Skip to content

Commit 27c4836

Browse files
committed
refactor: simplify the schedule
Signed-off-by: Umberto Sgueglia <usgueglia@contractor.linuxfoundation.org>
1 parent 624d8ed commit 27c4836

19 files changed

Lines changed: 1211 additions & 434 deletions

File tree

backend/.env.dist.local

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -185,9 +185,9 @@ ENRICHER_BATCH_SIZE=100
185185
ENRICHER_REPO_UPDATE_INTERVAL_HOURS=24
186186
ENRICHER_IDLE_SLEEP_SEC=60
187187

188-
OSSPCKGS_GCP_PROJECT=
189-
OSSPCKGS_GCS_BUCKET=
190-
OSSPCKGS_GCP_CREDENTIALS_B64=
188+
OSSPCKGS_GCP_PROJECT=local-dev
189+
OSSPCKGS_GCS_BUCKET=local-dev
190+
OSSPCKGS_GCP_CREDENTIALS_B64=e30=
191191

192192
# osv-sync (Temporal-scheduled; see services/apps/packages_worker/src/osv/schedule.ts)
193193
# OSV_ECOSYSTEMS uses OSV's canonical bucket case (npm lowercase, Maven titlecase) because
@@ -200,13 +200,16 @@ OSV_TMP_DIR=/tmp/osv
200200
OSV_BATCH_SIZE=500
201201
OSV_DERIVE_BATCH_SIZE=1000
202202
# maven enricher
203-
POM_FETCHER_BATCH_SIZE=50
204-
POM_FETCHER_CONCURRENCY=3
203+
204+
POM_FETCHER_BATCH_SIZE=2000
205+
POM_FETCHER_CONCURRENCY=10
205206
POM_FETCHER_NON_CRITICAL_BATCH_SIZE=500
206207
POM_FETCHER_NON_CRITICAL_CONCURRENCY=20
207208
POM_FETCHER_REFRESH_DAYS=1
208-
POM_FETCHER_GROUP_DELAY_MS=500
209-
POM_FETCHER_IDLE_SLEEP_SEC=3600
209+
POM_FETCHER_GROUP_DELAY_MS=100
210210
# Set to 'true' on first run against a fresh/restored DB to skip the version-unchanged
211211
# optimisation and force full POM extraction. Set to 'false' after the first pass.
212-
POM_FETCHER_FORCE_FULL_EXTRACTION=false
212+
POM_FETCHER_FORCE_FULL_EXTRACTION=true
213+
POM_FETCHER_MAVEN_BASE_URL=https://maven-central.storage-download.googleapis.com/maven2
214+
MAVEN_SYNC_SOURCE=both
215+
MAVEN_DELTA_API_URL=https://maven-fetcher-production.up.railway.app

scripts/services/maven.yaml

Lines changed: 0 additions & 71 deletions
This file was deleted.

scripts/services/packages-worker.yaml

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,10 +8,6 @@ x-env-args: &env-args
88
CROWD_TEMPORAL_NAMESPACE: ${CROWD_PACKAGES_TEMPORAL_NAMESPACE}
99
SHELL: /bin/sh
1010
SUPPRESS_NO_CONFIG_WARNING: 'true'
11-
POM_FETCHER_BATCH_SIZE: '50'
12-
POM_FETCHER_CONCURRENCY: '3'
13-
POM_FETCHER_STALE_DAYS: '7'
14-
POM_FETCHER_IDLE_SLEEP_SEC: '3600'
1511

1612
services:
1713
packages-worker:

