From 6b0cdbcecc47e45aac9d296dfd09b0b4e8d74b8f Mon Sep 17 00:00:00 2001 From: boris Date: Wed, 9 Sep 2026 10:31:25 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E5=85=B1=E4=BA=AB=E5=9B=A0?= =?UTF-8?q?=E5=AD=90=E4=BA=8B=E4=BB=B6=E8=A1=A8=E8=BE=BE=E5=BC=8F=E4=B8=8E?= =?UTF-8?q?=E5=AE=8C=E6=95=B4=E6=88=AA=E9=9D=A2=E7=AE=97=E5=AD=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Cargo.lock | 14 + crates/fidc-core/Cargo.toml | 1 + .../fidc-core/examples/factor-event-kernel.rs | 35 + crates/fidc-core/src/daily_patterns.rs | 133 ++- crates/fidc-core/src/factor_cross_section.rs | 190 +++ crates/fidc-core/src/factor_events.rs | 1047 +++++++++++++++++ crates/fidc-core/src/lib.rs | 2 + 7 files changed, 1421 insertions(+), 1 deletion(-) create mode 100644 crates/fidc-core/examples/factor-event-kernel.rs create mode 100644 crates/fidc-core/src/factor_cross_section.rs create mode 100644 crates/fidc-core/src/factor_events.rs diff --git a/Cargo.lock b/Cargo.lock index ef0b58b..49082cc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -153,6 +153,7 @@ dependencies = [ "rhai", "serde", "serde_json", + "ta-lib", "thiserror", ] @@ -477,6 +478,19 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "ta-lib" +version = "0.8.1" +source = "git+https://github.com/TA-Lib/ta-lib.git?rev=dd5a90259a3f9e04e2da9f38bf0719a841b40108#dd5a90259a3f9e04e2da9f38bf0719a841b40108" +dependencies = [ + "ta-lib-dispatch", +] + +[[package]] +name = "ta-lib-dispatch" +version = "0.1.2" +source = "git+https://github.com/TA-Lib/ta-lib.git?rev=dd5a90259a3f9e04e2da9f38bf0719a841b40108#dd5a90259a3f9e04e2da9f38bf0719a841b40108" + [[package]] name = "thin-vec" version = "0.2.16" diff --git a/crates/fidc-core/Cargo.toml b/crates/fidc-core/Cargo.toml index 0fcfb41..b1e1365 100644 --- a/crates/fidc-core/Cargo.toml +++ b/crates/fidc-core/Cargo.toml @@ -14,3 +14,4 @@ rhai.workspace = true serde.workspace = true serde_json.workspace = true thiserror.workspace = true +ta-lib = { git = "https://github.com/TA-Lib/ta-lib.git", rev = "dd5a90259a3f9e04e2da9f38bf0719a841b40108" } diff --git a/crates/fidc-core/examples/factor-event-kernel.rs b/crates/fidc-core/examples/factor-event-kernel.rs new file mode 100644 index 0000000..ad82788 --- /dev/null +++ b/crates/fidc-core/examples/factor-event-kernel.rs @@ -0,0 +1,35 @@ +use fidc_core::factor_events::{self, Expr, Frame}; +use serde::Deserialize; +use serde_json::{Value, json}; +use std::io::{self, Read}; + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Request { + expressions: std::collections::BTreeMap, + frame: Frame, +} + +fn main() -> Result<(), Box> { + let mut input = String::new(); + io::stdin().read_to_string(&mut input)?; + let output = if input.trim().is_empty() { + factor_events::catalog() + } else { + let request: Request = serde_json::from_str(&input)?; + let results = request + .expressions + .iter() + .map(|(id, expr)| { + let result = match factor_events::evaluate(expr, &request.frame) { + Ok(v) => json!({"result":v}), + Err(e) => json!({"error":e}), + }; + (id.clone(), result) + }) + .collect::>(); + json!({"contract":factor_events::CONTRACT,"results":results,"read_only":true}) + }; + println!("{}", serde_json::to_string(&output)?); + Ok(()) +} diff --git a/crates/fidc-core/src/daily_patterns.rs b/crates/fidc-core/src/daily_patterns.rs index 0530152..a49529f 100644 --- a/crates/fidc-core/src/daily_patterns.rs +++ b/crates/fidc-core/src/daily_patterns.rs @@ -1,6 +1,6 @@ //! Completed-session OHLCV rules shared by research and strategy execution. use crate::DataSet; -use chrono::NaiveDate; +use chrono::{FixedOffset, NaiveDate, TimeZone}; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; use std::collections::{BTreeMap, BTreeSet}; @@ -9,6 +9,7 @@ pub const CONTRACT: &str = "fidc_daily_ohlcv_pattern_v1"; pub fn catalog() -> Value { json!({"contract":CONTRACT,"templates":{ + "expression":{"label":"指标与事件条件","parameters":{"history_window":[300,2,3000]},"stages":["selection","buy","sell","position_management"],"method":"冻结历史窗口与表达式;预热不足或未定义值不产生信号。复用共享指标事件内核,不修改既有任务。"}, "strength":{"label":"趋势强势","parameters":{"momentum_window":[25,5,120],"fast_window":[20,2,60],"slow_window":[60,20,252]},"stages":["selection","buy"],"method":"收盘价>短均线>长均线,按区间动量排序;不是当日金叉。"}, "breakout":{"label":"前高突破","parameters":{"high_window":[60,5,252],"volume_window":[10,2,60],"volume_multiple":[1.3,1,10],"max_upper_shadow":[0.1,0,1]},"stages":["selection","buy"],"method":"收盘突破此前N日最高价,量达到此前M日均量倍数,上影比例受限;参考窗口不含当日。"}, "volume_spike":{"label":"放量上涨","parameters":{"volume_window":[5,2,60],"volume_multiple":[3.0,1,10]},"stages":["selection","buy"],"method":"当日上涨且量达到此前N日最大量的指定倍数;不等同价格创新高。"}, @@ -24,9 +25,14 @@ pub struct PatternSpec { pub template: String, #[serde(default)] pub parameters: BTreeMap, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub expression: Option, } impl PatternSpec { pub fn validate(mut self) -> Result { + if (self.template == "expression") != self.expression.is_some() { + return Err("expression_template_requires_expression_only".into()); + } let catalog = catalog(); let definition = catalog["templates"] .get(&self.template) @@ -66,6 +72,7 @@ impl PatternSpec { } pub fn history_len(&self) -> usize { match self.template.as_str() { + "expression" => self.n("history_window"), "strength" => self.n("slow_window").max(self.n("momentum_window") + 1), "breakout" => self.n("high_window").max(self.n("volume_window")) + 1, "volume_spike" | "volume_down" => self.n("volume_window") + 1, @@ -246,6 +253,80 @@ pub fn evaluate( result.anchor = json!({"date":days[len-1],"raw_close":by_day[&days[len-1]].close,"factor":by_day[&days[len-1]].adjustment_factor_backward1}); let mut score = None; match spec.template.as_str() { + "expression" => { + let zone = FixedOffset::east_opt(8 * 3600).unwrap(); + let timestamps = days + .iter() + .map(|d| { + zone.from_local_datetime(&d.and_hms_opt(16, 0, 0).unwrap()) + .single() + .unwrap() + }) + .collect::>(); + let anchor = by_day[&days[len - 1]].adjustment_factor_backward1.unwrap(); + let mut fields = BTreeMap::from([ + ( + "open".into(), + prices.iter().map(|b| Some(b.0 / anchor)).collect(), + ), + ( + "high".into(), + prices.iter().map(|b| Some(b.1 / anchor)).collect(), + ), + ( + "low".into(), + prices.iter().map(|b| Some(b.2 / anchor)).collect(), + ), + ( + "close".into(), + prices.iter().map(|b| Some(b.3 / anchor)).collect(), + ), + ("volume".into(), prices.iter().map(|b| Some(b.4)).collect()), + ]); + for (name, index) in [ + ("raw_open", 0), + ("raw_high", 1), + ("raw_low", 2), + ("raw_close", 3), + ] { + fields.insert( + name.into(), + days.iter() + .map(|d| { + let b = by_day[d]; + [b.open, b.high, b.low, b.close][index] + }) + .collect(), + ); + } + let frame = crate::factor_events::Frame { + symbol: series.symbol.clone(), + frequency: "1d".into(), + decision_at: *timestamps.last().unwrap(), + available_at: timestamps.clone(), + timestamps, + fields, + }; + let values = crate::factor_events::evaluate(spec.expression.as_ref().unwrap(), &frame)?; + let latest = values.values.last().copied().flatten(); + result.values["expression"] = json!(values); + result.values["expression_contract"] = json!(crate::factor_events::CONTRACT); + result.values["price_policy"] = json!("backward1_anchored_to_decision_close"); + result.score = latest; + if latest.is_none() { + result.exclusion = Some( + json!({"reason":"expression_undefined_or_warmup","signal_date":days.last()}), + ); + } else if values.value_type == crate::factor_events::ValueType::Boolean { + result.matched = latest == Some(1.0); + result + .checks + .push(json!({"label":"组合条件","passed":result.matched})); + } else { + return Err("expression_signal_requires_boolean: 数值因子必须显式比较或组合,不能自动视为买卖信号".into()); + } + return Ok(result); + } "strength" => { let fast = mean(prices[len - spec.n("fast_window")..].iter().map(|b| b.3))?; let slow = mean(prices[len - spec.n("slow_window")..].iter().map(|b| b.3))?; @@ -491,10 +572,60 @@ pub fn expression_specs(expression: &str) -> Result, String> { #[cfg(test)] mod tests { use super::*; + #[test] + fn expression_condition_preserves_native_types_and_rejects_numeric_as_signal() { + let make = |expression: Value| { + serde_json::from_value::(json!({"template":"expression","parameters":{"history_window":3},"expression":expression})).unwrap().validate().unwrap() + }; + let spec = make( + json!({"kind":"operator","name":"GT","args":[{"kind":"field","name":"close"},{"kind":"indicator","name":"SMA","inputs":[{"kind":"field","name":"close"}],"parameters":{"optInTimePeriod":2}}]}), + ); + let days = ["2026-09-04", "2026-09-07", "2026-09-08"] + .map(|d| NaiveDate::parse_from_str(d, "%Y-%m-%d").unwrap()); + let series = PatternSeries { + symbol: "TEST".into(), + name: None, + listed_at: None, + bars: days + .iter() + .enumerate() + .map(|(i, &date)| { + let p = 10.0 + i as f64; + PatternBar { + date, + open: Some(p), + high: Some(p), + low: Some(p), + close: Some(p), + volume: Some(100.0), + adjustment_factor_backward1: Some(1.0), + paused: Some(false), + source_path: None, + } + }) + .collect(), + }; + let result = evaluate(&spec, &days, &series).unwrap(); + assert!(result.matched); + assert_eq!(result.score, Some(1.0)); + assert!( + evaluate( + &make(json!({"kind":"field","name":"close"})), + &days, + &series + ) + .unwrap_err() + .contains("requires_boolean") + ); + let mut missing = series.clone(); + missing.bars[1].close = None; + assert!(evaluate(&spec, &days, &missing).is_err()); + } fn fixture(template: &str) -> (PatternSpec, Vec, PatternSeries) { let spec = PatternSpec { template: template.into(), parameters: BTreeMap::new(), + expression: None, } .validate() .unwrap(); diff --git a/crates/fidc-core/src/factor_cross_section.rs b/crates/fidc-core/src/factor_cross_section.rs new file mode 100644 index 0000000..87f4038 --- /dev/null +++ b/crates/fidc-core/src/factor_cross_section.rs @@ -0,0 +1,190 @@ +//! Cross-sectional operators require an explicit complete universe, never a UI page. +use serde::{Deserialize, Serialize}; +use std::collections::{BTreeMap, BTreeSet}; + +pub const OPERATORS: &[&str] = &[ + "RANK", + "PERCENTILE", + "TOP", + "BOTTOM", + "TOP_PERCENT", + "BOTTOM_PERCENT", + "WINSORIZE", + "INDUSTRY_NEUTRALIZE", + "SIZE_NEUTRALIZE", +]; + +#[derive(Clone, Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Observation { + pub symbol: String, + pub value: f64, + pub industry: Option, + pub market_cap: Option, +} +#[derive(Debug, Serialize)] +pub struct Output { + pub symbol: String, + pub value: f64, +} + +fn mean(values: &[f64]) -> f64 { + let base = values[0]; + base + values + .iter() + .skip(1) + .map(|v| (v - base) / values.len() as f64) + .sum::() +} +fn quantile(sorted: &[f64], p: f64) -> f64 { + let x = p * (sorted.len() - 1) as f64; + let l = x.floor() as usize; + let r = x.ceil() as usize; + sorted[l] + (sorted[r] - sorted[l]) * (x - l as f64) +} + +pub fn evaluate( + name: &str, + universe: &[String], + rows: &[Observation], + threshold: f64, +) -> Result, String> { + let expected = universe.iter().collect::>(); + if rows.is_empty() + || rows.len() > 20_000 + || expected.len() != universe.len() + || rows.len() != universe.len() + || rows.iter().map(|r| &r.symbol).collect::>() != expected + || rows.iter().any(|r| !r.value.is_finite()) + { + return Err("cross_section_incomplete_or_invalid_universe".into()); + } + if !OPERATORS.contains(&name) || !threshold.is_finite() { + return Err("cross_section_operator_invalid".into()); + } + if matches!(name, "TOP" | "BOTTOM") && (threshold < 1.0 || threshold.fract() != 0.0) + || matches!(name, "TOP_PERCENT" | "BOTTOM_PERCENT") && !(0.0..=1.0).contains(&threshold) + || name == "WINSORIZE" && !(0.0..0.5).contains(&threshold) + { + return Err("cross_section_threshold_invalid".into()); + } + let mut sorted = rows.iter().map(|r| r.value).collect::>(); + sorted.sort_by(f64::total_cmp); + let mut industry_values: BTreeMap<&str, Vec> = BTreeMap::new(); + if name == "INDUSTRY_NEUTRALIZE" { + for row in rows { + let industry = row + .industry + .as_deref() + .filter(|v| !v.trim().is_empty()) + .ok_or("cross_section_pit_industry_missing")?; + industry_values.entry(industry).or_default().push(row.value); + } + } + let size = if name == "SIZE_NEUTRALIZE" { + let x = rows + .iter() + .map(|r| { + r.market_cap + .filter(|v| v.is_finite() && *v > 0.0) + .map(f64::ln) + .ok_or("cross_section_market_cap_missing") + }) + .collect::, _>>()?; + let xm = mean(&x); + let ym = mean(&sorted); + let variance = x.iter().map(|v| (v - xm).powi(2)).sum::(); + if variance == 0.0 || rows.len() < 3 { + return Err("cross_section_size_regression_unidentified".into()); + } + let beta = x + .iter() + .zip(rows) + .map(|(x, y)| (x - xm) * (y.value - ym)) + .sum::() + / variance; + Some((x, xm, ym, beta)) + } else { + None + }; + rows.iter() + .enumerate() + .map(|(index, row)| { + let low = sorted.partition_point(|v| *v < row.value); + let high = sorted.partition_point(|v| *v <= row.value); + let rank = (low + 1 + high) as f64 / 2.0; + let descending = (rows.len() + 1) as f64 - rank; + let percentile = if rows.len() == 1 { + 0.5 + } else { + (rank - 1.0) / (rows.len() - 1) as f64 + }; + let value = match name { + "RANK" => descending, + "PERCENTILE" => percentile, + "TOP" => f64::from(descending <= threshold), + "BOTTOM" => f64::from(rank <= threshold), + "TOP_PERCENT" => f64::from(descending <= threshold * rows.len() as f64), + "BOTTOM_PERCENT" => f64::from(rank <= threshold * rows.len() as f64), + "WINSORIZE" => row.value.clamp( + quantile(&sorted, threshold), + quantile(&sorted, 1.0 - threshold), + ), + "INDUSTRY_NEUTRALIZE" => { + row.value - mean(&industry_values[row.industry.as_deref().unwrap()]) + } + "SIZE_NEUTRALIZE" => { + let (x, xm, ym, beta) = size.as_ref().unwrap(); + row.value - (ym + beta * (x[index] - xm)) + } + _ => unreachable!(), + }; + if !value.is_finite() { + return Err("cross_section_result_nonfinite".into()); + } + Ok(Output { + symbol: row.symbol.clone(), + value, + }) + }) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + fn rows() -> Vec { + [1.0, 3.0, 3.0, 4.0] + .iter() + .enumerate() + .map(|(i, &value)| Observation { + symbol: format!("S{i}"), + value, + industry: Some(if i < 2 { "A" } else { "B" }.into()), + market_cap: Some(10.0 + i as f64), + }) + .collect() + } + #[test] + fn ties_keep_equal_rank_and_missing_universe_rejects() { + let r = rows(); + let u = r.iter().map(|r| r.symbol.clone()).collect::>(); + let out = evaluate("RANK", &u, &r, 0.0).unwrap(); + assert_eq!( + out.iter().map(|r| r.value).collect::>(), + vec![4.0, 2.5, 2.5, 1.0] + ); + assert!(evaluate("RANK", &u, &r[..3], 0.0).is_err()); + } + #[test] + fn neutralization_preserves_input_order() { + let r = rows(); + let u = r.iter().map(|r| r.symbol.clone()).collect::>(); + let out = evaluate("INDUSTRY_NEUTRALIZE", &u, &r, 0.0).unwrap(); + assert_eq!( + out.iter().map(|r| r.value).collect::>(), + vec![-1.0, 1.0, -0.5, 0.5] + ); + assert!(evaluate("TOP_PERCENT", &u, &r, 20.0).is_err()); + } +} diff --git a/crates/fidc-core/src/factor_events.rs b/crates/fidc-core/src/factor_events.rs new file mode 100644 index 0000000..7330693 --- /dev/null +++ b/crates/fidc-core/src/factor_events.rs @@ -0,0 +1,1047 @@ +//! Causal, typed indicator/event expressions shared by research and trading. +use chrono::{DateTime, FixedOffset}; +use serde::{Deserialize, Serialize}; +use serde_json::{Value, json}; +use std::collections::BTreeMap; +use ta_lib::{ + Core, + abstract_api::{self, InputType, OptInputType, OutputType}, +}; + +pub const CONTRACT: &str = "fidc_factor_event_expression_v1"; +pub const TA_REV: &str = "dd5a90259a3f9e04e2da9f38bf0719a841b40108"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] +pub enum Expr { + Number { + value: f64, + }, + Field { + name: String, + }, + Indicator { + name: String, + #[serde(default)] + inputs: Vec, + #[serde(default)] + parameters: BTreeMap, + #[serde(default)] + output: usize, + }, + Operator { + name: String, + args: Vec, + #[serde(default)] + window: Option, + }, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct Frame { + pub symbol: String, + pub frequency: String, + pub decision_at: DateTime, + pub timestamps: Vec>, + pub available_at: Vec>, + pub fields: BTreeMap>>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum ValueType { + Number, + Boolean, +} + +#[derive(Debug, Clone, Serialize)] +pub struct Series { + pub value_type: ValueType, + pub values: Vec>, +} + +const OPERATORS: &[&str] = &[ + "GT", + "GTE", + "LT", + "LTE", + "EQ", + "NEQ", + "BETWEEN", + "OUTSIDE", + "CROSS_ABOVE", + "CROSS_BELOW", + "BREAK_ABOVE", + "BREAK_BELOW", + "BREAK_HIGH", + "BREAK_LOW", + "CHANGE", + "DIFF", + "DELTA", + "PCT_CHANGE", + "LOG_RETURN", + "RISING", + "FALLING", + "NON_DECREASING", + "NON_INCREASING", + "TURN_UP", + "TURN_DOWN", + "BOTTOM_REVERSAL", + "TOP_REVERSAL", + "SLOPE", + "SLOPE_CHANGE", + "ACCELERATION", + "HHV", + "LLV", + "ARGMAX", + "ARGMIN", + "DISTANCE_TO_HIGH", + "DISTANCE_TO_LOW", + "NEW_HIGH", + "NEW_LOW", + "NEAR_HIGH", + "NEAR_LOW", + "BULLISH_DIVERGENCE", + "BEARISH_DIVERGENCE", + "ZSCORE", + "MINMAX", + "STANDARDIZE", + "NORMALIZE", + "COUNT", + "COUNT_TRUE", + "CONSECUTIVE", + "BARS_SINCE", + "DURATION", + "DAYS_SINCE", + "TIME_SINCE", + "REF", + "LAG", + "PREV", + "SHIFT", + "ROLLING_MEAN", + "ROLLING_SUM", + "ROLLING_STD", + "ROLLING_MAX", + "ROLLING_MIN", + "ROLLING_MEDIAN", + "ROLLING_CORR", + "ROLLING_COV", + "AND", + "OR", + "NOT", + "XOR", + "ADD", + "SUB", + "MUL", + "DIV", + "ABS", + "MAX", + "MIN", + "LOG", + "SQRT", + "POWER", + "CUMMAX", + "CUMMIN", + "SIGN", + "IF", +]; + +pub fn catalog() -> Value { + let indicators: Vec = abstract_api::funcs().map(|f| json!({ + "name":f.name, "group":format!("{:?}",f.group), "description":f.hint, + "inputs":f.inputs.iter().map(|p|json!({"name":p.param_name,"kind":format!("{:?}",p.kind),"flags":p.flags.0})).collect::>(), + "parameters":f.opt_inputs.iter().map(|p|json!({"name":p.param_name,"label":p.display_name,"description":p.hint,"domain":format!("{:?}",p.kind)})).collect::>(), + "outputs":f.outputs.iter().enumerate().map(|(i,p)|json!({"index":i,"name":p.param_name,"kind":format!("{:?}",p.kind)})).collect::>(), + "unstable_period":format!("{:?}",f.unst_id), "production_eligible":false, + })).collect(); + json!({"contract":CONTRACT,"library":{"name":"TA-Lib native Rust","revision":TA_REV,"license":"BSD-3-Clause"}, + "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", + "breakout":"previous_window_excludes_current","boolean":"three_valued_logic","daily_execution":"next_completed_session", + "minute_execution":"strictly_after_completed_bar","cross_section":"requires_separate_complete_universe_contract"}}) +} + +impl Frame { + pub fn validate(&self) -> Result<(), String> { + let n = self.timestamps.len(); + if self.symbol.is_empty() + || n == 0 + || n > 200_000 + || self.available_at.len() != n + || self.fields.len() > 100 + || n.saturating_mul(self.fields.len()) > 1_000_000 + { + return Err("factor_frame_invalid: identity/shape/limit".into()); + } + if !["1d", "1w", "1m", "5m", "15m", "30m", "60m"].contains(&self.frequency.as_str()) { + return Err("factor_frame_invalid: unsupported_frequency".into()); + } + for i in 0..n { + if (i > 0 && self.timestamps[i] <= self.timestamps[i - 1]) + || self.available_at[i] < self.timestamps[i] + || self.available_at[i] > self.decision_at + { + return Err(format!( + "factor_input_not_visible: {} index={i}", + self.symbol + )); + } + } + for (field, values) in &self.fields { + if values.len() != n || values.iter().flatten().any(|v| !v.is_finite()) { + return Err(format!("factor_field_invalid: {} {field}", self.symbol)); + } + } + Ok(()) + } +} + +pub fn evaluate(expr: &Expr, frame: &Frame) -> Result { + frame.validate()?; + fn cost(expr: &Expr, depth: usize, nodes: &mut usize) -> Result { + *nodes += 1; + if depth > 24 || *nodes > 256 { + return Err("factor_expression_size_exceeded".into()); + } + let (children, own) = match expr { + Expr::Indicator { + inputs, parameters, .. + } => ( + inputs.as_slice(), + parameters + .values() + .filter_map(Value::as_u64) + .max() + .unwrap_or(30) + .min(1_000_000) as usize, + ), + Expr::Operator { args, window, .. } => (args.as_slice(), window.unwrap_or(1)), + _ => (&[][..], 1), + }; + children.iter().try_fold(own, |total, child| { + Ok(total.saturating_add(cost(child, depth + 1, nodes)?)) + }) + } + if frame + .timestamps + .len() + .saturating_mul(cost(expr, 0, &mut 0)?) + > 20_000_000 + { + return Err("factor_expression_compute_budget_exceeded".into()); + } + evaluate_inner(expr, frame, 0) +} + +fn evaluate_inner(expr: &Expr, frame: &Frame, depth: usize) -> Result { + if depth > 24 { + return Err("factor_expression_too_deep".into()); + } + match expr { + Expr::Number { value } if value.is_finite() => Ok(Series { + value_type: ValueType::Number, + values: vec![Some(*value); frame.timestamps.len()], + }), + Expr::Number { .. } => Err("factor_constant_nonfinite".into()), + Expr::Field { name } => Ok(Series { + value_type: ValueType::Number, + values: frame + .fields + .get(name) + .ok_or_else(|| format!("factor_source_field_missing: {} {name}", frame.symbol))? + .clone(), + }), + Expr::Indicator { + name, + inputs, + parameters, + output, + } => indicator(name, inputs, parameters, *output, frame, depth), + Expr::Operator { name, args, window } => { + if args.len() > 16 { + return Err("factor_operator_arity_exceeded".into()); + } + let args = args + .iter() + .map(|a| evaluate_inner(a, frame, depth + 1)) + .collect::, _>>()?; + operator(name, &args, *window, frame) + } + } +} + +fn indicator( + name: &str, + inputs: &[Expr], + parameters: &BTreeMap, + output: usize, + frame: &Frame, + depth: usize, +) -> Result { + let id = + abstract_api::get_func_handle(name).ok_or_else(|| format!("indicator_unknown: {name}"))?; + let info = id.info(); + if output >= info.outputs.len() { + return Err("indicator_output_invalid".into()); + } + let real_count = info + .inputs + .iter() + .filter(|i| i.kind == InputType::Real) + .count(); + if inputs.len() != real_count || info.inputs.iter().any(|i| i.kind == InputType::Integer) { + return Err(format!( + "indicator_inputs_invalid: {name} expects {real_count} real series" + )); + } + let mut data = inputs + .iter() + .map(|a| evaluate_inner(a, frame, depth + 1)) + .collect::, _>>()?; + if data.iter().any(|s| s.value_type != ValueType::Number) { + return Err("indicator_requires_numeric_input".into()); + } + let price_names = ["open", "high", "low", "close", "volume", "open_interest"]; + let flags = info + .inputs + .iter() + .filter(|i| i.kind == InputType::Price) + .fold(0, |v, i| v | i.flags.0); + let mut price_indices = [None; 6]; + for (i, field) in price_names.iter().enumerate() { + if flags & (1 << i) != 0 { + price_indices[i] = Some(data.len()); + data.push(evaluate_inner( + &Expr::Field { + name: (*field).into(), + }, + frame, + depth + 1, + )?); + } + } + let core = Core::new(); + let mut validation = id.new_call(&core); + for (key, v) in parameters { + let slot = info + .opt_inputs + .iter() + .position(|p| p.param_name == key) + .ok_or_else(|| format!("indicator_parameter_unknown: {name}.{key}"))?; + match info.opt_inputs[slot].kind { + OptInputType::IntegerRange { .. } | OptInputType::IntegerList { .. } => { + let v = v + .as_i64() + .and_then(|v| i32::try_from(v).ok()) + .ok_or("indicator_parameter_requires_integer")?; + validation.set_opt(slot, v).map_err(|e| format!("{e:?}"))?; + } + _ => { + validation + .set_opt( + slot, + v.as_f64() + .filter(|v| v.is_finite()) + .ok_or("indicator_parameter_requires_finite_number")?, + ) + .map_err(|e| format!("{e:?}"))?; + } + } + } + let lookback = validation + .lookback() + .map_err(|e| format!("indicator_parameter_invalid: {name} {e:?}"))?; + let n = frame.timestamps.len(); + let mut result = vec![None; n]; + let mut start = 0; + // Never bridge missing source observations. Recursive indicators rewarm after a gap. + while start < n { + if data.iter().any(|s| s.values[start].is_none()) { + start += 1; + continue; + } + let mut end = start + 1; + while end < n && data.iter().all(|s| s.values[end].is_some()) { + end += 1; + } + if end - start <= lookback { + start = end; + continue; + } + let arrays = data + .iter() + .map(|s| { + s.values[start..end] + .iter() + .map(|v| v.unwrap()) + .collect::>() + }) + .collect::>(); + let mut float_out = (0..info.outputs.len()) + .map(|_| vec![0.0; end - start]) + .collect::>(); + let mut int_out = (0..info.outputs.len()) + .map(|_| vec![0i32; end - start]) + .collect::>(); + let mut call = id.new_call(&core); + for (key, v) in parameters { + let slot = info + .opt_inputs + .iter() + .position(|p| p.param_name == key) + .unwrap(); + match info.opt_inputs[slot].kind { + OptInputType::IntegerRange { .. } | OptInputType::IntegerList { .. } => { + call.set_opt(slot, v.as_i64().unwrap() as i32) + .map_err(|e| format!("{e:?}"))?; + } + _ => { + call.set_opt(slot, v.as_f64().unwrap()) + .map_err(|e| format!("{e:?}"))?; + } + } + } + let mut real_slot = 0; + for (slot, i) in info.inputs.iter().enumerate() { + if i.kind == InputType::Real { + call.set_input(slot, &arrays[real_slot]) + .map_err(|e| format!("{e:?}"))?; + real_slot += 1; + } else { + let p = price_indices.map(|i| i.map(|i| arrays[i].as_slice())); + call.set_price_input(slot, p[0], p[1], p[2], p[3], p[4], p[5]) + .map_err(|e| format!("{e:?}"))?; + } + } + for (slot, (floats, ints)) in float_out.iter_mut().zip(int_out.iter_mut()).enumerate() { + if info.outputs[slot].kind == OutputType::Real { + call.set_output(slot, floats) + .map_err(|e| format!("{e:?}"))?; + } else { + call.set_int_output(slot, ints) + .map_err(|e| format!("{e:?}"))?; + } + } + let range = call + .call(0, end - start - 1) + .map_err(|e| format!("indicator_failed: {name} {e:?}"))?; + drop(call); + for j in 0..range.count { + let value = if info.outputs[output].kind == OutputType::Real { + float_out[output][j] + } else { + int_out[output][j] as f64 + }; + if !value.is_finite() { + return Err(format!( + "indicator_nonfinite: {name} index={}", + start + range.beg_idx + j + )); + } + result[start + range.beg_idx + j] = Some(value); + } + start = end; + } + Ok(Series { + value_type: ValueType::Number, + values: result, + }) +} + +fn average(v: &[f64]) -> f64 { + v[0] + v + .iter() + .skip(1) + .map(|x| (x - v[0]) / v.len() as f64) + .sum::() +} +fn slope(v: &[f64]) -> f64 { + let x = (v.len() - 1) as f64 / 2.0; + let y = average(v); + let num = v + .iter() + .enumerate() + .map(|(i, v)| (i as f64 - x) * (v - y)) + .sum::(); + let den = (0..v.len()).map(|i| (i as f64 - x).powi(2)).sum::(); + num / den +} +fn boolean(v: bool) -> Option { + Some(if v { 1.0 } else { 0.0 }) +} + +fn operator( + name: &str, + args: &[Series], + window: Option, + frame: &Frame, +) -> Result { + if !OPERATORS.contains(&name) { + return Err(format!("operator_not_registered: {name}")); + } + let bool_input = matches!( + name, + "AND" + | "OR" + | "NOT" + | "XOR" + | "COUNT" + | "COUNT_TRUE" + | "CONSECUTIVE" + | "BARS_SINCE" + | "DURATION" + | "DAYS_SINCE" + | "TIME_SINCE" + ); + let lag = matches!(name, "REF" | "LAG" | "PREV" | "SHIFT"); + if args.is_empty() + || (name == "IF" + && (args.len() != 3 + || args[0].value_type != ValueType::Boolean + || args[1].value_type != args[2].value_type)) + || (!lag + && name != "IF" + && args + .iter() + .any(|a| (a.value_type == ValueType::Boolean) != bool_input)) + { + return Err(format!("operator_input_type_invalid: {name}")); + } + let arity = match name { + "BETWEEN" | "OUTSIDE" | "IF" => 3, + "GT" | "GTE" | "LT" | "LTE" | "EQ" | "NEQ" | "CROSS_ABOVE" | "CROSS_BELOW" + | "BREAK_ABOVE" | "BREAK_BELOW" | "ADD" | "SUB" | "MUL" | "DIV" | "MAX" | "MIN" + | "POWER" | "XOR" | "ROLLING_CORR" | "ROLLING_COV" | "NEAR_HIGH" | "NEAR_LOW" + | "BULLISH_DIVERGENCE" | "BEARISH_DIVERGENCE" => 2, + "AND" | "OR" => args.len(), + _ => 1, + }; + if args.len() != arity { + return Err(format!("operator_arity_invalid: {name}")); + } + let windowed = matches!( + name, + "BREAK_HIGH" + | "BREAK_LOW" + | "RISING" + | "FALLING" + | "NON_DECREASING" + | "NON_INCREASING" + | "SLOPE" + | "SLOPE_CHANGE" + | "HHV" + | "LLV" + | "ARGMAX" + | "ARGMIN" + | "DISTANCE_TO_HIGH" + | "DISTANCE_TO_LOW" + | "NEW_HIGH" + | "NEW_LOW" + | "NEAR_HIGH" + | "NEAR_LOW" + | "BULLISH_DIVERGENCE" + | "BEARISH_DIVERGENCE" + | "ZSCORE" + | "STANDARDIZE" + | "MINMAX" + | "NORMALIZE" + | "COUNT" + | "COUNT_TRUE" + ) || name.starts_with("ROLLING_"); + let n = window.unwrap_or(1); + if n == 0 + || n > 10_000 + || (windowed && window.is_none()) + || (matches!( + name, + "SLOPE" + | "SLOPE_CHANGE" + | "ZSCORE" + | "STANDARDIZE" + | "ROLLING_STD" + | "ROLLING_CORR" + | "ROLLING_COV" + ) && n < 2) + { + return Err(format!("operator_window_invalid: {name}")); + } + let returns_bool = matches!( + name, + "GT" | "GTE" + | "LT" + | "LTE" + | "EQ" + | "NEQ" + | "BETWEEN" + | "OUTSIDE" + | "CROSS_ABOVE" + | "CROSS_BELOW" + | "BREAK_ABOVE" + | "BREAK_BELOW" + | "BREAK_HIGH" + | "BREAK_LOW" + | "RISING" + | "FALLING" + | "NON_DECREASING" + | "NON_INCREASING" + | "TURN_UP" + | "TURN_DOWN" + | "BOTTOM_REVERSAL" + | "TOP_REVERSAL" + | "NEW_HIGH" + | "NEW_LOW" + | "NEAR_HIGH" + | "NEAR_LOW" + | "BULLISH_DIVERGENCE" + | "BEARISH_DIVERGENCE" + | "AND" + | "OR" + | "NOT" + | "XOR" + ); + let len = frame.timestamps.len(); + let mut out = vec![None; len]; + let mut last_true = None; + let mut consecutive = Some(0usize); + let mut extreme: Option = None; + let mut cumulative_complete = true; + for i in 0..len { + let a = args[0].values[i]; + let b = args.get(1).and_then(|a| a.values[i]); + let at = |j: usize| args[0].values.get(j).copied().flatten(); + let history = |end: usize, count: usize| -> Option> { + if end < count { + None + } else { + args[0].values[end - count..end].iter().copied().collect() + } + }; + out[i] = match name { + "IF" => a.and_then(|a| { + if a == 1.0 { + args[1].values[i] + } else { + args[2].values[i] + } + }), + "SIGN" => a.map(|v| { + if v == 0.0 { + 0.0 + } else if v > 0.0 { + 1.0 + } else { + -1.0 + } + }), + "CUMMAX" | "CUMMIN" => { + cumulative_complete &= a.is_some(); + extreme = a.filter(|_| cumulative_complete).map(|v| { + extreme.map_or(v, |p| if name == "CUMMAX" { p.max(v) } else { p.min(v) }) + }); + extreme + } + "AND" => { + if args.iter().any(|a| a.values[i] == Some(0.0)) { + Some(0.0) + } else if args.iter().any(|a| a.values[i].is_none()) { + None + } else { + Some(1.0) + } + } + "OR" => { + if args.iter().any(|a| a.values[i] == Some(1.0)) { + Some(1.0) + } else if args.iter().any(|a| a.values[i].is_none()) { + None + } else { + Some(0.0) + } + } + "NOT" => a.map(|v| 1.0 - v), + "XOR" => a.zip(b).and_then(|(a, b)| boolean(a != b)), + "GT" | "GTE" | "LT" | "LTE" | "EQ" | "NEQ" => a.zip(b).and_then(|(a, b)| { + boolean(match name { + "GT" => a > b, + "GTE" => a >= b, + "LT" => a < b, + "LTE" => a <= b, + "EQ" => a == b, + _ => a != b, + }) + }), + "BETWEEN" | "OUTSIDE" => a.zip(b).zip(args[2].values[i]).and_then(|((a, b), c)| { + if b > c { + None + } else { + boolean((a >= b && a <= c) == (name == "BETWEEN")) + } + }), + "CROSS_ABOVE" | "CROSS_BELOW" | "BREAK_ABOVE" | "BREAK_BELOW" => { + if i == 0 { + None + } else { + a.zip(b).zip(at(i - 1).zip(args[1].values[i - 1])).and_then( + |((a, b), (p, q))| { + boolean(if name.ends_with("ABOVE") { + p <= q && a > b + } else { + p >= q && a < b + }) + }, + ) + } + } + "REF" | "LAG" | "PREV" | "SHIFT" => i.checked_sub(n).and_then(at), + "CHANGE" | "DIFF" | "DELTA" | "PCT_CHANGE" | "LOG_RETURN" => a + .zip(i.checked_sub(n).and_then(at)) + .and_then(|(a, p)| match name { + "PCT_CHANGE" => { + if p == 0.0 { + None + } else { + Some(a / p - 1.0) + } + } + "LOG_RETURN" => { + if a <= 0.0 || p <= 0.0 { + None + } else { + Some((a / p).ln()) + } + } + _ => Some(a - p), + }), + "ACCELERATION" => a + .zip(i.checked_sub(n).and_then(at)) + .zip(i.checked_sub(n * 2).and_then(at)) + .map(|((a, p), q)| a - 2.0 * p + q), + "BULLISH_DIVERGENCE" | "BEARISH_DIVERGENCE" => { + if i < n || n < 4 { + None + } else { + let price: Option> = + args[0].values[i - n..=i].iter().copied().collect(); + let indicator: Option> = + args[1].values[i - n..=i].iter().copied().collect(); + price.zip(indicator).and_then(|(price, indicator)| { + let low = name == "BULLISH_DIVERGENCE"; + let pivots = (1..n) + .filter(|&j| { + if low { + price[j] < price[j - 1] && price[j] < price[j + 1] + } else { + price[j] > price[j - 1] && price[j] > price[j + 1] + } + }) + .collect::>(); + if pivots.last() != Some(&(n - 1)) || pivots.len() < 2 { + return boolean(false); + } + let a = pivots[pivots.len() - 2]; + let b = n - 1; + boolean(if low { + price[b] < price[a] && indicator[b] > indicator[a] + } else { + price[b] > price[a] && indicator[b] < indicator[a] + }) + }) + } + } + "TURN_UP" | "TURN_DOWN" | "BOTTOM_REVERSAL" | "TOP_REVERSAL" => { + if i < 2 { + None + } else { + a.zip(at(i - 1)).zip(at(i - 2)).and_then(|((a, p), q)| { + if name == "ACCELERATION" { + Some(a - 2.0 * p + q) + } else { + boolean(if matches!(name, "TURN_UP" | "BOTTOM_REVERSAL") { + p < q && a > p + } else { + p > q && a < p + }) + } + }) + } + } + "ABS" => a.map(f64::abs), + "LOG" => a.filter(|v| *v > 0.0).map(f64::ln), + "SQRT" => a.filter(|v| *v >= 0.0).map(f64::sqrt), + "ADD" => a.zip(b).map(|(a, b)| a + b), + "SUB" => a.zip(b).map(|(a, b)| a - b), + "MUL" => a.zip(b).map(|(a, b)| a * b), + "DIV" => a.zip(b).filter(|(_, b)| *b != 0.0).map(|(a, b)| a / b), + "MAX" => a.zip(b).map(|(a, b)| a.max(b)), + "MIN" => a.zip(b).map(|(a, b)| a.min(b)), + "POWER" => a.zip(b).map(|(a, b)| a.powf(b)), + "BARS_SINCE" | "DAYS_SINCE" | "TIME_SINCE" => { + if a == Some(1.0) { + last_true = Some(i); + } + if a.is_none() { + last_true = None; + } + last_true.map(|t| { + if name == "BARS_SINCE" { + (i - t) as f64 + } else { + let secs = (frame.timestamps[i] - frame.timestamps[t]).num_seconds() as f64; + if name == "DAYS_SINCE" { + secs / 86400.0 + } else { + secs + } + } + }) + } + "CONSECUTIVE" | "DURATION" => { + consecutive = match a { + Some(1.0) => consecutive.map(|v| v + 1), + Some(_) => Some(0), + None => None, + }; + consecutive.map(|v| v as f64) + } + "BREAK_HIGH" | "NEW_HIGH" | "BREAK_LOW" | "NEW_LOW" => { + a.zip(history(i, n)).and_then(|(a, v)| { + boolean(if matches!(name, "BREAK_HIGH" | "NEW_HIGH") { + a > v.into_iter().fold(f64::NEG_INFINITY, f64::max) + } else { + a < v.into_iter().fold(f64::INFINITY, f64::min) + }) + }) + } + "RISING" | "FALLING" | "NON_DECREASING" | "NON_INCREASING" => history(i + 1, n + 1) + .and_then(|v| { + boolean(v.windows(2).all(|p| match name { + "RISING" => p[1] > p[0], + "FALLING" => p[1] < p[0], + "NON_DECREASING" => p[1] >= p[0], + _ => p[1] <= p[0], + })) + }), + "SLOPE_CHANGE" => history(i + 1, n) + .zip(history(i, n)) + .map(|(a, b)| slope(&a) - slope(&b)), + _ => history(i + 1, n).and_then(|mut v| { + let mean = average(&v); + let lo = v.iter().copied().fold(f64::INFINITY, f64::min); + let hi = v.iter().copied().fold(f64::NEG_INFINITY, f64::max); + let variance = v.iter().map(|v| (v - mean).powi(2)).sum::() / n as f64; + match name { + "HHV" | "ROLLING_MAX" => Some(hi), + "LLV" | "ROLLING_MIN" => Some(lo), + "ARGMAX" => v.iter().rposition(|x| *x == hi).map(|p| (n - 1 - p) as f64), + "ARGMIN" => v.iter().rposition(|x| *x == lo).map(|p| (n - 1 - p) as f64), + "DISTANCE_TO_HIGH" => { + if hi == 0.0 { + None + } else { + Some(v[n - 1] / hi - 1.0) + } + } + "DISTANCE_TO_LOW" => { + if lo == 0.0 { + None + } else { + Some(v[n - 1] / lo - 1.0) + } + } + "NEAR_HIGH" | "NEAR_LOW" => b.filter(|b| *b >= 0.0).and_then(|b| { + let base = if name == "NEAR_HIGH" { hi } else { lo }; + if base == 0.0 { + None + } else { + boolean((v[n - 1] / base - 1.0).abs() <= b) + } + }), + "ZSCORE" | "STANDARDIZE" => { + if variance == 0.0 { + None + } else { + Some((v[n - 1] - mean) / variance.sqrt()) + } + } + "MINMAX" | "NORMALIZE" => { + if hi == lo { + None + } else { + Some((v[n - 1] - lo) / (hi - lo)) + } + } + "ROLLING_MEAN" => Some(mean), + "ROLLING_SUM" | "COUNT" | "COUNT_TRUE" => Some(v.iter().sum()), + "ROLLING_STD" => Some(variance.sqrt()), + "ROLLING_MEDIAN" => { + v.sort_by(f64::total_cmp); + Some(if n % 2 == 1 { + v[n / 2] + } else { + (v[n / 2 - 1] + v[n / 2]) / 2.0 + }) + } + "SLOPE" => Some(slope(&v)), + "ROLLING_CORR" | "ROLLING_COV" => { + let b: Option> = + args[1].values[i + 1 - n..=i].iter().copied().collect(); + b.and_then(|b| { + let bm = average(&b); + let cov = v + .iter() + .zip(&b) + .map(|(a, b)| (a - mean) * (b - bm)) + .sum::() + / n as f64; + if name == "ROLLING_COV" { + Some(cov) + } else { + let bv = b.iter().map(|b| (b - bm).powi(2)).sum::() / n as f64; + let d = (variance * bv).sqrt(); + if d == 0.0 { None } else { Some(cov / d) } + } + }) + } + _ => None, + } + }), + } + .filter(|v| v.is_finite()); + } + Ok(Series { + value_type: if name == "IF" { + args[1].value_type + } else if lag { + args[0].value_type + } else if returns_bool { + ValueType::Boolean + } else { + ValueType::Number + }, + values: out, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + fn frame(values: Vec>) -> Frame { + let start = DateTime::parse_from_rfc3339("2026-09-01T15:30:00+08:00").unwrap(); + let times = (0..values.len()) + .map(|i| start + chrono::Duration::days(i as i64)) + .collect::>(); + Frame { + symbol: "TEST".into(), + frequency: "1d".into(), + decision_at: *times.last().unwrap(), + available_at: times.clone(), + timestamps: times, + fields: BTreeMap::from([("close".into(), values)]), + } + } + fn expr(v: Value) -> Expr { + serde_json::from_value(v).unwrap() + } + #[test] + fn ta_sma_real_values_and_parameter_validation() { + let frame = frame(vec![Some(1.0), Some(2.0), Some(3.0), Some(4.0)]); + let e = expr( + json!({"kind":"indicator","name":"SMA","inputs":[{"kind":"field","name":"close"}],"parameters":{"optInTimePeriod":3}}), + ); + assert_eq!( + evaluate(&e, &frame).unwrap().values, + vec![None, None, Some(2.0), Some(3.0)] + ); + let bad = expr( + json!({"kind":"indicator","name":"SMA","inputs":[{"kind":"field","name":"close"}],"parameters":{"period":3}}), + ); + assert!( + evaluate(&bad, &frame) + .unwrap_err() + .contains("parameter_unknown") + ); + } + #[test] + fn cross_is_event_not_state_and_never_uses_future() { + let f = frame(vec![ + Some(9.0), + Some(10.0), + Some(11.0), + Some(12.0), + Some(8.0), + ]); + let e = expr( + json!({"kind":"operator","name":"CROSS_ABOVE","args":[{"kind":"field","name":"close"},{"kind":"number","value":10.0}]}), + ); + assert_eq!( + evaluate(&e, &f).unwrap().values, + vec![None, Some(0.0), Some(1.0), Some(0.0), Some(0.0)] + ); + let mut invalid = f.clone(); + invalid.available_at[4] = invalid.decision_at + chrono::Duration::seconds(1); + assert!(evaluate(&e, &invalid).is_err()); + } + #[test] + fn missing_is_not_zero_and_breakout_excludes_current() { + let f = frame(vec![Some(1.0), Some(2.0), Some(3.0), None, Some(5.0)]); + let e = expr( + json!({"kind":"operator","name":"BREAK_HIGH","window":2,"args":[{"kind":"field","name":"close"}]}), + ); + assert_eq!( + evaluate(&e, &f).unwrap().values, + vec![None, None, Some(1.0), None, None] + ); + let zero = expr( + json!({"kind":"operator","name":"DIV","args":[{"kind":"field","name":"close"},{"kind":"number","value":0}]}), + ); + assert!( + evaluate(&zero, &f) + .unwrap() + .values + .iter() + .all(Option::is_none) + ); + } + #[test] + fn ta_rewarms_after_gap_and_const_zscore_is_unknown() { + let f = frame(vec![Some(1.0), Some(1.0), None, Some(2.0), Some(2.0)]); + let e = expr( + json!({"kind":"indicator","name":"SMA","inputs":[{"kind":"field","name":"close"}],"parameters":{"optInTimePeriod":2}}), + ); + assert_eq!( + evaluate(&e, &f).unwrap().values, + vec![None, Some(1.0), None, None, Some(2.0)] + ); + let e = expr( + json!({"kind":"operator","name":"ZSCORE","window":2,"args":[{"kind":"field","name":"close"}]}), + ); + assert!(evaluate(&e, &f).unwrap().values.iter().all(Option::is_none)); + } + #[test] + fn no_event_has_no_bars_since_and_type_errors_reject() { + let f = frame(vec![Some(1.0), Some(1.0), Some(1.0)]); + let state = json!({"kind":"operator","name":"GT","args":[{"kind":"field","name":"close"},{"kind":"number","value":5}]}); + let e = expr(json!({"kind":"operator","name":"BARS_SINCE","args":[state]})); + assert!(evaluate(&e, &f).unwrap().values.iter().all(Option::is_none)); + assert!( + evaluate( + &expr( + json!({"kind":"operator","name":"NOT","args":[{"kind":"field","name":"close"}]}) + ), + &f + ) + .is_err() + ); + } + #[test] + fn literal_unknown_fields_reject_and_catalog_is_not_trading_permission() { + assert!( + serde_json::from_value::(json!({"kind":"number","value":1,"account_id":2})) + .is_err() + ); + let c = catalog(); + assert!(c["indicators"].as_array().unwrap().len() > 190); + assert_eq!(c["live_routing"], false); + } +} diff --git a/crates/fidc-core/src/lib.rs b/crates/fidc-core/src/lib.rs index 7213797..859b921 100644 --- a/crates/fidc-core/src/lib.rs +++ b/crates/fidc-core/src/lib.rs @@ -3,6 +3,8 @@ pub mod calendar; pub mod cost; pub mod data; pub mod daily_patterns; +pub mod factor_events; +pub mod factor_cross_section; pub mod engine; pub mod event_bus; pub mod events;