Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -153,43 +154,66 @@ export async function runShadowDiffForChannel(
const unitDiffs: { unit: IShadowDiffUnit; result: IShadowDiffUnitResult }[] = []
let confirmationHttp: Promise<ConnectorHttp> | 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 {
http = await confirmationHttp
} 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)
Comment on lines +200 to +202
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 })
Expand Down
175 changes: 175 additions & 0 deletions services/apps/connectors_worker/src/deletedRecords.test.ts
Original file line number Diff line number Diff line change
@@ -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<unknown>): 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)
})
})
Loading
Loading