agentplane 0.23.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
//! What an HTTP failure means for a model call.
//!
//! Shared by every driver that speaks HTTP, because the mapping is **doctrine
//! rather than vendor detail** and two copies of it would drift. The question a
//! driver has to answer is not "what did the provider say" but:
//!
//! 1. did it reach them, and
//! 2. did it cost anything.
//!
//! Those two answers decide whether the runtime may ask again and whether the
//! budget ceiling is telling the truth. Everything else — field names, envelope
//! shapes, which key holds the token counts — is per-provider and stays in the
//! driver.

use super::{ModelError, ModelId};

/// Classify a non-success HTTP status.
///
/// Every driver routes its non-2xx responses through here, so the rules live in
/// one place:
///
/// * **429 and 529** are rate limiting. Separate from an ordinary refusal
///   because the response is different: this one is worth retrying, and it is
///   the one case where retrying is unambiguously safe *and* free. The
///   provider's `Retry-After` rides along, because it is the only number that
///   makes the retry useful — see [`ModelError::RateLimited`].
/// * **408 and 425** are the transient 4xx: a request the *server* timed out
///   or declined to process early, not one it judged wrong. Classed with the
///   retryable failures, because `Refused` means *repeating is pointless* and
///   these are the two 4xx codes for which repeating is the documented remedy.
/// * **every other 4xx** is a refusal before generating — bad request, unknown
///   model, bad key, content filtered on the way in. Nothing was metered, and
///   repeating is pointless rather than merely unsafe: the retry loop spends
///   no attempt on it.
/// * **anything else** reached the provider and did not say what it cost. See
///   [`ModelError::Unavailable`]: guessing "free" lets a retry loop spend
///   against a ceiling reading zero, and guessing "fatal" makes a transient blip
///   end a run.
pub fn classify_status(
    model: &ModelId,
    status: u16,
    headers: &reqwest::header::HeaderMap,
    body: &str,
) -> ModelError {
    let detail = format!("HTTP {status}: {}", trim(body));
    match status {
        429 | 529 => ModelError::RateLimited {
            model: model.clone(),
            detail,
            retry_after: retry_after(headers),
        },
        408 | 425 => ModelError::Unavailable {
            model: model.clone(),
            detail,
        },
        400..=499 => ModelError::Refused {
            model: model.clone(),
            detail,
        },
        _ => ModelError::Unavailable {
            model: model.clone(),
            detail,
        },
    }
}

/// The provider's `Retry-After`, in seconds, when it named one.
///
/// The parsing rule is [`core::retry_after_seconds`](crate::core::retry_after_seconds),
/// shared with every other wire this crate reads advice on: delta-seconds only,
/// because the HTTP-date form means trusting somebody else's clock against ours.
fn retry_after(headers: &reqwest::header::HeaderMap) -> Option<u64> {
    crate::core::retry_after_seconds(headers.get(reqwest::header::RETRY_AFTER)?.to_str().ok()?)
}

/// Classify a transport failure.
///
/// The distinction that matters is whether anything was written. A connection
/// that was never established sent nothing; one that failed later may have
/// delivered the request and generated an answer nobody will see.
pub fn classify_transport(model: &ModelId, e: &reqwest::Error) -> ModelError {
    if e.is_connect() {
        return ModelError::Unreachable {
            model: model.clone(),
            detail: format!("could not connect: {e}"),
        };
    }
    ModelError::Unavailable {
        model: model.clone(),
        detail: e.to_string(),
    }
}

/// Parse the answer when a schema was declared.
///
/// Provider constrained generation is the first line of defence. This also
/// parses and locally validates the answer against the exact requested schema,
/// so a provider bug or ignored constraint becomes a loud, metered `Unusable`
/// rather than invalid data reaching a later step. External schema reference
/// resolution is disabled, so validation cannot introduce hidden I/O.
///
/// # Errors
///
/// [`ModelError::Unusable`], carrying the usage — because a malformed answer was
/// still generated and still billed.
pub fn structured(
    schema: Option<&serde_json::Value>,
    text: &str,
    tool_calls: &[super::ToolCall],
    model: &ModelId,
    usage: super::Usage,
) -> Result<Option<serde_json::Value>, ModelError> {
    let Some(schema) = schema else {
        return Ok(None);
    };
    // A turn that asks for tools is not the final answer, and a schema binds
    // only the final answer — the same exemption `honour_declared_schema`
    // applies at the effect boundary. Spelled once here rather than once per
    // driver, because the copy that drifts is on whichever driver a
    // deployment does not exercise: failing a tool-asking turn does worse
    // than waste it — the error path carries no continuation, so a provider's
    // signed reasoning blocks are dropped from the retry, which the provider
    // then rejects. Emulated forced-tool answers pass an empty slice: there
    // the tool call *is* the answer, and its arguments must parse.
    if !tool_calls.is_empty() {
        return Ok(None);
    }
    let value: serde_json::Value =
        serde_json::from_str(text).map_err(|e| ModelError::Unusable {
            model: model.clone(),
            usage,
            detail: format!("a schema was required and the answer is not JSON: {e}"),
        })?;
    super::validate_schema(schema, &value).map_err(|detail| ModelError::Unusable {
        model: model.clone(),
        usage,
        detail,
    })?;
    Ok(Some(value))
}

