Skip to main content

kestrel_chartkit/
adapters.rs

1//! Provider-neutral integration contracts: a data-feed trait and a notification-sink trait, plus
2//! dependency-free reference implementations (an in-memory feed, a logging sink, and a
3//! transport-agnostic webhook-shaped sink). Tools embedding kestrel-chartkit implement these
4//! traits with their own broker/exchange/webhook specifics; this crate never depends on a
5//! specific provider, and ships no HTTP client — [`WebhookNotificationSink`] takes the actual
6//! transport as an injected closure instead.
7
8use std::collections::HashMap;
9
10use crate::model::Bar;
11use crate::timeframe::Timeframe;
12
13/// Historical/live bar data source contract. Provider clients (broker APIs, exchange
14/// WebSockets, vendor SDKs) live in consuming applications, not here.
15pub trait DataFeedAdapter {
16    type Error;
17
18    /// Fetches historical bars for `symbol` at `timeframe` within `[from, to]` (Unix seconds).
19    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    /// Registers interest in live updates for `symbol`/`timeframe`. A no-op for adapters that
28    /// only ever serve historical data.
29    fn subscribe_live(&mut self, symbol: &str, timeframe: Timeframe) -> Result<(), Self::Error>;
30
31    /// Returns any new bars received since the last call, for all subscribed symbols.
32    fn poll_live(&mut self) -> Result<Vec<Bar>, Self::Error>;
33}
34
35/// A fully generic, dependency-free [`DataFeedAdapter`]: serves bars from an in-memory table.
36/// Useful for backtests, replay, and tests — not tied to any provider.
37#[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    /// Loads (replacing any existing) bars for `symbol`. Bars are assumed sorted ascending by
50    /// timestamp.
51    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    /// "Live" bars for an in-memory feed are simply whatever was `load`ed beyond the cursor
92    /// established at subscription time — a deterministic stand-in for a real push feed, useful
93    /// for testing consumers of [`DataFeedAdapter`] without a network dependency.
94    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
124/// Notification-channel contract, decoupled from any specific transport (webhook, email, push,
125/// broker action). Deterministic chart-kit events (alerts, scenario transitions, fills) are the
126/// producer; the transport is entirely the consumer's concern.
127pub trait NotificationSink {
128    type Error;
129
130    fn notify(&mut self, event: &NotificationEvent) -> Result<(), Self::Error>;
131}
132
133/// A fully generic, dependency-free [`NotificationSink`]: appends every event to an in-memory
134/// log. Useful for tests, or as a base to wrap with a real transport.
135#[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
155/// A transport-agnostic webhook-shaped [`NotificationSink`]: holds a destination `url` and an
156/// injected `send` closure, so this crate needs no HTTP client dependency — works with any
157/// webhook-style endpoint (Slack, Discord, a generic HTTP receiver); the consumer supplies the
158/// actual transport.
159pub 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}