Class: ActiveSanction::Sync

Inherits:
Object
  • Object
show all
Extended by:
T::Sig
Defined in:
lib/active_sanction/sync.rb,
lib/active_sanction/sync/report.rb,
lib/active_sanction/sync/result.rb

Overview

Fetch, parse and store every configured list in one run, isolating each source from the others.

report = ActiveSanction.sync!                 # every configured source
report = ActiveSanction.sync!(:ofac_sdn)      # one
report = ActiveSanction.sync!(force: true)    # bypass conditional GET

report.failed?                                # => false
report[:ofac_sdn].record_count                # => 19015

One source at a time is Sources::Base#sync, which fetches, parses and checksums, and deliberately rescues nothing. This is the layer above it: what a run does, which is a different set of decisions.

One source failing must not abort the others

Government endpoints go down, change format without notice, and occasionally serve half a file. If a UN outage stopped OFAC from syncing, the library would fail exactly when it is most needed -- during an incident, which is when lists move. So every source runs inside its own rescue: the failure is captured into its Result, the remaining sources carry on, and the run ends with a summary that says which one broke.

StandardError and not Exception. An Interrupt or a SIGTERM is somebody stopping this run on purpose, and swallowing it to go on downloading three more lists is not isolation, it is a job that will not die.

A failed source keeps its previous snapshot

This is the single most important behaviour here, and it is a decision about what "no data" costs. Nothing clears a source's stored list on failure -- not a 500, not a parse error, not a publisher that started serving HTML where XML used to be. Screening against yesterday's OFAC list produces a report with a known, visible age on it; screening against an empty list produces a clean report for every customer, which is the most expensive thing this library can get wrong.

That trade is only safe while the age is visible, so every Result carries the age and record count of the snapshot that source is being screened against now -- see Sync::Result, which is where the failure ends up.

Unchanged sources cost nothing

