use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
use std::fmt;
use bytes::Bytes;
use web_time::Instant;
use crate::relay_reducer::RelayReducer;
use crate::route_control::Establishment;
use crate::session::{SessionEffect, SessionInput, SessionReducer};
use crate::{
ApplicationFrame, CoreError, DiscoverEvent, DiscoverPlan, DiscoverWalk, Envelope, ErrorCode,
Kind, NodeCore, NodeIdentity, RouteAck, RouteAckStatus, RouteAdvertisement, RouteDelta,
RouteError, RouteSnapshot, RouteWithdrawal, SessionClass, TargetPath, WalkInput, WalkOutput,
};
const MAX_EFFECTS: usize = 64;
const MAX_EFFECTS_PER_INPUT: usize = 4;
const MAX_DISCOVERY_EFFECTS_PER_INPUT: usize = 4;
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct SessionId(String);
impl SessionId {
pub fn as_str(&self) -> &str {
&self.0
}
}
impl From<&str> for SessionId {
fn from(value: &str) -> Self {
Self(value.to_owned())
}
}
impl From<String> for SessionId {
fn from(value: String) -> Self {
Self(value)
}
}
impl fmt::Display for SessionId {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(formatter)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct CorrelationId(String);
impl CorrelationId {
pub fn as_str(&self) -> &str {
&self.0
}
}
impl From<&str> for CorrelationId {
fn from(value: &str) -> Self {
Self(value.to_owned())
}
}
impl From<String> for CorrelationId {
fn from(value: String) -> Self {
Self(value)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct ClientOperationId(CorrelationId);
impl ClientOperationId {
pub fn as_str(&self) -> &str {
self.0.as_str()
}
}
impl From<CorrelationId> for ClientOperationId {
fn from(value: CorrelationId) -> Self {
Self(value)
}
}
impl From<&str> for ClientOperationId {
fn from(value: &str) -> Self {
Self(CorrelationId::from(value))
}
}
impl From<String> for ClientOperationId {
fn from(value: String) -> Self {
Self(CorrelationId::from(value))
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum ClientDelivery {
Item(ApplicationFrame),
Terminal(ApplicationFrame),
Cancelled,
TimedOut,
SessionClosed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct EffectId(u64);
impl EffectId {
pub fn new(value: u64) -> Self {
Self(value)
}
pub fn get(self) -> u64 {
self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct StreamKey {
pub session: SessionId,
pub corr: CorrelationId,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ApplicationInvocation {
pub stream: StreamKey,
pub reservation: EffectId,
pub origin: ApplicationOrigin,
pub frame: ApplicationFrame,
}
#[derive(Debug, Clone, PartialEq)]
pub enum ApplicationOrigin {
Client {
session: SessionId,
},
Peer {
session: SessionId,
peer: NodeIdentity,
},
}
#[derive(Debug, Clone, PartialEq)]
pub enum PeerAdmission {
Admitted(NodeIdentity),
Rejected(String),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CapacityResult {
Available,
Busy,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RelayOpenResult {
Opened(StreamKey),
Failed(ApplicationFailure),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TargetReadinessResult {
Ready,
Unavailable { message: String },
}
#[derive(Debug, Clone)]
pub struct ApplicationResponse {
pub head: http::Response<()>,
pub body: Option<crate::BodyId>,
}
impl PartialEq for ApplicationResponse {
fn eq(&self, other: &Self) -> bool {
self.head.status() == other.head.status()
&& self.head.version() == other.head.version()
&& self.head.headers() == other.head.headers()
&& self.body == other.body
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum ApplicationResult {
Response(ApplicationResponse),
Event(ApplicationResponse),
Finished(ApplicationResponse),
Bridged,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApplicationFailure {
pub code: ErrorCode,
pub message: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SendResult {
Reserved,
Written,
ReservationTimedOut,
Closed,
Cancelled,
WriteFailed(String),
Refused { code: ErrorCode, message: String },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OperationStreamDirection {
Opening,
Return,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperationStreamOutcome {
Clean,
Cancelled,
TimedOut,
Busy,
PayloadTooLarge(String),
Truncated,
Protocol(String),
Transport(String),
}
#[derive(Debug, Clone, PartialEq)]
pub enum OperationInput {
Body(Option<crate::BodyId>),
Discovery(DiscoverPlan),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RetirementReason {
EstablishmentTimeout,
SessionClosed,
MalformedIdentity,
DuplicateIdentity,
AdmissionRejected,
AdmissionIdentityMismatch,
SelfConnection,
UnexpectedPeer,
IdentityAcceptanceOrder,
SnapshotSequence,
MalformedSnapshot,
SnapshotBeforeIdentity,
SnapshotRejected,
SnapshotOrder,
MalformedAcknowledgement,
AcknowledgementOrder,
MalformedRouteControl,
RouteUpdateRejected,
DuplicateSession,
DuplicateSessionReplaced,
SendClosed,
SendTransport,
TransportFailed,
}
#[derive(Debug, Clone, PartialEq)]
pub enum CoreInput {
SessionOpened {
session: SessionId,
initiator: bool,
establish_peer: bool,
expected_peer: Option<String>,
},
FrameReceived {
session: SessionId,
envelope: Envelope,
},
ApplicationFrameReceived {
session: SessionId,
frame: ApplicationFrame,
},
OperationStreamEnded {
session: SessionId,
corr: CorrelationId,
direction: OperationStreamDirection,
outcome: OperationStreamOutcome,
},
StartClientOperation {
session: SessionId,
target_path: String,
kind: Kind,
input: OperationInput,
hops: Option<u8>,
headers: serde_json::Map<String, serde_json::Value>,
timeout: Option<std::time::Duration>,
},
CancelClientOperation {
session: SessionId,
operation: ClientOperationId,
},
ClientOperationTimeout {
session: SessionId,
operation: ClientOperationId,
},
SessionTimeout {
session: SessionId,
},
PeerAdmissionCompleted {
effect: EffectId,
result: PeerAdmission,
},
CapacityChecked {
effect: EffectId,
result: CapacityResult,
},
DispatchCompleted {
effect: EffectId,
result: Result<ApplicationResult, ApplicationFailure>,
},
RelayOpenCompleted {
effect: EffectId,
result: RelayOpenResult,
},
TargetReadinessCompleted {
effect: EffectId,
result: TargetReadinessResult,
},
RelayForwardCompleted {
effect: EffectId,
result: SendResult,
},
SendCompleted {
effect: EffectId,
result: SendResult,
},
EstablishmentTimeout {
session: SessionId,
},
DiscoveryNeighborEvent {
stream: StreamKey,
peer: String,
event: DiscoverEvent,
},
DiscoveryNeighborDone {
stream: StreamKey,
peer: String,
},
DiscoveryNeighborTimeout {
stream: StreamKey,
peer: String,
},
DiscoveryTargetEvent {
stream: StreamKey,
event: DiscoverEvent,
},
DiscoveryTargetFailed {
stream: StreamKey,
message: String,
},
SessionClosed {
session: SessionId,
},
SessionFailed {
session: SessionId,
},
ContinueDiscovery {
stream: StreamKey,
},
LocalCapabilitiesInstalled {
capabilities: BTreeMap<String, serde_json::Value>,
},
}
#[derive(Debug, Clone, PartialEq)]
pub enum CoreEffect {
SendFrame {
session: SessionId,
envelope: Envelope,
},
HandshakeEstablished {
session: SessionId,
version: u16,
},
DeliverClient {
session: SessionId,
operation: ClientOperationId,
delivery: ClientDelivery,
},
RegisterClient {
session: SessionId,
operation: ClientOperationId,
},
ScheduleClientDeadline {
session: SessionId,
operation: ClientOperationId,
deadline: Instant,
},
Deliver {
session: SessionId,
envelope: Envelope,
},
ReleaseBody {
session: SessionId,
body: crate::BodyId,
},
StreamClosed {
session: SessionId,
operation: ClientOperationId,
},
CloseTransport {
session: SessionId,
code: ErrorCode,
message: String,
},
ScheduleSessionDeadline {
session: SessionId,
deadline: Instant,
},
RequestPeerAdmission {
effect: EffectId,
session: SessionId,
remote: NodeIdentity,
},
AbortDispatch {
session: SessionId,
effect: EffectId,
},
CheckDispatchCapacity {
effect: EffectId,
stream: StreamKey,
frame: ApplicationFrame,
},
InvokeApplication {
effect: EffectId,
invocation: ApplicationInvocation,
},
OpenRelay {
effect: EffectId,
source: StreamKey,
peer: String,
frame: ApplicationFrame,
},
AwaitTargetReadiness {
effect: EffectId,
stream: StreamKey,
target: String,
},
ForwardRelay {
effect: EffectId,
source: StreamKey,
target: StreamKey,
frame: ApplicationFrame,
terminal: bool,
},
Send {
effect: EffectId,
session: SessionId,
frame: ApplicationFrame,
},
SendProtocol {
effect: EffectId,
session: SessionId,
envelope: Envelope,
},
SessionEstablished {
session: SessionId,
peer: NodeIdentity,
},
SessionRetired {
session: SessionId,
reason: RetirementReason,
},
RouteSnapshotApplied {
session: SessionId,
peer: String,
snapshot: RouteSnapshot,
changed: bool,
},
RouteDeltaApplied {
session: SessionId,
peer: String,
delta: RouteDelta,
changed: bool,
},
RouteSessionWithdrawn {
session: SessionId,
changed: bool,
},
RouteExportAcked {
session: SessionId,
},
QueryDiscoveryNeighbor {
stream: StreamKey,
peer: String,
plan: DiscoverPlan,
},
QueryDiscoveryTarget {
stream: StreamKey,
peer: String,
target_path: String,
plan: DiscoverPlan,
},
ContinueDiscovery {
stream: StreamKey,
},
}
impl CoreEffect {
pub fn session(&self) -> &SessionId {
match self {
Self::SendFrame { session, .. }
| Self::HandshakeEstablished { session, .. }
| Self::DeliverClient { session, .. }
| Self::RegisterClient { session, .. }
| Self::ScheduleClientDeadline { session, .. }
| Self::Deliver { session, .. }
| Self::ReleaseBody { session, .. }
| Self::StreamClosed { session, .. }
| Self::CloseTransport { session, .. }
| Self::ScheduleSessionDeadline { session, .. }
| Self::RequestPeerAdmission { session, .. }
| Self::AbortDispatch { session, .. }
| Self::Send { session, .. }
| Self::SendProtocol { session, .. }
| Self::SessionEstablished { session, .. }
| Self::SessionRetired { session, .. }
| Self::RouteSnapshotApplied { session, .. }
| Self::RouteDeltaApplied { session, .. }
| Self::RouteSessionWithdrawn { session, .. }
| Self::RouteExportAcked { session } => session,
Self::CheckDispatchCapacity { stream, .. }
| Self::InvokeApplication {
invocation: ApplicationInvocation { stream, .. },
..
}
| Self::OpenRelay { source: stream, .. }
| Self::AwaitTargetReadiness { stream, .. }
| Self::ForwardRelay { target: stream, .. }
| Self::QueryDiscoveryNeighbor { stream, .. }
| Self::QueryDiscoveryTarget { stream, .. }
| Self::ContinueDiscovery { stream } => &stream.session,
}
}
}
enum PendingEffect {
PeerAdmission(SessionId),
Capacity(Box<ApplicationInvocation>),
Dispatch(StreamKey),
RelayOpen {
stream: StreamKey,
body: Option<crate::BodyId>,
},
TargetReadiness {
stream: StreamKey,
frame: ApplicationFrame,
},
RelayForward {
source: StreamKey,
target: SessionId,
terminal: bool,
},
Send {
stream: StreamKey,
terminal: bool,
},
}
struct PeerSession {
initiator: bool,
establish_peer: bool,
expected_peer: Option<String>,
establishment: Establishment,
remote: Option<NodeIdentity>,
held_snapshot: Option<RouteSnapshot>,
deferred_ack: Option<RouteAck>,
ready: bool,
retired: bool,
export: Option<RouteExportState>,
}
struct RouteExportState {
generation: u64,
applied_ack: u64,
handled_resync: Option<u64>,
routes: Vec<RouteAdvertisement>,
}
impl PeerSession {
fn new(initiator: bool, establish_peer: bool, expected_peer: Option<String>) -> Self {
Self {
initiator,
establish_peer,
expected_peer,
establishment: Establishment::default(),
remote: None,
held_snapshot: None,
deferred_ack: None,
ready: false,
retired: false,
export: None,
}
}
}
pub(crate) const CLIENT_TOMBSTONE_LIMIT: usize = 4096;
pub struct ProtocolCore {
node: NodeCore,
sessions: HashMap<SessionId, SessionReducer>,
peers: HashMap<SessionId, PeerSession>,
active_peers: HashMap<String, SessionId>,
discoveries: HashMap<StreamKey, DiscoverWalk>,
target_discoveries: HashSet<StreamKey>,
client_operations: HashMap<StreamKey, Option<Instant>>,
completed_client_operations: HashSet<StreamKey>,
completed_client_order: VecDeque<StreamKey>,
relays: RelayReducer,
effects: VecDeque<CoreEffect>,
pending: HashMap<EffectId, PendingEffect>,
next_effect: u64,
}
impl ProtocolCore {
pub fn new(node: &str) -> Self {
let identity = NodeIdentity {
node_id: node.to_string(),
instance_id: node.to_string(),
epoch: 0,
proof: serde_json::Value::Null,
};
Self::with_identity(identity)
}
#[cfg(feature = "fuzzing")]
pub fn effect_len(&self) -> usize {
self.effects.len()
}
pub fn with_identity(identity: NodeIdentity) -> Self {
let mut node = NodeCore::new(&identity.node_id);
node.set_node_identity(identity);
Self::with_node(node)
}
pub fn with_node(node: NodeCore) -> Self {
Self {
node,
sessions: HashMap::new(),
peers: HashMap::new(),
active_peers: HashMap::new(),
discoveries: HashMap::new(),
target_discoveries: HashSet::new(),
client_operations: HashMap::new(),
completed_client_operations: HashSet::new(),
completed_client_order: VecDeque::new(),
relays: RelayReducer::default(),
effects: VecDeque::new(),
pending: HashMap::new(),
next_effect: 0,
}
}
pub fn node(&self) -> &NodeCore {
&self.node
}
pub fn handle(&mut self, now: Instant, input: CoreInput) -> Result<(), CoreError> {
if self.effects.len() > MAX_EFFECTS - MAX_EFFECTS_PER_INPUT {
return Err(CoreError::EffectQueueFull);
}
let session = match input {
CoreInput::SessionOpened {
session,
initiator,
establish_peer,
expected_peer,
} => {
if self.sessions.contains_key(&session) {
return Err(CoreError::DuplicateSession(session.to_string()));
}
let mut core = SessionReducer::authoritative();
if !establish_peer {
core.mark_peer_ready();
}
core.handle_input(now, SessionInput::Start { initiator })?;
self.sessions.insert(session.clone(), core);
self.peers.insert(
session.clone(),
PeerSession::new(initiator, establish_peer, expected_peer),
);
session
}
CoreInput::FrameReceived { session, envelope } => {
if (envelope.kind.is_application_request() && envelope.kind != Kind::Discover)
|| (envelope.corr.is_some()
&& (envelope.kind.is_application_response()
|| envelope.kind == Kind::Cancel))
{
return Err(CoreError::Malformed(
"application frames require the payload-opaque core input".into(),
));
}
self.session_mut(&session)?
.handle_input(now, SessionInput::FrameReceived(envelope))?;
session
}
CoreInput::ApplicationFrameReceived { session, frame } => {
if let Some(corr) = frame.head.corr.as_deref() {
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr),
};
if self.completed_client_operations.contains(&stream) {
self.release_body(&session, frame.body);
return Ok(());
}
}
self.session_mut(&session)?
.handle_input(now, SessionInput::FrameReceived(frame.into_envelope()))?;
session
}
CoreInput::OperationStreamEnded {
session,
corr,
direction: _,
outcome,
} => {
let stream = StreamKey {
session: session.clone(),
corr: corr.clone(),
};
if self.completed_client_operations.contains(&stream) {
return Ok(());
}
if outcome == OperationStreamOutcome::Clean {
return Ok(());
}
let (code, message) = match outcome {
OperationStreamOutcome::Clean => unreachable!(),
OperationStreamOutcome::Cancelled => (
ErrorCode::Cancelled,
"operation stream cancelled".to_owned(),
),
OperationStreamOutcome::TimedOut => (
ErrorCode::PeerUnreachable,
"operation stream timed out".to_owned(),
),
OperationStreamOutcome::Busy => (
ErrorCode::Busy,
"operation stream capacity unavailable".to_owned(),
),
OperationStreamOutcome::PayloadTooLarge(message) => {
(ErrorCode::PayloadTooLarge, message)
}
OperationStreamOutcome::Truncated => (
ErrorCode::Protocol,
"operation stream ended with a truncated record".to_owned(),
),
OperationStreamOutcome::Protocol(message) => (ErrorCode::Protocol, message),
OperationStreamOutcome::Transport(message) => {
(ErrorCode::PeerUnreachable, message)
}
};
self.session_mut(&session)?
.operation_stream_failed(corr.as_str(), code, &message);
session
}
CoreInput::StartClientOperation {
session,
target_path,
kind,
input,
hops,
headers,
timeout,
} => {
let (payload, body) = match input {
OperationInput::Body(body) if kind != Kind::Discover => (Bytes::new(), body),
OperationInput::Discovery(plan) if kind == Kind::Discover => (
Envelope::encode_payload(
&serde_json::to_value(plan)
.map_err(|error| CoreError::Malformed(error.to_string()))?,
),
None,
),
_ => {
return Err(CoreError::Malformed(
"operation input does not match its application kind".into(),
))
}
};
let corr = if kind == Kind::Discover {
self.session_mut(&session)?.open_discovery(
&target_path,
payload,
hops,
headers,
)?
} else {
self.session_mut(&session)?.open_stream_from_body(
&target_path,
kind,
body.as_ref().map(ToString::to_string),
hops,
headers,
)?
};
let operation = ClientOperationId::from(corr.clone());
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr),
};
let deadline = timeout.map(|timeout| now + timeout);
self.client_operations.insert(stream.clone(), deadline);
self.effects.push_back(CoreEffect::RegisterClient {
session: session.clone(),
operation: operation.clone(),
});
if let Some(deadline) = deadline {
self.effects.push_back(CoreEffect::ScheduleClientDeadline {
session: session.clone(),
operation,
deadline,
});
}
let SessionEffect::SendFrame(envelope) = self.session_mut(&session)?.poll_effect()
else {
return Err(CoreError::Malformed(
"client opening frame missing".to_owned(),
));
};
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::Send {
stream,
terminal: false,
},
);
if envelope.kind == Kind::Discover {
self.effects.push_back(CoreEffect::SendProtocol {
effect,
session,
envelope,
});
} else {
self.effects.push_back(CoreEffect::Send {
effect,
session,
frame: ApplicationFrame::from_envelope(&envelope)?,
});
}
return Ok(());
}
CoreInput::CancelClientOperation { session, operation } => {
self.complete_client_operation(
&session,
&operation,
ClientDelivery::Cancelled,
true,
);
return Ok(());
}
CoreInput::ClientOperationTimeout { session, operation } => {
let stream = StreamKey {
session: session.clone(),
corr: operation.0.clone(),
};
if self
.client_operations
.get(&stream)
.is_some_and(|deadline| deadline.is_some_and(|deadline| now >= deadline))
{
self.complete_client_operation(
&session,
&operation,
ClientDelivery::TimedOut,
true,
);
}
return Ok(());
}
CoreInput::SessionTimeout { session } => {
self.session_mut(&session)?
.handle_input(now, SessionInput::Deadline(now))?;
session
}
CoreInput::PeerAdmissionCompleted { effect, result } => {
let Some(PendingEffect::PeerAdmission(session)) = self.pending.remove(&effect)
else {
return Ok(());
};
self.complete_admission(session, result)?;
return Ok(());
}
CoreInput::CapacityChecked { effect, result } => {
let Some(PendingEffect::Capacity(invocation)) = self.pending.remove(&effect) else {
return Ok(());
};
let invocation = *invocation;
match result {
CapacityResult::Available => {
let effect = self.effect_id();
self.pending
.insert(effect, PendingEffect::Dispatch(invocation.stream.clone()));
self.effects
.push_back(CoreEffect::InvokeApplication { effect, invocation });
}
CapacityResult::Busy => {
self.release_body(&invocation.stream.session, invocation.frame.body);
self.fail_stream(
&invocation.stream,
ErrorCode::Busy,
"node at activation capacity",
)?;
}
}
return Ok(());
}
CoreInput::TargetReadinessCompleted { effect, result } => {
let Some(PendingEffect::TargetReadiness { stream, frame }) =
self.pending.remove(&effect)
else {
return Ok(());
};
match result {
TargetReadinessResult::Ready => {
let envelope = frame.clone().into_envelope();
if !self.resolve_relay_opening_with_wait(&stream, &envelope, false) {
self.release_body(&stream.session, frame.body);
self.fail_stream(
&stream,
ErrorCode::PeerUnreachable,
"target readiness completed without an authoritative route",
)?;
}
}
TargetReadinessResult::Unavailable { message } => {
self.release_body(&stream.session, frame.body);
self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message)?;
}
}
return Ok(());
}
CoreInput::RelayOpenCompleted { effect, result } => {
let Some(PendingEffect::RelayOpen { stream: source, .. }) =
self.pending.remove(&effect)
else {
return Ok(());
};
match result {
RelayOpenResult::Opened(target) => {
if !self.relays.pair(source.clone(), target) {
self.fail_stream(
&source,
ErrorCode::Conflict,
"relay stream is already paired",
)?;
}
}
RelayOpenResult::Failed(failure) => {
self.fail_stream(&source, failure.code, &failure.message)?;
}
}
return Ok(());
}
CoreInput::RelayForwardCompleted { effect, result } => {
if result == SendResult::Reserved {
return Ok(());
}
let Some(PendingEffect::RelayForward {
source,
target,
terminal,
}) = self.pending.remove(&effect)
else {
return Ok(());
};
match result {
SendResult::Written => {
if terminal {
self.relays.remove_pair(&source);
}
}
SendResult::ReservationTimedOut => self.fail_relay_forward(&source),
SendResult::Refused { code, message } => {
self.fail_relay_forward_with(&source, code, &message)
}
SendResult::Closed | SendResult::Cancelled => {
self.fail_relay_forward(&source);
self.retire(target, RetirementReason::SendClosed);
}
SendResult::WriteFailed(_) => {
self.fail_relay_forward(&source);
self.retire(target, RetirementReason::SendTransport);
}
SendResult::Reserved => unreachable!(),
}
return Ok(());
}
CoreInput::DispatchCompleted { effect, result } => {
let Some(PendingEffect::Dispatch(stream)) = self.pending.remove(&effect) else {
return Ok(());
};
match result {
Ok(ApplicationResult::Response(payload))
| Ok(ApplicationResult::Finished(payload)) => {
self.pending.remove(&effect);
self.queue_application_response(stream, payload, true)?;
}
Ok(ApplicationResult::Event(payload)) => {
self.pending
.insert(effect, PendingEffect::Dispatch(stream.clone()));
self.queue_application_response(stream, payload, false)?;
}
Ok(ApplicationResult::Bridged) => {}
Err(failure) => {
self.pending.remove(&effect);
self.fail_stream(&stream, failure.code, &failure.message)?;
}
}
return Ok(());
}
CoreInput::SendCompleted { effect, result } => {
if result == SendResult::Reserved {
return Ok(());
}
let Some(PendingEffect::Send { stream, terminal }) = self.pending.remove(&effect)
else {
return Ok(());
};
let relay_leg = self.relays.peer(&stream).is_some();
match result {
SendResult::Written => {}
SendResult::ReservationTimedOut if relay_leg => {
self.fail_relay_forward(&stream);
}
SendResult::ReservationTimedOut if !terminal => {
self.fail_stream(&stream, ErrorCode::Busy, "send capacity unavailable")?;
}
SendResult::ReservationTimedOut => {}
SendResult::Refused { code, message } if relay_leg => {
self.fail_relay_forward_with(&stream, code, &message);
}
SendResult::Refused { code, message } if !terminal => {
self.fail_stream(&stream, code, &message)?;
}
SendResult::Refused { code, message } => {
if let Ok(session) = self.session_mut(&stream.session) {
session.fail_closed(stream.corr.as_str(), code, &message);
}
self.drain_session(stream.session.clone());
}
SendResult::Closed | SendResult::Cancelled => {
self.retire(stream.session, RetirementReason::SendClosed);
}
SendResult::WriteFailed(_) => {
self.retire(stream.session, RetirementReason::SendTransport);
}
SendResult::Reserved => unreachable!(),
}
return Ok(());
}
CoreInput::EstablishmentTimeout { session } => {
if self.peers.get(&session).is_some_and(|peer| peer.ready)
|| self.session_class(&session)? == SessionClass::Client
{
return Ok(());
}
self.retire(session, RetirementReason::EstablishmentTimeout);
return Ok(());
}
CoreInput::DiscoveryNeighborEvent {
stream,
peer,
event,
} => {
self.handle_discovery(stream, WalkInput::NeighborEvent { peer, event });
return Ok(());
}
CoreInput::DiscoveryNeighborDone { stream, peer } => {
self.handle_discovery(stream, WalkInput::NeighborDone { peer });
return Ok(());
}
CoreInput::DiscoveryNeighborTimeout { stream, peer } => {
self.handle_discovery(stream, WalkInput::NeighborTimeout { peer });
return Ok(());
}
CoreInput::DiscoveryTargetEvent { stream, event } => {
if !self.target_discoveries.contains(&stream) {
return Ok(());
}
let terminal = matches!(event, DiscoverEvent::Done { .. });
let payload = Envelope::encode_payload(
&serde_json::to_value(event)
.map_err(|error| CoreError::Malformed(error.to_string()))?,
);
let kind = if terminal {
Kind::Response
} else {
Kind::Event
};
if terminal {
self.target_discoveries.remove(&stream);
}
let _ = self.queue_send(stream, kind, payload);
return Ok(());
}
CoreInput::DiscoveryTargetFailed { stream, message } => {
if self.target_discoveries.remove(&stream) {
let _ = self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message);
}
return Ok(());
}
CoreInput::SessionClosed { session } => {
self.close_session(&session, RetirementReason::SessionClosed);
return Ok(());
}
CoreInput::SessionFailed { session } => {
self.close_session(&session, RetirementReason::TransportFailed);
return Ok(());
}
CoreInput::ContinueDiscovery { stream } => {
self.drain_discovery(stream);
return Ok(());
}
CoreInput::LocalCapabilitiesInstalled { capabilities } => {
self.node.install_local_capabilities(capabilities);
return Ok(());
}
};
self.drain_session(session);
Ok(())
}
pub fn poll_effect(&mut self) -> Option<CoreEffect> {
self.effects.pop_front()
}
pub fn session_deadline(&self, session: &SessionId) -> Result<Option<Instant>, CoreError> {
self.sessions
.get(session)
.map(SessionReducer::deadline)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
}
pub fn session_class(&self, session: &SessionId) -> Result<SessionClass, CoreError> {
self.sessions
.get(session)
.map(SessionReducer::class)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
}
pub fn session_ready(&self, session: &SessionId) -> Result<bool, CoreError> {
self.peers
.get(session)
.map(|peer| peer.ready)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
}
pub fn relay_peer(&self, stream: &StreamKey) -> Option<&StreamKey> {
self.relays.peer(stream)
}
pub fn open_stream_body(
&mut self,
session: &SessionId,
target_path: &str,
kind: Kind,
body: Option<crate::BodyId>,
hops: Option<u8>,
headers: serde_json::Map<String, serde_json::Value>,
) -> Result<String, CoreError> {
if kind == Kind::Discover {
return Err(CoreError::BadKind(format!("{kind:?}")));
}
let corr = self.session_mut(session)?.open_stream_from_body(
target_path,
kind,
body.map(|body| body.to_string()),
hops,
headers,
)?;
self.drain_session(session.clone());
Ok(corr)
}
pub fn send_body(
&mut self,
session: &SessionId,
corr: &str,
body: Option<crate::BodyId>,
) -> Result<(), CoreError> {
self.session_mut(session)?.send_body(
corr,
body.map(|body| body.to_string()),
Default::default(),
)?;
self.drain_session(session.clone());
Ok(())
}
pub fn respond_body(
&mut self,
session: &SessionId,
corr: &str,
body: Option<crate::BodyId>,
headers: serde_json::Map<String, serde_json::Value>,
) -> Result<(), CoreError> {
self.session_mut(session)?.respond_body(
corr,
body.map(|body| body.to_string()),
headers,
)?;
self.drain_session(session.clone());
Ok(())
}
pub fn respond_terminal_body(
&mut self,
session: &SessionId,
corr: &str,
body: Option<crate::BodyId>,
path: Vec<String>,
) -> Result<(), CoreError> {
self.session_mut(session)?.respond_terminal_body(
corr,
body.map(|body| body.to_string()),
path,
Default::default(),
)?;
self.drain_session(session.clone());
Ok(())
}
pub fn control(
&mut self,
session: &SessionId,
kind: Kind,
payload: Bytes,
) -> Result<(), CoreError> {
self.session_mut(session)?.control(kind, payload)?;
self.drain_session(session.clone());
Ok(())
}
pub fn fail(
&mut self,
session: &SessionId,
corr: &str,
code: ErrorCode,
message: &str,
) -> Result<(), CoreError> {
self.session_mut(session)?.fail(corr, code, message)?;
self.drain_session(session.clone());
Ok(())
}
pub fn cancel(&mut self, session: &SessionId, corr: &str) -> Result<(), CoreError> {
self.session_mut(session)?.cancel(corr)?;
self.drain_session(session.clone());
Ok(())
}
fn session_mut(&mut self, session: &SessionId) -> Result<&mut SessionReducer, CoreError> {
self.sessions
.get_mut(session)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
}
fn drain_session(&mut self, session: SessionId) {
loop {
let output = self.sessions.get_mut(&session).unwrap().poll_effect();
match output {
SessionEffect::HandshakeEstablished { version } => {
if self
.peers
.get(&session)
.is_some_and(|peer| peer.initiator && peer.establish_peer)
{
let _ = self.send_identify(&session);
}
self.effects.push_back(CoreEffect::HandshakeEstablished {
session: session.clone(),
version,
});
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::Identify => {
if self.peers[&session].establish_peer {
self.on_identify(&session, envelope);
} else {
self.deliver_client(session.clone(), envelope);
}
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::IdentityAccepted => {
if self.peers[&session].establish_peer {
self.on_identity_accepted(&session);
} else {
self.deliver_client(session.clone(), envelope);
}
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteSnapshot => {
if !self.peers[&session].establish_peer {
self.deliver_client(session.clone(), envelope);
} else if self.peers.get(&session).is_some_and(|peer| peer.ready) {
self.on_route_snapshot(&session, envelope);
} else {
self.on_initial_snapshot(&session, envelope);
}
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteDelta => {
if self.peers.get(&session).is_some_and(|peer| peer.ready) {
self.on_route_delta(&session, envelope);
} else {
self.retire(session.clone(), RetirementReason::SnapshotOrder);
}
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::RouteAck => {
if !self.peers[&session].establish_peer {
self.deliver_client(session.clone(), envelope);
} else if self.peers.get(&session).is_some_and(|peer| peer.ready) {
self.on_route_ack(&session, envelope);
} else {
self.on_initial_ack(&session, envelope);
}
}
SessionEffect::Deliver(envelope) if envelope.kind == Kind::Discover => {
self.start_discovery(&session, envelope);
}
SessionEffect::Deliver(envelope) if envelope.kind.is_application_request() => {
if !self.peers[&session].establish_peer {
self.deliver_client(session.clone(), envelope);
continue;
}
if self
.peers
.get(&session)
.is_some_and(|peer| peer.expected_peer.is_some() && !peer.ready)
{
self.retire(session.clone(), RetirementReason::SessionClosed);
continue;
}
if let Some(corr) = envelope.corr.clone() {
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr),
};
if self.resolve_relay_opening(&stream, &envelope) {
continue;
}
let effect = self.effect_id();
let frame = match ApplicationFrame::from_envelope(&envelope) {
Ok(frame) => frame,
Err(_) => {
self.retire(session.clone(), RetirementReason::TransportFailed);
continue;
}
};
let invocation = ApplicationInvocation {
stream: stream.clone(),
reservation: effect,
origin: self.application_origin(&session),
frame: frame.clone(),
};
self.pending
.insert(effect, PendingEffect::Capacity(Box::new(invocation)));
self.effects.push_back(CoreEffect::CheckDispatchCapacity {
effect,
stream,
frame,
});
}
}
SessionEffect::Deliver(envelope) => {
let stream = envelope.corr.as_ref().map(|corr| StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr.clone()),
});
let forwarded =
ApplicationFrame::from_envelope(&envelope)
.ok()
.and_then(|frame| {
stream.as_ref().and_then(|stream| {
self.relays.forward(stream, frame, self.node.node())
})
});
if let (Some(source), Some(mut forwarded)) = (stream, forwarded) {
if forwarded.terminal {
if let Some(target) = self.sessions.get_mut(&forwarded.target.session) {
let forwarded_envelope = forwarded.frame.clone().into_envelope();
let result = target.respond_terminal_with(
forwarded.target.corr.as_str(),
forwarded_envelope.kind,
forwarded_envelope.payload,
forwarded_envelope.path,
forwarded_envelope.headers,
);
if result.is_ok() {
if let SessionEffect::SendFrame(mut envelope) =
target.poll_effect()
{
envelope.body_token =
forwarded.frame.body.as_ref().map(ToString::to_string);
if let Ok(frame) =
ApplicationFrame::from_envelope(&envelope)
{
forwarded.frame = frame;
}
}
}
}
}
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::RelayForward {
source: source.clone(),
target: forwarded.target.session.clone(),
terminal: forwarded.terminal,
},
);
self.effects.push_back(CoreEffect::ForwardRelay {
effect,
source,
target: forwarded.target,
frame: forwarded.frame,
terminal: forwarded.terminal,
});
} else {
self.deliver_client(session.clone(), envelope);
}
}
SessionEffect::Closed { code, message } => {
self.effects.push_back(CoreEffect::CloseTransport {
session: session.clone(),
code,
message,
});
self.retire(session.clone(), RetirementReason::SessionClosed);
}
SessionEffect::SendFrame(envelope)
if envelope.kind.is_application_request()
|| (envelope.corr.is_some()
&& (envelope.kind.is_application_response()
|| envelope.kind == Kind::Cancel)) =>
{
if envelope.kind == Kind::Discover {
let effect = self.effect_id();
self.effects.push_back(CoreEffect::SendProtocol {
effect,
session: session.clone(),
envelope,
});
continue;
}
let Some(corr) = envelope.corr.clone() else {
self.retire(session.clone(), RetirementReason::TransportFailed);
continue;
};
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr),
};
let terminal =
matches!(envelope.kind, Kind::Response | Kind::Error | Kind::Cancel);
let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
self.retire(session.clone(), RetirementReason::TransportFailed);
continue;
};
let effect = self.effect_id();
self.pending
.insert(effect, PendingEffect::Send { stream, terminal });
self.effects.push_back(CoreEffect::Send {
effect,
session: session.clone(),
frame,
});
}
SessionEffect::SendFrame(envelope) => {
self.effects.push_back(CoreEffect::SendFrame {
session: session.clone(),
envelope,
})
}
SessionEffect::StreamClosed { corr } => self.stream_closed(session.clone(), corr),
SessionEffect::Deadline(Some(deadline)) => {
self.effects.push_back(CoreEffect::ScheduleSessionDeadline {
session: session.clone(),
deadline,
});
break;
}
SessionEffect::Deadline(None) => break,
}
}
}
fn effect_id(&mut self) -> EffectId {
let effect = EffectId::new(self.next_effect);
self.next_effect = self.next_effect.checked_add(1).expect("effect id overflow");
effect
}
fn resolve_relay_opening(&mut self, stream: &StreamKey, envelope: &Envelope) -> bool {
self.resolve_relay_opening_with_wait(stream, envelope, true)
}
fn resolve_relay_opening_with_wait(
&mut self,
stream: &StreamKey,
envelope: &Envelope,
allow_wait: bool,
) -> bool {
if envelope.target.is_empty() {
return false;
}
match self.node.resolve(&envelope.target) {
crate::Resolution::Route(_) => match self.node.forward(envelope.clone()) {
Ok((peer, envelope)) => {
let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
let _ = self.fail_stream(
stream,
ErrorCode::Protocol,
"relay application frame is malformed",
);
return true;
};
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::RelayOpen {
stream: stream.clone(),
body: frame.body.clone(),
},
);
self.effects.push_back(CoreEffect::OpenRelay {
effect,
source: stream.clone(),
peer,
frame,
});
}
Err(crate::RouteError::HopLimitExceeded) => {
self.release_envelope_body(&stream.session, envelope);
let payload = Envelope::encode_payload(&serde_json::json!({
"code": ErrorCode::HopLimitExceeded,
"message": "hop limit exceeded",
}));
let node = self.node.node().to_string();
let _ = self.session_mut(&stream.session).and_then(|session| {
session.respond_terminal(
stream.corr.as_str(),
Kind::Error,
payload,
vec![node],
)
});
self.drain_session(stream.session.clone());
}
Err(error) => {
self.release_envelope_body(&stream.session, envelope);
let _ = self.fail_stream(stream, ErrorCode::Internal, &error.to_string());
}
},
crate::Resolution::Conflicted { owners } => {
self.release_envelope_body(&stream.session, envelope);
let _ = self.fail_stream(
stream,
ErrorCode::PeerUnreachable,
&format!(
"destination node {:?} has multiple live incarnations: {}",
envelope.target,
owners.join(", ")
),
);
}
crate::Resolution::Unknown => {
if allow_wait {
let Ok(frame) = ApplicationFrame::from_envelope(envelope) else {
let _ = self.fail_stream(
stream,
ErrorCode::Protocol,
"relay application frame is malformed",
);
return true;
};
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::TargetReadiness {
stream: stream.clone(),
frame,
},
);
self.effects.push_back(CoreEffect::AwaitTargetReadiness {
effect,
stream: stream.clone(),
target: envelope.target.clone(),
});
} else {
self.release_envelope_body(&stream.session, envelope);
let known = self.node.reachable_names();
let known = known.iter().map(String::as_str).collect::<Vec<_>>();
let message = crate::teach_unknown("node", &envelope.target, &known);
let _ = self.fail_stream(stream, ErrorCode::PeerUnreachable, &message);
}
}
crate::Resolution::Local => return false,
}
true
}
fn release_envelope_body(&mut self, session: &SessionId, envelope: &Envelope) {
let body = ApplicationFrame::from_envelope(envelope)
.ok()
.and_then(|frame| frame.body);
self.release_body(session, body);
}
fn start_discovery(&mut self, session: &SessionId, envelope: Envelope) {
let Some(corr) = envelope.corr else {
return;
};
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr),
};
let Ok(plan) = DiscoverPlan::decode(&envelope.payload) else {
let _ = self.fail_stream(&stream, ErrorCode::InvalidInput, "malformed discovery plan");
return;
};
match self.node.resolve(&envelope.target) {
crate::Resolution::Route(peer) => {
let hops = envelope.hops.unwrap_or(crate::DEFAULT_HOPS);
if hops == 0 {
let _ = self.fail_stream(
&stream,
ErrorCode::HopLimitExceeded,
"hop limit exceeded",
);
return;
}
let Ok(target_path) = TargetPath::discovery(&envelope.target) else {
let _ = self.fail_stream(
&stream,
ErrorCode::InvalidInput,
"invalid discovery target",
);
return;
};
let mut plan = plan;
plan.hops = plan.hops.min(hops - 1);
self.target_discoveries.insert(stream.clone());
self.effects.push_back(CoreEffect::QueryDiscoveryTarget {
stream,
peer,
target_path: target_path.to_string(),
plan,
});
return;
}
crate::Resolution::Conflicted { owners } => {
let _ = self.fail_stream(
&stream,
ErrorCode::PeerUnreachable,
&format!(
"destination node {:?} has multiple live incarnations: {}",
envelope.target,
owners.join(", ")
),
);
return;
}
crate::Resolution::Unknown => {
let known = self.node.reachable_names();
let known = known.iter().map(String::as_str).collect::<Vec<_>>();
let message = crate::teach_unknown("node", &envelope.target, &known);
let _ = self.fail_stream(&stream, ErrorCode::PeerUnreachable, &message);
return;
}
crate::Resolution::Local => {}
}
let candidates = self.active_peers.keys().cloned().collect();
let snapshot = self.node.catalog_snapshot(plan.detail.is_full());
self.discoveries.insert(
stream.clone(),
DiscoverWalk::start(snapshot, plan, candidates),
);
self.drain_discovery(stream);
}
fn handle_discovery(&mut self, stream: StreamKey, input: WalkInput) {
let Some(walk) = self.discoveries.get_mut(&stream) else {
return;
};
walk.handle(input);
self.drain_discovery(stream);
}
fn drain_discovery(&mut self, stream: StreamKey) {
for _ in 0..MAX_DISCOVERY_EFFECTS_PER_INPUT - 1 {
let output = self
.discoveries
.get_mut(&stream)
.and_then(DiscoverWalk::drain);
match output {
Some(WalkOutput::Emit(event)) => {
let payload = Envelope::encode_payload(&serde_json::to_value(event).unwrap());
let _ = self.queue_send(stream.clone(), Kind::Event, payload);
}
Some(WalkOutput::AskNeighbor { peer, plan }) => {
self.effects.push_back(CoreEffect::QueryDiscoveryNeighbor {
stream: stream.clone(),
peer,
plan,
});
}
Some(WalkOutput::Finish) => {
let discover_id = self
.discoveries
.get(&stream)
.map(|walk| walk.discover_id().to_string())
.unwrap_or_default();
let failure = self
.discoveries
.get(&stream)
.and_then(DiscoverWalk::strict_failure)
.map(str::to_owned);
self.discoveries.remove(&stream);
if let Some(peer) = failure {
let _ = self.fail_stream(
&stream,
ErrorCode::PeerUnreachable,
&format!("discovery timed out at {peer}"),
);
} else {
let done = DiscoverEvent::Done { discover_id };
let payload =
Envelope::encode_payload(&serde_json::to_value(done).unwrap());
let _ = self.queue_send(stream, Kind::Response, payload);
}
return;
}
None => return,
}
}
if self
.discoveries
.get(&stream)
.is_some_and(DiscoverWalk::has_output)
{
self.effects
.push_back(CoreEffect::ContinueDiscovery { stream });
}
}
fn close_session(&mut self, session: &SessionId, reason: RetirementReason) {
let affected: Vec<_> = self
.discoveries
.keys()
.filter(|stream| &stream.session == session)
.cloned()
.collect();
for stream in affected {
self.discoveries.remove(&stream);
}
self.target_discoveries
.retain(|stream| &stream.session != session);
for (peer, peer_is_source) in self.relays.drain_session(session) {
if peer_is_source {
let payload = Envelope::encode_payload(&serde_json::json!({
"code": ErrorCode::PeerUnreachable,
"message": "the peer serving this stream disconnected",
}));
let node = self.node.node().to_string();
if let Ok(target) = self.session_mut(&peer.session) {
let _ = target.respond_terminal(
peer.corr.as_str(),
Kind::Error,
payload,
vec![node],
);
}
} else if let Ok(target) = self.session_mut(&peer.session) {
let _ = target.cancel(peer.corr.as_str());
}
self.drain_session(peer.session);
}
self.retire(session.clone(), reason);
}
fn fail_relay_forward(&mut self, source: &StreamKey) {
self.fail_relay_forward_with(
source,
ErrorCode::Busy,
"relay dropped a slow consumer stream",
);
}
fn fail_relay_forward_with(&mut self, source: &StreamKey, code: ErrorCode, message: &str) {
let Some((peer, source_is_return)) = self.relays.remove_pair(source) else {
return;
};
if source_is_return {
let payload = Envelope::encode_payload(&serde_json::json!({
"code": code,
"message": message,
}));
let node = self.node.node().to_string();
if let Ok(session) = self.session_mut(&peer.session) {
let _ =
session.respond_terminal(peer.corr.as_str(), Kind::Error, payload, vec![node]);
}
self.drain_session(peer.session.clone());
if let Ok(session) = self.session_mut(&source.session) {
let _ = session.cancel(source.corr.as_str());
}
self.drain_session(source.session.clone());
} else {
if let Ok(session) = self.session_mut(&source.session) {
let _ = session.cancel(source.corr.as_str());
}
self.drain_session(source.session.clone());
if let Ok(session) = self.session_mut(&peer.session) {
let _ = session.cancel(peer.corr.as_str());
}
self.drain_session(peer.session);
}
}
fn on_identify(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(remote) = envelope.parse_payload::<NodeIdentity>() else {
self.retire(session.clone(), RetirementReason::MalformedIdentity);
return;
};
let Some(peer) = self.peers.get_mut(session) else {
return;
};
if peer.establishment.on_identify(remote.clone()).is_err() {
self.retire(session.clone(), RetirementReason::DuplicateIdentity);
return;
}
let effect = self.effect_id();
self.pending
.insert(effect, PendingEffect::PeerAdmission(session.clone()));
self.effects.push_back(CoreEffect::RequestPeerAdmission {
effect,
session: session.clone(),
remote,
});
}
fn complete_admission(
&mut self,
session: SessionId,
result: PeerAdmission,
) -> Result<(), CoreError> {
let PeerAdmission::Admitted(verified) = result else {
self.retire(session, RetirementReason::AdmissionRejected);
return Ok(());
};
let local_node = self.node.node().to_string();
let Some(peer) = self.peers.get_mut(&session) else {
return Ok(());
};
if peer.retired {
return Ok(());
}
let Some(declared) = peer.establishment.remote_identity() else {
return Ok(());
};
if verified.node_id != declared.node_id || verified.instance_id != declared.instance_id {
self.retire(session, RetirementReason::AdmissionIdentityMismatch);
return Ok(());
}
if verified.node_id == local_node {
self.retire(session, RetirementReason::SelfConnection);
return Ok(());
}
if peer
.expected_peer
.as_deref()
.is_some_and(|expected| expected != verified.node_id)
{
self.retire(session, RetirementReason::UnexpectedPeer);
return Ok(());
}
peer.establishment.local_accept()?;
peer.remote = Some(verified);
if !peer.initiator {
self.send_identify(&session)?;
}
self.control(&session, Kind::IdentityAccepted, Bytes::new())?;
self.maybe_send_snapshot(&session)?;
Ok(())
}
fn on_identity_accepted(&mut self, session: &SessionId) {
let result = self
.peers
.get_mut(session)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
.and_then(|peer| peer.establishment.on_identity_accepted());
if result.is_err() {
self.retire(session.clone(), RetirementReason::IdentityAcceptanceOrder);
return;
}
if self.maybe_send_snapshot(session).is_err() {
self.retire(session.clone(), RetirementReason::SnapshotSequence);
return;
}
self.try_promote(session.clone());
}
fn on_initial_snapshot(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(snapshot) = envelope.parse_payload::<RouteSnapshot>() else {
self.retire(session.clone(), RetirementReason::MalformedSnapshot);
return;
};
let Some(remote) = self
.peers
.get(session)
.and_then(|peer| peer.remote.as_ref())
.cloned()
else {
self.retire(session.clone(), RetirementReason::SnapshotBeforeIdentity);
return;
};
let mut probe = self.node.clone();
if probe
.apply_snapshot(session.as_str(), &remote.node_id, &snapshot)
.is_err()
{
self.retire(session.clone(), RetirementReason::SnapshotRejected);
return;
}
let Some(peer) = self.peers.get_mut(session) else {
return;
};
if peer.establishment.on_snapshot_applied().is_err() {
self.retire(session.clone(), RetirementReason::SnapshotOrder);
return;
}
peer.held_snapshot = Some(snapshot.clone());
let ack = RouteAck {
generation: snapshot.generation,
status: RouteAckStatus::Applied,
};
if peer.initiator {
let _ = self.control(
session,
Kind::RouteAck,
Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
);
} else {
peer.deferred_ack = Some(ack);
}
self.try_promote(session.clone());
}
fn on_initial_ack(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(ack) = envelope.parse_payload::<RouteAck>() else {
self.retire(session.clone(), RetirementReason::MalformedAcknowledgement);
return;
};
let result = self
.peers
.get_mut(session)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))
.and_then(|peer| peer.establishment.on_route_ack(&ack));
if result.is_err() {
self.retire(session.clone(), RetirementReason::AcknowledgementOrder);
return;
}
self.effects.push_back(CoreEffect::RouteExportAcked {
session: session.clone(),
});
self.try_promote(session.clone());
}
fn on_route_snapshot(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(snapshot) = envelope.parse_payload::<RouteSnapshot>() else {
self.retire(session.clone(), RetirementReason::MalformedRouteControl);
return;
};
let Some(peer) = self.peer_node(session) else {
return;
};
match self.node.apply_snapshot(session.as_str(), &peer, &snapshot) {
Ok(changed) => {
self.route_ack(session, snapshot.generation, RouteAckStatus::Applied);
if !changed.is_empty() {
let _ = self.recompute_route_exports();
}
self.effects.push_back(CoreEffect::RouteSnapshotApplied {
session: session.clone(),
peer,
snapshot,
changed: !changed.is_empty(),
});
}
Err(RouteError::StaleUpdate { .. }) => {
if let Some(generation) = self.node.applied_generation(session.as_str()) {
self.route_ack(session, generation, RouteAckStatus::Applied);
}
}
Err(_) => self.retire(session.clone(), RetirementReason::RouteUpdateRejected),
}
}
fn on_route_delta(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(delta) = envelope.parse_payload::<RouteDelta>() else {
self.retire(session.clone(), RetirementReason::MalformedRouteControl);
return;
};
let Some(peer) = self.peer_node(session) else {
return;
};
match self.node.apply_delta(session.as_str(), &peer, &delta) {
Ok(changed) => {
self.route_ack(session, delta.generation, RouteAckStatus::Applied);
if !changed.is_empty() {
let _ = self.recompute_route_exports();
}
self.effects.push_back(CoreEffect::RouteDeltaApplied {
session: session.clone(),
peer,
delta,
changed: !changed.is_empty(),
});
}
Err(RouteError::StaleUpdate { .. }) => {
if let Some(generation) = self.node.applied_generation(session.as_str()) {
self.route_ack(session, generation, RouteAckStatus::Applied);
}
}
Err(RouteError::GenerationGap { expected, .. }) => self.route_ack(
session,
expected.saturating_sub(1),
RouteAckStatus::ResyncRequired,
),
Err(_) => self.retire(session.clone(), RetirementReason::RouteUpdateRejected),
}
}
fn on_route_ack(&mut self, session: &SessionId, envelope: Envelope) {
let Ok(ack) = envelope.parse_payload::<RouteAck>() else {
self.retire(session.clone(), RetirementReason::MalformedAcknowledgement);
return;
};
let Some(export) = self
.peers
.get_mut(session)
.and_then(|peer| peer.export.as_mut())
else {
return;
};
if ack.generation > export.generation {
self.retire(session.clone(), RetirementReason::AcknowledgementOrder);
return;
}
match ack.status {
RouteAckStatus::Applied => {
export.applied_ack = export.applied_ack.max(ack.generation);
}
RouteAckStatus::ResyncRequired => {
if ack.generation < export.applied_ack
|| export
.handled_resync
.is_some_and(|handled| ack.generation <= handled)
{
return;
}
export.handled_resync = Some(ack.generation);
export.generation = export.generation.saturating_add(1);
let snapshot = RouteSnapshot::canonical(export.generation, export.routes.clone());
let _ = self.control(
session,
Kind::RouteSnapshot,
Envelope::encode_payload(&serde_json::to_value(snapshot).unwrap()),
);
}
}
}
fn on_route_exports_changed(
&mut self,
session: &SessionId,
mut routes: Vec<RouteAdvertisement>,
) -> Result<(), CoreError> {
routes.sort();
let Some(export) = self
.peers
.get_mut(session)
.and_then(|peer| peer.export.as_mut())
else {
return Ok(());
};
if routes == export.routes {
return Ok(());
}
let previous: BTreeMap<&str, &RouteAdvertisement> = export
.routes
.iter()
.map(|route| (route.destination.as_str(), route))
.collect();
let current: BTreeMap<&str, &RouteAdvertisement> = routes
.iter()
.map(|route| (route.destination.as_str(), route))
.collect();
let upsert = routes
.iter()
.filter(|route| previous.get(route.destination.as_str()) != Some(route))
.cloned()
.collect();
let withdraw = export
.routes
.iter()
.filter(|route| !current.contains_key(route.destination.as_str()))
.map(|route| RouteWithdrawal {
destination: route.destination.clone(),
owner: route.owner.clone(),
owner_instance: route.owner_instance.clone(),
owner_epoch: route.owner_epoch,
owner_revision: route.owner_revision,
})
.collect();
export.generation = export.generation.saturating_add(1);
let delta = RouteDelta {
generation: export.generation,
upsert,
withdraw,
};
export.routes = routes;
self.control(
session,
Kind::RouteDelta,
Envelope::encode_payload(&serde_json::to_value(delta).unwrap()),
)
}
fn recompute_route_exports(&mut self) -> Result<(), CoreError> {
let peers: Vec<_> = self
.peers
.iter()
.filter_map(|(session, peer)| {
peer.ready
.then(|| {
peer.remote
.as_ref()
.map(|remote| (session.clone(), remote.node_id.clone()))
})
.flatten()
})
.collect();
for (session, peer) in peers {
self.on_route_exports_changed(&session, self.node.export_for(&peer))?;
}
Ok(())
}
fn peer_node(&self, session: &SessionId) -> Option<String> {
self.peers
.get(session)
.and_then(|peer| peer.remote.as_ref())
.map(|peer| peer.node_id.clone())
}
fn route_ack(&mut self, session: &SessionId, generation: u64, status: RouteAckStatus) {
let ack = RouteAck { generation, status };
let _ = self.control(
session,
Kind::RouteAck,
Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
);
}
fn send_identify(&mut self, session: &SessionId) -> Result<(), CoreError> {
let peer = self
.peers
.get_mut(session)
.ok_or_else(|| CoreError::UnknownSession(session.to_string()))?;
peer.establishment.identify_sent()?;
let payload =
Envelope::encode_payload(&serde_json::to_value(self.node.identity()).unwrap());
self.control(session, Kind::Identify, payload)
}
fn maybe_send_snapshot(&mut self, session: &SessionId) -> Result<(), CoreError> {
let Some(remote) = self
.peers
.get(session)
.and_then(|peer| peer.remote.as_ref())
.cloned()
else {
return Ok(());
};
if !self.peers[session].establishment.identities_accepted() {
return Ok(());
}
let snapshot = RouteSnapshot::canonical(1, self.node.export_for(&remote.node_id));
self.peers.get_mut(session).unwrap().export = Some(RouteExportState {
generation: snapshot.generation,
applied_ack: 0,
handled_resync: None,
routes: snapshot.routes.clone(),
});
self.peers
.get_mut(session)
.unwrap()
.establishment
.snapshot_sent(snapshot.generation)?;
self.control(
session,
Kind::RouteSnapshot,
Envelope::encode_payload(&serde_json::to_value(snapshot).unwrap()),
)
}
fn try_promote(&mut self, session: SessionId) {
let Some(peer) = self.peers.get(&session) else {
return;
};
if peer.ready || !peer.establishment.ready() {
return;
}
let Some(remote) = peer.remote.clone() else {
return;
};
if let Some(active) = self.active_peers.get(&remote.node_id).cloned() {
if active != session {
let old = &self.peers[&active];
if old
.remote
.as_ref()
.is_some_and(|id| id.instance_id == remote.instance_id)
{
let smaller_initiates = self.node.node() < remote.node_id.as_str();
let new_preferred = peer.initiator == smaller_initiates;
let old_preferred = old.initiator == smaller_initiates;
if !new_preferred || old_preferred {
self.retire(session, RetirementReason::DuplicateSession);
return;
}
}
self.retire(active, RetirementReason::DuplicateSessionReplaced);
}
}
let peer = self.peers.get_mut(&session).unwrap();
peer.ready = true;
self.sessions.get_mut(&session).unwrap().mark_peer_ready();
let snapshot = peer.held_snapshot.take();
let deferred_ack = peer.deferred_ack.take();
let changed = snapshot.as_ref().and_then(|snapshot| {
self.node
.apply_snapshot(session.as_str(), &remote.node_id, snapshot)
.ok()
});
self.active_peers
.insert(remote.node_id.clone(), session.clone());
self.effects.push_back(CoreEffect::SessionEstablished {
session: session.clone(),
peer: remote.clone(),
});
let peer_node = remote.node_id.clone();
let changed_routes = changed.as_ref().is_some_and(|changed| !changed.is_empty());
if let (Some(snapshot), Some(_)) = (snapshot, changed) {
self.effects.push_back(CoreEffect::RouteSnapshotApplied {
session: session.clone(),
peer: peer_node,
snapshot,
changed: changed_routes,
});
}
if let Some(ack) = deferred_ack {
let _ = self.control(
&session,
Kind::RouteAck,
Envelope::encode_payload(&serde_json::to_value(ack).unwrap()),
);
}
if changed_routes {
let _ = self.recompute_route_exports();
}
}
fn retire(&mut self, session: SessionId, reason: RetirementReason) {
let reason = self.converged_retirement(&session, reason);
let mut withdrew_routes = false;
let mut withdrew_session = false;
if let Some(peer) = self.peers.get_mut(&session) {
if peer.retired {
return;
}
peer.retired = true;
withdrew_session = peer.ready;
peer.ready = false;
if let Some(remote) = &peer.remote {
if self.active_peers.get(&remote.node_id) == Some(&session) {
self.active_peers.remove(&remote.node_id);
withdrew_routes = !self.node.leave(session.as_str()).is_empty();
}
}
}
if withdrew_routes {
let _ = self.recompute_route_exports();
}
if withdrew_session {
self.effects.push_back(CoreEffect::RouteSessionWithdrawn {
session: session.clone(),
changed: withdrew_routes,
});
}
self.abort_pending_dispatch(&session, None);
self.pending.retain(|_, pending| match pending {
PendingEffect::PeerAdmission(target) => target != &session,
PendingEffect::Capacity(invocation) => invocation.stream.session != session,
PendingEffect::Dispatch(stream) | PendingEffect::Send { stream, .. } => {
stream.session != session
}
PendingEffect::RelayOpen { stream, .. } => stream.session != session,
PendingEffect::TargetReadiness { stream, .. } => stream.session != session,
PendingEffect::RelayForward { source, target, .. } => {
source.session != session && target != &session
}
});
self.discoveries
.retain(|stream, _| stream.session != session);
self.completed_client_operations
.retain(|stream| stream.session != session);
self.completed_client_order
.retain(|stream| stream.session != session);
let operations = self
.client_operations
.keys()
.filter(|stream| stream.session == session)
.cloned()
.collect::<Vec<_>>();
for stream in operations {
self.client_operations.remove(&stream);
self.effects.push_back(CoreEffect::DeliverClient {
session: session.clone(),
operation: ClientOperationId::from(stream.corr),
delivery: ClientDelivery::SessionClosed,
});
}
self.effects
.push_back(CoreEffect::SessionRetired { session, reason });
}
fn converged_retirement(
&self,
session: &SessionId,
reason: RetirementReason,
) -> RetirementReason {
if !matches!(reason, RetirementReason::SessionClosed) {
return reason;
}
let Some(peer) = self.peers.get(session) else {
return reason;
};
if peer.ready || peer.retired {
return reason;
}
let Some(remote) = peer.remote.as_ref() else {
return reason;
};
match self.active_peers.get(&remote.node_id) {
Some(active)
if active != session && self.peers.get(active).is_some_and(|peer| peer.ready) =>
{
RetirementReason::DuplicateSession
}
_ => reason,
}
}
fn fail_stream(
&mut self,
stream: &StreamKey,
code: ErrorCode,
message: &str,
) -> Result<(), CoreError> {
self.session_mut(&stream.session)?
.fail(stream.corr.as_str(), code, message)?;
self.drain_session(stream.session.clone());
Ok(())
}
fn queue_send(
&mut self,
stream: StreamKey,
kind: Kind,
payload: Bytes,
) -> Result<(), CoreError> {
let terminal = match kind {
Kind::Response => self
.session_mut(&stream.session)?
.respond_discovery(stream.corr.as_str(), payload)
.map(|()| true)?,
Kind::Event => self
.session_mut(&stream.session)?
.send_discovery_event(stream.corr.as_str(), payload)
.map(|()| false)?,
_ => unreachable!(),
};
let outputs = {
let session = self.session_mut(&stream.session)?;
let mut outputs = Vec::new();
loop {
let output = session.poll_effect();
if matches!(output, SessionEffect::Deadline(_)) {
break;
}
outputs.push(output);
}
outputs
};
for output in outputs {
if let SessionEffect::SendFrame(envelope) = output {
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::Send {
stream: stream.clone(),
terminal,
},
);
self.effects.push_back(CoreEffect::SendProtocol {
effect,
session: stream.session.clone(),
envelope,
});
} else {
match output {
SessionEffect::Deliver(envelope) => {
self.deliver_client(stream.session.clone(), envelope)
}
SessionEffect::StreamClosed { corr } => {
self.stream_closed(stream.session.clone(), corr)
}
SessionEffect::Closed { code, message } => {
self.effects.push_back(CoreEffect::CloseTransport {
session: stream.session.clone(),
code,
message,
})
}
SessionEffect::HandshakeEstablished { version } => {
self.effects.push_back(CoreEffect::HandshakeEstablished {
session: stream.session.clone(),
version,
})
}
SessionEffect::Deadline(Some(deadline)) => {
self.effects.push_back(CoreEffect::ScheduleSessionDeadline {
session: stream.session.clone(),
deadline,
})
}
SessionEffect::Deadline(None) | SessionEffect::SendFrame(_) => {}
}
}
}
Ok(())
}
fn queue_application_response(
&mut self,
stream: StreamKey,
response: ApplicationResponse,
terminal: bool,
) -> Result<(), CoreError> {
let (parts, ()) = response.head.into_parts();
let mut envelope =
Envelope::from_response(http::Response::from_parts(parts, Bytes::new()))?;
envelope.body_token = response.body.map(|body| body.to_string());
let body = envelope.body_token.take();
if terminal {
self.session_mut(&stream.session)?.respond_terminal_body(
stream.corr.as_str(),
body,
envelope.path,
envelope.headers,
)?;
} else {
let session = self.session_mut(&stream.session)?;
session.send_body(stream.corr.as_str(), body, envelope.headers)?;
}
self.drain_send_outputs(stream, terminal)
}
fn drain_send_outputs(&mut self, stream: StreamKey, terminal: bool) -> Result<(), CoreError> {
let outputs = {
let session = self.session_mut(&stream.session)?;
let mut outputs = Vec::new();
loop {
let output = session.poll_effect();
if matches!(output, SessionEffect::Deadline(_)) {
break;
}
outputs.push(output);
}
outputs
};
for output in outputs {
if let SessionEffect::SendFrame(envelope) = output {
let effect = self.effect_id();
self.pending.insert(
effect,
PendingEffect::Send {
stream: stream.clone(),
terminal,
},
);
let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
return Err(CoreError::Malformed(
"application send frame is malformed".into(),
));
};
self.effects.push_back(CoreEffect::Send {
effect,
session: stream.session.clone(),
frame,
});
} else {
self.queue_session_effect(stream.session.clone(), output);
}
}
Ok(())
}
fn application_origin(&self, session: &SessionId) -> ApplicationOrigin {
self.peers
.get(session)
.and_then(|peer| peer.ready.then(|| peer.remote.clone()).flatten())
.map_or_else(
|| ApplicationOrigin::Client {
session: session.clone(),
},
|peer| ApplicationOrigin::Peer {
session: session.clone(),
peer,
},
)
}
fn deliver_client(&mut self, session: SessionId, envelope: Envelope) {
if let Some(corr) = envelope.corr.clone() {
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr.clone()),
};
if self.completed_client_operations.contains(&stream) {
self.release_body(
&session,
envelope.body_token.as_deref().map(crate::BodyId::from),
);
return;
}
let terminal = matches!(envelope.kind, Kind::Response | Kind::Error | Kind::Cancel);
if !self.client_operations.contains_key(&stream) {
self.effects
.push_back(CoreEffect::Deliver { session, envelope });
return;
}
if terminal {
self.client_operations.remove(&stream);
self.tombstone_client_operation(stream);
}
let Ok(frame) = ApplicationFrame::from_envelope(&envelope) else {
self.effects
.push_back(CoreEffect::Deliver { session, envelope });
return;
};
self.effects.push_back(CoreEffect::DeliverClient {
session,
operation: ClientOperationId::from(corr),
delivery: if terminal {
ClientDelivery::Terminal(frame)
} else {
ClientDelivery::Item(frame)
},
});
} else {
self.effects
.push_back(CoreEffect::Deliver { session, envelope });
}
}
fn queue_session_effect(&mut self, session: SessionId, effect: SessionEffect) {
let effect = match effect {
SessionEffect::SendFrame(envelope) => CoreEffect::SendFrame { session, envelope },
SessionEffect::HandshakeEstablished { version } => {
CoreEffect::HandshakeEstablished { session, version }
}
SessionEffect::Deliver(envelope) => {
self.deliver_client(session, envelope);
return;
}
SessionEffect::StreamClosed { corr } => {
self.stream_closed(session, corr);
return;
}
SessionEffect::Closed { code, message } => CoreEffect::CloseTransport {
session,
code,
message,
},
SessionEffect::Deadline(Some(deadline)) => {
CoreEffect::ScheduleSessionDeadline { session, deadline }
}
SessionEffect::Deadline(None) => return,
};
self.effects.push_back(effect);
}
fn complete_client_operation(
&mut self,
session: &SessionId,
operation: &ClientOperationId,
delivery: ClientDelivery,
cancel: bool,
) {
let stream = StreamKey {
session: session.clone(),
corr: operation.0.clone(),
};
if self.client_operations.remove(&stream).is_none() {
return;
}
self.tombstone_client_operation(stream);
if cancel {
let _ = self
.session_mut(session)
.and_then(|core| core.cancel(operation.as_str()));
self.drain_session(session.clone());
}
self.effects.push_back(CoreEffect::DeliverClient {
session: session.clone(),
operation: operation.clone(),
delivery,
});
}
fn abort_pending_dispatch(&mut self, session: &SessionId, corr: Option<&str>) {
let aborted: Vec<(EffectId, Option<crate::BodyId>)> = self
.pending
.iter()
.filter_map(|(effect, entry)| match entry {
PendingEffect::Capacity(invocation)
if invocation.stream.session == *session
&& corr.is_none_or(|corr| invocation.stream.corr.as_str() == corr) =>
{
Some((*effect, invocation.frame.body.clone()))
}
PendingEffect::Dispatch(stream)
if stream.session == *session
&& corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
{
Some((*effect, None))
}
PendingEffect::RelayOpen { stream, body }
if stream.session == *session
&& corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
{
Some((*effect, body.clone()))
}
PendingEffect::TargetReadiness { stream, frame }
if stream.session == *session
&& corr.is_none_or(|corr| stream.corr.as_str() == corr) =>
{
Some((*effect, frame.body.clone()))
}
_ => None,
})
.collect();
for (effect, body) in aborted {
self.pending.remove(&effect);
self.release_body(session, body);
self.effects.push_back(CoreEffect::AbortDispatch {
session: session.clone(),
effect,
});
}
}
fn release_body(&mut self, session: &SessionId, body: Option<crate::BodyId>) {
if let Some(body) = body {
self.effects.push_back(CoreEffect::ReleaseBody {
session: session.clone(),
body,
});
}
}
fn tombstone_client_operation(&mut self, stream: StreamKey) {
if !self.completed_client_operations.insert(stream.clone()) {
return;
}
self.completed_client_order.push_back(stream);
while self.completed_client_order.len() > CLIENT_TOMBSTONE_LIMIT {
if let Some(oldest) = self.completed_client_order.pop_front() {
self.completed_client_operations.remove(&oldest);
}
}
}
fn stream_closed(&mut self, session: SessionId, corr: String) {
self.abort_pending_dispatch(&session, Some(&corr));
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from(corr.clone()),
};
self.discoveries.remove(&stream);
self.target_discoveries.remove(&stream);
if self.completed_client_operations.contains(&stream) {
return;
}
if self.client_operations.remove(&stream).is_some() {
self.tombstone_client_operation(stream);
self.effects.push_back(CoreEffect::DeliverClient {
session,
operation: ClientOperationId::from(corr),
delivery: ClientDelivery::Cancelled,
});
return;
}
self.effects.push_back(CoreEffect::StreamClosed {
session,
operation: ClientOperationId::from(corr),
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ApplicationHead;
fn established_client_core() -> (ProtocolCore, SessionId) {
let mut core = ProtocolCore::new("node");
let session = SessionId::from("client");
let now = web_time::Instant::now();
core.handle(
now,
CoreInput::SessionOpened {
session: session.clone(),
initiator: true,
establish_peer: false,
expected_peer: None,
},
)
.unwrap();
while core.poll_effect().is_some() {}
core.handle(
now,
CoreInput::FrameReceived {
session: session.clone(),
envelope: Envelope {
v: crate::PROTOCOL_VERSION,
id: "welcome".into(),
target: String::new(),
subject: String::new(),
kind: Kind::Welcome,
corr: None,
seq: None,
hops: None,
body_token: None,
payload: Envelope::encode_payload(&serde_json::json!({"version": 1})),
path: Vec::new(),
headers: Default::default(),
},
},
)
.unwrap();
while core.poll_effect().is_some() {}
(core, session)
}
#[test]
fn retiring_session_removes_completed_client_operation_tombstones() {
let session = SessionId::from("session");
let stream = StreamKey {
session: session.clone(),
corr: CorrelationId::from("corr"),
};
let mut core = ProtocolCore::new("node");
core.completed_client_operations.insert(stream);
core.peers
.insert(session.clone(), PeerSession::new(false, false, None));
core.retire(session, RetirementReason::SessionClosed);
assert!(core.completed_client_operations.is_empty());
}
#[test]
fn client_operation_tombstones_are_bounded() {
let mut core = ProtocolCore::new("node");
for index in 0..(CLIENT_TOMBSTONE_LIMIT + 10) {
core.tombstone_client_operation(StreamKey {
session: SessionId::from("session"),
corr: CorrelationId::from(format!("corr-{index}")),
});
}
assert_eq!(
core.completed_client_operations.len(),
CLIENT_TOMBSTONE_LIMIT
);
assert_eq!(core.completed_client_order.len(), CLIENT_TOMBSTONE_LIMIT);
assert!(!core.completed_client_operations.contains(&StreamKey {
session: SessionId::from("session"),
corr: CorrelationId::from("corr-0"),
}));
}
#[test]
fn operation_stream_fault_is_scoped_and_late_duplicates_are_idempotent() {
let (mut core, session) = established_client_core();
let now = web_time::Instant::now();
core.handle(
now,
CoreInput::StartClientOperation {
session: session.clone(),
target_path: "/node/echo".into(),
kind: Kind::Request,
input: OperationInput::Body(None),
hops: None,
headers: Default::default(),
timeout: None,
},
)
.unwrap();
let mut corr = None;
while let Some(effect) = core.poll_effect() {
if let CoreEffect::Send { frame, .. } = effect {
corr = frame.head.corr;
}
}
let corr = CorrelationId::from(corr.unwrap());
let fault = CoreInput::OperationStreamEnded {
session: session.clone(),
corr: corr.clone(),
direction: OperationStreamDirection::Return,
outcome: OperationStreamOutcome::Truncated,
};
core.handle(now, fault.clone()).unwrap();
let effects = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
assert!(effects.iter().any(|effect| matches!(
effect,
CoreEffect::DeliverClient {
delivery: ClientDelivery::Terminal(frame),
..
} if frame.head.error.as_ref().is_some_and(|error| error.code == ErrorCode::Protocol)
)));
assert!(!effects
.iter()
.any(|effect| matches!(effect, CoreEffect::SessionRetired { .. })));
core.handle(now, fault).unwrap();
let late = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
assert!(!late.iter().any(|effect| matches!(
effect,
CoreEffect::DeliverClient { .. } | CoreEffect::SessionRetired { .. }
)));
}
#[test]
fn terminal_response_wins_over_a_late_body_fault() {
let (mut core, session) = established_client_core();
let now = web_time::Instant::now();
core.handle(
now,
CoreInput::StartClientOperation {
session: session.clone(),
target_path: "/node/echo".into(),
kind: Kind::Request,
input: OperationInput::Body(None),
hops: None,
headers: Default::default(),
timeout: None,
},
)
.unwrap();
let corr = std::iter::from_fn(|| core.poll_effect())
.find_map(|effect| match effect {
CoreEffect::Send { frame, .. } => frame.head.corr,
_ => None,
})
.unwrap();
core.handle(
now,
CoreInput::ApplicationFrameReceived {
session: session.clone(),
frame: ApplicationFrame {
head: ApplicationHead {
kind: Kind::Response,
corr: Some(corr.clone()),
..ApplicationHead::default()
},
body: Some(crate::BodyId::from("response-body")),
},
},
)
.unwrap();
let terminal = std::iter::from_fn(|| core.poll_effect()).collect::<Vec<_>>();
assert!(terminal.iter().any(|effect| matches!(
effect,
CoreEffect::DeliverClient {
delivery: ClientDelivery::Terminal(_),
..
}
)));
core.handle(
now,
CoreInput::OperationStreamEnded {
session,
corr: CorrelationId::from(corr),
direction: OperationStreamDirection::Return,
outcome: OperationStreamOutcome::Truncated,
},
)
.unwrap();
assert!(core.poll_effect().is_none());
}
#[test]
fn cancelled_or_timed_out_operation_releases_a_late_response_body() {
let (mut core, session) = established_client_core();
let now = web_time::Instant::now();
core.handle(
now,
CoreInput::StartClientOperation {
session: session.clone(),
target_path: "/node/echo".into(),
kind: Kind::Request,
input: OperationInput::Body(None),
hops: None,
headers: Default::default(),
timeout: None,
},
)
.unwrap();
let corr = std::iter::from_fn(|| core.poll_effect())
.find_map(|effect| match effect {
CoreEffect::Send { frame, .. } => frame.head.corr,
_ => None,
})
.unwrap();
core.handle(
now,
CoreInput::CancelClientOperation {
session: session.clone(),
operation: ClientOperationId::from(corr.clone()),
},
)
.unwrap();
while core.poll_effect().is_some() {}
core.handle(
now,
CoreInput::ApplicationFrameReceived {
session: session.clone(),
frame: ApplicationFrame {
head: ApplicationHead {
kind: Kind::Response,
corr: Some(corr),
..ApplicationHead::default()
},
body: Some(crate::BodyId::from("late-body")),
},
},
)
.unwrap();
assert!(matches!(
core.poll_effect(),
Some(CoreEffect::ReleaseBody { session: owner, body })
if owner == session && body.as_str() == "late-body"
));
assert!(core.poll_effect().is_none());
core.handle(
now,
CoreInput::StartClientOperation {
session: session.clone(),
target_path: "/node/echo".into(),
kind: Kind::Request,
input: OperationInput::Body(None),
hops: None,
headers: Default::default(),
timeout: Some(std::time::Duration::ZERO),
},
)
.unwrap();
let timed_out_corr = std::iter::from_fn(|| core.poll_effect())
.find_map(|effect| match effect {
CoreEffect::Send { frame, .. } => frame.head.corr,
_ => None,
})
.unwrap();
core.handle(
now,
CoreInput::ClientOperationTimeout {
session: session.clone(),
operation: ClientOperationId::from(timed_out_corr.clone()),
},
)
.unwrap();
while core.poll_effect().is_some() {}
core.handle(
now,
CoreInput::ApplicationFrameReceived {
session: session.clone(),
frame: ApplicationFrame {
head: ApplicationHead {
kind: Kind::Response,
corr: Some(timed_out_corr),
..ApplicationHead::default()
},
body: Some(crate::BodyId::from("timed-out-body")),
},
},
)
.unwrap();
assert!(matches!(
core.poll_effect(),
Some(CoreEffect::ReleaseBody { session: owner, body })
if owner == session && body.as_str() == "timed-out-body"
));
assert!(core.poll_effect().is_none());
}
}