kestrel_chartkit/indicator/
anchored_vwap.rs1use 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#[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#[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 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}