diff --git a/scripts/builders/packages.env b/scripts/builders/packages.env index b8d1a41a68..03cc6369b7 100644 --- a/scripts/builders/packages.env +++ b/scripts/builders/packages.env @@ -1,4 +1,4 @@ DOCKERFILE="./services/docker/Dockerfile.packages" CONTEXT="../" REPO="sjc.ocir.io/axbydjxa5zuh/packages" -SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker" +SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker go-worker" diff --git a/scripts/services/go-worker.yaml b/scripts/services/go-worker.yaml new file mode 100644 index 0000000000..a0c2a56c06 --- /dev/null +++ b/scripts/services/go-worker.yaml @@ -0,0 +1,66 @@ +version: '3.1' + +x-env-args: &env-args + DOCKER_BUILDKIT: 1 + NODE_ENV: docker + SERVICE: go-worker + CROWD_TEMPORAL_TASKQUEUE: go-worker + SHELL: /bin/sh + SUPPRESS_NO_CONFIG_WARNING: 'true' + +services: + go-worker: + build: + context: ../../ + dockerfile: ./scripts/services/docker/Dockerfile.packages + command: 'pnpm run start:go-worker' + working_dir: /usr/crowd/app/services/apps/packages_worker + env_file: + - ../../backend/.env.dist.local + - ../../backend/.env.dist.composed + - ../../backend/.env.override.local + - ../../backend/.env.override.composed + environment: + <<: *env-args + restart: always + networks: + - crowd-bridge + + go-worker-dev: + build: + context: ../../ + dockerfile: ./scripts/services/docker/Dockerfile.packages + command: 'pnpm run dev:go-worker' + working_dir: /usr/crowd/app/services/apps/packages_worker + # user: '${USER_ID}:${GROUP_ID}' + env_file: + - ../../backend/.env.dist.local + - ../../backend/.env.dist.composed + - ../../backend/.env.override.local + - ../../backend/.env.override.composed + environment: + <<: *env-args + hostname: go-worker + networks: + - crowd-bridge + volumes: + - ../../services/libs/audit-logs/src:/usr/crowd/app/services/libs/audit-logs/src + - ../../services/libs/common/src:/usr/crowd/app/services/libs/common/src + - ../../services/libs/common_services/src:/usr/crowd/app/services/libs/common_services/src + - ../../services/libs/data-access-layer/src:/usr/crowd/app/services/libs/data-access-layer/src + - ../../services/libs/database/src:/usr/crowd/app/services/libs/database/src + - ../../services/libs/integrations/src:/usr/crowd/app/services/libs/integrations/src + - ../../services/libs/logging/src:/usr/crowd/app/services/libs/logging/src + - ../../services/libs/nango/src:/usr/crowd/app/services/libs/nango/src + - ../../services/libs/opensearch/src:/usr/crowd/app/services/libs/opensearch/src + - ../../services/libs/queue/src:/usr/crowd/app/services/libs/queue/src + - ../../services/libs/redis/src:/usr/crowd/app/services/libs/redis/src + - ../../services/libs/snowflake/src:/usr/crowd/app/services/libs/snowflake/src + - ../../services/libs/telemetry/src:/usr/crowd/app/services/libs/telemetry/src + - ../../services/libs/temporal/src:/usr/crowd/app/services/libs/temporal/src + - ../../services/libs/types/src:/usr/crowd/app/services/libs/types/src + - ../../services/apps/packages_worker/src:/usr/crowd/app/services/apps/packages_worker/src + +networks: + crowd-bridge: + external: true diff --git a/services/apps/packages_worker/package.json b/services/apps/packages_worker/package.json index f030b1af59..d49e5c2f34 100644 --- a/services/apps/packages_worker/package.json +++ b/services/apps/packages_worker/package.json @@ -21,6 +21,8 @@ "export-to-bucket:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=bq-dataset-ingest tsx src/scripts/exportToBucket.ts", "trigger-bootstrap": "SERVICE=bq-dataset-ingest tsx src/scripts/triggerBootstrap.ts", "trigger-bootstrap:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=bq-dataset-ingest tsx src/scripts/triggerBootstrap.ts", + "trigger-go": "SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts", + "trigger-go:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts", "start:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker SERVICE=npm-worker tsx src/bin/npm-worker.ts", "dev:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker SERVICE=npm-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/npm-worker.ts", "dev:npm-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=npm-worker SERVICE=npm-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/npm-worker.ts", @@ -33,6 +35,9 @@ "start:cargo-worker": "CROWD_TEMPORAL_TASKQUEUE=cargo-worker SERVICE=cargo-worker tsx src/bin/cargo-worker.ts", "dev:cargo-worker": "CROWD_TEMPORAL_TASKQUEUE=cargo-worker SERVICE=cargo-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9240 src/bin/cargo-worker.ts", "dev:cargo-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=cargo-worker SERVICE=cargo-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9240 src/bin/cargo-worker.ts", + "start:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker tsx src/bin/go-worker.ts", + "dev:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9241 src/bin/go-worker.ts", + "dev:go-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9241 src/bin/go-worker.ts", "backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts", "backfill:maven:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=maven LOG_LEVEL=info tsx src/bin/maven-backfill.ts", "backfill:stewardship": "SERVICE=stewardship-backfill tsx src/bin/stewardship-backfill.ts", diff --git a/services/apps/packages_worker/src/activities.ts b/services/apps/packages_worker/src/activities.ts index 09a591c869..acd9079f24 100644 --- a/services/apps/packages_worker/src/activities.ts +++ b/services/apps/packages_worker/src/activities.ts @@ -24,3 +24,4 @@ export { cargoFlushAudit, cargoCleanup, } from './cargo/activities' +export { enrichGoVersionsBatch, enrichGoStatusBatch } from './go/activities' diff --git a/services/apps/packages_worker/src/bin/go-worker.ts b/services/apps/packages_worker/src/bin/go-worker.ts new file mode 100644 index 0000000000..15af3b6819 --- /dev/null +++ b/services/apps/packages_worker/src/bin/go-worker.ts @@ -0,0 +1,9 @@ +import { scheduleGoStatus, scheduleGoVersions } from '../go/schedule' +import { svc } from '../service' + +setImmediate(async () => { + await svc.init() + await scheduleGoVersions() + await scheduleGoStatus() + await svc.start() +}) diff --git a/services/apps/packages_worker/src/config.ts b/services/apps/packages_worker/src/config.ts index b524205ba5..3fb773261c 100644 --- a/services/apps/packages_worker/src/config.ts +++ b/services/apps/packages_worker/src/config.ts @@ -63,6 +63,13 @@ export function getCargoConfig() { } } +export function getGoConfig() { + return { + fetchTimeoutMs: parseInt(process.env.GO_FETCH_TIMEOUT_MS ?? '15000', 10), + proxyConcurrency: Math.max(1, parseInt(process.env.GO_PROXY_CONCURRENCY ?? '10', 10)), + } +} + export function getDockerhubConfig() { return { hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'), diff --git a/services/apps/packages_worker/src/go/activities.ts b/services/apps/packages_worker/src/go/activities.ts new file mode 100644 index 0000000000..0ff111bd7e --- /dev/null +++ b/services/apps/packages_worker/src/go/activities.ts @@ -0,0 +1,141 @@ +import { Context } from '@temporalio/activity' + +import { logAuditFieldChanges } from '@crowd/data-access-layer/src/packages' +import type { QueryExecutor } from '@crowd/data-access-layer/src/queryExecutor' +import { getServiceChildLogger } from '@crowd/logging' + +import { getGoConfig } from '../config' +import { getPackagesDb } from '../db' + +import { fetchStatus } from './pkgGoDevClient' +import { fetchLatest } from './proxyClient' +import { isFetchError } from './types' + +const log = getServiceChildLogger('go') + +const PROXY_SOURCE = 'go-proxy' +const PKGGODEV_SOURCE = 'pkg-go-dev' + +// TODO: filter to critical packages once computed +async function getGoBatch( + qx: QueryExecutor, + afterPurl: string, + batchSize: number, +): Promise> { + return qx.select( + `SELECT purl, name FROM packages + WHERE ecosystem = 'go' AND purl > $(after) + ORDER BY last_synced_at ASC NULLS FIRST, purl ASC + LIMIT $(limit)`, + { after: afterPurl, limit: batchSize }, + ) +} + +export async function enrichGoVersionsBatch( + afterPurl: string, + batchSize: number, +): Promise { + const qx = await getPackagesDb() + const rows = await getGoBatch(qx, afterPurl, batchSize) + if (rows.length === 0) return null + + const { fetchTimeoutMs, proxyConcurrency } = getGoConfig() + + const enrichOne = async (row: { purl: string; name: string }): Promise => { + Context.current().heartbeat(row.purl) + const result = await fetchLatest(row.name, fetchTimeoutMs) + if (isFetchError(result)) { + log.warn( + { purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode }, + 'go proxy fetch failed — skipping package', + ) + return + } + const changed = await qx.selectOne( + `WITH old AS ( + SELECT latest_version AS v, latest_release_at AS t, repository_url AS r + FROM packages WHERE purl = $(purl) + ), + upd AS ( + UPDATE packages p SET + latest_version = $(version), + latest_release_at = $(releaseAt), + repository_url = COALESCE($(repoUrl), p.repository_url), + last_synced_at = NOW() + WHERE p.purl = $(purl) + RETURNING latest_version AS v, latest_release_at AS t, repository_url AS r + ) + SELECT + (SELECT v FROM old) IS DISTINCT FROM (SELECT v FROM upd) AS v_changed, + (SELECT t FROM old) IS DISTINCT FROM (SELECT t FROM upd) AS t_changed, + (SELECT r FROM old) IS DISTINCT FROM (SELECT r FROM upd) AS r_changed`, + { + version: result.version, + releaseAt: result.releaseAt, + repoUrl: result.repoUrl, + purl: row.purl, + }, + ) + const changedFields = [ + changed?.v_changed ? 'packages.latest_version' : null, + changed?.t_changed ? 'packages.latest_release_at' : null, + changed?.r_changed ? 'packages.repository_url' : null, + ].filter(Boolean) as string[] + await logAuditFieldChanges(qx, PROXY_SOURCE, row.purl, changedFields) + } + + for (let i = 0; i < rows.length; i += proxyConcurrency) { + await Promise.all(rows.slice(i, i + proxyConcurrency).map(enrichOne)) + } + + log.info({ count: rows.length, concurrency: proxyConcurrency }, 'Enriched go versions batch') + return rows[rows.length - 1].purl +} + +export async function enrichGoStatusBatch( + afterPurl: string, + batchSize: number, +): Promise { + const qx = await getPackagesDb() + const rows = await getGoBatch(qx, afterPurl, batchSize) + if (rows.length === 0) return null + + const { fetchTimeoutMs } = getGoConfig() + for (const row of rows) { + const result = await fetchStatus(row.name, fetchTimeoutMs, () => + Context.current().heartbeat(row.purl), + ) + if (isFetchError(result)) { + log.warn( + { purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode }, + 'pkg.go.dev fetch failed — skipping package', + ) + continue + } + const changed = await qx.selectOne( + `WITH old AS ( + SELECT status AS s, versions_count AS vc FROM packages WHERE purl = $(purl) + ), + upd AS ( + UPDATE packages p SET + status = $(status), + versions_count = COALESCE($(versionsCount), p.versions_count), + last_synced_at = NOW() + WHERE p.purl = $(purl) + RETURNING status AS s, versions_count AS vc + ) + SELECT + (SELECT s FROM old) IS DISTINCT FROM (SELECT s FROM upd) AS s_changed, + (SELECT vc FROM old) IS DISTINCT FROM (SELECT vc FROM upd) AS vc_changed`, + { status: result.status, versionsCount: result.versionsCount, purl: row.purl }, + ) + const changedFields = [ + changed?.s_changed ? 'packages.status' : null, + changed?.vc_changed ? 'packages.versions_count' : null, + ].filter(Boolean) as string[] + await logAuditFieldChanges(qx, PKGGODEV_SOURCE, row.purl, changedFields) + } + + log.info({ count: rows.length }, 'Enriched go status batch') + return rows[rows.length - 1].purl +} diff --git a/services/apps/packages_worker/src/go/pkgGoDevClient.ts b/services/apps/packages_worker/src/go/pkgGoDevClient.ts new file mode 100644 index 0000000000..c9581cfac4 --- /dev/null +++ b/services/apps/packages_worker/src/go/pkgGoDevClient.ts @@ -0,0 +1,105 @@ +import { FetchError, GoStatusResult, isFetchError } from './types' + +const BASE = process.env.PKGGODEV_BASE_URL ?? 'https://pkg.go.dev' +// 40 QPS per IP. Keep a single in-process gap under that ceiling. +const MIN_INTERVAL_MS = parseInt(process.env.PKGGODEV_MIN_INTERVAL_MS ?? '30', 10) +const MAX_429_RETRIES = 5 +const MAX_PAGES = 20 + +interface VersionItem { + version: string + deprecated?: boolean + retracted?: boolean + latestVersion?: string +} +interface VersionsPage { + items?: VersionItem[] + total?: number + nextPageToken?: string +} + +let lastRequestAt = 0 + +function sleep(ms: number): Promise { + return new Promise((r) => setTimeout(r, ms)) +} + +async function throttle(): Promise { + const wait = lastRequestAt + MIN_INTERVAL_MS - Date.now() + if (wait > 0) await sleep(wait) + lastRequestAt = Date.now() +} + +async function getPage(url: string, timeoutMs: number): Promise { + for (let attempt = 0; attempt <= MAX_429_RETRIES; attempt++) { + await throttle() + const controller = new AbortController() + const timer = setTimeout(() => controller.abort(), timeoutMs) + let res: Response + try { + res = await fetch(url, { signal: controller.signal }) + } catch (e) { + clearTimeout(timer) + return { kind: 'TRANSIENT', message: `network error: ${(e as Error).message}` } + } + clearTimeout(timer) + + if (res.status === 429) { + const reset = parseInt(res.headers.get('x-ratelimit-reset') ?? '0', 10) + const waitMs = reset ? Math.max(1000, reset * 1000 - Date.now() + 500) : 2000 + await sleep(waitMs) + continue + } + // Any other 4xx is permanent (e.g. 400 for submodule/non-module-root paths) — skip, don't retry. + if (res.status >= 400 && res.status < 500) { + return { kind: 'NOT_FOUND', statusCode: res.status, message: `${res.status}` } + } + if (res.status !== 200) { + return { + kind: 'TRANSIENT', + statusCode: res.status, + message: `unexpected status ${res.status}`, + } + } + try { + return (await res.json()) as VersionsPage + } catch { + return { kind: 'MALFORMED', message: 'invalid json' } + } + } + return { kind: 'RATE_LIMIT', statusCode: 429, message: '429 after retries' } +} + +// Module status from /v1beta/versions: 'deprecated' if the latest version is +// deprecated/retracted, else 'active'. Pages are newest-first and every item carries +// latestVersion, so the match is normally on page 1; paginate (token query param, +// request otherwise verbatim) only if it isn't. +export async function fetchStatus( + module: string, + timeoutMs: number, + onHeartbeat?: () => void, +): Promise { + let token: string | undefined + let versionsCount: number | null = null + for (let page = 0; page < MAX_PAGES; page++) { + const url = `${BASE}/v1beta/versions/${module}${token ? `?token=${encodeURIComponent(token)}` : ''}` + const result = await getPage(url, timeoutMs) + onHeartbeat?.() + if (isFetchError(result)) return result + + if (versionsCount === null && typeof result.total === 'number') versionsCount = result.total + const items = result.items ?? [] + const latestVersion = items.find((i) => i.latestVersion)?.latestVersion + const match = latestVersion ? items.find((i) => i.version === latestVersion) : undefined + if (match) { + return { + status: match.deprecated || match.retracted ? 'deprecated' : 'active', + versionsCount, + } + } + + if (!result.nextPageToken) break + token = result.nextPageToken + } + return { kind: 'NOT_FOUND', message: 'latest version entry not found' } +} diff --git a/services/apps/packages_worker/src/go/proxyClient.ts b/services/apps/packages_worker/src/go/proxyClient.ts new file mode 100644 index 0000000000..4df6bdce34 --- /dev/null +++ b/services/apps/packages_worker/src/go/proxyClient.ts @@ -0,0 +1,52 @@ +import { FetchError, GoProxyLatest } from './types' + +const BASE = process.env.GO_PROXY_BASE_URL ?? 'https://proxy.golang.org' +const ZERO_TIME = '0001-01-01T00:00:00Z' + +// GOPROXY spec: uppercase letters in a module path are escaped as '!' + lowercase. +export function escapeModulePath(module: string): string { + return module.replace(/!/g, '!!').replace(/[A-Z]/g, (c) => '!' + c.toLowerCase()) +} + +export async function fetchLatest( + module: string, + timeoutMs: number, +): Promise { + const url = `${BASE}/${escapeModulePath(module)}/@latest` + const controller = new AbortController() + const timer = setTimeout(() => controller.abort(), timeoutMs) + + let res: Response + try { + res = await fetch(url, { signal: controller.signal }) + } catch (e) { + return { kind: 'TRANSIENT', message: `network error: ${(e as Error).message}` } + } finally { + clearTimeout(timer) + } + + if (res.status === 429) { + return { kind: 'RATE_LIMIT', statusCode: 429, message: 'rate limited' } + } + // Any other 4xx is permanent (unknown/invalid module path) — skip, don't retry. + if (res.status >= 400 && res.status < 500) { + return { kind: 'NOT_FOUND', statusCode: res.status, message: `${res.status}` } + } + if (res.status !== 200) { + return { kind: 'TRANSIENT', statusCode: res.status, message: `unexpected status ${res.status}` } + } + + let body: { Version?: string; Time?: string; Origin?: { URL?: string } } + try { + body = (await res.json()) as { Version?: string; Time?: string; Origin?: { URL?: string } } + } catch { + return { kind: 'MALFORMED', message: 'invalid json' } + } + if (!body.Version) return { kind: 'MALFORMED', message: 'missing Version' } + + return { + version: body.Version, + releaseAt: body.Time && body.Time !== ZERO_TIME ? body.Time : null, + repoUrl: body.Origin?.URL || null, + } +} diff --git a/services/apps/packages_worker/src/go/schedule.ts b/services/apps/packages_worker/src/go/schedule.ts new file mode 100644 index 0000000000..c0918c15e5 --- /dev/null +++ b/services/apps/packages_worker/src/go/schedule.ts @@ -0,0 +1,76 @@ +import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client' + +import { svc } from '../service' +import { enrichGoStatus, enrichGoVersions } from '../workflows' + +export async function scheduleGoVersions(): Promise { + const { temporal } = svc + if (!temporal) throw new Error('Temporal client not initialized') + + try { + await temporal.schedule.create({ + scheduleId: 'go-version-enrich', + spec: { + cronExpressions: ['0 2 * * *'], + }, + policies: { + overlap: ScheduleOverlapPolicy.SKIP, + catchupWindow: '1 hour', + }, + action: { + type: 'startWorkflow', + workflowType: enrichGoVersions, + taskQueue: 'go-worker', + workflowRunTimeout: '24 hours', + retry: { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, + }, + args: [], + }, + }) + } catch (err) { + if (err instanceof ScheduleAlreadyRunning) { + svc.log.info('Schedule go-version-enrich already registered.') + } else { + throw err + } + } +} + +export async function scheduleGoStatus(): Promise { + const { temporal } = svc + if (!temporal) throw new Error('Temporal client not initialized') + + try { + await temporal.schedule.create({ + scheduleId: 'go-status-enrich', + spec: { + cronExpressions: ['0 5 * * *'], + }, + policies: { + overlap: ScheduleOverlapPolicy.SKIP, + catchupWindow: '1 hour', + }, + action: { + type: 'startWorkflow', + workflowType: enrichGoStatus, + taskQueue: 'go-worker', + workflowRunTimeout: '24 hours', + retry: { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, + }, + args: [], + }, + }) + } catch (err) { + if (err instanceof ScheduleAlreadyRunning) { + svc.log.info('Schedule go-status-enrich already registered.') + } else { + throw err + } + } +} diff --git a/services/apps/packages_worker/src/go/types.ts b/services/apps/packages_worker/src/go/types.ts new file mode 100644 index 0000000000..0ce7534a72 --- /dev/null +++ b/services/apps/packages_worker/src/go/types.ts @@ -0,0 +1,22 @@ +export interface GoProxyLatest { + version: string + releaseAt: string | null + repoUrl: string | null +} + +export type GoStatus = 'active' | 'deprecated' + +export interface GoStatusResult { + status: GoStatus + versionsCount: number | null +} + +export interface FetchError { + kind: 'NOT_FOUND' | 'RATE_LIMIT' | 'TRANSIENT' | 'MALFORMED' + statusCode?: number + message: string +} + +export function isFetchError(v: unknown): v is FetchError { + return typeof v === 'object' && v !== null && 'kind' in v && 'message' in v +} diff --git a/services/apps/packages_worker/src/go/workflows.ts b/services/apps/packages_worker/src/go/workflows.ts new file mode 100644 index 0000000000..efa93af1c5 --- /dev/null +++ b/services/apps/packages_worker/src/go/workflows.ts @@ -0,0 +1,40 @@ +import { continueAsNew, proxyActivities } from '@temporalio/workflow' + +import type * as activities from './activities' + +const acts = proxyActivities({ + startToCloseTimeout: '15 minutes', + heartbeatTimeout: '2 minutes', + retry: { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 5, + }, +}) + +const BATCH = 100 +const ROUNDS_PER_RUN = 200 + +interface ScanState { + cursor: string +} + +export async function enrichGoVersions(state: ScanState = { cursor: '' }): Promise { + let { cursor } = state + for (let r = 0; r < ROUNDS_PER_RUN; r++) { + const next = await acts.enrichGoVersionsBatch(cursor, BATCH) + if (next === null) return + cursor = next + } + await continueAsNew({ cursor }) +} + +export async function enrichGoStatus(state: ScanState = { cursor: '' }): Promise { + let { cursor } = state + for (let r = 0; r < ROUNDS_PER_RUN; r++) { + const next = await acts.enrichGoStatusBatch(cursor, BATCH) + if (next === null) return + cursor = next + } + await continueAsNew({ cursor }) +} diff --git a/services/apps/packages_worker/src/scripts/triggerGoEnrich.ts b/services/apps/packages_worker/src/scripts/triggerGoEnrich.ts new file mode 100644 index 0000000000..1b2b2154d3 --- /dev/null +++ b/services/apps/packages_worker/src/scripts/triggerGoEnrich.ts @@ -0,0 +1,65 @@ +import { TEMPORAL_CONFIG, getTemporalClient } from '@crowd/temporal' + +import { enrichGoStatus, enrichGoVersions } from '../go/workflows' + +const HELP = ` +Usage: trigger-go [versions|status|both] + +Arguments: + versions Enrich latest_version + latest_release_at from proxy.golang.org + status Enrich status from pkg.go.dev + both Trigger both (default) + +Examples: + pnpm trigger-go:local + pnpm trigger-go:local versions + pnpm trigger-go:local status +` + +async function main(): Promise { + const args = process.argv.slice(2) + if (args.includes('--help') || args.includes('-h')) { + console.log(HELP) + process.exit(0) + } + + const target = (args[0] ?? 'both') as 'versions' | 'status' | 'both' + if (!['versions', 'status', 'both'].includes(target)) { + console.error(`Unknown target "${target}". Use "versions", "status", or "both".`) + process.exit(1) + } + + const cfg = TEMPORAL_CONFIG() + if (!cfg.serverUrl || !cfg.namespace) { + console.error('Missing CROWD_TEMPORAL_SERVER_URL or CROWD_TEMPORAL_NAMESPACE') + process.exit(1) + } + + const client = await getTemporalClient(cfg) + const now = Date.now() + + if (target === 'versions' || target === 'both') { + const handle = await client.workflow.start(enrichGoVersions, { + taskQueue: 'go-worker', + workflowId: `go-version-enrich-manual-${now}`, + args: [], + }) + console.log(`Started workflow ${handle.workflowId}`) + } + + if (target === 'status' || target === 'both') { + const handle = await client.workflow.start(enrichGoStatus, { + taskQueue: 'go-worker', + workflowId: `go-status-enrich-manual-${now}`, + args: [], + }) + console.log(`Started workflow ${handle.workflowId}`) + } +} + +main() + .then(() => process.exit(0)) + .catch((err) => { + console.error('Failed to trigger go enrich:', err) + process.exit(1) + }) diff --git a/services/apps/packages_worker/src/workflows/index.ts b/services/apps/packages_worker/src/workflows/index.ts index 667ffa9180..244bb26b6c 100644 --- a/services/apps/packages_worker/src/workflows/index.ts +++ b/services/apps/packages_worker/src/workflows/index.ts @@ -19,3 +19,4 @@ export { mavenCriticalWorkflow, mavenNonCriticalWorkflow } from '../maven/workfl export { ingestScorecard } from '../scorecard/workflows' export { rankPackagesWorkflow } from '../criticality/workflow' export { cargoSyncWorkflow } from '../cargo/workflows' +export { enrichGoVersions, enrichGoStatus } from '../go/workflows'