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 {
29 lookback: usize,
30 num_bins: usize,
31 hvn_multiplier: f64,
32 lvn_multiplier: f64,
33 records: VecDeque<Bar>,
34 alerts: Vec<IndicatorAlert>,
35}
36
37impl ExtendedVolumeProfileEngine {
38 pub fn new(lookback: usize, num_bins: usize) -> Self {
39 Self {
40 lookback: lookback.max(1),
41 num_bins: num_bins.max(1),
42 hvn_multiplier: 1.5,
43 lvn_multiplier: 0.5,
44 records: VecDeque::new(),
45 alerts: Vec::new(),
46 }
47 }
48
49 pub fn with_thresholds(mut self, hvn_multiplier: f64, lvn_multiplier: f64) -> Self {
51 self.hvn_multiplier = hvn_multiplier;
52 self.lvn_multiplier = lvn_multiplier;
53 self
54 }
55
56 fn push_record(&mut self, bar: Bar) {
57 if self.records.len() >= self.lookback {
58 self.records.pop_front();
59 }
60 self.records.push_back(bar);
61 }
62
63 pub fn on_intrabar_group(&mut self, group: &IntrabarGroup) -> Option<IndicatorOutput> {
68 for child in &group.children {
69 self.push_record(child.clone());
70 }
71 self.compute()
72 }
73
74 fn compute(&mut self) -> Option<IndicatorOutput> {
75 self.alerts.clear();
76 if self.records.len() < self.lookback {
77 return None;
78 }
79
80 let mut min_p = f64::MAX;
81 let mut max_p = f64::MIN;
82 for b in &self.records {
83 min_p = min_p.min(b.low);
84 max_p = max_p.max(b.high);
85 }
86 if (max_p - min_p).abs() < 1e-8 {
87 return None;
88 }
89
90 let step = (max_p - min_p) / self.num_bins as f64;
91 let mut volumes = vec![0.0f64; self.num_bins];
92 let mut buy_volumes = vec![0.0f64; self.num_bins];
93 let mut sell_volumes = vec![0.0f64; self.num_bins];
94
95 for b in &self.records {
96 let bar_vol = if b.volume > 0.0 {
97 b.volume
98 } else {
99 b.high - b.low
100 };
101 let aggressor = estimate_aggressor_from_ohlc(b);
102 let (buy_frac, sell_frac) = if b.volume > 0.0 {
103 (
104 aggressor.buy_volume / b.volume,
105 aggressor.sell_volume / b.volume,
106 )
107 } else {
108 (0.5, 0.5)
109 };
110
111 let b_start = (((b.low - min_p) / step).floor() as usize).min(self.num_bins - 1);
112 let b_end = (((b.high - min_p) / step).floor() as usize).min(self.num_bins - 1);
113 let bin_count = (b_end - b_start + 1) as f64;
114 let vol_per_bin = bar_vol / bin_count;
115
116 for bin_idx in b_start..=b_end {
117 volumes[bin_idx] += vol_per_bin;
118 buy_volumes[bin_idx] += vol_per_bin * buy_frac;
119 sell_volumes[bin_idx] += vol_per_bin * sell_frac;
120 }
121 }
122
123 let total_vol: f64 = volumes.iter().sum();
124 let mean_bin_vol = total_vol / self.num_bins as f64;
125 let classify = |v: f64| -> VolumeNodeClass {
126 if v >= mean_bin_vol * self.hvn_multiplier {
127 VolumeNodeClass::Hvn
128 } else if v <= mean_bin_vol * self.lvn_multiplier {
129 VolumeNodeClass::Lvn
130 } else {
131 VolumeNodeClass::Avn
132 }
133 };
134 let classes: Vec<VolumeNodeClass> = volumes.iter().map(|&v| classify(v)).collect();
135
136 let (poc_idx, _) =
137 volumes
138 .iter()
139 .enumerate()
140 .fold(
141 (0usize, 0.0f64),
142 |(bi, bv), (i, &v)| {
143 if v > bv {
144 (i, v)
145 } else {
146 (bi, bv)
147 }
148 },
149 );
150 let poc_price = min_p + (poc_idx as f64 + 0.5) * step;
151
152 let target_vol = total_vol * 0.70;
153 let mut accumulated = volumes[poc_idx];
154 let mut val_idx = poc_idx;
155 let mut vah_idx = poc_idx;
156 while accumulated < target_vol && (val_idx > 0 || vah_idx < self.num_bins - 1) {
157 let next_down = if val_idx > 0 {
158 volumes[val_idx - 1]
159 } else {
160 -1.0
161 };
162 let next_up = if vah_idx < self.num_bins - 1 {
163 volumes[vah_idx + 1]
164 } else {
165 -1.0
166 };
167 if next_up >= next_down && vah_idx < self.num_bins - 1 {
168 vah_idx += 1;
169 accumulated += volumes[vah_idx];
170 } else if val_idx > 0 {
171 val_idx -= 1;
172 accumulated += volumes[val_idx];
173 } else if vah_idx < self.num_bins - 1 {
174 vah_idx += 1;
175 accumulated += volumes[vah_idx];
176 }
177 }
178 let vah_price = min_p + (vah_idx as f64 + 1.0) * step;
179 let val_price = min_p + val_idx as f64 * step;
180
181 let bin_price = |i: usize| ProfileBin {
182 price_low: min_p + i as f64 * step,
183 price_high: min_p + (i as f64 + 1.0) * step,
184 value: volumes[i],
185 };
186 let bins: Vec<ProfileBin> = (0..self.num_bins).map(bin_price).collect();
187
188 let delta_bins: Vec<ProfileBin> = (0..self.num_bins)
189 .map(|i| ProfileBin {
190 price_low: min_p + i as f64 * step,
191 price_high: min_p + (i as f64 + 1.0) * step,
192 value: buy_volumes[i] - sell_volumes[i],
193 })
194 .collect();
195
196 let mut zones = Vec::new();
198 let mut i = 0;
199 while i < self.num_bins {
200 let class = classes[i];
201 if class == VolumeNodeClass::Avn {
202 i += 1;
203 continue;
204 }
205 let start = i;
206 while i < self.num_bins && classes[i] == class {
207 i += 1;
208 }
209 let end = i;
210 let zone_vol: f64 = volumes[start..end].iter().sum();
211 zones.push(ZoneArtifact {
212 kind: match class {
213 VolumeNodeClass::Hvn => "hvn_zone".to_string(),
214 VolumeNodeClass::Lvn => "lvn_zone".to_string(),
215 VolumeNodeClass::Avn => unreachable!("filtered above"),
216 },
217 price_top: min_p + end as f64 * step,
218 price_bottom: min_p + start as f64 * step,
219 strength: if total_vol > 0.0 {
220 zone_vol / total_vol
221 } else {
222 0.0
223 },
224 touches: 0,
225 });
226 }
227
228 let profile_artifact = ProfileArtifact {
229 kind: "volume_profile".to_string(),
230 bins,
231 poc: poc_price,
232 value_area_high: vah_price,
233 value_area_low: val_price,
234 };
235 let delta_artifact = ProfileArtifact {
236 kind: "delta_profile".to_string(),
237 bins: delta_bins,
238 poc: poc_price,
239 value_area_high: vah_price,
240 value_area_low: val_price,
241 };
242
243 let mut output = IndicatorOutput::new(poc_price)
244 .with_artifact(profile_artifact)
245 .with_artifact(delta_artifact);
246 for zone in zones {
247 output = output.with_artifact(zone);
248 }
249 Some(output)
250 }
251}
252
253impl Indicator for ExtendedVolumeProfileEngine {
254 fn name(&self) -> &str {
255 "extended_volume_profile"
256 }
257
258 fn warmup_period(&self) -> usize {
259 self.lookback
260 }
261
262 fn reset(&mut self) {
263 self.records.clear();
264 self.alerts.clear();
265 }
266
267 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
268 self.push_record(bar.clone());
269 self.compute()
270 }
271
272 fn alerts(&self) -> Vec<IndicatorAlert> {
273 self.alerts.clone()
274 }
275}
276
277#[cfg(test)]
278mod tests {
279 use super::*;
280
281 fn bar_at(price: f64, volume: f64) -> Bar {
282 Bar::new(0, price, price + 1.0, price - 1.0, price, volume)
283 }
284
285 fn narrow_bar(price: f64, volume: f64) -> Bar {
288 Bar::new(0, price, price + 0.4, price - 0.4, price, volume)
289 }
290
291 #[test]
292 fn test_classifies_hvn_and_lvn_bins() {
293 let mut engine = ExtendedVolumeProfileEngine::new(10, 4).with_thresholds(1.5, 0.5);
294 for _ in 0..8 {
297 engine.on_bar(&narrow_bar(95.0, 500.0));
298 }
299 engine.on_bar(&bar_at(80.0, 1.0));
300 let out = engine.on_bar(&bar_at(120.0, 1.0)).unwrap();
301
302 let profile = out
303 .artifacts
304 .iter()
305 .find_map(|a| match a {
306 crate::artifact::Artifact::Profile(p) if p.kind == "volume_profile" => Some(p),
307 _ => None,
308 })
309 .expect("volume_profile artifact must be present");
310
311 assert_eq!(profile.bins.len(), 4);
312 let max_bin_volume = profile.bins.iter().map(|b| b.value).fold(0.0, f64::max);
314 let total: f64 = profile.bins.iter().map(|b| b.value).sum();
315 assert!(max_bin_volume / total > 0.5);
316 }
317
318 #[test]
319 fn test_delta_profile_reflects_buy_sell_split() {
320 let mut engine = ExtendedVolumeProfileEngine::new(3, 2);
321 engine.on_bar(&Bar::new(0, 100.0, 102.0, 98.0, 102.0, 100.0));
323 engine.on_bar(&Bar::new(60, 100.0, 102.0, 98.0, 102.0, 100.0));
324 let out = engine
325 .on_bar(&Bar::new(120, 100.0, 102.0, 98.0, 102.0, 100.0))
326 .unwrap();
327
328 let delta = out
329 .artifacts
330 .iter()
331 .find_map(|a| match a {
332 crate::artifact::Artifact::Profile(p) if p.kind == "delta_profile" => Some(p),
333 _ => None,
334 })
335 .unwrap();
336 let total_delta: f64 = delta.bins.iter().map(|b| b.value).sum();
337 assert!(
338 total_delta > 0.0,
339 "close-at-high bars must skew delta positive"
340 );
341 }
342
343 #[test]
344 fn test_zones_formed_from_contiguous_hvn_bins() {
345 let mut engine = ExtendedVolumeProfileEngine::new(8, 5).with_thresholds(1.5, 0.5);
346 engine.on_bar(&bar_at(50.0, 1.0));
350 for _ in 0..6 {
351 engine.on_bar(&narrow_bar(100.0, 1000.0));
352 }
353 let out = engine.on_bar(&bar_at(150.0, 1.0)).unwrap();
354 let zones: Vec<_> = out
355 .artifacts
356 .iter()
357 .filter_map(|a| match a {
358 crate::artifact::Artifact::Zone(z) => Some(z),
359 _ => None,
360 })
361 .collect();
362 assert!(
363 !zones.is_empty(),
364 "a tight volume cluster must form at least one zone"
365 );
366 assert!(zones.iter().any(|z| z.kind == "hvn_zone"));
367 }
368
369 #[test]
370 fn test_none_until_lookback_filled() {
371 let mut engine = ExtendedVolumeProfileEngine::new(5, 4);
372 for _ in 0..4 {
373 assert!(engine.on_bar(&bar_at(100.0, 100.0)).is_none());
374 }
375 assert!(engine.on_bar(&bar_at(100.0, 100.0)).is_some());
376 }
377}