diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 4eb8f742ab..7541586a7a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1348,6 +1348,9 @@ importers: jsonwebtoken: specifier: ^9.0.0 version: 9.0.3 + pg-copy-streams: + specifier: ^7.0.0 + version: 7.0.0 semver: specifier: ^7.6.0 version: 7.6.0 @@ -1370,6 +1373,9 @@ importers: '@types/node': specifier: ^20.8.2 version: 20.12.7 + '@types/pg-copy-streams': + specifier: ^1.2.5 + version: 1.2.5 '@types/semver': specifier: ^7.5.8 version: 7.5.8 @@ -4884,6 +4890,12 @@ packages: '@types/node@22.19.10': resolution: {integrity: sha512-tF5VOugLS/EuDlTBijk0MqABfP8UxgYazTLo3uIn3b4yJgg26QRbVYJYsDtHrjdDUIRfP70+VfhTTc+CE1yskw==} + '@types/pg-copy-streams@1.2.5': + resolution: {integrity: sha512-7D6/GYW2uHIaVU6S/5omI+6RZnwlZBpLQDZAH83xX1rjxAOK0f6/deKyyUTewxqts145VIGn6XWYz1YGf50G5g==} + + '@types/pg@8.20.0': + resolution: {integrity: sha512-bEPFOaMAHTEP1EzpvHTbmwR8UsFyHSKsRisLIHVMXnpNefSbGA1bD6CVy+qKjGSqmZqNqBDV2azOBo8TgkcVow==} + '@types/q@1.5.8': resolution: {integrity: sha512-hroOstUScF6zhIi+5+x0dzqrHA1EJi+Irri6b1fxolMTqqHIV/Cg77EtnQcZqZCu8hR3mX2BzIxN4/GzI68Kfw==} @@ -8700,6 +8712,9 @@ packages: pg-connection-string@2.6.4: resolution: {integrity: sha512-v+Z7W/0EO707aNMaAEfiGnGL9sxxumwLl2fJvCQtMn9Fxsg+lPpPkdcyBSv/KFgpGdYkMfn+EI1Or2EHjpgLCA==} + pg-copy-streams@7.0.0: + resolution: {integrity: sha512-zBvnY6wtaBRE2ae2xXWOOGMaNVPkXh1vhypAkNSKgMdciJeTyIQAHZaEeRAxUjs/p1El5jgzYmwG5u871Zj3dQ==} + pg-cursor@2.12.0: resolution: {integrity: sha512-rppw54OnuYZfMUjiJI2zJMwAjjt2V9EtLUb+t7V5tqwSE5Jxod+7vA7Y0FI6Nq976jNLciA0hoVkwvjjB8qzEw==} peerDependencies: @@ -10026,7 +10041,7 @@ packages: uuid@3.3.2: resolution: {integrity: sha512-yXJmeNaw3DnnKAOKJE51sL/ZaYfWJRl1pK9dr19YFCu0ObS231AB1/LbqTKRAQ5kw8A90rA6fr4riOUpTZvQZA==} - deprecated: Please upgrade to version 7 or higher. Older versions may use Math.random() in certain circumstances, which is known to be problematic. See https://v8.dev/blog/math-random for details. + deprecated: uuid@10 and below is no longer supported. For ESM codebases, update to uuid@latest. For CommonJS codebases, use uuid@11 (but be aware this version will likely be deprecated in 2028). hasBin: true uuid@8.3.2: @@ -10472,8 +10487,8 @@ snapshots: dependencies: '@aws-crypto/sha256-browser': 3.0.0 '@aws-crypto/sha256-js': 3.0.0 - '@aws-sdk/client-sso-oidc': 3.572.0 - '@aws-sdk/client-sts': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0) + '@aws-sdk/client-sso-oidc': 3.572.0(@aws-sdk/client-sts@3.572.0) + '@aws-sdk/client-sts': 3.572.0 '@aws-sdk/core': 3.572.0 '@aws-sdk/credential-provider-node': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0)(@aws-sdk/client-sts@3.572.0) '@aws-sdk/middleware-host-header': 3.567.0 @@ -10667,11 +10682,11 @@ snapshots: transitivePeerDependencies: - aws-crt - '@aws-sdk/client-sso-oidc@3.572.0': + '@aws-sdk/client-sso-oidc@3.572.0(@aws-sdk/client-sts@3.572.0)': dependencies: '@aws-crypto/sha256-browser': 3.0.0 '@aws-crypto/sha256-js': 3.0.0 - '@aws-sdk/client-sts': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0) + '@aws-sdk/client-sts': 3.572.0 '@aws-sdk/core': 3.572.0 '@aws-sdk/credential-provider-node': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0)(@aws-sdk/client-sts@3.572.0) '@aws-sdk/middleware-host-header': 3.567.0 @@ -10710,6 +10725,7 @@ snapshots: '@smithy/util-utf8': 2.3.0 tslib: 2.6.2 transitivePeerDependencies: + - '@aws-sdk/client-sts' - aws-crt '@aws-sdk/client-sso@3.556.0': @@ -10885,11 +10901,11 @@ snapshots: transitivePeerDependencies: - aws-crt - '@aws-sdk/client-sts@3.572.0(@aws-sdk/client-sso-oidc@3.572.0)': + '@aws-sdk/client-sts@3.572.0': dependencies: '@aws-crypto/sha256-browser': 3.0.0 '@aws-crypto/sha256-js': 3.0.0 - '@aws-sdk/client-sso-oidc': 3.572.0 + '@aws-sdk/client-sso-oidc': 3.572.0(@aws-sdk/client-sts@3.572.0) '@aws-sdk/core': 3.572.0 '@aws-sdk/credential-provider-node': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0)(@aws-sdk/client-sts@3.572.0) '@aws-sdk/middleware-host-header': 3.567.0 @@ -10928,7 +10944,6 @@ snapshots: '@smithy/util-utf8': 2.3.0 tslib: 2.6.2 transitivePeerDependencies: - - '@aws-sdk/client-sso-oidc' - aws-crt '@aws-sdk/client-sts@3.985.0': @@ -11094,7 +11109,7 @@ snapshots: '@aws-sdk/credential-provider-ini@3.572.0(@aws-sdk/client-sso-oidc@3.572.0)(@aws-sdk/client-sts@3.572.0)': dependencies: - '@aws-sdk/client-sts': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0) + '@aws-sdk/client-sts': 3.572.0 '@aws-sdk/credential-provider-env': 3.568.0 '@aws-sdk/credential-provider-process': 3.572.0 '@aws-sdk/credential-provider-sso': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0) @@ -11271,7 +11286,7 @@ snapshots: '@aws-sdk/credential-provider-web-identity@3.568.0(@aws-sdk/client-sts@3.572.0)': dependencies: - '@aws-sdk/client-sts': 3.572.0(@aws-sdk/client-sso-oidc@3.572.0) + '@aws-sdk/client-sts': 3.572.0 '@aws-sdk/types': 3.567.0 '@smithy/property-provider': 2.2.0 '@smithy/types': 2.12.0 @@ -11583,7 +11598,7 @@ snapshots: '@aws-sdk/token-providers@3.572.0(@aws-sdk/client-sso-oidc@3.572.0)': dependencies: - '@aws-sdk/client-sso-oidc': 3.572.0 + '@aws-sdk/client-sso-oidc': 3.572.0(@aws-sdk/client-sts@3.572.0) '@aws-sdk/types': 3.567.0 '@smithy/property-provider': 2.2.0 '@smithy/shared-ini-file-loader': 2.4.0 @@ -14029,6 +14044,17 @@ snapshots: dependencies: undici-types: 6.21.0 + '@types/pg-copy-streams@1.2.5': + dependencies: + '@types/node': 20.12.7 + '@types/pg': 8.20.0 + + '@types/pg@8.20.0': + dependencies: + '@types/node': 20.12.7 + pg-protocol: 1.6.1 + pg-types: 2.2.0 + '@types/q@1.5.8': {} '@types/qs@6.9.15': {} @@ -18374,6 +18400,8 @@ snapshots: pg-connection-string@2.6.4: {} + pg-copy-streams@7.0.0: {} + pg-cursor@2.12.0(pg@8.11.5): dependencies: pg: 8.11.5 diff --git a/scripts/builders/packages.env b/scripts/builders/packages.env index 38006ad26a..b8d1a41a68 100644 --- a/scripts/builders/packages.env +++ b/scripts/builders/packages.env @@ -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" +SERVICES="github-repos-enricher bq-dataset-ingest npm-worker maven-worker osv-worker dockerhub-sync cargo-worker" diff --git a/scripts/services/cargo-worker.yaml b/scripts/services/cargo-worker.yaml new file mode 100644 index 0000000000..1dba91921a --- /dev/null +++ b/scripts/services/cargo-worker.yaml @@ -0,0 +1,65 @@ +version: '3.1' + +x-env-args: &env-args + DOCKER_BUILDKIT: 1 + NODE_ENV: docker + SERVICE: cargo-worker + SHELL: /bin/sh + SUPPRESS_NO_CONFIG_WARNING: 'true' + CROWD_TEMPORAL_TASKQUEUE: cargo-worker + +services: + cargo-worker: + build: + context: ../../ + dockerfile: ./scripts/services/docker/Dockerfile.packages + command: 'pnpm run start:cargo-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 + + cargo-worker-dev: + build: + context: ../../ + dockerfile: ./scripts/services/docker/Dockerfile.packages + command: 'pnpm run dev:cargo-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 + hostname: cargo-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 diff --git a/services/apps/packages_worker/package.json b/services/apps/packages_worker/package.json index 52abd880b0..04a6f35efd 100644 --- a/services/apps/packages_worker/package.json +++ b/services/apps/packages_worker/package.json @@ -30,6 +30,9 @@ "start:maven-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=maven-worker tsx src/bin/maven-worker.ts", "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", "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", + "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", "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", @@ -61,6 +64,7 @@ "semver": "^7.6.0", "axios": "^1.16.1", "fast-xml-parser": "^5.8.0", + "pg-copy-streams": "^7.0.0", "tsx": "^4.7.1", "typescript": "^5.6.3", "undici": "^5.29.0", @@ -69,6 +73,7 @@ "devDependencies": { "@types/jsonwebtoken": "^9.0.0", "@types/node": "^20.8.2", + "@types/pg-copy-streams": "^1.2.5", "@types/semver": "^7.5.8", "@types/unzipper": "^0.10.10", "nodemon": "^3.0.1", diff --git a/services/apps/packages_worker/src/activities.ts b/services/apps/packages_worker/src/activities.ts index 41ee885e50..09a591c869 100644 --- a/services/apps/packages_worker/src/activities.ts +++ b/services/apps/packages_worker/src/activities.ts @@ -14,3 +14,13 @@ export * from './deps-dev/activities' export { osvSyncEcosystem, osvDeriveCriticalFlag } from './osv/activities' export { processMavenCriticalBatch, processMavenNonCriticalBatch } from './maven/activities' export { criticalityComputePageRank, rankPackages } from './criticality/activities' +export { + cargoDownloadAndLoad, + cargoEnrichPackages, + cargoEnrichVersions, + cargoEnrichRepos, + cargoEnrichMaintainers, + cargoEnrichDownloadsDaily, + cargoFlushAudit, + cargoCleanup, +} from './cargo/activities' diff --git a/services/apps/packages_worker/src/bin/cargo-worker.ts b/services/apps/packages_worker/src/bin/cargo-worker.ts new file mode 100644 index 0000000000..32c92f54a6 --- /dev/null +++ b/services/apps/packages_worker/src/bin/cargo-worker.ts @@ -0,0 +1,8 @@ +import { scheduleCargoSync } from '../cargo/schedule' +import { svc } from '../service' + +setImmediate(async () => { + await svc.init() + await scheduleCargoSync() + await svc.start() +}) diff --git a/services/apps/packages_worker/src/cargo/activities.ts b/services/apps/packages_worker/src/cargo/activities.ts new file mode 100644 index 0000000000..d5e3439021 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/activities.ts @@ -0,0 +1,70 @@ +import { rm } from 'node:fs/promises' + +import { getServiceChildLogger } from '@crowd/logging' + +import { getCargoConfig } from '../config' +import { getPackagesDb, getPackagesDbConnection } from '../db' + +import { DUMP_DIR, downloadAndExtractDump } from './dump' +import { + enrichDownloadsDaily, + enrichMaintainers, + enrichPackages, + enrichRepos, + enrichVersions, + flushAudit, +} from './enrich' +import { STAGING_SCHEMA, loadDump } from './loadDump' +import { + EnrichDownloadsDailyResult, + EnrichMaintainersResult, + EnrichPackagesResult, + EnrichReposResult, + EnrichVersionsResult, + LoadResult, +} from './types' + +const log = getServiceChildLogger('cargo-activity') + +// Config read at call time — workflow bundle imports this file for type discovery without env. +export async function cargoDownloadAndLoad(): Promise { + const { dumpUrl } = getCargoConfig() + const dumpDir = await downloadAndExtractDump(dumpUrl) + const qx = await getPackagesDb() + const conn = await getPackagesDbConnection() + const result = await loadDump(qx, conn, dumpDir) + log.info({ ...result }, 'cargo dump loaded') + return result +} + +export async function cargoEnrichPackages(): Promise { + return enrichPackages(await getPackagesDb()) +} + +export async function cargoEnrichVersions(): Promise { + return enrichVersions(await getPackagesDb()) +} + +export async function cargoEnrichRepos(): Promise { + return enrichRepos(await getPackagesDb()) +} + +export async function cargoEnrichMaintainers(): Promise { + return enrichMaintainers(await getPackagesDb()) +} + +export async function cargoEnrichDownloadsDaily(): Promise { + return enrichDownloadsDaily(await getPackagesDb()) +} + +export async function cargoFlushAudit(): Promise { + return flushAudit(await getPackagesDb()) +} + +// Best-effort: a crashed run self-heals on next run (schema rebuilt, DUMP_DIR cleared). +export async function cargoCleanup(): Promise { + const qx = await getPackagesDb() + await qx.result(`DROP SCHEMA IF EXISTS ${STAGING_SCHEMA} CASCADE`) + await rm(DUMP_DIR, { recursive: true, force: true }) + log.info('cargo cleanup complete') +} diff --git a/services/apps/packages_worker/src/cargo/dump.ts b/services/apps/packages_worker/src/cargo/dump.ts new file mode 100644 index 0000000000..21f6fbd7c9 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/dump.ts @@ -0,0 +1,57 @@ +import { execFile } from 'node:child_process' +import { createWriteStream } from 'node:fs' +import { mkdir, rm } from 'node:fs/promises' +import * as path from 'node:path' +import { Readable } from 'node:stream' +import { pipeline } from 'node:stream/promises' +import { promisify } from 'node:util' + +const execFileAsync = promisify(execFile) + +export const DUMP_DIR = '/tmp/cargo-dump' +const DOWNLOAD_TIMEOUT_MS = 60 * 60 * 1000 + +async function downloadTarball(url: string, target: string): Promise { + const ac = new AbortController() + const timer = setTimeout(() => ac.abort(), DOWNLOAD_TIMEOUT_MS) + + try { + let response: Response + try { + response = await fetch(url, { signal: ac.signal }) + } catch (err) { + throw new Error(`Failed to GET ${url}: ${(err as Error).message}`) + } + + if (!response.ok) { + throw new Error(`Unexpected response from ${url}: HTTP ${response.status}`) + } + + if (!response.body) { + throw new Error(`Empty body from ${url}`) + } + + try { + await pipeline(Readable.fromWeb(response.body), createWriteStream(target)) + } catch (err) { + throw new Error(`Stream failed for ${url}: ${(err as Error).message}`) + } + } finally { + clearTimeout(timer) + } +} + +export async function downloadAndExtractDump(url: string): Promise { + await rm(DUMP_DIR, { recursive: true, force: true }) + await mkdir(DUMP_DIR, { recursive: true }) + + const tarPath = path.join(DUMP_DIR, 'db-dump.tar.gz') + + await downloadTarball(url, tarPath) + + await execFileAsync('tar', ['-xzf', tarPath, '-C', DUMP_DIR, '--strip-components=1']) + + await rm(tarPath, { force: true }) + + return DUMP_DIR +} diff --git a/services/apps/packages_worker/src/cargo/enrich.ts b/services/apps/packages_worker/src/cargo/enrich.ts new file mode 100644 index 0000000000..58e1115103 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/enrich.ts @@ -0,0 +1,354 @@ +import { QueryExecutor } from '@crowd/data-access-layer' +import { getServiceChildLogger } from '@crowd/logging' + +import { STAGING_SCHEMA } from './loadDump' +import { + EnrichDownloadsDailyResult, + EnrichMaintainersResult, + EnrichPackagesResult, + EnrichReposResult, + EnrichVersionsResult, +} from './types' + +const log = getServiceChildLogger('cargo-enrich') + +export const AUDIT_WORKER = 'cargo-registry' +const INGESTION_SOURCE = 'cargo-registry' +const REPO_LINK_SOURCE = 'declared' // same convention as npm/maven for manifest-declared repo URLs +const REPO_LINK_CONFIDENCE = 0.8 + +// synchronous_commit off: skip WAL fsync on bulk writes — job is idempotent. +const WORK_MEM = '512MB' +const SYNC_COMMIT = 'off' + +// Reads settings back after applying — fails loudly if they didn't take. +async function withTunedSession( + qx: QueryExecutor, + phase: string, + fn: (tx: QueryExecutor) => Promise, +): Promise { + return qx.tx(async (tx) => { + await tx.result(`SET LOCAL work_mem = '${WORK_MEM}'`) + await tx.result(`SET LOCAL synchronous_commit = ${SYNC_COMMIT}`) + const s = await tx.selectOne( + `SELECT current_setting('work_mem') AS work_mem, + current_setting('synchronous_commit') AS synchronous_commit`, + ) + if (s.work_mem !== WORK_MEM || s.synchronous_commit !== SYNC_COMMIT) { + throw new Error(`cargo enrich (${phase}) session settings not applied: ${JSON.stringify(s)}`) + } + log.info({ phase, ...s }, 'enrich session tuned') + return fn(tx) + }) +} + +// Nullable fields are COALESCEd — dump nulls don't wipe existing values. +export async function enrichPackages(qx: QueryExecutor): Promise { + const row = await withTunedSession(qx, 'packages', (tx) => + tx.selectOne( + `WITH snap AS ( + SELECT p.id, p.status, p.description, p.homepage, p.declared_repository_url, + p.licenses, p.licenses_raw, p.keywords, p.versions_count, p.latest_version, + p.first_release_at, p.latest_release_at, p.dependent_count, + p.dependent_repos_count, p.downloads_last_30d + FROM packages p + JOIN ${STAGING_SCHEMA}.enrich_packages e ON e.package_id = p.id + ), + upd AS ( + UPDATE packages p SET + status = e.status, + description = COALESCE(e.description, p.description), + homepage = COALESCE(e.homepage, p.homepage), + declared_repository_url = COALESCE(e.declared_repository_url, p.declared_repository_url), + licenses = COALESCE(e.licenses, p.licenses), + licenses_raw = COALESCE(e.licenses_raw, p.licenses_raw), + keywords = COALESCE(e.keywords, p.keywords), + versions_count = e.versions_count, + latest_version = e.latest_version, + first_release_at = e.first_release_at, + latest_release_at = e.latest_release_at, + dependent_count = e.dependent_count, + dependent_repos_count = e.dependent_repos_count, + downloads_last_30d = e.downloads_last_30d, + ingestion_source = $(ingestionSource), + last_synced_at = NOW() + FROM ${STAGING_SCHEMA}.enrich_packages e + WHERE p.id = e.package_id + RETURNING p.id + ), + diff AS ( + SELECT s.id AS package_id, f.field + FROM snap s + JOIN ${STAGING_SCHEMA}.enrich_packages e ON e.package_id = s.id + CROSS JOIN LATERAL (VALUES + ('packages.status', s.status IS DISTINCT FROM e.status), + ('packages.description', s.description IS DISTINCT FROM COALESCE(e.description, s.description)), + ('packages.homepage', s.homepage IS DISTINCT FROM COALESCE(e.homepage, s.homepage)), + ('packages.declared_repository_url', s.declared_repository_url IS DISTINCT FROM COALESCE(e.declared_repository_url, s.declared_repository_url)), + ('packages.licenses', s.licenses IS DISTINCT FROM COALESCE(e.licenses, s.licenses)), + ('packages.licenses_raw', s.licenses_raw IS DISTINCT FROM COALESCE(e.licenses_raw, s.licenses_raw)), + ('packages.keywords', s.keywords IS DISTINCT FROM COALESCE(e.keywords, s.keywords)), + ('packages.versions_count', s.versions_count IS DISTINCT FROM e.versions_count), + ('packages.latest_version', s.latest_version IS DISTINCT FROM e.latest_version), + ('packages.first_release_at', s.first_release_at IS DISTINCT FROM e.first_release_at), + ('packages.latest_release_at', s.latest_release_at IS DISTINCT FROM e.latest_release_at), + ('packages.dependent_count', s.dependent_count IS DISTINCT FROM e.dependent_count), + ('packages.dependent_repos_count', s.dependent_repos_count IS DISTINCT FROM e.dependent_repos_count), + ('packages.downloads_last_30d', s.downloads_last_30d IS DISTINCT FROM e.downloads_last_30d) + ) AS f(field, changed) + WHERE f.changed + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT package_id, field FROM diff RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM upd)::int AS updated`, + { ingestionSource: INGESTION_SOURCE }, + ), + ) + return { updated: row.updated } +} + +// namespace/name from the package row; license stored as ARRAY[spdx_string]. +export async function enrichVersions(qx: QueryExecutor): Promise { + const row = await withTunedSession(qx, 'versions', (tx) => + tx.selectOne( + `WITH old AS ( + SELECT v.package_id, v.number, v.published_at, v.is_latest, v.is_prerelease, v.licenses + FROM versions v + WHERE v.package_id IN (SELECT package_id FROM ${STAGING_SCHEMA}.enrich_versions) + ), + ins AS ( + INSERT INTO versions ( + package_id, ecosystem, namespace, name, number, + published_at, is_latest, is_prerelease, licenses, last_synced_at, created_at + ) + SELECT e.package_id, 'cargo', p.namespace, p.name, e.number, + e.published_at, e.is_latest, e.is_prerelease, + CASE WHEN e.license IS NOT NULL THEN ARRAY[e.license] END, NOW(), NOW() + FROM ${STAGING_SCHEMA}.enrich_versions e + JOIN packages p ON p.id = e.package_id + ON CONFLICT (package_id, number) DO UPDATE SET + published_at = COALESCE(EXCLUDED.published_at, versions.published_at), + is_latest = EXCLUDED.is_latest, + is_prerelease = EXCLUDED.is_prerelease, + licenses = COALESCE(EXCLUDED.licenses, versions.licenses), + last_synced_at = NOW() + RETURNING package_id, number, published_at, is_latest, is_prerelease, licenses + ), + diff AS ( + SELECT ins.package_id, f.field + FROM ins + LEFT JOIN old o ON o.package_id = ins.package_id AND o.number = ins.number + CROSS JOIN LATERAL (VALUES + ('versions.number', o.number IS NULL), + ('versions.published_at', o.number IS NULL OR o.published_at IS DISTINCT FROM ins.published_at), + ('versions.is_latest', o.number IS NULL OR o.is_latest IS DISTINCT FROM ins.is_latest), + ('versions.is_prerelease', o.number IS NULL OR o.is_prerelease IS DISTINCT FROM ins.is_prerelease), + ('versions.licenses', o.number IS NULL OR o.licenses IS DISTINCT FROM ins.licenses) + ) AS f(field, changed) + WHERE f.changed + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT DISTINCT package_id, field FROM diff RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM ins)::int AS upserted`, + ), + ) + return { upserted: row.upserted } +} + +// Writes only url + host — other repo fields belong to the GitHub enricher. +export async function enrichRepos(qx: QueryExecutor): Promise { + return withTunedSession(qx, 'repos', async (tx) => { + const repoRow = await tx.selectOne( + `WITH new_repos AS ( + INSERT INTO repos (url, host, updated_at) + SELECT DISTINCT e.declared_repository_url, + CASE + WHEN e.declared_repository_url ~* '://([^/]+\\.)?github\\.com(/|$)' THEN 'github' + WHEN e.declared_repository_url ~* '://[^/]*gitlab' THEN 'gitlab' + WHEN e.declared_repository_url ~* '://([^/]+\\.)?bitbucket\\.org(/|$)' THEN 'bitbucket' + ELSE 'other' + END, + NOW() + FROM ${STAGING_SCHEMA}.enrich_packages e + WHERE e.declared_repository_url IS NOT NULL AND e.declared_repository_url LIKE 'http%' + ON CONFLICT (url) DO NOTHING + RETURNING url + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT e.package_id, f.field + FROM ${STAGING_SCHEMA}.enrich_packages e + JOIN new_repos nr ON nr.url = e.declared_repository_url + CROSS JOIN LATERAL (VALUES ('repos.url'), ('repos.host')) AS f(field) + RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM new_repos)::int AS repos`, + ) + + const linkRow = await tx.selectOne( + `WITH old AS ( + SELECT pr.package_id, pr.repo_id, pr.source, pr.confidence + FROM package_repos pr + WHERE pr.package_id IN ( + SELECT package_id FROM ${STAGING_SCHEMA}.enrich_packages WHERE declared_repository_url IS NOT NULL + ) + ), + ins AS ( + INSERT INTO package_repos (package_id, repo_id, source, confidence, created_at, verified_at) + SELECT e.package_id, r.id, $(source), $(confidence), NOW(), NOW() + FROM ${STAGING_SCHEMA}.enrich_packages e + JOIN repos r ON r.url = e.declared_repository_url + WHERE e.declared_repository_url IS NOT NULL + ON CONFLICT (package_id, repo_id) DO UPDATE SET + source = EXCLUDED.source, + confidence = EXCLUDED.confidence, + verified_at = NOW() + RETURNING package_id, repo_id, source, confidence + ), + diff AS ( + SELECT ins.package_id, f.field + FROM ins + LEFT JOIN old o ON o.package_id = ins.package_id AND o.repo_id = ins.repo_id + CROSS JOIN LATERAL (VALUES + ('package_repos.repo_id', o.repo_id IS NULL), + ('package_repos.source', o.repo_id IS NULL OR o.source IS DISTINCT FROM ins.source), + ('package_repos.confidence', o.repo_id IS NULL OR o.confidence IS DISTINCT FROM ins.confidence) + ) AS f(field, changed) + WHERE f.changed + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT package_id, field FROM diff RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM ins)::int AS links`, + { source: REPO_LINK_SOURCE, confidence: REPO_LINK_CONFIDENCE }, + ) + + return { repos: repoRow.repos, links: linkRow.links } + }) +} + +// Fully replaces each package's maintainer links (role='owner'). +export async function enrichMaintainers(qx: QueryExecutor): Promise { + return withTunedSession(qx, 'maintainers', async (tx) => { + await tx.result( + `DROP TABLE IF EXISTS ${STAGING_SCHEMA}.mnt_before; + CREATE TABLE ${STAGING_SCHEMA}.mnt_before AS + SELECT pm.package_id, pm.maintainer_id + FROM package_maintainers pm + WHERE pm.package_id IN (SELECT package_id FROM ${STAGING_SCHEMA}.enrich_maintainers)`, + ) + + const mntRow = await tx.selectOne( + `WITH old AS (SELECT username, github_login, url FROM maintainers WHERE ecosystem = 'cargo'), + ins AS ( + INSERT INTO maintainers (ecosystem, username, github_login, url, created_at, updated_at) + SELECT DISTINCT 'cargo', em.github_login, em.github_login, + 'https://github.com/' || em.github_login, NOW(), NOW() + FROM ${STAGING_SCHEMA}.enrich_maintainers em + ON CONFLICT (ecosystem, username) DO UPDATE SET + github_login = COALESCE(EXCLUDED.github_login, maintainers.github_login), + url = COALESCE(EXCLUDED.url, maintainers.url), + updated_at = NOW() + RETURNING username, github_login, url + ), + changed AS ( + SELECT ins.username, f.field + FROM ins + LEFT JOIN old o ON o.username = ins.username + CROSS JOIN LATERAL (VALUES + ('maintainers.github_login', o.username IS NULL OR o.github_login IS DISTINCT FROM ins.github_login), + ('maintainers.url', o.username IS NULL OR o.url IS DISTINCT FROM ins.url) + ) AS f(field, ch) + WHERE f.ch + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT DISTINCT em.package_id, c.field + FROM changed c + JOIN ${STAGING_SCHEMA}.enrich_maintainers em ON em.github_login = c.username + RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM ins)::int AS n`, + ) + + await tx.result( + `DELETE FROM package_maintainers + WHERE package_id IN (SELECT package_id FROM ${STAGING_SCHEMA}.enrich_maintainers)`, + ) + + const links = await tx.result( + `INSERT INTO package_maintainers (package_id, maintainer_id, role, created_at, updated_at) + SELECT em.package_id, m.id, 'owner', NOW(), NOW() + FROM ${STAGING_SCHEMA}.enrich_maintainers em + JOIN maintainers m ON m.ecosystem = 'cargo' AND m.username = em.github_login + ON CONFLICT (package_id, maintainer_id) DO NOTHING`, + ) + + // Membership delta (symmetric difference of before vs after) → audit. + await tx.result( + `INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT DISTINCT package_id, 'package_maintainers.maintainer_id' + FROM ( + (SELECT package_id, maintainer_id FROM ${STAGING_SCHEMA}.mnt_before + EXCEPT + SELECT package_id, maintainer_id FROM package_maintainers + WHERE package_id IN (SELECT package_id FROM ${STAGING_SCHEMA}.enrich_maintainers)) + UNION + (SELECT package_id, maintainer_id FROM package_maintainers + WHERE package_id IN (SELECT package_id FROM ${STAGING_SCHEMA}.enrich_maintainers) + EXCEPT + SELECT package_id, maintainer_id FROM ${STAGING_SCHEMA}.mnt_before) + ) d`, + ) + + return { maintainers: mntRow.n, links } + }) +} + +// Batched by date — downloads_daily is range-partitioned. DO NOTHING preserves history. +export async function enrichDownloadsDaily(qx: QueryExecutor): Promise { + return withTunedSession(qx, 'downloadsDaily', async (tx) => { + const dates: Array<{ date: string }> = await tx.select( + `SELECT DISTINCT date::text AS date FROM ${STAGING_SCHEMA}.enrich_downloads_daily ORDER BY date`, + ) + let inserted = 0 + for (const { date } of dates) { + const row = await tx.selectOne( + `WITH ins AS ( + INSERT INTO downloads_daily (package_id, date, count, created_at, updated_at) + SELECT package_id, date, downloads, NOW(), NOW() + FROM ${STAGING_SCHEMA}.enrich_downloads_daily + WHERE date = $(date)::date + ON CONFLICT (package_id, date) DO NOTHING + RETURNING package_id + ), + ins_audit AS ( + INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field) + SELECT package_id, f.field FROM ins + CROSS JOIN LATERAL (VALUES ('downloads_daily.date'), ('downloads_daily.count')) AS f(field) + RETURNING 1 + ) + SELECT (SELECT COUNT(*) FROM ins)::int AS inserted`, + { date }, + ) + inserted += row.inserted + } + return { inserted } + }) +} + +export async function flushAudit(qx: QueryExecutor): Promise { + return qx.result( + `INSERT INTO audit_field_changes (worker, purl, changed_fields) + SELECT $(worker), p.purl, array_agg(DISTINCT ac.field) + FROM ${STAGING_SCHEMA}.audit_changes ac + JOIN packages p ON p.id = ac.package_id + GROUP BY p.purl`, + { worker: AUDIT_WORKER }, + ) +} diff --git a/services/apps/packages_worker/src/cargo/loadDump.ts b/services/apps/packages_worker/src/cargo/loadDump.ts new file mode 100644 index 0000000000..bda080f3ac --- /dev/null +++ b/services/apps/packages_worker/src/cargo/loadDump.ts @@ -0,0 +1,324 @@ +import { createReadStream } from 'node:fs' +import * as path from 'node:path' +import { Transform } from 'node:stream' +import { pipeline } from 'node:stream/promises' +import { CopyStreamQuery, from as copyFrom } from 'pg-copy-streams' + +import { QueryExecutor } from '@crowd/data-access-layer' +import { DbConnection } from '@crowd/database' +import { getServiceChildLogger } from '@crowd/logging' + +import { LoadResult } from './types' + +const log = getServiceChildLogger('cargo-load') + +export const STAGING_SCHEMA = 'cargo_sync' + +// Column order matches each CSV header exactly so COPY FROM STDIN CSV HEADER maps by position. +const STAGING_DDL = ` + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.crates CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.crates ( + created_at text, description text, documentation text, homepage text, + id integer, max_features text, max_upload_size text, name text, + readme text, repository text, trustpub_only text, updated_at text + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.versions CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.versions ( + bin_names text, categories text, checksum text, crate_id integer, + crate_size text, created_at text, description text, documentation text, + downloads text, edition text, features text, has_lib text, homepage text, + id integer, keywords text, license text, links text, num text, + num_no_build text, published_by text, repository text, rust_version text, + updated_at text, yanked text + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.default_versions CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.default_versions ( + crate_id integer, num_versions integer, version_id integer + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.version_downloads CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.version_downloads ( + date date, downloads integer, version_id integer + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.dependencies CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.dependencies ( + crate_id integer, default_features text, explicit_name text, features text, + id integer, kind integer, optional text, req text, target text, + version_id integer + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.keywords CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.keywords ( + crates_cnt integer, created_at text, id integer, keyword text + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.crates_keywords CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.crates_keywords ( + crate_id integer, keyword_id integer + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.crate_owners CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.crate_owners ( + crate_id integer, created_at text, created_by text, owner_id integer, + owner_kind integer + ); + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.oauth_github CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.oauth_github ( + account_id text, avatar text, login text, user_id integer + ); +` + +// crates.csv needs NUL stripping — readme blobs contain 0x00 bytes that COPY rejects. +const CSV_FILES: Array<{ table: string; file: string; stripNul?: boolean }> = [ + { table: 'crates', file: 'crates.csv', stripNul: true }, + { table: 'versions', file: 'versions.csv' }, + { table: 'default_versions', file: 'default_versions.csv' }, + { table: 'version_downloads', file: 'version_downloads.csv' }, + { table: 'dependencies', file: 'dependencies.csv' }, + { table: 'keywords', file: 'keywords.csv' }, + { table: 'crates_keywords', file: 'crates_keywords.csv' }, + { table: 'crate_owners', file: 'crate_owners.csv' }, + { table: 'oauth_github', file: 'oauth_github.csv' }, +] + +const stripNul = (): Transform => + new Transform({ + transform(chunk: Buffer, _enc, cb) { + cb(null, chunk.includes(0) ? chunk.filter((b) => b !== 0) : chunk) + }, + }) + +async function copyCsv( + db: DbConnection, + table: string, + csvPath: string, + doStripNul: boolean, +): Promise { + const con = await db.connect() + try { + // @types/pg-copy-streams Submittable doesn't match @types/pg's query overload. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const stream: CopyStreamQuery = (con.client as any).query( + copyFrom(`COPY ${STAGING_SCHEMA}.${table} FROM STDIN CSV HEADER`), + ) + const source = createReadStream(csvPath) + await (doStripNul ? pipeline(source, stripNul(), stream) : pipeline(source, stream)) + const rows = stream.rowCount + con.done() + return rows + } catch (err) { + con.done(true) // destroy — a failed COPY can leave the client unusable + throw new Error(`Failed to COPY ${table} from ${csvPath}: ${(err as Error).message}`) + } +} + +const STAGING_INDEXES = ` + CREATE INDEX ON ${STAGING_SCHEMA}.versions (crate_id); + CREATE INDEX ON ${STAGING_SCHEMA}.versions (id); + CREATE INDEX ON ${STAGING_SCHEMA}.version_downloads (version_id); + CREATE INDEX ON ${STAGING_SCHEMA}.dependencies (version_id); + CREATE INDEX ON ${STAGING_SCHEMA}.default_versions (version_id); + CREATE INDEX ON ${STAGING_SCHEMA}.default_versions (crate_id); + CREATE INDEX ON ${STAGING_SCHEMA}.crates (id); + CREATE INDEX ON ${STAGING_SCHEMA}.crate_owners (crate_id); + CREATE INDEX ON ${STAGING_SCHEMA}.oauth_github (user_id); + CREATE INDEX ON ${STAGING_SCHEMA}.crates_keywords (crate_id); +` + +const AGGREGATIONS = ` + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_downloads_30d CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_downloads_30d AS + SELECT v.crate_id, SUM(vd.downloads)::bigint AS downloads_30d + FROM ${STAGING_SCHEMA}.version_downloads vd + JOIN ${STAGING_SCHEMA}.versions v ON v.id = vd.version_id + WHERE vd.date >= CURRENT_DATE - INTERVAL '30 days' + GROUP BY v.crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_downloads_30d (crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_downloads_daily CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_downloads_daily AS + SELECT v.crate_id, vd.date, SUM(vd.downloads)::bigint AS downloads + FROM ${STAGING_SCHEMA}.version_downloads vd + JOIN ${STAGING_SCHEMA}.versions v ON v.id = vd.version_id + WHERE vd.date >= CURRENT_DATE - INTERVAL '30 days' + GROUP BY v.crate_id, vd.date; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_downloads_daily (crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_dep_packages CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_dep_packages AS + SELECT d.crate_id AS target_crate_id, COUNT(DISTINCT dv.crate_id)::bigint AS dependent_packages + FROM ${STAGING_SCHEMA}.default_versions dv + JOIN ${STAGING_SCHEMA}.versions v ON v.id = dv.version_id AND v.yanked = 'f' + JOIN ${STAGING_SCHEMA}.dependencies d ON d.version_id = dv.version_id AND d.kind = 0 + GROUP BY d.crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_dep_packages (target_crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_dep_repos CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_dep_repos AS + SELECT d.crate_id AS target_crate_id, COUNT(DISTINCT c.repository)::bigint AS dependent_repos + FROM ${STAGING_SCHEMA}.default_versions dv + JOIN ${STAGING_SCHEMA}.versions v ON v.id = dv.version_id AND v.yanked = 'f' + JOIN ${STAGING_SCHEMA}.dependencies d ON d.version_id = dv.version_id AND d.kind = 0 + JOIN ${STAGING_SCHEMA}.crates c ON c.id = dv.crate_id + AND c.repository IS NOT NULL AND c.repository <> '' + GROUP BY d.crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_dep_repos (target_crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_keywords CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_keywords AS + SELECT ck.crate_id, array_agg(k.keyword ORDER BY k.keyword) AS keywords + FROM ${STAGING_SCHEMA}.crates_keywords ck + JOIN ${STAGING_SCHEMA}.keywords k ON k.id = ck.keyword_id + GROUP BY ck.crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_keywords (crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.agg_version_stats CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.agg_version_stats AS + SELECT crate_id, + COUNT(*)::integer AS versions_count, + MIN(created_at)::timestamptz AS first_release_at, + MAX(created_at)::timestamptz AS latest_release_at + FROM ${STAGING_SCHEMA}.versions + GROUP BY crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.agg_version_stats (crate_id); +` + +const DENORMALIZE = ` + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.crate_package CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.crate_package AS + SELECT c.id AS crate_id, p.id AS package_id + FROM ${STAGING_SCHEMA}.crates c + JOIN packages p ON p.ecosystem = 'cargo' + AND p.purl = 'pkg:cargo/' || LOWER(REPLACE(c.name, '-', '_')); + CREATE INDEX ON ${STAGING_SCHEMA}.crate_package (crate_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.enrich_packages CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.enrich_packages AS + SELECT + cp.package_id, + CASE WHEN v.yanked = 't' THEN 'deprecated' ELSE 'active' END AS status, + NULLIF(c.description, '') AS description, + NULLIF(c.homepage, '') AS homepage, + NULLIF(c.repository, '') AS declared_repository_url, + CASE WHEN NULLIF(v.license, '') IS NOT NULL THEN ARRAY[v.license] END AS licenses, + NULLIF(v.license, '') AS licenses_raw, + kw.keywords AS keywords, + vs.versions_count, + v.num AS latest_version, + vs.first_release_at, + vs.latest_release_at, + COALESCE(dp.dependent_packages, 0) AS dependent_count, + COALESCE(dr.dependent_repos, 0) AS dependent_repos_count, + COALESCE(d30.downloads_30d, 0) AS downloads_last_30d + FROM ${STAGING_SCHEMA}.crate_package cp + JOIN ${STAGING_SCHEMA}.crates c ON c.id = cp.crate_id + JOIN ${STAGING_SCHEMA}.default_versions dv ON dv.crate_id = c.id + JOIN ${STAGING_SCHEMA}.versions v ON v.id = dv.version_id + JOIN ${STAGING_SCHEMA}.agg_version_stats vs ON vs.crate_id = c.id + LEFT JOIN ${STAGING_SCHEMA}.agg_downloads_30d d30 ON d30.crate_id = c.id + LEFT JOIN ${STAGING_SCHEMA}.agg_dep_packages dp ON dp.target_crate_id = c.id + LEFT JOIN ${STAGING_SCHEMA}.agg_dep_repos dr ON dr.target_crate_id = c.id + LEFT JOIN ${STAGING_SCHEMA}.agg_keywords kw ON kw.crate_id = c.id; + CREATE INDEX ON ${STAGING_SCHEMA}.enrich_packages (package_id); + + -- Incremental: only crates newer than packages.latest_release_at (read before enrichPackages overwrites it). + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.enrich_versions CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.enrich_versions AS + SELECT + cp.package_id, + v.num AS number, + NULLIF(v.created_at, '')::timestamptz AS published_at, + (v.id = dv.version_id) AS is_latest, + (v.num LIKE '%-%') AS is_prerelease, + NULLIF(v.license, '') AS license + FROM ${STAGING_SCHEMA}.versions v + JOIN ${STAGING_SCHEMA}.crate_package cp ON cp.crate_id = v.crate_id + JOIN packages p ON p.id = cp.package_id + JOIN ${STAGING_SCHEMA}.agg_version_stats vs ON vs.crate_id = v.crate_id + LEFT JOIN ${STAGING_SCHEMA}.default_versions dv ON dv.crate_id = v.crate_id + WHERE p.latest_release_at IS NULL OR vs.latest_release_at > p.latest_release_at; + CREATE INDEX ON ${STAGING_SCHEMA}.enrich_versions (package_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.enrich_maintainers CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.enrich_maintainers AS + SELECT DISTINCT cp.package_id, og.login AS github_login + FROM ${STAGING_SCHEMA}.crate_owners co + JOIN ${STAGING_SCHEMA}.crate_package cp ON cp.crate_id = co.crate_id + JOIN ${STAGING_SCHEMA}.oauth_github og ON og.user_id = co.owner_id + WHERE co.owner_kind = 0; + CREATE INDEX ON ${STAGING_SCHEMA}.enrich_maintainers (package_id); + + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.enrich_downloads_daily CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.enrich_downloads_daily AS + SELECT cp.package_id, dd.date, dd.downloads + FROM ${STAGING_SCHEMA}.agg_downloads_daily dd + JOIN ${STAGING_SCHEMA}.crate_package cp ON cp.crate_id = dd.crate_id; + CREATE INDEX ON ${STAGING_SCHEMA}.enrich_downloads_daily (date); + + -- Audit scratch: enrich phases append (package_id, field); flushAudit aggregates per purl. + DROP TABLE IF EXISTS ${STAGING_SCHEMA}.audit_changes CASCADE; + CREATE TABLE ${STAGING_SCHEMA}.audit_changes (package_id bigint, field text); +` + +export async function loadDump( + qx: QueryExecutor, + db: DbConnection, + dumpDir: string, +): Promise { + const startedAt = Date.now() + const dataDir = path.join(dumpDir, 'data') + + log.info('Creating staging schema...') + await qx.result( + `CREATE SCHEMA IF NOT EXISTS ${STAGING_SCHEMA}; + ${STAGING_DDL}`, + ) + + const counts: Record = {} + for (const { table, file, stripNul: doStripNul } of CSV_FILES) { + counts[table] = await copyCsv(db, table, path.join(dataDir, file), doStripNul ?? false) + log.info({ table, rows: counts[table] }, 'Loaded staging table') + } + + // Guard against a truncated/empty dump — COPY can succeed with zero rows. + if (!counts.crates || !counts.versions) { + throw new Error( + `cargo dump load produced empty core tables (crates=${counts.crates}, versions=${counts.versions}) — dump may be truncated`, + ) + } + + log.info('Building staging indexes...') + await qx.result(STAGING_INDEXES) + + // Prevent autoanalyze mid-aggregation and give the planner accurate stats. + log.info('Analyzing staging tables...') + await qx.result(`ANALYZE ${STAGING_SCHEMA}.crates, ${STAGING_SCHEMA}.versions, + ${STAGING_SCHEMA}.version_downloads, ${STAGING_SCHEMA}.dependencies, + ${STAGING_SCHEMA}.default_versions, ${STAGING_SCHEMA}.crate_owners, + ${STAGING_SCHEMA}.oauth_github, ${STAGING_SCHEMA}.crates_keywords, + ${STAGING_SCHEMA}.keywords`) + + // work_mem: hash joins stay in RAM. max_parallel_workers=0: parallel workers + // contend on dsm over fresh tables and fail non-deterministically. + log.info('Aggregating and denormalizing...') + let matched = 0 + await qx.tx(async (tx) => { + await tx.result(`SET LOCAL work_mem = '512MB'`) + await tx.result(`SET LOCAL max_parallel_workers_per_gather = 0`) + await tx.result(AGGREGATIONS) + await tx.result(DENORMALIZE) + const row = await tx.selectOne( + `SELECT COUNT(*)::int AS n FROM ${STAGING_SCHEMA}.enrich_packages`, + ) + matched = row.n + }) + + const durationMs = Date.now() - startedAt + log.info({ matched, durationMs }, 'Dump load complete') + + return { + crates: counts.crates ?? 0, + versions: counts.versions ?? 0, + dependencies: counts.dependencies ?? 0, + versionDownloads: counts.version_downloads ?? 0, + owners: counts.crate_owners ?? 0, + matched, + durationMs, + } +} diff --git a/services/apps/packages_worker/src/cargo/schedule.ts b/services/apps/packages_worker/src/cargo/schedule.ts new file mode 100644 index 0000000000..3983ddb7a3 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/schedule.ts @@ -0,0 +1,41 @@ +import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client' + +import { svc } from '../service' +import { cargoSyncWorkflow } from '../workflows' + +const SCHEDULE_ID = 'cargo-registry-sync' + +// 06:00 UTC: after crates.io publishes (~02:00 UTC), clear of 03:00–05:00 ingest window. +export async function scheduleCargoSync(): Promise { + const { temporal } = svc + if (!temporal) throw new Error('Temporal client not initialized') + + try { + await temporal.schedule.create({ + scheduleId: SCHEDULE_ID, + spec: { cronExpressions: ['0 6 * * *'] }, + policies: { + overlap: ScheduleOverlapPolicy.SKIP, + catchupWindow: '1 hour', + }, + action: { + type: 'startWorkflow', + workflowType: cargoSyncWorkflow, + taskQueue: 'cargo-worker', + workflowExecutionTimeout: '3 hours', + retry: { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, + }, + args: [], + }, + }) + } catch (err) { + if (err instanceof ScheduleAlreadyRunning) { + svc.log.info(`Schedule ${SCHEDULE_ID} already registered.`) + } else { + throw err + } + } +} diff --git a/services/apps/packages_worker/src/cargo/types.ts b/services/apps/packages_worker/src/cargo/types.ts new file mode 100644 index 0000000000..bd1ac08d06 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/types.ts @@ -0,0 +1,35 @@ +export interface CargoConfig { + dumpUrl: string +} + +export interface LoadResult { + crates: number + versions: number + dependencies: number + versionDownloads: number + owners: number + matched: number + durationMs: number +} + +export interface EnrichPackagesResult { + updated: number +} + +export interface EnrichVersionsResult { + upserted: number +} + +export interface EnrichReposResult { + repos: number + links: number +} + +export interface EnrichMaintainersResult { + maintainers: number + links: number +} + +export interface EnrichDownloadsDailyResult { + inserted: number +} diff --git a/services/apps/packages_worker/src/cargo/workflows.ts b/services/apps/packages_worker/src/cargo/workflows.ts new file mode 100644 index 0000000000..1edb6dbf81 --- /dev/null +++ b/services/apps/packages_worker/src/cargo/workflows.ts @@ -0,0 +1,51 @@ +import { log, proxyActivities } from '@temporalio/workflow' + +import type * as activities from './activities' + +const RETRY = { + initialInterval: '30 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, +} + +const { cargoDownloadAndLoad } = proxyActivities({ + startToCloseTimeout: '90 minutes', + retry: RETRY, +}) + +const { + cargoEnrichPackages, + cargoEnrichVersions, + cargoEnrichRepos, + cargoEnrichMaintainers, + cargoEnrichDownloadsDaily, + cargoFlushAudit, + cargoCleanup, +} = proxyActivities({ + startToCloseTimeout: '30 minutes', + retry: RETRY, +}) + +// Sequential: phases touch disjoint tables; flushAudit before cleanup. +export async function cargoSyncWorkflow(): Promise { + const load = await cargoDownloadAndLoad() + log.info('cargoSync loaded dump', { ...load }) + + const packages = await cargoEnrichPackages() + const versions = await cargoEnrichVersions() + const repos = await cargoEnrichRepos() + const maintainers = await cargoEnrichMaintainers() + const downloads = await cargoEnrichDownloadsDaily() + const auditRows = await cargoFlushAudit() + await cargoCleanup() + + log.info('cargoSync complete', { + matched: load.matched, + packages, + versions, + repos, + maintainers, + downloads, + auditRows, + }) +} diff --git a/services/apps/packages_worker/src/config.ts b/services/apps/packages_worker/src/config.ts index cf96e0cb14..b524205ba5 100644 --- a/services/apps/packages_worker/src/config.ts +++ b/services/apps/packages_worker/src/config.ts @@ -57,6 +57,12 @@ export function getMavenConfig() { } } +export function getCargoConfig() { + return { + dumpUrl: process.env.CARGO_DUMP_URL ?? 'https://static.crates.io/db-dump.tar.gz', + } +} + export function getDockerhubConfig() { return { hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'), diff --git a/services/apps/packages_worker/src/db.ts b/services/apps/packages_worker/src/db.ts index f314861d8b..eb96d1d9b1 100644 --- a/services/apps/packages_worker/src/db.ts +++ b/services/apps/packages_worker/src/db.ts @@ -1,5 +1,5 @@ import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor' -import { getDbConnection } from '@crowd/database' +import { DbConnection, getDbConnection } from '@crowd/database' import { getPackagesDbConfig } from './config' @@ -7,3 +7,10 @@ export async function getPackagesDb() { const conn = await getDbConnection(getPackagesDbConfig()) return pgpQx(conn) } + +// Raw pg-promise connection for the same pool getPackagesDb() wraps +// (getDbConnection caches per host:database). Needed for COPY FROM STDIN via +// pg-copy-streams, which requires direct access to the underlying pg client. +export async function getPackagesDbConnection(): Promise { + return getDbConnection(getPackagesDbConfig()) +} diff --git a/services/apps/packages_worker/src/workflows/index.ts b/services/apps/packages_worker/src/workflows/index.ts index 3e4e1b8113..667ffa9180 100644 --- a/services/apps/packages_worker/src/workflows/index.ts +++ b/services/apps/packages_worker/src/workflows/index.ts @@ -18,3 +18,4 @@ export { osvSync } from '../osv/workflows' export { mavenCriticalWorkflow, mavenNonCriticalWorkflow } from '../maven/workflows' export { ingestScorecard } from '../scorecard/workflows' export { rankPackagesWorkflow } from '../criticality/workflow' +export { cargoSyncWorkflow } from '../cargo/workflows'