Skip to content

Commit 7fe5c53

Browse files
authored
feat: cargo registry implementation [CM-1264] (#4236)
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent b1795fa commit 7fe5c53

16 files changed

Lines changed: 1075 additions & 13 deletions

File tree

pnpm-lock.yaml

Lines changed: 39 additions & 11 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

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"
4+
SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker"

scripts/services/cargo-worker.yaml

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

services/apps/packages_worker/package.json

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,9 @@
3030
"start:maven-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=maven-worker tsx src/bin/maven-worker.ts",
3131
"dev:maven-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=maven-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/maven-worker.ts",
3232
"dev:maven-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=maven-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/maven-worker.ts",
33+
"start:cargo-worker": "CROWD_TEMPORAL_TASKQUEUE=cargo-worker SERVICE=cargo-worker tsx src/bin/cargo-worker.ts",
34+
"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",
35+
"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",
3336
"backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts",
3437
"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",
3538
"backfill:stewardship": "SERVICE=stewardship-backfill tsx src/bin/stewardship-backfill.ts",
@@ -61,6 +64,7 @@
6164
"semver": "^7.6.0",
6265
"axios": "^1.16.1",
6366
"fast-xml-parser": "^5.8.0",
67+
"pg-copy-streams": "^7.0.0",
6468
"tsx": "^4.7.1",
6569
"typescript": "^5.6.3",
6670
"undici": "^5.29.0",
@@ -69,6 +73,7 @@
6973
"devDependencies": {
7074
"@types/jsonwebtoken": "^9.0.0",
7175
"@types/node": "^20.8.2",
76+
"@types/pg-copy-streams": "^1.2.5",
7277
"@types/semver": "^7.5.8",
7378
"@types/unzipper": "^0.10.10",
7479
"nodemon": "^3.0.1",

services/apps/packages_worker/src/activities.ts

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,3 +14,13 @@ export * from './deps-dev/activities'
1414
export { osvSyncEcosystem, osvDeriveCriticalFlag } from './osv/activities'
1515
export { processMavenCriticalBatch, processMavenNonCriticalBatch } from './maven/activities'
1616
export { criticalityComputePageRank, rankPackages } from './criticality/activities'
17+
export {
18+
cargoDownloadAndLoad,
19+
cargoEnrichPackages,
20+
cargoEnrichVersions,
21+
cargoEnrichRepos,
22+
cargoEnrichMaintainers,
23+
cargoEnrichDownloadsDaily,
24+
cargoFlushAudit,
25+
cargoCleanup,
26+
} from './cargo/activities'
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
import { scheduleCargoSync } from '../cargo/schedule'
2+
import { svc } from '../service'
3+
4+
setImmediate(async () => {
5+
await svc.init()
6+
await scheduleCargoSync()
7+
await svc.start()
8+
})
Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
import { rm } from 'node:fs/promises'
2+
3+
import { getServiceChildLogger } from '@crowd/logging'
4+
5+
import { getCargoConfig } from '../config'
6+
import { getPackagesDb, getPackagesDbConnection } from '../db'
7+
8+
import { DUMP_DIR, downloadAndExtractDump } from './dump'
9+
import {
10+
enrichDownloadsDaily,
11+
enrichMaintainers,
12+
enrichPackages,
13+
enrichRepos,
14+
enrichVersions,
15+
flushAudit,
16+
} from './enrich'
17+
import { STAGING_SCHEMA, loadDump } from './loadDump'
18+
import {
19+
EnrichDownloadsDailyResult,
20+
EnrichMaintainersResult,
21+
EnrichPackagesResult,
22+
EnrichReposResult,
23+
EnrichVersionsResult,
24+
LoadResult,
25+
} from './types'
26+
27+
const log = getServiceChildLogger('cargo-activity')
28+
29+
// Config read at call time — workflow bundle imports this file for type discovery without env.
30+
export async function cargoDownloadAndLoad(): Promise<LoadResult> {
31+
const { dumpUrl } = getCargoConfig()
32+
const dumpDir = await downloadAndExtractDump(dumpUrl)
33+
const qx = await getPackagesDb()
34+
const conn = await getPackagesDbConnection()
35+
const result = await loadDump(qx, conn, dumpDir)
36+
log.info({ ...result }, 'cargo dump loaded')
37+
return result
38+
}
39+
40+
export async function cargoEnrichPackages(): Promise<EnrichPackagesResult> {
41+
return enrichPackages(await getPackagesDb())
42+
}
43+
44+
export async function cargoEnrichVersions(): Promise<EnrichVersionsResult> {
45+
return enrichVersions(await getPackagesDb())
46+
}
47+
48+
export async function cargoEnrichRepos(): Promise<EnrichReposResult> {
49+
return enrichRepos(await getPackagesDb())
50+
}
51+
52+
export async function cargoEnrichMaintainers(): Promise<EnrichMaintainersResult> {
53+
return enrichMaintainers(await getPackagesDb())
54+
}
55+
56+
export async function cargoEnrichDownloadsDaily(): Promise<EnrichDownloadsDailyResult> {
57+
return enrichDownloadsDaily(await getPackagesDb())
58+
}
59+
60+
export async function cargoFlushAudit(): Promise<number> {
61+
return flushAudit(await getPackagesDb())
62+
}
63+
64+
// Best-effort: a crashed run self-heals on next run (schema rebuilt, DUMP_DIR cleared).
65+
export async function cargoCleanup(): Promise<void> {
66+
const qx = await getPackagesDb()
67+
await qx.result(`DROP SCHEMA IF EXISTS ${STAGING_SCHEMA} CASCADE`)
68+
await rm(DUMP_DIR, { recursive: true, force: true })
69+
log.info('cargo cleanup complete')
70+
}

0 commit comments

Comments
 (0)