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    /// First bar that touched this bin.
49    pub first_touched_ts: i64,
50    /// Most recent bar that touched this bin.
51    pub last_touched_ts: i64,
52}
53
54/// Persistent volume profile over the last `lookback` bars on a fixed price grid.
55///
56/// A bin is `[k * bin_width, (k + 1) * bin_width)` with `k = floor(price / bin_width)`. Each bar's
57/// volume — its range when it carries no volume — is spread evenly over the bins from its low's
58/// key to its high's key; when the bar leaves the window exactly that contribution is taken back,
59/// and a bin with no touches or no volume left is removed. `value` is the POC: the centre of the
60/// first (lowest-priced) live bin with the largest volume. The profile and a per-bin absorption
61/// profile come as [`ProfileArtifact`]s; bins whose volume per touch is a robust outlier (median
62/// plus `absorption_k` scaled MAD, `2.5` by default) additionally as [`ZoneArtifact`]s.
63///
64/// First output: with the `lookback`-th bar. [`Indicator::reset`] clears bins and window.
65pub 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    /// `bin_width` is a fixed price-grid resolution (not derived from the rolling window's
76    /// min/max), the mechanism that gives bins a stable identity across updates.
77    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    /// Currently live bins (born and not yet expired), price-ascending.
105    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    /// Per-bin absorption view: a bin is flagged when its volume-per-touch is a robust outlier
113    /// (median + `absorption_k` scaled-MAD) across the currently live bins — a small number of
114    /// bars depositing disproportionate volume at one price level without it rolling away.
115    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                    // Note: when a majority of bins share the same ratio, MAD is 0 and the
143                    // threshold collapses to the median itself — any value strictly above it is
144                    // still a real outlier (a tight majority plus one clear outsider), so this
145                    // does not require `mad > 0.0` as an extra gate.
146                    is_absorption: ratio > threshold,
147                }
148            })
149            .collect()
150    }
151
152    fn record_bar(&mut self, bar: &Bar) {
153        // A non-finite or inverted range cannot be mapped onto the fixed price grid (e.g.
154        // `high = inf` yields `bin_key` = `i64::MAX`, overflowing the span arithmetic below).
155        // Skip such a bar rather than panic; it simply contributes nothing this call.
156        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        // Die Bin-Zustände führen ihre Berührungszeiten bereits — daraus ergibt sich
270        // das Fenster, über das dieses Profil gewachsen ist.
271        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                // Die Zone besteht, seit ihr Bin zuerst berührt wurde.
329                .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    /// A narrow-range bar (high-low = 0.1) so it always lands inside a single 1.0-wide bin
371    /// regardless of where `price` falls relative to a bin boundary.
372    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        // Same bin (same 1.0-wide grid cell around 100), volume accumulated rather than replaced.
385        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        // Two more bars at a distant price roll the original bar out of the lookback=2 window.
398        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    // Baseline/spike prices deliberately offset from whole numbers (bin_width = 1.0) so their
409    // narrow +/-0.05 range never straddles a bin boundary.
410    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        // Baseline: modest, evenly distributed volume across several distinct price levels.
417        for price in BASELINE_PRICES {
418            engine.on_bar(&bar_at(price, 50.0));
419        }
420        // One outlier bar dumps a huge amount of volume into a single new bin in one touch.
421        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        // Large enough lookback that none of the 19 bars fed below roll out of the window.
437        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        // Same per-touch volume (50) as the baseline bins, just reached over 10 touches at one
442        // level instead of a single one -> normal participation, not absorption.
443        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}