44 lines
2.4 KiB
Rust
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)
|
|
}
|