1use std::collections::VecDeque;
12
13use crate::model::Bar;
14
15use super::{Indicator, IndicatorAlert, IndicatorOutput};
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub enum LiquidityPoolKind {
25 Bsl,
26 Ssl,
27}
28
29#[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 pub touches: u32,
47 pub formed_at: i64,
48 pub state: LiquidityPoolState,
49}
50
51pub 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#[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#[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 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 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#[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#[derive(Debug, Clone, PartialEq)]
348pub struct LinkedStructureEvent {
349 pub timestamp: i64,
350 pub categories: Vec<StructureCategory>,
351 pub confluence_score: f64,
353 pub kinds: Vec<String>,
354}
355
356pub 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 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 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 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), 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 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 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 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 tracker.on_bar(&Bar::new(60, 110.0, 112.0, 108.0, 111.0, 10.0));
520 assert!(!tracker.zones()[0].filled);
521
522 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 linker.observe(10, &[IndicatorAlert::new("structure_break", "BOS", 0.9)]);
534 assert!(linker.check_confluence(10).is_none());
535
536 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 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}