Skip to main content

relay_knowledge/interfaces/web/
mod.rs

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