use std::time::{Duration, Instant};
use serde::Serialize;
use serde::de::DeserializeOwned;
use liminal_protocol::wire::RecordAdmissionAttemptToken;
use crate::envelope::Envelope;
use crate::error::{CallError, PublishError};
use crate::id::{MessageKind, PatternId, PublicationId};
use crate::outcome::{PublicationItem, PublishReceipt};
use crate::store::ResumeStore;
use super::pump::PumpStep;
use super::state::{ConversationHandle, QueuedItem};
impl<S: ResumeStore> ConversationHandle<S> {
pub fn publish<P: Serialize>(
&mut self,
id: PublicationId,
body: &P,
) -> Result<PublishReceipt, PublishError> {
let envelope = Envelope::from_typed(
PatternId::PubSub,
id.to_correlation(),
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_with_token(bytes, RecordAdmissionAttemptToken::new(id.to_bytes()))
}
pub fn next_publication<P: DeserializeOwned>(
&mut self,
wait: Duration,
) -> Result<Option<PublicationItem<P>>, CallError> {
let started = Instant::now();
loop {
while let Some(item) = self.publications_inbox.pop_front() {
if let QueuedItem::Event { seq, .. } = &item {
let value = seq.value();
if self
.last_presented_publication
.is_some_and(|last| value <= last)
{
self.counters.duplicate_publications += 1;
continue;
}
self.last_presented_publication = Some(value);
}
return Ok(Some(convert_publication(item)));
}
if started.elapsed() >= wait {
return Ok(None);
}
match self.pump_step() {
Ok(PumpStep::Classified | PumpStep::Quiet) => {}
Err(detail) => return Err(CallError::ConnectionLost { detail }),
}
}
}
}
fn convert_publication<P: DeserializeOwned>(item: QueuedItem) -> PublicationItem<P> {
match item {
QueuedItem::Event {
envelope,
publisher,
seq,
} => {
let id = PublicationId::from_correlation(envelope.correlation);
match envelope.decode_payload::<P>() {
Ok(body) => PublicationItem::Publication {
id,
body,
publisher,
seq,
},
Err(refusal) => PublicationItem::SchemaInvalid {
id,
publisher,
seq,
detail: refusal.detail,
},
}
}
QueuedItem::Joined { peer, seq } => PublicationItem::PeerJoined { peer, seq },
QueuedItem::Departed { peer, seq, reason } => {
PublicationItem::PeerDeparted { peer, seq, reason }
}
QueuedItem::Failed { peer, seq, failure } => {
PublicationItem::PeerFailed { peer, seq, failure }
}
QueuedItem::Compacted { seq } => PublicationItem::HistoryCompacted { seq },
QueuedItem::Gap { expected, observed } => PublicationItem::Gap { expected, observed },
}
}