1#[path = "web_model_config.rs"]
4mod web_model_config;
5
6use std::{
7 path::{Component, Path, PathBuf},
8 sync::Arc,
9};
10
11use axum::{
12 Json, Router,
13 body::Body,
14 extract::{Path as AxumPath, Query, State},
15 http::{StatusCode, header},
16 response::{IntoResponse, Response},
17 routing::{get, post},
18};
19use serde::{Deserialize, Serialize};
20use serde_json::{Value, json};
21use tower_http::limit::RequestBodyLimitLayer;
22
23use crate::{
24 api::{
25 ApiError, AuditQueryApiRequest, CodeRepositoryRegisterRequest, ErrorKind, FileIndexRequest,
26 FileQueryRequest, GRAPH_CANVAS_DEFAULT_LIMIT, GraphCanvasKind, GraphCanvasRequest,
27 GraphInspectionRequest, HybridRetrievalRequest, IndexRefreshRequest, IngestEvidence,
28 IngestRequest, InterfaceKind, ProposalDecisionApiRequest, ProposalListApiRequest,
29 RequestContext, WorkerRunRequest, WorkerStatusRequest,
30 },
31 application::RelayKnowledgeService,
32 domain::{
33 CodeImpactRequest, CodeIndexMode, CodeIndexRequest, CodeQueryKind, CodeRepositorySelector,
34 CodeRetrievalRequest, FreshnessPolicy, IndexKind, ProposalState, WorkerKind,
35 },
36};
37
38pub fn router(service: RelayKnowledgeService, max_request_body_bytes: u64) -> Router {
40 router_with_assets(service, default_web_dist(), max_request_body_bytes)
41}
42
43fn router_with_assets(
44 service: RelayKnowledgeService,
45 asset_root: PathBuf,
46 max_request_body_bytes: u64,
47) -> Router {
48 let state = WebState {
49 service,
50 asset_root: Arc::new(asset_root),
51 };
52 let body_limit = usize::try_from(max_request_body_bytes).unwrap_or(usize::MAX);
53
54 Router::new()
55 .route("/api/project/status", get(project_status))
56 .route("/api/health", get(health))
57 .route("/api/service/status", get(service_status))
58 .route("/api/web/graph/canvas", get(graph_canvas))
59 .route("/api/web/operations/execute", post(execute_operation))
60 .merge(web_model_config::routes())
61 .route("/", get(index))
62 .route("/{*path}", get(asset_or_index))
63 .with_state(state)
64 .layer(RequestBodyLimitLayer::new(body_limit))
65}
66
67async fn project_status(State(state): State<WebState>) -> Response {
68 match state
69 .service
70 .project_status(RequestContext::for_interface(InterfaceKind::Web))
71 .await
72 {
73 Ok(response) => Json(response).into_response(),
74 Err(error) => api_error_response(error),
75 }
76}
77
78async fn health(State(state): State<WebState>) -> Response {
79 match state
80 .service
81 .health(RequestContext::for_interface(InterfaceKind::Web))
82 .await
83 {
84 Ok(response) => Json(response).into_response(),
85 Err(error) => api_error_response(error),
86 }
87}
88
89async fn service_status(State(state): State<WebState>) -> Response {
90 match state
91 .service
92 .service_status(RequestContext::for_interface(InterfaceKind::Web))
93 .await
94 {
95 Ok(response) => Json(response).into_response(),
96 Err(error) => api_error_response(error),
97 }
98}
99
100async fn graph_canvas(
101 State(state): State<WebState>,
102 Query(query): Query<GraphCanvasQuery>,
103) -> Response {
104 let kind = match query
105 .kind
106 .as_deref()
107 .map(GraphCanvasKind::parse)
108 .transpose()
109 {
110 Ok(kind) => kind.unwrap_or(GraphCanvasKind::Knowledge),
111 Err(message) => return WebError::bad_request(message).into_response(),
112 };
113 let request = GraphCanvasRequest {
114 kind,
115 source_scope: query.scope.and_then(non_empty_query_value),
116 query: query.query.and_then(non_empty_query_value),
117 limit: query.limit.unwrap_or(GRAPH_CANVAS_DEFAULT_LIMIT),
118 };
119
120 match state
121 .service
122 .graph_canvas(request, RequestContext::for_interface(InterfaceKind::Web))
123 .await
124 {
125 Ok(response) => Json(response).into_response(),
126 Err(error) => api_error_response(error),
127 }
128}
129
130async fn execute_operation(
131 State(state): State<WebState>,
132 Json(request): Json<ExecuteOperationRequest>,
133) -> Result<Response, WebError> {
134 let operation = string_field(&request.snapshot.payload, "operation")?;
135 let context = RequestContext::for_interface(InterfaceKind::Web);
136 let (metadata, result) = dispatch_operation(
137 &state.service,
138 operation,
139 &request.snapshot.payload,
140 context,
141 )
142 .await?;
143 let response = ExecuteOperationResponse {
144 metadata,
145 operation: operation.to_owned(),
146 name: request.snapshot.name,
147 command: request.snapshot.command,
148 result,
149 };
150
151 Ok(Json(response).into_response())
152}
153
154async fn dispatch_operation(
155 service: &RelayKnowledgeService,
156 operation: &str,
157 payload: &Value,
158 context: RequestContext,
159) -> Result<(crate::api::ApiMetadata, Value), WebError> {
160 match operation {
161 "retrieve.context" => {
162 let response = service
163 .retrieve_context(retrieve_request(payload)?, context)
164 .await?;
165 Ok((response.metadata.clone(), json!(response)))
166 }
167 "graph.ingest" => {
168 let response = service.ingest(ingest_request(payload)?, context).await?;
169 Ok((response.metadata.clone(), json!(response)))
170 }
171 "graph.inspect" => {
172 let response = service
173 .inspect_graph(graph_request(payload), context)
174 .await?;
175 Ok((response.metadata.clone(), json!(response)))
176 }
177 "index.refresh" => {
178 let response = service
179 .refresh_indexes(index_request(payload)?, context)
180 .await?;
181 Ok((response.metadata.clone(), json!(response)))
182 }
183 "files.index" => {
184 let response = service
185 .index_files(file_index_request(payload)?, context)
186 .await?;
187 Ok((response.metadata.clone(), json!(response)))
188 }
189 "files.query" => {
190 let response = service
191 .query_files(file_query_request(payload)?, context)
192 .await?;
193 Ok((response.metadata.clone(), json!(response)))
194 }
195 "service.doctor" | "service.run.streamable_http" => {
196 let response = service.service_status(context).await?;
197 Ok((response.metadata.clone(), json!(response)))
198 }
199 "provider.embedding.probe" => {
200 let response = service.probe_embedding_provider(context).await?;
201 Ok((response.metadata.clone(), json!(response)))
202 }
203 "worker.status" => {
204 let response = service
205 .worker_status(
206 WorkerStatusRequest {
207 kind: optional_worker_kind(payload)?,
208 },
209 context,
210 )
211 .await?;
212 Ok((response.metadata.clone(), json!(response)))
213 }
214 "worker.run-once" => {
215 let response = service
216 .run_worker_once(
217 WorkerRunRequest {
218 kind: optional_worker_kind(payload)?,
219 },
220 context,
221 )
222 .await?;
223 Ok((response.metadata.clone(), json!(response)))
224 }
225 "proposal.list" => {
226 let response = service
227 .list_proposals(
228 ProposalListApiRequest {
229 state: optional_proposal_state(payload)?,
230 limit: usize_field(payload, "limit")?,
231 },
232 context,
233 )
234 .await?;
235 Ok((response.metadata.clone(), json!(response)))
236 }
237 "proposal.show" => {
238 let response = service
239 .show_proposal(string_field(payload, "proposal_id")?.to_owned(), context)
240 .await?;
241 Ok((response.metadata.clone(), json!(response)))
242 }
243 "proposal.accept" => {
244 let response = service
245 .accept_proposal(
246 string_field(payload, "proposal_id")?.to_owned(),
247 proposal_decision_request(payload)?,
248 context,
249 )
250 .await?;
251 Ok((response.metadata.clone(), json!(response)))
252 }
253 "proposal.reject" => {
254 let response = service
255 .decide_proposal_without_commit(
256 string_field(payload, "proposal_id")?.to_owned(),
257 ProposalState::Rejected,
258 proposal_decision_request(payload)?,
259 context,
260 )
261 .await?;
262 Ok((response.metadata.clone(), json!(response)))
263 }
264 "proposal.supersede" => {
265 let response = service
266 .decide_proposal_without_commit(
267 string_field(payload, "proposal_id")?.to_owned(),
268 ProposalState::Superseded,
269 proposal_decision_request(payload)?,
270 context,
271 )
272 .await?;
273 Ok((response.metadata.clone(), json!(response)))
274 }
275 "audit.query" => {
276 let response = service
277 .query_audit(
278 AuditQueryApiRequest {
279 operation: optional_string_field(payload, "filter_operation"),
280 limit: usize_field(payload, "limit")?,
281 },
282 context,
283 )
284 .await?;
285 Ok((response.metadata.clone(), json!(response)))
286 }
287 "code.repo.register" => {
288 let response = service
289 .register_code_repository(code_register_request(payload)?, context)
290 .await?;
291 Ok((response.metadata.clone(), json!(response)))
292 }
293 "code.repo.index" => {
294 let response = service
295 .start_code_repository_index(
296 code_index_request(payload, CodeIndexMode::Full)?,
297 context,
298 )
299 .await?;
300 Ok((response.metadata.clone(), json!(response)))
301 }
302 "code.repo.update" => {
303 let mode = CodeIndexMode::incremental(
304 string_field(payload, "base_ref")?,
305 string_field(payload, "head_ref")?,
306 )
307 .map_err(|error| WebError::bad_request(error.to_string()))?;
308 let response = service
309 .index_code_repository(code_index_request(payload, mode)?, context)
310 .await?;
311 Ok((response.metadata.clone(), json!(response)))
312 }
313 "code.repo.query" => {
314 let response = service
315 .query_code_repository(code_query_request(payload)?, context)
316 .await?;
317 Ok((response.metadata.clone(), json!(response)))
318 }
319 "code.repo.impact" => {
320 let response = service
321 .impact_code_repository(code_impact_request(payload)?, context)
322 .await?;
323 Ok((response.metadata.clone(), json!(response)))
324 }
325 "code.repo.status" => {
326 let response = service
327 .code_repository_status(code_selector(payload)?, context)
328 .await?;
329 Ok((response.metadata.clone(), json!(response)))
330 }
331 other => Err(WebError::bad_request(format!(
332 "unsupported web operation '{other}'"
333 ))),
334 }
335}
336
337async fn index(State(state): State<WebState>) -> Response {
338 serve_file_or_status(index_path(&state.asset_root), StatusCode::NOT_FOUND).await
339}
340
341async fn asset_or_index(
342 State(state): State<WebState>,
343 AxumPath(path): AxumPath<String>,
344) -> Response {
345 if path.starts_with("api/") {
346 return (StatusCode::NOT_FOUND, Json(json!({"message": "not found"}))).into_response();
347 }
348
349 match sanitized_asset_path(&state.asset_root, &path) {
350 Some(asset_path)
351 if tokio::fs::metadata(&asset_path)
352 .await
353 .is_ok_and(|meta| meta.is_file()) =>
354 {
355 serve_file_or_status(asset_path, StatusCode::NOT_FOUND).await
356 }
357 _ => serve_file_or_status(index_path(&state.asset_root), StatusCode::NOT_FOUND).await,
358 }
359}
360
361async fn serve_file_or_status(path: PathBuf, missing_status: StatusCode) -> Response {
362 match tokio::fs::read(&path).await {
363 Ok(body) => (
364 StatusCode::OK,
365 [(header::CONTENT_TYPE, content_type(&path))],
366 Body::from(body),
367 )
368 .into_response(),
369 Err(_) => (
370 missing_status,
371 Json(json!({"message": "web assets are not built; run ./build.sh"})),
372 )
373 .into_response(),
374 }
375}
376
377fn api_error_response(error: ApiError) -> Response {
378 let status = match error.error_kind {
379 ErrorKind::InvalidArgument => StatusCode::BAD_REQUEST,
380 ErrorKind::StorageUnavailable => StatusCode::SERVICE_UNAVAILABLE,
381 ErrorKind::Timeout => StatusCode::REQUEST_TIMEOUT,
382 ErrorKind::Internal => StatusCode::INTERNAL_SERVER_ERROR,
383 };
384
385 (status, Json(error)).into_response()
386}
387
388fn sanitized_asset_path(root: &Path, requested: &str) -> Option<PathBuf> {
389 let mut path = root.to_path_buf();
390 for component in Path::new(requested).components() {
391 match component {
392 Component::Normal(segment) => path.push(segment),
393 Component::CurDir => {}
394 Component::ParentDir | Component::RootDir | Component::Prefix(_) => return None,
395 }
396 }
397
398 Some(path)
399}
400
401fn index_path(root: &Path) -> PathBuf {
402 root.join("index.html")
403}
404
405fn default_web_dist() -> PathBuf {
406 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
407 .join("web")
408 .join("dist")
409}
410
411fn content_type(path: &Path) -> &'static str {
412 match path.extension().and_then(|extension| extension.to_str()) {
413 Some("css") => "text/css; charset=utf-8",
414 Some("html") => "text/html; charset=utf-8",
415 Some("js") => "text/javascript; charset=utf-8",
416 Some("json") => "application/json",
417 Some("svg") => "image/svg+xml",
418 Some("wasm") => "application/wasm",
419 _ => "application/octet-stream",
420 }
421}
422
423fn retrieve_request(payload: &Value) -> Result<HybridRetrievalRequest, WebError> {
424 Ok(HybridRetrievalRequest {
425 query: string_field(payload, "query")?.to_owned(),
426 source_scope: optional_string_field(payload, "source_scope"),
427 freshness: parse_freshness(string_field(payload, "freshness")?)?,
428 limit: usize_field(payload, "limit")?,
429 })
430}
431
432fn ingest_request(payload: &Value) -> Result<IngestRequest, WebError> {
433 Ok(IngestRequest {
434 source_scope: string_field(payload, "source_scope")?.to_owned(),
435 evidence: vec![IngestEvidence {
436 id: None,
437 source_path: None,
438 span: None,
439 confidence: None,
440 status: None,
441 content: string_field(payload, "content")?.to_owned(),
442 entity_labels: string_array_field(payload, "entity_labels")?,
443 extraction: None,
444 }],
445 relations: Vec::new(),
446 claims: Vec::new(),
447 events: Vec::new(),
448 })
449}
450
451fn graph_request(payload: &Value) -> GraphInspectionRequest {
452 GraphInspectionRequest {
453 source_scope: optional_string_field(payload, "source_scope"),
454 }
455}
456
457fn index_request(payload: &Value) -> Result<IndexRefreshRequest, WebError> {
458 Ok(IndexRefreshRequest {
459 kinds: string_array_field(payload, "kinds")?
460 .into_iter()
461 .map(|kind| parse_index_kind(&kind))
462 .collect::<Result<Vec<_>, _>>()?,
463 })
464}
465
466fn code_register_request(payload: &Value) -> Result<CodeRepositoryRegisterRequest, WebError> {
467 Ok(CodeRepositoryRegisterRequest {
468 root_path: string_field(payload, "root_path")?.to_owned(),
469 alias: string_field(payload, "alias")?.to_owned(),
470 path_filters: optional_string_array_field(payload, "path_filters")?,
471 language_filters: optional_string_array_field(payload, "language_filters")?,
472 })
473}
474
475fn code_index_request(payload: &Value, mode: CodeIndexMode) -> Result<CodeIndexRequest, WebError> {
476 Ok(CodeIndexRequest {
477 repository: code_selector(payload)?,
478 mode,
479 freshness_policy: FreshnessPolicy::AllowStale,
480 })
481}
482
483fn file_index_request(payload: &Value) -> Result<FileIndexRequest, WebError> {
484 Ok(FileIndexRequest {
485 source_scope: optional_string_field(payload, "source_scope"),
486 roots: optional_string_array_field(payload, "roots")?,
487 })
488}
489
490fn file_query_request(payload: &Value) -> Result<FileQueryRequest, WebError> {
491 Ok(FileQueryRequest {
492 query: string_field(payload, "query")?.to_owned(),
493 source_scope: optional_string_field(payload, "source_scope"),
494 root_id: optional_string_field(payload, "root_id"),
495 limit: usize_field(payload, "limit")?,
496 })
497}
498
499fn code_query_request(payload: &Value) -> Result<CodeRetrievalRequest, WebError> {
500 CodeRetrievalRequest::new(
501 string_field(payload, "query")?,
502 code_selector(payload)?,
503 parse_code_query_kind(string_field(payload, "kind")?)?,
504 usize_field(payload, "limit")?,
505 parse_freshness(string_field(payload, "freshness")?)?,
506 )
507 .map_err(|error| WebError::bad_request(error.to_string()))
508}
509
510fn code_impact_request(payload: &Value) -> Result<CodeImpactRequest, WebError> {
511 CodeImpactRequest::new(
512 code_selector(payload)?,
513 string_field(payload, "base_ref")?,
514 string_field(payload, "head_ref")?,
515 usize_field(payload, "limit")?,
516 )
517 .map_err(|error| WebError::bad_request(error.to_string()))
518}
519
520fn code_selector(payload: &Value) -> Result<CodeRepositorySelector, WebError> {
521 CodeRepositorySelector::new(
522 string_field(payload, "alias")?,
523 optional_string_field(payload, "ref").unwrap_or_else(|| "HEAD".to_owned()),
524 optional_string_array_field(payload, "path_filters")?,
525 optional_string_array_field(payload, "language_filters")?,
526 )
527 .map_err(|error| WebError::bad_request(error.to_string()))
528}
529
530fn string_field<'a>(payload: &'a Value, field: &'static str) -> Result<&'a str, WebError> {
531 payload
532 .get(field)
533 .and_then(Value::as_str)
534 .filter(|value| !value.trim().is_empty())
535 .ok_or_else(|| WebError::bad_request(format!("{field} is required")))
536}
537
538fn optional_string_field(payload: &Value, field: &'static str) -> Option<String> {
539 payload
540 .get(field)
541 .and_then(Value::as_str)
542 .map(str::trim)
543 .filter(|value| !value.is_empty())
544 .map(ToOwned::to_owned)
545}
546
547fn string_array_field(payload: &Value, field: &'static str) -> Result<Vec<String>, WebError> {
548 payload
549 .get(field)
550 .and_then(Value::as_array)
551 .ok_or_else(|| WebError::bad_request(format!("{field} must be an array")))?
552 .iter()
553 .map(|item| {
554 item.as_str()
555 .map(str::trim)
556 .filter(|value| !value.is_empty())
557 .map(ToOwned::to_owned)
558 .ok_or_else(|| {
559 WebError::bad_request(format!("{field} contains a non-string value"))
560 })
561 })
562 .collect()
563}
564
565fn optional_string_array_field(
566 payload: &Value,
567 field: &'static str,
568) -> Result<Vec<String>, WebError> {
569 if payload.get(field).is_none() {
570 return Ok(Vec::new());
571 }
572
573 string_array_field(payload, field)
574}
575
576fn usize_field(payload: &Value, field: &'static str) -> Result<usize, WebError> {
577 payload
578 .get(field)
579 .and_then(Value::as_u64)
580 .and_then(|value| usize::try_from(value).ok())
581 .filter(|value| *value > 0)
582 .ok_or_else(|| WebError::bad_request(format!("{field} must be a positive integer")))
583}
584
585fn parse_freshness(value: &str) -> Result<FreshnessPolicy, WebError> {
586 match value {
587 "allow-stale" => Ok(FreshnessPolicy::AllowStale),
588 "wait-until-fresh" => Ok(FreshnessPolicy::WaitUntilFresh),
589 "graph-only" => Ok(FreshnessPolicy::GraphOnly),
590 other => Err(WebError::bad_request(format!(
591 "unsupported freshness '{other}'"
592 ))),
593 }
594}
595
596fn parse_index_kind(value: &str) -> Result<IndexKind, WebError> {
597 match value {
598 "bm25" => Ok(IndexKind::Bm25),
599 "semantic" => Ok(IndexKind::Semantic),
600 "vector" => Ok(IndexKind::Vector),
601 other => Err(WebError::bad_request(format!(
602 "unsupported index kind '{other}'"
603 ))),
604 }
605}
606
607fn parse_code_query_kind(value: &str) -> Result<CodeQueryKind, WebError> {
608 match value {
609 "hybrid" => Ok(CodeQueryKind::Hybrid),
610 "symbol" => Ok(CodeQueryKind::Symbol),
611 "definition" => Ok(CodeQueryKind::Definition),
612 "references" => Ok(CodeQueryKind::References),
613 "callers" => Ok(CodeQueryKind::Callers),
614 "callees" => Ok(CodeQueryKind::Callees),
615 "imports" => Ok(CodeQueryKind::Imports),
616 other => Err(WebError::bad_request(format!(
617 "unsupported code query kind '{other}'"
618 ))),
619 }
620}
621
622fn optional_worker_kind(payload: &Value) -> Result<Option<WorkerKind>, WebError> {
623 optional_string_field(payload, "kind")
624 .map(|kind| {
625 WorkerKind::parse(&kind)
626 .map_err(|_| WebError::bad_request(format!("unsupported worker kind '{kind}'")))
627 })
628 .transpose()
629}
630
631fn optional_proposal_state(payload: &Value) -> Result<Option<ProposalState>, WebError> {
632 optional_string_field(payload, "state")
633 .map(|state| {
634 ProposalState::parse(&state)
635 .map_err(|_| WebError::bad_request(format!("unsupported proposal state '{state}'")))
636 })
637 .transpose()
638}
639
640fn proposal_decision_request(payload: &Value) -> Result<ProposalDecisionApiRequest, WebError> {
641 Ok(ProposalDecisionApiRequest {
642 actor: string_field(payload, "actor")?.to_owned(),
643 reason: optional_string_field(payload, "reason"),
644 })
645}
646
647#[derive(Debug, Deserialize)]
648struct GraphCanvasQuery {
649 kind: Option<String>,
650 scope: Option<String>,
651 query: Option<String>,
652 limit: Option<usize>,
653}
654
655fn non_empty_query_value(value: String) -> Option<String> {
656 let trimmed = value.trim();
657
658 (!trimmed.is_empty()).then(|| trimmed.to_owned())
659}
660
661#[derive(Debug, Deserialize)]
662struct ExecuteOperationRequest {
663 snapshot: WebOperationSnapshot,
664}
665
666#[derive(Debug, Deserialize)]
667struct WebOperationSnapshot {
668 name: String,
669 command: String,
670 payload: Value,
671}
672
673#[derive(Debug, Serialize)]
674struct ExecuteOperationResponse {
675 metadata: crate::api::ApiMetadata,
676 operation: String,
677 name: String,
678 command: String,
679 result: Value,
680}
681
682#[derive(Debug)]
683struct WebError {
684 status: StatusCode,
685 message: String,
686}
687
688impl WebError {
689 fn bad_request(message: String) -> Self {
690 Self {
691 status: StatusCode::BAD_REQUEST,
692 message,
693 }
694 }
695}
696
697impl From<ApiError> for WebError {
698 fn from(error: ApiError) -> Self {
699 let status = match error.error_kind {
700 ErrorKind::InvalidArgument => StatusCode::BAD_REQUEST,
701 ErrorKind::StorageUnavailable => StatusCode::SERVICE_UNAVAILABLE,
702 ErrorKind::Timeout => StatusCode::GATEWAY_TIMEOUT,
703 ErrorKind::Internal => StatusCode::INTERNAL_SERVER_ERROR,
704 };
705
706 Self {
707 status,
708 message: error.message,
709 }
710 }
711}
712
713impl IntoResponse for WebError {
714 fn into_response(self) -> Response {
715 (self.status, Json(json!({ "error": self.message }))).into_response()
716 }
717}
718
719#[derive(Clone)]
720pub(super) struct WebState {
721 pub(super) service: RelayKnowledgeService,
722 asset_root: Arc<PathBuf>,
723}
724
725#[cfg(test)]
726#[path = "web_tests.rs"]
727mod tests;