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