1#![forbid(unsafe_code)]
32
33use std::{error::Error, fmt, path::PathBuf};
34
35use serde::{Deserialize, Serialize};
36
37pub mod frame;
38pub mod manifest;
39pub mod session;
40
41pub use frame::{Frame, FrameBuildError};
42
43#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
64pub struct BindIdentity {
65 pub project_root: PathBuf,
66 pub harness: String,
67 pub session: String,
68}
69
70#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
72#[serde(tag = "kind", rename_all = "snake_case")]
73pub enum Principal {
74 Reserved { module_id: String },
76 Direct,
78 Unverified,
80}
81
82#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
95#[serde(tag = "kind", rename_all = "snake_case")]
96pub enum RouteTarget {
97 ToolProvider {
98 module_id: String,
99 },
100 ManagementSurface {
101 module_id: String,
102 },
103 InternalService {
104 module_id: String,
105 service_id: String,
106 },
107}
108
109pub const PROTOCOL_VERSION: u8 = 2;
111
112pub const MIN_SUPPORTED_VERSION: u8 = 2;
114
115pub const SUBC_MODULE_ID_ENV: &str = "SUBC_MODULE_ID";
118
119pub const SUBC_LAUNCH_NONCE_ENV: &str = "SUBC_LAUNCH_NONCE";
124
125pub const HEADER_LEN: usize = 21;
127
128pub const FROZEN_PREFIX_LEN: usize = 5;
132
133pub const MAX_FRAME_BODY_LEN: u32 = 64 * 1024 * 1024;
139
140#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
142pub struct ErrorBody {
143 pub code: String,
144 pub message: String,
145}
146
147#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
149pub struct ModuleHelloBody {
150 pub manifest: manifest::ModuleManifest,
151 pub protocol_ver: u8,
152 #[serde(default)]
153 pub control_ops: Option<Vec<String>>,
154 #[serde(default, skip_serializing_if = "Option::is_none")]
162 pub launch_nonce: Option<String>,
163}
164
165#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
167pub struct ModuleHelloAckBody {
168 pub negotiated_ver: u8,
169 pub subc_ops: Vec<String>,
170 pub subc_capabilities: Vec<String>,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub storage: Option<serde_json::Value>,
179}
180
181#[derive(Debug, Clone, Copy, PartialEq, Eq)]
186#[repr(u8)]
187pub enum FrameType {
188 Request = 0,
189 Response = 1,
190 Push = 2,
191 StreamData = 3,
192 StreamEnd = 4,
193 Error = 5,
194 Cancel = 6,
195 Ping = 7,
196 Pong = 8,
197 Hello = 9,
198 HelloAck = 10,
199 Goodbye = 11,
200}
201
202impl FrameType {
203 pub fn from_u8(b: u8) -> Option<Self> {
205 Some(match b {
206 0 => Self::Request,
207 1 => Self::Response,
208 2 => Self::Push,
209 3 => Self::StreamData,
210 4 => Self::StreamEnd,
211 5 => Self::Error,
212 6 => Self::Cancel,
213 7 => Self::Ping,
214 8 => Self::Pong,
215 9 => Self::Hello,
216 10 => Self::HelloAck,
217 11 => Self::Goodbye,
218 _ => return None,
219 })
220 }
221
222 pub fn is_pure_header(self) -> bool {
223 matches!(self, Self::Cancel | Self::Ping | Self::Pong | Self::Goodbye)
224 }
225}
226
227#[derive(Debug, Clone, Copy, PartialEq, Eq)]
230#[repr(u8)]
231pub enum Priority {
232 Passive = 0,
233 Interactive = 1,
234 Background = 2,
235}
236
237impl Priority {
238 fn from_bits(bits: u8) -> Option<Self> {
239 Some(match bits {
240 0 => Self::Passive,
241 1 => Self::Interactive,
242 2 => Self::Background,
243 _ => return None,
244 })
245 }
246}
247
248#[derive(Debug, Clone, Copy, PartialEq, Eq)]
250#[repr(u8)]
251pub enum AdmissionClass {
252 Normal = 0,
253 Expedite = 1,
254 Sheddable = 2,
255}
256
257impl AdmissionClass {
258 fn from_bits(bits: u8) -> Option<Self> {
259 Some(match bits {
260 0 => Self::Normal,
261 1 => Self::Expedite,
262 2 => Self::Sheddable,
263 _ => return None,
264 })
265 }
266}
267
268const FLAG_BINARY: u8 = 0b0000_0001; const FLAG_PRIORITY_MASK: u8 = 0b0000_0110; const FLAG_PRIORITY_SHIFT: u8 = 1;
271const FLAG_LAST: u8 = 0b0000_1000; const FLAG_ADMISSION_MASK: u8 = 0b0011_0000; const FLAG_ADMISSION_SHIFT: u8 = 4;
274const FLAG_RESERVED_MASK: u8 = 0b1100_0000; #[derive(Debug, Clone, Copy, PartialEq, Eq)]
278pub struct Flags(pub u8);
279
280impl Flags {
281 pub fn new(binary: bool, priority: Priority, last: bool) -> Self {
283 let mut b = 0u8;
284 if binary {
285 b |= FLAG_BINARY;
286 }
287 b |= (priority as u8) << FLAG_PRIORITY_SHIFT;
288 if last {
289 b |= FLAG_LAST;
290 }
291 Flags(b)
292 }
293
294 pub fn with_admission_class(mut self, admission_class: AdmissionClass) -> Self {
296 self.0 =
297 (self.0 & !FLAG_ADMISSION_MASK) | ((admission_class as u8) << FLAG_ADMISSION_SHIFT);
298 self
299 }
300
301 pub fn is_binary(self) -> bool {
303 self.0 & FLAG_BINARY != 0
304 }
305
306 pub fn is_last(self) -> bool {
308 self.0 & FLAG_LAST != 0
309 }
310
311 pub fn priority(self) -> Option<Priority> {
313 Priority::from_bits((self.0 & FLAG_PRIORITY_MASK) >> FLAG_PRIORITY_SHIFT)
314 }
315
316 pub fn admission_class(self) -> Option<AdmissionClass> {
318 AdmissionClass::from_bits((self.0 & FLAG_ADMISSION_MASK) >> FLAG_ADMISSION_SHIFT)
319 }
320
321 pub fn has_reserved_bits(self) -> bool {
323 self.0 & FLAG_RESERVED_MASK != 0
324 }
325}
326
327#[derive(Debug, Clone, Copy, PartialEq, Eq)]
329pub struct EnvelopeHeader {
330 pub len: u32,
332 pub ver: u8,
334 pub ty: FrameType,
336 pub flags: Flags,
338 pub channel: u16,
340 pub epoch: u32,
342 pub corr: u64,
344}
345
346impl EnvelopeHeader {
347 pub fn encode(&self) -> [u8; HEADER_LEN] {
349 let mut buf = [0u8; HEADER_LEN];
350 buf[0..4].copy_from_slice(&self.len.to_le_bytes());
351 buf[4] = self.ver;
352 buf[5] = self.ty as u8;
353 buf[6] = self.flags.0;
354 buf[7..9].copy_from_slice(&self.channel.to_le_bytes());
355 buf[9..13].copy_from_slice(&self.epoch.to_le_bytes());
356 buf[13..21].copy_from_slice(&self.corr.to_le_bytes());
357 buf
358 }
359}
360
361#[derive(Debug, Clone, Copy, PartialEq, Eq)]
363pub enum DecodeError {
364 TooShortForPrefix { have: usize },
366 UnsupportedVersion { ver: u8 },
368 TooShortForHeader { have: usize, need: usize },
370 UnknownFrameType { byte: u8 },
372 ReservedFlagBits { flags: u8 },
374 ReservedPriorityBits { flags: u8 },
376 ReservedAdmissionClass { flags: u8 },
378 SheddableIllegalFrameType { ty: FrameType, flags: u8 },
380 NonzeroEpochOnControlChannel { epoch: u32 },
382 PureHeaderFrameWithBody { ty: FrameType, len: u32 },
384}
385
386impl fmt::Display for DecodeError {
387 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
388 match self {
389 Self::TooShortForPrefix { have } => {
390 write!(f, "header shorter than frozen prefix: have {have} bytes")
391 }
392 Self::UnsupportedVersion { ver } => write!(f, "unsupported envelope version {ver}"),
393 Self::TooShortForHeader { have, need } => {
394 write!(
395 f,
396 "header too short for version: have {have} bytes, need {need}"
397 )
398 }
399 Self::UnknownFrameType { byte } => write!(f, "unknown frame type byte {byte}"),
400 Self::ReservedFlagBits { flags } => {
401 write!(f, "reserved flag bits set in flags 0b{flags:08b}")
402 }
403 Self::ReservedPriorityBits { flags } => {
404 write!(f, "reserved priority bits set in flags 0b{flags:08b}")
405 }
406 Self::ReservedAdmissionClass { flags } => {
407 write!(f, "reserved admission class set in flags 0b{flags:08b}")
408 }
409 Self::SheddableIllegalFrameType { ty, flags } => write!(
410 f,
411 "SHEDDABLE admission class is illegal on {ty:?} in flags 0b{flags:08b}"
412 ),
413 Self::NonzeroEpochOnControlChannel { epoch } => {
414 write!(f, "control channel carried nonzero epoch {epoch}")
415 }
416 Self::PureHeaderFrameWithBody { ty, len } => {
417 write!(
418 f,
419 "pure-header frame {ty:?} declared non-zero body length {len}"
420 )
421 }
422 }
423 }
424}
425
426impl Error for DecodeError {}
427
428fn header_len_for_version(ver: u8) -> Option<usize> {
431 match ver {
432 PROTOCOL_VERSION => Some(HEADER_LEN),
433 _ => None,
434 }
435}
436
437pub fn decode_header(bytes: &[u8]) -> Result<EnvelopeHeader, DecodeError> {
445 if bytes.len() < FROZEN_PREFIX_LEN {
446 return Err(DecodeError::TooShortForPrefix { have: bytes.len() });
447 }
448 let ver = bytes[4];
449 let need = header_len_for_version(ver).ok_or(DecodeError::UnsupportedVersion { ver })?;
450 if bytes.len() < need {
451 return Err(DecodeError::TooShortForHeader {
452 have: bytes.len(),
453 need,
454 });
455 }
456
457 let len = u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
458 let ty =
459 FrameType::from_u8(bytes[5]).ok_or(DecodeError::UnknownFrameType { byte: bytes[5] })?;
460 let flags = Flags(bytes[6]);
461 if flags.has_reserved_bits() {
462 return Err(DecodeError::ReservedFlagBits { flags: bytes[6] });
463 }
464 if flags.priority().is_none() {
465 return Err(DecodeError::ReservedPriorityBits { flags: bytes[6] });
466 }
467 let admission_class = flags
468 .admission_class()
469 .ok_or(DecodeError::ReservedAdmissionClass { flags: bytes[6] })?;
470 if admission_class == AdmissionClass::Sheddable
471 && !matches!(ty, FrameType::Push | FrameType::StreamData)
472 {
473 return Err(DecodeError::SheddableIllegalFrameType {
474 ty,
475 flags: bytes[6],
476 });
477 }
478 if ty.is_pure_header() && len != 0 {
479 return Err(DecodeError::PureHeaderFrameWithBody { ty, len });
480 }
481 let channel = u16::from_le_bytes([bytes[7], bytes[8]]);
482 let epoch = u32::from_le_bytes([bytes[9], bytes[10], bytes[11], bytes[12]]);
483 if channel == 0 && epoch != 0 {
484 return Err(DecodeError::NonzeroEpochOnControlChannel { epoch });
485 }
486 let corr = u64::from_le_bytes([
487 bytes[13], bytes[14], bytes[15], bytes[16], bytes[17], bytes[18], bytes[19], bytes[20],
488 ]);
489
490 Ok(EnvelopeHeader {
491 len,
492 ver,
493 ty,
494 flags,
495 channel,
496 epoch,
497 corr,
498 })
499}
500
501#[cfg(test)]
502mod tests {
503 use super::*;
504
505 fn hdr(len: u32, ty: FrameType, flags: Flags, channel: u16, corr: u64) -> EnvelopeHeader {
506 hdr_with_epoch(len, ty, flags, channel, u32::from(channel != 0), corr)
507 }
508
509 fn hdr_with_epoch(
510 len: u32,
511 ty: FrameType,
512 flags: Flags,
513 channel: u16,
514 epoch: u32,
515 corr: u64,
516 ) -> EnvelopeHeader {
517 EnvelopeHeader {
518 len,
519 ver: PROTOCOL_VERSION,
520 ty,
521 flags,
522 channel,
523 epoch,
524 corr,
525 }
526 }
527
528 #[test]
529 fn bind_identity_round_trips_json() {
530 let identity = BindIdentity {
531 project_root: PathBuf::from("/tmp/project"),
532 harness: "opencode".to_string(),
533 session: "session-1".to_string(),
534 };
535
536 let encoded = serde_json::to_vec(&identity).unwrap();
537 let decoded: BindIdentity = serde_json::from_slice(&encoded).unwrap();
538
539 assert_eq!(decoded, identity);
540 }
541
542 #[test]
543 fn route_target_variants_round_trip_json() {
544 let targets = [
545 RouteTarget::ToolProvider {
546 module_id: "aft".to_string(),
547 },
548 RouteTarget::ManagementSurface {
549 module_id: "memory".to_string(),
550 },
551 RouteTarget::InternalService {
552 module_id: "bus".to_string(),
553 service_id: "dm".to_string(),
554 },
555 ];
556
557 for target in targets {
558 let encoded = serde_json::to_vec(&target).unwrap();
559 let decoded: RouteTarget = serde_json::from_slice(&encoded).unwrap();
560 assert_eq!(decoded, target);
561 }
562 }
563
564 #[test]
565 fn error_body_round_trips_json() {
566 let body = ErrorBody {
567 code: "config_divergence".to_string(),
568 message: "active config differs".to_string(),
569 };
570
571 let encoded = serde_json::to_vec(&body).unwrap();
572 let decoded: ErrorBody = serde_json::from_slice(&encoded).unwrap();
573
574 assert_eq!(decoded, body);
575 }
576
577 #[test]
578 fn round_trip_request() {
579 let h = hdr(
580 1234,
581 FrameType::Request,
582 Flags::new(false, Priority::Interactive, false),
583 42,
584 0xDEAD_BEEF_0000_0001,
585 );
586 let decoded = decode_header(&h.encode()).unwrap();
587 assert_eq!(h, decoded);
588 }
589
590 #[test]
591 fn round_trip_all_frame_types() {
592 for b in 0u8..=11 {
593 let ty = FrameType::from_u8(b).unwrap();
594 let h = hdr(0, ty, Flags::new(false, Priority::Passive, false), 0, 0);
595 assert_eq!(decode_header(&h.encode()).unwrap().ty, ty);
596 }
597 }
598
599 #[test]
600 fn pure_header_frame_has_zero_len() {
601 let h = hdr(
603 0,
604 FrameType::Cancel,
605 Flags::new(false, Priority::Passive, false),
606 7,
607 99,
608 );
609 let d = decode_header(&h.encode()).unwrap();
610 assert_eq!(d.len, 0);
611 assert_eq!(d.corr, 99);
612 }
613
614 #[test]
615 fn flags_round_trip() {
616 let f = Flags::new(true, Priority::Background, true)
617 .with_admission_class(AdmissionClass::Expedite);
618 assert!(f.is_binary());
619 assert!(f.is_last());
620 assert_eq!(f.priority(), Some(Priority::Background));
621 assert_eq!(f.admission_class(), Some(AdmissionClass::Expedite));
622 let h = hdr(8, FrameType::StreamData, f, 1, 1);
623 assert_eq!(decode_header(&h.encode()).unwrap().flags, f);
624 }
625
626 #[test]
627 fn little_endian_and_frozen_prefix_layout() {
628 let h = hdr_with_epoch(
629 0x0403_0201,
630 FrameType::Request,
631 Flags(0),
632 0x0605,
633 0x0a09_0807,
634 0x1211_100f_0e0d_0c0b,
635 );
636 let buf = h.encode();
637 assert_eq!(&buf[0..4], &[1, 2, 3, 4]);
638 assert_eq!(buf[4], PROTOCOL_VERSION);
639 assert_eq!(&buf[7..9], &[5, 6]);
640 assert_eq!(&buf[9..13], &[7, 8, 9, 10]);
641 assert_eq!(&buf[13..21], &[11, 12, 13, 14, 15, 16, 17, 18]);
642 assert_eq!(buf.len(), HEADER_LEN);
643 }
644
645 #[test]
646 fn reject_too_short_for_prefix() {
647 assert_eq!(
648 decode_header(&[0, 0, 0, 0]),
649 Err(DecodeError::TooShortForPrefix { have: 4 })
650 );
651 }
652
653 #[test]
654 fn reject_too_short_for_header() {
655 let mut b = [0u8; 10];
657 b[4] = PROTOCOL_VERSION;
658 assert_eq!(
659 decode_header(&b),
660 Err(DecodeError::TooShortForHeader {
661 have: 10,
662 need: HEADER_LEN
663 })
664 );
665 }
666
667 #[test]
668 fn reject_unsupported_version() {
669 let mut b = [0u8; HEADER_LEN];
670 b[4] = 1;
671 assert_eq!(
672 decode_header(&b),
673 Err(DecodeError::UnsupportedVersion { ver: 1 })
674 );
675 }
676
677 #[test]
678 fn reject_unknown_frame_type() {
679 let mut b = [0u8; HEADER_LEN];
680 b[4] = PROTOCOL_VERSION;
681 b[5] = 99;
682 assert_eq!(
683 decode_header(&b),
684 Err(DecodeError::UnknownFrameType { byte: 99 })
685 );
686 }
687
688 #[test]
689 fn reject_reserved_flag_bits() {
690 let mut b = [0u8; HEADER_LEN];
691 b[4] = PROTOCOL_VERSION;
692 b[5] = FrameType::Request as u8;
693 b[6] = 0b1000_0000; assert_eq!(
695 decode_header(&b),
696 Err(DecodeError::ReservedFlagBits { flags: 0b1000_0000 })
697 );
698 }
699
700 #[test]
701 fn reject_reserved_priority_bits() {
702 let mut b = [0u8; HEADER_LEN];
703 b[4] = PROTOCOL_VERSION;
704 b[5] = FrameType::Request as u8;
705 b[6] = 0b0000_0110; assert_eq!(
707 decode_header(&b),
708 Err(DecodeError::ReservedPriorityBits { flags: 0b0000_0110 })
709 );
710 }
711
712 #[test]
713 fn reject_pure_header_frame_with_body_len() {
714 let h = hdr(
715 1,
716 FrameType::Ping,
717 Flags::new(false, Priority::Passive, false),
718 0,
719 1,
720 );
721 assert_eq!(
722 decode_header(&h.encode()),
723 Err(DecodeError::PureHeaderFrameWithBody {
724 ty: FrameType::Ping,
725 len: 1
726 })
727 );
728 }
729
730 #[test]
731 fn epoch_boundaries_round_trip() {
732 for (channel, epoch) in [(0, 0), (1, 1), (u16::MAX, u32::MAX)] {
733 let h = hdr_with_epoch(
734 0,
735 FrameType::Request,
736 Flags::new(false, Priority::Passive, false),
737 channel,
738 epoch,
739 9,
740 );
741 assert_eq!(decode_header(&h.encode()).unwrap(), h);
742 }
743 }
744
745 #[test]
746 fn admission_classes_accept_three_values_and_reject_reserved_value() {
747 for (ty, admission_class) in [
748 (FrameType::Request, AdmissionClass::Normal),
749 (FrameType::Request, AdmissionClass::Expedite),
750 (FrameType::Push, AdmissionClass::Sheddable),
751 (FrameType::StreamData, AdmissionClass::Sheddable),
752 ] {
753 let flags = Flags::new(false, Priority::Interactive, false)
754 .with_admission_class(admission_class);
755 let h = hdr(0, ty, flags, 1, 2);
756 assert_eq!(decode_header(&h.encode()).unwrap().flags, flags);
757 }
758
759 let mut h = hdr(
760 0,
761 FrameType::Push,
762 Flags::new(false, Priority::Passive, false),
763 1,
764 2,
765 )
766 .encode();
767 h[6] |= 0b0011_0000;
768 assert_eq!(
769 decode_header(&h),
770 Err(DecodeError::ReservedAdmissionClass { flags: h[6] })
771 );
772 }
773
774 #[test]
775 fn sheddable_rejected_on_every_illegal_frame_type() {
776 let flags = Flags::new(false, Priority::Passive, false)
777 .with_admission_class(AdmissionClass::Sheddable);
778 for ty in [
779 FrameType::Request,
780 FrameType::Response,
781 FrameType::StreamEnd,
782 FrameType::Error,
783 FrameType::Cancel,
784 FrameType::Ping,
785 FrameType::Pong,
786 FrameType::Hello,
787 FrameType::HelloAck,
788 FrameType::Goodbye,
789 ] {
790 let h = hdr(0, ty, flags, 1, 2);
791 assert_eq!(
792 decode_header(&h.encode()),
793 Err(DecodeError::SheddableIllegalFrameType { ty, flags: flags.0 })
794 );
795 }
796 }
797
798 #[test]
799 fn nonzero_epoch_on_control_channel_is_rejected() {
800 let h = hdr_with_epoch(
801 0,
802 FrameType::Request,
803 Flags::new(false, Priority::Passive, false),
804 0,
805 u32::MAX,
806 2,
807 );
808 assert_eq!(
809 decode_header(&h.encode()),
810 Err(DecodeError::NonzeroEpochOnControlChannel { epoch: u32::MAX })
811 );
812 }
813}