From c5a1a4f9b38a0db742acb50ea7ff23059284c738 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 15:22:52 +0000 Subject: [PATCH 1/7] feat: add member reference, commit count and sync state to repo contributors (CM-1824) Signed-off-by: Mouad BANI --- .../V1790954467__repo_contributors_git_activity.sql | 9 +++++++++ 1 file changed, 9 insertions(+) create mode 100644 backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql diff --git a/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql b/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql new file mode 100644 index 0000000000..951ba0eeef --- /dev/null +++ b/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql @@ -0,0 +1,9 @@ +ALTER TABLE repo_contributors + ADD COLUMN IF NOT EXISTS cdp_member_id UUID, + ADD COLUMN IF NOT EXISTS commit_count INTEGER; + +CREATE TABLE IF NOT EXISTS repo_contributors_sync_state ( + source TEXT PRIMARY KEY, + watermark TIMESTAMPTZ NOT NULL, + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); From 3cb264d70fe1d231326d95b6470b287be1b80155 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 15:47:42 +0000 Subject: [PATCH 2/7] feat: sync repo contributors from Tinybird git activity with a watermark (CM-1824) Signed-off-by: Mouad BANI --- .../git-activity/mapRows.ts | 47 ++++ .../git-activity/readCommitContributors.ts | 73 ++++++ .../syncGitActivityContributors.ts | 234 ++++++++++++++++++ 3 files changed, 354 insertions(+) create mode 100644 services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts create mode 100644 services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts create mode 100644 services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts b/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts new file mode 100644 index 0000000000..e6745bd8d2 --- /dev/null +++ b/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts @@ -0,0 +1,47 @@ +import { RepoContributorRow, classifyIdentity } from '../governance/mapRows' + +export const GIT_ACTIVITY_SOURCE = 'git_activity' +export const GIT_ACTIVITY_ROLE = 'commit-author' +export const GIT_ACTIVITY_ROLE_KIND = 'contributor' + +export interface TinybirdCommitContributorRow { + channel: string + memberId: string + platform: string + username: string + commitCount: number | string + firstCommitAt: string + lastCommitAt: string + lastUpdatedAt: string +} + +export interface GitActivityContributorRow extends RepoContributorRow { + cdpMemberId: string + commitCount: number +} + +export function parseTinybirdDateTime(value: string): Date { + const parsed = new Date(`${value.replace(' ', 'T')}Z`) + if (Number.isNaN(parsed.getTime())) { + throw new Error(`Unparseable Tinybird datetime: ${value}`) + } + return parsed +} + +export function toGitActivityContributorRow( + row: TinybirdCommitContributorRow, + repoId: string, +): GitActivityContributorRow { + return { + repoId, + source: GIT_ACTIVITY_SOURCE, + role: GIT_ACTIVITY_ROLE, + roleKind: GIT_ACTIVITY_ROLE_KIND, + ...classifyIdentity(row.platform, 'username', row.username), + firstSeenAt: parseTinybirdDateTime(row.firstCommitAt), + lastSeenAt: parseTinybirdDateTime(row.lastCommitAt), + endedAt: null, + cdpMemberId: row.memberId, + commitCount: Number(row.commitCount), + } +} diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts new file mode 100644 index 0000000000..e1d532fa32 --- /dev/null +++ b/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts @@ -0,0 +1,73 @@ +import { TinybirdClient } from '@crowd/database' + +import { TinybirdCommitContributorRow, parseTinybirdDateTime } from './mapRows' + +const DATASOURCE = 'repo_commit_contributors_copy_ds' + +interface TinybirdResult { + data: T[] +} + +export type CommitContributorPageKey = Pick< + TinybirdCommitContributorRow, + 'channel' | 'memberId' | 'platform' | 'username' +> + +export function pageKeyOf(row: TinybirdCommitContributorRow): CommitContributorPageKey { + return { + channel: row.channel, + memberId: row.memberId, + platform: row.platform, + username: row.username, + } +} + +export function formatTinybirdDateTime(value: Date): string { + return value.toISOString().replace('T', ' ').replace('Z', '') +} + +export async function readSnapshotComputedAt(tb: TinybirdClient): Promise { + const result = await tb.executeSql>( + `SELECT toString(max(computedAt)) AS computedAt FROM ${DATASOURCE} FORMAT JSON`, + ) + const computedAt = result.data[0]?.computedAt + if (!computedAt || computedAt === '1970-01-01 00:00:00') return null + return parseTinybirdDateTime(computedAt) +} + +export async function readCommitContributorPage( + tb: TinybirdClient, + after: CommitContributorPageKey | null, + updatedSince: Date | null, + limit: number, +): Promise { + const query = `% + SELECT channel, memberId, platform, username, commitCount, + toString(firstCommitAt) AS firstCommitAt, + toString(lastCommitAt) AS lastCommitAt, + toString(lastUpdatedAt) AS lastUpdatedAt + FROM ${DATASOURCE} + WHERE 1 = 1 + {% if defined(afterChannel) %} + AND (channel, memberId, platform, username) > + ({{String(afterChannel)}}, {{String(afterMemberId)}}, {{String(afterPlatform)}}, {{String(afterUsername)}}) + {% end %} + {% if defined(updatedSince) %} + AND lastUpdatedAt > parseDateTime64BestEffort({{String(updatedSince)}}) + {% end %} + ORDER BY channel, memberId, platform, username + LIMIT {{Int32(limit)}} + FORMAT JSON` + + const params: Record = { limit } + if (after) { + params.afterChannel = after.channel + params.afterMemberId = after.memberId + params.afterPlatform = after.platform + params.afterUsername = after.username + } + if (updatedSince) params.updatedSince = formatTinybirdDateTime(updatedSince) + + const result = await tb.executeSql>(query, params) + return result.data +} diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts new file mode 100644 index 0000000000..74ee29ef6c --- /dev/null +++ b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts @@ -0,0 +1,234 @@ +import { heartbeat } from '@temporalio/activity' + +import { QueryExecutor } from '@crowd/data-access-layer/src/queryExecutor' +import { prepareBulkInsert } from '@crowd/data-access-layer/src/utils' +import { TinybirdClient } from '@crowd/database' +import { getServiceChildLogger } from '@crowd/logging' + +import { getPackagesDb } from '../../db' +import { canonicalGovernanceRepoUrl } from '../governance/mapRows' +import { findRepoIdsByUrl, readRunStart } from '../governance/syncGovernanceContributors' +import { + GIT_ACTIVITY_SOURCE, + GitActivityContributorRow, + TinybirdCommitContributorRow, + toGitActivityContributorRow, +} from './mapRows' +import { + CommitContributorPageKey, + pageKeyOf, + readCommitContributorPage, + readSnapshotComputedAt, +} from './readCommitContributors' + +const log = getServiceChildLogger('syncGitActivityContributors') + +const PAGE_SIZE = 5000 +const WATERMARK_MARGIN_MS = 24 * 60 * 60 * 1000 + +export interface GitActivitySyncOptions { + full: boolean +} + +export interface GitActivitySyncCounts { + full: boolean + read: number + upserted: number + removed: number + skippedUnknownRepo: number +} + +export type GitActivityPageCounts = Pick + +const CONTRIBUTOR_COLUMNS = [ + 'repo_id', + 'source', + 'role', + 'role_kind', + 'identity_type', + 'identity_value', + 'first_seen_at', + 'last_seen_at', + 'ended_at', + 'cdp_member_id', + 'commit_count', +] + +const UPSERT_CLAUSE = `(repo_id, source, identity_type, identity_value, role) DO UPDATE SET + role_kind = EXCLUDED.role_kind, + first_seen_at = EXCLUDED.first_seen_at, + last_seen_at = EXCLUDED.last_seen_at, + ended_at = EXCLUDED.ended_at, + cdp_member_id = EXCLUDED.cdp_member_id, + commit_count = EXCLUDED.commit_count, + updated_at = NOW()` + +export async function readWatermark(pkgsQx: QueryExecutor): Promise { + const rows: Array<{ watermarkMs: number }> = await pkgsQx.select( + `SELECT FLOOR(EXTRACT(EPOCH FROM watermark) * 1000) AS "watermarkMs" + FROM repo_contributors_sync_state WHERE source = $(source)`, + { source: GIT_ACTIVITY_SOURCE }, + ) + return rows.length === 0 ? null : new Date(rows[0].watermarkMs) +} + +export async function writeWatermark(pkgsQx: QueryExecutor, watermark: Date): Promise { + await pkgsQx.result( + `INSERT INTO repo_contributors_sync_state (source, watermark, updated_at) + VALUES ($(source), $(watermark), NOW()) + ON CONFLICT (source) DO UPDATE SET watermark = EXCLUDED.watermark, updated_at = NOW()`, + { source: GIT_ACTIVITY_SOURCE, watermark }, + ) +} + +export function resolveUpdatedSince(watermark: Date | null, full: boolean): Date | null { + if (full || watermark === null) return null + return new Date(watermark.getTime() - WATERMARK_MARGIN_MS) +} + +export function assertSnapshotIsFresh(computedAt: Date | null, watermark: Date | null): Date { + if (computedAt === null) { + throw new Error('Tinybird repo_commit_contributors_copy_ds is empty') + } + if (watermark !== null && computedAt.getTime() <= watermark.getTime()) { + throw new Error( + `Tinybird repo_commit_contributors snapshot (${computedAt.toISOString()}) is not newer than the watermark (${watermark.toISOString()})`, + ) + } + return computedAt +} + +function contributorRowKey(row: GitActivityContributorRow): string { + return [row.repoId, row.identityType, row.identityValue].join('\u0000') +} + +export function mergeContributorRows( + rows: GitActivityContributorRow[], +): GitActivityContributorRow[] { + const byKey = new Map() + for (const row of rows) { + const key = contributorRowKey(row) + const existing = byKey.get(key) + if (!existing) { + byKey.set(key, row) + continue + } + byKey.set(key, { + ...existing, + firstSeenAt: existing.firstSeenAt < row.firstSeenAt ? existing.firstSeenAt : row.firstSeenAt, + lastSeenAt: existing.lastSeenAt > row.lastSeenAt ? existing.lastSeenAt : row.lastSeenAt, + commitCount: existing.commitCount + row.commitCount, + }) + } + return [...byKey.values()] +} + +function toDbRow(row: GitActivityContributorRow): Record { + return { + repo_id: row.repoId, + source: row.source, + role: row.role, + role_kind: row.roleKind, + identity_type: row.identityType, + identity_value: row.identityValue, + first_seen_at: row.firstSeenAt, + last_seen_at: row.lastSeenAt, + ended_at: row.endedAt, + cdp_member_id: row.cdpMemberId, + commit_count: row.commitCount, + } +} + +export async function upsertGitActivityContributors( + pkgsQx: QueryExecutor, + rows: GitActivityContributorRow[], +): Promise { + if (rows.length === 0) return 0 + return pkgsQx.result( + prepareBulkInsert('repo_contributors', CONTRIBUTOR_COLUMNS, rows.map(toDbRow), UPSERT_CLAUSE), + ) +} + +export async function removeGitActivityContributorsUntouchedSince( + pkgsQx: QueryExecutor, + runStartedAt: Date, +): Promise { + return pkgsQx.result( + 'DELETE FROM repo_contributors WHERE source = $(source) AND updated_at < $(runStartedAt)', + { source: GIT_ACTIVITY_SOURCE, runStartedAt }, + ) +} + +export async function applyGitActivityPage( + pkgsQx: QueryExecutor, + page: TinybirdCommitContributorRow[], +): Promise { + const canonicalUrls = page.map((row) => canonicalGovernanceRepoUrl(row.channel)) + const knownUrls = canonicalUrls.filter((url): url is string => url !== null) + const repoIdByUrl = await findRepoIdsByUrl(pkgsQx, [...new Set(knownUrls)]) + + const rows: GitActivityContributorRow[] = [] + let skippedUnknownRepo = 0 + page.forEach((row, i) => { + const url = canonicalUrls[i] + const repoId = url === null ? undefined : repoIdByUrl.get(url) + if (repoId === undefined) { + skippedUnknownRepo += 1 + return + } + rows.push(toGitActivityContributorRow(row, repoId)) + }) + + const upserted = await upsertGitActivityContributors(pkgsQx, mergeContributorRows(rows)) + return { upserted, skippedUnknownRepo } +} + +export async function syncGitActivityContributors( + options: GitActivitySyncOptions, +): Promise { + const pkgsQx = await getPackagesDb() + const tb = new TinybirdClient() + + const runStartedAt = await readRunStart(pkgsQx) + const watermark = await readWatermark(pkgsQx) + const snapshotComputedAt = assertSnapshotIsFresh(await readSnapshotComputedAt(tb), watermark) + const full = options.full || watermark === null + const updatedSince = resolveUpdatedSince(watermark, full) + + const counts: GitActivitySyncCounts = { + full, + read: 0, + upserted: 0, + removed: 0, + skippedUnknownRepo: 0, + } + + let after: CommitContributorPageKey | null = null + for (;;) { + const page = await readCommitContributorPage(tb, after, updatedSince, PAGE_SIZE) + if (page.length === 0) break + + const pageCounts = await applyGitActivityPage(pkgsQx, page) + counts.read += page.length + counts.upserted += pageCounts.upserted + counts.skippedUnknownRepo += pageCounts.skippedUnknownRepo + + after = pageKeyOf(page[page.length - 1]) + heartbeat(counts) + } + + if (counts.read === 0 || counts.upserted === 0) { + throw new Error( + `Refusing to finish git-activity contributors sync: read ${counts.read} rows, upserted ${counts.upserted}`, + ) + } + + if (full) { + counts.removed = await removeGitActivityContributorsUntouchedSince(pkgsQx, runStartedAt) + } + + await writeWatermark(pkgsQx, snapshotComputedAt) + + log.info(counts, 'Git-activity contributors synced from Tinybird') + return counts +} From 95ea8df86a9c354ac88459d96ed9446dc2a6d959 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 15:49:08 +0000 Subject: [PATCH 3/7] feat: schedule daily and biweekly git-activity contributors syncs (CM-1824) Signed-off-by: Mouad BANI --- .../apps/packages_worker/src/activities.ts | 5 +- .../src/bin/security-contacts-worker.ts | 6 +- .../src/member-contributors/activities.ts | 1 + .../src/member-contributors/schedule.ts | 91 ++++++++++++++----- .../src/member-contributors/workflows.ts | 28 +++++- .../packages_worker/src/workflows/index.ts | 5 +- 6 files changed, 107 insertions(+), 29 deletions(-) diff --git a/services/apps/packages_worker/src/activities.ts b/services/apps/packages_worker/src/activities.ts index ec3c83ecd2..24d6eb2b24 100644 --- a/services/apps/packages_worker/src/activities.ts +++ b/services/apps/packages_worker/src/activities.ts @@ -69,7 +69,10 @@ export { blastRadiusReachability, blastRadiusReport, } from './blast-radius/activities' -export { syncGovernanceContributors } from './member-contributors/activities' +export { + syncGitActivityContributors, + syncGovernanceContributors, +} from './member-contributors/activities' export { slackNotify } from './activities/index' export { syncGithubRepos } from './scorecard/activities' export { diff --git a/services/apps/packages_worker/src/bin/security-contacts-worker.ts b/services/apps/packages_worker/src/bin/security-contacts-worker.ts index 1236c7b61e..db1b7be9e5 100644 --- a/services/apps/packages_worker/src/bin/security-contacts-worker.ts +++ b/services/apps/packages_worker/src/bin/security-contacts-worker.ts @@ -1,4 +1,7 @@ -import { scheduleGovernanceFileContributorsSync } from '../member-contributors/schedule' +import { + scheduleGitActivityContributorsSync, + scheduleGovernanceFileContributorsSync, +} from '../member-contributors/schedule' import { scheduleReportingProtocolIngestion } from '../security-contacts/protocol/schedule' import { scheduleSecurityContactsIngestion } from '../security-contacts/schedule' import { svc } from '../service' @@ -8,5 +11,6 @@ setImmediate(async () => { await scheduleSecurityContactsIngestion() await scheduleReportingProtocolIngestion() await scheduleGovernanceFileContributorsSync() + await scheduleGitActivityContributorsSync() await svc.start() }) diff --git a/services/apps/packages_worker/src/member-contributors/activities.ts b/services/apps/packages_worker/src/member-contributors/activities.ts index 6946fe444c..828389c18d 100644 --- a/services/apps/packages_worker/src/member-contributors/activities.ts +++ b/services/apps/packages_worker/src/member-contributors/activities.ts @@ -1 +1,2 @@ +export { syncGitActivityContributors } from './git-activity/syncGitActivityContributors' export { syncGovernanceContributors } from './governance/syncGovernanceContributors' diff --git a/services/apps/packages_worker/src/member-contributors/schedule.ts b/services/apps/packages_worker/src/member-contributors/schedule.ts index 1e90e4f24e..05e6e33bf0 100644 --- a/services/apps/packages_worker/src/member-contributors/schedule.ts +++ b/services/apps/packages_worker/src/member-contributors/schedule.ts @@ -1,57 +1,106 @@ import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client' import { svc } from '../service' -import { syncGovernanceFileContributors } from '../workflows' +import { syncGovernanceFileContributors, syncRepoContributorsFromGitActivity } from '../workflows' +import { GitActivitySyncOptions } from './git-activity/syncGitActivityContributors' -const SCHEDULE_ID = 'governance-file-contributors-sync' -const WORKFLOW_EXECUTION_TIMEOUT = '6 hours' +const TASK_QUEUE = 'security-contacts-worker' +const RETRY = { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, +} + +interface ContributorsSchedule { + scheduleId: string + cron: string + action: { + type: 'startWorkflow' + workflowType: typeof syncGovernanceFileContributors | typeof syncRepoContributorsFromGitActivity + workflowId: string + taskQueue: string + workflowExecutionTimeout: string + retry: typeof RETRY + args: [] | [GitActivitySyncOptions] + } +} -function scheduleAction() { +function governanceFileSchedule(): ContributorsSchedule { return { - type: 'startWorkflow' as const, - workflowType: syncGovernanceFileContributors, - workflowId: 'governance-file-contributors-daily', - taskQueue: 'security-contacts-worker', - workflowExecutionTimeout: WORKFLOW_EXECUTION_TIMEOUT, - retry: { - initialInterval: '30 seconds', - backoffCoefficient: 2, - maximumAttempts: 3, + scheduleId: 'governance-file-contributors-sync', + cron: '0 5 * * *', + action: { + type: 'startWorkflow', + workflowType: syncGovernanceFileContributors, + workflowId: 'governance-file-contributors-daily', + taskQueue: TASK_QUEUE, + workflowExecutionTimeout: '6 hours', + retry: RETRY, + args: [], }, - args: [] as [], } } -export async function scheduleGovernanceFileContributorsSync(): Promise { +function gitActivitySchedule(full: boolean): ContributorsSchedule { + const mode = full ? 'full' : 'incremental' + return { + scheduleId: `git-activity-contributors-${mode}-sync`, + cron: full ? '0 6 1,15 * *' : '0 6 2-14,16-31 * *', + action: { + type: 'startWorkflow', + workflowType: syncRepoContributorsFromGitActivity, + workflowId: `git-activity-contributors-${mode}`, + taskQueue: TASK_QUEUE, + workflowExecutionTimeout: '8 hours', + retry: RETRY, + args: [{ full }], + }, + } +} + +async function createOrReconcileSchedule(schedule: ContributorsSchedule): Promise { const { temporal } = svc if (!temporal) throw new Error('Temporal client not initialized') try { await temporal.schedule.create({ - scheduleId: SCHEDULE_ID, + scheduleId: schedule.scheduleId, spec: { - cronExpressions: ['0 5 * * *'], + cronExpressions: [schedule.cron], }, policies: { overlap: ScheduleOverlapPolicy.SKIP, catchupWindow: '1 hour', }, - action: scheduleAction(), + action: schedule.action, }) } catch (err) { if (err instanceof ScheduleAlreadyRunning) { - svc.log.info(`Schedule ${SCHEDULE_ID} already exists, reconciling action.`) - const handle = temporal.schedule.getHandle(SCHEDULE_ID) + svc.log.info(`Schedule ${schedule.scheduleId} already exists, reconciling action.`) + const handle = temporal.schedule.getHandle(schedule.scheduleId) await handle.update((prev) => ({ ...prev, + spec: { + ...prev.spec, + cronExpressions: [schedule.cron], + }, policies: { ...prev.policies, overlap: ScheduleOverlapPolicy.SKIP, }, - action: scheduleAction(), + action: schedule.action, })) } else { throw err } } } + +export async function scheduleGovernanceFileContributorsSync(): Promise { + await createOrReconcileSchedule(governanceFileSchedule()) +} + +export async function scheduleGitActivityContributorsSync(): Promise { + await createOrReconcileSchedule(gitActivitySchedule(false)) + await createOrReconcileSchedule(gitActivitySchedule(true)) +} diff --git a/services/apps/packages_worker/src/member-contributors/workflows.ts b/services/apps/packages_worker/src/member-contributors/workflows.ts index f633a0e2af..56d6e129ce 100644 --- a/services/apps/packages_worker/src/member-contributors/workflows.ts +++ b/services/apps/packages_worker/src/member-contributors/workflows.ts @@ -1,18 +1,36 @@ import { proxyActivities } from '@temporalio/workflow' import type * as memberContributorActivities from './activities' +import type { + GitActivitySyncCounts, + GitActivitySyncOptions, +} from './git-activity/syncGitActivityContributors' import type { GovernanceSyncCounts } from './governance/syncGovernanceContributors' +const RETRY = { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, +} + const { syncGovernanceContributors } = proxyActivities({ startToCloseTimeout: '1 hour', heartbeatTimeout: '5 minutes', - retry: { - initialInterval: '30 seconds', - backoffCoefficient: 2, - maximumAttempts: 3, - }, + retry: RETRY, +}) + +const { syncGitActivityContributors } = proxyActivities({ + startToCloseTimeout: '4 hours', + heartbeatTimeout: '10 minutes', + retry: RETRY, }) export async function syncGovernanceFileContributors(): Promise { return syncGovernanceContributors() } + +export async function syncRepoContributorsFromGitActivity( + options: GitActivitySyncOptions, +): Promise { + return syncGitActivityContributors(options) +} diff --git a/services/apps/packages_worker/src/workflows/index.ts b/services/apps/packages_worker/src/workflows/index.ts index 30f4ffc503..663c716e37 100644 --- a/services/apps/packages_worker/src/workflows/index.ts +++ b/services/apps/packages_worker/src/workflows/index.ts @@ -39,6 +39,9 @@ export { ingestSecurityContactsForPurlWorkflow, ingestReportingProtocols, } from '../security-contacts/workflows' -export { syncGovernanceFileContributors } from '../member-contributors/workflows' +export { + syncGovernanceFileContributors, + syncRepoContributorsFromGitActivity, +} from '../member-contributors/workflows' export { analyzeBlastRadius } from '../blast-radius/workflows' export { sweepPackageRepoConfidence } from '../package-repos/workflows' From 2ccbe53640921a2c77a311a798fa94e55251e4e8 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 15:51:03 +0000 Subject: [PATCH 4/7] test: cover the git-activity contributors sync (CM-1824) Signed-off-by: Mouad BANI --- .../syncGitActivityContributors.test.ts | 291 ++++++++++++++++++ 1 file changed, 291 insertions(+) create mode 100644 services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts diff --git a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts new file mode 100644 index 0000000000..9374d9987d --- /dev/null +++ b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts @@ -0,0 +1,291 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' + +import { QueryExecutor } from '@crowd/data-access-layer/src/queryExecutor' +import { TinybirdClient } from '@crowd/database' + +import { getPackagesDb } from '../../db' +import { + TinybirdCommitContributorRow, + parseTinybirdDateTime, + toGitActivityContributorRow, +} from '../git-activity/mapRows' +import { readCommitContributorPage } from '../git-activity/readCommitContributors' +import { + assertSnapshotIsFresh, + mergeContributorRows, + resolveUpdatedSince, + syncGitActivityContributors, +} from '../git-activity/syncGitActivityContributors' + +vi.mock('@temporalio/activity', () => ({ heartbeat: vi.fn() })) +vi.mock('../../db', () => ({ getPackagesDb: vi.fn() })) +vi.mock('@crowd/database', () => ({ TinybirdClient: vi.fn() })) + +const runStart = new Date('2026-10-03T06:00:00Z') +const computedAt = '2026-10-03 03:45:00' +const watermark = new Date('2026-10-02T03:45:00Z') + +function tbRow( + overrides: Partial = {}, +): TinybirdCommitContributorRow { + return { + channel: 'https://github.com/org/repo', + memberId: 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', + platform: 'github', + username: 'JaneDoe', + commitCount: '12', + firstCommitAt: '2020-01-01 10:00:00.000', + lastCommitAt: '2026-09-30 12:30:00.000', + lastUpdatedAt: '2026-10-02 20:00:00.000', + ...overrides, + } +} + +function stubTinybird(pages: TinybirdCommitContributorRow[][], snapshotComputedAt = computedAt) { + const executeSql = vi.fn().mockImplementation(async (query: string) => { + if (query.includes('max(computedAt)')) return { data: [{ computedAt: snapshotComputedAt }] } + return { data: pages.shift() ?? [] } + }) + vi.mocked(TinybirdClient).mockImplementation(function () { + return { executeSql } as unknown as TinybirdClient + }) + return executeSql +} + +function stubPackagesQx( + repos: Array<{ id: string; url: string }>, + storedWatermark: Date | null, + removed = 0, +) { + const select = vi + .fn() + .mockResolvedValueOnce([{ nowMs: runStart.getTime() }]) + .mockResolvedValueOnce(storedWatermark ? [{ watermarkMs: storedWatermark.getTime() }] : []) + .mockResolvedValue(repos) + const result = vi + .fn() + .mockImplementation((sql: string) => Promise.resolve(sql.startsWith('DELETE') ? removed : 1)) + const qx = { select, result } as unknown as QueryExecutor + vi.mocked(getPackagesDb).mockResolvedValue(qx) + return { select, result } +} + +function sqlCalls(result: ReturnType, prefix: string): string[] { + return result.mock.calls.map((c) => String(c[0])).filter((sql) => sql.startsWith(prefix)) +} + +beforeEach(() => { + vi.clearAllMocks() +}) + +describe('parseTinybirdDateTime', () => { + it('reads Tinybird datetimes as UTC', () => { + expect(parseTinybirdDateTime('2026-10-03 03:45:00').toISOString()).toBe( + '2026-10-03T03:45:00.000Z', + ) + expect(parseTinybirdDateTime('2026-10-03 03:45:00.123').toISOString()).toBe( + '2026-10-03T03:45:00.123Z', + ) + }) + + it('throws on garbage', () => { + expect(() => parseTinybirdDateTime('not a date')).toThrow('Unparseable Tinybird datetime') + }) +}) + +describe('toGitActivityContributorRow', () => { + it('maps a github login to a lowercase github-login identity with member and count', () => { + const row = toGitActivityContributorRow(tbRow(), '42') + expect(row).toMatchObject({ + repoId: '42', + source: 'git_activity', + role: 'commit-author', + roleKind: 'contributor', + identityType: 'github-login', + identityValue: 'janedoe', + endedAt: null, + cdpMemberId: 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', + commitCount: 12, + }) + expect(row.firstSeenAt.toISOString()).toBe('2020-01-01T10:00:00.000Z') + expect(row.lastSeenAt.toISOString()).toBe('2026-09-30T12:30:00.000Z') + }) + + it('maps a git author email to an email identity', () => { + const row = toGitActivityContributorRow( + tbRow({ platform: 'git', username: 'Jane@Example.com' }), + '42', + ) + expect(row.identityType).toBe('email') + expect(row.identityValue).toBe('jane@example.com') + }) +}) + +describe('mergeContributorRows', () => { + it('sums counts and widens the seen window for rows sharing repo and identity', () => { + const a = toGitActivityContributorRow(tbRow({ commitCount: 3 }), '1') + const b = toGitActivityContributorRow( + tbRow({ + commitCount: 4, + firstCommitAt: '2019-01-01 00:00:00.000', + lastCommitAt: '2026-10-01 00:00:00.000', + }), + '1', + ) + const merged = mergeContributorRows([a, b]) + expect(merged).toHaveLength(1) + expect(merged[0].commitCount).toBe(7) + expect(merged[0].firstSeenAt).toEqual(b.firstSeenAt) + expect(merged[0].lastSeenAt).toEqual(b.lastSeenAt) + }) + + it('keeps rows that differ in repo or identity', () => { + const rows = [ + toGitActivityContributorRow(tbRow(), '1'), + toGitActivityContributorRow(tbRow(), '2'), + toGitActivityContributorRow(tbRow({ username: 'other' }), '1'), + ] + expect(mergeContributorRows(rows)).toHaveLength(3) + }) +}) + +describe('resolveUpdatedSince', () => { + it('is null for full runs and first runs', () => { + expect(resolveUpdatedSince(null, false)).toBeNull() + expect(resolveUpdatedSince(watermark, true)).toBeNull() + }) + + it('subtracts a one-day margin from the watermark', () => { + expect(resolveUpdatedSince(watermark, false)?.toISOString()).toBe('2026-10-01T03:45:00.000Z') + }) +}) + +describe('assertSnapshotIsFresh', () => { + it('throws on an empty snapshot', () => { + expect(() => assertSnapshotIsFresh(null, null)).toThrow('is empty') + }) + + it('throws when the snapshot is not newer than the watermark', () => { + expect(() => assertSnapshotIsFresh(watermark, watermark)).toThrow('not newer') + }) + + it('returns the snapshot time otherwise', () => { + const fresh = new Date('2026-10-03T03:45:00Z') + expect(assertSnapshotIsFresh(fresh, watermark)).toBe(fresh) + expect(assertSnapshotIsFresh(fresh, null)).toBe(fresh) + }) +}) + +describe('readCommitContributorPage', () => { + it('passes the keyset cursor and watermark as template params', async () => { + const executeSql = vi.fn().mockResolvedValue({ data: [] }) + const tb = { executeSql } as unknown as TinybirdClient + const after = { channel: 'c', memberId: 'm', platform: 'github', username: 'u' } + + await readCommitContributorPage(tb, after, watermark, 10) + + const [query, params] = executeSql.mock.calls[0] + expect(query).toContain('{% if defined(afterChannel) %}') + expect(params).toEqual({ + limit: 10, + afterChannel: 'c', + afterMemberId: 'm', + afterPlatform: 'github', + afterUsername: 'u', + updatedSince: '2026-10-02 03:45:00.000', + }) + }) + + it('omits cursor and watermark params on a first full page', async () => { + const executeSql = vi.fn().mockResolvedValue({ data: [] }) + const tb = { executeSql } as unknown as TinybirdClient + + await readCommitContributorPage(tb, null, null, 10) + + expect(executeSql.mock.calls[0][1]).toEqual({ limit: 10 }) + }) +}) + +describe('syncGitActivityContributors', () => { + const repo = { id: '1', url: 'https://github.com/org/repo' } + + it('runs a full backfill when no watermark is stored and writes the snapshot time', async () => { + const executeSql = stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], null, 3) + + const counts = await syncGitActivityContributors({ full: false }) + + expect(counts).toEqual({ full: true, read: 1, upserted: 1, removed: 3, skippedUnknownRepo: 0 }) + expect(executeSql.mock.calls[1][1]).toEqual({ limit: 5000 }) + expect(sqlCalls(result, 'DELETE')).toHaveLength(1) + const [watermarkSql, watermarkParams] = result.mock.calls[result.mock.calls.length - 1] + expect(watermarkSql).toContain('repo_contributors_sync_state') + expect(watermarkParams.watermark.toISOString()).toBe('2026-10-03T03:45:00.000Z') + }) + + it('runs incrementally from the watermark minus a day and does not delete', async () => { + const executeSql = stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], watermark) + + const counts = await syncGitActivityContributors({ full: false }) + + expect(counts.full).toBe(false) + expect(counts.removed).toBe(0) + expect(executeSql.mock.calls[1][1]).toMatchObject({ updatedSince: '2026-10-01 03:45:00.000' }) + expect(sqlCalls(result, 'DELETE')).toHaveLength(0) + }) + + it('ignores the watermark and reconciles when full is requested', async () => { + const executeSql = stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], watermark, 2) + + const counts = await syncGitActivityContributors({ full: true }) + + expect(counts.full).toBe(true) + expect(counts.removed).toBe(2) + expect(executeSql.mock.calls[1][1]).toEqual({ limit: 5000 }) + expect(sqlCalls(result, 'DELETE')).toHaveLength(1) + }) + + it('advances the keyset cursor between pages and skips unknown repos', async () => { + const executeSql = stubTinybird([ + [tbRow(), tbRow({ channel: 'https://github.com/org/unknown' })], + [tbRow({ username: 'second' })], + [], + ]) + stubPackagesQx([repo], watermark) + + const counts = await syncGitActivityContributors({ full: false }) + + expect(counts.read).toBe(3) + expect(counts.skippedUnknownRepo).toBe(1) + expect(executeSql.mock.calls[2][1]).toMatchObject({ + afterChannel: 'https://github.com/org/unknown', + afterUsername: 'JaneDoe', + }) + }) + + it('throws when the snapshot is stale and touches nothing', async () => { + stubTinybird([[tbRow()], []], '2026-10-02 03:45:00') + const { result } = stubPackagesQx([repo], watermark) + + await expect(syncGitActivityContributors({ full: false })).rejects.toThrow('not newer') + expect(result).not.toHaveBeenCalled() + }) + + it('throws when nothing was read and leaves the watermark untouched', async () => { + stubTinybird([[]]) + const { result } = stubPackagesQx([repo], watermark) + + await expect(syncGitActivityContributors({ full: true })).rejects.toThrow('read 0 rows') + expect(result).not.toHaveBeenCalled() + }) + + it('throws when rows were read but none upserted', async () => { + stubTinybird([[tbRow({ channel: 'https://github.com/org/unknown' })], []]) + const { result } = stubPackagesQx([repo], watermark) + + await expect(syncGitActivityContributors({ full: true })).rejects.toThrow('upserted 0') + expect(sqlCalls(result, 'INSERT INTO repo_contributors_sync_state')).toHaveLength(0) + }) +}) From 06283169ad8d759b9e88f6e647392072b335cdd1 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 16:46:07 +0000 Subject: [PATCH 5/7] refactor: drop the commit count from git-activity contributors (CM-1824) Signed-off-by: Mouad BANI --- .../V1790954467__repo_contributors_git_activity.sql | 3 +-- .../__tests__/syncGitActivityContributors.test.ts | 10 +++------- .../src/member-contributors/git-activity/mapRows.ts | 3 --- .../git-activity/readCommitContributors.ts | 2 +- .../git-activity/syncGitActivityContributors.ts | 4 ---- 5 files changed, 5 insertions(+), 17 deletions(-) diff --git a/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql b/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql index 951ba0eeef..3200224517 100644 --- a/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql +++ b/backend/src/osspckgs/migrations/V1790954467__repo_contributors_git_activity.sql @@ -1,6 +1,5 @@ ALTER TABLE repo_contributors - ADD COLUMN IF NOT EXISTS cdp_member_id UUID, - ADD COLUMN IF NOT EXISTS commit_count INTEGER; + ADD COLUMN IF NOT EXISTS cdp_member_id UUID; CREATE TABLE IF NOT EXISTS repo_contributors_sync_state ( source TEXT PRIMARY KEY, diff --git a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts index 9374d9987d..41da4be6fa 100644 --- a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts +++ b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts @@ -33,7 +33,6 @@ function tbRow( memberId: 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', platform: 'github', username: 'JaneDoe', - commitCount: '12', firstCommitAt: '2020-01-01 10:00:00.000', lastCommitAt: '2026-09-30 12:30:00.000', lastUpdatedAt: '2026-10-02 20:00:00.000', @@ -94,7 +93,7 @@ describe('parseTinybirdDateTime', () => { }) describe('toGitActivityContributorRow', () => { - it('maps a github login to a lowercase github-login identity with member and count', () => { + it('maps a github login to a lowercase github-login identity with its member', () => { const row = toGitActivityContributorRow(tbRow(), '42') expect(row).toMatchObject({ repoId: '42', @@ -105,7 +104,6 @@ describe('toGitActivityContributorRow', () => { identityValue: 'janedoe', endedAt: null, cdpMemberId: 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', - commitCount: 12, }) expect(row.firstSeenAt.toISOString()).toBe('2020-01-01T10:00:00.000Z') expect(row.lastSeenAt.toISOString()).toBe('2026-09-30T12:30:00.000Z') @@ -122,11 +120,10 @@ describe('toGitActivityContributorRow', () => { }) describe('mergeContributorRows', () => { - it('sums counts and widens the seen window for rows sharing repo and identity', () => { - const a = toGitActivityContributorRow(tbRow({ commitCount: 3 }), '1') + it('widens the seen window for rows sharing repo and identity', () => { + const a = toGitActivityContributorRow(tbRow(), '1') const b = toGitActivityContributorRow( tbRow({ - commitCount: 4, firstCommitAt: '2019-01-01 00:00:00.000', lastCommitAt: '2026-10-01 00:00:00.000', }), @@ -134,7 +131,6 @@ describe('mergeContributorRows', () => { ) const merged = mergeContributorRows([a, b]) expect(merged).toHaveLength(1) - expect(merged[0].commitCount).toBe(7) expect(merged[0].firstSeenAt).toEqual(b.firstSeenAt) expect(merged[0].lastSeenAt).toEqual(b.lastSeenAt) }) diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts b/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts index e6745bd8d2..7e3c94e1e4 100644 --- a/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts +++ b/services/apps/packages_worker/src/member-contributors/git-activity/mapRows.ts @@ -9,7 +9,6 @@ export interface TinybirdCommitContributorRow { memberId: string platform: string username: string - commitCount: number | string firstCommitAt: string lastCommitAt: string lastUpdatedAt: string @@ -17,7 +16,6 @@ export interface TinybirdCommitContributorRow { export interface GitActivityContributorRow extends RepoContributorRow { cdpMemberId: string - commitCount: number } export function parseTinybirdDateTime(value: string): Date { @@ -42,6 +40,5 @@ export function toGitActivityContributorRow( lastSeenAt: parseTinybirdDateTime(row.lastCommitAt), endedAt: null, cdpMemberId: row.memberId, - commitCount: Number(row.commitCount), } } diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts index e1d532fa32..7ad34995d0 100644 --- a/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts +++ b/services/apps/packages_worker/src/member-contributors/git-activity/readCommitContributors.ts @@ -42,7 +42,7 @@ export async function readCommitContributorPage( limit: number, ): Promise { const query = `% - SELECT channel, memberId, platform, username, commitCount, + SELECT channel, memberId, platform, username, toString(firstCommitAt) AS firstCommitAt, toString(lastCommitAt) AS lastCommitAt, toString(lastUpdatedAt) AS lastUpdatedAt diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts index 74ee29ef6c..f6fba8c26f 100644 --- a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts +++ b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts @@ -51,7 +51,6 @@ const CONTRIBUTOR_COLUMNS = [ 'last_seen_at', 'ended_at', 'cdp_member_id', - 'commit_count', ] const UPSERT_CLAUSE = `(repo_id, source, identity_type, identity_value, role) DO UPDATE SET @@ -60,7 +59,6 @@ const UPSERT_CLAUSE = `(repo_id, source, identity_type, identity_value, role) DO last_seen_at = EXCLUDED.last_seen_at, ended_at = EXCLUDED.ended_at, cdp_member_id = EXCLUDED.cdp_member_id, - commit_count = EXCLUDED.commit_count, updated_at = NOW()` export async function readWatermark(pkgsQx: QueryExecutor): Promise { @@ -117,7 +115,6 @@ export function mergeContributorRows( ...existing, firstSeenAt: existing.firstSeenAt < row.firstSeenAt ? existing.firstSeenAt : row.firstSeenAt, lastSeenAt: existing.lastSeenAt > row.lastSeenAt ? existing.lastSeenAt : row.lastSeenAt, - commitCount: existing.commitCount + row.commitCount, }) } return [...byKey.values()] @@ -135,7 +132,6 @@ function toDbRow(row: GitActivityContributorRow): Record { last_seen_at: row.lastSeenAt, ended_at: row.endedAt, cdp_member_id: row.cdpMemberId, - commit_count: row.commitCount, } } From d0e97f5932bf4ff7b6888713782774a21eec1f83 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 16:47:01 +0000 Subject: [PATCH 6/7] fix: widen the git-activity seen window across channel variants (CM-1824) Signed-off-by: Mouad BANI --- .../__tests__/syncGitActivityContributors.test.ts | 13 +++++++++++++ .../git-activity/syncGitActivityContributors.ts | 4 ++-- 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts index 41da4be6fa..8690312d21 100644 --- a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts +++ b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts @@ -219,6 +219,19 @@ describe('syncGitActivityContributors', () => { expect(watermarkParams.watermark.toISOString()).toBe('2026-10-03T03:45:00.000Z') }) + it('widens the stored seen window instead of overwriting it', async () => { + stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], watermark) + + await syncGitActivityContributors({ full: false }) + + const upsertSql = result.mock.calls + .map((c) => String(c[0])) + .find((sql) => sql.includes('INSERT INTO "repo_contributors"')) + expect(upsertSql).toContain('LEAST(repo_contributors.first_seen_at, EXCLUDED.first_seen_at)') + expect(upsertSql).toContain('GREATEST(repo_contributors.last_seen_at, EXCLUDED.last_seen_at)') + }) + it('runs incrementally from the watermark minus a day and does not delete', async () => { const executeSql = stubTinybird([[tbRow()], []]) const { result } = stubPackagesQx([repo], watermark) diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts index f6fba8c26f..b61facbbbd 100644 --- a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts +++ b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts @@ -55,8 +55,8 @@ const CONTRIBUTOR_COLUMNS = [ const UPSERT_CLAUSE = `(repo_id, source, identity_type, identity_value, role) DO UPDATE SET role_kind = EXCLUDED.role_kind, - first_seen_at = EXCLUDED.first_seen_at, - last_seen_at = EXCLUDED.last_seen_at, + first_seen_at = LEAST(repo_contributors.first_seen_at, EXCLUDED.first_seen_at), + last_seen_at = GREATEST(repo_contributors.last_seen_at, EXCLUDED.last_seen_at), ended_at = EXCLUDED.ended_at, cdp_member_id = EXCLUDED.cdp_member_id, updated_at = NOW()` From 362bc8ae0eeebf965b72cd9e07bb5c12b6f23cc8 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 2 Oct 2026 17:04:05 +0000 Subject: [PATCH 7/7] fix: refuse the full git-activity reconcile on an unexpectedly small snapshot (CM-1824) Signed-off-by: Mouad BANI --- .../syncGitActivityContributors.test.ts | 35 ++++++++++++++++--- .../syncGitActivityContributors.ts | 21 +++++++++++ 2 files changed, 51 insertions(+), 5 deletions(-) diff --git a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts index 8690312d21..0bbe40449d 100644 --- a/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts +++ b/services/apps/packages_worker/src/member-contributors/__tests__/syncGitActivityContributors.test.ts @@ -55,12 +55,16 @@ function stubPackagesQx( repos: Array<{ id: string; url: string }>, storedWatermark: Date | null, removed = 0, + existing = { total: '0', untouched: '0' }, ) { - const select = vi - .fn() - .mockResolvedValueOnce([{ nowMs: runStart.getTime() }]) - .mockResolvedValueOnce(storedWatermark ? [{ watermarkMs: storedWatermark.getTime() }] : []) - .mockResolvedValue(repos) + const select = vi.fn().mockImplementation(async (sql: string) => { + if (sql.includes('"nowMs"')) return [{ nowMs: runStart.getTime() }] + if (sql.includes('"watermarkMs"')) { + return storedWatermark ? [{ watermarkMs: storedWatermark.getTime() }] : [] + } + if (sql.includes('count(*)')) return [existing] + return repos + }) const result = vi .fn() .mockImplementation((sql: string) => Promise.resolve(sql.startsWith('DELETE') ? removed : 1)) @@ -256,6 +260,27 @@ describe('syncGitActivityContributors', () => { expect(sqlCalls(result, 'DELETE')).toHaveLength(1) }) + it('reconciles when the untouched share stays within the threshold', async () => { + stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], watermark, 10, { total: '100', untouched: '10' }) + + const counts = await syncGitActivityContributors({ full: true }) + + expect(counts.removed).toBe(10) + expect(sqlCalls(result, 'DELETE')).toHaveLength(1) + }) + + it('refuses to reconcile when the snapshot would remove too many existing rows', async () => { + stubTinybird([[tbRow()], []]) + const { result } = stubPackagesQx([repo], watermark, 90, { total: '100', untouched: '90' }) + + await expect(syncGitActivityContributors({ full: true })).rejects.toThrow( + 'would remove 90 of 100 existing rows', + ) + expect(sqlCalls(result, 'DELETE')).toHaveLength(0) + expect(sqlCalls(result, 'INSERT INTO repo_contributors_sync_state')).toHaveLength(0) + }) + it('advances the keyset cursor between pages and skips unknown repos', async () => { const executeSql = stubTinybird([ [tbRow(), tbRow({ channel: 'https://github.com/org/unknown' })], diff --git a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts index b61facbbbd..84aa011783 100644 --- a/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts +++ b/services/apps/packages_worker/src/member-contributors/git-activity/syncGitActivityContributors.ts @@ -25,6 +25,7 @@ const log = getServiceChildLogger('syncGitActivityContributors') const PAGE_SIZE = 5000 const WATERMARK_MARGIN_MS = 24 * 60 * 60 * 1000 +const MAX_RECONCILE_REMOVAL_SHARE = 0.2 export interface GitActivitySyncOptions { full: boolean @@ -145,6 +146,25 @@ export async function upsertGitActivityContributors( ) } +export async function assertReconcileRemovalIsPlausible( + pkgsQx: QueryExecutor, + runStartedAt: Date, +): Promise { + const rows: Array<{ total: string; untouched: string }> = await pkgsQx.select( + `SELECT count(*) AS total, + count(*) FILTER (WHERE updated_at < $(runStartedAt)) AS untouched + FROM repo_contributors WHERE source = $(source)`, + { source: GIT_ACTIVITY_SOURCE, runStartedAt }, + ) + const total = Number(rows[0].total) + const untouched = Number(rows[0].untouched) + if (total > 0 && untouched / total > MAX_RECONCILE_REMOVAL_SHARE) { + throw new Error( + `Refusing to reconcile git-activity contributors: snapshot would remove ${untouched} of ${total} existing rows`, + ) + } +} + export async function removeGitActivityContributorsUntouchedSince( pkgsQx: QueryExecutor, runStartedAt: Date, @@ -220,6 +240,7 @@ export async function syncGitActivityContributors( } if (full) { + await assertReconcileRemovalIsPlausible(pkgsQx, runStartedAt) counts.removed = await removeGitActivityContributorsUntouchedSince(pkgsQx, runStartedAt) }