Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
2 changes: 1 addition & 1 deletion scripts/builders/packages.env
Original file line number Diff line number Diff line change
@@ -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"
66 changes: 66 additions & 0 deletions scripts/services/go-worker.yaml
Original file line number Diff line number Diff line change
@@ -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
5 changes: 5 additions & 0 deletions services/apps/packages_worker/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
1 change: 1 addition & 0 deletions services/apps/packages_worker/src/activities.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,3 +24,4 @@ export {
cargoFlushAudit,
cargoCleanup,
} from './cargo/activities'
export { enrichGoVersionsBatch, enrichGoStatusBatch } from './go/activities'
9 changes: 9 additions & 0 deletions services/apps/packages_worker/src/bin/go-worker.ts
Original file line number Diff line number Diff line change
@@ -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()
})
7 changes: 7 additions & 0 deletions services/apps/packages_worker/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'),
Expand Down
141 changes: 141 additions & 0 deletions services/apps/packages_worker/src/go/activities.ts
Original file line number Diff line number Diff line change
@@ -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<Array<{ purl: string; name: string }>> {
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)`,
Comment on lines +26 to +29
{ after: afterPurl, limit: batchSize },
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cursor mismatches batch ordering

High Severity

getGoBatch orders rows by last_synced_at then purl, but pagination only applies purl > $(after) and the workflow cursor is the last row’s purl. After a batch, every package with a smaller purl that was not in that batch is excluded for the rest of the run, so the scan can finish early and leave Go packages unenriched until the next scheduled workflow.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 8f85149. Configure here.

}

export async function enrichGoVersionsBatch(
afterPurl: string,
batchSize: number,
): Promise<string | null> {
const qx = await getPackagesDb()
const rows = await getGoBatch(qx, afterPurl, batchSize)
if (rows.length === 0) return null

const { fetchTimeoutMs, proxyConcurrency } = getGoConfig()

Comment thread
mbani01 marked this conversation as resolved.
const enrichOne = async (row: { purl: string; name: string }): Promise<void> => {
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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Transient fetch errors swallowed

Medium Severity

Version and status enrichment treat every FetchError, including RATE_LIMIT and TRANSIENT, as a warn-and-skip. Unlike the npm worker, those cases never throw, so Temporal does not retry the activity and the batch cursor still advances past packages that were not updated.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 0946688. Configure here.

}
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<string | null> {
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
}
105 changes: 105 additions & 0 deletions services/apps/packages_worker/src/go/pkgGoDevClient.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
return new Promise((r) => setTimeout(r, ms))
}

async function throttle(): Promise<void> {
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<VersionsPage | FetchError> {
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' }
}
Comment thread
mbani01 marked this conversation as resolved.

// 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<GoStatusResult | FetchError> {
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' }
}
Loading
Loading