Skip to main content

kestrel_chartkit/
event.rs

1//! Event and alert enrichment model.
2//!
3//! [`crate::indicator::IndicatorAlert`] stays a lightweight per-bar signal (kind/note/strength)
4//! so every existing indicator keeps constructing it unchanged. [`AlertEvent`] is the layer a
5//! runner/composition graph wraps around it once timestamp, instrument, and timeframe context is
6//! known, adding a stable event ID, an explicit state phase, and deduplication support.
7
8use 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/// Lifecycle phase of a tracked setup/alert.
18#[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/// An [`IndicatorAlert`] enriched with timestamp, instrument, timeframe, a stable event ID, and
48/// an explicit lifecycle phase, suitable for deduplication and cross-bar tracking.
49#[derive(Debug, Clone, PartialEq)]
50#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
51pub struct AlertEvent {
52    pub alert: IndicatorAlert,
53    pub timestamp: i64,
54    pub phase: EventPhase,
55    pub instrument: Option<String>,
56    pub timeframe: Option<Timeframe>,
57    pub event_id: String,
58}
59
60impl AlertEvent {
61    /// Builds an event with a deterministic ID derived from the alert kind, phase, and timestamp,
62    /// so re-emitting the same logical event on the same bar always yields the same ID.
63    pub fn new(alert: IndicatorAlert, timestamp: i64, phase: EventPhase) -> Self {
64        let event_id = format!("{}:{}:{}", alert.kind, phase, timestamp);
65        Self {
66            alert,
67            timestamp,
68            phase,
69            instrument: None,
70            timeframe: None,
71            event_id,
72        }
73    }
74
75    pub fn with_instrument(mut self, instrument: impl Into<String>) -> Self {
76        self.instrument = Some(instrument.into());
77        self
78    }
79
80    pub fn with_timeframe(mut self, timeframe: Timeframe) -> Self {
81        self.timeframe = Some(timeframe);
82        self
83    }
84}
85
86/// Tracks event IDs already admitted so repeated recalculation (e.g. rollback/idempotent replay
87/// of the same bar) does not re-emit duplicate alerts downstream.
88#[derive(Debug, Clone, Default)]
89pub struct AlertDeduplicator {
90    seen: HashSet<String>,
91}
92
93impl AlertDeduplicator {
94    pub fn new() -> Self {
95        Self::default()
96    }
97
98    /// Returns `true` and records the event if its ID has not been seen before; returns `false`
99    /// without side effects if it is a duplicate.
100    pub fn admit(&mut self, event: &AlertEvent) -> bool {
101        self.seen.insert(event.event_id.clone())
102    }
103
104    pub fn reset(&mut self) {
105        self.seen.clear();
106    }
107
108    pub fn len(&self) -> usize {
109        self.seen.len()
110    }
111
112    pub fn is_empty(&self) -> bool {
113        self.seen.is_empty()
114    }
115}
116
117#[cfg(test)]
118mod tests {
119    use super::*;
120
121    #[test]
122    fn test_alert_event_deterministic_id() {
123        let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
124        let event_a = AlertEvent::new(alert.clone(), 1_000, EventPhase::Trigger);
125        let event_b = AlertEvent::new(alert, 1_000, EventPhase::Trigger);
126        assert_eq!(event_a.event_id, event_b.event_id);
127    }
128
129    #[test]
130    fn test_alert_event_builders() {
131        let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
132        let event = AlertEvent::new(alert, 1_000, EventPhase::Watch)
133            .with_instrument("GENERIC")
134            .with_timeframe(Timeframe::Minute(5));
135        assert_eq!(event.instrument.as_deref(), Some("GENERIC"));
136        assert_eq!(event.timeframe, Some(Timeframe::Minute(5)));
137    }
138
139    #[test]
140    fn test_alert_deduplicator() {
141        let alert = IndicatorAlert::new("cross_up", "RSI crossed above 70", 0.8);
142        let event = AlertEvent::new(alert.clone(), 1_000, EventPhase::Trigger);
143        let mut dedup = AlertDeduplicator::new();
144
145        assert!(dedup.admit(&event));
146        assert!(
147            !dedup.admit(&event),
148            "duplicate event must not be re-admitted"
149        );
150        assert_eq!(dedup.len(), 1);
151
152        let other_bar = AlertEvent::new(alert, 1_060, EventPhase::Trigger);
153        assert!(dedup.admit(&other_bar), "distinct timestamp is a new event");
154        assert_eq!(dedup.len(), 2);
155
156        dedup.reset();
157        assert!(dedup.is_empty());
158    }
159}