Skip to main content

kestrel_chartkit/
cross_asset.rs

1//! Pairwise return correlation and relative-strength ranking over close samples.
2//! Inputs must be finite, chronological, and use a common timestamp convention.
3
4use std::collections::BTreeMap;
5
6#[cfg(feature = "serde")]
7use serde::Serialize;
8
9/// Timestamp is an opaque alignment key. All series must use the same unit and bar boundary.
10#[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    /// Pearson correlation of the two instruments' aligned percentage returns,
22    /// in −1..1.
23    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    /// Percentage change of close over the lookback window.
31    pub change_pct: f64,
32}
33
34/// Pairwise return correlation across `series` (each `(epic, bars)`, bars
35/// oldest first). Pairs are aligned on their common bar timestamps — FX /
36/// commodity instruments share the same session grid but can differ at the
37/// edges (gaps, differing seed depth), so aligning on `ts` avoids correlating
38/// misaligned rows. `lookback` caps how many of the most recent aligned
39/// returns are used. A pair with fewer than 3 common returns is omitted.
40/// Returns are `(c_t - c_{t-1}) / c_{t-1}` over consecutive common timestamps, a pair of closes
41/// with a zero base skipped; `corr` is the Pearson correlation of the last `lookback` of them,
42/// clamped to `-1..=1`, and a pair with a flat return series is omitted.
43/// For compatibility, lookbacks below 3 leave the full aligned history in use.
44pub fn correlation_matrix(
45    series: &[(String, Vec<CloseSample>)],
46    lookback: usize,
47) -> Vec<CorrelationCell> {
48    // Pre-index each instrument's closes by timestamp once.
49    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
78/// Ranks instruments by percentage change of close over the last `lookback`
79/// bars (each `(epic, bars)`, bars oldest first), strongest first:
80/// `change_pct = 100 · (close_last - close_base) / close_base` with the base `lookback` bars
81/// before the last. Instruments without enough bars, or with a zero base close, are dropped;
82/// a `lookback` of 0 drops all.
83pub 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
108/// Pearson correlation of percentage returns over the common timestamps of
109/// two close series, using at most the last `lookback` returns.
110fn pair_correlation(
111    a: &BTreeMap<i64, f64>,
112    b: &BTreeMap<i64, f64>,
113    lookback: usize,
114) -> Option<f64> {
115    // Common timestamps, ascending, with both closes.
116    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; // < 3 returns
122    }
123    // Percentage returns from consecutive common closes.
124    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; // a flat series has no defined correlation
163    }
164    Some((cov / denom).clamp(-1.0, 1.0))
165}
166
167/// Evaluated rolling beta, alpha, and regression fit of an asset against a benchmark series.
168#[derive(Debug, Clone, PartialEq)]
169#[cfg_attr(feature = "serde", derive(Serialize))]
170pub struct RollingBetaResult {
171    /// Covariance(asset, benchmark) / Variance(benchmark).
172    pub beta: f64,
173    /// Mean(asset) - beta * Mean(benchmark).
174    pub alpha: f64,
175    /// Coefficient of determination (R^2), in 0.0..=1.0.
176    pub r_squared: f64,
177    /// Correlation between asset and benchmark returns.
178    pub correlation: f64,
179    /// Number of aligned return periods evaluated.
180    pub periods: usize,
181}
182
183/// Computes rolling beta, alpha, and R^2 of `asset_returns` against `benchmark_returns`.
184///
185/// Both slices must have the same length and at least 3 observations.
186pub 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        // Benchmark is flat; beta is undefined
215        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/// A snapshot observation of a constituent member within a market universe at an as-of timestamp.
239#[derive(Debug, Clone, PartialEq)]
240#[cfg_attr(feature = "serde", derive(Serialize))]
241pub struct UniverseMemberObservation {
242    pub symbol: String,
243    /// Return over the relevant period (e.g. 1-day return).
244    pub period_return: f64,
245    /// Current price level.
246    pub current_price: f64,
247    /// Moving average reference level (e.g. SMA 20 or SMA 50), if available.
248    pub ma_reference: Option<f64>,
249}
250
251/// Cross-sectional market breadth summary computed across active universe members.
252#[derive(Debug, Clone, PartialEq)]
253#[cfg_attr(feature = "serde", derive(Serialize))]
254pub struct MarketBreadthSnapshot {
255    /// Total number of active constituents with valid data at this as-of timestamp.
256    pub total_active: usize,
257    /// Number of advancing constituents (period_return > 0).
258    pub advancing: usize,
259    /// Number of declining constituents (period_return < 0).
260    pub declining: usize,
261    /// Number of unchanged constituents (period_return == 0).
262    pub unchanged: usize,
263    /// Ratio of advancing to declining constituents: `advancing / max(1, declining)`.
264    pub advance_decline_ratio: f64,
265    /// Net advance percentage: `(advancing - declining) / total_active`.
266    pub net_advancing_pct: f64,
267    /// Percentage of constituents trading above their moving average reference level, in 0.0..=1.0.
268    pub pct_above_ma: Option<f64>,
269}
270
271/// Computes cross-sectional market breadth across active universe constituents.
272///
273/// Missing or inactive constituents are excluded, correctly adjusting the denominator.
274pub 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/// Evaluated descriptive pair spread, hedge ratio, and residual z-score.
334#[derive(Debug, Clone, PartialEq)]
335#[cfg_attr(feature = "serde", derive(Serialize))]
336pub struct PairSpreadResult {
337    /// OLS hedge ratio $\gamma$ (Asset A price vs. Asset B price).
338    pub hedge_ratio: f64,
339    /// Latest spread level: $P_A - \gamma \cdot P_B$.
340    pub current_spread: f64,
341    /// Historical mean of the spread over the window.
342    pub mean_spread: f64,
343    /// Historical standard deviation of the spread.
344    pub std_spread: f64,
345    /// Current normalized residual z-score: `(current_spread - mean_spread) / std_spread`.
346    pub residual_z_score: f64,
347}
348
349/// Computes descriptive pair spread, OLS hedge ratio, and current residual z-score between two price series.
350pub 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    // Compute spread series: S = P_a - hedge_ratio * P_b
379    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/// Pairwise correlation between two indicator/signal time series to detect collinearity and double counting.
411#[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    /// Pearson correlation of the two signal score series, in -1.0..=1.0.
417    pub correlation: f64,
418    /// True if absolute correlation exceeds the redundancy threshold (e.g. >= 0.85).
419    pub is_redundant: bool,
420}
421
422/// Computes pairwise signal score correlations across multiple named signal output streams.
423pub 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}