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:

  1. Connector (Pds#full_dataset_files) skips shards already ingested when resuming the same snapshot.
  2. SyncService marker lifecycle pins/clears full_sync_target_sequence so first_sync? keeps routing to /all for 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)
  end

Shards 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 until on 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.rb
  • malware_pds_spec / malware_offline_spec: resuming an interrupted first sync returns only shards above checkpoint.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
end

PDS (online) — needs a short-lived staging IJWT:

  1. On the staging rails console, mint the token and copy it:
    puts CloudConnector::Tokens.get(unit_primitive: :malware_advisories, resource: :instance)
    Save it to /tmp/pds_ijwt.txt locally (a credential — do not commit).
  2. 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:

  1. dest=vendor/package_metadata/malware_advisories/v3/npm
    mkdir -p "$dest" && tar xf /path/to/full_dataset.tar -C "$dest"
  2. 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 SyncService run on staging awaits malware ingestion (#602431 (closed)) — the marker lifecycle is unit-tested and the connector resume-skip is staging-verified.

References

Edited by Bala Kumar

Merge request reports

Loading
Loading