1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
//! The orchestration runtime behind [`script!`](crate::script!).
//!
//! `script!` is a thin syntax layer: every block it accepts desugars into one
//! call in this module, and this module is where the semantics live. That split
//! is deliberate. Syntax is read by people and models; semantics are read by
//! the compiler, the tests, and anyone running `cargo expand`. Keeping the
//! semantics in ordinary, documented functions means the expansion of a script
//! is a short list of calls an auditor can follow, not a generated state
//! machine nobody can read.
//!
//! # What a flow cannot get wrong
//!
//! The estate's fix history (45 repositories, 2,766 fix and revert commits as
//! of 2026-09-30) is dominated by the same few orchestration defects: work with
//! no ceiling, waits with no deadline, retries that duplicate an effect,
//! identities that leak across tenants, errors that vanish, and tasks nobody
//! owns. Each primitive here makes one of those unrepresentable rather than
//! documented:
//!
//! | Block | Primitive | What it rules out |
//! |---|---|---|
//! | `each x in xs, at most N at once:` | [`each`] | unbounded fan-out; lost results; orphaned siblings after a failure |
//! | `FanOut::new(xs).at_most(N).within(d).run(body)` | [`FanOut`] | the same, with no scope to open and your own error type back, with the failing item's position |
//! | `within 2s:` | [`within`] | a wait with no deadline; a wait that ignores cancellation |
//! | `retry up to 3 times, waiting 100ms:` | [`retry`] | retry storms; retrying a permanent failure; a new identity per attempt |
//! | `together:` | `try_join!` | sequential awaits that should overlap; a failed branch that keeps running |
//! | `step name:` | [`Scope::enter`] | anonymous work with no stable key or error location |
//! | `await ready(&db, 30s)` | [`Readiness`] | a readiness guessed by sleeping; a poll loop; a stale instance releasing dependants |
//! | `admit the model's plan:` | [`admit`] | an untrusted payload becoming an instruction; a refusal with no location, no provenance or no ceiling |
//!
//! [`admit`] is the one way a task body crosses untrusted model or tool output
//! into a run. It enters its step, charges one run-scoped [`Gate`] — a decoder, a
//! surface, an admission budget and a repair ledger — turns a
//! [`proposal::Refusal`] or a
//! [`proposal::Intervention`] into its own
//! [`FlowError`] arm carrying the [`Provenance`](crate::proposal::Provenance)
//! of the exact bytes, and records the refusal through the run store so a run
//! resumed on a fresh host reads back what this run refused. It admits a
//! [`Plan`](crate::proposal::Plan) of operation *names*; it performs nothing and
//! is not a fifth verb.
//!
//! Every flow takes a [`Scope`] first. The scope carries the [`Tenant`], the
//! path of steps that led here, and the cancellation token, so identity,
//! location and cancellation reach every step without being threaded by hand,
//! and a [`StepKey`] is always the hash of *tenant and path*: two tenants
//! running the same flow over the same input never share a key.
//!
//! # Orchestration is the architecture
//!
//! A `script!` block also emits a constant `ARCHITECTURE` ([`Architecture`]):
//! the flows it declares and the tree of blocks inside each, with bounds,
//! deadlines, retry budgets and the flows each one runs. It is compiled from the
//! same tokens as the code, so it cannot drift from it, and it is the thing to
//! read (or diff) instead of re-deriving the orchestration from source.
//!
//! # Boundary
//!
//! [`each`] drives its bodies concurrently **on the task that awaits it**. That
//! is what lets a body borrow the flow's locals the way a Python loop body
//! would, with no `Arc`, no `clone`, and no `'static` bound. A flow is `Send`
//! exactly when what its steps hold is: bodies are stored as their own future
//! types, never erased, so a flow can be spawned onto a multi-threaded runtime
//! and each tenant's flow runs on whichever worker is free. The bodies of one
//! `each` share that one task. Concurrency is for waiting on many things at
//! once. CPU parallelism across cores is a different tool:
//! [`join_all_bounded`](crate::rt::task::join_all_bounded) on owned inputs.
//!
//! [`each`]: crate::script::each
//! [`within`]: crate::script::within
//! [`retry`]: crate::script::retry
//! [`admit`]: crate::script::admit
//! [`Gate`]: crate::script::Gate
//! [`FlowError`]: crate::script::FlowError
//! [`Scope`]: crate::script::Scope
//! [`Scope::enter`]: crate::script::Scope::enter
//! [`Tenant`]: crate::script::Tenant
//! [`StepKey`]: crate::script::StepKey
//! [`Architecture`]: crate::script::Architecture
//! [`Readiness`]: crate::script::Readiness
//! [`FanOut`]: crate::script::FanOut
pub
pub
pub
use Duration;
// The clock is re-exported here rather than reached through `rt::clock` because
// this is the module that *consumes* it: a flow author writing `within` needs
// the type, and the two paths to it is one more name to keep straight. The
// underlying module stays the single definition.
use crateclock as rt_clock;
pub use ;
pub use ;
pub use each;
pub use ;
pub use ;
pub use ;
pub use ;
pub use Clock;
pub use ;
pub use ;
/// How many step paths a scope's trail retains unless a host says otherwise.
///
/// A reported bound rather than an implicit one: a run that entered thousands
/// of `each` items and retained them all would make the report the largest
/// object in the process, which is the opposite of what a report is for. See
/// [`crate::task::Report::steps`].
pub const DEFAULT_TRAIL_STEPS: usize = 64;
/// The longest tenant name a [`Tenant`] accepts, in bytes.
///
/// A tenant name is part of every key and every error location, so it is
/// bounded like any other payload that is copied into each of them.
pub const MAX_TENANT_BYTES: usize = 128;
/// How deeply steps may nest inside one flow run.
///
/// A flow that runs itself, directly or through another flow, grows its path
/// by one segment per call. The limit turns that into a typed
/// [`FlowError::TooDeep`] instead of an ever-longer path and an eventual stack
/// overflow.
pub const MAX_DEPTH: u16 = 64;
/// The most bodies one [`each`] may have in flight at once.
///
/// `at most N at once` is required, and this caps what `N` may say: a bound of
/// a million is a bound in name only.
pub const MAX_IN_FLIGHT: usize = 65_536;
/// The most attempts one [`retry`] may make.
pub const MAX_ATTEMPTS: u32 = 1_000;
/// The longest single wait between two [`retry`] attempts, before jitter.
pub const MAX_BACKOFF: Duration = from_secs;
/// Bodies one poll of [`each`] may poll before it hands the executor back.
///
/// Without a ceiling, a fan-out whose bodies keep waking themselves is one
/// endless poll. That is not hypothetical: the async engine's cooperative
/// budget answers a task that has done enough work in one poll with `Pending`
/// **and an immediate wake**, so a driver that re-polls every woken body before
/// returning spins on that wake forever. It was observed as a hang of a
/// 10,000-item fan-out over timers. Returning after this many polls, with the
/// task re-woken, lets the engine reset its budget; the woken bodies stay
/// queued and are polled on the next turn. 128 matches the engine's own
/// per-poll budget.
pub const POLL_BUDGET: usize = 128;