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
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
//! A model provider with no model behind it.
//!
//! For tests, examples, and local runs where the point is the *plane* rather
//! than the answer. Lives in `testkit` — off by default, never in a production
//! build — for the same reason [`StubSigner`](super::StubSigner) does: something
//! that stands in for a real component must not be reachable by accident.
//!
//! # Two traps, both of which make a fake worse than useless
//!
//! **A fake that is not deterministic destroys the property under test.** This
//! crate exists so a run replays to the same answer; a fake returning arbitrary
//! text would make every example non-replayable and every replay test a
//! coin-toss that mostly passes. So the default answer is a pure function of the
//! prompt, and scripted answers are consumed in a fixed order.
//!
//! **A fake that reports zero usage makes every budget test vacuous.** Token
//! ceilings, cost ceilings, the metered-failure path — all of them read
//! [`Usage`], and a provider that always answers "free" lets them pass over a
//! runtime that has stopped counting. So usage is derived from the prompt, and
//! scripted failures can carry usage of their own.
//!
//! # What it does not fake
//!
//! Intelligence. The default answer is an echo, not a plausible completion. A
//! fake that produced convincing prose would invite tests that assert on
//! *content*, which is the one thing a real provider will never reproduce.

use std::sync::{Arc, Mutex};

use async_trait::async_trait;
use serde_json::{Value, json};

#[cfg(test)]
use crate::model::ModelCall;
use crate::model::{
    Completion, ModelError, ModelId, ModelProvider, ModelStreamEvent, ReasoningEffort, Request,
    Usage,
};

/// What the fake was asked.
#[derive(Debug, Clone, PartialEq)]
pub struct Ask {
    pub model: ModelId,
    pub prompt: Value,
    pub max_output_tokens: u32,
    pub reasoning_effort: Option<ReasoningEffort>,
    pub schema: Option<Value>,
    /// The exact tool surface the model was offered this turn.
    pub tools: Vec<crate::model::ToolDeclaration>,
    /// What this turn was told about the tools the last turn asked for.
    ///
    /// Recorded because it is the only place a *refusal* reaches a model, and
    /// without it no test can assert what the model was told — which is how a
    /// loop handing back the precise policy message passed every test in the
    /// suite. `PolicyError::for_model` existed, was tested, and had no callers.
    pub exchanges: Vec<crate::model::ToolExchange>,
}

/// A `ModelProvider` that answers without a model.
#[derive(Debug, Default)]
pub struct FakeProvider {
    /// Answers handed out in order, before the default takes over.
    scripted: Mutex<std::collections::VecDeque<Result<Completion, ModelError>>>,
    /// Every call, in order.
    asked: Mutex<Vec<Ask>>,
    /// Whether to emit text deltas to a caller's observer before answering.
    streaming: Mutex<bool>,
}

impl FakeProvider {
    /// A provider that always echoes.
    #[must_use]
    pub fn new() -> Arc<Self> {
        Arc::new(Self::default())
    }

    /// Queue one answer. Consumed before the default echo.
    ///
    /// Takes `&self` rather than `self` so a test can arrange a provider it has
    /// already handed to a runtime — which is the usual shape, because the
    /// runtime wants an `Arc` at construction.
    pub fn will_answer(&self, completion: Completion) -> &Self {
        self.scripted
            .lock()
            .expect("fake")
            .push_back(Ok(completion));
        self
    }

    /// Queue one failure.
    ///
    /// Use the metered variants — `Unusable`, `Interrupted` — to exercise the
    /// path that matters: a call that generated, cost money, and produced
    /// nothing usable. A fake that can only fail for free cannot test the
    /// ceiling that exists for exactly that case.
    pub fn will_fail(&self, error: ModelError) -> &Self {
        self.scripted.lock().expect("fake").push_back(Err(error));
        self
    }

    /// Queue a plain text answer with usage derived from its length.
    pub fn will_say(&self, text: impl Into<String>) -> &Self {
        let text = text.into();
        let usage = usage_for(&json!(&text));
        self.will_answer(Completion {
            tool_calls: Vec::new(),
            text,
            usage,
            stop_reason: Some("end_turn".to_owned()),
            truncated: false,
            structured: None,
            continuation: None,
        })
    }

