agentplane 0.4.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
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
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
//! Plans — compiled artifacts, and the authorization graph.
//!
//! Whatever produced a plan — a manifest, capability routing, or a model — the
//! output is one validated [`PlanIR`], frozen and content-addressed before
//! anything runs.
//!
//! # The plan is an authorization graph
//!
//! Because a plan is compiled from *trusted* input only and frozen before any
//! untrusted data is touched, it is more than a schedule: it is a statement of
//! what this run is permitted to do, made before anything could have influenced
//! it. The journal that follows can then be checked against it.
//!
//! That check goes further than "which tools may run". Every argument declares
//! where it comes from — a named upstream node, the run's input, or a constant —
//! and the executor refuses an argument whose actual provenance does not match.
//! Labels say *how much to trust* a value; source bindings say *where it came
//! from*. Both are needed: a label alone permits substituting one untrusted
//! value for another.

use std::collections::{BTreeMap, BTreeSet};

use serde::{Deserialize, Serialize};
use serde_json::Value;

use crate::core::{Capability, Digest, StepId, canon};

/// How many agents contribute to one task.
///
/// Declared rather than emergent, because the choice is expensive in a way that
/// is easy to make by accident: coordination between agents is the largest
/// measured failure category in multi-agent systems after specification, and it
/// is a category that *only exists if you opt into it*.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Topology {
    /// One agent, one context, many tools. No inter-agent surface at all.
    #[default]
    Single,
    /// Several agents contribute to one task.
    ///
    /// Requires a justification (see [`Collaboration`]) that the plan contract
    /// checks, because the cost is real and the benefit is often assumed.
    Collaborative(Collaboration),
}

/// Why collaboration is worth its cost here.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum Collaboration {
    // `ContextOverflow` used to sit here and was removed rather than kept as a
    // documented no-op. Whether work exceeds a context window is not a property
    // of the graph, so the contract could not check it — and an unchecked
    // justification is not a weak control, it is an *escape hatch*: a plan
    // refused as false parallelism was approved by editing one word, which made
    // the two real checks optional for anyone who noticed. I12 leaves two
    // choices for a control nothing enforces, and this is the other one.
    /// Sub-tasks operate on disjoint inputs.
    ///
    /// Checked, not taken on trust: overlapping inputs are *false parallelism*,
    /// where the coordination cost is paid and no parallelism is obtained.
    ParallelDisjoint,
    /// Sub-tasks need strictly different authority.
    ///
    /// The best reason to split agents, and the one least often named: if a
    /// sub-task needs credentials the parent should not hold, delegating to a
    /// narrower agent buys least privilege rather than hypothetical speed.
    DistinctAuthority,
}

/// Where one argument comes from.
///
/// The whole point of naming this: an argument whose actual provenance differs
/// from its declaration is refused before the effect is dispatched.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "from", rename_all = "snake_case")]
pub enum ArgSource {
    /// The run's own input, optionally a field of it.
    RunInput {
        #[serde(default, skip_serializing_if = "Option::is_none")]
        field: Option<String>,
    },
    /// An upstream node's output, optionally a field of it.
    Node {
        step: StepId,
        #[serde(default, skip_serializing_if = "Option::is_none")]
        field: Option<String>,
    },
    /// A value fixed in the plan. Trusted by construction: it was frozen before
    /// anything untrusted was read.
    Const { value: Value },
}

impl ArgSource {
    #[must_use]
    pub fn run_input() -> Self {
        Self::RunInput { field: None }
    }

    #[must_use]
    pub fn input_field(f: impl Into<String>) -> Self {
        Self::RunInput {
            field: Some(f.into()),
        }
    }

    #[must_use]
    pub fn node(step: StepId) -> Self {
        Self::Node { step, field: None }
    }

    #[must_use]
    pub fn node_field(step: StepId, f: impl Into<String>) -> Self {
        Self::Node {
            step,
            field: Some(f.into()),
        }
    }

    #[must_use]
    pub fn constant(v: Value) -> Self {
        Self::Const { value: v }
    }

    /// The node this argument depends on, if any.
    #[must_use]
    pub fn depends_on(&self) -> Option<StepId> {
        match self {
            Self::Node { step, .. } => Some(*step),
            _ => None,
        }
    }
}

/// One unit of work in a plan.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PlanNode {
    pub id: StepId,
    /// What this node needs done. Resolved against registered skills.
    pub capability: Capability,
    /// Nodes that must complete first.
    ///
    /// Structural, not conversational: a node cannot run before its
    /// predecessors' outputs are bound to its inputs, which is how "one agent
    /// ignored another's result" stops being possible rather than discouraged.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub depends_on: Vec<StepId>,
    /// Each argument and where it comes from.
    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
    pub args: BTreeMap<String, ArgSource>,
    /// Whether the plan is complete once this node is done.
    #[serde(default)]
    pub terminal: bool,
    /// Whether this node checks another's work.
    ///
    /// Named so the contract can require one: nothing checking the work is a
    /// fifth of observed multi-agent failures.
    #[serde(default)]
    pub verifies: bool,
    /// Judge this node's work more than once, from declared angles.
    ///
    /// For steps where a single execution is not adequate evidence — the
    /// pass^k collapse. `None` is one judgement, which is the right default:
    /// a panel on every node would pay the cost everywhere and mean nothing
    /// anywhere.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub quorum: Option<crate::core::Quorum>,
}