services/apps/packages_worker/package.json

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -6,23 +6,22 @@
66
"start:github-repos-enricher": "SERVICE=github-repos-enricher tsx src/bin/github-repos-enricher.ts",
77
"start:packages-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=packages-worker tsx src/bin/packages-worker.ts",
88
"start:pom-fetcher": "SERVICE=pom-fetcher tsx src/bin/pom-fetcher.ts",
9+
"backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts",
910
"dev:packages-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=packages-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9233 src/bin/packages-worker.ts",
1011
"dev:deps-dev-ingest": "CROWD_TEMPORAL_TASKQUEUE=deps-dev-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=deps-dev-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/deps-dev-ingest.ts",
1112
"start:maven": "SERVICE=maven tsx src/bin/maven.ts",
1213
"dev:github-repos-enricher": "SERVICE=github-repos-enricher LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9234 src/bin/github-repos-enricher.ts",
13-
"dev:maven": "SERVICE=maven LOG_LEVEL=info nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/maven.ts",
1414
"dev:packages-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=packages-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9233 src/bin/packages-worker.ts",
1515
"dev:deps-dev-ingest:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=deps-dev-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=deps-dev-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/deps-dev-ingest.ts",
1616
"dev:github-repos-enricher:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=github-repos-enricher LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9234 src/bin/github-repos-enricher.ts",
1717
"export-to-bucket": "SERVICE=deps-dev-ingest tsx src/scripts/exportToBucket.ts",
1818
"export-to-bucket:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=deps-dev-ingest tsx src/scripts/exportToBucket.ts",
19-
"monitor:osspckgs": "SERVICE=monitor tsx src/scripts/monitorOsspckgs.ts",
20-
"monitor:osspckgs:local": "bash -c 'set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=monitor tsx src/scripts/monitorOsspckgs.ts'",
21-
"trigger-bootstrap": "SERVICE=deps-dev-ingest tsx src/scripts/triggerBootstrap.ts",
22-
"trigger-bootstrap:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=deps-dev-ingest tsx src/scripts/triggerBootstrap.ts",
23-
"dev:pom-fetcher": "SERVICE=pom-fetcher LOG_LEVEL=info nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/pom-fetcher.ts",
24-
"dev:pom-fetcher:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=pom-fetcher LOG_LEVEL=info nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/pom-fetcher.ts",
25-
"dev:maven:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=maven LOG_LEVEL=info nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/maven.ts",
19+
"monitor:osspckgs:local": "bash -c 'set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && node ../../../scripts/monitor-osspckgs.mjs'",
20+
"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",
21+
"benchmark:maven-delta": "SERVICE=maven tsx src/scripts/benchmarkMavenDelta.ts",
22+
"benchmark:maven-delta:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=maven LOG_LEVEL=info tsx src/scripts/benchmarkMavenDelta.ts",
23+
"validate:maven-quality": "SERVICE=maven tsx src/scripts/validateDataQuality.ts",
24+
"validate:maven-quality:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=maven LOG_LEVEL=info tsx src/scripts/validateDataQuality.ts",
2625
"lint": "npx eslint --ext .ts src --max-warnings=0",
2726
"format": "npx prettier --write \"src/**/*.ts\"",
2827
"format-check": "npx prettier --check .",
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
import { getServiceLogger } from '@crowd/logging'
2+
3+
import { getMavenConfig } from '../config'
4+
import { getPackagesDb } from '../db'
5+
import { runMavenCriticalBackfill } from '../maven/runMavenEnrichmentLoop'
6+
7+
const log = getServiceLogger()
8+
9+
let shuttingDown = false
10+
11+
// Graceful stop: finish the in-flight batch, then exit. Safe to interrupt — every
12+
// write is an idempotent upsert and the DB state is the cursor, so re-running
13+
// resumes where it left off.
14+
const shutdown = () => {
15+
if (shuttingDown) return
16+
shuttingDown = true
17+
log.info('Shutting down maven backfill (stopping after the current batch)...')
18+
}
19+
20+
process.on('SIGINT', shutdown)
21+
process.on('SIGTERM', shutdown)
22+
23+
const main = async () => {
24+
log.info('maven backfill starting (one-shot, full extraction)...')
25+
26+
const config = getMavenConfig()
27+
log.info(config, 'Config loaded')
28+
29+
const qx = await getPackagesDb()
30+
await qx.selectOne('SELECT 1')
31+
log.info('Connected to packages-db.')
32+
33+
const totals = await runMavenCriticalBackfill(qx, config, () => shuttingDown)
34+
35+
log.info({ ...totals }, 'maven backfill complete')
36+
process.exit(0)
37+
}
38+
39+
main().catch((err) => {
40+
log.error({ err }, 'maven backfill fatal error')
41+
process.exit(1)
42+
})

