Class: ActiveSanction::Sync
- Inherits:
-
Object
- Object
- ActiveSanction::Sync
- 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
-
#concurrency ⇒ Integer
readonly
How many publishers to fetch from at once.
- #force ⇒ Boolean readonly
-
#instrumenter ⇒ T.untyped
readonly
Where the
:syncand:storeevents go, or nil for nothing listening. - #keys ⇒ Array<Symbol> readonly
- #logger ⇒ T.untyped readonly
-
#sources ⇒ Array<T.untyped>
readonly
The adapters this run covers: classes as registered, or instances a caller passed in.
- #store ⇒ T.untyped readonly
Class Method Summary collapse
Instance Method Summary collapse
-
#call(&block) ⇒ Report
Runs the sync and returns the Report.
-
#initialize(sources: nil, store: nil, force: false, concurrency: nil, logger: ActiveSanction.config.logger, instrumenter: ActiveSanction.config.instrumenter) ⇒ void
constructor
sources:takes source keys, adapter classes, adapter instances, or nil for whateverconfig.sourcesnames. - #inspect ⇒ String
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.
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.
148 149 150 |
# File 'lib/active_sanction/sync.rb', line 148 def concurrency @concurrency end |
#force ⇒ Boolean (readonly)
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.
162 163 164 |
# File 'lib/active_sanction/sync.rb', line 162 def instrumenter @instrumenter end |
#keys ⇒ Array<Symbol> (readonly)
138 139 140 |
# File 'lib/active_sanction/sync.rb', line 138 def keys @keys end |
#logger ⇒ T.untyped (readonly)
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.
135 136 137 |
# File 'lib/active_sanction/sync.rb', line 135 def sources @sources end |
#store ⇒ T.untyped (readonly)
141 142 143 |
# File 'lib/active_sanction/sync.rb', line 141 def store @store end |
Class Method Details
.call(**options, &block) ⇒ Report
165 |
# File 'lib/active_sanction/sync.rb', line 165 def self.call(**, &block) = T.unsafe(self).new(**).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 }
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
228 |
# File 'lib/active_sanction/sync.rb', line 228 def inspect = "#<#{self.class} #{keys.join(", ")}#{" forced" if force} concurrency=#{concurrency}>" |