Files
fidc-backtest-engine/crates/fidc-signal-client/src/lib.rs
T

44 lines
2.4 KiB
Rust

//! 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<Arc<ValidatedSignalBook>,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)
}