kestrel_chartkit/indicator/
volume_profile_persistent.rs1use std::collections::{HashMap, VecDeque};
11
12use crate::artifact::{ProfileArtifact, ProfileBin, ZoneArtifact};
13use crate::model::Bar;
14use crate::stats::rolling_median;
15
16use super::volume_flow_hires::estimate_aggressor_from_ohlc;
17use super::{Indicator, IndicatorAlert, IndicatorOutput};
18
19const MAD_CONSISTENCY_CONSTANT: f64 = 1.482_602_218_505_602;
22
23#[derive(Debug, Clone, Copy, PartialEq)]
24struct BinLifecycle {
25 volume: f64,
26 buy_volume: f64,
27 sell_volume: f64,
28 touches: u32,
29 first_touched_ts: i64,
30 last_touched_ts: i64,
31}
32
33struct RecordContribution {
36 per_bin: Vec<(i64, f64, f64, f64)>,
37}
38
39#[derive(Debug, Clone, Copy, PartialEq)]
41pub struct AbsorptionBin {
42 pub price_low: f64,
43 pub price_high: f64,
44 pub volume: f64,
45 pub touches: u32,
46 pub volume_per_touch: f64,
47 pub is_absorption: bool,
48}
49
50pub struct PersistentVolumeProfileEngine {
51 lookback: usize,
52 bin_width: f64,
53 absorption_k: f64,
54 bins: HashMap<i64, BinLifecycle>,
55 window: VecDeque<RecordContribution>,
56 alerts: Vec<IndicatorAlert>,
57}
58
59impl PersistentVolumeProfileEngine {
60 pub fn new(lookback: usize, bin_width: f64) -> Self {
63 Self {
64 lookback: lookback.max(1),
65 bin_width: bin_width.max(1e-9),
66 absorption_k: 2.5,
67 bins: HashMap::new(),
68 window: VecDeque::new(),
69 alerts: Vec::new(),
70 }
71 }
72
73 pub fn with_absorption_k(mut self, k: f64) -> Self {
74 self.absorption_k = k;
75 self
76 }
77
78 fn bin_key(&self, price: f64) -> i64 {
79 (price / self.bin_width).floor() as i64
80 }
81
82 fn bin_price_range(&self, key: i64) -> (f64, f64) {
83 (
84 key as f64 * self.bin_width,
85 (key as f64 + 1.0) * self.bin_width,
86 )
87 }
88
89 fn live_bins(&self) -> Vec<(i64, BinLifecycle)> {
91 let mut entries: Vec<(i64, BinLifecycle)> =
92 self.bins.iter().map(|(&k, &v)| (k, v)).collect();
93 entries.sort_by_key(|(k, _)| *k);
94 entries
95 }
96
97 pub fn absorption_profile(&self) -> Vec<AbsorptionBin> {
101 let live = self.live_bins();
102 if live.is_empty() {
103 return Vec::new();
104 }
105
106 let ratios: Vec<f64> = live
107 .iter()
108 .map(|(_, b)| b.volume / b.touches.max(1) as f64)
109 .collect();
110 let median = rolling_median(&ratios);
111 let abs_dev: Vec<f64> = ratios.iter().map(|r| (r - median).abs()).collect();
112 let mad = rolling_median(&abs_dev) * MAD_CONSISTENCY_CONSTANT;
113 let threshold = median + self.absorption_k * mad;
114
115 live.into_iter()
116 .zip(ratios)
117 .map(|((key, bin), ratio)| {
118 let (price_low, price_high) = self.bin_price_range(key);
119 AbsorptionBin {
120 price_low,
121 price_high,
122 volume: bin.volume,
123 touches: bin.touches,
124 volume_per_touch: ratio,
125 is_absorption: ratio > threshold,
130 }
131 })
132 .collect()
133 }
134
135 fn record_bar(&mut self, bar: &Bar) {
136 if !bar.low.is_finite()
140 || !bar.high.is_finite()
141 || !bar.volume.is_finite()
142 || bar.high < bar.low
143 {
144 return;
145 }
146
147 let bar_vol = if bar.volume > 0.0 {
148 bar.volume
149 } else {
150 bar.high - bar.low
151 };
152 let aggressor = estimate_aggressor_from_ohlc(bar);
153 let (buy_frac, sell_frac) = if bar.volume > 0.0 {
154 (
155 aggressor.buy_volume / bar.volume,
156 aggressor.sell_volume / bar.volume,
157 )
158 } else {
159 (0.5, 0.5)
160 };
161
162 let start_key = self.bin_key(bar.low);
163 let end_key = self.bin_key(bar.high).max(start_key);
164 let bin_count = (end_key - start_key + 1) as f64;
165 let vol_per_bin = bar_vol / bin_count;
166 let buy_per_bin = vol_per_bin * buy_frac;
167 let sell_per_bin = vol_per_bin * sell_frac;
168
169 let mut contribution = Vec::with_capacity((end_key - start_key + 1) as usize);
170 for key in start_key..=end_key {
171 let entry = self.bins.entry(key).or_insert(BinLifecycle {
172 volume: 0.0,
173 buy_volume: 0.0,
174 sell_volume: 0.0,
175 touches: 0,
176 first_touched_ts: bar.timestamp,
177 last_touched_ts: bar.timestamp,
178 });
179 entry.volume += vol_per_bin;
180 entry.buy_volume += buy_per_bin;
181 entry.sell_volume += sell_per_bin;
182 entry.touches += 1;
183 entry.last_touched_ts = bar.timestamp;
184 contribution.push((key, vol_per_bin, buy_per_bin, sell_per_bin));
185 }
186
187 self.window.push_back(RecordContribution {
188 per_bin: contribution,
189 });
190 if self.window.len() > self.lookback {
191 let evicted = self
192 .window
193 .pop_front()
194 .expect("just checked len > lookback");
195 for (key, vol, buy, sell) in evicted.per_bin {
196 let expired = if let Some(entry) = self.bins.get_mut(&key) {
197 entry.volume -= vol;
198 entry.buy_volume -= buy;
199 entry.sell_volume -= sell;
200 entry.touches = entry.touches.saturating_sub(1);
201 entry.touches == 0 || entry.volume <= 1e-9
202 } else {
203 false
204 };
205 if expired {
206 self.bins.remove(&key);
207 let (price_low, price_high) = self.bin_price_range(key);
208 self.alerts.push(IndicatorAlert::new(
209 "bin_expired",
210 format!("Price bin [{:.4}, {:.4}) expired", price_low, price_high),
211 0.3,
212 ));
213 }
214 }
215 }
216 }
217
218 fn build_output(&self) -> Option<IndicatorOutput> {
219 let live = self.live_bins();
220 if live.is_empty() {
221 return None;
222 }
223
224 let bins: Vec<ProfileBin> = live
225 .iter()
226 .map(|(key, b)| {
227 let (price_low, price_high) = self.bin_price_range(*key);
228 ProfileBin {
229 price_low,
230 price_high,
231 value: b.volume,
232 }
233 })
234 .collect();
235
236 let (poc_pos, poc_volume) =
237 live.iter()
238 .enumerate()
239 .fold((0usize, f64::MIN), |(bi, bv), (i, (_, b))| {
240 if b.volume > bv {
241 (i, b.volume)
242 } else {
243 (bi, bv)
244 }
245 });
246 let _ = poc_volume;
247 let poc_key = live[poc_pos].0;
248 let (poc_low, poc_high) = self.bin_price_range(poc_key);
249 let poc_price = (poc_low + poc_high) / 2.0;
250
251 let profile_artifact = ProfileArtifact {
252 kind: "persistent_volume_profile".to_string(),
253 bins,
254 poc: poc_price,
255 value_area_high: poc_high,
256 value_area_low: poc_low,
257 };
258
259 let absorption = self.absorption_profile();
260 let absorption_bins: Vec<ProfileBin> = absorption
261 .iter()
262 .map(|a| ProfileBin {
263 price_low: a.price_low,
264 price_high: a.price_high,
265 value: a.volume_per_touch,
266 })
267 .collect();
268 let absorption_artifact = ProfileArtifact {
269 kind: "absorption_profile".to_string(),
270 bins: absorption_bins,
271 poc: poc_price,
272 value_area_high: poc_high,
273 value_area_low: poc_low,
274 };
275
276 let mut output = IndicatorOutput::new(poc_price)
277 .with_artifact(profile_artifact)
278 .with_artifact(absorption_artifact);
279
280 for a in absorption.iter().filter(|a| a.is_absorption) {
281 output = output.with_artifact(ZoneArtifact {
282 kind: "absorption_zone".to_string(),
283 price_top: a.price_high,
284 price_bottom: a.price_low,
285 strength: (a.volume_per_touch).min(1.0),
286 touches: a.touches,
287 });
288 }
289
290 Some(output)
291 }
292}
293
294impl Indicator for PersistentVolumeProfileEngine {
295 fn name(&self) -> &str {
296 "persistent_volume_profile"
297 }
298
299 fn warmup_period(&self) -> usize {
300 self.lookback
301 }
302
303 fn reset(&mut self) {
304 self.bins.clear();
305 self.window.clear();
306 self.alerts.clear();
307 }
308
309 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
310 self.alerts.clear();
311 self.record_bar(bar);
312 if self.window.len() < self.lookback {
313 return None;
314 }
315 self.build_output()
316 }
317
318 fn alerts(&self) -> Vec<IndicatorAlert> {
319 self.alerts.clone()
320 }
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326
327 fn bar_at(price: f64, volume: f64) -> Bar {
330 Bar::new(0, price, price + 0.05, price - 0.05, price, volume)
331 }
332
333 #[test]
334 fn test_bin_persists_and_grows_across_updates() {
335 let mut engine = PersistentVolumeProfileEngine::new(3, 1.0);
336 engine.on_bar(&bar_at(100.2, 100.0));
337 let key = engine.bin_key(100.2);
338 assert_eq!(engine.bins.get(&key).unwrap().volume, 100.0);
339
340 engine.on_bar(&bar_at(100.3, 50.0));
341 assert_eq!(engine.bin_key(100.3), key);
343 assert_eq!(engine.bins.get(&key).unwrap().volume, 150.0);
344 assert_eq!(engine.bins.get(&key).unwrap().touches, 2);
345 }
346
347 #[test]
348 fn test_bin_dies_once_its_contributing_bars_roll_out() {
349 let mut engine = PersistentVolumeProfileEngine::new(2, 1.0);
350 let key = engine.bin_key(50.0);
351 engine.on_bar(&bar_at(50.0, 100.0));
352 assert!(engine.bins.contains_key(&key));
353
354 engine.on_bar(&bar_at(200.0, 10.0));
356 engine.on_bar(&bar_at(200.0, 10.0));
357
358 assert!(
359 !engine.bins.contains_key(&key),
360 "bin must expire once its only contributing bar leaves the window"
361 );
362 assert!(engine.alerts().iter().any(|a| a.kind == "bin_expired"));
363 }
364
365 const BASELINE_PRICES: [f64; 9] = [90.3, 92.3, 94.3, 96.3, 98.3, 102.3, 104.3, 106.3, 108.3];
368 const SPIKE_PRICE: f64 = 100.3;
369
370 #[test]
371 fn test_absorption_flags_concentrated_single_bar_volume() {
372 let mut engine = PersistentVolumeProfileEngine::new(10, 1.0).with_absorption_k(1.5);
373 for price in BASELINE_PRICES {
375 engine.on_bar(&bar_at(price, 50.0));
376 }
377 engine.on_bar(&bar_at(SPIKE_PRICE, 5000.0));
379
380 let absorption = engine.absorption_profile();
381 let flagged = absorption.iter().find(|a| a.is_absorption);
382 assert!(
383 flagged.is_some(),
384 "a single-touch volume spike must be flagged as absorption"
385 );
386 assert!(
387 flagged.unwrap().price_low <= SPIKE_PRICE && flagged.unwrap().price_high > SPIKE_PRICE
388 );
389 }
390
391 #[test]
392 fn test_evenly_touched_bin_is_not_absorption() {
393 let mut engine = PersistentVolumeProfileEngine::new(19, 1.0).with_absorption_k(1.5);
395 for price in BASELINE_PRICES {
396 engine.on_bar(&bar_at(price, 50.0));
397 }
398 for _ in 0..10 {
401 engine.on_bar(&bar_at(SPIKE_PRICE, 50.0));
402 }
403
404 let absorption = engine.absorption_profile();
405 let bin_spike = absorption
406 .iter()
407 .find(|a| a.price_low <= SPIKE_PRICE && a.price_high > SPIKE_PRICE)
408 .unwrap();
409 assert!(!bin_spike.is_absorption);
410 }
411
412 #[test]
413 fn test_none_until_lookback_filled() {
414 let mut engine = PersistentVolumeProfileEngine::new(4, 1.0);
415 for _ in 0..3 {
416 assert!(engine.on_bar(&bar_at(100.0, 10.0)).is_none());
417 }
418 assert!(engine.on_bar(&bar_at(100.0, 10.0)).is_some());
419 }
420}