Skip to main content

kestrel_chartkit/indicator/
smart_money_structure.rs

1//! Linked smart-money structure: named BSL/SSL liquidity pools (from Equal High/Low clusters)
2//! with an explicit Stop-Hunt-vs-Breakout-vs-Reclaim classification, persistent FVG zones with
3//! fill tracking, and a lightweight correlator that links BOS/CHOCH, liquidity, order-block, and
4//! FVG-fill events across the crate's *existing*, independently-running detectors
5//! ([`super::bos_choch::BosChochEngine`], [`super::liquidity_sweeps::LiquiditySweepEngine`],
6//! [`super::order_block::OrderBlockEngine`], [`super::liquidity_fvg::LiquidityFvgEngine`]) —
7//! rather than reimplementing their detection logic, this observes their
8//! [`super::IndicatorAlert`] streams (every indicator already exposes one) and reports when
9//! several independently corroborate the same bar.
10
11use std::collections::VecDeque;
12
13use crate::model::Bar;
14
15use super::{Indicator, IndicatorAlert, IndicatorOutput};
16
17// ---------------------------------------------------------------------------------------------
18// BSL/SSL liquidity pools
19// ---------------------------------------------------------------------------------------------
20
21/// Buy-side liquidity (resting above an Equal-High cluster) or sell-side liquidity (resting below
22/// an Equal-Low cluster).
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub enum LiquidityPoolKind {
25    Bsl,
26    Ssl,
27}
28
29/// Lifecycle state of a [`LiquidityPool`]: explicitly distinguishes a stop hunt (price pierced
30/// the pool then closed back inside — liquidity taken, no sustained follow-through) from a
31/// breakout (price pierced and closed beyond, i.e. sustained), and a subsequent reclaim (a prior
32/// breakout later reverses back through the level).
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34pub enum LiquidityPoolState {
35    Active,
36    StopHunted,
37    BrokenThrough,
38    Reclaimed,
39}
40
41#[derive(Debug, Clone, PartialEq)]
42pub struct LiquidityPool {
43    pub kind: LiquidityPoolKind,
44    pub price: f64,
45    /// How many near-equal pivots contributed to this pool (Equal-High/Low cluster size).
46    pub touches: u32,
47    pub formed_at: i64,
48    pub state: LiquidityPoolState,
49}
50
51/// Detects BSL/SSL pools from Equal High/Low pivot clusters and classifies every interaction as a
52/// stop hunt, a breakout, or (for a previously broken pool) a reclaim.
53pub struct LiquidityPoolEngine {
54    pivot_len: usize,
55    tolerance_pct: f64,
56    bars: VecDeque<Bar>,
57    pools: Vec<LiquidityPool>,
58    alerts: Vec<IndicatorAlert>,
59}
60
61impl LiquidityPoolEngine {
62    pub fn new(pivot_len: usize, tolerance_pct: f64) -> Self {
63        let pivot_len = pivot_len.max(2);
64        Self {
65            pivot_len,
66            tolerance_pct: tolerance_pct.max(0.001),
67            bars: VecDeque::with_capacity(pivot_len * 2 + 1),
68            pools: Vec::new(),
69            alerts: Vec::new(),
70        }
71    }
72
73    pub fn with_defaults() -> Self {
74        Self::new(5, 0.2)
75    }
76
77    pub fn pools(&self) -> &[LiquidityPool] {
78        &self.pools
79    }
80
81    fn register_pivot(&mut self, kind: LiquidityPoolKind, price: f64, timestamp: i64) {
82        let tol = self.tolerance_pct / 100.0;
83        let existing = self.pools.iter_mut().find(|p| {
84            p.kind == kind
85                && p.state == LiquidityPoolState::Active
86                && p.price != 0.0
87                && (p.price - price).abs() / p.price.abs() <= tol
88        });
89        match existing {
90            Some(pool) => {
91                pool.touches += 1;
92                pool.price = (pool.price + price) / 2.0;
93            }
94            None => self.pools.push(LiquidityPool {
95                kind,
96                price,
97                touches: 1,
98                formed_at: timestamp,
99                state: LiquidityPoolState::Active,
100            }),
101        }
102    }
103}
104
105impl Indicator for LiquidityPoolEngine {
106    fn name(&self) -> &str {
107        "liquidity_pools"
108    }
109
110    fn warmup_period(&self) -> usize {
111        self.pivot_len * 2 + 1
112    }
113
114    fn reset(&mut self) {
115        self.bars.clear();
116        self.pools.clear();
117        self.alerts.clear();
118    }
119
120    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
121        self.alerts.clear();
122
123        self.bars.push_back(bar.clone());
124        if self.bars.len() > self.pivot_len * 2 + 1 {
125            self.bars.pop_front();
126        }
127        if self.bars.len() < self.pivot_len * 2 + 1 {
128            return None;
129        }
130
131        let mid_idx = self.pivot_len;
132        let mid_bar = self.bars[mid_idx].clone();
133
134        let is_pivot_high = self
135            .bars
136            .iter()
137            .enumerate()
138            .all(|(i, b)| i == mid_idx || b.high <= mid_bar.high);
139        let is_pivot_low = self
140            .bars
141            .iter()
142            .enumerate()
143            .all(|(i, b)| i == mid_idx || b.low >= mid_bar.low);
144
145        if is_pivot_high {
146            self.register_pivot(LiquidityPoolKind::Bsl, mid_bar.high, mid_bar.timestamp);
147        }
148        if is_pivot_low {
149            self.register_pivot(LiquidityPoolKind::Ssl, mid_bar.low, mid_bar.timestamp);
150        }
151
152        for pool in &mut self.pools {
153            match (pool.kind, pool.state) {
154                (LiquidityPoolKind::Bsl, LiquidityPoolState::Active) if bar.high > pool.price => {
155                    if bar.close < pool.price {
156                        pool.state = LiquidityPoolState::StopHunted;
157                        self.alerts.push(IndicatorAlert::new(
158                            "liquidity_pool_stop_hunt",
159                            format!(
160                                "BSL pool at {:.4} swept and reclaimed (stop hunt)",
161                                pool.price
162                            ),
163                            0.85,
164                        ));
165                    } else {
166                        pool.state = LiquidityPoolState::BrokenThrough;
167                        self.alerts.push(IndicatorAlert::new(
168                            "liquidity_pool_breakout",
169                            format!(
170                                "BSL pool at {:.4} broken through (sustained breakout)",
171                                pool.price
172                            ),
173                            0.7,
174                        ));
175                    }
176                }
177                (LiquidityPoolKind::Ssl, LiquidityPoolState::Active) if bar.low < pool.price => {
178                    if bar.close > pool.price {
179                        pool.state = LiquidityPoolState::StopHunted;
180                        self.alerts.push(IndicatorAlert::new(
181                            "liquidity_pool_stop_hunt",
182                            format!(
183                                "SSL pool at {:.4} swept and reclaimed (stop hunt)",
184                                pool.price
185                            ),
186                            0.85,
187                        ));
188                    } else {
189                        pool.state = LiquidityPoolState::BrokenThrough;
190                        self.alerts.push(IndicatorAlert::new(
191                            "liquidity_pool_breakout",
192                            format!(
193                                "SSL pool at {:.4} broken through (sustained breakout)",
194                                pool.price
195                            ),
196                            0.7,
197                        ));
198                    }
199                }
200                (LiquidityPoolKind::Bsl, LiquidityPoolState::BrokenThrough)
201                    if bar.close < pool.price =>
202                {
203                    pool.state = LiquidityPoolState::Reclaimed;
204                    self.alerts.push(IndicatorAlert::new(
205                        "liquidity_pool_reclaim",
206                        format!(
207                            "BSL breakout at {:.4} reclaimed (failed breakout)",
208                            pool.price
209                        ),
210                        0.75,
211                    ));
212                }
213                (LiquidityPoolKind::Ssl, LiquidityPoolState::BrokenThrough)
214                    if bar.close > pool.price =>
215                {
216                    pool.state = LiquidityPoolState::Reclaimed;
217                    self.alerts.push(IndicatorAlert::new(
218                        "liquidity_pool_reclaim",
219                        format!(
220                            "SSL breakout at {:.4} reclaimed (failed breakout)",
221                            pool.price
222                        ),
223                        0.75,
224                    ));
225                }
226                _ => {}
227            }
228        }
229
230        let active_count = self
231            .pools
232            .iter()
233            .filter(|p| p.state == LiquidityPoolState::Active)
234            .count();
235        Some(IndicatorOutput::new(active_count as f64))
236    }
237
238    fn alerts(&self) -> Vec<IndicatorAlert> {
239        self.alerts.clone()
240    }
241}
242
243// ---------------------------------------------------------------------------------------------
244// FVG zones with fill tracking
245// ---------------------------------------------------------------------------------------------
246
247#[derive(Debug, Clone, PartialEq)]
248pub struct FvgZone {
249    pub is_bullish: bool,
250    pub top: f64,
251    pub bottom: f64,
252    pub formed_at: i64,
253    pub filled: bool,
254}
255
256/// Tracks Fair Value Gap zones as persistent objects and marks them filled once price trades back
257/// through them — the "FVG-Fill" lifecycle no existing FVG detector in this crate tracks (they
258/// only emit a one-shot creation alert).
259#[derive(Debug, Clone, Default)]
260pub struct FvgZoneTracker {
261    zones: Vec<FvgZone>,
262}
263
264impl FvgZoneTracker {
265    pub fn new() -> Self {
266        Self::default()
267    }
268
269    /// Registers a newly detected FVG (e.g. from [`super::liquidity_fvg::LiquidityFvgEngine`]'s
270    /// per-bar gap output).
271    pub fn register(&mut self, is_bullish: bool, top: f64, bottom: f64, formed_at: i64) {
272        self.zones.push(FvgZone {
273            is_bullish,
274            top,
275            bottom,
276            formed_at,
277            filled: false,
278        });
279    }
280
281    pub fn zones(&self) -> &[FvgZone] {
282        &self.zones
283    }
284
285    pub fn reset(&mut self) {
286        self.zones.clear();
287    }
288
289    /// Feeds a bar, marking any unfilled zone the bar traded back into as filled. Returns the
290    /// zones newly filled this bar.
291    pub fn on_bar(&mut self, bar: &Bar) -> Vec<&FvgZone> {
292        let mut newly_filled_indices = Vec::new();
293        for (i, zone) in self.zones.iter_mut().enumerate() {
294            if zone.filled {
295                continue;
296            }
297            let touched = if zone.is_bullish {
298                bar.low <= zone.top
299            } else {
300                bar.high >= zone.bottom
301            };
302            if touched {
303                zone.filled = true;
304                newly_filled_indices.push(i);
305            }
306        }
307        newly_filled_indices
308            .iter()
309            .map(|&i| &self.zones[i])
310            .collect()
311    }
312}
313
314// ---------------------------------------------------------------------------------------------
315// Cross-detector event linker
316// ---------------------------------------------------------------------------------------------
317
318/// Which structural family an [`IndicatorAlert::kind`] belongs to, for correlation purposes.
319#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
320pub enum StructureCategory {
321    Break,
322    Liquidity,
323    OrderBlock,
324    Fvg,
325}
326
327fn categorize(kind: &str) -> Option<StructureCategory> {
328    match kind {
329        "structure_break" => Some(StructureCategory::Break),
330        "sweep"
331        | "liquidity_pool_stop_hunt"
332        | "liquidity_pool_breakout"
333        | "liquidity_pool_reclaim"
334        | "bullish_liquidity_sweep"
335        | "bearish_liquidity_sweep" => Some(StructureCategory::Liquidity),
336        "bullish_order_block"
337        | "bearish_order_block"
338        | "ob_retest_bullish"
339        | "ob_retest_bearish" => Some(StructureCategory::OrderBlock),
340        "bullish_fvg" | "bearish_fvg" | "fvg_filled" => Some(StructureCategory::Fvg),
341        _ => None,
342    }
343}
344
345/// A confirmed cross-detector confluence: a structure break (BOS/CHOCH) co-occurring, within the
346/// linker's window, with corroborating evidence from at least one other family.
347#[derive(Debug, Clone, PartialEq)]
348pub struct LinkedStructureEvent {
349    pub timestamp: i64,
350    pub categories: Vec<StructureCategory>,
351    /// `categories.len() / 4.0`, capped at `1.0`: how many of the four families corroborate.
352    pub confluence_score: f64,
353    pub kinds: Vec<String>,
354}
355
356/// Correlates [`IndicatorAlert`] streams from independently-running detectors within a trailing
357/// bar window, without depending on their internal types — any indicator's alerts can feed this
358/// via [`SmartMoneyStructureLinker::observe`].
359pub struct SmartMoneyStructureLinker {
360    window_bars: i64,
361    events: VecDeque<(i64, String)>,
362}
363
364impl SmartMoneyStructureLinker {
365    pub fn new(window_bars: i64) -> Self {
366        Self {
367            window_bars: window_bars.max(1),
368            events: VecDeque::new(),
369        }
370    }
371
372    pub fn reset(&mut self) {
373        self.events.clear();
374    }
375
376    /// Records `alerts` at `timestamp` and prunes anything older than the window.
377    pub fn observe(&mut self, timestamp: i64, alerts: &[IndicatorAlert]) {
378        for alert in alerts {
379            self.events.push_back((timestamp, alert.kind.clone()));
380        }
381        while self
382            .events
383            .front()
384            .map(|(ts, _)| timestamp - ts > self.window_bars)
385            .unwrap_or(false)
386        {
387            self.events.pop_front();
388        }
389    }
390
391    /// Call after [`SmartMoneyStructureLinker::observe`] for `timestamp`: if a structure-break
392    /// event occurred exactly at `timestamp`, checks whether other families co-occurred within
393    /// the trailing window and returns the linked confluence event.
394    pub fn check_confluence(&self, timestamp: i64) -> Option<LinkedStructureEvent> {
395        let break_now = self.events.iter().any(|(ts, kind)| {
396            *ts == timestamp && categorize(kind) == Some(StructureCategory::Break)
397        });
398        if !break_now {
399            return None;
400        }
401
402        let mut categories: Vec<StructureCategory> = self
403            .events
404            .iter()
405            .filter_map(|(_, kind)| categorize(kind))
406            .collect();
407        categories.sort();
408        categories.dedup();
409
410        if categories.len() < 2 {
411            return None;
412        }
413
414        let kinds: Vec<String> = self.events.iter().map(|(_, k)| k.clone()).collect();
415        let confluence_score = (categories.len() as f64 / 4.0).min(1.0);
416
417        Some(LinkedStructureEvent {
418            timestamp,
419            categories,
420            confluence_score,
421            kinds,
422        })
423    }
424}
425
426#[cfg(test)]
427mod tests {
428    use super::*;
429
430    fn trending_bars(n: usize, step: f64) -> Vec<Bar> {
431        (0..n)
432            .map(|i| {
433                let base = 100.0 + i as f64 * step;
434                Bar::new(
435                    i as i64 * 60,
436                    base,
437                    base + 3.0,
438                    base - 3.0,
439                    base + 1.0,
440                    100.0,
441                )
442            })
443            .collect()
444    }
445
446    #[test]
447    fn test_liquidity_pool_stop_hunt_vs_breakout_are_distinguished() {
448        let mut engine = LiquidityPoolEngine::new(2, 0.1);
449        // Build up bars to form a swing high pivot around 110.
450        let bars = vec![
451            Bar::new(0, 100.0, 105.0, 99.0, 100.0, 10.0),
452            Bar::new(60, 100.0, 108.0, 99.0, 100.0, 10.0),
453            Bar::new(120, 100.0, 110.0, 99.0, 100.0, 10.0), // pivot high
454            Bar::new(180, 100.0, 106.0, 99.0, 100.0, 10.0),
455            Bar::new(240, 100.0, 104.0, 99.0, 100.0, 10.0),
456        ];
457        for bar in &bars {
458            engine.on_bar(bar);
459        }
460        assert!(
461            !engine.pools().is_empty(),
462            "a BSL pool must have formed at the swing high"
463        );
464
465        // Stop hunt: pierce above 110 but close back below it.
466        let hunt = engine.on_bar(&Bar::new(300, 100.0, 111.0, 99.0, 105.0, 10.0));
467        assert!(hunt.is_some());
468        assert!(engine
469            .pools()
470            .iter()
471            .any(|p| p.state == LiquidityPoolState::StopHunted));
472        assert!(engine
473            .alerts()
474            .iter()
475            .any(|a| a.kind == "liquidity_pool_stop_hunt"));
476    }
477
478    #[test]
479    fn test_liquidity_pool_breakout_and_reclaim() {
480        let mut engine = LiquidityPoolEngine::new(2, 0.1);
481        let bars = vec![
482            Bar::new(0, 100.0, 105.0, 99.0, 100.0, 10.0),
483            Bar::new(60, 100.0, 108.0, 99.0, 100.0, 10.0),
484            Bar::new(120, 100.0, 110.0, 99.0, 100.0, 10.0),
485            Bar::new(180, 100.0, 106.0, 99.0, 100.0, 10.0),
486            Bar::new(240, 100.0, 104.0, 99.0, 100.0, 10.0),
487        ];
488        for bar in &bars {
489            engine.on_bar(bar);
490        }
491
492        // Breakout: pierce and close above 110.
493        engine.on_bar(&Bar::new(300, 100.0, 112.0, 99.0, 111.0, 10.0));
494        assert!(engine
495            .pools()
496            .iter()
497            .any(|p| p.state == LiquidityPoolState::BrokenThrough));
498
499        // Reclaim: price reverses back below the broken level.
500        let reclaim = engine.on_bar(&Bar::new(360, 111.0, 111.5, 108.0, 109.0, 10.0));
501        assert!(reclaim.is_some());
502        assert!(engine
503            .pools()
504            .iter()
505            .any(|p| p.state == LiquidityPoolState::Reclaimed));
506        assert!(engine
507            .alerts()
508            .iter()
509            .any(|a| a.kind == "liquidity_pool_reclaim"));
510    }
511
512    #[test]
513    fn test_fvg_zone_tracker_marks_fill() {
514        let mut tracker = FvgZoneTracker::new();
515        tracker.register(true, 105.0, 100.0, 0);
516        assert!(!tracker.zones()[0].filled);
517
518        // Price stays above the gap: not filled yet.
519        tracker.on_bar(&Bar::new(60, 110.0, 112.0, 108.0, 111.0, 10.0));
520        assert!(!tracker.zones()[0].filled);
521
522        // Price trades back down into the gap zone [100, 105].
523        let filled = tracker.on_bar(&Bar::new(120, 106.0, 107.0, 102.0, 103.0, 10.0));
524        assert_eq!(filled.len(), 1);
525        assert!(tracker.zones()[0].filled);
526    }
527
528    #[test]
529    fn test_linker_requires_break_plus_corroboration() {
530        let mut linker = SmartMoneyStructureLinker::new(5);
531
532        // A structure break alone, no corroboration, must not link.
533        linker.observe(10, &[IndicatorAlert::new("structure_break", "BOS", 0.9)]);
534        assert!(linker.check_confluence(10).is_none());
535
536        // A liquidity sweep shortly before a later break must corroborate it.
537        let mut linker = SmartMoneyStructureLinker::new(5);
538        linker.observe(8, &[IndicatorAlert::new("sweep", "swept", 0.85)]);
539        linker.observe(10, &[IndicatorAlert::new("structure_break", "BOS", 0.9)]);
540
541        let event = linker.check_confluence(10).unwrap();
542        assert!(event.categories.contains(&StructureCategory::Break));
543        assert!(event.categories.contains(&StructureCategory::Liquidity));
544        assert!(event.confluence_score > 0.0);
545    }
546
547    #[test]
548    fn test_linker_ignores_break_outside_current_bar() {
549        let mut linker = SmartMoneyStructureLinker::new(5);
550        linker.observe(8, &[IndicatorAlert::new("sweep", "swept", 0.85)]);
551        linker.observe(9, &[IndicatorAlert::new("structure_break", "BOS", 0.9)]);
552        // Querying a bar where no break occurred must return None even if other events exist.
553        assert!(linker.check_confluence(10).is_none());
554    }
555
556    #[test]
557    fn test_smoke_no_panic_across_trending_bars() {
558        let mut engine = LiquidityPoolEngine::with_defaults();
559        for bar in trending_bars(60, 1.5) {
560            engine.on_bar(&bar);
561        }
562    }
563}