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}