services/apps/packages_worker/src/bin/maven.ts

Lines changed: 0 additions & 39 deletions
This file was deleted.
Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,12 @@
1+
import { scheduleMavenCritical } from '../maven/schedule'
12
import { scheduleNpmIngest } from '../npm/schedule'
23
import { scheduleOsvSync } from '../osv/schedule'
3-
import { scheduleMavenCritical, scheduleMavenNonCritical } from '../maven/schedule'
44
import { svc } from '../service'
55

66
setImmediate(async () => {
77
await svc.init()
88
await scheduleNpmIngest()
99
await scheduleOsvSync()
1010
await scheduleMavenCritical()
11-
await scheduleMavenNonCritical()
1211
await svc.start()
1312
})

services/apps/packages_worker/src/config.ts

Lines changed: 31 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,15 +33,44 @@ export function getEnricherConfig() {
3333
}
3434
}
3535

36+
// Which source drives the critical Maven sync:
37+
// 'maven' → poll packages_universe by staleness (current behaviour, default/fallback)
38+
// 'api' → only enrich what our delta feed reports as changed
39+
// 'both' → run both passes in the same Temporal tick
40+
export type MavenSyncSource = 'api' | 'maven' | 'both'
41+
42+
function parseMavenSyncSource(raw: string | undefined): MavenSyncSource {
43+
if (raw === 'api' || raw === 'both') return raw
44+
// Anything else (unset, typo, legacy value) falls back to the current behaviour.
45+
return 'maven'
46+
}
47+
3648
export function getMavenConfig() {
49+
const syncSource = parseMavenSyncSource(process.env.MAVEN_SYNC_SOURCE)
50+
const deltaApiBaseUrl = (process.env.MAVEN_DELTA_API_URL ?? '').replace(/\/+$/, '')
51+
52+
if (syncSource !== 'maven' && !deltaApiBaseUrl) {
53+
throw new Error(`MAVEN_SYNC_SOURCE='${syncSource}' requires MAVEN_DELTA_API_URL to be set`)
54+
}
55+
3756
return {
3857
batchSize: requireEnvInt('POM_FETCHER_BATCH_SIZE'),
3958
concurrency: requireEnvInt('POM_FETCHER_CONCURRENCY'),
4059
nonCriticalBatchSize: requireEnvInt('POM_FETCHER_NON_CRITICAL_BATCH_SIZE'),
4160
nonCriticalConcurrency: requireEnvInt('POM_FETCHER_NON_CRITICAL_CONCURRENCY'),
4261
refreshDays: requireEnvInt('POM_FETCHER_REFRESH_DAYS'),
4362
groupDelayMs: requireEnvInt('POM_FETCHER_GROUP_DELAY_MS'),
44-
idleSleepSec: requireEnvInt('POM_FETCHER_IDLE_SLEEP_SEC'),
45-
forceFullExtraction: requireEnv('POM_FETCHER_FORCE_FULL_EXTRACTION') === 'true',
63+
syncSource,
64+
deltaApi: {
65+
baseUrl: deltaApiBaseUrl,
66+
token: process.env.MAVEN_DELTA_API_TOKEN || undefined,
67+
pageSize: process.env.MAVEN_DELTA_API_PAGE_SIZE
68+
? parseInt(process.env.MAVEN_DELTA_API_PAGE_SIZE, 10)
69+
: 100,
70+
lookbackMinutes: process.env.MAVEN_DELTA_API_LOOKBACK_MINUTES
71+
? parseInt(process.env.MAVEN_DELTA_API_LOOKBACK_MINUTES, 10)
72+
: 15,
73+
includePrerelease: process.env.MAVEN_DELTA_API_INCLUDE_PRERELEASE === 'true',
74+
},
4675
}
4776
}

0 commit comments

Comments
 (0)