Resume malware advisory /all first sync from the last ingested shard
What does this MR do?
Wires the resumable-first-sync checkpoint marker into PackageMetadata::MalwareAdvisorySyncService so an interrupted malware advisory first sync resumes from the last fully ingested shard instead of restarting the whole /all snapshot from shard 00.
Why
Merge request !244523 (merged) added resumability via a full_sync_target_sequence checkpoint marker, keeping checkpoint.first_sync? true across an interruption so subsequent runs skip already-ingested shards. That marker lifecycle was implemented in PackageMetadata::SyncService, which handles only advisories, licenses, and cve_enrichment — never malware_advisories. MalwareAdvisorySyncService, the actual entry point for malware syncs, never sets or clears the marker, so resumability is dead code for malware.
Correctness is preserved today (an interrupted first sync restarts /all rather than dropping shards), but it restarts from shard 00 and re-ingests everything. That is a latent livelock risk as the corpus grows: a bootstrap that cannot finish in one run's budget re-ingests the same prefix every run and never completes. This is a robustness fix for a latent issue, not a response to a live incident. Per-/all resumability is the goal of #603628 (closed).
How
The change is confined to ee/app/services/package_metadata/malware_advisory_sync_service.rb:
- Capture
full_sync = checkpoint.first_sync?at the start ofingest_stream, before any checkpoint mutation. - Delta path (unchanged): advance the checkpoint per archive, on the sequence boundary.
- Full-sync path (new): after each shard is ingested, advance the checkpoint and set
full_sync_target_sequence, sofirst_sync?stays true across runs and the connector resumes/allfrom the last completed shard. - On clean, uninterrupted completion of a full sync, clear the marker so later runs take
/delta. - Nothing is committed on an interruption or persistence failure, so the run always resumes from the last fully ingested shard.
The delta path and the bulk /delta path are unchanged.
Testing
Unit specs added to ee/spec/services/package_metadata/malware_advisory_sync_service_spec.rb cover a full first sync (per-shard checkpoint advance with the marker set, marker cleared on completion) and an interrupted first sync (progress and marker preserved so the next run resumes).
Local testing
/all resumability — offline and PDS, before/after
Verified on a local GDK with the offline dataset vendored at vendor/package_metadata/malware_advisories/v3/npm/full_dataset/ (64 shards 00–3f + checkpoint.json), and against staging PDS. Each run resets the npm checkpoint to a fresh first sync, interrupts after 3 shards, then reports the checkpoint and which shards the next run would process.
Steps (offline shown; for PDS build a :pds config against PDS_MALWARE_STAGING_ENDPOINT and supply a staging instance token, overriding MalwarePds.instance_token / request_headers to match the token's gitlab_realm and sub):
# gitlab-rails runner
ENV['PM_SYNC_IN_DEV'] = 'true'
config = PackageMetadata::SyncConfiguration.configs_for('malware_advisories').find { |c| c.purl_type.to_s == 'npm' }
cp = PackageMetadata::Checkpoint.for_malware_advisories.find_by!(purl_type: :npm)
cp.update!(sequence: 0, chunk: 0, full_sync_target_sequence: nil) # fresh first sync
signal = Object.new
signal.define_singleton_method(:stop?) { @n = (@n || 0) + 1; @n >= 3 } # interrupt after 3 shards
PackageMetadata::MalwareAdvisorySyncService.new(config, signal, checkpoint: cp).execute
cp.reload
puts "after interrupt: seq=#{cp.sequence} chunk=#{cp.chunk} marker=#{cp.full_sync_target_sequence.inspect} first_sync?=#{cp.first_sync?}"
conn = Gitlab::PackageMetadata::Connector::MalwareOffline.new(config)
puts "next run processes chunks: #{conn.data_after(cp).map(&:chunk).first(6).inspect}"Observations — the interrupted first sync resumes instead of restarting:
| Path | Before (master) |
After (this MR) |
|---|---|---|
| Offline | seq=0 chunk=0 marker=nil → next run [0, 1, 2, 3, 4, 5] (restarts from shard 00) |
seq=1783005651 chunk=1 marker=1783005651 → next run [2, 3, 4, 5, 6, 7] (resumes) |
| PDS (staging) | seq=0 chunk=0 marker=nil → next run [0, 1, 2, 3] (restarts from shard 00) |
seq=1783077798 chunk=1 marker=1783077798 → next run [2, 3, 4, 5] (resumes) |
In both cases first_sync? stays true across the interruption (via full_sync_target_sequence), so the connector re-selects /all and skips the already-ingested shards rather than dropping them or flipping to /delta.
Resume / interruption logging
The started and interrupted events (event: malware_advisory_sync, Gitlab::AppJsonLogger) now carry the shard (chunk) and sequence, and a resumed run sets resuming: true. From a local offline interrupt → resume:
{"phase":"started","sync_mode":"full","resuming":false,"from_sequence":0,"from_chunk":0}
{"phase":"interrupted","sync_mode":"full","to_sequence":1783005651,"to_chunk":1}
{"phase":"started","sync_mode":"full","resuming":true,"from_sequence":1783005651,"from_chunk":1}
{"phase":"interrupted","sync_mode":"full","to_sequence":1783005651,"to_chunk":2}from_* is where a run picks up (resuming: true marks a continued bootstrap); to_* is where an interrupted run stopped.
Database
The full /all sync now calls checkpoint.update(...) once per shard (previously once at the end). Each is a single-row UPDATE of pm_checkpoints by primary key:
UPDATE "pm_checkpoints"
SET "sequence" = $1, "chunk" = $2, "full_sync_target_sequence" = $3, "updated_at" = $4
WHERE "pm_checkpoints"."id" = $5- Bounded by the shard count of one snapshot (≤ 64 per
purl_type) and spaced by the inter-slice ingest throttle. - Same per-file checkpoint cadence
PackageMetadata::SyncServicealready uses for advisories/licenses. - No scope, bulk operation, or new/expensive query. A formal database review is likely unnecessary, but flagging per process.
References
- Resolves #617851 (closed)
- Marker introduced in !244523 (merged)
- Resumable first-sync design: #603628 (closed)
- Epic: &20876