/// The name of the single tool used when emulating structured output.
///
/// Fixed rather than caller-chosen: it goes into the request, and a name that
/// varied per call would change the request bytes without changing the
/// question, which is noise in anything that diffs them.
pub const RESPOND_TOOL: &str = "agentplane_respond";

/// Why a schema cannot be used with strict constrained decoding, if it cannot.
///
/// `OpenAI`'s strict mode accepts a **subset** of JSON Schema, and a schema that
/// is perfectly valid elsewhere is rejected with a 400 that does not say which
/// rule it broke. Checking here turns that into a refusal naming the exact
/// problem, before anything is sent and before anything is billed.
///
/// Deliberately **not** auto-corrected. Rewriting the caller's schema would mean
/// the effect key records one shape and the wire carries another — and a run
/// whose journal disagrees with what it asked for is exactly the class of quiet
/// divergence this crate exists to prevent. The caller fixes the schema.
pub fn strict_schema_problem(schema: &serde_json::Value) -> Option<String> {
    fn walk(node: &serde_json::Value, path: &str, out: &mut Vec<String>) {
        let Some(obj) = node.as_object() else { return };

        if obj.contains_key("default") {
            out.push(format!(
                "`{path}` uses `default`, which strict mode rejects"
            ));
        }

        if obj.get("type").and_then(|t| t.as_str()) == Some("object") {
            if obj.get("additionalProperties") != Some(&serde_json::Value::Bool(false)) {
                out.push(format!(
                    "`{path}` is an object without `additionalProperties: false`"
                ));
            }
            let properties = obj.get("properties").and_then(|p| p.as_object());
            if let Some(properties) = properties {
                let required: Vec<&str> = obj
                    .get("required")
                    .and_then(|r| r.as_array())
                    .map(|r| r.iter().filter_map(|v| v.as_str()).collect())
                    .unwrap_or_default();
                for key in properties.keys() {
                    if !required.contains(&key.as_str()) {
                        out.push(format!(
                            "`{path}.{key}` is optional; strict mode requires every \
                             property to be listed in `required`"
                        ));
                    }
                }
            }
        }

        for (key, child) in obj {
            match key.as_str() {
                "properties" | "$defs" | "definitions" => {
                    if let Some(map) = child.as_object() {
                        for (name, sub) in map {
                            walk(sub, &format!("{path}.{name}"), out);
                        }
                    }
                }
                "items" | "not" => walk(child, &format!("{path}.{key}"), out),
                "anyOf" | "oneOf" | "allOf" => {
                    if let Some(list) = child.as_array() {
                        for (i, sub) in list.iter().enumerate() {
                            walk(sub, &format!("{path}.{key}[{i}]"), out);
                        }
                    }
                }
                _ => {}
            }
        }
    }

    let mut problems = Vec::new();
    walk(schema, "schema", &mut problems);
    if problems.is_empty() {
        return None;
    }
    Some(problems.join("; "))
}

