|
1 | 1 | import { Context } from '@temporalio/activity' |
2 | 2 |
|
| 3 | +import { createIngestJob, markJobStatus } from '@crowd/data-access-layer' |
3 | 4 | import { getServiceChildLogger } from '@crowd/logging' |
4 | 5 |
|
5 | 6 | import { getPackagesDb } from '../db' |
@@ -63,3 +64,34 @@ export async function criticalityComputePageRank( |
63 | 64 |
|
64 | 65 | return { ecosystem, nodeCount: graph.N, edgeCount, iterations, durationMs: Date.now() - start } |
65 | 66 | } |
| 67 | + |
| 68 | +export async function rankPackages(): Promise<{ scoredRows: number; rankedRows: number }> { |
| 69 | + const qx = await getPackagesDb() |
| 70 | + |
| 71 | + // On retry, a pending row from the prior attempt may already exist — reuse it. |
| 72 | + // Do NOT reuse a done row: it belongs to a previous bootstrap run and ranking must re-execute. |
| 73 | + const existing = await qx.selectOneOrNone( |
| 74 | + `SELECT id FROM osspckgs_ingest_jobs |
| 75 | + WHERE job_kind = 'ranking' AND status = 'pending' |
| 76 | + ORDER BY id DESC LIMIT 1`, |
| 77 | + ) |
| 78 | + |
| 79 | + const jobId = existing?.id ?? (await createIngestJob(qx, 'ranking', 'ranking', null)) |
| 80 | + try { |
| 81 | + const [result] = await qx.select(`SELECT * FROM rank_packages()`) |
| 82 | + const scoredRows = Number(result.scored_rows ?? 0) |
| 83 | + const rankedRows = Number(result.ranked_rows ?? 0) |
| 84 | + await markJobStatus(qx, jobId, 'done', { |
| 85 | + rowCountPg: scoredRows, |
| 86 | + tableRowCounts: { scored: scoredRows, ranked: rankedRows }, |
| 87 | + finishedAt: new Date(), |
| 88 | + }) |
| 89 | + return { scoredRows, rankedRows } |
| 90 | + } catch (err) { |
| 91 | + await markJobStatus(qx, jobId, 'failed', { |
| 92 | + errorMessage: (err as Error).message, |
| 93 | + finishedAt: new Date(), |
| 94 | + }) |
| 95 | + throw err |
| 96 | + } |
| 97 | +} |
0 commit comments