diff --git a/backend/src/database/migrations/V1789112777__skip-reason-to-project-catalog.sql b/backend/src/database/migrations/V1789112777__skip-reason-to-project-catalog.sql new file mode 100644 index 0000000000..f7d80e4b34 --- /dev/null +++ b/backend/src/database/migrations/V1789112777__skip-reason-to-project-catalog.sql @@ -0,0 +1,2 @@ +ALTER TABLE "projectCatalog" + ADD COLUMN IF NOT EXISTS "skipReason" TEXT; diff --git a/services/apps/automatic_onboarding_worker/src/activities/activities.ts b/services/apps/automatic_onboarding_worker/src/activities/activities.ts index 1463ca00e3..43964ffb54 100644 --- a/services/apps/automatic_onboarding_worker/src/activities/activities.ts +++ b/services/apps/automatic_onboarding_worker/src/activities/activities.ts @@ -2,14 +2,22 @@ import { findProjectCatalogById, findProjectCatalogPendingOnboarding, markProjectCatalogOnboardingFailed, + markProjectCatalogOnboardingSkipped, updateProjectCatalog, } from '@crowd/data-access-layer' +import { + InsightsProjectField, + findInsightsProjectBySlugIncludingDeleted, +} from '@crowd/data-access-layer/src/collections' import { IDbProjectCatalog } from '@crowd/data-access-layer/src/project-catalog/types' import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor' import { getServiceLogger } from '@crowd/logging' import { svc } from '../main' -import { onboardProject } from '../onboarder/onboarder' +import { deriveProjectSlug, onboardProject } from '../onboarder/onboarder' +import { OnboardAndUpdateProjectOutcome } from '../types' + +import { buildInsightsProjectSkipReason } from './insightsProjectSkip' const log = getServiceLogger() @@ -33,7 +41,17 @@ async function findAlreadyOnboarded( return fresh?.onboardedAt ? fresh : null } -export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promise { +async function findDeletedInsightsProjectBySlug(qx: ReturnType, projectSlug: string) { + const slug = deriveProjectSlug(projectSlug) + const insightsProject = await findInsightsProjectBySlugIncludingDeleted(qx, slug, [ + InsightsProjectField.DELETED_AT, + ]) + return insightsProject?.deletedAt ? insightsProject : null +} + +export async function onboardAndUpdateProject( + project: IDbProjectCatalog, +): Promise { const qx = pgpQx(svc.postgres.writer.connection()) const startTime = Date.now() @@ -44,7 +62,38 @@ export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promi { id: project.id, repoUrl: project.repoUrl, onboardedAt: fresh.onboardedAt }, 'Project already onboarded, skipping API call.', ) - return + return 'already-onboarded' + } + + const deletedInsightsProject = await findDeletedInsightsProjectBySlug(qx, project.projectSlug) + if (deletedInsightsProject) { + const reason = buildInsightsProjectSkipReason( + project.projectSlug, + deletedInsightsProject.deletedAt, + ) + const updatedRows = await markProjectCatalogOnboardingSkipped(qx, project.id, reason) + if (updatedRows > 0) { + log.info({ id: project.id, repoUrl: project.repoUrl, reason }, 'Onboarding skipped.') + return 'skipped' + } + + const current = await findProjectCatalogById(qx, project.id) + if (current?.onboardedAt) { + log.info( + { id: project.id, repoUrl: project.repoUrl, onboardedAt: current.onboardedAt }, + 'Project already onboarded, skipping API call.', + ) + return 'already-onboarded' + } + if (current?.action === 'skip') { + log.info({ id: project.id, repoUrl: project.repoUrl }, 'Project already skipped.') + return 'skipped' + } + log.info( + { id: project.id, repoUrl: project.repoUrl, action: current?.action }, + 'Skip guard was a no-op; project catalog row changed concurrently, not calling onboarding API.', + ) + return 'catalog-changed' } log.info({ id: project.id, repoUrl: project.repoUrl }, 'Starting onboarding.') @@ -72,6 +121,8 @@ export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promi { id: project.id, repoUrl: project.repoUrl, segmentId: result.segmentId, elapsedSeconds }, 'Onboarding complete.', ) + + return 'onboarded' } export async function markProjectOnboardingFailed( diff --git a/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.test.ts b/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.test.ts new file mode 100644 index 0000000000..80e7d69858 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.test.ts @@ -0,0 +1,16 @@ +import { describe, expect, it } from 'vitest' + +import { buildInsightsProjectSkipReason } from './insightsProjectSkip' + +describe('buildInsightsProjectSkipReason', () => { + it('normalizes the project slug and includes the deletion date', () => { + const reason = buildInsightsProjectSkipReason( + 'nonlf_gerritcodereview-gerrit', + '2026-04-10T00:00:00.000Z', + ) + + expect(reason).toBe( + "Insights project 'nonlf-gerritcodereview-gerrit' was deleted on 2026-04-10T00:00:00.000Z; onboarding skipped for manual review", + ) + }) +}) diff --git a/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.ts b/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.ts new file mode 100644 index 0000000000..f8fddd6810 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/activities/insightsProjectSkip.ts @@ -0,0 +1,8 @@ +import { deriveProjectSlug } from '../onboarder/onboarder' + +// A soft-deleted insightsProjects row still owns its slug (the unique index can't be made +// partial on deletedAt: three FKs reference it), so segment creation would 500 on it. +export function buildInsightsProjectSkipReason(projectSlug: string, deletedAt: string): string { + const slug = deriveProjectSlug(projectSlug) + return `Insights project '${slug}' was deleted on ${deletedAt}; onboarding skipped for manual review` +} diff --git a/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.test.ts b/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.test.ts new file mode 100644 index 0000000000..3b87aaae33 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.test.ts @@ -0,0 +1,27 @@ +import { describe, expect, it } from 'vitest' + +import { readErrorBody } from './onboarder' + +describe('readErrorBody', () => { + it('returns the response body text', async () => { + const response = new Response('{"error":"insightsProjects slug already exists"}') + + expect(await readErrorBody(response)).toBe('{"error":"insightsProjects slug already exists"}') + }) + + it('truncates a body longer than 500 characters', async () => { + const response = new Response('a'.repeat(600)) + + const body = await readErrorBody(response) + + expect(body).toBe(`${'a'.repeat(500)}…`) + }) + + it('returns an empty string when the body cannot be read', async () => { + const response = new Response(null) + // Consuming the body once locks the stream, so a second read fails. + await response.text() + + expect(await readErrorBody(response)).toBe('') + }) +}) diff --git a/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.ts b/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.ts index c82c49340a..ec822065f0 100644 --- a/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.ts +++ b/services/apps/automatic_onboarding_worker/src/onboarder/onboarder.ts @@ -12,6 +12,16 @@ interface ISegmentQueryResponse { const SEGMENT_QUERY_PAGE_SIZE = 20 const BACKEND_REQUEST_TIMEOUT_MS = 30_000 const GITHUB_REQUEST_TIMEOUT_MS = 10_000 +const ERROR_BODY_MAX_LENGTH = 500 + +export async function readErrorBody(response: Response): Promise { + try { + const body = await response.text() + return body.length > ERROR_BODY_MAX_LENGTH ? `${body.slice(0, ERROR_BODY_MAX_LENGTH)}…` : body + } catch { + return '' + } +} export function deriveProjectName(repoName: string): string { return repoName @@ -99,7 +109,9 @@ async function queryProjectByName( }) if (!response.ok) { - throw new Error(`Segment query returned HTTP ${response.status}: ${response.statusText}`) + throw new Error( + `Segment query returned HTTP ${response.status}: ${response.statusText} - ${await readErrorBody(response)}`, + ) } const body = (await response.json()) as ISegmentQueryResponse @@ -139,7 +151,9 @@ async function createProjectSegment( }) if (!response.ok) { - throw new Error(`Segment creation returned HTTP ${response.status}: ${response.statusText}`) + throw new Error( + `Segment creation returned HTTP ${response.status}: ${response.statusText} - ${await readErrorBody(response)}`, + ) } // POST /segment/project does not return the created segment's id in its response body; re-query by name to get it. @@ -211,7 +225,9 @@ async function createGithubIntegration( }) if (!response.ok) { - throw new Error(`GitHub integration returned HTTP ${response.status}: ${response.statusText}`) + throw new Error( + `GitHub integration returned HTTP ${response.status}: ${response.statusText} - ${await readErrorBody(response)}`, + ) } } diff --git a/services/apps/automatic_onboarding_worker/src/types.ts b/services/apps/automatic_onboarding_worker/src/types.ts index cb0f970d90..59786e75ca 100644 --- a/services/apps/automatic_onboarding_worker/src/types.ts +++ b/services/apps/automatic_onboarding_worker/src/types.ts @@ -1,3 +1,9 @@ export interface IOnboardProjectsInput { batchSize?: number } + +export type OnboardAndUpdateProjectOutcome = + | 'onboarded' + | 'skipped' + | 'already-onboarded' + | 'catalog-changed' diff --git a/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts b/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts index 48fa3587e5..890dd1c307 100644 --- a/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts +++ b/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts @@ -36,6 +36,8 @@ export async function onboardProjects(input: IOnboardProjectsInput = {}): Promis log.info(`Onboarding ${projects.length} project(s) (batch size: ${batchSize}).`) let succeeded = 0 + let skipped = 0 + let racedOut = 0 let failed = 0 for (let i = 0; i < projects.length; i++) { @@ -43,8 +45,14 @@ export async function onboardProjects(input: IOnboardProjectsInput = {}): Promis log.info(`[${i + 1}/${projects.length}] Onboarding: ${project.repoUrl}`) try { - await onboardActivities.onboardAndUpdateProject(project) - succeeded++ + const outcome = await onboardActivities.onboardAndUpdateProject(project) + if (outcome === 'skipped') { + skipped++ + } else if (outcome === 'catalog-changed') { + racedOut++ + } else { + succeeded++ + } } catch (err) { // Activity-level retries are already exhausted at this point — mark as a // terminal error so the daily schedule stops retrying this project forever. @@ -64,6 +72,6 @@ export async function onboardProjects(input: IOnboardProjectsInput = {}): Promis } log.info( - `Batch onboarding complete. total=${projects.length} succeeded=${succeeded} failed=${failed}`, + `Batch onboarding complete. total=${projects.length} succeeded=${succeeded} skipped=${skipped} racedOut=${racedOut} failed=${failed}`, ) } diff --git a/services/libs/data-access-layer/src/collections/index.test.ts b/services/libs/data-access-layer/src/collections/index.test.ts new file mode 100644 index 0000000000..e526a28e1a --- /dev/null +++ b/services/libs/data-access-layer/src/collections/index.test.ts @@ -0,0 +1,55 @@ +import { test as base, describe, expect } from 'vitest' + +import { withQx } from '@crowd/test-kit/db' + +import { + InsightsProjectField, + createInsightsProject, + deleteInsightsProject, + findInsightsProjectBySlugIncludingDeleted, +} from './index' + +const test = withQx(base) + +describe('findInsightsProjectBySlugIncludingDeleted', () => { + test('finds a soft-deleted row by slug, with deletedAt set', async ({ qx }) => { + const created = await createInsightsProject(qx, { + name: 'Gerrit', + slug: 'gerritcodereview-gerrit', + isLF: false, + }) + await deleteInsightsProject(qx, created.id) + + const found = await findInsightsProjectBySlugIncludingDeleted(qx, 'gerritcodereview-gerrit', [ + InsightsProjectField.ID, + InsightsProjectField.DELETED_AT, + ]) + + expect(found?.id).toBe(created.id) + expect(found?.deletedAt).not.toBeNull() + }) + + test('finds a live row by slug, with deletedAt null', async ({ qx }) => { + const created = await createInsightsProject(qx, { + name: 'Kubernetes', + slug: 'kubernetes-kubernetes', + isLF: true, + }) + + const found = await findInsightsProjectBySlugIncludingDeleted(qx, 'kubernetes-kubernetes', [ + InsightsProjectField.ID, + InsightsProjectField.DELETED_AT, + ]) + + expect(found?.id).toBe(created.id) + expect(found?.deletedAt).toBeNull() + }) + + test('returns null when no row has the slug', async ({ qx }) => { + const found = await findInsightsProjectBySlugIncludingDeleted(qx, 'does-not-exist', [ + InsightsProjectField.ID, + ]) + + expect(found).toBeNull() + }) +}) diff --git a/services/libs/data-access-layer/src/collections/index.ts b/services/libs/data-access-layer/src/collections/index.ts index be5bc75aa9..d62968dbf9 100644 --- a/services/libs/data-access-layer/src/collections/index.ts +++ b/services/libs/data-access-layer/src/collections/index.ts @@ -196,6 +196,20 @@ export async function queryInsightsProjects( return queryTable(qx, 'insightsProjects', Object.values(InsightsProjectField), opts) } +export async function findInsightsProjectBySlugIncludingDeleted( + qx: QueryExecutor, + slug: string, + fields: T[], +): Promise | null> { + const rows = await queryTable(qx, 'insightsProjects', Object.values(InsightsProjectField), { + fields, + filter: { slug: { eq: slug } }, + limit: 1, + }) + + return rows.length > 0 ? rows[0] : null +} + export async function createInsightsProject( qx: QueryExecutor, insightProject: Partial, diff --git a/services/libs/data-access-layer/src/project-catalog/projectCatalog.test.ts b/services/libs/data-access-layer/src/project-catalog/projectCatalog.test.ts new file mode 100644 index 0000000000..f0bfcfec1a --- /dev/null +++ b/services/libs/data-access-layer/src/project-catalog/projectCatalog.test.ts @@ -0,0 +1,73 @@ +import { test as base, describe, expect } from 'vitest' + +import { withQx } from '@crowd/test-kit/db' + +import { + findProjectCatalogById, + insertProjectCatalog, + markProjectCatalogOnboardingSkipped, + updateProjectCatalog, +} from './projectCatalog' + +const test = withQx(base) + +function catalogRow(overrides: Partial[1]> = {}) { + return { + projectSlug: 'gerritcodereview-gerrit', + repoName: 'gerrit', + repoUrl: 'https://github.com/gerritcodereview/gerrit', + action: 'onboard' as const, + ...overrides, + } +} + +describe('markProjectCatalogOnboardingSkipped', () => { + test('transitions a pending row to skip with the reason, and clears a prior onboarding error', async ({ + qx, + }) => { + const inserted = await insertProjectCatalog(qx, catalogRow()) + await updateProjectCatalog(qx, inserted.id, { + onboardingError: 'Segment creation returned HTTP 500: Internal Server Error', + }) + + const updatedRows = await markProjectCatalogOnboardingSkipped( + qx, + inserted.id, + "Insights project 'gerritcodereview-gerrit' was deleted on 2026-04-10; onboarding skipped for manual review", + ) + + const row = await findProjectCatalogById(qx, inserted.id) + expect(updatedRows).toBe(1) + expect(row?.action).toBe('skip') + expect(row?.skipReason).toBe( + "Insights project 'gerritcodereview-gerrit' was deleted on 2026-04-10; onboarding skipped for manual review", + ) + expect(row?.onboardingError).toBeNull() + }) + + test('does not touch a row whose action is no longer onboard', async ({ qx }) => { + const inserted = await insertProjectCatalog(qx, catalogRow({ action: 'error' })) + + const updatedRows = await markProjectCatalogOnboardingSkipped(qx, inserted.id, 'some reason') + + const row = await findProjectCatalogById(qx, inserted.id) + expect(updatedRows).toBe(0) + expect(row?.action).toBe('error') + expect(row?.skipReason).toBeNull() + }) + + test('does not touch a row already onboarded', async ({ qx }) => { + const inserted = await insertProjectCatalog(qx, catalogRow()) + await updateProjectCatalog(qx, inserted.id, { + action: 'onboarded', + onboardedAt: new Date().toISOString(), + }) + + const updatedRows = await markProjectCatalogOnboardingSkipped(qx, inserted.id, 'some reason') + + const row = await findProjectCatalogById(qx, inserted.id) + expect(updatedRows).toBe(0) + expect(row?.action).toBe('onboarded') + expect(row?.skipReason).toBeNull() + }) +}) diff --git a/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts b/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts index 63c72b42a9..f4768f8262 100644 --- a/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts +++ b/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts @@ -23,6 +23,7 @@ const PROJECT_CATALOG_COLUMNS = [ 'evaluatedAt', 'onboardedAt', 'onboardingError', + 'skipReason', 'syncedAt', 'createdAt', 'updatedAt', @@ -518,6 +519,10 @@ export async function updateProjectCatalog( setClauses.push('"onboardingError" = $(onboardingError)') params.onboardingError = data.onboardingError } + if (data.skipReason !== undefined) { + setClauses.push('"skipReason" = $(skipReason)') + params.skipReason = data.skipReason + } if (setClauses.length === 0) { return findProjectCatalogById(qx, id) @@ -551,6 +556,23 @@ export async function markProjectCatalogOnboardingFailed( ) } +// Guarded like markProjectCatalogOnboardingFailed: a concurrent manual action wins. +// onboardingError is cleared in case this row was previously failed and requeued. +export async function markProjectCatalogOnboardingSkipped( + qx: QueryExecutor, + id: string, + reason: string, +): Promise { + return qx.result( + ` + UPDATE "projectCatalog" + SET "action" = 'skip', "skipReason" = $(reason), "onboardingError" = NULL, "updatedAt" = NOW() + WHERE id = $(id) AND "action" = 'onboard' AND "onboardedAt" IS NULL + `, + { id, reason }, + ) +} + // Guarded like markProjectCatalogOnboardingFailed above: a manual request // (POST /project-catalog) may have moved the row out of 'evaluate' while this // evaluation was in flight — in that case the manual action wins and this diff --git a/services/libs/data-access-layer/src/project-catalog/types.ts b/services/libs/data-access-layer/src/project-catalog/types.ts index f3822f3bb0..abfb1af250 100644 --- a/services/libs/data-access-layer/src/project-catalog/types.ts +++ b/services/libs/data-access-layer/src/project-catalog/types.ts @@ -25,6 +25,7 @@ export interface IDbProjectCatalog { evaluatedAt: string | null onboardedAt: string | null onboardingError: string | null + skipReason: string | null syncedAt: string | null createdAt: string | null updatedAt: string | null @@ -41,6 +42,7 @@ type ProjectCatalogWritable = Pick< | 'evaluationResult' | 'evaluationReason' | 'onboardingError' + | 'skipReason' > export type IDbProjectCatalogCreate = Omit< @@ -51,6 +53,7 @@ export type IDbProjectCatalogCreate = Omit< | 'evaluationResult' | 'evaluationReason' | 'onboardingError' + | 'skipReason' > & { source?: string | null action?: ProjectCatalogAction @@ -58,6 +61,7 @@ export type IDbProjectCatalogCreate = Omit< evaluationResult?: string | null evaluationReason?: string | null onboardingError?: string | null + skipReason?: string | null } export type IDbProjectCatalogUpdate = Partial & {