    /// Queue a turn in which the model asks for a tool.
    ///
    /// The id is the one a real provider would issue and the loop must echo
    /// back, so a test exercises the pairing rather than assuming it.
    pub fn will_call_tool(
        &self,
        id: impl Into<String>,
        name: impl Into<String>,
        arguments: serde_json::Value,
    ) -> &Self {
        self.will_answer(Completion {
            tool_calls: vec![crate::model::ToolCall {
                id: id.into(),
                name: name.into(),
                arguments,
            }],
            text: String::new(),
            usage: Usage::default(),
            stop_reason: Some("tool_use".to_owned()),
            truncated: false,
            structured: None,
            continuation: None,
        })
    }

    /// Emit the answer as text deltas before returning it whole.
    ///
    /// Streaming is the one provider behaviour a fake could not reach, and its
    /// absence had a cost: an embedder building a live view had nothing to test
    /// against, so the only exercise of the observer path in this repository ran
    /// against a stub HTTP server — which tests the SSE parser, not the seam.
    ///
    /// The split is deliberately **not** a tokenizer. Whitespace with the
    /// separator kept on the preceding chunk means the concatenation of every
    /// delta is byte-identical to [`Completion::text`], which is the property an
    /// observer actually depends on and the one a clever split would break.
    ///
    /// What this does not fake is timing. Deltas are delivered synchronously,
    /// before the completion returns, because the alternative is a fake whose
    /// output order depends on the scheduler — and a non-deterministic fake
    /// destroys the property under test.
    pub fn streaming(&self) -> &Self {
        *self.streaming.lock().expect("fake") = true;
        self
    }

    /// Everything it was asked, in order.
    #[must_use]
    pub fn asked(&self) -> Vec<Ask> {
        self.asked.lock().expect("fake").clone()
    }

    /// How many times it was called.
    #[must_use]
    pub fn calls(&self) -> usize {
        self.asked.lock().expect("fake").len()
    }

    /// Whether every scripted answer was used.
    ///
    /// Worth asserting at the end of a test: leftover answers mean the run made
    /// fewer calls than the test believed, and a test that scripts three
    /// responses and checks the result of one is not testing what it thinks.
    #[must_use]
    pub fn script_exhausted(&self) -> bool {
        self.scripted.lock().expect("fake").is_empty()
    }
}

/// Chunk text so that concatenating every chunk reproduces it exactly.
///
/// The separator stays on the chunk before it. That is the whole trick: an
/// observer's job is usually to append deltas into a buffer, so a split that
/// dropped or duplicated a space would make the buffer disagree with the
/// canonical `Completion::text` in a way no assertion on chunk *count* notices.
fn split_for_stream(text: &str) -> Vec<String> {
    if text.is_empty() {
        return Vec::new();
    }
    let mut chunks = Vec::new();
    let mut current = String::new();
    for ch in text.chars() {
        current.push(ch);
        if ch.is_whitespace() {
            chunks.push(std::mem::take(&mut current));
        }
    }
    if !current.is_empty() {
        chunks.push(current);
    }
    chunks
}

/// Deterministic token counts, so budgets are exercised rather than bypassed.
///
/// Four bytes to a token is wrong for every real tokenizer and right for this
/// purpose: it is stable, monotonic in prompt size, and never zero — which are
/// the three properties a budget test actually depends on.
fn usage_for(prompt: &Value) -> Usage {
    let len = prompt.to_string().len() as u64;
    Usage {
        input_tokens: (len / 4).max(1),
        output_tokens: (len / 8).max(1),
        cache_write_tokens: 0,
        cache_read_tokens: 0,
        minor_units: 0,
    }
}

