pub mod result;
use crate::ipc_types::{IpcHttpRequest, IpcHttpResponse};
use crate::communication::ipc_port_negotiation::PortNegotiationManager;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use tokio::sync::broadcast;
use tokio::sync::Mutex;
use tracing::{debug, error, info, warn, trace, Level};
use once_cell::sync::Lazy;
use uuid::Uuid;
pub use self::result::*;
static HTTP_REQUEST_CHANNEL: Lazy<(broadcast::Sender<IpcHttpRequest>, Mutex<Option<broadcast::Receiver<IpcHttpRequest>>>)> =
Lazy::new(|| {
let (tx, rx) = broadcast::channel(256); (tx, Mutex::new(Some(rx)))
});
static HTTP_RESPONSE_CHANNEL: Lazy<(broadcast::Sender<IpcHttpResponse>, Mutex<Option<broadcast::Receiver<IpcHttpResponse>>>)> =
Lazy::new(|| {
let (tx, rx) = broadcast::channel(256); (tx, Mutex::new(Some(rx)))
});
#[derive(Debug, Default)]
pub struct HttpIpcMetrics {
pub requests_received: std::sync::atomic::AtomicUsize,
pub responses_sent: std::sync::atomic::AtomicUsize,
pub errors_encountered: std::sync::atomic::AtomicUsize,
pub avg_response_time_ms: std::sync::atomic::AtomicUsize,
}
static HTTP_IPC_METRICS: Lazy<HttpIpcMetrics> = Lazy::new(|| HttpIpcMetrics::default());
pub fn metrics() -> &'static HttpIpcMetrics {
&*HTTP_IPC_METRICS
}
pub fn subscribe_http_requests() -> broadcast::Receiver<IpcHttpRequest> {
let receiver = HTTP_REQUEST_CHANNEL.0.subscribe();
debug!("New subscriber for HTTP IPC requests (total receivers: {})",
HTTP_REQUEST_CHANNEL.0.receiver_count());
receiver
}
pub async fn send_http_response(response: IpcHttpResponse) -> std::result::Result<(), crate::Error> {
debug!("Sending HTTP response - request_id: {}, status: {}",
response.request_id, response.status_code);
let body_size = response.body.as_ref().map_or(0, |b| b.len());
trace!("Response body size: {} bytes", body_size);
if Level::TRACE <= tracing::level_filters::LevelFilter::current() {
for (name, value) in &response.headers {
trace!("Response header: {}: {}", name, value);
}
}
let mut retry_count = 0;
let max_retries = 2;
let mut last_error = None;
while retry_count <= max_retries {
match HTTP_RESPONSE_CHANNEL.0.send(response.clone()) {
Ok(_) => {
HTTP_IPC_METRICS.responses_sent.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
debug!("HTTP response sent successfully (request_id: {})", response.request_id);
return Ok(());
}
Err(e) => {
warn!("Failed to send HTTP response (attempt {}/{}): {}",
retry_count + 1, max_retries + 1, e);
if retry_count == max_retries {
error!("Failed to send HTTP response after {} attempts: {}",
max_retries + 1, e);
last_error = Some(e);
break;
}
tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
retry_count += 1;
}
}
}
HTTP_IPC_METRICS.errors_encountered.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
if let Some(error) = last_error {
Err(crate::Error::Config(crate::error::ConfigError::Invalid(
format!("Failed to send HTTP response after {} attempts: {}",
max_retries + 1, error)
)))
} else {
Err(crate::Error::Config(crate::error::ConfigError::Invalid(
"Failed to send HTTP response after all attempts".to_string()
)))
}
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ApiResponse<T> {
pub status: String,
pub data: Option<T>,
pub message: Option<String>,
}
pub type HandlerFn<S> = Arc<
dyn Fn(
IpcHttpRequest,
Arc<S>,
) -> Pin<Box<dyn Future<Output = HttpResult<IpcHttpResponse>> + Send>>
+ Send
+ Sync,
>;
#[derive(Clone)]
struct Route<S> {
method: String,
path: String,
handler: HandlerFn<S>,
}
#[derive(Clone)]
pub struct HttpIpcRouter<S: Send + Sync + Clone + 'static> {
routes: Arc<Vec<Route<S>>>,
not_found_handler: Option<HandlerFn<S>>,
middleware: Vec<Arc<dyn Fn(IpcHttpRequest) -> Pin<Box<dyn Future<Output = IpcHttpRequest> + Send>> + Send + Sync>>,
}
impl<S: Send + Sync + Clone + 'static> Default for HttpIpcRouter<S> {
fn default() -> Self {
Self::new()
}
}
impl<S: Send + Sync + Clone + 'static> HttpIpcRouter<S> {
pub fn new() -> Self {
Self {
routes: Arc::new(Vec::new()),
not_found_handler: None,
middleware: Vec::new(),
}
}
pub fn route<F, Fut>(mut self, method: &str, path: &str, handler: F) -> Self
where
F: Fn(IpcHttpRequest, Arc<S>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = HttpResult<IpcHttpResponse>> + Send + 'static,
{
let handler = Arc::new(move |req, state| {
let fut = handler(req, state);
Box::pin(fut) as Pin<Box<dyn Future<Output = HttpResult<IpcHttpResponse>> + Send>>
});
let mut routes = (*self.routes).clone();
routes.push(Route {
method: method.to_uppercase(),
path: path.to_string(),
handler,
});
self.routes = Arc::new(routes);
debug!("Added route: {} {}", method.to_uppercase(), path);
self
}
pub fn middleware<F, Fut>(mut self, middleware_fn: F) -> Self
where
F: Fn(IpcHttpRequest) -> Fut + Send + Sync + 'static,
Fut: Future<Output = IpcHttpRequest> + Send + 'static,
{
let middleware = Arc::new(move |req: IpcHttpRequest| {
let fut = middleware_fn(req);
Box::pin(fut) as Pin<Box<dyn Future<Output = IpcHttpRequest> + Send>>
});
self.middleware.push(middleware);
debug!("Added middleware to router");
self
}
pub fn not_found_handler<F, Fut>(mut self, handler: F) -> Self
where
F: Fn(IpcHttpRequest, Arc<S>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = HttpResult<IpcHttpResponse>> + Send + 'static,
{
let handler = Arc::new(move |req, state| {
let fut = handler(req, state);
Box::pin(fut) as Pin<Box<dyn Future<Output = HttpResult<IpcHttpResponse>> + Send>>
});
self.not_found_handler = Some(handler);
debug!("Custom not-found handler set");
self
}
pub async fn handle_request(&self, mut request: IpcHttpRequest, state: Arc<S>) -> IpcHttpResponse {
let start_time = std::time::Instant::now();
info!(
"HTTP_IPC_ROUTER: Processing request - ID: {}, Method: {}, URI: {}",
request.request_id, request.method, request.uri
);
for middleware in &self.middleware {
request = middleware(request).await;
}
let method = request.method.to_uppercase();
let path = request
.uri
.split('?')
.next()
.unwrap_or(&request.uri)
.to_string();
debug!("Routing request: {} {} (request_id: {})", method, path, request.request_id);
let routes_count = self.routes.len();
if routes_count > 10 {
debug!("Searching through {} routes for match", routes_count);
}
for route in self.routes.as_ref() {
if route.method == method && route.path == path {
debug!("Found matching route: {} {}", route.method, route.path);
let handler_result = match (route.handler)(request.clone(), state.clone()).await {
Ok(response) => {
let elapsed = start_time.elapsed();
let elapsed_ms = elapsed.as_millis() as usize;
let current_avg = HTTP_IPC_METRICS.avg_response_time_ms.load(std::sync::atomic::Ordering::Relaxed);
let count = HTTP_IPC_METRICS.responses_sent.load(std::sync::atomic::Ordering::Relaxed);
let new_avg = if count == 0 {
elapsed_ms
} else {
((current_avg * count) + elapsed_ms) / (count + 1)
};
HTTP_IPC_METRICS.avg_response_time_ms.store(new_avg, std::sync::atomic::Ordering::Relaxed);
debug!("Handler completed in {}ms (request_id: {})", elapsed_ms, request.request_id);
Ok(response)
},
Err(err) => {
error!("Handler error for {} {}: {} (request_id: {})",
method, path, err, request.request_id);
HTTP_IPC_METRICS.errors_encountered.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
Err(err)
}
};
match handler_result {
Ok(mut response) => {
response.request_id = request.request_id;
return response;
}
Err(err) => {
return result::error_to_response(err, &request.request_id);
}
}
}
}
if let Some(handler) = &self.not_found_handler {
match handler(request.clone(), state).await {
Ok(mut response) => {
response.request_id = request.request_id;
info!("Not found handler completed for {} {} (request_id: {})",
method, path, request.request_id);
response
}
Err(err) => {
error!("Not found handler error: {} (request_id: {})",
err, request.request_id);
HTTP_IPC_METRICS.errors_encountered.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
result::error_to_response(err, &request.request_id)
}
}
} else {
let response = ApiResponse::<()> {
status: "error".to_string(),
data: None,
message: Some(format!("Endpoint not found: {} {}", method, path)),
};
warn!("No route found for {} {} (request_id: {})",
method, path, request.request_id);
let body = match serde_json::to_vec(&response) {
Ok(body) => Some(body),
Err(e) => {
error!("Failed to serialize not found response: {} (request_id: {})",
e, request.request_id);
Some(format!("{{\"status\":\"error\",\"message\":\"Endpoint not found: {} {}\"}}",
method, path).into_bytes())
}
};
let mut headers = HashMap::new();
headers.insert("Content-Type".to_string(), "application/json".to_string());
IpcHttpResponse {
request_id: request.request_id,
status_code: 404,
headers,
body,
}
}
}
}
#[derive(Debug, Clone)]
pub struct HttpIpcServerConfig {
pub max_concurrent_requests: usize,
pub trace_requests: bool,
pub request_id_logger: Option<Arc<dyn Fn(&str) + Send + Sync>>,
}
impl Default for HttpIpcServerConfig {
fn default() -> Self {
Self {
max_concurrent_requests: 100,
trace_requests: false,
request_id_logger: None,
}
}
}
pub async fn start_http_ipc_server<S>(
router: HttpIpcRouter<S>,
state: S,
shutdown_signal: Option<broadcast::Receiver<()>>,
config: Option<HttpIpcServerConfig>,
) -> std::result::Result<(), crate::Error>
where
S: Send + Sync + Clone + 'static,
{
let server_config = config.unwrap_or_default();
info!("Starting HTTP IPC server (max_concurrent_requests: {})",
server_config.max_concurrent_requests);
let mut rx = subscribe_http_requests();
let semaphore = Arc::new(tokio::sync::Semaphore::new(server_config.max_concurrent_requests));
let state = Arc::new(state);
let server_task = tokio::spawn(async move {
info!("HTTP IPC server task started");
if let Some(port) = PortNegotiationManager::get_allocated_port() {
info!("HTTP IPC server using allocated port: {}", port);
} else {
info!("HTTP IPC server running in pure IPC mode (no TCP port allocated)");
}
loop {
match rx.recv().await {
Ok(request) => {
HTTP_IPC_METRICS.requests_received.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
info!("HTTP IPC server received request - ID: {}, Method: {}, URI: {}",
request.request_id, request.method, request.uri);
if server_config.trace_requests && Level::TRACE <= tracing::level_filters::LevelFilter::current() {
trace!("Request headers: {:?}", request.headers);
if let Some(body) = &request.body {
if body.len() < 1024 {
match std::str::from_utf8(body) {
Ok(body_str) => trace!("Request body (string): {}", body_str),
Err(_) => trace!("Request body (bytes): {:?}", body),
}
} else {
trace!("Request body size: {} bytes", body.len());
}
} else {
trace!("Request has no body");
}
}
if let Some(id_logger) = &server_config.request_id_logger {
id_logger(&request.request_id);
}
let permit = match semaphore.clone().acquire_owned().await {
Ok(permit) => permit,
Err(e) => {
error!("Failed to acquire semaphore permit: {}", e);
continue;
}
};
let router = router.clone();
let state = Arc::clone(&state);
let request_id = request.request_id.clone();
tokio::spawn(async move {
debug!("Processing request in handler task (request_id: {})", request_id);
let response = router.handle_request(request, state).await;
if let Err(e) = send_http_response(response).await {
error!("Failed to send HTTP response: {} (request_id: {})",
e, request_id);
}
drop(permit);
debug!("Handler task completed (request_id: {})", request_id);
});
}
Err(e) => {
match e {
broadcast::error::RecvError::Closed => {
error!("HTTP request channel closed. Shutting down HTTP IPC server.");
break;
}
broadcast::error::RecvError::Lagged(count) => {
warn!("HTTP IPC server lagged behind by {} messages. Consider increasing channel capacity.", count);
rx = subscribe_http_requests();
}
}
}
}
}
});
if let Some(mut shutdown_rx) = shutdown_signal {
tokio::select! {
_ = shutdown_rx.recv() => {
info!("Received shutdown signal. Shutting down HTTP IPC server.");
}
_ = server_task => {
info!("HTTP IPC server task completed unexpectedly.");
}
}
} else {
server_task.await.map_err(|e| {
crate::Error::Config(crate::error::ConfigError::Invalid(
format!("HTTP IPC server task failed: {}", e)
))
})?;
}
info!("HTTP IPC server stopped.");
Ok(())
}