ruststream_sea_file/
stream.rs1use ruststream::SubscriptionSource;
4
5use crate::error::SeaFileError;
6use crate::file::ConnectedFileBroker;
7use crate::subscriber::FileSubscriber;
8
9#[derive(Debug, Clone, PartialEq, Eq)]
28#[must_use]
29pub struct FileStream {
30 stream: String,
31 replay: bool,
32}
33
34impl FileStream {
35 pub fn new(stream: impl Into<String>) -> Self {
37 Self {
38 stream: stream.into(),
39 replay: false,
40 }
41 }
42
43 pub fn replay(mut self) -> Self {
46 self.replay = true;
47 self
48 }
49
50 #[must_use]
52 pub fn stream(&self) -> &str {
53 &self.stream
54 }
55
56 pub(crate) fn replay_value(&self) -> bool {
57 self.replay
58 }
59
60 pub(crate) fn validate(&self) -> Result<(), SeaFileError> {
62 if self.stream.is_empty() {
63 return Err(SeaFileError::Invalid("stream key must be non-empty".into()));
64 }
65 Ok(())
66 }
67}
68
69impl SubscriptionSource<ConnectedFileBroker> for FileStream {
70 type Subscriber = FileSubscriber;
71
72 fn name(&self) -> &str {
73 self.stream()
74 }
75
76 async fn subscribe(
77 self,
78 connected: &ConnectedFileBroker,
79 ) -> Result<FileSubscriber, SeaFileError> {
80 connected.subscribe_stream(self).await
81 }
82}
83
84#[cfg(test)]
85mod tests {
86 use super::*;
87
88 #[test]
89 fn empty_stream_keys_are_rejected_before_io() {
90 assert!(FileStream::new("").validate().is_err());
91 }
92
93 #[test]
94 fn replay_reads_the_retained_file() {
95 assert!(FileStream::new("orders").replay().replay_value());
96 assert!(!FileStream::new("orders").replay_value());
97 }
98}