Skip to main content

kestrel_chartkit/indicator/
anchored_vwap.rs

1use super::{Indicator, IndicatorAlert, IndicatorOutput};
2use crate::model::Bar;
3use crate::session::{SessionConfig, SessionConfigError, SessionTracker};
4use crate::timeframe::Timeframe;
5use std::collections::HashMap;
6
7#[cfg(feature = "serde")]
8use serde::{Deserialize, Serialize};
9
10/// Anchor trigger condition for resetting cumulative VWAP calculations.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
12#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
13pub enum VwapAnchorKind {
14    #[default]
15    Session,
16    Day,
17    Week,
18    Month,
19    ManualTimestamp(i64),
20    External,
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
24#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
25pub enum ZeroVolumePolicy {
26    Skip,
27    #[default]
28    EqualWeight,
29}
30
31/// Anchored VWAP: a cumulative volume-weighted average price since the last anchor, with
32/// volume-weighted standard-deviation bands.
33///
34/// From the anchor on, `VWAP = sum(typical * volume) / sum(volume)` with
35/// `typical = (high + low + close) / 3`, and
36/// `stddev = sqrt(sum(typical^2 * volume) / sum(volume) - VWAP^2)`, floored at zero; the bands sit
37/// at `mult1` and `mult2` standard deviations. A bar without volume counts with weight `1` under
38/// [`ZeroVolumePolicy::EqualWeight`] (the default) and is skipped under
39/// [`ZeroVolumePolicy::Skip`].
40///
41/// The anchor is chosen by [`VwapAnchorKind`]. The default, `Session`, restarts at the open of
42/// each session of the configured [`SessionConfig`]; the default config has the same start and end
43/// (`00:00`), which makes the session the whole day — every bar is in session and the VWAP restarts
44/// at 00:00 UTC. Bars outside a configured session produce no output. `Day`, `Week` and `Month`
45/// restart at the calendar boundary (shifted by the UTC offset), `ManualTimestamp` starts with the
46/// first bar at or after its timestamp, `External` on the caller's anchor events.
47///
48/// `value` and `extra["vwap"]`: the VWAP; `extra["stddev"]` and the four band edges. First output:
49/// with the first bar inside an active anchor. [`Indicator::reset`] clears the accumulation.
50#[derive(Debug, Clone)]
51pub struct AnchoredVwapEngine {
52    anchor_kind: VwapAnchorKind,
53    cum_pv: f64,
54    cum_vol: f64,
55    cum_pv2: f64,
56    prev_timestamp: Option<i64>,
57    utc_offset_seconds: i32,
58    active: bool,
59    zero_volume_policy: ZeroVolumePolicy,
60    session_tracker: Option<SessionTracker>,
61    stddev_mult1: f64,
62    stddev_mult2: f64,
63}
64
65impl AnchoredVwapEngine {
66    pub fn new(anchor_kind: VwapAnchorKind, stddev_mult1: f64, stddev_mult2: f64) -> Self {
67        Self {
68            anchor_kind,
69            cum_pv: 0.0,
70            cum_vol: 0.0,
71            cum_pv2: 0.0,
72            prev_timestamp: None,
73            utc_offset_seconds: 0,
74            active: !matches!(anchor_kind, VwapAnchorKind::ManualTimestamp(_)),
75            zero_volume_policy: ZeroVolumePolicy::EqualWeight,
76            session_tracker: matches!(anchor_kind, VwapAnchorKind::Session).then(|| {
77                SessionTracker::new(SessionConfig::default()).expect("default session is valid")
78            }),
79            stddev_mult1,
80            stddev_mult2,
81        }
82    }
83
84    pub fn with_defaults() -> Self {
85        Self::new(VwapAnchorKind::Session, 1.0, 2.0)
86    }
87
88    pub fn with_utc_offset(mut self, utc_offset_seconds: i32) -> Self {
89        self.utc_offset_seconds = utc_offset_seconds;
90        self
91    }
92
93    pub fn with_zero_volume_policy(mut self, policy: ZeroVolumePolicy) -> Self {
94        self.zero_volume_policy = policy;
95        self
96    }
97
98    pub fn with_session_config(
99        mut self,
100        config: SessionConfig,
101    ) -> Result<Self, SessionConfigError> {
102        self.anchor_kind = VwapAnchorKind::Session;
103        self.session_tracker = Some(SessionTracker::new(config)?);
104        self.active = true;
105        Ok(self)
106    }
107
108    fn check_anchor_reset(&self, current_ts: i64) -> bool {
109        let prev_ts = match self.prev_timestamp {
110            Some(ts) => ts,
111            None => return false,
112        };
113
114        let timeframe = match self.anchor_kind {
115            VwapAnchorKind::Session => return false,
116            VwapAnchorKind::Day => Timeframe::Day(1),
117            VwapAnchorKind::Week => Timeframe::Week(1),
118            VwapAnchorKind::Month => Timeframe::Month(1),
119            VwapAnchorKind::ManualTimestamp(_) | VwapAnchorKind::External => return false,
120        };
121        timeframe.bucket_start(prev_ts, self.utc_offset_seconds)
122            != timeframe.bucket_start(current_ts, self.utc_offset_seconds)
123    }
124
125    /// Processes a bar and optionally resets an externally/pivot-anchored VWAP.
126    pub fn on_bar_with_anchor(&mut self, bar: &Bar, anchor_event: bool) -> Option<IndicatorOutput> {
127        let session_reset = if let Some(tracker) = &mut self.session_tracker {
128            tracker.on_bar(bar);
129            if !tracker.in_session() {
130                self.prev_timestamp = Some(bar.timestamp);
131                return None;
132            }
133            tracker.is_new_session()
134        } else {
135            false
136        };
137        let manual_activated = match self.anchor_kind {
138            VwapAnchorKind::ManualTimestamp(timestamp) => {
139                !self.active && bar.timestamp >= timestamp
140            }
141            VwapAnchorKind::External => anchor_event,
142            _ => false,
143        };
144        if session_reset
145            || manual_activated
146            || anchor_event
147            || self.check_anchor_reset(bar.timestamp)
148        {
149            self.cum_pv = 0.0;
150            self.cum_vol = 0.0;
151            self.cum_pv2 = 0.0;
152            self.active = true;
153        }
154        self.prev_timestamp = Some(bar.timestamp);
155        if !self.active {
156            return None;
157        }
158
159        let volume = if bar.volume > 0.0 {
160            bar.volume
161        } else if self.zero_volume_policy == ZeroVolumePolicy::EqualWeight {
162            1.0
163        } else {
164            return None;
165        };
166        let price = bar.typical_price();
167        self.cum_pv += price * volume;
168        self.cum_vol += volume;
169        self.cum_pv2 += price * price * volume;
170
171        let vwap = self.cum_pv / self.cum_vol;
172        let variance = (self.cum_pv2 / self.cum_vol - vwap * vwap).max(0.0);
173        let stddev = variance.sqrt();
174        let mut extra = HashMap::new();
175        extra.insert("vwap".to_string(), vwap);
176        extra.insert("stddev".to_string(), stddev);
177        extra.insert("band1_upper".to_string(), vwap + self.stddev_mult1 * stddev);
178        extra.insert("band1_lower".to_string(), vwap - self.stddev_mult1 * stddev);
179        extra.insert("band2_upper".to_string(), vwap + self.stddev_mult2 * stddev);
180        extra.insert("band2_lower".to_string(), vwap - self.stddev_mult2 * stddev);
181        Some(IndicatorOutput::with_extra(vwap, extra))
182    }
183}
184
185impl Indicator for AnchoredVwapEngine {
186    fn name(&self) -> &str {
187        "anchored_vwap"
188    }
189
190    fn warmup_period(&self) -> usize {
191        1
192    }
193
194    fn reset(&mut self) {
195        self.cum_pv = 0.0;
196        self.cum_vol = 0.0;
197        self.cum_pv2 = 0.0;
198        self.prev_timestamp = None;
199        self.active = !matches!(self.anchor_kind, VwapAnchorKind::ManualTimestamp(_));
200        if let Some(tracker) = &mut self.session_tracker {
201            tracker.reset();
202        }
203    }
204
205    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
206        self.on_bar_with_anchor(bar, false)
207    }
208
209    fn alerts(&self) -> Vec<IndicatorAlert> {
210        Vec::new()
211    }
212}
213
214#[cfg(test)]
215mod tests {
216    use super::*;
217
218    #[test]
219    fn test_anchored_vwap_reset() {
220        let mut avwap = AnchoredVwapEngine::with_defaults();
221        let bar1 = Bar::new(0, 100.0, 105.0, 95.0, 100.0, 1000.0);
222        let bar2 = Bar::new(86400, 200.0, 205.0, 195.0, 200.0, 1000.0);
223
224        let out1 = avwap.on_bar(&bar1).unwrap();
225        assert_eq!(out1.value, 100.0);
226
227        let out2 = avwap.on_bar(&bar2).unwrap();
228        assert_eq!(out2.value, 200.0);
229    }
230
231    #[test]
232    fn manual_anchor_emits_only_from_anchor_timestamp() {
233        let mut avwap = AnchoredVwapEngine::new(VwapAnchorKind::ManualTimestamp(10), 1.0, 2.0);
234        assert!(avwap
235            .on_bar(&Bar::new(9, 100.0, 100.0, 100.0, 100.0, 1.0))
236            .is_none());
237        assert_eq!(
238            avwap
239                .on_bar(&Bar::new(10, 110.0, 110.0, 110.0, 110.0, 1.0))
240                .unwrap()
241                .value,
242            110.0
243        );
244    }
245
246    #[test]
247    fn external_anchor_resets_accumulation() {
248        let mut avwap = AnchoredVwapEngine::new(VwapAnchorKind::External, 1.0, 2.0);
249        avwap.on_bar(&Bar::new(1, 100.0, 100.0, 100.0, 100.0, 1.0));
250        let output = avwap
251            .on_bar_with_anchor(&Bar::new(2, 200.0, 200.0, 200.0, 200.0, 1.0), true)
252            .unwrap();
253        assert_eq!(output.value, 200.0);
254    }
255
256    #[test]
257    fn session_anchor_observes_configured_session() {
258        let config = SessionConfig {
259            start_hour: 9,
260            start_minute: 0,
261            end_hour: 10,
262            end_minute: 0,
263            orb_duration_mins: 30,
264            utc_offset_seconds: 0,
265        };
266        let mut avwap = AnchoredVwapEngine::with_defaults()
267            .with_session_config(config)
268            .unwrap();
269        assert!(avwap
270            .on_bar(&Bar::new(8 * 3_600, 100.0, 100.0, 100.0, 100.0, 1.0))
271            .is_none());
272        assert_eq!(
273            avwap
274                .on_bar(&Bar::new(9 * 3_600, 110.0, 110.0, 110.0, 110.0, 1.0))
275                .unwrap()
276                .value,
277            110.0
278        );
279        assert!(avwap
280            .on_bar(&Bar::new(10 * 3_600, 120.0, 120.0, 120.0, 120.0, 1.0))
281            .is_none());
282    }
283}