diff --git a/backend/src/database/migrations/V1787916364__onboarding-error-to-project-catalog.sql b/backend/src/database/migrations/V1787916364__onboarding-error-to-project-catalog.sql new file mode 100644 index 0000000000..a8ea7bc2ef --- /dev/null +++ b/backend/src/database/migrations/V1787916364__onboarding-error-to-project-catalog.sql @@ -0,0 +1,2 @@ +ALTER TABLE "projectCatalog" + ADD COLUMN IF NOT EXISTS "onboardingError" TEXT; diff --git a/services/apps/automatic_onboarding_worker/src/activities.ts b/services/apps/automatic_onboarding_worker/src/activities.ts index 336ce12bb9..3662234550 100644 --- a/services/apps/automatic_onboarding_worker/src/activities.ts +++ b/services/apps/automatic_onboarding_worker/src/activities.ts @@ -1 +1 @@ -export {} +export * from './activities/activities' diff --git a/services/apps/automatic_onboarding_worker/src/activities/activities.ts b/services/apps/automatic_onboarding_worker/src/activities/activities.ts new file mode 100644 index 0000000000..9236359346 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/activities/activities.ts @@ -0,0 +1,93 @@ +import { + findProjectCatalogById, + findProjectCatalogPendingOnboarding, + markProjectCatalogOnboardingFailed, + updateProjectCatalog, +} from '@crowd/data-access-layer' +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' + +const log = getServiceLogger() + +export async function fetchProjectsPendingOnboarding( + batchSize: number, +): Promise { + const qx = pgpQx(svc.postgres.reader.connection()) + + const projects = await findProjectCatalogPendingOnboarding(qx, { limit: batchSize }) + + log.info({ count: projects.length, batchSize }, 'Fetched projects pending onboarding.') + + return projects +} + +async function findAlreadyOnboarded( + qx: ReturnType, + projectId: string, +): Promise { + const fresh = await findProjectCatalogById(qx, projectId) + return fresh?.onboardedAt ? fresh : null +} + +export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promise { + const qx = pgpQx(svc.postgres.writer.connection()) + const startTime = Date.now() + + // Guard: uses the writer connection to avoid replica lag missing a just-written onboardedAt. + const fresh = await findAlreadyOnboarded(qx, project.id) + if (fresh) { + log.info( + { id: project.id, repoUrl: project.repoUrl, onboardedAt: fresh.onboardedAt }, + 'Project already onboarded, skipping API call.', + ) + return + } + + log.info({ id: project.id, repoUrl: project.repoUrl }, 'Starting onboarding.') + + const result = await onboardProject({ + id: project.id, + repoUrl: project.repoUrl, + repoName: project.repoName, + projectSlug: project.projectSlug, + }) + + if (result.outcome === 'error') { + throw new Error(result.error ?? 'Unknown onboarding error') + } + + await updateProjectCatalog(qx, project.id, { + onboardedAt: new Date().toISOString(), + onboardingError: null, + }) + + const elapsedSeconds = ((Date.now() - startTime) / 1000).toFixed(1) + + log.info( + { id: project.id, repoUrl: project.repoUrl, segmentId: result.segmentId, elapsedSeconds }, + 'Onboarding complete.', + ) +} + +export async function markProjectOnboardingFailed( + projectId: string, + reason: string, +): Promise { + const qx = pgpQx(svc.postgres.writer.connection()) + + const updatedRows = await markProjectCatalogOnboardingFailed(qx, projectId, reason) + + if (updatedRows === 0) { + log.info( + { id: projectId }, + 'Project was already onboarded or no longer pending, not marking as error.', + ) + return + } + + log.error({ id: projectId, reason }, 'Onboarding permanently failed, marked as error.') +} diff --git a/services/apps/automatic_onboarding_worker/src/main.ts b/services/apps/automatic_onboarding_worker/src/main.ts index 21df335993..0bf7fa94d8 100644 --- a/services/apps/automatic_onboarding_worker/src/main.ts +++ b/services/apps/automatic_onboarding_worker/src/main.ts @@ -1,6 +1,8 @@ import { Config } from '@crowd/archetype-standard' import { Options, ServiceWorker } from '@crowd/archetype-worker' +import { scheduleProjectsOnboarding } from './schedules/scheduleProjectsOnboarding' + const config: Config = { envvars: [ 'CROWD_API_SERVICE_URL', @@ -34,5 +36,9 @@ setImmediate(async () => { svc.log.info('Automatic onboarding worker starting up.') + await scheduleProjectsOnboarding() + + svc.log.info('Automatic onboarding worker running — schedule registered, waiting for Temporal.') + await svc.start() }) diff --git a/services/apps/automatic_onboarding_worker/src/schedules/scheduleProjectsOnboarding.ts b/services/apps/automatic_onboarding_worker/src/schedules/scheduleProjectsOnboarding.ts new file mode 100644 index 0000000000..48adbdaba9 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/schedules/scheduleProjectsOnboarding.ts @@ -0,0 +1,45 @@ +import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client' + +import { svc } from '../main' +import { IOnboardProjectsInput, onboardProjects } from '../workflows' + +const ONBOARDING_ARGS: IOnboardProjectsInput = { + batchSize: 20, +} + +export const scheduleProjectsOnboarding = async () => { + svc.log.info('Scheduling projects onboarding') + + try { + await svc.temporal.schedule.create({ + scheduleId: 'automaticOnboarding', + spec: { + // Daily: catches up on whatever landed in 'onboard' state, independent of the evaluation schedule's timing. + cronExpressions: ['0 8 * * *'], + }, + policies: { + overlap: ScheduleOverlapPolicy.SKIP, + catchupWindow: '1 hour', + }, + action: { + type: 'startWorkflow', + workflowType: onboardProjects, + taskQueue: 'automatic-onboarding', + args: [ONBOARDING_ARGS], + workflowExecutionTimeout: '6 hours', + retry: { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, + }, + }, + }) + } catch (err) { + if (err instanceof ScheduleAlreadyRunning) { + svc.log.info('Schedule already registered in Temporal.') + svc.log.info('Configuration may have changed since. Please make sure they are in sync.') + } else { + throw new Error(err) + } + } +} diff --git a/services/apps/automatic_onboarding_worker/src/types.ts b/services/apps/automatic_onboarding_worker/src/types.ts new file mode 100644 index 0000000000..cb0f970d90 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/types.ts @@ -0,0 +1,3 @@ +export interface IOnboardProjectsInput { + batchSize?: number +} diff --git a/services/apps/automatic_onboarding_worker/src/workflows.ts b/services/apps/automatic_onboarding_worker/src/workflows.ts new file mode 100644 index 0000000000..98607ffff8 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/workflows.ts @@ -0,0 +1,5 @@ +import type { IOnboardProjectsInput } from './types' +import { onboardProjects } from './workflows/onboardProjects' + +export { onboardProjects } +export type { IOnboardProjectsInput } diff --git a/services/apps/automatic_onboarding_worker/src/workflows/index.ts b/services/apps/automatic_onboarding_worker/src/workflows/index.ts deleted file mode 100644 index 336ce12bb9..0000000000 --- a/services/apps/automatic_onboarding_worker/src/workflows/index.ts +++ /dev/null @@ -1 +0,0 @@ -export {} diff --git a/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts b/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts new file mode 100644 index 0000000000..48fa3587e5 --- /dev/null +++ b/services/apps/automatic_onboarding_worker/src/workflows/onboardProjects.ts @@ -0,0 +1,69 @@ +import { log, proxyActivities, rootCause } from '@temporalio/workflow' + +import type * as activities from '../activities' +import type { IOnboardProjectsInput } from '../types' + +// Short timeout: just a DB read. +const fetchActivities = proxyActivities({ + startToCloseTimeout: '2 minutes', + retry: { maximumAttempts: 3 }, +}) + +// Each onboarding call chains a segment create/query plus GitHub enrichment and integration calls, +// each of which can individually approach a ~30s backend timeout; give generous headroom per project. +const onboardActivities = proxyActivities({ + startToCloseTimeout: '5 minutes', + retry: { maximumAttempts: 2 }, +}) + +const failureActivities = proxyActivities({ + startToCloseTimeout: '2 minutes', + retry: { maximumAttempts: 2 }, +}) + +export async function onboardProjects(input: IOnboardProjectsInput = {}): Promise { + const { batchSize = 20 } = input + + log.info('onboardProjects workflow started.') + + const projects = await fetchActivities.fetchProjectsPendingOnboarding(batchSize) + + if (projects.length === 0) { + log.info('No projects pending onboarding. Nothing to do.') + return + } + + log.info(`Onboarding ${projects.length} project(s) (batch size: ${batchSize}).`) + + let succeeded = 0 + let failed = 0 + + for (let i = 0; i < projects.length; i++) { + const project = projects[i] + log.info(`[${i + 1}/${projects.length}] Onboarding: ${project.repoUrl}`) + + try { + await onboardActivities.onboardAndUpdateProject(project) + 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. + failed++ + const reason = rootCause(err) ?? String(err) + log.error( + `Onboarding failed for project id=${project.id} repoUrl=${project.repoUrl}: ${reason}`, + ) + + try { + await failureActivities.markProjectOnboardingFailed(project.id, reason) + } catch (markErr) { + // Don't let a failure to record the error state abort the rest of the batch. + log.error(`Failed to mark project id=${project.id} as errored: ${String(markErr)}`) + } + } + } + + log.info( + `Batch onboarding complete. total=${projects.length} succeeded=${succeeded} failed=${failed}`, + ) +} 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 6e2ad4c6d9..922918049c 100644 --- a/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts +++ b/services/libs/data-access-layer/src/project-catalog/projectCatalog.ts @@ -20,6 +20,7 @@ const PROJECT_CATALOG_COLUMNS = [ 'evaluationReason', 'evaluatedAt', 'onboardedAt', + 'onboardingError', 'syncedAt', 'createdAt', 'updatedAt', @@ -335,7 +336,7 @@ export async function upsertProjectCatalog( "repoName" = EXCLUDED."repoName", "source" = COALESCE(EXCLUDED."source", "projectCatalog"."source"), "action" = CASE - WHEN "projectCatalog"."action" IN ('onboard', 'skip', 'unsure') THEN "projectCatalog"."action" + WHEN "projectCatalog"."action" IN ('onboard', 'skip', 'unsure', 'error') THEN "projectCatalog"."action" WHEN EXCLUDED.action = 'evaluate' THEN 'evaluate' ELSE "projectCatalog"."action" END, @@ -408,7 +409,7 @@ export async function bulkUpsertProjectCatalog( "repoName" = EXCLUDED."repoName", "source" = COALESCE(EXCLUDED."source", "projectCatalog"."source"), "action" = CASE - WHEN "projectCatalog"."action" IN ('onboard', 'skip', 'unsure') THEN "projectCatalog"."action" + WHEN "projectCatalog"."action" IN ('onboard', 'skip', 'unsure', 'error') THEN "projectCatalog"."action" WHEN EXCLUDED.action = 'evaluate' THEN 'evaluate' ELSE "projectCatalog"."action" END, @@ -472,6 +473,10 @@ export async function updateProjectCatalog( setClauses.push('"onboardedAt" = $(onboardedAt)') params.onboardedAt = data.onboardedAt } + if (data.onboardingError !== undefined) { + setClauses.push('"onboardingError" = $(onboardingError)') + params.onboardingError = data.onboardingError + } if (setClauses.length === 0) { return findProjectCatalogById(qx, id) @@ -490,6 +495,21 @@ export async function updateProjectCatalog( ) } +export async function markProjectCatalogOnboardingFailed( + qx: QueryExecutor, + id: string, + reason: string, +): Promise { + return qx.result( + ` + UPDATE "projectCatalog" + SET "action" = 'error', "onboardingError" = $(reason), "updatedAt" = NOW() + WHERE id = $(id) AND "action" = 'onboard' AND "onboardedAt" IS NULL + `, + { id, reason }, + ) +} + export async function updateProjectCatalogSyncedAt(qx: QueryExecutor, id: string): Promise { await qx.selectNone( ` 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 69d93254be..e54b214b04 100644 --- a/services/libs/data-access-layer/src/project-catalog/types.ts +++ b/services/libs/data-access-layer/src/project-catalog/types.ts @@ -1,4 +1,4 @@ -export type ProjectCatalogAction = 'auto' | 'evaluate' | 'onboard' | 'skip' | 'unsure' +export type ProjectCatalogAction = 'auto' | 'evaluate' | 'onboard' | 'skip' | 'unsure' | 'error' export interface IDbProjectCatalog { id: string @@ -12,6 +12,7 @@ export interface IDbProjectCatalog { evaluationReason: string | null evaluatedAt: string | null onboardedAt: string | null + onboardingError: string | null syncedAt: string | null createdAt: string | null updatedAt: string | null @@ -27,17 +28,24 @@ type ProjectCatalogWritable = Pick< | 'lfCriticalityScore' | 'evaluationResult' | 'evaluationReason' + | 'onboardingError' > export type IDbProjectCatalogCreate = Omit< ProjectCatalogWritable, - 'source' | 'action' | 'lfCriticalityScore' | 'evaluationResult' | 'evaluationReason' + | 'source' + | 'action' + | 'lfCriticalityScore' + | 'evaluationResult' + | 'evaluationReason' + | 'onboardingError' > & { source?: string | null action?: ProjectCatalogAction lfCriticalityScore?: number evaluationResult?: string | null evaluationReason?: string | null + onboardingError?: string | null } export type IDbProjectCatalogUpdate = Partial & {