Skip to main content

relay_knowledge/interfaces/
web.rs

1//! Web HTTP adapter for same-origin diagnostics and static assets.
2
3#[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
38/// Builds the Web router without opening sockets.
39pub 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;