kestrel_chartkit/indicator/
liquidity_sweeps.rs1use super::{Indicator, IndicatorAlert, IndicatorOutput};
2use crate::model::Bar;
3use std::collections::{HashMap, VecDeque};
4
5#[derive(Debug, Clone)]
7pub struct LiquiditySweepEngine {
8 pivot_len: usize,
9 tolerance_pct: f64,
10 bars: VecDeque<Bar>,
11 pivot_highs: Vec<f64>,
12 pivot_lows: Vec<f64>,
13 sweep_detected: i8, }
15
16impl LiquiditySweepEngine {
17 pub fn new(pivot_len: usize, tolerance_pct: f64) -> Self {
18 Self {
19 pivot_len: pivot_len.max(2),
20 tolerance_pct: tolerance_pct.max(0.01),
21 bars: VecDeque::with_capacity(pivot_len * 2 + 1),
22 pivot_highs: Vec::new(),
23 pivot_lows: Vec::new(),
24 sweep_detected: 0,
25 }
26 }
27
28 pub fn with_defaults() -> Self {
29 Self::new(5, 0.2) }
31}
32
33impl Indicator for LiquiditySweepEngine {
34 fn name(&self) -> &str {
35 "liquidity_sweeps"
36 }
37
38 fn warmup_period(&self) -> usize {
39 self.pivot_len * 2 + 1
40 }
41
42 fn reset(&mut self) {
43 self.bars.clear();
44 self.pivot_highs.clear();
45 self.pivot_lows.clear();
46 self.sweep_detected = 0;
47 }
48
49 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
50 self.bars.push_back(bar.clone());
51 if self.bars.len() > self.pivot_len * 2 + 1 {
52 self.bars.pop_front();
53 }
54
55 if self.bars.len() < self.pivot_len * 2 + 1 {
56 return None;
57 }
58
59 let mid_idx = self.pivot_len;
60 let mid_bar = &self.bars[mid_idx];
61
62 let is_pivot_high = self
63 .bars
64 .iter()
65 .enumerate()
66 .all(|(i, b)| i == mid_idx || b.high <= mid_bar.high);
67 let is_pivot_low = self
68 .bars
69 .iter()
70 .enumerate()
71 .all(|(i, b)| i == mid_idx || b.low >= mid_bar.low);
72
73 if is_pivot_high {
74 self.pivot_highs.push(mid_bar.high);
75 if self.pivot_highs.len() > 20 {
76 self.pivot_highs.remove(0);
77 }
78 }
79 if is_pivot_low {
80 self.pivot_lows.push(mid_bar.low);
81 if self.pivot_lows.len() > 20 {
82 self.pivot_lows.remove(0);
83 }
84 }
85
86 self.sweep_detected = 0;
87
88 for &ph in &self.pivot_highs {
90 if bar.high > ph && bar.close < ph {
91 self.sweep_detected = -1;
92 break;
93 }
94 }
95
96 if self.sweep_detected == 0 {
98 for &pl in &self.pivot_lows {
99 if bar.low < pl && bar.close > pl {
100 self.sweep_detected = 1;
101 break;
102 }
103 }
104 }
105
106 let eqh_count = self
108 .pivot_highs
109 .windows(2)
110 .filter(|w| (w[0] - w[1]).abs() / w[0] * 100.0 <= self.tolerance_pct)
111 .count();
112 let eql_count = self
113 .pivot_lows
114 .windows(2)
115 .filter(|w| (w[0] - w[1]).abs() / w[0] * 100.0 <= self.tolerance_pct)
116 .count();
117
118 let mut extra = HashMap::new();
119 extra.insert("sweep".to_string(), self.sweep_detected as f64);
120 extra.insert("eqh_count".to_string(), eqh_count as f64);
121 extra.insert("eql_count".to_string(), eql_count as f64);
122
123 Some(IndicatorOutput::with_extra(
124 self.sweep_detected as f64,
125 extra,
126 ))
127 }
128
129 fn alerts(&self) -> Vec<IndicatorAlert> {
130 let mut alerts = Vec::new();
131 if self.sweep_detected == 1 {
132 alerts.push(IndicatorAlert::new(
133 "sweep",
134 "Bullish Liquidity Sweep & Reclaim",
135 0.85,
136 ));
137 } else if self.sweep_detected == -1 {
138 alerts.push(IndicatorAlert::new(
139 "sweep",
140 "Bearish Liquidity Sweep & Reclaim",
141 0.85,
142 ));
143 }
144 alerts
145 }
146}
147
148#[cfg(test)]
149mod tests {
150 use super::*;
151
152 #[test]
153 fn test_liquidity_sweep_detection() {
154 let mut sweep = LiquiditySweepEngine::with_defaults();
155 for i in 0..20 {
156 let b = Bar::new(i, 100.0, 105.0, 95.0, 100.0, 1000.0);
157 sweep.on_bar(&b);
158 }
159 let sweep_bar = Bar::new(20, 99.0, 101.0, 90.0, 98.0, 1000.0);
161 let out = sweep.on_bar(&sweep_bar);
162 assert!(out.is_some());
163 }
164}