Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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
@@ -0,0 +1,2 @@
ALTER TABLE "projectCatalog"
ADD COLUMN IF NOT EXISTS "onboardingError" TEXT;
Original file line number Diff line number Diff line change
@@ -1 +1 @@
export {}
export * from './activities/activities'
Original file line number Diff line number Diff line change
@@ -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<IDbProjectCatalog[]> {
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<typeof pgpQx>,
projectId: string,
): Promise<IDbProjectCatalog | null> {
const fresh = await findProjectCatalogById(qx, projectId)
return fresh?.onboardedAt ? fresh : null
}

export async function onboardAndUpdateProject(project: IDbProjectCatalog): Promise<void> {
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)
Comment thread
ulemons marked this conversation as resolved.
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,
})
Comment thread
ulemons marked this conversation as resolved.
Comment thread
ulemons marked this conversation as resolved.

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<void> {
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.')
}
6 changes: 6 additions & 0 deletions services/apps/automatic_onboarding_worker/src/main.ts
Original file line number Diff line number Diff line change
@@ -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',
Expand Down Expand Up @@ -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()
})
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client'

import { svc } from '../main'
import { IOnboardProjectsInput, onboardProjects } from '../workflows'

const ONBOARDING_ARGS: IOnboardProjectsInput = {
batchSize: 50,
}

export const scheduleProjectsOnboarding = async () => {
svc.log.info('Scheduling projects onboarding')

try {
await svc.temporal.schedule.create({
scheduleId: 'projectsOnboarding',
spec: {
// Daily: catches up on whatever landed in 'onboard' state, independent of the evaluation schedule's timing.
cronExpressions: ['0 8 * * *'],
Comment thread
ulemons marked this conversation as resolved.
},
policies: {
overlap: ScheduleOverlapPolicy.SKIP,
catchupWindow: '1 hour',
},
action: {
type: 'startWorkflow',
workflowType: onboardProjects,
taskQueue: 'automatic-onboarding',
args: [ONBOARDING_ARGS],
workflowExecutionTimeout: '14 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)
}
}
}
3 changes: 3 additions & 0 deletions services/apps/automatic_onboarding_worker/src/types.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
export interface IOnboardProjectsInput {
batchSize?: number
}
5 changes: 5 additions & 0 deletions services/apps/automatic_onboarding_worker/src/workflows.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
import type { IOnboardProjectsInput } from './types'
import { onboardProjects } from './workflows/onboardProjects'

export { onboardProjects }
export type { IOnboardProjectsInput }

This file was deleted.

Original file line number Diff line number Diff line change
@@ -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<typeof activities>({
Comment thread
ulemons marked this conversation as resolved.
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<typeof activities>({
startToCloseTimeout: '5 minutes',
retry: { maximumAttempts: 2 },
})

const failureActivities = proxyActivities<typeof activities>({
startToCloseTimeout: '2 minutes',
retry: { maximumAttempts: 2 },
})
Comment thread
ulemons marked this conversation as resolved.

export async function onboardProjects(input: IOnboardProjectsInput = {}): Promise<void> {
const { batchSize = 50 } = 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}`,
)
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ const PROJECT_CATALOG_COLUMNS = [
'evaluationReason',
'evaluatedAt',
'onboardedAt',
'onboardingError',
'syncedAt',
'createdAt',
'updatedAt',
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand All @@ -490,6 +495,21 @@ export async function updateProjectCatalog(
)
}

export async function markProjectCatalogOnboardingFailed(
qx: QueryExecutor,
id: string,
reason: string,
): Promise<number> {
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<void> {
await qx.selectNone(
`
Expand Down
12 changes: 10 additions & 2 deletions services/libs/data-access-layer/src/project-catalog/types.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
export type ProjectCatalogAction = 'auto' | 'evaluate' | 'onboard' | 'skip' | 'unsure'
export type ProjectCatalogAction = 'auto' | 'evaluate' | 'onboard' | 'skip' | 'unsure' | 'error'
Comment thread
ulemons marked this conversation as resolved.

export interface IDbProjectCatalog {
id: string
Expand All @@ -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
Expand All @@ -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
Comment thread
ulemons marked this conversation as resolved.
}

export type IDbProjectCatalogUpdate = Partial<ProjectCatalogWritable> & {
Expand Down
Loading