/// The default answer: a pure function of the request.
///
/// An echo rather than plausible prose, deliberately. A fake that produced
/// convincing text would invite assertions on content, and content is the one
/// thing a real provider will never reproduce.
fn echo(request: &Request<'_>) -> Completion {
    let usage = usage_for(request.prompt);
    match request.schema {
        // A schema was asked for, so the answer must satisfy the *shape*
        // contract: valid JSON. Built from the schema's declared properties so
        // it is at least plausibly conformant, without this file becoming a
        // JSON Schema implementation.
        Some(schema) => {
            let value = sample(schema);
            Completion {
                tool_calls: Vec::new(),
                text: value.to_string(),
                usage,
                stop_reason: Some("end_turn".to_owned()),
                truncated: false,
                structured: Some(value),
                continuation: None,
            }
        }
        None => Completion {
            tool_calls: Vec::new(),
            text: format!("fake answer to {}", request.prompt),
            usage,
            stop_reason: Some("end_turn".to_owned()),
            truncated: false,
            structured: None,
            continuation: None,
        },
    }
}

/// A minimal value satisfying a schema's declared types.
///
/// Handles the shapes a test is likely to declare and falls back to `null`
/// elsewhere. Not a JSON Schema implementation and not trying to be — the
/// crate deliberately does not validate schemas, so a fake that did would be
/// asserting a contract the real path never checks.
fn sample(schema: &Value) -> Value {
    match schema.get("type").and_then(Value::as_str) {
        Some("object") => {
            let mut out = serde_json::Map::new();
            if let Some(props) = schema.get("properties").and_then(Value::as_object) {
                for (name, sub) in props {
                    out.insert(name.clone(), sample(sub));
                }
            }
            Value::Object(out)
        }
        Some("array") => match schema.get("items") {
            Some(items) => json!([sample(items)]),
            None => json!([]),
        },
        Some("string") => json!("fake"),
        Some("number" | "integer") => json!(0),
        Some("boolean") => json!(false),
        _ => Value::Null,
    }
}

#[async_trait]
impl ModelProvider for FakeProvider {
    async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
        self.asked.lock().expect("fake").push(Ask {
            model: request.model.clone(),
            prompt: request.prompt.clone(),
            max_output_tokens: request.max_output_tokens,
            reasoning_effort: request.reasoning_effort,
            schema: request.schema.cloned(),
            tools: request.tools.to_vec(),
            exchanges: request.exchanges.to_vec(),
        });

        // Scoped so the guard is gone before anything else happens: a lock held
        // across a suspension is held on the thread, and this one is taken on
        // every model call in the suite.
        let scripted = self.scripted.lock().expect("fake").pop_front();
        let answer = scripted.unwrap_or_else(|| Ok(echo(&request)));

        // Deltas first, then the whole answer — the order a real driver
        // produces, so an observer that assumes it is exercised rather than
        // assumed. Only on the success path: a call that failed before
        // generating has no text to have streamed.
        if *self.streaming.lock().expect("fake")
            && let (Ok(completion), Some((observer, label))) = (&answer, request.stream)
        {
            for delta in split_for_stream(&completion.text) {
                observer.event(crate::core::Tainted::with_label(
                    ModelStreamEvent::TextDelta(delta),
                    label.clone(),
                ));
            }
            observer.event(crate::core::Tainted::with_label(
                ModelStreamEvent::Usage(completion.usage),
                label.clone(),
            ));
        }
        answer
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn model() -> ModelId {
        ModelId::new("fake", "m")
    }

