1#![forbid(unsafe_code)]
9
10use std::path::PathBuf;
11
12use serde::{
13 de::{Error as _, MapAccess, SeqAccess, Visitor},
14 ser::SerializeMap,
15 Deserialize, Deserializer, Serialize, Serializer,
16};
17use subc_protocol::{
18 manifest::{CapabilityDeclarations, ManifestProvenance, ProviderRole, SelfSignalDeclaration},
19 session::HealthStatus,
20 BindIdentity, RouteTarget,
21};
22
23pub use subc_protocol::RouteCloseReason;
24
25macro_rules! open_string_enum {
26 (
27 $(#[$meta:meta])*
28 $name:ident {
29 $( $variant:ident => $wire_name:literal ),+ $(,)?
30 }
31 ) => {
32 $(#[$meta])*
33 #[derive(Debug, Clone, PartialEq, Eq)]
34 pub enum $name {
35 $( $variant, )+
36 Unknown(String),
37 }
38
39 impl $name {
40 fn wire_name(&self) -> &str {
41 match self {
42 $( Self::$variant => $wire_name, )+
43 Self::Unknown(value) => value,
44 }
45 }
46 }
47
48 impl Serialize for $name {
49 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
50 where
51 S: serde::Serializer,
52 {
53 serializer.serialize_str(self.wire_name())
54 }
55 }
56
57 impl<'de> Deserialize<'de> for $name {
58 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
59 where
60 D: serde::Deserializer<'de>,
61 {
62 let value = String::deserialize(deserializer)?;
63 Ok(match value.as_str() {
64 $( $wire_name => Self::$variant, )+
65 _ => Self::Unknown(value),
66 })
67 }
68 }
69 };
70}
71
72#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
74pub struct ConsumerIdentity {
75 pub module_id: String,
76 pub launch_nonce: String,
77}
78
79pub mod ops {
90 pub const SERVER: &str = "server.";
91 pub const CATALOG: &str = "catalog.";
92 pub const ROUTE: &str = "route.";
93 pub const SUPERVISOR: &str = "supervisor.";
94 pub const CONFIG: &str = "config.";
95
96 pub const SERVER_DESCRIBE: &str = "server.describe";
97 pub const CATALOG_LIST: &str = "catalog.list";
98 pub const ROUTE_OPEN: &str = "route.open";
99 pub const ROUTE_POLL: &str = "route.poll";
100 pub const ROUTE_CLOSING: &str = "route.closing";
101 pub const ROUTE_CLOSED: &str = "route.closed";
102 pub const SUPERVISOR_LIST: &str = "supervisor.list";
103 pub const SUPERVISOR_RESTART: &str = "supervisor.restart";
104 pub const SUPERVISOR_RELOAD: &str = "supervisor.reload";
105 pub const SUPERVISOR_RESCAN: &str = "supervisor.rescan";
106 pub const SUPERVISOR_RELEASE_RESERVED: &str = "supervisor.release_reserved";
107 pub const SUPERVISOR_SET_ENABLED: &str = "supervisor.set_enabled";
108 pub const SUPERVISOR_HEALTH_PROBE: &str = "supervisor.health_probe";
109 pub const SUPERVISOR_HEALTH: &str = "supervisor.health";
110 pub const SUPERVISOR_STDERR_TAIL: &str = "supervisor.stderr_tail";
111 pub const SUPERVISOR_TERMINALS: &str = "supervisor.terminals";
112 pub const SUPERVISOR_ROUTES: &str = "supervisor.routes";
113 pub const SUPERVISOR_PROVENANCE: &str = "supervisor.provenance";
114}
115
116#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
118#[serde(tag = "op")]
119#[allow(clippy::large_enum_variant)]
122pub enum ClientControlRequest {
123 #[serde(rename = "server.describe")]
124 ServerDescribe {},
125 #[serde(rename = "catalog.list")]
126 CatalogList {
127 #[serde(default)]
132 module_id: Option<String>,
133 },
134 #[serde(rename = "route.open")]
135 RouteOpen {
136 target: RouteTarget,
137 identity: BindIdentity,
138 #[serde(default, skip_serializing_if = "Option::is_none")]
147 consumer_identity: Option<ConsumerIdentity>,
148 #[serde(default, skip_serializing_if = "Option::is_none")]
156 consumer_capabilities: Option<Vec<String>>,
157 #[serde(default, skip_serializing_if = "Option::is_none")]
159 admission_facts: Option<serde_json::Value>,
160 },
161 #[serde(rename = "route.poll")]
162 RoutePoll {
163 route_channel: u16,
164 route_epoch: u32,
165 kind: PollKind,
166 },
167 #[serde(rename = "supervisor.list")]
168 SupervisorList {},
169 #[serde(rename = "supervisor.restart")]
170 SupervisorRestart {
171 module_id: String,
172 #[serde(default, skip_serializing_if = "Option::is_none")]
180 drain_timeout_ms: Option<u64>,
181 },
182 #[serde(rename = "supervisor.reload")]
183 SupervisorReload { module_id: String },
184 #[serde(rename = "supervisor.rescan")]
185 SupervisorRescan {
186 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
206 preview: bool,
207 },
208 #[serde(rename = "supervisor.release_reserved")]
212 SupervisorReleaseReserved { module_id: String },
213 #[serde(rename = "supervisor.set_enabled")]
214 SupervisorSetEnabled { module_id: String, enabled: bool },
215 #[serde(rename = "supervisor.health_probe")]
216 SupervisorHealthProbe { module_id: String },
217 #[serde(rename = "supervisor.health")]
218 SupervisorHealth {},
219 #[serde(rename = "supervisor.routes")]
232 SupervisorRoutes {
233 #[serde(default, skip_serializing_if = "Option::is_none")]
234 module_id: Option<String>,
235 },
236 #[serde(rename = "supervisor.provenance")]
239 SupervisorProvenance {
240 #[serde(default, skip_serializing_if = "Option::is_none")]
241 module_id: Option<String>,
242 },
243 #[serde(rename = "supervisor.stderr_tail")]
251 SupervisorStderrTail {
252 module_id: String,
253 #[serde(default, skip_serializing_if = "Option::is_none")]
254 max_lines: Option<u32>,
255 #[serde(default, skip_serializing_if = "Option::is_none")]
256 max_bytes: Option<u32>,
257 },
258 #[serde(rename = "supervisor.terminals")]
271 SupervisorTerminals { module_id: String },
272}
273
274#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
276#[serde(tag = "op")]
277pub enum ClientControlResponse {
278 #[serde(rename = "server.describe")]
279 ServerDescribe {
280 protocol_ver: u8,
281 subc_ops: Vec<String>,
282 capabilities: Vec<String>,
283 connected_clients: u64,
284 #[serde(default, skip_serializing_if = "Option::is_none")]
285 counters: Option<serde_json::Value>,
286 #[serde(default, skip_serializing_if = "Option::is_none")]
293 build_git_sha: Option<String>,
294 #[serde(default, skip_serializing_if = "Option::is_none")]
300 build_lock_digest: Option<String>,
301 #[serde(default, skip_serializing_if = "Vec::is_empty")]
305 capability_requirements: Vec<CapabilityRequirementStatus>,
306 },
307 #[serde(rename = "catalog.list")]
308 CatalogList {
309 generation: u64,
310 modules: Vec<CatalogEntry>,
311 subc_ops: Vec<String>,
312 },
313 #[serde(rename = "route.open")]
314 RouteOpen {
315 route_channel: u16,
316 route_epoch: u32,
317 },
318 #[serde(rename = "route.poll")]
319 RoutePoll {
320 route_channel: u16,
321 route_epoch: u32,
322 status: Option<String>,
323 live: Option<bool>,
324 },
325 #[serde(rename = "supervisor.list")]
326 SupervisorList {
327 generation: u64,
328 modules: Vec<SupervisorEntry>,
329 },
330 #[serde(rename = "supervisor.ack")]
331 SupervisorAck { module_id: String, applied: bool },
332 #[serde(rename = "supervisor.rescan")]
333 SupervisorRescan {
334 #[serde(flatten)]
335 result: SupervisorRescanResult,
336 },
337 #[serde(rename = "supervisor.health_probe")]
338 SupervisorHealthProbe {
339 module_id: String,
340 status: HealthStatus,
341 #[serde(default, skip_serializing_if = "Option::is_none")]
342 detail: Option<String>,
343 #[serde(default, skip_serializing_if = "Option::is_none")]
344 metrics: Option<serde_json::Value>,
345 },
346 #[serde(rename = "supervisor.health")]
347 SupervisorHealth {
348 generation: u64,
349 modules: Vec<SupervisorHealthEntry>,
350 },
351 #[serde(rename = "supervisor.routes")]
352 SupervisorRoutes { modules: Vec<SupervisorRouteModule> },
353 #[serde(rename = "supervisor.provenance")]
354 SupervisorProvenance {
355 daemon: SupervisorDaemonProvenance,
356 modules: Vec<SupervisorModuleProvenance>,
357 },
358 #[serde(rename = "supervisor.stderr_tail")]
359 SupervisorStderrTail {
360 module_id: String,
361 #[serde(flatten)]
362 tail: StderrTail,
363 },
364 #[serde(rename = "supervisor.terminals")]
365 SupervisorTerminals {
366 module_id: String,
367 #[serde(flatten)]
368 terminals: TerminalHistory,
369 },
370}
371
372#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
377#[serde(tag = "op")]
378pub enum ClientControlPush {
379 #[serde(rename = "route.closing")]
380 RouteClosing {
381 module_id: String,
382 reason: RouteCloseReason,
383 },
384 #[serde(rename = "route.closed")]
385 RouteClosed {
386 module_id: String,
387 reason: RouteCloseReason,
388 drained: bool,
390 abandoned: u32,
393 #[serde(default)]
395 excluded_subscriptions: u32,
396 #[serde(default, skip_serializing_if = "Option::is_none")]
402 terminal: Option<bool>,
403 },
404}
405
406#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
408pub struct StderrTail {
409 pub capture: StderrCaptureState,
410 pub entries: Vec<StderrTailEntry>,
411 #[serde(default, skip_serializing_if = "is_zero_u64")]
420 pub dropped_lines: u64,
421}
422
423#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
425pub struct SupervisorRouteModule {
426 pub module_id: String,
427 pub routes: Vec<SupervisorRoute>,
428}
429
430#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
432pub struct SupervisorRoute {
433 pub consumer: SupervisorRouteConsumer,
434 pub age_ms: u64,
436 pub draining: bool,
439 #[serde(default, skip_serializing_if = "Option::is_none")]
445 pub drain_reason: Option<RouteCloseReason>,
446}
447
448#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
450pub struct SupervisorModuleProvenance {
451 pub module_id: String,
452 pub module_declared: ModuleDeclaredProvenance,
453 pub daemon_observed: SupervisorObservedProcess,
454}
455
456#[derive(Debug, Clone, PartialEq)]
458pub enum ModuleDeclaredProvenance {
459 Reported {
460 build: ManifestProvenance,
461 },
462 Unverifiable,
463 Unknown {
466 tag: String,
467 body: OrderedJsonObject,
468 },
469}
470
471#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
476pub struct SupervisorObservedProcess {
477 #[serde(default, skip_serializing_if = "Option::is_none")]
478 pub pid: Option<u32>,
479 #[serde(default, skip_serializing_if = "Option::is_none")]
480 pub spawned_at_ms: Option<u64>,
481 #[serde(default, skip_serializing_if = "Option::is_none")]
482 pub spawned_from: Option<PathBuf>,
483 pub running_image: RunningImageAgreement,
484}
485
486#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
488pub struct SupervisorDaemonProvenance {
489 pub daemon_build: DaemonBuildProvenance,
490 pub daemon_observed: DaemonObservedProcess,
491}
492
493#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
495pub struct DaemonBuildProvenance {
496 #[serde(default, skip_serializing_if = "Option::is_none")]
497 pub build_git_sha: Option<String>,
498 #[serde(default, skip_serializing_if = "Option::is_none")]
499 pub build_lock_digest: Option<String>,
500}
501
502#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
504pub struct DaemonObservedProcess {
505 #[serde(default, skip_serializing_if = "Option::is_none")]
506 pub pid: Option<u32>,
507 #[serde(default, skip_serializing_if = "Option::is_none")]
508 pub started_at_ms: Option<u64>,
509 pub running_image: RunningImageAgreement,
510}
511
512#[derive(Debug, Clone, PartialEq)]
514pub enum RunningImageAgreement {
515 Match {
516 evidence: RunningImageEvidence,
517 },
518 Mismatch {
519 running: RunningImageEvidence,
520 disk: RunningImageEvidence,
521 },
522 Unavailable {
523 reason: RunningImageUnavailableReason,
524 },
525 Unknown {
528 tag: String,
529 body: OrderedJsonObject,
530 },
531}
532
533#[derive(Debug, Clone, PartialEq)]
535pub enum RunningImageEvidence {
536 LinuxProcSha256 {
537 digest: String,
538 },
539 MacosSpawnInode {
540 device: u64,
541 inode: u64,
542 },
543 Unknown {
546 tag: String,
547 body: OrderedJsonObject,
548 },
549}
550
551open_string_enum! {
552 RunningImageUnavailableReason {
554 NotRunning => "not_running",
555 UnsupportedPlatform => "unsupported_platform",
556 RunningExecutableUnreadable => "running_executable_unreadable",
557 SpawnedPathUnreadable => "spawned_path_unreadable",
558 HashFailed => "hash_failed",
559 ProcessIdentityUnconfirmed => "process_identity_unconfirmed",
560 }
561}
562
563#[derive(Debug, Clone, PartialEq)]
569pub enum SupervisorRouteConsumer {
570 Reserved {
571 module_id: String,
572 },
573 Direct {
574 connection_id: u64,
575 },
576 Unknown {
579 tag: String,
580 body: OrderedJsonObject,
581 },
582}
583
584#[derive(Debug, Clone, PartialEq)]
591pub enum StderrCaptureState {
592 Captured,
595 Incomplete { reason: String },
597 NotCaptured { reason: String },
599 Unknown {
602 tag: String,
603 body: OrderedJsonObject,
604 },
605}
606
607#[derive(Debug, Clone, PartialEq)]
608pub enum StderrTailEntry {
609 Line {
610 text: String,
611 truncated: bool,
616 },
617 ProcessStart,
622 Unknown {
625 tag: String,
626 body: OrderedJsonObject,
627 },
628}
629
630#[derive(Debug, Serialize, Deserialize)]
631#[serde(tag = "status", rename_all = "snake_case")]
632enum ModuleDeclaredProvenanceWire {
633 Reported { build: ManifestProvenance },
634 Unverifiable,
635}
636
637#[derive(Debug, Serialize, Deserialize)]
638#[serde(tag = "status", rename_all = "snake_case")]
639enum RunningImageAgreementWire {
640 Match {
641 evidence: RunningImageEvidence,
642 },
643 Mismatch {
644 running: RunningImageEvidence,
645 disk: RunningImageEvidence,
646 },
647 Unavailable {
648 reason: RunningImageUnavailableReason,
649 },
650}
651
652#[derive(Debug, Serialize, Deserialize)]
653#[serde(tag = "method", rename_all = "snake_case")]
654enum RunningImageEvidenceWire {
655 LinuxProcSha256 { digest: String },
656 MacosSpawnInode { device: u64, inode: u64 },
657}
658
659#[derive(Debug, Serialize, Deserialize)]
660#[serde(tag = "kind", rename_all = "snake_case")]
661enum SupervisorRouteConsumerWire {
662 Reserved { module_id: String },
663 Direct { connection_id: u64 },
664}
665
666#[derive(Debug, Serialize, Deserialize)]
667#[serde(tag = "state", rename_all = "snake_case")]
668enum StderrCaptureStateWire {
669 Captured,
670 Incomplete { reason: String },
671 NotCaptured { reason: String },
672}
673
674#[derive(Debug, Serialize, Deserialize)]
675#[serde(tag = "kind", rename_all = "snake_case")]
676enum StderrTailEntryWire {
677 Line {
678 text: String,
679 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
680 truncated: bool,
681 },
682 ProcessStart,
683}
684
685#[derive(Debug, Clone, PartialEq)]
687pub enum OrderedJsonValue {
688 Null,
689 Bool(bool),
690 Number(serde_json::Number),
691 String(String),
692 Array(Vec<Self>),
693 Object(OrderedJsonObject),
694}
695
696#[derive(Debug, Clone, PartialEq)]
698pub struct OrderedJsonObject(Vec<(String, OrderedJsonValue)>);
699
700impl OrderedJsonObject {
701 pub fn as_entries(&self) -> &[(String, OrderedJsonValue)] {
703 &self.0
704 }
705
706 fn into_value(self) -> serde_json::Value {
707 serde_json::Value::Object(
708 self.0
709 .into_iter()
710 .map(|(key, value)| (key, value.into_value()))
711 .collect(),
712 )
713 }
714}
715
716impl OrderedJsonValue {
717 fn into_value(self) -> serde_json::Value {
718 match self {
719 Self::Null => serde_json::Value::Null,
720 Self::Bool(value) => serde_json::Value::Bool(value),
721 Self::Number(value) => serde_json::Value::Number(value),
722 Self::String(value) => serde_json::Value::String(value),
723 Self::Array(values) => {
724 serde_json::Value::Array(values.into_iter().map(Self::into_value).collect())
725 }
726 Self::Object(value) => value.into_value(),
727 }
728 }
729}
730
731impl Serialize for OrderedJsonValue {
732 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
733 where
734 S: Serializer,
735 {
736 match self {
737 Self::Null => serializer.serialize_unit(),
738 Self::Bool(value) => serializer.serialize_bool(*value),
739 Self::Number(value) => value.serialize(serializer),
740 Self::String(value) => serializer.serialize_str(value),
741 Self::Array(values) => values.serialize(serializer),
742 Self::Object(value) => value.serialize(serializer),
743 }
744 }
745}
746
747impl<'de> Deserialize<'de> for OrderedJsonValue {
748 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
749 where
750 D: Deserializer<'de>,
751 {
752 struct OrderedValueVisitor;
753
754 impl<'de> Visitor<'de> for OrderedValueVisitor {
755 type Value = OrderedJsonValue;
756
757 fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
758 formatter.write_str("a JSON value with ordered object members")
759 }
760
761 fn visit_unit<E>(self) -> Result<Self::Value, E>
762 where
763 E: serde::de::Error,
764 {
765 Ok(OrderedJsonValue::Null)
766 }
767
768 fn visit_none<E>(self) -> Result<Self::Value, E>
769 where
770 E: serde::de::Error,
771 {
772 Ok(OrderedJsonValue::Null)
773 }
774
775 fn visit_some<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
776 where
777 D: Deserializer<'de>,
778 {
779 OrderedJsonValue::deserialize(deserializer)
780 }
781
782 fn visit_bool<E>(self, value: bool) -> Result<Self::Value, E>
783 where
784 E: serde::de::Error,
785 {
786 Ok(OrderedJsonValue::Bool(value))
787 }
788
789 fn visit_i64<E>(self, value: i64) -> Result<Self::Value, E>
790 where
791 E: serde::de::Error,
792 {
793 Ok(OrderedJsonValue::Number(value.into()))
794 }
795
796 fn visit_u64<E>(self, value: u64) -> Result<Self::Value, E>
797 where
798 E: serde::de::Error,
799 {
800 Ok(OrderedJsonValue::Number(value.into()))
801 }
802
803 fn visit_f64<E>(self, value: f64) -> Result<Self::Value, E>
804 where
805 E: serde::de::Error,
806 {
807 serde_json::Number::from_f64(value)
808 .map(OrderedJsonValue::Number)
809 .ok_or_else(|| E::custom("non-finite JSON number"))
810 }
811
812 fn visit_str<E>(self, value: &str) -> Result<Self::Value, E>
813 where
814 E: serde::de::Error,
815 {
816 Ok(OrderedJsonValue::String(value.to_owned()))
817 }
818
819 fn visit_string<E>(self, value: String) -> Result<Self::Value, E>
820 where
821 E: serde::de::Error,
822 {
823 Ok(OrderedJsonValue::String(value))
824 }
825
826 fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
827 where
828 A: SeqAccess<'de>,
829 {
830 let mut values = Vec::new();
831 while let Some(value) = sequence.next_element()? {
832 values.push(value);
833 }
834 Ok(OrderedJsonValue::Array(values))
835 }
836
837 fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
838 where
839 A: MapAccess<'de>,
840 {
841 let mut entries = Vec::new();
842 while let Some((key, value)) = map.next_entry()? {
843 entries.push((key, value));
844 }
845 Ok(OrderedJsonValue::Object(OrderedJsonObject(entries)))
846 }
847 }
848
849 deserializer.deserialize_any(OrderedValueVisitor)
850 }
851}
852
853impl Serialize for OrderedJsonObject {
854 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
855 where
856 S: Serializer,
857 {
858 let mut map = serializer.serialize_map(Some(self.0.len()))?;
859 for (key, value) in &self.0 {
860 map.serialize_entry(key, value)?;
861 }
862 map.end()
863 }
864}
865
866impl<'de> Deserialize<'de> for OrderedJsonObject {
867 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
868 where
869 D: Deserializer<'de>,
870 {
871 struct OrderedObjectVisitor;
872
873 impl<'de> Visitor<'de> for OrderedObjectVisitor {
874 type Value = OrderedJsonObject;
875
876 fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
877 formatter.write_str("an object with ordered JSON members")
878 }
879
880 fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
881 where
882 A: MapAccess<'de>,
883 {
884 let mut entries = Vec::new();
885 while let Some((key, value)) = map.next_entry()? {
886 entries.push((key, value));
887 }
888 Ok(OrderedJsonObject(entries))
889 }
890 }
891
892 deserializer.deserialize_map(OrderedObjectVisitor)
893 }
894}
895
896fn read_tagged<'de, D>(
897 deserializer: D,
898 field: &'static str,
899) -> Result<(String, OrderedJsonObject), D::Error>
900where
901 D: Deserializer<'de>,
902{
903 let body = OrderedJsonObject::deserialize(deserializer)?;
904 let mut tag = None;
905 for (key, value) in body.as_entries() {
906 if key != field {
907 continue;
908 }
909 if tag.is_some() {
910 return Err(D::Error::custom(format!(
911 "tagged object has duplicate `{field}` field"
912 )));
913 }
914 let OrderedJsonValue::String(value) = value else {
915 return Err(D::Error::custom(format!(
916 "tagged object has no string `{field}` field"
917 )));
918 };
919 tag = Some(value);
920 }
921 let Some(tag) = tag else {
922 return Err(D::Error::custom(format!(
923 "tagged object has no string `{field}` field"
924 )));
925 };
926 Ok((tag.to_string(), body))
927}
928
929fn read_ordered_tagged(
930 value: OrderedJsonValue,
931 field: &'static str,
932) -> Result<(String, OrderedJsonObject), String> {
933 let OrderedJsonValue::Object(body) = value else {
934 return Err(format!("expected tagged object with `{field}` field"));
935 };
936 let mut tag = None;
937 for (key, value) in body.as_entries() {
938 if key != field {
939 continue;
940 }
941 if tag.is_some() {
942 return Err(format!("tagged object has duplicate `{field}` field"));
943 }
944 let OrderedJsonValue::String(value) = value else {
945 return Err(format!("tagged object has no string `{field}` field"));
946 };
947 tag = Some(value);
948 }
949 let Some(tag) = tag else {
950 return Err(format!("tagged object has no string `{field}` field"));
951 };
952 Ok((tag.to_string(), body))
953}
954
955fn ordered_field<'a>(body: &'a OrderedJsonObject, field: &str) -> Option<&'a OrderedJsonValue> {
956 body.as_entries()
957 .iter()
958 .find_map(|(key, value)| (key == field).then_some(value))
959}
960
961fn ordered_string(body: &OrderedJsonObject, field: &str) -> Result<String, String> {
962 match ordered_field(body, field) {
963 Some(OrderedJsonValue::String(value)) => Ok(value.clone()),
964 Some(_) => Err(format!("tagged object field `{field}` is not a string")),
965 None => Err(format!("tagged object has no `{field}` field")),
966 }
967}
968
969fn decode_running_image_evidence(value: OrderedJsonValue) -> Result<RunningImageEvidence, String> {
970 let (tag, body) = read_ordered_tagged(value, "method")?;
971 match tag.as_str() {
972 "linux_proc_sha256" => Ok(RunningImageEvidence::LinuxProcSha256 {
973 digest: ordered_string(&body, "digest")?,
974 }),
975 "macos_spawn_inode" => {
976 let device = ordered_field(&body, "device")
977 .and_then(|value| match value {
978 OrderedJsonValue::Number(number) => number.as_u64(),
979 _ => None,
980 })
981 .ok_or_else(|| "tagged object has no unsigned `device` field".to_string())?;
982 let inode = ordered_field(&body, "inode")
983 .and_then(|value| match value {
984 OrderedJsonValue::Number(number) => number.as_u64(),
985 _ => None,
986 })
987 .ok_or_else(|| "tagged object has no unsigned `inode` field".to_string())?;
988 Ok(RunningImageEvidence::MacosSpawnInode { device, inode })
989 }
990 _ => Ok(RunningImageEvidence::Unknown { tag, body }),
991 }
992}
993
994impl Serialize for ModuleDeclaredProvenance {
995 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
996 where
997 S: Serializer,
998 {
999 match self {
1000 Self::Reported { build } => ModuleDeclaredProvenanceWire::Reported {
1001 build: build.clone(),
1002 }
1003 .serialize(serializer),
1004 Self::Unverifiable => ModuleDeclaredProvenanceWire::Unverifiable.serialize(serializer),
1005 Self::Unknown { body, .. } => body.serialize(serializer),
1006 }
1007 }
1008}
1009
1010impl<'de> Deserialize<'de> for ModuleDeclaredProvenance {
1011 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1012 where
1013 D: serde::Deserializer<'de>,
1014 {
1015 let (tag, value) = read_tagged(deserializer, "status")?;
1016 match tag.as_str() {
1017 "reported" => match serde_json::from_value(value.into_value())
1018 .map_err(D::Error::custom)?
1019 {
1020 ModuleDeclaredProvenanceWire::Reported { build } => Ok(Self::Reported { build }),
1021 ModuleDeclaredProvenanceWire::Unverifiable => unreachable!(),
1022 },
1023 "unverifiable" => {
1024 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1025 ModuleDeclaredProvenanceWire::Unverifiable => Ok(Self::Unverifiable),
1026 ModuleDeclaredProvenanceWire::Reported { .. } => unreachable!(),
1027 }
1028 }
1029 _ => Ok(Self::Unknown { tag, body: value }),
1030 }
1031 }
1032}
1033
1034impl Serialize for RunningImageAgreement {
1035 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1036 where
1037 S: Serializer,
1038 {
1039 match self {
1040 Self::Match { evidence } => RunningImageAgreementWire::Match {
1041 evidence: evidence.clone(),
1042 }
1043 .serialize(serializer),
1044 Self::Mismatch { running, disk } => RunningImageAgreementWire::Mismatch {
1045 running: running.clone(),
1046 disk: disk.clone(),
1047 }
1048 .serialize(serializer),
1049 Self::Unavailable { reason } => RunningImageAgreementWire::Unavailable {
1050 reason: reason.clone(),
1051 }
1052 .serialize(serializer),
1053 Self::Unknown { body, .. } => body.serialize(serializer),
1054 }
1055 }
1056}
1057
1058impl<'de> Deserialize<'de> for RunningImageAgreement {
1059 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1060 where
1061 D: serde::Deserializer<'de>,
1062 {
1063 let (tag, value) = read_tagged(deserializer, "status")?;
1064 match tag.as_str() {
1065 "match" => Ok(Self::Match {
1066 evidence: decode_running_image_evidence(
1067 ordered_field(&value, "evidence")
1068 .cloned()
1069 .ok_or_else(|| D::Error::custom("tagged object has no `evidence` field"))?,
1070 )
1071 .map_err(D::Error::custom)?,
1072 }),
1073 "mismatch" => Ok(Self::Mismatch {
1074 running: decode_running_image_evidence(
1075 ordered_field(&value, "running")
1076 .cloned()
1077 .ok_or_else(|| D::Error::custom("tagged object has no `running` field"))?,
1078 )
1079 .map_err(D::Error::custom)?,
1080 disk: decode_running_image_evidence(
1081 ordered_field(&value, "disk")
1082 .cloned()
1083 .ok_or_else(|| D::Error::custom("tagged object has no `disk` field"))?,
1084 )
1085 .map_err(D::Error::custom)?,
1086 }),
1087 "unavailable" => Ok(Self::Unavailable {
1088 reason: serde_json::from_value(
1089 ordered_field(&value, "reason")
1090 .cloned()
1091 .ok_or_else(|| D::Error::custom("tagged object has no `reason` field"))?
1092 .into_value(),
1093 )
1094 .map_err(D::Error::custom)?,
1095 }),
1096 _ => Ok(Self::Unknown { tag, body: value }),
1097 }
1098 }
1099}
1100
1101impl Serialize for RunningImageEvidence {
1102 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1103 where
1104 S: Serializer,
1105 {
1106 match self {
1107 Self::LinuxProcSha256 { digest } => RunningImageEvidenceWire::LinuxProcSha256 {
1108 digest: digest.clone(),
1109 }
1110 .serialize(serializer),
1111 Self::MacosSpawnInode { device, inode } => RunningImageEvidenceWire::MacosSpawnInode {
1112 device: *device,
1113 inode: *inode,
1114 }
1115 .serialize(serializer),
1116 Self::Unknown { body, .. } => body.serialize(serializer),
1117 }
1118 }
1119}
1120
1121impl<'de> Deserialize<'de> for RunningImageEvidence {
1122 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1123 where
1124 D: serde::Deserializer<'de>,
1125 {
1126 let (tag, value) = read_tagged(deserializer, "method")?;
1127 match tag.as_str() {
1128 "linux_proc_sha256" => {
1129 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1130 RunningImageEvidenceWire::LinuxProcSha256 { digest } => {
1131 Ok(Self::LinuxProcSha256 { digest })
1132 }
1133 _ => unreachable!(),
1134 }
1135 }
1136 "macos_spawn_inode" => {
1137 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1138 RunningImageEvidenceWire::MacosSpawnInode { device, inode } => {
1139 Ok(Self::MacosSpawnInode { device, inode })
1140 }
1141 _ => unreachable!(),
1142 }
1143 }
1144 _ => Ok(Self::Unknown { tag, body: value }),
1145 }
1146 }
1147}
1148
1149impl Serialize for SupervisorRouteConsumer {
1150 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1151 where
1152 S: Serializer,
1153 {
1154 match self {
1155 Self::Reserved { module_id } => SupervisorRouteConsumerWire::Reserved {
1156 module_id: module_id.clone(),
1157 }
1158 .serialize(serializer),
1159 Self::Direct { connection_id } => SupervisorRouteConsumerWire::Direct {
1160 connection_id: *connection_id,
1161 }
1162 .serialize(serializer),
1163 Self::Unknown { body, .. } => body.serialize(serializer),
1164 }
1165 }
1166}
1167
1168impl<'de> Deserialize<'de> for SupervisorRouteConsumer {
1169 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1170 where
1171 D: serde::Deserializer<'de>,
1172 {
1173 let (tag, value) = read_tagged(deserializer, "kind")?;
1174 match tag.as_str() {
1175 "reserved" => {
1176 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1177 SupervisorRouteConsumerWire::Reserved { module_id } => {
1178 Ok(Self::Reserved { module_id })
1179 }
1180 _ => unreachable!(),
1181 }
1182 }
1183 "direct" => {
1184 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1185 SupervisorRouteConsumerWire::Direct { connection_id } => {
1186 Ok(Self::Direct { connection_id })
1187 }
1188 _ => unreachable!(),
1189 }
1190 }
1191 _ => Ok(Self::Unknown { tag, body: value }),
1192 }
1193 }
1194}
1195
1196impl Serialize for StderrCaptureState {
1197 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1198 where
1199 S: Serializer,
1200 {
1201 match self {
1202 Self::Captured => StderrCaptureStateWire::Captured.serialize(serializer),
1203 Self::Incomplete { reason } => StderrCaptureStateWire::Incomplete {
1204 reason: reason.clone(),
1205 }
1206 .serialize(serializer),
1207 Self::NotCaptured { reason } => StderrCaptureStateWire::NotCaptured {
1208 reason: reason.clone(),
1209 }
1210 .serialize(serializer),
1211 Self::Unknown { body, .. } => body.serialize(serializer),
1212 }
1213 }
1214}
1215
1216impl<'de> Deserialize<'de> for StderrCaptureState {
1217 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1218 where
1219 D: serde::Deserializer<'de>,
1220 {
1221 let (tag, value) = read_tagged(deserializer, "state")?;
1222 match tag.as_str() {
1223 "captured" => {
1224 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1225 StderrCaptureStateWire::Captured => Ok(Self::Captured),
1226 _ => unreachable!(),
1227 }
1228 }
1229 "incomplete" => match serde_json::from_value(value.into_value())
1230 .map_err(D::Error::custom)?
1231 {
1232 StderrCaptureStateWire::Incomplete { reason } => Ok(Self::Incomplete { reason }),
1233 _ => unreachable!(),
1234 },
1235 "not_captured" => match serde_json::from_value(value.into_value())
1236 .map_err(D::Error::custom)?
1237 {
1238 StderrCaptureStateWire::NotCaptured { reason } => Ok(Self::NotCaptured { reason }),
1239 _ => unreachable!(),
1240 },
1241 _ => Ok(Self::Unknown { tag, body: value }),
1242 }
1243 }
1244}
1245
1246impl Serialize for StderrTailEntry {
1247 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1248 where
1249 S: Serializer,
1250 {
1251 match self {
1252 Self::Line { text, truncated } => StderrTailEntryWire::Line {
1253 text: text.clone(),
1254 truncated: *truncated,
1255 }
1256 .serialize(serializer),
1257 Self::ProcessStart => StderrTailEntryWire::ProcessStart.serialize(serializer),
1258 Self::Unknown { body, .. } => body.serialize(serializer),
1259 }
1260 }
1261}
1262
1263impl<'de> Deserialize<'de> for StderrTailEntry {
1264 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1265 where
1266 D: serde::Deserializer<'de>,
1267 {
1268 let (tag, value) = read_tagged(deserializer, "kind")?;
1269 match tag.as_str() {
1270 "line" => match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1271 StderrTailEntryWire::Line { text, truncated } => Ok(Self::Line { text, truncated }),
1272 _ => unreachable!(),
1273 },
1274 "process_start" => {
1275 match serde_json::from_value(value.into_value()).map_err(D::Error::custom)? {
1276 StderrTailEntryWire::ProcessStart => Ok(Self::ProcessStart),
1277 _ => unreachable!(),
1278 }
1279 }
1280 _ => Ok(Self::Unknown { tag, body: value }),
1281 }
1282 }
1283}
1284
1285fn is_zero_u64(value: &u64) -> bool {
1286 *value == 0
1287}
1288
1289#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1291pub struct TerminalHistory {
1292 pub daemon_started_at_ms: u64,
1295 pub entries: Vec<TerminalEntry>,
1296 #[serde(default, skip_serializing_if = "is_zero_u64")]
1298 pub dropped: u64,
1299}
1300
1301#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1303pub struct TerminalEntry {
1304 #[serde(default, skip_serializing_if = "Option::is_none")]
1305 pub exit_code: Option<i32>,
1306 #[serde(default, skip_serializing_if = "Option::is_none")]
1307 pub exit_signal: Option<i32>,
1308 pub at_ms: u64,
1309 pub disposition: TerminalDisposition,
1310 #[serde(default, skip_serializing_if = "Option::is_none")]
1314 pub exit_kind: Option<TerminalExitKind>,
1315 #[serde(default, skip_serializing_if = "Option::is_none")]
1322 pub disposition_detail: Option<String>,
1323}
1324
1325#[derive(Debug, Clone, PartialEq, Eq)]
1330pub enum TerminalExitKind {
1331 Clean,
1332 Crash,
1333 DeliberateSeverance,
1334 Unknown(String),
1335}
1336
1337impl TerminalExitKind {
1338 fn wire_name(&self) -> &str {
1339 match self {
1340 Self::Clean => "clean",
1341 Self::Crash => "crash",
1342 Self::DeliberateSeverance => "deliberate_severance",
1343 Self::Unknown(value) => value,
1344 }
1345 }
1346}
1347
1348impl Serialize for TerminalExitKind {
1349 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1350 where
1351 S: serde::Serializer,
1352 {
1353 serializer.serialize_str(self.wire_name())
1354 }
1355}
1356
1357impl<'de> Deserialize<'de> for TerminalExitKind {
1358 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1359 where
1360 D: serde::Deserializer<'de>,
1361 {
1362 let value = String::deserialize(deserializer)?;
1363 Ok(match value.as_str() {
1364 "clean" => Self::Clean,
1365 "crash" => Self::Crash,
1366 "deliberate_severance" => Self::DeliberateSeverance,
1367 _ => Self::Unknown(value),
1368 })
1369 }
1370}
1371
1372open_string_enum! {
1373 TerminalDisposition {
1375 Stopped => "stopped",
1376 Disabled => "disabled",
1377 Failed => "failed",
1378 Restarting => "restarting",
1379 }
1380}
1381
1382#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1383#[serde(rename_all = "snake_case")]
1384pub enum PollKind {
1385 Status,
1386 Liveness,
1387}
1388
1389#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1390pub struct CatalogEntry {
1391 pub module_id: String,
1392 #[serde(default, skip_serializing_if = "Option::is_none")]
1413 pub module_version: Option<String>,
1414 pub roles: Vec<ProviderRole>,
1415 pub control_ops: Vec<String>,
1416 #[serde(default, skip_serializing_if = "Option::is_none")]
1421 pub capabilities: Option<CapabilityDeclarations>,
1422 #[serde(default, skip_serializing_if = "Option::is_none")]
1425 pub self_signals: Option<Vec<SelfSignalDeclaration>>,
1426}
1427
1428#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1429pub struct CapabilityRequirementStatus {
1430 pub consumer: String,
1431 pub capability: String,
1432 pub need: String,
1433 pub verdict: String,
1434 pub episode_seq: u64,
1435 pub config_satisfiable: bool,
1436 pub runtime_available: bool,
1437 pub detail: String,
1438}
1439
1440#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1441pub struct SupervisorRescanResult {
1442 pub added: Vec<String>,
1443 pub removed: Vec<String>,
1444 pub changed_pending_reload: Vec<String>,
1445 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1458 pub enabled_changes: Vec<String>,
1459 pub unchanged: u32,
1460 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
1468 pub preview: bool,
1469 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1487 pub restart_required: Vec<String>,
1488 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1492 pub capability_warnings: Vec<String>,
1493}
1494
1495#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1496pub struct SupervisorEntry {
1497 pub module_id: String,
1498 pub state: String,
1499 pub enabled: bool,
1500 pub live: bool,
1501 pub health: SupervisorHealthStatus,
1502 #[serde(default)]
1508 pub last_probe_ms: Option<u64>,
1509 #[serde(default, skip_serializing_if = "Option::is_none")]
1513 pub last_exit_code: Option<i32>,
1514 #[serde(default, skip_serializing_if = "Option::is_none")]
1518 pub last_exit_signal: Option<i32>,
1519 #[serde(default, skip_serializing_if = "Option::is_none")]
1523 pub last_exit_ms: Option<u64>,
1524 #[serde(default, skip_serializing_if = "Option::is_none")]
1527 pub last_exit_kind: Option<TerminalExitKind>,
1528 #[serde(default, skip_serializing_if = "Option::is_none")]
1545 pub restart_count: Option<u32>,
1546 #[serde(default, skip_serializing_if = "Option::is_none")]
1549 pub max_restarts: Option<u32>,
1550 #[serde(default, skip_serializing_if = "Option::is_none")]
1553 pub lifetime_restarts: Option<u32>,
1554 #[serde(default, skip_serializing_if = "Option::is_none")]
1564 pub restart_window_secs: Option<u64>,
1565 #[serde(default, skip_serializing_if = "Option::is_none")]
1569 pub drain_timeout_ms: Option<u64>,
1570 #[serde(default, skip_serializing_if = "Option::is_none")]
1573 pub restart_backoff_ms: Option<u64>,
1574 #[serde(default, skip_serializing_if = "Option::is_none")]
1577 pub restart_max_backoff_ms: Option<u64>,
1578}
1579
1580#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1581#[serde(rename_all = "snake_case")]
1582pub enum SupervisorHealthStatus {
1583 Ok,
1584 Degraded,
1585 Failing,
1586 Unresponsive,
1587 Unknown,
1588}
1589
1590#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1591pub struct SupervisorHealthEntry {
1592 pub module_id: String,
1593 pub status: SupervisorHealthStatus,
1594 #[serde(default, skip_serializing_if = "Option::is_none")]
1600 pub detail: Option<String>,
1601 #[serde(default, skip_serializing_if = "Option::is_none")]
1606 pub metrics: Option<serde_json::Value>,
1607 pub consecutive_failures: u32,
1608 #[serde(default)]
1611 pub late_answer_count: u64,
1612 #[serde(default, skip_serializing_if = "Option::is_none")]
1614 pub last_late_answer_latency_ms: Option<u64>,
1615 #[serde(default)]
1620 pub last_action: Option<String>,
1621 #[serde(default)]
1624 pub last_action_ms: Option<u64>,
1625 #[serde(default, skip_serializing_if = "Option::is_none")]
1638 pub last_probe_ms: Option<u64>,
1639}
1640
1641#[cfg(test)]
1642mod tests {
1643 use super::*;
1644 use subc_protocol::{BindIdentity, RouteTarget};
1645
1646 #[test]
1647 fn legacy_terminal_decoder_ignores_deliberate_severance_kind() {
1648 let entry = TerminalEntry {
1649 exit_code: Some(1),
1650 exit_signal: None,
1651 at_ms: 1_700_000_000_123,
1652 disposition: TerminalDisposition::Restarting,
1653 exit_kind: Some(TerminalExitKind::DeliberateSeverance),
1654 disposition_detail: None,
1655 };
1656 let wire = serde_json::to_string(&entry).expect("terminal entry serializes");
1657 assert_eq!(
1658 serde_json::from_str::<serde_json::Value>(&wire).expect("terminal entry is JSON")
1659 ["exit_kind"],
1660 "deliberate_severance"
1661 );
1662
1663 #[derive(serde::Deserialize)]
1664 struct LegacyTerminalEntry {
1665 exit_code: Option<i32>,
1666 exit_signal: Option<i32>,
1667 at_ms: u64,
1668 disposition: TerminalDisposition,
1669 }
1670
1671 let decoded: LegacyTerminalEntry =
1672 serde_json::from_str(&wire).expect("legacy decoder keeps the terminal record");
1673 assert_eq!(decoded.exit_code, Some(1));
1674 assert_eq!(decoded.exit_signal, None);
1675 assert_eq!(decoded.at_ms, 1_700_000_000_123);
1676 assert_eq!(decoded.disposition, TerminalDisposition::Restarting);
1677
1678 let future_wire = wire.replace("deliberate_severance", "future_exit_kind");
1679 let future: TerminalEntry =
1680 serde_json::from_str(&future_wire).expect("new decoder keeps a future terminal kind");
1681 assert_eq!(
1682 future.exit_kind,
1683 Some(TerminalExitKind::Unknown("future_exit_kind".to_string()))
1684 );
1685 }
1686
1687 #[test]
1688 fn route_poll_uses_kind_field() {
1689 let body = serde_json::to_value(ClientControlRequest::RoutePoll {
1690 route_channel: 7,
1691 route_epoch: 11,
1692 kind: PollKind::Status,
1693 })
1694 .unwrap();
1695
1696 assert_eq!(body["op"], "route.poll");
1697 assert_eq!(body["route_epoch"], 11);
1698 assert_eq!(body["kind"], "status");
1699 assert!(body.get("op").is_some());
1700 }
1701
1702 #[test]
1703 fn route_open_is_internally_tagged() {
1704 let request = ClientControlRequest::RouteOpen {
1705 target: RouteTarget::ToolProvider {
1706 module_id: "aft".to_string(),
1707 },
1708 identity: BindIdentity::new("/tmp/project", "opencode", "session-1"),
1709 consumer_identity: None,
1710 consumer_capabilities: None,
1711 admission_facts: None,
1712 };
1713
1714 let body = serde_json::to_value(request).unwrap();
1715 assert_eq!(body["op"], "route.open");
1716 assert_eq!(body["target"]["kind"], "tool_provider");
1717 assert!(body.get("consumer_identity").is_none());
1718 assert!(body.get("consumer_capabilities").is_none());
1719 }
1720
1721 #[test]
1722 fn route_open_without_optional_fields_still_decodes() {
1723 let body = serde_json::json!({
1724 "op": "route.open",
1725 "target": { "kind": "tool_provider", "module_id": "aft" },
1726 "identity": {
1727 "project_root": "/tmp/project",
1728 "harness": "opencode",
1729 "session": "session-1"
1730 }
1731 });
1732
1733 let decoded: ClientControlRequest = serde_json::from_value(body).unwrap();
1734 let ClientControlRequest::RouteOpen {
1735 consumer_identity,
1736 consumer_capabilities,
1737 admission_facts,
1738 ..
1739 } = decoded
1740 else {
1741 panic!("decoded wrong request variant");
1742 };
1743 assert_eq!(consumer_identity, None);
1744 assert_eq!(consumer_capabilities, None);
1745 assert_eq!(admission_facts, None);
1746 }
1747
1748 #[test]
1749 fn new_route_closed_decoder_defaults_fields_absent_from_old_daemon() {
1750 let old_wire = r#"{"op":"route.closed","module_id":"aft-tools","reason":"crash","drained":false,"abandoned":0}"#;
1751 let decoded: ClientControlPush = serde_json::from_str(old_wire).unwrap();
1752 match decoded {
1753 ClientControlPush::RouteClosed {
1754 excluded_subscriptions,
1755 terminal,
1756 ..
1757 } => {
1758 assert_eq!(excluded_subscriptions, 0);
1759 assert_eq!(terminal, None);
1760 }
1761 other => panic!("unexpected push: {other:?}"),
1762 }
1763 assert!(!serde_json::to_string(&decoded)
1764 .unwrap()
1765 .contains("terminal"));
1766 }
1767
1768 #[test]
1769 fn old_route_closed_decoder_ignores_new_terminal_field() {
1770 #[derive(serde::Deserialize)]
1771 #[serde(tag = "op")]
1772 enum LegacyClientControlPush {
1773 #[serde(rename = "route.closed")]
1774 RouteClosed {
1775 module_id: String,
1776 reason: RouteCloseReason,
1777 drained: bool,
1778 abandoned: u32,
1779 },
1780 }
1781
1782 let wire = r#"{"op":"route.closed","module_id":"aft-tools","reason":"crash","drained":false,"abandoned":0,"excluded_subscriptions":3,"terminal":true}"#;
1783 let decoded: LegacyClientControlPush = serde_json::from_str(wire).unwrap();
1784 match decoded {
1785 LegacyClientControlPush::RouteClosed {
1786 module_id,
1787 reason,
1788 drained,
1789 abandoned,
1790 } => {
1791 assert_eq!(module_id, "aft-tools");
1792 assert_eq!(reason, RouteCloseReason::Crash);
1793 assert!(!drained);
1794 assert_eq!(abandoned, 0);
1795 }
1796 }
1797 }
1798
1799 #[test]
1800 fn supervisor_routes_is_a_control_plane_request() {
1801 let body = serde_json::json!({
1802 "op": "supervisor.routes",
1803 "module_id": "aft"
1804 });
1805
1806 let request: ClientControlRequest = serde_json::from_value(body.clone()).unwrap();
1807 assert_eq!(serde_json::to_value(request).unwrap(), body);
1808 }
1809
1810 #[test]
1811 fn diagnostic_string_enums_retain_unknown_wire_values() {
1812 let reason: RunningImageUnavailableReason =
1813 serde_json::from_str("\"future_reason\"").unwrap();
1814 let disposition: TerminalDisposition =
1815 serde_json::from_str("\"future_disposition\"").unwrap();
1816
1817 assert_eq!(
1818 reason,
1819 RunningImageUnavailableReason::Unknown("future_reason".to_string())
1820 );
1821 assert_eq!(
1822 disposition,
1823 TerminalDisposition::Unknown("future_disposition".to_string())
1824 );
1825 }
1826
1827 #[test]
1828 fn diagnostic_string_enums_preserve_existing_wire_names() {
1829 let names = [
1830 (RunningImageUnavailableReason::NotRunning, "not_running"),
1831 (
1832 RunningImageUnavailableReason::UnsupportedPlatform,
1833 "unsupported_platform",
1834 ),
1835 (
1836 RunningImageUnavailableReason::RunningExecutableUnreadable,
1837 "running_executable_unreadable",
1838 ),
1839 (
1840 RunningImageUnavailableReason::SpawnedPathUnreadable,
1841 "spawned_path_unreadable",
1842 ),
1843 (RunningImageUnavailableReason::HashFailed, "hash_failed"),
1844 (
1845 RunningImageUnavailableReason::ProcessIdentityUnconfirmed,
1846 "process_identity_unconfirmed",
1847 ),
1848 ];
1849 for (value, expected) in names {
1850 let wire = serde_json::to_string(&value).unwrap();
1851 assert_eq!(wire, format!("\"{expected}\""));
1852 let decoded: RunningImageUnavailableReason = serde_json::from_str(&wire).unwrap();
1853 assert_eq!(decoded, value);
1854 }
1855
1856 for (value, expected) in [
1857 (TerminalDisposition::Stopped, "stopped"),
1858 (TerminalDisposition::Disabled, "disabled"),
1859 (TerminalDisposition::Failed, "failed"),
1860 (TerminalDisposition::Restarting, "restarting"),
1861 ] {
1862 let wire = serde_json::to_string(&value).unwrap();
1863 assert_eq!(wire, format!("\"{expected}\""));
1864 let decoded: TerminalDisposition = serde_json::from_str(&wire).unwrap();
1865 assert_eq!(decoded, value);
1866 }
1867 }
1868
1869 #[test]
1870 fn diagnostic_string_enums_reject_non_string_bodies() {
1871 assert!(serde_json::from_str::<RunningImageUnavailableReason>("42").is_err());
1872 assert!(serde_json::from_str::<TerminalDisposition>("{\"value\":\"failed\"}").is_err());
1873 }
1874
1875 #[test]
1876 fn unknown_provenance_reason_does_not_discard_healthy_siblings() {
1877 let body = serde_json::json!({
1878 "op": "supervisor.provenance",
1879 "daemon": {
1880 "daemon_build": {},
1881 "daemon_observed": {
1882 "running_image": {
1883 "status": "unavailable",
1884 "reason": "not_running"
1885 }
1886 }
1887 },
1888 "modules": [
1889 {
1890 "module_id": "future",
1891 "module_declared": { "status": "unverifiable" },
1892 "daemon_observed": {
1893 "running_image": {
1894 "status": "unavailable",
1895 "reason": "future_reason"
1896 }
1897 }
1898 },
1899 {
1900 "module_id": "healthy-a",
1901 "module_declared": { "status": "unverifiable" },
1902 "daemon_observed": {
1903 "running_image": {
1904 "status": "match",
1905 "evidence": {
1906 "method": "linux_proc_sha256",
1907 "digest": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
1908 }
1909 }
1910 }
1911 },
1912 {
1913 "module_id": "healthy-b",
1914 "module_declared": { "status": "unverifiable" },
1915 "daemon_observed": {
1916 "running_image": {
1917 "status": "unavailable",
1918 "reason": "unsupported_platform"
1919 }
1920 }
1921 }
1922 ]
1923 });
1924
1925 let decoded: ClientControlResponse = serde_json::from_value(body).unwrap();
1926 let ClientControlResponse::SupervisorProvenance { modules, .. } = decoded else {
1927 panic!("decoded wrong response variant");
1928 };
1929 assert_eq!(modules.len(), 3);
1930 assert_eq!(modules[0].module_id, "future");
1931 assert_eq!(
1932 modules[0].daemon_observed.running_image,
1933 RunningImageAgreement::Unavailable {
1934 reason: RunningImageUnavailableReason::Unknown("future_reason".to_string())
1935 }
1936 );
1937 assert_eq!(modules[1].module_id, "healthy-a");
1938 assert_eq!(modules[2].module_id, "healthy-b");
1939 }
1940
1941 #[test]
1942 fn tagged_unknown_values_retain_tag_and_body() {
1943 macro_rules! assert_unknown_round_trip {
1944 ($ty:ident, $field:literal, $value:expr) => {
1945 let value = $value;
1946 let wire = serde_json::to_string(&value).unwrap();
1947 let decoded: $ty = serde_json::from_str(&wire).unwrap();
1948 match decoded {
1949 $ty::Unknown { tag, body } => {
1950 assert_eq!(tag, value[$field].as_str().unwrap());
1951 assert_eq!(serde_json::to_value(&body).unwrap(), value);
1952 }
1953 _ => panic!("decoded known variant"),
1954 }
1955 };
1956 }
1957
1958 assert_unknown_round_trip!(
1959 ModuleDeclaredProvenance,
1960 "status",
1961 serde_json::json!({"status": "future", "build": {"version": 7}})
1962 );
1963 assert_unknown_round_trip!(
1964 RunningImageAgreement,
1965 "status",
1966 serde_json::json!({"status": "future", "evidence": {"digest": "abc"}})
1967 );
1968 assert_unknown_round_trip!(
1969 RunningImageEvidence,
1970 "method",
1971 serde_json::json!({"method": "future", "digest": "abc"})
1972 );
1973 assert_unknown_round_trip!(
1974 SupervisorRouteConsumer,
1975 "kind",
1976 serde_json::json!({"kind": "future", "module_id": "m"})
1977 );
1978 assert_unknown_round_trip!(
1979 StderrCaptureState,
1980 "state",
1981 serde_json::json!({"state": "future", "reason": "because"})
1982 );
1983 assert_unknown_round_trip!(
1984 StderrTailEntry,
1985 "kind",
1986 serde_json::json!({"kind": "future", "text": "line"})
1987 );
1988 }
1989
1990 #[test]
1991 fn tagged_unknown_values_round_trip_the_original_json() {
1992 let wire = r#"{"kind":"future_consumer","detail":{"z":1}}"#;
1993 let decoded: SupervisorRouteConsumer = serde_json::from_str(wire).unwrap();
1994 assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
1995 }
1996
1997 #[test]
1998 fn tagged_unknown_values_round_trip_trailing_tag() {
1999 let route_wire = r#"{"detail":{"z":1},"kind":"future_consumer"}"#;
2000 let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2001 assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2002
2003 let stderr_wire = r#"{"reason":"because","state":"future_state"}"#;
2004 let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2005 assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2006 }
2007
2008 #[test]
2009 fn tagged_unknown_values_round_trip_middle_tag() {
2010 let route_wire = r#"{"a":1,"kind":"future_x","b":2}"#;
2011 let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2012 assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2013
2014 let stderr_wire = r#"{"a":1,"state":"future_state","b":2}"#;
2015 let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2016 assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2017 }
2018
2019 #[test]
2020 fn tagged_unknown_values_round_trip_deep_payload() {
2021 let route_wire = r#"{"a":{"n":[1,2]},"kind":"future_x","zz":"s","b":null}"#;
2022 let route: SupervisorRouteConsumer = serde_json::from_str(route_wire).unwrap();
2023 assert_eq!(serde_json::to_string(&route).unwrap(), route_wire);
2024
2025 let stderr_wire = r#"{"a":{"n":[1,2]},"state":"future_state","zz":"s","b":null}"#;
2026 let stderr: StderrCaptureState = serde_json::from_str(stderr_wire).unwrap();
2027 assert_eq!(serde_json::to_string(&stderr).unwrap(), stderr_wire);
2028 }
2029
2030 #[test]
2031 fn tagged_unknown_values_reject_non_object_bodies() {
2032 for wire in ["42", r#""future""#, "[]"] {
2033 assert!(serde_json::from_str::<SupervisorRouteConsumer>(wire).is_err());
2034 assert!(serde_json::from_str::<StderrCaptureState>(wire).is_err());
2035 }
2036 }
2037
2038 #[test]
2039 fn duplicate_discriminators_reject_without_panicking() {
2040 assert_eq!(
2041 serde_json::from_str::<ModuleDeclaredProvenance>(r#"{"status":"unverifiable"}"#)
2042 .unwrap(),
2043 ModuleDeclaredProvenance::Unverifiable
2044 );
2045 match serde_json::from_str::<ModuleDeclaredProvenance>(r#"{"status":"future_thing"}"#)
2046 .unwrap()
2047 {
2048 ModuleDeclaredProvenance::Unknown { tag, .. } => assert_eq!(tag, "future_thing"),
2049 _ => panic!("future discriminator decoded as a known variant"),
2050 }
2051
2052 let wires = [
2053 r#"{"status":"reported","status":"unverifiable"}"#,
2054 r#"{"status":"unverifiable","status":"reported"}"#,
2055 r#"{"status":"reported","build":{},"status":"unverifiable"}"#,
2056 r#"{"status":"unverifiable","build":{},"status":"reported"}"#,
2057 ];
2058
2059 for wire in wires {
2060 let result =
2061 std::panic::catch_unwind(|| serde_json::from_str::<ModuleDeclaredProvenance>(wire));
2062 assert!(result.is_ok(), "duplicate discriminator panicked: {wire}");
2063 assert!(
2064 result.unwrap().is_err(),
2065 "duplicate discriminator decoded: {wire}"
2066 );
2067 }
2068
2069 let wire = r#"{"state":"captured","state":"incomplete","reason":"x"}"#;
2070 let result = std::panic::catch_unwind(|| serde_json::from_str::<StderrCaptureState>(wire));
2071 assert!(result.is_ok(), "duplicate discriminator panicked: {wire}");
2072 assert!(
2073 result.unwrap().is_err(),
2074 "duplicate discriminator decoded: {wire}"
2075 );
2076 }
2077
2078 #[test]
2079 fn nested_unknown_values_round_trip_without_normalizing_member_order() {
2080 let known_wire =
2081 r#"{"status":"match","evidence":{"method":"linux_proc_sha256","digest":"abc"}}"#;
2082 let known: RunningImageAgreement = serde_json::from_str(known_wire).unwrap();
2083 assert_eq!(serde_json::to_string(&known).unwrap(), known_wire);
2084
2085 for wire in [
2086 r#"{"kind":"future_x","detail":{"zeta":1,"alpha":2}}"#,
2087 r#"{"kind":"future_x","d":{"b":{"zz":1,"aa":2}}}"#,
2088 ] {
2089 let decoded: SupervisorRouteConsumer = serde_json::from_str(wire).unwrap();
2090 assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2091 }
2092
2093 for wire in [
2094 r#"{"status":"match","evidence":{"method":"future_probe","zz":1,"aa":2}}"#,
2095 r#"{"status":"match","evidence":{"method":"future_probe","d":{"zz":1,"aa":2}}}"#,
2096 ] {
2097 let decoded: RunningImageAgreement = serde_json::from_str(wire).unwrap();
2098 assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2099 }
2100
2101 let wire = r#"{"status":"mismatch","running":{"detail":{"z":1},"method":"future_running"},"disk":{"method":"future_disk","detail":{"z":1}}}"#;
2102 let decoded: RunningImageAgreement = serde_json::from_str(wire).unwrap();
2103 assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2104
2105 let wire = r#"{"capture":{"state":"captured"},"entries":[{"detail":{"z":1,"a":2},"kind":"future_line"},{"kind":"future_restart","meta":{"b":{"zz":1,"aa":2}}}]}"#;
2106 let decoded: StderrTail = serde_json::from_str(wire).unwrap();
2107 assert_eq!(serde_json::to_string(&decoded).unwrap(), wire);
2108 }
2109
2110 #[test]
2111 fn tagged_unknown_member_does_not_discard_known_siblings() {
2112 let body = serde_json::json!({
2113 "modules": [{
2114 "module_id": "target",
2115 "routes": [
2116 {"consumer": {"kind": "future_consumer", "module_id": "m", "detail": {"retry": true}}, "age_ms": 0, "draining": false},
2117 {"consumer": {"kind": "direct", "connection_id": 7}, "age_ms": 0, "draining": false}
2118 ]
2119 }]
2120 });
2121 let decoded: ClientControlResponse = serde_json::from_value(
2122 serde_json::json!({"op": "supervisor.routes", "modules": body["modules"]}),
2123 )
2124 .unwrap();
2125 let ClientControlResponse::SupervisorRoutes { modules } = decoded else {
2126 panic!("decoded wrong response variant");
2127 };
2128 assert_eq!(modules[0].routes.len(), 2);
2129 assert_eq!(
2130 modules[0].routes[1].consumer,
2131 SupervisorRouteConsumer::Direct { connection_id: 7 }
2132 );
2133 }
2134}