feat: share verified signal books across backtest and trading clients
This commit is contained in:
@@ -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<StrategyDecision, BacktestError> {
|
||||
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
|
||||
|
||||
@@ -19,6 +19,8 @@ use crate::{
|
||||
pub struct StrategyRuntimeSpec {
|
||||
#[serde(default)]
|
||||
pub signal_book: Option<crate::signal_contract::SignalBook>,
|
||||
#[serde(default)]
|
||||
pub signal_book_ref: Option<crate::signal_contract::SignalBookReference>,
|
||||
#[serde(default, alias = "strategy_id")]
|
||||
pub strategy_id: Option<String>,
|
||||
#[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());
|
||||
}
|
||||
|
||||
@@ -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<String,Weak<ValidatedSignalBook>>,
|
||||
retained: std::collections::VecDeque<(String,Arc<ValidatedSignalBook>,usize)>,
|
||||
}
|
||||
|
||||
fn signal_cache() -> &'static Mutex<SignalCache> {
|
||||
static CACHE: OnceLock<Mutex<SignalCache>> = OnceLock::new();
|
||||
CACHE.get_or_init(||Mutex::new(SignalCache::default()))
|
||||
}
|
||||
|
||||
pub fn cached_signal_book(reference: &SignalBookReference) -> Result<Option<Arc<ValidatedSignalBook>>,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<Arc<ValidatedSignalBook>,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::<usize>()>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<Utc>,
|
||||
pub decision_at: DateTime<Utc>,
|
||||
pub input_as_of: DateTime<Utc>,
|
||||
pub input_available_at: DateTime<Utc>,
|
||||
@@ -64,7 +127,8 @@ pub struct SignalBook {
|
||||
pub schema: String,
|
||||
pub version_sha256: String,
|
||||
pub generator_sha256: String,
|
||||
pub knowledge_cutoff: DateTime<Utc>,
|
||||
pub model_sha256: Option<String>,
|
||||
pub knowledge_cutoff: Option<DateTime<Utc>>,
|
||||
pub provenance: SignalProvenance,
|
||||
pub frequency: SignalFrequency,
|
||||
pub expected_decisions: Vec<DateTime<Utc>>,
|
||||
@@ -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<Utc> = "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"));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user