Skip to main content

agentplane/model/
mod.rs

1//! Calling a model.
2//!
3//! A completion is an effect like any other — journaled once, replayed from the
4//! record, untrusted on the way back. What is different is the meter.
5//!
6//! # A failed completion is not a free one
7//!
8//! Every other outward call this crate makes either happens or does not. A model
9//! call has a third state: it ran, generated four hundred tokens, and then the
10//! stream died. The provider will bill those tokens. The answer is unusable.
11//!
12//! Two consequences, and both are easy to get backwards.
13//!
14//! **The spend must be reported anyway.** [`EffectError::Metered`] carries what
15//! was consumed, and the runtime bills it on the failure path. Without that, the
16//! token and cost ceilings — the ones that exist to bound exactly this — count
17//! zero while a retry loop against a flaky provider spends real money.
18//!
19//! **A died-mid-stream call is [`Disposition::Landed`], not `InDoubt`.** The
20//! usual reasoning about reaching the peer is inverted here: we know perfectly
21//! well that it reached the provider, because we watched it generate. What we
22//! lack is the *answer*, and repeating the call buys a second bill for the same
23//! question. `InDoubt` would invite [`Recovery`] to resolve an outcome that is
24//! not in doubt at all.
25//!
26//! # Determinism
27//!
28//! A model is the least deterministic thing a run touches, which is exactly why
29//! the completion is journaled: replay reads the recorded answer rather than
30//! asking again. The prompt is part of the effect key, so a changed prompt is a
31//! changed effect and shows up as divergence rather than as a quietly different
32//! run.
33
34#[cfg(feature = "providers")]
35pub mod anthropic;
36#[cfg(feature = "providers")]
37mod anthropic_stream;
38#[cfg(feature = "bedrock")]
39pub mod bedrock;
40#[cfg(feature = "bedrock")]
41mod bedrock_stream;
42#[cfg(feature = "providers")]
43pub mod chat_completions;
44#[cfg(feature = "providers")]
45mod chat_completions_stream;
46/// The embeddings wire, beside the drivers it shares a transport with.
47// `any(..)`, not `providers`: `BedrockEmbedder` lives in here and is the driver
48// a plane whose data may not leave one AWS account needs, so gating the module
49// on `providers` made the `bedrock` feature pay for the AWS SDK and expose no
50// embedder at all — the "a feature that builds is not a feature that delivers"
51// shape the store gate already had. A `const _` in `lib.rs` names each
52// embedder's type so a gate that configures one out fails *this* crate's build.
53#[cfg(any(feature = "providers", feature = "bedrock"))]
54pub mod embeddings;
55#[cfg(feature = "providers")]
56pub mod gemini;
57#[cfg(feature = "providers")]
58mod gemini_stream;
59#[cfg(feature = "providers")]
60pub mod openai;
61#[cfg(feature = "providers")]
62mod openai_stream;
63#[cfg(feature = "providers")]
64mod sse;
65#[cfg(feature = "providers")]
66mod wire;
67
68use std::fmt::Debug;
69use std::sync::Arc;
70
71use async_trait::async_trait;
72#[cfg(feature = "media")]
73use base64::Engine as _;
74use serde::{Deserialize, Serialize};
75use serde_json::Value;
76
77use crate::core::{
78    Disposition, Effect, EffectDescriptor, EffectError, Recovery, RetryPolicy, Sensitivity, Spend,
79    Trust,
80};
81
82#[cfg(any(feature = "manifest", feature = "providers", feature = "bedrock"))]
83pub(crate) fn validate_schema(schema: &Value, value: &Value) -> Result<(), String> {
84    let validator = jsonschema::validator_for(schema)
85        .map_err(|error| format!("the declared JSON Schema is invalid: {error}"))?;
86    validator
87        .validate(value)
88        .map_err(|error| format!("value does not satisfy the declared JSON Schema: {error}"))
89}
90
91/// A provider-side media reference hidden inside a provider-native prompt.
92///
93/// Deliberately structural rather than a search for strings that look like
94/// URLs. A user may ask a model to discuss a URL; the dangerous forms are the
95/// content blocks that instruct the provider to dereference one. The built-in
96/// drivers accept provider-native JSON, so both providers' spellings are
97/// recognized here and the runtime applies the same hard cut to custom drivers.
98fn provider_side_media_reference(value: &Value) -> Option<&'static str> {
99    fn remote_url(value: Option<&Value>) -> bool {
100        let url = value.and_then(|value| {
101            value
102                .as_str()
103                .or_else(|| value.get("url").and_then(Value::as_str))
104        });
105        url.is_some_and(|url| !url.starts_with("data:"))
106    }
107
108    match value {
109        Value::Array(values) => values.iter().find_map(provider_side_media_reference),
110        Value::Object(object) => {
111            let kind = object.get("type").and_then(Value::as_str);
112
113            // Anthropic Messages: image/document source { type: "url", url: ... }.
114            if matches!(kind, Some("image" | "document"))
115                && object
116                    .get("source")
117                    .and_then(Value::as_object)
118                    .is_some_and(|source| {
119                        source.get("type").and_then(Value::as_str) == Some("url")
120                            && source.get("url").and_then(Value::as_str).is_some()
121                    })
122            {
123                return Some("an Anthropic image/document URL source");
124            }
125
126            // OpenAI Responses, plus the older image_url content-block spelling
127            // accepted by compatible endpoints. A data URL carries bytes in the
128            // request and is not a provider-side network fetch.
129            if matches!(kind, Some("input_image" | "image_url"))
130                && remote_url(object.get("image_url"))
131            {
132                return Some("an OpenAI image URL");
133            }
134            if kind == Some("input_file") && remote_url(object.get("file_url")) {
135                return Some("an OpenAI file URL");
136            }
137
138            // Gemini: `fileData { fileUri }`, the form that tells Google to
139            // fetch the bytes itself — a Files API URI or a plain remote URL,
140            // and both are a fetch from the provider's network rather than
141            // this plane's. Google's REST surface accepts camelCase and
142            // snake_case interchangeably, so both spellings are checked: a
143            // control that only knows one of two accepted spellings is one an
144            // author bypasses by writing the other, without meaning to.
145            for key in ["fileData", "file_data"] {
146                if let Some(file) = object.get(key).and_then(Value::as_object)
147                    && (remote_url(file.get("fileUri")) || remote_url(file.get("file_uri")))
148                {
149                    return Some("a Gemini fileData URI");
150                }
151            }
152
153            object.values().find_map(provider_side_media_reference)
154        }
155        _ => None,
156    }
157}
158
159fn provider_side_media_refusal(kind: &str) -> String {
160    format!(
161        "{kind} was refused before dispatch: the model provider would fetch it outside \
162         this plane's egress policy and journal; inline the media bytes, or fetch them \
163         through an explicit governed effect first"
164    )
165}
166
167pub(crate) fn refuse_provider_side_media(
168    prompt: &Value,
169    model: &ModelId,
170) -> Result<(), ModelError> {
171    let Some(kind) = provider_side_media_reference(prompt) else {
172        return Ok(());
173    };
174    Err(ModelError::Refused {
175        model: model.clone(),
176        detail: provider_side_media_refusal(kind),
177    })
178}
179
180/// Hold a completion to the schema its request declared.
181///
182/// The built-in drivers already validate: [`wire::structured`] parses and checks
183/// every answer they return. But the [`ModelProvider`] trait is public, and the
184/// runtime treats `structured` as *guaranteed to match the schema* — it indexes
185/// into the value it asked for rather than re-checking a contract it believes is
186/// already held. A provider this crate cannot inspect makes that belief a
187/// third party's promise.
188///
189/// So the check is repeated at the effect boundary, exactly as
190/// [`refuse_provider_side_media`] is: a control implemented once per driver is
191/// one a driver written elsewhere does not have. The cost is one validation of a
192/// value that is almost always already valid; the alternative is that a
193/// malformed answer reaches code with no way to reject it.
194///
195/// A completion carrying **tool calls** is exempt: choosing a tool is a
196/// legitimate answer to a schema-bearing request, and the schema governs the
197/// turn that finally answers.
198///
199/// Failure is [`ModelError::Unusable`] — **metered**, because the tokens were
200/// generated and the provider will bill for them however unusable the answer is.
201fn honour_declared_schema(
202    completion: &Completion,
203    schema: Option<&Value>,
204    model: &ModelId,
205) -> Result<(), ModelError> {
206    let Some(schema) = schema else {
207        return Ok(());
208    };
209    if !completion.tool_calls.is_empty() {
210        return Ok(());
211    }
212    let Some(value) = completion.structured.as_ref() else {
213        return Err(ModelError::Unusable {
214            model: model.clone(),
215            usage: completion.usage,
216            detail: "a schema was declared and the provider returned no structured value; \
217                     a driver must parse its own answer into `Completion::structured`, or \
218                     report the call unusable itself"
219                .to_owned(),
220        });
221    };
222    #[cfg(any(feature = "manifest", feature = "providers", feature = "bedrock"))]
223    validate_schema(schema, value).map_err(|detail| ModelError::Unusable {
224        model: model.clone(),
225        usage: completion.usage,
226        detail,
227    })?;
228    // Without `jsonschema` in the build there is no validator to run. The
229    // presence check above still holds, and it is the half that keeps the
230    // runtime from indexing into an absent value. Named rather than silenced,
231    // so a reader can see which half is missing and why.
232    #[cfg(not(any(feature = "manifest", feature = "providers", feature = "bedrock")))]
233    let _ = (schema, value);
234    Ok(())
235}
236
237/// Which model, from which provider.
238#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
239pub struct ModelId {
240    pub provider: String,
241    pub model: String,
242}
243
244impl ModelId {
245    pub fn new(provider: impl Into<String>, model: impl Into<String>) -> Self {
246        Self {
247            provider: provider.into(),
248            model: model.into(),
249        }
250    }
251}
252
253impl std::fmt::Display for ModelId {
254    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
255        write!(f, "{}/{}", self.provider, self.model)
256    }
257}
258
259/// What a completion cost.
260///
261/// # Cached tokens are the trap
262///
263/// Prompt caching is the feature most likely to make a token ceiling lie, and it
264/// lies in *both* directions depending on the provider:
265///
266/// * **Anthropic** reports `cache_creation_input_tokens` and
267///   `cache_read_input_tokens` **alongside** `input_tokens`, which excludes
268///   them. A driver that reads only `input_tokens` bills a cached call at close
269///   to nothing while the provider charges a premium for the write and a tenth
270///   of the rate for the read.
271/// * **`OpenAI`** reports `input_tokens_details.cached_tokens` as a **subset** of
272///   `input_tokens`. Adding it would double-count.
273///
274/// Same words, opposite arithmetic. So this type keeps the cached counts in
275/// their own fields, `input_tokens` always means *everything sent*, and each
276/// driver is responsible for normalising into that. The alternative — a bare
277/// `input_tokens` each driver fills differently — is a budget whose meaning
278/// depends on which provider a run happened to use.
279#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
280pub struct Usage {
281    /// Every input token the provider processed, cached ones included.
282    pub input_tokens: u64,
283    pub output_tokens: u64,
284    /// Of `input_tokens`, how many were written into a cache.
285    ///
286    /// Billed at a premium over ordinary input. Reported separately because the
287    /// *rate* differs, and a deployment pricing its own runs needs the split.
288    #[serde(default)]
289    pub cache_write_tokens: u64,
290    /// Of `input_tokens`, how many were served from a cache.
291    ///
292    /// Billed at roughly a tenth of the input rate. A run that reported these as
293    /// ordinary input would over-state its cost by an order of magnitude on the
294    /// cached portion — which is the opposite failure to omitting them, and just
295    /// as wrong.
296    #[serde(default)]
297    pub cache_read_tokens: u64,
298    /// Money in minor units, if the provider reports it.
299    ///
300    /// Priced by the driver rather than derived here: rates change, differ per
301    /// model, and are a deployment's contract with its provider, not this
302    /// crate's guess.
303    pub minor_units: u64,
304}
305
306impl Usage {
307    /// What this counts against a run's ceilings.
308    ///
309    /// Tokens are summed flat — input plus output — because a *ceiling* is about
310    /// bounding how much work a run may cause, and a cached input token is still
311    /// a token the provider processed. Cost weighting belongs in `minor_units`,
312    /// which the driver prices; conflating the two would make the token ceiling
313    /// mean something different for every provider.
314    #[must_use]
315    pub const fn spend(&self) -> Spend {
316        Spend {
317            tokens: self.input_tokens + self.output_tokens,
318            minor_units: self.minor_units,
319        }
320    }
321
322    /// Input tokens that were neither written to nor read from a cache.
323    #[must_use]
324    pub const fn uncached_input_tokens(&self) -> u64 {
325        self.input_tokens
326            .saturating_sub(self.cache_write_tokens)
327            .saturating_sub(self.cache_read_tokens)
328    }
329}
330
331/// A tool the model asked to call.
332///
333/// A **request, never an instruction**. Each one still has to pass the gate: the
334/// agent's manifest must grant the tool, policy must allow the call, and the
335/// budget must have room. Model output is a proposal, and this is the most
336/// literal case of that rule.
337#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
338pub struct ToolCall {
339    /// The provider's own identifier for this call.
340    ///
341    /// Load-bearing in a loop: each result must go back under the id the model
342    /// used, or the model cannot tell which answer belongs to which question —
343    /// and providers reject a result carrying an id they never issued.
344    pub id: String,
345    /// The tool's name, as the model wrote it.
346    ///
347    /// Untrusted like everything else a model emits. A name matching no grant is
348    /// refused, never resolved to a near neighbour.
349    pub name: String,
350    /// The arguments, decoded.
351    ///
352    /// Normalised across providers — Anthropic sends an object, `OpenAI` a JSON
353    /// string — so a caller need not know which driver answered in order to read
354    /// them.
355    pub arguments: Value,
356}
357
358/// What came back.
359#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
360pub struct Completion {
361    pub text: String,
362    /// Tools the model asked to call.
363    ///
364    /// Empty for an ordinary answer, and empty for the forced-tool path used to
365    /// obtain structured output: that tool is this crate's mechanism for "answer
366    /// in this shape", not a request for the runtime to *do* anything, and
367    /// surfacing it would make every schema-shaped completion look like a tool
368    /// invocation.
369    ///
370    /// Filled identically by the buffered and streaming paths. Streaming is the
371    /// default, so a field populated only when buffering would be silently empty
372    /// in most deployments — which is worse than absent, because callers would
373    /// build loops on it and see them never fire.
374    #[serde(default, skip_serializing_if = "Vec::is_empty")]
375    pub tool_calls: Vec<ToolCall>,
376    pub usage: Usage,
377    /// Why generation stopped, in the provider's words.
378    ///
379    /// Passed through unnormalised on purpose: `end_turn`, `max_tokens`,
380    /// `incomplete:max_output_tokens` and `stop` mean subtly different things to
381    /// the providers that emit them, and flattening them into a shared
382    /// vocabulary would lose exactly the detail a caller debugging a truncated
383    /// answer needs. The one thing a caller must not have to *parse* out of it
384    /// is whether the answer is complete — see [`truncated`](Self::truncated).
385    pub stop_reason: Option<String>,
386    /// Whether the answer was cut short.
387    ///
388    /// A typed field rather than a string a caller has to recognise, and its own
389    /// field rather than an error, for the same reason the worklist reports
390    /// `truncated` beside its page: **a partial answer returned as a whole one is
391    /// a silent truncation**, which this crate refuses everywhere else (P7).
392    ///
393    /// It is not an error because a cut-off answer is often still useful — prose
394    /// that stops early is readable, and the caller is the only one who knows
395    /// whether they were parsing JSON. What they must not be able to do is
396    /// *overlook* it, and a `bool` in the struct they already destructure is
397    /// harder to overlook than a stop reason they have to compare against a
398    /// provider-specific string.
399    pub truncated: bool,
400    /// The answer parsed as JSON, when a schema was asked for.
401    ///
402    /// `None` when no schema was declared. When one was, this is the parsed
403    /// value and [`text`](Self::text) still holds the raw string.
404    ///
405    /// Provider constrained decoding prevents malformed output before tokens
406    /// are emitted; this crate then validates the parsed value locally as
407    /// defense in depth. Provider bugs and forced-tool best-effort behavior are
408    /// therefore loud, metered `Unusable` responses rather than malformed data
409    /// reaching downstream code. External schema references are not resolved:
410    /// validation performs no hidden file or network I/O.
411    #[serde(default, skip_serializing_if = "Option::is_none")]
412    pub structured: Option<Value>,
413    /// Provider-owned response items required to continue this exact turn.
414    ///
415    /// This is deliberately opaque to the runtime. `OpenAI` uses complete
416    /// Responses output items (including encrypted reasoning); Anthropic uses
417    /// the complete assistant content blocks (including signed thinking). The
418    /// next request returns the value only to the provider that issued it.
419    /// Keeping it in the journal makes continuation independent of expiring
420    /// provider-side conversation state and reproducible on replay.
421    #[serde(default, skip_serializing_if = "Option::is_none")]
422    pub continuation: Option<ProviderContinuation>,
423}
424
425/// Opaque, self-contained provider state for one continuation turn.
426#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
427pub struct ProviderContinuation {
428    /// Provider name that owns [`state`](Self::state).
429    pub provider: String,
430    /// Exact provider-native items emitted by the preceding response.
431    pub state: Value,
432}
433
434/// A live, non-durable model-stream event.
435#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
436#[serde(rename_all = "snake_case", tag = "type", content = "value")]
437pub enum ModelStreamEvent {
438    /// Visible answer text only. Opaque reasoning is never exposed here.
439    TextDelta(String),
440    /// The latest provider-reported usage snapshot.
441    Usage(Usage),
442}
443
444/// Receives live model progress while one terminal completion remains canonical.
445pub trait ModelStreamObserver: Send + Sync + Debug {
446    /// Delivery is advisory and must not block provider consumption. A caller
447    /// needing network backpressure should enqueue into its own bounded channel.
448    fn event(&self, event: crate::core::Tainted<ModelStreamEvent>);
449}
450
451impl ProviderContinuation {
452    #[must_use]
453    pub fn new(provider: impl Into<String>, state: Value) -> Self {
454        Self {
455            provider: provider.into(),
456            state,
457        }
458    }
459}
460
461/// Why a completion failed.
462#[derive(Debug, thiserror::Error)]
463pub enum ModelError {
464    /// Never reached the provider.
465    #[error("could not reach '{model}': {detail}")]
466    Unreachable { model: ModelId, detail: String },
467
468    /// The provider refused before generating: bad request, unknown model, a
469    /// content filter on the *input*. Nothing was metered.
470    #[error("'{model}' refused the request: {detail}")]
471    Refused { model: ModelId, detail: String },
472
473    /// Rate-limited before generating.
474    ///
475    /// Separate from [`Refused`](ModelError::Refused) because the response is
476    /// different: this one is worth retrying, and it is the one case here where
477    /// retrying is unambiguously safe.
478    #[error("'{model}' is rate limiting: {detail}")]
479    RateLimited { model: ModelId, detail: String },
480
481    /// It generated, and then the stream died.
482    ///
483    /// The expensive case. The tokens counted here have been spent whatever
484    /// happens next.
485    #[error("'{model}' stopped mid-response after {} token(s): {detail}", usage.input_tokens + usage.output_tokens)]
486    Interrupted {
487        model: ModelId,
488        usage: Usage,
489        detail: String,
490    },
491
492    /// It reached the provider, and nothing came back that says whether it
493    /// generated.
494    ///
495    /// A non-streaming 5xx, or a response that could not be read. The honest
496    /// position is that this is *unknowable* from here, and both guesses are
497    /// wrong in a different way: calling it `Interrupted` makes a transient blip
498    /// fatal, and calling it free lets a retry loop spend real money against a
499    /// ceiling that reads zero.
500    ///
501    /// Treated as safe to repeat, because a completion does not change the
502    /// world — so repeating is a correctness no-op and only a cost. The
503    /// documented price is that the spend ceiling may under-count by at most one
504    /// call per occurrence.
505    ///
506    /// A driver that *can* see partial usage must report
507    /// [`Interrupted`](ModelError::Interrupted) instead — which is what both
508    /// shipped drivers do when streaming, and why they stream by default. Where
509    /// the provider makes even that impossible, the answer is
510    /// [`Unaccounted`](ModelError::Unaccounted), not this.
511    #[error("'{model}' did not say whether it generated: {detail}")]
512    Unavailable { model: ModelId, detail: String },
513
514    /// It generated, the stream died, and the cost is unknowable.
515    ///
516    /// The state `OpenAI`'s Responses stream can produce and Anthropic's cannot.
517    /// Usage appears there only in the terminal event, so a connection cut after
518    /// four hundred tokens of deltas leaves the driver *certain* that generation
519    /// happened and *ignorant* of what it cost.
520    ///
521    /// Neither neighbour says that, which is why this variant exists rather than
522    /// being folded into one of them:
523    ///
524    /// * [`Unavailable`](ModelError::Unavailable) means it may never have
525    ///   generated, and is therefore safe to repeat. Here we watched it generate;
526    ///   asking again buys a second bill for the same question.
527    /// * [`Interrupted`](ModelError::Interrupted) carries a [`Usage`], and
528    ///   filling it with zeroes is the "guess free" failure this crate refuses
529    ///   everywhere else — it reads as *this cost nothing* rather than as
530    ///   *nobody knows*.
531    ///
532    /// So it is [`Disposition::Landed`] with no usage, and the under-count is
533    /// admitted rather than hidden: the budget will be short by whatever this
534    /// call generated. What the variant buys is that the runtime stops paying
535    /// **twice** for it. A caller who needs the true figure has the provider's
536    /// response id and a [`Recovery`] policy to reconcile with; a driver quietly
537    /// making a second unjournaled request to find out is not the answer.
538    #[error("'{model}' generated and then died without saying what it cost: {detail}")]
539    Unaccounted { model: ModelId, detail: String },
540
541    /// It answered, and the answer was not usable — truncated JSON, a refusal
542    /// where a tool call was required. Metered, because it generated.
543    #[error("'{model}' returned an unusable answer: {detail}")]
544    Unusable {
545        model: ModelId,
546        usage: Usage,
547        detail: String,
548    },
549}
550
551impl ModelError {
552    /// What this failure says about whether the call reached the provider.
553    #[must_use]
554    pub const fn disposition(&self) -> Disposition {
555        match self {
556            Self::Unreachable { .. }
557            | Self::Refused { .. }
558            | Self::RateLimited { .. }
559            // Safe to repeat despite having reached the provider: a completion
560            // is the one outward call here that does not change the world, so
561            // the only cost of asking again is money — which the budget bounds.
562            | Self::Unavailable { .. } => Disposition::DidNotHappen,
563            // We watched it generate. There is nothing in doubt: it happened,
564            // it was billed, and repeating it buys a second bill for the same
565            // question. `Unaccounted` belongs here for exactly that reason and
566            // despite reporting no usage — what is unknown is the *amount*, not
567            // whether it happened.
568            Self::Interrupted { .. } | Self::Unusable { .. } | Self::Unaccounted { .. } => {
569                Disposition::Landed
570            }
571        }
572    }
573
574    /// What was consumed before the failure.
575    #[must_use]
576    pub const fn usage(&self) -> Usage {
577        match self {
578            Self::Interrupted { usage, .. } | Self::Unusable { usage, .. } => *usage,
579            _ => Usage {
580                input_tokens: 0,
581                output_tokens: 0,
582                cache_write_tokens: 0,
583                cache_read_tokens: 0,
584                minor_units: 0,
585            },
586        }
587    }
588}
589
590/// One model role, as a declaration resolves it: which model, and the
591/// ceilings the reviewer put beside it.
592///
593/// A manifest declares `max_tokens` and `reasoning_effort` **per role**
594/// because the model is only half the decision — the role that reads hostile
595/// content is the one whose output most needs a ceiling somebody reviewed.
596/// Passing the id alone is how those two fields get parsed into the digest and
597/// then dropped, so the runtime's seams carry the whole role and
598/// [`applied_to`](Self::applied_to) puts it on a call in one motion.
599#[derive(Debug, Clone, PartialEq, Eq)]
600pub struct ModelRole {
601    pub model: ModelId,
602    /// Cap on generated tokens, when the declaration set one.
603    pub max_output_tokens: Option<u32>,
604    /// Requested reasoning depth, when the declaration set one.
605    pub reasoning_effort: Option<ReasoningEffort>,
606}
607
608impl ModelRole {
609    /// A role with no declared ceilings — the driver defaults apply.
610    #[must_use]
611    pub fn new(model: ModelId) -> Self {
612        Self {
613            model,
614            max_output_tokens: None,
615            reasoning_effort: None,
616        }
617    }
618
619    /// Put this role's declared ceilings on a call.
620    ///
621    /// The model itself is *not* applied here — a [`ModelCall`] is constructed
622    /// with its model, so applying it afterwards would be a second place the
623    /// choice is made.
624    #[must_use]
625    pub fn applied_to(&self, mut call: ModelCall) -> ModelCall {
626        if let Some(max_output_tokens) = self.max_output_tokens {
627            call = call.with_max_output_tokens(max_output_tokens);
628        }
629        if let Some(effort) = self.reasoning_effort {
630            call = call.with_reasoning_effort(effort);
631        }
632        call
633    }
634}
635
636/// One request to a provider.
637///
638/// A struct rather than a widening argument list, because what a model call
639/// carries is the part of this seam most likely to grow — and every growth would
640/// otherwise be a breaking change to every driver.
641#[derive(Debug, Clone)]
642pub struct Request<'a> {
643    pub model: &'a ModelId,
644    pub prompt: &'a Value,
645    /// Maximum tokens the provider may generate for this call.
646    ///
647    /// Provider-neutral because both shipped APIs expose the same control, and
648    /// per-call because a manifest declares it per model role. Keeping it on a
649    /// driver silently discarded that declaration and kept the real request
650    /// limit out of the effect key.
651    pub max_output_tokens: u32,
652    /// How much internal reasoning to request, when explicitly configured.
653    pub reasoning_effort: Option<ReasoningEffort>,
654    /// A JSON Schema the answer must conform to, if one was declared.
655    ///
656    /// Passed straight through to the provider's own structured-output mode —
657    /// `text.format` on `OpenAI` Responses, `output_config.format` on Anthropic —
658    /// where the constraint is *enforced during generation* rather than checked
659    /// afterwards. That is the whole reason to use it: a schema applied after
660    /// the fact rejects a bad answer you have already paid for.
661    pub schema: Option<&'a Value>,
662    /// The tools the model may ask for.
663    ///
664    /// Empty means the model is told of none, which is not the same as being
665    /// forbidden: authorization happens when a call comes back, against the
666    /// operator's grants. Declaring nothing simply gives it nothing to choose.
667    pub tools: &'a [ToolDeclaration],
668    /// Tools already run this turn, and what they returned.
669    pub exchanges: &'a [ToolExchange],
670    /// Exact provider-native state emitted beside those tool calls.
671    pub continuation: Option<&'a ProviderContinuation>,
672    /// Live observer. Not provider-visible and therefore not effect identity.
673    /// Strict replay never calls it because replay never performs the provider.
674    pub stream: Option<(&'a dyn ModelStreamObserver, &'a crate::core::Label)>,
675}
676
677/// Provider-neutral reasoning depth.
678///
679/// Providers and models support different subsets. An explicit unsupported
680/// value is refused before dispatch rather than silently downgraded.
681#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
682#[serde(rename_all = "kebab-case")]
683pub enum ReasoningEffort {
684    None,
685    Minimal,
686    Low,
687    Medium,
688    High,
689    XHigh,
690    Max,
691}
692
693impl ReasoningEffort {
694    #[must_use]
695    pub const fn as_str(self) -> &'static str {
696        match self {
697            Self::None => "none",
698            Self::Minimal => "minimal",
699            Self::Low => "low",
700            Self::Medium => "medium",
701            Self::High => "high",
702            Self::XHigh => "xhigh",
703            Self::Max => "max",
704        }
705    }
706}
707
708#[cfg(test)]
709mod reasoning_effort_tests {
710    use super::ReasoningEffort;
711
712    #[test]
713    fn every_reasoning_effort_has_a_pinned_wire_spelling() {
714        for (effort, wire) in [
715            (ReasoningEffort::None, "none"),
716            (ReasoningEffort::Minimal, "minimal"),
717            (ReasoningEffort::Low, "low"),
718            (ReasoningEffort::Medium, "medium"),
719            (ReasoningEffort::High, "high"),
720            (ReasoningEffort::XHigh, "xhigh"),
721            (ReasoningEffort::Max, "max"),
722        ] {
723            assert_eq!(effort.as_str(), wire);
724        }
725    }
726}
727
728/// How a driver should obtain a schema-conforming answer.
729///
730/// **Native structured output is not universally available**, and that is the
731/// whole reason this is a choice rather than an implementation detail. Anthropic
732/// gates grammar-constrained generation on particular models; `OpenAI`'s strict
733/// mode is only on newer ones, with older models offering a JSON *mode* that
734/// guarantees valid JSON and nothing about its shape.
735///
736/// Which mode a given model supports is a fact the deployment knows and the
737/// crate cannot discover — asking would be a network call on a path that must
738/// not make one, and guessing from a model-name pattern is a lookup table that
739/// is wrong the week a model ships.
740#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
741pub enum SchemaMode {
742    /// The provider's own constrained decoding.
743    ///
744    /// Strongest where it exists: the schema is enforced token by token during
745    /// generation, so a non-conforming answer is not merely rejected but
746    /// unproducible. The default, because a deployment that has not thought
747    /// about this should get the strong thing and a loud failure, not a silent
748    /// downgrade.
749    #[default]
750    Native,
751    /// A single forced tool whose input schema *is* the desired output schema.
752    ///
753    /// The universal fallback, and older than native support: define one tool,
754    /// force the model to call it, and read the arguments it was obliged to
755    /// construct. Works on any model that can call tools at all, which is a much
756    /// wider set than those with constrained decoding.
757    ///
758    /// Weaker in one specific way worth knowing: the model is constrained to
759    /// *produce a tool call*, and providers vary in how strictly they validate
760    /// its arguments against the declared schema. Native mode makes a malformed
761    /// answer impossible; this makes it unlikely.
762    ForcedTool,
763}
764
765/// Talks to a provider.
766#[async_trait]
767pub trait ModelProvider: Send + Sync + Debug {
768    /// Stable, non-secret configuration that changes the provider wire request.
769    ///
770    /// It is part of [`ModelCall`]'s effect identity. A provider switching from
771    /// native schema enforcement to a forced tool, changing endpoint/API
772    /// version, or changing streaming behavior must not replay an answer
773    /// produced under the old transport contract. API keys never belong here.
774    fn request_profile(&self, _model: &ModelId) -> Value {
775        Value::Null
776    }
777
778    /// Complete a prompt.
779    ///
780    /// # Errors
781    ///
782    /// A [`ModelError`] that states both what is known about reaching the
783    /// provider *and* what was consumed. A driver that reports zero usage for an
784    /// interrupted stream is telling the budget the call was free.
785    async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError>;
786}
787
788/// A tool the model may ask for, as the request declares it.
789///
790/// Provider-neutral, because the two shapes differ in ways that are easy to get
791/// subtly wrong: Anthropic takes `{name, description, input_schema}` at the top
792/// level, while `OpenAI` wraps it as `{type: "function", function: {name,
793/// description, parameters}}`. A driver renders this into whichever it speaks,
794/// so a caller writes the declaration once.
795///
796/// # What a declaration is *not*
797///
798/// It is not a grant. Declaring a tool tells the model the tool exists; it does
799/// not authorize the call. The model's choice of tool and its arguments come
800/// back **untrusted**, are matched against the operator's grants exactly, and
801/// are dispatched through `cx.sink` where field provenance and the egress
802/// ceiling apply. A framework that executes what the model asked for has
803/// authorized the model; this one authorizes the operator's declaration and
804/// treats the model's request as a suggestion.
805#[derive(Debug, Clone, PartialEq, Eq)]
806pub struct ToolDeclaration {
807    /// The name the model will use when it asks for this tool.
808    pub name: String,
809    /// What it does, for the model's benefit.
810    pub description: String,
811    /// JSON Schema for the arguments.
812    ///
813    /// Sent with provider-side strict enforcement when the provider accepts the
814    /// schema. `OpenAI` supports only a subset for strict tools, so a valid
815    /// schema with optional fields is sent non-strict rather than rejected by
816    /// the API; typed local tools still deserialize the result exactly. Either
817    /// way this is not a security control: a well-formed argument is still an
818    /// untrusted one, and the field-provenance check is what authorizes it.
819    pub parameters: Value,
820}
821
822impl ToolDeclaration {
823    /// Declare a tool.
824    #[must_use]
825    pub fn new(name: impl Into<String>, description: impl Into<String>, parameters: Value) -> Self {
826        Self {
827            name: name.into(),
828            description: description.into(),
829            parameters,
830        }
831    }
832}
833
834/// A tool the model asked for, and what came back.
835///
836/// Handed to the next request so the model can see the result of what it asked
837/// for. Provider-neutral because the continuation shapes differ more than the
838/// declarations do, and in ways that fail loudly at the API rather than quietly
839/// in the answer.
840///
841/// # Both halves travel, not just the result
842///
843/// The **call** is echoed back alongside its output. That is not redundancy: a
844/// provider matches a result to the request that produced it by id, and one sent
845/// without its call is rejected — `OpenAI` answers *"No tool call found for
846/// function call output with `call_id`"*. Carrying the pair makes that
847/// unrepresentable.
848///
849/// # Why the transcript is passed rather than referenced
850///
851/// `OpenAI` will hold the conversation for you behind `previous_response_id`.
852/// This crate does not use it, and will not: replay would then depend on state
853/// a provider holds, expires and can lose — so a run that replayed correctly
854/// today would diverge when that state aged out, for a reason nothing in the
855/// journal could explain. Everything needed to continue is in the request.
856#[derive(Debug, Clone, PartialEq)]
857pub struct ToolExchange {
858    /// What the model asked for, including the id it issued.
859    pub call: ToolCall,
860    /// What came back, as the tool produced it.
861    pub output: Value,
862    /// Whether the tool failed.
863    ///
864    /// Sent as Anthropic's `is_error`, so the model is told the difference
865    /// between a tool that answered and one that could not. A failure rendered
866    /// as an ordinary result teaches it that the operation succeeded and
867    /// returned something strange.
868    pub failed: bool,
869}
870
871impl ToolExchange {
872    /// A tool that answered.
873    #[must_use]
874    pub fn ok(call: ToolCall, output: Value) -> Self {
875        Self {
876            call,
877            output,
878            failed: false,
879        }
880    }
881
882    /// A tool that failed, with what to tell the model.
883    #[must_use]
884    pub fn failed(call: ToolCall, detail: impl Into<String>) -> Self {
885        Self {
886            call,
887            output: Value::String(detail.into()),
888            failed: true,
889        }
890    }
891}
892
893/// One completion.
894#[derive(Debug)]
895pub struct ModelCall {
896    model: ModelId,
897    prompt: Value,
898    schema: Option<Value>,
899    tools: Vec<ToolDeclaration>,
900    exchanges: Vec<ToolExchange>,
901    continuation: Option<ProviderContinuation>,
902    stream: Option<Arc<dyn ModelStreamObserver>>,
903    max_output_tokens: u32,
904    reasoning_effort: Option<ReasoningEffort>,
905    provider: Arc<dyn ModelProvider>,
906    max_sensitivity: Sensitivity,
907    output_sensitivity: Sensitivity,
908    retry: RetryPolicy,
909    #[cfg(feature = "media")]
910    media: Option<Arc<dyn crate::blob::BlobStore>>,
911    #[cfg(feature = "media")]
912    media_grants: std::collections::BTreeSet<(crate::core::Digest, String)>,
913    /// `/system` when the prompt has one, empty otherwise.
914    ///
915    /// Computed at construction rather than returned fresh, because
916    /// [`Effect::protected_fields`] hands back a borrowed slice — and because a
917    /// declared field is *mandatory*, so declaring `/system` unconditionally
918    /// would refuse every prompt that legitimately has no instruction.
919    protected: Vec<crate::core::ProtectedField>,
920}
921
922impl ModelCall {
923    /// Conservative per-call output ceiling used when the caller does not set
924    /// one explicitly.
925    pub const DEFAULT_MAX_OUTPUT_TOKENS: u32 = 4096;
926
927    /// A completion from this provider.
928    #[must_use]
929    pub fn new(provider: Arc<dyn ModelProvider>, model: ModelId, prompt: Value) -> Self {
930        Self {
931            model,
932            prompt,
933            schema: None,
934            tools: Vec::new(),
935            exchanges: Vec::new(),
936            continuation: None,
937            stream: None,
938            max_output_tokens: Self::DEFAULT_MAX_OUTPUT_TOKENS,
939            reasoning_effort: None,
940            provider,
941            max_sensitivity: Sensitivity::Public,
942            output_sensitivity: Sensitivity::Public,
943            retry: RetryPolicy::never(),
944            #[cfg(feature = "media")]
945            media: None,
946            #[cfg(feature = "media")]
947            media_grants: std::collections::BTreeSet::new(),
948            protected: Vec::new(),
949        }
950        .with_protected_instruction()
951    }
952
953    /// Require the instruction to be trusted, when there is one.
954    ///
955    /// # The instruction slot carries authority; the content does not
956    ///
957    /// A model reads its instruction and its data as the same undifferentiated
958    /// text, so text that *arrives as data* and reads like an instruction is
959    /// obeyed like one. The usual defence — label the data, gate the sinks —
960    /// contains what the model may then *do*, and this crate does that. It does
961    /// not answer the prior question of who was allowed to give the order.
962    ///
963    /// So `/system` is protected: if the prompt has an instruction, it must be
964    /// trusted. Untrusted material belongs in `messages`, where it is content
965    /// the model reasons *about* rather than a directive it reasons *under*.
966    ///
967    /// The consequence is deliberate and it will be met immediately. Building a
968    /// prompt with `untrusted.map(|d| json!({"system": "…", "messages": [d]}))`
969    /// is refused, because `map` cannot prove how a closure reshaped a value and
970    /// so conservatively taints the whole thing — instruction included.
971    /// [`Tainted::object`](crate::core::Tainted::object) keeps the two apart,
972    /// which is what it is for.
973    fn with_protected_instruction(mut self) -> Self {
974        self.protected = if self.prompt.get("system").is_some_and(|s| !s.is_null()) {
975            vec![crate::core::ProtectedField::trusted("/system")]
976        } else {
977            Vec::new()
978        };
979        self
980    }
981
982    /// Tell the model which tools it may ask for.
983    ///
984    /// What comes back is a *request*, not an action: the chosen name is matched
985    /// against the operator's grants exactly — never resolved to a near
986    /// neighbour — and the arguments stay untrusted until they pass the sink's
987    /// field-provenance rules. Declaring is telling; authorizing is separate.
988    ///
989    /// Declare only what is granted. Offering the model a tool the manifest does
990    /// not grant produces a call that is refused after the model has been paid
991    /// for choosing it, and teaches nobody anything.
992    #[must_use]
993    pub fn with_tools(mut self, tools: impl IntoIterator<Item = ToolDeclaration>) -> Self {
994        self.tools = tools.into_iter().collect();
995        self
996    }
997
998    /// Continue after tools ran, showing the model what came back.
999    ///
1000    /// Each exchange carries the call *and* its output: a provider matches them
1001    /// by the id it issued, and an output without its call is rejected.
1002    ///
1003    /// The prompt stays what it was. A continuation is the same question with
1004    /// more known, so re-stating it would change the effect key and make each
1005    /// turn of a loop a different call for replay purposes.
1006    #[must_use]
1007    pub fn continuing(mut self, exchanges: impl IntoIterator<Item = ToolExchange>) -> Self {
1008        self.exchanges = exchanges.into_iter().collect();
1009        self
1010    }
1011
1012    /// Continue with exact provider-native state from the preceding response.
1013    ///
1014    /// This state is not interpreted, synthesized, or fetched by id. It is
1015    /// journaled as part of this call's identity and returned only to the
1016    /// provider that produced it.
1017    #[must_use]
1018    pub fn with_continuation(mut self, continuation: ProviderContinuation) -> Self {
1019        self.continuation = Some(continuation);
1020        self
1021    }
1022
1023    /// Observe visible model text as it arrives during live execution.
1024    #[must_use]
1025    pub fn streaming_to(mut self, observer: Arc<dyn ModelStreamObserver>) -> Self {
1026        self.stream = Some(observer);
1027        self
1028    }
1029
1030    /// Bound how many tokens this call may generate.
1031    #[must_use]
1032    pub const fn with_max_output_tokens(mut self, max_output_tokens: u32) -> Self {
1033        self.max_output_tokens = max_output_tokens;
1034        self
1035    }
1036
1037    /// Request an explicit reasoning depth.
1038    #[must_use]
1039    pub const fn with_reasoning_effort(mut self, effort: ReasoningEffort) -> Self {
1040        self.reasoning_effort = Some(effort);
1041        self
1042    }
1043
1044    /// The highest sensitivity this model may be shown.
1045    ///
1046    /// The control that matters for a hosted model: a prompt assembled from a
1047    /// secret is an exfiltration whether or not anyone meant it.
1048    ///
1049    /// It is also the **floor of the completion's own label** — see
1050    /// [`with_output_sensitivity`](Self::with_output_sensitivity) for why the
1051    /// two are one decision.
1052    #[must_use]
1053    pub const fn with_max_sensitivity(mut self, s: Sensitivity) -> Self {
1054        self.max_sensitivity = s;
1055        self
1056    }
1057
1058    /// Declare the completion's sensitivity floor, above the derived one.
1059    ///
1060    /// The derived floor is [`max_sensitivity`](Self::with_max_sensitivity),
1061    /// always: a caller who raised the egress ceiling to show the model a
1062    /// confidential prompt has told the runtime what class of data the answer
1063    /// was derived from, and a completion labelled below its own prompt is a
1064    /// laundering primitive — ask the model to restate the secret and read it
1065    /// back a level down. So this setter can only *raise* the floor; a value
1066    /// below the ceiling is kept and simply loses to it at
1067    /// [`Effect::output_sensitivity`], where the two are joined in one place
1068    /// rather than at every call site that used to compensate by hand.
1069    #[must_use]
1070    pub const fn with_output_sensitivity(mut self, s: Sensitivity) -> Self {
1071        self.output_sensitivity = s;
1072        self
1073    }
1074
1075    /// The one floor both the terminal completion and the live stream carry:
1076    /// the declared output sensitivity, never below the egress ceiling the
1077    /// prompt was shown under.
1078    ///
1079    /// One function on purpose — the stream label computed its own copy once,
1080    /// and a second implementation of a floor is the pair that drifts at the
1081    /// boundary nobody probed.
1082    fn declared_output_floor(&self) -> Sensitivity {
1083        self.output_sensitivity.max(self.max_sensitivity)
1084    }
1085
1086    #[must_use]
1087    pub const fn with_retry(mut self, r: RetryPolicy) -> Self {
1088        self.retry = r;
1089        self
1090    }
1091
1092    /// Permit these exact [`FetchedMedia`](crate::media::FetchedMedia) artifacts
1093    /// to materialize from this blob store immediately before live dispatch.
1094    ///
1095    /// The prompt and effect key contain only media digests. Strict replay does
1096    /// not execute `perform`, so it reads neither blob storage nor the network.
1097    /// A prompt marker without a matching digest/type grant is refused;
1098    /// knowing another case's digest is not authority to read that blob.
1099    /// Model output remains journaled for replay and may itself reproduce media
1100    /// content; digest-only input storage is not an output-redaction promise.
1101    #[cfg(feature = "media")]
1102    #[must_use]
1103    pub fn with_media<'a>(
1104        mut self,
1105        media: Arc<dyn crate::blob::BlobStore>,
1106        artifacts: impl IntoIterator<Item = &'a crate::media::FetchedMedia>,
1107    ) -> Self {
1108        self.media = Some(media);
1109        for artifact in artifacts {
1110            self.media_grants
1111                .insert((artifact.digest, artifact.media_type.clone()));
1112        }
1113        self
1114    }
1115
1116    /// Require the answer to conform to a JSON Schema.
1117    ///
1118    /// The schema goes into the **effect key**, which is the point: editing a
1119    /// schema changes the effect, so a replayed run reports divergence instead
1120    /// of quietly reading back an answer shaped to different rules. A schema
1121    /// that lived outside the key would let today's shape re-interpret last
1122    /// year's stored answer.
1123    ///
1124    /// Enforcement happens at the provider, during generation, where a
1125    /// constraint can prevent a malformed answer rather than reject one already
1126    /// paid for. What this crate adds on top is the parse — see
1127    /// [`Completion::structured`].
1128    #[must_use]
1129    pub fn expecting(mut self, schema: Value) -> Self {
1130        self.schema = Some(schema);
1131        self
1132    }
1133}
1134
1135#[async_trait]
1136impl Effect for ModelCall {
1137    type Output = Completion;
1138
1139    fn gen_ai_operation(&self) -> Option<&'static str> {
1140        Some(crate::runtime::telemetry::GEN_AI_CHAT)
1141    }
1142
1143    fn descriptor(&self) -> EffectDescriptor {
1144        // Every provider-visible input is in the key. A changed prompt, schema,
1145        // offered tool, or continuation transcript is a changed effect, so an
1146        // edit shows up on replay as divergence rather than reading an answer
1147        // produced for a request nobody made this time.
1148        EffectDescriptor::new(
1149            "model.complete",
1150            serde_json::json!({
1151                "provider": self.model.provider,
1152                "model": self.model.model,
1153                "provider_profile": self.provider.request_profile(&self.model),
1154                "prompt": self.prompt,
1155                // In the key for the same reason the prompt is: a changed
1156                // schema is a changed question, and a replay that read back an
1157                // answer shaped to the old one would be answering a question
1158                // nobody asked.
1159                "schema": self.schema,
1160                "max_output_tokens": self.max_output_tokens,
1161                "reasoning_effort": self.reasoning_effort,
1162                // Tool descriptions and schemas steer generation just as the
1163                // prompt does. Omitting them would let strict replay consume a
1164                // completion produced while a different capability surface was
1165                // offered.
1166                "tools": self.tools.iter().map(|tool| serde_json::json!({
1167                    "name": tool.name,
1168                    "description": tool.description,
1169                    "parameters": tool.parameters,
1170                })).collect::<Vec<_>>(),
1171                // A continuation is the original question plus the exact calls
1172                // and results already observed. IDs, arguments, outputs and the
1173                // failure bit all affect the next provider response.
1174                "exchanges": self.exchanges.iter().map(|exchange| serde_json::json!({
1175                    "call": {
1176                        "id": exchange.call.id,
1177                        "name": exchange.call.name,
1178                        "arguments": exchange.call.arguments,
1179                    },
1180                    "output": exchange.output,
1181                    "failed": exchange.failed,
1182                })).collect::<Vec<_>>(),
1183                "continuation": self.continuation,
1184            }),
1185        )
1186    }
1187
1188    /// A completion does not change the world.
1189    ///
1190    /// Said plainly because it is the one outward call in this crate where that
1191    /// is true, and it is what makes retrying a rate-limit sane. It is *not* free
1192    /// — see `spend` — but a second completion does not move money twice.
1193    fn mutates(&self) -> bool {
1194        false
1195    }
1196
1197    /// `/system` when the prompt has an instruction.
1198    ///
1199    /// A model reads its instruction and its data as the same undifferentiated
1200    /// text, so text arriving as *data* that reads like a directive is obeyed
1201    /// like one. Every other control here bounds what the model may then **do**;
1202    /// this is the only one that asks who was allowed to give the order.
1203    /// Untrusted material belongs in `messages`.
1204    fn protected_fields(&self) -> &[crate::core::ProtectedField] {
1205        &self.protected
1206    }
1207
1208    fn recovery(&self) -> Recovery {
1209        Recovery::Retry
1210    }
1211
1212    fn retry(&self) -> RetryPolicy {
1213        self.retry
1214    }
1215
1216    fn max_sensitivity(&self) -> Sensitivity {
1217        self.max_sensitivity
1218    }
1219
1220    /// Never below [`max_sensitivity`](Effect::max_sensitivity).
1221    ///
1222    /// A sink's result inherits the class of what was sent into it: a caller
1223    /// who raised the egress ceiling to get a confidential prompt out has
1224    /// declared what the answer derives from, and returning that answer at a
1225    /// lower label would make one round trip through a model a release nobody
1226    /// authorized. `with_output_sensitivity` still raises the floor further;
1227    /// nothing lowers it below the ceiling.
1228    fn output_sensitivity(&self) -> Sensitivity {
1229        self.declared_output_floor()
1230    }
1231
1232    fn sink_arguments(&self) -> Option<&Value> {
1233        Some(&self.prompt)
1234    }
1235
1236    /// Model output is untrusted, and this is the case the rule was written for.
1237    ///
1238    /// A completion is a plausible-sounding string produced from whatever was in
1239    /// the context window — including anything untrusted that got there. It is
1240    /// the canonical prompt-injection carrier.
1241    fn trust(&self) -> Trust {
1242        Trust::Untrusted
1243    }
1244
1245    /// Which model answered — `model:{provider}/{model}` — the same spelling
1246    /// the live stream label has always carried, now derived in one place so
1247    /// the terminal completion and the stream cannot name one answer two ways.
1248    fn source(&self) -> crate::core::SourceId {
1249        crate::core::SourceId::new(format!("model:{}", self.model))
1250    }
1251
1252    fn spend(&self, output: &Completion) -> Spend {
1253        output.usage.spend()
1254    }
1255
1256    async fn perform(&self) -> Result<Completion, EffectError> {
1257        #[cfg(feature = "media")]
1258        let prompt =
1259            materialize_media(&self.prompt, self.media.as_ref(), &self.media_grants).await?;
1260        #[cfg(not(feature = "media"))]
1261        let prompt = self.prompt.clone();
1262
1263        // The trait is public, so an embedder may supply a provider whose wire
1264        // implementation this crate cannot inspect. Apply the hard cut at the
1265        // effect boundary as well as inside the built-in drivers: no provider
1266        // reached through the runtime receives a remote media URL.
1267        refuse_provider_side_media(&prompt, &self.model)
1268            .map_err(|error| EffectError::Rejected(error.to_string()))?;
1269
1270        // The same identity the terminal completion is labelled with, from the
1271        // same derivation — `Effect::source` — so the two spellings cannot
1272        // drift apart.
1273        let mut stream_label = crate::core::Label::untrusted(Effect::source(self));
1274        // Raised, never assigned: the untrusted label already carries the
1275        // `Internal` floor every model answer has, and the terminal completion
1276        // is floored the same way at the effect boundary. A plain `=` here
1277        // once *lowered* the stream below that — the same bytes left the
1278        // plane twice, once labelled and once laundered, differing only in
1279        // whether the caller read them live.
1280        stream_label.sensitivity = stream_label.sensitivity.max(self.declared_output_floor());
1281        let answered = self
1282            .provider
1283            .complete(Request {
1284                model: &self.model,
1285                prompt: &prompt,
1286                max_output_tokens: self.max_output_tokens,
1287                reasoning_effort: self.reasoning_effort,
1288                schema: self.schema.as_ref(),
1289                tools: &self.tools,
1290                exchanges: &self.exchanges,
1291                continuation: self.continuation.as_ref(),
1292                stream: self
1293                    .stream
1294                    .as_deref()
1295                    .map(|observer| (observer, &stream_label)),
1296            })
1297            .await;
1298
1299        // The declared schema is a contract on the *answer*, and the runtime
1300        // relies on it: `structured` is indexed into rather than re-checked.
1301        // Held here as well as in each driver, for the reason stated above the
1302        // media check — the provider trait is public.
1303        let answered = answered.and_then(|completion| {
1304            honour_declared_schema(&completion, self.schema.as_ref(), &self.model)?;
1305            Ok(completion)
1306        });
1307
1308        answered.map_err(|e| {
1309            let detail = e.to_string();
1310            let spend = e.usage().spend();
1311            // A failure that consumed nothing is an ordinary failure. One
1312            // that generated tokens has to carry them, or the ceiling that
1313            // exists to bound a runaway provider counts zero.
1314            if spend.is_zero() {
1315                match e.disposition() {
1316                    // A provider's refusal is an answer, not a fault: the
1317                    // request is *wrong* — unknown model, malformed schema,
1318                    // input filtered — and asking again asks the same rule the
1319                    // same question. Carried as `Refused` so the retry loop
1320                    // spends no attempt on it, where a rate limit or an
1321                    // outage stays `Rejected` and retries under policy.
1322                    Disposition::DidNotHappen if matches!(e, ModelError::Refused { .. }) => {
1323                        EffectError::Refused(detail)
1324                    }
1325                    Disposition::DidNotHappen => EffectError::Rejected(detail),
1326                    Disposition::InDoubt => EffectError::Interrupted {
1327                        driver: self.model.to_string(),
1328                        detail,
1329                    },
1330                    Disposition::Landed => EffectError::Performed(detail),
1331                }
1332            } else {
1333                EffectError::Metered {
1334                    detail,
1335                    spend,
1336                    disposition: e.disposition(),
1337                }
1338            }
1339        })
1340    }
1341}
1342
1343#[cfg(feature = "media")]
1344#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1345struct MediaMaterialization {
1346    digest: crate::core::Digest,
1347    media_type: String,
1348    encoding: String,
1349}
1350
1351#[cfg(feature = "media")]
1352async fn materialize_media(
1353    prompt: &Value,
1354    store: Option<&Arc<dyn crate::blob::BlobStore>>,
1355    grants: &std::collections::BTreeSet<(crate::core::Digest, String)>,
1356) -> Result<Value, EffectError> {
1357    let mut references = Vec::new();
1358    collect_media_references(prompt, &mut references)?;
1359    if references.is_empty() {
1360        return Ok(prompt.clone());
1361    }
1362    let store = store.ok_or_else(|| {
1363        EffectError::Rejected(
1364            "the prompt contains governed-media references but ModelCall has no media store"
1365                .to_owned(),
1366        )
1367    })?;
1368    references.sort();
1369    references.dedup();
1370
1371    let mut replacements = std::collections::BTreeMap::new();
1372    for reference in references {
1373        if !grants.contains(&(reference.digest, reference.media_type.clone())) {
1374            return Err(EffectError::Rejected(format!(
1375                "governed media {} with type '{}' is not explicitly granted to this model call",
1376                reference.digest, reference.media_type
1377            )));
1378        }
1379        let bytes = store.get(reference.digest).await.map_err(|error| {
1380            EffectError::Rejected(format!(
1381                "governed media {} could not be materialized: {error}",
1382                reference.digest
1383            ))
1384        })?;
1385        crate::media::verify_materialized(&reference.media_type, &bytes)
1386            .map_err(EffectError::Rejected)?;
1387        let encoded = base64::engine::general_purpose::STANDARD.encode(bytes);
1388        let value = match reference.encoding.as_str() {
1389            "base64" => encoded,
1390            "data_url" => format!("data:{};base64,{encoded}", reference.media_type),
1391            other => {
1392                return Err(EffectError::Rejected(format!(
1393                    "unknown governed-media encoding '{other}'"
1394                )));
1395            }
1396        };
1397        replacements.insert(reference, Value::String(value));
1398    }
1399
1400    replace_media_references(prompt, &replacements)
1401}
1402
1403#[cfg(feature = "media")]
1404fn collect_media_references(
1405    value: &Value,
1406    out: &mut Vec<MediaMaterialization>,
1407) -> Result<(), EffectError> {
1408    match value {
1409        Value::Array(values) => {
1410            for value in values {
1411                collect_media_references(value, out)?;
1412            }
1413        }
1414        Value::Object(object) => {
1415            if object.contains_key("$agentplane_media") {
1416                out.push(parse_media_reference(value)?);
1417            } else {
1418                for value in object.values() {
1419                    collect_media_references(value, out)?;
1420                }
1421            }
1422        }
1423        _ => {}
1424    }
1425    Ok(())
1426}
1427
1428#[cfg(feature = "media")]
1429fn replace_media_references(
1430    value: &Value,
1431    replacements: &std::collections::BTreeMap<MediaMaterialization, Value>,
1432) -> Result<Value, EffectError> {
1433    match value {
1434        Value::Array(values) => values
1435            .iter()
1436            .map(|value| replace_media_references(value, replacements))
1437            .collect::<Result<Vec<_>, _>>()
1438            .map(Value::Array),
1439        Value::Object(object) if object.contains_key("$agentplane_media") => replacements
1440            .get(&parse_media_reference(value)?)
1441            .cloned()
1442            .ok_or_else(|| {
1443                EffectError::Rejected("governed-media replacement is missing".to_owned())
1444            }),
1445        Value::Object(object) => object
1446            .iter()
1447            .map(|(key, value)| Ok((key.clone(), replace_media_references(value, replacements)?)))
1448            .collect::<Result<serde_json::Map<_, _>, EffectError>>()
1449            .map(Value::Object),
1450        _ => Ok(value.clone()),
1451    }
1452}
1453
1454#[cfg(feature = "media")]
1455fn parse_media_reference(value: &Value) -> Result<MediaMaterialization, EffectError> {
1456    let outer = value.as_object().ok_or_else(|| {
1457        EffectError::Rejected("governed-media marker must be an object".to_owned())
1458    })?;
1459    if outer.len() != 1 {
1460        return Err(EffectError::Rejected(
1461            "governed-media marker may not contain sibling fields".to_owned(),
1462        ));
1463    }
1464    let marker = outer
1465        .get("$agentplane_media")
1466        .and_then(Value::as_object)
1467        .ok_or_else(|| {
1468            EffectError::Rejected("governed-media marker body must be an object".to_owned())
1469        })?;
1470    if marker.len() != 3 {
1471        return Err(EffectError::Rejected(
1472            "governed-media marker must contain exactly digest, media_type, and encoding"
1473                .to_owned(),
1474        ));
1475    }
1476    let digest = marker
1477        .get("digest")
1478        .and_then(Value::as_str)
1479        .ok_or_else(|| EffectError::Rejected("governed-media digest is missing".to_owned()))?;
1480    let digest = crate::core::Digest::from_hex(digest).map_err(|error| {
1481        EffectError::Rejected(format!("invalid governed-media digest: {error}"))
1482    })?;
1483    let media_type = marker
1484        .get("media_type")
1485        .and_then(Value::as_str)
1486        .ok_or_else(|| EffectError::Rejected("governed-media media_type is missing".to_owned()))?
1487        .to_owned();
1488    let encoding = marker
1489        .get("encoding")
1490        .and_then(Value::as_str)
1491        .ok_or_else(|| EffectError::Rejected("governed-media encoding is missing".to_owned()))?
1492        .to_owned();
1493    Ok(MediaMaterialization {
1494        digest,
1495        media_type,
1496        encoding,
1497    })
1498}
1499
1500#[cfg(test)]
1501mod tests {
1502    #[cfg(feature = "media")]
1503    use std::sync::Mutex;
1504    use std::sync::atomic::{AtomicUsize, Ordering};
1505
1506    use serde_json::json;
1507
1508    use super::*;
1509
1510    #[derive(Debug)]
1511    struct RecordingProvider(Arc<AtomicUsize>);
1512
1513    #[async_trait]
1514    impl ModelProvider for RecordingProvider {
1515        async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
1516            self.0.fetch_add(1, Ordering::Relaxed);
1517            Err(ModelError::Unavailable {
1518                model: request.model.clone(),
1519                detail: "recording provider was called".to_owned(),
1520            })
1521        }
1522    }
1523
1524    /// `with_retry` replaces the policy, and the default declines to retry.
1525    ///
1526    /// The default is `never` on purpose: a model call that reached the provider
1527    /// and died mid-stream is `Landed`, so repeating it buys a second bill for
1528    /// the same question. A deployment that has decided otherwise sets its own
1529    /// policy here — and the builder had no caller and no test, so a `with_retry`
1530    /// that dropped the value on the floor would have left every such deployment
1531    /// silently on the default.
1532    #[test]
1533    fn with_retry_replaces_a_deliberately_unretrying_default() {
1534        let provider: Arc<dyn ModelProvider> =
1535            Arc::new(RecordingProvider(Arc::new(AtomicUsize::new(0))));
1536        let plain = ModelCall::new(
1537            Arc::clone(&provider),
1538            ModelId::new("custom", "m"),
1539            json!({"q": "hi"}),
1540        );
1541        assert_eq!(
1542            Effect::retry(&plain).max_attempts,
1543            RetryPolicy::never().max_attempts,
1544            "a model call must not retry by default — a died-mid-stream call already landed"
1545        );
1546
1547        let insistent = ModelCall::new(provider, ModelId::new("custom", "m"), json!({"q": "hi"}))
1548            .with_retry(RetryPolicy::default());
1549        assert_eq!(
1550            Effect::retry(&insistent).max_attempts,
1551            RetryPolicy::default().max_attempts
1552        );
1553    }
1554
1555    /// A completion's floor derives from its egress ceiling.
1556    ///
1557    /// A caller who raised `max_sensitivity` to show the model a confidential
1558    /// prompt has said what the answer derives from; a completion labelled
1559    /// below that is a laundering primitive — ask the model to restate the
1560    /// secret and read it back a level down. Every direction is asserted:
1561    /// the derived default, an explicit value below the ceiling losing to it,
1562    /// and an explicit value above it still raising the floor.
1563    #[test]
1564    fn a_completions_floor_derives_from_its_egress_ceiling() {
1565        let provider = || -> Arc<dyn ModelProvider> {
1566            Arc::new(RecordingProvider(Arc::new(AtomicUsize::new(0))))
1567        };
1568        let call = |p| ModelCall::new(p, ModelId::new("custom", "m"), json!({"q": "hi"}));
1569
1570        let derived = call(provider()).with_max_sensitivity(Sensitivity::Confidential);
1571        assert_eq!(
1572            Effect::output_sensitivity(&derived),
1573            Sensitivity::Confidential,
1574            "raising the egress ceiling without setting an output sensitivity \
1575             must floor the completion at the ceiling"
1576        );
1577
1578        let lowered = call(provider())
1579            .with_max_sensitivity(Sensitivity::Confidential)
1580            .with_output_sensitivity(Sensitivity::Public);
1581        assert_eq!(
1582            Effect::output_sensitivity(&lowered),
1583            Sensitivity::Confidential,
1584            "an explicit output sensitivity below the ceiling must lose to it"
1585        );
1586
1587        let raised = call(provider())
1588            .with_max_sensitivity(Sensitivity::Internal)
1589            .with_output_sensitivity(Sensitivity::Secret);
1590        assert_eq!(
1591            Effect::output_sensitivity(&raised),
1592            Sensitivity::Secret,
1593            "an explicit output sensitivity above the ceiling still raises the floor"
1594        );
1595    }
1596
1597    /// Captures the label the runtime hands a driver for live stream delivery.
1598    #[derive(Debug, Default)]
1599    struct StreamLabelProbe(std::sync::Mutex<Option<crate::core::Label>>);
1600
1601    #[async_trait]
1602    impl ModelProvider for StreamLabelProbe {
1603        async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
1604            *self.0.lock().unwrap() = request.stream.map(|(_, label)| label.clone());
1605            Ok(Completion {
1606                text: "ok".to_owned(),
1607                tool_calls: Vec::new(),
1608                usage: Usage::default(),
1609                stop_reason: Some("end_turn".to_owned()),
1610                truncated: false,
1611                structured: None,
1612                continuation: None,
1613            })
1614        }
1615    }
1616
1617    #[derive(Debug)]
1618    struct DropsEvents;
1619    impl ModelStreamObserver for DropsEvents {
1620        fn event(&self, _event: crate::core::Tainted<ModelStreamEvent>) {}
1621    }
1622
1623    /// Stream delivery never carries less than the terminal completion's floor.
1624    ///
1625    /// The terminal answer is floored at the effect boundary — untrusted, so
1626    /// `Internal` at least, raised to the declared output sensitivity. The
1627    /// stream label used to be *assigned* from the declared value instead,
1628    /// so the same bytes left the plane twice: once labelled, once lowered to
1629    /// the `Public` default, differing only in whether the caller read them
1630    /// live.
1631    #[tokio::test]
1632    async fn a_stream_label_never_dips_below_the_terminal_floor() {
1633        let probe = Arc::new(StreamLabelProbe::default());
1634        let call = ModelCall::new(
1635            Arc::clone(&probe) as Arc<dyn ModelProvider>,
1636            ModelId::new("custom", "m"),
1637            json!({"q": "hi"}),
1638        )
1639        .streaming_to(Arc::new(DropsEvents));
1640        call.perform().await.expect("completes");
1641        let label = probe.0.lock().unwrap().clone().expect("a stream label");
1642        assert_eq!(
1643            label.sensitivity,
1644            Sensitivity::Internal,
1645            "with nothing declared, stream delivery dipped below the Internal \
1646             floor every untrusted completion carries"
1647        );
1648
1649        let probe = Arc::new(StreamLabelProbe::default());
1650        let call = ModelCall::new(
1651            Arc::clone(&probe) as Arc<dyn ModelProvider>,
1652            ModelId::new("custom", "m"),
1653            json!({"q": "hi"}),
1654        )
1655        .with_max_sensitivity(Sensitivity::Confidential)
1656        .streaming_to(Arc::new(DropsEvents));
1657        call.perform().await.expect("completes");
1658        let label = probe.0.lock().unwrap().clone().expect("a stream label");
1659        assert_eq!(
1660            label.sensitivity,
1661            Sensitivity::Confidential,
1662            "stream delivery must carry the ceiling-derived floor the terminal \
1663             completion carries"
1664        );
1665    }
1666
1667    /// The effect boundary protects custom providers, not only the built-ins.
1668    #[tokio::test]
1669    async fn a_model_call_refuses_provider_side_media_before_any_provider() {
1670        let calls = Arc::new(AtomicUsize::new(0));
1671        let call = ModelCall::new(
1672            Arc::new(RecordingProvider(Arc::clone(&calls))),
1673            ModelId::new("custom", "vision"),
1674            json!({
1675                "input": [{
1676                    "role": "user",
1677                    "content": [{
1678                        "type": "input_image",
1679                        "image_url": "https://media.example/private.png"
1680                    }]
1681                }]
1682            }),
1683        );
1684
1685        let error = call
1686            .perform()
1687            .await
1688            .expect_err("remote media must be refused");
1689        assert!(matches!(error, EffectError::Rejected(_)), "{error}");
1690        assert_eq!(calls.load(Ordering::Relaxed), 0, "the provider was called");
1691    }
1692
1693    #[cfg(feature = "media")]
1694    #[derive(Debug, Default)]
1695    struct CapturingProvider(Mutex<Option<Value>>);
1696
1697    #[cfg(feature = "media")]
1698    #[async_trait]
1699    impl ModelProvider for CapturingProvider {
1700        async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
1701            *self.0.lock().unwrap() = Some(request.prompt.clone());
1702            Ok(Completion {
1703                tool_calls: Vec::new(),
1704                text: "described".to_owned(),
1705                usage: Usage::default(),
1706                stop_reason: Some("stop".to_owned()),
1707                truncated: false,
1708                structured: None,
1709                continuation: None,
1710            })
1711        }
1712    }
1713
1714    #[cfg(feature = "media")]
1715    #[tokio::test]
1716    async fn governed_media_is_digest_only_until_live_model_dispatch() {
1717        use crate::blob::{BlobStore, MemoryBlobs};
1718        use crate::media::{FetchedMedia, MediaRetention};
1719
1720        let blobs = Arc::new(MemoryBlobs::new());
1721        let bytes = b"\x89PNG\r\n\x1a\nbody";
1722        let digest = blobs.put(bytes).await.unwrap();
1723        let fetched = FetchedMedia {
1724            digest,
1725            media_type: "image/png".to_owned(),
1726            bytes: bytes.len(),
1727            source_url: "https://media.example/a.png".to_owned(),
1728            final_url: "https://media.example/a.png".to_owned(),
1729            redirects: 0,
1730            validated_by: Vec::new(),
1731            hops: Vec::new(),
1732            retention: MediaRetention::External {
1733                policy: "test".to_owned(),
1734            },
1735        };
1736        let provider = Arc::new(CapturingProvider::default());
1737        let call = ModelCall::new(
1738            Arc::clone(&provider) as Arc<dyn ModelProvider>,
1739            ModelId::new("openai", "vision"),
1740            json!({ "input": [{ "content": [fetched.openai_image()] }] }),
1741        )
1742        .with_media(blobs as Arc<dyn BlobStore>, [&fetched]);
1743
1744        let identity = serde_json::to_string(&call.descriptor()).unwrap();
1745        assert!(identity.contains(&digest.to_hex()));
1746        assert!(
1747            !identity.contains("iVBOR"),
1748            "media bytes entered the effect key"
1749        );
1750
1751        call.perform().await.unwrap();
1752        let prompt = provider.0.lock().unwrap().clone().unwrap();
1753        let data_url = prompt["input"][0]["content"][0]["image_url"]
1754            .as_str()
1755            .unwrap();
1756        assert!(data_url.starts_with("data:image/png;base64,iVBOR"));
1757    }
1758
1759    #[cfg(feature = "media")]
1760    #[tokio::test]
1761    async fn knowing_a_media_digest_is_not_authority_to_materialize_its_blob() {
1762        use crate::blob::{BlobStore, MemoryBlobs};
1763        use crate::media::FetchedMedia;
1764
1765        let blobs = Arc::new(MemoryBlobs::new());
1766        let bytes = b"\x89PNG\r\n\x1a\nprivate";
1767        let digest = blobs.put(bytes).await.unwrap();
1768        let provider = Arc::new(RecordingProvider(Arc::new(AtomicUsize::new(0))));
1769        let calls = Arc::clone(&provider.0);
1770        let call = ModelCall::new(
1771            provider as Arc<dyn ModelProvider>,
1772            ModelId::new("openai", "vision"),
1773            json!({
1774                "input": [{
1775                    "content": [{
1776                        "type": "input_image",
1777                        "image_url": {
1778                            "$agentplane_media": {
1779                                "digest": digest,
1780                                "media_type": "image/png",
1781                                "encoding": "data_url"
1782                            }
1783                        }
1784                    }]
1785                }]
1786            }),
1787        )
1788        .with_media(
1789            blobs as Arc<dyn BlobStore>,
1790            std::iter::empty::<&FetchedMedia>(),
1791        );
1792
1793        let error = call.perform().await.expect_err("ungranted digest");
1794        assert!(
1795            error.to_string().contains("not explicitly granted"),
1796            "{error}"
1797        );
1798        assert_eq!(calls.load(Ordering::Relaxed), 0, "the provider was called");
1799    }
1800
1801    #[test]
1802    fn provider_side_media_url_shapes_are_classified_structurally() {
1803        for remote in [
1804            json!({
1805                "type": "image",
1806                "source": { "type": "url", "url": "https://media.example/image.png" }
1807            }),
1808            json!({
1809                "type": "document",
1810                "source": { "type": "url", "url": "https://media.example/document.pdf" }
1811            }),
1812            json!({ "type": "input_image", "image_url": "https://media.example/image.png" }),
1813            json!({ "type": "image_url", "image_url": { "url": "https://media.example/image.png" } }),
1814            json!({ "type": "input_file", "file_url": "https://media.example/document.pdf" }),
1815        ] {
1816            assert!(provider_side_media_reference(&remote).is_some(), "{remote}");
1817        }
1818
1819        for inline_or_text in [
1820            json!({
1821                "type": "image",
1822                "source": { "type": "base64", "media_type": "image/png", "data": "iVBORw0KGgo=" }
1823            }),
1824            json!({
1825                "type": "document",
1826                "source": { "type": "base64", "media_type": "application/pdf", "data": "JVBERi0=" }
1827            }),
1828            json!({
1829                "type": "input_image",
1830                "image_url": "data:image/png;base64,iVBORw0KGgo="
1831            }),
1832            json!({
1833                "type": "input_file",
1834                "filename": "document.pdf",
1835                "file_data": "data:application/pdf;base64,JVBERi0="
1836            }),
1837            json!({ "type": "input_text", "text": "Discuss https://example.com/image.png" }),
1838            json!("https://example.com/image.png"),
1839        ] {
1840            assert!(
1841                provider_side_media_reference(&inline_or_text).is_none(),
1842                "{inline_or_text}"
1843            );
1844        }
1845    }
1846}