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///
50/// `instrument` and `timeframe` are private (with [`AlertEvent::instrument`]/
51/// [`AlertEvent::timeframe`] accessors) so that [`AlertEvent::with_instrument`]/
52/// [`AlertEvent::with_timeframe`] are the only way to change them, guaranteeing `event_id` is
53/// always recomputed from the complete, current identity context and can never go stale relative
54/// to it — regardless of the order those builders are called in.
55#[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
66/// Encodes `s` with an explicit byte-length prefix so that a delimiter character occurring inside
67/// `s` (e.g. a `:` in an instrument or alert-kind name) cannot be mistaken for a component
68/// boundary when concatenated with other encoded components.
69fn encode_component(s: &str) -> String {
70    format!("{}#{}", s.len(), s)
71}
72
73/// Computes the full identity-based event ID from every field that distinguishes one logical
74/// event from another: instrument, timeframe, alert kind, lifecycle phase, and bar timestamp. A
75/// missing instrument/timeframe is encoded as an explicit empty component (distinct from any
76/// non-empty value), so "no instrument" is a stable, singular identity rather than being
77/// indistinguishable from some in-band placeholder string.
78fn 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    /// Builds an event with a deterministic ID derived from the alert kind, phase, and timestamp
97    /// (instrument/timeframe default to absent; attach them via [`AlertEvent::with_instrument`]/
98    /// [`AlertEvent::with_timeframe`]), so re-emitting the same logical event on the same bar
99    /// always yields the same ID.
100    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    /// Attaches an instrument and recomputes `event_id` from the complete identity context, so
121    /// two events that only differ by instrument never collide.
122    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    /// Attaches a timeframe and recomputes `event_id` from the complete identity context, so two
135    /// events that only differ by timeframe never collide.
136    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/// Tracks event IDs already admitted so repeated recalculation (e.g. rollback/idempotent replay
150/// of the same bar) does not re-emit duplicate alerts downstream.
151#[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    /// Returns `true` and records the event if its ID has not been seen before; returns `false`
162    /// without side effects if it is a duplicate.
163    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    /// Finding 03: events that differ only by instrument must both be admitted by a shared
203    /// deduplicator, not collide.
204    #[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    /// Finding 03: events that differ only by timeframe must both be admitted.
225    #[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    /// A fully identical context (including instrument and timeframe) must still be deduplicated.
246    #[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    /// The final ID must not depend on the order `with_instrument`/`with_timeframe` were called
264    /// in.
265    #[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    /// A delimiter character occurring inside an instrument or alert-kind name must not create an
283    /// ID collision with an otherwise-different event, since components are length-prefixed
284    /// rather than joined with a bare separator.
285    #[test]
286    fn test_event_id_special_characters_do_not_collide() {
287        // Without length-prefixing, instrument "A:B" + kind "C" could collide with instrument "A"
288        // + kind "B:C" once joined with ':'.
289        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}