Skip to content

Commit 027f453

Browse files
committed
feat: go packages worker
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent 43395bb commit 027f453

14 files changed

Lines changed: 575 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: 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: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
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+
async function getGoBatch(
20+
qx: QueryExecutor,
21+
afterPurl: string,
22+
batchSize: number,
23+
): Promise<Array<{ purl: string; name: string }>> {
24+
return qx.select(
25+
`SELECT purl, name FROM packages
26+
WHERE ecosystem = 'go' AND purl > $(after)
27+
ORDER BY purl ASC
28+
LIMIT $(limit)`,
29+
{ after: afterPurl, limit: batchSize },
30+
)
31+
}
32+
33+
export async function enrichGoVersionsBatch(
34+
afterPurl: string,
35+
batchSize: number,
36+
): Promise<string | null> {
37+
const qx = await getPackagesDb()
38+
const rows = await getGoBatch(qx, afterPurl, batchSize)
39+
if (rows.length === 0) return null
40+
41+
const { fetchTimeoutMs, proxyConcurrency } = getGoConfig()
42+
43+
const enrichOne = async (row: { purl: string; name: string }): Promise<void> => {
44+
Context.current().heartbeat(row.purl)
45+
const result = await fetchLatest(row.name, fetchTimeoutMs)
46+
if (isFetchError(result)) {
47+
log.warn(
48+
{ purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode },
49+
'go proxy fetch failed — skipping package',
50+
)
51+
return
52+
}
53+
const changed = await qx.selectOne(
54+
`WITH old AS (
55+
SELECT latest_version AS v, latest_release_at AS t, repository_url AS r
56+
FROM packages WHERE purl = $(purl)
57+
),
58+
upd AS (
59+
UPDATE packages p SET
60+
latest_version = $(version),
61+
latest_release_at = $(releaseAt),
62+
repository_url = COALESCE($(repoUrl), p.repository_url),
63+
last_synced_at = NOW()
64+
WHERE p.purl = $(purl)
65+
RETURNING latest_version AS v, latest_release_at AS t, repository_url AS r
66+
)
67+
SELECT
68+
(SELECT v FROM old) IS DISTINCT FROM (SELECT v FROM upd) AS v_changed,
69+
(SELECT t FROM old) IS DISTINCT FROM (SELECT t FROM upd) AS t_changed,
70+
(SELECT r FROM old) IS DISTINCT FROM (SELECT r FROM upd) AS r_changed`,
71+
{ version: result.version, releaseAt: result.releaseAt, repoUrl: result.repoUrl, purl: row.purl },
72+
)
73+
const changedFields = [
74+
changed?.v_changed ? 'packages.latest_version' : null,
75+
changed?.t_changed ? 'packages.latest_release_at' : null,
76+
changed?.r_changed ? 'packages.repository_url' : null,
77+
].filter(Boolean) as string[]
78+
await logAuditFieldChanges(qx, PROXY_SOURCE, row.purl, changedFields)
79+
}
80+
81+
for (let i = 0; i < rows.length; i += proxyConcurrency) {
82+
await Promise.all(rows.slice(i, i + proxyConcurrency).map(enrichOne))
83+
}
84+
85+
log.info({ count: rows.length, concurrency: proxyConcurrency }, 'Enriched go versions batch')
86+
return rows[rows.length - 1].purl
87+
}
88+
89+
export async function enrichGoStatusBatch(
90+
afterPurl: string,
91+
batchSize: number,
92+
): Promise<string | null> {
93+
const qx = await getPackagesDb()
94+
const rows = await getGoBatch(qx, afterPurl, batchSize)
95+
if (rows.length === 0) return null
96+
97+
const { fetchTimeoutMs } = getGoConfig()
98+
for (const row of rows) {
99+
Context.current().heartbeat(row.purl)
100+
const result = await fetchStatus(row.name, fetchTimeoutMs)
101+
if (isFetchError(result)) {
102+
log.warn(
103+
{ purl: row.purl, name: row.name, kind: result.kind, statusCode: result.statusCode },
104+
'pkg.go.dev fetch failed — skipping package',
105+
)
106+
continue
107+
}
108+
const changed = await qx.selectOne(
109+
`WITH old AS (
110+
SELECT status AS s, versions_count AS vc FROM packages WHERE purl = $(purl)
111+
),
112+
upd AS (
113+
UPDATE packages p SET
114+
status = $(status),
115+
versions_count = COALESCE($(versionsCount), p.versions_count),
116+
last_synced_at = NOW()
117+
WHERE p.purl = $(purl)
118+
RETURNING status AS s, versions_count AS vc
119+
)
120+
SELECT
121+
(SELECT s FROM old) IS DISTINCT FROM (SELECT s FROM upd) AS s_changed,
122+
(SELECT vc FROM old) IS DISTINCT FROM (SELECT vc FROM upd) AS vc_changed`,
123+
{ status: result.status, versionsCount: result.versionsCount, purl: row.purl },
124+
)
125+
const changedFields = [
126+
changed?.s_changed ? 'packages.status' : null,
127+
changed?.vc_changed ? 'packages.versions_count' : null,
128+
].filter(Boolean) as string[]
129+
await logAuditFieldChanges(qx, PKGGODEV_SOURCE, row.purl, changedFields)
130+
}
131+
132+
log.info({ count: rows.length }, 'Enriched go status batch')
133+
return rows[rows.length - 1].purl
134+
}
Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,96 @@
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 { kind: 'TRANSIENT', statusCode: res.status, message: `unexpected status ${res.status}` }
59+
}
60+
try {
61+
return (await res.json()) as VersionsPage
62+
} catch {
63+
return { kind: 'MALFORMED', message: 'invalid json' }
64+
}
65+
}
66+
return { kind: 'RATE_LIMIT', statusCode: 429, message: '429 after retries' }
67+
}
68+
69+
// Module status from /v1beta/versions: 'deprecated' if the latest version is
70+
// deprecated/retracted, else 'active'. Pages are newest-first and every item carries
71+
// latestVersion, so the match is normally on page 1; paginate (token query param,
72+
// request otherwise verbatim) only if it isn't.
73+
export async function fetchStatus(
74+
module: string,
75+
timeoutMs: number,
76+
): Promise<GoStatusResult | FetchError> {
77+
let token: string | undefined
78+
let versionsCount: number | null = null
79+
for (let page = 0; page < MAX_PAGES; page++) {
80+
const url = `${BASE}/v1beta/versions/${module}${token ? `?token=${encodeURIComponent(token)}` : ''}`
81+
const result = await getPage(url, timeoutMs)
82+
if (isFetchError(result)) return result
83+
84+
if (versionsCount === null && typeof result.total === 'number') versionsCount = result.total
85+
const items = result.items ?? []
86+
const latestVersion = items.find((i) => i.latestVersion)?.latestVersion
87+
const match = latestVersion ? items.find((i) => i.version === latestVersion) : undefined
88+
if (match) {
89+
return { status: match.deprecated || match.retracted ? 'deprecated' : 'active', versionsCount }
90+
}
91+
92+
if (!result.nextPageToken) break
93+
token = result.nextPageToken
94+
}
95+
return { kind: 'NOT_FOUND', message: 'latest version entry not found' }
96+
}

0 commit comments

Comments
 (0)