use kotoba_core::types::{TxId, ContentHash, Result, KotobaError, Value, Properties};
use crate::http::ir::*;
use kotoba_graph::prelude::*;
use kotoba_rewrite::prelude::*;
use kotoba_security::{SecurityService, AuditResult};
use kotoba_core::ir::rule::{RuleIR, Match};
use kotoba_core::ir::strategy::{StrategyIR, StrategyOp};
use kotoba_core::ir::patch::Patch;
pub use kotoba_security::{JwtClaims, User, AuthResult, AuthzResult, Principal, Resource};
use std::collections::HashMap;
use std::sync::Arc;
#[derive(Clone)]
pub struct HttpRequestProcessor {
rewrite_engine: Arc<RewriteEngine>,
mvcc: Arc<MVCCManager>,
merkle: Arc<MerkleDAG>,
security: Arc<SecurityService>,
}
impl HttpRequestProcessor {
pub fn new(
rewrite_engine: Arc<RewriteEngine>,
mvcc: Arc<MVCCManager>,
merkle: Arc<MerkleDAG>,
security: Arc<SecurityService>,
) -> Self {
Self {
rewrite_engine,
mvcc,
merkle,
security,
}
}
pub async fn process_request(&self, request: HttpRequest) -> Result<HttpResponse> {
self.process_request_simple(request).await
}
async fn process_request_simple(&self, request: HttpRequest) -> Result<HttpResponse> {
match request.path.as_str() {
"/ping" => {
let mut headers = HttpHeaders::new();
headers.set("content-type".to_string(), "application/json".to_string());
Ok(HttpResponse::new(
request.id,
HttpStatus::ok(),
headers,
Some(ContentHash::sha256([0; 32])), ))
},
"/health" => {
let mut headers = HttpHeaders::new();
headers.set("content-type".to_string(), "application/json".to_string());
Ok(HttpResponse::new(
request.id,
HttpStatus::ok(),
headers,
Some(ContentHash::sha256([1; 32])), ))
},
"/graphql" => {
let mut headers = HttpHeaders::new();
headers.set("content-type".to_string(), "application/json".to_string());
Ok(HttpResponse::new(
request.id,
HttpStatus::ok(),
headers,
Some(ContentHash::sha256([2; 32])),
))
},
_ => {
Ok(HttpResponse::new(
request.id,
HttpStatus::not_found(),
HttpHeaders::new(),
None,
))
}
}
}
}
pub struct HttpRewriteExterns;
impl HttpRewriteExterns {
pub fn new() -> Self {
Self
}
}
impl RewriteExterns for HttpRewriteExterns {
fn deg_ge(&self, _v: crate::types::VertexId, _k: u32) -> bool {
true
}
fn edge_count_nonincreasing(&self, _g0: &GraphRef, _g1: &GraphRef) -> bool {
true
}
fn custom_measure(&self, _name: &str, _args: &[kotoba_core::types::Value]) -> f64 {
0.0
}
}
#[derive(Clone)]
pub struct MiddlewareProcessor {
middlewares: Vec<HttpMiddleware>,
security: Arc<SecurityService>,
}
impl MiddlewareProcessor {
pub fn new(middlewares: Vec<HttpMiddleware>, security: Arc<SecurityService>) -> Self {
Self { middlewares, security }
}
pub async fn process(&self, request: &mut HttpRequest) -> Result<()> {
let mut sorted_middlewares = self.middlewares.clone();
sorted_middlewares.sort_by_key(|mw| mw.order);
for middleware in sorted_middlewares {
self.execute_middleware(&middleware, request).await?;
}
Ok(())
}
async fn execute_middleware(&self, middleware: &HttpMiddleware, request: &mut HttpRequest) -> Result<()> {
match middleware.name.as_str() {
"request_id" => {
let request_id = format!("req_{}", request.id);
request.headers.set("x-request-id".to_string(), request_id);
},
"logger" => {
println!("Request: {} {} {}", request.method, request.path, request.id);
},
"cors" => {
},
"jwt_auth" => {
self.execute_jwt_auth_middleware(request).await?;
},
"authorization" => {
self.execute_authorization_middleware(request).await?;
},
"rate_limit" => {
self.execute_rate_limit_middleware(request).await?;
},
"csrf" => {
self.execute_csrf_middleware(request).await?;
},
_ => {
println!("Executing custom middleware: {}", middleware.name);
}
}
Ok(())
}
async fn execute_jwt_auth_middleware(&self, request: &mut HttpRequest) -> Result<()> {
let auth_header = request.headers.get("authorization");
if let Some(auth_value) = auth_header {
if auth_value.starts_with("Bearer ") {
let token = &auth_value[7..];
match self.security.validate_token(token) {
Ok(claims) => {
println!("AUDIT: JWT authentication successful for user: {}", claims.sub);
request.attributes.insert(
"user_id".to_string(),
Value::String(claims.sub)
);
request.attributes.insert(
"roles".to_string(),
Value::Array(claims.roles)
);
return Ok(());
}
Err(_) => {
println!("AUDIT: JWT authentication failed - invalid token");
return Err(KotobaError::Security("Invalid JWT token".to_string()));
}
}
}
}
println!("AUDIT: Missing authentication token");
Err(KotobaError::Security("Authentication required".to_string()))
}
async fn execute_authorization_middleware(&self, request: &mut HttpRequest) -> Result<()> {
if request.attributes.contains_key("user_id") {
println!("AUDIT: Authorization successful for path: {}", request.path);
Ok(())
} else {
println!("AUDIT: Authorization failed - user not authenticated");
Err(KotobaError::Security("Authorization required".to_string()))
}
}
fn determine_resource_from_request(&self, request: &HttpRequest) -> (kotoba_security::ResourceType, kotoba_security::Action, Option<String>) {
use kotoba_security::{ResourceType, Action};
match request.path.as_str() {
path if path.starts_with("/api/graph/") => {
let resource_id = if path.len() > "/api/graph/".len() {
Some(path["/api/graph/".len()..].to_string())
} else {
None
};
match request.method {
HttpMethod::GET => (ResourceType::Graph, Action::Read, resource_id),
HttpMethod::POST => (ResourceType::Graph, Action::Create, resource_id),
HttpMethod::PUT => (ResourceType::Graph, Action::Update, resource_id),
HttpMethod::DELETE => (ResourceType::Graph, Action::Delete, resource_id),
_ => (ResourceType::Graph, Action::Read, resource_id),
}
},
path if path.starts_with("/api/query") => {
(ResourceType::Query, Action::Execute, None)
},
path if path.starts_with("/api/admin") => {
(ResourceType::Admin, Action::Admin, None)
},
path if path.starts_with("/api/user") => {
(ResourceType::User, Action::Read, None)
},
_ => {
(ResourceType::FileSystem, Action::Read, Some(request.path.clone()))
}
}
}
async fn execute_rate_limit_middleware(&self, _request: &mut HttpRequest) -> Result<()> {
Ok(())
}
async fn execute_csrf_middleware(&self, request: &mut HttpRequest) -> Result<()> {
let is_state_changing = matches!(request.method,
crate::http::ir::HttpMethod::POST |
crate::http::ir::HttpMethod::PUT |
crate::http::ir::HttpMethod::DELETE);
if is_state_changing {
if request.headers.get("x-csrf-token").is_some() {
Ok(())
} else {
Err(KotobaError::Security("CSRF token required".to_string()))
}
} else {
Ok(())
}
}
}
#[derive(Clone)]
pub struct HandlerProcessor;
impl HandlerProcessor {
pub fn new() -> Self {
Self
}
pub async fn process(&self, route: &HttpRoute, request: &HttpRequest) -> Result<HttpResponse> {
match route.pattern.as_str() {
"/ping" => {
let mut headers = HttpHeaders::new();
headers.set("content-type".to_string(), "application/json".to_string());
Ok(HttpResponse::new(
request.id.clone(),
HttpStatus::ok(),
headers,
Some(ContentHash::sha256([0; 32])), ))
},
"/health" => {
let mut headers = HttpHeaders::new();
headers.set("content-type".to_string(), "application/json".to_string());
Ok(HttpResponse::new(
request.id.clone(),
HttpStatus::ok(),
headers,
Some(ContentHash::sha256([1; 32])), ))
},
_ => {
Ok(HttpResponse::new(
request.id.clone(),
HttpStatus::not_found(),
HttpHeaders::new(),
None,
))
}
}
}
}