Skip to content

Commit 77e6b6c

Browse files
committed
fix: guard packagist watermark against re-processing
Signed-off-by: anilb <epipav@gmail.com>
1 parent 3a9d2bc commit 77e6b6c

4 files changed

Lines changed: 120 additions & 62 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ describe('markPackagistMetadataScanned', () => {
7272
asQx(qx),
7373
'pkg:composer/a/x',
7474
{ status: 'success', attempts: 1 },
75-
'Wed, 01 Jul 2026 00:00:00 GMT',
75+
{ metadataLastModified: 'Wed, 01 Jul 2026 00:00:00 GMT' },
7676
)
7777

7878
const [sql, params] = qx.result.mock.calls[0]

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

Lines changed: 22 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@ const mockSeedInsert = vi.mocked(insertPackagistPackages)
7575
const qx = {} as QueryExecutor
7676
const PURL = 'pkg:composer/monolog/monolog'
7777
const RUN_DATE = '2026-07-15'
78+
const SCHEDULED_AT = '2026-07-15T00:00:00.000Z'
7879

7980
const statsJson = {
8081
package: { name: 'monolog/monolog', downloads: { total: 1000, monthly: 300, daily: 10 } },
@@ -113,7 +114,7 @@ describe('ingestOnePackagistMetadata', () => {
113114
it('audits phase 1 immediately and phase 2 separately, then marks scanned with the fresh Last-Modified', async () => {
114115
happyMocks()
115116

116-
await ingestOnePackagistMetadata(qx, candidate)
117+
await ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)
117118

118119
expect(mockPersistInfo).toHaveBeenCalledWith(qx, PURL, expect.anything())
119120
expect(mockFetchP2).toHaveBeenCalledWith('monolog/monolog', 'Tue, 30 Jun 2026 00:00:00 GMT')
@@ -127,15 +128,15 @@ describe('ingestOnePackagistMetadata', () => {
127128
qx,
128129
PURL,
129130
expect.objectContaining({ status: 'success', attempts: 1 }),
130-
'Wed, 01 Jul 2026 00:00:00 GMT',
131+
{ metadataLastModified: 'Wed, 01 Jul 2026 00:00:00 GMT' },
131132
)
132133
})
133134

134135
it('records a p2 NOT_MODIFIED as success: info persisted, versions skipped', async () => {
135136
happyMocks()
136137
mockFetchP2.mockResolvedValue({ kind: 'NOT_MODIFIED' } as never)
137138

138-
await ingestOnePackagistMetadata(qx, candidate)
139+
await ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)
139140

140141
expect(mockPersistInfo).toHaveBeenCalled()
141142
expect(mockPersistMetadata).not.toHaveBeenCalled()
@@ -144,7 +145,7 @@ describe('ingestOnePackagistMetadata', () => {
144145
qx,
145146
PURL,
146147
expect.objectContaining({ status: 'success' }),
147-
null,
148+
{ metadataLastModified: null },
148149
)
149150
})
150151

@@ -156,7 +157,7 @@ describe('ingestOnePackagistMetadata', () => {
156157
message: 'not found',
157158
} as never)
158159

159-
const p = ingestOnePackagistMetadata(qx, candidate)
160+
const p = ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)
160161
await vi.runAllTimersAsync()
161162
await p
162163

@@ -167,6 +168,7 @@ describe('ingestOnePackagistMetadata', () => {
167168
qx,
168169
PURL,
169170
expect.objectContaining({ status: 'error', errorKind: 'NOT_FOUND', httpStatus: 404 }),
171+
{ notBefore: SCHEDULED_AT },
170172
)
171173
})
172174

@@ -177,7 +179,7 @@ describe('ingestOnePackagistMetadata', () => {
177179
message: 'HTTP 500',
178180
} as never)
179181

180-
await expect(ingestOnePackagistMetadata(qx, candidate)).rejects.toThrow()
182+
await expect(ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)).rejects.toThrow()
181183
expect(mockMarkMetadata).not.toHaveBeenCalled()
182184
expect(mockFetchStats).toHaveBeenCalledTimes(1)
183185
})
@@ -192,7 +194,7 @@ describe('ingestOnePackagistMetadata', () => {
192194
message: 'not found',
193195
} as never)
194196

195-
const p = ingestOnePackagistMetadata(qx, candidate)
197+
const p = ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)
196198
await vi.runAllTimersAsync()
197199
await p
198200

