Skip to main content

kestrel_chartkit/indicator/
volume_indicators.rs

1use std::collections::{HashMap, VecDeque};
2
3use crate::indicator::{Indicator, IndicatorAlert, IndicatorOutput};
4use crate::model::Bar;
5
6/// Bar volume with its moving average.
7///
8/// `value` is the bar's own volume. Once `ma_period` volumes have been seen, `extra` adds
9/// `avg_volume` — the plain mean of the last `ma_period` volumes, this bar included — and
10/// `volume_ratio = volume / avg_volume` (`1` for a zero average); a volume above twice the average
11/// raises an alert. Before that, only the volume itself is published.
12///
13/// First output: with the first bar. [`Indicator::reset`] clears the window.
14pub struct VolumeEngine {
15    ma_period: usize,
16    volumes: VecDeque<f64>,
17    alerts: Vec<IndicatorAlert>,
18}
19
20impl VolumeEngine {
21    pub fn new(ma_period: usize) -> Self {
22        Self {
23            ma_period,
24            volumes: VecDeque::new(),
25            alerts: Vec::new(),
26        }
27    }
28}
29
30impl Indicator for VolumeEngine {
31    fn name(&self) -> &str {
32        "volume"
33    }
34
35    fn warmup_period(&self) -> usize {
36        self.ma_period
37    }
38
39    fn reset(&mut self) {
40        self.volumes.clear();
41        self.alerts.clear();
42    }
43
44    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
45        let vol = bar.volume;
46        self.volumes.push_back(vol);
47        if self.volumes.len() > self.ma_period {
48            self.volumes.pop_front();
49        }
50
51        self.alerts.clear();
52        if self.volumes.len() < self.ma_period {
53            return Some(IndicatorOutput::new(vol));
54        }
55
56        let avg_vol: f64 = self.volumes.iter().sum::<f64>() / self.ma_period as f64;
57        let mut extra = HashMap::new();
58        extra.insert("avg_volume".to_string(), avg_vol);
59        extra.insert(
60            "volume_ratio".to_string(),
61            if avg_vol > 0.0 { vol / avg_vol } else { 1.0 },
62        );
63
64        if avg_vol > 0.0 && vol > 2.0 * avg_vol {
65            self.alerts.push(IndicatorAlert::new(
66                "high_volume_spike",
67                format!("High Volume Spike: {:.0} (>2.0x avg {:.0})", vol, avg_vol),
68                0.80,
69            ));
70        }
71
72        Some(IndicatorOutput::with_extra(vol, extra))
73    }
74
75    fn alerts(&self) -> Vec<IndicatorAlert> {
76        self.alerts.clone()
77    }
78}
79
80/// Relative Volume (RVOL): this bar's volume against the recent average.
81///
82/// `RVOL = volume / avg`, where `avg` is the plain mean of the last `period` volumes **including
83/// this bar's own**, and `1` for a zero average. Including the current bar damps the ratio: a
84/// spike raises its own reference. For the comparison against the same time of day on earlier
85/// days see [`super::rvat::RelativeVolumeAtTime`].
86///
87/// First output: with the `period`-th bar. [`Indicator::reset`] clears the window.
88#[derive(Debug, Clone)]
89pub struct RvolEngine {
90    period: usize,
91    volumes: VecDeque<f64>,
92    alerts: Vec<IndicatorAlert>,
93}
94
95impl RvolEngine {
96    pub fn new(period: usize) -> Self {
97        Self {
98            period,
99            volumes: VecDeque::new(),
100            alerts: Vec::new(),
101        }
102    }
103}
104
105impl Indicator for RvolEngine {
106    fn name(&self) -> &str {
107        "rvol"
108    }
109
110    fn warmup_period(&self) -> usize {
111        self.period
112    }
113
114    fn reset(&mut self) {
115        self.volumes.clear();
116        self.alerts.clear();
117    }
118
119    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
120        let vol = bar.volume;
121        self.volumes.push_back(vol);
122        if self.volumes.len() > self.period {
123            self.volumes.pop_front();
124        }
125
126        self.alerts.clear();
127        if self.volumes.len() < self.period {
128            return None;
129        }
130
131        let avg_vol: f64 = self.volumes.iter().sum::<f64>() / self.period as f64;
132        let rvol = if avg_vol > 0.0 { vol / avg_vol } else { 1.0 };
133
134        if rvol >= 2.5 {
135            self.alerts.push(IndicatorAlert::new(
136                "extreme_rvol",
137                format!("Extreme Relative Volume: {:.2}x", rvol),
138                0.90,
139            ));
140        }
141
142        Some(IndicatorOutput::new(rvol))
143    }
144
145    fn alerts(&self) -> Vec<IndicatorAlert> {
146        self.alerts.clone()
147    }
148}
149
150/// On-Balance Volume (OBV): a running total of volume signed by the direction of the close.
151///
152/// Starting at `0`, each bar adds its volume when the close rose against the previous close,
153/// subtracts it when the close fell, and leaves the total unchanged when the close is equal. The
154/// first bar has no previous close and contributes nothing.
155///
156/// First output: with the first bar. [`Indicator::reset`] returns the total to zero.
157pub struct ObvEngine {
158    prev_close: Option<f64>,
159    cum_obv: f64,
160    alerts: Vec<IndicatorAlert>,
161}
162
163impl ObvEngine {
164    pub fn new() -> Self {
165        Self {
166            prev_close: None,
167            cum_obv: 0.0,
168            alerts: Vec::new(),
169        }
170    }
171}
172
173impl Default for ObvEngine {
174    fn default() -> Self {
175        Self::new()
176    }
177}
178
179impl Indicator for ObvEngine {
180    fn name(&self) -> &str {
181        "obv"
182    }
183
184    fn warmup_period(&self) -> usize {
185        1
186    }
187
188    fn reset(&mut self) {
189        self.prev_close = None;
190        self.cum_obv = 0.0;
191        self.alerts.clear();
192    }
193
194    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
195        if let Some(prev) = self.prev_close {
196            if bar.close > prev {
197                self.cum_obv += bar.volume;
198            } else if bar.close < prev {
199                self.cum_obv -= bar.volume;
200            }
201        }
202        self.prev_close = Some(bar.close);
203
204        Some(IndicatorOutput::new(self.cum_obv))
205    }
206
207    fn alerts(&self) -> Vec<IndicatorAlert> {
208        self.alerts.clone()
209    }
210}
211
212/// Chaikin Money Flow (CMF) over `period` bars.
213///
214/// Each bar's money-flow multiplier places the close in its range,
215/// `MFM = ((close - low) - (high - close)) / (high - low)` (`0` for a range below `1e-8`), and its
216/// money-flow volume is `MFM * volume`. CMF is the sum of the last `period` money-flow volumes over
217/// the sum of their volumes (`0` without volume), clamped to `-1..=1`.
218///
219/// First output: with the `period`-th bar. [`Indicator::reset`] clears both windows.
220pub struct CmfEngine {
221    period: usize,
222    mf_volumes: VecDeque<f64>,
223    volumes: VecDeque<f64>,
224    alerts: Vec<IndicatorAlert>,
225}
226
227impl CmfEngine {
228    pub fn new(period: usize) -> Self {
229        Self {
230            period,
231            mf_volumes: VecDeque::new(),
232            volumes: VecDeque::new(),
233            alerts: Vec::new(),
234        }
235    }
236}
237
238impl Indicator for CmfEngine {
239    fn name(&self) -> &str {
240        "cmf"
241    }
242
243    fn warmup_period(&self) -> usize {
244        self.period
245    }
246
247    fn reset(&mut self) {
248        self.mf_volumes.clear();
249        self.volumes.clear();
250        self.alerts.clear();
251    }
252
253    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
254        let high_low = bar.high - bar.low;
255        let mfm = if high_low > 1e-8 {
256            ((bar.close - bar.low) - (bar.high - bar.close)) / high_low
257        } else {
258            0.0
259        };
260        let mfv = mfm * bar.volume;
261
262        self.mf_volumes.push_back(mfv);
263        self.volumes.push_back(bar.volume);
264
265        if self.mf_volumes.len() > self.period {
266            self.mf_volumes.pop_front();
267            self.volumes.pop_front();
268        }
269
270        self.alerts.clear();
271        if self.mf_volumes.len() < self.period {
272            return None;
273        }
274
275        let sum_mfv: f64 = self.mf_volumes.iter().sum();
276        let sum_vol: f64 = self.volumes.iter().sum();
277
278        let cmf = if sum_vol > 0.0 {
279            sum_mfv / sum_vol
280        } else {
281            0.0
282        };
283
284        if cmf > 0.20 {
285            self.alerts.push(IndicatorAlert::new(
286                "cmf_bullish",
287                format!("Strong Buying Pressure (CMF: {:.2})", cmf),
288                0.80,
289            ));
290        } else if cmf < -0.20 {
291            self.alerts.push(IndicatorAlert::new(
292                "cmf_bearish",
293                format!("Strong Selling Pressure (CMF: {:.2})", cmf),
294                0.80,
295            ));
296        }
297
298        Some(IndicatorOutput::new(cmf.clamp(-1.0, 1.0)))
299    }
300
301    fn alerts(&self) -> Vec<IndicatorAlert> {
302        self.alerts.clone()
303    }
304}
305
306/// Accumulation/Distribution Line (A/D): a running total of money-flow volume.
307///
308/// Each bar adds `MFM * volume` with the money-flow multiplier of [`CmfEngine`],
309/// `((close - low) - (high - close)) / (high - low)` (`0` for a range below `1e-8`). A close in
310/// the middle of its range therefore adds nothing, however large the volume.
311///
312/// First output: with the first bar. [`Indicator::reset`] returns the total to zero.
313#[derive(Debug, Clone)]
314pub struct AccDistEngine {
315    cum_ad: f64,
316    alerts: Vec<IndicatorAlert>,
317}
318
319impl AccDistEngine {
320    pub fn new() -> Self {
321        Self {
322            cum_ad: 0.0,
323            alerts: Vec::new(),
324        }
325    }
326}
327
328impl Default for AccDistEngine {
329    fn default() -> Self {
330        Self::new()
331    }
332}
333
334impl Indicator for AccDistEngine {
335    fn name(&self) -> &str {
336        "acc_dist"
337    }
338
339    fn warmup_period(&self) -> usize {
340        1
341    }
342
343    fn reset(&mut self) {
344        self.cum_ad = 0.0;
345        self.alerts.clear();
346    }
347
348    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
349        let high_low = bar.high - bar.low;
350        let mfm = if high_low > 1e-8 {
351            ((bar.close - bar.low) - (bar.high - bar.close)) / high_low
352        } else {
353            0.0
354        };
355        let mfv = mfm * bar.volume;
356        self.cum_ad += mfv;
357
358        Some(IndicatorOutput::new(self.cum_ad))
359    }
360
361    fn alerts(&self) -> Vec<IndicatorAlert> {
362        self.alerts.clone()
363    }
364}