/// Keep an error body short enough to log.
///
/// A provider's error payload can carry the echoed prompt, and a prompt can
/// carry whatever the run was working on. Truncating here keeps a failure from
/// becoming an exfiltration channel into the operator's log aggregator.
fn trim(body: &str) -> String {
    const LIMIT: usize = 400;
    if body.len() <= LIMIT {
        return body.to_owned();
    }
    let mut cut = LIMIT;
    while cut > 0 && !body.is_char_boundary(cut) {
        cut -= 1;
    }
    format!("{}… ({} bytes)", &body[..cut], body.len())
}

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

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

    fn no_headers() -> reqwest::header::HeaderMap {
        reqwest::header::HeaderMap::new()
    }

    fn advising(value: &str) -> reqwest::header::HeaderMap {
        let mut headers = reqwest::header::HeaderMap::new();
        headers.insert(
            reqwest::header::RETRY_AFTER,
            reqwest::header::HeaderValue::from_str(value).expect("a header value"),
        );
        headers
    }

    /// The window is the whole point of the classification: without it the
    /// retry loop computes a schedule in hundreds of milliseconds against a
    /// limit measured in tens of seconds, spends every permitted attempt
    /// inside the window, and reports the provider as down.
    #[test]
    fn a_named_rate_limit_window_survives_classification() {
        let e = classify_status(&model(), 429, &advising("42"), "");
        assert!(
            matches!(
                e,
                ModelError::RateLimited {
                    retry_after: Some(42),
                    ..
                }
            ),
            "the provider named its window and the classification dropped it: {e}"
        );
    }

    /// A provider that throttles without saying when to come back is ordinary,
    /// and the effect's own schedule applies. What must not happen is a
    /// fabricated window: `None` is *no advice*, not zero seconds.
    #[test]
    fn an_unnamed_window_is_absent_rather_than_invented() {
        for value in [
            "",
            "  ",
            "0",
            "later",
            "-5",
            "Wed, 21 Oct 2026 07:28:00 GMT",
        ] {
            let e = classify_status(&model(), 429, &advising(value), "");
            assert!(
                matches!(
                    e,
                    ModelError::RateLimited {
                        retry_after: None,
                        ..
                    }
                ),
                "'{value}' is not advice this crate can act on, and reading it as \
                 one would replace a real backoff with a made-up schedule: {e}"
            );
        }
        assert!(matches!(
            classify_status(&model(), 429, &no_headers(), ""),
            ModelError::RateLimited {
                retry_after: None,
                ..
            }
        ));
    }

    #[test]
    fn rate_limiting_is_told_apart_from_refusal() {
        for s in [429u16, 529] {
            assert!(matches!(
                classify_status(&model(), s, &no_headers(), ""),
                ModelError::RateLimited { .. }
            ));
        }
    }

    #[test]
    fn a_client_error_did_not_generate() {
        for s in [400u16, 401, 403, 404, 422] {
            let e = classify_status(&model(), s, &no_headers(), "");
            assert_eq!(e.disposition(), Disposition::DidNotHappen);
            assert_eq!(e.usage().spend().tokens, 0);
            assert!(
                matches!(e, ModelError::Refused { .. }),
                "HTTP {s} is a judgement about the request, and repeating a \
                 judged request asks the same rule the same question"
            );
        }
    }

    /// 408 and 425 are the transient 4xx: the server timed out or declined to
    /// process *early*, not judged the request wrong. Classing them as
    /// `Refused` would make a hiccup terminal — the retry loop spends no
    /// attempt on a refusal, and these are the two 4xx codes whose documented
    /// remedy is the retry.
    #[test]
    fn the_transient_4xx_are_not_judgements() {
        for s in [408u16, 425] {
            let e = classify_status(&model(), s, &no_headers(), "");
            assert_eq!(e.disposition(), Disposition::DidNotHappen);
            assert!(
                matches!(e, ModelError::Unavailable { .. }),
                "HTTP {s} is transient and must stay retryable, got: {e}"
            );
        }
    }

    #[test]
    fn a_server_error_says_it_does_not_know() {
        assert!(matches!(
            classify_status(&model(), 500, &no_headers(), ""),
            ModelError::Unavailable { .. }
        ));
    }

    /// A provider's error body can echo the prompt back.
    #[test]
    fn a_long_error_body_is_trimmed() {
        let secret = "x".repeat(5_000);
        let e = classify_status(&model(), 400, &no_headers(), &secret);
        let rendered = e.to_string();
        assert!(
            rendered.len() < 600,
            "an error body went into the log at full length ({} chars), and a \
             provider echoes the prompt back in it",
            rendered.len()
        );
        assert!(rendered.contains("5000 bytes"), "{rendered}");
    }

    /// Multi-byte characters must not be cut through.
    #[test]
    fn trimming_respects_character_boundaries() {
        let body = "ü".repeat(1_000);
        let _ = trim(&body);
    }
}

#[cfg(test)]
mod schema_validation_tests {
    use serde_json::json;

    use super::*;

    #[test]
    fn parseable_but_nonconforming_structured_output_is_unusable() {
        let model = ModelId::new("test", "structured");
        let usage = super::super::Usage {
            output_tokens: 5,
            ..Default::default()
        };
        let error = structured(
            Some(&json!({
                "type": "object",
                "properties": {"id": {"type": "string", "minLength": 5}},
                "required": ["id"]
            })),
            r#"{"id":"abc"}"#,
            &[],
            &model,
            usage,
        )
        .expect_err("provider-constrained output still needs defense-in-depth validation");
        assert!(matches!(error, ModelError::Unusable { usage: u, .. } if u == usage));
    }
}