Skip to content

Limits: how much a pipeline may consume

Building one?

Step 4 of the backend track fills these meters in from real stages — which ones reserve, which only charge, and what the debug panel makes of it.

A policy answers what a pipeline may be made of. This answers how much of anything it may use, and the two are separate because a graph built entirely of allowed stages can still loop, fan out into thousands of tasks, or run until the process dies.

from stageflow import Limits, Policy, Session

plan = Policy(
    stages={"LoadTicket", "ClassifyByRules", "LlmReply"},
    limits=Limits(
        counters={"seconds": 120, "steps": 5_000, "iterations": 2_000,
                  "tokens": 200_000, "llm_calls": 100},
        gauges={"concurrency": 8, "depth": 4, "frame_bytes": 4_000_000},
        max_retries=5,
        max_delay_seconds=30,
    ),
)

session = Session(id="run-1", pipeline=pipeline, context=ctx, policy=plan)
result = await session.run()
result.meters   # {"steps": 412, "seconds": 8.2, "tokens": 5120, "llm_calls": 9}

One mechanism: meters

Not three systems for time, cost and size — one, with two kinds of meter:

Kind Behaviour Examples Checked
counter grows, never falls seconds, steps, iterations, tokens on every charge
gauge an instantaneous value concurrency, depth, frame_bytes when it is taken up

The runtime charges steps and iterations and holds the gauges; stages charge whatever they spent. concurrency is the one gauge that waits rather than failing: a hundred items to process is a good pipeline that may not have a hundred calls in the air at once, and refusing it for being wide would be refusing the work rather than bounding it. Depth and size do fail — too deep is too deep, and waiting would not help. Meter names are not fixed by the core. A host counts what is scarce for it, and a meter nobody limits simply accumulates and comes back in result.meters, which is what a bill is made of.

seconds is the one hybrid, deliberately. It is reported like a counter so that it shows up in the same table as everything else, but it is enforced as a deadline — a counter is only checked when something charges it, and a stage that hangs charges nothing.

Nothing is measured unless it is limited

Every check costs something and frame_bytes costs the most: measuring it means serialising the frame. A meter absent from the limits is not computed at all, so the price of all this to a host that sets no limits is one if.

What a stage spends

A stage charges the units it consumed — tokens, llm_calls, http_calls, rows. What a unit is worth is a price list; price lists change without any code changing, and they belong to the host. A stage knows how much it used, not what that costs.

@register_stage("LlmTriageStage")
class LlmTriageStage(BaseStage):
    """
    description: "Classify a ticket with a model"
    reserve:
      llm_calls: 1
      tokens: "args.max_tokens + size(args.text) / 3"
    """

    timeout = 60

    async def run(self):
        answer = await self.client.chat(self.get_arguments()["text"])
        self.set_outputs({"topic": answer.topic})
        self.charge(tokens=answer.usage.total, llm_calls=1)

Two mechanisms, and they do not overlap:

Where Sees Answers
reserve the spec, declarative args may this be attempted at all
charge() the stage, in code everything the stage knows what it actually spent

reserve is declarative because it is read before anything runs: an editor can show "up to 4000 tokens a call", and a host can refuse a graph without executing it. Its values are numbers or CEL over args, the arguments as the stage will receive them.

The order is reserve → run → settle:

  1. the reservation is evaluated and held. It does not fit what is left → the stage never starts. An expensive call is better not begun than cut off after the money is gone;
  2. the stage runs;
  3. charge() replaces the reservation for the meters it names and adds the ones it does not — an amount both reserved and charged is not counted twice. A stage that charges nothing is settled at what it reserved, and that is true of a failure too: a call that went out and then timed out still spent what it spent.

A stage that reserves and charges nothing is free.

A loop is not reserved as a whole

Reservations are per stage, taken as that stage is about to run — nothing sums up what a map will cost before entering it, because the reservation of an iteration depends on its arguments, which depend on the item, which comes from the data. What is checked up front is the iteration count: iterations is charged before the first pass, so a loop over a thousand items on a plan allowing fifty is refused without running any of them.

A loop therefore stops as soon as it cannot afford the next stage, with the passes it did make already paid for. In parallel mode the overshoot is bounded by concurrency rather than by one stage: calls already in flight have already been made. A host that needs a hard ceiling for the whole loop gets it by arithmetic rather than prediction — cap iterations so that the count times the worst case fits the budget.

What validation refuses before the run

Pipeline.validate(policy) checks what is soundly knowable from the JSON, so that a tenant saving a graph is told what to change rather than finding out on the tenth element:

Checked Because
the shortest way through the graph against steps every run passes at least that many nodes, so a graph whose cheapest path does not fit cannot finish at all
how deep the declared subpipelines nest against depth the nesting is written down, not computed
what each retry asks for against max_retries and max_delay_seconds the numbers are in the JSON

The bound on steps is a lower one, which is what makes it safe to refuse by: a graph with a long road and a short one is judged by the short one.

What is not checked is anything that would have to be guessed. A loop's cost is unknowable before the run — the number of passes comes from the data and the reservation of a pass from its arguments — and multiplying something here would refuse pipelines that fit. A false refusal on a working graph is worse than the ceiling it was meant to save.

Running out

The ceiling is hard. BudgetExceeded inherits BaseException, deliberately outside the StageFlowError tree, because try blocks catch Exception and so does the retry machinery — a tenant able to wrap a graph in try and swallow this would make every limit decoration.

Session.run is the one place that catches it, and it ends the run the way a stop does, rather than throwing an exception through work that did happen:

result.result   # {"status": "budget_exceeded", "meter": "tokens",
                #  "limit": 200000, "spent": 200512}
result.meters   # everything spent, for the host to bill by

A child session of a subpipeline shares the parent's budget — a fresh one per child would multiply the allowance by the depth of nesting — and lets the stop travel up rather than swallowing it.

The tenant's knobs are requests

retry and the pauses between attempts come from the pipeline, which is to say from the less-trusted side, so max_retries and max_delay_seconds deal with them twice.

A node declaring more than the plan allows is refused at validation, like anything else written in the JSON: the tenant is told what to change rather than quietly given something other than what they asked for. A retry that reaches a session another way — a graph assembled in Python and never validated — is clamped instead: fifty attempts on a plan allowing three run three times, and the event says three.

Waiting for a person is not work

A pipeline that asks a question and waits is not consuming the host, so the deadline does not tick while it waits. Without that, a deadline would be a limit on how fast somebody reads.

This also fixes an older trap: a stage that waits for input used to die at its own 30-second timeout. A stage timeout is now measured in busy time as well, and clamped by what is left of the deadline — otherwise a stage allowed sixty seconds and started with two left would run for sixty, and the budget would leak by the length of its last stage.

What limits do not do

concurrency: 8 bounds one run. A hundred runs are eight hundred tasks: the core cannot know how many runs exist, so admission control is the host's — a queue and a pool per tenant. Likewise a monthly quota: the core enforces the ceiling of one run and reports what it spent; subtracting that from an allowance for the month, and turning units into money, is the platform's accounting.