Skip to content

Commit 2bb3a23

Browse files
committed
feat: adding onboarding worker as standard pattern, and onboarding error reason on db
Signed-off-by: Umberto Sgueglia <usgueglia@contractor.linuxfoundation.org>
1 parent 144926e commit 2bb3a23

11 files changed

Lines changed: 220 additions & 4 deletions

File tree

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
ALTER TABLE "projectCatalog"
2+
ADD COLUMN IF NOT EXISTS "onboardingError" TEXT;
Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
export {}
1+
export * from './activities/activities'
Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
import {
2+
findProjectCatalogById,
3+
findProjectCatalogPendingOnboarding,
4+
updateProjectCatalog,
5+
} from '@crowd/data-access-layer'
6+
import { IDbProjectCatalog } from '@crowd/data-access-layer/src/project-catalog/types'
7+
import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor'
8+
import { getServiceLogger } from '@crowd/logging'
9+
10+
import { svc } from '../main'
11+
import { onboardProject } from '../onboarder/onboarder'
12+
13+
const log = getServiceLogger()
14+
15+
export async function fetchProjectsPendingOnboarding(
16+
batchSize: number,
17+
): Promise<IDbProjectCatalog[]> {
18+
const qx = pgpQx(svc.postgres.reader.connection())
19+
20+
const projects = await findProjectCatalogPendingOnboarding(qx, { limit: batchSize })
21+
22+
log.info({ count: projects.length, batchSize }, 'Fetched projects pending onboarding.')
23+
24+
return projects
25+
}
26+
27+
export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promise<void> {
28+
const qx = pgpQx(svc.postgres.writer.connection())
29+
const startTime = Date.now()
30+
31+
// Guard: fetch fresh state to ensure the API is called at most once per project.
32+
// Uses the writer connection to avoid replica lag missing a just-written onboardedAt.
33+
const fresh = await findProjectCatalogById(qx, project.id)
34+
if (fresh?.onboardedAt) {
35+
log.info(
36+
{ id: project.id, repoUrl: project.repoUrl, onboardedAt: fresh.onboardedAt },
37+
'Project already onboarded, skipping API call.',
38+
)
39+
return
40+
}
41+
42+
log.info({ id: project.id, repoUrl: project.repoUrl }, 'Starting onboarding.')
43+
44+
const result = await onboardProject({
45+
id: project.id,
46+
repoUrl: project.repoUrl,
47+
repoName: project.repoName,
48+
projectSlug: project.projectSlug,
49+
})
50+
51+
if (result.outcome === 'error') {
52+
throw new Error(result.error ?? 'Unknown onboarding error')
53+
}
54+
55+
await updateProjectCatalog(qx, project.id, {
56+
onboardedAt: new Date().toISOString(),
57+
})
58+
59+
const elapsedSeconds = ((Date.now() - startTime) / 1000).toFixed(1)
60+
61+
log.info(
62+
{ id: project.id, repoUrl: project.repoUrl, segmentId: result.segmentId, elapsedSeconds },
63+
'Onboarding complete.',
64+
)
65+
}
66+
67+
export async function markProjectOnboardingFailed(
68+
projectId: string,
69+
reason: string,
70+
): Promise<void> {
71+
const qx = pgpQx(svc.postgres.writer.connection())
72+
73+
await updateProjectCatalog(qx, projectId, {
74+
action: 'error',
75+
onboardingError: reason,
76+
})
77+
78+
log.error({ id: projectId, reason }, 'Onboarding permanently failed, marked as error.')
79+
}

services/apps/automatic_onboarding_worker/src/main.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
import { Config } from '@crowd/archetype-standard'
22
import { Options, ServiceWorker } from '@crowd/archetype-worker'
33

