//! Completed-session OHLCV rules shared by research and strategy execution. use crate::DataSet; use chrono::{FixedOffset, NaiveDate, TimeZone}; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; use std::collections::{BTreeMap, BTreeSet}; 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":"冻结历史窗口与表达式;预热不足或未定义值不产生信号。复用共享指标事件内核,不修改既有任务。"}, "session_event":{"label":"已完成分钟事件","parameters":{"opening_minutes":[30,1,120],"volume_window":[5,2,120],"volume_multiple":[3.0,1,20]},"stages":["selection","buy","sell"],"method":"仅本交易日完整分钟OHLCVA,信号K线必须早于执行时点;不使用盘口快照伪造K线。"}, "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日最大量的指定倍数;不等同价格创新高。"}, "mean_volume_spike":{"label":"均量倍增","parameters":{"volume_window":[5,2,60],"volume_multiple":[3.0,1,10]},"stages":["selection","buy"],"method":"量达到此前N个交易日均量的M倍且当日上涨。分母不含当日;保留与最大量规则的区别。"}, "mean_shrink_breakout":{"label":"倍量后缩量阳线突破","parameters":{"spike_lookback":[5,2,30],"volume_window":[5,2,60],"volume_multiple":[3.0,1,10],"shrink_ratio":[0.5,0.01,1]},"stages":["selection","buy"],"method":"此前出现N日均量M倍放量,当前缩量阳线收盘突破该放量日最高价。"}, "breakout_retest":{"label":"突破回踩站回","parameters":{"high_window":[60,5,252],"retest_lookback":[10,2,30],"price_tolerance":[0.02,0,0.2],"shrink_ratio":[0.8,0.01,1]},"stages":["selection","buy"],"method":"观察窗先收盘突破此前N日最高价,随后低点回踩突破位容差区,今日收盘站回该位且不低于昨日、成交量收缩。突破与回踩不得同日。"}, "limit_consolidation":{"label":"涨停后整理(日线)","parameters":{"anchor_lag":[4,2,30],"price_band":[0.05,0,0.3],"volume_band":[0.15,0,2],"ma_window":[5,2,60]},"stages":["selection","buy"],"method":"明确T-i日按真实涨停价收盘,后续收盘和量相对锚日偏离受限,今日收盘低于完整日线均线;不是盘中动态MA条件。"}, "shrink_breakout":{"label":"缩量突破","parameters":{"spike_lookback":[5,2,30],"volume_window":[5,2,60],"volume_multiple":[3.0,1,10],"shrink_ratio":[0.5,0.01,1]},"stages":["selection","buy"],"method":"此前观察窗有放量日,今日收盘超过该日最高价,成交量不超过其指定比例。"}, "ma_below":{"label":"均线下方","parameters":{"ma_window":[20,2,252]},"stages":["sell"],"method":"完整收盘价低于含当日的N日均线;独立卖出条件。"}, "volume_down":{"label":"放量下跌","parameters":{"volume_window":[5,2,60],"volume_multiple":[3.0,1,10]},"stages":["sell"],"method":"当日下跌且量达到此前N日最大量的指定倍数。"} },"data_frequency":"1d","execution_policies":["next_session_open"]}) } #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct PatternSpec { pub template: String, #[serde(default)] pub parameters: BTreeMap, #[serde(default, skip_serializing_if = "Option::is_none")] pub expression: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub execution_context: Option, #[serde(default,skip_serializing_if="Option::is_none")] pub session_event:Option, } impl PatternSpec { pub fn validate(self) -> Result { if self.template=="session_event" { if !self.session_event.as_deref().is_some_and(|id|crate::session_events::EVENTS.contains(&id)) || self.execution_context.is_some() {return Err("session_event_contract_invalid".into());} } else if self.session_event.is_some() {return Err("unexpected_session_event_id".into());} let allowed = if let Some(context) = &self.execution_context { context.validate(self.expression.as_ref().ok_or("pattern_context_requires_expression")?)?; crate::pattern_context::CONTEXT_FIELDS } else { &[] }; let spec = self.validate_with_context(allowed)?; if let Some(context) = &spec.execution_context { if context.rank_universe.len().saturating_mul(spec.history_len()) > 2_000_000 { return Err("pattern_rank_window_budget_exceeded: 完整截面不得截断".into()); } } Ok(spec) } fn validate_with_context(mut self, context_fields: &[&str]) -> Result { if (self.template == "expression") != self.expression.is_some() { return Err("expression_template_requires_expression_only".into()); } if let Some(expr) = &self.expression { let supported = [ "open", "high", "low", "close", "volume", "raw_open", "raw_high", "raw_low", "raw_close", "prev_close", "amount", ]; let missing = crate::factor_events::field_dependencies(expr) .into_iter() .filter(|f| !supported.contains(&f.as_str()) && !context_fields.contains(&f.as_str())) .collect::>(); if !missing.is_empty() { return Err(format!( "expression_source_mapping_required: {}", missing.join(",") )); } } let catalog = catalog(); let definition = catalog["templates"] .get(&self.template) .ok_or("未登记的量价模板")?; let parameters = definition["parameters"].as_object().unwrap(); if self.parameters.keys().any(|k| !parameters.contains_key(k)) { return Err("模板包含未知参数".into()); } for (key, bounds) in parameters { let value = self.parameters.get(key).unwrap_or(&bounds[0]); let number = value .as_f64() .filter(|v| v.is_finite()) .ok_or_else(|| format!("{key}必须为有限数值"))?; if number < bounds[1].as_f64().unwrap() || number > bounds[2].as_f64().unwrap() { return Err(format!("{key}超出允许范围")); } if key.ends_with("window") || key.ends_with("lookback") || key == "anchor_lag" || key=="opening_minutes" { if number.fract() != 0.0 { return Err(format!("{key}必须是整数")); } self.parameters.insert(key.clone(), json!(number as usize)); } else { self.parameters.insert(key.clone(), json!(number)); } } if self.template == "strength" && self.n("fast_window") >= self.n("slow_window") { return Err("短均线必须小于长均线".into()); } Ok(self) } pub fn n(&self, key: &str) -> usize { self.parameters[key].as_u64().unwrap() as usize } pub fn v(&self, key: &str) -> f64 { self.parameters[key].as_f64().unwrap() } pub fn history_len(&self) -> usize { match self.template.as_str() { "session_event"=>1, "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" | "mean_volume_spike" => self.n("volume_window") + 1, "breakout_retest" => self.n("high_window") + self.n("retest_lookback") + 1, "limit_consolidation" => (self.n("anchor_lag")+1).max(self.n("ma_window")), "ma_below" => self.n("ma_window").max(2), "shrink_breakout" | "mean_shrink_breakout" => self.n("spike_lookback") + self.n("volume_window") + 1, _ => unreachable!(), } } } #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct PatternBar { pub date: NaiveDate, pub open: Option, pub high: Option, pub low: Option, pub close: Option, pub volume: Option, #[serde(default)] pub prev_close: Option, #[serde(default)] pub amount: Option, #[serde(default)] pub upper_limit: Option, #[serde(default)] pub no_limit: Option, pub adjustment_factor_backward1: Option, pub paused: Option, #[serde(default)] pub source_path: Option, } #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct PatternSeries { pub symbol: String, #[serde(default)] pub name: Option, #[serde(default)] pub listed_at: Option, pub bars: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct PatternResult { pub symbol: String, pub name: Option, pub matched: bool, pub score: Option, pub checks: Vec, pub values: Value, pub anchor: Value, pub exclusion: Option, } fn number(v: Option, symbol: &str, day: NaiveDate, field: &str) -> Result { v.filter(|v|v.is_finite()).ok_or_else(||format!("pattern_input_invalid: symbol={symbol}, date={day}, field={field}, reason=missing_or_nonfinite")) } fn check(checks: &mut Vec, label: &str, actual: f64, operator: &str, threshold: f64) { let passed = match operator { ">" => actual > threshold, "<" => actual < threshold, ">=" => actual >= threshold, "<=" => actual <= threshold, _ => false, }; checks.push(json!({"label":label,"actual":actual,"operator":operator,"threshold":threshold,"passed":passed})); } fn mean(mut values: impl ExactSizeIterator) -> Result { let count = values.len(); let first = values.next().ok_or("pattern_mean_empty")?; // Center before summation so an unchanged decimal price stays exactly unchanged. let result = first + values .map(|value| (value - first) / count as f64) .sum::(); if !result.is_finite() { return Err("pattern_mean_nonfinite".into()); } Ok(result) } /// No calendar compression, fill-forward prices or numerical substitutes. pub fn evaluate( spec: &PatternSpec, days: &[NaiveDate], series: &PatternSeries, ) -> Result { evaluate_with_context(spec, days, series, &BTreeMap::new(), false) } pub(crate) fn evaluate_with_context( spec: &PatternSpec, days: &[NaiveDate], series: &PatternSeries, context: &BTreeMap>>, numeric_output: bool, ) -> Result { if spec.template=="session_event" {return Err("session_event_requires_completed_minute_endpoint".into());} if days.len() != spec.history_len() || days.windows(2).any(|w| w[0] >= w[1]) { return Err("pattern_calendar_incomplete: 需要完整、唯一且递增的真实交易日窗口".into()); } let by_day = series .bars .iter() .map(|b| (b.date, b)) .collect::>(); if by_day.len() != series.bars.len() || series.bars.iter().any(|b| days.binary_search(&b.date).is_err()) { return Err(format!( "pattern_input_invalid: symbol={}, reason=duplicate_or_out_of_scope", series.symbol )); } let mut unavailable = Vec::new(); let mut prices = Vec::new(); for &day in days { let Some(b) = by_day.get(&day) else { if series.listed_at.is_some_and(|listed| day < listed) { unavailable.push( json!({"date":day,"reason":"before_listing","listed_at":series.listed_at}), ); continue; } return Err(format!( "pattern_input_invalid: symbol={}, date={day}, reason=missing_market_row", series.symbol )); }; let factor = if b.close.is_some_and(|c| c.is_finite() && c > 0.0) { let factor = number( b.adjustment_factor_backward1, &series.symbol, day, "adjustment_factor_backward1", )?; if factor <= 0.0 { return Err(format!( "pattern_input_invalid: symbol={}, date={day}, field=adjustment_factor_backward1, reason=nonpositive", series.symbol )); } factor } else { 1.0 }; let paused = b.paused.ok_or_else(|| { format!( "pattern_input_invalid: symbol={}, date={day}, field=paused", series.symbol ) })?; if paused { unavailable.push( json!({"date":day,"reason":"confirmed_suspension","source_path":b.source_path}), ); continue; } if series.listed_at.is_some_and(|listed| day < listed) { return Err(format!( "pattern_input_invalid: symbol={}, date={day}, reason=price_before_listing", series.symbol )); } let o = number(b.open, &series.symbol, day, "open")?; let h = number(b.high, &series.symbol, day, "high")?; let l = number(b.low, &series.symbol, day, "low")?; let c = number(b.close, &series.symbol, day, "close")?; let v = number(b.volume, &series.symbol, day, "volume")?; if o <= 0.0 || l <= 0.0 || c <= 0.0 || h < o.max(c) || l > o.min(c) || v < 0.0 { return Err(format!( "pattern_input_invalid: symbol={}, date={day}, reason=invalid_ohlcv", series.symbol )); } prices.push((o * factor, h * factor, l * factor, c * factor, v)); } let mut result = PatternResult { symbol: series.symbol.clone(), name: series.name.clone(), matched: false, score: None, checks: vec![], values: json!({}), anchor: Value::Null, exclusion: None, }; if !unavailable.is_empty() { result.exclusion = Some( json!({"reason":"proven_incomplete_window","signal_date":days.last(),"evidence":unavailable}), ); return Ok(result); } let len = prices.len(); let (o, h, l, c, v) = prices[len - 1]; let change = c / prices[len - 2].3 - 1.0; result.values = json!({"close":by_day[&days[len-1]].close,"daily_return":change}); 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 needed = crate::factor_events::field_dependencies(spec.expression.as_ref().unwrap()); for name in ["prev_close", "amount"] { if !needed.contains(name) { continue; } let values = days .iter() .map(|d| { let b = by_day[d]; let value = number( if name == "prev_close" { b.prev_close } else { b.amount }, &series.symbol, *d, name, )?; if value < 0.0 || (name == "prev_close" && value == 0.0) { return Err(format!( "pattern_input_invalid: symbol={}, date={d}, field={name}, reason=invalid_value", series.symbol )); } Ok(Some(value)) }) .collect::, String>>()?; fields.insert(name.into(), values); } for (name, values) in context { if fields.contains_key(name) || values.len() != days.len() || values.iter().flatten().any(|v| !v.is_finite()) { return Err(format!("research_context_invalid: {} {name}", series.symbol)); } fields.insert(name.clone(), values.clone()); } 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 numeric_output { if values.value_type != crate::factor_events::ValueType::Number { return Err("research_rank_input_requires_numeric_expression".into()); } return Ok(result); } 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":"组合条件","actual":latest,"operator":"==","threshold":1,"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))?; let momentum = c / prices[len - 1 - spec.n("momentum_window")].3 - 1.0; score = Some(momentum); result.values["momentum"] = json!(momentum); result.values["fast_ma"] = json!(fast); result.values["slow_ma"] = json!(slow); check(&mut result.checks, "收盘高于短均线", c, ">", fast); check(&mut result.checks, "短均线高于长均线", fast, ">", slow); } "breakout" => { let prior_high = prices[len - 1 - spec.n("high_window")..len - 1] .iter() .map(|b| b.1) .fold(f64::NEG_INFINITY, f64::max); let avg = mean( prices[len - 1 - spec.n("volume_window")..len - 1] .iter() .map(|b| b.4), )?; if avg <= 0.0 { return Err(format!( "pattern_input_invalid: symbol={}, reason=zero_reference_volume", series.symbol )); } let shadow = if h > l { (h - o.max(c)) / (h - l) } else { 0.0 }; score = Some(c / prior_high - 1.0); result.values["volume_ratio"] = json!(v / avg); result.values["upper_shadow"] = json!(shadow); check(&mut result.checks, "收盘突破前高", c, ">", prior_high); check( &mut result.checks, "均量倍数", v / avg, ">=", spec.v("volume_multiple"), ); check( &mut result.checks, "上影比例", shadow, "<=", spec.v("max_upper_shadow"), ); } "volume_spike" | "volume_down" | "mean_volume_spike" => { let reference = &prices[len - 1 - spec.n("volume_window")..len - 1]; let high = if spec.template=="mean_volume_spike" {mean(reference.iter().map(|b|b.4))?} else {reference .iter() .map(|b| b.4) .fold(0.0, f64::max)}; if high <= 0.0 { return Err(format!( "pattern_input_invalid: symbol={}, reason=zero_reference_volume", series.symbol )); } score = Some(v / high); result.values["volume_ratio"] = json!(v / high); check( &mut result.checks, if spec.template=="mean_volume_spike" {"均量倍数"} else {"最大量倍数"}, v / high, ">=", spec.v("volume_multiple"), ); check( &mut result.checks, if spec.template != "volume_down" { "当日上涨" } else { "当日下跌" }, change, if spec.template != "volume_down" { ">" } else { "<" }, 0.0, ); } "ma_below" => { let avg = mean(prices[len - spec.n("ma_window")..].iter().map(|b| b.3))?; score = Some(avg / c - 1.0); result.values["ma"] = json!(avg); check(&mut result.checks, "收盘低于均线", c, "<", avg); } "shrink_breakout" | "mean_shrink_breakout" => { let mut spikes = Vec::new(); let mut eligible = Vec::new(); for i in len - 1 - spec.n("spike_lookback")..len - 1 { let reference=&prices[i - spec.n("volume_window")..i]; let prior = if spec.template=="mean_shrink_breakout"{mean(reference.iter().map(|b|b.4))?}else{reference .iter() .map(|b| b.4) .fold(0.0, f64::max)}; if prior <= 0.0 { return Err(format!( "pattern_input_invalid: symbol={}, date={}, reason=zero_reference_volume", series.symbol, days[i] )); } if prices[i].4 >= prior * spec.v("volume_multiple") { spikes.push(i); if c > prices[i].1 && v <= prices[i].4 * spec.v("shrink_ratio") { eligible.push(i); } } } check( &mut result.checks, "观察窗存在放量日", spikes.len() as f64, ">", 0.0, ); if let Some(&i) = eligible.last().or_else(|| spikes.last()) { score = Some(c / prices[i].1 - 1.0); result.values["spike_date"] = json!(days[i]); result.values["volume_ratio"] = json!(v / prices[i].4); check( &mut result.checks, "收盘突破放量日高点", c, ">", prices[i].1, ); check( &mut result.checks, "缩量比例", v / prices[i].4, "<=", spec.v("shrink_ratio"), ); } if spec.template=="mean_shrink_breakout" {check(&mut result.checks,"当前为阳线",c,">",o);} } "breakout_retest" => { let mut anchors=Vec::new();let mut eligible=Vec::new(); for i in len-1-spec.n("retest_lookback")..len-1 { let level=prices[i-spec.n("high_window")..i].iter().map(|b|b.1).fold(f64::NEG_INFINITY,f64::max); if prices[i].3<=level {continue;} let retraced=prices[i+1..].iter().any(|b| b.2 <= level*(1.0+spec.v("price_tolerance"))); anchors.push((i,level,retraced)); if retraced&&c>=level&&c>=prices[len-2].3&&prices[i].4>0.0&&v<=prices[i].4*spec.v("shrink_ratio") {eligible.push((i,level,retraced));} } check(&mut result.checks,"观察窗存在先前突破",anchors.len() as f64,">",0.0); if let Some(&(i,level,retraced))=eligible.last().or_else(||anchors.last()) { result.values["breakout_date"]=json!(days[i]);result.values["breakout_level"]=json!(level);result.values["days_since_breakout"]=json!(len-1-i); check(&mut result.checks,"突破后曾回踩",if retraced{1.0}else{0.0},">",0.0); check(&mut result.checks,"收盘重新站回突破位",c,">=",level); check(&mut result.checks,"收盘不低于昨日",c,">=",prices[len-2].3); if prices[i].4<=0.0{return Err("突破锚日成交量为零,不能计算缩量比例".into());} check(&mut result.checks,"相对突破日缩量",v/prices[i].4,"<=",spec.v("shrink_ratio")); score=Some(c/level-1.0); } } "limit_consolidation" => { let i=len-1-spec.n("anchor_lag");let anchor=by_day[&days[i]]; let is_limit=if anchor.no_limit==Some(true){false}else{ let upper=number(anchor.upper_limit,&series.symbol,days[i],"upper_limit")?; if upper<=0.0||upper>=99999.0{return Err("涨停事件缺少有效源涨停价或无涨跌幅限制证据,禁止按比例推算".into());} (anchor.close.unwrap()/upper-1.0).abs()<=1e-8 }; if prices[i].4<=0.0{return Err("涨停锚日成交量为零".into());} let price_gap=prices[i+1..].iter().map(|b|(b.3/prices[i].3-1.0).abs()).fold(0.0,f64::max); let volume_gap=prices[i+1..].iter().map(|b|(b.4/prices[i].4-1.0).abs()).fold(0.0,f64::max); let avg=mean(prices[len-spec.n("ma_window")..].iter().map(|b|b.3))?; result.values["limit_date"]=json!(days[i]);result.values["price_deviation"]=json!(price_gap);result.values["volume_deviation"]=json!(volume_gap); check(&mut result.checks,"锚日真实涨停收盘",if is_limit{1.0}else{0.0},">",0.0); check(&mut result.checks,"后续收盘最大偏离",price_gap,"<=",spec.v("price_band")); check(&mut result.checks,"后续成交量最大偏离",volume_gap,"<=",spec.v("volume_band")); check(&mut result.checks,"收盘低于日线均线",c,"<",avg);score=Some(avg/c-1.0); } _ => unreachable!(), } if score.is_some_and(|v| !v.is_finite()) { return Err("pattern_result_nonfinite".into()); } result.score = score; result.matched = result.checks.iter().all(|c| c["passed"] == true); Ok(result) } pub fn evaluate_dataset( spec: &PatternSpec, data: &DataSet, date: NaiveDate, symbol: &str, ) -> Result { let context = crate::pattern_context::build_dataset_context(spec, data, date)?; evaluate_dataset_context(spec, data, date, symbol, &context) } pub fn dataset_series(data: &DataSet, days: &[NaiveDate], symbol: &str) -> PatternSeries { let bars = days.iter().filter_map(|&d| data.market(d, symbol).map(|b| PatternBar { date:d, open:Some(b.open), high:Some(b.high), low:Some(b.low), close:Some(b.close), volume:Some(b.volume as f64), prev_close:Some(b.prev_close), amount:data.factor_numeric_value(d,symbol,"amount"),upper_limit:Some(b.upper_limit), no_limit:data.factor_numeric_value(d,symbol,"no_limit").map(|v|v==1.0), adjustment_factor_backward1:data.factor(d,symbol).and_then(|f|f.adjustment_factor_backward1), paused:Some(b.paused), source_path:None, })).collect(); PatternSeries{symbol:symbol.into(),name:data.instrument(symbol).map(|i|i.name.clone()), listed_at:data.instrument(symbol).and_then(|i|i.listed_at),bars} } pub fn evaluate_dataset_context( spec: &PatternSpec, data: &DataSet, date: NaiveDate, symbol: &str, context: &ResearchContext, ) -> Result { let days = data.calendar().trailing_days(date, spec.history_len()); let mut fields = context.common.clone(); fields.extend(context.by_symbol.get(symbol).cloned().unwrap_or_default()); let outside = spec.execution_context.as_ref().is_some_and(|c| c.rank_expression.is_some() && !c.rank_universe.iter().any(|s|s==symbol)); if outside { for name in ["scope_rank","scope_percentile"] {fields.insert(name.into(),vec![None;days.len()]);} fields.insert("scope_size".into(),vec![Some(spec.execution_context.as_ref().unwrap().rank_universe.len() as f64);days.len()]); } let mut result = evaluate_with_context(spec,&days,&dataset_series(data,&days,symbol),&fields,false)?; if outside && result.score.is_none() { result.exclusion=Some(json!({"reason":"outside_frozen_rank_universe","symbol":symbol,"signal_date":date})); } result.values["execution_context_latest"]=json!(fields.iter().map(|(k,v)|(k,v.last().copied().flatten())).collect::>()); Ok(result) } pub fn evaluate_batch( spec: PatternSpec, days: &[NaiveDate], series: &[PatternSeries], ) -> Result { evaluate_batch_with_policy(spec, days, series, false) } /// Partial results are research diagnostics, never strategy execution inputs. pub fn evaluate_batch_with_policy( spec: PatternSpec, days: &[NaiveDate], series: &[PatternSeries], isolate_data_errors: bool, ) -> Result { let spec = spec.validate()?; if series.is_empty() || series.len() > 200 || series .iter() .map(|s| &s.symbol) .collect::>() .len() != series.len() { return Err("pattern_batch_invalid: 需要1至200只唯一证券".into()); } let rows = series .iter() .map(|s| research_row(evaluate(&spec, days, s), s, isolate_data_errors)) .collect::, _>>()?; Ok( json!({"contract":CONTRACT,"spec":spec,"required_history":spec.history_len(),"rows":rows,"read_only":true}), ) } fn research_row(result: Result, series: &PatternSeries, isolate: bool) -> Result { match result { Ok(row) => Ok(json!(row)), Err(detail) if isolate && detail.starts_with(&format!("pattern_input_invalid: symbol={},", series.symbol)) => { let fields = detail.split(", ").filter_map(|p| p.split_once('=')).collect::>(); Ok(json!({"symbol":series.symbol,"name":series.name,"matched":null,"score":null, "checks":[],"values":{},"anchor":null,"exclusion":null, "data_issue":{"reason":fields.get("reason"),"date":fields.get("date"),"field":fields.get("field"),"detail":detail}})) }, Err(error) => Err(error), } } /// Values are supplied only by the verified research transport or dataset context builder. #[derive(Debug, Clone, Default, Deserialize)] #[serde(deny_unknown_fields)] pub struct ResearchContext { #[serde(default)] pub common: BTreeMap>>, #[serde(default)] pub by_symbol: BTreeMap>>>, } pub fn evaluate_research_batch( spec: PatternSpec, days: &[NaiveDate], series: &[PatternSeries], context: &ResearchContext, numeric_output: bool, ) -> Result { evaluate_research_batch_with_policy(spec, days, series, context, numeric_output, false) } 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"].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())) || context.by_symbol.iter().any(|(s, fields)| !series.iter().any(|row| &row.symbol == s) || fields.keys().any(|k| !symbol_fields.contains(&k.as_str()))) { return Err("research_context_scope_or_fields_invalid".into()); } for (name, values) in &context.common { 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}")); } } let allowed = common_fields.into_iter().chain(symbol_fields).collect::>(); let spec = spec.validate_with_context(&allowed)?; let dependencies = crate::factor_events::field_dependencies(spec.expression.as_ref().unwrap()); let mut rows = Vec::with_capacity(series.len()); for item in series { let mut fields = context.common.clone(); fields.extend(context.by_symbol.get(&item.symbol).cloned().unwrap_or_default()); if allowed.iter().any(|f| dependencies.contains(*f) && !fields.contains_key(*f)) { return Err(format!("research_context_missing: {}", item.symbol)); } for (name, values) in &fields { if values.len() != days.len() || values.iter().flatten().any(|v| !v.is_finite() || (name == "scope_percentile" && !(0.0..=1.0).contains(v)) || (matches!(name.as_str(), "scope_rank" | "scope_size") && *v < 1.0)) { return Err(format!("research_context_invalid: {} {name}", item.symbol)); } } let mut result = research_row(evaluate_with_context(&spec, days, item, &fields, numeric_output), item, isolate_data_errors)?; result["values"]["research_context_latest"] = json!(fields.iter().map(|(k,v)|(k,v.last().copied().flatten())).collect::>()); rows.push(result); } Ok(json!({"contract":CONTRACT,"context_contract":"fidc_research_event_context_v1","spec":spec, "required_history":spec.history_len(),"rows":rows,"read_only":true, "source_evidence_verified":false,"live_routing":false,"rule_backtest_supported":false})) } pub fn expression_specs(expression: &str) -> Result, String> { let mut specs = Vec::new(); for helper in ["pattern_signal", "pattern_score"] { for (index, _) in expression.match_indices(helper) { if index > 0 && expression[..index] .chars() .next_back() .is_some_and(|c| c.is_alphanumeric() || c == '_') { continue; } let rest = expression[index + helper.len()..].trim_start(); let Some(rest) = rest.strip_prefix('(') else { continue; }; let rest = rest.trim_start(); let mut stream = serde_json::Deserializer::from_str(rest).into_iter::(); let text = stream .next() .ok_or("missing pattern JSON")? .map_err(|e| e.to_string())?; if !rest[stream.byte_offset()..].trim_start().starts_with(')') { return Err("pattern helper takes one JSON string".into()); } let spec: PatternSpec = serde_json::from_str(&text).map_err(|e| e.to_string())?; specs.push(spec.validate()?); } } Ok(specs) } #[cfg(test)] mod tests { use super::*; #[test] fn research_isolates_missing_listing_day_without_weakening_execution() { let days = ["2026-06-11", "2026-06-12"].map(|d|d.parse::().unwrap()); let spec: PatternSpec = serde_json::from_value(json!({"template":"ma_below","parameters":{"ma_window":2}})).unwrap(); let make = |symbol: &str| -> PatternSeries { serde_json::from_value(json!({"symbol":symbol,"listed_at":"2026-06-11","bars":days.map(|d|json!({"date":d,"open":10.,"high":11.,"low":9.,"close":10.,"volume":100.,"adjustment_factor_backward1":1.,"paused":false,"source_path":"/controlled/source.parquet"}))})).unwrap() }; let complete=make("300395.SZ");let mut missing=make("920083.BJ");missing.bars.remove(0); let members=[complete.clone(),missing]; assert!(evaluate_batch(spec.clone(),&days,&members).unwrap_err().contains("missing_market_row")); let partial=evaluate_batch_with_policy(spec.clone(),&days,&members,true).unwrap(); assert_eq!(partial["rows"][0],json!(evaluate(&spec.validate().unwrap(),&days,&complete).unwrap())); assert!(partial["rows"][1]["matched"].is_null()); assert_eq!(partial["rows"][1]["data_issue"]["date"],"2026-06-11"); assert_eq!(partial["rows"][1]["data_issue"]["reason"],"missing_market_row"); let invalid:PatternSpec=serde_json::from_value(json!({"template":"not-a-template"})).unwrap(); assert!(evaluate_batch_with_policy(invalid,&days,&members,true).is_err()); } #[test] fn research_index_and_ranking_context_never_unlock_strategy_mapping() { let days=["2026-09-04","2026-09-07","2026-09-08"].map(|s|s.parse::().unwrap()); let spec:PatternSpec=serde_json::from_value(json!({"template":"expression","parameters":{"history_window":3}, "expression":{"kind":"operator","name":"CROSS_ABOVE","args":[{"kind":"field","name":"close"},{"kind":"field","name":"index_close"}]}})).unwrap(); assert!(spec.clone().validate().unwrap_err().contains("mapping_required")); let series:PatternSeries=serde_json::from_value(json!({"symbol":"TEST","bars":days.iter().zip([9.0,10.0,11.0]).map(|(d,c)|json!({"date":d,"open":c,"high":c,"low":c,"close":c,"volume":100.0,"adjustment_factor_backward1":1.0,"paused":false})).collect::>()})).unwrap(); let mut context=ResearchContext{common:BTreeMap::from([("index_close".into(),vec![Some(10.0);3])]),..Default::default()}; let result=evaluate_research_batch(spec.clone(),&days,&[series.clone()],&context,false).unwrap(); assert_eq!(result["rows"][0]["matched"],true); assert_eq!(result["source_evidence_verified"],false); assert_eq!(result["rule_backtest_supported"],false); context.common.get_mut("index_close").unwrap()[1]=None; assert!(evaluate_research_batch(spec.clone(),&days,&[series.clone()],&context,false).unwrap_err().contains("index_window_incomplete")); context.common=BTreeMap::from([("close".into(),vec![Some(10.0);3])]); assert!(evaluate_research_batch(spec,&days,&[series],&context,false).is_err()); } #[test] fn research_numeric_output_keeps_warmup_unknown_without_a_false_signal() { let days=["2026-09-04","2026-09-07","2026-09-08"].map(|s|s.parse::().unwrap()); let spec:PatternSpec=serde_json::from_value(json!({"template":"expression","parameters":{"history_window":3}, "expression":{"kind":"operator","name":"PCT_CHANGE","window":2,"args":[{"kind":"field","name":"close"}]}})).unwrap(); let series:PatternSeries=serde_json::from_value(json!({"symbol":"TEST","bars":days.iter().zip([10.0,10.5,11.0]).map(|(d,c)|json!({"date":d,"open":c,"high":c,"low":c,"close":c,"volume":100.0,"adjustment_factor_backward1":1.0,"paused":false})).collect::>()})).unwrap(); assert!(evaluate_batch(spec.clone(),&days,&[series.clone()]).is_err()); let result=evaluate_research_batch(spec,&days,&[series],&ResearchContext::default(),true).unwrap(); let values=&result["rows"][0]["values"]["expression"]["values"]; assert!(values[0].is_null() && values[1].is_null()); assert!((values[2].as_f64().unwrap()-0.1).abs()<1e-12); assert_eq!(result["rows"][0]["matched"],false); } #[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), prev_close: Some(p - 1.0), amount: Some(p * 100.0), upper_limit: None, no_limit: None, 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)); let vwap_spec = make( json!({"kind":"operator","name":"GT","args":[{"kind":"operator","name":"DIV","args":[{"kind":"field","name":"amount"},{"kind":"field","name":"volume"}]},{"kind":"field","name":"prev_close"}]}), ); assert!(evaluate(&vwap_spec, &days, &series).unwrap().matched); let mut missing_amount = series.clone(); missing_amount.bars[1].amount = None; assert!( evaluate(&vwap_spec, &days, &missing_amount) .unwrap_err() .contains("amount") ); let mut missing_previous=series.clone();missing_previous.bars[1].prev_close=None; assert!(evaluate(&vwap_spec,&days,&missing_previous).unwrap_err().contains("prev_close")); 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, execution_context: None, session_event: None, } .validate() .unwrap(); let days = (0..spec.history_len()) .map(|n| { NaiveDate::from_ymd_opt(2025, 1, 1).unwrap() + chrono::Duration::days(n as i64) }) .collect::>(); let bars = days .iter() .enumerate() .map(|(n, &date)| { let c = 10.0 + n as f64; PatternBar { date, open: Some(c), high: Some(c), low: Some(c), close: Some(c), volume: Some(1000.0), prev_close: Some(c - 1.0), amount: Some(c * 1000.0), upper_limit: None, no_limit: None, adjustment_factor_backward1: Some(1.0), paused: Some(false), source_path: Some("fixture.parquet".into()), } }) .collect(); ( spec, days, PatternSeries { symbol: "000001.SZ".into(), name: None, listed_at: Some(NaiveDate::from_ymd_opt(1991, 4, 3).unwrap()), bars, }, ) } #[test] fn daily_patterns_all_templates_and_score_absence() { for template in [ "strength", "breakout", "volume_spike", "shrink_breakout", "ma_below", "volume_down", ] { let (spec, days, series) = fixture(template); let result = evaluate(&spec, &days, &series).unwrap(); assert_eq!(result.matched, template == "strength"); assert_eq!(result.score.is_none(), template == "shrink_breakout"); } } #[test] fn daily_patterns_adjusts_all_prices_not_volume() { let (spec, days, series) = fixture("strength"); let a = evaluate(&spec, &days, &series).unwrap(); let mut split = series.clone(); for b in &mut split.bars { b.open = b.open.map(|p| p / 2.0); b.high = b.high.map(|p| p / 2.0); b.low = b.low.map(|p| p / 2.0); b.close = b.close.map(|p| p / 2.0); b.adjustment_factor_backward1 = Some(2.0); } let b = evaluate(&spec, &days, &split).unwrap(); assert_eq!(a.score, b.score); assert_eq!(a.checks, b.checks); } #[test] fn mean_volume_is_not_prior_max_and_excludes_current_bar() { let (spec,days,mut series)=fixture("mean_volume_spike"); for (bar,volume) in series.bars.iter_mut().zip([10.,10.,10.,10.,100.,100.]) {bar.volume=Some(volume);} assert!(evaluate(&spec,&days,&series).unwrap().matched); let mut old=spec.clone();old.template="volume_spike".into(); assert!(!evaluate(&old,&days,&series).unwrap().matched); assert_eq!(evaluate(&spec,&days,&series).unwrap().values["volume_ratio"],json!(100./28.)); } #[test] fn mean_volume_followup_requires_bullish_breakout_and_shrink() { let (spec,days,mut series)=fixture("mean_shrink_breakout"); series.bars[6].volume=Some(4000.); let last=series.bars.last_mut().unwrap();last.open=Some(19.);last.low=Some(19.); assert!(evaluate(&spec,&days,&series).unwrap().matched); series.bars.last_mut().unwrap().volume=Some(3000.); assert!(!evaluate(&spec,&days,&series).unwrap().matched); } #[test] fn breakout_retest_needs_a_later_retest_not_the_breakout_candle_itself() { let (spec,days,mut series)=fixture("breakout_retest"); for b in &mut series.bars {b.open=Some(10.);b.high=Some(10.);b.low=Some(10.);b.close=Some(10.);} let anchor=series.bars.len()-11; let b=&mut series.bars[anchor];b.open=Some(11.);b.high=Some(12.1);b.low=Some(9.9);b.close=Some(12.);b.volume=Some(2000.); for b in &mut series.bars[anchor+1..] {b.open=Some(10.3);b.high=Some(10.4);b.low=Some(10.3);b.close=Some(10.4);} assert!(!evaluate(&spec,&days,&series).unwrap().matched); series.bars[anchor+1].low=Some(9.95); assert!(evaluate(&spec,&days,&series).unwrap().matched); } #[test] fn limit_consolidation_requires_real_limit_and_never_infers_ten_percent() { let (spec,days,mut series)=fixture("limit_consolidation"); for (b,c) in series.bars.iter_mut().zip([10.,10.1,10.2,10.1,9.9]) {b.open=Some(c);b.high=Some(c);b.low=Some(c);b.close=Some(c);} assert!(evaluate(&spec,&days,&series).unwrap_err().contains("upper_limit")); series.bars[0].upper_limit=Some(10.); assert!(evaluate(&spec,&days,&series).unwrap().matched); series.bars[0].no_limit=Some(true); assert!(!evaluate(&spec,&days,&series).unwrap().matched); series.bars[0].no_limit=Some(false);series.bars[0].upper_limit=Some(0.); assert!(evaluate(&spec,&days,&series).is_err()); } #[test] fn daily_patterns_flat_decimal_prices_do_not_create_a_sell_signal() { let (mut spec, _, mut series) = fixture("strength"); spec.template = "ma_below".into(); spec.parameters = BTreeMap::from([("ma_window".into(), json!(60))]); for bar in &mut series.bars { bar.open = Some(10.1); bar.high = Some(10.1); bar.low = Some(10.1); bar.close = Some(10.1); } let days = series.bars.iter().map(|bar| bar.date).collect::>(); let result = evaluate(&spec, &days, &series).unwrap(); assert!( !result.matched, "unchanged decimal prices must not trigger a below-MA sell: {:?}", result.checks ); assert_eq!(result.values["ma"], 10.1); } #[test] fn daily_patterns_no_missing_data_fallback() { let (spec, days, mut series) = fixture("strength"); series.bars[0].adjustment_factor_backward1 = None; assert!( evaluate(&spec, &days, &series) .unwrap_err() .contains("adjustment_factor") ); series.bars[0].paused = Some(true); assert!(evaluate(&spec, &days, &series).is_err()); series.bars[0].adjustment_factor_backward1 = Some(1.0); let excluded = evaluate(&spec, &days, &series).unwrap(); assert!(excluded.exclusion.is_some()); assert!(!excluded.matched); series.bars.remove(0); assert!(evaluate(&spec, &days, &series).is_err()); series.listed_at = Some(days[1]); assert!(evaluate(&spec, &days, &series).unwrap().exclusion.is_some()); } #[test] fn daily_patterns_rejects_future_and_duplicate_bars() { let (spec, days, mut series) = fixture("strength"); series.bars.push(series.bars[0].clone()); assert!(evaluate(&spec, &days, &series).is_err()); series.bars.last_mut().unwrap().date = *days.last().unwrap() + chrono::Duration::days(1); assert!(evaluate(&spec, &days, &series).is_err()); } #[test] fn daily_patterns_helper_literal_preserves_parameters() { let text = serde_json::to_string(&json!({"template":"breakout","parameters":{"high_window":252}})) .unwrap(); let expression = format!("pattern_signal({})", serde_json::to_string(&text).unwrap()); assert_eq!(expression_specs(&expression).unwrap()[0].history_len(), 253); assert!(expression_specs("pattern_signal(\"{}\")").is_err()); } }