Skip to main content

unb_core/
protocol_core.rs

1use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
2use std::fmt;
3
4use bytes::Bytes;
5use web_time::Instant;
6
7use crate::relay_reducer::RelayReducer;
8use crate::route_control::Establishment;
9use crate::session::{SessionEffect, SessionInput, SessionReducer};
10use crate::{
11    ApplicationFrame, CoreError, DiscoverEvent, DiscoverPlan, DiscoverWalk, Envelope, ErrorCode,
12    Kind, NodeCore, NodeIdentity, RouteAck, RouteAckStatus, RouteAdvertisement, RouteDelta,
13    RouteError, RouteSnapshot, RouteWithdrawal, SessionClass, TargetPath, WalkInput, WalkOutput,
14};
15
16const MAX_EFFECTS: usize = 64;
17const MAX_EFFECTS_PER_INPUT: usize = 4;
18const MAX_DISCOVERY_EFFECTS_PER_INPUT: usize = 4;
19
20#[derive(Debug, Clone, PartialEq, Eq, Hash)]
21pub struct SessionId(String);
22
23impl SessionId {
24    pub fn as_str(&self) -> &str {
25        &self.0
26    }
27}
28
29impl From<&str> for SessionId {
30    fn from(value: &str) -> Self {
31        Self(value.to_owned())
32    }
33}
34
35impl From<String> for SessionId {
36    fn from(value: String) -> Self {
37        Self(value)
38    }
39}
40
41impl fmt::Display for SessionId {
42    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
43        self.0.fmt(formatter)
44    }
45}
46
47#[derive(Debug, Clone, PartialEq, Eq, Hash)]
48pub struct CorrelationId(String);
49
50impl CorrelationId {
51    pub fn as_str(&self) -> &str {
52        &self.0
53    }
54}
55
56impl From<&str> for CorrelationId {
57    fn from(value: &str) -> Self {
58        Self(value.to_owned())
59    }
60}
61
62impl From<String> for CorrelationId {
63    fn from(value: String) -> Self {
64        Self(value)
65    }
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, Hash)]
69pub struct ClientOperationId(CorrelationId);
70
71impl ClientOperationId {
72    pub fn as_str(&self) -> &str {
73        self.0.as_str()
74    }
75}
76
77impl From<CorrelationId> for ClientOperationId {
78    fn from(value: CorrelationId) -> Self {
79        Self(value)
80    }
81}
82
83impl From<&str> for ClientOperationId {
84    fn from(value: &str) -> Self {
85        Self(CorrelationId::from(value))
86    }
87}
88
89impl From<String> for ClientOperationId {
90    fn from(value: String) -> Self {
91        Self(CorrelationId::from(value))
92    }
93}
94
95#[derive(Debug, Clone, PartialEq)]
96pub enum ClientDelivery {
97    Item(ApplicationFrame),
98    Terminal(ApplicationFrame),
99    Cancelled,
100    TimedOut,
101    SessionClosed,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
105pub struct EffectId(u64);
106
107impl EffectId {
108    pub fn new(value: u64) -> Self {
109        Self(value)
110    }
111
112    pub fn get(self) -> u64 {
113        self.0
114    }
115}
116
117#[derive(Debug, Clone, PartialEq, Eq, Hash)]
118pub struct StreamKey {
119    pub session: SessionId,
120    pub corr: CorrelationId,
121}
122
123#[derive(Debug, Clone, PartialEq)]
124pub struct ApplicationInvocation {
125    pub stream: StreamKey,
126    pub reservation: EffectId,
127    pub origin: ApplicationOrigin,
128    pub frame: ApplicationFrame,
129}
130
131#[derive(Debug, Clone, PartialEq)]
132pub enum ApplicationOrigin {
133    Client {
134        session: SessionId,
135    },
136    Peer {
137        session: SessionId,
138        peer: NodeIdentity,
139    },
140}
141
142#[derive(Debug, Clone, PartialEq)]
143pub enum PeerAdmission {
144    Admitted(NodeIdentity),
145    Rejected(String),
146}
147
148#[derive(Debug, Clone, Copy, PartialEq, Eq)]
149pub enum CapacityResult {
150    Available,
151    Busy,
152}
153
154#[derive(Debug, Clone, PartialEq, Eq)]
155pub enum RelayOpenResult {
156    Opened(StreamKey),
157    Failed(ApplicationFailure),
158}
159
160#[derive(Debug, Clone, PartialEq, Eq)]
161pub enum TargetReadinessResult {
162    Ready,
163    Unavailable { message: String },
164}
165
166#[derive(Debug, Clone)]
167pub struct ApplicationResponse {
168    pub head: http::Response<()>,
169    pub body: Option<crate::BodyId>,
170}
171
172impl PartialEq for ApplicationResponse {
173    fn eq(&self, other: &Self) -> bool {
174        self.head.status() == other.head.status()
175            && self.head.version() == other.head.version()
176            && self.head.headers() == other.head.headers()
177            && self.body == other.body
178    }
179}
180
181#[derive(Debug, Clone, PartialEq)]
182pub enum ApplicationResult {
183    Response(ApplicationResponse),
184    Event(ApplicationResponse),
185    Finished(ApplicationResponse),
186    Bridged,
187}
188
189#[derive(Debug, Clone, PartialEq, Eq)]
190pub struct ApplicationFailure {
191    pub code: ErrorCode,
192    pub message: String,
193}
194
195#[derive(Debug, Clone, PartialEq, Eq)]
196pub enum SendResult {
197    Reserved,
198    Written,
199    ReservationTimedOut,
200    Closed,
201    Cancelled,
202    WriteFailed(String),
203    Refused { code: ErrorCode, message: String },
204}
205
206#[derive(Debug, Clone, Copy, PartialEq, Eq)]
207pub enum OperationStreamDirection {
208    Opening,
209    Return,
210}
211
212#[derive(Debug, Clone, PartialEq, Eq)]
213pub enum OperationStreamOutcome {
214    Clean,
215    Cancelled,
216    TimedOut,
217    Busy,
218    PayloadTooLarge(String),
219    Truncated,
220    Protocol(String),
221    Transport(String),
222}
223
224#[derive(Debug, Clone, PartialEq)]
225pub enum OperationInput {
226    Body(Option<crate::BodyId>),
227    Discovery(DiscoverPlan),
228}
229
230#[derive(Debug, Clone, Copy, PartialEq, Eq)]
231pub enum RetirementReason {
232    EstablishmentTimeout,
233    SessionClosed,
234    MalformedIdentity,
235    DuplicateIdentity,
236    AdmissionRejected,
237    AdmissionIdentityMismatch,
238    SelfConnection,
239    UnexpectedPeer,
240    IdentityAcceptanceOrder,
241    SnapshotSequence,
242    MalformedSnapshot,
243    SnapshotBeforeIdentity,
244    SnapshotRejected,
245    SnapshotOrder,
246    MalformedAcknowledgement,
247    AcknowledgementOrder,
248    MalformedRouteControl,
249    RouteUpdateRejected,
250    DuplicateSession,
251    DuplicateSessionReplaced,
252    SendClosed,
253    SendTransport,
254    TransportFailed,
255}
256
257#[derive(Debug, Clone, PartialEq)]
258pub enum CoreInput {
259    SessionOpened {
260        session: SessionId,
261        initiator: bool,
262        establish_peer: bool,
263        expected_peer: Option<String>,
264    },
265    FrameReceived {
266        session: SessionId,
267        envelope: Envelope,
268    },
269    ApplicationFrameReceived {
270        session: SessionId,
271        frame: ApplicationFrame,
272    },
273    OperationStreamEnded {
274        session: SessionId,
275        corr: CorrelationId,
276        direction: OperationStreamDirection,
277        outcome: OperationStreamOutcome,
278    },
279    StartClientOperation {
280        session: SessionId,
281        target_path: String,
282        kind: Kind,
283        input: OperationInput,
284        hops: Option<u8>,
285        headers: serde_json::Map<String, serde_json::Value>,
286        timeout: Option<std::time::Duration>,
287    },
288    CancelClientOperation {
289        session: SessionId,
290        operation: ClientOperationId,
291    },
292    ClientOperationTimeout {
293        session: SessionId,
294        operation: ClientOperationId,
295    },
296    SessionTimeout {
297        session: SessionId,
298    },
299    PeerAdmissionCompleted {
300        effect: EffectId,
301        result: PeerAdmission,
302    },
303    CapacityChecked {
304        effect: EffectId,
305        result: CapacityResult,
306    },
307    DispatchCompleted {
308        effect: EffectId,
309        result: Result<ApplicationResult, ApplicationFailure>,
310    },
311    RelayOpenCompleted {
312        effect: EffectId,
313        result: RelayOpenResult,
314    },
315    TargetReadinessCompleted {
316        effect: EffectId,
317        result: TargetReadinessResult,
318    },
319    RelayForwardCompleted {
320        effect: EffectId,
321        result: SendResult,
322    },
323    SendCompleted {
324        effect: EffectId,
325        result: SendResult,
326    },
327    EstablishmentTimeout {
328        session: SessionId,
329    },
330    DiscoveryNeighborEvent {
331        stream: StreamKey,
332        peer: String,
333        event: DiscoverEvent,
334    },
335    DiscoveryNeighborDone {
336        stream: StreamKey,
337        peer: String,
338    },
339    DiscoveryNeighborTimeout {
340        stream: StreamKey,
341        peer: String,
342    },
343    DiscoveryTargetEvent {
344        stream: StreamKey,
345        event: DiscoverEvent,
346    },
347    DiscoveryTargetFailed {
348        stream: StreamKey,
349        message: String,
350    },
351    SessionClosed {
352        session: SessionId,
353    },
354    SessionFailed {
355        session: SessionId,
356    },
357    ContinueDiscovery {
358        stream: StreamKey,
359    },
360    LocalCapabilitiesInstalled {
361        capabilities: BTreeMap<String, serde_json::Value>,
362    },
363}
364
365#[derive(Debug, Clone, PartialEq)]
366pub enum CoreEffect {
367    SendFrame {
368        session: SessionId,
369        envelope: Envelope,
370    },
371    HandshakeEstablished {
372        session: SessionId,
373        version: u16,
374    },
375    DeliverClient {
376        session: SessionId,
377        operation: ClientOperationId,
378        delivery: ClientDelivery,
379    },
380    RegisterClient {
381        session: SessionId,
382        operation: ClientOperationId,
383    },
384    ScheduleClientDeadline {
385        session: SessionId,
386        operation: ClientOperationId,
387        deadline: Instant,
388    },
389    Deliver {
390        session: SessionId,
391        envelope: Envelope,
392    },
393    ReleaseBody {
394        session: SessionId,
395        body: crate::BodyId,
396    },
397    StreamClosed {
398        session: SessionId,
399        operation: ClientOperationId,
400    },
401    CloseTransport {
402        session: SessionId,
403        code: ErrorCode,
404        message: String,
405    },
406    ScheduleSessionDeadline {
407        session: SessionId,
408        deadline: Instant,
409    },
410    RequestPeerAdmission {
411        effect: EffectId,
412        session: SessionId,
413        remote: NodeIdentity,
414    },
415    AbortDispatch {
416        session: SessionId,
417        effect: EffectId,
418    },
419    CheckDispatchCapacity {
420        effect: EffectId,
421        stream: StreamKey,
422        frame: ApplicationFrame,
423    },
424    InvokeApplication {
425        effect: EffectId,
426        invocation: ApplicationInvocation,
427    },
428    OpenRelay {
429        effect: EffectId,
430        source: StreamKey,
431        peer: String,
432        frame: ApplicationFrame,
433    },
434    AwaitTargetReadiness {
435        effect: EffectId,
436        stream: StreamKey,
437        target: String,
438    },
439    ForwardRelay {
440        effect: EffectId,
441        source: StreamKey,
442        target: StreamKey,
443        frame: ApplicationFrame,
444        terminal: bool,
445    },
446    Send {
447        effect: EffectId,
448        session: SessionId,
449        frame: ApplicationFrame,
450    },
451    SendProtocol {
452        effect: EffectId,
453        session: SessionId,
454        envelope: Envelope,
455    },
456    SessionEstablished {
457        session: SessionId,
458        peer: NodeIdentity,
459    },
460    SessionRetired {
461        session: SessionId,
462        reason: RetirementReason,
463    },
464    RouteSnapshotApplied {
465        session: SessionId,
466        peer: String,
467        snapshot: RouteSnapshot,
468        changed: bool,
469    },
470    RouteDeltaApplied {
471        session: SessionId,
472        peer: String,
473        delta: RouteDelta,
474        changed: bool,
475    },
476    RouteSessionWithdrawn {
477        session: SessionId,
478        changed: bool,
479    },
480    RouteExportAcked {
481        session: SessionId,
482    },
483    QueryDiscoveryNeighbor {
484        stream: StreamKey,
485        peer: String,
486        plan: DiscoverPlan,
487    },
488    QueryDiscoveryTarget {
489        stream: StreamKey,
490        peer: String,
491        target_path: String,
492        plan: DiscoverPlan,
493    },
494    ContinueDiscovery {
495        stream: StreamKey,
496    },
497}
498
499impl CoreEffect {
500    pub fn session(&self) -> &SessionId {
501        match self {
502            Self::SendFrame { session, .. }
503            | Self::HandshakeEstablished { session, .. }
504            | Self::DeliverClient { session, .. }
505            | Self::RegisterClient { session, .. }
506            | Self::ScheduleClientDeadline { session, .. }
507            | Self::Deliver { session, .. }
508            | Self::ReleaseBody { session, .. }
509            | Self::StreamClosed { session, .. }
510            | Self::CloseTransport { session, .. }
511            | Self::ScheduleSessionDeadline { session, .. }
512            | Self::RequestPeerAdmission { session, .. }
513            | Self::AbortDispatch { session, .. }
514            | Self::Send { session, .. }
515            | Self::SendProtocol { session, .. }
516            | Self::SessionEstablished { session, .. }
517            | Self::SessionRetired { session, .. }
518            | Self::RouteSnapshotApplied { session, .. }
519            | Self::RouteDeltaApplied { session, .. }
520            | Self::RouteSessionWithdrawn { session, .. }
521            | Self::RouteExportAcked { session } => session,
522            Self::CheckDispatchCapacity { stream, .. }
523            | Self::InvokeApplication {
524                invocation: ApplicationInvocation { stream, .. },
525                ..
526            }
527            | Self::OpenRelay { source: stream, .. }
528            | Self::AwaitTargetReadiness { stream, .. }
529            | Self::ForwardRelay { target: stream, .. }
530            | Self::QueryDiscoveryNeighbor { stream, .. }
531            | Self::QueryDiscoveryTarget { stream, .. }
532            | Self::ContinueDiscovery { stream } => &stream.session,
533        }
534    }
535}
536
537enum PendingEffect {
538    PeerAdmission(SessionId),
539    Capacity(Box<ApplicationInvocation>),
540    Dispatch(StreamKey),
541    RelayOpen {
542        stream: StreamKey,
543        body: Option<crate::BodyId>,
544    },
545    TargetReadiness {
546        stream: StreamKey,
547        frame: ApplicationFrame,
548    },
549    RelayForward {
550        source: StreamKey,
551        target: SessionId,
552        terminal: bool,
553    },
554    Send {
555        stream: StreamKey,
556        terminal: bool,
557    },
558}
559
560struct PeerSession {
561    initiator: bool,
562    establish_peer: bool,
563    expected_peer: Option<String>,
564    establishment: Establishment,
565    remote: Option<NodeIdentity>,
566    held_snapshot: Option<RouteSnapshot>,
567    deferred_ack: Option<RouteAck>,
568    ready: bool,
569    retired: bool,
570    export: Option<RouteExportState>,
571}
572
573struct RouteExportState {
574    generation: u64,
575    applied_ack: u64,
576    handled_resync: Option<u64>,
577    routes: Vec<RouteAdvertisement>,
578}
579
580impl PeerSession {
581    fn new(initiator: bool, establish_peer: bool, expected_peer: Option<String>) -> Self {
582        Self {
583            initiator,
584            establish_peer,
585            expected_peer,
586            establishment: Establishment::default(),
587            remote: None,
588            held_snapshot: None,
589            deferred_ack: None,
590            ready: false,
591            retired: false,
592            export: None,
593        }
594    }
595}
596
597pub(crate) const CLIENT_TOMBSTONE_LIMIT: usize = 4096;
598
599pub struct ProtocolCore {
600    node: NodeCore,
601    sessions: HashMap<SessionId, SessionReducer>,
602    peers: HashMap<SessionId, PeerSession>,
603    active_peers: HashMap<String, SessionId>,
604    discoveries: HashMap<StreamKey, DiscoverWalk>,
605    target_discoveries: HashSet<StreamKey>,
606    client_operations: HashMap<StreamKey, Option<Instant>>,
607    completed_client_operations: HashSet<StreamKey>,
608    completed_client_order: VecDeque<StreamKey>,
609    relays: RelayReducer,
610    effects: VecDeque<CoreEffect>,
611    pending: HashMap<EffectId, PendingEffect>,
612    next_effect: u64,
613}
614
615impl ProtocolCore {
616    pub fn new(node: &str) -> Self {
617        let identity = NodeIdentity {
618            node_id: node.to_string(),
619            instance_id: node.to_string(),
620            epoch: 0,
621            proof: serde_json::Value::Null,
622        };
623        Self::with_identity(identity)
624    }
625
626    #[cfg(feature = "fuzzing")]
627    pub fn effect_len(&self) -> usize {
628        self.effects.len()
629    }
630
631    pub fn with_identity(identity: NodeIdentity) -> Self {
632        let mut node = NodeCore::new(&identity.node_id);
633        node.set_node_identity(identity);
634        Self::with_node(node)
635    }
636
637    pub fn with_node(node: NodeCore) -> Self {
638        Self {
639            node,
640            sessions: HashMap::new(),
641            peers: HashMap::new(),
642            active_peers: HashMap::new(),
643            discoveries: HashMap::new(),
644            target_discoveries: HashSet::new(),
645            client_operations: HashMap::new(),
646            completed_client_operations: HashSet::new(),
647            completed_client_order: VecDeque::new(),
648            relays: RelayReducer::default(),
649            effects: VecDeque::new(),
650            pending: HashMap::new(),
651            next_effect: 0,
652        }
653    }
654
655    pub fn node(&self) -> &NodeCore {
656        &self.node
657    }
658
659    pub fn handle(&mut self, now: Instant, input: CoreInput) -> Result<(), CoreError> {
660        if self.effects.len() > MAX_EFFECTS - MAX_EFFECTS_PER_INPUT {
661            return Err(CoreError::EffectQueueFull);
662        }
663        let session = match input {
664            CoreInput::SessionOpened {
665                session,
666                initiator,
667                establish_peer,
668                expected_peer,
669            } => {
670                if self.sessions.contains_key(&session) {
671                    return Err(CoreError::DuplicateSession(session.to_string()));
672                }
673                let mut core = SessionReducer::authoritative();
674                if !establish_peer {
675                    core.mark_peer_ready();
676                }
677                core.handle_input(now, SessionInput::Start { initiator })?;
678                self.sessions.insert(session.clone(), core);
679                self.peers.insert(
680                    session.clone(),
681                    PeerSession::new(initiator, establish_peer, expected_peer),
682                );
683                session
684            }
685            CoreInput::FrameReceived { session, envelope } => {
686                if (envelope.kind.is_application_request() && envelope.kind != Kind::Discover)
687                    || (envelope.corr.is_some()
688                        && (envelope.kind.is_application_response()
689                            || envelope.kind == Kind::Cancel))
690                {
691                    return Err(CoreError::Malformed(
692                        "application frames require the payload-opaque core input".into(),
693                    ));
694                }
695                self.session_mut(&session)?
696                    .handle_input(now, SessionInput::FrameReceived(envelope))?;
697                session
698            }
699            CoreInput::ApplicationFrameReceived { session, frame } => {
700                if let Some(corr) = frame.head.corr.as_deref() {
701                    let stream = StreamKey {
702                        session: session.clone(),
703                        corr: CorrelationId::from(corr),
704                    };
705                    if self.completed_client_operations.contains(&stream) {
706                        self.release_body(&session, frame.body);
707                        return Ok(());
708                    }
709                }
710                self.session_mut(&session)?
711                    .handle_input(now, SessionInput::FrameReceived(frame.into_envelope()))?;
712                session
713            }
714            CoreInput::OperationStreamEnded {
715                session,
716                corr,
717                direction: _,
718                outcome,
719            } => {
720                let stream = StreamKey {
721                    session: session.clone(),
722                    corr: corr.clone(),
723                };
724                if self.completed_client_operations.contains(&stream) {
725                    return Ok(());
726                }
727                if outcome == OperationStreamOutcome::Clean {
728                    return Ok(());
729                }
730                let (code, message) = match outcome {
731                    OperationStreamOutcome::Clean => unreachable!(),
732                    OperationStreamOutcome::Cancelled => (
733                        ErrorCode::Cancelled,
734                        "operation stream cancelled".to_owned(),
735                    ),
736                    OperationStreamOutcome::TimedOut => (
737                        ErrorCode::PeerUnreachable,
738                        "operation stream timed out".to_owned(),
739                    ),
740                    OperationStreamOutcome::Busy => (
741                        ErrorCode::Busy,
742                        "operation stream capacity unavailable".to_owned(),
743                    ),
744                    OperationStreamOutcome::PayloadTooLarge(message) => {
745                        (ErrorCode::PayloadTooLarge, message)
746                    }
747                    OperationStreamOutcome::Truncated => (
748                        ErrorCode::Protocol,
749                        "operation stream ended with a truncated record".to_owned(),
750                    ),
751                    OperationStreamOutcome::Protocol(message) => (ErrorCode::Protocol, message),
752                    OperationStreamOutcome::Transport(message) => {
753                        (ErrorCode::PeerUnreachable, message)
754                    }
755                };
756                self.session_mut(&session)?
757                    .operation_stream_failed(corr.as_str(), code, &message);
758                session
759            }
760            CoreInput::StartClientOperation {
761                session,
762                target_path,
763                kind,
764                input,
765                hops,
766                headers,
767                timeout,
768            } => {
769                let (payload, body) = match input {
770                    OperationInput::Body(body) if kind != Kind::Discover => (Bytes::new(), body),
771                    OperationInput::Discovery(plan) if kind == Kind::Discover => (
772                        Envelope::encode_payload(
773                            &serde_json::to_value(plan)
774                                .map_err(|error| CoreError::Malformed(error.to_string()))?,
775                        ),
776                        None,
777                    ),
778                    _ => {
779                        return Err(CoreError::Malformed(
780                            "operation input does not match its application kind".into(),
781                        ))
782                    }
783                };
784                let corr = if kind == Kind::Discover {
785                    self.session_mut(&session)?.open_discovery(
786                        &target_path,
787                        payload,
788                        hops,
789                        headers,
790                    )?
791                } else {
792                    self.session_mut(&session)?.open_stream_from_body(
793                        &target_path,
794                        kind,
795                        body.as_ref().map(ToString::to_string),
796                        hops,
797                        headers,
798                    )?
799                };
800                let operation = ClientOperationId::from(corr.clone());
801                let stream = StreamKey {
802                    session: session.clone(),
803                    corr: CorrelationId::from(corr),
804                };
805                let deadline = timeout.map(|timeout| now + timeout);
806                self.client_operations.insert(stream.clone(), deadline);
807                self.effects.push_back(CoreEffect::RegisterClient {
808                    session: session.clone(),
809                    operation: operation.clone(),
810                });
811                if let Some(deadline) = deadline {
812                    self.effects.push_back(CoreEffect::ScheduleClientDeadline {
813                        session: session.clone(),
814                        operation,
815                        deadline,
816                    });
817                }
818                let SessionEffect::SendFrame(envelope) = self.session_mut(&session)?.poll_effect()
819                else {
820                    return Err(CoreError::Malformed(
821                        "client opening frame missing".to_owned(),
822                    ));
823                };
824                let effect = self.effect_id();
825                self.pending.insert(
826                    effect,
827                    PendingEffect::Send {
828                        stream,
829                        terminal: false,
830                    },
831                );
832                if envelope.kind == Kind::Discover {
833                    self.effects.push_back(CoreEffect::SendProtocol {
834                        effect,
835                        session,
836                        envelope,
837                    });
838                } else {
839                    self.effects.push_back(CoreEffect::Send {
840                        effect,
841                        session,
842                        frame: ApplicationFrame::from_envelope(&envelope)?,
843                    });
844                }
845                return Ok(());
846            }
847            CoreInput::CancelClientOperation { session, operation } => {
848                self.complete_client_operation(
849                    &session,
850                    &operation,
851                    ClientDelivery::Cancelled,
852                    true,
853                );
854                return Ok(());
855            }
856            CoreInput::ClientOperationTimeout { session, operation } => {
857                let stream = StreamKey {
858                    session: session.clone(),
859                    corr: operation.0.clone(),
860                };
861                if self
862                    .client_operations
863                    .get(&stream)
864                    .is_some_and(|deadline| deadline.is_some_and(|deadline| now >= deadline))
865                {
866                    self.complete_client_operation(
867                        &session,
868                        &operation,
869                        ClientDelivery::TimedOut,
870                        true,
871                    );
872                }
873                return Ok(());
874            }
875            CoreInput::SessionTimeout { session } => {
876                self.session_mut(&session)?
877                    .handle_input(now, SessionInput::Deadline(now))?;
878                session
879            }
880            CoreInput::PeerAdmissionCompleted { effect, result } => {
881                let Some(PendingEffect::PeerAdmission(session)) = self.pending.remove(&effect)
882                else {
883                    return Ok(());
884                };
885                self.complete_admission(session, result)?;
886                return Ok(());
887            }
888            CoreInput::CapacityChecked { effect, result } => {
889                let Some(PendingEffect::Capacity(invocation)) = self.pending.remove(&effect) else {
890                    return Ok(());
891                };
892                let invocation = *invocation;
893                match result {
894                    CapacityResult::Available => {
895                        let effect = self.effect_id();
896                        self.pending
897                            .insert(effect, PendingEffect::Dispatch(invocation.stream.clone()));
898                        self.effects
899                            .push_back(CoreEffect::InvokeApplication { effect, invocation });
900                    }
901                    CapacityResult::Busy => {
902                        self.release_body(&invocation.stream.session, invocation.frame.body);
903                        self.fail_stream(
904                            &invocation.stream,
905                            ErrorCode::Busy,
906                            "node at activation capacity",
907                        )?;
908                    }
909                }
910                return Ok(());
911            }
912            CoreInput::TargetReadinessCompleted { effect, result } => {
913                let Some(PendingEffect::TargetReadiness { stream, frame }) =
914                    self.pending.remove(&effect)
915                else {
916                    return Ok(());
917                };
918                match result {
919                    TargetReadinessResult::Ready => {
920                        let envelope = frame.clone().into_envelope();
921                        if !self.resolve_relay_opening_with_wait(&stream, &envelope, false) {
922                            self.release_body(&stream.session, frame.body);
923                            self.fail_stream(
924                                &stream,
925                                ErrorCode::PeerUnreachable,
926                                "target readiness completed without an authoritative route",
927                            )?;
928                        }
929                    }
930                    TargetReadinessResult::Unavailable { message } => {
931                        self.release_body(&stream.session, frame.body);
932                        self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message)?;
933                    }
934                }
935                return Ok(());
936            }
937            CoreInput::RelayOpenCompleted { effect, result } => {
938                let Some(PendingEffect::RelayOpen { stream: source, .. }) =
939                    self.pending.remove(&effect)
940                else {
941                    return Ok(());
942                };
943                match result {
944                    RelayOpenResult::Opened(target) => {
945                        if !self.relays.pair(source.clone(), target) {
946                            self.fail_stream(
947                                &source,
948                                ErrorCode::Conflict,
949                                "relay stream is already paired",
950                            )?;
951                        }
952                    }
953                    RelayOpenResult::Failed(failure) => {
954                        self.fail_stream(&source, failure.code, &failure.message)?;
955                    }
956                }
957                return Ok(());
958            }
959            CoreInput::RelayForwardCompleted { effect, result } => {
960                if result == SendResult::Reserved {
961                    return Ok(());
962                }
963                let Some(PendingEffect::RelayForward {
964                    source,
965                    target,
966                    terminal,
967                }) = self.pending.remove(&effect)
968                else {
969                    return Ok(());
970                };
971                match result {
972                    SendResult::Written => {
973                        if terminal {
974                            self.relays.remove_pair(&source);
975                        }
976                    }
977                    SendResult::ReservationTimedOut => self.fail_relay_forward(&source),
978                    SendResult::Refused { code, message } => {
979                        self.fail_relay_forward_with(&source, code, &message)
980                    }
981                    SendResult::Closed | SendResult::Cancelled => {
982                        self.fail_relay_forward(&source);
983                        self.retire(target, RetirementReason::SendClosed);
984                    }
985                    SendResult::WriteFailed(_) => {
986                        self.fail_relay_forward(&source);
987                        self.retire(target, RetirementReason::SendTransport);
988                    }
989                    SendResult::Reserved => unreachable!(),
990                }
991                return Ok(());
992            }
993            CoreInput::DispatchCompleted { effect, result } => {
994                let Some(PendingEffect::Dispatch(stream)) = self.pending.remove(&effect) else {
995                    return Ok(());
996                };
997                match result {
998                    Ok(ApplicationResult::Response(payload))
999                    | Ok(ApplicationResult::Finished(payload)) => {
1000                        self.pending.remove(&effect);
1001                        self.queue_application_response(stream, payload, true)?;
1002                    }
1003                    Ok(ApplicationResult::Event(payload)) => {
1004                        self.pending
1005                            .insert(effect, PendingEffect::Dispatch(stream.clone()));
1006                        self.queue_application_response(stream, payload, false)?;
1007                    }
1008                    Ok(ApplicationResult::Bridged) => {}
1009                    Err(failure) => {
1010                        self.pending.remove(&effect);
1011                        self.fail_stream(&stream, failure.code, &failure.message)?;
1012                    }
1013                }
1014                return Ok(());
1015            }
1016            CoreInput::SendCompleted { effect, result } => {
1017                if result == SendResult::Reserved {
1018                    return Ok(());
1019                }
1020                let Some(PendingEffect::Send { stream, terminal }) = self.pending.remove(&effect)
1021                else {
1022                    return Ok(());
1023                };
1024                let relay_leg = self.relays.peer(&stream).is_some();
1025                match result {
1026                    SendResult::Written => {}
1027                    SendResult::ReservationTimedOut if relay_leg => {
1028                        self.fail_relay_forward(&stream);
1029                    }
1030                    SendResult::ReservationTimedOut if !terminal => {
1031                        self.fail_stream(&stream, ErrorCode::Busy, "send capacity unavailable")?;
1032                    }
1033                    SendResult::ReservationTimedOut => {}
1034                    SendResult::Refused { code, message } if relay_leg => {
1035                        self.fail_relay_forward_with(&stream, code, &message);
1036                    }
1037                    SendResult::Refused { code, message } if !terminal => {
1038                        self.fail_stream(&stream, code, &message)?;
1039                    }
1040                    SendResult::Refused { code, message } => {
1041                        if let Ok(session) = self.session_mut(&stream.session) {
1042                            session.fail_closed(stream.corr.as_str(), code, &message);
1043                        }
1044                        self.drain_session(stream.session.clone());
1045                    }
1046                    SendResult::Closed | SendResult::Cancelled => {
1047                        self.retire(stream.session, RetirementReason::SendClosed);
1048                    }
1049                    SendResult::WriteFailed(_) => {
1050                        self.retire(stream.session, RetirementReason::SendTransport);
1051                    }
1052                    SendResult::Reserved => unreachable!(),
1053                }
1054                return Ok(());
1055            }
1056            CoreInput::EstablishmentTimeout { session } => {
1057                if self.peers.get(&session).is_some_and(|peer| peer.ready)
1058                    || self.session_class(&session)? == SessionClass::Client
1059                {
1060                    return Ok(());
1061                }
1062                self.retire(session, RetirementReason::EstablishmentTimeout);
1063                return Ok(());
1064            }
1065            CoreInput::DiscoveryNeighborEvent {
1066                stream,
1067                peer,
1068                event,
1069            } => {
1070                self.handle_discovery(stream, WalkInput::NeighborEvent { peer, event });
1071                return Ok(());
1072            }
1073            CoreInput::DiscoveryNeighborDone { stream, peer } => {
1074                self.handle_discovery(stream, WalkInput::NeighborDone { peer });
1075                return Ok(());
1076            }
1077            CoreInput::DiscoveryNeighborTimeout { stream, peer } => {
1078                self.handle_discovery(stream, WalkInput::NeighborTimeout { peer });
1079                return Ok(());
1080            }
1081            CoreInput::DiscoveryTargetEvent { stream, event } => {
1082                if !self.target_discoveries.contains(&stream) {
1083                    return Ok(());
1084                }
1085                let terminal = matches!(event, DiscoverEvent::Done { .. });
1086                let payload = Envelope::encode_payload(
1087                    &serde_json::to_value(event)
1088                        .map_err(|error| CoreError::Malformed(error.to_string()))?,
1089                );
1090                let kind = if terminal {
1091                    Kind::Response
1092                } else {
1093                    Kind::Event
1094                };
1095                if terminal {
1096                    self.target_discoveries.remove(&stream);
1097                }
1098                let _ = self.queue_send(stream, kind, payload);
1099                return Ok(());
1100            }
1101            CoreInput::DiscoveryTargetFailed { stream, message } => {
1102                if self.target_discoveries.remove(&stream) {
1103                    let _ = self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message);
1104                }
1105                return Ok(());
1106            }
1107            CoreInput::SessionClosed { session } => {
1108                self.close_session(&session, RetirementReason::SessionClosed);
1109                return Ok(());
1110            }
1111            CoreInput::SessionFailed { session } => {
1112                self.close_session(&session, RetirementReason::TransportFailed);
1113                return Ok(());
1114            }
1115            CoreInput::ContinueDiscovery { stream } => {
1116                self.drain_discovery(stream);
1117                return Ok(());
1118            }
1119            CoreInput::LocalCapabilitiesInstalled { capabilities } => {
1120                self.node.install_local_capabilities(capabilities);
1121                return Ok(());
1122            }
1123        };
1124        self.drain_session(session);
1125        Ok(())
1126    }
1127
1128    pub fn poll_effect(&mut self) -> Option<CoreEffect> {
1129        self.effects.pop_front()
1130    }
1131
1132    pub fn session_deadline(&self, session: &SessionId) -> Result<Option<Instant>, CoreError> {
1133        self.sessions
1134            .get(session)
1135            .map(SessionReducer::deadline)
1136            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1137    }
1138
1139    pub fn session_class(&self, session: &SessionId) -> Result<SessionClass, CoreError> {
1140        self.sessions
1141            .get(session)
1142            .map(SessionReducer::class)
1143            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1144    }
1145
1146    pub fn session_ready(&self, session: &SessionId) -> Result<bool, CoreError> {
1147        self.peers
1148            .get(session)
1149            .map(|peer| peer.ready)
1150            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1151    }
1152
1153    pub fn relay_peer(&self, stream: &StreamKey) -> Option<&StreamKey> {
1154        self.relays.peer(stream)
1155    }
1156
1157    pub fn open_stream_body(
1158        &mut self,
1159        session: &SessionId,
1160        target_path: &str,
1161        kind: Kind,
1162        body: Option<crate::BodyId>,
1163        hops: Option<u8>,
1164        headers: serde_json::Map<String, serde_json::Value>,
1165    ) -> Result<String, CoreError> {
1166        if kind == Kind::Discover {
1167            return Err(CoreError::BadKind(format!("{kind:?}")));
1168        }
1169        let corr = self.session_mut(session)?.open_stream_from_body(
1170            target_path,
1171            kind,
1172            body.map(|body| body.to_string()),
1173            hops,
1174            headers,
1175        )?;
1176        self.drain_session(session.clone());
1177        Ok(corr)
1178    }
1179
1180    pub fn send_body(
1181        &mut self,
1182        session: &SessionId,
1183        corr: &str,
1184        body: Option<crate::BodyId>,
1185    ) -> Result<(), CoreError> {
1186        self.session_mut(session)?.send_body(
1187            corr,
1188            body.map(|body| body.to_string()),
1189            Default::default(),
1190        )?;
1191        self.drain_session(session.clone());
1192        Ok(())
1193    }
1194
1195    pub fn respond_body(
1196        &mut self,
1197        session: &SessionId,
1198        corr: &str,
1199        body: Option<crate::BodyId>,
1200        headers: serde_json::Map<String, serde_json::Value>,
1201    ) -> Result<(), CoreError> {
1202        self.session_mut(session)?.respond_body(
1203            corr,
1204            body.map(|body| body.to_string()),
1205            headers,
1206        )?;
1207        self.drain_session(session.clone());
1208        Ok(())
1209    }
1210
1211    pub fn respond_terminal_body(
1212        &mut self,
1213        session: &SessionId,
1214        corr: &str,
1215        body: Option<crate::BodyId>,
1216        path: Vec<String>,
1217    ) -> Result<(), CoreError> {
1218        self.session_mut(session)?.respond_terminal_body(
1219            corr,
1220            body.map(|body| body.to_string()),
1221            path,
1222            Default::default(),
1223        )?;
1224        self.drain_session(session.clone());
1225        Ok(())
1226    }
1227
1228    pub fn control(
1229        &mut self,
1230        session: &SessionId,
1231        kind: Kind,
1232        payload: Bytes,
1233    ) -> Result<(), CoreError> {
1234        self.session_mut(session)?.control(kind, payload)?;
1235        self.drain_session(session.clone());
1236        Ok(())
1237    }
1238
1239    pub fn fail(
1240        &mut self,
1241        session: &SessionId,
1242        corr: &str,
1243        code: ErrorCode,
1244        message: &str,
1245    ) -> Result<(), CoreError> {
1246        self.session_mut(session)?.fail(corr, code, message)?;
1247        self.drain_session(session.clone());
1248        Ok(())
1249    }
1250
1251    pub fn cancel(&mut self, session: &SessionId, corr: &str) -> Result<(), CoreError> {
1252        self.session_mut(session)?.cancel(corr)?;
1253        self.drain_session(session.clone());
1254        Ok(())
1255    }
1256
1257    fn session_mut(&mut self, session: &SessionId) -> Result<&mut SessionReducer, CoreError> {
1258        self.sessions
1259            .get_mut(session)
1260            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1261    }
1262
1263    fn drain_session(&mut self, session: SessionId) {
1264        loop {
1265            let output = self.sessions.get_mut(&session).unwrap().poll_effect();
1266            match output {
1267                SessionEffect::HandshakeEstablished { version } => {
1268                    if self
1269                        .peers
1270                        .get(&session)
1271                        .is_some_and(|peer| peer.initiator && peer.establish_peer)
1272                    {
1273                        let _ = self.send_identify(&session);
1274                    }
1275                    self.effects.push_back(CoreEffect::HandshakeEstablished {
1276                        session: session.clone(),
1277                        version,
1278                    });
1279                }
1280                SessionEffect::Deliver(envelope) if envelope.kind == Kind::Identify => {
1281                    if self.peers[&session].establish_peer {
1282                        self.on_identify(&session, envelope);
1283                    } else {
1284                        self.deliver_client(session.clone(), envelope);
1285                    }
1286                }
1287                SessionEffect::Deliver(envelope) if envelope.kind == Kind::IdentityAccepted => {
1288                    if self.peers[&session].establish_peer {
1289                        self.on_identity_accepted(&session);
1290                    } else {
1291                        self.deliver_client(session.clone(), envelope);
1292                    }
1293                }
1294                SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteSnapshot => {
1295                    if !self.peers[&session].establish_peer {
1296                        self.deliver_client(session.clone(), envelope);
1297                    } else if self.peers.get(&session).is_some_and(|peer| peer.ready) {
1298                        self.on_route_snapshot(&session, envelope);
1299                    } else {
1300                        self.on_initial_snapshot(&session, envelope);
1301                    }
1302                }
1303                SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteDelta => {
1304                    if self.peers.get(&session).is_some_and(|peer| peer.ready) {
1305                        self.on_route_delta(&session, envelope);
1306                    } else {
1307                        self.retire(session.clone(), RetirementReason::SnapshotOrder);
1308                    }
1309                }
1310                SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteAck => {
1311                    if !self.peers[&session].establish_peer {
1312                        self.deliver_client(session.clone(), envelope);
1313                    } else if self.peers.get(&session).is_some_and(|peer| peer.ready) {
1314                        self.on_route_ack(&session, envelope);
1315                    } else {
1316                        self.on_initial_ack(&session, envelope);
1317                    }
1318                }
1319                SessionEffect::Deliver(envelope) if envelope.kind == Kind::Discover => {
1320                    self.start_discovery(&session, envelope);
1321                }
1322                SessionEffect::Deliver(envelope) if envelope.kind.is_application_request() => {
1323                    if !self.peers[&session].establish_peer {
1324                        self.deliver_client(session.clone(), envelope);
1325                        continue;
1326                    }
1327                    if self
1328                        .peers
1329                        .get(&session)
1330                        .is_some_and(|peer| peer.expected_peer.is_some() && !peer.ready)
1331                    {
1332                        self.retire(session.clone(), RetirementReason::SessionClosed);
1333                        continue;
1334                    }
1335                    if let Some(corr) = envelope.corr.clone() {
1336                        let stream = StreamKey {
1337                            session: session.clone(),
1338                            corr: CorrelationId::from(corr),
1339                        };
1340                        if self.resolve_relay_opening(&stream, &envelope) {
1341                            continue;
1342                        }
1343                        let effect = self.effect_id();
1344                        let frame = match ApplicationFrame::from_envelope(&envelope) {
1345                            Ok(frame) => frame,
1346                            Err(_) => {
1347                                self.retire(session.clone(), RetirementReason::TransportFailed);
1348                                continue;
1349                            }
1350                        };
1351                        let invocation = ApplicationInvocation {
1352                            stream: stream.clone(),
1353                            reservation: effect,
1354                            origin: self.application_origin(&session),
1355                            frame: frame.clone(),
1356                        };
1357                        self.pending
1358                            .insert(effect, PendingEffect::Capacity(Box::new(invocation)));
1359                        self.effects.push_back(CoreEffect::CheckDispatchCapacity {
1360                            effect,
1361                            stream,
1362                            frame,
1363                        });
1364                    }
1365                }
1366                SessionEffect::Deliver(envelope) => {
1367                    let stream = envelope.corr.as_ref().map(|corr| StreamKey {
1368                        session: session.clone(),
1369                        corr: CorrelationId::from(corr.clone()),
1370                    });
1371                    let forwarded =
1372                        ApplicationFrame::from_envelope(&envelope)
1373                            .ok()
1374                            .and_then(|frame| {
1375                                stream.as_ref().and_then(|stream| {
1376                                    self.relays.forward(stream, frame, self.node.node())
1377                                })
1378                            });
1379                    if let (Some(source), Some(mut forwarded)) = (stream, forwarded) {
1380                        if forwarded.terminal {
1381                            if let Some(target) = self.sessions.get_mut(&forwarded.target.session) {
1382                                let forwarded_envelope = forwarded.frame.clone().into_envelope();
1383                                let result = target.respond_terminal_with(
1384                                    forwarded.target.corr.as_str(),
1385                                    forwarded_envelope.kind,
1386                                    forwarded_envelope.payload,
1387                                    forwarded_envelope.path,
1388                                    forwarded_envelope.headers,
1389                                );
1390                                if result.is_ok() {
1391                                    if let SessionEffect::SendFrame(mut envelope) =
1392                                        target.poll_effect()
1393                                    {
1394                                        envelope.body_token =
1395                                            forwarded.frame.body.as_ref().map(ToString::to_string);
1396                                        if let Ok(frame) =
1397                                            ApplicationFrame::from_envelope(&envelope)
1398                                        {
1399                                            forwarded.frame = frame;
1400                                        }
1401                                    }
1402                                }
1403                            }
1404                        }
1405                        let effect = self.effect_id();
1406                        self.pending.insert(
1407                            effect,
1408                            PendingEffect::RelayForward {
1409                                source: source.clone(),
1410                                target: forwarded.target.session.clone(),
1411                                terminal: forwarded.terminal,
1412                            },
1413                        );
1414                        self.effects.push_back(CoreEffect::ForwardRelay {
1415                            effect,
1416                            source,
1417                            target: forwarded.target,
1418                            frame: forwarded.frame,
1419                            terminal: forwarded.terminal,
1420                        });
1421                    } else {
1422                        self.deliver_client(session.clone(), envelope);
1423                    }
1424                }
1425                SessionEffect::Closed { code, message } => {
1426                    self.effects.push_back(CoreEffect::CloseTransport {
1427                        session: session.clone(),
1428                        code,
1429                        message,
1430                    });
1431                    self.retire(session.clone(), RetirementReason::SessionClosed);
1432                }
1433                SessionEffect::SendFrame(envelope)
1434                    if envelope.kind.is_application_request()
1435                        || (envelope.corr.is_some()
1436                            && (envelope.kind.is_application_response()
1437                                || envelope.kind == Kind::Cancel)) =>
1438                {
1439                    if envelope.kind == Kind::Discover {
1440                        let effect = self.effect_id();
1441                        self.effects.push_back(CoreEffect::SendProtocol {
1442                            effect,
1443                            session: session.clone(),
1444                            envelope,
1445                        });
1446                        continue;
1447                    }
1448                    let Some(corr) = envelope.corr.clone() else {
1449                        self.retire(session.clone(), RetirementReason::TransportFailed);
1450                        continue;
1451                    };
1452                    let stream = StreamKey {
1453                        session: session.clone(),
1454                        corr: CorrelationId::from(corr),
1455                    };
1456                    let terminal =
1457                        matches!(envelope.kind, Kind::Response | Kind::Error | Kind::Cancel);
1458                    let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
1459                        self.retire(session.clone(), RetirementReason::TransportFailed);
1460                        continue;
1461                    };
1462                    let effect = self.effect_id();
1463                    self.pending
1464                        .insert(effect, PendingEffect::Send { stream, terminal });
1465                    self.effects.push_back(CoreEffect::Send {
1466                        effect,
1467                        session: session.clone(),
1468                        frame,
1469                    });
1470                }
1471                SessionEffect::SendFrame(envelope) => {
1472                    self.effects.push_back(CoreEffect::SendFrame {
1473                        session: session.clone(),
1474                        envelope,
1475                    })
1476                }
1477                SessionEffect::StreamClosed { corr } => self.stream_closed(session.clone(), corr),
1478                SessionEffect::Deadline(Some(deadline)) => {
1479                    self.effects.push_back(CoreEffect::ScheduleSessionDeadline {
1480                        session: session.clone(),
1481                        deadline,
1482                    });
1483                    break;
1484                }
1485                SessionEffect::Deadline(None) => break,
1486            }
1487        }
1488    }
1489
1490    fn effect_id(&mut self) -> EffectId {
1491        let effect = EffectId::new(self.next_effect);
1492        self.next_effect = self.next_effect.checked_add(1).expect("effect id overflow");
1493        effect
1494    }
1495
1496    fn resolve_relay_opening(&mut self, stream: &StreamKey, envelope: &Envelope) -> bool {
1497        self.resolve_relay_opening_with_wait(stream, envelope, true)
1498    }
1499
1500    fn resolve_relay_opening_with_wait(
1501        &mut self,
1502        stream: &StreamKey,
1503        envelope: &Envelope,
1504        allow_wait: bool,
1505    ) -> bool {
1506        if envelope.target.is_empty() {
1507            return false;
1508        }
1509        match self.node.resolve(&envelope.target) {
1510            crate::Resolution::Route(_) => match self.node.forward(envelope.clone()) {
1511                Ok((peer, envelope)) => {
1512                    let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
1513                        let _ = self.fail_stream(
1514                            stream,
1515                            ErrorCode::Protocol,
1516                            "relay application frame is malformed",
1517                        );
1518                        return true;
1519                    };
1520                    let effect = self.effect_id();
1521                    self.pending.insert(
1522                        effect,
1523                        PendingEffect::RelayOpen {
1524                            stream: stream.clone(),
1525                            body: frame.body.clone(),
1526                        },
1527                    );
1528                    self.effects.push_back(CoreEffect::OpenRelay {
1529                        effect,
1530                        source: stream.clone(),
1531                        peer,
1532                        frame,
1533                    });
1534                }
1535                Err(crate::RouteError::HopLimitExceeded) => {
1536                    self.release_envelope_body(&stream.session, envelope);
1537                    let payload = Envelope::encode_payload(&serde_json::json!({
1538                        "code": ErrorCode::HopLimitExceeded,
1539                        "message": "hop limit exceeded",
1540                    }));
1541                    let node = self.node.node().to_string();
1542                    let _ = self.session_mut(&stream.session).and_then(|session| {
1543                        session.respond_terminal(
1544                            stream.corr.as_str(),
1545                            Kind::Error,
1546                            payload,
1547                            vec![node],
1548                        )
1549                    });
1550                    self.drain_session(stream.session.clone());
1551                }
1552                Err(error) => {
1553                    self.release_envelope_body(&stream.session, envelope);
1554                    let _ = self.fail_stream(stream, ErrorCode::Internal, &error.to_string());
1555                }
1556            },
1557            crate::Resolution::Conflicted { owners } => {
1558                self.release_envelope_body(&stream.session, envelope);
1559                let _ = self.fail_stream(
1560                    stream,
1561                    ErrorCode::PeerUnreachable,
1562                    &format!(
1563                        "destination node {:?} has multiple live incarnations: {}",
1564                        envelope.target,
1565                        owners.join(", ")
1566                    ),
1567                );
1568            }
1569            crate::Resolution::Unknown => {
1570                if allow_wait {
1571                    let Ok(frame) = ApplicationFrame::from_envelope(envelope) else {
1572                        let _ = self.fail_stream(
1573                            stream,
1574                            ErrorCode::Protocol,
1575                            "relay application frame is malformed",
1576                        );
1577                        return true;
1578                    };
1579                    let effect = self.effect_id();
1580                    self.pending.insert(
1581                        effect,
1582                        PendingEffect::TargetReadiness {
1583                            stream: stream.clone(),
1584                            frame,
1585                        },
1586                    );
1587                    self.effects.push_back(CoreEffect::AwaitTargetReadiness {
1588                        effect,
1589                        stream: stream.clone(),
1590                        target: envelope.target.clone(),
1591                    });
1592                } else {
1593                    self.release_envelope_body(&stream.session, envelope);
1594                    let known = self.node.reachable_names();
1595                    let known = known.iter().map(String::as_str).collect::<Vec<_>>();
1596                    let message = crate::teach_unknown("node", &envelope.target, &known);
1597                    let _ = self.fail_stream(stream, ErrorCode::PeerUnreachable, &message);
1598                }
1599            }
1600            crate::Resolution::Local => return false,
1601        }
1602        true
1603    }
1604
1605    fn release_envelope_body(&mut self, session: &SessionId, envelope: &Envelope) {
1606        let body = ApplicationFrame::from_envelope(envelope)
1607            .ok()
1608            .and_then(|frame| frame.body);
1609        self.release_body(session, body);
1610    }
1611
1612    fn start_discovery(&mut self, session: &SessionId, envelope: Envelope) {
1613        let Some(corr) = envelope.corr else {
1614            return;
1615        };
1616        let stream = StreamKey {
1617            session: session.clone(),
1618            corr: CorrelationId::from(corr),
1619        };
1620        let Ok(plan) = DiscoverPlan::decode(&envelope.payload) else {
1621            let _ = self.fail_stream(&stream, ErrorCode::InvalidInput, "malformed discovery plan");
1622            return;
1623        };
1624        match self.node.resolve(&envelope.target) {
1625            crate::Resolution::Route(peer) => {
1626                let hops = envelope.hops.unwrap_or(crate::DEFAULT_HOPS);
1627                if hops == 0 {
1628                    let _ = self.fail_stream(
1629                        &stream,
1630                        ErrorCode::HopLimitExceeded,
1631                        "hop limit exceeded",
1632                    );
1633                    return;
1634                }
1635                let Ok(target_path) = TargetPath::discovery(&envelope.target) else {
1636                    let _ = self.fail_stream(
1637                        &stream,
1638                        ErrorCode::InvalidInput,
1639                        "invalid discovery target",
1640                    );
1641                    return;
1642                };
1643                let mut plan = plan;
1644                plan.hops = plan.hops.min(hops - 1);
1645                self.target_discoveries.insert(stream.clone());
1646                self.effects.push_back(CoreEffect::QueryDiscoveryTarget {
1647                    stream,
1648                    peer,
1649                    target_path: target_path.to_string(),
1650                    plan,
1651                });
1652                return;
1653            }
1654            crate::Resolution::Conflicted { owners } => {
1655                let _ = self.fail_stream(
1656                    &stream,
1657                    ErrorCode::PeerUnreachable,
1658                    &format!(
1659                        "destination node {:?} has multiple live incarnations: {}",
1660                        envelope.target,
1661                        owners.join(", ")
1662                    ),
1663                );
1664                return;
1665            }
1666            crate::Resolution::Unknown => {
1667                let known = self.node.reachable_names();
1668                let known = known.iter().map(String::as_str).collect::<Vec<_>>();
1669                let message = crate::teach_unknown("node", &envelope.target, &known);
1670                let _ = self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message);
1671                return;
1672            }
1673            crate::Resolution::Local => {}
1674        }
1675        let candidates = self.active_peers.keys().cloned().collect();
1676        let snapshot = self.node.catalog_snapshot(plan.detail.is_full());
1677        self.discoveries.insert(
1678            stream.clone(),
1679            DiscoverWalk::start(snapshot, plan, candidates),
1680        );
1681        self.drain_discovery(stream);
1682    }
1683
1684    fn handle_discovery(&mut self, stream: StreamKey, input: WalkInput) {
1685        let Some(walk) = self.discoveries.get_mut(&stream) else {
1686            return;
1687        };
1688        walk.handle(input);
1689        self.drain_discovery(stream);
1690    }
1691
1692    fn drain_discovery(&mut self, stream: StreamKey) {
1693        for _ in 0..MAX_DISCOVERY_EFFECTS_PER_INPUT - 1 {
1694            let output = self
1695                .discoveries
1696                .get_mut(&stream)
1697                .and_then(DiscoverWalk::drain);
1698            match output {
1699                Some(WalkOutput::Emit(event)) => {
1700                    let payload = Envelope::encode_payload(&serde_json::to_value(event).unwrap());
1701                    let _ = self.queue_send(stream.clone(), Kind::Event, payload);
1702                }
1703                Some(WalkOutput::AskNeighbor { peer, plan }) => {
1704                    self.effects.push_back(CoreEffect::QueryDiscoveryNeighbor {
1705                        stream: stream.clone(),
1706                        peer,
1707                        plan,
1708                    });
1709                }
1710                Some(WalkOutput::Finish) => {
1711                    let discover_id = self
1712                        .discoveries
1713                        .get(&stream)
1714                        .map(|walk| walk.discover_id().to_string())
1715                        .unwrap_or_default();
1716                    let failure = self
1717                        .discoveries
1718                        .get(&stream)
1719                        .and_then(DiscoverWalk::strict_failure)
1720                        .map(str::to_owned);
1721                    self.discoveries.remove(&stream);
1722                    if let Some(peer) = failure {
1723                        let _ = self.fail_stream(
1724                            &stream,
1725                            ErrorCode::PeerUnreachable,
1726                            &format!("discovery timed out at {peer}"),
1727                        );
1728                    } else {
1729                        let done = DiscoverEvent::Done { discover_id };
1730                        let payload =
1731                            Envelope::encode_payload(&serde_json::to_value(done).unwrap());
1732                        let _ = self.queue_send(stream, Kind::Response, payload);
1733                    }
1734                    return;
1735                }
1736                None => return,
1737            }
1738        }
1739        if self
1740            .discoveries
1741            .get(&stream)
1742            .is_some_and(DiscoverWalk::has_output)
1743        {
1744            self.effects
1745                .push_back(CoreEffect::ContinueDiscovery { stream });
1746        }
1747    }
1748
1749    fn close_session(&mut self, session: &SessionId, reason: RetirementReason) {
1750        let affected: Vec<_> = self
1751            .discoveries
1752            .keys()
1753            .filter(|stream| &stream.session == session)
1754            .cloned()
1755            .collect();
1756        for stream in affected {
1757            self.discoveries.remove(&stream);
1758        }
1759        self.target_discoveries
1760            .retain(|stream| &stream.session != session);
1761        for (peer, peer_is_source) in self.relays.drain_session(session) {
1762            if peer_is_source {
1763                let payload = Envelope::encode_payload(&serde_json::json!({
1764                    "code": ErrorCode::PeerUnreachable,
1765                    "message": "the peer serving this stream disconnected",
1766                }));
1767                let node = self.node.node().to_string();
1768                if let Ok(target) = self.session_mut(&peer.session) {
1769                    let _ = target.respond_terminal(
1770                        peer.corr.as_str(),
1771                        Kind::Error,
1772                        payload,
1773                        vec![node],
1774                    );
1775                }
1776            } else if let Ok(target) = self.session_mut(&peer.session) {
1777                let _ = target.cancel(peer.corr.as_str());
1778            }
1779            self.drain_session(peer.session);
1780        }
1781        self.retire(session.clone(), reason);
1782    }
1783
1784    fn fail_relay_forward(&mut self, source: &StreamKey) {
1785        self.fail_relay_forward_with(
1786            source,
1787            ErrorCode::Busy,
1788            "relay dropped a slow consumer stream",
1789        );
1790    }
1791
1792    fn fail_relay_forward_with(&mut self, source: &StreamKey, code: ErrorCode, message: &str) {
1793        let Some((peer, source_is_return)) = self.relays.remove_pair(source) else {
1794            return;
1795        };
1796        if source_is_return {
1797            let payload = Envelope::encode_payload(&serde_json::json!({
1798                "code": code,
1799                "message": message,
1800            }));
1801            let node = self.node.node().to_string();
1802            if let Ok(session) = self.session_mut(&peer.session) {
1803                let _ =
1804                    session.respond_terminal(peer.corr.as_str(), Kind::Error, payload, vec![node]);
1805            }
1806            self.drain_session(peer.session.clone());
1807            if let Ok(session) = self.session_mut(&source.session) {
1808                let _ = session.cancel(source.corr.as_str());
1809            }
1810            self.drain_session(source.session.clone());
1811        } else {
1812            if let Ok(session) = self.session_mut(&source.session) {
1813                let _ = session.cancel(source.corr.as_str());
1814            }
1815            self.drain_session(source.session.clone());
1816            if let Ok(session) = self.session_mut(&peer.session) {
1817                let _ = session.cancel(peer.corr.as_str());
1818            }
1819            self.drain_session(peer.session);
1820        }
1821    }
1822
1823    fn on_identify(&mut self, session: &SessionId, envelope: Envelope) {
1824        let Ok(remote) = envelope.parse_payload::<NodeIdentity>() else {
1825            self.retire(session.clone(), RetirementReason::MalformedIdentity);
1826            return;
1827        };
1828        let Some(peer) = self.peers.get_mut(session) else {
1829            return;
1830        };
1831        if peer.establishment.on_identify(remote.clone()).is_err() {
1832            self.retire(session.clone(), RetirementReason::DuplicateIdentity);
1833            return;
1834        }
1835        let effect = self.effect_id();
1836        self.pending
1837            .insert(effect, PendingEffect::PeerAdmission(session.clone()));
1838        self.effects.push_back(CoreEffect::RequestPeerAdmission {
1839            effect,
1840            session: session.clone(),
1841            remote,
1842        });
1843    }
1844
1845    fn complete_admission(
1846        &mut self,
1847        session: SessionId,
1848        result: PeerAdmission,
1849    ) -> Result<(), CoreError> {
1850        let PeerAdmission::Admitted(verified) = result else {
1851            self.retire(session, RetirementReason::AdmissionRejected);
1852            return Ok(());
1853        };
1854        let local_node = self.node.node().to_string();
1855        let Some(peer) = self.peers.get_mut(&session) else {
1856            return Ok(());
1857        };
1858        if peer.retired {
1859            return Ok(());
1860        }
1861        let Some(declared) = peer.establishment.remote_identity() else {
1862            return Ok(());
1863        };
1864        if verified.node_id != declared.node_id || verified.instance_id != declared.instance_id {
1865            self.retire(session, RetirementReason::AdmissionIdentityMismatch);
1866            return Ok(());
1867        }
1868        if verified.node_id == local_node {
1869            self.retire(session, RetirementReason::SelfConnection);
1870            return Ok(());
1871        }
1872        if peer
1873            .expected_peer
1874            .as_deref()
1875            .is_some_and(|expected| expected != verified.node_id)
1876        {
1877            self.retire(session, RetirementReason::UnexpectedPeer);
1878            return Ok(());
1879        }
1880        peer.establishment.local_accept()?;
1881        peer.remote = Some(verified);
1882        if !peer.initiator {
1883            self.send_identify(&session)?;
1884        }
1885        self.control(&session, Kind::IdentityAccepted, Bytes::new())?;
1886        self.maybe_send_snapshot(&session)?;
1887        Ok(())
1888    }
1889
1890    fn on_identity_accepted(&mut self, session: &SessionId) {
1891        let result = self
1892            .peers
1893            .get_mut(session)
1894            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1895            .and_then(|peer| peer.establishment.on_identity_accepted());
1896        if result.is_err() {
1897            self.retire(session.clone(), RetirementReason::IdentityAcceptanceOrder);
1898            return;
1899        }
1900        if self.maybe_send_snapshot(session).is_err() {
1901            self.retire(session.clone(), RetirementReason::SnapshotSequence);
1902            return;
1903        }
1904        self.try_promote(session.clone());
1905    }
1906
1907    fn on_initial_snapshot(&mut self, session: &SessionId, envelope: Envelope) {
1908        let Ok(snapshot) = envelope.parse_payload::<RouteSnapshot>() else {
1909            self.retire(session.clone(), RetirementReason::MalformedSnapshot);
1910            return;
1911        };
1912        let Some(remote) = self
1913            .peers
1914            .get(session)
1915            .and_then(|peer| peer.remote.as_ref())
1916            .cloned()
1917        else {
1918            self.retire(session.clone(), RetirementReason::SnapshotBeforeIdentity);
1919            return;
1920        };
1921        let mut probe = self.node.clone();
1922        if probe
1923            .apply_snapshot(session.as_str(), &remote.node_id, &snapshot)
1924            .is_err()
1925        {
1926            self.retire(session.clone(), RetirementReason::SnapshotRejected);
1927            return;
1928        }
1929        let Some(peer) = self.peers.get_mut(session) else {
1930            return;
1931        };
1932        if peer.establishment.on_snapshot_applied().is_err() {
1933            self.retire(session.clone(), RetirementReason::SnapshotOrder);
1934            return;
1935        }
1936        peer.held_snapshot = Some(snapshot.clone());
1937        let ack = RouteAck {
1938            generation: snapshot.generation,
1939            status: RouteAckStatus::Applied,
1940        };
1941        if peer.initiator {
1942            let _ = self.control(
1943                session,
1944                Kind::RouteAck,
1945                Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
1946            );
1947        } else {
1948            peer.deferred_ack = Some(ack);
1949        }
1950        self.try_promote(session.clone());
1951    }
1952
1953    fn on_initial_ack(&mut self, session: &SessionId, envelope: Envelope) {
1954        let Ok(ack) = envelope.parse_payload::<RouteAck>() else {
1955            self.retire(session.clone(), RetirementReason::MalformedAcknowledgement);
1956            return;
1957        };
1958        let result = self
1959            .peers
1960            .get_mut(session)
1961            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))
1962            .and_then(|peer| peer.establishment.on_route_ack(&ack));
1963        if result.is_err() {
1964            self.retire(session.clone(), RetirementReason::AcknowledgementOrder);
1965            return;
1966        }
1967        self.effects.push_back(CoreEffect::RouteExportAcked {
1968            session: session.clone(),
1969        });
1970        self.try_promote(session.clone());
1971    }
1972
1973    fn on_route_snapshot(&mut self, session: &SessionId, envelope: Envelope) {
1974        let Ok(snapshot) = envelope.parse_payload::<RouteSnapshot>() else {
1975            self.retire(session.clone(), RetirementReason::MalformedRouteControl);
1976            return;
1977        };
1978        let Some(peer) = self.peer_node(session) else {
1979            return;
1980        };
1981        match self.node.apply_snapshot(session.as_str(), &peer, &snapshot) {
1982            Ok(changed) => {
1983                self.route_ack(session, snapshot.generation, RouteAckStatus::Applied);
1984                if !changed.is_empty() {
1985                    let _ = self.recompute_route_exports();
1986                }
1987                self.effects.push_back(CoreEffect::RouteSnapshotApplied {
1988                    session: session.clone(),
1989                    peer,
1990                    snapshot,
1991                    changed: !changed.is_empty(),
1992                });
1993            }
1994            Err(RouteError::StaleUpdate { .. }) => {
1995                if let Some(generation) = self.node.applied_generation(session.as_str()) {
1996                    self.route_ack(session, generation, RouteAckStatus::Applied);
1997                }
1998            }
1999            Err(_) => self.retire(session.clone(), RetirementReason::RouteUpdateRejected),
2000        }
2001    }
2002
2003    fn on_route_delta(&mut self, session: &SessionId, envelope: Envelope) {
2004        let Ok(delta) = envelope.parse_payload::<RouteDelta>() else {
2005            self.retire(session.clone(), RetirementReason::MalformedRouteControl);
2006            return;
2007        };
2008        let Some(peer) = self.peer_node(session) else {
2009            return;
2010        };
2011        match self.node.apply_delta(session.as_str(), &peer, &delta) {
2012            Ok(changed) => {
2013                self.route_ack(session, delta.generation, RouteAckStatus::Applied);
2014                if !changed.is_empty() {
2015                    let _ = self.recompute_route_exports();
2016                }
2017                self.effects.push_back(CoreEffect::RouteDeltaApplied {
2018                    session: session.clone(),
2019                    peer,
2020                    delta,
2021                    changed: !changed.is_empty(),
2022                });
2023            }
2024            Err(RouteError::StaleUpdate { .. }) => {
2025                if let Some(generation) = self.node.applied_generation(session.as_str()) {
2026                    self.route_ack(session, generation, RouteAckStatus::Applied);
2027                }
2028            }
2029            Err(RouteError::GenerationGap { expected, .. }) => self.route_ack(
2030                session,
2031                expected.saturating_sub(1),
2032                RouteAckStatus::ResyncRequired,
2033            ),
2034            Err(_) => self.retire(session.clone(), RetirementReason::RouteUpdateRejected),
2035        }
2036    }
2037
2038    fn on_route_ack(&mut self, session: &SessionId, envelope: Envelope) {
2039        let Ok(ack) = envelope.parse_payload::<RouteAck>() else {
2040            self.retire(session.clone(), RetirementReason::MalformedAcknowledgement);
2041            return;
2042        };
2043        let Some(export) = self
2044            .peers
2045            .get_mut(session)
2046            .and_then(|peer| peer.export.as_mut())
2047        else {
2048            return;
2049        };
2050        if ack.generation > export.generation {
2051            self.retire(session.clone(), RetirementReason::AcknowledgementOrder);
2052            return;
2053        }
2054        match ack.status {
2055            RouteAckStatus::Applied => {
2056                export.applied_ack = export.applied_ack.max(ack.generation);
2057            }
2058            RouteAckStatus::ResyncRequired => {
2059                if ack.generation < export.applied_ack
2060                    || export
2061                        .handled_resync
2062                        .is_some_and(|handled| ack.generation <= handled)
2063                {
2064                    return;
2065                }
2066                export.handled_resync = Some(ack.generation);
2067                export.generation = export.generation.saturating_add(1);
2068                let snapshot = RouteSnapshot::canonical(export.generation, export.routes.clone());
2069                let _ = self.control(
2070                    session,
2071                    Kind::RouteSnapshot,
2072                    Envelope::encode_payload(&serde_json::to_value(snapshot).unwrap()),
2073                );
2074            }
2075        }
2076    }
2077
2078    fn on_route_exports_changed(
2079        &mut self,
2080        session: &SessionId,
2081        mut routes: Vec<RouteAdvertisement>,
2082    ) -> Result<(), CoreError> {
2083        routes.sort();
2084        let Some(export) = self
2085            .peers
2086            .get_mut(session)
2087            .and_then(|peer| peer.export.as_mut())
2088        else {
2089            return Ok(());
2090        };
2091        if routes == export.routes {
2092            return Ok(());
2093        }
2094        let previous: BTreeMap<&str, &RouteAdvertisement> = export
2095            .routes
2096            .iter()
2097            .map(|route| (route.destination.as_str(), route))
2098            .collect();
2099        let current: BTreeMap<&str, &RouteAdvertisement> = routes
2100            .iter()
2101            .map(|route| (route.destination.as_str(), route))
2102            .collect();
2103        let upsert = routes
2104            .iter()
2105            .filter(|route| previous.get(route.destination.as_str()) != Some(route))
2106            .cloned()
2107            .collect();
2108        let withdraw = export
2109            .routes
2110            .iter()
2111            .filter(|route| !current.contains_key(route.destination.as_str()))
2112            .map(|route| RouteWithdrawal {
2113                destination: route.destination.clone(),
2114                owner: route.owner.clone(),
2115                owner_instance: route.owner_instance.clone(),
2116                owner_epoch: route.owner_epoch,
2117                owner_revision: route.owner_revision,
2118            })
2119            .collect();
2120        export.generation = export.generation.saturating_add(1);
2121        let delta = RouteDelta {
2122            generation: export.generation,
2123            upsert,
2124            withdraw,
2125        };
2126        export.routes = routes;
2127        self.control(
2128            session,
2129            Kind::RouteDelta,
2130            Envelope::encode_payload(&serde_json::to_value(delta).unwrap()),
2131        )
2132    }
2133
2134    fn recompute_route_exports(&mut self) -> Result<(), CoreError> {
2135        let peers: Vec<_> = self
2136            .peers
2137            .iter()
2138            .filter_map(|(session, peer)| {
2139                peer.ready
2140                    .then(|| {
2141                        peer.remote
2142                            .as_ref()
2143                            .map(|remote| (session.clone(), remote.node_id.clone()))
2144                    })
2145                    .flatten()
2146            })
2147            .collect();
2148        for (session, peer) in peers {
2149            self.on_route_exports_changed(&session, self.node.export_for(&peer))?;
2150        }
2151        Ok(())
2152    }
2153
2154    fn peer_node(&self, session: &SessionId) -> Option<String> {
2155        self.peers
2156            .get(session)
2157            .and_then(|peer| peer.remote.as_ref())
2158            .map(|peer| peer.node_id.clone())
2159    }
2160
2161    fn route_ack(&mut self, session: &SessionId, generation: u64, status: RouteAckStatus) {
2162        let ack = RouteAck { generation, status };
2163        let _ = self.control(
2164            session,
2165            Kind::RouteAck,
2166            Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
2167        );
2168    }
2169
2170    fn send_identify(&mut self, session: &SessionId) -> Result<(), CoreError> {
2171        let peer = self
2172            .peers
2173            .get_mut(session)
2174            .ok_or_else(|| CoreError::UnknownSession(session.to_string()))?;
2175        peer.establishment.identify_sent()?;
2176        let payload =
2177            Envelope::encode_payload(&serde_json::to_value(self.node.identity()).unwrap());
2178        self.control(session, Kind::Identify, payload)
2179    }
2180
2181    fn maybe_send_snapshot(&mut self, session: &SessionId) -> Result<(), CoreError> {
2182        let Some(remote) = self
2183            .peers
2184            .get(session)
2185            .and_then(|peer| peer.remote.as_ref())
2186            .cloned()
2187        else {
2188            return Ok(());
2189        };
2190        if !self.peers[session].establishment.identities_accepted() {
2191            return Ok(());
2192        }
2193        let snapshot = RouteSnapshot::canonical(1, self.node.export_for(&remote.node_id));
2194        self.peers.get_mut(session).unwrap().export = Some(RouteExportState {
2195            generation: snapshot.generation,
2196            applied_ack: 0,
2197            handled_resync: None,
2198            routes: snapshot.routes.clone(),
2199        });
2200        self.peers
2201            .get_mut(session)
2202            .unwrap()
2203            .establishment
2204            .snapshot_sent(snapshot.generation)?;
2205        self.control(
2206            session,
2207            Kind::RouteSnapshot,
2208            Envelope::encode_payload(&serde_json::to_value(snapshot).unwrap()),
2209        )
2210    }
2211
2212    fn try_promote(&mut self, session: SessionId) {
2213        let Some(peer) = self.peers.get(&session) else {
2214            return;
2215        };
2216        if peer.ready || !peer.establishment.ready() {
2217            return;
2218        }
2219        let Some(remote) = peer.remote.clone() else {
2220            return;
2221        };
2222        if let Some(active) = self.active_peers.get(&remote.node_id).cloned() {
2223            if active != session {
2224                let old = &self.peers[&active];
2225                if old
2226                    .remote
2227                    .as_ref()
2228                    .is_some_and(|id| id.instance_id == remote.instance_id)
2229                {
2230                    let smaller_initiates = self.node.node() < remote.node_id.as_str();
2231                    let new_preferred = peer.initiator == smaller_initiates;
2232                    let old_preferred = old.initiator == smaller_initiates;
2233                    if !new_preferred || old_preferred {
2234                        self.retire(session, RetirementReason::DuplicateSession);
2235                        return;
2236                    }
2237                }
2238                self.retire(active, RetirementReason::DuplicateSessionReplaced);
2239            }
2240        }
2241        let peer = self.peers.get_mut(&session).unwrap();
2242        peer.ready = true;
2243        self.sessions.get_mut(&session).unwrap().mark_peer_ready();
2244        let snapshot = peer.held_snapshot.take();
2245        let deferred_ack = peer.deferred_ack.take();
2246        let changed = snapshot.as_ref().and_then(|snapshot| {
2247            self.node
2248                .apply_snapshot(session.as_str(), &remote.node_id, snapshot)
2249                .ok()
2250        });
2251        self.active_peers
2252            .insert(remote.node_id.clone(), session.clone());
2253        self.effects.push_back(CoreEffect::SessionEstablished {
2254            session: session.clone(),
2255            peer: remote.clone(),
2256        });
2257        let peer_node = remote.node_id.clone();
2258        let changed_routes = changed.as_ref().is_some_and(|changed| !changed.is_empty());
2259        if let (Some(snapshot), Some(_)) = (snapshot, changed) {
2260            self.effects.push_back(CoreEffect::RouteSnapshotApplied {
2261                session: session.clone(),
2262                peer: peer_node,
2263                snapshot,
2264                changed: changed_routes,
2265            });
2266        }
2267        if let Some(ack) = deferred_ack {
2268            let _ = self.control(
2269                &session,
2270                Kind::RouteAck,
2271                Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
2272            );
2273        }
2274        if changed_routes {
2275            let _ = self.recompute_route_exports();
2276        }
2277    }
2278
2279    fn retire(&mut self, session: SessionId, reason: RetirementReason) {
2280        let reason = self.converged_retirement(&session, reason);
2281        let mut withdrew_routes = false;
2282        let mut withdrew_session = false;
2283        if let Some(peer) = self.peers.get_mut(&session) {
2284            if peer.retired {
2285                return;
2286            }
2287            peer.retired = true;
2288            withdrew_session = peer.ready;
2289            peer.ready = false;
2290            if let Some(remote) = &peer.remote {
2291                if self.active_peers.get(&remote.node_id) == Some(&session) {
2292                    self.active_peers.remove(&remote.node_id);
2293                    withdrew_routes = !self.node.leave(session.as_str()).is_empty();
2294                }
2295            }
2296        }
2297        if withdrew_routes {
2298            let _ = self.recompute_route_exports();
2299        }
2300        if withdrew_session {
2301            self.effects.push_back(CoreEffect::RouteSessionWithdrawn {
2302                session: session.clone(),
2303                changed: withdrew_routes,
2304            });
2305        }
2306        self.abort_pending_dispatch(&session, None);
2307        self.pending.retain(|_, pending| match pending {
2308            PendingEffect::PeerAdmission(target) => target != &session,
2309            PendingEffect::Capacity(invocation) => invocation.stream.session != session,
2310            PendingEffect::Dispatch(stream) | PendingEffect::Send { stream, .. } => {
2311                stream.session != session
2312            }
2313            PendingEffect::RelayOpen { stream, .. } => stream.session != session,
2314            PendingEffect::TargetReadiness { stream, .. } => stream.session != session,
2315            PendingEffect::RelayForward { source, target, .. } => {
2316                source.session != session && target != &session
2317            }
2318        });
2319        self.discoveries
2320            .retain(|stream, _| stream.session != session);
2321        self.completed_client_operations
2322            .retain(|stream| stream.session != session);
2323        self.completed_client_order
2324            .retain(|stream| stream.session != session);
2325        let operations = self
2326            .client_operations
2327            .keys()
2328            .filter(|stream| stream.session == session)
2329            .cloned()
2330            .collect::<Vec<_>>();
2331        for stream in operations {
2332            self.client_operations.remove(&stream);
2333            self.effects.push_back(CoreEffect::DeliverClient {
2334                session: session.clone(),
2335                operation: ClientOperationId::from(stream.corr),
2336                delivery: ClientDelivery::SessionClosed,
2337            });
2338        }
2339        self.effects
2340            .push_back(CoreEffect::SessionRetired { session, reason });
2341    }
2342
2343    fn converged_retirement(
2344        &self,
2345        session: &SessionId,
2346        reason: RetirementReason,
2347    ) -> RetirementReason {
2348        if !matches!(reason, RetirementReason::SessionClosed) {
2349            return reason;
2350        }
2351        let Some(peer) = self.peers.get(session) else {
2352            return reason;
2353        };
2354        if peer.ready || peer.retired {
2355            return reason;
2356        }
2357        let Some(remote) = peer.remote.as_ref() else {
2358            return reason;
2359        };
2360        match self.active_peers.get(&remote.node_id) {
2361            Some(active)
2362                if active != session && self.peers.get(active).is_some_and(|peer| peer.ready) =>
2363            {
2364                RetirementReason::DuplicateSession
2365            }
2366            _ => reason,
2367        }
2368    }
2369
2370    fn fail_stream(
2371        &mut self,
2372        stream: &StreamKey,
2373        code: ErrorCode,
2374        message: &str,
2375    ) -> Result<(), CoreError> {
2376        self.session_mut(&stream.session)?
2377            .fail(stream.corr.as_str(), code, message)?;
2378        self.drain_session(stream.session.clone());
2379        Ok(())
2380    }
2381
2382    fn queue_send(
2383        &mut self,
2384        stream: StreamKey,
2385        kind: Kind,
2386        payload: Bytes,
2387    ) -> Result<(), CoreError> {
2388        let terminal = match kind {
2389            Kind::Response => self
2390                .session_mut(&stream.session)?
2391                .respond_discovery(stream.corr.as_str(), payload)
2392                .map(|()| true)?,
2393            Kind::Event => self
2394                .session_mut(&stream.session)?
2395                .send_discovery_event(stream.corr.as_str(), payload)
2396                .map(|()| false)?,
2397            _ => unreachable!(),
2398        };
2399        let outputs = {
2400            let session = self.session_mut(&stream.session)?;
2401            let mut outputs = Vec::new();
2402            loop {
2403                let output = session.poll_effect();
2404                if matches!(output, SessionEffect::Deadline(_)) {
2405                    break;
2406                }
2407                outputs.push(output);
2408            }
2409            outputs
2410        };
2411        for output in outputs {
2412            if let SessionEffect::SendFrame(envelope) = output {
2413                let effect = self.effect_id();
2414                self.pending.insert(
2415                    effect,
2416                    PendingEffect::Send {
2417                        stream: stream.clone(),
2418                        terminal,
2419                    },
2420                );
2421                self.effects.push_back(CoreEffect::SendProtocol {
2422                    effect,
2423                    session: stream.session.clone(),
2424                    envelope,
2425                });
2426            } else {
2427                match output {
2428                    SessionEffect::Deliver(envelope) => {
2429                        self.deliver_client(stream.session.clone(), envelope)
2430                    }
2431                    SessionEffect::StreamClosed { corr } => {
2432                        self.stream_closed(stream.session.clone(), corr)
2433                    }
2434                    SessionEffect::Closed { code, message } => {
2435                        self.effects.push_back(CoreEffect::CloseTransport {
2436                            session: stream.session.clone(),
2437                            code,
2438                            message,
2439                        })
2440                    }
2441                    SessionEffect::HandshakeEstablished { version } => {
2442                        self.effects.push_back(CoreEffect::HandshakeEstablished {
2443                            session: stream.session.clone(),
2444                            version,
2445                        })
2446                    }
2447                    SessionEffect::Deadline(Some(deadline)) => {
2448                        self.effects.push_back(CoreEffect::ScheduleSessionDeadline {
2449                            session: stream.session.clone(),
2450                            deadline,
2451                        })
2452                    }
2453                    SessionEffect::Deadline(None) | SessionEffect::SendFrame(_) => {}
2454                }
2455            }
2456        }
2457        Ok(())
2458    }
2459
2460    fn queue_application_response(
2461        &mut self,
2462        stream: StreamKey,
2463        response: ApplicationResponse,
2464        terminal: bool,
2465    ) -> Result<(), CoreError> {
2466        let (parts, ()) = response.head.into_parts();
2467        let mut envelope =
2468            Envelope::from_response(http::Response::from_parts(parts, Bytes::new()))?;
2469        envelope.body_token = response.body.map(|body| body.to_string());
2470        let body = envelope.body_token.take();
2471        if terminal {
2472            self.session_mut(&stream.session)?.respond_terminal_body(
2473                stream.corr.as_str(),
2474                body,
2475                envelope.path,
2476                envelope.headers,
2477            )?;
2478        } else {
2479            let session = self.session_mut(&stream.session)?;
2480            session.send_body(stream.corr.as_str(), body, envelope.headers)?;
2481        }
2482        self.drain_send_outputs(stream, terminal)
2483    }
2484
2485    fn drain_send_outputs(&mut self, stream: StreamKey, terminal: bool) -> Result<(), CoreError> {
2486        let outputs = {
2487            let session = self.session_mut(&stream.session)?;
2488            let mut outputs = Vec::new();
2489            loop {
2490                let output = session.poll_effect();
2491                if matches!(output, SessionEffect::Deadline(_)) {
2492                    break;
2493                }
2494                outputs.push(output);
2495            }
2496            outputs
2497        };
2498        for output in outputs {
2499            if let SessionEffect::SendFrame(envelope) = output {
2500                let effect = self.effect_id();
2501                self.pending.insert(
2502                    effect,
2503                    PendingEffect::Send {
2504                        stream: stream.clone(),
2505                        terminal,
2506                    },
2507                );
2508                let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
2509                    return Err(CoreError::Malformed(
2510                        "application send frame is malformed".into(),
2511                    ));
2512                };
2513                self.effects.push_back(CoreEffect::Send {
2514                    effect,
2515                    session: stream.session.clone(),
2516                    frame,
2517                });
2518            } else {
2519                self.queue_session_effect(stream.session.clone(), output);
2520            }
2521        }
2522        Ok(())
2523    }
2524
2525    fn application_origin(&self, session: &SessionId) -> ApplicationOrigin {
2526        self.peers
2527            .get(session)
2528            .and_then(|peer| peer.ready.then(|| peer.remote.clone()).flatten())
2529            .map_or_else(
2530                || ApplicationOrigin::Client {
2531                    session: session.clone(),
2532                },
2533                |peer| ApplicationOrigin::Peer {
2534                    session: session.clone(),
2535                    peer,
2536                },
2537            )
2538    }
2539
2540    fn deliver_client(&mut self, session: SessionId, envelope: Envelope) {
2541        if let Some(corr) = envelope.corr.clone() {
2542            let stream = StreamKey {
2543                session: session.clone(),
2544                corr: CorrelationId::from(corr.clone()),
2545            };
2546            if self.completed_client_operations.contains(&stream) {
2547                self.release_body(
2548                    &session,
2549                    envelope.body_token.as_deref().map(crate::BodyId::from),
2550                );
2551                return;
2552            }
2553            let terminal = matches!(envelope.kind, Kind::Response | Kind::Error | Kind::Cancel);
2554            if !self.client_operations.contains_key(&stream) {
2555                self.effects
2556                    .push_back(CoreEffect::Deliver { session, envelope });
2557                return;
2558            }
2559            if terminal {
2560                self.client_operations.remove(&stream);
2561                self.tombstone_client_operation(stream);
2562            }
2563            let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
2564                self.effects
2565                    .push_back(CoreEffect::Deliver { session, envelope });
2566                return;
2567            };
2568            self.effects.push_back(CoreEffect::DeliverClient {
2569                session,
2570                operation: ClientOperationId::from(corr),
2571                delivery: if terminal {
2572                    ClientDelivery::Terminal(frame)
2573                } else {
2574                    ClientDelivery::Item(frame)
2575                },
2576            });
2577        } else {
2578            self.effects
2579                .push_back(CoreEffect::Deliver { session, envelope });
2580        }
2581    }
2582
2583    fn queue_session_effect(&mut self, session: SessionId, effect: SessionEffect) {
2584        let effect = match effect {
2585            SessionEffect::SendFrame(envelope) => CoreEffect::SendFrame { session, envelope },
2586            SessionEffect::HandshakeEstablished { version } => {
2587                CoreEffect::HandshakeEstablished { session, version }
2588            }
2589            SessionEffect::Deliver(envelope) => {
2590                self.deliver_client(session, envelope);
2591                return;
2592            }
2593            SessionEffect::StreamClosed { corr } => {
2594                self.stream_closed(session, corr);
2595                return;
2596            }
2597            SessionEffect::Closed { code, message } => CoreEffect::CloseTransport {
2598                session,
2599                code,
2600                message,
2601            },
2602            SessionEffect::Deadline(Some(deadline)) => {
2603                CoreEffect::ScheduleSessionDeadline { session, deadline }
2604            }
2605            SessionEffect::Deadline(None) => return,
2606        };
2607        self.effects.push_back(effect);
2608    }
2609
2610    fn complete_client_operation(
2611        &mut self,
2612        session: &SessionId,
2613        operation: &ClientOperationId,
2614        delivery: ClientDelivery,
2615        cancel: bool,
2616    ) {
2617        let stream = StreamKey {
2618            session: session.clone(),
2619            corr: operation.0.clone(),
2620        };
2621        if self.client_operations.remove(&stream).is_none() {
2622            return;
2623        }
2624        self.tombstone_client_operation(stream);
2625        if cancel {
2626            let _ = self
2627                .session_mut(session)
2628                .and_then(|core| core.cancel(operation.as_str()));
2629            self.drain_session(session.clone());
2630        }
2631        self.effects.push_back(CoreEffect::DeliverClient {
2632            session: session.clone(),
2633            operation: operation.clone(),
2634            delivery,
2635        });
2636    }
2637
2638    fn abort_pending_dispatch(&mut self, session: &SessionId, corr: Option<&str>) {
2639        let aborted: Vec<(EffectId, Option<crate::BodyId>)> = self
2640            .pending
2641            .iter()
2642            .filter_map(|(effect, entry)| match entry {
2643                PendingEffect::Capacity(invocation)
2644                    if invocation.stream.session == *session
2645                        && corr.is_none_or(|corr| invocation.stream.corr.as_str() == corr) =>
2646                {
2647                    Some((*effect, invocation.frame.body.clone()))
2648                }
2649                PendingEffect::Dispatch(stream)
2650                    if stream.session == *session
2651                        && corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
2652                {
2653                    Some((*effect, None))
2654                }
2655                PendingEffect::RelayOpen { stream, body }
2656                    if stream.session == *session
2657                        && corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
2658                {
2659                    Some((*effect, body.clone()))
2660                }
2661                PendingEffect::TargetReadiness { stream, frame }
2662                    if stream.session == *session
2663                        && corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
2664                {
2665                    Some((*effect, frame.body.clone()))
2666                }
2667                _ => None,
2668            })
2669            .collect();
2670        for (effect, body) in aborted {
2671            self.pending.remove(&effect);
2672            self.release_body(session, body);
2673            self.effects.push_back(CoreEffect::AbortDispatch {
2674                session: session.clone(),
2675                effect,
2676            });
2677        }
2678    }
2679
2680    fn release_body(&mut self, session: &SessionId, body: Option<crate::BodyId>) {
2681        if let Some(body) = body {
2682            self.effects.push_back(CoreEffect::ReleaseBody {
2683                session: session.clone(),
2684                body,
2685            });
2686        }
2687    }
2688
2689    fn tombstone_client_operation(&mut self, stream: StreamKey) {
2690        if !self.completed_client_operations.insert(stream.clone()) {
2691            return;
2692        }
2693        self.completed_client_order.push_back(stream);
2694        while self.completed_client_order.len() > CLIENT_TOMBSTONE_LIMIT {
2695            if let Some(oldest) = self.completed_client_order.pop_front() {
2696                self.completed_client_operations.remove(&oldest);
2697            }
2698        }
2699    }
2700
2701    fn stream_closed(&mut self, session: SessionId, corr: String) {
2702        self.abort_pending_dispatch(&session, Some(&corr));
2703        let stream = StreamKey {
2704            session: session.clone(),
2705            corr: CorrelationId::from(corr.clone()),
2706        };
2707        self.discoveries.remove(&stream);
2708        self.target_discoveries.remove(&stream);
2709        if self.completed_client_operations.contains(&stream) {
2710            return;
2711        }
2712        if self.client_operations.remove(&stream).is_some() {
2713            self.tombstone_client_operation(stream);
2714            self.effects.push_back(CoreEffect::DeliverClient {
2715                session,
2716                operation: ClientOperationId::from(corr),
2717                delivery: ClientDelivery::Cancelled,
2718            });
2719            return;
2720        }
2721        self.effects.push_back(CoreEffect::StreamClosed {
2722            session,
2723            operation: ClientOperationId::from(corr),
2724        });
2725    }
2726}
2727
2728#[cfg(test)]
2729mod tests {
2730    use super::*;
2731    use crate::ApplicationHead;
2732
2733    fn established_client_core() -> (ProtocolCore, SessionId) {
2734        let mut core = ProtocolCore::new("node");
2735        let session = SessionId::from("client");
2736        let now = web_time::Instant::now();
2737        core.handle(
2738            now,
2739            CoreInput::SessionOpened {
2740                session: session.clone(),
2741                initiator: true,
2742                establish_peer: false,
2743                expected_peer: None,
2744            },
2745        )
2746        .unwrap();
2747        while core.poll_effect().is_some() {}
2748        core.handle(
2749            now,
2750            CoreInput::FrameReceived {
2751                session: session.clone(),
2752                envelope: Envelope {
2753                    v: crate::PROTOCOL_VERSION,
2754                    id: "welcome".into(),
2755                    target: String::new(),
2756                    subject: String::new(),
2757                    kind: Kind::Welcome,
2758                    corr: None,
2759                    seq: None,
2760                    hops: None,
2761                    body_token: None,
2762                    payload: Envelope::encode_payload(&serde_json::json!({"version": 1})),
2763                    path: Vec::new(),
2764                    headers: Default::default(),
2765                },
2766            },
2767        )
2768        .unwrap();
2769        while core.poll_effect().is_some() {}
2770        (core, session)
2771    }
2772
2773    #[test]
2774    fn retiring_session_removes_completed_client_operation_tombstones() {
2775        let session = SessionId::from("session");
2776        let stream = StreamKey {
2777            session: session.clone(),
2778            corr: CorrelationId::from("corr"),
2779        };
2780        let mut core = ProtocolCore::new("node");
2781        core.completed_client_operations.insert(stream);
2782        core.peers
2783            .insert(session.clone(), PeerSession::new(false, false, None));
2784
2785        core.retire(session, RetirementReason::SessionClosed);
2786
2787        assert!(core.completed_client_operations.is_empty());
2788    }
2789
2790    #[test]
2791    fn client_operation_tombstones_are_bounded() {
2792        let mut core = ProtocolCore::new("node");
2793        for index in 0..(CLIENT_TOMBSTONE_LIMIT + 10) {
2794            core.tombstone_client_operation(StreamKey {
2795                session: SessionId::from("session"),
2796                corr: CorrelationId::from(format!("corr-{index}")),
2797            });
2798        }
2799
2800        assert_eq!(
2801            core.completed_client_operations.len(),
2802            CLIENT_TOMBSTONE_LIMIT
2803        );
2804        assert_eq!(core.completed_client_order.len(), CLIENT_TOMBSTONE_LIMIT);
2805        assert!(!core.completed_client_operations.contains(&StreamKey {
2806            session: SessionId::from("session"),
2807            corr: CorrelationId::from("corr-0"),
2808        }));
2809    }
2810
2811    #[test]
2812    fn operation_stream_fault_is_scoped_and_late_duplicates_are_idempotent() {
2813        let (mut core, session) = established_client_core();
2814        let now = web_time::Instant::now();
2815        core.handle(
2816            now,
2817            CoreInput::StartClientOperation {
2818                session: session.clone(),
2819                target_path: "/node/echo".into(),
2820                kind: Kind::Request,
2821                input: OperationInput::Body(None),
2822                hops: None,
2823                headers: Default::default(),
2824                timeout: None,
2825            },
2826        )
2827        .unwrap();
2828        let mut corr = None;
2829        while let Some(effect) = core.poll_effect() {
2830            if let CoreEffect::Send { frame, .. } = effect {
2831                corr = frame.head.corr;
2832            }
2833        }
2834        let corr = CorrelationId::from(corr.unwrap());
2835        let fault = CoreInput::OperationStreamEnded {
2836            session: session.clone(),
2837            corr: corr.clone(),
2838            direction: OperationStreamDirection::Return,
2839            outcome: OperationStreamOutcome::Truncated,
2840        };
2841        core.handle(now, fault.clone()).unwrap();
2842
2843        let effects = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
2844        assert!(effects.iter().any(|effect| matches!(
2845            effect,
2846            CoreEffect::DeliverClient {
2847                delivery: ClientDelivery::Terminal(frame),
2848                ..
2849            } if frame.head.error.as_ref().is_some_and(|error| error.code == ErrorCode::Protocol)
2850        )));
2851        assert!(!effects
2852            .iter()
2853            .any(|effect| matches!(effect, CoreEffect::SessionRetired { .. })));
2854
2855        core.handle(now, fault).unwrap();
2856        let late = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
2857        assert!(!late.iter().any(|effect| matches!(
2858            effect,
2859            CoreEffect::DeliverClient { .. } | CoreEffect::SessionRetired { .. }
2860        )));
2861    }
2862
2863    #[test]
2864    fn terminal_response_wins_over_a_late_body_fault() {
2865        let (mut core, session) = established_client_core();
2866        let now = web_time::Instant::now();
2867        core.handle(
2868            now,
2869            CoreInput::StartClientOperation {
2870                session: session.clone(),
2871                target_path: "/node/echo".into(),
2872                kind: Kind::Request,
2873                input: OperationInput::Body(None),
2874                hops: None,
2875                headers: Default::default(),
2876                timeout: None,
2877            },
2878        )
2879        .unwrap();
2880        let corr = std::iter::from_fn(|| core.poll_effect())
2881            .find_map(|effect| match effect {
2882                CoreEffect::Send { frame, .. } => frame.head.corr,
2883                _ => None,
2884            })
2885            .unwrap();
2886        core.handle(
2887            now,
2888            CoreInput::ApplicationFrameReceived {
2889                session: session.clone(),
2890                frame: ApplicationFrame {
2891                    head: ApplicationHead {
2892                        kind: Kind::Response,
2893                        corr: Some(corr.clone()),
2894                        ..ApplicationHead::default()
2895                    },
2896                    body: Some(crate::BodyId::from("response-body")),
2897                },
2898            },
2899        )
2900        .unwrap();
2901        let terminal = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
2902        assert!(terminal.iter().any(|effect| matches!(
2903            effect,
2904            CoreEffect::DeliverClient {
2905                delivery: ClientDelivery::Terminal(_),
2906                ..
2907            }
2908        )));
2909
2910        core.handle(
2911            now,
2912            CoreInput::OperationStreamEnded {
2913                session,
2914                corr: CorrelationId::from(corr),
2915                direction: OperationStreamDirection::Return,
2916                outcome: OperationStreamOutcome::Truncated,
2917            },
2918        )
2919        .unwrap();
2920        assert!(core.poll_effect().is_none());
2921    }
2922
2923    #[test]
2924    fn cancelled_or_timed_out_operation_releases_a_late_response_body() {
2925        let (mut core, session) = established_client_core();
2926        let now = web_time::Instant::now();
2927        core.handle(
2928            now,
2929            CoreInput::StartClientOperation {
2930                session: session.clone(),
2931                target_path: "/node/echo".into(),
2932                kind: Kind::Request,
2933                input: OperationInput::Body(None),
2934                hops: None,
2935                headers: Default::default(),
2936                timeout: None,
2937            },
2938        )
2939        .unwrap();
2940        let corr = std::iter::from_fn(|| core.poll_effect())
2941            .find_map(|effect| match effect {
2942                CoreEffect::Send { frame, .. } => frame.head.corr,
2943                _ => None,
2944            })
2945            .unwrap();
2946        core.handle(
2947            now,
2948            CoreInput::CancelClientOperation {
2949                session: session.clone(),
2950                operation: ClientOperationId::from(corr.clone()),
2951            },
2952        )
2953        .unwrap();
2954        while core.poll_effect().is_some() {}
2955
2956        core.handle(
2957            now,
2958            CoreInput::ApplicationFrameReceived {
2959                session: session.clone(),
2960                frame: ApplicationFrame {
2961                    head: ApplicationHead {
2962                        kind: Kind::Response,
2963                        corr: Some(corr),
2964                        ..ApplicationHead::default()
2965                    },
2966                    body: Some(crate::BodyId::from("late-body")),
2967                },
2968            },
2969        )
2970        .unwrap();
2971        assert!(matches!(
2972            core.poll_effect(),
2973            Some(CoreEffect::ReleaseBody { session: owner, body })
2974                if owner == session && body.as_str() == "late-body"
2975        ));
2976        assert!(core.poll_effect().is_none());
2977
2978        core.handle(
2979            now,
2980            CoreInput::StartClientOperation {
2981                session: session.clone(),
2982                target_path: "/node/echo".into(),
2983                kind: Kind::Request,
2984                input: OperationInput::Body(None),
2985                hops: None,
2986                headers: Default::default(),
2987                timeout: Some(std::time::Duration::ZERO),
2988            },
2989        )
2990        .unwrap();
2991        let timed_out_corr = std::iter::from_fn(|| core.poll_effect())
2992            .find_map(|effect| match effect {
2993                CoreEffect::Send { frame, .. } => frame.head.corr,
2994                _ => None,
2995            })
2996            .unwrap();
2997        core.handle(
2998            now,
2999            CoreInput::ClientOperationTimeout {
3000                session: session.clone(),
3001                operation: ClientOperationId::from(timed_out_corr.clone()),
3002            },
3003        )
3004        .unwrap();
3005        while core.poll_effect().is_some() {}
3006        core.handle(
3007            now,
3008            CoreInput::ApplicationFrameReceived {
3009                session: session.clone(),
3010                frame: ApplicationFrame {
3011                    head: ApplicationHead {
3012                        kind: Kind::Response,
3013                        corr: Some(timed_out_corr),
3014                        ..ApplicationHead::default()
3015                    },
3016                    body: Some(crate::BodyId::from("timed-out-body")),
3017                },
3018            },
3019        )
3020        .unwrap();
3021        assert!(matches!(
3022            core.poll_effect(),
3023            Some(CoreEffect::ReleaseBody { session: owner, body })
3024                if owner == session && body.as_str() == "timed-out-body"
3025        ));
3026        assert!(core.poll_effect().is_none());
3027    }
3028}