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