use std::collections::BTreeSet;
use std::pin::Pin;
use std::task::{Context, Poll};
use futures_util::Stream;
use lash_core::{
LiveReplayGap, SessionCursor, SessionObservationEvent, SessionObservationEventPayload,
SessionReadView,
};
use crate::Result;
use crate::session::{ObservableSession, SessionObservationStream, SessionObservationStreamItem};
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct RecoverableChatEventId {
pub session_id: String,
pub cursor: String,
}
impl RecoverableChatEventId {
fn from_event(event: &SessionObservationEvent) -> Self {
Self {
session_id: event.session_id.clone(),
cursor: event.cursor.to_string(),
}
}
}
#[derive(Clone, Debug)]
pub struct RecoverableChatSnapshot {
pub read_view: SessionReadView,
pub cursor: SessionCursor,
}
impl RecoverableChatSnapshot {
pub fn capture(observable: &ObservableSession) -> Self {
let observation = observable.current_observation();
Self {
read_view: observation.read_view,
cursor: observation.cursor,
}
}
}
#[derive(Clone, Debug)]
pub enum RecoverableChatUpdate {
Event {
id: RecoverableChatEventId,
event: std::sync::Arc<SessionObservationEvent>,
},
ReplayGap {
snapshot: RecoverableChatSnapshot,
gap: LiveReplayGap,
},
TerminalReplacement {
id: RecoverableChatEventId,
event: std::sync::Arc<SessionObservationEvent>,
snapshot: RecoverableChatSnapshot,
},
}
pub struct RecoverableChatSubscription {
inner: SessionObservationStream,
applied: BTreeSet<RecoverableChatEventId>,
}
impl RecoverableChatSubscription {
pub(crate) fn new(inner: SessionObservationStream) -> Self {
Self {
inner,
applied: BTreeSet::new(),
}
}
pub fn with_applied_event_ids(
mut self,
ids: impl IntoIterator<Item = RecoverableChatEventId>,
) -> Self {
self.applied.extend(ids);
self
}
pub fn cursor(&self) -> &SessionCursor {
self.inner.cursor()
}
}
impl Stream for RecoverableChatSubscription {
type Item = Result<RecoverableChatUpdate>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
loop {
let item = match Pin::new(&mut self.inner).poll_next(cx) {
Poll::Pending => return Poll::Pending,
Poll::Ready(None) => return Poll::Ready(None),
Poll::Ready(Some(Err(err))) => return Poll::Ready(Some(Err(err))),
Poll::Ready(Some(Ok(item))) => item,
};
match item {
SessionObservationStreamItem::Gap { observation, gap } => {
self.applied.clear();
return Poll::Ready(Some(Ok(RecoverableChatUpdate::ReplayGap {
snapshot: RecoverableChatSnapshot {
read_view: observation.read_view,
cursor: observation.cursor,
},
gap,
})));
}
SessionObservationStreamItem::Event(event) => {
let id = RecoverableChatEventId::from_event(&event);
if !self.applied.insert(id.clone()) {
continue;
}
if let SessionObservationEventPayload::Committed { read_view } = &event.payload
{
self.applied.clear();
self.applied.insert(id.clone());
return Poll::Ready(Some(Ok(RecoverableChatUpdate::TerminalReplacement {
id,
snapshot: RecoverableChatSnapshot {
read_view: read_view.clone(),
cursor: event.cursor.clone(),
},
event,
})));
}
return Poll::Ready(Some(Ok(RecoverableChatUpdate::Event { id, event })));
}
}
}
}
}
impl ObservableSession {
pub fn recoverable_chat_snapshot(&self) -> RecoverableChatSnapshot {
RecoverableChatSnapshot::capture(self)
}
pub fn subscribe_recoverable_chat(&self, cursor: SessionCursor) -> RecoverableChatSubscription {
RecoverableChatSubscription::new(self.subscribe_and_recover(cursor))
}
}