use std::time::{Duration, Instant};
use liminal_protocol::wire::{ClientRequest, ParticipantAck, ServerValue};
use serde::Serialize;
use serde::de::DeserializeOwned;
use crate::envelope::Envelope;
use crate::error::{CallError, CursorCommit, PublishError};
use crate::id::{ConversationSeq, CorrelationId, MessageKind, PatternId};
use crate::outcome::{PublishReceipt, SubscriptionItem};
use crate::seam::classify_refusal;
use crate::store::ResumeStore;
use super::pump::{PumpStep, submit_to_publish};
use super::state::{ConversationHandle, QueuedItem};
impl<S: ResumeStore> ConversationHandle<S> {
pub fn publish_event<E: Serialize>(
&mut self,
body: &E,
) -> Result<PublishReceipt, PublishError> {
let envelope = Envelope::from_typed(
PatternId::Subscription,
CorrelationId::mint(),
MessageKind::Event,
body,
)
.map_err(|refusal| PublishError::SchemaInvalid {
detail: refusal.detail,
})?;
let bytes = envelope
.encode()
.map_err(|refusal| PublishError::Protocol {
detail: refusal.detail,
})?;
self.admit(bytes)
}
pub fn next_event<E: DeserializeOwned>(
&mut self,
wait: Duration,
) -> Result<Option<SubscriptionItem<E>>, CallError> {
let started = Instant::now();
loop {
if let Some(item) = self.events_inbox.pop_front() {
return Ok(Some(convert_item(item)));
}
if started.elapsed() >= wait {
return Ok(None);
}
match self.pump_step() {
Ok(PumpStep::Classified | PumpStep::Quiet) => {}
Err(detail) => return Err(CallError::ConnectionLost { detail }),
}
}
}
pub fn commit_cursor(
&mut self,
through: ConversationSeq,
) -> Result<CursorCommit, PublishError> {
let answer = self
.submit(ClientRequest::ParticipantAck(ParticipantAck {
conversation_id: self.conversation.to_wire(),
participant_id: self.me.to_wire(),
capability_generation: self.generation,
through_seq: through.value(),
}))
.map_err(submit_to_publish)?;
match answer {
ServerValue::AckCommitted(committed) => Ok(CursorCommit::Committed {
through: ConversationSeq::new(committed.current_cursor()),
}),
ServerValue::AckNoOp(noop) => Ok(CursorCommit::NoOp {
current: ConversationSeq::new(noop.current_cursor()),
}),
ServerValue::AckGap(gap) => Ok(CursorCommit::Gap {
requested: ConversationSeq::new(gap.request().through_seq),
current: ConversationSeq::new(gap.current_cursor()),
}),
ServerValue::AckRegression(regression) => Ok(CursorCommit::Regression {
requested: ConversationSeq::new(regression.request().through_seq),
}),
other => {
let (class, detail) = classify_refusal(&other);
Err(PublishError::Refused { class, detail })
}
}
}
}
fn convert_item<E: DeserializeOwned>(item: QueuedItem) -> SubscriptionItem<E> {
match item {
QueuedItem::Event {
envelope,
publisher,
seq,
} => match envelope.decode_payload::<E>() {
Ok(body) => SubscriptionItem::Event {
body,
publisher,
seq,
},
Err(refusal) => SubscriptionItem::SchemaInvalid {
publisher,
seq,
detail: refusal.detail,
},
},
QueuedItem::Joined { peer, seq } => SubscriptionItem::PeerJoined { peer, seq },
QueuedItem::Departed { peer, seq, reason } => {
SubscriptionItem::PeerDeparted { peer, seq, reason }
}
QueuedItem::Failed { peer, seq, failure } => {
SubscriptionItem::PeerFailed { peer, seq, failure }
}
QueuedItem::Compacted { seq } => SubscriptionItem::HistoryCompacted { seq },
QueuedItem::Gap { expected, observed } => SubscriptionItem::Gap { expected, observed },
}
}