Skip to main content

meerkat_comms/
io_task.rs

1//! IO task for handling incoming connections.
2//!
3//! Each incoming connection spawns an IO task that:
4//! 1. Reads envelope (with length-prefix framing)
5//! 2. Verifies signature (optional)
6//! 3. Runs ingress admission through the inbox seam
7//! 4. If admitted: sends Ack (unless it's an Ack or Response)
8//! 5. Closes connection (or keeps alive)
9
10use std::sync::Arc;
11
12use bytes::Bytes;
13use futures::{SinkExt, StreamExt};
14use tokio::io::{AsyncRead, AsyncWrite};
15use tokio_util::codec::Framed;
16
17use crate::identity::{Keypair, Signature};
18use crate::inbox::{AdmissionOutcome, DropReason, InboxSender};
19use crate::transport::TransportError;
20use crate::transport::codec::{EnvelopeFrame, TransportCodec};
21use crate::types::{Envelope, MessageKind};
22
23/// Handle an incoming connection.
24///
25/// Reads envelopes, validates them, admits them through the inbox seam,
26/// then sends acks for admitted ingress as appropriate.
27///
28/// # Arguments
29/// * `stream` - The async read/write stream (e.g., TcpStream or UnixStream)
30/// * `keypair` - Our keypair for signing acks
31/// * `require_peer_auth` - Whether to enforce signature+trusted-peer validation
32/// * `inbox_sender` - Channel to send validated messages to the inbox
33pub async fn handle_connection<S>(
34    stream: S,
35    require_peer_auth: bool,
36    keypair: &Keypair,
37    inbox_sender: &InboxSender,
38) -> Result<(), IoTaskError>
39where
40    S: AsyncRead + AsyncWrite + Unpin,
41{
42    let mut framed = Framed::new(
43        stream,
44        TransportCodec::new(crate::transport::MAX_PAYLOAD_SIZE),
45    );
46    let envelope = match framed.next().await {
47        Some(Ok(frame)) => frame.envelope,
48        Some(Err(err)) => return Err(IoTaskError::Io(err)),
49        None => {
50            return Err(IoTaskError::Io(std::io::Error::new(
51                std::io::ErrorKind::UnexpectedEof,
52                "connection closed",
53            )));
54        }
55    };
56
57    // Verify signature (when peer auth is enabled).
58    //
59    // A rejected envelope is a *failed* admission, not successful handling.
60    // Return a typed auth/address fault (mirroring the `IngressDropped` arm
61    // below) so the listener's `Err -> warn` arm records the rejection fact
62    // rather than treating the silent `Ok(())` as a clean connection.
63    if require_peer_auth && !envelope.verify() {
64        return Err(IoTaskError::InvalidSignature {
65            envelope_id: envelope.id,
66        });
67    }
68
69    // Verify envelope is addressed to us
70    if envelope.to != keypair.public_key() {
71        return Err(IoTaskError::Misaddressed {
72            envelope_id: envelope.id,
73        });
74    }
75
76    // Admit through the inbox seam first. Typed admission outcome: explicit
77    // drops are surfaced as `IoTaskError` so the IO task can react (close
78    // connection, log, etc.) rather than silently returning `Ok(())`.
79    match inbox_sender.send_connection_ingress(envelope.clone(), require_peer_auth) {
80        AdmissionOutcome::Admitted => {
81            if should_ack(&envelope.kind) {
82                let ack = create_ack(&envelope, keypair);
83                let frame = EnvelopeFrame {
84                    envelope: ack,
85                    raw: Arc::new(Bytes::new()),
86                };
87                framed.send(frame).await?;
88            }
89            Ok(())
90        }
91        AdmissionOutcome::Dropped { reason } => Err(match reason {
92            DropReason::SessionClosed => IoTaskError::InboxClosed,
93            DropReason::InboxFull => IoTaskError::InboxFull,
94            DropReason::UntrustedSender | DropReason::ClassificationRejected => {
95                IoTaskError::IngressDropped(reason)
96            }
97        }),
98    }
99}
100
101/// Determine if we should send an ack for this message kind.
102///
103/// Per spec:
104/// - Message: Yes
105/// - Request: Yes
106/// - Response: No
107/// - Ack: Never (would cause infinite loop)
108fn should_ack(kind: &MessageKind) -> bool {
109    matches!(
110        kind,
111        MessageKind::Message { .. } | MessageKind::Request { .. }
112    )
113}
114
115/// Create an Ack envelope in reply to the given envelope.
116fn create_ack(original: &Envelope, keypair: &Keypair) -> Envelope {
117    let mut ack = Envelope {
118        id: uuid::Uuid::new_v4(),
119        from: keypair.public_key(),
120        to: original.from,
121        kind: MessageKind::Ack {
122            in_reply_to: original.id,
123        },
124        sig: Signature::new([0u8; 64]),
125    };
126    ack.sign(keypair);
127    ack
128}
129
130/// Errors that can occur in IO task operations.
131#[derive(Debug, thiserror::Error)]
132pub enum IoTaskError {
133    #[error("IO error: {0}")]
134    Io(#[from] std::io::Error),
135    #[error("Transport error: {0}")]
136    Transport(#[from] TransportError),
137    #[error("CBOR error: {0}")]
138    Cbor(String),
139    #[error("Inbox closed")]
140    InboxClosed,
141    #[error("Inbox full")]
142    InboxFull,
143    #[error("Ingress dropped: {0:?}")]
144    IngressDropped(DropReason),
145    #[error("Rejected envelope {envelope_id}: invalid signature")]
146    InvalidSignature { envelope_id: uuid::Uuid },
147    #[error("Rejected envelope {envelope_id}: misaddressed (not addressed to us)")]
148    Misaddressed { envelope_id: uuid::Uuid },
149}
150
151impl IoTaskError {
152    /// Whether this error represents a *rejected admission* (auth/address/policy
153    /// fault) rather than a transport/IO failure. Lets the listener distinguish
154    /// "we refused this peer" from "the connection broke" for metrics.
155    pub fn is_admission_rejection(&self) -> bool {
156        matches!(
157            self,
158            Self::InvalidSignature { .. } | Self::Misaddressed { .. } | Self::IngressDropped(_)
159        )
160    }
161}
162
163#[cfg(test)]
164#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
165mod tests {
166    use super::*;
167    use crate::classify::test_support;
168    use crate::identity::PubKey;
169    use crate::inbox::Inbox;
170    use crate::trust::{TrustEntry, TrustStore};
171    use crate::types::InboxItem;
172    use futures::StreamExt;
173    use parking_lot::RwLock;
174    use tokio::io::{AsyncReadExt, AsyncWriteExt};
175    use tokio_util::codec::FramedRead;
176    use uuid::Uuid;
177
178    fn make_keypair() -> Keypair {
179        Keypair::generate()
180    }
181
182    fn test_trust_entry(name: &str, pubkey: &PubKey, addr: &str) -> TrustEntry {
183        TrustEntry {
184            peer_id: pubkey.to_peer_id(),
185            name: meerkat_core::comms::PeerName::new(name).expect("valid peer name"),
186            pubkey: *pubkey,
187            address: meerkat_core::comms::PeerAddress::parse(addr).expect("valid peer address"),
188            meta: crate::PeerMeta::default(),
189        }
190    }
191
192    fn make_trusted_peers(pubkey: &PubKey) -> Arc<RwLock<TrustStore>> {
193        let mut store = TrustStore::new();
194        store
195            .insert(test_trust_entry(
196                "test-peer",
197                pubkey,
198                "tcp://127.0.0.1:4200",
199            ))
200            .expect("trusted test peer should insert");
201        Arc::new(RwLock::new(store))
202    }
203
204    fn classified_inbox_for_trust(
205        trusted: &Arc<RwLock<TrustStore>>,
206        require_peer_auth: bool,
207    ) -> (Inbox, InboxSender) {
208        Inbox::new_classified(test_support::classification_context_shared(
209            trusted.clone(),
210            require_peer_auth,
211        ))
212    }
213
214    fn make_signed_envelope(from_keypair: &Keypair, to: PubKey, kind: MessageKind) -> Envelope {
215        let mut envelope = Envelope {
216            id: Uuid::new_v4(),
217            from: from_keypair.public_key(),
218            to,
219            kind,
220            sig: Signature::new([0u8; 64]),
221        };
222        envelope.sign(from_keypair);
223        envelope
224    }
225
226    async fn envelope_to_bytes(envelope: &Envelope) -> Vec<u8> {
227        let mut payload = Vec::new();
228        ciborium::into_writer(envelope, &mut payload).unwrap();
229        let len = payload.len() as u32;
230        let mut bytes = Vec::new();
231        bytes.extend_from_slice(&len.to_be_bytes());
232        bytes.extend_from_slice(&payload);
233        bytes
234    }
235
236    async fn read_one_envelope<R>(reader: &mut R) -> Result<Envelope, std::io::Error>
237    where
238        R: tokio::io::AsyncRead + Unpin,
239    {
240        let mut framed = FramedRead::new(
241            reader,
242            TransportCodec::new(crate::transport::MAX_PAYLOAD_SIZE),
243        );
244        match framed.next().await {
245            Some(Ok(frame)) => Ok(frame.envelope),
246            Some(Err(err)) => Err(err),
247            None => Err(std::io::Error::new(
248                std::io::ErrorKind::UnexpectedEof,
249                "connection closed",
250            )),
251        }
252    }
253
254    #[test]
255    fn test_handle_connection_compiles() {
256        // This test just verifies the function signature compiles
257        fn _check_signature<S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin>(
258            _stream: S,
259            _keypair: &Keypair,
260            _trusted: &Arc<RwLock<TrustStore>>,
261            _inbox_sender: &InboxSender,
262        ) {
263            // The handle_connection function exists with correct signature
264        }
265    }
266
267    #[tokio::test]
268    async fn test_io_task_reads_envelope() {
269        let sender_keypair = make_keypair();
270        let receiver_keypair = make_keypair();
271        let _trusted = make_trusted_peers(&sender_keypair.public_key());
272        let (_inbox, _inbox_sender) = Inbox::new();
273
274        let envelope = make_signed_envelope(
275            &sender_keypair,
276            receiver_keypair.public_key(),
277            MessageKind::Message {
278                content_taint: None,
279                blocks: None,
280                body: "hello".to_string(),
281                handling_mode: None,
282            },
283        );
284        let envelope_id = envelope.id;
285        let _bytes = envelope_to_bytes(&envelope).await;
286
287        // We need to handle that Cursor doesn't really support async write back
288        // For this test, we'll use a duplex stream instead
289        let (client, server) = tokio::io::duplex(4096);
290        let (mut server_read, _server_write) = tokio::io::split(server);
291        let (_client_read, mut client_write) = tokio::io::split(client);
292
293        // Write envelope from client
294        let bytes = envelope_to_bytes(&envelope).await;
295        tokio::spawn(async move {
296            client_write.write_all(&bytes).await.unwrap();
297        });
298
299        // Read just the envelope (not the full handle_connection)
300        let received = read_one_envelope(&mut server_read).await.unwrap();
301        assert_eq!(received.id, envelope_id);
302    }
303
304    #[tokio::test]
305    async fn test_io_task_verifies_signature() {
306        let sender_keypair = make_keypair();
307        let receiver_keypair = make_keypair();
308        let trusted = make_trusted_peers(&sender_keypair.public_key());
309        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
310
311        // Create envelope with invalid signature
312        let envelope = Envelope {
313            id: Uuid::new_v4(),
314            from: sender_keypair.public_key(),
315            to: receiver_keypair.public_key(),
316            kind: MessageKind::Message {
317                content_taint: None,
318                blocks: None,
319                body: "hello".to_string(),
320                handling_mode: None,
321            },
322            sig: Signature::new([0u8; 64]), // Invalid signature
323        };
324        let bytes = envelope_to_bytes(&envelope).await;
325
326        let (client, server) = tokio::io::duplex(4096);
327        let (_client_read, mut client_write) = tokio::io::split(client);
328
329        tokio::spawn(async move {
330            client_write.write_all(&bytes).await.unwrap();
331        });
332
333        let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
334        // ROW #300: an invalid signature is a rejected admission, surfaced as a
335        // typed auth fault — not a silent `Ok(())` that the listener would treat
336        // as successful handling.
337        assert!(matches!(result, Err(IoTaskError::InvalidSignature { .. })));
338        assert!(
339            result
340                .as_ref()
341                .err()
342                .is_some_and(IoTaskError::is_admission_rejection),
343            "invalid signature must classify as an admission rejection"
344        );
345
346        // No item in inbox
347        let items = inbox.try_drain_classified();
348        assert!(items.is_empty());
349    }
350
351    /// ROW #300 gate: an envelope addressed to a different recipient is rejected
352    /// with a typed `Misaddressed` fault, not silently swallowed as `Ok(())`.
353    #[tokio::test]
354    async fn test_io_task_misaddressed_returns_typed_fault() {
355        let sender_keypair = make_keypair();
356        let receiver_keypair = make_keypair();
357        let other_keypair = make_keypair();
358        let trusted = make_trusted_peers(&sender_keypair.public_key());
359        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
360
361        // Signed correctly, but addressed to `other_keypair`, not the receiver.
362        let envelope = make_signed_envelope(
363            &sender_keypair,
364            other_keypair.public_key(),
365            MessageKind::Message {
366                content_taint: None,
367                blocks: None,
368                body: "hello".to_string(),
369                handling_mode: None,
370            },
371        );
372        let bytes = envelope_to_bytes(&envelope).await;
373
374        let (client, server) = tokio::io::duplex(4096);
375        let (_client_read, mut client_write) = tokio::io::split(client);
376        tokio::spawn(async move {
377            client_write.write_all(&bytes).await.unwrap();
378        });
379
380        let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
381        assert!(matches!(result, Err(IoTaskError::Misaddressed { .. })));
382        assert!(
383            result
384                .as_ref()
385                .err()
386                .is_some_and(IoTaskError::is_admission_rejection)
387        );
388
389        let items = inbox.try_drain_classified();
390        assert!(items.is_empty());
391    }
392
393    #[tokio::test]
394    async fn test_io_task_checks_trust() {
395        let sender_keypair = make_keypair();
396        let receiver_keypair = make_keypair();
397        let untrusted_keypair = make_keypair();
398        let trusted = make_trusted_peers(&sender_keypair.public_key()); // Only trust sender_keypair
399        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
400
401        // Create envelope from untrusted peer
402        let envelope = make_signed_envelope(
403            &untrusted_keypair, // Not in trusted list
404            receiver_keypair.public_key(),
405            MessageKind::Message {
406                content_taint: None,
407                blocks: None,
408                body: "hello".to_string(),
409                handling_mode: None,
410            },
411        );
412        let bytes = envelope_to_bytes(&envelope).await;
413
414        let (client, server) = tokio::io::duplex(4096);
415        let (_client_read, mut client_write) = tokio::io::split(client);
416
417        tokio::spawn(async move {
418            client_write.write_all(&bytes).await.unwrap();
419        });
420
421        let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
422        assert!(matches!(
423            result,
424            Err(IoTaskError::IngressDropped(DropReason::UntrustedSender))
425        ));
426
427        // No item in inbox
428        let items = inbox.try_drain_classified();
429        assert!(items.is_empty());
430    }
431
432    #[tokio::test]
433    async fn test_io_task_accepts_invalid_signature_when_auth_disabled() {
434        let sender_keypair = make_keypair();
435        let receiver_keypair = make_keypair();
436        let trusted = make_trusted_peers(&make_keypair().public_key());
437        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, false);
438
439        let envelope = Envelope {
440            id: Uuid::new_v4(),
441            from: sender_keypair.public_key(),
442            to: receiver_keypair.public_key(),
443            kind: MessageKind::Message {
444                content_taint: None,
445                blocks: None,
446                body: "hello".to_string(),
447                handling_mode: None,
448            },
449            sig: Signature::new([0u8; 64]), // Invalid signature
450        };
451        let bytes = envelope_to_bytes(&envelope).await;
452        let expected_id = envelope.id;
453
454        let (client, server) = tokio::io::duplex(4096);
455        let (mut client_read, mut client_write) = tokio::io::split(client);
456
457        tokio::spawn(async move {
458            client_write.write_all(&bytes).await.unwrap();
459        });
460
461        handle_connection(server, false, &receiver_keypair, &inbox_sender)
462            .await
463            .unwrap();
464
465        let ack = read_one_envelope(&mut client_read).await.unwrap();
466        match ack.kind {
467            MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, expected_id),
468            _ => panic!("expected Ack"),
469        }
470
471        let items = inbox.try_drain_classified();
472        assert_eq!(items.len(), 1);
473        match &items[0].item {
474            InboxItem::External { envelope } => assert_eq!(envelope.id, expected_id),
475            _ => panic!("expected External"),
476        }
477    }
478
479    #[tokio::test]
480    async fn test_io_task_accepts_untrusted_sender_when_auth_disabled() {
481        let sender_keypair = make_keypair();
482        let receiver_keypair = make_keypair();
483        let untrusted_keypair = make_keypair();
484        let trusted = make_trusted_peers(&untrusted_keypair.public_key()); // not relevant in no-auth mode
485        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, false);
486
487        let envelope = make_signed_envelope(
488            &sender_keypair, // not in trusted list
489            receiver_keypair.public_key(),
490            MessageKind::Message {
491                content_taint: None,
492                blocks: None,
493                body: "hello".to_string(),
494                handling_mode: None,
495            },
496        );
497        let bytes = envelope_to_bytes(&envelope).await;
498        let expected_id = envelope.id;
499
500        let (client, server) = tokio::io::duplex(4096);
501        let (mut client_read, mut client_write) = tokio::io::split(client);
502
503        tokio::spawn(async move {
504            client_write.write_all(&bytes).await.unwrap();
505        });
506
507        handle_connection(server, false, &receiver_keypair, &inbox_sender)
508            .await
509            .unwrap();
510
511        let ack = read_one_envelope(&mut client_read).await.unwrap();
512        match ack.kind {
513            MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, expected_id),
514            _ => panic!("expected Ack"),
515        }
516
517        let items = inbox.try_drain_classified();
518        assert_eq!(items.len(), 1);
519        match &items[0].item {
520            InboxItem::External { envelope } => assert_eq!(envelope.id, expected_id),
521            _ => panic!("expected External"),
522        }
523    }
524
525    #[tokio::test]
526    async fn test_io_task_sends_ack() {
527        let sender_keypair = make_keypair();
528        let receiver_keypair = make_keypair();
529        let trusted = make_trusted_peers(&sender_keypair.public_key());
530        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
531
532        let envelope = make_signed_envelope(
533            &sender_keypair,
534            receiver_keypair.public_key(),
535            MessageKind::Message {
536                content_taint: None,
537                blocks: None,
538                body: "hello".to_string(),
539                handling_mode: None,
540            },
541        );
542        let original_id = envelope.id;
543        let bytes = envelope_to_bytes(&envelope).await;
544
545        let (client, server) = tokio::io::duplex(4096);
546        let (mut client_read, mut client_write) = tokio::io::split(client);
547
548        // Send envelope
549        tokio::spawn(async move {
550            client_write.write_all(&bytes).await.unwrap();
551        });
552
553        // Handle connection
554        let handle = tokio::spawn(async move {
555            handle_connection(server, true, &receiver_keypair, &inbox_sender).await
556        });
557
558        // Read ack from client side
559        let ack = read_one_envelope(&mut client_read).await.unwrap();
560        handle.await.unwrap().unwrap();
561
562        // Verify ack
563        match ack.kind {
564            MessageKind::Ack { in_reply_to } => {
565                assert_eq!(in_reply_to, original_id);
566            }
567            _ => panic!("expected Ack"),
568        }
569        assert!(ack.verify());
570    }
571
572    #[tokio::test]
573    async fn test_io_task_enqueues_to_inbox() {
574        let sender_keypair = make_keypair();
575        let receiver_keypair = make_keypair();
576        let trusted = make_trusted_peers(&sender_keypair.public_key());
577        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
578
579        let envelope = make_signed_envelope(
580            &sender_keypair,
581            receiver_keypair.public_key(),
582            MessageKind::Message {
583                content_taint: None,
584                blocks: None,
585                body: "hello".to_string(),
586                handling_mode: None,
587            },
588        );
589        let envelope_id = envelope.id;
590        let bytes = envelope_to_bytes(&envelope).await;
591
592        let (client, server) = tokio::io::duplex(4096);
593        let (mut client_read, mut client_write) = tokio::io::split(client);
594
595        tokio::spawn(async move {
596            client_write.write_all(&bytes).await.unwrap();
597            // Read the ack to prevent blocking
598            let mut buf = vec![0u8; 1024];
599            let _ = client_read.read(&mut buf).await;
600        });
601
602        handle_connection(server, true, &receiver_keypair, &inbox_sender)
603            .await
604            .unwrap();
605
606        // Check inbox
607        let items = inbox.try_drain_classified();
608        assert_eq!(items.len(), 1);
609        match &items[0].item {
610            InboxItem::External { envelope } => {
611                assert_eq!(envelope.id, envelope_id);
612            }
613            _ => panic!("expected External"),
614        }
615    }
616
617    #[tokio::test]
618    async fn test_ack_for_message() {
619        let sender_keypair = make_keypair();
620        let receiver_keypair = make_keypair();
621        let trusted = make_trusted_peers(&sender_keypair.public_key());
622        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
623
624        let envelope = make_signed_envelope(
625            &sender_keypair,
626            receiver_keypair.public_key(),
627            MessageKind::Message {
628                content_taint: None,
629                blocks: None,
630                body: "hello".to_string(),
631                handling_mode: None,
632            },
633        );
634        let original_id = envelope.id;
635        let bytes = envelope_to_bytes(&envelope).await;
636
637        let (client, server) = tokio::io::duplex(4096);
638        let (mut client_read, mut client_write) = tokio::io::split(client);
639
640        tokio::spawn(async move {
641            client_write.write_all(&bytes).await.unwrap();
642        });
643
644        let handle = tokio::spawn(async move {
645            handle_connection(server, true, &receiver_keypair, &inbox_sender).await
646        });
647
648        // Should receive an ack
649        let ack = read_one_envelope(&mut client_read).await.unwrap();
650        handle.await.unwrap().unwrap();
651
652        match ack.kind {
653            MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
654            _ => panic!("expected Ack for Message"),
655        }
656    }
657
658    #[tokio::test]
659    async fn test_ack_for_multimodal_message() {
660        let sender_keypair = make_keypair();
661        let receiver_keypair = make_keypair();
662        let trusted = make_trusted_peers(&sender_keypair.public_key());
663        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
664
665        let envelope = make_signed_envelope(
666            &sender_keypair,
667            receiver_keypair.public_key(),
668            MessageKind::Message {
669                content_taint: None,
670                blocks: Some(vec![meerkat_core::ContentBlock::Image {
671                    media_type: "image/png".to_string(),
672                    data: "abc".into(),
673                }]),
674                body: "hello".to_string(),
675                handling_mode: None,
676            },
677        );
678        let original_id = envelope.id;
679        let bytes = envelope_to_bytes(&envelope).await;
680
681        let (client, server) = tokio::io::duplex(4096);
682        let (mut client_read, mut client_write) = tokio::io::split(client);
683
684        tokio::spawn(async move {
685            client_write.write_all(&bytes).await.unwrap();
686        });
687
688        let handle = tokio::spawn(async move {
689            handle_connection(server, true, &receiver_keypair, &inbox_sender).await
690        });
691
692        let ack = read_one_envelope(&mut client_read).await.unwrap();
693        handle.await.unwrap().unwrap();
694
695        match ack.kind {
696            MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
697            _ => panic!("expected Ack for multimodal Message"),
698        }
699    }
700
701    #[tokio::test]
702    async fn test_ack_for_request() {
703        let sender_keypair = make_keypair();
704        let receiver_keypair = make_keypair();
705        let trusted = make_trusted_peers(&sender_keypair.public_key());
706        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
707
708        let envelope = make_signed_envelope(
709            &sender_keypair,
710            receiver_keypair.public_key(),
711            MessageKind::Request {
712                content_taint: None,
713                intent: "test".to_string(),
714                params: serde_json::json!({}),
715                blocks: None,
716                handling_mode: None,
717            },
718        );
719        let original_id = envelope.id;
720        let bytes = envelope_to_bytes(&envelope).await;
721
722        let (client, server) = tokio::io::duplex(4096);
723        let (mut client_read, mut client_write) = tokio::io::split(client);
724
725        tokio::spawn(async move {
726            client_write.write_all(&bytes).await.unwrap();
727        });
728
729        let handle = tokio::spawn(async move {
730            handle_connection(server, true, &receiver_keypair, &inbox_sender).await
731        });
732
733        // Should receive an ack
734        let ack = read_one_envelope(&mut client_read).await.unwrap();
735        handle.await.unwrap().unwrap();
736
737        match ack.kind {
738            MessageKind::Ack { in_reply_to } => assert_eq!(in_reply_to, original_id),
739            _ => panic!("expected Ack for Request"),
740        }
741    }
742
743    #[tokio::test]
744    async fn test_no_ack_for_ack() {
745        let sender_keypair = make_keypair();
746        let receiver_keypair = make_keypair();
747        let trusted = make_trusted_peers(&sender_keypair.public_key());
748        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
749
750        let envelope = make_signed_envelope(
751            &sender_keypair,
752            receiver_keypair.public_key(),
753            MessageKind::Ack {
754                in_reply_to: Uuid::new_v4(),
755            },
756        );
757        let bytes = envelope_to_bytes(&envelope).await;
758
759        let (client, server) = tokio::io::duplex(4096);
760        let (mut client_read, mut client_write) = tokio::io::split(client);
761
762        tokio::spawn(async move {
763            client_write.write_all(&bytes).await.unwrap();
764        });
765
766        handle_connection(server, true, &receiver_keypair, &inbox_sender)
767            .await
768            .unwrap();
769
770        // Should NOT receive an ack - connection closes without data
771        let result = read_one_envelope(&mut client_read).await;
772        assert!(result.is_err(), "Should not receive ack for Ack message");
773    }
774
775    #[tokio::test]
776    async fn test_no_ack_for_response() {
777        let sender_keypair = make_keypair();
778        let receiver_keypair = make_keypair();
779        let trusted = make_trusted_peers(&sender_keypair.public_key());
780        let (_inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
781
782        let envelope = make_signed_envelope(
783            &sender_keypair,
784            receiver_keypair.public_key(),
785            MessageKind::Response {
786                content_taint: None,
787                in_reply_to: Uuid::new_v4(),
788                status: crate::types::Status::Completed,
789                result: serde_json::json!({}),
790                blocks: None,
791                handling_mode: None,
792            },
793        );
794        let bytes = envelope_to_bytes(&envelope).await;
795
796        let (client, server) = tokio::io::duplex(4096);
797        let (mut client_read, mut client_write) = tokio::io::split(client);
798
799        tokio::spawn(async move {
800            client_write.write_all(&bytes).await.unwrap();
801        });
802
803        handle_connection(server, true, &receiver_keypair, &inbox_sender)
804            .await
805            .unwrap();
806
807        // Should NOT receive an ack - connection closes without data
808        let result = read_one_envelope(&mut client_read).await;
809        assert!(
810            result.is_err(),
811            "Should not receive ack for Response message"
812        );
813    }
814
815    #[tokio::test]
816    async fn test_drop_invalid_signature() {
817        let sender_keypair = make_keypair();
818        let receiver_keypair = make_keypair();
819        let trusted = make_trusted_peers(&sender_keypair.public_key());
820        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
821
822        // Create envelope with invalid signature
823        let envelope = Envelope {
824            id: Uuid::new_v4(),
825            from: sender_keypair.public_key(),
826            to: receiver_keypair.public_key(),
827            kind: MessageKind::Message {
828                content_taint: None,
829                blocks: None,
830                body: "hello".to_string(),
831                handling_mode: None,
832            },
833            sig: Signature::new([0u8; 64]), // Invalid
834        };
835        let bytes = envelope_to_bytes(&envelope).await;
836
837        let (client, server) = tokio::io::duplex(4096);
838        let (mut client_read, mut client_write) = tokio::io::split(client);
839
840        tokio::spawn(async move {
841            client_write.write_all(&bytes).await.unwrap();
842        });
843
844        // ROW #300: an invalid signature is a rejected admission surfaced as a
845        // typed auth fault, not a silent `Ok(())` the listener would treat as a
846        // clean connection.
847        let outcome = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
848        assert!(
849            matches!(outcome, Err(IoTaskError::InvalidSignature { .. })),
850            "invalid signature must classify as a typed admission rejection, got {outcome:?}"
851        );
852
853        // No ack sent - the connection is dropped on the rejection rather than
854        // acknowledged.
855        let result = read_one_envelope(&mut client_read).await;
856        assert!(result.is_err(), "Should not send ack for invalid signature");
857
858        // No inbox item
859        let items = inbox.try_drain_classified();
860        assert!(items.is_empty());
861    }
862
863    #[tokio::test]
864    async fn test_drop_untrusted_sender() {
865        let sender_keypair = make_keypair();
866        let receiver_keypair = make_keypair();
867        let other_keypair = make_keypair();
868        let trusted = make_trusted_peers(&other_keypair.public_key()); // sender NOT trusted
869        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
870
871        let envelope = make_signed_envelope(
872            &sender_keypair, // Not trusted
873            receiver_keypair.public_key(),
874            MessageKind::Message {
875                content_taint: None,
876                blocks: None,
877                body: "hello".to_string(),
878                handling_mode: None,
879            },
880        );
881        let bytes = envelope_to_bytes(&envelope).await;
882
883        let (client, server) = tokio::io::duplex(4096);
884        let (mut client_read, mut client_write) = tokio::io::split(client);
885
886        tokio::spawn(async move {
887            client_write.write_all(&bytes).await.unwrap();
888        });
889
890        let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
891        assert!(matches!(
892            result,
893            Err(IoTaskError::IngressDropped(DropReason::UntrustedSender))
894        ));
895
896        // No ack sent - connection closes without data
897        let result = read_one_envelope(&mut client_read).await;
898        assert!(result.is_err(), "Should not send ack for untrusted sender");
899
900        // No inbox item
901        let items = inbox.try_drain_classified();
902        assert!(items.is_empty());
903    }
904
905    #[tokio::test]
906    async fn test_ack_waits_for_final_admission_outcome() {
907        // DOGMA-12 defensive scan: if admission rejects the ingress item,
908        // the transport must not send an Ack first.
909        let sender_keypair = make_keypair();
910        let receiver_keypair = make_keypair();
911        let trusted = make_trusted_peers(&sender_keypair.public_key());
912        let (inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
913        drop(inbox);
914
915        let envelope = make_signed_envelope(
916            &sender_keypair,
917            receiver_keypair.public_key(),
918            MessageKind::Message {
919                content_taint: None,
920                blocks: None,
921                body: "hello".to_string(),
922                handling_mode: None,
923            },
924        );
925        let bytes = envelope_to_bytes(&envelope).await;
926
927        let (client, server) = tokio::io::duplex(4096);
928        let (mut client_read, mut client_write) = tokio::io::split(client);
929
930        tokio::spawn(async move {
931            client_write.write_all(&bytes).await.unwrap();
932        });
933
934        let result = handle_connection(server, true, &receiver_keypair, &inbox_sender).await;
935        assert!(matches!(result, Err(IoTaskError::InboxClosed)));
936
937        let read_result = read_one_envelope(&mut client_read).await;
938        assert!(
939            read_result.is_err(),
940            "admission rejection must not leak an Ack before the final outcome"
941        );
942    }
943
944    /// Regression: the IO task must read trust through the shared
945    /// `Arc<RwLock<TrustStore>>` handle — not a snapshot — so that a
946    /// peer added to the router *after* the connection is accepted is
947    /// still admitted. This locks in the Wave 3 D Row 20 invariant: one
948    /// trust authority, one read path, no snapshot divergence.
949    #[tokio::test]
950    async fn test_io_task_reads_live_trust_after_listener_spawn() {
951        let sender_keypair = make_keypair();
952        let receiver_keypair = make_keypair();
953
954        // Start with an EMPTY trust set. If the IO task snapshotted at
955        // spawn time, the subsequent add below would not be visible and
956        // the envelope would be silently dropped.
957        let trusted = Arc::new(RwLock::new(TrustStore::new()));
958        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
959
960        let envelope = make_signed_envelope(
961            &sender_keypair,
962            receiver_keypair.public_key(),
963            MessageKind::Message {
964                content_taint: None,
965                blocks: None,
966                body: "hello".to_string(),
967                handling_mode: None,
968            },
969        );
970        let envelope_id = envelope.id;
971        let bytes = envelope_to_bytes(&envelope).await;
972
973        let (client, server) = tokio::io::duplex(4096);
974        let (mut client_read, mut client_write) = tokio::io::split(client);
975
976        // Mutate the trust set AFTER the IO task would normally have
977        // snapshotted, but BEFORE the envelope is delivered. The live
978        // read must observe this mutation.
979        trusted
980            .write()
981            .insert(test_trust_entry(
982                "sender",
983                &sender_keypair.public_key(),
984                "tcp://127.0.0.1:0",
985            ))
986            .expect("live trust insert should succeed");
987
988        tokio::spawn(async move {
989            client_write.write_all(&bytes).await.unwrap();
990            // Read the ack to avoid blocking.
991            let mut buf = vec![0u8; 1024];
992            let _ = client_read.read(&mut buf).await;
993        });
994
995        handle_connection(server, true, &receiver_keypair, &inbox_sender)
996            .await
997            .unwrap();
998
999        let items = inbox.try_drain_classified();
1000        assert_eq!(items.len(), 1, "envelope should be admitted via live trust");
1001        match &items[0].item {
1002            InboxItem::External { envelope } => assert_eq!(envelope.id, envelope_id),
1003            _ => panic!("expected External"),
1004        }
1005    }
1006}