kestrel_chartkit/indicator/
volume_profile_extended.rs1use std::collections::VecDeque;
8
9use crate::artifact::{ProfileArtifact, ProfileBin, ZoneArtifact};
10use crate::intrabar::IntrabarGroup;
11use crate::model::Bar;
12
13use super::volume_flow_hires::estimate_aggressor_from_ohlc;
14use super::{Indicator, IndicatorAlert, IndicatorOutput};
15
16#[derive(Debug, Clone, Copy, PartialEq, Eq)]
18pub enum VolumeNodeClass {
19 Hvn,
21 Avn,
23 Lvn,
25}
26
27pub struct ExtendedVolumeProfileEngine {
41 lookback: usize,
42 num_bins: usize,
43 hvn_multiplier: f64,
44 lvn_multiplier: f64,
45 records: VecDeque<Bar>,
46 alerts: Vec<IndicatorAlert>,
47}
48
49impl ExtendedVolumeProfileEngine {
50 pub fn new(lookback: usize, num_bins: usize) -> Self {
51 Self {
52 lookback: lookback.max(1),
53 num_bins: num_bins.max(1),
54 hvn_multiplier: 1.5,
55 lvn_multiplier: 0.5,
56 records: VecDeque::new(),
57 alerts: Vec::new(),
58 }
59 }
60
61 pub fn with_thresholds(mut self, hvn_multiplier: f64, lvn_multiplier: f64) -> Self {
63 self.hvn_multiplier = hvn_multiplier;
64 self.lvn_multiplier = lvn_multiplier;
65 self
66 }
67
68 fn push_record(&mut self, bar: Bar) {
69 if self.records.len() >= self.lookback {
70 self.records.pop_front();
71 }
72 self.records.push_back(bar);
73 }
74
75 pub fn on_intrabar_group(&mut self, group: &IntrabarGroup) -> Option<IndicatorOutput> {
80 for child in &group.children {
81 self.push_record(child.clone());
82 }
83 self.compute()
84 }
85
86 fn compute(&mut self) -> Option<IndicatorOutput> {
87 self.alerts.clear();
88 if self.records.len() < self.lookback {
89 return None;
90 }
91
92 let mut min_p = f64::MAX;
93 let mut max_p = f64::MIN;
94 for b in &self.records {
95 min_p = min_p.min(b.low);
96 max_p = max_p.max(b.high);
97 }
98 if (max_p - min_p).abs() < 1e-8 {
99 return None;
100 }
101
102 let step = (max_p - min_p) / self.num_bins as f64;
103 let mut volumes = vec![0.0f64; self.num_bins];
104 let mut buy_volumes = vec![0.0f64; self.num_bins];
105 let mut sell_volumes = vec![0.0f64; self.num_bins];
106
107 for b in &self.records {
108 let bar_vol = if b.volume > 0.0 {
109 b.volume
110 } else {
111 b.high - b.low
112 };
113 let aggressor = estimate_aggressor_from_ohlc(b);
114 let (buy_frac, sell_frac) = if b.volume > 0.0 {
115 (
116 aggressor.buy_volume / b.volume,
117 aggressor.sell_volume / b.volume,
118 )
119 } else {
120 (0.5, 0.5)
121 };
122
123 let raw_start = ((b.low - min_p) / step).floor();
124 let b_start = if raw_start.is_finite() && raw_start >= 0.0 {
125 (raw_start as usize).min(self.num_bins.saturating_sub(1))
126 } else {
127 0
128 };
129 let raw_end = ((b.high - min_p) / step).floor();
130 let b_end = if raw_end.is_finite() && raw_end >= 0.0 {
131 (raw_end as usize).min(self.num_bins.saturating_sub(1))
132 } else {
133 0
134 };
135 let b_end = b_end.max(b_start);
136 let bin_count = (b_end - b_start + 1) as f64;
137 let vol_per_bin = bar_vol / bin_count;
138
139 for bin_idx in b_start..=b_end {
140 if let Some(vol) = volumes.get_mut(bin_idx) {
141 *vol += vol_per_bin;
142 }
143 if let Some(buy_vol) = buy_volumes.get_mut(bin_idx) {
144 *buy_vol += vol_per_bin * buy_frac;
145 }
146 if let Some(sell_vol) = sell_volumes.get_mut(bin_idx) {
147 *sell_vol += vol_per_bin * sell_frac;
148 }
149 }
150 }
151
152 let total_vol: f64 = volumes.iter().sum();
153 let mean_bin_vol = total_vol / self.num_bins as f64;
154 let classify = |v: f64| -> VolumeNodeClass {
155 if v >= mean_bin_vol * self.hvn_multiplier {
156 VolumeNodeClass::Hvn
157 } else if v <= mean_bin_vol * self.lvn_multiplier {
158 VolumeNodeClass::Lvn
159 } else {
160 VolumeNodeClass::Avn
161 }
162 };
163 let classes: Vec<VolumeNodeClass> = volumes.iter().map(|&v| classify(v)).collect();
164
165 let (poc_idx, _) =
166 volumes
167 .iter()
168 .enumerate()
169 .fold(
170 (0usize, 0.0f64),
171 |(bi, bv), (i, &v)| {
172 if v > bv {
173 (i, v)
174 } else {
175 (bi, bv)
176 }
177 },
178 );
179 let poc_price = min_p + (poc_idx as f64 + 0.5) * step;
180
181 let target_vol = total_vol * 0.70;
182 let mut accumulated = volumes[poc_idx];
183 let mut val_idx = poc_idx;
184 let mut vah_idx = poc_idx;
185 while accumulated < target_vol && (val_idx > 0 || vah_idx < self.num_bins - 1) {
186 let next_down = if val_idx > 0 {
187 volumes[val_idx - 1]
188 } else {
189 -1.0
190 };
191 let next_up = if vah_idx < self.num_bins - 1 {
192 volumes[vah_idx + 1]
193 } else {
194 -1.0
195 };
196 if next_up >= next_down && vah_idx < self.num_bins - 1 {
197 vah_idx += 1;
198 accumulated += volumes[vah_idx];
199 } else if val_idx > 0 {
200 val_idx -= 1;
201 accumulated += volumes[val_idx];
202 } else if vah_idx < self.num_bins - 1 {
203 vah_idx += 1;
204 accumulated += volumes[vah_idx];
205 }
206 }
207 let vah_price = min_p + (vah_idx as f64 + 1.0) * step;
208 let val_price = min_p + val_idx as f64 * step;
209
210 let bin_price = |i: usize| ProfileBin {
211 price_low: min_p + i as f64 * step,
212 price_high: min_p + (i as f64 + 1.0) * step,
213 value: volumes[i],
214 };
215 let bins: Vec<ProfileBin> = (0..self.num_bins).map(bin_price).collect();
216
217 let delta_bins: Vec<ProfileBin> = (0..self.num_bins)
218 .map(|i| ProfileBin {
219 price_low: min_p + i as f64 * step,
220 price_high: min_p + (i as f64 + 1.0) * step,
221 value: buy_volumes[i] - sell_volumes[i],
222 })
223 .collect();
224
225 let window = match (self.records.front(), self.records.back()) {
227 (Some(first), Some(last)) => Some((first.timestamp, last.timestamp)),
228 _ => None,
229 };
230
231 let mut zones = Vec::new();
233 let mut i = 0;
234 while i < self.num_bins {
235 let class = classes[i];
236 if class == VolumeNodeClass::Avn {
237 i += 1;
238 continue;
239 }
240 let start = i;
241 while i < self.num_bins && classes[i] == class {
242 i += 1;
243 }
244 let end = i;
245 let zone_vol: f64 = volumes[start..end].iter().sum();
246 let mut zone = ZoneArtifact {
247 kind: match class {
248 VolumeNodeClass::Hvn => "hvn_zone".to_string(),
249 VolumeNodeClass::Lvn => "lvn_zone".to_string(),
250 VolumeNodeClass::Avn => unreachable!("filtered above"),
251 },
252 price_top: min_p + end as f64 * step,
253 price_bottom: min_p + start as f64 * step,
254 strength: if total_vol > 0.0 {
255 zone_vol / total_vol
256 } else {
257 0.0
258 },
259 touches: 0,
260 from_ts: None,
261 to_ts: None,
262 };
263 if let Some((from, to)) = window {
264 zone = zone.spanning(from, to);
265 }
266 zones.push(zone);
267 }
268
269 let mut profile_artifact = ProfileArtifact {
270 kind: "volume_profile".to_string(),
271 bins,
272 poc: poc_price,
273 value_area_high: vah_price,
274 value_area_low: val_price,
275 from_ts: None,
276 to_ts: None,
277 };
278 let mut delta_artifact = ProfileArtifact {
279 kind: "delta_profile".to_string(),
280 bins: delta_bins,
281 poc: poc_price,
282 value_area_high: vah_price,
283 value_area_low: val_price,
284 from_ts: None,
285 to_ts: None,
286 };
287 if let Some((from, to)) = window {
288 profile_artifact = profile_artifact.spanning(from, to);
289 delta_artifact = delta_artifact.spanning(from, to);
290 }
291
292 let mut output = IndicatorOutput::new(poc_price)
293 .with_artifact(profile_artifact)
294 .with_artifact(delta_artifact);
295 for zone in zones {
296 output = output.with_artifact(zone);
297 }
298 Some(output)
299 }
300}
301
302impl Indicator for ExtendedVolumeProfileEngine {
303 fn name(&self) -> &str {
304 "extended_volume_profile"
305 }
306
307 fn warmup_period(&self) -> usize {
308 self.lookback
309 }
310
311 fn reset(&mut self) {
312 self.records.clear();
313 self.alerts.clear();
314 }
315
316 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
317 self.push_record(bar.clone());
318 self.compute()
319 }
320
321 fn alerts(&self) -> Vec<IndicatorAlert> {
322 self.alerts.clone()
323 }
324}
325
326#[cfg(test)]
327mod tests {
328 use super::*;
329
330 fn bar_at(price: f64, volume: f64) -> Bar {
331 Bar::new(0, price, price + 1.0, price - 1.0, price, volume)
332 }
333
334 fn narrow_bar(price: f64, volume: f64) -> Bar {
337 Bar::new(0, price, price + 0.4, price - 0.4, price, volume)
338 }
339
340 #[test]
341 fn test_classifies_hvn_and_lvn_bins() {
342 let mut engine = ExtendedVolumeProfileEngine::new(10, 4).with_thresholds(1.5, 0.5);
343 for _ in 0..8 {
346 engine.on_bar(&narrow_bar(95.0, 500.0));
347 }
348 engine.on_bar(&bar_at(80.0, 1.0));
349 let out = engine.on_bar(&bar_at(120.0, 1.0)).unwrap();
350
351 let profile = out
352 .artifacts
353 .iter()
354 .find_map(|a| match a {
355 crate::artifact::Artifact::Profile(p) if p.kind == "volume_profile" => Some(p),
356 _ => None,
357 })
358 .expect("volume_profile artifact must be present");
359
360 assert_eq!(profile.bins.len(), 4);
361 let max_bin_volume = profile.bins.iter().map(|b| b.value).fold(0.0, f64::max);
363 let total: f64 = profile.bins.iter().map(|b| b.value).sum();
364 assert!(max_bin_volume / total > 0.5);
365 }
366
367 #[test]
368 fn test_delta_profile_reflects_buy_sell_split() {
369 let mut engine = ExtendedVolumeProfileEngine::new(3, 2);
370 engine.on_bar(&Bar::new(0, 100.0, 102.0, 98.0, 102.0, 100.0));
372 engine.on_bar(&Bar::new(60, 100.0, 102.0, 98.0, 102.0, 100.0));
373 let out = engine
374 .on_bar(&Bar::new(120, 100.0, 102.0, 98.0, 102.0, 100.0))
375 .unwrap();
376
377 let delta = out
378 .artifacts
379 .iter()
380 .find_map(|a| match a {
381 crate::artifact::Artifact::Profile(p) if p.kind == "delta_profile" => Some(p),
382 _ => None,
383 })
384 .unwrap();
385 let total_delta: f64 = delta.bins.iter().map(|b| b.value).sum();
386 assert!(
387 total_delta > 0.0,
388 "close-at-high bars must skew delta positive"
389 );
390 }
391
392 #[test]
393 fn test_zones_formed_from_contiguous_hvn_bins() {
394 let mut engine = ExtendedVolumeProfileEngine::new(8, 5).with_thresholds(1.5, 0.5);
395 engine.on_bar(&bar_at(50.0, 1.0));
399 for _ in 0..6 {
400 engine.on_bar(&narrow_bar(100.0, 1000.0));
401 }
402 let out = engine.on_bar(&bar_at(150.0, 1.0)).unwrap();
403 let zones: Vec<_> = out
404 .artifacts
405 .iter()
406 .filter_map(|a| match a {
407 crate::artifact::Artifact::Zone(z) => Some(z),
408 _ => None,
409 })
410 .collect();
411 assert!(
412 !zones.is_empty(),
413 "a tight volume cluster must form at least one zone"
414 );
415 assert!(zones.iter().any(|z| z.kind == "hvn_zone"));
416 }
417
418 #[test]
419 fn test_none_until_lookback_filled() {
420 let mut engine = ExtendedVolumeProfileEngine::new(5, 4);
421 for _ in 0..4 {
422 assert!(engine.on_bar(&bar_at(100.0, 100.0)).is_none());
423 }
424 assert!(engine.on_bar(&bar_at(100.0, 100.0)).is_some());
425 }
426}