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::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
24/// A resolved inbound recipient identity: the addressed identity's signing
25/// keypair and its inbox.
26///
27/// The keypair is load-bearing for acks, not cosmetic: the sender's router
28/// accepts an ack only when `ack.from` equals the envelope's original `to`
29/// (see `Router::send_on_stream`), so an ack signed by any other local
30/// keypair surfaces sender-side as `PeerOffline`.
31type ResolvedInboundIdentity = (Arc<Keypair>, InboxSender);
32
33/// Handle an incoming connection.
34///
35/// Reads envelopes, validates them, admits them through the inbox seam,
36/// then sends acks for admitted ingress as appropriate.
37///
38/// # Arguments
39/// * `stream` - The async read/write stream (e.g., TcpStream or UnixStream)
40/// * `keypair` - Our keypair for signing acks
41/// * `require_peer_auth` - Whether to enforce signature+trusted-peer validation
42/// * `inbox_sender` - Channel to send validated messages to the inbox
43pub 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    // The degenerate one-identity resolver: the single runtime identity is
53    // addressable, everything else is misaddressed.
54    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
60/// TCP-specific single-identity handler carrying the kernel-observed remote
61/// socket address into ingress classification.
62///
63/// Keep [`handle_connection`] as the transport-generic wrapper for UDS and
64/// duplex-based tests. Production TCP listeners must use this entry point.
65pub(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/// Handle an incoming connection on the multi-identity host acceptor path,
82/// demultiplexing by `envelope.to` against the registered identity set.
83///
84/// Peer auth is deliberately NOT a parameter: the host path verifies
85/// envelope signatures by construction (D1) — no field or flag exists that
86/// could disable it.
87#[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
98/// TCP-specific multi-identity host-demux handler carrying the observed
99/// remote socket address into the addressed member's classifier.
100pub(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
114/// Shared per-connection body, parameterized by recipient-identity
115/// resolution. Order is identical to the pre-demux single-identity handler:
116/// frame read, signature verify, identity gate, inbox admission, ack.
117async 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    // Verify signature (when peer auth is enabled).
143    //
144    // A rejected envelope is a *failed* admission, not successful handling.
145    // Return a typed auth/address fault (mirroring the `IngressDropped` arm
146    // below) so the listener's `Err -> warn` arm records the rejection fact
147    // rather than treating the silent `Ok(())` as a clean connection.
148    if require_peer_auth && !envelope.verify() {
149        return Err(IoTaskError::InvalidSignature {
150            envelope_id: envelope.id,
151        });
152    }
153
154    // Verify the envelope is addressed to a resolvable local identity. One
155    // typed error covers both "never registered" and "registered then
156    // removed" recipients.
157    let Some((keypair, inbox_sender)) = resolve(&envelope.to) else {
158        return Err(IoTaskError::Misaddressed {
159            envelope_id: envelope.id,
160        });
161    };
162
163    // Admit through the inbox seam first. Typed admission outcome: explicit
164    // drops are surfaced as `IoTaskError` so the IO task can react (close
165    // connection, log, etc.) rather than silently returning `Ok(())`.
166    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
194/// Determine if we should send an ack for this message kind.
195///
196/// Per spec:
197/// - Message: Yes
198/// - Request: Yes
199/// - Response: No
200/// - Ack: Never (would cause infinite loop)
201fn should_ack(kind: &MessageKind) -> bool {
202    matches!(
203        kind,
204        MessageKind::Message { .. }
205            | MessageKind::IncarnationFencedMessage { .. }
206            | MessageKind::Request { .. }
207    )
208}
209
210/// Create an Ack envelope in reply to the given envelope.
211fn 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/// Errors that can occur in IO task operations.
226#[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    /// Whether this error represents a *rejected admission* (auth/address/policy
248    /// fault) rather than a transport/IO failure. Lets the listener distinguish
249    /// "we refused this peer" from "the connection broke" for metrics.
250    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        // This test just verifies the function signature compiles
352        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            // The handle_connection function exists with correct signature
359        }
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        // We need to handle that Cursor doesn't really support async write back
384        // For this test, we'll use a duplex stream instead
385        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        // Write envelope from client
390        let bytes = envelope_to_bytes(&envelope).await;
391        tokio::spawn(async move {
392            client_write.write_all(&bytes).await.unwrap();
393        });
394
395        // Read just the envelope (not the full handle_connection)
396        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        // Create envelope with invalid signature
408        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]), // Invalid signature
420        };
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        // ROW #300: an invalid signature is a rejected admission, surfaced as a
432        // typed auth fault — not a silent `Ok(())` that the listener would treat
433        // as successful handling.
434        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        // No item in inbox
444        let items = inbox.try_drain_classified();
445        assert!(items.is_empty());
446    }
447
448    /// ROW #300 gate: an envelope addressed to a different recipient is rejected
449    /// with a typed `Misaddressed` fault, not silently swallowed as `Ok(())`.
450    #[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        // Signed correctly, but addressed to `other_keypair`, not the receiver.
459        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()); // Only trust sender_keypair
497        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
498
499        // Create envelope from untrusted peer
500        let envelope = make_signed_envelope(
501            &untrusted_keypair, // Not in trusted list
502            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        // No item in inbox
527        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]), // Invalid signature
550        };
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()); // not relevant in no-auth mode
585        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, false);
586
587        let envelope = make_signed_envelope(
588            &sender_keypair, // not in trusted list
589            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        // Send envelope
651        tokio::spawn(async move {
652            client_write.write_all(&bytes).await.unwrap();
653        });
654
655        // Handle connection
656        let handle = tokio::spawn(async move {
657            handle_connection(server, true, &receiver_keypair, &inbox_sender).await
658        });
659
660        // Read ack from client side
661        let ack = read_one_envelope(&mut client_read).await.unwrap();
662        handle.await.unwrap().unwrap();
663
664        // Verify ack
665        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            // Read the ack to prevent blocking
701            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        // Check inbox
710        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        // Should receive an ack
753        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        // Should receive an ack
841        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        // Should NOT receive an ack - connection closes without data
921        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        // Should NOT receive an ack - connection closes without data
959        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        // Create envelope with invalid signature
974        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]), // Invalid
986        };
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        // ROW #300: an invalid signature is a rejected admission surfaced as a
997        // typed auth fault, not a silent `Ok(())` the listener would treat as a
998        // clean connection.
999        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        // No ack sent - the connection is dropped on the rejection rather than
1006        // acknowledged.
1007        let result = read_one_envelope(&mut client_read).await;
1008        assert!(result.is_err(), "Should not send ack for invalid signature");
1009
1010        // No inbox item
1011        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()); // sender NOT trusted
1021        let (mut inbox, inbox_sender) = classified_inbox_for_trust(&trusted, true);
1022
1023        let envelope = make_signed_envelope(
1024            &sender_keypair, // Not trusted
1025            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        // No ack sent - connection closes without data
1050        let result = read_one_envelope(&mut client_read).await;
1051        assert!(result.is_err(), "Should not send ack for untrusted sender");
1052
1053        // No inbox item
1054        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        // DOGMA-12 defensive scan: if admission rejects the ingress item,
1061        // the transport must not send an Ack first.
1062        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    /// Regression: the IO task must read trust through the shared
1099    /// `Arc<RwLock<TrustStore>>` handle — not a snapshot — so that a
1100    /// peer added to the router *after* the connection is accepted is
1101    /// still admitted. This locks in the Wave 3 D Row 20 invariant: one
1102    /// trust authority, one read path, no snapshot divergence.
1103    #[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        // Start with an EMPTY trust set. If the IO task snapshotted at
1109        // spawn time, the subsequent add below would not be visible and
1110        // the envelope would be silently dropped.
1111        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        // Mutate the trust set AFTER the IO task would normally have
1132        // snapshotted, but BEFORE the envelope is delivered. The live
1133        // read must observe this mutation.
1134        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            // Read the ack to avoid blocking.
1146            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    // ---- host acceptor demux path (D1) ----
1163
1164    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    /// §11 demux row: an envelope addressed to member B lands in B's inbox
1189    /// and the ack is signed by B's keypair — the ADDRESSED identity, not
1190    /// any other registered identity.
1191    #[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, &registry).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(), &registry)
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    /// §11 misaddressed row: an envelope addressed to a never-registered key
1295    /// is rejected with the typed `Misaddressed` fault, no ack sent.
1296    #[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, &registry).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    /// §11 unregistered-identity row: register → remove → send is the same
1333    /// typed reject as never-registered; removal requires the installed
1334    /// owner (a wrong owner gets a typed error and the registry is never
1335    /// mutated).
1336    #[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        // Wrong owner: typed error, registry never mutated.
1355        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        // Installed owner: removal succeeds, and the identity now rejects.
1367        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, &registry).await;
1393        assert!(matches!(result, Err(IoTaskError::Misaddressed { .. })));
1394        assert!(inbox_a.try_drain_classified().is_empty());
1395    }
1396
1397    /// §11 mandatory-auth row: the demux entry point has no peer-auth
1398    /// parameter (type-level absence — compare `handle_connection`'s
1399    /// `require_peer_auth`), so an unsigned envelope is rejected even when
1400    /// the registered identity's own inbox context was built from a
1401    /// `require_peer_auth = false` configuration.
1402    #[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        // The context a `require_peer_auth = false` CoreCommsConfig would
1408        // produce — the acceptor never reads it.
1409        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]), // unsigned
1424        };
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, &registry).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    /// §11 byte-compat row: the demux path accepts the SAME hand-written
1442    /// wire bytes (u32-BE length prefix + ciborium CBOR envelope — the
1443    /// `envelope_to_bytes` recipe) as the single-identity path, and the
1444    /// signature verifies — no signed-region or codec drift across the
1445    /// refactor. The exact-byte signable-region fixtures in types.rs stand
1446    /// untouched as the primary pin.
1447    #[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        // Hand-encode the frame: 4-byte big-endian length prefix followed by
1469        // the ciborium CBOR envelope, no codec involved on the write side.
1470        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, &registry).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}