Skip to main content

kestrel_chartkit/indicator/
money_flow_profile.rs

1//! Money Flow Profile: a rolling-window Volume-by-Price profile that bins by **dollar volume**
2//! (`volume × row-mid-price`) instead of raw volume, plus an aggregate bull%/bear% flow-bias
3//! scalar with debounced (50%-crossing) bias-flip alerts and distance-scaled value-area breakout
4//! alerts.
5//!
6//! Chartkit's existing [`super::volume_profile`]/[`super::volume_profile_extended`]/
7//! [`super::volume_profile_persistent`] family bins by raw `bar.volume` and is, taken together,
8//! more feature-rich overall (HVN/LVN zones, absorption, per-bin delta). What none of them do is
9//! dollar-volume weighting — for a wide-range instrument this genuinely shifts where POC/VAH/VAL
10//! land, since a bar with modest volume at a high price level can outweigh a bar with heavy
11//! volume at a low price level once weighted by price. Kept as a separate engine rather than a
12//! mode flag on the existing family, matching the precedent of that family already being three
13//! separate engines rather than one engine with every mode as a flag.
14//!
15//! This is a **windowed recompute**, not an incremental running average: row boundaries are
16//! re-derived from the window's current high/low every bar (there is no meaningful "add one, drop
17//! one" incremental step once bin edges themselves move), same as
18//! [`super::volume_profile::VolumeProfileEngine`].
19//!
20//! Scope: the numeric substance only — POC, Value Area High/Low, Delta POC and the overall bull%
21//! flow bias. Drawing, HVN/LVN zone overlays, absorption, intrabar delta and alternative
22//! money-flow/price modes are not part of this engine.
23//!
24//! Bars with `volume <= 0.0` are treated as `1.0` (synthetic volume) rather than dropped, a
25//! deliberate default that matters for feeds that report `volume = 0`/unavailable.
26
27use std::collections::VecDeque;
28
29use crate::model::Bar;
30
31use super::smoothing::{crossed_over, crossed_under};
32use super::{Indicator, IndicatorAlert, IndicatorOutput};
33
34use std::collections::HashMap;
35
36/// Money Flow Profile engine — see the module doc comment for how this relates to the
37/// [`super::volume_profile`] family.
38///
39/// Per bar, over the last `lookback` bars (at least two): the window's lowest low `L` and highest
40/// high `H` are split into `rows` bins of height `step = (H - L) / rows`. Every bar with a range
41/// spreads its flow over the bins it overlaps, in proportion to the overlap — the part of its
42/// `high - low` inside the bin over `high - low`. A bin receives `volume · overlap · mid` from the
43/// bar, `mid` the bin's middle price and a non-positive volume counted as 1; the bullish part of
44/// that is the same times the bar's `clamp((close - low) / (high - low), 0, 1)`.
45///
46/// `value` is the POC, the middle of the bin with the largest flow (the upper bin on a tie).
47/// `extra["delta_poc"]` is the middle of the first bin with the largest `|2 · bullish - flow|`.
48/// The value area grows from the POC bin one neighbour at a time — the one with more flow, the
49/// lower one on a tie — until it holds `va_pct` of the total flow; `extra["vah"]`/`extra["val"]`
50/// are its outer bin edges. `extra["bull_pct"] = 100 · sum bullish / sum flow`.
51///
52/// From the second output on, alerts fire when the close crosses above the value-area high or
53/// below the low — at or below (above) the previous bar's level before, beyond the current level
54/// now — with the distance beyond it over the value-area width, clamped to `0..=1`, as strength;
55/// and when `bull_pct` crosses 50 in the same sense. `None` for a single bar, a window without
56/// range, or no flow at all.
57pub struct MoneyFlowProfileEngine {
58    lookback: usize,
59    rows: usize,
60    va_pct: f64,
61
62    window: VecDeque<Bar>,
63    prev_close: Option<f64>,
64    prev_vah: Option<f64>,
65    prev_val: Option<f64>,
66    prev_bull_pct: Option<f64>,
67
68    alerts: Vec<IndicatorAlert>,
69}
70
71impl MoneyFlowProfileEngine {
72    /// `lookback` bars form the rolling window; `rows` bins the window's high/low range;
73    /// `va_pct` is the fraction of total flow the value area expands to capture from POC
74    /// outward (default `0.70`).
75    pub fn new(lookback: usize, rows: usize, va_pct: f64) -> Self {
76        let lookback = lookback.max(1);
77        let rows = rows.max(1);
78        Self {
79            lookback,
80            rows,
81            va_pct,
82            window: VecDeque::with_capacity(lookback),
83            prev_close: None,
84            prev_vah: None,
85            prev_val: None,
86            prev_bull_pct: None,
87            alerts: Vec::new(),
88        }
89    }
90
91    /// Defaults: `lookback=200`, `rows=25`, `va_pct=0.70`.
92    pub fn with_defaults() -> Self {
93        Self::new(200, 25, 0.70)
94    }
95}
96
97impl Indicator for MoneyFlowProfileEngine {
98    fn name(&self) -> &str {
99        "money_flow_profile"
100    }
101
102    fn warmup_period(&self) -> usize {
103        // Technical minimum for a non-degenerate high/low range: the profile is computed over
104        // whatever history exists rather than waiting for the full lookback, and only becomes
105        // operationally meaningful once the window has accumulated close to `lookback` bars.
106        2
107    }
108
109    fn reset(&mut self) {
110        self.window.clear();
111        self.prev_close = None;
112        self.prev_vah = None;
113        self.prev_val = None;
114        self.prev_bull_pct = None;
115        self.alerts.clear();
116    }
117
118    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
119        self.alerts.clear();
120
121        if self.window.len() == self.lookback {
122            self.window.pop_front();
123        }
124        self.window.push_back(bar.clone());
125        if self.window.len() < 2 {
126            return None;
127        }
128
129        let p_lo = self
130            .window
131            .iter()
132            .map(|b| b.low)
133            .fold(f64::INFINITY, f64::min);
134        let p_hi = self
135            .window
136            .iter()
137            .map(|b| b.high)
138            .fold(f64::NEG_INFINITY, f64::max);
139        if p_hi <= p_lo {
140            return None;
141        }
142        let p_step = (p_hi - p_lo) / self.rows as f64;
143
144        let mut total_flow = vec![0.0_f64; self.rows];
145        let mut bull_flow = vec![0.0_f64; self.rows];
146
147        for b in &self.window {
148            let (h, l, c) = (b.high, b.low, b.close);
149            if h <= l {
150                continue;
151            }
152            let v = if b.volume > 0.0 { b.volume } else { 1.0 };
153            let buy_ratio = ((c - l) / (h - l)).clamp(0.0, 1.0);
154
155            for r in 0..self.rows {
156                let row_lo = p_lo + r as f64 * p_step;
157                let row_hi = row_lo + p_step;
158                if h < row_lo || l >= row_hi {
159                    continue;
160                }
161                let overlap = if l >= row_lo && h > row_hi {
162                    (row_hi - l) / (h - l)
163                } else if h <= row_hi && l < row_lo {
164                    (h - row_lo) / (h - l)
165                } else if l >= row_lo && h <= row_hi {
166                    1.0
167                } else {
168                    p_step / (h - l)
169                };
170
171                let mf_price = p_lo + (r as f64 + 0.5) * p_step; // Row Mid (default)
172                let flow = v * overlap * mf_price; // Money Flow source (default): dollar volume.
173                total_flow[r] += flow;
174                bull_flow[r] += flow * buy_ratio;
175            }
176        }
177
178        let tot_max = total_flow.iter().cloned().fold(0.0_f64, f64::max);
179        if tot_max <= 0.0 {
180            return None;
181        }
182        let tot_sum: f64 = total_flow.iter().sum();
183        let poc_idx = total_flow
184            .iter()
185            .enumerate()
186            .max_by(|a, b| a.1.partial_cmp(b.1).unwrap())
187            .map(|(i, _)| i)
188            .unwrap();
189        let poc_price = p_lo + (poc_idx as f64 + 0.5) * p_step;
190
191        let mut delta_poc_idx = 0usize;
192        let mut delta_poc_abs_max = 0.0_f64;
193        for r in 0..self.rows {
194            let d = (2.0 * bull_flow[r] - total_flow[r]).abs();
195            if d > delta_poc_abs_max {
196                delta_poc_abs_max = d;
197                delta_poc_idx = r;
198            }
199        }
200        let delta_poc_price = p_lo + (delta_poc_idx as f64 + 0.5) * p_step;
201
202        // Value Area: expand from POC toward the highest adjacent row until `va_pct` of total
203        // flow is captured.
204        let mut va_lo = poc_idx;
205        let mut va_hi = poc_idx;
206        let mut va_acc = total_flow[poc_idx];
207        let va_tgt = tot_sum * self.va_pct;
208        while va_acc < va_tgt {
209            let add_lo = if va_lo > 0 {
210                total_flow[va_lo - 1]
211            } else {
212                -1.0
213            };
214            let add_hi = if va_hi < self.rows - 1 {
215                total_flow[va_hi + 1]
216            } else {
217                -1.0
218            };
219            if add_lo < 0.0 && add_hi < 0.0 {
220                break;
221            }
222            if add_lo >= add_hi {
223                va_lo -= 1;
224                va_acc += add_lo;
225            } else {
226                va_hi += 1;
227                va_acc += add_hi;
228            }
229        }
230        let vah_price = p_lo + (va_hi + 1) as f64 * p_step;
231        let val_price = p_lo + va_lo as f64 * p_step;
232
233        let bull_pct = bull_flow.iter().sum::<f64>() / tot_sum * 100.0;
234
235        let mut vah_breakout = false;
236        let mut val_breakdown = false;
237        let mut bull_bias = false;
238        let mut bear_bias = false;
239        let mut vah_breakout_strength = 0.0;
240        let mut val_breakdown_strength = 0.0;
241
242        if let (Some(prev_close), Some(prev_vah), Some(prev_val), Some(prev_bull_pct)) = (
243            self.prev_close,
244            self.prev_vah,
245            self.prev_val,
246            self.prev_bull_pct,
247        ) {
248            vah_breakout = crossed_over(prev_close, prev_vah, bar.close, vah_price);
249            val_breakdown = crossed_under(prev_close, prev_val, bar.close, val_price);
250            bull_bias = crossed_over(prev_bull_pct, 50.0, bull_pct, 50.0);
251            bear_bias = crossed_under(prev_bull_pct, 50.0, bull_pct, 50.0);
252
253            let va_width = vah_price - val_price;
254            vah_breakout_strength = if va_width > 0.0 {
255                ((bar.close - vah_price) / va_width).clamp(0.0, 1.0)
256            } else {
257                1.0
258            };
259            val_breakdown_strength = if va_width > 0.0 {
260                ((val_price - bar.close) / va_width).clamp(0.0, 1.0)
261            } else {
262                1.0
263            };
264        }
265
266        if vah_breakout {
267            self.alerts.push(IndicatorAlert::new(
268                "vah_breakout",
269                "Money Flow Profile: close crossed above the Value Area High",
270                vah_breakout_strength,
271            ));
272        }
273        if val_breakdown {
274            self.alerts.push(IndicatorAlert::new(
275                "val_breakdown",
276                "Money Flow Profile: close crossed below the Value Area Low",
277                val_breakdown_strength,
278            ));
279        }
280        if bull_bias {
281            self.alerts.push(IndicatorAlert::new(
282                "bull_bias",
283                "Money Flow Profile: flow bias turned bullish",
284                1.0,
285            ));
286        }
287        if bear_bias {
288            self.alerts.push(IndicatorAlert::new(
289                "bear_bias",
290                "Money Flow Profile: flow bias turned bearish",
291                1.0,
292            ));
293        }
294
295        self.prev_close = Some(bar.close);
296        self.prev_vah = Some(vah_price);
297        self.prev_val = Some(val_price);
298        self.prev_bull_pct = Some(bull_pct);
299
300        let mut extra = HashMap::new();
301        extra.insert("vah".to_string(), vah_price);
302        extra.insert("val".to_string(), val_price);
303        extra.insert("delta_poc".to_string(), delta_poc_price);
304        extra.insert("bull_pct".to_string(), bull_pct);
305
306        Some(IndicatorOutput::with_extra(poc_price, extra))
307    }
308
309    fn alerts(&self) -> Vec<IndicatorAlert> {
310        self.alerts.clone()
311    }
312}
313
314#[cfg(test)]
315mod tests {
316    use super::*;
317
318    // Reference values: `tests/golden_reference_volume.rs`, derived in `reference/`.
319
320    #[test]
321    fn synthetic_volume_substitutes_for_non_positive_volume() {
322        let mut mfp = MoneyFlowProfileEngine::new(2, 5, 0.70);
323        let bars = [
324            Bar::new(1, 100.0, 101.0, 99.0, 100.0, 0.0),
325            Bar::new(2, 101.0, 102.0, 100.0, 101.0, 0.0),
326        ];
327        // Both bars report zero volume; the engine must still produce a profile (treating volume
328        // as 1.0) instead of degenerating to an all-zero, POC-less window.
329        mfp.on_bar(&bars[0]);
330        let out = mfp.on_bar(&bars[1]);
331        assert!(
332            out.is_some(),
333            "zero-volume bars must fall back to synthetic volume, not None"
334        );
335    }
336
337    #[test]
338    fn reset_clears_window_and_cross_state() {
339        let mut mfp = MoneyFlowProfileEngine::new(2, 10, 0.70);
340        mfp.on_bar(&Bar::new(1, 10.0, 10.1, 9.9, 10.0, 1000.0));
341        mfp.on_bar(&Bar::new(2, 49.9, 50.0, 49.8, 49.9, 300.0));
342        mfp.reset();
343        assert!(mfp
344            .on_bar(&Bar::new(1, 10.0, 10.1, 9.9, 10.0, 1000.0))
345            .is_none());
346    }
347}