Skip to main content

fin_primitives/alternative_data/
mod.rs

1//! Alternative data integration: social sentiment, web traffic, satellite imagery,
2//! credit card data, job postings, and patent filings.
3//! Provides `AltDataAggregator` with composite signals, Pearson correlation, z-score anomaly
4//! detection, and staleness checking.
5//!
6//! Alternative data integration: sentiment signals, web traffic proxies,
7//! satellite imagery, credit card data, job postings, and patent filings.
8//!
9//! ## Key Types
10//!
11//! - [`AltDataSource`] — enum of alternative data source categories
12//! - [`AltDataPoint`] — a single timestamped observation from a source
13//! - [`SentimentSignal`] — aggregated social sentiment for a symbol
14//! - [`AltDataAggregator`] — ingestion, lookup, composite scoring, and anomaly detection
15
16use std::collections::HashMap;
17
18// ─────────────────────────────────────────
19//  AltDataSource
20// ─────────────────────────────────────────
21
22/// Categories of alternative data sources.
23#[derive(Debug, Clone, PartialEq, Eq, Hash)]
24pub enum AltDataSource {
25    /// Social media and news sentiment (Twitter, Reddit, StockTwits, etc.).
26    SocialSentiment,
27    /// Web traffic metrics (SimilarWeb, Alexa rank, search trends).
28    WebTraffic,
29    /// Satellite imagery-derived signals (retail parking lots, oil tank levels, etc.).
30    SatelliteImagery,
31    /// Credit card transaction data aggregates.
32    CreditCardData,
33    /// Job posting counts and hiring velocity.
34    JobPostings,
35    /// Patent application filings.
36    PatentFilings,
37}
38
39impl AltDataSource {
40    /// Typical refresh cadence for this source type (hours).
41    pub fn refresh_frequency_hours(&self) -> u8 {
42        match self {
43            AltDataSource::SocialSentiment => 1,
44            AltDataSource::WebTraffic => 24,
45            AltDataSource::SatelliteImagery => 24,
46            AltDataSource::CreditCardData => 24,
47            AltDataSource::JobPostings => 24,
48            AltDataSource::PatentFilings => 168, // weekly
49        }
50    }
51}
52
53// ─────────────────────────────────────────
54//  AltDataPoint
55// ─────────────────────────────────────────
56
57/// A single timestamped alternative data observation.
58#[derive(Debug, Clone)]
59pub struct AltDataPoint {
60    /// Source category.
61    pub source: AltDataSource,
62    /// Ticker / symbol this data pertains to.
63    pub symbol: String,
64    /// Numeric signal value (interpretation is source-dependent).
65    pub value: f64,
66    /// Confidence in this data point (0.0–1.0).
67    pub confidence: f64,
68    /// Unix timestamp (seconds) when this data was recorded.
69    pub timestamp: u64,
70    /// Arbitrary key-value metadata (e.g. `"provider"`, `"region"`).
71    pub metadata: HashMap<String, String>,
72}
73
74// ─────────────────────────────────────────
75//  SentimentSignal
76// ─────────────────────────────────────────
77
78/// Aggregated social sentiment for a given symbol.
79#[derive(Debug, Clone)]
80pub struct SentimentSignal {
81    /// Ticker symbol.
82    pub symbol: String,
83    /// Composite sentiment score in `[-1.0, 1.0]` (negative = bearish, positive = bullish).
84    pub score: f64,
85    /// Total message / post volume used to compute the score.
86    pub volume: u64,
87    /// Number of distinct sources contributing.
88    pub source_count: u32,
89    /// Fraction of bullish messages (0.0–1.0).
90    pub bullish_pct: f64,
91}
92
93// ─────────────────────────────────────────
94//  AltDataAggregator
95// ─────────────────────────────────────────
96
97/// Ingests and queries alternative data across sources and symbols.
98///
99/// Internally keyed by `(symbol, source)`. Only the most-recent data point per
100/// `(symbol, source)` pair is retained for `latest()` queries; all historical
101/// values are retained for anomaly scoring.
102pub struct AltDataAggregator {
103    /// Latest data point per (symbol, source).
104    latest: HashMap<(String, String), AltDataPoint>,
105    /// Full history per (symbol, source) for statistics.
106    history: HashMap<(String, String), Vec<f64>>,
107}
108
109impl AltDataAggregator {
110    /// Create an empty aggregator.
111    pub fn new() -> Self {
112        Self {
113            latest: HashMap::new(),
114            history: HashMap::new(),
115        }
116    }
117
118    fn key(symbol: &str, source: &AltDataSource) -> (String, String) {
119        (symbol.to_string(), format!("{source:?}"))
120    }
121
122    /// Ingest a new data point, updating the latest observation and history.
123    pub fn ingest(&mut self, point: AltDataPoint) {
124        let k = Self::key(&point.symbol, &point.source);
125        self.history
126            .entry(k.clone())
127            .or_default()
128            .push(point.value);
129        self.latest.insert(k, point);
130    }
131
132    /// Return the most-recent data point for `(symbol, source)`, if any.
133    pub fn latest(&self, symbol: &str, source: &AltDataSource) -> Option<&AltDataPoint> {
134        let k = Self::key(symbol, source);
135        self.latest.get(&k)
136    }
137
138    /// Compute a composite signal for `symbol` as a weighted average of all
139    /// available source values.
140    ///
141    /// `weights` maps source debug names (e.g. `"SocialSentiment"`) to their weight.
142    /// Sources absent from the map receive weight 1.0.
143    /// Returns `0.0` if no data is available for the symbol.
144    pub fn composite_signal(&self, symbol: &str, weights: &HashMap<String, f64>) -> f64 {
145        let mut weighted_sum = 0.0_f64;
146        let mut total_weight = 0.0_f64;
147
148        for ((sym, src_key), point) in &self.latest {
149            if sym != symbol {
150                continue;
151            }
152            let w = weights.get(src_key).copied().unwrap_or(1.0);
153            weighted_sum += point.value * w;
154            total_weight += w;
155        }
156
157        if total_weight == 0.0 {
158            0.0
159        } else {
160            weighted_sum / total_weight
161        }
162    }
163
164    /// Compute the Pearson correlation between the composite signals of two symbols
165    /// across matching `(symbol, source)` time series.
166    ///
167    /// Uses all historical values for each source shared by both symbols.
168    /// Returns `None` if there are fewer than 2 aligned observations.
169    pub fn signal_correlation(&self, symbol_a: &str, symbol_b: &str) -> Option<f64> {
170        // Collect per-source history for both symbols.
171        let mut pairs: Vec<(f64, f64)> = Vec::new();
172
173        // Gather all source keys present for symbol_a.
174        for (sym_src, vals_a) in &self.history {
175            if sym_src.0 != symbol_a {
176                continue;
177            }
178            let src_key = &sym_src.1;
179            let key_b = (symbol_b.to_string(), src_key.clone());
180            if let Some(vals_b) = self.history.get(&key_b) {
181                // Align by taking the shorter length.
182                let n = vals_a.len().min(vals_b.len());
183                for i in 0..n {
184                    pairs.push((vals_a[i], vals_b[i]));
185                }
186            }
187        }
188
189        let n = pairs.len();
190        if n < 2 {
191            return None;
192        }
193
194        let n_f = n as f64;
195        let mean_a = pairs.iter().map(|p| p.0).sum::<f64>() / n_f;
196        let mean_b = pairs.iter().map(|p| p.1).sum::<f64>() / n_f;
197
198        let (cov, var_a, var_b) = pairs.iter().fold((0.0, 0.0, 0.0), |(c, va, vb), &(a, b)| {
199            let da = a - mean_a;
200            let db = b - mean_b;
201            (c + da * db, va + da * da, vb + db * db)
202        });
203
204        let denom = (var_a * var_b).sqrt();
205        if denom == 0.0 {
206            None
207        } else {
208            Some(cov / denom)
209        }
210    }
211
212    /// Compute the z-score of the latest value for `(symbol, source)` relative to
213    /// the historical distribution for that source.
214    ///
215    /// Returns `None` if fewer than 2 historical values exist.
216    pub fn anomaly_score(&self, symbol: &str, source: &AltDataSource) -> Option<f64> {
217        let k = Self::key(symbol, source);
218        let history = self.history.get(&k)?;
219        let latest = self.latest.get(&k)?;
220        let n = history.len();
221        if n < 2 {
222            return None;
223        }
224        let mean = history.iter().sum::<f64>() / n as f64;
225        let variance =
226            history.iter().map(|v| (v - mean).powi(2)).sum::<f64>() / (n - 1) as f64;
227        let std = variance.sqrt();
228        if std == 0.0 {
229            return Some(0.0);
230        }
231        Some((latest.value - mean) / std)
232    }
233
234    /// Return all `(symbol, source)` pairs whose latest observation is older than
235    /// `max_age_hours` relative to `now` (Unix seconds).
236    pub fn freshness_check(
237        &self,
238        max_age_hours: u8,
239        now: u64,
240    ) -> Vec<(String, AltDataSource)> {
241        let max_age_secs = max_age_hours as u64 * 3600;
242        let mut stale = Vec::new();
243
244        for point in self.latest.values() {
245            let age = now.saturating_sub(point.timestamp);
246            if age > max_age_secs {
247                stale.push((point.symbol.clone(), point.source.clone()));
248            }
249        }
250
251        stale
252    }
253}
254
255impl Default for AltDataAggregator {
256    fn default() -> Self {
257        Self::new()
258    }
259}
260
261// ─────────────────────────────────────────
262//  Tests
263// ─────────────────────────────────────────
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268
269    fn make_point(symbol: &str, source: AltDataSource, value: f64, timestamp: u64) -> AltDataPoint {
270        AltDataPoint {
271            source,
272            symbol: symbol.to_string(),
273            value,
274            confidence: 0.9,
275            timestamp,
276            metadata: HashMap::new(),
277        }
278    }
279
280    // ── AltDataSource ──────────────────────────────────────────────────────
281
282    #[test]
283    fn refresh_frequency_social() {
284        assert_eq!(AltDataSource::SocialSentiment.refresh_frequency_hours(), 1);
285    }
286
287    #[test]
288    fn refresh_frequency_patent() {
289        assert_eq!(AltDataSource::PatentFilings.refresh_frequency_hours(), 168);
290    }
291
292    // ── AltDataAggregator::ingest / latest ────────────────────────────────
293
294    #[test]
295    fn latest_none_before_ingest() {
296        let agg = AltDataAggregator::new();
297        assert!(agg.latest("AAPL", &AltDataSource::SocialSentiment).is_none());
298    }
299
300    #[test]
301    fn latest_after_ingest() {
302        let mut agg = AltDataAggregator::new();
303        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.7, 1_000));
304        let p = agg.latest("AAPL", &AltDataSource::SocialSentiment).unwrap();
305        assert!((p.value - 0.7).abs() < 1e-12);
306    }
307
308    #[test]
309    fn latest_updated_on_reingest() {
310        let mut agg = AltDataAggregator::new();
311        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.5, 1_000));
312        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.9, 2_000));
313        let p = agg.latest("AAPL", &AltDataSource::SocialSentiment).unwrap();
314        assert!((p.value - 0.9).abs() < 1e-12);
315    }
316
317    // ── composite_signal ──────────────────────────────────────────────────
318
319    #[test]
320    fn composite_signal_no_data() {
321        let agg = AltDataAggregator::new();
322        let weights = HashMap::new();
323        assert_eq!(agg.composite_signal("AAPL", &weights), 0.0);
324    }
325
326    #[test]
327    fn composite_signal_single_source() {
328        let mut agg = AltDataAggregator::new();
329        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.6, 1_000));
330        let weights = HashMap::new();
331        let cs = agg.composite_signal("AAPL", &weights);
332        assert!((cs - 0.6).abs() < 1e-12);
333    }
334
335    #[test]
336    fn composite_signal_weighted() {
337        let mut agg = AltDataAggregator::new();
338        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.8, 1_000));
339        agg.ingest(make_point("AAPL", AltDataSource::WebTraffic, 0.4, 1_000));
340        let mut weights = HashMap::new();
341        weights.insert("SocialSentiment".to_string(), 2.0);
342        weights.insert("WebTraffic".to_string(), 1.0);
343        let cs = agg.composite_signal("AAPL", &weights);
344        // (0.8*2 + 0.4*1) / 3 = 2.0/3 ≈ 0.6667
345        assert!((cs - 2.0 / 3.0).abs() < 1e-12);
346    }
347
348    // ── signal_correlation ────────────────────────────────────────────────
349
350    #[test]
351    fn signal_correlation_insufficient() {
352        let mut agg = AltDataAggregator::new();
353        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.5, 1_000));
354        assert!(agg.signal_correlation("AAPL", "GOOG").is_none());
355    }
356
357    #[test]
358    fn signal_correlation_perfect_positive() {
359        let mut agg = AltDataAggregator::new();
360        for i in 0..5u64 {
361            agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, i as f64, i));
362            agg.ingest(make_point("GOOG", AltDataSource::SocialSentiment, i as f64 * 2.0, i));
363        }
364        let corr = agg.signal_correlation("AAPL", "GOOG").unwrap();
365        assert!((corr - 1.0).abs() < 1e-10);
366    }
367
368    // ── anomaly_score ─────────────────────────────────────────────────────
369
370    #[test]
371    fn anomaly_score_none_before_ingest() {
372        let agg = AltDataAggregator::new();
373        assert!(agg.anomaly_score("AAPL", &AltDataSource::SocialSentiment).is_none());
374    }
375
376    #[test]
377    fn anomaly_score_computed() {
378        let mut agg = AltDataAggregator::new();
379        // Mean = 1.0, std ~= 1.0 (sample); latest = 3.0 → z ≈ 2.0
380        agg.ingest(make_point("AAPL", AltDataSource::WebTraffic, 0.0, 1));
381        agg.ingest(make_point("AAPL", AltDataSource::WebTraffic, 1.0, 2));
382        agg.ingest(make_point("AAPL", AltDataSource::WebTraffic, 3.0, 3));
383        let z = agg.anomaly_score("AAPL", &AltDataSource::WebTraffic).unwrap();
384        // mean([0,1,3])=4/3, sample std of [0,1,3]
385        let vals = [0.0_f64, 1.0, 3.0];
386        let mean = vals.iter().sum::<f64>() / 3.0;
387        let std = (vals.iter().map(|v| (v - mean).powi(2)).sum::<f64>() / 2.0).sqrt();
388        let expected = (3.0 - mean) / std;
389        assert!((z - expected).abs() < 1e-10);
390    }
391
392    // ── freshness_check ───────────────────────────────────────────────────
393
394    #[test]
395    fn freshness_check_empty() {
396        let agg = AltDataAggregator::new();
397        assert!(agg.freshness_check(24, 100_000).is_empty());
398    }
399
400    #[test]
401    fn freshness_check_fresh() {
402        let mut agg = AltDataAggregator::new();
403        let now = 10_000_u64;
404        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.5, now - 60));
405        assert!(agg.freshness_check(1, now).is_empty());
406    }
407
408    #[test]
409    fn freshness_check_stale() {
410        let mut agg = AltDataAggregator::new();
411        let now = 100_000_u64;
412        // Timestamp 3600 seconds ago but max_age is 0 hours → immediately stale
413        agg.ingest(make_point("AAPL", AltDataSource::SocialSentiment, 0.5, now - 3_601));
414        let stale = agg.freshness_check(1, now);
415        assert_eq!(stale.len(), 1);
416        assert_eq!(stale[0].0, "AAPL");
417    }
418}