Skip to content

Commit 80d1eec

Browse files
authored
feat: go packages worker [CM-1275] (#4256)
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent 4ad17ce commit 80d1eec

14 files changed

Lines changed: 591 additions & 1 deletion

File tree

scripts/builders/packages.env

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
11
DOCKERFILE="./services/docker/Dockerfile.packages"
22
CONTEXT="../"
33
REPO="sjc.ocir.io/axbydjxa5zuh/packages"
4-
SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker"
4+
SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker go-worker"

scripts/services/go-worker.yaml

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,66 @@
1+
version: '3.1'
2+
3+
x-env-args: &env-args
4+
DOCKER_BUILDKIT: 1
5+
NODE_ENV: docker
6+
SERVICE: go-worker
7+
CROWD_TEMPORAL_TASKQUEUE: go-worker
8+
SHELL: /bin/sh
9+
SUPPRESS_NO_CONFIG_WARNING: 'true'
10+
11+
services:
12+
go-worker:
13+
build:
14+
context: ../../
15+
dockerfile: ./scripts/services/docker/Dockerfile.packages
16+
command: 'pnpm run start:go-worker'
17+
working_dir: /usr/crowd/app/services/apps/packages_worker
18+
env_file:
19+
- ../../backend/.env.dist.local
20+
- ../../backend/.env.dist.composed
21+
- ../../backend/.env.override.local
22+
- ../../backend/.env.override.composed
23+
environment:
24+
<<: *env-args
25+
restart: always
26+
networks:
27+
- crowd-bridge
28+
29+
go-worker-dev:
30+
build:
31+
context: ../../
32+
dockerfile: ./scripts/services/docker/Dockerfile.packages
33+
command: 'pnpm run dev:go-worker'
34+
working_dir: /usr/crowd/app/services/apps/packages_worker
35+
# user: '${USER_ID}:${GROUP_ID}'
36+
env_file:
37+
- ../../backend/.env.dist.local
38+
- ../../backend/.env.dist.composed
39+
- ../../backend/.env.override.local
40+
- ../../backend/.env.override.composed
41+
environment:
42+
<<: *env-args
43+
hostname: go-worker
44+
networks:
45+
- crowd-bridge
46+
volumes:
47+
- ../../services/libs/audit-logs/src:/usr/crowd/app/services/libs/audit-logs/src
48+
- ../../services/libs/common/src:/usr/crowd/app/services/libs/common/src
49+
- ../../services/libs/common_services/src:/usr/crowd/app/services/libs/common_services/src
50+
- ../../services/libs/data-access-layer/src:/usr/crowd/app/services/libs/data-access-layer/src
51+
- ../../services/libs/database/src:/usr/crowd/app/services/libs/database/src
52+
- ../../services/libs/integrations/src:/usr/crowd/app/services/libs/integrations/src
53+
- ../../services/libs/logging/src:/usr/crowd/app/services/libs/logging/src
54+
- ../../services/libs/nango/src:/usr/crowd/app/services/libs/nango/src
55+
- ../../services/libs/opensearch/src:/usr/crowd/app/services/libs/opensearch/src
56+
- ../../services/libs/queue/src:/usr/crowd/app/services/libs/queue/src
57+
- ../../services/libs/redis/src:/usr/crowd/app/services/libs/redis/src
58+
- ../../services/libs/snowflake/src:/usr/crowd/app/services/libs/snowflake/src
59+
- ../../services/libs/telemetry/src:/usr/crowd/app/services/libs/telemetry/src
60+
- ../../services/libs/temporal/src:/usr/crowd/app/services/libs/temporal/src
61+
- ../../services/libs/types/src:/usr/crowd/app/services/libs/types/src
62+
- ../../services/apps/packages_worker/src:/usr/crowd/app/services/apps/packages_worker/src
63+
64+
networks:
65+
crowd-bridge:
66+
external: true

services/apps/packages_worker/package.json

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
"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",
2222
"trigger-bootstrap": "SERVICE=bq-dataset-ingest tsx src/scripts/triggerBootstrap.ts",
2323
"trigger-bootstrap:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=bq-dataset-ingest tsx src/scripts/triggerBootstrap.ts",
24+
"trigger-go": "SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts",
25+
"trigger-go:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts",
2426
"start:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker SERVICE=npm-worker tsx src/bin/npm-worker.ts",
2527
"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",
2628
"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 @@
3335
"start:cargo-worker": "CROWD_TEMPORAL_TASKQUEUE=cargo-worker SERVICE=cargo-worker tsx src/bin/cargo-worker.ts",
3436
"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",
3537
"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",
38+
"start:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker tsx src/bin/go-worker.ts",
39+
"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",
40+
"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",
3641
"backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts",
3742
"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",
3843
"backfill:stewardship": "SERVICE=stewardship-backfill tsx src/bin/stewardship-backfill.ts",

services/apps/packages_worker/src/activities.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,3 +24,4 @@ export {
2424
cargoFlushAudit,
2525
cargoCleanup,
2626
} from './cargo/activities'
27+
export { enrichGoVersionsBatch, enrichGoStatusBatch } from './go/activities'
Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
import { scheduleGoStatus, scheduleGoVersions } from '../go/schedule'
2+
import { svc } from '../service'
3+
4+
setImmediate(async () => {
5+
await svc.init()
6+
await scheduleGoVersions()
7+
await scheduleGoStatus()
8+
await svc.start()
9+
})

services/apps/packages_worker/src/config.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,13 @@ export function getCargoConfig() {
6363
}
6464
}
6565

