kestrel_chartkit/indicator/
money_flow_profile.rs1use 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
36pub 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 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 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 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; let flow = v * overlap * mf_price; 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 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 #[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 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}