Resume PDS /all full sync from the last ingested shard
What does this MR do?
Makes the PDS /all (full-dataset) malware sync resume from the last ingested
shard after an interruption, instead of restarting from the first shard. Two
parts:
- Connector (
Pds#full_dataset_files) skips shards already ingested when resuming the same snapshot. SyncServicemarker lifecycle pins/clearsfull_sync_target_sequencesofirst_sync?keeps routing to/allfor the whole first sync — the interrupted run resumes there instead of flipping to/delta.
Stacked on !243687 (merged) (the PDS / offline connectors) — target that branch, review after it merges. Part of #603628 (closed).
Background
The connector MR downloads every /all shard on each call. On a resumed first
sync it re-downloads and re-parses shards already ingested — idempotent (upserts)
but wasteful, and for a ~215k-advisory / 64-shard snapshot that is a lot of
repeated network work each time a job is interrupted by its time budget.
Change
Connector: skip ingested shards on resume
full_dataset_files now takes the checkpoint and skips done shards when
resuming the same snapshot:
snapshot_until = body['until']
resuming = checkpoint.sequence.to_i == snapshot_until
Array.wrap(body['shards'])
.sort_by { |entry| entry['shard'].to_i(16) }
.lazy.flat_map do |entry|
shard = entry['shard'].to_i(16)
next [] if resuming && shard <= checkpoint.chunk.to_i
files_from(entry['signed_url'], sequence: snapshot_until, chunk: shard)
endShards are ordered by id so checkpoint.chunk cleanly partitions done vs.
remaining, and skipped shards are never downloaded. A fresh sync
(sequence 0) or a changed snapshot (different until) re-fetches every
shard, so we never skip against a stale cursor.
Connector::MalwareOffline applies the same skip, so the air-gapped sync resumes
identically — reading shards from disk instead of signed URLs.
SyncService: full-sync marker lifecycle
first_sync? is full_sync_target_sequence.present? || sequence == 0. Without a
marker, after the first shard the checkpoint's sequence becomes until, so the
next run sees first_sync? == false, flips to /delta, and drops the rest of
the dataset. SyncService now manages the marker for malware syncs:
def execute
full_sync = full_dataset_sync? # malware_advisories? && checkpoint.first_sync?
connector.data_after(checkpoint).each do |file|
ingest_file(file)
checkpoint.update(**checkpoint_position(file, full_sync)) # pins the marker on the first shard
return log_stop_signal if signal.stop?
end
finalize_full_sync if full_sync # clears the marker once the loop drains
end- Pinned to the snapshot
untilon the first shard, cleared only when the loop drains (the whole dataset is ingested). An interrupt returns before finalize, so the marker persists and the next run resumes via/all. - Gated on
malware_advisories?, so advisories / licenses / cve are unaffected. - Also fixes the offline malware first sync, which previously stopped after
the first shard once
first_sync?turned false.
Verification
Unit specs
bundle exec rspec \
ee/spec/lib/gitlab/package_metadata/connector/malware_pds_spec.rb \
ee/spec/lib/gitlab/package_metadata/connector/malware_offline_spec.rb \
ee/spec/services/package_metadata/sync_service_spec.rbmalware_pds_spec/malware_offline_spec: resuming an interrupted first sync returns only shards abovecheckpoint.chunk; already-ingested shards are not re-fetched / re-read.sync_service_spec: marker pinned on the first shard and cleared on completion; retained on interrupt (first_sync?stays true); cleared when a marker-set checkpoint is resumed to completion.
Reproduce the resumable sync locally (both connectors)
Both connectors resume from the last ingested shard. This helper simulates
SyncService: drive data_after with a fresh first-sync checkpoint, stop after a
few shards, then re-drive with the full_sync_target_sequence marker set (as the
sync flow does) and confirm the remaining shards are yielded.
# Stands in for PackageMetadata::Checkpoint (first_sync? / sequence / chunk).
class Cp
def initialize(seq, chunk, first) = (@seq = seq; @chunk = chunk; @first = first)
def first_sync? = @first
def sequence = @seq
def chunk = @chunk
end
def resume_demo(connector, interrupt: 3)
processed = []
until_seq = nil
connector.data_after(Cp.new(0, 0, true)).each do |f|
until_seq = f.sequence; processed << f.chunk; break if processed.size >= interrupt
end
resumed = []
connector.data_after(Cp.new(until_seq, processed.last, true)).each do |f|
resumed << f.chunk; break if resumed.size >= interrupt
end
{ processed: processed, resumed: resumed } # resumed should start at processed.last + 1
endPDS (online) — needs a short-lived staging IJWT:
- On the staging rails console, mint the token and copy it:
Save it to
puts CloudConnector::Tokens.get(unit_primitive: :malware_advisories, resource: :instance)/tmp/pds_ijwt.txtlocally (a credential — do not commit). - In a local
rails console, drive the real connector against staging, matching the identity headers to the token's claims:require 'jwt' token = File.read('/tmp/pds_ijwt.txt').strip claims = JWT.decode(token, nil, false).first cfg = PackageMetadata::SyncConfiguration.new( 'malware_advisories', :pds, PackageMetadata::SyncConfiguration::Location::PDS_MALWARE_STAGING_ENDPOINT, 'v3', :npm) connector = Gitlab::PackageMetadata::Connector::MalwarePds.new(cfg) connector.define_singleton_method(:request_headers) do CloudConnector.headers(nil).merge( 'x-gitlab-realm' => claims['gitlab_realm'], 'x-gitlab-instance-id' => claims['sub'], 'Authorization' => "Bearer #{token}") end resume_demo(connector) # => { processed: [0, 1, 2], resumed: [3, 4, 5] }
Offline (air-gapped) — no token; unpack the vendor full_dataset:
-
dest=vendor/package_metadata/malware_advisories/v3/npm mkdir -p "$dest" && tar xf /path/to/full_dataset.tar -C "$dest" - In
rails console:cfg = PackageMetadata::SyncConfiguration.new( 'malware_advisories', :offline, PackageMetadata::SyncConfiguration::Location::MALWARE_ADVISORIES_PATH, 'v3', :npm) resume_demo(Gitlab::PackageMetadata::Connector::MalwareOffline.new(cfg)) # => { processed: [0, 1, 2], resumed: [3, 4, 5] }
Observed
- PDS (staging,
until 1783077798, 64 shards, ~215k advisories): fresh run processed[0, 1, 2]; resume yielded[3, 4, 5]— skipped the ingested shards. - Offline (vendor
full_dataset,until 1783005651): fresh[0, 1, 2]; resume[3, 4, 5].
Not in this MR
- End-to-end
SyncServicerun on staging awaits malware ingestion (#602431 (closed)) — the marker lifecycle is unit-tested and the connector resume-skip is staging-verified.
References
- Issue: #603628 (closed)
- Connector MR (stacked on): !243687 (merged)
- Epic: &20876