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