1pub use meerkat_core::comms::{
11 CommsCommandError, CommsPeerRequestIntent, InputSource, InputStreamMode, PeerAddress,
12 PeerCapabilitySet, PeerDirectoryEntry, PeerDirectoryListing, PeerDirectorySource, PeerId,
13 PeerLifecycleKind, PeerName, PeerSendability, PeerTransport, SendTaintOverride,
14 SenderContentTaint,
15};
16pub use meerkat_core::interaction::ResponseStatus;
17pub use meerkat_core::types::HandlingMode;
18
19use super::supervisor_bridge::{BridgeCommand, BridgePeerSpec, BridgeReply};
20use serde::{Deserialize, Serialize};
21
22#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
29#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
30#[serde(deny_unknown_fields)]
31pub struct CommsPeerLifecycleParams {
32 pub peer: String,
33 #[serde(default, skip_serializing_if = "Option::is_none")]
34 pub role: Option<String>,
35 #[serde(default, skip_serializing_if = "Option::is_none")]
36 pub description: Option<String>,
37 #[serde(default, skip_serializing_if = "Option::is_none")]
38 pub peer_spec: Option<BridgePeerSpec>,
39}
40
41impl CommsPeerLifecycleParams {
42 fn into_json_value(self) -> Result<serde_json::Value, serde_json::Error> {
43 serde_json::to_value(self)
44 }
45}
46
47#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
50#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
51#[serde(deny_unknown_fields)]
52pub struct CommsChecksumTokenParams {
53 pub subject: String,
54}
55
56#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
58#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
59pub enum CommsChecksumTokenResultIntent {
60 #[serde(rename = "checksum_token")]
61 ChecksumToken,
62}
63
64#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
66#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
67#[serde(deny_unknown_fields)]
68pub struct CommsChecksumTokenResult {
69 pub request_intent: CommsChecksumTokenResultIntent,
70 pub request_subject: String,
71 pub token: String,
72}
73
74#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
76#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
77#[serde(untagged)]
78pub enum CommsPeerRequestParams {
79 SupervisorBridge(Box<BridgeCommand>),
80 ChecksumToken(CommsChecksumTokenParams),
81}
82
83impl CommsPeerRequestParams {
84 fn matches_intent(&self, intent: &CommsPeerRequestIntent) -> bool {
85 matches!(
86 (intent, self),
87 (
88 CommsPeerRequestIntent::SupervisorBridge,
89 Self::SupervisorBridge(_)
90 ) | (
91 CommsPeerRequestIntent::ChecksumToken,
92 Self::ChecksumToken(_)
93 )
94 )
95 }
96
97 fn into_json_value(self) -> Result<serde_json::Value, serde_json::Error> {
98 match self {
99 Self::SupervisorBridge(params) => serde_json::to_value(params),
100 Self::ChecksumToken(params) => serde_json::to_value(params),
101 }
102 }
103}
104
105#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
110#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
111#[serde(untagged)]
112pub enum CommsPeerResponseResult {
113 SupervisorBridge(BridgeReply),
114 ChecksumToken(CommsChecksumTokenResult),
115}
116
117impl CommsPeerResponseResult {
118 fn into_json_value(self) -> Result<serde_json::Value, serde_json::Error> {
119 match self {
120 Self::SupervisorBridge(result) => serde_json::to_value(result),
121 Self::ChecksumToken(result) => serde_json::to_value(result),
122 }
123 }
124}
125
126impl From<BridgeReply> for CommsPeerResponseResult {
127 fn from(reply: BridgeReply) -> Self {
128 Self::SupervisorBridge(reply)
129 }
130}
131
132#[derive(Debug, thiserror::Error)]
135pub enum CommsCommandProjectionError {
136 #[error(transparent)]
137 Command(#[from] CommsCommandError),
138 #[error("peer_request params do not match typed intent {intent}")]
139 IntentParamsMismatch { intent: &'static str },
140 #[error("failed to project typed comms {field} to compatibility JSON: {source}")]
141 CompatibilityJson {
142 field: &'static str,
143 #[source]
144 source: serde_json::Error,
145 },
146}
147
148impl CommsCommandProjectionError {
149 fn compatibility_json(field: &'static str, source: serde_json::Error) -> Self {
150 Self::CompatibilityJson { field, source }
151 }
152}
153
154fn single_authority_body(body: String, blocks: &Option<Vec<meerkat_core::ContentBlock>>) -> String {
168 match blocks {
169 Some(blocks) => blocks
170 .iter()
171 .filter_map(|block| match block {
172 meerkat_core::ContentBlock::Text { text } => Some(text.as_str()),
173 _ => None,
174 })
175 .collect::<Vec<_>>()
176 .join("\n"),
177 None => body,
178 }
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
183#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
184#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
185pub enum CommsCommandRequest {
186 Input {
187 body: String,
188 #[serde(default, skip_serializing_if = "Option::is_none")]
189 blocks: Option<Vec<meerkat_core::ContentBlock>>,
190 #[serde(default, skip_serializing_if = "Option::is_none")]
191 source: Option<InputSource>,
192 #[serde(default, skip_serializing_if = "Option::is_none")]
193 stream: Option<InputStreamMode>,
194 #[serde(default, skip_serializing_if = "Option::is_none")]
195 handling_mode: Option<HandlingMode>,
196 #[serde(default, skip_serializing_if = "Option::is_none")]
197 allow_self_session: Option<bool>,
198 },
199 PeerMessage {
200 to: PeerId,
201 body: String,
202 #[serde(default, skip_serializing_if = "Option::is_none")]
203 blocks: Option<Vec<meerkat_core::ContentBlock>>,
204 #[serde(default, skip_serializing_if = "Option::is_none")]
207 content_taint: Option<SendTaintOverride>,
208 #[serde(default, skip_serializing_if = "Option::is_none")]
209 handling_mode: Option<HandlingMode>,
210 },
211 PeerLifecycle {
212 to: PeerId,
213 lifecycle_kind: PeerLifecycleKind,
214 params: CommsPeerLifecycleParams,
215 },
216 PeerRequest {
217 to: PeerId,
218 intent: CommsPeerRequestIntent,
219 params: CommsPeerRequestParams,
220 #[serde(default, skip_serializing_if = "Option::is_none")]
221 blocks: Option<Vec<meerkat_core::ContentBlock>>,
222 #[serde(default, skip_serializing_if = "Option::is_none")]
225 content_taint: Option<SendTaintOverride>,
226 #[serde(default, skip_serializing_if = "Option::is_none")]
227 handling_mode: Option<HandlingMode>,
228 #[serde(default, skip_serializing_if = "Option::is_none")]
229 stream: Option<InputStreamMode>,
230 },
231 PeerResponse {
232 to: PeerId,
233 in_reply_to: meerkat_core::interaction::InteractionId,
234 status: ResponseStatus,
235 #[serde(default, skip_serializing_if = "Option::is_none")]
236 result: Option<CommsPeerResponseResult>,
237 #[serde(default, skip_serializing_if = "Option::is_none")]
238 blocks: Option<Vec<meerkat_core::ContentBlock>>,
239 #[serde(default, skip_serializing_if = "Option::is_none")]
242 content_taint: Option<SendTaintOverride>,
243 #[serde(default, skip_serializing_if = "Option::is_none")]
244 handling_mode: Option<HandlingMode>,
245 },
246}
247
248impl CommsCommandRequest {
249 pub fn peer_label(&self) -> Option<String> {
250 match self {
251 Self::Input { .. } => None,
252 Self::PeerMessage { to, .. }
253 | Self::PeerLifecycle { to, .. }
254 | Self::PeerRequest { to, .. }
255 | Self::PeerResponse { to, .. } => Some(to.to_string()),
256 }
257 }
258
259 pub fn into_core_request(
260 self,
261 ) -> Result<meerkat_core::comms::CommsCommandRequest, CommsCommandProjectionError> {
262 Ok(match self {
263 Self::Input {
264 body,
265 blocks,
266 source,
267 stream,
268 handling_mode,
269 allow_self_session,
270 } => meerkat_core::comms::CommsCommandRequest::Input {
271 body: single_authority_body(body, &blocks),
272 blocks,
273 source,
274 stream,
275 handling_mode,
276 allow_self_session,
277 },
278 Self::PeerMessage {
279 to,
280 body,
281 blocks,
282 content_taint,
283 handling_mode,
284 } => meerkat_core::comms::CommsCommandRequest::PeerMessage {
285 to,
286 body: single_authority_body(body, &blocks),
287 blocks,
288 content_taint,
289 handling_mode,
290 },
291 Self::PeerLifecycle {
292 to,
293 lifecycle_kind,
294 params,
295 } => meerkat_core::comms::CommsCommandRequest::PeerLifecycle {
296 to,
297 lifecycle_kind,
298 params: params.into_json_value().map_err(|source| {
299 CommsCommandProjectionError::compatibility_json("peer_lifecycle.params", source)
300 })?,
301 },
302 Self::PeerRequest {
303 to,
304 intent,
305 params,
306 blocks,
307 content_taint,
308 handling_mode,
309 stream,
310 } => {
311 if !params.matches_intent(&intent) {
312 return Err(CommsCommandProjectionError::IntentParamsMismatch {
313 intent: intent.as_str(),
314 });
315 }
316 meerkat_core::comms::CommsCommandRequest::PeerRequest {
317 to,
318 intent,
326 params: params.into_json_value().map_err(|source| {
327 CommsCommandProjectionError::compatibility_json(
328 "peer_request.params",
329 source,
330 )
331 })?,
332 blocks,
333 content_taint,
334 handling_mode,
335 stream,
336 }
337 }
338 Self::PeerResponse {
339 to,
340 in_reply_to,
341 status,
342 result,
343 blocks,
344 content_taint,
345 handling_mode,
346 } => meerkat_core::comms::CommsCommandRequest::PeerResponse {
347 to,
348 in_reply_to,
349 status,
350 result: match result {
351 Some(result) => result.into_json_value().map_err(|source| {
352 CommsCommandProjectionError::compatibility_json(
353 "peer_response.result",
354 source,
355 )
356 })?,
357 None => serde_json::Value::Null,
358 },
359 blocks,
360 content_taint,
361 handling_mode,
362 },
363 })
364 }
365
366 pub fn into_command(
367 self,
368 session_id: &meerkat_core::types::SessionId,
369 ) -> Result<meerkat_core::comms::CommsCommand, CommsCommandProjectionError> {
370 Ok(self.into_core_request()?.into_command(session_id)?)
371 }
372}
373
374#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
376#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
377#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
378pub enum CommsSendParams {
379 Input {
380 session_id: String,
381 body: String,
382 #[serde(default, skip_serializing_if = "Option::is_none")]
383 blocks: Option<Vec<meerkat_core::ContentBlock>>,
384 #[serde(default, skip_serializing_if = "Option::is_none")]
385 source: Option<InputSource>,
386 #[serde(default, skip_serializing_if = "Option::is_none")]
387 stream: Option<InputStreamMode>,
388 #[serde(default, skip_serializing_if = "Option::is_none")]
389 handling_mode: Option<HandlingMode>,
390 #[serde(default, skip_serializing_if = "Option::is_none")]
391 allow_self_session: Option<bool>,
392 },
393 PeerMessage {
394 session_id: String,
395 to: PeerId,
396 body: String,
397 #[serde(default, skip_serializing_if = "Option::is_none")]
398 blocks: Option<Vec<meerkat_core::ContentBlock>>,
399 #[serde(default, skip_serializing_if = "Option::is_none")]
402 content_taint: Option<SendTaintOverride>,
403 #[serde(default, skip_serializing_if = "Option::is_none")]
404 handling_mode: Option<HandlingMode>,
405 },
406 PeerLifecycle {
407 session_id: String,
408 to: PeerId,
409 lifecycle_kind: PeerLifecycleKind,
410 params: CommsPeerLifecycleParams,
411 },
412 PeerRequest {
413 session_id: String,
414 to: PeerId,
415 intent: CommsPeerRequestIntent,
416 params: CommsPeerRequestParams,
417 #[serde(default, skip_serializing_if = "Option::is_none")]
418 blocks: Option<Vec<meerkat_core::ContentBlock>>,
419 #[serde(default, skip_serializing_if = "Option::is_none")]
422 content_taint: Option<SendTaintOverride>,
423 #[serde(default, skip_serializing_if = "Option::is_none")]
424 handling_mode: Option<HandlingMode>,
425 #[serde(default, skip_serializing_if = "Option::is_none")]
426 stream: Option<InputStreamMode>,
427 },
428 PeerResponse {
429 session_id: String,
430 to: PeerId,
431 in_reply_to: meerkat_core::interaction::InteractionId,
432 status: ResponseStatus,
433 #[serde(default, skip_serializing_if = "Option::is_none")]
434 result: Option<CommsPeerResponseResult>,
435 #[serde(default, skip_serializing_if = "Option::is_none")]
436 blocks: Option<Vec<meerkat_core::ContentBlock>>,
437 #[serde(default, skip_serializing_if = "Option::is_none")]
440 content_taint: Option<SendTaintOverride>,
441 #[serde(default, skip_serializing_if = "Option::is_none")]
442 handling_mode: Option<HandlingMode>,
443 },
444}
445
446impl CommsSendParams {
447 pub fn session_id(&self) -> &str {
448 match self {
449 Self::Input { session_id, .. }
450 | Self::PeerMessage { session_id, .. }
451 | Self::PeerLifecycle { session_id, .. }
452 | Self::PeerRequest { session_id, .. }
453 | Self::PeerResponse { session_id, .. } => session_id,
454 }
455 }
456
457 pub fn peer_label(&self) -> Option<String> {
459 match self {
460 Self::Input { .. } => None,
461 Self::PeerMessage { to, .. }
462 | Self::PeerLifecycle { to, .. }
463 | Self::PeerRequest { to, .. }
464 | Self::PeerResponse { to, .. } => Some(to.to_string()),
465 }
466 }
467
468 pub fn into_command(self) -> CommsCommandRequest {
469 match self {
470 Self::Input {
471 body,
472 blocks,
473 source,
474 stream,
475 handling_mode,
476 allow_self_session,
477 ..
478 } => CommsCommandRequest::Input {
479 body,
480 blocks,
481 source,
482 stream,
483 handling_mode,
484 allow_self_session,
485 },
486 Self::PeerMessage {
487 to,
488 body,
489 blocks,
490 content_taint,
491 handling_mode,
492 ..
493 } => CommsCommandRequest::PeerMessage {
494 to,
495 body,
496 blocks,
497 content_taint,
498 handling_mode,
499 },
500 Self::PeerLifecycle {
501 to,
502 lifecycle_kind,
503 params,
504 ..
505 } => CommsCommandRequest::PeerLifecycle {
506 to,
507 lifecycle_kind,
508 params,
509 },
510 Self::PeerRequest {
511 to,
512 intent,
513 params,
514 blocks,
515 content_taint,
516 handling_mode,
517 stream,
518 ..
519 } => CommsCommandRequest::PeerRequest {
520 to,
521 intent,
522 params,
523 blocks,
524 content_taint,
525 handling_mode,
526 stream,
527 },
528 Self::PeerResponse {
529 to,
530 in_reply_to,
531 status,
532 result,
533 blocks,
534 content_taint,
535 handling_mode,
536 ..
537 } => CommsCommandRequest::PeerResponse {
538 to,
539 in_reply_to,
540 status,
541 result,
542 blocks,
543 content_taint,
544 handling_mode,
545 },
546 }
547 }
548}
549
550#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
552#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
553#[serde(deny_unknown_fields)]
554pub struct CommsPeersParams {
555 pub session_id: String,
556}
557
558#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
560#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
561#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
562pub enum CommsSendResult {
563 InputAccepted {
564 interaction_id: String,
565 stream_reserved: bool,
566 },
567 PeerMessageSent {
568 envelope_id: String,
569 acked: bool,
570 },
571 PeerLifecycleSent {
572 envelope_id: String,
573 },
574 PeerRequestSent {
575 envelope_id: String,
576 interaction_id: String,
577 request_id: String,
578 stream_reserved: bool,
579 },
580 PeerResponseSent {
581 envelope_id: String,
582 in_reply_to: String,
583 },
584}
585
586impl From<meerkat_core::comms::SendReceipt> for CommsSendResult {
587 fn from(receipt: meerkat_core::comms::SendReceipt) -> Self {
588 match receipt {
589 meerkat_core::comms::SendReceipt::InputAccepted {
590 interaction_id,
591 stream_reserved,
592 } => Self::InputAccepted {
593 interaction_id: interaction_id.0.to_string(),
594 stream_reserved,
595 },
596 meerkat_core::comms::SendReceipt::PeerMessageSent { envelope_id, acked } => {
597 Self::PeerMessageSent {
598 envelope_id: envelope_id.to_string(),
599 acked,
600 }
601 }
602 meerkat_core::comms::SendReceipt::PeerLifecycleSent { envelope_id } => {
603 Self::PeerLifecycleSent {
604 envelope_id: envelope_id.to_string(),
605 }
606 }
607 meerkat_core::comms::SendReceipt::PeerRequestSent {
608 envelope_id,
609 interaction_id,
610 stream_reserved,
611 } => {
612 let envelope_id = envelope_id.to_string();
613 Self::PeerRequestSent {
614 request_id: envelope_id.clone(),
615 envelope_id,
616 interaction_id: interaction_id.0.to_string(),
617 stream_reserved,
618 }
619 }
620 meerkat_core::comms::SendReceipt::PeerResponseSent {
621 envelope_id,
622 in_reply_to,
623 } => Self::PeerResponseSent {
624 envelope_id: envelope_id.to_string(),
625 in_reply_to: in_reply_to.0.to_string(),
626 },
627 }
628 }
629}
630
631#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
633#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
634#[serde(rename_all = "snake_case")]
635pub enum CommsPeerUnreachableReason {
636 OfflineOrNoAck,
637 TransportError,
638}
639
640#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
645#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
646#[serde(tag = "code", rename_all = "snake_case")]
647pub enum CommsSendErrorData {
648 PeerNotFoundOrNotTrusted {
649 peer: String,
650 message: String,
651 },
652 PeerUnreachable {
653 peer: String,
654 reason: CommsPeerUnreachableReason,
655 message: String,
656 #[serde(default, skip_serializing_if = "Option::is_none")]
657 details: Option<String>,
658 },
659 PeerAdmissionDropped {
660 peer: String,
661 reason: meerkat_core::comms::AdmissionDropReason,
662 message: String,
663 },
664 SendFailed {
665 message: String,
666 },
667 InvalidCommand {
668 message: String,
669 },
670}
671
672impl CommsSendErrorData {
673 #[must_use]
675 pub fn message(&self) -> &str {
676 match self {
677 Self::PeerNotFoundOrNotTrusted { message, .. }
678 | Self::PeerUnreachable { message, .. }
679 | Self::PeerAdmissionDropped { message, .. }
680 | Self::SendFailed { message }
681 | Self::InvalidCommand { message } => message,
682 }
683 }
684}
685
686pub type CommsPeerEntry = PeerDirectoryEntry;
687
688#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
690#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
691pub struct CommsPeersResult {
692 pub peers: Vec<PeerDirectoryEntry>,
693}
694
695impl CommsPeersResult {
696 pub fn from_entries(entries: &[meerkat_core::comms::PeerDirectoryEntry]) -> Self {
697 Self {
698 peers: entries.to_vec(),
699 }
700 }
701}
702
703impl From<PeerDirectoryListing> for CommsPeersResult {
704 fn from(listing: PeerDirectoryListing) -> Self {
705 Self {
706 peers: listing.peers,
707 }
708 }
709}
710
711#[cfg(test)]
712#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
713mod tests {
714 use super::super::supervisor_bridge::SUPERVISOR_BRIDGE_INTENT;
715 use super::*;
716 use serde_json::json;
717
718 fn peer_id() -> PeerId {
719 PeerId::new()
720 }
721
722 fn bridge_peer_spec() -> serde_json::Value {
723 let pubkey = [7u8; 32];
724 json!({
725 "name": "supervisor",
726 "peer_id": PeerId::from_ed25519_pubkey(&pubkey).to_string(),
727 "address": "inproc://supervisor",
728 "pubkey": pubkey,
729 })
730 }
731
732 fn supervisor_bridge_params() -> serde_json::Value {
733 json!({
734 "command": "observe_member",
735 "supervisor": bridge_peer_spec(),
736 "epoch": 1,
737 "protocol_version": 2,
738 })
739 }
740
741 #[test]
742 fn peer_request_unknown_intent_fails_closed() {
743 let err = serde_json::from_value::<CommsSendParams>(json!({
744 "session_id": "sid_1",
745 "kind": "peer_request",
746 "to": peer_id().to_string(),
747 "intent": "local.default",
748 "params": {}
749 }))
750 .expect_err("unknown public comms intent must fail at serde boundary");
751
752 let message = err.to_string();
753 assert!(
754 message.contains("local.default") || message.contains("variant"),
755 "error should mention the rejected intent, got: {message}"
756 );
757 }
758
759 #[test]
760 fn peer_request_malformed_params_cannot_be_typed_success() {
761 let err = serde_json::from_value::<CommsSendParams>(json!({
762 "session_id": "sid_1",
763 "kind": "peer_request",
764 "to": peer_id().to_string(),
765 "intent": "supervisor.bridge",
766 "params": "not-an-object"
767 }))
768 .expect_err("malformed supervisor bridge params must not deserialize");
769
770 let message = err.to_string();
771 assert!(
772 message.contains("invalid type")
773 || message.contains("params")
774 || message.contains("did not match any variant"),
775 "expected typed params error, got: {message}"
776 );
777 }
778
779 #[test]
780 fn peer_response_malformed_result_cannot_be_typed_success() {
781 let err = serde_json::from_value::<CommsSendParams>(json!({
782 "session_id": "sid_1",
783 "kind": "peer_response",
784 "to": peer_id().to_string(),
785 "in_reply_to": uuid::Uuid::new_v4().to_string(),
786 "status": "completed",
787 "result": {
788 "result": "ack",
789 "ok": "yes"
790 }
791 }))
792 .expect_err("malformed typed result must not deserialize");
793
794 let message = err.to_string();
795 assert!(
796 message.contains("ok")
797 || message.contains("invalid type")
798 || message.contains("did not match any variant"),
799 "expected typed result error, got: {message}"
800 );
801 }
802
803 #[test]
804 fn comms_send_result_unknown_field_fails_closed() {
805 let err = serde_json::from_value::<CommsSendResult>(json!({
806 "kind": "peer_request_sent",
807 "envelope_id": uuid::Uuid::new_v4().to_string(),
808 "interaction_id": uuid::Uuid::new_v4().to_string(),
809 "request_id": uuid::Uuid::new_v4().to_string(),
810 "stream_reserved": true,
811 "extra_behavior": true
812 }))
813 .expect_err("unknown result fields must fail at serde boundary");
814
815 let message = err.to_string();
816 assert!(
817 message.contains("extra_behavior") || message.contains("unknown field"),
818 "expected unknown result field error, got: {message}"
819 );
820 }
821
822 #[test]
823 fn comms_send_error_data_preserves_admission_drop_reason() {
824 let value = serde_json::to_value(CommsSendErrorData::PeerAdmissionDropped {
825 peer: "peer-a".to_string(),
826 reason: meerkat_core::comms::AdmissionDropReason::InboxFull,
827 message: "peer 'peer-a' rejected envelope at ingress: inbox_full".to_string(),
828 })
829 .expect("admission drop error data should serialize");
830
831 assert_eq!(value["code"], "peer_admission_dropped");
832 assert_eq!(value["peer"], "peer-a");
833 assert_eq!(value["reason"], "inbox_full");
834
835 let roundtrip: CommsSendErrorData =
836 serde_json::from_value(value).expect("admission drop error data should deserialize");
837 assert!(matches!(
838 roundtrip,
839 CommsSendErrorData::PeerAdmissionDropped {
840 reason: meerkat_core::comms::AdmissionDropReason::InboxFull,
841 ..
842 }
843 ));
844 }
845
846 #[test]
847 fn public_peer_request_projects_typed_intent_and_params_to_core() {
848 let params = serde_json::from_value::<CommsSendParams>(json!({
849 "session_id": "sid_1",
850 "kind": "peer_request",
851 "to": peer_id().to_string(),
852 "intent": "supervisor.bridge",
853 "params": supervisor_bridge_params(),
854 "handling_mode": "queue",
855 "stream": "reserve_interaction"
856 }))
857 .expect("typed supervisor bridge request should deserialize");
858
859 let session_id = meerkat_core::types::SessionId::new();
860 let command = params
861 .into_command()
862 .into_command(&session_id)
863 .expect("typed comms request should project to core command");
864
865 let meerkat_core::comms::CommsCommand::PeerRequest {
866 intent,
867 params,
868 handling_mode,
869 stream,
870 ..
871 } = command
872 else {
873 panic!("expected core peer request");
874 };
875
876 assert_eq!(intent, SUPERVISOR_BRIDGE_INTENT);
877 assert_eq!(params["command"], "observe_member");
878 assert_eq!(handling_mode, HandlingMode::Queue);
879 assert_eq!(stream, InputStreamMode::ReserveInteraction);
880 }
881
882 #[test]
883 fn public_checksum_token_request_projects_typed_intent_and_params_to_core() {
884 let params = serde_json::from_value::<CommsSendParams>(json!({
885 "session_id": "sid_1",
886 "kind": "peer_request",
887 "to": peer_id().to_string(),
888 "intent": "checksum_token",
889 "params": {
890 "subject": "alpha beta gamma"
891 },
892 "handling_mode": "steer",
893 "stream": "reserve_interaction"
894 }))
895 .expect("typed checksum token request should deserialize");
896
897 let session_id = meerkat_core::types::SessionId::new();
898 let command = params
899 .into_command()
900 .into_command(&session_id)
901 .expect("typed checksum token request should project to core command");
902
903 let meerkat_core::comms::CommsCommand::PeerRequest {
904 intent,
905 params,
906 handling_mode,
907 stream,
908 ..
909 } = command
910 else {
911 panic!("expected core peer request");
912 };
913
914 assert_eq!(intent, "checksum_token");
915 assert_eq!(params["subject"], "alpha beta gamma");
916 assert_eq!(handling_mode, HandlingMode::Steer);
917 assert_eq!(stream, InputStreamMode::ReserveInteraction);
918 }
919
920 #[test]
921 fn public_peer_request_rejects_intent_params_mismatch_before_dispatch() {
922 let params = serde_json::from_value::<CommsSendParams>(json!({
923 "session_id": "sid_1",
924 "kind": "peer_request",
925 "to": peer_id().to_string(),
926 "intent": "checksum_token",
927 "params": supervisor_bridge_params()
928 }))
929 .expect("mismatched typed params can deserialize but must not project");
930
931 let session_id = meerkat_core::types::SessionId::new();
932 let err = params
933 .into_command()
934 .into_command(&session_id)
935 .expect_err("intent/params mismatch must not become a core command");
936
937 assert!(
938 err.to_string().contains("checksum_token"),
939 "expected mismatch error to name intent, got: {err}"
940 );
941 }
942
943 #[test]
944 fn public_peer_response_result_projects_typed_bridge_reply_to_core() {
945 let params = serde_json::from_value::<CommsSendParams>(json!({
946 "session_id": "sid_1",
947 "kind": "peer_response",
948 "to": peer_id().to_string(),
949 "in_reply_to": uuid::Uuid::new_v4().to_string(),
950 "status": "completed",
951 "result": {
952 "result": "ack",
953 "ok": true
954 }
955 }))
956 .expect("typed bridge reply should deserialize");
957
958 let session_id = meerkat_core::types::SessionId::new();
959 let command = params
960 .into_command()
961 .into_command(&session_id)
962 .expect("typed comms response should project to core command");
963
964 let meerkat_core::comms::CommsCommand::PeerResponse { result, .. } = command else {
965 panic!("expected core peer response");
966 };
967
968 assert_eq!(result["result"], "ack");
969 assert_eq!(result["ok"], true);
970 }
971
972 #[test]
973 fn public_peer_response_result_projects_typed_checksum_token_to_core() {
974 let params = serde_json::from_value::<CommsSendParams>(json!({
975 "session_id": "sid_1",
976 "kind": "peer_response",
977 "to": peer_id().to_string(),
978 "in_reply_to": uuid::Uuid::new_v4().to_string(),
979 "status": "completed",
980 "result": {
981 "request_intent": "checksum_token",
982 "request_subject": "alpha beta gamma",
983 "token": "birch seventeen"
984 }
985 }))
986 .expect("typed checksum token reply should deserialize");
987
988 let session_id = meerkat_core::types::SessionId::new();
989 let command = params
990 .into_command()
991 .into_command(&session_id)
992 .expect("typed checksum token reply should project to core command");
993
994 let meerkat_core::comms::CommsCommand::PeerResponse { result, .. } = command else {
995 panic!("expected core peer response");
996 };
997
998 assert_eq!(result["request_intent"], "checksum_token");
999 assert_eq!(result["request_subject"], "alpha beta gamma");
1000 assert_eq!(result["token"], "birch seventeen");
1001 }
1002
1003 #[test]
1004 fn input_body_is_projected_from_blocks_single_authority() {
1005 let request = CommsCommandRequest::Input {
1009 body: "diverged caller body".to_string(),
1010 blocks: Some(vec![
1011 meerkat_core::ContentBlock::Text {
1012 text: "authoritative line one".to_string(),
1013 },
1014 meerkat_core::ContentBlock::Text {
1015 text: "authoritative line two".to_string(),
1016 },
1017 ]),
1018 source: None,
1019 stream: None,
1020 handling_mode: None,
1021 allow_self_session: None,
1022 };
1023
1024 let core = request
1025 .into_core_request()
1026 .expect("input request should project to core");
1027
1028 let meerkat_core::comms::CommsCommandRequest::Input { body, .. } = core else {
1029 panic!("expected core input request");
1030 };
1031 assert_eq!(body, "authoritative line one\nauthoritative line two");
1032 }
1033
1034 #[test]
1035 fn peer_message_body_stands_alone_when_no_blocks() {
1036 let request = CommsCommandRequest::PeerMessage {
1039 to: peer_id(),
1040 body: "lone body".to_string(),
1041 blocks: None,
1042 content_taint: None,
1043 handling_mode: None,
1044 };
1045
1046 let core = request
1047 .into_core_request()
1048 .expect("peer message should project to core");
1049
1050 let meerkat_core::comms::CommsCommandRequest::PeerMessage { body, .. } = core else {
1051 panic!("expected core peer message");
1052 };
1053 assert_eq!(body, "lone body");
1054 }
1055
1056 #[test]
1057 fn peer_message_content_taint_defaults_absent_and_threads_to_core() {
1058 let inherited = serde_json::from_value::<CommsSendParams>(json!({
1062 "session_id": "sid_1",
1063 "kind": "peer_message",
1064 "to": peer_id().to_string(),
1065 "body": "hello"
1066 }))
1067 .expect("peer message without content_taint should deserialize");
1068 let command = inherited
1069 .into_command()
1070 .into_command(&meerkat_core::types::SessionId::new())
1071 .expect("peer message should project to core command");
1072 let meerkat_core::comms::CommsCommand::PeerMessage { content_taint, .. } = command else {
1073 panic!("expected core peer message");
1074 };
1075 assert_eq!(content_taint, None);
1076
1077 for (wire, expected) in [
1078 (
1079 json!({"declare": "tainted"}),
1080 SendTaintOverride::Declare(SenderContentTaint::Tainted),
1081 ),
1082 (
1083 json!({"declare": "clean"}),
1084 SendTaintOverride::Declare(SenderContentTaint::Clean),
1085 ),
1086 (json!("undeclared"), SendTaintOverride::Undeclared),
1087 ] {
1088 let params = serde_json::from_value::<CommsSendParams>(json!({
1089 "session_id": "sid_1",
1090 "kind": "peer_message",
1091 "to": peer_id().to_string(),
1092 "body": "hello",
1093 "content_taint": wire
1094 }))
1095 .expect("peer message with content_taint should deserialize");
1096 let command = params
1097 .into_command()
1098 .into_command(&meerkat_core::types::SessionId::new())
1099 .expect("peer message should project to core command");
1100 let meerkat_core::comms::CommsCommand::PeerMessage { content_taint, .. } = command
1101 else {
1102 panic!("expected core peer message");
1103 };
1104 assert_eq!(content_taint, Some(expected));
1105 }
1106 }
1107
1108 #[test]
1109 fn public_peer_response_result_rejects_checksum_token_without_subject() {
1110 serde_json::from_value::<CommsSendParams>(json!({
1111 "session_id": "sid_1",
1112 "kind": "peer_response",
1113 "to": peer_id().to_string(),
1114 "in_reply_to": uuid::Uuid::new_v4().to_string(),
1115 "status": "completed",
1116 "result": {
1117 "request_intent": "checksum_token",
1118 "token": "birch seventeen"
1119 }
1120 }))
1121 .expect_err("checksum token replies must carry the request subject discriminator");
1122 }
1123}