use std::pin::Pin;
use std::task::{Context, Poll};
use ciborium::Value as CborValue;
use futures::Stream;
use tokio::sync::mpsc;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Action {
Create,
Update,
Delete,
}
impl Action {
pub fn parse(s: &str) -> Option<Self> {
match s.to_ascii_uppercase().as_str() {
"CREATE" => Some(Action::Create),
"UPDATE" => Some(Action::Update),
"DELETE" => Some(Action::Delete),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub struct Notification {
pub query_id: String,
pub action: Action,
pub record_id: CborValue,
pub data: CborValue,
}
pub struct LiveStream {
pub(crate) query_id: String,
pub(crate) rx: mpsc::UnboundedReceiver<Notification>,
}
impl LiveStream {
pub fn query_id(&self) -> &str {
&self.query_id
}
pub async fn recv(&mut self) -> Option<Notification> {
self.rx.recv().await
}
}
impl Stream for LiveStream {
type Item = Notification;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.rx.poll_recv(cx)
}
}