impl PlanNode {
    pub fn new(id: u32, capability: impl Into<Capability>) -> Self {
        Self {
            quorum: None,
            id: StepId(id),
            capability: capability.into(),
            depends_on: Vec::new(),
            args: BTreeMap::new(),
            terminal: false,
            verifies: false,
        }
    }

    /// Bind an argument, recording the dependency it implies.
    #[must_use]
    pub fn arg(mut self, name: impl Into<String>, source: ArgSource) -> Self {
        if let Some(dep) = source.depends_on()
            && !self.depends_on.contains(&dep)
        {
            self.depends_on.push(dep);
        }
        self.args.insert(name.into(), source);
        self
    }

    #[must_use]
    pub fn after(mut self, step: u32) -> Self {
        let s = StepId(step);
        if !self.depends_on.contains(&s) {
            self.depends_on.push(s);
        }
        self
    }

    #[must_use]
    pub fn terminal(mut self) -> Self {
        self.terminal = true;
        self
    }

    #[must_use]
    pub fn verifies(mut self) -> Self {
        self.verifies = true;
        self
    }

    /// Judge this node from several declared angles, requiring agreement.
    ///
    /// For steps where one execution is not adequate evidence. Failure to reach
    /// the quorum escalates; there is deliberately no way to resolve a split
    /// panel to a majority — see [`Quorum`](crate::core::Quorum).
    #[must_use]
    pub fn with_quorum(mut self, quorum: crate::core::Quorum) -> Self {
        self.quorum = Some(quorum);
        self
    }
}

/// A frozen, content-addressed plan.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PlanIR {
    pub version: u32,
    /// The plan this one replaced, if any.
    ///
    /// Replanning produces a *new version* rather than mutating in place, so the
    /// audit trail shows what the run intended before it changed its mind —
    /// usually the interesting part, and structurally absent from any system
    /// that edits a plan.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub derived_from: Option<Digest>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub reason: Option<String>,
    pub topology: Topology,
    pub nodes: Vec<PlanNode>,
}

impl PlanIR {
    #[must_use]
    pub fn new(nodes: Vec<PlanNode>) -> Self {
        Self {
            version: 1,
            derived_from: None,
            reason: None,
            topology: Topology::Single,
            nodes,
        }
    }

    /// A one-step plan, for the common case of "just run this".
    #[must_use]
    pub fn single(capability: impl Into<Capability>) -> Self {
        Self::new(vec![
            PlanNode::new(0, capability)
                .arg("input", ArgSource::run_input())
                .terminal(),
        ])
    }

    /// Run several capabilities on the same input, concurrently, then feed
    /// every result to one aggregator.
    ///
    /// The shape people mean by "fan out to all the matching specialists and
    /// combine what they say", written once instead of by hand. The branches
    /// have no edge between them, so they are one ready set and are dispatched
    /// **concurrently**; the aggregator depends on all of them, so it runs when
    /// the last finishes.
    ///
    /// ```
    /// # use agentplane::core::PlanIR;
    /// let plan = PlanIR::fan_out(
    ///     ["billing.anomaly", "billing.regulatory"],
    ///     "billing.decide",
    /// );
    /// assert_eq!(plan.nodes.len(), 3);
    /// ```
    ///
    /// The aggregator receives one argument per branch, named for the
    /// capability that produced it — so adding a specialist does not silently
    /// renumber what the aggregator reads, which hand-wiring by `StepId` does.
    ///
    /// # Why this is a plan rather than an effect
    ///
    /// A `race`-style primitive — dispatch N, take the first, abandon the rest —
    /// is deliberately **not** offered. Abandoning an in-flight branch
    /// manufactures exactly the unknown outcome the effect protocol exists to
    /// prevent: the loser was announced, it may have reached a model or a tool,
    /// and cancelling it mid-flight leaves a started effect with no terminal
    /// record. Every branch here therefore runs to completion and every outcome
    /// is on the record, which costs more and is the only version that can be
    /// replayed or recovered from a crash.
    ///
    /// # Panics
    ///
    /// If `branches` is empty. A fan-out with nothing to fan out to is a
    /// one-node plan written the long way, and accepting it would produce an
    /// aggregator with no subject — the shape [`validate`](crate::plan::validate)
    /// refuses for verifiers, for the same reason.
    #[must_use]
    pub fn fan_out(
        branches: impl IntoIterator<Item = impl Into<Capability>>,
        aggregate: impl Into<Capability>,
    ) -> Self {
        let branches: Vec<Capability> = branches.into_iter().map(Into::into).collect();
        assert!(
            !branches.is_empty(),
            "a fan-out needs at least one branch; with none the aggregator has \
             nothing to aggregate"
        );

        let mut nodes: Vec<PlanNode> = branches
            .iter()
            .enumerate()
            .map(|(i, capability)| {
                PlanNode::new(u32::try_from(i).unwrap_or(u32::MAX), capability.clone())
                    .arg("input", ArgSource::run_input())
            })
            .collect();

        let join_id = u32::try_from(branches.len()).unwrap_or(u32::MAX);
        let mut join = PlanNode::new(join_id, aggregate).terminal();
        for (i, capability) in branches.iter().enumerate() {
            let step = StepId(u32::try_from(i).unwrap_or(u32::MAX));
            join = join.arg(&capability.0, ArgSource::node(step));
        }
        nodes.push(join);
        Self::new(nodes)
    }

