1use std::collections::HashSet;
9use std::fmt;
10
11#[cfg(feature = "serde")]
12use serde::{Deserialize, Serialize};
13
14use crate::indicator::IndicatorAlert;
15use crate::timeframe::Timeframe;
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
19#[cfg_attr(
20 feature = "serde",
21 derive(Serialize, Deserialize),
22 serde(rename_all = "snake_case")
23)]
24pub enum EventPhase {
25 Setup,
26 Watch,
27 Trigger,
28 Invalidation,
29 Expiry,
30 TargetHit,
31}
32
33impl fmt::Display for EventPhase {
34 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
35 let s = match self {
36 EventPhase::Setup => "setup",
37 EventPhase::Watch => "watch",
38 EventPhase::Trigger => "trigger",
39 EventPhase::Invalidation => "invalidation",
40 EventPhase::Expiry => "expiry",
41 EventPhase::TargetHit => "target_hit",
42 };
43 f.write_str(s)
44 }
45}
46
47#[derive(Debug, Clone, PartialEq)]
56#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
57pub struct AlertEvent {
58 pub alert: IndicatorAlert,
59 pub timestamp: i64,
60 pub phase: EventPhase,
61 instrument: Option<String>,
62 timeframe: Option<Timeframe>,
63 pub event_id: String,
64}
65
66fn encode_component(s: &str) -> String {
70 format!("{}#{}", s.len(), s)
71}
72
73fn compute_event_id(
79 kind: &str,
80 phase: EventPhase,
81 timestamp: i64,
82 instrument: Option<&str>,
83 timeframe: Option<Timeframe>,
84) -> String {
85 let instrument_component = encode_component(instrument.unwrap_or(""));
86 let timeframe_component =
87 encode_component(&timeframe.map(|tf| tf.to_string()).unwrap_or_default());
88 let kind_component = encode_component(kind);
89 let phase_component = encode_component(&phase.to_string());
90 format!(
91 "{instrument_component}:{timeframe_component}:{kind_component}:{phase_component}:{timestamp}"
92 )
93}
94
95impl AlertEvent {
96 pub fn new(alert: IndicatorAlert, timestamp: i64, phase: EventPhase) -> Self {
101 let event_id = compute_event_id(&alert.kind, phase, timestamp, None, None);
102 Self {
103 alert,
104 timestamp,
105 phase,
106 instrument: None,
107 timeframe: None,
108 event_id,
109 }
110 }
111
112 pub fn instrument(&self) -> Option<&str> {
113 self.instrument.as_deref()
114 }
115
116 pub fn timeframe(&self) -> Option<Timeframe> {
117 self.timeframe
118 }
119
120 pub fn with_instrument(mut self, instrument: impl Into<String>) -> Self {
123 self.instrument = Some(instrument.into());
124 self.event_id = compute_event_id(
125 &self.alert.kind,
126 self.phase,
127 self.timestamp,
128 self.instrument.as_deref(),
129 self.timeframe,
130 );
131 self
132 }
133
134 pub fn with_timeframe(mut self, timeframe: Timeframe) -> Self {
137 self.timeframe = Some(timeframe);
138 self.event_id = compute_event_id(
139 &self.alert.kind,
140 self.phase,
141 self.timestamp,
142 self.instrument.as_deref(),
143 self.timeframe,
144 );
145 self
146 }
147}
148
149#[derive(Debug, Clone, Default)]
152pub struct AlertDeduplicator {
153 seen: HashSet<String>,
154}
155
156impl AlertDeduplicator {
157 pub fn new() -> Self {
158 Self::default()
159 }
160
161 pub fn admit(&mut self, event: &AlertEvent) -> bool {
164 self.seen.insert(event.event_id.clone())
165 }
166
167 pub fn reset(&mut self) {
168 self.seen.clear();
169 }
170
171 pub fn len(&self) -> usize {
172 self.seen.len()
173 }
174
175 pub fn is_empty(&self) -> bool {
176 self.seen.is_empty()
177 }
178}
179
180#[cfg(test)]
181mod tests {
182 use super::*;
183
184 #[test]
185 fn test_alert_event_deterministic_id() {
186 let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
187 let event_a = AlertEvent::new(alert.clone(), 1_000, EventPhase::Trigger);
188 let event_b = AlertEvent::new(alert, 1_000, EventPhase::Trigger);
189 assert_eq!(event_a.event_id, event_b.event_id);
190 }
191
192 #[test]
193 fn test_alert_event_builders() {
194 let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
195 let event = AlertEvent::new(alert, 1_000, EventPhase::Watch)
196 .with_instrument("GENERIC")
197 .with_timeframe(Timeframe::Minute(5));
198 assert_eq!(event.instrument(), Some("GENERIC"));
199 assert_eq!(event.timeframe(), Some(Timeframe::Minute(5)));
200 }
201
202 #[test]
205 fn test_event_id_distinguishes_instrument() {
206 let alert = IndicatorAlert::new("cross_up", "crossed above trigger", 0.8);
207 let dax = AlertEvent::new(alert.clone(), 1_700_000_000, EventPhase::Trigger)
208 .with_instrument("DAX")
209 .with_timeframe(Timeframe::Minute(5));
210 let es = AlertEvent::new(alert, 1_700_000_000, EventPhase::Trigger)
211 .with_instrument("ES")
212 .with_timeframe(Timeframe::Minute(5));
213
214 assert_ne!(dax.event_id, es.event_id);
215
216 let mut dedup = AlertDeduplicator::new();
217 assert!(dedup.admit(&dax));
218 assert!(
219 dedup.admit(&es),
220 "same kind/phase/timestamp on a different instrument must not be treated as a duplicate"
221 );
222 }
223
224 #[test]
226 fn test_event_id_distinguishes_timeframe() {
227 let alert = IndicatorAlert::new("cross_up", "crossed above trigger", 0.8);
228 let m5 = AlertEvent::new(alert.clone(), 1_700_000_000, EventPhase::Trigger)
229 .with_instrument("DAX")
230 .with_timeframe(Timeframe::Minute(5));
231 let m15 = AlertEvent::new(alert, 1_700_000_000, EventPhase::Trigger)
232 .with_instrument("DAX")
233 .with_timeframe(Timeframe::Minute(15));
234
235 assert_ne!(m5.event_id, m15.event_id);
236
237 let mut dedup = AlertDeduplicator::new();
238 assert!(dedup.admit(&m5));
239 assert!(
240 dedup.admit(&m15),
241 "same instrument on a different timeframe must not be treated as a duplicate"
242 );
243 }
244
245 #[test]
247 fn test_event_id_same_full_context_is_duplicate() {
248 let alert = IndicatorAlert::new("cross_up", "crossed above trigger", 0.8);
249 let first = AlertEvent::new(alert.clone(), 1_700_000_000, EventPhase::Trigger)
250 .with_instrument("DAX")
251 .with_timeframe(Timeframe::Minute(5));
252 let second = AlertEvent::new(alert, 1_700_000_000, EventPhase::Trigger)
253 .with_instrument("DAX")
254 .with_timeframe(Timeframe::Minute(5));
255
256 assert_eq!(first.event_id, second.event_id);
257
258 let mut dedup = AlertDeduplicator::new();
259 assert!(dedup.admit(&first));
260 assert!(!dedup.admit(&second));
261 }
262
263 #[test]
266 fn test_event_id_independent_of_builder_call_order() {
267 let alert = IndicatorAlert::new("cross_up", "crossed above trigger", 0.8);
268 let instrument_then_timeframe =
269 AlertEvent::new(alert.clone(), 1_700_000_000, EventPhase::Trigger)
270 .with_instrument("DAX")
271 .with_timeframe(Timeframe::Minute(5));
272 let timeframe_then_instrument = AlertEvent::new(alert, 1_700_000_000, EventPhase::Trigger)
273 .with_timeframe(Timeframe::Minute(5))
274 .with_instrument("DAX");
275
276 assert_eq!(
277 instrument_then_timeframe.event_id,
278 timeframe_then_instrument.event_id
279 );
280 }
281
282 #[test]
286 fn test_event_id_special_characters_do_not_collide() {
287 let alert_ab_c = IndicatorAlert::new("C", "note", 0.5);
290 let event_1 =
291 AlertEvent::new(alert_ab_c, 1_000, EventPhase::Trigger).with_instrument("A:B");
292
293 let alert_b_c = IndicatorAlert::new("B:C", "note", 0.5);
294 let event_2 = AlertEvent::new(alert_b_c, 1_000, EventPhase::Trigger).with_instrument("A");
295
296 assert_ne!(event_1.event_id, event_2.event_id);
297 }
298
299 #[test]
300 fn test_alert_deduplicator() {
301 let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
302 let event = AlertEvent::new(alert.clone(), 1_000, EventPhase::Trigger);
303 let mut dedup = AlertDeduplicator::new();
304
305 assert!(dedup.admit(&event));
306 assert!(
307 !dedup.admit(&event),
308 "duplicate event must not be re-admitted"
309 );
310 assert_eq!(dedup.len(), 1);
311
312 let other_bar = AlertEvent::new(alert, 1_060, EventPhase::Trigger);
313 assert!(dedup.admit(&other_bar), "distinct timestamp is a new event");
314 assert_eq!(dedup.len(), 2);
315
316 dedup.reset();
317 assert!(dedup.is_empty());
318 }
319}