use ruststream::SubscriptionSource;
use crate::error::SeaFileError;
use crate::file::ConnectedFileBroker;
use crate::subscriber::FileSubscriber;
#[derive(Debug, Clone, PartialEq, Eq)]
#[must_use]
pub struct FileStream {
stream: String,
replay: bool,
}
impl FileStream {
pub fn new(stream: impl Into<String>) -> Self {
Self {
stream: stream.into(),
replay: false,
}
}
pub fn replay(mut self) -> Self {
self.replay = true;
self
}
#[must_use]
pub fn stream(&self) -> &str {
&self.stream
}
pub(crate) fn replay_value(&self) -> bool {
self.replay
}
pub(crate) fn validate(&self) -> Result<(), SeaFileError> {
if self.stream.is_empty() {
return Err(SeaFileError::Invalid("stream key must be non-empty".into()));
}
Ok(())
}
}
impl SubscriptionSource<ConnectedFileBroker> for FileStream {
type Subscriber = FileSubscriber;
fn name(&self) -> &str {
self.stream()
}
async fn subscribe(
self,
connected: &ConnectedFileBroker,
) -> Result<FileSubscriber, SeaFileError> {
connected.subscribe_stream(self).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn empty_stream_keys_are_rejected_before_io() {
assert!(FileStream::new("").validate().is_err());
}
#[test]
fn replay_reads_the_retained_file() {
assert!(FileStream::new("orders").replay().replay_value());
assert!(!FileStream::new("orders").replay_value());
}
}