Skip to content

How the strategist plans a flow

For each asked question the strategist picks targets:

  1. there is a rule: the rule's arguments;
  2. a head was trained with fit: the features the head selected;
  3. uses is set: those facts;
  4. 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 def parts (fn, extract, check, rule, alternative producers) are awaited — also when marked blocking=True, which only matters for sync parts; a sync part marked blocking=True runs in a worker thread (asyncio.to_thread); any other sync part runs inline, as in ask.
  • 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 that ask would not make, and a failed hard check stops the paid lookups behind it. speculate=True starts 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). Cancelling aask itself cancels every pending call. Under a learned order (System(order="learned"), learn_order()) the hard checks run one at a time in that order, so speculate=True is ignored, with a UserWarning.
  • Timeouts. A call that takes longer than its part's timeout= (or aask(timeout=), or System(timeout=)) fails with timed out after 2 s: the fact is missing and the questions that need it abstain with guard timeout (a hard check that times out: "could not be evaluated", as for any error). The safeguard timeout is in res.safeguards, the audit and system.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 it blocking=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 of ask on the same input, whatever finished first (tested on every gallery case and on the examples, with and without speculate). 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 / Route work as under ask. Concurrent aask calls on one System are safe on one event loop: its costs, stats and storage are updated between awaits.
  • ask still works on a catalog with async def parts: each such call runs in an event loop of its own, one after another (on a worker thread when ask is called from a running loop). system.is_async says whether a catalog has parts that aask awaits; solvi serve answers such systems with aask.

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.