Skip to main content

kestrel_chartkit/indicator/
volume_flow_hires.rs

1//! High-resolution volume flow: direct aggressor/delta inputs, intrabar delta from grouped child
2//! bars, absorption detection, and explicitly quality-tagged fallbacks — the capabilities
3//! [`super::volume_flow::CvdEngine`]'s single OHLC-close-location heuristic per bar cannot offer.
4
5use std::collections::HashMap;
6
7use crate::clustering::RollingRobustThreshold;
8use crate::intrabar::IntrabarGroup;
9use crate::model::{Bar, BarQuality};
10
11use super::{Indicator, IndicatorAlert, IndicatorOutput};
12
13/// Directly known aggressor-side (taker buy vs. taker sell) volume for a bar, from real trade/
14/// tick data rather than an OHLC-inferred estimate.
15#[derive(Debug, Clone, Copy, PartialEq)]
16pub struct AggressorVolume {
17    pub buy_volume: f64,
18    pub sell_volume: f64,
19}
20
21impl AggressorVolume {
22    pub fn delta(&self) -> f64 {
23        self.buy_volume - self.sell_volume
24    }
25}
26
27/// Estimates a bar's aggressor split from its OHLC close-location within its range — the same
28/// heuristic [`super::volume_flow::CvdEngine`] uses, kept here as the explicit, tagged fallback
29/// path when no direct aggressor/tick data is available. Shared with
30/// [`super::volume_profile_extended`] for its delta-profile bins.
31pub(crate) fn estimate_aggressor_from_ohlc(bar: &Bar) -> AggressorVolume {
32    let range = (bar.high - bar.low).max(1e-8);
33    let buy_pct = ((bar.close - bar.low) / range).clamp(0.0, 1.0);
34    AggressorVolume {
35        buy_volume: bar.volume * buy_pct,
36        sell_volume: bar.volume * (1.0 - buy_pct),
37    }
38}
39
40#[derive(Debug, Clone, Copy, PartialEq)]
41pub struct HiResVolumeFlowOutput {
42    pub delta: f64,
43    pub cumulative_delta: f64,
44    pub quality: BarQuality,
45    /// `true` when this bar's volume is an outlier-high spike (per the rolling robust volume
46    /// threshold) while its price range stayed unremarkable — high participation without
47    /// proportional price displacement, the classic order-flow absorption signature.
48    pub absorption: bool,
49}
50
51/// Streaming high-resolution volume flow engine.
52pub struct HiResVolumeFlowEngine {
53    cumulative_delta: f64,
54    volume_threshold: RollingRobustThreshold,
55    range_threshold: RollingRobustThreshold,
56    alerts: Vec<IndicatorAlert>,
57}
58
59impl HiResVolumeFlowEngine {
60    pub fn new(window_len: usize) -> Self {
61        Self {
62            cumulative_delta: 0.0,
63            volume_threshold: RollingRobustThreshold::new(window_len, 2.5),
64            range_threshold: RollingRobustThreshold::new(window_len, 2.5),
65            alerts: Vec::new(),
66        }
67    }
68
69    pub fn reset(&mut self) {
70        self.cumulative_delta = 0.0;
71        self.volume_threshold.reset();
72        self.range_threshold.reset();
73        self.alerts.clear();
74    }
75
76    fn absorb_step(&mut self, bar: &Bar, delta: f64, quality: BarQuality) -> HiResVolumeFlowOutput {
77        self.alerts.clear();
78        self.cumulative_delta += delta;
79
80        let volume_band = self.volume_threshold.update(bar.volume);
81        let range_band = self.range_threshold.update(bar.high - bar.low);
82
83        let absorption = match (volume_band, range_band) {
84            (Some(vb), Some(rb)) => bar.volume > vb.upper && (bar.high - bar.low) <= rb.median,
85            _ => false,
86        };
87
88        if absorption {
89            self.alerts.push(IndicatorAlert::new(
90                "volume_absorption",
91                "Outlier volume with unremarkable price range: possible absorption",
92                0.7,
93            ));
94        }
95
96        HiResVolumeFlowOutput {
97            delta,
98            cumulative_delta: self.cumulative_delta,
99            quality,
100            absorption,
101        }
102    }
103
104    /// High-resolution path: feeds a bar with directly known aggressor volume (e.g. aggregated
105    /// from trade-tape data).
106    pub fn on_bar_with_aggressor(
107        &mut self,
108        bar: &Bar,
109        aggressor: AggressorVolume,
110    ) -> HiResVolumeFlowOutput {
111        self.absorb_step(bar, aggressor.delta(), BarQuality::observed())
112    }
113
114    /// High-resolution path: feeds a confirmed [`IntrabarGroup`] (a parent bar's full ordered
115    /// child-bar sequence), summing each child's own OHLC-inferred delta instead of applying the
116    /// heuristic once to the aggregate parent bar — preserves the intrabar price path the
117    /// aggregate alone loses.
118    pub fn on_intrabar_group(&mut self, group: &IntrabarGroup) -> HiResVolumeFlowOutput {
119        let mut delta = 0.0;
120        let mut open = f64::NAN;
121        let mut high = f64::NEG_INFINITY;
122        let mut low = f64::INFINITY;
123        let mut close = f64::NAN;
124        let mut volume = 0.0;
125
126        for (i, child) in group.children.iter().enumerate() {
127            let aggressor = estimate_aggressor_from_ohlc(child);
128            delta += aggressor.delta();
129            if i == 0 {
130                open = child.open;
131            }
132            high = high.max(child.high);
133            low = low.min(child.low);
134            close = child.close;
135            volume += child.volume;
136        }
137
138        let parent_bar = Bar::new(group.parent_timestamp, open, high, low, close, volume);
139        let mut quality = BarQuality::observed();
140        quality.is_forward_filled = false;
141        self.absorb_step(&parent_bar, delta, quality)
142    }
143
144    /// Fallback path: OHLC-inferred estimate when no direct aggressor or intrabar data is
145    /// available, explicitly tagged as such via [`BarQuality::volume_available`] staying accurate
146    /// but the estimate not being a genuine tick-level measurement.
147    pub fn on_bar_estimated(&mut self, bar: &Bar) -> HiResVolumeFlowOutput {
148        let aggressor = estimate_aggressor_from_ohlc(bar);
149        let mut quality = BarQuality::observed();
150        quality.is_synthetic = true; // the buy/sell split itself is inferred, not observed
151        self.absorb_step(bar, aggressor.delta(), quality)
152    }
153}
154
155impl Indicator for HiResVolumeFlowEngine {
156    fn name(&self) -> &str {
157        "hires_volume_flow"
158    }
159
160    fn reset(&mut self) {
161        HiResVolumeFlowEngine::reset(self)
162    }
163
164    /// Delegates to the OHLC-estimated fallback path; use
165    /// [`HiResVolumeFlowEngine::on_bar_with_aggressor`] or
166    /// [`HiResVolumeFlowEngine::on_intrabar_group`] directly for the high-resolution paths, which
167    /// this trait's single-`Bar` signature cannot express.
168    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
169        let out = self.on_bar_estimated(bar);
170        let mut extra = HashMap::new();
171        extra.insert("delta".to_string(), out.delta);
172        extra.insert(
173            "is_estimated".to_string(),
174            if out.quality.is_synthetic { 1.0 } else { 0.0 },
175        );
176        Some(
177            IndicatorOutput::with_extra(out.cumulative_delta, extra).with_state(
178                if out.absorption {
179                    "absorption"
180                } else {
181                    "normal"
182                },
183            ),
184        )
185    }
186
187    fn alerts(&self) -> Vec<IndicatorAlert> {
188        self.alerts.clone()
189    }
190}
191
192#[cfg(test)]
193mod tests {
194    use super::*;
195    use crate::intrabar::IntrabarGrouper;
196    use crate::timeframe::Timeframe;
197
198    #[test]
199    fn test_direct_aggressor_delta_matches_input_exactly() {
200        let mut engine = HiResVolumeFlowEngine::new(5);
201        let bar = Bar::new(0, 100.0, 101.0, 99.0, 100.5, 1000.0);
202        let out = engine.on_bar_with_aggressor(
203            &bar,
204            AggressorVolume {
205                buy_volume: 700.0,
206                sell_volume: 300.0,
207            },
208        );
209        assert_eq!(out.delta, 400.0);
210        assert_eq!(out.quality, BarQuality::observed());
211    }
212
213    #[test]
214    fn test_estimated_fallback_is_tagged_synthetic() {
215        let mut engine = HiResVolumeFlowEngine::new(5);
216        let bar = Bar::new(0, 100.0, 101.0, 99.0, 100.5, 1000.0);
217        let out = engine.on_bar_estimated(&bar);
218        assert!(
219            out.quality.is_synthetic,
220            "estimated delta must be tagged as such"
221        );
222    }
223
224    #[test]
225    fn test_intrabar_group_sums_child_deltas() {
226        let mut grouper = IntrabarGrouper::new(Timeframe::Minute(5)).unwrap();
227        let children = [
228            Bar::new(0, 100.0, 101.0, 100.0, 101.0, 100.0),
229            Bar::new(60, 101.0, 102.0, 100.5, 101.5, 100.0),
230        ];
231        for child in &children {
232            grouper.on_child_bar(child);
233        }
234        // Starts a new parent bucket, completing the first with the two children above.
235        let completed = grouper.on_child_bar(&Bar::new(300, 101.5, 102.0, 101.0, 101.8, 50.0));
236
237        let mut engine = HiResVolumeFlowEngine::new(5);
238        let out = engine.on_intrabar_group(&completed.unwrap());
239
240        let expected: f64 = children
241            .iter()
242            .map(|c| estimate_aggressor_from_ohlc(c).delta())
243            .sum();
244        assert!((out.delta - expected).abs() < 1e-9);
245    }
246
247    #[test]
248    fn test_absorption_flags_outlier_volume_with_tight_range() {
249        let mut engine = HiResVolumeFlowEngine::new(6);
250        // Establish a normal volume/range baseline.
251        for _ in 0..5 {
252            engine.on_bar_estimated(&Bar::new(0, 100.0, 101.0, 99.0, 100.5, 100.0));
253        }
254        // One bar with a volume spike but a tight (unremarkable) range.
255        let out = engine.on_bar_estimated(&Bar::new(60, 100.0, 100.3, 99.8, 100.1, 5000.0));
256        assert!(
257            out.absorption,
258            "large volume with tight range must flag absorption"
259        );
260        assert!(!engine.alerts().is_empty());
261    }
262
263    #[test]
264    fn test_normal_bar_does_not_flag_absorption() {
265        let mut engine = HiResVolumeFlowEngine::new(6);
266        let mut last = HiResVolumeFlowOutput {
267            delta: 0.0,
268            cumulative_delta: 0.0,
269            quality: BarQuality::observed(),
270            absorption: false,
271        };
272        for _ in 0..6 {
273            last = engine.on_bar_estimated(&Bar::new(0, 100.0, 101.0, 99.0, 100.5, 100.0));
274        }
275        assert!(!last.absorption);
276    }
277}