1use 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
42use 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 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
286pub(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
351async 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 #[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
408async 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
419async 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 crate::core::pathjail::warn_if_relaxed();
991
992 crate::core::plugins::PluginManager::init();
993 crate::core::savings_autopush::spawn_if_enabled();
994
995 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
1034pub(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
1067pub 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 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 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 *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 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}