4+
import { scheduleProjectsOnboarding } from './schedules/scheduleProjectsOnboarding'
5+
46
const config: Config = {
57
envvars: [
68
'CROWD_API_SERVICE_URL',
@@ -34,5 +36,9 @@ setImmediate(async () => {
3436

3537
svc.log.info('Automatic onboarding worker starting up.')
3638

39+
await scheduleProjectsOnboarding()
40+
41+
svc.log.info('Automatic onboarding worker running — schedule registered, waiting for Temporal.')
42+
3743
await svc.start()
3844
})
Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,46 @@
1+
import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client'
2+
3+
import { svc } from '../main'
4+
import { IOnboardProjectsInput, onboardProjects } from '../workflows'
5+
6+
const ONBOARDING_ARGS: IOnboardProjectsInput = {
7+
batchSize: 50,
8+
}
9+
10+
export const scheduleProjectsOnboarding = async () => {
11+
svc.log.info('Scheduling projects onboarding')
12+
13+
try {
14+
await svc.temporal.schedule.create({
15+
scheduleId: 'projectsOnboarding',
16+
spec: {
17+
// Daily: catches up on whatever landed in 'onboard' state, independent of the evaluation schedule's timing.
18+
cronExpressions: ['0 8 * * *'],
19+
},
20+
policies: {
21+
overlap: ScheduleOverlapPolicy.SKIP,
22+
catchupWindow: '1 minute',
23+
},
24+
action: {
25+
type: 'startWorkflow',
26+
workflowType: onboardProjects,
27+
taskQueue: 'automatic-onboarding',
28+
args: [ONBOARDING_ARGS],
29+
// 50 projects × up to ~2min per attempt, up to 2 attempts each = ~3.3h worst case; set ceiling with margin.
30+
workflowExecutionTimeout: '6 hours',
31+
retry: {
32+
initialInterval: '30 seconds',
33+
backoffCoefficient: 2,
34+
maximumAttempts: 3,
35+
},
36+
},
37+
})
38+
} catch (err) {
39+
if (err instanceof ScheduleAlreadyRunning) {
40+
svc.log.info('Schedule already registered in Temporal.')
41+
svc.log.info('Configuration may have changed since. Please make sure they are in sync.')
42+
} else {
43+
throw new Error(err)
44+
}
45+
}
46+
}
Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
export interface IOnboardProjectsInput {
2+
batchSize?: number
3+
}
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
import type { IOnboardProjectsInput } from './types'
2+
import { onboardProjects } from './workflows/onboardProjects'
3+
4+
export { onboardProjects }
5+
export type { IOnboardProjectsInput }

services/apps/automatic_onboarding_worker/src/workflows/index.ts

Lines changed: 0 additions & 1 deletion
This file was deleted.
Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,63 @@
1+
import { log, proxyActivities, rootCause } from '@temporalio/workflow'
2+
3+
import type * as activities from '../activities'
4+
import type { IOnboardProjectsInput } from '../types'
5+
6+
// Short timeout: just a DB read.
7+
const fetchActivities = proxyActivities<typeof activities>({
8+
startToCloseTimeout: '2 minutes',
9+
retry: { maximumAttempts: 3 },
10+
})
11+
12+
// Each onboarding call chains a segment create/query plus GitHub enrichment and integration calls; give generous headroom per project.
13+
const onboardActivities = proxyActivities<typeof activities>({
14+
startToCloseTimeout: '2 minutes',
15+
retry: { maximumAttempts: 2 },
16+
})
17+
18+
export async function onboardProjects(input: IOnboardProjectsInput = {}): Promise<void> {
19+
const { batchSize = 50 } = input
20+
21+
log.info('onboardProjects workflow started.')
22+
23+
const projects = await fetchActivities.fetchProjectsPendingOnboarding(batchSize)
24+
25+
if (projects.length === 0) {
26+
log.info('No projects pending onboarding. Nothing to do.')
27+
return
28+
}
29+
30+
log.info(`Onboarding ${projects.length} project(s) (batch size: ${batchSize}).`)
31+
32+
let succeeded = 0
33+
let failed = 0
34+
35+
for (let i = 0; i < projects.length; i++) {
36+
const project = projects[i]
37+
log.info(`[${i + 1}/${projects.length}] Onboarding: ${project.repoUrl}`)
38+
39+
try {
40+
await onboardActivities.onboardAndUpdateProject(project)
41+
succeeded++
42+
} catch (err) {
43+
// Activity-level retries are already exhausted at this point — mark as a
44+
// terminal error so the daily schedule stops retrying this project forever.
45+
failed++
46+
const reason = rootCause(err) ?? String(err)
47+
log.error(
48+
`Onboarding failed for project id=${project.id} repoUrl=${project.repoUrl}: ${reason}`,
49+
)
50+
51+
try {
52+
await onboardActivities.markProjectOnboardingFailed(project.id, reason)
53+
} catch (markErr) {
54+
// Don't let a failure to record the error state abort the rest of the batch.
55+
log.error(`Failed to mark project id=${project.id} as errored: ${String(markErr)}`)
56+
}
57+
}
58+
}
59+
60+
log.info(
61+
`Batch onboarding complete. total=${projects.length} succeeded=${succeeded} failed=${failed}`,
62+
)
63+
}

services/libs/data-access-layer/src/project-catalog/projectCatalog.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ const PROJECT_CATALOG_COLUMNS = [
2020
'evaluationReason',
2121
'evaluatedAt',
2222
'onboardedAt',
23+
'onboardingError',
2324
'syncedAt',
2425
'createdAt',
2526
'updatedAt',
@@ -472,6 +473,10 @@ export async function updateProjectCatalog(
472473
setClauses.push('"onboardedAt" = $(onboardedAt)')
473474
params.onboardedAt = data.onboardedAt
474475
}
476+
if (data.onboardingError !== undefined) {
477+
setClauses.push('"onboardingError" = $(onboardingError)')
478+
params.onboardingError = data.onboardingError
479+
}
475480

476481
if (setClauses.length === 0) {
477482
return findProjectCatalogById(qx, id)

0 commit comments

Comments
 (0)