1use std::net::SocketAddr;
11use std::sync::Arc;
12
13use bytes::Bytes;
14use futures::{SinkExt, StreamExt};
15use tokio::io::{AsyncRead, AsyncWrite};
16use tokio_util::codec::Framed;
17
18use crate::identity::{Keypair, PubKey, Signature};
19use crate::inbox::{AdmissionOutcome, DropReason, InboxSender};
20use crate::transport::TransportError;
21use crate::transport::codec::{EnvelopeFrame, TransportCodec};
22use crate::types::{Envelope, MessageKind};
23
24type ResolvedInboundIdentity = (Arc<Keypair>, InboxSender);
32
33pub async fn handle_connection<S>(
44 stream: S,
45 require_peer_auth: bool,
46 keypair: &Keypair,
47 inbox_sender: &InboxSender,
48) -> Result<(), IoTaskError>
49where
50 S: AsyncRead + AsyncWrite + Unpin,
51{
52 handle_connection_inner(stream, require_peer_auth, None, |to| {
55 (*to == keypair.public_key()).then(|| (Arc::new(keypair.clone()), inbox_sender.clone()))
56 })
57 .await
58}
59
60pub(crate) async fn handle_tcp_connection<S>(
66 stream: S,
67 observed_source: SocketAddr,
68 require_peer_auth: bool,
69 keypair: &Keypair,
70 inbox_sender: &InboxSender,
71) -> Result<(), IoTaskError>
72where
73 S: AsyncRead + AsyncWrite + Unpin,
74{
75 handle_connection_inner(stream, require_peer_auth, Some(observed_source), |to| {
76 (*to == keypair.public_key()).then(|| (Arc::new(keypair.clone()), inbox_sender.clone()))
77 })
78 .await
79}
80
81#[cfg(test)]
88pub(crate) async fn handle_connection_demux<S>(
89 stream: S,
90 registry: &crate::host_acceptor::HostAcceptorIdentityRegistry,
91) -> Result<(), IoTaskError>
92where
93 S: AsyncRead + AsyncWrite + Unpin,
94{
95 handle_connection_inner(stream, true, None, |to| registry.resolve(to)).await
96}
97
98pub(crate) async fn handle_tcp_connection_demux<S>(
101 stream: S,
102 observed_source: SocketAddr,
103 registry: &crate::host_acceptor::HostAcceptorIdentityRegistry,
104) -> Result<(), IoTaskError>
105where
106 S: AsyncRead + AsyncWrite + Unpin,
107{
108 handle_connection_inner(stream, true, Some(observed_source), |to| {
109 registry.resolve(to)
110 })
111 .await
112}
113
114async fn handle_connection_inner<S, R>(
118 stream: S,
119 require_peer_auth: bool,
120 observed_tcp_source: Option<SocketAddr>,
121 resolve: R,
122) -> Result<(), IoTaskError>
123where
124 S: AsyncRead + AsyncWrite + Unpin,
125 R: FnOnce(&PubKey) -> Option<ResolvedInboundIdentity>,
126{
127 let mut framed = Framed::new(
128 stream,
129 TransportCodec::new(crate::transport::MAX_PAYLOAD_SIZE),
130 );
131 let envelope = match framed.next().await {
132 Some(Ok(frame)) => frame.envelope,
133 Some(Err(err)) => return Err(IoTaskError::Io(err)),
134 None => {
135 return Err(IoTaskError::Io(std::io::Error::new(
136 std::io::ErrorKind::UnexpectedEof,
137 "connection closed",
138 )));
139 }
140 };
141
142 if require_peer_auth && !envelope.verify() {
149 return Err(IoTaskError::InvalidSignature {
150 envelope_id: envelope.id,
151 });
152 }
153
154 let Some((keypair, inbox_sender)) = resolve(&envelope.to) else {
158 return Err(IoTaskError::Misaddressed {
159 envelope_id: envelope.id,
160 });
161 };
162
163 let admission = match observed_tcp_source {
167 Some(source) => {
168 inbox_sender.send_tcp_connection_ingress(envelope.clone(), require_peer_auth, source)
169 }
170 None => inbox_sender.send_connection_ingress(envelope.clone(), require_peer_auth),
171 };
172 match admission {
173 AdmissionOutcome::Admitted => {
174 if should_ack(&envelope.kind) {
175 let ack = create_ack(&envelope, &keypair);
176 let frame = EnvelopeFrame {
177 envelope: ack,
178 raw: Arc::new(Bytes::new()),
179 };
180 framed.send(frame).await?;
181 }
182 Ok(())
183 }
184 AdmissionOutcome::Dropped { reason } => Err(match reason {
185 DropReason::SessionClosed => IoTaskError::InboxClosed,
186 DropReason::InboxFull => IoTaskError::InboxFull,
187 DropReason::UntrustedSender | DropReason::ClassificationRejected => {
188 IoTaskError::IngressDropped(reason)
189 }
190 }),
191 }
192}
193
194fn should_ack(kind: &MessageKind) -> bool {
202 matches!(
203 kind,
204 MessageKind::Message { .. }
205 | MessageKind::IncarnationFencedMessage { .. }
206 | MessageKind::Request { .. }
207 )
208}
209
210fn create_ack(original: &Envelope, keypair: &Keypair) -> Envelope {
212 let mut ack = Envelope {
213 id: uuid::Uuid::new_v4(),
214 from: keypair.public_key(),
215 to: original.from,
216 kind: MessageKind::Ack {
217 in_reply_to: original.id,
218 },
219 sig: Signature::new([0u8; 64]),
220 };
221 ack.sign(keypair);
222 ack
223}
224
225#[derive(Debug, thiserror::Error)]
227pub enum IoTaskError {
228 #[error("IO error: {0}")]
229 Io(#[from] std::io::Error),
230 #[error("Transport error: {0}")]
231 Transport(#[from] TransportError),
232 #[error("CBOR error: {0}")]
233 Cbor(String),
234 #[error("Inbox closed")]
235 InboxClosed,
236 #[error("Inbox full")]
237 InboxFull,
238 #[error("Ingress dropped: {0:?}")]
239 IngressDropped(DropReason),
240 #[error("Rejected envelope {envelope_id}: invalid signature")]
241 InvalidSignature { envelope_id: uuid::Uuid },
242 #[error("Rejected envelope {envelope_id}: misaddressed (not addressed to us)")]
243 Misaddressed { envelope_id: uuid::Uuid },
244}
245
246impl IoTaskError {
247 pub fn is_admission_rejection(&self) -> bool {
251 matches!(
252 self,
253 Self::InvalidSignature { .. } | Self::Misaddressed { .. } | Self::IngressDropped(_)
254 )
255 }
256}
257
258#[cfg(test)]
259#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
260mod tests {
261 use super::*;
262 use crate::classify::test_support;
263 use crate::identity::PubKey;
264 use crate::inbox::Inbox;
265 use crate::trust::{TrustEntry, TrustStore};
266 use crate::types::InboxItem;
267 use futures::StreamExt;
268 use parking_lot::RwLock;
269 use tokio::io::{AsyncReadExt, AsyncWriteExt};
270 use tokio_util::codec::FramedRead;
271 use uuid::Uuid;
272
273 fn make_keypair() -> Keypair {
274 Keypair::generate()
275 }
276
277 fn test_trust_entry(name: &str, pubkey: &PubKey, addr: &str) -> TrustEntry {
278 TrustEntry {
279 peer_id: pubkey.to_peer_id(),
280 name: meerkat_core::comms::PeerName::new(name).expect("valid peer name"),
281 pubkey: *pubkey,
282 address: meerkat_core::comms::PeerAddress::parse(addr).expect("valid peer address"),
283 meta: crate::PeerMeta::default(),
284 }
285 }
286
287 fn make_trusted_peers(pubkey: &PubKey) -> Arc<RwLock<TrustStore>> {
288 let mut store = TrustStore::new();
289 store
290 .insert(test_trust_entry(
291 "test-peer",
292 pubkey,
293 "tcp://127.0.0.1:4200",
294 ))
295 .expect("trusted test peer should insert");
296 Arc::new(RwLock::new(store))
297 }
298
299 fn classified_inbox_for_trust(
300 trusted: &Arc<RwLock<TrustStore>>,
301 require_peer_auth: bool,
302 ) -> (Inbox, InboxSender) {
303 Inbox::new_classified(test_support::classification_context_shared(
304 trusted.clone(),
305 require_peer_auth,
306 ))
307 }
308
309 fn make_signed_envelope(from_keypair: &Keypair, to: PubKey, kind: MessageKind) -> Envelope {
310 let mut envelope = Envelope {
311 id: Uuid::new_v4(),
312 from: from_keypair.public_key(),
313 to,
314 kind,
315 sig: Signature::new([0u8; 64]),
316 };
317 envelope.sign(from_keypair);
318 envelope
319 }
320
321 async fn envelope_to_bytes(envelope: &Envelope) -> Vec<u8> {
322 let mut payload = Vec::new();
323 ciborium::into_writer(envelope, &mut payload).unwrap();
324 let len = payload.len() as u32;
325 let mut bytes = Vec::new();
326 bytes.extend_from_slice(&len.to_be_bytes());
327 bytes.extend_from_slice(&payload);
328 bytes
329 }
330
331 async fn read_one_envelope<R>(reader: &mut R) -> Result<Envelope, std::io::Error>
332 where
333 R: tokio::io::AsyncRead + Unpin,
334 {
335 let mut framed = FramedRead::new(
336 reader,
337 TransportCodec::new(crate::transport::MAX_PAYLOAD_SIZE),
338 );
339 match framed.next().await {
340 Some(Ok(frame)) => Ok(frame.envelope),
341 Some(Err(err)) => Err(err),
342 None => Err(std::io::Error::new(
343 std::io::ErrorKind::UnexpectedEof,
344 "connection closed",
345 )),
346 }
347 }
348
349 #[test]
350 fn test_handle_connection_compiles() {
351 fn _check_signature<S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin>(
353 _stream: S,
354 _keypair: &Keypair,
355 _trusted: &Arc<RwLock<TrustStore>>,
356 _inbox_sender: &InboxSender,
357 ) {
358 }
360 }
361
362 #[tokio::test]
363 async fn test_io_task_reads_envelope() {
364 let sender_keypair = make_keypair();
365 let receiver_keypair = make_keypair();
366 let _trusted = make_trusted_peers(&sender_keypair.public_key());
367 let (_inbox, _inbox_sender) = Inbox::new();
368
369 let envelope = make_signed_envelope(
370 &sender_keypair,
371 receiver_keypair.public_key(),
372 MessageKind::Message {
373 objective_id: None,
374 content_taint: None,
375 blocks: None,
376 body: "hello".to_string(),
377 handling_mode: None,
378 },
379 );
380 let envelope_id = envelope.id;
381 let _bytes = envelope_to_bytes(&envelope).await;
382
383 let (client, server) = tokio::io::duplex(4096);
386 let (mut server_read, _server_write) = tokio::io::split(server);
387 let (_client_read, mut client_write) = tokio::io::split(client);
388
389 let bytes = envelope_to_bytes(&envelope).await;
391 tokio::spawn(async move {
392 client_write.write_all(&bytes).await.unwrap();
393 });
394
395 let received = read_one_envelope(&mut server_read).await.unwrap();
397 assert_eq!(received.id, envelope_id);
398 }
399
400 #[tokio::test]
401 async fn test_io_task_verifies_signature() {
402 let sender_keypair = make_keypair();
403 let receiver_keypair = make_keypair();
404 let trusted = make_trusted_peers(&sender_keypair.public_key());
405 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
406
407 let envelope = Envelope {
409 id: Uuid::new_v4(),
410 from: sender_keypair.public_key(),
411 to: receiver_keypair.public_key(),
412 kind: MessageKind::Message {
413 objective_id: None,
414 content_taint: None,
415 blocks: None,
416 body: "hello".to_string(),
417 handling_mode: None,
418 },
419 sig: Signature::new([0u8; 64]), };
421 let bytes = envelope_to_bytes(&envelope).await;
422
423 let (client, server) = tokio::io::duplex(4096);
424 let (_client_read, mut client_write) = tokio::io::split(client);
425
426 tokio::spawn(async move {
427 client_write.write_all(&bytes).await.unwrap();
428 });
429
430 let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
431 assert!(matches!(result, Err(IoTaskError::InvalidSignature { .. })));
435 assert!(
436 result
437 .as_ref()
438 .err()
439 .is_some_and(IoTaskError::is_admission_rejection),
440 "invalid signature must classify as an admission rejection"
441 );
442
443 let items = inbox.try_drain_classified();
445 assert!(items.is_empty());
446 }
447
448 #[tokio::test]
451 async fn test_io_task_misaddressed_returns_typed_fault() {
452 let sender_keypair = make_keypair();
453 let receiver_keypair = make_keypair();
454 let other_keypair = make_keypair();
455 let trusted = make_trusted_peers(&sender_keypair.public_key());
456 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
457
458 let envelope = make_signed_envelope(
460 &sender_keypair,
461 other_keypair.public_key(),
462 MessageKind::Message {
463 objective_id: None,
464 content_taint: None,
465 blocks: None,
466 body: "hello".to_string(),
467 handling_mode: None,
468 },
469 );
470 let bytes = envelope_to_bytes(&envelope).await;
471
472 let (client, server) = tokio::io::duplex(4096);
473 let (_client_read, mut client_write) = tokio::io::split(client);
474 tokio::spawn(async move {
475 client_write.write_all(&bytes).await.unwrap();
476 });
477
478 let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
479 assert!(matches!(result, Err(IoTaskError::Misaddressed { .. })));
480 assert!(
481 result
482 .as_ref()
483 .err()
484 .is_some_and(IoTaskError::is_admission_rejection)
485 );
486
487 let items = inbox.try_drain_classified();
488 assert!(items.is_empty());
489 }
490
491 #[tokio::test]
492 async fn test_io_task_checks_trust() {
493 let sender_keypair = make_keypair();
494 let receiver_keypair = make_keypair();
495 let untrusted_keypair = make_keypair();
496 let trusted = make_trusted_peers(&sender_keypair.public_key()); let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
498
499 let envelope = make_signed_envelope(
501 &untrusted_keypair, receiver_keypair.public_key(),
503 MessageKind::Message {
504 objective_id: None,
505 content_taint: None,
506 blocks: None,
507 body: "hello".to_string(),
508 handling_mode: None,
509 },
510 );
511 let bytes = envelope_to_bytes(&envelope).await;
512
513 let (client, server) = tokio::io::duplex(4096);
514 let (_client_read, mut client_write) = tokio::io::split(client);
515
516 tokio::spawn(async move {
517 client_write.write_all(&bytes).await.unwrap();
518 });
519
520 let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
521 assert!(matches!(
522 result,
523 Err(IoTaskError::IngressDropped(DropReason::UntrustedSender))
524 ));
525
526 let items = inbox.try_drain_classified();
528 assert!(items.is_empty());
529 }
530
531 #[tokio::test]
532 async fn test_io_task_accepts_invalid_signature_when_auth_disabled() {
533 let sender_keypair = make_keypair();
534 let receiver_keypair = make_keypair();
535 let trusted = make_trusted_peers(&make_keypair().public_key());
536 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, false);
537
538 let envelope = Envelope {
539 id: Uuid::new_v4(),
540 from: sender_keypair.public_key(),
541 to: receiver_keypair.public_key(),
542 kind: MessageKind::Message {
543 objective_id: None,
544 content_taint: None,
545 blocks: None,
546 body: "hello".to_string(),
547 handling_mode: None,
548 },
549 sig: Signature::new([0u8; 64]), };
551 let bytes = envelope_to_bytes(&envelope).await;
552 let expected_id = envelope.id;
553
554 let (client, server) = tokio::io::duplex(4096);
555 let (mut client_read, mut client_write) = tokio::io::split(client);
556
557 tokio::spawn(async move {
558 client_write.write_all(&bytes).await.unwrap();
559 });
560
561 handle_connection(server, false, &receiver_keypair, &inbox_sender)
562 .await
563 .unwrap();
564
565 let ack = read_one_envelope(&mut client_read).await.unwrap();
566 match ack.kind {
567 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, expected_id),
568 _ => panic!("expected Ack"),
569 }
570
571 let items = inbox.try_drain_classified();
572 assert_eq!(items.len(), 1);
573 match &items[0].item {
574 InboxItem::External { envelope } => assert_eq!(envelope.id, expected_id),
575 _ => panic!("expected External"),
576 }
577 }
578
579 #[tokio::test]
580 async fn test_io_task_accepts_untrusted_sender_when_auth_disabled() {
581 let sender_keypair = make_keypair();
582 let receiver_keypair = make_keypair();
583 let untrusted_keypair = make_keypair();
584 let trusted = make_trusted_peers(&untrusted_keypair.public_key()); let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, false);
586
587 let envelope = make_signed_envelope(
588 &sender_keypair, receiver_keypair.public_key(),
590 MessageKind::Message {
591 objective_id: None,
592 content_taint: None,
593 blocks: None,
594 body: "hello".to_string(),
595 handling_mode: None,
596 },
597 );
598 let bytes = envelope_to_bytes(&envelope).await;
599 let expected_id = envelope.id;
600
601 let (client, server) = tokio::io::duplex(4096);
602 let (mut client_read, mut client_write) = tokio::io::split(client);
603
604 tokio::spawn(async move {
605 client_write.write_all(&bytes).await.unwrap();
606 });
607
608 handle_connection(server, false, &receiver_keypair, &inbox_sender)
609 .await
610 .unwrap();
611
612 let ack = read_one_envelope(&mut client_read).await.unwrap();
613 match ack.kind {
614 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, expected_id),
615 _ => panic!("expected Ack"),
616 }
617
618 let items = inbox.try_drain_classified();
619 assert_eq!(items.len(), 1);
620 match &items[0].item {
621 InboxItem::External { envelope } => assert_eq!(envelope.id, expected_id),
622 _ => panic!("expected External"),
623 }
624 }
625
626 #[tokio::test]
627 async fn test_io_task_sends_ack() {
628 let sender_keypair = make_keypair();
629 let receiver_keypair = make_keypair();
630 let trusted = make_trusted_peers(&sender_keypair.public_key());
631 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
632
633 let envelope = make_signed_envelope(
634 &sender_keypair,
635 receiver_keypair.public_key(),
636 MessageKind::Message {
637 objective_id: None,
638 content_taint: None,
639 blocks: None,
640 body: "hello".to_string(),
641 handling_mode: None,
642 },
643 );
644 let original_id = envelope.id;
645 let bytes = envelope_to_bytes(&envelope).await;
646
647 let (client, server) = tokio::io::duplex(4096);
648 let (mut client_read, mut client_write) = tokio::io::split(client);
649
650 tokio::spawn(async move {
652 client_write.write_all(&bytes).await.unwrap();
653 });
654
655 let handle = tokio::spawn(async move {
657 handle_connection(server, true, &receiver_keypair, &inbox_sender).await
658 });
659
660 let ack = read_one_envelope(&mut client_read).await.unwrap();
662 handle.await.unwrap().unwrap();
663
664 match ack.kind {
666 MessageKind::Ack { in_reply_to } => {
667 assert_eq!(in_reply_to, original_id);
668 }
669 _ => panic!("expected Ack"),
670 }
671 assert!(ack.verify());
672 }
673
674 #[tokio::test]
675 async fn test_io_task_enqueues_to_inbox() {
676 let sender_keypair = make_keypair();
677 let receiver_keypair = make_keypair();
678 let trusted = make_trusted_peers(&sender_keypair.public_key());
679 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
680
681 let envelope = make_signed_envelope(
682 &sender_keypair,
683 receiver_keypair.public_key(),
684 MessageKind::Message {
685 objective_id: None,
686 content_taint: None,
687 blocks: None,
688 body: "hello".to_string(),
689 handling_mode: None,
690 },
691 );
692 let envelope_id = envelope.id;
693 let bytes = envelope_to_bytes(&envelope).await;
694
695 let (client, server) = tokio::io::duplex(4096);
696 let (mut client_read, mut client_write) = tokio::io::split(client);
697
698 tokio::spawn(async move {
699 client_write.write_all(&bytes).await.unwrap();
700 let mut buf = vec![0u8; 1024];
702 let _ = client_read.read(&mut buf).await;
703 });
704
705 handle_connection(server, true, &receiver_keypair, &inbox_sender)
706 .await
707 .unwrap();
708
709 let items = inbox.try_drain_classified();
711 assert_eq!(items.len(), 1);
712 match &items[0].item {
713 InboxItem::External { envelope } => {
714 assert_eq!(envelope.id, envelope_id);
715 }
716 _ => panic!("expected External"),
717 }
718 }
719
720 #[tokio::test]
721 async fn test_ack_for_message() {
722 let sender_keypair = make_keypair();
723 let receiver_keypair = make_keypair();
724 let trusted = make_trusted_peers(&sender_keypair.public_key());
725 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
726
727 let envelope = make_signed_envelope(
728 &sender_keypair,
729 receiver_keypair.public_key(),
730 MessageKind::Message {
731 objective_id: None,
732 content_taint: None,
733 blocks: None,
734 body: "hello".to_string(),
735 handling_mode: None,
736 },
737 );
738 let original_id = envelope.id;
739 let bytes = envelope_to_bytes(&envelope).await;
740
741 let (client, server) = tokio::io::duplex(4096);
742 let (mut client_read, mut client_write) = tokio::io::split(client);
743
744 tokio::spawn(async move {
745 client_write.write_all(&bytes).await.unwrap();
746 });
747
748 let handle = tokio::spawn(async move {
749 handle_connection(server, true, &receiver_keypair, &inbox_sender).await
750 });
751
752 let ack = read_one_envelope(&mut client_read).await.unwrap();
754 handle.await.unwrap().unwrap();
755
756 match ack.kind {
757 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
758 _ => panic!("expected Ack for Message"),
759 }
760 }
761
762 #[tokio::test]
763 async fn test_ack_for_multimodal_message() {
764 let sender_keypair = make_keypair();
765 let receiver_keypair = make_keypair();
766 let trusted = make_trusted_peers(&sender_keypair.public_key());
767 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
768
769 let envelope = make_signed_envelope(
770 &sender_keypair,
771 receiver_keypair.public_key(),
772 MessageKind::Message {
773 objective_id: None,
774 content_taint: None,
775 blocks: Some(vec![meerkat_core::ContentBlock::Image {
776 media_type: "image/png".to_string(),
777 data: "abc".into(),
778 }]),
779 body: "hello".to_string(),
780 handling_mode: None,
781 },
782 );
783 let original_id = envelope.id;
784 let bytes = envelope_to_bytes(&envelope).await;
785
786 let (client, server) = tokio::io::duplex(4096);
787 let (mut client_read, mut client_write) = tokio::io::split(client);
788
789 tokio::spawn(async move {
790 client_write.write_all(&bytes).await.unwrap();
791 });
792
793 let handle = tokio::spawn(async move {
794 handle_connection(server, true, &receiver_keypair, &inbox_sender).await
795 });
796
797 let ack = read_one_envelope(&mut client_read).await.unwrap();
798 handle.await.unwrap().unwrap();
799
800 match ack.kind {
801 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
802 _ => panic!("expected Ack for multimodal Message"),
803 }
804 }
805
806 #[tokio::test]
807 async fn test_ack_for_request() {
808 let sender_keypair = make_keypair();
809 let receiver_keypair = make_keypair();
810 let trusted = make_trusted_peers(&sender_keypair.public_key());
811 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
812
813 let envelope = make_signed_envelope(
814 &sender_keypair,
815 receiver_keypair.public_key(),
816 MessageKind::Request {
817 objective_id: None,
818 content_taint: None,
819 intent: "test".to_string(),
820 params: serde_json::json!({}),
821 blocks: None,
822 reply_endpoint: None,
823 handling_mode: None,
824 },
825 );
826 let original_id = envelope.id;
827 let bytes = envelope_to_bytes(&envelope).await;
828
829 let (client, server) = tokio::io::duplex(4096);
830 let (mut client_read, mut client_write) = tokio::io::split(client);
831
832 tokio::spawn(async move {
833 client_write.write_all(&bytes).await.unwrap();
834 });
835
836 let handle = tokio::spawn(async move {
837 handle_connection(server, true, &receiver_keypair, &inbox_sender).await
838 });
839
840 let ack = read_one_envelope(&mut client_read).await.unwrap();
842 handle.await.unwrap().unwrap();
843
844 match ack.kind {
845 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
846 _ => panic!("expected Ack for Request"),
847 }
848 }
849
850 #[tokio::test]
851 async fn tcp_handler_threads_observed_source_into_reply_endpoint_classification() {
852 let sender_keypair = make_keypair();
853 let receiver_keypair = make_keypair();
854 let trusted = make_trusted_peers(&sender_keypair.public_key());
855 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
856 let declared = meerkat_core::comms::PeerAddress::parse("tcp://203.0.113.88:4311").unwrap();
857 let envelope = make_signed_envelope(
858 &sender_keypair,
859 receiver_keypair.public_key(),
860 MessageKind::Request {
861 content_taint: None,
862 intent: "probe".to_string(),
863 params: serde_json::json!({}),
864 blocks: None,
865 reply_endpoint: Some(declared),
866 handling_mode: None,
867 objective_id: None,
868 },
869 );
870 let bytes = envelope_to_bytes(&envelope).await;
871
872 let (client, server) = tokio::io::duplex(4096);
873 let (mut client_read, mut client_write) = tokio::io::split(client);
874 tokio::spawn(async move {
875 client_write.write_all(&bytes).await.unwrap();
876 });
877
878 let source = "192.0.2.77:51900".parse().unwrap();
879 let handle = tokio::spawn(async move {
880 handle_tcp_connection(server, source, true, &receiver_keypair, &inbox_sender).await
881 });
882 let _ack = read_one_envelope(&mut client_read).await.unwrap();
883 handle.await.unwrap().unwrap();
884
885 let entries = inbox.try_drain_classified();
886 assert_eq!(entries.len(), 1);
887 assert_eq!(
888 entries[0].ingress_fact.declared_reply_endpoint,
889 Some(meerkat_core::comms::PeerAddress::parse("tcp://192.0.2.77:4311").unwrap())
890 );
891 }
892
893 #[tokio::test]
894 async fn test_no_ack_for_ack() {
895 let sender_keypair = make_keypair();
896 let receiver_keypair = make_keypair();
897 let trusted = make_trusted_peers(&sender_keypair.public_key());
898 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
899
900 let envelope = make_signed_envelope(
901 &sender_keypair,
902 receiver_keypair.public_key(),
903 MessageKind::Ack {
904 in_reply_to: Uuid::new_v4(),
905 },
906 );
907 let bytes = envelope_to_bytes(&envelope).await;
908
909 let (client, server) = tokio::io::duplex(4096);
910 let (mut client_read, mut client_write) = tokio::io::split(client);
911
912 tokio::spawn(async move {
913 client_write.write_all(&bytes).await.unwrap();
914 });
915
916 handle_connection(server, true, &receiver_keypair, &inbox_sender)
917 .await
918 .unwrap();
919
920 let result = read_one_envelope(&mut client_read).await;
922 assert!(result.is_err(), "Should not receive ack for Ack message");
923 }
924
925 #[tokio::test]
926 async fn test_no_ack_for_response() {
927 let sender_keypair = make_keypair();
928 let receiver_keypair = make_keypair();
929 let trusted = make_trusted_peers(&sender_keypair.public_key());
930 let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
931
932 let envelope = make_signed_envelope(
933 &sender_keypair,
934 receiver_keypair.public_key(),
935 MessageKind::Response {
936 objective_id: None,
937 content_taint: None,
938 in_reply_to: Uuid::new_v4(),
939 status: crate::types::Status::Completed,
940 result: serde_json::json!({}),
941 blocks: None,
942 handling_mode: None,
943 },
944 );
945 let bytes = envelope_to_bytes(&envelope).await;
946
947 let (client, server) = tokio::io::duplex(4096);
948 let (mut client_read, mut client_write) = tokio::io::split(client);
949
950 tokio::spawn(async move {
951 client_write.write_all(&bytes).await.unwrap();
952 });
953
954 handle_connection(server, true, &receiver_keypair, &inbox_sender)
955 .await
956 .unwrap();
957
958 let result = read_one_envelope(&mut client_read).await;
960 assert!(
961 result.is_err(),
962 "Should not receive ack for Response message"
963 );
964 }
965
966 #[tokio::test]
967 async fn test_drop_invalid_signature() {
968 let sender_keypair = make_keypair();
969 let receiver_keypair = make_keypair();
970 let trusted = make_trusted_peers(&sender_keypair.public_key());
971 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
972
973 let envelope = Envelope {
975 id: Uuid::new_v4(),
976 from: sender_keypair.public_key(),
977 to: receiver_keypair.public_key(),
978 kind: MessageKind::Message {
979 objective_id: None,
980 content_taint: None,
981 blocks: None,
982 body: "hello".to_string(),
983 handling_mode: None,
984 },
985 sig: Signature::new([0u8; 64]), };
987 let bytes = envelope_to_bytes(&envelope).await;
988
989 let (client, server) = tokio::io::duplex(4096);
990 let (mut client_read, mut client_write) = tokio::io::split(client);
991
992 tokio::spawn(async move {
993 client_write.write_all(&bytes).await.unwrap();
994 });
995
996 let outcome = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
1000 assert!(
1001 matches!(outcome, Err(IoTaskError::InvalidSignature { .. })),
1002 "invalid signature must classify as a typed admission rejection, got {outcome:?}"
1003 );
1004
1005 let result = read_one_envelope(&mut client_read).await;
1008 assert!(result.is_err(), "Should not send ack for invalid signature");
1009
1010 let items = inbox.try_drain_classified();
1012 assert!(items.is_empty());
1013 }
1014
1015 #[tokio::test]
1016 async fn test_drop_untrusted_sender() {
1017 let sender_keypair = make_keypair();
1018 let receiver_keypair = make_keypair();
1019 let other_keypair = make_keypair();
1020 let trusted = make_trusted_peers(&other_keypair.public_key()); let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
1022
1023 let envelope = make_signed_envelope(
1024 &sender_keypair, receiver_keypair.public_key(),
1026 MessageKind::Message {
1027 objective_id: None,
1028 content_taint: None,
1029 blocks: None,
1030 body: "hello".to_string(),
1031 handling_mode: None,
1032 },
1033 );
1034 let bytes = envelope_to_bytes(&envelope).await;
1035
1036 let (client, server) = tokio::io::duplex(4096);
1037 let (mut client_read, mut client_write) = tokio::io::split(client);
1038
1039 tokio::spawn(async move {
1040 client_write.write_all(&bytes).await.unwrap();
1041 });
1042
1043 let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
1044 assert!(matches!(
1045 result,
1046 Err(IoTaskError::IngressDropped(DropReason::UntrustedSender))
1047 ));
1048
1049 let result = read_one_envelope(&mut client_read).await;
1051 assert!(result.is_err(), "Should not send ack for untrusted sender");
1052
1053 let items = inbox.try_drain_classified();
1055 assert!(items.is_empty());
1056 }
1057
1058 #[tokio::test]
1059 async fn test_ack_waits_for_final_admission_outcome() {
1060 let sender_keypair = make_keypair();
1063 let receiver_keypair = make_keypair();
1064 let trusted = make_trusted_peers(&sender_keypair.public_key());
1065 let (inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
1066 drop(inbox);
1067
1068 let envelope = make_signed_envelope(
1069 &sender_keypair,
1070 receiver_keypair.public_key(),
1071 MessageKind::Message {
1072 objective_id: None,
1073 content_taint: None,
1074 blocks: None,
1075 body: "hello".to_string(),
1076 handling_mode: None,
1077 },
1078 );
1079 let bytes = envelope_to_bytes(&envelope).await;
1080
1081 let (client, server) = tokio::io::duplex(4096);
1082 let (mut client_read, mut client_write) = tokio::io::split(client);
1083
1084 tokio::spawn(async move {
1085 client_write.write_all(&bytes).await.unwrap();
1086 });
1087
1088 let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
1089 assert!(matches!(result, Err(IoTaskError::InboxClosed)));
1090
1091 let read_result = read_one_envelope(&mut client_read).await;
1092 assert!(
1093 read_result.is_err(),
1094 "admission rejection must not leak an Ack before the final outcome"
1095 );
1096 }
1097
1098 #[tokio::test]
1104 async fn test_io_task_reads_live_trust_after_listener_spawn() {
1105 let sender_keypair = make_keypair();
1106 let receiver_keypair = make_keypair();
1107
1108 let trusted = Arc::new(RwLock::new(TrustStore::new()));
1112 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
1113
1114 let envelope = make_signed_envelope(
1115 &sender_keypair,
1116 receiver_keypair.public_key(),
1117 MessageKind::Message {
1118 objective_id: None,
1119 content_taint: None,
1120 blocks: None,
1121 body: "hello".to_string(),
1122 handling_mode: None,
1123 },
1124 );
1125 let envelope_id = envelope.id;
1126 let bytes = envelope_to_bytes(&envelope).await;
1127
1128 let (client, server) = tokio::io::duplex(4096);
1129 let (mut client_read, mut client_write) = tokio::io::split(client);
1130
1131 trusted
1135 .write()
1136 .insert(test_trust_entry(
1137 "sender",
1138 &sender_keypair.public_key(),
1139 "tcp://127.0.0.1:0",
1140 ))
1141 .expect("live trust insert should succeed");
1142
1143 tokio::spawn(async move {
1144 client_write.write_all(&bytes).await.unwrap();
1145 let mut buf = vec![0u8; 1024];
1147 let _ = client_read.read(&mut buf).await;
1148 });
1149
1150 handle_connection(server, true, &receiver_keypair, &inbox_sender)
1151 .await
1152 .unwrap();
1153
1154 let items = inbox.try_drain_classified();
1155 assert_eq!(items.len(), 1, "envelope should be admitted via live trust");
1156 match &items[0].item {
1157 InboxItem::External { envelope } => assert_eq!(envelope.id, envelope_id),
1158 _ => panic!("expected External"),
1159 }
1160 }
1161
1162 use crate::host_acceptor::HostAcceptorIdentityRegistry;
1165
1166 fn demux_registry(
1167 identities: &[(&Keypair, &InboxSender)],
1168 ) -> (
1169 HostAcceptorIdentityRegistry,
1170 Arc<dyn std::any::Any + Send + Sync>,
1171 ) {
1172 let registry = HostAcceptorIdentityRegistry::new();
1173 let owner: Arc<dyn std::any::Any + Send + Sync> = Arc::new(());
1174 registry.install_owner(owner.clone()).unwrap();
1175 for (keypair, inbox_sender) in identities {
1176 registry
1177 .register_identity(
1178 &owner,
1179 keypair.public_key(),
1180 Arc::new((*keypair).clone()),
1181 (*inbox_sender).clone(),
1182 )
1183 .unwrap();
1184 }
1185 (registry, owner)
1186 }
1187
1188 #[tokio::test]
1192 async fn test_demux_routes_by_to_and_acks_with_addressed_keypair() {
1193 let sender_keypair = make_keypair();
1194 let member_a = make_keypair();
1195 let member_b = make_keypair();
1196 let trusted = make_trusted_peers(&sender_keypair.public_key());
1197 let (mut inbox_a, inbox_sender_a) = classified_inbox_for_trust(&trusted, true);
1198 let (mut inbox_b, inbox_sender_b) = classified_inbox_for_trust(&trusted, true);
1199 let (registry, _owner) =
1200 demux_registry(&[(&member_a, &inbox_sender_a), (&member_b, &inbox_sender_b)]);
1201
1202 let envelope = make_signed_envelope(
1203 &sender_keypair,
1204 member_b.public_key(),
1205 MessageKind::Message {
1206 content_taint: None,
1207 blocks: None,
1208 body: "hello".to_string(),
1209 handling_mode: None,
1210 objective_id: None,
1211 },
1212 );
1213 let original_id = envelope.id;
1214 let bytes = envelope_to_bytes(&envelope).await;
1215
1216 let (client, server) = tokio::io::duplex(4096);
1217 let (mut client_read, mut client_write) = tokio::io::split(client);
1218 tokio::spawn(async move {
1219 client_write.write_all(&bytes).await.unwrap();
1220 });
1221
1222 let handle = tokio::spawn(async move { handle_connection_demux(server, ®istry).await });
1223
1224 let ack = read_one_envelope(&mut client_read).await.unwrap();
1225 handle.await.unwrap().unwrap();
1226
1227 match ack.kind {
1228 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
1229 _ => panic!("expected Ack"),
1230 }
1231 assert_eq!(
1232 ack.from,
1233 member_b.public_key(),
1234 "ack must be signed by the ADDRESSED member's keypair"
1235 );
1236 assert!(ack.verify(), "ack must carry a valid member-B signature");
1237
1238 let items_b = inbox_b.try_drain_classified();
1239 assert_eq!(items_b.len(), 1, "envelope should land in B's inbox");
1240 match &items_b[0].item {
1241 InboxItem::External { envelope } => assert_eq!(envelope.id, original_id),
1242 _ => panic!("expected External"),
1243 }
1244 assert!(
1245 inbox_a.try_drain_classified().is_empty(),
1246 "member A must not receive an envelope addressed to B"
1247 );
1248 }
1249
1250 #[tokio::test]
1251 async fn tcp_demux_threads_observed_source_to_addressed_member_classifier() {
1252 let sender_keypair = make_keypair();
1253 let member = make_keypair();
1254 let trusted = make_trusted_peers(&sender_keypair.public_key());
1255 let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
1256 let (registry, _owner) = demux_registry(&[(&member, &inbox_sender)]);
1257 let envelope = make_signed_envelope(
1258 &sender_keypair,
1259 member.public_key(),
1260 MessageKind::Request {
1261 content_taint: None,
1262 intent: "probe".to_string(),
1263 params: serde_json::json!({}),
1264 blocks: None,
1265 reply_endpoint: Some(
1266 meerkat_core::comms::PeerAddress::parse("tcp://198.51.100.90:4311").unwrap(),
1267 ),
1268 handling_mode: None,
1269 objective_id: None,
1270 },
1271 );
1272 let bytes = envelope_to_bytes(&envelope).await;
1273 let (client, server) = tokio::io::duplex(4096);
1274 let (mut client_read, mut client_write) = tokio::io::split(client);
1275 tokio::spawn(async move {
1276 client_write.write_all(&bytes).await.unwrap();
1277 });
1278
1279 let handle = tokio::spawn(async move {
1280 handle_tcp_connection_demux(server, "192.0.2.91:52000".parse().unwrap(), ®istry)
1281 .await
1282 });
1283 let _ack = read_one_envelope(&mut client_read).await.unwrap();
1284 handle.await.unwrap().unwrap();
1285
1286 let entries = inbox.try_drain_classified();
1287 assert_eq!(entries.len(), 1);
1288 assert_eq!(
1289 entries[0].ingress_fact.declared_reply_endpoint,
1290 Some(meerkat_core::comms::PeerAddress::parse("tcp://192.0.2.91:4311").unwrap())
1291 );
1292 }
1293
1294 #[tokio::test]
1297 async fn test_demux_misaddressed_unregistered_identity_rejected() {
1298 let sender_keypair = make_keypair();
1299 let member_a = make_keypair();
1300 let stranger = make_keypair();
1301 let trusted = make_trusted_peers(&sender_keypair.public_key());
1302 let (mut inbox_a, inbox_sender_a) = classified_inbox_for_trust(&trusted, true);
1303 let (registry, _owner) = demux_registry(&[(&member_a, &inbox_sender_a)]);
1304
1305 let envelope = make_signed_envelope(
1306 &sender_keypair,
1307 stranger.public_key(),
1308 MessageKind::Message {
1309 content_taint: None,
1310 blocks: None,
1311 body: "hello".to_string(),
1312 handling_mode: None,
1313 objective_id: None,
1314 },
1315 );
1316 let bytes = envelope_to_bytes(&envelope).await;
1317
1318 let (client, server) = tokio::io::duplex(4096);
1319 let (mut client_read, mut client_write) = tokio::io::split(client);
1320 tokio::spawn(async move {
1321 client_write.write_all(&bytes).await.unwrap();
1322 });
1323
1324 let result = handle_connection_demux(server, ®istry).await;
1325 assert!(matches!(result, Err(IoTaskError::Misaddressed { .. })));
1326
1327 let read_result = read_one_envelope(&mut client_read).await;
1328 assert!(read_result.is_err(), "no ack for a misaddressed envelope");
1329 assert!(inbox_a.try_drain_classified().is_empty());
1330 }
1331
1332 #[tokio::test]
1337 async fn test_demux_registered_then_removed_identity_rejected() {
1338 let sender_keypair = make_keypair();
1339 let member_a = make_keypair();
1340 let trusted = make_trusted_peers(&sender_keypair.public_key());
1341 let (mut inbox_a, inbox_sender_a) = classified_inbox_for_trust(&trusted, true);
1342 let registry = HostAcceptorIdentityRegistry::new();
1343 let owner: Arc<dyn std::any::Any + Send + Sync> = Arc::new(());
1344 registry.install_owner(owner.clone()).unwrap();
1345 registry
1346 .register_identity(
1347 &owner,
1348 member_a.public_key(),
1349 Arc::new(member_a.clone()),
1350 inbox_sender_a.clone(),
1351 )
1352 .unwrap();
1353
1354 let intruder: Arc<dyn std::any::Any + Send + Sync> = Arc::new(());
1356 let wrong = registry.remove_identity(&intruder, &member_a.public_key());
1357 assert!(matches!(
1358 wrong,
1359 Err(crate::host_acceptor::HostAcceptorError::OwnerMismatch)
1360 ));
1361 assert!(
1362 registry.resolve(&member_a.public_key()).is_some(),
1363 "wrong-owner removal must not mutate the registry"
1364 );
1365
1366 assert!(
1368 registry
1369 .remove_identity(&owner, &member_a.public_key())
1370 .unwrap()
1371 );
1372
1373 let envelope = make_signed_envelope(
1374 &sender_keypair,
1375 member_a.public_key(),
1376 MessageKind::Message {
1377 content_taint: None,
1378 blocks: None,
1379 body: "hello".to_string(),
1380 handling_mode: None,
1381 objective_id: None,
1382 },
1383 );
1384 let bytes = envelope_to_bytes(&envelope).await;
1385
1386 let (client, server) = tokio::io::duplex(4096);
1387 let (_client_read, mut client_write) = tokio::io::split(client);
1388 tokio::spawn(async move {
1389 client_write.write_all(&bytes).await.unwrap();
1390 });
1391
1392 let result = handle_connection_demux(server, ®istry).await;
1393 assert!(matches!(result, Err(IoTaskError::Misaddressed { .. })));
1394 assert!(inbox_a.try_drain_classified().is_empty());
1395 }
1396
1397 #[tokio::test]
1403 async fn test_demux_rejects_unsigned_envelope_despite_no_auth_context() {
1404 let sender_keypair = make_keypair();
1405 let member_a = make_keypair();
1406 let trusted = make_trusted_peers(&sender_keypair.public_key());
1407 let (mut inbox_a, inbox_sender_a) = classified_inbox_for_trust(&trusted, false);
1410 let (registry, _owner) = demux_registry(&[(&member_a, &inbox_sender_a)]);
1411
1412 let envelope = Envelope {
1413 id: Uuid::new_v4(),
1414 from: sender_keypair.public_key(),
1415 to: member_a.public_key(),
1416 kind: MessageKind::Message {
1417 content_taint: None,
1418 blocks: None,
1419 body: "hello".to_string(),
1420 handling_mode: None,
1421 objective_id: None,
1422 },
1423 sig: Signature::new([0u8; 64]), };
1425 let bytes = envelope_to_bytes(&envelope).await;
1426
1427 let (client, server) = tokio::io::duplex(4096);
1428 let (mut client_read, mut client_write) = tokio::io::split(client);
1429 tokio::spawn(async move {
1430 client_write.write_all(&bytes).await.unwrap();
1431 });
1432
1433 let result = handle_connection_demux(server, ®istry).await;
1434 assert!(matches!(result, Err(IoTaskError::InvalidSignature { .. })));
1435
1436 let read_result = read_one_envelope(&mut client_read).await;
1437 assert!(read_result.is_err(), "no ack for an unsigned envelope");
1438 assert!(inbox_a.try_drain_classified().is_empty());
1439 }
1440
1441 #[tokio::test]
1448 async fn test_demux_accepts_hand_encoded_envelope_bytes() {
1449 let sender_keypair = make_keypair();
1450 let member_a = make_keypair();
1451 let trusted = make_trusted_peers(&sender_keypair.public_key());
1452 let (mut inbox_a, inbox_sender_a) = classified_inbox_for_trust(&trusted, true);
1453 let (registry, _owner) = demux_registry(&[(&member_a, &inbox_sender_a)]);
1454
1455 let envelope = make_signed_envelope(
1456 &sender_keypair,
1457 member_a.public_key(),
1458 MessageKind::Message {
1459 content_taint: None,
1460 blocks: None,
1461 body: "hello".to_string(),
1462 handling_mode: None,
1463 objective_id: None,
1464 },
1465 );
1466 let original_id = envelope.id;
1467
1468 let mut payload = Vec::new();
1471 ciborium::into_writer(&envelope, &mut payload).unwrap();
1472 let mut bytes = Vec::new();
1473 bytes.extend_from_slice(&(payload.len() as u32).to_be_bytes());
1474 bytes.extend_from_slice(&payload);
1475
1476 let (client, server) = tokio::io::duplex(4096);
1477 let (mut client_read, mut client_write) = tokio::io::split(client);
1478 tokio::spawn(async move {
1479 client_write.write_all(&bytes).await.unwrap();
1480 });
1481
1482 let handle = tokio::spawn(async move { handle_connection_demux(server, ®istry).await });
1483
1484 let ack = read_one_envelope(&mut client_read).await.unwrap();
1485 handle.await.unwrap().unwrap();
1486 match ack.kind {
1487 MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
1488 _ => panic!("expected Ack"),
1489 }
1490 assert_eq!(ack.from, member_a.public_key());
1491
1492 let items = inbox_a.try_drain_classified();
1493 assert_eq!(items.len(), 1);
1494 match &items[0].item {
1495 InboxItem::External { envelope } => assert_eq!(envelope.id, original_id),
1496 _ => panic!("expected External"),
1497 }
1498 }
1499}