Skip to main content

lean_ctx/http_server/team/
server.rs

1#[allow(clippy::wildcard_imports)]
2use super::*;
3
4pub async fn serve_team(cfg: TeamServerConfig) -> Result<()> {
5    cfg.validate_for_serve()?;
6
7    let addr: std::net::SocketAddr = format!("{}:{}", cfg.host, cfg.port)
8        .parse()
9        .context("invalid host/port")?;
10
11    let team_server = TeamCtxServer {
12        default_workspace_id: cfg.default_workspace_id.clone(),
13        roots: Arc::new(
14            cfg.workspaces
15                .iter()
16                .map(|w| (w.id.clone(), w.root.to_string_lossy().to_string()))
17                .collect(),
18        ),
19    };
20    let engine = Arc::new(TeamContextEngine::new(team_server.clone()));
21
22    let audit_file = tokio::fs::OpenOptions::new()
23        .create(true)
24        .append(true)
25        .open(&cfg.audit_log_path)
26        .await
27        .with_context(|| format!("open audit log {}", cfg.audit_log_path.display()))?;
28
29    let savings_dir = cfg
30        .audit_log_path
31        .parent()
32        .unwrap_or(std::path::Path::new("."))
33        .join("savings");
34    let workspace_roots: Vec<(String, std::path::PathBuf)> = cfg
35        .workspaces
36        .iter()
37        .map(|w| (w.id.clone(), w.root.clone()))
38        .collect();
39    let storage_roots = crate::http_server::team_billing::storage_roots_from_config(
40        &cfg.audit_log_path,
41        &workspace_roots,
42        cfg.storage_quota_bytes,
43    );
44    // Connector run state lives next to the audit log / savings store on the
45    // persistent data volume, so per-connector cursors survive redeploys.
46    let connectors_state_dir = cfg
47        .audit_log_path
48        .parent()
49        .unwrap_or(std::path::Path::new("."))
50        .join("connectors");
51    let connectors = Arc::new(cfg.connectors.clone());
52
53    // Hosted managed-connector scheduler (#281): scheduled in-process syncs of
54    // each configured source into the workspace's BM25/graph/knowledge stores,
55    // paused by the storage quota backstop (#282). A no-op with no connectors.
56    connectors::spawn_scheduler(
57        connectors.clone(),
58        team_server.roots.clone(),
59        cfg.default_workspace_id.clone(),
60        connectors_state_dir.clone(),
61        storage_roots.data_root.clone(),
62        storage_roots.quota_bytes,
63        Duration::from_mins(1),
64    );
65
66    let team = Arc::new(TeamState {
67        auth: Arc::new(cfg.tokens.clone()),
68        engine,
69        audit: Arc::new(tokio::sync::Mutex::new(audit_file)),
70        savings_store_dir: Arc::new(tokio::sync::Mutex::new(savings_dir)),
71        storage_roots,
72        storage_cache: Arc::new(tokio::sync::Mutex::new(
73            crate::http_server::team_billing::StorageCache::default(),
74        )),
75        connectors,
76        connectors_state_dir: Arc::new(connectors_state_dir),
77    });
78
79    let state = TeamAppState {
80        concurrency: Arc::new(tokio::sync::Semaphore::new(cfg.max_concurrency.max(1))),
81        rate: Arc::new(crate::http_server::RateLimiter::new(
82            cfg.max_rps,
83            cfg.rate_burst,
84        )),
85        timeout: Duration::from_millis(cfg.request_timeout_ms.max(1)),
86        team,
87        max_body_bytes: cfg.max_body_bytes,
88    };
89
90    let service_factory =
91        move || -> std::result::Result<TeamCtxServer, std::io::Error> { Ok(team_server.clone()) };
92    let mcp_http = StreamableHttpService::new(
93        service_factory,
94        Arc::new(
95            rmcp::transport::streamable_http_server::session::local::LocalSessionManager::default(),
96        ),
97        streamable_http_config(&cfg),
98    );
99
100    // Weekly team-ROI webhook (GL #388): validated at boot so a bad URL is a
101    // loud startup error, not a silent weekly no-op.
102    if let Some(url) = &cfg.roi_webhook_url {
103        crate::http_server::roi_webhook::validate_webhook_url(url)
104            .map_err(|e| anyhow!("invalid roiWebhookUrl in team config: {e}"))?;
105        crate::http_server::roi_webhook::spawn_weekly_roi_webhook(state.clone(), url.clone());
106        tracing::info!("team ROI webhook enabled (weekly)");
107    }
108
109    let app = Router::new()
110        .route("/health", get(crate::http_server::health))
111        .route("/v1/manifest", get(v1_manifest))
112        .route("/v1/tools", get(v1_tools))
113        .route("/v1/tools/call", axum::routing::post(v1_tool_call))
114        .route("/v1/events", get(v1_events))
115        .route(
116            "/v1/context/summary",
117            get(crate::http_server::context_views::v1_context_summary),
118        )
119        .route(
120            "/v1/events/search",
121            get(crate::http_server::context_views::v1_events_search),
122        )
123        .route(
124            "/v1/events/lineage",
125            get(crate::http_server::context_views::v1_event_lineage),
126        )
127        .route("/v1/metrics", get(v1_team_metrics))
128        .route(
129            "/v1/savings/summary",
130            get(crate::http_server::savings_summary::v1_savings_summary),
131        )
132        .route(
133            "/v1/savings/member/{signer}",
134            get(crate::http_server::savings_summary::v1_savings_member),
135        )
136        .route(
137            "/v1/storage",
138            get(crate::http_server::team_billing::v1_storage),
139        )
140        .route("/v1/usage", get(crate::http_server::team_billing::v1_usage))
141        .route("/v1/connectors", get(connectors::v1_connectors))
142        .route(
143            "/api/v1/savings/ingest",
144            axum::routing::post(crate::http_server::savings_ingest::v1_savings_ingest),
145        )
146        .fallback_service(mcp_http)
147        .layer(axum::extract::DefaultBodyLimit::max(cfg.max_body_bytes))
148        .layer(middleware::from_fn_with_state(
149            state.clone(),
150            team_rate_limit_middleware,
151        ))
152        .layer(middleware::from_fn_with_state(
153            state.clone(),
154            team_concurrency_middleware,
155        ))
156        .layer(middleware::from_fn_with_state(
157            state.clone(),
158            team_auth_middleware,
159        ))
160        // Outermost: SLO measurement sees the full client-observed latency.
161        .layer(middleware::from_fn(team_slo_middleware))
162        .with_state(state);
163
164    crate::core::team_slo::global().mark_started();
165
166    let listener = tokio::net::TcpListener::bind(addr)
167        .await
168        .with_context(|| format!("bind {addr}"))?;
169
170    tracing::info!(
171        "lean-ctx TEAM server listening on http://{addr} (workspaces={}, audit={})",
172        cfg.workspaces.len(),
173        cfg.audit_log_path.display()
174    );
175
176    axum::serve(listener, app)
177        .with_graceful_shutdown(async move {
178            let _ = tokio::signal::ctrl_c().await;
179        })
180        .await
181        .context("team http server")?;
182    Ok(())
183}
184
185pub fn create_token() -> Result<(String, String)> {
186    let mut bytes = [0u8; 32];
187    getrandom::fill(&mut bytes).map_err(|e| anyhow!("getrandom: {e}"))?;
188    let token = hex_lower(&bytes);
189    let hash = sha256_hex(token.as_bytes());
190    Ok((token, hash))
191}