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