use bytes::Bytes;
use tokio::sync::oneshot;
use crate::broker::broker::PublishOutcome;
use crate::broker::fanout::{ConnectionSink, SubscribeIntent, SubscriptionId};
use crate::error::Result;
use crate::frame::Frame;
use crate::storage::StoredSnapshot;
pub enum TopicMsg {
Publish {
frame: Frame,
reply_to: oneshot::Sender<Result<PublishOutcome>>,
},
Subscribe {
sink: ConnectionSink,
intent: SubscribeIntent,
reply_to: oneshot::Sender<Result<SubscriptionId>>,
},
Unsubscribe {
id: SubscriptionId,
reply_to: oneshot::Sender<Result<bool>>,
},
Replay {
from: i64,
to: i64,
reply_to: oneshot::Sender<Result<Vec<Bytes>>>,
},
Snapshot {
reply_to: oneshot::Sender<Result<Option<StoredSnapshot>>>,
},
HeadOffset { reply_to: oneshot::Sender<i64> },
DropSink {
sink_id: u64,
reply_to: oneshot::Sender<usize>,
},
Shutdown { reply_to: oneshot::Sender<()> },
}
impl std::fmt::Debug for TopicMsg {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
TopicMsg::Publish { .. } => f.write_str("Publish"),
TopicMsg::Subscribe { .. } => f.write_str("Subscribe"),
TopicMsg::Unsubscribe { .. } => f.write_str("Unsubscribe"),
TopicMsg::Replay { .. } => f.write_str("Replay"),
TopicMsg::Snapshot { .. } => f.write_str("Snapshot"),
TopicMsg::HeadOffset { .. } => f.write_str("HeadOffset"),
TopicMsg::DropSink { .. } => f.write_str("DropSink"),
TopicMsg::Shutdown { .. } => f.write_str("Shutdown"),
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum WireTopicMsg {
Publish {
request_id: u32,
frame: Frame,
},
Subscribe {
request_id: u32,
intent: SubscribeIntent,
sink_id: u64,
},
Unsubscribe {
request_id: u32,
id: u64,
},
Replay {
request_id: u32,
from: i64,
to: i64,
},
Snapshot {
request_id: u32,
},
HeadOffset {
request_id: u32,
},
DropSink {
request_id: u32,
sink_id: u64,
},
Shutdown {
request_id: u32,
},
}