Skip to main content

lean_ctx/http_server/
mod.rs

1//! HTTP server combining Streamable-HTTP MCP transport, REST APIs, and
2//! optional local HTTP surfaces.
3//!
4//! # Pillar mapping
5//!
6//! - **Engine:** HTTP MCP transport, Context OS event bus (SSE), A2A handoffs,
7//!   agent registry, capabilities/manifest/openapi endpoints.
8
9use std::net::SocketAddr;
10use std::path::PathBuf;
11use std::sync::Arc;
12
13use anyhow::{Context, Result, anyhow};
14use axum::{
15    Router,
16    extract::Json,
17    extract::Query,
18    extract::State,
19    http::{Request, StatusCode, header},
20    middleware::{self, Next},
21    response::sse::{Event as SseEvent, KeepAlive, Sse},
22    response::{IntoResponse, Response},
23    routing::get,
24};
25use futures::Stream;
26use rmcp::transport::{StreamableHttpServerConfig, StreamableHttpService};
27use serde::Deserialize;
28use serde_json::Value;
29use tokio::sync::broadcast;
30use tokio::time::{Duration, Instant};
31
32use crate::core::context_os::ContextOsMetrics;
33use crate::engine::ContextEngine;
34use crate::tools::LeanCtxServer;
35
36mod handlers;
37#[allow(clippy::wildcard_imports)]
38use handlers::*;
39
40pub mod kernel_api;
41
42/// Wrapper stream that calls `record_sse_disconnect` on drop.
43use std::pin::Pin;
44
45pub(crate) struct SseDisconnectGuard<I> {
46    pub(crate) inner: Pin<Box<dyn Stream<Item = I> + Send>>,
47    pub(crate) metrics: Arc<ContextOsMetrics>,
48}
49
50impl<I> Stream for SseDisconnectGuard<I> {
51    type Item = I;
52
53    fn poll_next(
54        mut self: Pin<&mut Self>,
55        cx: &mut std::task::Context<'_>,
56    ) -> std::task::Poll<Option<Self::Item>> {
57        self.inner.as_mut().poll_next(cx)
58    }
59}
60
61impl<I> Drop for SseDisconnectGuard<I> {
62    fn drop(&mut self) {
63        self.metrics.record_sse_disconnect();
64    }
65}
66
67const MAX_ID_LEN: usize = 64;
68
69fn sanitize_id(raw: &str) -> String {
70    let trimmed = raw.trim();
71    if trimmed.is_empty() {
72        return "default".to_string();
73    }
74    let cleaned: String = trimmed
75        .chars()
76        .filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_' || *c == '.')
77        .take(MAX_ID_LEN)
78        .collect();
79    if cleaned.is_empty() {
80        "default".to_string()
81    } else {
82        cleaned
83    }
84}
85
86#[derive(Clone, Debug)]
87pub struct HttpServerConfig {
88    pub host: String,
89    pub port: u16,
90    pub project_root: PathBuf,
91    pub auth_token: Option<String>,
92    pub stateful_mode: bool,
93    pub json_response: bool,
94    pub disable_host_check: bool,
95    pub allowed_hosts: Vec<String>,
96    pub max_body_bytes: usize,
97    pub max_concurrency: usize,
98    pub max_rps: u32,
99    pub rate_burst: u32,
100    pub request_timeout_ms: u64,
101}
102
103impl Default for HttpServerConfig {
104    fn default() -> Self {
105        let project_root = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
106        Self {
107            host: "127.0.0.1".to_string(),
108            port: 8080,
109            project_root,
110            auth_token: None,
111            stateful_mode: false,
112            json_response: true,
113            disable_host_check: false,
114            allowed_hosts: Vec::new(),
115            max_body_bytes: 2 * 1024 * 1024,
116            max_concurrency: 32,
117            max_rps: 50,
118            rate_burst: 100,
119            request_timeout_ms: 30_000,
120        }
121    }
122}
123
124impl HttpServerConfig {
125    pub fn validate(&self) -> Result<()> {
126        let host = self.host.trim().to_lowercase();
127        let is_loopback = host == "127.0.0.1" || host == "localhost" || host == "::1";
128        if !is_loopback && self.auth_token.as_deref().unwrap_or("").is_empty() {
129            return Err(anyhow!(
130                "Refusing to bind to host='{host}' without auth. Provide --auth-token (or bind to 127.0.0.1)."
131            ));
132        }
133        Ok(())
134    }
135
136    pub fn effective_auth_token(&self) -> Option<String> {
137        if let Some(ref token) = self.auth_token
138            && !token.is_empty()
139        {
140            return Some(token.clone());
141        }
142        let host = self.host.trim().to_lowercase();
143        let is_loopback = host == "127.0.0.1" || host == "localhost" || host == "::1";
144        if is_loopback {
145            let auto_token = crate::core::session_token::generate_token();
146            eprintln!(
147                "[lean-ctx] Auto-generated auth token for loopback: {auto_token}\n\
148                 Pass as Bearer token or set --auth-token explicitly."
149            );
150            Some(auto_token)
151        } else {
152            None
153        }
154    }
155
156    fn mcp_http_config(&self) -> StreamableHttpServerConfig {
157        let mut cfg = StreamableHttpServerConfig::default()
158            .with_stateful_mode(self.stateful_mode)
159            .with_json_response(self.json_response);
160
161        if self.disable_host_check {
162            tracing::warn!(
163                "⚠ --disable-host-check is active: DNS rebinding protection is OFF. \
164                 Do NOT use this in production or on non-loopback interfaces."
165            );
166            cfg = cfg.disable_allowed_hosts();
167            return cfg;
168        }
169
170        if !self.allowed_hosts.is_empty() {
171            cfg = cfg.with_allowed_hosts(self.allowed_hosts.clone());
172            return cfg;
173        }
174
175        // Keep rmcp's secure loopback defaults; also allow the configured host (if it's loopback).
176        let host = self.host.trim();
177        if host == "127.0.0.1" || host == "localhost" || host == "::1" {
178            cfg.allowed_hosts.push(host.to_string());
179        }
180
181        cfg
182    }
183}
184
185#[derive(Clone)]
186struct AppState {
187    token: Option<String>,
188    concurrency: Arc<tokio::sync::Semaphore>,
189    rate: Arc<RateLimiter>,
190    project_root: String,
191    timeout: Duration,
192    server: LeanCtxServer,
193}
194
195#[derive(Debug)]
196struct RateLimiter {
197    max_rps: f64,
198    burst: f64,
199    state: tokio::sync::Mutex<RateState>,
200}
201
202#[derive(Debug, Clone, Copy)]
203struct RateState {
204    tokens: f64,
205    last: Instant,
206}
207
208impl RateLimiter {
209    fn new(max_rps: u32, burst: u32) -> Self {
210        let now = Instant::now();
211        Self {
212            max_rps: (max_rps.max(1)) as f64,
213            burst: (burst.max(1)) as f64,
214            state: tokio::sync::Mutex::new(RateState {
215                tokens: (burst.max(1)) as f64,
216                last: now,
217            }),
218        }
219    }
220
221    async fn allow(&self) -> bool {
222        let mut s = self.state.lock().await;
223        let now = Instant::now();
224        let elapsed = now.saturating_duration_since(s.last);
225        let refill = elapsed.as_secs_f64() * self.max_rps;
226        s.tokens = (s.tokens + refill).min(self.burst);
227        s.last = now;
228        if s.tokens >= 1.0 {
229            s.tokens -= 1.0;
230            true
231        } else {
232            false
233        }
234    }
235}
236
237async fn auth_middleware(
238    State(state): State<AppState>,
239    req: Request<axum::body::Body>,
240    next: Next,
241) -> Response {
242    if state.token.is_none() {
243        return next.run(req).await;
244    }
245
246    if req.uri().path() == "/health" {
247        return next.run(req).await;
248    }
249
250    let expected = state.token.as_deref().unwrap_or("");
251    let Some(h) = req.headers().get(header::AUTHORIZATION) else {
252        return json_error(
253            StatusCode::UNAUTHORIZED,
254            "unauthorized",
255            "missing Authorization header",
256        );
257    };
258    let Ok(s) = h.to_str() else {
259        return json_error(
260            StatusCode::UNAUTHORIZED,
261            "unauthorized",
262            "malformed Authorization header",
263        );
264    };
265    let Some(token) = s
266        .strip_prefix("Bearer ")
267        .or_else(|| s.strip_prefix("bearer "))
268    else {
269        return json_error(
270            StatusCode::UNAUTHORIZED,
271            "unauthorized",
272            "Authorization must use the Bearer scheme",
273        );
274    };
275    if !constant_time_eq(token.as_bytes(), expected.as_bytes()) {
276        return json_error(
277            StatusCode::UNAUTHORIZED,
278            "unauthorized",
279            "invalid bearer token",
280        );
281    }
282
283    next.run(req).await
284}
285
286/// Structured REST error envelope: `{ "error": <human message>, "error_code": <stable code> }`.
287///
288/// `error_code` is the stable, machine-readable string SDKs switch on; `error` carries the
289/// human-facing message. Used for every REST (non-A2A) error so clients branch on a code
290/// instead of parsing prose. The A2A JSON-RPC surface keeps its own `-32xxx` envelope.
291pub(crate) fn json_error(status: StatusCode, error_code: &str, message: &str) -> Response {
292    (
293        status,
294        Json(serde_json::json!({ "error": message, "error_code": error_code })),
295    )
296        .into_response()
297}
298
299fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
300    use subtle::ConstantTimeEq;
301    if a.len() != b.len() {
302        return false;
303    }
304    bool::from(a.ct_eq(b))
305}
306
307async fn rate_limit_middleware(
308    State(state): State<AppState>,
309    req: Request<axum::body::Body>,
310    next: Next,
311) -> Response {
312    if !state.rate.allow().await {
313        return StatusCode::TOO_MANY_REQUESTS.into_response();
314    }
315    next.run(req).await
316}
317
318async fn concurrency_middleware(
319    State(state): State<AppState>,
320    req: Request<axum::body::Body>,
321    next: Next,
322) -> Response {
323    let Ok(permit) = state.concurrency.clone().try_acquire_owned() else {
324        return StatusCode::TOO_MANY_REQUESTS.into_response();
325    };
326    let resp = next.run(req).await;
327    drop(permit);
328    resp
329}
330
331async fn health() -> impl IntoResponse {
332    (StatusCode::OK, "ok\n")
333}
334
335async fn v1_shutdown() -> impl IntoResponse {
336    tokio::spawn(async {
337        tokio::time::sleep(Duration::from_millis(100)).await;
338        std::process::exit(0);
339    });
340    (StatusCode::OK, "shutting down\n")
341}
342
343#[derive(Debug, Deserialize)]
344#[serde(rename_all = "camelCase")]
345struct IndexEnsureBody {
346    root: String,
347    #[serde(default)]
348    extra_roots: Vec<String>,
349}
350
351/// Daemon-side index delegation (#460). A thin-client session POSTs the repo it
352/// needs warmed and the daemon — the single long-lived indexer — builds it once
353/// in the background (deduped per root). Every other session for the same root
354/// then load-shares the on-disk result via the `graph-idx`/`bm25-idx`
355/// cross-process locks instead of running its own scan, so N concurrent sessions
356/// cost ~one index pass machine-wide instead of N. Returns immediately; the
357/// build runs in the orchestrator's own worker thread.
358async fn v1_index_ensure(Json(body): Json<IndexEnsureBody>) -> impl IntoResponse {
359    if body.root.trim().is_empty() {
360        return (StatusCode::BAD_REQUEST, "root is required\n");
361    }
362    let root = body.root;
363    let extra = body.extra_roots;
364    tokio::task::spawn_blocking(move || {
365        crate::core::index_orchestrator::ensure_all_background(&root);
366        if !extra.is_empty() {
367            crate::core::index_orchestrator::ensure_extra_roots_background(&root, &extra);
368        }
369    });
370    (StatusCode::OK, "{\"status\":\"ok\"}\n")
371}
372
373#[derive(Debug, Deserialize)]
374#[serde(rename_all = "camelCase")]
375struct ToolCallBody {
376    name: String,
377    #[serde(default)]
378    arguments: Option<Value>,
379    #[serde(default)]
380    _workspace_id: Option<String>,
381    #[serde(default)]
382    _channel_id: Option<String>,
383}
384
385#[derive(Debug, Deserialize)]
386#[serde(rename_all = "camelCase")]
387struct EventsQuery {
388    #[serde(default)]
389    workspace_id: Option<String>,
390    #[serde(default)]
391    channel_id: Option<String>,
392    #[serde(default)]
393    since: Option<i64>,
394    #[serde(default)]
395    limit: Option<usize>,
396    /// Comma-separated event kind filter (e.g. `tool_call,session_start`).
397    /// When set, only matching events are delivered via SSE.
398    #[serde(default)]
399    kind: Option<String>,
400}
401
402async fn v1_manifest(State(state): State<AppState>) -> impl IntoResponse {
403    let _ = state;
404    let v = crate::core::mcp_manifest::manifest_value();
405    (StatusCode::OK, Json(v))
406}
407
408/// `GET /v1/capabilities` — discovery document describing what this instance
409/// supports (presets, tools, read modes, features, extensions, contract
410/// versions). See `docs/contracts/capabilities-contract-v1.md`.
411async fn v1_capabilities(State(state): State<AppState>) -> impl IntoResponse {
412    let _ = state;
413    (
414        StatusCode::OK,
415        Json(crate::core::server_capabilities::capabilities_value()),
416    )
417}
418
419/// `GET /v1/openapi.json` — OpenAPI 3.0 document for the public `/v1` surface,
420/// generated from the in-code endpoint inventory (`core::openapi`).
421async fn v1_openapi(State(state): State<AppState>) -> impl IntoResponse {
422    let _ = state;
423    (StatusCode::OK, Json(crate::core::openapi::openapi_value()))
424}
425
426#[derive(Debug, Deserialize)]
427#[serde(rename_all = "camelCase")]
428struct ToolsQuery {
429    #[serde(default)]
430    offset: Option<usize>,
431    #[serde(default)]
432    limit: Option<usize>,
433}
434
435async fn v1_tools(State(state): State<AppState>, Query(q): Query<ToolsQuery>) -> impl IntoResponse {
436    let _ = state;
437    let v = crate::core::mcp_manifest::manifest_value();
438    let tools = v
439        .get("tools")
440        .and_then(|t| t.get("granular"))
441        .cloned()
442        .unwrap_or(Value::Array(vec![]));
443
444    let all = tools.as_array().cloned().unwrap_or_default();
445    let total = all.len();
446    let offset = q.offset.unwrap_or(0).min(total);
447    let limit = q.limit.unwrap_or(200).min(500);
448    let page = all.into_iter().skip(offset).take(limit).collect::<Vec<_>>();
449
450    (
451        StatusCode::OK,
452        Json(serde_json::json!({
453            "tools": page,
454            "total": total,
455            "offset": offset,
456            "limit": limit,
457        })),
458    )
459}
460
461async fn v1_tool_call(
462    State(state): State<AppState>,
463    Json(body): Json<ToolCallBody>,
464) -> impl IntoResponse {
465    let engine = ContextEngine::from_server(state.server.clone());
466    match tokio::time::timeout(
467        state.timeout,
468        engine.call_tool_value(&body.name, body.arguments),
469    )
470    .await
471    {
472        Ok(Ok(v)) => (StatusCode::OK, Json(serde_json::json!({ "result": v }))).into_response(),
473        Ok(Err(e)) => {
474            tracing::warn!("tool call error: {e}");
475            json_error(
476                StatusCode::BAD_REQUEST,
477                "tool_error",
478                "tool execution failed",
479            )
480        }
481        Err(_) => json_error(
482            StatusCode::GATEWAY_TIMEOUT,
483            "request_timeout",
484            "tool call timed out",
485        ),
486    }
487}
488
489async fn v1_events(
490    State(state): State<AppState>,
491    Query(q): Query<EventsQuery>,
492) -> Sse<impl Stream<Item = Result<SseEvent, std::convert::Infallible>>> {
493    use crate::core::context_os::{ContextEventV1, RedactionLevel, redact_event_payload};
494
495    let ws = sanitize_id(&q.workspace_id.unwrap_or_else(|| "default".to_string()));
496    let ch = sanitize_id(&q.channel_id.unwrap_or_else(|| "default".to_string()));
497    let _ = &state.project_root;
498    let since = q.since.unwrap_or(0);
499    let limit = q.limit.unwrap_or(200).min(1000);
500    let redaction = RedactionLevel::RefsOnly;
501
502    let kind_filter: Option<Vec<String>> = q
503        .kind
504        .as_deref()
505        .map(|k| k.split(',').map(|s| s.trim().to_string()).collect());
506
507    let rt = crate::core::context_os::runtime();
508    let replay = rt.bus.read(&ws, &ch, since, limit);
509
510    let replay = if let Some(ref kinds) = kind_filter {
511        replay
512            .into_iter()
513            .filter(|ev| kinds.contains(&ev.kind))
514            .collect()
515    } else {
516        replay
517    };
518
519    let rx = if let Some(ref kinds) = kind_filter {
520        let kind_refs: Vec<&str> = kinds.iter().map(String::as_str).collect();
521        let filter = crate::core::context_os::TopicFilter::kinds(&kind_refs);
522        if let Some(sub) = rt.bus.subscribe_filtered(&ws, &ch, filter) {
523            crate::core::context_os::SubscriptionKind::Filtered(sub)
524        } else {
525            tracing::warn!("SSE subscriber limit reached for {ws}/{ch}");
526            let (_, rx) = broadcast::channel::<ContextEventV1>(1);
527            crate::core::context_os::SubscriptionKind::Unfiltered(rx)
528        }
529    } else if let Some(sub) = rt.bus.subscribe(&ws, &ch) {
530        crate::core::context_os::SubscriptionKind::Unfiltered(sub)
531    } else {
532        tracing::warn!("SSE subscriber limit reached for {ws}/{ch}");
533        let (_, rx) = broadcast::channel::<ContextEventV1>(1);
534        crate::core::context_os::SubscriptionKind::Unfiltered(rx)
535    };
536
537    rt.metrics.record_sse_connect();
538    rt.metrics.record_events_replayed(replay.len() as u64);
539    rt.metrics.record_workspace_active(&ws);
540
541    let bus = rt.bus.clone();
542    let metrics = rt.metrics.clone();
543    let pending: std::collections::VecDeque<ContextEventV1> = replay.into();
544
545    let stream = futures::stream::unfold(
546        (
547            pending,
548            rx,
549            ws.clone(),
550            ch.clone(),
551            since,
552            redaction,
553            bus,
554            metrics,
555        ),
556        |(mut pending, mut rx, ws, ch, mut last_id, redaction, bus, metrics)| async move {
557            if let Some(mut ev) = pending.pop_front() {
558                last_id = ev.id;
559                redact_event_payload(&mut ev, redaction);
560                let data = serde_json::to_string(&ev).unwrap_or_else(|_| "{}".to_string());
561                let evt = SseEvent::default()
562                    .id(ev.id.to_string())
563                    .event(ev.kind)
564                    .data(data);
565                return Some((
566                    Ok(evt),
567                    (pending, rx, ws, ch, last_id, redaction, bus, metrics),
568                ));
569            }
570
571            loop {
572                match rx.recv().await {
573                    Ok(mut ev) if ev.id > last_id => {
574                        last_id = ev.id;
575                        redact_event_payload(&mut ev, redaction);
576                        let data = serde_json::to_string(&ev).unwrap_or_else(|_| "{}".to_string());
577                        let evt = SseEvent::default()
578                            .id(ev.id.to_string())
579                            .event(ev.kind)
580                            .data(data);
581                        return Some((
582                            Ok(evt),
583                            (pending, rx, ws, ch, last_id, redaction, bus, metrics),
584                        ));
585                    }
586                    Ok(_) => {}
587                    Err(broadcast::error::RecvError::Closed) => return None,
588                    Err(broadcast::error::RecvError::Lagged(skipped)) => {
589                        let missed = bus.read(&ws, &ch, last_id, skipped as usize);
590                        metrics.record_events_replayed(missed.len() as u64);
591                        for ev in missed {
592                            last_id = last_id.max(ev.id);
593                            pending.push_back(ev);
594                        }
595                    }
596                }
597            }
598        },
599    );
600
601    let metrics_ref = rt.metrics.clone();
602    let guarded = SseDisconnectGuard {
603        inner: Box::pin(stream),
604        metrics: metrics_ref,
605    };
606
607    Sse::new(guarded).keep_alive(KeepAlive::new().interval(Duration::from_secs(15)))
608}
609
610#[derive(Debug, Deserialize)]
611struct AuditEventsQuery {
612    #[serde(default = "default_audit_limit")]
613    limit: usize,
614}
615
616fn default_audit_limit() -> usize {
617    100
618}
619
620async fn v1_audit_events(Query(q): Query<AuditEventsQuery>) -> impl IntoResponse {
621    let capped = q.limit.min(1000);
622    let boundary_events = crate::core::memory_boundary::load_audit_events(capped);
623    let trail_events = crate::core::audit_trail::load_recent(capped);
624
625    Json(serde_json::json!({
626        "cross_project_events": boundary_events,
627        "audit_trail": trail_events,
628    }))
629}
630
631#[derive(Deserialize)]
632struct ContextSummaryQuery {
633    #[serde(default)]
634    workspace_id: Option<String>,
635    #[serde(default)]
636    limit: Option<usize>,
637}
638
639async fn v1_context_summary(
640    State(state): State<AppState>,
641    Query(q): Query<ContextSummaryQuery>,
642) -> impl IntoResponse {
643    let ws = sanitize_id(&q.workspace_id.unwrap_or_else(|| "default".to_string()));
644    let limit = q.limit.unwrap_or(100).min(1000);
645    let rt = crate::core::context_os::runtime();
646    let events = rt.bus.read(&ws, "default", 0, limit);
647    let mut counts_by_kind = serde_json::Map::new();
648    for ev in &events {
649        let counter = counts_by_kind
650            .entry(ev.kind.clone())
651            .or_insert(serde_json::Value::from(0u64));
652        if let Some(n) = counter.as_u64() {
653            *counter = serde_json::Value::from(n + 1);
654        }
655    }
656    Json(serde_json::json!({
657        "workspaceId": ws,
658        "channelId": "default",
659        "projectRoot": state.project_root,
660        "totalEvents": events.len(),
661        "eventCountsByKind": counts_by_kind,
662        "limit": limit,
663    }))
664}
665
666#[derive(Deserialize)]
667struct EventsSearchQuery {
668    #[serde(default)]
669    q: Option<String>,
670    #[serde(default)]
671    limit: Option<usize>,
672}
673
674async fn v1_events_search(Query(q): Query<EventsSearchQuery>) -> impl IntoResponse {
675    let query = q.q.unwrap_or_default();
676    let limit = q.limit.unwrap_or(50).min(500);
677    let rt = crate::core::context_os::runtime();
678    let all_events = rt.bus.read("default", "default", 0, limit * 10);
679    let results: Vec<serde_json::Value> = all_events
680        .into_iter()
681        .filter(|ev| {
682            let payload = serde_json::to_string(ev).unwrap_or_default();
683            query.is_empty() || payload.contains(&query)
684        })
685        .take(limit)
686        .map(|ev| serde_json::to_value(&ev).unwrap_or_default())
687        .collect();
688    Json(serde_json::json!({
689        "query": query,
690        "results": results,
691        "count": results.len(),
692    }))
693}
694
695#[derive(Deserialize)]
696struct EventLineageQuery {
697    #[serde(default)]
698    id: Option<i64>,
699    #[serde(default)]
700    depth: Option<usize>,
701}
702
703async fn v1_events_lineage(Query(q): Query<EventLineageQuery>) -> impl IntoResponse {
704    let event_id = q.id.unwrap_or(0);
705    let depth = q.depth.unwrap_or(10).min(100);
706    let rt = crate::core::context_os::runtime();
707    let events = rt.bus.read("default", "default", 0, depth);
708    let chain: Vec<serde_json::Value> = events
709        .into_iter()
710        .filter(|ev| ev.id >= event_id)
711        .take(depth)
712        .map(|ev| serde_json::to_value(&ev).unwrap_or_default())
713        .collect();
714    Json(serde_json::json!({
715        "eventId": event_id,
716        "depth": depth,
717        "chain": chain,
718    }))
719}
720
721async fn v1_metrics(State(_state): State<AppState>) -> impl IntoResponse {
722    let rt = crate::core::context_os::runtime();
723    let snap = rt.metrics.snapshot();
724    (
725        StatusCode::OK,
726        Json(serde_json::to_value(snap).unwrap_or_default()),
727    )
728}
729
730async fn a2a_jsonrpc(Json(body): Json<Value>) -> impl IntoResponse {
731    let req: crate::core::a2a::a2a_compat::JsonRpcRequest = match serde_json::from_value(body) {
732        Ok(r) => r,
733        Err(e) => {
734            tracing::debug!("a2a JSON-RPC parse error: {e}");
735            return (
736                StatusCode::BAD_REQUEST,
737                Json(serde_json::json!({
738                    "jsonrpc": "2.0",
739                    "id": null,
740                    "error": {"code": -32700, "message": "invalid request"}
741                })),
742            );
743        }
744    };
745    let resp = crate::core::a2a::a2a_compat::handle_a2a_jsonrpc(&req);
746    let json = serde_json::to_value(resp).unwrap_or_default();
747    (StatusCode::OK, Json(json))
748}
749
750async fn v1_a2a_agent_card(State(state): State<AppState>) -> impl IntoResponse {
751    let card = crate::core::a2a::agent_card::build_agent_card(&state.project_root);
752    (
753        StatusCode::OK,
754        [(header::CONTENT_TYPE, "application/json")],
755        Json(card),
756    )
757}
758
759async fn mcp_server_card() -> impl IntoResponse {
760    let card = serde_json::json!({
761        "name": "lean-ctx",
762        "version": env!("CARGO_PKG_VERSION"),
763        "description": "Context Infrastructure Layer — compression, caching, governance for AI agents",
764        "capabilities": {
765            "tools": true,
766            "resources": false,
767            "prompts": false,
768            "sampling": false
769        },
770        "tool_categories": [
771            {"name": "file_operations", "tools": ["ctx_read", "ctx_search", "ctx_tree", "ctx_edit"], "avg_token_cost": 150},
772            {"name": "session_management", "tools": ["ctx_session", "ctx_compress", "ctx_dedup", "ctx_preload"], "avg_token_cost": 80},
773            {"name": "intelligence", "tools": ["ctx_knowledge", "ctx_semantic_search", "ctx_graph", "ctx_overview"], "avg_token_cost": 200},
774            {"name": "agent_ops", "tools": ["ctx_agent", "ctx_handoff", "ctx_task", "ctx_share"], "avg_token_cost": 120}
775        ],
776        "features": {
777            "compression": "deterministic AST-based, 40-70% token reduction",
778            "caching": "session-scoped with zstd, re-reads ~13 tokens",
779            "audit_trail": "SHA-256 chained JSONL",
780            "rbac": "5 built-in roles with capability-based access",
781            "sandboxing": "Level 0 (subprocess) + Level 1 (OS-level)",
782            "secret_detection": "8 regex patterns + custom"
783        },
784        "security": {
785            "path_jail": true,
786            "rate_limiting": true,
787            "budget_tracking": true,
788            "signed_handoffs": true,
789            "timing_safe_auth": true
790        }
791    });
792    Json(card)
793}
794
795async fn v1_agents_register(
796    State(state): State<AppState>,
797    Json(body): Json<Value>,
798) -> impl IntoResponse {
799    let agent_type = body
800        .get("agent_type")
801        .and_then(|v| v.as_str())
802        .unwrap_or("unknown");
803    let role = body.get("role").and_then(|v| v.as_str());
804    let project_root = body
805        .get("project_root")
806        .and_then(|v| v.as_str())
807        .unwrap_or(&state.project_root);
808
809    let agent_id = crate::core::agents::AgentRegistry::mutate_locked(|registry| {
810        registry.register(agent_type, role, project_root)
811    })
812    .map(|(_, id)| id)
813    .unwrap_or_default();
814
815    Json(serde_json::json!({
816        "agent_id": agent_id,
817        "status": "registered"
818    }))
819}
820
821async fn v1_agents_heartbeat(Json(body): Json<Value>) -> impl IntoResponse {
822    let agent_id = body.get("agent_id").and_then(|v| v.as_str()).unwrap_or("");
823    let _ = crate::core::agents::AgentRegistry::mutate_locked(|registry| {
824        registry.update_heartbeat(agent_id);
825    });
826    Json(serde_json::json!({"status": "ok"}))
827}
828
829async fn v1_agents_list() -> impl IntoResponse {
830    let registry = crate::core::agents::AgentRegistry::load_or_create();
831    let active = registry.list_active(None);
832    Json(serde_json::json!({
833        "agents": active.iter().map(|a| serde_json::json!({
834            "agent_id": a.agent_id,
835            "agent_type": a.agent_type,
836            "role": a.role,
837            "status": a.status.to_string(),
838            "last_active": a.last_active.to_rfc3339(),
839        })).collect::<Vec<_>>()
840    }))
841}
842
843async fn v1_agents_deregister(Json(body): Json<Value>) -> impl IntoResponse {
844    let agent_id = body.get("agent_id").and_then(|v| v.as_str()).unwrap_or("");
845    let _ = crate::core::agents::AgentRegistry::mutate_locked(|registry| {
846        registry.set_status(
847            agent_id,
848            crate::core::agents::AgentStatus::Finished,
849            Some("deregistered via API"),
850        );
851    });
852    Json(serde_json::json!({"status": "deregistered"}))
853}
854
855async fn v1_agents_events_sse()
856-> Sse<impl Stream<Item = Result<SseEvent, std::convert::Infallible>>> {
857    let stream = futures::stream::unfold(0usize, |last_count| async move {
858        loop {
859            tokio::time::sleep(Duration::from_secs(5)).await;
860            let registry = crate::core::agents::AgentRegistry::load_or_create();
861            let active = registry.list_active(None);
862            let count = active.len();
863            if count != last_count {
864                let data = serde_json::json!({
865                    "type": "agents_changed",
866                    "active_count": count,
867                    "agents": active.iter().map(|a| &a.agent_id).collect::<Vec<_>>(),
868                });
869                return Some((
870                    Ok::<_, std::convert::Infallible>(SseEvent::default().data(data.to_string())),
871                    count,
872                ));
873            }
874        }
875    });
876
877    Sse::new(stream).keep_alive(KeepAlive::new().interval(Duration::from_secs(15)))
878}
879
880fn build_app_router(cfg: &HttpServerConfig) -> Router {
881    build_app_router_with_auth(cfg, true)
882}
883
884fn build_app_router_with_auth(cfg: &HttpServerConfig, require_auth: bool) -> Router {
885    let project_root = cfg.project_root.to_string_lossy().to_string();
886    let service_project_root = project_root.clone();
887    let service_factory = move || -> Result<LeanCtxServer, std::io::Error> {
888        Ok(LeanCtxServer::new_shared_with_context(
889            &service_project_root,
890            "default",
891            "default",
892        ))
893    };
894    let mcp_http = StreamableHttpService::new(
895        service_factory,
896        Arc::new(
897            rmcp::transport::streamable_http_server::session::local::LocalSessionManager::default(),
898        ),
899        cfg.mcp_http_config(),
900    );
901
902    let rest_server = LeanCtxServer::new_shared_with_context(&project_root, "default", "default");
903
904    let state = AppState {
905        token: if require_auth {
906            cfg.effective_auth_token()
907        } else {
908            None
909        },
910        concurrency: Arc::new(tokio::sync::Semaphore::new(cfg.max_concurrency.max(1))),
911        rate: Arc::new(RateLimiter::new(cfg.max_rps, cfg.rate_burst)),
912        project_root,
913        timeout: Duration::from_millis(cfg.request_timeout_ms.max(1)),
914        server: rest_server,
915    };
916
917    Router::new()
918        .route("/health", get(health))
919        .route("/v1/shutdown", axum::routing::post(v1_shutdown))
920        .route("/v1/index/ensure", axum::routing::post(v1_index_ensure))
921        .route("/v1/manifest", get(v1_manifest))
922        .route("/v1/capabilities", get(v1_capabilities))
923        .route("/v1/openapi.json", get(v1_openapi))
924        .route("/v1/cache/stats", get(v1_cache_stats))
925        .route("/v1/tools", get(v1_tools))
926        .route("/v1/tools/call", axum::routing::post(v1_tool_call))
927        .route("/v1/events", get(v1_events))
928        .route("/v1/metrics", get(v1_metrics))
929        .route("/v1/context/summary", get(v1_context_summary))
930        .route("/v1/events/search", get(v1_events_search))
931        .route("/v1/events/lineage", get(v1_events_lineage))
932        .route("/v1/audit/events", get(v1_audit_events))
933        .route("/v1/a2a/handoff", axum::routing::post(v1_a2a_handoff))
934        .route("/v1/a2a/agent-card", get(v1_a2a_agent_card))
935        .route("/.well-known/agent.json", get(v1_a2a_agent_card))
936        .route("/.well-known/mcp-server.json", get(mcp_server_card))
937        .route("/a2a", axum::routing::post(a2a_jsonrpc))
938        .route(
939            "/v1/agents/register",
940            axum::routing::post(v1_agents_register),
941        )
942        .route(
943            "/v1/agents/heartbeat",
944            axum::routing::post(v1_agents_heartbeat),
945        )
946        .route("/v1/agents/list", get(v1_agents_list))
947        .route(
948            "/v1/agents/deregister",
949            axum::routing::post(v1_agents_deregister),
950        )
951        .route("/v1/agents/events", get(v1_agents_events_sse))
952        .route("/v1/kernel/dashboard", get(kernel_api::dashboard))
953        .route("/v1/kernel/etpao", get(kernel_api::etpao))
954        .route("/v1/kernel/config", get(kernel_api::get_config))
955        .route(
956            "/v1/kernel/config",
957            axum::routing::post(kernel_api::set_config),
958        )
959        .route("/v1/kernel/evidence", get(kernel_api::evidence))
960        .route("/v1/kernel/health", get(kernel_api::health))
961        .route("/v1/kernel/report", get(kernel_api::report))
962        .route(
963            "/v1/kernel/reset",
964            axum::routing::post(kernel_api::reset_state),
965        )
966        .merge(crate::core::ocla::wire_api::ocla_router().with_state(()))
967        .fallback_service(mcp_http)
968        .layer(axum::extract::DefaultBodyLimit::max(cfg.max_body_bytes))
969        .layer(middleware::from_fn_with_state(
970            state.clone(),
971            rate_limit_middleware,
972        ))
973        .layer(middleware::from_fn_with_state(
974            state.clone(),
975            concurrency_middleware,
976        ))
977        .layer(middleware::from_fn_with_state(
978            state.clone(),
979            auth_middleware,
980        ))
981        .with_state(state)
982}
983
984pub async fn serve(cfg: HttpServerConfig) -> Result<()> {
985    crate::core::protocol::set_mcp_context(true);
986    cfg.validate()?;
987
988    // Surface any path-jail relaxation inherited from the launch env or config,
989    // so a loosened boundary is never silent (GH security audit, finding 3).
990    crate::core::pathjail::warn_if_relaxed();
991
992    crate::core::plugins::PluginManager::init();
993    crate::core::savings_autopush::spawn_if_enabled();
994
995    // Pre-warm the project indices in the background for this long-lived HTTP
996    // server. The stdio path deliberately stays lazy — short-lived respawns must
997    // not each pay a full graph + BM25 scan (#453) — but `serve` is a single,
998    // persistent process: one background build gives the first heavy/search tool
999    // call a warm index instead of racing a cold scan of a large project root
1000    // against the per-request timeout (the SDK-conformance regression, GL #395).
1001    // The build is deduped per root and idle CPU settles flat once it completes
1002    // (the memory guard backs off), so #453 idle hygiene is preserved.
1003    let warm_root = cfg.project_root.to_string_lossy().to_string();
1004    if !warm_root.is_empty() {
1005        crate::core::index_orchestrator::ensure_all_background(&warm_root);
1006    }
1007
1008    let addr: SocketAddr = format!("{}:{}", cfg.host, cfg.port)
1009        .parse()
1010        .context("invalid host/port")?;
1011
1012    let app = build_app_router(&cfg);
1013
1014    let listener = tokio::net::TcpListener::bind(addr)
1015        .await
1016        .with_context(|| format!("bind {addr}"))?;
1017
1018    tracing::info!(
1019        "lean-ctx Streamable HTTP server listening on http://{addr} (project_root={})",
1020        cfg.project_root.display()
1021    );
1022
1023    axum::serve(listener, app)
1024        .with_graceful_shutdown(async move {
1025            let _ = tokio::signal::ctrl_c().await;
1026        })
1027        .await
1028        .context("http server")?;
1029
1030    fire_session_end();
1031    Ok(())
1032}
1033
1034/// Fire the `on_session_end` plugin hook synchronously (best-effort, bounded by
1035/// each plugin's own timeout) so listeners run before the process exits. A
1036/// no-op unless a plugin declares the hook.
1037pub(crate) fn fire_session_end() {
1038    if crate::core::plugins::PluginManager::has_listener("on_session_end") {
1039        let _ = crate::core::plugins::PluginManager::fire_hook(
1040            &crate::core::plugins::executor::HookPoint::OnSessionEnd,
1041        );
1042    }
1043}
1044
1045#[cfg(windows)]
1046impl axum::serve::Listener for crate::ipc::NamedPipeListener {
1047    type Io = tokio::net::windows::named_pipe::NamedPipeServer;
1048    type Addr = String;
1049
1050    async fn accept(&mut self) -> (Self::Io, Self::Addr) {
1051        loop {
1052            match self.accept_pipe().await {
1053                Ok(pipe) => return (pipe, self.name().to_string()),
1054                Err(e) => {
1055                    tracing::error!("named pipe accept error: {e}");
1056                    tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1057                }
1058            }
1059        }
1060    }
1061
1062    fn local_addr(&self) -> std::io::Result<Self::Addr> {
1063        Ok(self.name().to_string())
1064    }
1065}
1066
1067/// Serve the daemon over a platform-independent IPC channel (UDS on Unix,
1068/// Named Pipes on Windows).
1069pub async fn serve_ipc(cfg: HttpServerConfig, addr: crate::ipc::DaemonAddr) -> Result<()> {
1070    cfg.validate()?;
1071
1072    crate::core::plugins::PluginManager::init();
1073    crate::core::savings_autopush::spawn_if_enabled();
1074
1075    match addr {
1076        #[cfg(unix)]
1077        crate::ipc::DaemonAddr::Unix(ref path) => {
1078            let app = build_app_router_with_auth(&cfg, false);
1079            let listener = crate::ipc::bind_listener(&addr)?;
1080
1081            tracing::info!(
1082                "lean-ctx daemon listening on {} (project_root={})",
1083                path.display(),
1084                cfg.project_root.display()
1085            );
1086
1087            axum::serve(listener, app.into_make_service())
1088                .with_graceful_shutdown(async move {
1089                    let _ = tokio::signal::ctrl_c().await;
1090                })
1091                .await
1092                .context("ipc server")?;
1093            Ok(())
1094        }
1095        #[cfg(windows)]
1096        crate::ipc::DaemonAddr::NamedPipe(ref name) => {
1097            let app = build_app_router_with_auth(&cfg, false);
1098            let listener = crate::ipc::bind_listener(&addr)?;
1099
1100            tracing::info!(
1101                "lean-ctx daemon listening on {} (project_root={})",
1102                name,
1103                cfg.project_root.display()
1104            );
1105
1106            axum::serve(listener, app.into_make_service())
1107                .with_graceful_shutdown(async move {
1108                    let _ = tokio::signal::ctrl_c().await;
1109                })
1110                .await
1111                .context("ipc server")?;
1112            Ok(())
1113        }
1114    }
1115}
1116
1117#[cfg(test)]
1118mod tests {
1119    use super::*;
1120    use axum::body::Body;
1121    use axum::http::Request;
1122    use futures::StreamExt;
1123    use rmcp::transport::{StreamableHttpServerConfig, StreamableHttpService};
1124    use serde_json::json;
1125    use tower::ServiceExt;
1126
1127    async fn read_first_sse_message(body: Body) -> String {
1128        let mut stream = body.into_data_stream();
1129        let mut buf: Vec<u8> = Vec::new();
1130        for _ in 0..32 {
1131            let next = tokio::time::timeout(Duration::from_secs(2), stream.next()).await;
1132            let Ok(Some(Ok(bytes))) = next else {
1133                break;
1134            };
1135            buf.extend_from_slice(&bytes);
1136            if buf.windows(2).any(|w| w == b"\n\n") {
1137                break;
1138            }
1139        }
1140        String::from_utf8_lossy(&buf).to_string()
1141    }
1142
1143    #[test]
1144    fn index_ensure_body_parses_root_and_optional_extra_roots() {
1145        // Wire contract for the #460 daemon delegation endpoint: camelCase
1146        // `extraRoots`, optional and defaulting to empty. daemon_client serializes
1147        // exactly this shape, so a drift here silently breaks delegation.
1148        let full: IndexEnsureBody =
1149            serde_json::from_str(r#"{"root":"/a","extraRoots":["/b","/c"]}"#).unwrap();
1150        assert_eq!(full.root, "/a");
1151        assert_eq!(full.extra_roots, vec!["/b".to_string(), "/c".to_string()]);
1152
1153        let minimal: IndexEnsureBody = serde_json::from_str(r#"{"root":"/a"}"#).unwrap();
1154        assert_eq!(minimal.root, "/a");
1155        assert!(minimal.extra_roots.is_empty());
1156    }
1157
1158    #[tokio::test]
1159    async fn ipc_router_allows_local_tools_without_bearer_header() {
1160        let dir = tempfile::tempdir().expect("tempdir");
1161        let cfg = HttpServerConfig {
1162            project_root: dir.path().to_path_buf(),
1163            auth_token: Some("secret".to_string()),
1164            ..HttpServerConfig::default()
1165        };
1166        let app = build_app_router_with_auth(&cfg, false);
1167
1168        let body = json!({
1169            "name": "ctx_cache",
1170            "arguments": { "action": "stats" }
1171        })
1172        .to_string();
1173        let req = Request::builder()
1174            .method("POST")
1175            .uri("/v1/tools/call")
1176            .header("Host", "localhost")
1177            .header("Content-Type", "application/json")
1178            .body(Body::from(body))
1179            .expect("request");
1180
1181        let resp = app.oneshot(req).await.expect("resp");
1182        assert_ne!(resp.status(), StatusCode::UNAUTHORIZED);
1183    }
1184
1185    #[tokio::test]
1186    async fn auth_token_blocks_requests_without_bearer_header() {
1187        let dir = tempfile::tempdir().expect("tempdir");
1188        let root_str = dir.path().to_string_lossy().to_string();
1189        let service_project_root = root_str.clone();
1190        let service_factory = move || -> Result<LeanCtxServer, std::io::Error> {
1191            Ok(LeanCtxServer::new_shared_with_context(
1192                &service_project_root,
1193                "default",
1194                "default",
1195            ))
1196        };
1197        let cfg = StreamableHttpServerConfig::default()
1198            .with_stateful_mode(false)
1199            .with_json_response(true);
1200
1201        let mcp_http = StreamableHttpService::new(
1202            service_factory,
1203            Arc::new(
1204                rmcp::transport::streamable_http_server::session::local::LocalSessionManager::default(),
1205            ),
1206            cfg,
1207        );
1208
1209        let state = AppState {
1210            token: Some("secret".to_string()),
1211            concurrency: Arc::new(tokio::sync::Semaphore::new(4)),
1212            rate: Arc::new(RateLimiter::new(50, 100)),
1213            project_root: root_str.clone(),
1214            timeout: Duration::from_secs(30),
1215            server: LeanCtxServer::new_shared_with_context(&root_str, "default", "default"),
1216        };
1217
1218        let app = Router::new()
1219            .fallback_service(mcp_http)
1220            .layer(middleware::from_fn_with_state(
1221                state.clone(),
1222                auth_middleware,
1223            ))
1224            .with_state(state);
1225
1226        let body = json!({
1227            "jsonrpc": "2.0",
1228            "id": 1,
1229            "method": "tools/list",
1230            "params": {}
1231        })
1232        .to_string();
1233
1234        let req = Request::builder()
1235            .method("POST")
1236            .uri("/")
1237            .header("Host", "localhost")
1238            .header("Accept", "application/json, text/event-stream")
1239            .header("Content-Type", "application/json")
1240            .body(Body::from(body))
1241            .expect("request");
1242
1243        let resp = app.clone().oneshot(req).await.expect("resp");
1244        assert_eq!(resp.status(), StatusCode::UNAUTHORIZED);
1245    }
1246
1247    #[tokio::test]
1248    async fn mcp_service_factory_isolates_per_client_state() {
1249        let dir = tempfile::tempdir().expect("tempdir");
1250        let root_str = dir.path().to_string_lossy().to_string();
1251
1252        // Mirrors the serve() setup: service_factory must create a fresh server per MCP session.
1253        let service_project_root = root_str.clone();
1254        let service_factory = move || -> Result<LeanCtxServer, std::convert::Infallible> {
1255            Ok(LeanCtxServer::new_shared_with_context(
1256                &service_project_root,
1257                "default",
1258                "default",
1259            ))
1260        };
1261
1262        let s1 = service_factory().expect("server 1");
1263        let s2 = service_factory().expect("server 2");
1264
1265        // If the two servers accidentally share the same Arc-backed fields, these writes would
1266        // clobber each other. This test stays independent of rmcp's InitializeRequestParams API.
1267        *s1.client_name.write().await = "client-a".to_string();
1268        *s2.client_name.write().await = "client-b".to_string();
1269
1270        let a = s1.client_name.read().await.clone();
1271        let b = s2.client_name.read().await.clone();
1272        assert_eq!(a, "client-a");
1273        assert_eq!(b, "client-b");
1274    }
1275
1276    #[tokio::test]
1277    async fn rate_limit_returns_429_when_exhausted() {
1278        let state = AppState {
1279            token: None,
1280            concurrency: Arc::new(tokio::sync::Semaphore::new(16)),
1281            rate: Arc::new(RateLimiter::new(1, 1)),
1282            project_root: ".".to_string(),
1283            timeout: Duration::from_secs(30),
1284            server: LeanCtxServer::new_shared_with_context(".", "default", "default"),
1285        };
1286
1287        let app = Router::new()
1288            .route("/limited", get(|| async { (StatusCode::OK, "ok\n") }))
1289            .layer(middleware::from_fn_with_state(
1290                state.clone(),
1291                rate_limit_middleware,
1292            ))
1293            .with_state(state);
1294
1295        let req1 = Request::builder()
1296            .method("GET")
1297            .uri("/limited")
1298            .header("Host", "localhost")
1299            .body(Body::empty())
1300            .expect("req1");
1301        let resp1 = app.clone().oneshot(req1).await.expect("resp1");
1302        assert_eq!(resp1.status(), StatusCode::OK);
1303
1304        let req2 = Request::builder()
1305            .method("GET")
1306            .uri("/limited")
1307            .header("Host", "localhost")
1308            .body(Body::empty())
1309            .expect("req2");
1310        let resp2 = app.clone().oneshot(req2).await.expect("resp2");
1311        assert_eq!(resp2.status(), StatusCode::TOO_MANY_REQUESTS);
1312    }
1313
1314    #[tokio::test]
1315    async fn audit_events_endpoint_returns_json() {
1316        let dir = tempfile::tempdir().expect("tempdir");
1317        let root_str = dir.path().to_string_lossy().to_string();
1318
1319        let state = AppState {
1320            token: None,
1321            concurrency: Arc::new(tokio::sync::Semaphore::new(16)),
1322            rate: Arc::new(RateLimiter::new(50, 100)),
1323            project_root: root_str.clone(),
1324            timeout: Duration::from_secs(30),
1325            server: LeanCtxServer::new_shared_with_context(&root_str, "default", "default"),
1326        };
1327
1328        let app = Router::new()
1329            .route("/v1/audit/events", get(v1_audit_events))
1330            .with_state(state);
1331
1332        let req = Request::builder()
1333            .method("GET")
1334            .uri("/v1/audit/events?limit=10")
1335            .header("Host", "localhost")
1336            .body(Body::empty())
1337            .unwrap();
1338
1339        let resp = app.oneshot(req).await.unwrap();
1340        assert_eq!(resp.status(), StatusCode::OK);
1341
1342        let body = axum::body::to_bytes(resp.into_body(), 1_000_000)
1343            .await
1344            .unwrap();
1345        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1346        assert!(json.get("cross_project_events").unwrap().is_array());
1347        assert!(json.get("audit_trail").unwrap().is_array());
1348    }
1349
1350    #[tokio::test]
1351    async fn capabilities_endpoint_returns_contract() {
1352        let dir = tempfile::tempdir().expect("tempdir");
1353        let root_str = dir.path().to_string_lossy().to_string();
1354
1355        let state = AppState {
1356            token: None,
1357            concurrency: Arc::new(tokio::sync::Semaphore::new(16)),
1358            rate: Arc::new(RateLimiter::new(50, 100)),
1359            project_root: root_str.clone(),
1360            timeout: Duration::from_secs(30),
1361            server: LeanCtxServer::new_shared_with_context(&root_str, "default", "default"),
1362        };
1363
1364        let app = Router::new()
1365            .route("/v1/capabilities", get(v1_capabilities))
1366            .with_state(state);
1367
1368        let req = Request::builder()
1369            .method("GET")
1370            .uri("/v1/capabilities")
1371            .header("Host", "localhost")
1372            .body(Body::empty())
1373            .unwrap();
1374
1375        let resp = app.oneshot(req).await.unwrap();
1376        assert_eq!(resp.status(), StatusCode::OK);
1377
1378        let body = axum::body::to_bytes(resp.into_body(), 1_000_000)
1379            .await
1380            .unwrap();
1381        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1382        assert_eq!(json["contract_version"], json!(1));
1383        assert!(json["tools"]["total"].as_u64().unwrap() > 0);
1384        assert!(json["features"]["compression"].as_bool().unwrap());
1385        assert!(json["contracts"].is_object());
1386    }
1387
1388    #[tokio::test]
1389    async fn cache_stats_endpoint_returns_live_shape() {
1390        let app = Router::new().route("/v1/cache/stats", get(v1_cache_stats));
1391        let request = Request::builder()
1392            .method("GET")
1393            .uri("/v1/cache/stats")
1394            .body(Body::empty())
1395            .unwrap();
1396        let response = app.oneshot(request).await.unwrap();
1397        assert_eq!(response.status(), StatusCode::OK);
1398        let body = axum::body::to_bytes(response.into_body(), 1_000_000)
1399            .await
1400            .unwrap();
1401        let value: serde_json::Value = serde_json::from_slice(&body).unwrap();
1402        assert!(value["l1"]["entries"].is_u64());
1403        assert!(value["l2"]["hit_rate"].is_number());
1404        assert!(value["l3"]["bytes"].is_u64());
1405        assert!(value["delivery"]["references_served"].is_u64());
1406        assert!(value["by_kind"]["shell_command"].is_object());
1407    }
1408
1409    #[tokio::test]
1410    async fn openapi_endpoint_returns_spec() {
1411        let dir = tempfile::tempdir().expect("tempdir");
1412        let root_str = dir.path().to_string_lossy().to_string();
1413
1414        let state = AppState {
1415            token: None,
1416            concurrency: Arc::new(tokio::sync::Semaphore::new(16)),
1417            rate: Arc::new(RateLimiter::new(50, 100)),
1418            project_root: root_str.clone(),
1419            timeout: Duration::from_secs(30),
1420            server: LeanCtxServer::new_shared_with_context(&root_str, "default", "default"),
1421        };
1422
1423        let app = Router::new()
1424            .route("/v1/openapi.json", get(v1_openapi))
1425            .with_state(state);
1426
1427        let req = Request::builder()
1428            .method("GET")
1429            .uri("/v1/openapi.json")
1430            .header("Host", "localhost")
1431            .body(Body::empty())
1432            .unwrap();
1433
1434        let resp = app.oneshot(req).await.unwrap();
1435        assert_eq!(resp.status(), StatusCode::OK);
1436
1437        let body = axum::body::to_bytes(resp.into_body(), 1_000_000)
1438            .await
1439            .unwrap();
1440        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1441        assert_eq!(json["openapi"], json!("3.0.3"));
1442        assert!(json["paths"]["/v1/capabilities"]["get"].is_object());
1443        assert!(json["paths"]["/v1/openapi.json"]["get"].is_object());
1444    }
1445
1446    #[tokio::test]
1447    async fn events_endpoint_replays_tool_call_event() {
1448        use crate::core::context_os::{self, ContextEventKindV1};
1449
1450        let dir = tempfile::tempdir().expect("tempdir");
1451        std::fs::create_dir_all(dir.path().join(".git")).expect("git marker");
1452        std::fs::write(dir.path().join("a.txt"), "ok").expect("file");
1453        let root_str = dir.path().to_string_lossy().to_string();
1454
1455        let state = AppState {
1456            token: None,
1457            concurrency: Arc::new(tokio::sync::Semaphore::new(16)),
1458            rate: Arc::new(RateLimiter::new(50, 100)),
1459            project_root: root_str.clone(),
1460            timeout: Duration::from_secs(30),
1461            server: LeanCtxServer::new_shared_with_context(&root_str, "default", "default"),
1462        };
1463
1464        let app = Router::new()
1465            .route("/v1/events", get(v1_events))
1466            .with_state(state);
1467
1468        // Directly append an event to the bus — no fire-and-forget timing dependency.
1469        let rt = context_os::runtime();
1470        rt.bus.append(
1471            "ws1",
1472            "ch1",
1473            &ContextEventKindV1::ToolCallRecorded,
1474            Some("test-agent"),
1475            json!({"tool": "ctx_session", "action": "status"}),
1476        );
1477
1478        let req = Request::builder()
1479            .method("GET")
1480            .uri("/v1/events?workspaceId=ws1&channelId=ch1&since=0&limit=1")
1481            .header("Host", "localhost")
1482            .header("Accept", "text/event-stream")
1483            .body(Body::empty())
1484            .expect("req");
1485        let resp = app.clone().oneshot(req).await.expect("events");
1486        assert_eq!(resp.status(), StatusCode::OK);
1487
1488        let msg = read_first_sse_message(resp.into_body()).await;
1489        assert!(msg.contains("event: tool_call_recorded"), "msg={msg:?}");
1490        assert!(msg.contains("\"ws1\""), "msg={msg:?}");
1491        assert!(msg.contains("\"ch1\""), "msg={msg:?}");
1492    }
1493}