diff --git a/Cargo.toml b/Cargo.toml index 4593a77..3985559 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,7 @@ [workspace] members = [ "crates/fidc-core", + "crates/fidc-signal-client", ] resolver = "2" diff --git a/crates/fidc-core/src/platform_expr_strategy.rs b/crates/fidc-core/src/platform_expr_strategy.rs index a914719..d7a234e 100644 --- a/crates/fidc-core/src/platform_expr_strategy.rs +++ b/crates/fidc-core/src/platform_expr_strategy.rs @@ -10100,7 +10100,7 @@ impl PlatformExprStrategy { buy_denials: Default::default(), rebalance: false, target_weights: BTreeMap::new(), - exit_symbols: BTreeSet::new(), + exit_symbols: if self.config.signal_book.is_some() { exit_symbols } else { BTreeSet::new() }, order_intents, notes: Vec::new(), diagnostics, @@ -12437,6 +12437,9 @@ impl Strategy for PlatformExprStrategy { impl PlatformExprStrategy { fn attach_buy_denials(&self, ctx: &StrategyContext<'_>, decision: &mut StrategyDecision) -> Result<(), BacktestError> { + if self.config.signal_book.is_none() && self.config.buy_filter_expr.trim().is_empty() { + return Ok(()); + } let symbols = decision.potential_buy_symbols(ctx.open_orders); if symbols.is_empty() { return Ok(()); @@ -12498,6 +12501,9 @@ impl PlatformExprStrategy { } fn compute_day_decision(&mut self, ctx: &StrategyContext<'_>) -> Result { + if self.config.signal_book.is_some() && self.config.explicit_action_schedule.is_some() { + return Ok(StrategyDecision::default()); + } if self.config.rotation_enabled && self .config diff --git a/crates/fidc-core/src/platform_strategy_spec.rs b/crates/fidc-core/src/platform_strategy_spec.rs index ecc75ad..24c824c 100644 --- a/crates/fidc-core/src/platform_strategy_spec.rs +++ b/crates/fidc-core/src/platform_strategy_spec.rs @@ -19,6 +19,8 @@ use crate::{ pub struct StrategyRuntimeSpec { #[serde(default)] pub signal_book: Option, + #[serde(default)] + pub signal_book_ref: Option, #[serde(default, alias = "strategy_id")] pub strategy_id: Option, #[serde(default)] @@ -2601,8 +2603,13 @@ pub fn platform_expr_config_from_spec( } cfg.strict_value_budget = true; - if let Some(raw) = &spec.signal_book { - let book = raw.clone().validate()?; + let signal_book = match (&spec.signal_book,&spec.signal_book_ref) { + (Some(_),Some(_)) => return Err("inline_and_registered_signal_book_are_mutually_exclusive".into()), + (Some(raw),None) => Some(std::sync::Arc::new(raw.clone().validate()?)), + (None,Some(reference)) => crate::signal_contract::cached_signal_book(reference)?, + (None,None) => None, + }; + if let Some(book) = signal_book { if cfg.explicit_actions.len() != 1 || !matches!(cfg.explicit_actions[0], PlatformTradeAction::ConsumeSignal) { return Err("signal_book_requires_one_consume_signal_action".into()); } @@ -2612,7 +2619,12 @@ pub fn platform_expr_config_from_spec( cfg.rotation_enabled = false; cfg.signal_rebalance_dates = book.decision_dates(); cfg.initial_subscriptions.extend(book.symbols()); - cfg.signal_book = Some(std::sync::Arc::new(book)); + cfg.signal_book = Some(book); + } else if spec.signal_book_ref.is_some() { + if cfg.explicit_actions.len()!=1 || !matches!(cfg.explicit_actions[0],PlatformTradeAction::ConsumeSignal) { + return Err("signal_book_requires_one_consume_signal_action".into()); + } + cfg.rotation_enabled=false; } else if cfg.explicit_actions.iter().any(|action| matches!(action, PlatformTradeAction::ConsumeSignal)) { return Err("consume_signal_requires_verified_signal_book".into()); } diff --git a/crates/fidc-core/src/signal_contract.rs b/crates/fidc-core/src/signal_contract.rs index 5784018..f1dd256 100644 --- a/crates/fidc-core/src/signal_contract.rs +++ b/crates/fidc-core/src/signal_contract.rs @@ -2,6 +2,7 @@ //! prices are intentionally absent; the existing broker owns those decisions. use std::collections::{BTreeMap, BTreeSet}; +use std::sync::{Arc, Mutex, OnceLock, Weak}; use chrono::{DateTime, FixedOffset, NaiveDate, NaiveDateTime, NaiveTime, Utc}; use serde::{Deserialize, Serialize}; @@ -11,6 +12,67 @@ use crate::portfolio::PortfolioState; pub const SIGNAL_BOOK_SCHEMA: &str = "fidc.signal-book/v1"; +#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub struct SignalBookReference { + pub book_id: String, + pub version_sha256: String, + pub artifact_sha256: String, +} + +impl SignalBookReference { + pub fn validate(&self) -> Result<(), String> { + if !valid_sha(&self.version_sha256) || !valid_sha(&self.artifact_sha256) + || self.book_id != format!("signal_book_{}",self.version_sha256) + { return Err("signal_book_reference_invalid".into()); } + Ok(()) + } +} + +#[derive(Default)] +struct SignalCache { + entries: BTreeMap>, + retained: std::collections::VecDeque<(String,Arc,usize)>, +} + +fn signal_cache() -> &'static Mutex { + static CACHE: OnceLock> = OnceLock::new(); + CACHE.get_or_init(||Mutex::new(SignalCache::default())) +} + +pub fn cached_signal_book(reference: &SignalBookReference) -> Result>,String> { + reference.validate()?; + let cache=signal_cache().lock().map_err(|_|"signal_cache_lock_failed")?; + let book=cache.entries.get(&reference.artifact_sha256).and_then(Weak::upgrade); + if book.as_ref().is_some_and(|book|book.version_sha256()!=reference.version_sha256) { + return Err("signal_book_cached_version_mismatch".into()); + } + Ok(book) +} + +pub fn register_signal_book(reference: &SignalBookReference, body: &[u8]) -> Result,String> { + use sha2::{Digest,Sha256}; + reference.validate()?; + if body.len()>64*1024*1024 || format!("{:x}",Sha256::digest(body))!=reference.artifact_sha256 { + return Err("signal_book_artifact_hash_or_size_invalid".into()); + } + let raw:SignalBook=serde_json::from_slice(body).map_err(|error|format!("signal_book_decode_failed: {error}"))?; + if raw.version_sha256!=reference.version_sha256 { return Err("signal_book_version_mismatch".into()); } + let book=Arc::new(raw.validate()?); + let mut cache=signal_cache().lock().map_err(|_|"signal_cache_lock_failed")?; + cache.entries.retain(|_,value|value.strong_count()>0); + if let Some(existing)=cache.entries.get(&reference.artifact_sha256).and_then(Weak::upgrade) { return Ok(existing); } + cache.entries.insert(reference.artifact_sha256.clone(),Arc::downgrade(&book)); + let estimated=body.len().saturating_mul(4); + if estimated<=128*1024*1024 { + cache.retained.push_back((reference.artifact_sha256.clone(),book.clone(),estimated)); + while cache.retained.len()>4 || cache.retained.iter().map(|entry|entry.2).sum::()>128*1024*1024 { + cache.retained.pop_front(); + } + } + Ok(book) +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize, Serialize)] #[serde(rename_all = "snake_case")] pub enum SignalProvenance { @@ -48,6 +110,7 @@ impl SignalAction { #[derive(Debug, Clone, Deserialize, Serialize)] #[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct SignalSnapshot { + pub signal_at: DateTime, pub decision_at: DateTime, pub input_as_of: DateTime, pub input_available_at: DateTime, @@ -64,7 +127,8 @@ pub struct SignalBook { pub schema: String, pub version_sha256: String, pub generator_sha256: String, - pub knowledge_cutoff: DateTime, + pub model_sha256: Option, + pub knowledge_cutoff: Option>, pub provenance: SignalProvenance, pub frequency: SignalFrequency, pub expected_decisions: Vec>, @@ -92,6 +156,9 @@ impl SignalBook { { return Err("signal_book_identity_invalid".into()); } + if self.model_sha256.as_ref().is_some_and(|value| !valid_sha(value)) + || self.model_sha256.is_some() != self.knowledge_cutoff.is_some() + { return Err("signal_model_training_identity_incomplete".into()); } if self.expected_decisions.is_empty() || self.expected_decisions.len() > 100_000 || self.expected_decisions.len() != self.snapshots.len() { @@ -105,11 +172,11 @@ impl SignalBook { return Err("signal_book_decisions_duplicate_or_unordered".into()); } previous = Some(*expected); - if self.knowledge_cutoff >= *expected || snapshot.input_as_of > *expected - || snapshot.input_available_at > *expected || snapshot.input_as_of > snapshot.input_available_at + if self.knowledge_cutoff.is_some_and(|cutoff| cutoff > snapshot.signal_at) || snapshot.signal_at > *expected + || snapshot.input_available_at > snapshot.signal_at || snapshot.input_as_of > snapshot.input_available_at || snapshot.published_at < snapshot.generated_at || !valid_sha(&snapshot.input_sha256) || snapshot.generated_at < snapshot.input_available_at - || snapshot.generated_at < self.knowledge_cutoff + || self.knowledge_cutoff.is_some_and(|cutoff| snapshot.generated_at < cutoff) { return Err("signal_book_future_or_invalid_input".into()); } @@ -179,7 +246,9 @@ impl ValidatedSignalBook { pub fn snapshot_for(&self, ctx: &StrategyContext<'_>) -> Result<&SignalSnapshot, String> { let snapshot = self.snapshot_at(ctx.execution_date, ctx.current_time(), ctx.is_lagged_execution())?; - if ctx.is_lagged_execution() && shanghai(snapshot.input_as_of) > ctx.decision_date.and_hms_opt(15,0,0).expect("completed decision session") { + let logical_clock=ctx.current_datetime().filter(|at|at.date()==ctx.decision_date) + .unwrap_or(ctx.decision_date.and_hms_opt(15,0,0).expect("completed decision session")); + if shanghai(snapshot.signal_at)>logical_clock || (ctx.is_lagged_execution() && shanghai(snapshot.input_as_of).date()>ctx.decision_date) { return Err("next_open_signal_contains_execution_session_inputs".into()); } Ok(snapshot) @@ -262,9 +331,11 @@ mod tests { let source: DateTime = "2025-01-06T15:00:00+08:00".parse().unwrap(); SignalBook { schema: SIGNAL_BOOK_SCHEMA.into(), version_sha256: "a".repeat(64), generator_sha256: "b".repeat(64), - knowledge_cutoff: "2024-12-31T15:00:00+08:00".parse().unwrap(), + model_sha256: Some("d".repeat(64)), + knowledge_cutoff: Some("2024-12-31T15:00:00+08:00".parse().unwrap()), provenance: SignalProvenance::Reconstructed, frequency: SignalFrequency::Daily, expected_decisions: vec![decision], snapshots: vec![SignalSnapshot { + signal_at: source, decision_at: decision, input_as_of: source, input_available_at: source, generated_at: decision + Duration::days(10), published_at: decision + Duration::days(10), input_sha256: "c".repeat(64), complete_targets: true, @@ -293,7 +364,7 @@ mod tests { match field { 0 => value.snapshots[0].input_as_of = future, 1 => value.snapshots[0].input_available_at = future, - _ => value.knowledge_cutoff = future, + _ => value.knowledge_cutoff = Some(future), } assert!(value.validate().unwrap_err().contains("future")); } diff --git a/crates/fidc-signal-client/Cargo.toml b/crates/fidc-signal-client/Cargo.toml new file mode 100644 index 0000000..d91515c --- /dev/null +++ b/crates/fidc-signal-client/Cargo.toml @@ -0,0 +1,10 @@ +[package] +name = "fidc-signal-client" +version.workspace = true +edition.workspace = true +license.workspace = true + +[dependencies] +fidc-core = { path = "../fidc-core" } +reqwest.workspace = true +serde_json.workspace = true diff --git a/crates/fidc-signal-client/src/lib.rs b/crates/fidc-signal-client/src/lib.rs new file mode 100644 index 0000000..b98279b --- /dev/null +++ b/crates/fidc-signal-client/src/lib.rs @@ -0,0 +1,43 @@ +//! Shared signal transport for FIDC backtest and trading services. + +use std::sync::Arc; +use fidc_core::signal_contract::{SignalBookReference,ValidatedSignalBook,cached_signal_book,register_signal_book}; +use reqwest::Client; +use serde_json::{Value,json}; + +#[derive(Clone,Copy)] +pub enum Purpose { Backtest, Online } + +pub async fn load(client:&Client, source_url:&str, token:&str, reference:&SignalBookReference, purpose:Purpose) + -> Result,String> +{ + reference.validate()?; + if token.len()<32 {return Err("signal_service_auth_not_configured".into());} + let purpose_name=match purpose {Purpose::Backtest=>"backtest",Purpose::Online=>"online"}; + let payload=json!({"reference":reference,"purpose":purpose_name}); + let root=format!("{}/api/strategy-signals/internal",source_url.trim_end_matches('/')); + // Registration/purpose validation always precedes a process-cache hit. + let response=client.post(format!("{root}/validate")) + .header("X-FIDC-Lifecycle-Token",token).json(&payload).send().await + .map_err(|_|"signal_validation_service_unavailable")?; + if !response.status().is_success() {return Err(format!("signal_validation_rejected_http_{}",response.status()));} + let validation:Value=response.json().await.map_err(|_|"signal_validation_response_invalid")?; + if validation.get("ok")!=Some(&Value::Bool(true)) || validation.get("reference")!=Some(&json!(reference)) { + return Err("signal_validation_identity_mismatch".into()); + } + let book=if let Some(book)=cached_signal_book(reference)? {book} else { + let mut response=client.post(format!("{root}/book")) + .header("X-FIDC-Lifecycle-Token",token).json(&payload).send().await + .map_err(|_|"signal_book_service_unavailable")?; + if !response.status().is_success() {return Err(format!("signal_book_rejected_http_{}",response.status()));} + if response.content_length().is_some_and(|bytes|bytes>64*1024*1024) {return Err("signal_book_transport_size_exceeded".into());} + let mut bytes=Vec::new(); + while let Some(chunk)=response.chunk().await.map_err(|_|"signal_book_transport_incomplete")? { + if bytes.len().saturating_add(chunk.len())>64*1024*1024 {return Err("signal_book_transport_size_exceeded".into());} + bytes.extend_from_slice(&chunk); + } + register_signal_book(reference,&bytes)? + }; + if matches!(purpose,Purpose::Online) {book.require_observed()?;} + Ok(book) +}