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 {
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#[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#[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 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 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#[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#[derive(Debug, Clone, PartialEq)]
361pub struct LinkedStructureEvent {
362 pub timestamp: i64,
363 pub categories: Vec<StructureCategory>,
364 pub confluence_score: f64,
366 pub kinds: Vec<String>,
367}
368
369pub 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 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 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 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), 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 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 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 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 tracker.on_bar(&Bar::new(60, 110.0, 112.0, 108.0, 111.0, 10.0));
533 assert!(!tracker.zones()[0].filled);
534
535 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 linker.observe(10, &[IndicatorAlert::new("structure_break", "BOS", 0.9)]);
547 assert!(linker.check_confluence(10).is_none());
548
549 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 assert!(linker.check_confluence(10).is_none());
567 }
568
569 #[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}