use liminal_protocol::client::{
DetachReplayRefusalReason, DetachTransportAttemptDecision, DetachTransportFate,
DetachTransportFateDecision, ExplicitReconnectAction, LostAuthorityKind,
LostOperationAuthorityDecision, LostReconnectAuthorityDecision, ProvedOnlineTransition,
ReconnectAttemptDecision, ReconnectAttemptFate, ReconnectAttemptFateDecision,
ReconnectAttemptFateRefusalReason, ReconnectAttemptRefusalReason, ReconnectPermitDecision,
RecoveredExpectedOperationDecision, RecoveredReconnectPermitDecision, record_attempt_fate,
record_explicit_reconnect, record_online_transition, recover_expected_operation,
recover_reconnect_permit, redeem_attempt, resolve_lost_operation_authority,
resolve_lost_reconnect_authority, transport_attempt_started, transport_fate,
};
use liminal_protocol::outcome::ReconnectState;
use liminal_protocol::wire::{ClientRequest, DetachRequest};
use super::{
OperationDurability, ParticipantResumeStore, RemoteOperationTransportFate,
RemoteParticipantError, RemoteParticipantHandle, RemoteParticipantOperation,
RemoteParticipantSendOutcome, RemoteReconnectPermit, RemoteReconnectPermitOutcome, persist,
record_connection_fate, record_operation_transport_fate, take_aggregate,
};
#[derive(Debug)]
pub enum RemoteExpectedOperationRecovery {
Recovered(RemoteParticipantOperation),
NotAvailable {
already_issued: bool,
},
}
#[derive(Debug, PartialEq, Eq)]
pub enum RemoteLostOperationResolution {
Recorded {
request: ClientRequest,
testimony: LostAuthorityKind,
},
DetachParked {
request: ClientRequest,
testimony: LostAuthorityKind,
},
Refused {
reason: liminal_protocol::client::LostAuthorityResolutionRefusalReason,
},
}
#[derive(Debug, PartialEq, Eq)]
pub enum RemoteLostReconnectResolution {
Recorded {
testimony: LostAuthorityKind,
},
Refused {
reason: liminal_protocol::client::LostAuthorityResolutionRefusalReason,
},
}
#[derive(Debug)]
pub enum RemoteReconnectPermitRecovery {
Recovered(RemoteReconnectPermit),
NotAvailable {
state: ReconnectState,
},
}
#[derive(Debug)]
pub enum RemoteReconnectAttemptOutcome {
Connected {
provenance: super::ParticipantResponseProvenance,
},
Failed {
error: crate::SdkError,
},
Refused {
permit: RemoteReconnectPermit,
reason: ReconnectAttemptRefusalReason,
},
FateRefused {
reason: ReconnectAttemptFateRefusalReason,
error: Option<crate::SdkError>,
provenance: Option<super::ParticipantResponseProvenance>,
},
}
#[derive(Debug)]
pub enum RemoteDetachReplayOutcome {
Send(RemoteParticipantSendOutcome),
Refused {
reason: DetachReplayRefusalReason,
},
}
#[derive(Debug, PartialEq, Eq)]
pub enum RemoteReplayApplyOutcome<T> {
Applied,
Refused {
input: T,
reason: DetachReplayRefusalReason,
},
}
#[derive(Debug)]
pub struct RemoteTransportLossOutcome {
pub operation_fate: RemoteOperationTransportFate,
pub reconnect: RemoteReconnectPermitOutcome,
}
impl<S: ParticipantResumeStore> RemoteParticipantHandle<S> {
pub fn recover_expected_operation(
&self,
) -> Result<RemoteExpectedOperationRecovery, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
match recover_expected_operation(aggregate) {
RecoveredExpectedOperationDecision::Recovered {
aggregate,
operation,
} => {
state.aggregate = Some(aggregate);
Ok(RemoteExpectedOperationRecovery::Recovered(
RemoteParticipantOperation {
operation,
durability: OperationDurability::WriteAhead,
},
))
}
RecoveredExpectedOperationDecision::NotAvailable {
aggregate,
already_issued,
} => {
state.aggregate = Some(aggregate);
Ok(RemoteExpectedOperationRecovery::NotAvailable { already_issued })
}
}
}
pub fn resolve_lost_operation_authority(
&self,
) -> Result<RemoteLostOperationResolution, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
let outcome = match resolve_lost_operation_authority(aggregate) {
LostOperationAuthorityDecision::Recorded {
aggregate,
request,
testimony,
} => {
state.aggregate = Some(aggregate);
RemoteLostOperationResolution::Recorded {
request,
testimony: testimony.kind(),
}
}
LostOperationAuthorityDecision::DetachParked {
aggregate,
request,
testimony,
} => {
state.aggregate = Some(aggregate);
RemoteLostOperationResolution::DetachParked {
request,
testimony: testimony.kind(),
}
}
LostOperationAuthorityDecision::Refused { aggregate, reason } => {
state.aggregate = Some(aggregate);
RemoteLostOperationResolution::Refused { reason }
}
};
checkpoint_state(&mut state)?;
Ok(outcome)
}
pub fn take_restored_operation_abandonment(
&self,
) -> Result<
Option<liminal_protocol::client::RestoredExpectedOperationAbandonment>,
RemoteParticipantError,
> {
let mut state = self.state.lock();
let mut aggregate = take_aggregate(&mut state)?;
let abandonment = aggregate.take_restored_operation_abandonment();
if abandonment.is_some() {
persist(&mut state.store, &aggregate)?;
}
state.aggregate = Some(aggregate);
Ok(abandonment)
}
pub fn record_transport_fate(
&self,
) -> Result<RemoteReconnectPermitOutcome, RemoteParticipantError> {
let mut state = self.state.lock();
record_connection_fate(&mut state)
}
pub fn record_online_transition(
&self,
) -> Result<RemoteReconnectPermitOutcome, RemoteParticipantError> {
self.record_fresh_reconnect(|aggregate| {
record_online_transition(aggregate, ProvedOnlineTransition::ProvedOnline)
})
}
pub fn record_explicit_reconnect(
&self,
) -> Result<RemoteReconnectPermitOutcome, RemoteParticipantError> {
self.record_fresh_reconnect(|aggregate| {
record_explicit_reconnect(aggregate, ExplicitReconnectAction::ReconnectNow)
})
}
fn record_fresh_reconnect(
&self,
decide: impl FnOnce(
liminal_protocol::client::ClientParticipantAggregate,
) -> ReconnectPermitDecision,
) -> Result<RemoteReconnectPermitOutcome, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
let outcome = match decide(aggregate) {
ReconnectPermitDecision::Permitted {
aggregate,
permit,
result,
} => {
state.aggregate = Some(aggregate);
RemoteReconnectPermitOutcome::Permitted {
permit: RemoteReconnectPermit { permit },
result,
}
}
ReconnectPermitDecision::Refused(refusal) => {
let reason = refusal.reason();
let result = refusal.result();
let (aggregate, _) = refusal.into_parts();
state.aggregate = Some(aggregate);
RemoteReconnectPermitOutcome::Refused { reason, result }
}
};
checkpoint_state(&mut state)?;
Ok(outcome)
}
pub fn recover_reconnect_permit(
&self,
) -> Result<RemoteReconnectPermitRecovery, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
match recover_reconnect_permit(aggregate) {
RecoveredReconnectPermitDecision::Recovered { aggregate, permit } => {
state.aggregate = Some(aggregate);
Ok(RemoteReconnectPermitRecovery::Recovered(
RemoteReconnectPermit { permit },
))
}
RecoveredReconnectPermitDecision::NotAvailable {
aggregate,
state: value,
} => {
state.aggregate = Some(aggregate);
Ok(RemoteReconnectPermitRecovery::NotAvailable { state: value })
}
}
}
pub fn resolve_lost_reconnect_authority(
&self,
) -> Result<RemoteLostReconnectResolution, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
let outcome = match resolve_lost_reconnect_authority(aggregate) {
LostReconnectAuthorityDecision::Recorded {
aggregate,
testimony,
} => {
state.aggregate = Some(aggregate);
RemoteLostReconnectResolution::Recorded {
testimony: testimony.kind(),
}
}
LostReconnectAuthorityDecision::Refused { aggregate, reason } => {
state.aggregate = Some(aggregate);
RemoteLostReconnectResolution::Refused { reason }
}
};
checkpoint_state(&mut state)?;
Ok(outcome)
}
pub fn reconnect(
&self,
permit: RemoteReconnectPermit,
) -> Result<RemoteReconnectAttemptOutcome, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
let (aggregate, attempt) = match redeem_attempt(aggregate, permit.permit) {
ReconnectAttemptDecision::Started { aggregate, attempt } => (aggregate, attempt),
ReconnectAttemptDecision::Refused {
aggregate,
permit,
reason,
} => {
state.aggregate = Some(aggregate);
return Ok(RemoteReconnectAttemptOutcome::Refused {
permit: RemoteReconnectPermit { permit },
reason,
});
}
};
persist(&mut state.store, &aggregate)?;
state.aggregate = Some(aggregate);
drop(state);
let transport_result = self.transport.reconnect_participant(&self.server_address);
let fate = if transport_result.is_ok() {
ReconnectAttemptFate::Connected
} else {
ReconnectAttemptFate::Failed
};
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
match record_attempt_fate(aggregate, attempt, fate) {
ReconnectAttemptFateDecision::Recorded(aggregate) => {
persist(&mut state.store, &aggregate)?;
state.aggregate = Some(aggregate);
match transport_result {
Ok(provenance) => Ok(RemoteReconnectAttemptOutcome::Connected { provenance }),
Err(error) => Ok(RemoteReconnectAttemptOutcome::Failed { error }),
}
}
ReconnectAttemptFateDecision::Refused {
aggregate,
attempt,
reason,
..
} => {
state.aggregate = Some(aggregate);
state.reconnect_attempt = Some(attempt);
let (provenance, error) = match transport_result {
Ok(value) => (Some(value), None),
Err(value) => (None, Some(value)),
};
Ok(RemoteReconnectAttemptOutcome::FateRefused {
reason,
error,
provenance,
})
}
}
}
pub fn record_established_transport_loss(
&self,
) -> Result<RemoteTransportLossOutcome, RemoteParticipantError> {
let mut state = self.state.lock();
let operation_fate = if let Some(correlation) = state.correlation.take() {
let aggregate = take_aggregate(&mut state)?;
record_operation_transport_fate(&mut state, aggregate, correlation)
} else {
RemoteOperationTransportFate::NotOutstanding
};
let reconnect = record_connection_fate(&mut state)?;
Ok(RemoteTransportLossOutcome {
operation_fate,
reconnect,
})
}
pub fn replay_detach(&self) -> Result<RemoteDetachReplayOutcome, RemoteParticipantError> {
let mut state = self.state.lock();
let aggregate = take_aggregate(&mut state)?;
let (aggregate, attempt) = match transport_attempt_started(aggregate) {
DetachTransportAttemptDecision::Started { aggregate, attempt } => (aggregate, attempt),
DetachTransportAttemptDecision::Refused(refusal) => {
let reason = refusal.reason();
let (aggregate, ()) = refusal.into_parts();
state.aggregate = Some(aggregate);
return Ok(RemoteDetachReplayOutcome::Refused { reason });
}
};
persist(&mut state.store, &aggregate)?;
let (request, correlation) = attempt.into_request();
let request = ClientRequest::Detach(DetachRequest {
conversation_id: request.conversation_id,
participant_id: request.participant_id,
capability_generation: request.capability_generation,
detach_attempt_token: request.detach_attempt_token,
});
match self
.transport
.send_participant(&self.server_address, &request)
{
Ok(provenance) => {
state.aggregate = Some(aggregate);
state.correlation = Some(correlation);
Ok(RemoteDetachReplayOutcome::Send(
RemoteParticipantSendOutcome::Sent { provenance },
))
}
Err(error) => {
let operation_fate = match transport_fate(
aggregate,
correlation,
DetachTransportFate::ResponseUnavailable,
) {
DetachTransportFateDecision::Parked(applied) => {
state.aggregate = Some(applied.into_aggregate());
RemoteOperationTransportFate::DetachParked
}
DetachTransportFateDecision::Refused(refusal) => {
let (aggregate, (correlation, _)) = refusal.into_parts();
state.aggregate = Some(aggregate);
state.correlation = Some(correlation);
RemoteOperationTransportFate::Refused {
reason: liminal_protocol::client::ExpectedOperationFateRefusalReason::DetachUsesReplayFate,
}
}
};
let reconnect = record_connection_fate(&mut state)?;
Ok(RemoteDetachReplayOutcome::Send(
RemoteParticipantSendOutcome::TransportLost {
error,
operation_fate,
reconnect,
},
))
}
}
}
}
fn checkpoint_state<S: ParticipantResumeStore>(
state: &mut super::RemoteParticipantState<S>,
) -> Result<(), RemoteParticipantError> {
let aggregate = take_aggregate(state)?;
persist(&mut state.store, &aggregate)?;
state.aggregate = Some(aggregate);
Ok(())
}