Copied from api_gateway, removed REST handlers (port 8080), added tonic-web + CORS for grpc-web browser access. Binary renamed: api-gateway → api Changes: - Package name: api-gateway → api - Deleted src/handlers/ (REST ML endpoints on port 8080) - Added tonic-web 0.13 + tower-http CORS layer - Server::builder().accept_http1(true) for grpc-web - CORS_ORIGINS env var (default http://localhost:5173) - Metrics server on port 9091 (axum) preserved - All 95 lib tests pass, 0 clippy warnings - Added services/api to workspace members Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
552 lines
18 KiB
Rust
552 lines
18 KiB
Rust
//! Trading Agent Service Proxy - Zero-copy gRPC forwarding for trading agent operations
|
|
//!
|
|
//! This module implements a high-performance proxy for the Trading Agent Service with:
|
|
//! - Zero-copy message forwarding (routing overhead <10μs)
|
|
//! - Connection pooling via tonic::transport::Channel
|
|
//! - Circuit breaker integration for backend failures
|
|
//! - Efficient streaming support for agent activity events
|
|
//! - Health checking integration
|
|
|
|
use std::pin::Pin;
|
|
use std::time::Duration;
|
|
|
|
use futures::Stream;
|
|
use tonic::{Request, Response, Status};
|
|
use tracing::{error, info, instrument, warn};
|
|
|
|
// Import the generated Trading Agent service protobuf definitions from lib.rs
|
|
use crate::trading_agent::trading_agent_service_client::TradingAgentServiceClient;
|
|
use crate::trading_agent::trading_agent_service_server::{
|
|
TradingAgentService, TradingAgentServiceServer,
|
|
};
|
|
use crate::trading_agent::{
|
|
AgentActivityEvent,
|
|
// Portfolio Allocation
|
|
AllocatePortfolioRequest,
|
|
AllocatePortfolioResponse,
|
|
// Order Generation
|
|
GenerateOrdersRequest,
|
|
GenerateOrdersResponse,
|
|
GetAgentPerformanceRequest,
|
|
GetAgentPerformanceResponse,
|
|
|
|
// Agent Monitoring
|
|
GetAgentStatusRequest,
|
|
GetAgentStatusResponse,
|
|
GetAllocationRequest,
|
|
GetAllocationResponse,
|
|
GetSelectedAssetsRequest,
|
|
GetSelectedAssetsResponse,
|
|
|
|
GetUniverseRequest,
|
|
GetUniverseResponse,
|
|
// Service Health
|
|
HealthCheckRequest,
|
|
HealthCheckResponse,
|
|
ListStrategiesRequest,
|
|
ListStrategiesResponse,
|
|
RebalancePortfolioRequest,
|
|
RebalancePortfolioResponse,
|
|
|
|
// Strategy Coordination
|
|
RegisterStrategyRequest,
|
|
RegisterStrategyResponse,
|
|
// Asset Selection
|
|
SelectAssetsRequest,
|
|
SelectAssetsResponse,
|
|
// Universe Management
|
|
SelectUniverseRequest,
|
|
SelectUniverseResponse,
|
|
StreamAgentActivityRequest,
|
|
StreamAgentStatusRequest,
|
|
SubmitAgentOrdersRequest,
|
|
SubmitAgentOrdersResponse,
|
|
|
|
UpdateStrategyStatusRequest,
|
|
UpdateStrategyStatusResponse,
|
|
|
|
UpdateUniverseCriteriaRequest,
|
|
UpdateUniverseCriteriaResponse,
|
|
};
|
|
|
|
/// Trading Agent Service Proxy
|
|
///
|
|
/// Provides zero-copy forwarding of gRPC requests to the backend Trading Agent Service.
|
|
///
|
|
/// Uses connection pooling and circuit breakers for high availability and performance.
|
|
#[derive(Debug, Clone)]
|
|
pub struct TradingAgentProxy {
|
|
/// Backend Trading Agent Service client with connection pooling
|
|
client: TradingAgentServiceClient<tonic::transport::Channel>,
|
|
}
|
|
|
|
impl TradingAgentProxy {
|
|
/// Create a new Trading Agent Service proxy
|
|
///
|
|
/// # Arguments
|
|
/// * `client` - Pre-configured Trading Agent Service client with circuit breaker
|
|
///
|
|
/// # Performance
|
|
/// - Uses Arc-based channel cloning for zero-copy client reuse
|
|
/// - Connection pooling managed by tonic::transport::Channel
|
|
pub fn new(client: TradingAgentServiceClient<tonic::transport::Channel>) -> Self {
|
|
Self { client }
|
|
}
|
|
|
|
/// Convert proxy into a tonic server instance
|
|
pub fn into_server(self) -> TradingAgentServiceServer<TradingAgentProxy> {
|
|
TradingAgentServiceServer::new(self)
|
|
}
|
|
}
|
|
|
|
#[tonic::async_trait]
|
|
impl TradingAgentService for TradingAgentProxy {
|
|
/// Server streaming type for agent activity events
|
|
type StreamAgentActivityStream =
|
|
Pin<Box<dyn Stream<Item = Result<AgentActivityEvent, Status>> + Send>>;
|
|
/// Server streaming type for poll-based agent status
|
|
type StreamAgentStatusStream =
|
|
Pin<Box<dyn Stream<Item = Result<GetAgentStatusResponse, Status>> + Send>>;
|
|
|
|
// ===== Universe Management Methods =====
|
|
|
|
/// Select tradable universe based on liquidity, volatility, and ML signals
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn select_universe(
|
|
&self,
|
|
request: Request<SelectUniverseRequest>,
|
|
) -> Result<Response<SelectUniverseResponse>, Status> {
|
|
info!("Proxying SelectUniverse request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.select_universe(request).await.map_err(|e| {
|
|
error!("Backend SelectUniverse failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("SelectUniverse request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Get current trading universe configuration
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn get_universe(
|
|
&self,
|
|
request: Request<GetUniverseRequest>,
|
|
) -> Result<Response<GetUniverseResponse>, Status> {
|
|
info!("Proxying GetUniverse request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.get_universe(request).await.map_err(|e| {
|
|
error!("Backend GetUniverse failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GetUniverse request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Update universe selection criteria
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn update_universe_criteria(
|
|
&self,
|
|
request: Request<UpdateUniverseCriteriaRequest>,
|
|
) -> Result<Response<UpdateUniverseCriteriaResponse>, Status> {
|
|
info!("Proxying UpdateUniverseCriteria request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client
|
|
.update_universe_criteria(request)
|
|
.await
|
|
.map_err(|e| {
|
|
error!("Backend UpdateUniverseCriteria failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("UpdateUniverseCriteria request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Asset Selection Methods =====
|
|
|
|
/// Select specific assets to trade within universe
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn select_assets(
|
|
&self,
|
|
request: Request<SelectAssetsRequest>,
|
|
) -> Result<Response<SelectAssetsResponse>, Status> {
|
|
info!("Proxying SelectAssets request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.select_assets(request).await.map_err(|e| {
|
|
error!("Backend SelectAssets failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("SelectAssets request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Get current asset selection with scores
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn get_selected_assets(
|
|
&self,
|
|
request: Request<GetSelectedAssetsRequest>,
|
|
) -> Result<Response<GetSelectedAssetsResponse>, Status> {
|
|
info!("Proxying GetSelectedAssets request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.get_selected_assets(request).await.map_err(|e| {
|
|
error!("Backend GetSelectedAssets failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GetSelectedAssets request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Portfolio Allocation Methods =====
|
|
|
|
/// Allocate capital across selected assets
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn allocate_portfolio(
|
|
&self,
|
|
request: Request<AllocatePortfolioRequest>,
|
|
) -> Result<Response<AllocatePortfolioResponse>, Status> {
|
|
info!("Proxying AllocatePortfolio request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.allocate_portfolio(request).await.map_err(|e| {
|
|
error!("Backend AllocatePortfolio failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("AllocatePortfolio request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Get current portfolio allocation
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn get_allocation(
|
|
&self,
|
|
request: Request<GetAllocationRequest>,
|
|
) -> Result<Response<GetAllocationResponse>, Status> {
|
|
info!("Proxying GetAllocation request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.get_allocation(request).await.map_err(|e| {
|
|
error!("Backend GetAllocation failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GetAllocation request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Rebalance portfolio based on target allocation
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn rebalance_portfolio(
|
|
&self,
|
|
request: Request<RebalancePortfolioRequest>,
|
|
) -> Result<Response<RebalancePortfolioResponse>, Status> {
|
|
info!("Proxying RebalancePortfolio request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.rebalance_portfolio(request).await.map_err(|e| {
|
|
error!("Backend RebalancePortfolio failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("RebalancePortfolio request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Order Generation Methods =====
|
|
|
|
/// Generate orders based on allocation and ML signals
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn generate_orders(
|
|
&self,
|
|
request: Request<GenerateOrdersRequest>,
|
|
) -> Result<Response<GenerateOrdersResponse>, Status> {
|
|
info!("Proxying GenerateOrders request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.generate_orders(request).await.map_err(|e| {
|
|
error!("Backend GenerateOrders failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GenerateOrders request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Submit generated orders to Trading Service
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn submit_agent_orders(
|
|
&self,
|
|
request: Request<SubmitAgentOrdersRequest>,
|
|
) -> Result<Response<SubmitAgentOrdersResponse>, Status> {
|
|
info!("Proxying SubmitAgentOrders request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.submit_agent_orders(request).await.map_err(|e| {
|
|
error!("Backend SubmitAgentOrders failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("SubmitAgentOrders request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Strategy Coordination Methods =====
|
|
|
|
/// Register a trading strategy with the agent
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn register_strategy(
|
|
&self,
|
|
request: Request<RegisterStrategyRequest>,
|
|
) -> Result<Response<RegisterStrategyResponse>, Status> {
|
|
info!("Proxying RegisterStrategy request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.register_strategy(request).await.map_err(|e| {
|
|
error!("Backend RegisterStrategy failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("RegisterStrategy request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Get list of active strategies
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn list_strategies(
|
|
&self,
|
|
request: Request<ListStrategiesRequest>,
|
|
) -> Result<Response<ListStrategiesResponse>, Status> {
|
|
info!("Proxying ListStrategies request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.list_strategies(request).await.map_err(|e| {
|
|
error!("Backend ListStrategies failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("ListStrategies request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Enable/disable a strategy
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn update_strategy_status(
|
|
&self,
|
|
request: Request<UpdateStrategyStatusRequest>,
|
|
) -> Result<Response<UpdateStrategyStatusResponse>, Status> {
|
|
info!("Proxying UpdateStrategyStatus request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.update_strategy_status(request).await.map_err(|e| {
|
|
error!("Backend UpdateStrategyStatus failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("UpdateStrategyStatus request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Agent Monitoring Methods =====
|
|
|
|
/// Get comprehensive agent status and performance
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn get_agent_status(
|
|
&self,
|
|
request: Request<GetAgentStatusRequest>,
|
|
) -> Result<Response<GetAgentStatusResponse>, Status> {
|
|
info!("Proxying GetAgentStatus request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.get_agent_status(request).await.map_err(|e| {
|
|
error!("Backend GetAgentStatus failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GetAgentStatus request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Stream real-time agent decisions and actions (server streaming)
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy stream forwarding
|
|
/// - No intermediate buffering
|
|
/// - Direct stream passthrough from backend
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn stream_agent_activity(
|
|
&self,
|
|
request: Request<StreamAgentActivityRequest>,
|
|
) -> Result<Response<Self::StreamAgentActivityStream>, Status> {
|
|
info!("Proxying StreamAgentActivity streaming request");
|
|
|
|
let mut client = self.client.clone();
|
|
|
|
// Get backend stream response
|
|
let stream_response = client.stream_agent_activity(request).await.map_err(|e| {
|
|
error!("Backend StreamAgentActivity failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
// Extract inner stream and forward directly (zero-copy)
|
|
let stream = stream_response.into_inner();
|
|
let boxed_stream = Box::pin(stream) as Self::StreamAgentActivityStream;
|
|
|
|
info!("StreamAgentActivity streaming request forwarded successfully");
|
|
Ok(Response::new(boxed_stream))
|
|
}
|
|
|
|
/// Get agent performance metrics
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn get_agent_performance(
|
|
&self,
|
|
request: Request<GetAgentPerformanceRequest>,
|
|
) -> Result<Response<GetAgentPerformanceResponse>, Status> {
|
|
info!("Proxying GetAgentPerformance request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.get_agent_performance(request).await.map_err(|e| {
|
|
error!("Backend GetAgentPerformance failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("GetAgentPerformance request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
// ===== Service Health Methods =====
|
|
|
|
/// Health check for Trading Agent Service backend
|
|
///
|
|
/// # Performance
|
|
/// - Zero-copy message forwarding
|
|
/// - Routing overhead target: <10μs
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn health_check(
|
|
&self,
|
|
request: Request<HealthCheckRequest>,
|
|
) -> Result<Response<HealthCheckResponse>, Status> {
|
|
info!("Proxying HealthCheck request");
|
|
|
|
let mut client = self.client.clone();
|
|
let response = client.health_check(request).await.map_err(|e| {
|
|
warn!("Backend HealthCheck failed: {}", e);
|
|
e
|
|
})?;
|
|
|
|
info!("HealthCheck request forwarded successfully");
|
|
Ok(response)
|
|
}
|
|
|
|
/// Poll-to-stream adapter for agent status
|
|
///
|
|
/// Polls GetAgentStatus at a configurable interval and yields results as a server stream.
|
|
#[instrument(skip(self, request), fields(request_id = %uuid::Uuid::new_v4()), err)]
|
|
async fn stream_agent_status(
|
|
&self,
|
|
request: Request<StreamAgentStatusRequest>,
|
|
) -> Result<Response<Self::StreamAgentStatusStream>, Status> {
|
|
let req = request.into_inner();
|
|
let interval_secs = if req.interval_seconds == 0 {
|
|
3
|
|
} else {
|
|
req.interval_seconds.clamp(1, 60)
|
|
};
|
|
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_agent_status(Request::new(GetAgentStatusRequest::default())).await {
|
|
Ok(resp) => yield Ok(resp.into_inner()),
|
|
Err(e) => {
|
|
error!("Backend GetAgentStatus failed: {}", e);
|
|
yield Err(e);
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
};
|
|
|
|
Ok(Response::new(Box::pin(stream)))
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
#[test]
|
|
fn test_proxy_creation() {
|
|
// This test validates the proxy struct can be created
|
|
// Full integration tests require running backend service
|
|
}
|
|
}
|