use alloc::vec::Vec;
use super::{
ClientBindingState, ClientParticipantAggregate, DetachReplayStatus, DetachReplayTerminal,
ExpectedOperationState, LostAuthorityKind, LostAuthorityTestimony, ReconnectAggregate,
RestoredExpectedOperationAbandonment, RestoredExpectedOperationAbandonmentReason,
SdkDetachReplayAggregate, reconnect::ReconnectMachineState, replay::DetachReplayState,
};
use super::{resume_decode::decode_facts, resume_encode::encode_aggregate};
use crate::wire::{ClientRequest, CodecError};
pub(super) const MAGIC: [u8; 4] = *b"LPCR";
pub(super) const VERSION: u16 = 1;
pub(super) const HEADER_LEN: usize = 14;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ClientResumeRecordSection {
Binding,
ExpectedOperation,
DetachReplay,
Reconnect,
Abandonment,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ClientResumeRecordEncodeError {
NestedCodec {
section: ClientResumeRecordSection,
source: CodecError,
},
LengthOverflow,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ClientResumeRecordDecodeError {
Truncated {
needed: usize,
remaining: usize,
},
InvalidMagic {
presented: [u8; 4],
},
UnsupportedVersion {
presented: u16,
},
LengthMismatch {
declared: u64,
actual: usize,
},
InvalidTag {
section: ClientResumeRecordSection,
tag: u8,
},
NestedCodec {
section: ClientResumeRecordSection,
source: Option<CodecError>,
},
InvalidAbandonmentRequest {
request: crate::wire::ClientDiscriminant,
},
TrailingBytes {
remaining: usize,
},
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ClientResumeRestoreError {
BindingGenerationMismatch,
ContinuousAckOutstanding,
ReplayTerminalMismatch,
InvalidOperationAuthorization,
ExpectedBindingMismatch,
ActiveReplayExpectedDetachMismatch,
ExpectedDetachActiveReplayMismatch,
InvalidReconnectAuthorization,
LostAuthorityTestimonyMismatch,
PendingAbandonmentConflict,
CorruptRecord(ClientResumeRecordDecodeError),
}
#[derive(Debug, PartialEq, Eq)]
pub struct ClientResumeRecord {
canonical: Vec<u8>,
}
impl ClientResumeRecord {
#[must_use]
pub fn encode_canonical(&self) -> Vec<u8> {
self.canonical.clone()
}
pub fn decode_canonical(input: &[u8]) -> Result<Self, ClientResumeRecordDecodeError> {
let _ = decode_facts(input)?;
Ok(Self {
canonical: input.to_vec(),
})
}
pub fn restore(self) -> Result<ClientParticipantAggregate, ClientResumeRestoreError> {
let facts =
decode_facts(&self.canonical).map_err(ClientResumeRestoreError::CorruptRecord)?;
validate_facts(&facts)?;
let mut expected = facts.expected;
let tokenless = expected
.as_ref()
.is_some_and(|expected| matches!(expected.request, ClientRequest::ObserverRecovery(_)));
let restored_abandonment = if tokenless {
expected
.take()
.map(|expected| RestoredExpectedOperationAbandonment {
request: expected.request,
reason: RestoredExpectedOperationAbandonmentReason::TokenlessAfterCrash,
was_issued: expected.issued,
})
} else {
facts.abandonment
};
if let Some(expected) = expected.as_mut()
&& expected.issued
&& expected.lost.is_none()
{
let kind = if matches!(expected.request, ClientRequest::Detach(_)) {
LostAuthorityKind::DetachTransportAttempt
} else {
LostAuthorityKind::IssuedOperationCorrelation
};
expected.lost = Some(LostAuthorityTestimony::mint(kind));
}
let mut reconnect_lost = facts.reconnect_lost;
if reconnect_lost.is_none() {
reconnect_lost = match facts.reconnect_state {
ReconnectMachineState::Permit { issued: true, .. } => Some(
LostAuthorityTestimony::mint(LostAuthorityKind::ReconnectPermit),
),
ReconnectMachineState::Attempt { .. } => Some(LostAuthorityTestimony::mint(
LostAuthorityKind::ReconnectAttempt,
)),
ReconnectMachineState::Parked
| ReconnectMachineState::Permit { issued: false, .. }
| ReconnectMachineState::Online => None,
};
}
Ok(ClientParticipantAggregate {
binding: facts.binding,
expected,
next_operation_authorization: facts.next_operation_authorization,
detach_replay: SdkDetachReplayAggregate {
state: facts.replay,
},
reconnect: ReconnectAggregate {
state: facts.reconnect_state,
next_authorization: facts.next_authorization,
lost: reconnect_lost,
},
restored_abandonment,
})
}
}
impl ClientParticipantAggregate {
pub fn resume_record(&self) -> Result<ClientResumeRecord, ClientResumeRecordEncodeError> {
Ok(ClientResumeRecord {
canonical: encode_aggregate(self)?,
})
}
}
impl super::ClientOperationCommit {
pub fn resume_record(&self) -> Result<ClientResumeRecord, ClientResumeRecordEncodeError> {
Ok(ClientResumeRecord {
canonical: encode_aggregate(&self.aggregate)?,
})
}
}
pub(super) struct DecodedFacts {
pub(super) binding: ClientBindingState,
pub(super) next_operation_authorization: u64,
pub(super) expected: Option<ExpectedOperationState>,
pub(super) replay: DetachReplayState,
pub(super) reconnect_state: ReconnectMachineState,
pub(super) next_authorization: u64,
pub(super) reconnect_lost: Option<LostAuthorityTestimony>,
pub(super) abandonment: Option<RestoredExpectedOperationAbandonment>,
}
fn validate_facts(facts: &DecodedFacts) -> Result<(), ClientResumeRestoreError> {
if let ClientBindingState::Bound {
generation,
binding_epoch,
..
} = facts.binding
&& generation != binding_epoch.capability_generation
{
return Err(ClientResumeRestoreError::BindingGenerationMismatch);
}
if matches!(
facts.expected,
Some(ExpectedOperationState {
request: ClientRequest::ParticipantAck(_),
..
})
) {
return Err(ClientResumeRestoreError::ContinuousAckOutstanding);
}
if facts.expected.as_ref().is_some_and(|expected| {
expected.authorization == 0 || expected.authorization > facts.next_operation_authorization
}) {
return Err(ClientResumeRestoreError::InvalidOperationAuthorization);
}
if facts
.expected
.as_ref()
.is_some_and(|expected| !facts.binding.accepts_request(&expected.request))
{
return Err(ClientResumeRestoreError::ExpectedBindingMismatch);
}
let active_replay = match &facts.replay {
DetachReplayState::Recorded { request, status }
if matches!(
status,
DetachReplayStatus::Parked | DetachReplayStatus::InFlight
) =>
{
Some((request, status))
}
DetachReplayState::Empty | DetachReplayState::Recorded { .. } => None,
};
let expected_detach = facts.expected.as_ref().and_then(|expected| {
let ClientRequest::Detach(value) = &expected.request else {
return None;
};
Some((value, expected.issued))
});
match (active_replay, expected_detach) {
(Some((request, status)), Some((value, issued)))
if value.conversation_id == request.conversation_id
&& value.participant_id == request.participant_id
&& value.capability_generation == request.capability_generation
&& value.detach_attempt_token == request.detach_attempt_token
&& ((matches!(status, DetachReplayStatus::Parked) && !issued)
|| (matches!(status, DetachReplayStatus::InFlight) && issued)) => {}
(Some(_), _) => {
return Err(ClientResumeRestoreError::ActiveReplayExpectedDetachMismatch);
}
(None, Some(_)) => {
return Err(ClientResumeRestoreError::ExpectedDetachActiveReplayMismatch);
}
(None, None) => {}
}
if let DetachReplayState::Recorded {
request,
status: DetachReplayStatus::Terminal(terminal),
} = &facts.replay
&& !terminal_matches(request, terminal)
{
return Err(ClientResumeRestoreError::ReplayTerminalMismatch);
}
let authorization = match facts.reconnect_state {
ReconnectMachineState::Permit { authorization, .. }
| ReconnectMachineState::Attempt { authorization, .. } => Some(authorization),
ReconnectMachineState::Parked | ReconnectMachineState::Online => None,
};
if authorization.is_some_and(|value| value == 0 || value > facts.next_authorization) {
return Err(ClientResumeRestoreError::InvalidReconnectAuthorization);
}
validate_testimony_coupling(facts)?;
Ok(())
}
fn validate_testimony_coupling(facts: &DecodedFacts) -> Result<(), ClientResumeRestoreError> {
if let Some(expected) = facts.expected.as_ref()
&& let Some(testimony) = expected.lost.as_ref()
{
let tokenless = matches!(expected.request, ClientRequest::ObserverRecovery(_));
let expected_kind = if matches!(expected.request, ClientRequest::Detach(_)) {
LostAuthorityKind::DetachTransportAttempt
} else {
LostAuthorityKind::IssuedOperationCorrelation
};
if !expected.issued || tokenless || testimony.kind() != expected_kind {
return Err(ClientResumeRestoreError::LostAuthorityTestimonyMismatch);
}
}
if let Some(testimony) = facts.reconnect_lost.as_ref() {
let state_kind = match facts.reconnect_state {
ReconnectMachineState::Permit { issued: true, .. } => {
Some(LostAuthorityKind::ReconnectPermit)
}
ReconnectMachineState::Attempt { .. } => Some(LostAuthorityKind::ReconnectAttempt),
ReconnectMachineState::Parked
| ReconnectMachineState::Permit { issued: false, .. }
| ReconnectMachineState::Online => None,
};
if state_kind != Some(testimony.kind()) {
return Err(ClientResumeRestoreError::LostAuthorityTestimonyMismatch);
}
}
if facts.abandonment.is_some()
&& facts
.expected
.as_ref()
.is_some_and(|expected| matches!(expected.request, ClientRequest::ObserverRecovery(_)))
{
return Err(ClientResumeRestoreError::PendingAbandonmentConflict);
}
Ok(())
}
fn terminal_matches(
request: &crate::wire::DetachEnvelope,
terminal: &DetachReplayTerminal,
) -> bool {
match terminal {
DetachReplayTerminal::DetachCommitted(value) => {
value.conversation_id() == request.conversation_id
&& value.participant_id() == request.participant_id
&& value.capability_generation() == request.capability_generation
&& value.detach_attempt_token() == request.detach_attempt_token
}
DetachReplayTerminal::DetachInProgress(value) => {
let expected_generation = request.capability_generation;
let presented_generation = value.presented_generation;
let expected_token = request.detach_attempt_token;
let presented_token = value.presented_token;
value.conversation_id == request.conversation_id
&& value.participant_id == request.participant_id
&& presented_generation == expected_generation
&& presented_token == expected_token
}
DetachReplayTerminal::TerminalizedDetachCell(value) => {
value.conversation_id() == request.conversation_id
&& value.participant_id() == request.participant_id
&& value.capability_generation() == request.capability_generation
&& value.detach_attempt_token() == request.detach_attempt_token
}
}
}