66+
export function getGoConfig() {
67+
return {
68+
fetchTimeoutMs: parseInt(process.env.GO_FETCH_TIMEOUT_MS ?? '15000', 10),
69+
proxyConcurrency: Math.max(1, parseInt(process.env.GO_PROXY_CONCURRENCY ?? '10', 10)),
70+
}
71+
}
72+
6673
export function getDockerhubConfig() {
6774
return {
6875
hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'),
Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
import { Context } from '@temporalio/activity'
2+
3+
import { logAuditFieldChanges } from '@crowd/data-access-layer/src/packages'
4+
import type { QueryExecutor } from '@crowd/data-access-layer/src/queryExecutor'
5+
import { getServiceChildLogger } from '@crowd/logging'
6+
7+
import { getGoConfig } from '../config'
8+
import { getPackagesDb } from '../db'
9+
10+
import { fetchStatus } from './pkgGoDevClient'
11+
import { fetchLatest } from './proxyClient'
12+
import { isFetchError } from './types'
13+
14+
const log = getServiceChildLogger('go')
15+
16+
const PROXY_SOURCE = 'go-proxy'
17+
const PKGGODEV_SOURCE = 'pkg-go-dev'
18+
19+
// TODO: filter to critical packages once computed
20+
async function getGoBatch(
21+
qx: QueryExecutor,
22+
afterPurl: string,
23+
batchSize: number,
24+
): Promise<Array<{ purl: string; name: string }>> {
25+
return qx.select(
26+
`SELECT purl, name FROM packages
27+
WHERE ecosystem = 'go' AND purl > $(after)
28+
ORDER BY last_synced_at ASC NULLS FIRST, purl ASC
29+
LIMIT $(limit)`,
30+
{ after: afterPurl, limit: batchSize },
31+
)
32+
}
33+
34+
export async function enrichGoVersionsBatch(
35+
afterPurl: string,
36+
batchSize: number,
37+
): Promise<string | null> {
38+
const qx = await getPackagesDb()
39+
const rows = await getGoBatch(qx, afterPurl, batchSize)
40+
if (rows.length === 0) return null
41+
42+
const { fetchTimeoutMs, proxyConcurrency } = getGoConfig()
43+
44+
const enrichOne = async (row: { purl: string; name: string }): Promise<void> => {
45+
Context.current().heartbeat(row.purl)
46+
const result = await fetchLatest(row.name, fetchTimeoutMs)
47+
if (isFetchError(result)) {
48+
log.warn(
49+
{ purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode },
50+
'go proxy fetch failed — skipping package',
51+
)
52+
return
53+
}
54+
const changed = await qx.selectOne(
55+
`WITH old AS (
56+
SELECT latest_version AS v, latest_release_at AS t, repository_url AS r
57+
FROM packages WHERE purl = $(purl)
58+
),
59+
upd AS (
60+
UPDATE packages p SET
61+
latest_version = $(version),
62+
latest_release_at = $(releaseAt),
63+
repository_url = COALESCE($(repoUrl), p.repository_url),
64+
last_synced_at = NOW()
65+
WHERE p.purl = $(purl)
66+
RETURNING latest_version AS v, latest_release_at AS t, repository_url AS r
67+
)
68+
SELECT
69+
(SELECT v FROM old) IS DISTINCT FROM (SELECT v FROM upd) AS v_changed,
70+
(SELECT t FROM old) IS DISTINCT FROM (SELECT t FROM upd) AS t_changed,
71+
(SELECT r FROM old) IS DISTINCT FROM (SELECT r FROM upd) AS r_changed`,
72+
{
73+
version: result.version,
74+
releaseAt: result.releaseAt,
75+
repoUrl: result.repoUrl,
76+
purl: row.purl,
77+
},
78+
)
79+
const changedFields = [
80+
changed?.v_changed ? 'packages.latest_version' : null,
81+
changed?.t_changed ? 'packages.latest_release_at' : null,
82+
changed?.r_changed ? 'packages.repository_url' : null,
83+
].filter(Boolean) as string[]
84+
await logAuditFieldChanges(qx, PROXY_SOURCE, row.purl, changedFields)
85+
}
86+
87+
for (let i = 0; i < rows.length; i += proxyConcurrency) {
88+
await Promise.all(rows.slice(i, i + proxyConcurrency).map(enrichOne))
89+
}
90+
91+
log.info({ count: rows.length, concurrency: proxyConcurrency }, 'Enriched go versions batch')
92+
return rows[rows.length - 1].purl
93+
}
94+
95+
export async function enrichGoStatusBatch(
96+
afterPurl: string,
97+
batchSize: number,
98+
): Promise<string | null> {
99+
const qx = await getPackagesDb()
100+
const rows = await getGoBatch(qx, afterPurl, batchSize)
101+
if (rows.length === 0) return null
102+
103+
const { fetchTimeoutMs } = getGoConfig()
104+
for (const row of rows) {
105+
const result = await fetchStatus(row.name, fetchTimeoutMs, () =>
106+
Context.current().heartbeat(row.purl),
107+
)
108+
if (isFetchError(result)) {
109+
log.warn(
110+
{ purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode },
111+
'pkg.go.dev fetch failed — skipping package',
112+
)
113+
continue
114+
}
115+
const changed = await qx.selectOne(
116+
`WITH old AS (
117+
SELECT status AS s, versions_count AS vc FROM packages WHERE purl = $(purl)
118+
),
119+
upd AS (
120+
UPDATE packages p SET
121+
status = $(status),
122+
versions_count = COALESCE($(versionsCount), p.versions_count),
123+
last_synced_at = NOW()
124+
WHERE p.purl = $(purl)
125+
RETURNING status AS s, versions_count AS vc
126+
)
127+
SELECT
128+
(SELECT s FROM old) IS DISTINCT FROM (SELECT s FROM upd) AS s_changed,
129+
(SELECT vc FROM old) IS DISTINCT FROM (SELECT vc FROM upd) AS vc_changed`,
130+
{ status: result.status, versionsCount: result.versionsCount, purl: row.purl },
131+
)
132+
const changedFields = [
133+
changed?.s_changed ? 'packages.status' : null,
134+
changed?.vc_changed ? 'packages.versions_count' : null,
135+
].filter(Boolean) as string[]
136+
await logAuditFieldChanges(qx, PKGGODEV_SOURCE, row.purl, changedFields)
137+
}
138+
139+
log.info({ count: rows.length }, 'Enriched go status batch')
140+
return rows[rows.length - 1].purl
141+
}
Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,105 @@
1+
import { FetchError, GoStatusResult, isFetchError } from './types'
2+
3+
const BASE = process.env.PKGGODEV_BASE_URL ?? 'https://pkg.go.dev'
4+
// 40 QPS per IP. Keep a single in-process gap under that ceiling.
5+
const MIN_INTERVAL_MS = parseInt(process.env.PKGGODEV_MIN_INTERVAL_MS ?? '30', 10)
6+
const MAX_429_RETRIES = 5
7+
const MAX_PAGES = 20
8+
9+
interface VersionItem {
10+
version: string
11+
deprecated?: boolean
12+
retracted?: boolean
13+
latestVersion?: string
14+
}
15+
interface VersionsPage {
16+
items?: VersionItem[]
17+
total?: number
18+
nextPageToken?: string
19+
}
20+
21+
let lastRequestAt = 0
22+
23+
function sleep(ms: number): Promise<void> {
24+
return new Promise((r) => setTimeout(r, ms))
25+
}
26+
27+
async function throttle(): Promise<void> {
28+
const wait = lastRequestAt + MIN_INTERVAL_MS - Date.now()
29+
if (wait > 0) await sleep(wait)
30+
lastRequestAt = Date.now()
31+
}
32+
33+
async function getPage(url: string, timeoutMs: number): Promise<VersionsPage | FetchError> {
34+
for (let attempt = 0; attempt <= MAX_429_RETRIES; attempt++) {
35+
await throttle()
36+
const controller = new AbortController()
37+
const timer = setTimeout(() => controller.abort(), timeoutMs)
38+
let res: Response
39+
try {
40+
res = await fetch(url, { signal: controller.signal })
41+
} catch (e) {
42+
clearTimeout(timer)
43+
return { kind: 'TRANSIENT', message: `network error: ${(e as Error).message}` }
44+
}
45+
clearTimeout(timer)
46+
47+
if (res.status === 429) {
48+
const reset = parseInt(res.headers.get('x-ratelimit-reset') ?? '0', 10)
49+
const waitMs = reset ? Math.max(1000, reset * 1000 - Date.now() + 500) : 2000
50+
await sleep(waitMs)
51+
continue
52+
}
53+
// Any other 4xx is permanent (e.g. 400 for submodule/non-module-root paths) — skip, don't retry.
54+
if (res.status >= 400 && res.status < 500) {
55+
return { kind: 'NOT_FOUND', statusCode: res.status, message: `${res.status}` }
56+
}
57+
if (res.status !== 200) {
58+
return {
59+
kind: 'TRANSIENT',
60+
statusCode: res.status,
61+
message: `unexpected status ${res.status}`,
62+
}
63+
}
64+
try {
65+
return (await res.json()) as VersionsPage
66+
} catch {
67+
return { kind: 'MALFORMED', message: 'invalid json' }
68+
}
69+
}
70+
return { kind: 'RATE_LIMIT', statusCode: 429, message: '429 after retries' }
71+
}
72+
73+
// Module status from /v1beta/versions: 'deprecated' if the latest version is
74+
// deprecated/retracted, else 'active'. Pages are newest-first and every item carries
75+
// latestVersion, so the match is normally on page 1; paginate (token query param,
76+
// request otherwise verbatim) only if it isn't.
77+
export async function fetchStatus(
78+
module: string,
79+
timeoutMs: number,
80+
onHeartbeat?: () => void,
81+
): Promise<GoStatusResult | FetchError> {
82+
let token: string | undefined
83+
let versionsCount: number | null = null
84+
for (let page = 0; page < MAX_PAGES; page++) {
85+
const url = `${BASE}/v1beta/versions/${module}${token ? `?token=${encodeURIComponent(token)}` : ''}`
86+
const result = await getPage(url, timeoutMs)
87+
onHeartbeat?.()
88+
if (isFetchError(result)) return result
89+
90+
if (versionsCount === null && typeof result.total === 'number') versionsCount = result.total
91+
const items = result.items ?? []
92+
const latestVersion = items.find((i) => i.latestVersion)?.latestVersion
93+
const match = latestVersion ? items.find((i) => i.version === latestVersion) : undefined
94+
if (match) {
95+
return {
96+
status: match.deprecated || match.retracted ? 'deprecated' : 'active',
97+
versionsCount,
98+
}
99+
}
100+
101+
if (!result.nextPageToken) break
102+
token = result.nextPageToken
103+
}
104+
return { kind: 'NOT_FOUND', message: 'latest version entry not found' }
105+
}

0 commit comments

Comments
 (0)