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}