kestrel_chartkit/indicator/
volume_flow_hires.rs1use 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#[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
27pub(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 pub absorption: bool,
49}
50
51pub struct HiResVolumeFlowEngine {
66 cumulative_delta: f64,
67 volume_threshold: RollingRobustThreshold,
68 range_threshold: RollingRobustThreshold,
69 alerts: Vec<IndicatorAlert>,
70}
71
72impl HiResVolumeFlowEngine {
73 pub fn new(window_len: usize) -> Self {
74 Self {
75 cumulative_delta: 0.0,
76 volume_threshold: RollingRobustThreshold::new(window_len, 2.5),
77 range_threshold: RollingRobustThreshold::new(window_len, 2.5),
78 alerts: Vec::new(),
79 }
80 }
81
82 pub fn reset(&mut self) {
83 self.cumulative_delta = 0.0;
84 self.volume_threshold.reset();
85 self.range_threshold.reset();
86 self.alerts.clear();
87 }
88
89 fn absorb_step(&mut self, bar: &Bar, delta: f64, quality: BarQuality) -> HiResVolumeFlowOutput {
90 self.alerts.clear();
91 self.cumulative_delta += delta;
92
93 let volume_band = self.volume_threshold.update(bar.volume);
94 let range_band = self.range_threshold.update(bar.high - bar.low);
95
96 let absorption = match (volume_band, range_band) {
97 (Some(vb), Some(rb)) => bar.volume > vb.upper && (bar.high - bar.low) <= rb.median,
98 _ => false,
99 };
100
101 if absorption {
102 self.alerts.push(IndicatorAlert::new(
103 "volume_absorption",
104 "Outlier volume with unremarkable price range: possible absorption",
105 0.7,
106 ));
107 }
108
109 HiResVolumeFlowOutput {
110 delta,
111 cumulative_delta: self.cumulative_delta,
112 quality,
113 absorption,
114 }
115 }
116
117 pub fn on_bar_with_aggressor(
120 &mut self,
121 bar: &Bar,
122 aggressor: AggressorVolume,
123 ) -> HiResVolumeFlowOutput {
124 self.absorb_step(bar, aggressor.delta(), BarQuality::observed())
125 }
126
127 pub fn on_intrabar_group(&mut self, group: &IntrabarGroup) -> HiResVolumeFlowOutput {
132 let mut delta = 0.0;
133 let mut open = f64::NAN;
134 let mut high = f64::NEG_INFINITY;
135 let mut low = f64::INFINITY;
136 let mut close = f64::NAN;
137 let mut volume = 0.0;
138
139 for (i, child) in group.children.iter().enumerate() {
140 let aggressor = estimate_aggressor_from_ohlc(child);
141 delta += aggressor.delta();
142 if i == 0 {
143 open = child.open;
144 }
145 high = high.max(child.high);
146 low = low.min(child.low);
147 close = child.close;
148 volume += child.volume;
149 }
150
151 let parent_bar = Bar::new(group.parent_timestamp, open, high, low, close, volume);
152 let mut quality = BarQuality::observed();
153 quality.is_forward_filled = false;
154 self.absorb_step(&parent_bar, delta, quality)
155 }
156
157 pub fn on_bar_estimated(&mut self, bar: &Bar) -> HiResVolumeFlowOutput {
161 let aggressor = estimate_aggressor_from_ohlc(bar);
162 let mut quality = BarQuality::observed();
163 quality.is_synthetic = true; self.absorb_step(bar, aggressor.delta(), quality)
165 }
166}
167
168impl Indicator for HiResVolumeFlowEngine {
169 fn name(&self) -> &str {
170 "hires_volume_flow"
171 }
172
173 fn reset(&mut self) {
174 HiResVolumeFlowEngine::reset(self)
175 }
176
177 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
182 let out = self.on_bar_estimated(bar);
183 let mut extra = HashMap::new();
184 extra.insert("delta".to_string(), out.delta);
185 extra.insert(
186 "is_estimated".to_string(),
187 if out.quality.is_synthetic { 1.0 } else { 0.0 },
188 );
189 Some(
190 IndicatorOutput::with_extra(out.cumulative_delta, extra).with_state(
191 if out.absorption {
192 "absorption"
193 } else {
194 "normal"
195 },
196 ),
197 )
198 }
199
200 fn alerts(&self) -> Vec<IndicatorAlert> {
201 self.alerts.clone()
202 }
203}
204
205#[cfg(test)]
206mod tests {
207 use super::*;
208 use crate::intrabar::IntrabarGrouper;
209 use crate::timeframe::Timeframe;
210
211 #[test]
212 fn test_direct_aggressor_delta_matches_input_exactly() {
213 let mut engine = HiResVolumeFlowEngine::new(5);
214 let bar = Bar::new(0, 100.0, 101.0, 99.0, 100.5, 1000.0);
215 let out = engine.on_bar_with_aggressor(
216 &bar,
217 AggressorVolume {
218 buy_volume: 700.0,
219 sell_volume: 300.0,
220 },
221 );
222 assert_eq!(out.delta, 400.0);
223 assert_eq!(out.quality, BarQuality::observed());
224 }
225
226 #[test]
227 fn test_estimated_fallback_is_tagged_synthetic() {
228 let mut engine = HiResVolumeFlowEngine::new(5);
229 let bar = Bar::new(0, 100.0, 101.0, 99.0, 100.5, 1000.0);
230 let out = engine.on_bar_estimated(&bar);
231 assert!(
232 out.quality.is_synthetic,
233 "estimated delta must be tagged as such"
234 );
235 }
236
237 #[test]
238 fn test_intrabar_group_sums_child_deltas() {
239 let mut grouper = IntrabarGrouper::new(Timeframe::Minute(5)).unwrap();
240 let children = [
241 Bar::new(0, 100.0, 101.0, 100.0, 101.0, 100.0),
242 Bar::new(60, 101.0, 102.0, 100.5, 101.5, 100.0),
243 ];
244 for child in &children {
245 grouper.on_child_bar(child);
246 }
247 let completed = grouper.on_child_bar(&Bar::new(300, 101.5, 102.0, 101.0, 101.8, 50.0));
249
250 let mut engine = HiResVolumeFlowEngine::new(5);
251 let out = engine.on_intrabar_group(&completed.unwrap());
252
253 let expected: f64 = children
254 .iter()
255 .map(|c| estimate_aggressor_from_ohlc(c).delta())
256 .sum();
257 assert!((out.delta - expected).abs() < 1e-9);
258 }
259
260 #[test]
261 fn test_absorption_flags_outlier_volume_with_tight_range() {
262 let mut engine = HiResVolumeFlowEngine::new(6);
263 for _ in 0..5 {
265 engine.on_bar_estimated(&Bar::new(0, 100.0, 101.0, 99.0, 100.5, 100.0));
266 }
267 let out = engine.on_bar_estimated(&Bar::new(60, 100.0, 100.3, 99.8, 100.1, 5000.0));
269 assert!(
270 out.absorption,
271 "large volume with tight range must flag absorption"
272 );
273 assert!(!engine.alerts().is_empty());
274 }
275
276 #[test]
277 fn test_normal_bar_does_not_flag_absorption() {
278 let mut engine = HiResVolumeFlowEngine::new(6);
279 let mut last = HiResVolumeFlowOutput {
280 delta: 0.0,
281 cumulative_delta: 0.0,
282 quality: BarQuality::observed(),
283 absorption: false,
284 };
285 for _ in 0..6 {
286 last = engine.on_bar_estimated(&Bar::new(0, 100.0, 101.0, 99.0, 100.5, 100.0));
287 }
288 assert!(!last.absorption);
289 }
290}