From 3b493a7912585d539365609fac9c359833a62043 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 25 Sep 2026 16:17:29 +0000 Subject: [PATCH 1/3] fix: skip github-deleted records in shadow diff (CM-1473) Signed-off-by: Mouad BANI --- .../src/activities/shadowDiffActivities.ts | 68 +++++++--- .../connectors_worker/src/deletedRecords.ts | 128 ++++++++++++++++++ 2 files changed, 174 insertions(+), 22 deletions(-) create mode 100644 services/apps/connectors_worker/src/deletedRecords.ts diff --git a/services/apps/connectors_worker/src/activities/shadowDiffActivities.ts b/services/apps/connectors_worker/src/activities/shadowDiffActivities.ts index 1ecb7b2649..c4169753d0 100644 --- a/services/apps/connectors_worker/src/activities/shadowDiffActivities.ts +++ b/services/apps/connectors_worker/src/activities/shadowDiffActivities.ts @@ -11,6 +11,7 @@ import { import { getNangoMappingForRepo } from '@crowd/data-access-layer/src/integrations' import { dbStoreQx } from '@crowd/data-access-layer/src/queryExecutor' +import { dropConfirmedDeletedRecords, hasDeletedRecordCandidates } from '../deletedRecords' import { dropConfirmedForcePushedCommits, hasForcePushCandidates } from '../forcePushedCommits' import { svc } from '../main' import { @@ -153,9 +154,12 @@ export async function runShadowDiffForChannel( const unitDiffs: { unit: IShadowDiffUnit; result: IShadowDiffUnitResult }[] = [] let confirmationHttp: Promise | null = null for (const unit of pendingUnits) { - const result = await diffUnit(qx, unit, mapping.connectionId, windowStart, windowEnd) + let result = await diffUnit(qx, unit, mapping.connectionId, windowStart, windowEnd) - if (hasForcePushCandidates(unit.syncName, result.mismatches)) { + const needsForcePushCheck = hasForcePushCandidates(unit.syncName, result.mismatches) + const needsDeletedRecordCheck = hasDeletedRecordCandidates(unit.syncName, result.mismatches) + + if (needsForcePushCheck || needsDeletedRecordCheck) { confirmationHttp ??= createGithubConfirmationHttp(qx, channel.integrationId) let http: ConnectorHttp try { @@ -163,33 +167,53 @@ export async function runShadowDiffForChannel( } catch (err) { svc.log.warn( { err, unitId: unit.id, day, channelName: channel.channelName }, - 'failed to set up github client for force-push confirmation, keeping candidates as missing_in_shadow', + 'failed to set up github client for missing_in_shadow confirmation, keeping candidates as missing_in_shadow', ) unitDiffs.push({ unit, result }) continue } - const { mismatches, skippedCount, failedCount } = await dropConfirmedForcePushedCommits( - unit.syncName, - result.mismatches, - owner, - repo, - http, - svc.log, - ) - if (skippedCount > 0) { - svc.log.info( - { unitId: unit.id, day, channelName: channel.channelName, skippedCount }, - 'skipped force-pushed commits confirmed missing from github', + + if (needsForcePushCheck) { + const { mismatches, skippedCount, failedCount } = await dropConfirmedForcePushedCommits( + unit.syncName, + result.mismatches, + owner, + repo, + http, + svc.log, ) + if (skippedCount > 0) { + svc.log.info( + { unitId: unit.id, day, channelName: channel.channelName, skippedCount }, + 'skipped force-pushed commits confirmed missing from github', + ) + } + if (failedCount > 0) { + svc.log.warn( + { unitId: unit.id, day, channelName: channel.channelName, failedCount }, + 'failed to confirm force-pushed commit candidates against github, keeping them as missing_in_shadow', + ) + } + result = { ...result, mismatches } } - if (failedCount > 0) { - svc.log.warn( - { unitId: unit.id, day, channelName: channel.channelName, failedCount }, - 'failed to confirm force-pushed commit candidates against github, keeping them as missing_in_shadow', - ) + + if (hasDeletedRecordCandidates(unit.syncName, result.mismatches)) { + const { mismatches, confirmedDeletedCount, keptUnconfirmedCount } = + await dropConfirmedDeletedRecords(unit.syncName, result.mismatches, http, svc.log) + if (confirmedDeletedCount > 0) { + svc.log.info( + { unitId: unit.id, day, channelName: channel.channelName, confirmedDeletedCount }, + 'skipped records confirmed deleted on github', + ) + } + if (keptUnconfirmedCount > 0) { + svc.log.warn( + { unitId: unit.id, day, channelName: channel.channelName, keptUnconfirmedCount }, + 'failed to confirm deleted-record candidates against github, keeping them as missing_in_shadow', + ) + } + result = { ...result, mismatches } } - unitDiffs.push({ unit, result: { ...result, mismatches } }) - continue } unitDiffs.push({ unit, result }) diff --git a/services/apps/connectors_worker/src/deletedRecords.ts b/services/apps/connectors_worker/src/deletedRecords.ts new file mode 100644 index 0000000000..6154c045b6 --- /dev/null +++ b/services/apps/connectors_worker/src/deletedRecords.ts @@ -0,0 +1,128 @@ +import type { ConnectorHttp } from '@crowd/connectors' +import { mapWithConcurrency } from '@crowd/connectors' +import type { Logger } from '@crowd/logging' + +import { IShadowDiffMismatch } from './shadowDiff' + +const COMMIT_SYNC_NAME = 'pull-request-commits' +const SYNTHETIC_SOURCE_ID_PREFIX = 'gen-' +const NODES_QUERY_BATCH_SIZE = 100 +const CONFIRM_CONCURRENCY = 3 +const CONFIRM_BUDGET_MS = 60_000 +const CONFIRM_REQUEST_TIMEOUT_MS = 10_000 +const CONFIRM_REQUEST_MAX_ATTEMPTS = 1 + +const NODES_QUERY = `query($ids: [ID!]!) { nodes(ids: $ids) { id } }` + +interface GraphqlEnvelope { + data?: T + errors?: { type?: string; message?: string }[] +} + +interface NodesQueryResult { + nodes: ({ id: string } | null)[] +} + +function isNodeIdCandidate(mismatch: IShadowDiffMismatch): boolean { + return !mismatch.sourceId.startsWith(SYNTHETIC_SOURCE_ID_PREFIX) +} + +export function hasDeletedRecordCandidates( + syncName: string, + mismatches: IShadowDiffMismatch[], +): boolean { + return ( + syncName !== COMMIT_SYNC_NAME && + mismatches.some((m) => m.kind === 'missing_in_shadow' && isNodeIdCandidate(m)) + ) +} + +export interface DeletedRecordFilterResult { + mismatches: IShadowDiffMismatch[] + confirmedDeletedCount: number + keptUnconfirmedCount: number +} + +function toBatches(items: T[], size: number): T[][] { + const batches: T[][] = [] + for (let i = 0; i < items.length; i += size) { + batches.push(items.slice(i, i + size)) + } + return batches +} + +async function fetchExistingNodeIds( + http: ConnectorHttp, + ids: string[], + log: Logger, +): Promise | null> { + try { + const body = await http.request>( + { + method: 'post', + url: 'https://api.github.com/graphql', + data: { query: NODES_QUERY, variables: { ids } }, + timeout: CONFIRM_REQUEST_TIMEOUT_MS, + }, + log, + CONFIRM_REQUEST_MAX_ATTEMPTS, + ) + if (!body.data) { + return null + } + return new Set(body.data.nodes.filter((n): n is { id: string } => n !== null).map((n) => n.id)) + } catch { + return null + } +} + +export async function dropConfirmedDeletedRecords( + syncName: string, + mismatches: IShadowDiffMismatch[], + http: ConnectorHttp, + log: Logger, +): Promise { + if (!hasDeletedRecordCandidates(syncName, mismatches)) { + return { mismatches, confirmedDeletedCount: 0, keptUnconfirmedCount: 0 } + } + + const candidates = mismatches.filter( + (m) => m.kind === 'missing_in_shadow' && isNodeIdCandidate(m), + ) + const candidateIds = [...new Set(candidates.map((m) => m.sourceId))] + const batches = toBatches(candidateIds, NODES_QUERY_BATCH_SIZE) + const deadline = Date.now() + CONFIRM_BUDGET_MS + + const confirmedDeletedIds = new Set() + const uncheckedIds = new Set() + + await mapWithConcurrency(batches, CONFIRM_CONCURRENCY, async (batch) => { + if (Date.now() >= deadline) { + for (const id of batch) uncheckedIds.add(id) + return + } + const existingIds = await fetchExistingNodeIds(http, batch, log) + if (existingIds === null) { + for (const id of batch) uncheckedIds.add(id) + return + } + for (const id of batch) { + if (!existingIds.has(id)) { + confirmedDeletedIds.add(id) + } + } + }) + + return { + mismatches: mismatches.filter( + (m) => + !( + m.kind === 'missing_in_shadow' && + isNodeIdCandidate(m) && + confirmedDeletedIds.has(m.sourceId) + ), + ), + confirmedDeletedCount: confirmedDeletedIds.size, + keptUnconfirmedCount: uncheckedIds.size, + } +} From 3a57bfb12da9a8cd8843189bd853e4122e7a107f Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 25 Sep 2026 16:29:53 +0000 Subject: [PATCH 2/3] fix: confirm node deletion only on NOT_FOUND errors (CM-1473) Signed-off-by: Mouad BANI --- .../connectors_worker/src/deletedRecords.ts | 44 ++++++++++++++----- 1 file changed, 32 insertions(+), 12 deletions(-) diff --git a/services/apps/connectors_worker/src/deletedRecords.ts b/services/apps/connectors_worker/src/deletedRecords.ts index 6154c045b6..6533b311a2 100644 --- a/services/apps/connectors_worker/src/deletedRecords.ts +++ b/services/apps/connectors_worker/src/deletedRecords.ts @@ -16,7 +16,7 @@ const NODES_QUERY = `query($ids: [ID!]!) { nodes(ids: $ids) { id } }` interface GraphqlEnvelope { data?: T - errors?: { type?: string; message?: string }[] + errors?: { type?: string; message?: string; path?: (string | number)[] }[] } interface NodesQueryResult { @@ -51,11 +51,16 @@ function toBatches(items: T[], size: number): T[][] { return batches } -async function fetchExistingNodeIds( +interface BatchConfirmation { + deletedIds: Set + unconfirmedIds: Set +} + +async function confirmDeletedNodeIds( http: ConnectorHttp, ids: string[], log: Logger, -): Promise | null> { +): Promise { try { const body = await http.request>( { @@ -67,10 +72,28 @@ async function fetchExistingNodeIds( log, CONFIRM_REQUEST_MAX_ATTEMPTS, ) - if (!body.data) { + if (!body.data || body.data.nodes.length !== ids.length) { return null } - return new Set(body.data.nodes.filter((n): n is { id: string } => n !== null).map((n) => n.id)) + const notFoundIndexes = new Set( + (body.errors ?? []) + .filter((e) => e.type === 'NOT_FOUND' && e.path?.[0] === 'nodes') + .map((e) => e.path?.[1]) + .filter((i): i is number => typeof i === 'number'), + ) + const deletedIds = new Set() + const unconfirmedIds = new Set() + body.data.nodes.forEach((node, i) => { + if (node !== null) { + return + } + if (notFoundIndexes.has(i)) { + deletedIds.add(ids[i]) + } else { + unconfirmedIds.add(ids[i]) + } + }) + return { deletedIds, unconfirmedIds } } catch { return null } @@ -101,16 +124,13 @@ export async function dropConfirmedDeletedRecords( for (const id of batch) uncheckedIds.add(id) return } - const existingIds = await fetchExistingNodeIds(http, batch, log) - if (existingIds === null) { + const confirmation = await confirmDeletedNodeIds(http, batch, log) + if (confirmation === null) { for (const id of batch) uncheckedIds.add(id) return } - for (const id of batch) { - if (!existingIds.has(id)) { - confirmedDeletedIds.add(id) - } - } + for (const id of confirmation.deletedIds) confirmedDeletedIds.add(id) + for (const id of confirmation.unconfirmedIds) uncheckedIds.add(id) }) return { From 19113ba2d307b1de99a40342a342ac7f850e9856 Mon Sep 17 00:00:00 2001 From: Mouad BANI Date: Fri, 25 Sep 2026 16:31:23 +0000 Subject: [PATCH 3/3] test: cover deleted record confirmation in shadow diff (CM-1473) Signed-off-by: Mouad BANI --- .../src/deletedRecords.test.ts | 175 ++++++++++++++++++ 1 file changed, 175 insertions(+) create mode 100644 services/apps/connectors_worker/src/deletedRecords.test.ts diff --git a/services/apps/connectors_worker/src/deletedRecords.test.ts b/services/apps/connectors_worker/src/deletedRecords.test.ts new file mode 100644 index 0000000000..d2a2674fa0 --- /dev/null +++ b/services/apps/connectors_worker/src/deletedRecords.test.ts @@ -0,0 +1,175 @@ +import { describe, expect, it, vi } from 'vitest' + +import type { ConnectorHttp } from '@crowd/connectors' +import type { Logger } from '@crowd/logging' + +import { dropConfirmedDeletedRecords, hasDeletedRecordCandidates } from './deletedRecords' +import { IShadowDiffMismatch } from './shadowDiff' + +const log = { info: vi.fn(), warn: vi.fn(), error: vi.fn() } as unknown as Logger + +function missingMismatch(sourceId: string, type = 'issue-comment'): IShadowDiffMismatch { + return { sourceId, type, kind: 'missing_in_shadow', severity: 'high' } +} + +function httpWithHandler(handler: (ids: string[]) => Promise): ConnectorHttp { + return { + request: async (config: { data: { variables: { ids: string[] } } }) => + handler(config.data.variables.ids), + requestCount: () => 0, + } as unknown as ConnectorHttp +} + +describe('hasDeletedRecordCandidates', () => { + it('is false for the pull-request-commits sync regardless of candidates', () => { + expect(hasDeletedRecordCandidates('pull-request-commits', [missingMismatch('IC_abc')])).toBe( + false, + ) + }) + + it('is false when the only missing_in_shadow records have synthetic sourceIds', () => { + expect( + hasDeletedRecordCandidates('pull-requests', [missingMismatch('gen-CE_PR_1_bob_2026')]), + ).toBe(false) + }) + + it('is true when a real node-id missing_in_shadow record exists', () => { + expect(hasDeletedRecordCandidates('issue-comments', [missingMismatch('IC_abc')])).toBe(true) + }) +}) + +describe('dropConfirmedDeletedRecords', () => { + it('drops records confirmed deleted (null node with NOT_FOUND) and keeps records still alive', async () => { + const mismatches = [missingMismatch('IC_deleted'), missingMismatch('IC_alive')] + const http = httpWithHandler(async (ids) => ({ + data: { nodes: ids.map((id) => (id === 'IC_alive' ? { id } : null)) }, + errors: ids.flatMap((id, i) => + id === 'IC_alive' ? [] : [{ type: 'NOT_FOUND', path: ['nodes', i] }], + ), + })) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual([missingMismatch('IC_alive')]) + expect(result.confirmedDeletedCount).toBe(1) + expect(result.keptUnconfirmedCount).toBe(0) + }) + + it('keeps records when the graphql request errors', async () => { + const mismatches = [missingMismatch('IC_unknown')] + const http = httpWithHandler(async () => { + throw new Error('network exploded') + }) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual(mismatches) + expect(result.confirmedDeletedCount).toBe(0) + expect(result.keptUnconfirmedCount).toBe(1) + }) + + it('keeps records when the response has no top-level data', async () => { + const mismatches = [missingMismatch('IC_unknown')] + const http = httpWithHandler(async () => ({ errors: [{ type: 'FORBIDDEN' }] })) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual(mismatches) + expect(result.confirmedDeletedCount).toBe(0) + expect(result.keptUnconfirmedCount).toBe(1) + }) + + it('treats NOT_FOUND per-id errors alongside data as a valid confirmation', async () => { + const mismatches = [missingMismatch('IC_deleted')] + const http = httpWithHandler(async (ids) => ({ + data: { nodes: ids.map(() => null) }, + errors: [{ type: 'NOT_FOUND', message: 'Could not resolve to a node', path: ['nodes', 0] }], + })) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual([]) + expect(result.confirmedDeletedCount).toBe(1) + }) + + it('keeps null nodes whose error is not NOT_FOUND (e.g. FORBIDDEN)', async () => { + const mismatches = [missingMismatch('IC_forbidden'), missingMismatch('IC_deleted')] + const http = httpWithHandler(async (ids) => ({ + data: { nodes: ids.map(() => null) }, + errors: [ + { type: 'FORBIDDEN', path: ['nodes', 0] }, + { type: 'NOT_FOUND', path: ['nodes', 1] }, + ], + })) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual([missingMismatch('IC_forbidden')]) + expect(result.confirmedDeletedCount).toBe(1) + expect(result.keptUnconfirmedCount).toBe(1) + }) + + it('keeps null nodes that have no matching error entry at all', async () => { + const mismatches = [missingMismatch('IC_unresolved')] + const http = httpWithHandler(async (ids) => ({ + data: { nodes: ids.map(() => null) }, + })) + + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + + expect(result.mismatches).toEqual(mismatches) + expect(result.confirmedDeletedCount).toBe(0) + expect(result.keptUnconfirmedCount).toBe(1) + }) + + it('never checks pull-request-commits candidates and leaves them untouched', async () => { + const mismatches = [missingMismatch('deadbeef', 'authored-commit')] + const requestSpy = vi.fn() + const http = { request: requestSpy, requestCount: () => 0 } as unknown as ConnectorHttp + + const result = await dropConfirmedDeletedRecords('pull-request-commits', mismatches, http, log) + + expect(result.mismatches).toEqual(mismatches) + expect(requestSpy).not.toHaveBeenCalled() + }) + + it('ignores synthetic gen- sourceIds since they are not graphql node ids', async () => { + const genMismatch = missingMismatch('gen-CE_PR_1_bob_2026', 'pull_request-closed') + const nodeMismatch = missingMismatch('PR_1', 'pull_request-opened') + const requestSpy = vi.fn(async (config: { data: { variables: { ids: string[] } } }) => ({ + data: { nodes: config.data.variables.ids.map(() => null) }, + errors: config.data.variables.ids.map((_, i) => ({ + type: 'NOT_FOUND', + path: ['nodes', i], + })), + })) + const http = { request: requestSpy, requestCount: () => 0 } as unknown as ConnectorHttp + + const result = await dropConfirmedDeletedRecords( + 'pull-requests', + [genMismatch, nodeMismatch], + http, + log, + ) + + expect(result.mismatches).toEqual([genMismatch]) + expect(requestSpy).toHaveBeenCalledTimes(1) + expect(requestSpy.mock.calls[0][0].data.variables.ids).toEqual(['PR_1']) + }) + + it('keeps every candidate unconfirmed once the time budget is already exhausted', async () => { + const manyIds = Array.from({ length: 150 }, (_, i) => `IC_${i}`) + const mismatches = manyIds.map((id) => missingMismatch(id)) + const requestSpy = vi.fn() + const http = { request: requestSpy, requestCount: () => 0 } as unknown as ConnectorHttp + + vi.spyOn(Date, 'now').mockReturnValueOnce(0).mockReturnValue(120_000) + const result = await dropConfirmedDeletedRecords('issue-comments', mismatches, http, log) + vi.restoreAllMocks() + + expect(requestSpy).not.toHaveBeenCalled() + expect(result.confirmedDeletedCount).toBe(0) + expect(result.keptUnconfirmedCount).toBe(manyIds.length) + expect(result.mismatches).toEqual(mismatches) + }) +})