//! Risk Service Proxy - Zero-copy gRPC forwarding for risk.RiskService //! //! Forwards to trading-service backend (risk runs in the same process). use futures::Stream; use std::pin::Pin; use std::time::Duration; use tonic::{Request, Response, Status}; use tracing::{error, instrument}; use crate::risk::risk_service_client::RiskServiceClient; use crate::risk::risk_service_server::RiskService; use crate::risk::{ EmergencyStopRequest, EmergencyStopResponse, GetCircuitBreakerStatusRequest, GetCircuitBreakerStatusResponse, GetPositionRiskRequest, GetPositionRiskResponse, GetRiskMetricsRequest, GetRiskMetricsResponse, GetVaRRequest, GetVaRResponse, RiskAlertEvent, StreamCircuitBreakerStatusRequest, StreamRiskAlertsRequest, StreamRiskMetricsRequest, StreamVaRRequest, VaREvent, ValidateOrderRequest, ValidateOrderResponse, }; #[derive(Debug, Clone)] pub struct RiskServiceProxy { client: RiskServiceClient, } impl RiskServiceProxy { pub fn new(client: RiskServiceClient) -> Self { Self { client } } } #[tonic::async_trait] impl RiskService for RiskServiceProxy { type StreamVaRUpdatesStream = Pin> + Send>>; type StreamRiskAlertsStream = Pin> + Send>>; type StreamCircuitBreakerStatusStream = Pin> + Send>>; type StreamRiskMetricsStream = Pin> + Send>>; #[instrument(skip(self, request), err)] async fn get_va_r( &self, request: Request, ) -> Result, Status> { self.client.clone().get_va_r(request).await.map_err(|e| { error!("Backend GetVaR failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn stream_va_r_updates( &self, request: Request, ) -> Result, Status> { let stream = self .client .clone() .stream_va_r_updates(request) .await .map_err(|e| { error!("Backend StreamVaRUpdates failed: {}", e); e })?; Ok(Response::new(Box::pin(stream.into_inner()))) } #[instrument(skip(self, request), err)] async fn get_position_risk( &self, request: Request, ) -> Result, Status> { self.client .clone() .get_position_risk(request) .await .map_err(|e| { error!("Backend GetPositionRisk failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn validate_order( &self, request: Request, ) -> Result, Status> { self.client .clone() .validate_order(request) .await .map_err(|e| { error!("Backend ValidateOrder failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn get_risk_metrics( &self, request: Request, ) -> Result, Status> { self.client .clone() .get_risk_metrics(request) .await .map_err(|e| { error!("Backend GetRiskMetrics failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn stream_risk_alerts( &self, request: Request, ) -> Result, Status> { let stream = self .client .clone() .stream_risk_alerts(request) .await .map_err(|e| { error!("Backend StreamRiskAlerts failed: {}", e); e })?; Ok(Response::new(Box::pin(stream.into_inner()))) } #[instrument(skip(self, request), err)] async fn emergency_stop( &self, request: Request, ) -> Result, Status> { self.client .clone() .emergency_stop(request) .await .map_err(|e| { error!("Backend EmergencyStop failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn get_circuit_breaker_status( &self, request: Request, ) -> Result, Status> { self.client .clone() .get_circuit_breaker_status(request) .await .map_err(|e| { error!("Backend GetCircuitBreakerStatus failed: {}", e); e }) } #[instrument(skip(self, request), err)] async fn stream_circuit_breaker_status( &self, request: Request, ) -> Result, Status> { let req = request.into_inner(); let interval_secs = if req.interval_seconds == 0 { 2 } else { req.interval_seconds.clamp(1, 60) }; let symbol = req.symbol; let mut client = self.client.clone(); let stream = async_stream::stream! { let mut interval = tokio::time::interval(Duration::from_secs(u64::from(interval_secs))); loop { interval.tick().await; match client.get_circuit_breaker_status(Request::new(GetCircuitBreakerStatusRequest { symbol: symbol.clone(), })).await { Ok(resp) => yield Ok(resp.into_inner()), Err(e) => { error!("Backend GetCircuitBreakerStatus failed: {}", e); yield Err(e); break; } } } }; Ok(Response::new(Box::pin(stream))) } #[instrument(skip(self, request), err)] async fn stream_risk_metrics( &self, request: Request, ) -> Result, Status> { let req = request.into_inner(); let interval_secs = if req.interval_seconds == 0 { 3 } else { req.interval_seconds.clamp(1, 60) }; let portfolio_id = req.portfolio_id; let mut client = self.client.clone(); let stream = async_stream::stream! { let mut interval = tokio::time::interval(Duration::from_secs(u64::from(interval_secs))); loop { interval.tick().await; match client.get_risk_metrics(Request::new(GetRiskMetricsRequest { portfolio_id: portfolio_id.clone(), })).await { Ok(resp) => yield Ok(resp.into_inner()), Err(e) => { error!("Backend GetRiskMetrics failed: {}", e); yield Err(e); break; } } } }; Ok(Response::new(Box::pin(stream))) } }