1use 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 pub first_touched_ts: i64,
50 pub last_touched_ts: i64,
52}
53
54pub struct PersistentVolumeProfileEngine {
66 lookback: usize,
67 bin_width: f64,
68 absorption_k: f64,
69 bins: HashMap<i64, BinLifecycle>,
70 window: VecDeque<RecordContribution>,
71 alerts: Vec<IndicatorAlert>,
72}
73
74impl PersistentVolumeProfileEngine {
75 pub fn new(lookback: usize, bin_width: f64) -> Self {
78 Self {
79 lookback: lookback.max(1),
80 bin_width: bin_width.max(1e-9),
81 absorption_k: 2.5,
82 bins: HashMap::new(),
83 window: VecDeque::new(),
84 alerts: Vec::new(),
85 }
86 }
87
88 pub fn with_absorption_k(mut self, k: f64) -> Self {
89 self.absorption_k = k;
90 self
91 }
92
93 fn bin_key(&self, price: f64) -> i64 {
94 (price / self.bin_width).floor() as i64
95 }
96
97 fn bin_price_range(&self, key: i64) -> (f64, f64) {
98 (
99 key as f64 * self.bin_width,
100 (key as f64 + 1.0) * self.bin_width,
101 )
102 }
103
104 fn live_bins(&self) -> Vec<(i64, BinLifecycle)> {
106 let mut entries: Vec<(i64, BinLifecycle)> =
107 self.bins.iter().map(|(&k, &v)| (k, v)).collect();
108 entries.sort_by_key(|(k, _)| *k);
109 entries
110 }
111
112 pub fn absorption_profile(&self) -> Vec<AbsorptionBin> {
116 let live = self.live_bins();
117 if live.is_empty() {
118 return Vec::new();
119 }
120
121 let ratios: Vec<f64> = live
122 .iter()
123 .map(|(_, b)| b.volume / b.touches.max(1) as f64)
124 .collect();
125 let median = rolling_median(&ratios);
126 let abs_dev: Vec<f64> = ratios.iter().map(|r| (r - median).abs()).collect();
127 let mad = rolling_median(&abs_dev) * MAD_CONSISTENCY_CONSTANT;
128 let threshold = median + self.absorption_k * mad;
129
130 live.into_iter()
131 .zip(ratios)
132 .map(|((key, bin), ratio)| {
133 let (price_low, price_high) = self.bin_price_range(key);
134 AbsorptionBin {
135 price_low,
136 price_high,
137 volume: bin.volume,
138 touches: bin.touches,
139 volume_per_touch: ratio,
140 first_touched_ts: bin.first_touched_ts,
141 last_touched_ts: bin.last_touched_ts,
142 is_absorption: ratio > threshold,
147 }
148 })
149 .collect()
150 }
151
152 fn record_bar(&mut self, bar: &Bar) {
153 if !bar.low.is_finite()
157 || !bar.high.is_finite()
158 || !bar.volume.is_finite()
159 || bar.high < bar.low
160 {
161 return;
162 }
163
164 let bar_vol = if bar.volume > 0.0 {
165 bar.volume
166 } else {
167 bar.high - bar.low
168 };
169 let aggressor = estimate_aggressor_from_ohlc(bar);
170 let (buy_frac, sell_frac) = if bar.volume > 0.0 {
171 (
172 aggressor.buy_volume / bar.volume,
173 aggressor.sell_volume / bar.volume,
174 )
175 } else {
176 (0.5, 0.5)
177 };
178
179 let start_key = self.bin_key(bar.low);
180 let end_key = self.bin_key(bar.high).max(start_key);
181 let bin_count = (end_key - start_key + 1) as f64;
182 let vol_per_bin = bar_vol / bin_count;
183 let buy_per_bin = vol_per_bin * buy_frac;
184 let sell_per_bin = vol_per_bin * sell_frac;
185
186 let cap = usize::try_from(end_key - start_key + 1).unwrap_or(0);
187 let mut contribution = Vec::with_capacity(cap);
188 for key in start_key..=end_key {
189 let entry = self.bins.entry(key).or_insert(BinLifecycle {
190 volume: 0.0,
191 buy_volume: 0.0,
192 sell_volume: 0.0,
193 touches: 0,
194 first_touched_ts: bar.timestamp,
195 last_touched_ts: bar.timestamp,
196 });
197 entry.volume += vol_per_bin;
198 entry.buy_volume += buy_per_bin;
199 entry.sell_volume += sell_per_bin;
200 entry.touches += 1;
201 entry.last_touched_ts = bar.timestamp;
202 contribution.push((key, vol_per_bin, buy_per_bin, sell_per_bin));
203 }
204
205 self.window.push_back(RecordContribution {
206 per_bin: contribution,
207 });
208 if self.window.len() > self.lookback {
209 let evicted = self
210 .window
211 .pop_front()
212 .expect("just checked len > lookback");
213 for (key, vol, buy, sell) in evicted.per_bin {
214 let expired = if let Some(entry) = self.bins.get_mut(&key) {
215 entry.volume -= vol;
216 entry.buy_volume -= buy;
217 entry.sell_volume -= sell;
218 entry.touches = entry.touches.saturating_sub(1);
219 entry.touches == 0 || entry.volume <= 1e-9
220 } else {
221 false
222 };
223 if expired {
224 self.bins.remove(&key);
225 let (price_low, price_high) = self.bin_price_range(key);
226 self.alerts.push(IndicatorAlert::new(
227 "bin_expired",
228 format!("Price bin [{:.4}, {:.4}) expired", price_low, price_high),
229 0.3,
230 ));
231 }
232 }
233 }
234 }
235
236 fn build_output(&self) -> Option<IndicatorOutput> {
237 let live = self.live_bins();
238 if live.is_empty() {
239 return None;
240 }
241
242 let bins: Vec<ProfileBin> = live
243 .iter()
244 .map(|(key, b)| {
245 let (price_low, price_high) = self.bin_price_range(*key);
246 ProfileBin {
247 price_low,
248 price_high,
249 value: b.volume,
250 }
251 })
252 .collect();
253
254 let (poc_pos, poc_volume) =
255 live.iter()
256 .enumerate()
257 .fold((0usize, f64::MIN), |(bi, bv), (i, (_, b))| {
258 if b.volume > bv {
259 (i, b.volume)
260 } else {
261 (bi, bv)
262 }
263 });
264 let _ = poc_volume;
265 let poc_key = live[poc_pos].0;
266 let (poc_low, poc_high) = self.bin_price_range(poc_key);
267 let poc_price = (poc_low + poc_high) / 2.0;
268
269 let window = live
272 .iter()
273 .fold(None::<(i64, i64)>, |acc, (_, b)| match acc {
274 None => Some((b.first_touched_ts, b.last_touched_ts)),
275 Some((from, to)) => Some((from.min(b.first_touched_ts), to.max(b.last_touched_ts))),
276 });
277
278 let mut profile_artifact = ProfileArtifact {
279 kind: "persistent_volume_profile".to_string(),
280 bins,
281 poc: poc_price,
282 value_area_high: poc_high,
283 value_area_low: poc_low,
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 }
290
291 let absorption = self.absorption_profile();
292 let absorption_bins: Vec<ProfileBin> = absorption
293 .iter()
294 .map(|a| ProfileBin {
295 price_low: a.price_low,
296 price_high: a.price_high,
297 value: a.volume_per_touch,
298 })
299 .collect();
300 let mut absorption_artifact = ProfileArtifact {
301 kind: "absorption_profile".to_string(),
302 bins: absorption_bins,
303 poc: poc_price,
304 value_area_high: poc_high,
305 value_area_low: poc_low,
306 from_ts: None,
307 to_ts: None,
308 };
309 if let Some((from, to)) = window {
310 absorption_artifact = absorption_artifact.spanning(from, to);
311 }
312
313 let mut output = IndicatorOutput::new(poc_price)
314 .with_artifact(profile_artifact)
315 .with_artifact(absorption_artifact);
316
317 for a in absorption.iter().filter(|a| a.is_absorption) {
318 output = output.with_artifact(
319 ZoneArtifact {
320 kind: "absorption_zone".to_string(),
321 price_top: a.price_high,
322 price_bottom: a.price_low,
323 strength: (a.volume_per_touch).min(1.0),
324 touches: a.touches,
325 from_ts: None,
326 to_ts: None,
327 }
328 .spanning(a.first_touched_ts, a.last_touched_ts),
330 );
331 }
332
333 Some(output)
334 }
335}
336
337impl Indicator for PersistentVolumeProfileEngine {
338 fn name(&self) -> &str {
339 "persistent_volume_profile"
340 }
341
342 fn warmup_period(&self) -> usize {
343 self.lookback
344 }
345
346 fn reset(&mut self) {
347 self.bins.clear();
348 self.window.clear();
349 self.alerts.clear();
350 }
351
352 fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
353 self.alerts.clear();
354 self.record_bar(bar);
355 if self.window.len() < self.lookback {
356 return None;
357 }
358 self.build_output()
359 }
360
361 fn alerts(&self) -> Vec<IndicatorAlert> {
362 self.alerts.clone()
363 }
364}
365
366#[cfg(test)]
367mod tests {
368 use super::*;
369
370 fn bar_at(price: f64, volume: f64) -> Bar {
373 Bar::new(0, price, price + 0.05, price - 0.05, price, volume)
374 }
375
376 #[test]
377 fn test_bin_persists_and_grows_across_updates() {
378 let mut engine = PersistentVolumeProfileEngine::new(3, 1.0);
379 engine.on_bar(&bar_at(100.2, 100.0));
380 let key = engine.bin_key(100.2);
381 assert_eq!(engine.bins.get(&key).unwrap().volume, 100.0);
382
383 engine.on_bar(&bar_at(100.3, 50.0));
384 assert_eq!(engine.bin_key(100.3), key);
386 assert_eq!(engine.bins.get(&key).unwrap().volume, 150.0);
387 assert_eq!(engine.bins.get(&key).unwrap().touches, 2);
388 }
389
390 #[test]
391 fn test_bin_dies_once_its_contributing_bars_roll_out() {
392 let mut engine = PersistentVolumeProfileEngine::new(2, 1.0);
393 let key = engine.bin_key(50.0);
394 engine.on_bar(&bar_at(50.0, 100.0));
395 assert!(engine.bins.contains_key(&key));
396
397 engine.on_bar(&bar_at(200.0, 10.0));
399 engine.on_bar(&bar_at(200.0, 10.0));
400
401 assert!(
402 !engine.bins.contains_key(&key),
403 "bin must expire once its only contributing bar leaves the window"
404 );
405 assert!(engine.alerts().iter().any(|a| a.kind == "bin_expired"));
406 }
407
408 const BASELINE_PRICES: [f64; 9] = [90.3, 92.3, 94.3, 96.3, 98.3, 102.3, 104.3, 106.3, 108.3];
411 const SPIKE_PRICE: f64 = 100.3;
412
413 #[test]
414 fn test_absorption_flags_concentrated_single_bar_volume() {
415 let mut engine = PersistentVolumeProfileEngine::new(10, 1.0).with_absorption_k(1.5);
416 for price in BASELINE_PRICES {
418 engine.on_bar(&bar_at(price, 50.0));
419 }
420 engine.on_bar(&bar_at(SPIKE_PRICE, 5000.0));
422
423 let absorption = engine.absorption_profile();
424 let flagged = absorption.iter().find(|a| a.is_absorption);
425 assert!(
426 flagged.is_some(),
427 "a single-touch volume spike must be flagged as absorption"
428 );
429 assert!(
430 flagged.unwrap().price_low <= SPIKE_PRICE && flagged.unwrap().price_high > SPIKE_PRICE
431 );
432 }
433
434 #[test]
435 fn test_evenly_touched_bin_is_not_absorption() {
436 let mut engine = PersistentVolumeProfileEngine::new(19, 1.0).with_absorption_k(1.5);
438 for price in BASELINE_PRICES {
439 engine.on_bar(&bar_at(price, 50.0));
440 }
441 for _ in 0..10 {
444 engine.on_bar(&bar_at(SPIKE_PRICE, 50.0));
445 }
446
447 let absorption = engine.absorption_profile();
448 let bin_spike = absorption
449 .iter()
450 .find(|a| a.price_low <= SPIKE_PRICE && a.price_high > SPIKE_PRICE)
451 .unwrap();
452 assert!(!bin_spike.is_absorption);
453 }
454
455 #[test]
456 fn test_none_until_lookback_filled() {
457 let mut engine = PersistentVolumeProfileEngine::new(4, 1.0);
458 for _ in 0..3 {
459 assert!(engine.on_bar(&bar_at(100.0, 10.0)).is_none());
460 }
461 assert!(engine.on_bar(&bar_at(100.0, 10.0)).is_some());
462 }
463}