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 of ingest_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, so first_sync? stays true across runs and the connector resumes /all from 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::SyncService already uses for advisories/licenses.
  • No scope, bulk operation, or new/expensive query. A formal database review is likely unnecessary, but flagging per process.

References

Edited by Bala Kumar

Merge request reports

Loading
Loading