Skip to content

Commit 5b609ca

Browse files
committed
fix: anchor packagist metadata due-selection to a fixed cutoff
Signed-off-by: anilb <epipav@gmail.com>
1 parent 3e42dbe commit 5b609ca

5 files changed

Lines changed: 64 additions & 18 deletions

File tree

services/apps/packages_worker/src/packagist/__tests__/dueSelection.test.ts

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -35,15 +35,20 @@ beforeEach(() => {
3535
})
3636

3737
// Metadata (merged enrichment) due-selection: keyset pagination, refresh window, and
38-
// the stored Last-Modified for If-Modified-Since replay on the p2 fetch.
38+
// the stored Last-Modified for If-Modified-Since replay on the p2 fetch. `cutoff` is a
39+
// pre-computed threshold (this run's fixed start time minus the refresh window), not a
40+
// live NOW() — a keyset scan only visits each purl once per drain, so due-selection
41+
// must stay anchored to one stable point in time for the whole run.
3942
describe('getPackagistMetadataDuePurls', () => {
40-
it('returns purls with their stored Last-Modified, scoped by refresh window', async () => {
43+
const CUTOFF = '2026-07-08T00:00:00.000Z'
44+
45+
it('returns purls with their stored Last-Modified, scoped by the fixed cutoff', async () => {
4146
qx.select.mockResolvedValue([
4247
{ purl: 'pkg:composer/a/x', metadata_last_modified: 'Tue, 30 Jun 2026 00:00:00 GMT' },
4348
{ purl: 'pkg:composer/b/y', metadata_last_modified: null },
4449
])
4550

46-
const candidates = await getPackagistMetadataDuePurls(asQx(qx), '', 50, 7, true)
51+
const candidates = await getPackagistMetadataDuePurls(asQx(qx), CUTOFF, '', 50, true)
4752

4853
expect(candidates).toEqual([
4954
{ purl: 'pkg:composer/a/x', metadataLastModified: 'Tue, 30 Jun 2026 00:00:00 GMT' },
@@ -54,11 +59,11 @@ describe('getPackagistMetadataDuePurls', () => {
5459
expect(sql).toMatch(/is_critical/)
5560
expect(sql).toMatch(/metadata_last_run_at/)
5661
expect(sql).toMatch(/ORDER BY\s+p\.purl/i)
57-
expect(params).toMatchObject({ batchSize: 50, refreshDays: 7, onlyCritical: true })
62+
expect(params).toMatchObject({ cutoff: CUTOFF, batchSize: 50, onlyCritical: true })
5863
})
5964

6065
it('selects across all packagist packages when onlyCritical is false (the all-packages default)', async () => {
61-
await getPackagistMetadataDuePurls(asQx(qx), '', 50, 7, false)
66+
await getPackagistMetadataDuePurls(asQx(qx), CUTOFF, '', 50, false)
6267

6368
const [sql, params] = qx.select.mock.calls[0]
6469
expect(sql).toMatch(/ecosystem = 'packagist'/)

services/apps/packages_worker/src/packagist/__tests__/ingest.test.ts

Lines changed: 29 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -414,39 +414,61 @@ describe('ingestPackagistItemsConcurrently', () => {
414414

415415
// All-packages metadata scope: the env default drives due-selection's onlyCritical arg.
416416
// Default is ALL packages (deps.dev has no Packagist data to fall back on); the env var
417-
// narrows back to the critical slice.
417+
// narrows back to the critical slice. The due-cutoff is derived once from the run's
418+
// fixed cutoff minus the refresh window, not a live NOW() re-evaluated per batch.
418419
describe('getPackagistMetadataBatch scope', () => {
419420
const mockMetadataDue = vi.mocked(getPackagistMetadataDuePurls)
421+
const RUN_CUTOFF = '2026-07-15T00:00:00.000Z'
422+
const DUE_CUTOFF = '2026-07-08T00:00:00.000Z' // RUN_CUTOFF minus the 7-day default
420423

421424
afterEach(() => {
422425
delete process.env.CROWD_PACKAGES_PACKAGIST_RUN_ONLY_FOR_CRITICAL
426+
delete process.env.CROWD_PACKAGES_PACKAGIST_METADATA_REFRESH_DAYS
427+
})
428+
429+
it('derives the due-cutoff from the fixed run cutoff, not a live clock', async () => {
430+
process.env.CROWD_PACKAGES_PACKAGIST_METADATA_REFRESH_DAYS = '3'
431+
mockMetadataDue.mockResolvedValue([])
432+
433+
await getPackagistMetadataBatch(RUN_CUTOFF, '', 50)
434+
435+
// 3 days before the run's own fixed cutoff — never "3 days before whenever this
436+
// particular batch happens to execute" — is what makes the keyset scan give every
437+
// purl exactly one refresh-eligibility check per drain.
438+
expect(mockMetadataDue).toHaveBeenCalledWith(
439+
expect.anything(),
440+
'2026-07-12T00:00:00.000Z',
441+
'',
442+
50,
443+
false,
444+
)
423445
})
424446

425447
it('selects across ALL packagist packages by default', async () => {
426448
mockMetadataDue.mockResolvedValue([])
427449

428-
await getPackagistMetadataBatch('', 50)
450+
await getPackagistMetadataBatch(RUN_CUTOFF, '', 50)
429451

430-
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), '', 50, 7, false)
452+
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), DUE_CUTOFF, '', 50, false)
431453
})
432454

433455
it('narrows to the critical slice when CROWD_PACKAGES_PACKAGIST_RUN_ONLY_FOR_CRITICAL=true', async () => {
434456
process.env.CROWD_PACKAGES_PACKAGIST_RUN_ONLY_FOR_CRITICAL = 'true'
435457
mockMetadataDue.mockResolvedValue([])
436458

437-
await getPackagistMetadataBatch('', 50)
459+
await getPackagistMetadataBatch(RUN_CUTOFF, '', 50)
438460

439-
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), '', 50, 7, true)
461+
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), DUE_CUTOFF, '', 50, true)
440462
})
441463

442464
it('keeps the all-packages default when the env value is unrecognized', async () => {
443465
// a typo must not silently flip the sweep to critical-only
444466
process.env.CROWD_PACKAGES_PACKAGIST_RUN_ONLY_FOR_CRITICAL = 'ture'
445467
mockMetadataDue.mockResolvedValue([])
446468

447-
await getPackagistMetadataBatch('', 50)
469+
await getPackagistMetadataBatch(RUN_CUTOFF, '', 50)
448470

449-
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), '', 50, 7, false)
471+
expect(mockMetadataDue).toHaveBeenCalledWith(expect.anything(), DUE_CUTOFF, '', 50, false)
450472
})
451473
})
452474

services/apps/packages_worker/src/packagist/activities.ts

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -323,16 +323,25 @@ export async function runPackagistPackageSeed(): Promise<{ discovered: number; i
323323
return { discovered: entries.length, invalid }
324324
}
325325

326+
// `cutoff` is the drain's own fixed start time (stable across every round via
327+
// continueAsNew), not a live NOW() — a keyset scan only visits each purl once per
328+
// drain, so re-deriving "7 days ago" fresh on every batch would silently skip a purl
329+
// that hasn't quite hit the refresh window yet when the cursor passes it, pushing its
330+
// effective cadence out toward two refresh cycles instead of one.
326331
export async function getPackagistMetadataBatch(
332+
cutoff: string,
327333
afterPurl: string,
328334
batchSize: number,
329335
): Promise<{ candidates: PackagistMetadataCandidate[]; nextCursor: string }> {
330336
const qx = await getPackagesDb()
337+
const dueCutoff = new Date(
338+
new Date(cutoff).getTime() - metadataRefreshDays() * 24 * 60 * 60 * 1000,
339+
).toISOString()
331340
const candidates = await getPackagistMetadataDuePurls(
332341
qx,
342+
dueCutoff,
333343
afterPurl,
334344
batchSize,
335-
metadataRefreshDays(),
336345
runOnlyForCritical(),
337346
)
338347
return {

services/apps/packages_worker/src/packagist/workflows.ts

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ const INGEST_BATCH = 50
2323
const ROUNDS_PER_RUN = 20
2424

2525
interface MetadataState {
26+
cutoff?: string
2627
cursor?: string
2728
}
2829

@@ -54,20 +55,29 @@ export async function seedPackagistPackages(): Promise<void> {
5455
}
5556
}
5657

58+
// The cutoff is fixed once per run (deterministic activity), same pattern as the
59+
// downloads-30d/daily lanes — a keyset scan only ever visits each purl once per drain,
60+
// so due-selection must be anchored to a stable point in time rather than a live NOW()
61+
// that would let a purl processed early in the run dodge this cycle's refresh window.
5762
export async function ingestPackagistMetadata(state: MetadataState = {}): Promise<void> {
63+
const cutoff = state.cutoff ?? (await acts.packagistCurrentTimestamp())
5864
let cursor = state.cursor || ''
5965
const stopAfterFirstPage = await acts.packagistStopAfterFirstPage()
6066

6167
for (let r = 0; r < ROUNDS_PER_RUN; r++) {
62-
const { candidates, nextCursor } = await acts.getPackagistMetadataBatch(cursor, INGEST_BATCH)
68+
const { candidates, nextCursor } = await acts.getPackagistMetadataBatch(
69+
cutoff,
70+
cursor,
71+
INGEST_BATCH,
72+
)
6373
if (candidates.length === 0) return
6474
await acts.ingestPackagistMetadataBatch(candidates)
6575
cursor = nextCursor
6676
if (stopAfterFirstPage) return
6777
if (candidates.length < INGEST_BATCH) return
6878
}
6979

70-
await continueAsNew<typeof ingestPackagistMetadata>({ cursor })
80+
await continueAsNew<typeof ingestPackagistMetadata>({ cutoff, cursor })
7181
}
7282

7383
// Monthly capture of the observed rolling 30d window for every packagist package.

services/libs/data-access-layer/src/packages/packagistPackageState.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -68,9 +68,9 @@ export async function markPackagistMetadataScanned(
6868

6969
export async function getPackagistMetadataDuePurls(
7070
qx: QueryExecutor,
71+
cutoff: string,
7172
afterPurl: string,
7273
batchSize: number,
73-
refreshDays: number,
7474
onlyCritical: boolean,
7575
): Promise<PackagistMetadataCandidate[]> {
7676
const rows: Array<{ purl: string; metadata_last_modified: string | null }> = await qx.select(
@@ -82,11 +82,11 @@ export async function getPackagistMetadataDuePurls(
8282
AND p.purl > $(afterPurl)
8383
AND (
8484
s.metadata_last_run_at IS NULL
85-
OR s.metadata_last_run_at < NOW() - ($(refreshDays) || ' days')::interval
85+
OR s.metadata_last_run_at < $(cutoff)::timestamptz
8686
)
8787
ORDER BY p.purl
8888
LIMIT $(batchSize)`,
89-
{ afterPurl, batchSize, refreshDays, onlyCritical },
89+
{ cutoff, afterPurl, batchSize, onlyCritical },
9090
)
9191
return rows.map((r) => ({
9292
purl: r.purl,

0 commit comments

Comments
 (0)