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