Skip to content

Commit f07fc1b

Browse files
authored
feat: implement nuget registry worker [CM-1276] (#4265)
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent 0210c52 commit f07fc1b

19 files changed

Lines changed: 1229 additions & 1 deletion

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
ALTER TABLE packages ADD COLUMN IF NOT EXISTS total_downloads bigint;

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

scripts/services/nuget-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: nuget-worker
7+
SHELL: /bin/sh
8+
SUPPRESS_NO_CONFIG_WARNING: 'true'
9+
CROWD_TEMPORAL_TASKQUEUE: nuget-worker
10+
11+
services:
12+
nuget-worker:
13+
build:
14+
context: ../../
15+
dockerfile: ./scripts/services/docker/Dockerfile.packages
16+
command: 'pnpm run start:nuget-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+
nuget-worker-dev:
30+
build:
31+
context: ../../
32+
dockerfile: ./scripts/services/docker/Dockerfile.packages
33+
command: 'pnpm run dev:nuget-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: nuget-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
@@ -23,6 +23,8 @@
2323
"trigger-bootstrap:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=bq-dataset-ingest tsx src/scripts/triggerBootstrap.ts",
2424
"trigger-go": "SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts",
2525
"trigger-go:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=go-worker tsx src/scripts/triggerGoEnrich.ts",
26+
"trigger-nuget": "SERVICE=nuget-worker tsx src/scripts/triggerNuGetSync.ts",
27+
"trigger-nuget:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=nuget-worker tsx src/scripts/triggerNuGetSync.ts",
2628
"start:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker SERVICE=npm-worker tsx src/bin/npm-worker.ts",
2729
"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",
2830
"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",
@@ -38,6 +40,9 @@
3840
"start:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker tsx src/bin/go-worker.ts",
3941
"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",
4042
"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",
43+
"start:nuget-worker": "CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker tsx src/bin/nuget-worker.ts",
44+
"dev:nuget-worker": "CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9242 src/bin/nuget-worker.ts",
45+
"dev:nuget-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9242 src/bin/nuget-worker.ts",
4146
"backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts",
4247
"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",
4348
"backfill:stewardship": "SERVICE=stewardship-backfill tsx src/bin/stewardship-backfill.ts",

services/apps/packages_worker/src/activities.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,3 +25,4 @@ export {
2525
cargoCleanup,
2626
} from './cargo/activities'
2727
export { enrichGoVersionsBatch, enrichGoStatusBatch } from './go/activities'
28+
export { processNuGetBatch } from './nuget/activities'
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
import { scheduleNuGetIngestion } from '../nuget/schedule'
2+
import { svc } from '../service'
3+
4+
setImmediate(async () => {
5+
await svc.init()
6+
await scheduleNuGetIngestion()
7+
await svc.start()
8+
})

services/apps/packages_worker/src/config.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,15 @@ export function getGoConfig() {
7070
}
7171
}
7272

