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,
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 sequence: u64,
136 pub context: InvocationContext,
138 pub observation: AgentObservation,
140}
141
142#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
144pub struct ExecutionFailure {
145 pub context: InvocationContext,
146 pub code: String,
147 pub message: String,
148}
149
150#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
152pub struct EventAck {
153 pub event_id: EventId,
154 pub checkpoint: CheckpointReference,
156}
157
158#[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#[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#[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
197pub use llm_api::GenerationSupport as LlmGenerationSupport;
200
201#[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 pub context_window_tokens: u32,
219 pub model_snapshot: LlmModelSnapshot,
221}
222
223#[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#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
235pub struct CommandRequest {
236 pub command: AgentCommand,
237 pub context_generation: Option<ContextGeneration>,
239 pub llm_plan: Option<LlmExecutionPlan>,
240 pub base_checkpoint: Option<CheckpointReference>,
242 pub abandoned_run_id: Option<RunId>,
244}
245
246#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
248pub struct EventMessage {
249 pub event: DurableEvent,
250 pub checkpoint: CheckpointReference,
251}
252
253#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
255#[serde(tag = "type", rename_all = "snake_case")]
256#[serde(deny_unknown_fields)]
257pub enum AppFacadeResponse {
258 Result {
260 result: FacadeResult,
262 },
263 Prepared {
265 action: PreparedAction,
267 },
268 #[serde(deserialize_with = "crate::execution::deserialize_empty_variant")]
270 Rejected,
271 Failed {
273 code: String,
275 message: String,
277 retry_after_ms: Option<u64>,
279 },
280}
281
282#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
284#[serde(tag = "family", content = "payload", rename_all = "snake_case")]
285pub enum Payload {
286 Command(Box<CommandRequest>),
288 CommandResponse(CommandResponse),
290 Event(EventMessage),
292 EventAck(EventAck),
294 CheckpointPurge(CheckpointPurgeRequest),
296 CheckpointPurgeAck(CheckpointPurgeAck),
298 Observation(ObservationMessage),
300 ExecutionFailure(ExecutionFailure),
302 AppFacadeRequest(AppFacadeRequest),
304 AppFacadeResponse(AppFacadeResponse),
306}
307
308#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
310pub struct Envelope {
311 pub metadata: Metadata,
313 pub payload: Payload,
315}
316
317impl Envelope {
318 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 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#[derive(Clone, Debug, Error, Eq, PartialEq)]
641pub enum ProtocolError {
642 #[error("invalid protocol JSON: {message}")]
644 InvalidJson {
645 message: String,
647 },
648 #[error("unknown protocol field: {path}")]
650 UnknownField {
651 path: String,
653 },
654 #[error("unsupported protocol version {actual}")]
656 UnsupportedVersion {
657 actual: u16,
659 },
660 #[error("missing required identifier: {0}")]
662 MissingIdentifier(&'static str),
663 #[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}