Skip to content

Commit 97d617b

Browse files
authored
feat: implement rubygems registry [CM-1241] (#4322)
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent 54e9449 commit 97d617b

22 files changed

Lines changed: 1152 additions & 7 deletions

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 pypi-worker maven-worker osv-worker dockerhub-sync cargo-worker go-worker nuget-worker security-contacts-worker"
4+
SERVICES="github-repos-enricher bq-dataset-ingest npm-worker pypi-worker maven-worker osv-worker dockerhub-sync cargo-worker go-worker nuget-worker security-contacts-worker rubygems-worker"
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: rubygems-worker
7+
SHELL: /bin/sh
8+
SUPPRESS_NO_CONFIG_WARNING: 'true'
9+
CROWD_TEMPORAL_TASKQUEUE: rubygems-worker
10+
11+
services:
12+
rubygems-worker:
13+
build:
14+
context: ../../
15+
dockerfile: ./scripts/services/docker/Dockerfile.packages
16+
command: 'pnpm run start:rubygems-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+
rubygems-worker-dev:
30+
build:
31+
context: ../../
32+
dockerfile: ./scripts/services/docker/Dockerfile.packages
33+
command: 'pnpm run dev:rubygems-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: rubygems-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: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,9 @@
4949
"start:nuget-worker": "CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker tsx src/bin/nuget-worker.ts",
5050
"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",
5151
"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",
52+
"start:rubygems-worker": "CROWD_TEMPORAL_TASKQUEUE=rubygems-worker SERVICE=rubygems-worker tsx src/bin/rubygems-worker.ts",
53+
"dev:rubygems-worker": "CROWD_TEMPORAL_TASKQUEUE=rubygems-worker SERVICE=rubygems-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9244 src/bin/rubygems-worker.ts",
54+
"dev:rubygems-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=rubygems-worker SERVICE=rubygems-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9244 src/bin/rubygems-worker.ts",
5255
"backfill:maven": "SERVICE=maven tsx src/bin/maven-backfill.ts",
5356
"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",
5457
"import:maven-maintainers": "SERVICE=maven tsx src/maven/scripts/importMaintainersFromCsv.ts",

services/apps/packages_worker/src/activities.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,11 @@ export {
3232
} from './pypi/activities'
3333
export { getCriticalPypiCount } from './pypi/downloads/getCriticalPypiCount'
3434
export { processNuGetBatch } from './nuget/activities'
35+
export {
36+
processRubyGemsCoreBatch,
37+
processRubyGemsCriticalBatch,
38+
processRubyGemsDependentsBatch,
39+
} from './rubygems/activities'
3540
export {
3641
processSecurityContactsBatch,
3742
ingestSecurityContactsForPurlActivity,
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import {
2+
scheduleRubyGemsCriticalIngestion,
3+
scheduleRubyGemsDependentsIngestion,
4+
scheduleRubyGemsIngestion,
5+
} from '../rubygems/schedule'
6+
import { svc } from '../service'
7+
8+
setImmediate(async () => {
9+
await svc.init()
10+
await scheduleRubyGemsIngestion()
11+
await scheduleRubyGemsCriticalIngestion()
12+
await scheduleRubyGemsDependentsIngestion()
13+
await svc.start()
14+
})

services/apps/packages_worker/src/config.ts

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,27 @@ export function getNuGetConfig() {
8686
}
8787
}
8888

89+
export function getRubyGemsConfig() {
90+
return {
91+
batchSize: parseInt(process.env.RUBYGEMS_FETCHER_BATCH_SIZE ?? '10000', 10),
92+
concurrency: parseInt(process.env.RUBYGEMS_FETCHER_CONCURRENCY ?? '8', 10),
93+
}
94+
}
95+
96+
export function getRubyGemsCriticalConfig() {
97+
return {
98+
batchSize: parseInt(process.env.RUBYGEMS_CRITICAL_FETCHER_BATCH_SIZE ?? '5000', 10),
99+
concurrency: parseInt(process.env.RUBYGEMS_CRITICAL_FETCHER_CONCURRENCY ?? '4', 10),
100+
}
101+
}
102+
103+
export function getRubyGemsDependentsConfig() {
104+
return {
105+
batchSize: parseInt(process.env.RUBYGEMS_DEPENDENTS_BATCH_SIZE ?? '10000', 10),
106+
concurrency: parseInt(process.env.RUBYGEMS_DEPENDENTS_CONCURRENCY ?? '50', 10),
107+
}
108+
}
109+
89110
export function getDockerhubConfig() {
90111
return {
91112
hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'),
Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
import { getServiceChildLogger } from '@crowd/logging'
2+
3+
import {
4+
getRubyGemsConfig,
5+
getRubyGemsCriticalConfig,
6+
getRubyGemsDependentsConfig,
7+
} from '../config'
8+
import { getPackagesDb } from '../db'
9+
10+
import { processBatch as processCoreBatch } from './runRubyGemsCoreLoop'
11+
import { processBatch as processCriticalBatch } from './runRubyGemsCriticalLoop'
12+
import {
13+
DependentsBatchResult,
14+
processBatch as processDependentsBatch,
15+
} from './runRubyGemsDependentsLoop'
16+
import { BatchResult } from './types'
17+
18+
const log = getServiceChildLogger('rubygems-activity')
19+
20+
export async function processRubyGemsCoreBatch(): Promise<BatchResult> {
21+
const config = getRubyGemsConfig()
22+
const qx = await getPackagesDb()
23+
const today = new Date().toISOString().split('T')[0]
24+
const result = await processCoreBatch(qx, config, today)
25+
log.info({ ...result }, 'RubyGems core batch complete')
26+
return result
27+
}
28+
29+
export async function processRubyGemsCriticalBatch(
30+
afterId = '0',
31+
): Promise<BatchResult & { lastId: string | null }> {
32+
const config = getRubyGemsCriticalConfig()
33+
const qx = await getPackagesDb()
34+
const result = await processCriticalBatch(qx, config, afterId)
35+
log.info({ ...result }, 'RubyGems critical batch complete')
36+
return result
37+
}
38+
39+
export async function processRubyGemsDependentsBatch(
40+
afterId = '0',
41+
): Promise<DependentsBatchResult> {
42+
const config = getRubyGemsDependentsConfig()
43+
const qx = await getPackagesDb()
44+
const result = await processDependentsBatch(qx, config, afterId)
45+
log.info({ ...result }, 'RubyGems dependents batch complete')
46+
return result
47+
}
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
import axios from 'axios'
2+
3+
import { acquireRubyGemsSlot, parseRetryAfterMs } from './rateLimiter'
4+
import {
5+
RubyGemsFetchResult,
6+
RubyGemsGemResponse,
7+
RubyGemsOwner,
8+
RubyGemsVersionItem,
9+
} from './types'
10+
11+
const MAX_RATE_LIMIT_RETRIES = 5
12+
13+
async function rubyGemsGet<T>(url: string): Promise<RubyGemsFetchResult<T>> {
14+
for (let attempt = 0; ; attempt++) {
15+
await acquireRubyGemsSlot()
16+
try {
17+
const resp = await axios.get<T>(url, { timeout: 15000 })
18+
return resp.data
19+
} catch (err) {
20+
if (axios.isAxiosError(err)) {
21+
const status = err.response?.status
22+
if (status === 404) return { kind: 'NOT_FOUND', status, message: err.message }
23+
if (status === 429) {
24+
if (attempt >= MAX_RATE_LIMIT_RETRIES) {
25+
return { kind: 'RATE_LIMIT', status, message: err.message }
26+
}
27+
await new Promise((r) =>
28+
setTimeout(r, parseRetryAfterMs(err.response?.headers['retry-after'])),
29+
)
30+
continue
31+
}
32+
}
33+
throw err
34+
}
35+
}
36+
}
37+
38+
export function fetchGem(name: string): Promise<RubyGemsFetchResult<RubyGemsGemResponse>> {
39+
return rubyGemsGet<RubyGemsGemResponse>(
40+
`https://rubygems.org/api/v1/gems/${encodeURIComponent(name)}.json`,
41+
)
42+
}
43+
44+
export function fetchVersions(name: string): Promise<RubyGemsFetchResult<RubyGemsVersionItem[]>> {
45+
return rubyGemsGet<RubyGemsVersionItem[]>(
46+
`https://rubygems.org/api/v1/versions/${encodeURIComponent(name)}.json`,
47+
)
48+
}
49+
50+
export function fetchOwners(name: string): Promise<RubyGemsFetchResult<RubyGemsOwner[]>> {
51+
return rubyGemsGet<RubyGemsOwner[]>(
52+
`https://rubygems.org/api/v1/gems/${encodeURIComponent(name)}/owners.json`,
53+
)
54+
}
55+
56+
export function fetchReverseDependencies(name: string): Promise<RubyGemsFetchResult<string[]>> {
57+
return rubyGemsGet<string[]>(
58+
`https://rubygems.org/api/v1/gems/${encodeURIComponent(name)}/reverse_dependencies.json`,
59+
)
60+
}
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
import {
2+
NormalizedRubyGemsOwner,
3+
NormalizedRubyGemsPackage,
4+
NormalizedRubyGemsVersion,
5+
RubyGemsGemResponse,
6+
RubyGemsOwner,
7+
RubyGemsVersionItem,
8+
} from './types'
9+
10+
function nonEmpty(value: string | null | undefined): string | null {
11+
if (!value) return null
12+
const trimmed = value.trim()
13+
return trimmed === '' ? null : trimmed
14+
}
15+
16+
export function normalizeRubyGemsPackage(doc: RubyGemsGemResponse): NormalizedRubyGemsPackage {
17+
const licenses = doc.licenses && doc.licenses.length > 0 ? doc.licenses : null
18+
return {
19+
description: nonEmpty(doc.info),
20+
homepage: nonEmpty(doc.homepage_uri),
21+
declaredRepositoryUrl: nonEmpty(doc.source_code_uri),
22+
licenses,
23+
licensesRaw: licenses ? licenses.join(', ') : null,
24+
latestVersion: nonEmpty(doc.version),
25+
totalDownloads: doc.downloads ?? 0,
26+
}
27+
}
28+
29+
function parseCreatedAt(value: string | undefined): Date | null {
30+
if (!value) return null
31+
const date = new Date(value)
32+
return isNaN(date.getTime()) ? null : date
33+
}
34+
35+
export function normalizeRubyGemsVersions(
36+
items: RubyGemsVersionItem[],
37+
): NormalizedRubyGemsVersion[] {
38+
return items.map((item) => ({
39+
number: item.number,
40+
publishedAt: parseCreatedAt(item.created_at),
41+
isPrerelease: item.prerelease ?? false,
42+
licenses: item.licenses && item.licenses.length > 0 ? item.licenses : null,
43+
}))
44+
}
45+
46+
export function pickLatestRubyGemsVersion(
47+
versions: NormalizedRubyGemsVersion[],
48+
): NormalizedRubyGemsVersion | null {
49+
if (versions.length === 0) return null
50+
const stable = versions.find((v) => !v.isPrerelease)
51+
return stable ?? versions[0]
52+
}
53+
54+
export function normalizeRubyGemsOwners(owners: RubyGemsOwner[]): NormalizedRubyGemsOwner[] {
55+
return owners
56+
.filter((o): o is RubyGemsOwner & { handle: string } => !!o.handle && o.handle.trim() !== '')
57+
.map((o) => ({ username: o.handle, email: nonEmpty(o.email) }))
58+
}
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
const MAX_RPS = Math.max(1, parseInt(process.env.RUBYGEMS_MAX_RPS ?? '10', 10))
2+
const INTERVAL_MS = 1000 / MAX_RPS
3+
4+
let nextSlot = 0
5+
6+
export async function acquireRubyGemsSlot(): Promise<void> {
7+
const now = Date.now()
8+
const slot = Math.max(now, nextSlot)
9+
nextSlot = slot + INTERVAL_MS
10+
const wait = slot - now
11+
if (wait > 0) await new Promise((r) => setTimeout(r, wait))
12+
}
13+
14+
export function parseRetryAfterMs(header: unknown): number {
15+
const FALLBACK_MS = 1000
16+
if (typeof header !== 'string') return FALLBACK_MS
17+
const seconds = Number(header)
18+
if (Number.isFinite(seconds)) return Math.max(0, seconds * 1000)
19+
const date = new Date(header)
20+
if (!Number.isNaN(date.getTime())) return Math.max(0, date.getTime() - Date.now())
21+
return FALLBACK_MS
22+
}

0 commit comments

Comments
 (0)