1mod assets;
4mod code;
5mod files;
6mod model_config;
7mod operation_request;
8
9use std::path::PathBuf;
10use std::sync::Arc;
11
12use axum::extract::{Query, State};
13use axum::http::StatusCode;
14use axum::response::{IntoResponse, Response};
15use axum::routing::{get, post};
16use axum::{Json, Router};
17use serde::{Deserialize, Serialize};
18use serde_json::{Value, json};
19use tower_http::limit::RequestBodyLimitLayer;
20
21use crate::{
22 api::{
23 ApiError, AuditQueryApiRequest, ErrorKind, GRAPH_CANVAS_DEFAULT_LIMIT, GraphCanvasKind,
24 GraphCanvasRequest, InterfaceKind, ProposalListApiRequest, RequestContext,
25 WorkerRunRequest, WorkerStatusRequest,
26 },
27 application::RelayKnowledgeService,
28 domain::{CodeIndexMode, ProposalState},
29};
30use assets::{asset_or_index, default_web_dist, index};
31use code::{code_index_request, code_view_request};
32use operation_request::{
33 code_context_request, code_feature_flag_request, code_impact_request, code_query_request,
34 code_register_request, code_repository_set_add_request, code_repository_set_create_request,
35 code_repository_set_query_request, code_repository_set_remove_request, code_selector,
36 code_software_request, graph_request, index_request, ingest_request, optional_bool_field,
37 optional_proposal_state, optional_string_array_field, optional_string_field,
38 optional_worker_kind, parse_freshness, proposal_decision_request, retrieve_request,
39 string_field, usize_field,
40};
41
42pub fn router(service: RelayKnowledgeService, max_request_body_bytes: u64) -> Router {
44 router_with_assets(service, default_web_dist(), max_request_body_bytes)
45}
46
47fn router_with_assets(
48 service: RelayKnowledgeService,
49 asset_root: PathBuf,
50 max_request_body_bytes: u64,
51) -> Router {
52 let state = WebState {
53 service,
54 asset_root: Arc::new(asset_root),
55 };
56 let body_limit = usize::try_from(max_request_body_bytes).unwrap_or(usize::MAX);
57
58 Router::new()
59 .route("/api/project/status", get(project_status))
60 .route("/api/health", get(health))
61 .route("/api/service/status", get(service_status))
62 .route("/api/v1/control/status", get(control_status))
63 .route("/api/v1/control/health", get(control_health))
64 .route(
65 "/api/v1/control/service/status",
66 get(read_only_service_status),
67 )
68 .route("/api/v1/control/storage/topology", get(storage_topology))
69 .merge(code::routes())
70 .route("/api/web/graph/canvas", get(graph_canvas))
71 .route("/api/web/operations/execute", post(execute_operation))
72 .merge(model_config::routes())
73 .route("/", get(index))
74 .route("/{*path}", get(asset_or_index))
75 .with_state(state)
76 .layer(RequestBodyLimitLayer::new(body_limit))
77}
78
79async fn project_status(State(state): State<WebState>) -> Response {
80 match state
81 .service
82 .project_status(RequestContext::for_interface(InterfaceKind::Web))
83 .await
84 {
85 Ok(response) => Json(response).into_response(),
86 Err(error) => api_error_response(error),
87 }
88}
89
90async fn control_status(State(state): State<WebState>) -> Response {
91 let (response, _) = state
92 .service
93 .runtime_diagnostics(RequestContext::for_interface(InterfaceKind::Web));
94
95 Json(response).into_response()
96}
97
98async fn health(State(state): State<WebState>) -> Response {
99 match state
100 .service
101 .health(RequestContext::for_interface(InterfaceKind::Web))
102 .await
103 {
104 Ok(response) => Json(response).into_response(),
105 Err(error) => api_error_response(error),
106 }
107}
108
109async fn control_health(State(state): State<WebState>) -> Response {
110 match state
111 .service
112 .read_only_health(RequestContext::for_interface(InterfaceKind::Web))
113 .await
114 {
115 Ok(response) => Json(response).into_response(),
116 Err(error) => api_error_response(error),
117 }
118}
119
120async fn service_status(State(state): State<WebState>) -> Response {
121 match state
122 .service
123 .service_status(RequestContext::for_interface(InterfaceKind::Web))
124 .await
125 {
126 Ok(response) => Json(response).into_response(),
127 Err(error) => api_error_response(error),
128 }
129}
130
131async fn read_only_service_status(State(state): State<WebState>) -> Response {
132 match state
133 .service
134 .read_only_service_status(RequestContext::for_interface(InterfaceKind::Web))
135 .await
136 {
137 Ok(response) => Json(response).into_response(),
138 Err(error) => api_error_response(error),
139 }
140}
141
142async fn storage_topology(State(state): State<WebState>) -> Response {
143 match state
144 .service
145 .storage_topology_status(RequestContext::for_interface(InterfaceKind::Web))
146 .await
147 {
148 Ok(response) => Json(response).into_response(),
149 Err(error) => api_error_response(error),
150 }
151}
152
153async fn graph_canvas(
154 State(state): State<WebState>,
155 Query(query): Query<GraphCanvasQuery>,
156) -> Response {
157 let kind = match query
158 .kind
159 .as_deref()
160 .map(GraphCanvasKind::parse)
161 .transpose()
162 {
163 Ok(kind) => kind.unwrap_or(GraphCanvasKind::Knowledge),
164 Err(message) => return WebError::bad_request(message).into_response(),
165 };
166 let request = GraphCanvasRequest {
167 kind,
168 source_scope: query.scope.and_then(non_empty_query_value),
169 query: query.query.and_then(non_empty_query_value),
170 limit: query.limit.unwrap_or(GRAPH_CANVAS_DEFAULT_LIMIT),
171 };
172
173 match state
174 .service
175 .graph_canvas(request, RequestContext::for_interface(InterfaceKind::Web))
176 .await
177 {
178 Ok(response) => Json(response).into_response(),
179 Err(error) => api_error_response(error),
180 }
181}
182
183async fn execute_operation(
184 State(state): State<WebState>,
185 Json(request): Json<ExecuteOperationRequest>,
186) -> Result<Response, WebError> {
187 let operation = string_field(&request.snapshot.payload, "operation")?;
188 let context = RequestContext::for_interface(InterfaceKind::Web);
189 let (metadata, result) = dispatch_operation(
190 &state.service,
191 operation,
192 &request.snapshot.payload,
193 context,
194 )
195 .await?;
196 let response = ExecuteOperationResponse {
197 metadata,
198 operation: operation.to_owned(),
199 name: request.snapshot.name,
200 command: request.snapshot.command,
201 result,
202 };
203
204 Ok(Json(response).into_response())
205}
206
207async fn dispatch_operation(
208 service: &RelayKnowledgeService,
209 operation: &str,
210 payload: &Value,
211 context: RequestContext,
212) -> Result<(crate::api::ApiMetadata, Value), WebError> {
213 match operation {
214 "retrieve.context" => {
215 let response = service
216 .retrieve_context(retrieve_request(payload)?, context)
217 .await?;
218 Ok((response.metadata.clone(), json!(response)))
219 }
220 "graph.ingest" => {
221 let response = service.ingest(ingest_request(payload)?, context).await?;
222 Ok((response.metadata.clone(), json!(response)))
223 }
224 "graph.inspect" => {
225 let response = service
226 .inspect_graph(graph_request(payload), context)
227 .await?;
228 Ok((response.metadata.clone(), json!(response)))
229 }
230 "index.refresh" => {
231 let response = service
232 .refresh_indexes(index_request(payload)?, context)
233 .await?;
234 Ok((response.metadata.clone(), json!(response)))
235 }
236 "files.index" | "files.query" | "files.content" => {
237 files::dispatch_file_operation(service, operation, payload, context).await
238 }
239 "service.doctor" | "service.run.streamable_http" => {
240 let response = service.service_status(context).await?;
241 Ok((response.metadata.clone(), json!(response)))
242 }
243 "provider.embedding.probe" => {
244 let response = service.probe_embedding_provider(context).await?;
245 Ok((response.metadata.clone(), json!(response)))
246 }
247 "worker.status" => {
248 let response = service
249 .worker_status(
250 WorkerStatusRequest {
251 kind: optional_worker_kind(payload)?,
252 },
253 context,
254 )
255 .await?;
256 Ok((response.metadata.clone(), json!(response)))
257 }
258 "worker.run-once" => {
259 let response = service
260 .run_worker_once(
261 WorkerRunRequest {
262 kind: optional_worker_kind(payload)?,
263 },
264 context,
265 )
266 .await?;
267 Ok((response.metadata.clone(), json!(response)))
268 }
269 "proposal.list" => {
270 let response = service
271 .list_proposals(
272 ProposalListApiRequest {
273 state: optional_proposal_state(payload)?,
274 limit: usize_field(payload, "limit")?,
275 },
276 context,
277 )
278 .await?;
279 Ok((response.metadata.clone(), json!(response)))
280 }
281 "proposal.show" => {
282 let response = service
283 .show_proposal(string_field(payload, "proposal_id")?.to_owned(), context)
284 .await?;
285 Ok((response.metadata.clone(), json!(response)))
286 }
287 "proposal.accept" => {
288 let response = service
289 .accept_proposal(
290 string_field(payload, "proposal_id")?.to_owned(),
291 proposal_decision_request(payload)?,
292 context,
293 )
294 .await?;
295 Ok((response.metadata.clone(), json!(response)))
296 }
297 "proposal.reject" => {
298 let response = service
299 .decide_proposal_without_commit(
300 string_field(payload, "proposal_id")?.to_owned(),
301 ProposalState::Rejected,
302 proposal_decision_request(payload)?,
303 context,
304 )
305 .await?;
306 Ok((response.metadata.clone(), json!(response)))
307 }
308 "proposal.supersede" => {
309 let response = service
310 .decide_proposal_without_commit(
311 string_field(payload, "proposal_id")?.to_owned(),
312 ProposalState::Superseded,
313 proposal_decision_request(payload)?,
314 context,
315 )
316 .await?;
317 Ok((response.metadata.clone(), json!(response)))
318 }
319 "audit.query" => {
320 let response = service
321 .query_audit(
322 AuditQueryApiRequest {
323 operation: optional_string_field(payload, "filter_operation"),
324 limit: usize_field(payload, "limit")?,
325 },
326 context,
327 )
328 .await?;
329 Ok((response.metadata.clone(), json!(response)))
330 }
331 "code.repo.register" => {
332 let response = service
333 .register_code_repository(code_register_request(payload)?, context)
334 .await?;
335 Ok((response.metadata.clone(), json!(response)))
336 }
337 "code.repo.list" => {
338 let response = service.list_indexed_code_repositories(context).await?;
339 Ok((response.metadata.clone(), json!(response)))
340 }
341 "code.repo.index" => {
342 let response = service
343 .start_code_repository_index(
344 code_index_request(payload, CodeIndexMode::Full)?,
345 context,
346 )
347 .await?;
348 Ok((response.metadata.clone(), json!(response)))
349 }
350 "code.repo.update" => {
351 let mode = CodeIndexMode::incremental(
352 string_field(payload, "base_ref")?,
353 string_field(payload, "head_ref")?,
354 )
355 .map_err(|error| WebError::bad_request(error.to_string()))?;
356 let response = service
357 .index_code_repository(code_index_request(payload, mode)?, context)
358 .await?;
359 Ok((response.metadata.clone(), json!(response)))
360 }
361 "code.repo.query" => {
362 let response = service
363 .query_code_repository(code_query_request(payload)?, context)
364 .await?;
365 Ok((response.metadata.clone(), json!(response)))
366 }
367 "code.repo.context" => {
368 let response = service
369 .codegraph_context(code_context_request(payload)?, context)
370 .await?;
371 Ok((response.metadata.clone(), json!(response)))
372 }
373 "code.repo.feature_flags" => {
374 let response = service
375 .query_code_repository_feature_flags(code_feature_flag_request(payload)?, context)
376 .await?;
377 Ok((response.metadata.clone(), json!(response)))
378 }
379 "code.repo.impact" => {
380 let response = service
381 .impact_code_repository(code_impact_request(payload)?, context)
382 .await?;
383 Ok((response.metadata.clone(), json!(response)))
384 }
385 "code.repo.view" => {
386 let response = service
387 .codebase_view(code_view_request(payload)?, context)
388 .await?;
389 Ok((response.metadata.clone(), json!(response)))
390 }
391 "code.repo.software" => {
392 let response = service
393 .software_global_projection(code_software_request(payload)?, context)
394 .await?;
395 Ok((response.metadata.clone(), json!(response)))
396 }
397 "code.repo.status" => {
398 let response = service
399 .code_repository_status(code_selector(payload)?, context)
400 .await?;
401 Ok((response.metadata.clone(), json!(response)))
402 }
403 "code.repo_set.create" => {
404 let response = service
405 .create_code_repository_set(code_repository_set_create_request(payload)?, context)
406 .await?;
407 Ok((response.metadata.clone(), json!(response)))
408 }
409 "code.repo_set.add" => {
410 let response = service
411 .add_code_repository_set_member(code_repository_set_add_request(payload)?, context)
412 .await?;
413 Ok((response.metadata.clone(), json!(response)))
414 }
415 "code.repo_set.remove" => {
416 let response = service
417 .remove_code_repository_set_member(
418 code_repository_set_remove_request(payload)?,
419 context,
420 )
421 .await?;
422 Ok((response.metadata.clone(), json!(response)))
423 }
424 "code.repo_set.query" => {
425 let response = service
426 .query_code_repository_set(code_repository_set_query_request(payload)?, context)
427 .await?;
428 Ok((response.metadata.clone(), json!(response)))
429 }
430 "code.repo_set.status" => {
431 let response = service
432 .code_repository_set_status(string_field(payload, "set_alias")?.to_owned(), context)
433 .await?;
434 Ok((response.metadata.clone(), json!(response)))
435 }
436 "code.repo_set.refresh" => {
437 let set_alias = string_field(payload, "set_alias")?.to_owned();
438 let response = if optional_bool_field(payload, "async")?.unwrap_or(false) {
439 service
440 .start_code_repository_set_refresh(set_alias, context)
441 .await?
442 } else {
443 service
444 .refresh_code_repository_set(set_alias, context)
445 .await?
446 };
447 Ok((response.metadata.clone(), json!(response)))
448 }
449 other => Err(WebError::bad_request(format!(
450 "unsupported web operation '{other}'"
451 ))),
452 }
453}
454
455pub(super) fn api_error_response(error: ApiError) -> Response {
456 let status = match error.error_kind {
457 ErrorKind::InvalidArgument => StatusCode::BAD_REQUEST,
458 ErrorKind::StorageUnavailable => StatusCode::SERVICE_UNAVAILABLE,
459 ErrorKind::QosRejected => StatusCode::TOO_MANY_REQUESTS,
460 ErrorKind::Timeout => StatusCode::REQUEST_TIMEOUT,
461 ErrorKind::Internal => StatusCode::INTERNAL_SERVER_ERROR,
462 };
463
464 (status, Json(error)).into_response()
465}
466
467#[derive(Debug, Deserialize)]
468struct GraphCanvasQuery {
469 kind: Option<String>,
470 scope: Option<String>,
471 query: Option<String>,
472 limit: Option<usize>,
473}
474
475fn non_empty_query_value(value: String) -> Option<String> {
476 let trimmed = value.trim();
477
478 (!trimmed.is_empty()).then(|| trimmed.to_owned())
479}
480
481#[derive(Debug, Deserialize)]
482struct ExecuteOperationRequest {
483 snapshot: WebOperationSnapshot,
484}
485
486#[derive(Debug, Deserialize)]
487struct WebOperationSnapshot {
488 name: String,
489 command: String,
490 payload: Value,
491}
492
493#[derive(Debug, Serialize)]
494struct ExecuteOperationResponse {
495 metadata: crate::api::ApiMetadata,
496 operation: String,
497 name: String,
498 command: String,
499 result: Value,
500}
501
502#[derive(Debug)]
503pub(in crate::interfaces) struct WebError {
504 status: StatusCode,
505 message: String,
506}
507
508impl WebError {
509 fn bad_request(message: String) -> Self {
510 Self {
511 status: StatusCode::BAD_REQUEST,
512 message,
513 }
514 }
515}
516
517impl From<ApiError> for WebError {
518 fn from(error: ApiError) -> Self {
519 let status = match error.error_kind {
520 ErrorKind::InvalidArgument => StatusCode::BAD_REQUEST,
521 ErrorKind::StorageUnavailable => StatusCode::SERVICE_UNAVAILABLE,
522 ErrorKind::QosRejected => StatusCode::TOO_MANY_REQUESTS,
523 ErrorKind::Timeout => StatusCode::GATEWAY_TIMEOUT,
524 ErrorKind::Internal => StatusCode::INTERNAL_SERVER_ERROR,
525 };
526
527 Self {
528 status,
529 message: error.message,
530 }
531 }
532}
533
534impl IntoResponse for WebError {
535 fn into_response(self) -> Response {
536 (self.status, Json(json!({ "error": self.message }))).into_response()
537 }
538}
539
540#[derive(Clone)]
541pub(super) struct WebState {
542 pub(super) service: RelayKnowledgeService,
543 asset_root: Arc<PathBuf>,
544}
545
546#[cfg(test)]
547#[path = "control_tests.rs"]
548mod control_tests;
549
550#[cfg(test)]
551#[path = "router_files_integration_tests.rs"]
552mod router_files_integration_tests;
553
554#[cfg(test)]
555#[path = "mod_tests.rs"]
556mod mod_tests;