use std::collections::{HashMap, VecDeque};
use std::time::Duration;
use frame_core::error::FailureReason;
use liminal_protocol::wire::Generation;
use liminal_sdk::remote::RemoteParticipantHandle;
use crate::anomaly::{Anomaly, AnomalyCounters};
use crate::envelope::Envelope;
use crate::id::{ConversationId, ConversationSeq, CorrelationId, ParticipantRef};
use crate::outcome::{DepartReason, JoinGrant};
use crate::seam::StoreAdapter;
use crate::store::ResumeStore;
use crate::track::{CursorTracker, ExchangeRegistry};
#[derive(Debug)]
pub(crate) struct HeldReply {
pub(crate) envelope: Envelope,
pub(crate) responder: ParticipantRef,
pub(crate) seq: ConversationSeq,
}
#[derive(Debug)]
pub(crate) struct HeldRequest {
pub(crate) envelope: Envelope,
pub(crate) requester: ParticipantRef,
pub(crate) seq: ConversationSeq,
}
#[derive(Debug, Clone)]
pub(crate) enum QueuedItem {
Event {
envelope: Envelope,
publisher: ParticipantRef,
seq: ConversationSeq,
},
Joined {
peer: ParticipantRef,
seq: ConversationSeq,
},
Departed {
peer: ParticipantRef,
seq: ConversationSeq,
reason: DepartReason,
},
Failed {
peer: ParticipantRef,
seq: ConversationSeq,
failure: FailureReason,
},
Compacted {
seq: ConversationSeq,
},
Gap {
expected: ConversationSeq,
observed: ConversationSeq,
},
}
#[derive(Debug, Clone)]
pub(crate) struct DeathNote {
pub(crate) peer: ParticipantRef,
pub(crate) failure: FailureReason,
pub(crate) seq: ConversationSeq,
}
pub struct ConversationHandle<S: ResumeStore> {
pub(crate) sdk: RemoteParticipantHandle<StoreAdapter<S>>,
pub(crate) conversation: ConversationId,
pub(crate) me: ParticipantRef,
pub(crate) grant: JoinGrant,
pub(crate) generation: Generation,
pub(crate) answer_window: Duration,
pub(crate) exchanges: ExchangeRegistry,
pub(crate) replies: HashMap<CorrelationId, HeldReply>,
pub(crate) requests_inbox: VecDeque<HeldRequest>,
pub(crate) events_inbox: VecDeque<QueuedItem>,
pub(crate) publications_inbox: VecDeque<QueuedItem>,
pub(crate) last_presented_publication: Option<u64>,
pub(crate) deaths: Vec<DeathNote>,
pub(crate) anomalies: Vec<Anomaly>,
pub(crate) counters: AnomalyCounters,
pub(crate) tracker: CursorTracker,
pub(crate) attached: bool,
}
impl<S: ResumeStore> std::fmt::Debug for ConversationHandle<S> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("ConversationHandle")
.field("conversation", &self.conversation)
.field("participant", &self.me)
.field("attached", &self.attached)
.finish_non_exhaustive()
}
}
impl<S: ResumeStore> ConversationHandle<S> {
#[must_use]
pub const fn conversation(&self) -> ConversationId {
self.conversation
}
#[must_use]
pub const fn participant(&self) -> ParticipantRef {
self.me
}
#[must_use]
pub const fn grant(&self) -> &JoinGrant {
&self.grant
}
#[must_use]
pub const fn attached(&self) -> bool {
self.attached
}
pub fn drain_anomalies(&mut self) -> Vec<Anomaly> {
std::mem::take(&mut self.anomalies)
}
#[must_use]
pub const fn anomaly_counters(&self) -> AnomalyCounters {
self.counters
}
}