1pub 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
17pub const CURRENT_VERSION: u16 = 1;
19
20#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
24pub struct DispatchBinding {
25 pub dispatch_id: String,
27 pub agent_session_id: String,
29}
30
31#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
33pub struct CheckpointReference {
34 pub object_key: String,
36 pub sha256: String,
38 pub format_version: u32,
40}
41
42#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
44pub struct Metadata {
45 pub version: u16,
47 pub message_id: String,
49 pub correlation_id: String,
51 pub causation_id: Option<String>,
53 pub sent_at_unix_ms: u64,
55 pub dispatch: DispatchBinding,
57}
58
59#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
61#[serde(tag = "phase", rename_all = "snake_case")]
62#[serde(deny_unknown_fields)]
63pub enum AppFacadeRequest {
64 Invoke {
66 context: InvocationContext,
68 request: FacadeRequest,
70 },
71 Prepare {
73 context: InvocationContext,
75 request: FacadeRequest,
77 },
78 Commit {
80 context: InvocationContext,
82 action: PreparedAction,
84 },
85 Reject {
87 context: InvocationContext,
89 action: PreparedAction,
91 },
92}
93
94#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
96#[serde(tag = "type", rename_all = "snake_case")]
97#[serde(deny_unknown_fields)]
98pub enum CommandResponse {
99 Accepted {
101 context: InvocationContext,
103 disposition: AdmissionDisposition,
105 queue_revision: Option<Revision>,
107 },
108 Rejected {
110 context: InvocationContext,
112 code: String,
114 message: String,
116 },
117}
118
119#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
121#[serde(rename_all = "snake_case")]
122pub enum AdmissionDisposition {
123 Queued,
125 InteractionReady,
127 CancelRequested,
129}
130
131#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
133pub struct ObservationMessage {
134 pub context: InvocationContext,
136 pub observation: AgentObservation,
138}
139
140#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
142pub struct ExecutionFailure {
143 pub context: InvocationContext,
144 pub code: String,
145 pub message: String,
146}
147
148#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
150pub struct EventAck {
151 pub event_id: EventId,
152 pub checkpoint: CheckpointReference,
154}
155
156#[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#[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#[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
195pub use llm_api::GenerationSupport as LlmGenerationSupport;
198
199#[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 pub context_window_tokens: u32,
217 pub model_snapshot: LlmModelSnapshot,
219}
220
221#[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#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
232pub struct LlmExecutionPlan {
233 pub routes: Vec<LlmRouteBinding>,
234}
235
236impl LlmExecutionPlan {
237 #[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#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
253pub struct CommandRequest {
254 pub command: AgentCommand,
255 pub context_generation: Option<ContextGeneration>,
257 pub llm_plan: Option<LlmExecutionPlan>,
258 pub base_checkpoint: Option<CheckpointReference>,
260 pub abandoned_run_id: Option<RunId>,
262}
263
264#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
266pub struct EventMessage {
267 pub event: DurableEvent,
268 pub checkpoint: CheckpointReference,
269}
270
271#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
273#[serde(tag = "type", rename_all = "snake_case")]
274#[serde(deny_unknown_fields)]
275pub enum AppFacadeResponse {
276 Result {
278 result: FacadeResult,
280 },
281 Prepared {
283 action: PreparedAction,
285 },
286 #[serde(deserialize_with = "crate::execution::deserialize_empty_variant")]
288 Rejected,
289 Failed {
291 code: String,
293 message: String,
295 retry_after_ms: Option<u64>,
297 },
298}
299
300#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
302#[serde(tag = "family", content = "payload", rename_all = "snake_case")]
303pub enum Payload {
304 Command(Box<CommandRequest>),
306 CommandResponse(CommandResponse),
308 Event(EventMessage),
310 EventAck(EventAck),
312 CheckpointPurge(CheckpointPurgeRequest),
314 CheckpointPurgeAck(CheckpointPurgeAck),
316 Observation(ObservationMessage),
318 ExecutionFailure(ExecutionFailure),
320 AppFacadeRequest(AppFacadeRequest),
322 AppFacadeResponse(AppFacadeResponse),
324}
325
326#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
328pub struct Envelope {
329 pub metadata: Metadata,
331 pub payload: Payload,
333}
334
335impl Envelope {
336 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 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#[derive(Clone, Debug, Error, Eq, PartialEq)]
671pub enum ProtocolError {
672 #[error("invalid protocol JSON: {message}")]
674 InvalidJson {
675 message: String,
677 },
678 #[error("unknown protocol field: {path}")]
680 UnknownField {
681 path: String,
683 },
684 #[error("unsupported protocol version {actual}")]
686 UnsupportedVersion {
687 actual: u16,
689 },
690 #[error("missing required identifier: {0}")]
692 MissingIdentifier(&'static str),
693 #[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}