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)]
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 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}