1use std::collections::BTreeMap;
5
6#[cfg(feature = "serde")]
7use serde::Serialize;
8
9#[derive(Debug, Clone, Copy)]
11pub struct CloseSample {
12 pub timestamp: i64,
13 pub close: f64,
14}
15
16#[derive(Debug, Clone, PartialEq)]
17#[cfg_attr(feature = "serde", derive(Serialize))]
18pub struct CorrelationCell {
19 pub a: String,
20 pub b: String,
21 pub corr: f64,
24}
25
26#[derive(Debug, Clone, PartialEq)]
27#[cfg_attr(feature = "serde", derive(Serialize))]
28pub struct RelativeStrength {
29 pub instrument: String,
30 pub change_pct: f64,
32}
33
34pub fn correlation_matrix(
45 series: &[(String, Vec<CloseSample>)],
46 lookback: usize,
47) -> Vec<CorrelationCell> {
48 let closes: Vec<(&str, BTreeMap<i64, f64>)> = series
50 .iter()
51 .map(|(epic, bars)| {
52 (
53 epic.as_str(),
54 bars.iter()
55 .map(|b| (b.timestamp, b.close))
56 .collect::<BTreeMap<_, _>>(),
57 )
58 })
59 .collect();
60
61 let mut out = Vec::new();
62 for i in 0..closes.len() {
63 for j in (i + 1)..closes.len() {
64 let (epic_a, map_a) = &closes[i];
65 let (epic_b, map_b) = &closes[j];
66 if let Some(corr) = pair_correlation(map_a, map_b, lookback) {
67 out.push(CorrelationCell {
68 a: epic_a.to_string(),
69 b: epic_b.to_string(),
70 corr,
71 });
72 }
73 }
74 }
75 out
76}
77
78pub fn relative_strength(
84 series: &[(String, Vec<CloseSample>)],
85 lookback: usize,
86) -> Vec<RelativeStrength> {
87 let mut out: Vec<RelativeStrength> = series
88 .iter()
89 .filter_map(|(epic, bars)| {
90 if bars.len() < lookback + 1 || lookback == 0 {
91 return None;
92 }
93 let last = bars[bars.len() - 1].close;
94 let base = bars[bars.len() - 1 - lookback].close;
95 if base == 0.0 {
96 return None;
97 }
98 Some(RelativeStrength {
99 instrument: epic.clone(),
100 change_pct: 100.0 * (last - base) / base,
101 })
102 })
103 .collect();
104 out.sort_by(|a, b| b.change_pct.total_cmp(&a.change_pct));
105 out
106}
107
108fn pair_correlation(
111 a: &BTreeMap<i64, f64>,
112 b: &BTreeMap<i64, f64>,
113 lookback: usize,
114) -> Option<f64> {
115 let common: Vec<(f64, f64)> = a
117 .iter()
118 .filter_map(|(ts, ca)| b.get(ts).map(|cb| (*ca, *cb)))
119 .collect();
120 if common.len() < 4 {
121 return None; }
123 let mut ra = Vec::with_capacity(common.len() - 1);
125 let mut rb = Vec::with_capacity(common.len() - 1);
126 for pair in common.windows(2) {
127 let (a0, b0) = pair[0];
128 let (a1, b1) = pair[1];
129 if a0 == 0.0 || b0 == 0.0 {
130 continue;
131 }
132 ra.push((a1 - a0) / a0);
133 rb.push((b1 - b0) / b0);
134 }
135 if ra.len() < 3 {
136 return None;
137 }
138 if ra.len() > lookback && lookback >= 3 {
139 let start = ra.len() - lookback;
140 ra.drain(0..start);
141 rb.drain(0..start);
142 }
143 pearson(&ra, &rb)
144}
145
146fn pearson(x: &[f64], y: &[f64]) -> Option<f64> {
147 let n = x.len() as f64;
148 let mean_x = x.iter().sum::<f64>() / n;
149 let mean_y = y.iter().sum::<f64>() / n;
150 let mut cov = 0.0;
151 let mut var_x = 0.0;
152 let mut var_y = 0.0;
153 for (xi, yi) in x.iter().zip(y) {
154 let dx = xi - mean_x;
155 let dy = yi - mean_y;
156 cov += dx * dy;
157 var_x += dx * dx;
158 var_y += dy * dy;
159 }
160 let denom = (var_x * var_y).sqrt();
161 if denom <= 0.0 {
162 return None; }
164 Some((cov / denom).clamp(-1.0, 1.0))
165}
166
167#[derive(Debug, Clone, PartialEq)]
169#[cfg_attr(feature = "serde", derive(Serialize))]
170pub struct RollingBetaResult {
171 pub beta: f64,
173 pub alpha: f64,
175 pub r_squared: f64,
177 pub correlation: f64,
179 pub periods: usize,
181}
182
183pub fn compute_rolling_beta(
187 asset_returns: &[f64],
188 benchmark_returns: &[f64],
189) -> Option<RollingBetaResult> {
190 if asset_returns.len() != benchmark_returns.len() || asset_returns.len() < 3 {
191 return None;
192 }
193
194 let n = asset_returns.len() as f64;
195 let mean_a = asset_returns.iter().sum::<f64>() / n;
196 let mean_b = benchmark_returns.iter().sum::<f64>() / n;
197
198 let mut cov = 0.0f64;
199 let mut var_b = 0.0f64;
200 let mut var_a = 0.0f64;
201
202 for (&ra, &rb) in asset_returns.iter().zip(benchmark_returns) {
203 if !ra.is_finite() || !rb.is_finite() {
204 return None;
205 }
206 let da = ra - mean_a;
207 let db = rb - mean_b;
208 cov += da * db;
209 var_a += da * da;
210 var_b += db * db;
211 }
212
213 if var_b <= 1e-14 {
214 return None;
216 }
217
218 let beta = cov / var_b;
219 let alpha = mean_a - beta * mean_b;
220
221 let denom = (var_a * var_b).sqrt();
222 let correlation = if denom > 1e-14 {
223 (cov / denom).clamp(-1.0, 1.0)
224 } else {
225 0.0
226 };
227 let r_squared = (correlation * correlation).clamp(0.0, 1.0);
228
229 Some(RollingBetaResult {
230 beta,
231 alpha,
232 r_squared,
233 correlation,
234 periods: asset_returns.len(),
235 })
236}
237
238#[derive(Debug, Clone, PartialEq)]
240#[cfg_attr(feature = "serde", derive(Serialize))]
241pub struct UniverseMemberObservation {
242 pub symbol: String,
243 pub period_return: f64,
245 pub current_price: f64,
247 pub ma_reference: Option<f64>,
249}
250
251#[derive(Debug, Clone, PartialEq)]
253#[cfg_attr(feature = "serde", derive(Serialize))]
254pub struct MarketBreadthSnapshot {
255 pub total_active: usize,
257 pub advancing: usize,
259 pub declining: usize,
261 pub unchanged: usize,
263 pub advance_decline_ratio: f64,
265 pub net_advancing_pct: f64,
267 pub pct_above_ma: Option<f64>,
269}
270
271pub fn compute_market_breadth(
275 members: &[UniverseMemberObservation],
276) -> Option<MarketBreadthSnapshot> {
277 if members.is_empty() {
278 return None;
279 }
280
281 let mut advancing = 0usize;
282 let mut declining = 0usize;
283 let mut unchanged = 0usize;
284 let mut above_ma_count = 0usize;
285 let mut total_with_ma = 0usize;
286
287 for m in members {
288 if !m.period_return.is_finite() || !m.current_price.is_finite() {
289 continue;
290 }
291 if m.period_return > 1e-9 {
292 advancing += 1;
293 } else if m.period_return < -1e-9 {
294 declining += 1;
295 } else {
296 unchanged += 1;
297 }
298
299 if let Some(ma) = m.ma_reference {
300 if ma.is_finite() && ma > 0.0 {
301 total_with_ma += 1;
302 if m.current_price > ma {
303 above_ma_count += 1;
304 }
305 }
306 }
307 }
308
309 let total_active = advancing + declining + unchanged;
310 if total_active == 0 {
311 return None;
312 }
313
314 let advance_decline_ratio = advancing as f64 / (declining.max(1) as f64);
315 let net_advancing_pct = (advancing as f64 - declining as f64) / total_active as f64;
316 let pct_above_ma = if total_with_ma > 0 {
317 Some(above_ma_count as f64 / total_with_ma as f64)
318 } else {
319 None
320 };
321
322 Some(MarketBreadthSnapshot {
323 total_active,
324 advancing,
325 declining,
326 unchanged,
327 advance_decline_ratio,
328 net_advancing_pct,
329 pct_above_ma,
330 })
331}
332
333#[derive(Debug, Clone, PartialEq)]
335#[cfg_attr(feature = "serde", derive(Serialize))]
336pub struct PairSpreadResult {
337 pub hedge_ratio: f64,
339 pub current_spread: f64,
341 pub mean_spread: f64,
343 pub std_spread: f64,
345 pub residual_z_score: f64,
347}
348
349pub fn compute_pair_spread(prices_a: &[f64], prices_b: &[f64]) -> Option<PairSpreadResult> {
351 if prices_a.len() != prices_b.len() || prices_a.len() < 3 {
352 return None;
353 }
354
355 let n = prices_a.len() as f64;
356 let mean_a = prices_a.iter().sum::<f64>() / n;
357 let mean_b = prices_b.iter().sum::<f64>() / n;
358
359 let mut cov = 0.0f64;
360 let mut var_b = 0.0f64;
361
362 for (&pa, &pb) in prices_a.iter().zip(prices_b) {
363 if !pa.is_finite() || !pb.is_finite() {
364 return None;
365 }
366 let da = pa - mean_a;
367 let db = pb - mean_b;
368 cov += da * db;
369 var_b += db * db;
370 }
371
372 if var_b <= 1e-14 {
373 return None;
374 }
375
376 let hedge_ratio = cov / var_b;
377
378 let mut spread_series = Vec::with_capacity(prices_a.len());
380 let mut sum_spread = 0.0f64;
381 for (&pa, &pb) in prices_a.iter().zip(prices_b) {
382 let s = pa - hedge_ratio * pb;
383 spread_series.push(s);
384 sum_spread += s;
385 }
386
387 let mean_spread = sum_spread / n;
388 let mut var_spread = 0.0f64;
389 for &s in &spread_series {
390 var_spread += (s - mean_spread).powi(2);
391 }
392 let std_spread = (var_spread / n).sqrt();
393
394 let current_spread = *spread_series.last()?;
395 let residual_z_score = if std_spread > 1e-14 {
396 (current_spread - mean_spread) / std_spread
397 } else {
398 0.0
399 };
400
401 Some(PairSpreadResult {
402 hedge_ratio,
403 current_spread,
404 mean_spread,
405 std_spread,
406 residual_z_score,
407 })
408}
409
410#[derive(Debug, Clone, PartialEq)]
412#[cfg_attr(feature = "serde", derive(Serialize))]
413pub struct SignalCorrelationCell {
414 pub signal_a: String,
415 pub signal_b: String,
416 pub correlation: f64,
418 pub is_redundant: bool,
420}
421
422pub fn compute_signal_correlation_matrix(
424 signals: &[(String, Vec<f64>)],
425 redundancy_threshold: f64,
426) -> Vec<SignalCorrelationCell> {
427 let mut out = Vec::new();
428 let num_signals = signals.len();
429
430 for i in 0..num_signals {
431 for j in (i + 1)..num_signals {
432 let (name_a, scores_a) = &signals[i];
433 let (name_b, scores_b) = &signals[j];
434
435 let min_len = scores_a.len().min(scores_b.len());
436 if min_len < 3 {
437 continue;
438 }
439
440 if let Some(corr) = pearson(&scores_a[..min_len], &scores_b[..min_len]) {
441 out.push(SignalCorrelationCell {
442 signal_a: name_a.clone(),
443 signal_b: name_b.clone(),
444 correlation: corr,
445 is_redundant: corr.abs() >= redundancy_threshold,
446 });
447 }
448 }
449 }
450
451 out
452}