73+
export function getNuGetConfig() {
74+
return {
75+
batchSize: parseInt(process.env.NUGET_FETCHER_BATCH_SIZE ?? '1000', 10),
76+
concurrency: parseInt(process.env.NUGET_FETCHER_CONCURRENCY ?? '20', 10),
77+
groupDelayMs: parseInt(process.env.NUGET_FETCHER_GROUP_DELAY_MS ?? '0', 10),
78+
isCritical: (process.env.NUGET_FETCHER_IS_CRITICAL ?? 'false') === 'true',
79+
}
80+
}
81+
7382
export function getDockerhubConfig() {
7483
return {
7584
hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'),
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import { getServiceChildLogger } from '@crowd/logging'
2+
3+
import { getNuGetConfig } from '../config'
4+
import { getPackagesDb } from '../db'
5+
6+
import { processBatch } from './runNuGetEnrichmentLoop'
7+
import { BatchResult } from './types'
8+
9+
const log = getServiceChildLogger('nuget-activity')
10+
11+
export async function processNuGetBatch(): Promise<BatchResult> {
12+
const config = getNuGetConfig()
13+
const qx = await getPackagesDb()
14+
15+
const today = new Date().toISOString().split('T')[0]
16+
const result = await processBatch(qx, config, today)
17+
log.info({ ...result }, 'NuGet batch complete')
18+
return result
19+
}
Lines changed: 146 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,146 @@
1+
import axios from 'axios'
2+
3+
import {
4+
NuGetFetchError,
5+
NuGetRegistrationIndex,
6+
NuGetRegistrationPage,
7+
NuGetSearchItem,
8+
} from './types'
9+
10+
const SERVICE_INDEX_URL = 'https://api.nuget.org/v3/index.json'
11+
12+
interface ServiceIndexResource {
13+
'@id': string
14+
'@type': string
15+
}
16+
17+
interface ServiceIndex {
18+
resources: ServiceIndexResource[]
19+
}
20+
21+
interface ResolvedEndpoints {
22+
searchBaseUrl: string
23+
registrationBaseUrl: string
24+
}
25+
26+
let cachedEndpoints: ResolvedEndpoints | null = null
27+
28+
async function resolveEndpoints(): Promise<ResolvedEndpoints> {
29+
if (cachedEndpoints) return cachedEndpoints
30+
31+
const resp = await axios.get<ServiceIndex>(SERVICE_INDEX_URL, {
32+
timeout: 10000,
33+
})
34+
35+
const resources = resp.data.resources
36+
37+
const findFirst = (...types: string[]): string | undefined => {
38+
for (const t of types) {
39+
const found = resources.find((r) => r['@type'] === t)
40+
if (found) return found['@id']
41+
}
42+
return undefined
43+
}
44+
45+
const searchBaseUrl = findFirst('SearchQueryService/3.5.0', 'SearchQueryService')
46+
const registrationBaseUrl = findFirst(
47+
'RegistrationsBaseUrl/3.6.0',
48+
'RegistrationsBaseUrl/3.5.0',
49+
'RegistrationsBaseUrl',
50+
)
51+
52+
if (!searchBaseUrl || !registrationBaseUrl) {
53+
throw new Error('NuGet service index missing required endpoints')
54+
}
55+
56+
cachedEndpoints = { searchBaseUrl, registrationBaseUrl }
57+
return cachedEndpoints
58+
}
59+
60+
function classifyError(err: unknown): NuGetFetchError | null {
61+
if (!axios.isAxiosError(err)) return null
62+
const status = err.response?.status
63+
if (status === 404) return { kind: 'NOT_FOUND', status, message: err.message }
64+
if (status === 429) return { kind: 'RATE_LIMIT', status, message: err.message }
65+
return null
66+
}
67+
68+
export async function fetchSearch(packageId: string): Promise<NuGetSearchItem | NuGetFetchError> {
69+
const { searchBaseUrl } = await resolveEndpoints()
70+
const lowerPackageId = packageId.toLowerCase()
71+
72+
try {
73+
const resp = await axios.get<{ totalHits: number; data: NuGetSearchItem[] }>(searchBaseUrl, {
74+
params: {
75+
q: `packageid:${packageId}`,
76+
prerelease: 'true',
77+
semVerLevel: '2.0.0',
78+
take: 20,
79+
},
80+
headers: {
81+
'Accept-Encoding': 'gzip',
82+
},
83+
timeout: 15000,
84+
})
85+
86+
const match = resp.data.data.find((item) => item.id.toLowerCase() === lowerPackageId)
87+
if (!match) {
88+
return { kind: 'NOT_FOUND', message: `Package ${packageId} not found in search results` }
89+
}
90+
return match
91+
} catch (err) {
92+
const classified = classifyError(err)
93+
if (classified) return classified
94+
throw err
95+
}
96+
}
97+
98+
async function fetchRegistrationPage(
99+
pageId: string,
100+
maxAttempts = 2,
101+
): Promise<NuGetRegistrationPage> {
102+
for (let attempt = 1; ; attempt++) {
103+
try {
104+
const resp = await axios.get<NuGetRegistrationPage>(pageId, {
105+
headers: { 'Accept-Encoding': 'gzip' },
106+
timeout: 15000,
107+
})
108+
return resp.data
109+
} catch (err) {
110+
if (attempt >= maxAttempts) throw err
111+
}
112+
}
113+
}
114+
115+
export async function fetchRegistration(
116+
packageId: string,
117+
): Promise<NuGetRegistrationIndex | NuGetFetchError> {
118+
const { registrationBaseUrl } = await resolveEndpoints()
119+
const lowerId = packageId.toLowerCase()
120+
121+
try {
122+
const resp = await axios.get<NuGetRegistrationIndex>(
123+
`${registrationBaseUrl}${lowerId}/index.json`,
124+
{
125+
headers: { 'Accept-Encoding': 'gzip' },
126+
timeout: 15000,
127+
},
128+
)
129+
130+
const index = resp.data
131+
132+
for (let i = 0; i < index.items.length; i++) {
133+
const page = index.items[i]
134+
if (!page.items) {
135+
const fullPage = await fetchRegistrationPage(page['@id'])
136+
index.items[i] = { ...page, items: fullPage.items ?? [] }
137+
}
138+
}
139+
140+
return index
141+
} catch (err) {
142+
const classified = classifyError(err)
143+
if (classified) return classified
144+
throw err
145+
}
146+
}

0 commit comments

Comments
 (0)