Skip to main content

kestrel_chartkit/indicator/
volume_profile_persistent.rs

1//! Persistent price/volume profile: bins keyed by a fixed price grid (not the rolling window's
2//! current min/max), so a specific price level's bin has a real lifecycle — it is born on first
3//! touch, grows/shrinks incrementally as bars enter and leave the trailing window, and is removed
4//! once its contributing bars have all rolled out — plus a dedicated per-bin absorption profile,
5//! rather than the bar-level absorption [`super::volume_flow_hires::HiResVolumeFlowEngine`]
6//! computes. Complements [`super::volume_profile_extended::ExtendedVolumeProfileEngine`], which
7//! recomputes its bins from scratch every call over the window's current price range and has no
8//! notion of a bin persisting across updates.
9
10use 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
19/// MAD-to-stddev consistency constant (see [`crate::clustering`]), reused here for the
20/// per-bin absorption threshold.
21const 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
33/// One bar's contribution to each bin it spanned, kept so evicting the bar from the trailing
34/// window can precisely reverse its effect on those bins.
35struct RecordContribution {
36    per_bin: Vec<(i64, f64, f64, f64)>,
37}
38
39/// A snapshot view of one persistent bin's current lifecycle state.
40#[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    /// `bin_width` is a fixed price-grid resolution (not derived from the rolling window's
61    /// min/max), the mechanism that gives bins a stable identity across updates.
62    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    /// Currently live bins (born and not yet expired), price-ascending.
90    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    /// Per-bin absorption view: a bin is flagged when its volume-per-touch is a robust outlier
98    /// (median + `absorption_k` scaled-MAD) across the currently live bins — a small number of
99    /// bars depositing disproportionate volume at one price level without it rolling away.
100    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                    // Note: when a majority of bins share the same ratio, MAD is 0 and the
126                    // threshold collapses to the median itself — any value strictly above it is
127                    // still a real outlier (a tight majority plus one clear outsider), so this
128                    // does not require `mad > 0.0` as an extra gate.
129                    is_absorption: ratio > threshold,
130                }
131            })
132            .collect()
133    }
134
135    fn record_bar(&mut self, bar: &Bar) {
136        // A non-finite or inverted range cannot be mapped onto the fixed price grid (e.g.
137        // `high = inf` yields `bin_key` = `i64::MAX`, overflowing the span arithmetic below).
138        // Skip such a bar rather than panic; it simply contributes nothing this call.
139        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    /// A narrow-range bar (high-low = 0.1) so it always lands inside a single 1.0-wide bin
328    /// regardless of where `price` falls relative to a bin boundary.
329    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        // Same bin (same 1.0-wide grid cell around 100), volume accumulated rather than replaced.
342        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        // Two more bars at a distant price roll the original bar out of the lookback=2 window.
355        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    // Baseline/spike prices deliberately offset from whole numbers (bin_width = 1.0) so their
366    // narrow +/-0.05 range never straddles a bin boundary.
367    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        // Baseline: modest, evenly distributed volume across several distinct price levels.
374        for price in BASELINE_PRICES {
375            engine.on_bar(&bar_at(price, 50.0));
376        }
377        // One outlier bar dumps a huge amount of volume into a single new bin in one touch.
378        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        // Large enough lookback that none of the 19 bars fed below roll out of the window.
394        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        // Same per-touch volume (50) as the baseline bins, just reached over 10 touches at one
399        // level instead of a single one -> normal participation, not absorption.
400        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}