    fn ask(schema: Option<&Value>) -> Request<'static> {
        // Leaked so the borrow is `'static` and the test reads as one line.
        // A test binary's leak is a test binary's problem.
        let prompt: &'static Value = Box::leak(Box::new(json!({"q": "what is the balance"})));
        let model: &'static ModelId = Box::leak(Box::new(model()));
        Request {
            model,
            prompt,
            max_output_tokens: ModelCall::DEFAULT_MAX_OUTPUT_TOKENS,
            reasoning_effort: None,
            schema: schema.map(|s| &*Box::leak(Box::new(s.clone()))),
            tools: &[],
            exchanges: &[],
            continuation: None,
            stream: None,
        }
    }

    /// The property the whole crate rests on. A fake that broke it would make
    /// every replay test a coin-toss that mostly passes.
    #[tokio::test]
    async fn the_same_question_gets_the_same_answer() {
        let p = FakeProvider::new();
        let a = p.complete(ask(None)).await.unwrap();
        let b = p.complete(ask(None)).await.unwrap();
        assert_eq!(a.text, b.text);
        assert_eq!(a.usage, b.usage);
    }

    /// Trap two: a provider that always answers "free" lets every budget test
    /// pass over a runtime that has stopped counting.
    #[tokio::test]
    async fn an_answer_is_never_free() {
        let p = FakeProvider::new();
        let c = p.complete(ask(None)).await.unwrap();
        assert!(
            c.usage.spend().tokens > 0,
            "a fake reporting zero usage makes every ceiling test vacuous"
        );
    }

    /// Longer prompt, more tokens — so a test can drive a run *over* a ceiling
    /// rather than merely up to a non-zero one.
    #[tokio::test]
    async fn usage_grows_with_the_prompt() {
        let p = FakeProvider::new();
        let short: &'static Value = Box::leak(Box::new(json!("hi")));
        let long: &'static Value = Box::leak(Box::new(json!("hi".repeat(500))));
        let m = model();
        let a = p
            .complete(Request {
                model: &m,
                prompt: short,
                max_output_tokens: ModelCall::DEFAULT_MAX_OUTPUT_TOKENS,
                reasoning_effort: None,
                schema: None,
                tools: &[],
                exchanges: &[],
                continuation: None,
                stream: None,
            })
            .await
            .unwrap();
        let b = p
            .complete(Request {
                model: &m,
                prompt: long,
                max_output_tokens: ModelCall::DEFAULT_MAX_OUTPUT_TOKENS,
                reasoning_effort: None,
                schema: None,
                tools: &[],
                exchanges: &[],
                continuation: None,
                stream: None,
            })
            .await
            .unwrap();
        assert!(b.usage.spend().tokens > a.usage.spend().tokens);
    }

    #[tokio::test]
    async fn scripted_answers_come_back_in_order_then_the_default_takes_over() {
        let p = FakeProvider::new();
        p.will_say("first").will_say("second");

        assert_eq!(p.complete(ask(None)).await.unwrap().text, "first");
        assert_eq!(p.complete(ask(None)).await.unwrap().text, "second");
        assert!(p.script_exhausted());
        assert!(
            p.complete(ask(None)).await.unwrap().text.contains("fake"),
            "past the script, the default echo answers"
        );
        assert_eq!(p.calls(), 3);
    }

    /// The path that matters: a call that generated, was billed, and produced
    /// nothing usable.
    #[tokio::test]
    async fn a_scripted_failure_can_carry_usage() {
        let p = FakeProvider::new();
        p.will_fail(ModelError::Interrupted {
            model: model(),
            usage: Usage {
                input_tokens: 100,
                output_tokens: 300,
                ..Usage::default()
            },
            detail: "reset".to_owned(),
        });
        let e = p.complete(ask(None)).await.expect_err("scripted");
        assert_eq!(e.usage().spend().tokens, 400);
    }

    #[tokio::test]
    async fn a_schema_gets_json_shaped_like_it() {
        let schema = json!({
            "type": "object",
            "properties": {
                "verdict": {"type": "string"},
                "score":   {"type": "integer"},
                "flags":   {"type": "array", "items": {"type": "boolean"}},
            },
        });
        let p = FakeProvider::new();
        let c = p.complete(ask(Some(&schema))).await.unwrap();
        let v = c.structured.expect("a schema was asked for");
        assert_eq!(v["verdict"], json!("fake"));
        assert_eq!(v["score"], json!(0));
        assert_eq!(v["flags"], json!([false]));
        assert_eq!(
            c.text,
            v.to_string(),
            "`text` holds the raw string even when a schema was parsed"
        );
    }

    #[tokio::test]
    async fn no_schema_means_no_structured_value() {
        let p = FakeProvider::new();
        assert!(p.complete(ask(None)).await.unwrap().structured.is_none());
    }

    #[tokio::test]
    async fn it_records_what_it_was_asked() {
        let schema = json!({"type": "string"});
        let p = FakeProvider::new();
        p.complete(ask(None)).await.unwrap();
        p.complete(ask(Some(&schema))).await.unwrap();

        let asked = p.asked();
        assert_eq!(asked.len(), 2);
        assert_eq!(asked[0].model, model());
        assert_eq!(asked[0].schema, None);
        assert_eq!(asked[1].schema, Some(schema));
    }
}