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