The launch lists change daily at most and every one of them serves ETag and Last-Modified, so an hourly sync should transfer bytes once a day. A source whose publisher answers 304 is never parsed and never stored: the whole saving of conditional GET (#10) is that the parse -- the expensive half for OFAC's three-file join -- is skipped along with the download.

A source whose bytes changed but whose content hashes to what is already stored is also reported unchanged and not rewritten. A publisher regenerating an identical file with a new timestamp is not a new list version, and rewriting tens of megabytes to say so would churn the checksum every audit record cites.

Optional parallelism, with a politeness limit

ActiveSanction.sync!(concurrency: 3)

Sources that share a publisher are never fetched at the same time. They are grouped by the host they download from and each group runs sequentially, so raising concurrency fetches from more governments at once and never harder from any one of them -- which matters because two of the built-in adapters (OFAC SDN and OFAC Consolidated) are the same file server. The default is 1 and a run of four lists takes about as long as its slowest list.

Each source is fetched by its own adapter, so nothing is shared between two sources in flight except the files underneath them: the validator store is one small JSON file, written by rename and resolving last-writer-wins, exactly as it already does for a sync running beside a CLI command. The cost of losing that race is one avoidable download on the next run, and it is the reason the default is sequential rather than the reason parallelism is unsafe.

Where this belongs

Syncing is a capability of the local backend (#56), not of every backend: a hosted one does not sync, because data freshness is exactly what its user is paying somebody else to handle, and it should say so through supports?(:sync). When that seam lands this becomes Backend::Local#sync unchanged -- which is why the report is a serializable object rather than console output, and why nothing here writes to $stdout.

Defined Under Namespace

Classes: Failed, Report, Result

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(sources: nil, store: nil, force: false, concurrency: nil, logger: ActiveSanction.config.logger, instrumenter: ActiveSanction.config.instrumenter) ⇒ void

sources: takes source keys, adapter classes, adapter instances, or nil for whatever config.sources names. A key nothing is registered under raises here, before the first list is downloaded, rather than after.

Parameters:

  • sources (T.untyped) (defaults to: nil)
  • store (T.untyped) (defaults to: nil)
  • force (Boolean) (defaults to: false)
  • concurrency (T.untyped) (defaults to: nil)
  • logger (T.untyped) (defaults to: ActiveSanction.config.logger)
  • instrumenter (T.untyped) (defaults to: ActiveSanction.config.instrumenter)


174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
# File 'lib/active_sanction/sync.rb', line 174

def initialize(sources: nil, store: nil, force: false, concurrency: nil, logger: ActiveSanction.config.logger,
               instrumenter: ActiveSanction.config.instrumenter)
  @sources = T.let(resolve(sources), T::Array[T.untyped])
  @keys = T.let(@sources.map { |source| Sources::Definition.key!(source.key) }, T::Array[Symbol])
  @store = T.let(store || ActiveSanction.storage, T.untyped)
  @force = T.let(force, T::Boolean)
  @concurrency = T.let(
    Configuration.sync_concurrency!(concurrency || ActiveSanction.config.sync_concurrency), Integer
  )
  @logger = T.let(logger, T.untyped)
  @instrumenter = T.let(instrumenter, T.untyped)
  @lock = T.let(Mutex.new, Mutex)
  # The settings this run was started under, so a worker thread reads them
  # rather than the default client's. A configuration is fiber-local and a
  # `Thread.new` does not inherit one -- see
  # ActiveSanction.with_configuration -- so a run through a Client with its
  # own User-Agent would otherwise identify itself as that client on the
  # first source and as the default on the next three, purely according to
  # `concurrency:`.
  @configuration = T.let(ActiveSanction.config, Configuration)
end

Instance Attribute Details

#concurrency ⇒ Integer (readonly)

How many publishers to fetch from at once. See the class comment.

Returns:

  • (Integer)


148
149
150
# File 'lib/active_sanction/sync.rb', line 148

def concurrency
  @concurrency
end

#force ⇒ Boolean (readonly)

Returns:

  • (Boolean)


144
145
146
# File 'lib/active_sanction/sync.rb', line 144

def force
  @force
end

#instrumenter ⇒ T.untyped (readonly)

Where the :sync and :store events go, or nil for nothing listening.

The :fetch and :parse events of the sources this run covers do not come from here: an adapter is constructed by the run and reads the configuration, exactly as it does for its logger and its User-Agent. So a host that instruments through ActiveSanction.configure sees all six events, and one that hands a run its own instrumenter sees the two this class emits. See Instrumentation.

Returns:

  • (T.untyped)


162
163
164
# File 'lib/active_sanction/sync.rb', line 162

def instrumenter
  @instrumenter
end

#keys ⇒ Array<Symbol> (readonly)

Returns:

  • (Array<Symbol>)


138
139
140
# File 'lib/active_sanction/sync.rb', line 138

def keys
  @keys
end

#logger ⇒ T.untyped (readonly)

Returns:

  • (T.untyped)


151
152
153
# File 'lib/active_sanction/sync.rb', line 151

def logger
  @logger
end

#sources ⇒ Array<T.untyped> (readonly)

The adapters this run covers: classes as registered, or instances a caller passed in.

Returns:

  • (Array<T.untyped>)


135
136
137
# File 'lib/active_sanction/sync.rb', line 135

def sources
  @sources
end

#store ⇒ T.untyped (readonly)

Returns:

  • (T.untyped)


141
142
143
# File 'lib/active_sanction/sync.rb', line 141

def store
  @store
end

Class Method Details

.call(**options, &block) ⇒ Report

Parameters:

  • options (T.untyped)
  • block (T.untyped)

Returns:



165
# File 'lib/active_sanction/sync.rb', line 165

def self.call(**options, &block) = T.unsafe(self).new(**options).call(&block)

Instance Method Details

#call(&block) ⇒ Report

Runs the sync and returns the Report. Never raises for a source that failed -- that is what the report is for -- and does raise for anything that makes the run itself impossible, such as a store that cannot be written to at all.

The optional block is the progress hook: it is called with each Result as that source finishes, so a long run says something before it ends. It is called under a lock, so a block that appends to an array or writes a line does not have to be thread-safe to be correct under concurrency:.

ActiveSanction.sync! { |result| puts result }

Parameters:

  • block (T.proc.params(result: Result).void, nil)

Returns:



208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
# File 'lib/active_sanction/sync.rb', line 208

def call(&block)
  started_at = Time.now.utc
  began = monotonic
  log_start
  # The run's own duration is the Report's, taken from the same clock
  # reading, so an event and the report a caller is holding never disagree
  # about how long a run took.
  report = Instrumentation.instrument(instrumenter, :sync,
                                      { sources: keys, forced: force, concurrency: concurrency }) do |event|
    work = keys.each_with_index.map { |key, at| [at, key, sources.fetch(at)] }
    results = run(work, &block).sort_by(&:first).map(&:last)
    finished = Report.new(results: results, started_at: started_at, duration: elapsed(began))
    summarize(event, finished)
    finished
  end
  log_finish(report)
  report
end

#inspect ⇒ String

Returns:

  • (String)


228
# File 'lib/active_sanction/sync.rb', line 228

def inspect = "#<#{self.class} #{keys.join(", ")}#{" forced" if force} concurrency=#{concurrency}>"