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 Engine with volume-weighted stddev bands.
32#[derive(Debug, Clone)]
33pub struct AnchoredVwapEngine {
34    anchor_kind: VwapAnchorKind,
35    cum_pv: f64,
36    cum_vol: f64,
37    cum_pv2: f64,
38    prev_timestamp: Option<i64>,
39    utc_offset_seconds: i32,
40    active: bool,
41    zero_volume_policy: ZeroVolumePolicy,
42    session_tracker: Option<SessionTracker>,
43    stddev_mult1: f64,
44    stddev_mult2: f64,
45}
46
47impl AnchoredVwapEngine {
48    pub fn new(anchor_kind: VwapAnchorKind, stddev_mult1: f64, stddev_mult2: f64) -> Self {
49        Self {
50            anchor_kind,
51            cum_pv: 0.0,
52            cum_vol: 0.0,
53            cum_pv2: 0.0,
54            prev_timestamp: None,
55            utc_offset_seconds: 0,
56            active: !matches!(anchor_kind, VwapAnchorKind::ManualTimestamp(_)),
57            zero_volume_policy: ZeroVolumePolicy::EqualWeight,
58            session_tracker: matches!(anchor_kind, VwapAnchorKind::Session).then(|| {
59                SessionTracker::new(SessionConfig::default()).expect("default session is valid")
60            }),
61            stddev_mult1,
62            stddev_mult2,
63        }
64    }
65
66    pub fn with_defaults() -> Self {
67        Self::new(VwapAnchorKind::Session, 1.0, 2.0)
68    }
69
70    pub fn with_utc_offset(mut self, utc_offset_seconds: i32) -> Self {
71        self.utc_offset_seconds = utc_offset_seconds;
72        self
73    }
74
75    pub fn with_zero_volume_policy(mut self, policy: ZeroVolumePolicy) -> Self {
76        self.zero_volume_policy = policy;
77        self
78    }
79
80    pub fn with_session_config(
81        mut self,
82        config: SessionConfig,
83    ) -> Result<Self, SessionConfigError> {
84        self.anchor_kind = VwapAnchorKind::Session;
85        self.session_tracker = Some(SessionTracker::new(config)?);
86        self.active = true;
87        Ok(self)
88    }
89
90    fn check_anchor_reset(&self, current_ts: i64) -> bool {
91        let prev_ts = match self.prev_timestamp {
92            Some(ts) => ts,
93            None => return false,
94        };
95
96        let timeframe = match self.anchor_kind {
97            VwapAnchorKind::Session => return false,
98            VwapAnchorKind::Day => Timeframe::Day(1),
99            VwapAnchorKind::Week => Timeframe::Week(1),
100            VwapAnchorKind::Month => Timeframe::Month(1),
101            VwapAnchorKind::ManualTimestamp(_) | VwapAnchorKind::External => return false,
102        };
103        timeframe.bucket_start(prev_ts, self.utc_offset_seconds)
104            != timeframe.bucket_start(current_ts, self.utc_offset_seconds)
105    }
106
107    /// Processes a bar and optionally resets an externally/pivot-anchored VWAP.
108    pub fn on_bar_with_anchor(&mut self, bar: &Bar, anchor_event: bool) -> Option<IndicatorOutput> {
109        let session_reset = if let Some(tracker) = &mut self.session_tracker {
110            tracker.on_bar(bar);
111            if !tracker.in_session() {
112                self.prev_timestamp = Some(bar.timestamp);
113                return None;
114            }
115            tracker.is_new_session()
116        } else {
117            false
118        };
119        let manual_activated = match self.anchor_kind {
120            VwapAnchorKind::ManualTimestamp(timestamp) => {
121                !self.active && bar.timestamp >= timestamp
122            }
123            VwapAnchorKind::External => anchor_event,
124            _ => false,
125        };
126        if session_reset
127            || manual_activated
128            || anchor_event
129            || self.check_anchor_reset(bar.timestamp)
130        {
131            self.cum_pv = 0.0;
132            self.cum_vol = 0.0;
133            self.cum_pv2 = 0.0;
134            self.active = true;
135        }
136        self.prev_timestamp = Some(bar.timestamp);
137        if !self.active {
138            return None;
139        }
140
141        let volume = if bar.volume > 0.0 {
142            bar.volume
143        } else if self.zero_volume_policy == ZeroVolumePolicy::EqualWeight {
144            1.0
145        } else {
146            return None;
147        };
148        let price = bar.typical_price();
149        self.cum_pv += price * volume;
150        self.cum_vol += volume;
151        self.cum_pv2 += price * price * volume;
152
153        let vwap = self.cum_pv / self.cum_vol;
154        let variance = (self.cum_pv2 / self.cum_vol - vwap * vwap).max(0.0);
155        let stddev = variance.sqrt();
156        let mut extra = HashMap::new();
157        extra.insert("vwap".to_string(), vwap);
158        extra.insert("stddev".to_string(), stddev);
159        extra.insert("band1_upper".to_string(), vwap + self.stddev_mult1 * stddev);
160        extra.insert("band1_lower".to_string(), vwap - self.stddev_mult1 * stddev);
161        extra.insert("band2_upper".to_string(), vwap + self.stddev_mult2 * stddev);
162        extra.insert("band2_lower".to_string(), vwap - self.stddev_mult2 * stddev);
163        Some(IndicatorOutput::with_extra(vwap, extra))
164    }
165}
166
167impl Indicator for AnchoredVwapEngine {
168    fn name(&self) -> &str {
169        "anchored_vwap"
170    }
171
172    fn warmup_period(&self) -> usize {
173        1
174    }
175
176    fn reset(&mut self) {
177        self.cum_pv = 0.0;
178        self.cum_vol = 0.0;
179        self.cum_pv2 = 0.0;
180        self.prev_timestamp = None;
181        self.active = !matches!(self.anchor_kind, VwapAnchorKind::ManualTimestamp(_));
182        if let Some(tracker) = &mut self.session_tracker {
183            tracker.reset();
184        }
185    }
186
187    fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
188        self.on_bar_with_anchor(bar, false)
189    }
190
191    fn alerts(&self) -> Vec<IndicatorAlert> {
192        Vec::new()
193    }
194}
195
196#[cfg(test)]
197mod tests {
198    use super::*;
199
200    #[test]
201    fn test_anchored_vwap_reset() {
202        let mut avwap = AnchoredVwapEngine::with_defaults();
203        let bar1 = Bar::new(0, 100.0, 105.0, 95.0, 100.0, 1000.0);
204        let bar2 = Bar::new(86400, 200.0, 205.0, 195.0, 200.0, 1000.0);
205
206        let out1 = avwap.on_bar(&bar1).unwrap();
207        assert_eq!(out1.value, 100.0);
208
209        let out2 = avwap.on_bar(&bar2).unwrap();
210        assert_eq!(out2.value, 200.0);
211    }
212
213    #[test]
214    fn manual_anchor_emits_only_from_anchor_timestamp() {
215        let mut avwap = AnchoredVwapEngine::new(VwapAnchorKind::ManualTimestamp(10), 1.0, 2.0);
216        assert!(avwap
217            .on_bar(&Bar::new(9, 100.0, 100.0, 100.0, 100.0, 1.0))
218            .is_none());
219        assert_eq!(
220            avwap
221                .on_bar(&Bar::new(10, 110.0, 110.0, 110.0, 110.0, 1.0))
222                .unwrap()
223                .value,
224            110.0
225        );
226    }
227
228    #[test]
229    fn external_anchor_resets_accumulation() {
230        let mut avwap = AnchoredVwapEngine::new(VwapAnchorKind::External, 1.0, 2.0);
231        avwap.on_bar(&Bar::new(1, 100.0, 100.0, 100.0, 100.0, 1.0));
232        let output = avwap
233            .on_bar_with_anchor(&Bar::new(2, 200.0, 200.0, 200.0, 200.0, 1.0), true)
234            .unwrap();
235        assert_eq!(output.value, 200.0);
236    }
237
238    #[test]
239    fn session_anchor_observes_configured_session() {
240        let config = SessionConfig {
241            start_hour: 9,
242            start_minute: 0,
243            end_hour: 10,
244            end_minute: 0,
245            orb_duration_mins: 30,
246            utc_offset_seconds: 0,
247        };
248        let mut avwap = AnchoredVwapEngine::with_defaults()
249            .with_session_config(config)
250            .unwrap();
251        assert!(avwap
252            .on_bar(&Bar::new(8 * 3_600, 100.0, 100.0, 100.0, 100.0, 1.0))
253            .is_none());
254        assert_eq!(
255            avwap
256                .on_bar(&Bar::new(9 * 3_600, 110.0, 110.0, 110.0, 110.0, 1.0))
257                .unwrap()
258                .value,
259            110.0
260        );
261        assert!(avwap
262            .on_bar(&Bar::new(10 * 3_600, 120.0, 120.0, 120.0, 120.0, 1.0))
263            .is_none());
264    }
265}