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 { frame, .. } => f
.debug_struct("Publish")
.field("topic", &frame.topic)
.field("message_id", &frame.message_id)
.finish(),
TopicMsg::Subscribe { .. } => f.write_str("Subscribe"),
TopicMsg::Unsubscribe { id, .. } => write!(f, "Unsubscribe {{ id: {:?} }}", id),
TopicMsg::Replay { from, to, .. } => {
write!(f, "Replay {{ from: {}, to: {} }}", from, to)
}
TopicMsg::Snapshot { .. } => f.write_str("Snapshot"),
TopicMsg::HeadOffset { .. } => f.write_str("HeadOffset"),
TopicMsg::DropSink { sink_id, .. } => {
write!(f, "DropSink {{ sink_id: {} }}", sink_id)
}
TopicMsg::Shutdown { .. } => f.write_str("Shutdown"),
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum WireTopicMsg {
Publish { request_id: u32, frame: Box<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 },
}