diff --git a/crates/fidc-core/examples/market-event-kernel.rs b/crates/fidc-core/examples/market-event-kernel.rs new file mode 100644 index 0000000..b1392fc --- /dev/null +++ b/crates/fidc-core/examples/market-event-kernel.rs @@ -0,0 +1,9 @@ +use std::io::{self, Read}; +fn main() { + let mut input=String::new();io::stdin().read_to_string(&mut input).unwrap(); + let request=serde_json::from_str(&input).unwrap(); + match fidc_core::market_event_context::aggregate(request) { + Ok(value)=>println!("{}",serde_json::to_string(&value).unwrap()), + Err(error)=>{eprintln!("{error}");std::process::exit(1);} + } +} diff --git a/crates/fidc-core/src/daily_patterns.rs b/crates/fidc-core/src/daily_patterns.rs index 9478bcc..13c8b72 100644 --- a/crates/fidc-core/src/daily_patterns.rs +++ b/crates/fidc-core/src/daily_patterns.rs @@ -737,8 +737,10 @@ pub fn evaluate_research_batch_with_policy( spec: PatternSpec, days: &[NaiveDate], series: &[PatternSeries], context: &ResearchContext, numeric_output: bool, isolate_data_errors: bool, ) -> Result { - let common_fields = ["index_open", "index_high", "index_low", "index_close"]; - let symbol_fields = ["scope_rank", "scope_percentile", "scope_size"]; + let common_fields = ["index_open", "index_high", "index_low", "index_close"].into_iter() + .chain(crate::market_event_context::COMMON_FIELDS.iter().copied()).collect::>(); + let symbol_fields = ["scope_rank", "scope_percentile", "scope_size"].into_iter() + .chain(crate::market_event_context::INDUSTRY_FIELDS.iter().copied()).collect::>(); if spec.template != "expression" || series.is_empty() || series.len() > 200 || series.iter().map(|s| &s.symbol).collect::>().len() != series.len() || context.common.keys().any(|k| !common_fields.contains(&k.as_str())) @@ -747,7 +749,9 @@ pub fn evaluate_research_batch_with_policy( return Err("research_context_scope_or_fields_invalid".into()); } for (name, values) in &context.common { - if values.len() != days.len() || values.iter().any(|v| !v.is_some_and(|x| x.is_finite() && x > 0.0)) { + if values.len() != days.len() || values.iter().any(|v| if name.starts_with("index_") { + !v.is_some_and(|x| x.is_finite() && x > 0.0) + } else { v.is_some_and(|x| !x.is_finite()) }) { return Err(format!("research_index_window_incomplete: {name}")); } } diff --git a/crates/fidc-core/src/factor_events.rs b/crates/fidc-core/src/factor_events.rs index f4c32ba..3539649 100644 --- a/crates/fidc-core/src/factor_events.rs +++ b/crates/fidc-core/src/factor_events.rs @@ -179,6 +179,10 @@ pub fn catalog() -> Value { json!({"contract":CONTRACT,"library":{"name":"TA-Lib native Rust","revision":TA_REV,"license":"BSD-3-Clause"}, "execution_context_contract":crate::pattern_context::CONTRACT, "execution_context_fields":crate::pattern_context::CONTEXT_FIELDS, + "market_event_context_contract":crate::market_event_context::CONTRACT, + "market_event_kernel_sha256":crate::market_event_context::implementation_sha256(), + "market_event_common_fields":crate::market_event_context::COMMON_FIELDS, + "market_event_industry_fields":crate::market_event_context::INDUSTRY_FIELDS, "session_events":crate::session_events::EVENTS,"session_event_contract":crate::session_events::CONTRACT, "indicators":indicators,"operators":OPERATORS,"cross_section_operators":crate::factor_cross_section::OPERATORS,"read_only":true,"live_routing":false, "policies":{"null":"unknown_not_false","warmup":"null_until_full_history","recursive_seed":"frozen_input_start", diff --git a/crates/fidc-core/src/lib.rs b/crates/fidc-core/src/lib.rs index 37592f6..5c8aacb 100644 --- a/crates/fidc-core/src/lib.rs +++ b/crates/fidc-core/src/lib.rs @@ -7,6 +7,7 @@ pub mod pattern_context; pub mod session_events; pub mod factor_events; pub mod factor_cross_section; +pub mod market_event_context; pub mod engine; pub mod event_bus; pub mod events; diff --git a/crates/fidc-core/src/market_event_context.rs b/crates/fidc-core/src/market_event_context.rs new file mode 100644 index 0000000..e6a0d3d --- /dev/null +++ b/crates/fidc-core/src/market_event_context.rs @@ -0,0 +1,239 @@ +//! Complete published daily cross sections, independent of trading candidates and accounts. +use chrono::NaiveDate; +use serde::{Deserialize, Serialize}; +use std::collections::{BTreeMap, BTreeSet}; + +pub const CONTRACT: &str = "fidc_market_event_context_v1"; +pub fn implementation_sha256() -> String { + use sha2::{Digest, Sha256}; + format!("{:x}", Sha256::digest(include_bytes!("market_event_context.rs"))) +} +pub const COMMON_FIELDS: &[&str] = &[ + "market_breadth", "market_return", "market_limit_up_count", "market_limit_down_count", + "market_limit_up_rate", "market_broken_limit_rate", "market_high_board", "market_profit_effect", +]; +pub const INDUSTRY_FIELDS: &[&str] = &[ + "industry_close", "industry_return_20", "industry_breadth", "industry_rank", "industry_size", +]; + +#[derive(Clone, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Observation { + pub symbol: String, + pub industry: Option, + pub close: Option, + pub high: Option, + pub previous_close: Option, + pub upper_limit: Option, + pub lower_limit: Option, + pub no_limit: Option, + pub paused: Option, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Day { + pub date: NaiveDate, + pub universe: Vec, + pub rows: Vec, +} + +#[derive(Default, Clone, Deserialize, Serialize)] +#[serde(default, deny_unknown_fields)] +pub struct State { + pub last_date: Option, + pub streaks: BTreeMap>, + pub limit_ups: BTreeSet, + pub industry_history: BTreeMap>, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Request { + pub days: Vec, + #[serde(default)] + pub previous: State, +} + +#[derive(Serialize)] +pub struct OutputDay { + pub date: NaiveDate, + pub common: BTreeMap>, + pub industries: BTreeMap>>, + pub members: BTreeMap>, + pub securities: usize, + pub active: usize, + pub paused: usize, + pub no_limit: usize, + pub profit_effect_members: Vec, + pub profit_effect_missing: Vec, +} + +#[derive(Serialize)] +pub struct Output { + pub contract: &'static str, + pub days: Vec, + pub state: State, +} + +fn positive(value: Option, symbol: &str, field: &str) -> Result { + value.filter(|v| v.is_finite() && *v > 0.0) + .ok_or_else(|| format!("market_event_input_invalid: {symbol} {field}")) +} +fn average(values: impl Iterator, n: usize) -> f64 { + values.map(|v| v / n as f64).sum() +} + +pub fn aggregate(request: Request) -> Result { + let mut state = request.previous; + if request.days.is_empty() || request.days.len() > 30 + || request.days.iter().map(|d| d.rows.len()).sum::() > 60_000 + || state.streaks.len() > 20_000 || state.limit_ups.len() > 20_000 + || state.industry_history.len() > 2000 + || state.industry_history.values().any(|v| v.is_empty() || v.len() > 21 + || v.iter().any(|x| !x.is_finite() || *x <= 0.0)) + || state.last_date.is_none() && (!state.streaks.is_empty() || !state.limit_ups.is_empty() || !state.industry_history.is_empty()) { + return Err("market_event_history_budget_or_state_invalid".into()); + } + let mut output = Vec::new(); + for day in request.days { + if state.last_date.is_some_and(|d| d >= day.date) + || day.universe.is_empty() || day.universe.len() > 20_000 + || day.universe.iter().collect::>().len() != day.universe.len() + || day.rows.len() != day.universe.len() + || day.rows.iter().map(|r| &r.symbol).collect::>() != day.universe.iter().collect::>() { + return Err(format!("market_event_incomplete_cross_section: {}", day.date)); + } + let mut returns = BTreeMap::new(); + let mut groups: BTreeMap> = BTreeMap::new(); + let mut members = BTreeMap::new(); + let mut streaks = BTreeMap::new(); + let mut ups = BTreeSet::new(); + let mut downs = 0; let mut touched = 0; let mut broken = 0; let mut paused = 0; let mut unlimited = 0; + for row in &day.rows { + let industry = row.industry.clone().filter(|s| !s.trim().is_empty()); + members.insert(row.symbol.clone(), industry.clone()); + match row.paused { + Some(true) => { + paused += 1; + streaks.insert(row.symbol.clone(), state.streaks.get(&row.symbol).copied().flatten()); + continue; + }, + Some(false) => {}, + None => return Err(format!("market_event_pause_state_missing: {} {}", day.date, row.symbol)), + } + let c = positive(row.close, &row.symbol, "close")?; + let h = positive(row.high, &row.symbol, "high")?; + let p = positive(row.previous_close, &row.symbol, "previous_close")?; + if h + 1e-8 < c { return Err(format!("market_event_high_below_close: {}", row.symbol)); } + let change = c / p - 1.0; + returns.insert(row.symbol.clone(), change); + if let Some(industry) = industry { groups.entry(industry).or_default().push(change); } + let is_up = match row.no_limit { + Some(true) => { unlimited += 1; false }, + Some(false) => { + let upper = positive(row.upper_limit, &row.symbol, "upper_limit")?; + let lower = positive(row.lower_limit, &row.symbol, "lower_limit")?; + if lower >= upper || c > upper + 1e-8 || c < lower - 1e-8 { + return Err(format!("market_event_limit_bounds_invalid: {} {}", day.date, row.symbol)); + } + let at_up = (c - upper).abs() <= 1e-8; + if (c - lower).abs() <= 1e-8 { downs += 1; } + if h >= upper - 1e-8 { touched += 1; if !at_up { broken += 1; } } + at_up + }, + None => return Err(format!("market_event_limit_policy_missing: {}", row.symbol)), + }; + if is_up { + ups.insert(row.symbol.clone()); + // The first observed limit-up may already be a continuing streak. + streaks.insert(row.symbol.clone(), state.streaks.get(&row.symbol).copied().flatten().map(|v| v + 1)); + } else { streaks.insert(row.symbol.clone(), Some(0)); } + } + let active = returns.len(); + if active == 0 { return Err(format!("market_event_no_active_market: {}", day.date)); } + let previous_ups = state.limit_ups.iter().cloned().collect::>(); + let profit_missing = previous_ups.iter().filter(|s| !returns.contains_key(*s)).cloned().collect::>(); + let profit = if previous_ups.is_empty() || !profit_missing.is_empty() { None } + else { Some(average(previous_ups.iter().map(|s| returns[s]), previous_ups.len())) }; + let board = if ups.iter().any(|s| streaks[s].is_none()) { None } + else { Some(ups.iter().map(|s| streaks[s].unwrap()).max().unwrap_or(0) as f64) }; + let common = BTreeMap::from([ + ("market_breadth".into(), Some(returns.values().filter(|v| **v > 0.0).count() as f64 / active as f64)), + ("market_return".into(), Some(average(returns.values().copied(), active))), + ("market_limit_up_count".into(), Some(ups.len() as f64)), + ("market_limit_down_count".into(), Some(downs as f64)), + ("market_limit_up_rate".into(), (active > unlimited).then(|| ups.len() as f64 / (active - unlimited) as f64)), + ("market_broken_limit_rate".into(), (touched > 0).then(|| broken as f64 / touched as f64)), + ("market_high_board".into(), board), + ("market_profit_effect".into(), profit), + ]); + let mut industries = BTreeMap::new(); + // A disappeared group breaks its continuous history; no stale NAV is carried forward. + state.industry_history.retain(|key, _| groups.contains_key(key)); + for (industry, values) in groups { + let history = state.industry_history.entry(industry.clone()).or_default(); + let nav = history.last().copied().unwrap_or(1.0) * (1.0 + average(values.iter().copied(), values.len())); + history.push(nav); + if history.len() > 21 { history.remove(0); } + let momentum = (history.len() == 21).then(|| nav / history[0] - 1.0); + industries.insert(industry, BTreeMap::from([ + ("industry_close".into(), Some(nav)), ("industry_return_20".into(), momentum), + ("industry_breadth".into(), Some(values.iter().filter(|v| **v > 0.0).count() as f64 / values.len() as f64)), + ])); + } + let universe = industries.keys().cloned().collect::>(); + let known = industries.values().all(|g| g["industry_return_20"].is_some()); + let ranks = if known && !universe.is_empty() { + crate::factor_cross_section::evaluate("RANK", &universe, &industries.iter().map(|(s,g)| + crate::factor_cross_section::Observation {symbol:s.clone(), value:g["industry_return_20"].unwrap(),industry:None,market_cap:None}).collect::>(),0.0)? + .into_iter().map(|r|(r.symbol,r.value)).collect::>() + } else { BTreeMap::new() }; + for (name, fields) in &mut industries { + fields.insert("industry_rank".into(), ranks.get(name).copied()); + fields.insert("industry_size".into(), Some(universe.len() as f64)); + } + output.push(OutputDay { date:day.date, common, industries, members, securities:day.rows.len(), active, paused, + no_limit:unlimited, profit_effect_members:previous_ups, profit_effect_missing:profit_missing }); + state.last_date = Some(day.date); state.streaks = streaks; state.limit_ups = ups; + } + Ok(Output {contract:CONTRACT, days:output, state}) +} + +#[cfg(test)] +mod tests { + use super::*; + fn day(n: u32, up: bool) -> Day { + Day {date:NaiveDate::from_ymd_opt(2026,9,n).unwrap(), universe:vec!["A".into(),"B".into()], rows:vec![ + Observation{symbol:"A".into(),industry:Some("I".into()),close:Some(if up {11.0}else{10.0}),high:Some(11.0),previous_close:Some(10.0),upper_limit:Some(11.0),lower_limit:Some(9.0),no_limit:Some(false),paused:Some(false)}, + Observation{symbol:"B".into(),industry:Some("J".into()),close:Some(9.0),high:Some(10.0),previous_close:Some(10.0),upper_limit:Some(11.0),lower_limit:Some(9.0),no_limit:Some(false),paused:Some(false)}]} + } + #[test] + fn formulas_use_real_limits_and_full_denominators() { + let r=aggregate(Request{days:vec![day(1,false),day(2,true),day(3,true)],previous:State::default()}).unwrap(); + let d=&r.days[1]; + assert_eq!(d.common["market_breadth"],Some(0.5)); + assert_eq!(d.common["market_limit_down_count"],Some(1.0)); + assert_eq!(r.days[0].common["market_broken_limit_rate"],Some(1.0)); + assert_eq!(r.days[2].common["market_high_board"],Some(2.0)); + assert!((r.days[2].common["market_profit_effect"].unwrap()-0.1).abs()<1e-12); + assert_eq!(r.days[0].common["market_profit_effect"],None); + } + #[test] + fn missing_duplicate_and_unproven_limit_states_fail() { + let mut d=day(1,true);d.rows.pop();assert!(aggregate(Request{days:vec![d],previous:State::default()}).is_err()); + let mut d=day(1,true);d.rows[0].upper_limit=None;assert!(aggregate(Request{days:vec![d],previous:State::default()}).is_err()); + let mut d=day(1,true);d.rows[0].no_limit=Some(true);d.rows[0].upper_limit=None; + assert_eq!(aggregate(Request{days:vec![d],previous:State::default()}).unwrap().days[0].no_limit,1); + } + #[test] + fn chunking_and_future_append_preserve_history() { + let first=aggregate(Request{days:vec![day(1,false),day(2,true)],previous:State::default()}).unwrap(); + let next=aggregate(Request{days:vec![day(3,true)],previous:first.state}).unwrap(); + let full=aggregate(Request{days:vec![day(1,false),day(2,true),day(3,true)],previous:State::default()}).unwrap(); + assert_eq!(serde_json::to_value(&first.days).unwrap(),serde_json::to_value(&full.days[..2]).unwrap()); + assert_eq!(serde_json::to_value(&next.days).unwrap(),serde_json::to_value(&full.days[2..]).unwrap()); + let unknown=aggregate(Request{days:vec![day(1,true)],previous:State::default()}).unwrap(); + assert_eq!(unknown.days[0].common["market_high_board"],None); + } +}