How the strategist plans a flow¶
For each asked question the strategist picks targets:
- there is a rule: the rule's arguments;
- a head was trained with
fit: the features the head selected; usesis set: those facts;- otherwise: everything computable from
init_state(the flow is marked as not narrowed).
It then walks backwards from the targets through the catalog signatures to the keys of init_state, and always adds:
- the question's required parts (
requires); - every check whose inputs are already available in the flow and which touches at least one computed (not input) fact, a "check of what was computed".
Each part runs once per request, even if several questions need it. Steps are executed in topological order, rules
last; hard checks and the facts they read are ordered first. Catalog parts that are not needed, or cannot run because their inputs are missing, are not executed, and
res.flow.skipped says which case applies. A fact on a cycle of the catalog can never be computed: the questions that
need it abstain (solvi check reports the cycle). PlanError is raised for a checkpoint that names no part.
The flow depends only on which keys init_state has, the catalog, the questions and the trained heads. It is the same for
every request with the same keys.
Other strategists: dead ends and costs¶
System(..., strategist=...) takes another planner. solvi.strategy.CostStrategist() builds the same plan with producers
whose inputs are never given dropped (the deterministic strategist needs the inputs of every producer of a fact);
CostStrategist(producers="equivalent") treats the producers of a fact as interchangeable and picks the cheapest verified
plan by declared cost=, keeping every hard check that governs a question (System(..., producers="equivalent") is a
shortcut for it). Both are code only. A segment model (ModelStrategist.load) and the name matcher of solvi.aliases are
experimental: no checkpoint is published for either. Details, the
trace record of a plan and what was measured: docs/strategist.md.
Costs from measurements. System(cat, questions, producers="equivalent", cost_policy="measured") plans with the run times
system.cost_book measures instead of declared costs: after a warm-up (each producer measured min_samples times; an
undeclared one is tried at 0 ms, a declared one keeps its cost= until measured) it picks the fastest of equivalent
producers — a local table over a 300 ms feed — and switches when that one slows down; a producer unused for recheck
asks gets one more trial. system.freeze_costs() stops the switching (unfreeze_costs() resumes). The plan record of each
trace says, per fact, which cost decided and where it came from (declared, warm-up, measured, recheck, frozen). Settings:
costs=solvi.costs.MeasuredCosts(min_samples=3, recheck=50, alpha=None); see
docs/strategist.md.
Early exit and parallel execution¶
At run time the executor first computes the hard checks and what they depend on. If a hard check fails, every question whose
flow contains it is settled (forced by then, or abstained), and the steps that only those questions needed are not run.
They are listed in res.trace.skipped with the check that made them unnecessary. This is the default, because the
skipped rest is often the expensive part (a model, an API). Its price: a decision forced by a hard check has no rule
values or downstream facts in its record, and a part listed in requires — it is in the flow, but the question was
settled before it ran — is missing from res.values.
When the record must hold everything — a scorecard whose points you want for every stored decision, a proposal to hand to a person when it is rejected — compute the whole flow anyway:
res = system.ask(state, early_exit=False) # this ask; System(cat, questions, early_exit=False): every ask
res["approve"].status # "forced": the failed hard check still decides
res.values["points"], res.trace.skipped # every fact and rule value is there; nothing was skipped
res.trace.early_exit # False: recorded in the trace (and in a stored response)
The answers are the same either way; only the steps that run differ. The trace records the switch, and a replay with
the flow checks that no planned step is missing from such a trace. aask, ask_text and aask_text take the same
argument; System.facts_for computes the whole flow for training.
System(catalog, questions, workers=8) (or system.ask(state, workers=8)) runs independent steps in parallel threads: a step
starts as soon as the steps it reads have finished. This pays off when parts wait on I/O — HTTP APIs, databases, model
inference that releases the GIL. Answers, records and hashes are identical to a sequential run, because records are written in
flow order after execution.
system = System(cat, QUESTIONS, workers=8)
res = system.ask(claim)
print(res.trace.skipped) # e.g. [("fraud_score", "not needed: hard check policy_in_force failed"), ...]
Async execution: aask¶
await system.aask(state) is ask on an event loop, for catalogs whose parts wait on the network — database lookups,
HTTP APIs, model servers — and for callers that are async themselves (web servers, agents; Pyodide in the browser):
@cat.fn(timeout=2.0) # seconds; else System(timeout=...), else no limit
async def credit_score(customer_id):
async with httpx.AsyncClient() as c:
return (await c.get(f"{BUREAU}/score/{customer_id}")).json()["score"]
@cat.fn(blocking=True) # a sync client: aask runs it in a worker thread
def sanctions_hit(name):
return screening.lookup(name)
system = System(cat, QUESTIONS, timeout=5.0)
res = await system.aask(application) # same Response as ask
res = await system.aask(application, speculate=True)
async defparts (fn, extract, check, rule, alternative producers) are awaited — also when markedblocking=True, which only matters for sync parts; a sync part markedblocking=Trueruns in a worker thread (asyncio.to_thread); any other sync part runs inline, as inask.- Steps run as soon as the steps they read have finished, all concurrently. By default in the phases of
ask: hard checks and what they read first, then what the open questions still need — so no call starts thataskwould not make, and a failed hard check stops the paid lookups behind it.speculate=Truestarts every step as soon as its inputs are ready and cancels the pending calls a failed hard check makes unnecessary (lower latency; some calls may start and be cancelled; steps that finished anyway are dropped). Cancellingaaskitself cancels every pending call. Under a learned order (System(order="learned"),learn_order()) the hard checks run one at a time in that order, sospeculate=Trueis ignored, with aUserWarning. - Timeouts. A call that takes longer than its part's
timeout=(oraask(timeout=), orSystem(timeout=)) fails withtimed out after 2 s: the fact is missing and the questions that need it abstain with guardtimeout(a hard check that times out: "could not be evaluated", as for any error). The safeguardtimeoutis inres.safeguards, the audit andsystem.stats["timeouts"]. A producer of a fact that times out is followed by the next producer. A plain sync part running inline cannot be interrupted: mark itblocking=True(the thread finishes in the background, its result is ignored) or make it async. - Same trace as
ask. Records are written in flow order after the run, so the answers, the records and every hash are those ofaskon the same input, whatever finished first (tested on every gallery case and on the examples, with and withoutspeculate). A replay re-runs async parts in an event loop of its own and does not re-run a step that timed out (a timeout depends on the moment, not on the inputs). - Storage, the audit, safeguards, batched decisions and
Cascade/Vote/Routework as underask. Concurrentaaskcalls on one System are safe on one event loop: its costs, stats and storage are updated between awaits. askstill works on a catalog withasync defparts: each such call runs in an event loop of its own, one after another (on a worker thread whenaskis called from a running loop).system.is_asyncsays whether a catalog has parts thataaskawaits;solvi serveanswers such systems withaask.
Plain CPU parts gain nothing from aask: for them the sync ask stays the default.
Learned order of hard checks¶
Every ask measures the run time of each part: system.cost_book keeps a moving average (ms) per part (cost= on a decorator
is the prior until a part has run; cost_policy="measured" also feeds it to the planner, see above). While learning is on —
System(order="learned"), producers="learned", learn=True, or after learn_order() — it also records which hard
checks failed on which input; a default System does not (learn=False: no work inside ask beyond the costs).
System(cat, questions, order="learned") — or system.learn_order(examples) on a list of init_states, which runs only the
hard checks and what they read and then switches the order — makes the executor evaluate hard checks one at a time, the
one with the highest expected saving first:
score = P(check fails | cheap facts) × cost of the steps its failure would skip ÷ cost of evaluating it
and stop as soon as the failed checks settle every question they govern. P(fail) comes from a small online model per hard
check (system.order_model: Laplace counts, then a FastHead refitted every 50 rows and updated in between) on the scalar
values of init_state (and one level of dicts, e.g. invoice.currency). learn_order(features=[...]) adds cheap computed
facts; they are computed before the hard checks.
Answers are identical to the default order. When several hard checks fail, the first one declared in the catalog decides;
so a question is settled by a failed check only after every hard check that governs it and is declared earlier has been
evaluated. Those earlier checks are scheduled next, since evaluating them settles the question whatever they return. Records
stay in flow order. What changes is only which steps run: res.trace.skipped lists the rest, and
res.trace.explain_order() (also printed by show) says why each hard check ran when it did:
1. policy_in_force: P(fail) 0.09 × saves 222.6 ms ÷ costs 0.0 ms = 712 → passed
2. fraud_ok: P(fail) 0.81 × saves 64.0 ms ÷ costs 170.2 ms = 0.305 → failed
3. no_litigation: ... [unblocks decision (already decided by a failed check)] → passed; settles decision, fast_track
ask(state, order=...) overrides the order for one request: "default", "learned", or any object with
p_fail(check, row) and row(vals, init_keys) (e.g. an oracle for experiments).
The learned order saves time only when hard checks fail often enough and their failure skips expensive work; when nothing
fails every hard check still runs. With workers > 1 the default order runs all hard checks at once, while the learned order
runs them one after another (their inputs still run in parallel) — measure before choosing it for parallel execution.