@@ -203,8 +205,7 @@ describe('ingestOnePackagistMetadata', () => {
203205
qx,
204206
PURL,
205207
expect.objectContaining({ status: 'error', errorKind: 'NOT_FOUND' }),
206-
undefined,
207-
false,
208+
{ bumpLastRunAt: false, notBefore: SCHEDULED_AT },
208209
)
209210
})
210211

@@ -219,7 +220,7 @@ describe('ingestOnePackagistMetadata', () => {
219220
message: 'not found',
220221
} as never)
221222

222-
const p = ingestOnePackagistMetadata(qx, candidate)
223+
const p = ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)
223224
await vi.runAllTimersAsync()
224225
await p
225226

@@ -228,8 +229,7 @@ describe('ingestOnePackagistMetadata', () => {
228229
qx,
229230
PURL,
230231
expect.objectContaining({ status: 'error', errorKind: 'NOT_FOUND' }),
231-
undefined,
232-
false,
232+
{ bumpLastRunAt: false, notBefore: SCHEDULED_AT },
233233
)
234234
})
235235

@@ -238,7 +238,7 @@ describe('ingestOnePackagistMetadata', () => {
238238
mockPersistInfo.mockResolvedValue({ found: true, changedFields: ['packages.description'] })
239239
mockFetchP2.mockResolvedValue({ kind: 'TRANSIENT', message: 'HTTP 502' } as never)
240240

241-
await expect(ingestOnePackagistMetadata(qx, candidate)).rejects.toThrow()
241+
await expect(ingestOnePackagistMetadata(qx, candidate, SCHEDULED_AT)).rejects.toThrow()
242242
expect(mockMarkMetadata).not.toHaveBeenCalled()
243243
// phase-1 writes are already committed when the throw happens — a retry re-runs
244244
// phase 1 idempotently and reports no changes, so this audit event can't be deferred
@@ -252,7 +252,7 @@ describe('ingestOnePackagist30dWindow', () => {
252252
mockFetchStats.mockResolvedValue(statsJson as never)
253253
mockPersist30d.mockResolvedValue(['downloads_last_30d.count', 'packages.downloads_last_30d'])
254254

255-
await ingestOnePackagist30dWindow(qx, PURL, RUN_DATE)
255+
await ingestOnePackagist30dWindow(qx, PURL, RUN_DATE, SCHEDULED_AT)
256256

257257
expect(mockPersist30d).toHaveBeenCalledWith(qx, PURL, 300, RUN_DATE)
258258
expect(mockAudit).toHaveBeenCalledWith(qx, 'packagist', PURL, [
@@ -274,7 +274,7 @@ describe('ingestOnePackagist30dWindow', () => {
274274
message: 'not found',
275275
} as never)
276276

277-
const p = ingestOnePackagist30dWindow(qx, PURL, RUN_DATE)
277+
const p = ingestOnePackagist30dWindow(qx, PURL, RUN_DATE, SCHEDULED_AT)
278278
await vi.runAllTimersAsync()
279279
await p
280280

@@ -283,13 +283,14 @@ describe('ingestOnePackagist30dWindow', () => {
283283
qx,
284284
PURL,
285285
expect.objectContaining({ status: 'error', errorKind: 'NOT_FOUND' }),
286+
SCHEDULED_AT,
286287
)
287288
})
288289

289290
it('throws on a transient result without marking processed', async () => {
290291
mockFetchStats.mockResolvedValue({ kind: 'TRANSIENT', message: 'HTTP 503' } as never)
291292

292-
await expect(ingestOnePackagist30dWindow(qx, PURL, RUN_DATE)).rejects.toThrow()
293+
await expect(ingestOnePackagist30dWindow(qx, PURL, RUN_DATE, SCHEDULED_AT)).rejects.toThrow()
293294
expect(mockMark30d).not.toHaveBeenCalled()
294295
})
295296
})
@@ -302,7 +303,7 @@ describe('ingestOnePackagistDailyDownload', () => {
302303
mockFetchStats.mockResolvedValue(statsJson as never)
303304
mockDaily.mockResolvedValue(['downloads_daily.date', 'downloads_daily.count'])
304305

305-
await ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE)
306+
await ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE, SCHEDULED_AT)
306307

307308
expect(mockDaily).toHaveBeenCalledWith(qx, '7', [{ day: RUN_DATE, downloads: 10 }])
308309
expect(mockAudit).toHaveBeenCalledWith(qx, 'packagist', PURL, [
@@ -319,7 +320,7 @@ describe('ingestOnePackagistDailyDownload', () => {
319320
it('marks success without inserting or auditing when the registry reports no daily count', async () => {
320321
mockFetchStats.mockResolvedValue({ package: { name: 'monolog/monolog' } } as never)
321322

322-
await ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE)
323+
await ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE, SCHEDULED_AT)
323324

324325
expect(mockDaily).not.toHaveBeenCalled()
325326
expect(mockAudit).not.toHaveBeenCalled()
@@ -333,7 +334,9 @@ describe('ingestOnePackagistDailyDownload', () => {
333334
it('throws on a transient result without marking processed', async () => {
334335
mockFetchStats.mockResolvedValue({ kind: 'TRANSIENT', message: 'HTTP 503' } as never)
335336

336-
await expect(ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE)).rejects.toThrow()
337+
await expect(
338+
ingestOnePackagistDailyDownload(qx, candidate, RUN_DATE, SCHEDULED_AT),
339+
).rejects.toThrow()
337340
expect(mockMarkDaily).not.toHaveBeenCalled()
338341
})
339342
})

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

Lines changed: 54 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,9 @@ function giveUpResult(error: FetchError, attempts: number): PackagistRunResult {
119119
export async function ingestOnePackagistMetadata(
120120
qx: QueryExecutor,
121121
candidate: PackagistMetadataCandidate,
122+
// This batch activity's own scheduledTimestampMs (stable across Temporal retries),
123+
// passed through to every give-up write below — see MarkMetadataScannedOptions.notBefore.
124+
scheduledAt: string,
122125
): Promise<void> {
123126
const name = packagistNameFromPurl(candidate.purl)
124127

@@ -132,7 +135,14 @@ export async function ingestOnePackagistMetadata(
132135
{ purl: candidate.purl, statusCode: info.error.statusCode, kind: info.error.kind },
133136
'packagist package info 4xx/malformed after fast retries — marking scanned and skipping',
134137
)
135-
await markPackagistMetadataScanned(qx, candidate.purl, giveUpResult(info.error, info.attempts))
138+
await markPackagistMetadataScanned(
139+
qx,
140+
candidate.purl,
141+
giveUpResult(info.error, info.attempts),
142+
{
143+
notBefore: scheduledAt,
144+
},
145+
)
136146
return
137147
}
138148

@@ -154,13 +164,10 @@ export async function ingestOnePackagistMetadata(
154164
// Phase 1 already succeeded — only p2 (versions/deps) failed to refresh. Don't push
155165
// metadata_last_run_at forward, or due-selection wrongly treats this package as
156166
// "recently scanned" and skips it for the full refresh window despite stale p2 data.
157-
await markPackagistMetadataScanned(
158-
qx,
159-
candidate.purl,
160-
giveUpResult(p2.error, p2.attempts),
161-
undefined,
162-
false,
163-
)
167+
await markPackagistMetadataScanned(qx, candidate.purl, giveUpResult(p2.error, p2.attempts), {
168+
bumpLastRunAt: false,
169+
notBefore: scheduledAt,
170+
})
164171
return
165172
}
166173

@@ -180,13 +187,13 @@ export async function ingestOnePackagistMetadata(
180187
}
181188

182189
await logAuditFieldChanges(qx, WORKER, candidate.purl, phase2ChangedFields)
183-
// markPackagistMetadataScanned treats a null 4th arg the same as omitting it
184-
// (COALESCEd against the stored value), so lastModified can always be passed.
185190
await markPackagistMetadataScanned(
186191
qx,
187192
candidate.purl,
188193
{ status: 'success', attempts: p2.attempts },
189-
lastModified,
194+
{
195+
metadataLastModified: lastModified,
196+
},
190197
)
191198
}
192199

@@ -195,6 +202,7 @@ export async function ingestOnePackagist30dWindow(
195202
qx: QueryExecutor,
196203
purl: string,
197204
runDate: string,
205+
scheduledAt: string,
198206
): Promise<void> {
199207
const name = packagistNameFromPurl(purl)
200208

@@ -207,7 +215,7 @@ export async function ingestOnePackagist30dWindow(
207215
{ purl, statusCode: info.error.statusCode, kind: info.error.kind },
208216
'packagist 30d downloads 4xx/malformed after fast retries — marking processed and skipping',
209217
)
210-
await markPackagist30dProcessed(qx, purl, giveUpResult(info.error, info.attempts))
218+
await markPackagist30dProcessed(qx, purl, giveUpResult(info.error, info.attempts), scheduledAt)
211219
return
212220
}
213221

@@ -226,6 +234,7 @@ export async function ingestOnePackagistDailyDownload(
226234
qx: QueryExecutor,
227235
candidate: PackagistDailyCandidate,
228236
runDate: string,
237+
scheduledAt: string,
229238
): Promise<void> {
230239
const name = packagistNameFromPurl(candidate.purl)
231240

@@ -238,7 +247,12 @@ export async function ingestOnePackagistDailyDownload(
238247
{ purl: candidate.purl, statusCode: info.error.statusCode, kind: info.error.kind },
239248
'packagist daily downloads 4xx/malformed after fast retries — marking processed and skipping',
240249
)
241-
await markPackagistDailyProcessed(qx, candidate.purl, giveUpResult(info.error, info.attempts))
250+
await markPackagistDailyProcessed(
251+
qx,
252+
candidate.purl,
253+
giveUpResult(info.error, info.attempts),
254+
scheduledAt,
255+
)
242256
return
243257
}
244258

@@ -332,6 +346,8 @@ export async function ingestPackagistMetadataBatch(
332346
if (candidates.length === 0) return
333347
const qx = await getPackagesDb()
334348
const attempt = Context.current().info.attempt
349+
// Stable across every Temporal retry of this same batch — see notBefore below.
350+
const scheduledAt = new Date(Context.current().info.scheduledTimestampMs).toISOString()
335351

336352
// The merged lane starts every ingest with a DYNAMIC-endpoint fetch, so it is
337353
// bounded by that endpoint's 10-concurrent limit — not p2's 20. Running hotter
@@ -340,13 +356,16 @@ export async function ingestPackagistMetadataBatch(
340356
candidates,
341357
attempt,
342358
statsConcurrency(),
343-
(candidate) => ingestOnePackagistMetadata(qx, candidate),
359+
(candidate) => ingestOnePackagistMetadata(qx, candidate, scheduledAt),
344360
(candidate, err) =>
345-
markPackagistMetadataScanned(qx, candidate.purl, {
346-
status: 'error',
347-
attempts: attempt,
348-
message: String(err),
349-
}),
361+
markPackagistMetadataScanned(
362+
qx,
363+
candidate.purl,
364+
{ status: 'error', attempts: attempt, message: String(err) },
365+
// An item that already succeeded earlier in this same batch's retry sequence
366+
// must not have that success overwritten by an unrelated re-processing failure.
367+
{ notBefore: scheduledAt },
368+
),
350369
)
351370

352371
log.info({ count: candidates.length }, 'Ingested Packagist metadata batch')
@@ -366,18 +385,20 @@ export async function ingestPackagist30dBatch(purls: string[], runDate: string):
366385
if (purls.length === 0) return
367386
const qx = await getPackagesDb()
368387
const attempt = Context.current().info.attempt
388+
const scheduledAt = new Date(Context.current().info.scheduledTimestampMs).toISOString()
369389

370390
await ingestPackagistItemsConcurrently(
371391
purls,
372392
attempt,
373393
statsConcurrency(),
374-
(purl) => ingestOnePackagist30dWindow(qx, purl, runDate),
394+
(purl) => ingestOnePackagist30dWindow(qx, purl, runDate, scheduledAt),
375395
(purl, err) =>
376-
markPackagist30dProcessed(qx, purl, {
377-
status: 'error',
378-
attempts: attempt,
379-
message: String(err),
380-
}),
396+
markPackagist30dProcessed(
397+
qx,
398+
purl,
399+
{ status: 'error', attempts: attempt, message: String(err) },
400+
scheduledAt,
401+
),
381402
)
382403

383404
log.info({ count: purls.length }, 'Ingested Packagist 30d downloads batch')
@@ -403,18 +424,20 @@ export async function ingestPackagistDailyBatch(
403424
if (candidates.length === 0) return
404425
const qx = await getPackagesDb()
405426
const attempt = Context.current().info.attempt
427+
const scheduledAt = new Date(Context.current().info.scheduledTimestampMs).toISOString()
406428

407429
await ingestPackagistItemsConcurrently(
408430
candidates,
409431
attempt,
410432
statsConcurrency(),
411-
(candidate) => ingestOnePackagistDailyDownload(qx, candidate, runDate),
433+
(candidate) => ingestOnePackagistDailyDownload(qx, candidate, runDate, scheduledAt),
412434
(candidate, err) =>
413-
markPackagistDailyProcessed(qx, candidate.purl, {
414-
status: 'error',
415-
attempts: attempt,
416-
message: String(err),
417-
}),
435+
markPackagistDailyProcessed(
436+
qx,
437+
candidate.purl,
438+
{ status: 'error', attempts: attempt, message: String(err) },
439+
scheduledAt,
440+
),
418441
)
419442

420443
log.info({ count: candidates.length }, 'Ingested Packagist daily downloads batch')

0 commit comments

Comments
 (0)