#[cfg(feature = "providers")]
pub mod anthropic;
#[cfg(feature = "providers")]
mod anthropic_stream;
#[cfg(feature = "bedrock")]
pub mod bedrock;
#[cfg(feature = "bedrock")]
mod bedrock_stream;
#[cfg(feature = "providers")]
pub mod chat_completions;
#[cfg(feature = "providers")]
mod chat_completions_stream;
#[cfg(feature = "providers")]
pub mod gemini;
#[cfg(feature = "providers")]
mod gemini_stream;
#[cfg(feature = "providers")]
pub mod openai;
#[cfg(feature = "providers")]
mod openai_stream;
#[cfg(feature = "providers")]
mod sse;
#[cfg(feature = "providers")]
mod wire;
use std::fmt::Debug;
use std::sync::Arc;
use async_trait::async_trait;
#[cfg(feature = "media")]
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{
Disposition, Effect, EffectDescriptor, EffectError, Recovery, RetryPolicy, Sensitivity, Spend,
Trust,
};
#[cfg(any(feature = "manifest", feature = "providers", feature = "bedrock"))]
pub(crate) fn validate_schema(schema: &Value, value: &Value) -> Result<(), String> {
let validator = jsonschema::validator_for(schema)
.map_err(|error| format!("the declared JSON Schema is invalid: {error}"))?;
validator
.validate(value)
.map_err(|error| format!("value does not satisfy the declared JSON Schema: {error}"))
}
fn provider_side_media_reference(value: &Value) -> Option<&'static str> {
fn remote_url(value: Option<&Value>) -> bool {
let url = value.and_then(|value| {
value
.as_str()
.or_else(|| value.get("url").and_then(Value::as_str))
});
url.is_some_and(|url| !url.starts_with("data:"))
}
match value {
Value::Array(values) => values.iter().find_map(provider_side_media_reference),
Value::Object(object) => {
let kind = object.get("type").and_then(Value::as_str);
if matches!(kind, Some("image" | "document"))
&& object
.get("source")
.and_then(Value::as_object)
.is_some_and(|source| {
source.get("type").and_then(Value::as_str) == Some("url")
&& source.get("url").and_then(Value::as_str).is_some()
})
{
return Some("an Anthropic image/document URL source");
}
if matches!(kind, Some("input_image" | "image_url"))
&& remote_url(object.get("image_url"))
{
return Some("an OpenAI image URL");
}
if kind == Some("input_file") && remote_url(object.get("file_url")) {
return Some("an OpenAI file URL");
}
for key in ["fileData", "file_data"] {
if let Some(file) = object.get(key).and_then(Value::as_object)
&& (remote_url(file.get("fileUri")) || remote_url(file.get("file_uri")))
{
return Some("a Gemini fileData URI");
}
}
object.values().find_map(provider_side_media_reference)
}
_ => None,
}
}
fn provider_side_media_refusal(kind: &str) -> String {
format!(
"{kind} was refused before dispatch: the model provider would fetch it outside \
this plane's egress policy and journal; inline the media bytes, or fetch them \
through an explicit governed effect first"
)
}
pub(crate) fn refuse_provider_side_media(
prompt: &Value,
model: &ModelId,
) -> Result<(), ModelError> {
let Some(kind) = provider_side_media_reference(prompt) else {
return Ok(());
};
Err(ModelError::Refused {
model: model.clone(),
detail: provider_side_media_refusal(kind),
})
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub struct ModelId {
pub provider: String,
pub model: String,
}
impl ModelId {
pub fn new(provider: impl Into<String>, model: impl Into<String>) -> Self {
Self {
provider: provider.into(),
model: model.into(),
}
}
}
impl std::fmt::Display for ModelId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}/{}", self.provider, self.model)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Usage {
pub input_tokens: u64,
pub output_tokens: u64,
#[serde(default)]
pub cache_write_tokens: u64,
#[serde(default)]
pub cache_read_tokens: u64,
pub minor_units: u64,
}
impl Usage {
#[must_use]
pub const fn spend(&self) -> Spend {
Spend {
tokens: self.input_tokens + self.output_tokens,
minor_units: self.minor_units,
}
}
#[must_use]
pub const fn uncached_input_tokens(&self) -> u64 {
self.input_tokens
.saturating_sub(self.cache_write_tokens)
.saturating_sub(self.cache_read_tokens)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCall {
pub id: String,
pub name: String,
pub arguments: Value,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Completion {
pub text: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub tool_calls: Vec<ToolCall>,
pub usage: Usage,
pub stop_reason: Option<String>,
pub truncated: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub structured: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub continuation: Option<ProviderContinuation>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ProviderContinuation {
pub provider: String,
pub state: Value,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case", tag = "type", content = "value")]
pub enum ModelStreamEvent {
TextDelta(String),
Usage(Usage),
}
pub trait ModelStreamObserver: Send + Sync + Debug {
fn event(&self, event: crate::core::Tainted<ModelStreamEvent>);
}
impl ProviderContinuation {
#[must_use]
pub fn new(provider: impl Into<String>, state: Value) -> Self {
Self {
provider: provider.into(),
state,
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum ModelError {
#[error("could not reach '{model}': {detail}")]
Unreachable { model: ModelId, detail: String },
#[error("'{model}' refused the request: {detail}")]
Refused { model: ModelId, detail: String },
#[error("'{model}' is rate limiting: {detail}")]
RateLimited { model: ModelId, detail: String },
#[error("'{model}' stopped mid-response after {} token(s): {detail}", usage.input_tokens + usage.output_tokens)]
Interrupted {
model: ModelId,
usage: Usage,
detail: String,
},
#[error("'{model}' did not say whether it generated: {detail}")]
Unavailable { model: ModelId, detail: String },
#[error("'{model}' generated and then died without saying what it cost: {detail}")]
Unaccounted { model: ModelId, detail: String },
#[error("'{model}' returned an unusable answer: {detail}")]
Unusable {
model: ModelId,
usage: Usage,
detail: String,
},
}
impl ModelError {
#[must_use]
pub const fn disposition(&self) -> Disposition {
match self {
Self::Unreachable { .. }
| Self::Refused { .. }
| Self::RateLimited { .. }
| Self::Unavailable { .. } => Disposition::DidNotHappen,
Self::Interrupted { .. } | Self::Unusable { .. } | Self::Unaccounted { .. } => {
Disposition::Landed
}
}
}
#[must_use]
pub const fn usage(&self) -> Usage {
match self {
Self::Interrupted { usage, .. } | Self::Unusable { usage, .. } => *usage,
_ => Usage {
input_tokens: 0,
output_tokens: 0,
cache_write_tokens: 0,
cache_read_tokens: 0,
minor_units: 0,
},
}
}
}
#[derive(Debug, Clone)]
pub struct Request<'a> {
pub model: &'a ModelId,
pub prompt: &'a Value,
pub max_output_tokens: u32,
pub reasoning_effort: Option<ReasoningEffort>,
pub schema: Option<&'a Value>,
pub tools: &'a [ToolDeclaration],
pub exchanges: &'a [ToolExchange],
pub continuation: Option<&'a ProviderContinuation>,
pub stream: Option<(&'a dyn ModelStreamObserver, &'a crate::core::Label)>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum ReasoningEffort {
None,
Minimal,
Low,
Medium,
High,
XHigh,
Max,
}
impl ReasoningEffort {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::None => "none",
Self::Minimal => "minimal",
Self::Low => "low",
Self::Medium => "medium",
Self::High => "high",
Self::XHigh => "xhigh",
Self::Max => "max",
}
}
}
#[cfg(test)]
mod reasoning_effort_tests {
use super::ReasoningEffort;
#[test]
fn every_reasoning_effort_has_a_pinned_wire_spelling() {
for (effort, wire) in [
(ReasoningEffort::None, "none"),
(ReasoningEffort::Minimal, "minimal"),
(ReasoningEffort::Low, "low"),
(ReasoningEffort::Medium, "medium"),
(ReasoningEffort::High, "high"),
(ReasoningEffort::XHigh, "xhigh"),
(ReasoningEffort::Max, "max"),
] {
assert_eq!(effort.as_str(), wire);
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SchemaMode {
#[default]
Native,
ForcedTool,
}
#[async_trait]
pub trait ModelProvider: Send + Sync + Debug {
fn request_profile(&self, _model: &ModelId) -> Value {
Value::Null
}
async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ToolDeclaration {
pub name: String,
pub description: String,
pub parameters: Value,
}
impl ToolDeclaration {
#[must_use]
pub fn new(name: impl Into<String>, description: impl Into<String>, parameters: Value) -> Self {
Self {
name: name.into(),
description: description.into(),
parameters,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct ToolExchange {
pub call: ToolCall,
pub output: Value,
pub failed: bool,
}
impl ToolExchange {
#[must_use]
pub fn ok(call: ToolCall, output: Value) -> Self {
Self {
call,
output,
failed: false,
}
}
#[must_use]
pub fn failed(call: ToolCall, detail: impl Into<String>) -> Self {
Self {
call,
output: Value::String(detail.into()),
failed: true,
}
}
}
#[derive(Debug)]
pub struct ModelCall {
model: ModelId,
prompt: Value,
schema: Option<Value>,
tools: Vec<ToolDeclaration>,
exchanges: Vec<ToolExchange>,
continuation: Option<ProviderContinuation>,
stream: Option<Arc<dyn ModelStreamObserver>>,
max_output_tokens: u32,
reasoning_effort: Option<ReasoningEffort>,
provider: Arc<dyn ModelProvider>,
max_sensitivity: Sensitivity,
output_sensitivity: Sensitivity,
retry: RetryPolicy,
#[cfg(feature = "media")]
media: Option<Arc<dyn crate::blob::BlobStore>>,
#[cfg(feature = "media")]
media_grants: std::collections::BTreeSet<(crate::core::Digest, String)>,
protected: Vec<crate::core::ProtectedField>,
}
impl ModelCall {
pub const DEFAULT_MAX_OUTPUT_TOKENS: u32 = 4096;
#[must_use]
pub fn new(provider: Arc<dyn ModelProvider>, model: ModelId, prompt: Value) -> Self {
Self {
model,
prompt,
schema: None,
tools: Vec::new(),
exchanges: Vec::new(),
continuation: None,
stream: None,
max_output_tokens: Self::DEFAULT_MAX_OUTPUT_TOKENS,
reasoning_effort: None,
provider,
max_sensitivity: Sensitivity::Public,
output_sensitivity: Sensitivity::Public,
retry: RetryPolicy::never(),
#[cfg(feature = "media")]
media: None,
#[cfg(feature = "media")]
media_grants: std::collections::BTreeSet::new(),
protected: Vec::new(),
}
.with_protected_instruction()
}
fn with_protected_instruction(mut self) -> Self {
self.protected = if self.prompt.get("system").is_some_and(|s| !s.is_null()) {
vec![crate::core::ProtectedField::trusted("/system")]
} else {
Vec::new()
};
self
}
#[must_use]
pub fn with_tools(mut self, tools: impl IntoIterator<Item = ToolDeclaration>) -> Self {
self.tools = tools.into_iter().collect();
self
}
#[must_use]
pub fn continuing(mut self, exchanges: impl IntoIterator<Item = ToolExchange>) -> Self {
self.exchanges = exchanges.into_iter().collect();
self
}
#[must_use]
pub fn with_continuation(mut self, continuation: ProviderContinuation) -> Self {
self.continuation = Some(continuation);
self
}
#[must_use]
pub fn streaming_to(mut self, observer: Arc<dyn ModelStreamObserver>) -> Self {
self.stream = Some(observer);
self
}
#[must_use]
pub const fn with_max_output_tokens(mut self, max_output_tokens: u32) -> Self {
self.max_output_tokens = max_output_tokens;
self
}
#[must_use]
pub const fn with_reasoning_effort(mut self, effort: ReasoningEffort) -> Self {
self.reasoning_effort = Some(effort);
self
}
#[must_use]
pub const fn with_max_sensitivity(mut self, s: Sensitivity) -> Self {
self.max_sensitivity = s;
self
}
#[must_use]
pub const fn with_output_sensitivity(mut self, s: Sensitivity) -> Self {
self.output_sensitivity = s;
self
}
#[must_use]
pub const fn with_retry(mut self, r: RetryPolicy) -> Self {
self.retry = r;
self
}
#[cfg(feature = "media")]
#[must_use]
pub fn with_media<'a>(
mut self,
media: Arc<dyn crate::blob::BlobStore>,
artifacts: impl IntoIterator<Item = &'a crate::media::FetchedMedia>,
) -> Self {
self.media = Some(media);
for artifact in artifacts {
self.media_grants
.insert((artifact.digest, artifact.media_type.clone()));
}
self
}
#[must_use]
pub fn expecting(mut self, schema: Value) -> Self {
self.schema = Some(schema);
self
}
}
#[async_trait]
impl Effect for ModelCall {
type Output = Completion;
fn gen_ai_operation(&self) -> Option<&'static str> {
Some(crate::runtime::telemetry::GEN_AI_CHAT)
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"model.complete",
serde_json::json!({
"provider": self.model.provider,
"model": self.model.model,
"provider_profile": self.provider.request_profile(&self.model),
"prompt": self.prompt,
"schema": self.schema,
"max_output_tokens": self.max_output_tokens,
"reasoning_effort": self.reasoning_effort,
"tools": self.tools.iter().map(|tool| serde_json::json!({
"name": tool.name,
"description": tool.description,
"parameters": tool.parameters,
})).collect::<Vec<_>>(),
"exchanges": self.exchanges.iter().map(|exchange| serde_json::json!({
"call": {
"id": exchange.call.id,
"name": exchange.call.name,
"arguments": exchange.call.arguments,
},
"output": exchange.output,
"failed": exchange.failed,
})).collect::<Vec<_>>(),
"continuation": self.continuation,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn protected_fields(&self) -> &[crate::core::ProtectedField] {
&self.protected
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn retry(&self) -> RetryPolicy {
self.retry
}
fn max_sensitivity(&self) -> Sensitivity {
self.max_sensitivity
}
fn output_sensitivity(&self) -> Sensitivity {
self.output_sensitivity
}
fn sink_arguments(&self) -> Option<&Value> {
Some(&self.prompt)
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn spend(&self, output: &Completion) -> Spend {
output.usage.spend()
}
async fn perform(&self) -> Result<Completion, EffectError> {
#[cfg(feature = "media")]
let prompt =
materialize_media(&self.prompt, self.media.as_ref(), &self.media_grants).await?;
#[cfg(not(feature = "media"))]
let prompt = self.prompt.clone();
refuse_provider_side_media(&prompt, &self.model)
.map_err(|error| EffectError::Rejected(error.to_string()))?;
let mut stream_label = crate::core::Label::untrusted(crate::core::SourceId::new(format!(
"model:{}",
self.model
)));
stream_label.sensitivity = self.output_sensitivity;
self.provider
.complete(Request {
model: &self.model,
prompt: &prompt,
max_output_tokens: self.max_output_tokens,
reasoning_effort: self.reasoning_effort,
schema: self.schema.as_ref(),
tools: &self.tools,
exchanges: &self.exchanges,
continuation: self.continuation.as_ref(),
stream: self
.stream
.as_deref()
.map(|observer| (observer, &stream_label)),
})
.await
.map_err(|e| {
let detail = e.to_string();
let spend = e.usage().spend();
if spend.is_zero() {
match e.disposition() {
Disposition::DidNotHappen => EffectError::Rejected(detail),
Disposition::InDoubt => EffectError::Interrupted {
driver: self.model.to_string(),
detail,
},
Disposition::Landed => EffectError::Performed(detail),
}
} else {
EffectError::Metered {
detail,
spend,
disposition: e.disposition(),
}
}
})
}
}
#[cfg(feature = "media")]
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
struct MediaMaterialization {
digest: crate::core::Digest,
media_type: String,
encoding: String,
}
#[cfg(feature = "media")]
async fn materialize_media(
prompt: &Value,
store: Option<&Arc<dyn crate::blob::BlobStore>>,
grants: &std::collections::BTreeSet<(crate::core::Digest, String)>,
) -> Result<Value, EffectError> {
let mut references = Vec::new();
collect_media_references(prompt, &mut references)?;
if references.is_empty() {
return Ok(prompt.clone());
}
let store = store.ok_or_else(|| {
EffectError::Rejected(
"the prompt contains governed-media references but ModelCall has no media store"
.to_owned(),
)
})?;
references.sort();
references.dedup();
let mut replacements = std::collections::BTreeMap::new();
for reference in references {
if !grants.contains(&(reference.digest, reference.media_type.clone())) {
return Err(EffectError::Rejected(format!(
"governed media {} with type '{}' is not explicitly granted to this model call",
reference.digest, reference.media_type
)));
}
let bytes = store.get(reference.digest).await.map_err(|error| {
EffectError::Rejected(format!(
"governed media {} could not be materialized: {error}",
reference.digest
))
})?;
crate::media::verify_materialized(&reference.media_type, &bytes)
.map_err(EffectError::Rejected)?;
let encoded = base64::engine::general_purpose::STANDARD.encode(bytes);
let value = match reference.encoding.as_str() {
"base64" => encoded,
"data_url" => format!("data:{};base64,{encoded}", reference.media_type),
other => {
return Err(EffectError::Rejected(format!(
"unknown governed-media encoding '{other}'"
)));
}
};
replacements.insert(reference, Value::String(value));
}
replace_media_references(prompt, &replacements)
}
#[cfg(feature = "media")]
fn collect_media_references(
value: &Value,
out: &mut Vec<MediaMaterialization>,
) -> Result<(), EffectError> {
match value {
Value::Array(values) => {
for value in values {
collect_media_references(value, out)?;
}
}
Value::Object(object) => {
if object.contains_key("$agentplane_media") {
out.push(parse_media_reference(value)?);
} else {
for value in object.values() {
collect_media_references(value, out)?;
}
}
}
_ => {}
}
Ok(())
}
#[cfg(feature = "media")]
fn replace_media_references(
value: &Value,
replacements: &std::collections::BTreeMap<MediaMaterialization, Value>,
) -> Result<Value, EffectError> {
match value {
Value::Array(values) => values
.iter()
.map(|value| replace_media_references(value, replacements))
.collect::<Result<Vec<_>, _>>()
.map(Value::Array),
Value::Object(object) if object.contains_key("$agentplane_media") => replacements
.get(&parse_media_reference(value)?)
.cloned()
.ok_or_else(|| {
EffectError::Rejected("governed-media replacement is missing".to_owned())
}),
Value::Object(object) => object
.iter()
.map(|(key, value)| Ok((key.clone(), replace_media_references(value, replacements)?)))
.collect::<Result<serde_json::Map<_, _>, EffectError>>()
.map(Value::Object),
_ => Ok(value.clone()),
}
}
#[cfg(feature = "media")]
fn parse_media_reference(value: &Value) -> Result<MediaMaterialization, EffectError> {
let outer = value.as_object().ok_or_else(|| {
EffectError::Rejected("governed-media marker must be an object".to_owned())
})?;
if outer.len() != 1 {
return Err(EffectError::Rejected(
"governed-media marker may not contain sibling fields".to_owned(),
));
}
let marker = outer
.get("$agentplane_media")
.and_then(Value::as_object)
.ok_or_else(|| {
EffectError::Rejected("governed-media marker body must be an object".to_owned())
})?;
if marker.len() != 3 {
return Err(EffectError::Rejected(
"governed-media marker must contain exactly digest, media_type, and encoding"
.to_owned(),
));
}
let digest = marker
.get("digest")
.and_then(Value::as_str)
.ok_or_else(|| EffectError::Rejected("governed-media digest is missing".to_owned()))?;
let digest = crate::core::Digest::from_hex(digest).map_err(|error| {
EffectError::Rejected(format!("invalid governed-media digest: {error}"))
})?;
let media_type = marker
.get("media_type")
.and_then(Value::as_str)
.ok_or_else(|| EffectError::Rejected("governed-media media_type is missing".to_owned()))?
.to_owned();
let encoding = marker
.get("encoding")
.and_then(Value::as_str)
.ok_or_else(|| EffectError::Rejected("governed-media encoding is missing".to_owned()))?
.to_owned();
Ok(MediaMaterialization {
digest,
media_type,
encoding,
})
}
#[cfg(test)]
mod tests {
#[cfg(feature = "media")]
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use serde_json::json;
use super::*;
#[derive(Debug)]
struct RecordingProvider(Arc<AtomicUsize>);
#[async_trait]
impl ModelProvider for RecordingProvider {
async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
self.0.fetch_add(1, Ordering::Relaxed);
Err(ModelError::Unavailable {
model: request.model.clone(),
detail: "recording provider was called".to_owned(),
})
}
}
#[test]
fn with_retry_replaces_a_deliberately_unretrying_default() {
let provider: Arc<dyn ModelProvider> =
Arc::new(RecordingProvider(Arc::new(AtomicUsize::new(0))));
let plain = ModelCall::new(
Arc::clone(&provider),
ModelId::new("custom", "m"),
json!({"q": "hi"}),
);
assert_eq!(
Effect::retry(&plain).max_attempts,
RetryPolicy::never().max_attempts,
"a model call must not retry by default — a died-mid-stream call already landed"
);
let insistent = ModelCall::new(provider, ModelId::new("custom", "m"), json!({"q": "hi"}))
.with_retry(RetryPolicy::default());
assert_eq!(
Effect::retry(&insistent).max_attempts,
RetryPolicy::default().max_attempts
);
}
#[tokio::test]
async fn a_model_call_refuses_provider_side_media_before_any_provider() {
let calls = Arc::new(AtomicUsize::new(0));
let call = ModelCall::new(
Arc::new(RecordingProvider(Arc::clone(&calls))),
ModelId::new("custom", "vision"),
json!({
"input": [{
"role": "user",
"content": [{
"type": "input_image",
"image_url": "https://media.example/private.png"
}]
}]
}),
);
let error = call
.perform()
.await
.expect_err("remote media must be refused");
assert!(matches!(error, EffectError::Rejected(_)), "{error}");
assert_eq!(calls.load(Ordering::Relaxed), 0, "the provider was called");
}
#[cfg(feature = "media")]
#[derive(Debug, Default)]
struct CapturingProvider(Mutex<Option<Value>>);
#[cfg(feature = "media")]
#[async_trait]
impl ModelProvider for CapturingProvider {
async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
*self.0.lock().unwrap() = Some(request.prompt.clone());
Ok(Completion {
tool_calls: Vec::new(),
text: "described".to_owned(),
usage: Usage::default(),
stop_reason: Some("stop".to_owned()),
truncated: false,
structured: None,
continuation: None,
})
}
}
#[cfg(feature = "media")]
#[tokio::test]
async fn governed_media_is_digest_only_until_live_model_dispatch() {
use crate::blob::{BlobStore, MemoryBlobs};
use crate::media::{FetchedMedia, MediaRetention};
let blobs = Arc::new(MemoryBlobs::new());
let bytes = b"\x89PNG\r\n\x1a\nbody";
let digest = blobs.put(bytes).await.unwrap();
let fetched = FetchedMedia {
digest,
media_type: "image/png".to_owned(),
bytes: bytes.len(),
source_url: "https://media.example/a.png".to_owned(),
final_url: "https://media.example/a.png".to_owned(),
redirects: 0,
validated_by: Vec::new(),
hops: Vec::new(),
retention: MediaRetention::External {
policy: "test".to_owned(),
},
};
let provider = Arc::new(CapturingProvider::default());
let call = ModelCall::new(
Arc::clone(&provider) as Arc<dyn ModelProvider>,
ModelId::new("openai", "vision"),
json!({ "input": [{ "content": [fetched.openai_image()] }] }),
)
.with_media(blobs as Arc<dyn BlobStore>, [&fetched]);
let identity = serde_json::to_string(&call.descriptor()).unwrap();
assert!(identity.contains(&digest.to_hex()));
assert!(
!identity.contains("iVBOR"),
"media bytes entered the effect key"
);
call.perform().await.unwrap();
let prompt = provider.0.lock().unwrap().clone().unwrap();
let data_url = prompt["input"][0]["content"][0]["image_url"]
.as_str()
.unwrap();
assert!(data_url.starts_with("data:image/png;base64,iVBOR"));
}
#[cfg(feature = "media")]
#[tokio::test]
async fn knowing_a_media_digest_is_not_authority_to_materialize_its_blob() {
use crate::blob::{BlobStore, MemoryBlobs};
use crate::media::FetchedMedia;
let blobs = Arc::new(MemoryBlobs::new());
let bytes = b"\x89PNG\r\n\x1a\nprivate";
let digest = blobs.put(bytes).await.unwrap();
let provider = Arc::new(RecordingProvider(Arc::new(AtomicUsize::new(0))));
let calls = Arc::clone(&provider.0);
let call = ModelCall::new(
provider as Arc<dyn ModelProvider>,
ModelId::new("openai", "vision"),
json!({
"input": [{
"content": [{
"type": "input_image",
"image_url": {
"$agentplane_media": {
"digest": digest,
"media_type": "image/png",
"encoding": "data_url"
}
}
}]
}]
}),
)
.with_media(
blobs as Arc<dyn BlobStore>,
std::iter::empty::<&FetchedMedia>(),
);
let error = call.perform().await.expect_err("ungranted digest");
assert!(
error.to_string().contains("not explicitly granted"),
"{error}"
);
assert_eq!(calls.load(Ordering::Relaxed), 0, "the provider was called");
}
#[test]
fn provider_side_media_url_shapes_are_classified_structurally() {
for remote in [
json!({
"type": "image",
"source": { "type": "url", "url": "https://media.example/image.png" }
}),
json!({
"type": "document",
"source": { "type": "url", "url": "https://media.example/document.pdf" }
}),
json!({ "type": "input_image", "image_url": "https://media.example/image.png" }),
json!({ "type": "image_url", "image_url": { "url": "https://media.example/image.png" } }),
json!({ "type": "input_file", "file_url": "https://media.example/document.pdf" }),
] {
assert!(provider_side_media_reference(&remote).is_some(), "{remote}");
}
for inline_or_text in [
json!({
"type": "image",
"source": { "type": "base64", "media_type": "image/png", "data": "iVBORw0KGgo=" }
}),
json!({
"type": "document",
"source": { "type": "base64", "media_type": "application/pdf", "data": "JVBERi0=" }
}),
json!({
"type": "input_image",
"image_url": "data:image/png;base64,iVBORw0KGgo="
}),
json!({
"type": "input_file",
"filename": "document.pdf",
"file_data": "data:application/pdf;base64,JVBERi0="
}),
json!({ "type": "input_text", "text": "Discuss https://example.com/image.png" }),
json!("https://example.com/image.png"),
] {
assert!(
provider_side_media_reference(&inline_or_text).is_none(),
"{inline_or_text}"
);
}
}
}