diff --git a/services/apps/packages_worker/src/activities.ts b/services/apps/packages_worker/src/activities.ts index 7375896cd9..0e3b4a5738 100644 --- a/services/apps/packages_worker/src/activities.ts +++ b/services/apps/packages_worker/src/activities.ts @@ -33,11 +33,7 @@ export { } from './pypi/activities' export { getCriticalPypiCount } from './pypi/downloads/getCriticalPypiCount' export { processNuGetBatch } from './nuget/activities' -export { - processRubyGemsCoreBatch, - processRubyGemsCriticalBatch, - processRubyGemsDependentsBatch, -} from './rubygems/activities' +export { processRubyGemsCoreBatch, processRubyGemsCriticalBatch } from './rubygems/activities' export { processSecurityContactsBatch, ingestSecurityContactsForPurlActivity, diff --git a/services/apps/packages_worker/src/bin/rubygems-worker.ts b/services/apps/packages_worker/src/bin/rubygems-worker.ts index f0d8332dbe..cd5415b53c 100644 --- a/services/apps/packages_worker/src/bin/rubygems-worker.ts +++ b/services/apps/packages_worker/src/bin/rubygems-worker.ts @@ -1,14 +1,9 @@ -import { - scheduleRubyGemsCriticalIngestion, - scheduleRubyGemsDependentsIngestion, - scheduleRubyGemsIngestion, -} from '../rubygems/schedule' +import { scheduleRubyGemsCriticalIngestion, scheduleRubyGemsIngestion } from '../rubygems/schedule' import { svc } from '../service' setImmediate(async () => { await svc.init() await scheduleRubyGemsIngestion() await scheduleRubyGemsCriticalIngestion() - await scheduleRubyGemsDependentsIngestion() await svc.start() }) diff --git a/services/apps/packages_worker/src/config.ts b/services/apps/packages_worker/src/config.ts index cbb9826ebb..6daaab0f10 100644 --- a/services/apps/packages_worker/src/config.ts +++ b/services/apps/packages_worker/src/config.ts @@ -100,13 +100,6 @@ export function getRubyGemsCriticalConfig() { } } -export function getRubyGemsDependentsConfig() { - return { - batchSize: parseInt(process.env.RUBYGEMS_DEPENDENTS_BATCH_SIZE ?? '10000', 10), - concurrency: parseInt(process.env.RUBYGEMS_DEPENDENTS_CONCURRENCY ?? '50', 10), - } -} - export function getDockerhubConfig() { return { hubBaseUrl: requireEnv('DOCKERHUB_API_BASE_URL'), diff --git a/services/apps/packages_worker/src/rubygems/activities.ts b/services/apps/packages_worker/src/rubygems/activities.ts index b124e45210..5f40ca2d58 100644 --- a/services/apps/packages_worker/src/rubygems/activities.ts +++ b/services/apps/packages_worker/src/rubygems/activities.ts @@ -1,18 +1,10 @@ import { getServiceChildLogger } from '@crowd/logging' -import { - getRubyGemsConfig, - getRubyGemsCriticalConfig, - getRubyGemsDependentsConfig, -} from '../config' +import { getRubyGemsConfig, getRubyGemsCriticalConfig } from '../config' import { getPackagesDb } from '../db' import { processBatch as processCoreBatch } from './runRubyGemsCoreLoop' import { processBatch as processCriticalBatch } from './runRubyGemsCriticalLoop' -import { - DependentsBatchResult, - processBatch as processDependentsBatch, -} from './runRubyGemsDependentsLoop' import { BatchResult } from './types' const log = getServiceChildLogger('rubygems-activity') @@ -35,13 +27,3 @@ export async function processRubyGemsCriticalBatch( log.info({ ...result }, 'RubyGems critical batch complete') return result } - -export async function processRubyGemsDependentsBatch( - afterId = '0', -): Promise { - const config = getRubyGemsDependentsConfig() - const qx = await getPackagesDb() - const result = await processDependentsBatch(qx, config, afterId) - log.info({ ...result }, 'RubyGems dependents batch complete') - return result -} diff --git a/services/apps/packages_worker/src/rubygems/client.ts b/services/apps/packages_worker/src/rubygems/client.ts index 3eb2adb5c2..4d48ba3266 100644 --- a/services/apps/packages_worker/src/rubygems/client.ts +++ b/services/apps/packages_worker/src/rubygems/client.ts @@ -52,9 +52,3 @@ export function fetchOwners(name: string): Promise> { - return rubyGemsGet( - `https://rubygems.org/api/v1/gems/${encodeURIComponent(name)}/reverse_dependencies.json`, - ) -} diff --git a/services/apps/packages_worker/src/rubygems/runRubyGemsDependentsLoop.ts b/services/apps/packages_worker/src/rubygems/runRubyGemsDependentsLoop.ts deleted file mode 100644 index 95a2888c86..0000000000 --- a/services/apps/packages_worker/src/rubygems/runRubyGemsDependentsLoop.ts +++ /dev/null @@ -1,87 +0,0 @@ -import { - QueryExecutor, - RubyGemsPackageForDependents, - listRubyGemsPackagesForDependents, - updateRubyGemsDependentCount, -} from '@crowd/data-access-layer' -import { getServiceChildLogger } from '@crowd/logging' - -import { fetchReverseDependencies } from './client' -import { isRubyGemsFetchError } from './types' - -const log = getServiceChildLogger('rubygems-dependents') - -export type RubyGemsDependentsConfig = { - batchSize: number - concurrency: number -} - -export type DependentsBatchResult = { - processed: number - notFound: number - error: number - lastId: string | null -} - -type PackageStatus = 'processed' | 'notFound' | 'error' - -async function processPackage( - qx: QueryExecutor, - pkg: RubyGemsPackageForDependents, -): Promise { - const result = await fetchReverseDependencies(pkg.name) - - if (isRubyGemsFetchError(result)) { - if (result.kind === 'NOT_FOUND') { - await updateRubyGemsDependentCount(qx, pkg.id, 0) - return 'notFound' - } - log.warn({ name: pkg.name }, 'Rate limited fetching reverse dependencies — will retry') - return 'error' - } - - await updateRubyGemsDependentCount(qx, pkg.id, result.length) - return 'processed' -} - -export async function processBatch( - qx: QueryExecutor, - config: RubyGemsDependentsConfig, - afterId: string, -): Promise { - const packages = await listRubyGemsPackagesForDependents(qx, { - limit: config.batchSize, - afterId, - }) - - if (packages.length === 0) return { processed: 0, notFound: 0, error: 0, lastId: null } - - log.info({ count: packages.length, afterId }, 'Dependents batch started') - - const counts: DependentsBatchResult = { processed: 0, notFound: 0, error: 0, lastId: null } - - for (let batchStart = 0; batchStart < packages.length; batchStart += config.concurrency) { - const group = packages.slice(batchStart, batchStart + config.concurrency) - - await Promise.all( - group.map(async (pkg) => { - try { - const status = await processPackage(qx, pkg) - counts[status]++ - } catch (err) { - const message = err instanceof Error ? err.message : String(err) - log.error({ name: pkg.name, error: message }, 'Unexpected error fetching dependents') - counts.error++ - } - }), - ) - - const done = batchStart + group.length - if (done % 1000 === 0 || done === packages.length) { - log.info({ done, total: packages.length, ...counts }, 'Dependents progress') - } - } - - counts.lastId = packages[packages.length - 1].id - return counts -} diff --git a/services/apps/packages_worker/src/rubygems/schedule.ts b/services/apps/packages_worker/src/rubygems/schedule.ts index 4567ad6bcd..1d06224a3c 100644 --- a/services/apps/packages_worker/src/rubygems/schedule.ts +++ b/services/apps/packages_worker/src/rubygems/schedule.ts @@ -1,11 +1,7 @@ import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client' import { svc } from '../service' -import { - ingestRubyGemsCriticalDetails, - ingestRubyGemsDependents, - ingestRubyGemsPackages, -} from '../workflows' +import { ingestRubyGemsCriticalDetails, ingestRubyGemsPackages } from '../workflows' export async function scheduleRubyGemsIngestion(): Promise { const { temporal } = svc @@ -80,40 +76,3 @@ export async function scheduleRubyGemsCriticalIngestion(): Promise { } } } - -export async function scheduleRubyGemsDependentsIngestion(): Promise { - const { temporal } = svc - if (!temporal) throw new Error('Temporal client not initialized') - - try { - await temporal.schedule.create({ - scheduleId: 'rubygems-dependents-ingest', - spec: { - cronExpressions: ['0 3 * * 0'], - }, - policies: { - overlap: ScheduleOverlapPolicy.SKIP, - catchupWindow: '1 hour', - }, - action: { - type: 'startWorkflow', - workflowType: ingestRubyGemsDependents, - workflowId: 'rubygems-weekly-dependents', - taskQueue: 'rubygems-worker', - workflowRunTimeout: '24 hours', - retry: { - initialInterval: '30 seconds', - backoffCoefficient: 2, - maximumAttempts: 5, - }, - args: [], - }, - }) - } catch (err) { - if (err instanceof ScheduleAlreadyRunning) { - svc.log.info('Schedule rubygems-dependents-ingest already exists, skipping creation.') - } else { - throw err - } - } -} diff --git a/services/apps/packages_worker/src/rubygems/workflows.ts b/services/apps/packages_worker/src/rubygems/workflows.ts index a6015f2c2a..d05db79369 100644 --- a/services/apps/packages_worker/src/rubygems/workflows.ts +++ b/services/apps/packages_worker/src/rubygems/workflows.ts @@ -24,12 +24,3 @@ export async function ingestRubyGemsCriticalDetails(afterId = '0'): Promise(result.lastId) } - -export async function ingestRubyGemsDependents(afterId = '0'): Promise { - const result = await acts.processRubyGemsDependentsBatch(afterId) - if (result.lastId === null) { - log.info('RubyGems dependents ingestion complete — no more work, exiting.', { ...result }) - return - } - await continueAsNew(result.lastId) -} diff --git a/services/apps/packages_worker/src/workflows/index.ts b/services/apps/packages_worker/src/workflows/index.ts index 465047c39c..a14aee8ed0 100644 --- a/services/apps/packages_worker/src/workflows/index.ts +++ b/services/apps/packages_worker/src/workflows/index.ts @@ -25,11 +25,7 @@ export { ingestPypiDownloadsDaily, } from '../pypi/downloads/ingestPypiDownloads' export { ingestNuGetPackages } from '../nuget/workflows' -export { - ingestRubyGemsCriticalDetails, - ingestRubyGemsDependents, - ingestRubyGemsPackages, -} from '../rubygems/workflows' +export { ingestRubyGemsCriticalDetails, ingestRubyGemsPackages } from '../rubygems/workflows' export { ingestSecurityContacts, ingestSecurityContactsForPurlWorkflow, diff --git a/services/libs/data-access-layer/src/osspckgs/rubygems.ts b/services/libs/data-access-layer/src/osspckgs/rubygems.ts index d72b65b8ee..831467d2f7 100644 --- a/services/libs/data-access-layer/src/osspckgs/rubygems.ts +++ b/services/libs/data-access-layer/src/osspckgs/rubygems.ts @@ -62,41 +62,3 @@ export async function listRubyGemsCriticalPackagesToSync( { limit, afterId }, ) } - -export type RubyGemsPackageForDependents = { - id: string - name: string -} - -export async function listRubyGemsPackagesForDependents( - qx: QueryExecutor, - options: { limit: number; afterId?: string }, -): Promise { - const { limit, afterId = '0' } = options - return qx.select( - ` - SELECT p.id, p.name - FROM packages p - WHERE p.ecosystem = 'rubygems' - AND p.id > $(afterId)::bigint - ORDER BY p.id ASC - LIMIT $(limit) - `, - { limit, afterId }, - ) -} - -export async function updateRubyGemsDependentCount( - qx: QueryExecutor, - packageId: string, - dependentCount: number, -): Promise { - await qx.result( - `UPDATE packages - SET dependent_count = $(dependentCount), - last_synced_at = NOW() - WHERE id = $(packageId)::bigint - AND dependent_count IS DISTINCT FROM $(dependentCount)`, - { packageId, dependentCount }, - ) -}