Skip to main content

conversation_api/execution/wire/
mod.rs

1//! Versioned transport contract for Agent commands, events, and App Facade calls.
2//!
3//! This crate defines envelopes but no WebSocket client, authentication, reconnect policy, or
4//! server. Transport implementations belong to their deployment repository.
5
6/// Golden enqueue envelope used by deployment and Runtime conformance tests.
7pub const ENQUEUE_FIXTURE: &str = include_str!("../../../fixtures/execution/enqueue.v1.json");
8
9use crate::execution::{
10    AgentCommand, AgentObservation, ContextGeneration, ConversationSurface, DurableEvent, EventId,
11    FacadeRequest, FacadeResult, InvocationContext, ModelMode, PreparedAction, Revision, RunId,
12    Scope, ThreadId,
13};
14use serde::{Deserialize, Serialize};
15use thiserror::Error;
16
17/// Current protocol major version.
18pub const CURRENT_VERSION: u16 = 1;
19
20/// Ephemeral cloud execution identity assigned by Lion.
21///
22/// Both values belong to the transport boundary and never enter the provider-neutral Runtime.
23#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
24pub struct DispatchBinding {
25    /// Unique identity for one execution attempt.
26    pub dispatch_id: String,
27    /// Authenticated Agent WebSocket session selected for the attempt.
28    pub agent_session_id: String,
29}
30
31/// Immutable object-store reference to one Runtime checkpoint generation.
32#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
33pub struct CheckpointReference {
34    /// Object key inside the configured checkpoint bucket.
35    pub object_key: String,
36    /// Lowercase hexadecimal SHA-256 of the checkpoint bytes.
37    pub sha256: String,
38    /// Runtime checkpoint format understood by `agent-rust`.
39    pub format_version: u32,
40}
41
42/// Metadata shared by every wire message.
43#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
44pub struct Metadata {
45    /// Protocol major version.
46    pub version: u16,
47    /// Globally unique message identifier.
48    pub message_id: String,
49    /// Identifier correlating a request and its response or emitted events.
50    pub correlation_id: String,
51    /// Optional identifier of the message that caused this one.
52    pub causation_id: Option<String>,
53    /// Sender timestamp in Unix milliseconds.
54    pub sent_at_unix_ms: u64,
55    /// Execution-attempt and Agent-session binding.
56    pub dispatch: DispatchBinding,
57}
58
59/// Request to execute an App Facade call over a transport.
60#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
61#[serde(tag = "phase", rename_all = "snake_case")]
62#[serde(deny_unknown_fields)]
63pub enum AppFacadeRequest {
64    /// Executes a direct App Facade request.
65    Invoke {
66        /// Authenticated invocation context.
67        context: InvocationContext,
68        /// Canonical Agent-selected request.
69        request: FacadeRequest,
70    },
71    /// Prepares a mutating operation.
72    Prepare {
73        /// Authenticated invocation context.
74        context: InvocationContext,
75        /// Canonical Agent-selected request.
76        request: FacadeRequest,
77    },
78    /// Commits a previously prepared operation.
79    Commit {
80        /// Authenticated invocation context.
81        context: InvocationContext,
82        /// Prepared operation and idempotency metadata.
83        action: PreparedAction,
84    },
85    /// Rejects a previously prepared operation.
86    Reject {
87        /// Authenticated invocation context.
88        context: InvocationContext,
89        /// Prepared operation and idempotency metadata.
90        action: PreparedAction,
91    },
92}
93
94/// Immediate response to command admission for the selected execution attempt.
95#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
96#[serde(tag = "type", rename_all = "snake_case")]
97#[serde(deny_unknown_fields)]
98pub enum CommandResponse {
99    /// Command was admitted; final state is delivered through events.
100    Accepted {
101        /// Authenticated context of the admitted command.
102        context: InvocationContext,
103        /// Stable admission disposition.
104        disposition: AdmissionDisposition,
105        /// Queue revision when the command mutated the admission queue.
106        queue_revision: Option<Revision>,
107    },
108    /// Command was rejected before admission.
109    Rejected {
110        /// Authenticated context of the rejected command.
111        context: InvocationContext,
112        /// Stable machine-readable error code.
113        code: String,
114        /// Safe diagnostic message.
115        message: String,
116    },
117}
118
119/// Stable wire representation of command admission.
120#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
121#[serde(rename_all = "snake_case")]
122pub enum AdmissionDisposition {
123    /// A Lion-dispatched request entered the Agent process's execution queue.
124    Queued,
125    /// Interaction decisions are durable and ready to resume.
126    InteractionReady,
127    /// Cancellation has been signaled to the live Agent process.
128    CancelRequested,
129}
130
131/// Routed best-effort observation sent outside the durable event outbox.
132#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
133pub struct ObservationMessage {
134    /// Monotonic within the bound execution; gaps are allowed and never replayed.
135    pub sequence: u64,
136    /// Authenticated invocation and routing context.
137    pub context: InvocationContext,
138    /// Request-scoped progress payload.
139    pub observation: AgentObservation,
140}
141
142/// Terminal transport failure for an admitted execution that could not produce a checkpoint event.
143#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
144pub struct ExecutionFailure {
145    pub context: InvocationContext,
146    pub code: String,
147    pub message: String,
148}
149
150/// Acknowledges one durable event after the receiver has applied or deduplicated it.
151#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
152pub struct EventAck {
153    pub event_id: EventId,
154    /// Checkpoint reference Lion actually selected while applying (or deduplicating) the event.
155    pub checkpoint: CheckpointReference,
156}
157
158/// Lion-owned archive GC request for all Runtime checkpoints under one private thread.
159#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
160pub struct CheckpointPurgeRequest {
161    pub purge_id: String,
162    pub conversation_key: String,
163    pub scope: Scope,
164    pub surface_id: ConversationSurface,
165    pub thread_id: ThreadId,
166}
167
168/// Result of one idempotent checkpoint-prefix purge.
169#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
170#[serde(tag = "status", rename_all = "snake_case")]
171#[serde(deny_unknown_fields)]
172pub enum CheckpointPurgeAck {
173    Deleted {
174        purge_id: String,
175        conversation_key: String,
176        thread_id: ThreadId,
177        deleted_objects: u32,
178    },
179    Failed {
180        purge_id: String,
181        conversation_key: String,
182        thread_id: ThreadId,
183        code: String,
184        message: String,
185    },
186}
187
188/// Model facts captured by the deployment with the desired route parameters.
189/// This snapshot belongs to admission, not to the provider-neutral LLM request.
190#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
191#[serde(deny_unknown_fields)]
192pub struct LlmModelSnapshot {
193    pub capabilities: llm_api::ModelCapabilities,
194    pub generation_support: LlmGenerationSupport,
195}
196
197/// Catalog parameter support travels with the plan so a run never resolves against a
198/// catalog that changed after admission. Resolution itself belongs to `llm_api::resolve`.
199pub use llm_api::GenerationSupport as LlmGenerationSupport;
200
201/// One concrete cloud model route frozen by the deployment before command admission.
202/// The generation fields are the deployment's desired preferences, not resolved values:
203/// the runtime resolves them once against `model_snapshot` and the backend protocol.
204#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
205pub struct LlmRoute {
206    pub backend: String,
207    pub profile: String,
208    pub model: String,
209    pub provider_preferences: Vec<String>,
210    pub temperature: Option<f64>,
211    pub cache_control: Option<bool>,
212    pub reasoning_effort: Option<String>,
213    #[serde(default, skip_serializing_if = "Option::is_none")]
214    pub thinking: Option<bool>,
215    #[serde(default, skip_serializing_if = "Option::is_none")]
216    pub fast_mode: Option<bool>,
217    /// Total context window of the exact frozen model route.
218    pub context_window_tokens: u32,
219    /// Required immutable facts; old plans must not silently resolve against a new catalog.
220    pub model_snapshot: LlmModelSnapshot,
221}
222
223/// One resolved model configuration frozen for the entire admitted run.
224/// `model_mode` records the entry-point selection; internal calls do not re-resolve it.
225#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
226pub struct LlmExecutionPlan {
227    pub model_mode: ModelMode,
228    pub revision: Option<String>,
229    pub route: LlmRoute,
230}
231
232/// Transport admission request. Environment bindings are consumed by the composition root and do
233/// not enter the provider-neutral Runtime command.
234#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
235pub struct CommandRequest {
236    pub command: AgentCommand,
237    /// Effective prompt-context generation selected by Lion for this dispatch.
238    pub context_generation: Option<ContextGeneration>,
239    pub llm_plan: Option<LlmExecutionPlan>,
240    /// Stable Runtime checkpoint selected by Lion, absent for a new thread.
241    pub base_checkpoint: Option<CheckpointReference>,
242    /// Attempt already terminalized by Lion and settled silently before new work starts.
243    pub abandoned_run_id: Option<RunId>,
244}
245
246/// Checkpoint-coupled Runtime event. Lion commits the generation before projecting the event.
247#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
248pub struct EventMessage {
249    pub event: DurableEvent,
250    pub checkpoint: CheckpointReference,
251}
252
253/// Result of an App Facade request transported back to Runtime.
254#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
255#[serde(tag = "type", rename_all = "snake_case")]
256#[serde(deny_unknown_fields)]
257pub enum AppFacadeResponse {
258    /// Read-only or committed operation result.
259    Result {
260        /// Structured App Facade result.
261        result: FacadeResult,
262    },
263    /// Prepared operation awaiting a commit or rejection.
264    Prepared {
265        /// Prepared action.
266        action: PreparedAction,
267    },
268    /// Rejection completed successfully.
269    #[serde(deserialize_with = "crate::execution::deserialize_empty_variant")]
270    Rejected,
271    /// Stable error safe to expose across the transport.
272    Failed {
273        /// Machine-readable category.
274        code: String,
275        /// Safe diagnostic message.
276        message: String,
277        /// Optional retry delay.
278        retry_after_ms: Option<u64>,
279    },
280}
281
282/// Payload families carried by the versioned envelope.
283#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
284#[serde(tag = "family", content = "payload", rename_all = "snake_case")]
285pub enum Payload {
286    /// Command entering Runtime.
287    Command(Box<CommandRequest>),
288    /// Immediate command-admission response.
289    CommandResponse(CommandResponse),
290    /// Event leaving Runtime.
291    Event(EventMessage),
292    /// Durable receiver acknowledgement used to clear the Runtime outbox.
293    EventAck(EventAck),
294    /// Lion asks the cloud composition to remove an archived thread's Runtime checkpoints.
295    CheckpointPurge(CheckpointPurgeRequest),
296    /// Agent-cloud reports the idempotent S3 purge result.
297    CheckpointPurgeAck(CheckpointPurgeAck),
298    /// Non-durable progress leaving Runtime.
299    Observation(ObservationMessage),
300    /// One dispatch failed outside the checkpointed Runtime state machine.
301    ExecutionFailure(ExecutionFailure),
302    /// Runtime request to the App Facade.
303    AppFacadeRequest(AppFacadeRequest),
304    /// App Facade response returned to Runtime.
305    AppFacadeResponse(AppFacadeResponse),
306}
307
308/// Complete transport message.
309#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
310pub struct Envelope {
311    /// Correlation and version metadata.
312    pub metadata: Metadata,
313    /// Typed payload.
314    pub payload: Payload,
315}
316
317impl Envelope {
318    /// Decodes a JSON envelope against the protocol types.
319    ///
320    /// Transport adapters should use this entry point rather than calling `serde_json` directly.
321    /// Unknown keys are rejected. Canonical emit omits skip-optional `None`; an explicit `null` on
322    /// those keys is accepted as absent so accept follows the types instead of a serialize roundtrip.
323    ///
324    /// # Errors
325    ///
326    /// Returns `ProtocolError` for malformed JSON, unknown fields, or failed envelope validation.
327    pub fn decode_json(input: &str) -> Result<Self, ProtocolError> {
328        let mut deserializer = serde_json::Deserializer::from_str(input);
329        let mut unknown = Vec::new();
330        let envelope: Self = serde_ignored::deserialize(&mut deserializer, |path| {
331            unknown.push(path.to_string());
332        })
333        .map_err(|error| ProtocolError::InvalidJson {
334            message: error.to_string(),
335        })?;
336        deserializer
337            .end()
338            .map_err(|error| ProtocolError::InvalidJson {
339                message: error.to_string(),
340            })?;
341        if let Some(path) = unknown.into_iter().next() {
342            return Err(ProtocolError::UnknownField { path });
343        }
344        envelope.validate()?;
345        Ok(envelope)
346    }
347
348    /// Validates transport-level invariants before dispatch.
349    ///
350    /// # Errors
351    ///
352    /// Returns `ProtocolError` when the major version is unsupported or a required identifier is
353    /// blank.
354    pub fn validate(&self) -> Result<(), ProtocolError> {
355        if self.metadata.version != CURRENT_VERSION {
356            return Err(ProtocolError::UnsupportedVersion {
357                actual: self.metadata.version,
358            });
359        }
360        if self.metadata.message_id.trim().is_empty() {
361            return Err(ProtocolError::MissingIdentifier("message_id"));
362        }
363        if self.metadata.correlation_id.trim().is_empty() {
364            return Err(ProtocolError::MissingIdentifier("correlation_id"));
365        }
366        validate_dispatch(&self.metadata.dispatch)?;
367        if matches!(&self.payload, Payload::Observation(message) if message.sequence == 0) {
368            return Err(ProtocolError::InvalidBinding(
369                "observation sequence must be positive".into(),
370            ));
371        }
372        if let Some(context) = self.context() {
373            validate_context(context)?;
374        }
375        if let Payload::Command(request) = &self.payload {
376            validate_command_binding(request)?;
377        }
378        if let Payload::CheckpointPurge(request) = &self.payload {
379            validate_purge_request(request)?;
380        }
381        if let Payload::CheckpointPurgeAck(ack) = &self.payload {
382            validate_purge_ack(ack)?;
383        }
384        if let Payload::EventAck(ack) = &self.payload {
385            validate_checkpoint_reference(&ack.checkpoint)?;
386        }
387        match &self.payload {
388            Payload::Command(request) => {
389                if let Some(reference) = &request.base_checkpoint {
390                    validate_checkpoint_reference(reference)?;
391                }
392                if matches!(&request.command, AgentCommand::Cancel { .. })
393                    && (request.base_checkpoint.is_some() || request.abandoned_run_id.is_some())
394                {
395                    return Err(ProtocolError::InvalidBinding(
396                        "cancel must target live execution without checkpoint recovery metadata"
397                            .to_owned(),
398                    ));
399                }
400                if request.abandoned_run_id.is_some() && request.base_checkpoint.is_none() {
401                    return Err(ProtocolError::InvalidBinding(
402                        "an abandoned run requires a base checkpoint".to_owned(),
403                    ));
404                }
405            }
406            Payload::Event(message) => validate_checkpoint_reference(&message.checkpoint)?,
407            _ => {}
408        }
409        Ok(())
410    }
411
412    fn context(&self) -> Option<&InvocationContext> {
413        match &self.payload {
414            Payload::Command(request) => match &request.command {
415                AgentCommand::Enqueue { context, .. }
416                | AgentCommand::SubmitInteraction { context, .. }
417                | AgentCommand::Cancel { context } => Some(context),
418            },
419            Payload::Event(message) => Some(&message.event.context),
420            Payload::Observation(message) => Some(&message.context),
421            Payload::ExecutionFailure(message) => Some(&message.context),
422            Payload::AppFacadeRequest(request) => match request {
423                AppFacadeRequest::Invoke { context, .. }
424                | AppFacadeRequest::Prepare { context, .. }
425                | AppFacadeRequest::Commit { context, .. }
426                | AppFacadeRequest::Reject { context, .. } => Some(context),
427            },
428            Payload::CommandResponse(response) => match response {
429                CommandResponse::Accepted { context, .. }
430                | CommandResponse::Rejected { context, .. } => Some(context),
431            },
432            Payload::EventAck(_)
433            | Payload::CheckpointPurge(_)
434            | Payload::CheckpointPurgeAck(_)
435            | Payload::AppFacadeResponse(_) => None,
436        }
437    }
438}
439
440fn validate_purge_request(request: &CheckpointPurgeRequest) -> Result<(), ProtocolError> {
441    if request.purge_id.trim().is_empty() {
442        return Err(ProtocolError::MissingIdentifier("purge_id"));
443    }
444    if request.conversation_key.trim().is_empty() {
445        return Err(ProtocolError::MissingIdentifier("conversation_key"));
446    }
447    if request.scope.scope_id.as_str().trim().is_empty() {
448        return Err(ProtocolError::MissingIdentifier("scope_id"));
449    }
450    if request.thread_id.as_str().trim().is_empty() {
451        return Err(ProtocolError::MissingIdentifier("thread_id"));
452    }
453    Ok(())
454}
455
456fn validate_purge_ack(ack: &CheckpointPurgeAck) -> Result<(), ProtocolError> {
457    let (purge_id, conversation_key, thread_id) = match ack {
458        CheckpointPurgeAck::Deleted {
459            purge_id,
460            conversation_key,
461            thread_id,
462            ..
463        }
464        | CheckpointPurgeAck::Failed {
465            purge_id,
466            conversation_key,
467            thread_id,
468            ..
469        } => (purge_id, conversation_key, thread_id),
470    };
471    if purge_id.trim().is_empty() {
472        return Err(ProtocolError::MissingIdentifier("purge_id"));
473    }
474    if conversation_key.trim().is_empty() {
475        return Err(ProtocolError::MissingIdentifier("conversation_key"));
476    }
477    if thread_id.as_str().trim().is_empty() {
478        return Err(ProtocolError::MissingIdentifier("thread_id"));
479    }
480    if let CheckpointPurgeAck::Failed { code, .. } = ack
481        && code.trim().is_empty()
482    {
483        return Err(ProtocolError::MissingIdentifier("code"));
484    }
485    Ok(())
486}
487
488fn validate_dispatch(dispatch: &DispatchBinding) -> Result<(), ProtocolError> {
489    if dispatch.dispatch_id.trim().is_empty() {
490        return Err(ProtocolError::MissingIdentifier("dispatch_id"));
491    }
492    if dispatch.agent_session_id.trim().is_empty() {
493        return Err(ProtocolError::MissingIdentifier("agent_session_id"));
494    }
495    Ok(())
496}
497
498fn validate_checkpoint_reference(reference: &CheckpointReference) -> Result<(), ProtocolError> {
499    if reference.object_key.trim().is_empty() {
500        return Err(ProtocolError::MissingIdentifier("checkpoint_object_key"));
501    }
502    if reference.format_version == 0
503        || reference.sha256.len() != 64
504        || !reference
505            .sha256
506            .bytes()
507            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
508    {
509        return Err(ProtocolError::InvalidBinding(
510            "checkpoint reference is invalid".to_owned(),
511        ));
512    }
513    Ok(())
514}
515
516fn validate_command_binding(request: &CommandRequest) -> Result<(), ProtocolError> {
517    match (&request.command, &request.context_generation) {
518        (AgentCommand::Enqueue { .. } | AgentCommand::SubmitInteraction { .. }, Some(value))
519            if value.is_complete() => {}
520        (AgentCommand::Enqueue { .. } | AgentCommand::SubmitInteraction { .. }, _) => {
521            return Err(ProtocolError::InvalidBinding(
522                "enqueue and submit_interaction require a context generation".to_owned(),
523            ));
524        }
525        (AgentCommand::Cancel { .. }, None) => {}
526        (AgentCommand::Cancel { .. }, Some(_)) => {
527            return Err(ProtocolError::InvalidBinding(
528                "cancel must not carry a context generation".to_owned(),
529            ));
530        }
531    }
532    match (&request.command, &request.llm_plan) {
533        (AgentCommand::Enqueue { options, .. }, Some(plan)) => {
534            validate_llm_plan(plan)?;
535            if plan.model_mode != options.model_mode {
536                return Err(ProtocolError::InvalidBinding(
537                    "LLM plan mode differs from the command's selected mode".to_owned(),
538                ));
539            }
540            Ok(())
541        }
542        (AgentCommand::Enqueue { .. }, None) => Err(ProtocolError::InvalidBinding(
543            "enqueue requires a frozen LLM plan".to_owned(),
544        )),
545        (AgentCommand::SubmitInteraction { .. }, Some(plan)) => validate_llm_plan(plan),
546        (AgentCommand::SubmitInteraction { .. }, None) => Err(ProtocolError::InvalidBinding(
547            "submit_interaction requires the run's frozen LLM plan".to_owned(),
548        )),
549        (_, Some(_)) => Err(ProtocolError::InvalidBinding(
550            "only commands that invoke the model may carry an LLM plan".to_owned(),
551        )),
552        (_, None) => Ok(()),
553    }
554}
555
556fn validate_llm_plan(plan: &LlmExecutionPlan) -> Result<(), ProtocolError> {
557    let route = &plan.route;
558    route
559        .model_snapshot
560        .generation_support
561        .validate()
562        .map_err(|error| ProtocolError::InvalidBinding(error.into()))?;
563    llm_api::GenerationParameters {
564        temperature: route.temperature,
565        reasoning_effort: route.reasoning_effort.clone(),
566        thinking: route.thinking,
567        fast_mode: route.fast_mode,
568    }
569    .validate()
570    .map_err(|error| ProtocolError::InvalidBinding(error.into()))?;
571    if plan.model_mode.0.trim().is_empty()
572        || route.backend.trim().is_empty()
573        || route.profile.trim().is_empty()
574        || route.model.trim().is_empty()
575    {
576        return Err(ProtocolError::InvalidBinding(
577            "LLM model mode, backend, profile, and model are required".to_owned(),
578        ));
579    }
580    if plan
581        .revision
582        .as_ref()
583        .is_some_and(|value| value.trim().is_empty())
584        || route
585            .provider_preferences
586            .iter()
587            .any(|value| value.trim().is_empty())
588        || route
589            .reasoning_effort
590            .as_ref()
591            .is_some_and(|value| value.trim().is_empty())
592        || route.temperature.is_some_and(|value| !value.is_finite())
593        || route.context_window_tokens <= 20_000
594    {
595        return Err(ProtocolError::InvalidBinding(
596            "LLM route contains invalid optional settings".to_owned(),
597        ));
598    }
599    Ok(())
600}
601
602fn validate_context(context: &InvocationContext) -> Result<(), ProtocolError> {
603    validate_scope(&context.scope)?;
604    for (name, value) in [
605        ("user_id", context.actor.user_id.as_str()),
606        ("thread_id", context.thread_id.as_str()),
607        ("run_id", context.run_id.as_str()),
608        ("operation_id", context.operation_id.as_str()),
609    ] {
610        if value.trim().is_empty() {
611            return Err(ProtocolError::MissingIdentifier(name));
612        }
613    }
614    if context
615        .actor
616        .client_id
617        .as_ref()
618        .is_some_and(|client| client.as_str().trim().is_empty())
619    {
620        return Err(ProtocolError::MissingIdentifier("client_id"));
621    }
622    Ok(())
623}
624
625fn validate_scope(scope: &Scope) -> Result<(), ProtocolError> {
626    if scope.scope_id.as_str().trim().is_empty() {
627        return Err(ProtocolError::MissingIdentifier("scope_id"));
628    }
629    if scope
630        .tenant_id
631        .as_ref()
632        .is_some_and(|tenant| tenant.as_str().trim().is_empty())
633    {
634        return Err(ProtocolError::MissingIdentifier("tenant_id"));
635    }
636    Ok(())
637}
638
639/// Invalid wire envelope.
640#[derive(Clone, Debug, Error, Eq, PartialEq)]
641pub enum ProtocolError {
642    /// JSON could not be decoded into an envelope.
643    #[error("invalid protocol JSON: {message}")]
644    InvalidJson {
645        /// Safe parser diagnostic.
646        message: String,
647    },
648    /// The payload contained a field outside the normative schema.
649    #[error("unknown protocol field: {path}")]
650    UnknownField {
651        /// Serde path to the unexpected field.
652        path: String,
653    },
654    /// The sender used a protocol major version this crate does not understand.
655    #[error("unsupported protocol version {actual}")]
656    UnsupportedVersion {
657        /// Received major version.
658        actual: u16,
659    },
660    /// A required identifier was blank.
661    #[error("missing required identifier: {0}")]
662    MissingIdentifier(&'static str),
663    /// Environment binding is absent or disagrees with its Runtime command.
664    #[error("invalid command binding: {0}")]
665    InvalidBinding(String),
666}
667
668#[cfg(test)]
669mod tests {
670    use crate::execution::{
671        AccessMode, Actor, AgentEvent, ClientId, ContentPart, EventId, Message, MessageRole,
672        ModelMode, OperationId, RunId, RunOptions, Scope, ScopeId, TenantId, ThreadId, UseCase,
673        UserId,
674    };
675
676    use super::*;
677
678    #[test]
679    fn rejects_unknown_protocol_version() {
680        let envelope = Envelope {
681            metadata: Metadata {
682                version: CURRENT_VERSION + 1,
683                message_id: "message-1".to_owned(),
684                correlation_id: "correlation-1".to_owned(),
685                causation_id: None,
686                sent_at_unix_ms: 0,
687                dispatch: dispatch(),
688            },
689            payload: Payload::Event(EventMessage {
690                event: DurableEvent {
691                    id: EventId::from("event-1"),
692                    sequence: 1,
693                    context: InvocationContext {
694                        scope: Scope {
695                            tenant_id: None,
696                            scope_id: ScopeId::from("scope-1"),
697                        },
698                        actor: Actor {
699                            user_id: UserId::from("user-1"),
700                            client_id: Some(ClientId::from("client-1")),
701                        },
702                        surface_id: ConversationSurface::client_personal("user-1").unwrap(),
703                        thread_id: ThreadId::from("thread-1"),
704                        run_id: RunId::from("run-1"),
705                        operation_id: OperationId::from("operation-1"),
706                        deadline_unix_ms: None,
707                        traceparent: None,
708                    },
709                    event: AgentEvent::Started,
710                },
711                checkpoint: checkpoint(),
712            }),
713        };
714
715        assert_eq!(
716            envelope.validate(),
717            Err(ProtocolError::UnsupportedVersion {
718                actual: CURRENT_VERSION + 1
719            })
720        );
721    }
722
723    #[test]
724    fn command_fixture_round_trips_without_shape_drift() {
725        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
726        let envelope = Envelope::decode_json(fixture).expect("valid golden fixture");
727
728        let expected = Envelope {
729            metadata: Metadata {
730                version: CURRENT_VERSION,
731                message_id: "message-1".to_owned(),
732                correlation_id: "correlation-1".to_owned(),
733                causation_id: None,
734                sent_at_unix_ms: 1_700_000_000_000,
735                dispatch: dispatch(),
736            },
737            payload: Payload::Command(Box::new(CommandRequest {
738                command: AgentCommand::Enqueue {
739                    context: InvocationContext {
740                        scope: Scope {
741                            tenant_id: Some(TenantId::from("tenant-1")),
742                            scope_id: ScopeId::from("scope-1"),
743                        },
744                        actor: Actor {
745                            user_id: UserId::from("user-1"),
746                            client_id: Some(ClientId::from("client-1")),
747                        },
748                        surface_id: ConversationSurface::client_personal("user-1").unwrap(),
749                        thread_id: ThreadId::from("thread-1"),
750                        run_id: RunId::from("run-1"),
751                        operation_id: OperationId::from("operation-1"),
752                        deadline_unix_ms: None,
753                        traceparent: None,
754                    },
755                    message: Message {
756                        continuation: None,
757                        role: MessageRole::User,
758                        content: vec![ContentPart::Text {
759                            text: "hello".to_owned(),
760                        }],
761                    },
762                    options: RunOptions {
763                        use_case: UseCase("chat".to_owned()),
764                        model_mode: ModelMode("auto".to_owned()),
765                        access_mode: AccessMode::Interactive,
766                        allow_tools: true,
767                        max_steps: 8,
768                        credit_budget: None,
769                        debug: false,
770                    },
771                },
772                context_generation: Some(ContextGeneration {
773                    user_scope: "scope-1".to_owned(),
774                    identity: "identity-1".to_owned(),
775                    memory: "memory-1".to_owned(),
776                    integration_guide: "guide-1".to_owned(),
777                    installed_integrations: "installed-1".to_owned(),
778                    scope_integrations: "scope-integrations-1".to_owned(),
779                }),
780                llm_plan: Some(LlmExecutionPlan {
781                    model_mode: ModelMode("auto".to_owned()),
782                    revision: Some("1700000000000".to_owned()),
783                    route: LlmRoute {
784                        backend: "openrouter".to_owned(),
785                        profile: "cloud_openrouter".to_owned(),
786                        model: "openai/gpt-5.4-mini".to_owned(),
787                        provider_preferences: vec!["OpenAI".to_owned()],
788                        temperature: Some(0.5),
789                        cache_control: Some(true),
790                        reasoning_effort: None,
791                        thinking: None,
792                        fast_mode: None,
793                        context_window_tokens: 128_000,
794                        model_snapshot: LlmModelSnapshot {
795                            capabilities: llm_api::ModelCapabilities {
796                                input: vec!["image".to_owned()],
797                                tool_calling: true,
798                                structured_output: true,
799                            },
800                            generation_support: LlmGenerationSupport {
801                                temperature: Some(true),
802                                max_tokens: Some(true),
803                                thinking: Some(true),
804                                reasoning_efforts: Some(vec!["high".into()]),
805                                temperature_with_reasoning: Some(true),
806                                ..LlmGenerationSupport::default()
807                            },
808                        },
809                    },
810                }),
811                base_checkpoint: None,
812                abandoned_run_id: None,
813            })),
814        };
815        assert_eq!(envelope, expected);
816        assert_eq!(
817            serde_json::to_value(envelope).expect("serializes"),
818            serde_json::from_str::<serde_json::Value>(fixture).expect("fixture JSON")
819        );
820    }
821
822    #[test]
823    fn snapshot_is_required_and_support_metadata_is_strict() {
824        let mut fixture: serde_json::Value = serde_json::from_str(ENQUEUE_FIXTURE).unwrap();
825        fixture["payload"]["payload"]["llm_plan"]["route"]
826            .as_object_mut()
827            .unwrap()
828            .remove("model_snapshot");
829        assert!(Envelope::decode_json(&fixture.to_string()).is_err());
830        let mut envelope = Envelope::decode_json(ENQUEUE_FIXTURE).unwrap();
831        let Payload::Command(command) = &mut envelope.payload else {
832            panic!("command");
833        };
834        command
835            .llm_plan
836            .as_mut()
837            .unwrap()
838            .route
839            .model_snapshot
840            .generation_support
841            .max_output_tokens = Some(0);
842        assert!(envelope.validate().is_err());
843    }
844
845    #[test]
846    fn generation_support_travels_in_catalog_spelling() {
847        let support = LlmGenerationSupport {
848            temperature: Some(true),
849            max_tokens: Some(true),
850            thinking: Some(true),
851            reasoning_efforts: Some(vec!["low".into(), "high".into()]),
852            reasoning_effort_default: Some("high".into()),
853            temperature_with_reasoning: Some(false),
854            temperature_max: Some(2.0),
855            max_output_tokens: Some(8_192),
856            fast_mode: Some(true),
857        };
858        let encoded = serde_json::to_value(&support).unwrap();
859        assert_eq!(
860            encoded["reasoningEfforts"],
861            serde_json::json!(["low", "high"])
862        );
863        assert_eq!(
864            encoded["temperatureWithReasoning"],
865            serde_json::json!(false)
866        );
867        assert_eq!(
868            serde_json::from_value::<LlmGenerationSupport>(encoded).unwrap(),
869            support
870        );
871        let ignored = serde_json::from_value::<LlmGenerationSupport>(serde_json::json!({
872            "temperature": true, "maxTokens": true, "reasoningEfforts": ["low"],
873            "temperatureWithReasoning": null, "temperatureMax": null, "maxOutputTokens": null,
874            "reasoningIntensity": "high"
875        }))
876        .unwrap();
877        assert_eq!(ignored.temperature, Some(true));
878        assert_eq!(
879            ignored.reasoning_efforts.as_deref(),
880            Some(["low".to_string()].as_slice())
881        );
882    }
883
884    fn dispatch() -> DispatchBinding {
885        DispatchBinding {
886            dispatch_id: "dispatch-1".to_owned(),
887            agent_session_id: "agent-session-1".to_owned(),
888        }
889    }
890
891    fn checkpoint() -> CheckpointReference {
892        CheckpointReference {
893            object_key: "checkpoints/abc.json".to_owned(),
894            sha256: "a".repeat(64),
895            format_version: 1,
896        }
897    }
898
899    #[test]
900    fn every_v1_golden_fixture_round_trips() {
901        let schema: serde_json::Value = serde_json::from_str(include_str!(
902            "../../../schema/execution/envelope.v1.schema.json"
903        ))
904        .expect("valid JSON Schema document");
905        let validator = jsonschema::validator_for(&schema).expect("valid JSON Schema semantics");
906        for fixture in [
907            include_str!("../../../fixtures/execution/enqueue.v1.json"),
908            include_str!("../../../fixtures/execution/completed-event.v1.json"),
909            include_str!("../../../fixtures/execution/interaction-required-event.v1.json"),
910            include_str!("../../../fixtures/execution/failed-event.v1.json"),
911            include_str!("../../../fixtures/execution/app-facade-prepare.v1.json"),
912            include_str!("../../../fixtures/execution/observation.v1.json"),
913            include_str!("../../../fixtures/execution/admission.v1.json"),
914            include_str!("../../../fixtures/execution/app-facade-prepared.v1.json"),
915            include_str!("../../../fixtures/execution/app-facade-user-action.v1.json"),
916            include_str!("../../../fixtures/execution/event-ack.v1.json"),
917            include_str!("../../../fixtures/execution/checkpoint-purge.v1.json"),
918            include_str!("../../../fixtures/execution/checkpoint-purge-ack.v1.json"),
919        ] {
920            let json: serde_json::Value =
921                serde_json::from_str(fixture).expect("valid golden fixture JSON");
922            validator
923                .validate(&json)
924                .expect("golden fixture matches the normative schema");
925            let envelope = Envelope::decode_json(fixture).expect("fixture matches Rust DTOs");
926            assert_eq!(serde_json::to_value(envelope).expect("serializes"), json);
927        }
928    }
929
930    #[test]
931    fn rejects_blank_scoped_identifiers() {
932        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
933        let mut envelope: Envelope = serde_json::from_str(fixture).expect("valid golden fixture");
934        let Payload::Command(request) = &mut envelope.payload else {
935            panic!("enqueue fixture changed family");
936        };
937        let CommandRequest {
938            command: AgentCommand::Enqueue { context, .. },
939            ..
940        } = request.as_mut()
941        else {
942            panic!("enqueue fixture changed family");
943        };
944        context.scope.scope_id = ScopeId::from(" ");
945        assert_eq!(
946            envelope.validate(),
947            Err(ProtocolError::MissingIdentifier("scope_id"))
948        );
949    }
950
951    #[test]
952    fn strict_decoder_checks_other_commands_and_message_parts() {
953        let golden: serde_json::Value = serde_json::from_str(ENQUEUE_FIXTURE).unwrap();
954        let context = golden["payload"]["payload"]["command"]["context"].clone();
955        for command in [
956            serde_json::json!({"type": "cancel", "context": context}),
957            serde_json::json!({"type": "submit_interaction", "context": context,
958                "interaction": {"batch_id": "batch-1", "decisions": [{"action_id": "action-1", "proceed": true}]}}),
959        ] {
960            let mut input = golden.clone();
961            if command["type"] == "cancel" {
962                input["payload"]["payload"]["llm_plan"] = serde_json::Value::Null;
963                input["payload"]["payload"]["context_generation"] = serde_json::Value::Null;
964            }
965            input["payload"]["payload"]["command"] = command;
966            Envelope::decode_json(&input.to_string()).expect("valid command");
967            input["payload"]["payload"]["command"]["unexpected"] = serde_json::Value::Null;
968            assert!(Envelope::decode_json(&input.to_string()).is_err());
969        }
970        for part in [
971            serde_json::json!({"type": "text", "text": "hello"}),
972            serde_json::json!({"type": "artifact", "uri": crate::execution::ArtifactReference::new("tenant-1", "scope-1", "a".repeat(64), crate::execution::ArtifactMetadata::image("image/png", 1, 1, 1).unwrap()).unwrap().uri().unwrap()}),
973            serde_json::json!({"type": "tool_call", "id": "call-1", "name": "test", "arguments": {"arbitrary": {"nested": null}}}),
974            serde_json::json!({"type": "tool_result", "call_id": "call-1", "result": {"arbitrary": true}, "is_error": false}),
975        ] {
976            let mut input = golden.clone();
977            input["payload"]["payload"]["command"]["message"]["content"][0] = part;
978            Envelope::decode_json(&input.to_string()).expect("valid content and opaque JSON");
979            input["payload"]["payload"]["command"]["message"]["content"][0]["unexpected"] =
980                serde_json::Value::Bool(true);
981            assert!(Envelope::decode_json(&input.to_string()).is_err());
982        }
983        assert!(
984            serde_json::from_value::<crate::execution::AgentEvent>(
985                serde_json::json!({"type": "started", "unexpected": true})
986            )
987            .is_err()
988        );
989        assert!(
990            serde_json::from_value::<AppFacadeResponse>(
991                serde_json::json!({"type": "rejected", "unexpected": null})
992            )
993            .is_err()
994        );
995        assert!(
996            serde_json::from_value::<crate::execution::AgentEvent>(
997                serde_json::json!({"type": "started"})
998            )
999            .is_ok()
1000        );
1001        assert!(
1002            serde_json::from_value::<AppFacadeResponse>(serde_json::json!({"type": "rejected"}))
1003                .is_ok()
1004        );
1005    }
1006
1007    #[test]
1008    fn decoder_checks_every_closed_object_in_golden_envelopes() {
1009        fn object_paths(value: &serde_json::Value, pointer: &str, paths: &mut Vec<String>) {
1010            match value {
1011                serde_json::Value::Object(fields) => {
1012                    paths.push(pointer.to_owned());
1013                    for (key, child) in fields {
1014                        let key = key.replace('~', "~0").replace('/', "~1");
1015                        object_paths(child, &format!("{pointer}/{key}"), paths);
1016                    }
1017                }
1018                serde_json::Value::Array(items) => {
1019                    for (index, child) in items.iter().enumerate() {
1020                        object_paths(child, &format!("{pointer}/{index}"), paths);
1021                    }
1022                }
1023                _ => {}
1024            }
1025        }
1026
1027        let schema: serde_json::Value = serde_json::from_str(include_str!(
1028            "../../../schema/execution/envelope.v1.schema.json"
1029        ))
1030        .unwrap();
1031        let validator = jsonschema::validator_for(&schema).unwrap();
1032        let fixtures = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("fixtures/execution");
1033        let mut accepted_unknown = Vec::new();
1034        let mut checked = 0;
1035        for file in std::fs::read_dir(fixtures).unwrap() {
1036            let file = file.unwrap().path();
1037            if file.extension().is_none_or(|extension| extension != "json") {
1038                continue;
1039            }
1040            let input = std::fs::read_to_string(&file).unwrap();
1041            let golden: serde_json::Value = serde_json::from_str(&input).unwrap();
1042            validator
1043                .validate(&golden)
1044                .expect("valid canonical fixture");
1045            Envelope::decode_json(&input).expect("valid fixture decodes");
1046            let mut paths = Vec::new();
1047            object_paths(&golden, "", &mut paths);
1048            for pointer in paths {
1049                for extra in [serde_json::Value::Bool(true), serde_json::Value::Null] {
1050                    let mut input = golden.clone();
1051                    input
1052                        .pointer_mut(&pointer)
1053                        .unwrap()
1054                        .as_object_mut()
1055                        .unwrap()
1056                        .insert("__unknown_probe".to_owned(), extra);
1057                    let decoded = Envelope::decode_json(&input.to_string());
1058                    if validator.is_valid(&input) {
1059                        decoded.expect("protocol-owned opaque JSON must remain extensible");
1060                    } else {
1061                        checked += 1;
1062                        if decoded.is_ok() {
1063                            accepted_unknown.push(format!(
1064                                "{}:{pointer}",
1065                                file.file_name().unwrap().to_string_lossy()
1066                            ));
1067                        }
1068                    }
1069                }
1070            }
1071        }
1072        assert!(checked > 100, "exercise the entire fixture graph");
1073        assert!(
1074            accepted_unknown.is_empty(),
1075            "unknown fields accepted at {accepted_unknown:#?}"
1076        );
1077    }
1078
1079    #[test]
1080    fn admission_requires_one_model_matching_the_selected_mode() {
1081        let mut input: serde_json::Value =
1082            serde_json::from_str(include_str!("../../../fixtures/execution/enqueue.v1.json"))
1083                .unwrap();
1084        input["payload"]["payload"]["llm_plan"]["model_mode"] = "smart".into();
1085        assert!(matches!(
1086            Envelope::decode_json(&input.to_string()),
1087            Err(ProtocolError::InvalidBinding(_))
1088        ));
1089        input["payload"]["payload"]["command"]["options"]["model_mode"] = "smart".into();
1090        assert!(Envelope::decode_json(&input.to_string()).is_ok());
1091        let plan = input["payload"]["payload"]["llm_plan"].take();
1092        input["payload"]["payload"]["llm_plan"] = serde_json::json!({"routes": [plan]});
1093        assert!(Envelope::decode_json(&input.to_string()).is_err());
1094    }
1095
1096    #[test]
1097    fn strict_decoder_rejects_unknown_nested_fields() {
1098        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
1099        let input = fixture.replacen(
1100            "\"sent_at_unix_ms\": 1700000000000",
1101            "\"sent_at_unix_ms\": 1700000000000, \"unexpected\": true",
1102            1,
1103        );
1104        assert!(matches!(
1105            Envelope::decode_json(&input),
1106            Err(ProtocolError::UnknownField { .. })
1107        ));
1108
1109        let mut input: serde_json::Value = serde_json::from_str(fixture).expect("fixture JSON");
1110        input["payload"]["payload"]["options"]["unexpected"] = serde_json::Value::Bool(true);
1111        assert!(matches!(
1112            Envelope::decode_json(&input.to_string()),
1113            Err(ProtocolError::UnknownField { .. })
1114        ));
1115    }
1116
1117    #[test]
1118    fn strict_decoder_treats_skip_optional_nulls_as_absent() {
1119        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
1120        let mut input: serde_json::Value = serde_json::from_str(fixture).expect("fixture JSON");
1121        input["payload"]["payload"]["llm_plan"]["route"]["fast_mode"] = serde_json::Value::Null;
1122        input["payload"]["payload"]["llm_plan"]["route"]["thinking"] = serde_json::Value::Null;
1123        input["payload"]["payload"]["llm_plan"]["route"]["model_snapshot"]["generation_support"]
1124            ["fastMode"] = serde_json::Value::Null;
1125        input["payload"]["payload"]["llm_plan"]["route"]["model_snapshot"]["generation_support"]
1126            ["reasoningEffortDefault"] = serde_json::Value::Null;
1127        input["payload"]["payload"]["command"]["message"]["continuation"] = serde_json::Value::Null;
1128
1129        let envelope = Envelope::decode_json(&input.to_string())
1130            .expect("null on skip-optional fields deserializes as absent");
1131        let canonical = serde_json::to_value(&envelope).expect("serializes");
1132        let canonical_route = &canonical["payload"]["payload"]["llm_plan"]["route"];
1133        assert!(canonical_route.get("fast_mode").is_none());
1134        assert!(canonical_route.get("thinking").is_none());
1135        let canonical_support = &canonical_route["model_snapshot"]["generation_support"];
1136        assert!(canonical_support.get("fastMode").is_none());
1137        assert!(canonical_support.get("reasoningEffortDefault").is_none());
1138        assert!(
1139            canonical["payload"]["payload"]["command"]["message"]
1140                .get("continuation")
1141                .is_none()
1142        );
1143
1144        let schema: serde_json::Value = serde_json::from_str(include_str!(
1145            "../../../schema/execution/envelope.v1.schema.json"
1146        ))
1147        .expect("schema");
1148        let validator = jsonschema::validator_for(&schema).expect("schema");
1149        assert!(
1150            validator.validate(&input).is_err(),
1151            "canonical emit omits skip-optional keys; schema must not advertise null"
1152        );
1153    }
1154
1155    #[test]
1156    fn schema_rejects_generation_support_and_continuation_nulls() {
1157        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
1158        let schema: serde_json::Value = serde_json::from_str(include_str!(
1159            "../../../schema/execution/envelope.v1.schema.json"
1160        ))
1161        .expect("schema");
1162        let validator = jsonschema::validator_for(&schema).expect("schema");
1163        let golden: serde_json::Value = serde_json::from_str(fixture).expect("fixture JSON");
1164        assert!(
1165            validator.validate(&golden).is_ok(),
1166            "golden enqueue fixture must remain emit-valid"
1167        );
1168
1169        let mut support_null = golden.clone();
1170        support_null["payload"]["payload"]["llm_plan"]["route"]["model_snapshot"]["generation_support"]
1171            ["fastMode"] = serde_json::Value::Null;
1172        assert!(
1173            validator.validate(&support_null).is_err(),
1174            "skip-optional generation_support.fastMode must omit, not null"
1175        );
1176
1177        let mut continuation_null = golden;
1178        continuation_null["payload"]["payload"]["command"]["message"]["continuation"] =
1179            serde_json::Value::Null;
1180        assert!(
1181            validator.validate(&continuation_null).is_err(),
1182            "skip-optional message.continuation must omit, not null"
1183        );
1184    }
1185
1186    #[test]
1187    fn strict_decoder_still_rejects_unknown_null_fields() {
1188        let fixture = include_str!("../../../fixtures/execution/enqueue.v1.json");
1189        let mut input: serde_json::Value = serde_json::from_str(fixture).expect("fixture JSON");
1190        input["payload"]["payload"]["llm_plan"]["route"]["unexpected"] = serde_json::Value::Null;
1191        assert!(matches!(
1192            Envelope::decode_json(&input.to_string()),
1193            Err(ProtocolError::UnknownField { .. })
1194        ));
1195    }
1196
1197    #[test]
1198    fn abandoned_run_requires_a_base_checkpoint() {
1199        let mut value: serde_json::Value =
1200            serde_json::from_str(include_str!("../../../fixtures/execution/enqueue.v1.json"))
1201                .unwrap();
1202        value["payload"]["payload"]["abandoned_run_id"] =
1203            serde_json::Value::String("run-abandoned".to_owned());
1204
1205        assert!(matches!(
1206            Envelope::decode_json(&value.to_string()),
1207            Err(ProtocolError::InvalidBinding(_))
1208        ));
1209    }
1210}