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, 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
43/// Builds the Web router without opening sockets.
44pub 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;