    #[must_use]
    pub fn topology(mut self, t: Topology) -> Self {
        self.topology = t;
        self
    }

    /// Content address over the canonical form.
    #[must_use]
    pub fn digest(&self) -> Digest {
        Digest::of(&canon::to_bytes(self).unwrap_or_default())
    }

    #[must_use]
    pub fn node(&self, id: StepId) -> Option<&PlanNode> {
        self.nodes.iter().find(|n| n.id == id)
    }

    /// Nodes whose dependencies are all satisfied and which have not run.
    ///
    /// Returned in a deterministic total order — topological rank, then id — so
    /// replay reproduces dispatch order exactly. Without that, a plan with any
    /// parallelism would replay differently every time.
    #[must_use]
    pub fn ready(&self, done: &BTreeSet<StepId>) -> Vec<StepId> {
        let mut ready: Vec<StepId> = self
            .nodes
            .iter()
            .filter(|n| !done.contains(&n.id))
            .filter(|n| n.depends_on.iter().all(|d| done.contains(d)))
            .map(|n| n.id)
            .collect();
        ready.sort_unstable();
        ready
    }

    /// Whether every terminal node has run.
    ///
    /// Completion is this, and only this. A workload asserting it finished is
    /// not evidence — agents announce success on unmet objectives often enough
    /// that self-report cannot be the signal.
    #[must_use]
    pub fn is_complete(&self, done: &BTreeSet<StepId>) -> bool {
        self.nodes
            .iter()
            .filter(|n| n.terminal)
            .all(|n| done.contains(&n.id))
    }
}

/// A plan that must not run, and why.
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum PlanError {
    #[error("the plan has no nodes")]
    Empty,

    #[error("step {0} appears more than once")]
    DuplicateStep(StepId),

    #[error("step {step} depends on {missing}, which is not in the plan")]
    MissingDependency { step: StepId, missing: StepId },

    #[error("the plan has a dependency cycle involving {0}")]
    Cycle(StepId),

    /// Without a terminal node nothing ever declares the plan finished, and the
    /// run would either loop or stop for no stated reason.
    #[error("the plan has no terminal node, so nothing marks it complete")]
    NoTerminal,

    #[error("step {step} is unreachable: nothing depends on it and it is not terminal")]
    Unreachable { step: StepId },

    #[error("no skill provides capability '{capability}' required by step {step}")]
    NoProvider { step: StepId, capability: String },

    // Not named `source`: `thiserror` reserves that for an error cause.
    #[error("step {step} takes argument '{arg}' from {from_step}, which is not an upstream node")]
    ArgumentNotUpstream {
        step: StepId,
        arg: String,
        from_step: StepId,
    },

    #[error("step {step} has no bound arguments, so its input is undefined")]
    NoArguments { step: StepId },

    /// A verifier that could not have seen the work it claims to check.
    #[error("step {step} verifies nothing: a verifier must depend on what it checks")]
    VerifierWithoutSubject { step: StepId },

    /// A panel was declared on a node that judges nothing.
    ///
    /// A quorum is several judgements of *one piece of work*. On a node with no
    /// subject it is several executions of the work itself, which is not a
    /// second opinion — it is the same act performed repeatedly, and for a
    /// mutating step it is the same act performed repeatedly *on the world*.
    #[error(
        "step {step} declares a quorum but judges nothing; a panel needs \
         something to judge, or it is repetition rather than review"
    )]
    QuorumWithoutSubject { step: StepId },

    #[error("this plan requires a verifier node and has none")]
    VerifierRequired,

    #[error(
        "collaboration claims parallel-disjoint, but steps {a} and {b} read the same source — \
         paying coordination cost for parallelism that is not there"
    )]
    FalseParallelism { a: StepId, b: StepId },

    #[error(
        "collaboration claims distinct-authority, but every step needs the same capability \
         '{capability}' — there is no authority to separate"
    )]
    NoAuthorityToSeparate { capability: String },

    #[error("the plan needs {steps} steps but the budget allows {allowed}")]
    TooManySteps { steps: usize, allowed: usize },
}