1use std::collections::HashMap;
9
10use crate::model::Bar;
11use crate::timeframe::Timeframe;
12
13pub trait DataFeedAdapter {
16 type Error;
17
18 fn fetch_historical(
20 &mut self,
21 symbol: &str,
22 timeframe: Timeframe,
23 from: i64,
24 to: i64,
25 ) -> Result<Vec<Bar>, Self::Error>;
26
27 fn subscribe_live(&mut self, symbol: &str, timeframe: Timeframe) -> Result<(), Self::Error>;
30
31 fn poll_live(&mut self) -> Result<Vec<Bar>, Self::Error>;
33}
34
35#[derive(Debug, Clone, Default)]
38pub struct InMemoryDataFeed {
39 bars_by_symbol: HashMap<String, Vec<Bar>>,
40 live_cursor: HashMap<String, usize>,
41 subscribed: Vec<String>,
42}
43
44impl InMemoryDataFeed {
45 pub fn new() -> Self {
46 Self::default()
47 }
48
49 pub fn load(&mut self, symbol: impl Into<String>, bars: Vec<Bar>) {
52 self.bars_by_symbol.insert(symbol.into(), bars);
53 }
54}
55
56impl DataFeedAdapter for InMemoryDataFeed {
57 type Error = String;
58
59 fn fetch_historical(
60 &mut self,
61 symbol: &str,
62 _timeframe: Timeframe,
63 from: i64,
64 to: i64,
65 ) -> Result<Vec<Bar>, Self::Error> {
66 let bars = self
67 .bars_by_symbol
68 .get(symbol)
69 .ok_or_else(|| format!("no bars loaded for symbol '{symbol}'"))?;
70 Ok(bars
71 .iter()
72 .filter(|b| b.timestamp >= from && b.timestamp <= to)
73 .cloned()
74 .collect())
75 }
76
77 fn subscribe_live(&mut self, symbol: &str, _timeframe: Timeframe) -> Result<(), Self::Error> {
78 if !self.subscribed.contains(&symbol.to_string()) {
79 self.subscribed.push(symbol.to_string());
80 self.live_cursor.insert(
81 symbol.to_string(),
82 self.bars_by_symbol
83 .get(symbol)
84 .map(|b| b.len())
85 .unwrap_or(0),
86 );
87 }
88 Ok(())
89 }
90
91 fn poll_live(&mut self) -> Result<Vec<Bar>, Self::Error> {
95 let mut new_bars = Vec::new();
96 for symbol in &self.subscribed {
97 let cursor = self.live_cursor.get(symbol).copied().unwrap_or(0);
98 if let Some(bars) = self.bars_by_symbol.get(symbol) {
99 if cursor < bars.len() {
100 new_bars.extend(bars[cursor..].iter().cloned());
101 self.live_cursor.insert(symbol.clone(), bars.len());
102 }
103 }
104 }
105 Ok(new_bars)
106 }
107}
108
109#[derive(Debug, Clone, Copy, PartialEq, Eq)]
110pub enum NotificationSeverity {
111 Info,
112 Warning,
113 Critical,
114}
115
116#[derive(Debug, Clone, PartialEq)]
117pub struct NotificationEvent {
118 pub timestamp: i64,
119 pub severity: NotificationSeverity,
120 pub title: String,
121 pub body: String,
122}
123
124pub trait NotificationSink {
128 type Error;
129
130 fn notify(&mut self, event: &NotificationEvent) -> Result<(), Self::Error>;
131}
132
133#[derive(Debug, Clone, Default)]
136pub struct LoggingNotificationSink {
137 pub events: Vec<NotificationEvent>,
138}
139
140impl LoggingNotificationSink {
141 pub fn new() -> Self {
142 Self::default()
143 }
144}
145
146impl NotificationSink for LoggingNotificationSink {
147 type Error = std::convert::Infallible;
148
149 fn notify(&mut self, event: &NotificationEvent) -> Result<(), Self::Error> {
150 self.events.push(event.clone());
151 Ok(())
152 }
153}
154
155pub struct WebhookNotificationSink<F>
160where
161 F: FnMut(&str, &NotificationEvent) -> Result<(), String>,
162{
163 pub url: String,
164 send: F,
165}
166
167impl<F> WebhookNotificationSink<F>
168where
169 F: FnMut(&str, &NotificationEvent) -> Result<(), String>,
170{
171 pub fn new(url: impl Into<String>, send: F) -> Self {
172 Self {
173 url: url.into(),
174 send,
175 }
176 }
177}
178
179impl<F> NotificationSink for WebhookNotificationSink<F>
180where
181 F: FnMut(&str, &NotificationEvent) -> Result<(), String>,
182{
183 type Error = String;
184
185 fn notify(&mut self, event: &NotificationEvent) -> Result<(), Self::Error> {
186 (self.send)(&self.url, event)
187 }
188}
189
190#[cfg(test)]
191mod tests {
192 use super::*;
193
194 fn bar(ts: i64) -> Bar {
195 Bar::new(ts, 100.0, 101.0, 99.0, 100.0, 10.0)
196 }
197
198 #[test]
199 fn test_in_memory_feed_fetch_historical_filters_range() {
200 let mut feed = InMemoryDataFeed::new();
201 feed.load("TEST", vec![bar(0), bar(60), bar(120), bar(180)]);
202
203 let result = feed
204 .fetch_historical("TEST", Timeframe::Minute(1), 60, 120)
205 .unwrap();
206 assert_eq!(result.len(), 2);
207 assert_eq!(result[0].timestamp, 60);
208 }
209
210 #[test]
211 fn test_in_memory_feed_fetch_unknown_symbol_errors() {
212 let mut feed = InMemoryDataFeed::new();
213 assert!(feed
214 .fetch_historical("NOPE", Timeframe::Minute(1), 0, 100)
215 .is_err());
216 }
217
218 #[test]
219 fn test_in_memory_feed_poll_live_only_returns_new_bars() {
220 let mut feed = InMemoryDataFeed::new();
221 feed.load("TEST", vec![bar(0), bar(60)]);
222 feed.subscribe_live("TEST", Timeframe::Minute(1)).unwrap();
223
224 let first_poll = feed.poll_live().unwrap();
225 assert!(
226 first_poll.is_empty(),
227 "no bars arrived after the subscription cursor yet"
228 );
229
230 feed.load("TEST", vec![bar(0), bar(60), bar(120)]);
231 let second_poll = feed.poll_live().unwrap();
232 assert_eq!(second_poll.len(), 1);
233 assert_eq!(second_poll[0].timestamp, 120);
234 }
235
236 #[test]
237 fn test_logging_sink_records_events() {
238 let mut sink = LoggingNotificationSink::new();
239 let event = NotificationEvent {
240 timestamp: 0,
241 severity: NotificationSeverity::Warning,
242 title: "test".to_string(),
243 body: "body".to_string(),
244 };
245 sink.notify(&event).unwrap();
246 assert_eq!(sink.events.len(), 1);
247 assert_eq!(sink.events[0].severity, NotificationSeverity::Warning);
248 }
249
250 #[test]
251 fn test_webhook_sink_delegates_to_injected_transport() {
252 let mut received: Vec<(String, String)> = Vec::new();
253 let mut sink = WebhookNotificationSink::new("https://example.test/hook", |url, event| {
254 received.push((url.to_string(), event.title.clone()));
255 Ok(())
256 });
257
258 let event = NotificationEvent {
259 timestamp: 0,
260 severity: NotificationSeverity::Critical,
261 title: "alert".to_string(),
262 body: "body".to_string(),
263 };
264 sink.notify(&event).unwrap();
265
266 assert_eq!(received.len(), 1);
267 assert_eq!(received[0].0, "https://example.test/hook");
268 assert_eq!(received[0].1, "alert");
269 }
270
271 #[test]
272 fn test_webhook_sink_propagates_transport_errors() {
273 let mut sink = WebhookNotificationSink::new("https://example.test/hook", |_url, _event| {
274 Err("network unreachable".to_string())
275 });
276 let event = NotificationEvent {
277 timestamp: 0,
278 severity: NotificationSeverity::Info,
279 title: "x".to_string(),
280 body: String::new(),
281 };
282 assert!(sink.notify(&event).is_err());
283 }
284}