Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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: parseInt(process.env.GO_PROXY_CONCURRENCY ?? '10', 10),
}
}

export function getDockerhubConfig() {
return {
hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'),
Expand Down
134 changes: 134 additions & 0 deletions services/apps/packages_worker/src/go/activities.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
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'

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 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 },

Check failure on line 71 in services/apps/packages_worker/src/go/activities.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Replace `·version:·result.version,·releaseAt:·result.releaseAt,·repoUrl:·result.repoUrl,·purl:·row.purl` with `⏎········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) {
Context.current().heartbeat(row.purl)
const result = await fetchStatus(row.name, fetchTimeoutMs)
Comment thread
cursor[bot] marked this conversation as resolved.
Outdated
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
}
96 changes: 96 additions & 0 deletions services/apps/packages_worker/src/go/pkgGoDevClient.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
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}` }

Check failure on line 58 in services/apps/packages_worker/src/go/pkgGoDevClient.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Replace `·kind:·'TRANSIENT',·statusCode:·res.status,·message:·`unexpected·status·${res.status}`` with `⏎········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,
): 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)
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 }

Check failure on line 89 in services/apps/packages_worker/src/go/pkgGoDevClient.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Replace `·status:·match.deprecated·||·match.retracted·?·'deprecated'·:·'active',·versionsCount` with `⏎········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