1use std::{
2 path::PathBuf,
3 sync::{Arc, OnceLock},
4 time::Instant,
5};
6
7use serde::Serialize;
8
9use crate::{
10 api::{
11 AgentProtocolStatus, ApiError, ApiMetadata, EmbeddingProviderProbeResponse,
12 GRAPH_CANVAS_MAX_LIMIT, GraphCanvasEdge, GraphCanvasKind, GraphCanvasNode,
13 GraphCanvasRequest, GraphCanvasResponse, GraphCanvasSummary, GraphInspectionRequest,
14 GraphInspectionResponse, HealthResponse, HybridRetrievalRequest, HybridRetrievalResponse,
15 IndexRefreshRequest, IndexRefreshResponse, IngestRequest, IngestResponse,
16 MultimodalExtractionRequest, MultimodalExtractionResponse, ProjectStatusResponse,
17 RequestContext, ServiceRecoveryReport, ServiceStatusResponse,
18 },
19 domain::{
20 AuditStatus, CodeParseStatusCounts, CodeRepositoryTotals, ContextGraphPath,
21 ContextPackItem, FreshnessPolicy, FusionDiagnostics, IndexKind, ProposalState,
22 RECIPROCAL_RANK_FUSION_K, RetrievalBackendStatus, RetrievalBudgetUsed, RetrievalHit,
23 RetrievalMode, RetrievedContextPack, RetrieverSource, SourceScope,
24 },
25 env::EnvironmentConfig,
26 model_provider::ModelProviderConfigService,
27 observability::ObservabilityRuntime,
28 project::{
29 DATABASE_FILE_NAME, LINUX_SERVICE_DEFINITION_FILE_NAME, MACOS_SERVICE_DEFINITION_FILE_NAME,
30 PROJECT_NAME, WINDOWS_SERVICE_DEFINITION_FILE_NAME,
31 },
32 retrieval::{
33 RetrievalPlan,
34 provider::{EmbeddingRequest, ProviderRetryClass, embedding_provider},
35 read_model_backend_statuses,
36 },
37 storage::{
38 FileIndexDiagnostics, GraphCanvasSelection, GraphCanvasStorageRequest, GraphInspection,
39 GraphSearchRequest, KnowledgeStore, NewAuditEvent, ProposalListRequest, SqliteGraphStore,
40 StorageError,
41 },
42};
43
44use super::{
45 RuntimeConfiguration, RuntimeConfigurationError,
46 index_refresh::{
47 filter_outcome_to_read_models, index_refresh_outcome, metadata_for_indexes,
48 reconcile_index_refreshes, recover_index_kinds, refresh_index_kinds,
49 },
50 ingest::mutation_batch_from_request,
51 multimodal::extraction_ingest_request,
52 status::{agent_protocol_status, runtime_status, runtime_status_with_model_profiles},
53};
54
55#[cfg(test)]
56use super::ingest::generated_evidence_id;
57
58#[derive(Clone)]
60pub struct RelayKnowledgeService {
61 pub(super) runtime: RuntimeConfiguration,
62 pub(super) storage: StorageProvider,
63}
64
65impl RelayKnowledgeService {
66 pub fn new(runtime: RuntimeConfiguration) -> Self {
68 let database_path = runtime.paths.data_dir.join(DATABASE_FILE_NAME);
69
70 Self {
71 runtime,
72 storage: StorageProvider::sqlite(database_path),
73 }
74 }
75
76 pub fn with_store(runtime: RuntimeConfiguration, store: Arc<dyn KnowledgeStore>) -> Self {
78 Self {
79 runtime,
80 storage: StorageProvider::ready(store),
81 }
82 }
83
84 pub async fn from_process_environment() -> Result<Self, RuntimeConfigurationError> {
86 RuntimeConfiguration::from_process_environment()
87 .await
88 .map(Self::new)
89 }
90
91 pub async fn from_environment(
93 environment: &EnvironmentConfig,
94 ) -> Result<Self, RuntimeConfigurationError> {
95 RuntimeConfiguration::from_environment(environment)
96 .await
97 .map(Self::new)
98 }
99
100 pub async fn refresh_network_from_environment(
102 &self,
103 environment: &EnvironmentConfig,
104 ) -> Result<(), RuntimeConfigurationError> {
105 self.runtime
106 .network
107 .refresh_from_environment(environment)
108 .map(|_| ())
109 .map_err(RuntimeConfigurationError::Network)
110 }
111
112 pub async fn refresh_network_from_process_environment(
114 &self,
115 ) -> Result<(), RuntimeConfigurationError> {
116 self.runtime
117 .network
118 .refresh_from_process_environment()
119 .map(|_| ())
120 .map_err(RuntimeConfigurationError::NetworkRuntime)
121 }
122
123 pub fn observability(&self) -> ObservabilityRuntime {
125 self.runtime.observability.clone()
126 }
127
128 pub fn model_provider_config(&self) -> ModelProviderConfigService {
130 ModelProviderConfigService::new(self.runtime.paths.clone())
131 }
132
133 pub async fn record_agent_audit(&self, event: AgentDurableAuditInput) -> Result<(), ApiError> {
135 let store = self.storage.get().await.map_err(storage_api_error)?;
136 store
137 .insert_audit_event(NewAuditEvent {
138 operation: event.operation,
139 interface: event.interface,
140 request_id: event.request_id,
141 trace_id: event.trace_id,
142 status: event.status,
143 actor: event.actor,
144 source_scope: event.source_scope,
145 graph_version: event.graph_version,
146 detail_json: event.detail_json,
147 message: event.message,
148 now_ms: current_time_millis(),
149 })
150 .await
151 .map(|_| ())
152 .map_err(storage_api_error)
153 }
154
155 pub async fn project_status(
157 &self,
158 context: RequestContext,
159 ) -> Result<ProjectStatusResponse, ApiError> {
160 let store = self.storage.get().await.map_err(storage_api_error)?;
161 let graph_version = store
162 .current_graph_version()
163 .await
164 .map_err(storage_api_error)?;
165
166 let model_profiles = self
167 .model_provider_config()
168 .profile_summary(&self.runtime.retrieval)
169 .await;
170
171 Ok(ProjectStatusResponse {
172 project_name: PROJECT_NAME.to_owned(),
173 metadata: ApiMetadata::graph_only(&context, graph_version),
174 runtime: runtime_status_with_model_profiles(&self.runtime, model_profiles),
175 })
176 }
177
178 pub fn runtime_diagnostics(
180 &self,
181 context: RequestContext,
182 ) -> (ProjectStatusResponse, AgentProtocolStatus) {
183 (
184 ProjectStatusResponse {
185 project_name: PROJECT_NAME.to_owned(),
186 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
187 runtime: runtime_status(&self.runtime),
188 },
189 agent_protocol_status(&self.runtime),
190 )
191 }
192
193 pub async fn ingest(
195 &self,
196 request: IngestRequest,
197 context: RequestContext,
198 ) -> Result<IngestResponse, ApiError> {
199 let batch = mutation_batch_from_request(request)
200 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
201 let worker_evidence = batch.evidence.clone();
202 let store = self.storage.get().await.map_err(storage_api_error)?;
203 let receipt = store
204 .commit_mutation_batch(batch)
205 .await
206 .map_err(storage_api_error)?;
207 self.queue_worker_tasks_for_evidence(&store, &worker_evidence, receipt.graph_version)
208 .await?;
209 let (indexes, metadata, index_refresh_error) = match refresh_index_kinds(
210 &store,
211 IndexKind::ALL,
212 receipt.graph_version,
213 &self.runtime.retrieval,
214 )
215 .await
216 {
217 Ok(outcome) => {
218 let metadata =
219 metadata_for_indexes(&context, receipt.graph_version, &outcome.indexes);
220
221 (outcome.indexes, metadata, None)
222 }
223 Err(error) => (
224 Vec::new(),
225 ApiMetadata::indexed(&context, receipt.graph_version, None, None, true),
226 Some(error.message),
227 ),
228 };
229
230 Ok(IngestResponse {
231 metadata,
232 receipt,
233 indexes,
234 index_refresh_error,
235 })
236 }
237
238 pub async fn commit_multimodal_extraction(
240 &self,
241 request: MultimodalExtractionRequest,
242 context: RequestContext,
243 ) -> Result<MultimodalExtractionResponse, ApiError> {
244 let converted = extraction_ingest_request(request).map_err(ApiError::invalid_argument)?;
245 let parent_evidence_id = converted.parent_evidence_id;
246 let derived_evidence_count = converted.derived_evidence_count;
247 let response = self.ingest(converted.ingest, context).await?;
248
249 Ok(MultimodalExtractionResponse {
250 metadata: response.metadata,
251 parent_evidence_id,
252 derived_evidence_count,
253 receipt: response.receipt,
254 indexes: response.indexes,
255 index_refresh_error: response.index_refresh_error,
256 })
257 }
258
259 pub async fn retrieve_context(
261 &self,
262 request: HybridRetrievalRequest,
263 context: RequestContext,
264 ) -> Result<HybridRetrievalResponse, ApiError> {
265 let source_scope = normalize_optional_source_scope(request.source_scope)
266 .map_err(ApiError::invalid_argument)?;
267 let plan = RetrievalPlan::new(
268 request.query,
269 source_scope,
270 request.limit,
271 request.freshness,
272 )
273 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
274 let store = self.storage.get().await.map_err(storage_api_error)?;
275 let graph_version = store
276 .current_graph_version()
277 .await
278 .map_err(storage_api_error)?;
279
280 let mut retrieval_mode = RetrievalMode::Hybrid;
281 let mut indexes = Vec::new();
282 let mut metadata = ApiMetadata::graph_only(&context, graph_version);
283 let mut degraded_reasons = Vec::new();
284 let backend_statuses = if plan.freshness == FreshnessPolicy::GraphOnly {
285 retrieval_mode = RetrievalMode::GraphOnly;
286 degraded_reasons.push("graph_only freshness policy selected".to_owned());
287 Vec::new()
288 } else {
289 indexes = store.index_statuses().await.map_err(storage_api_error)?;
290 let mut active_indexes = indexes
291 .iter()
292 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
293 .cloned()
294 .collect::<Vec<_>>();
295 if plan.freshness == FreshnessPolicy::WaitUntilFresh {
296 let stale_kinds = active_indexes
297 .iter()
298 .filter(|status| status.is_stale_for(graph_version))
299 .map(|status| status.kind)
300 .collect::<Vec<_>>();
301 if !stale_kinds.is_empty() {
302 refresh_index_kinds(
303 &store,
304 stale_kinds,
305 graph_version,
306 &self.runtime.retrieval,
307 )
308 .await?;
309 indexes = store.index_statuses().await.map_err(storage_api_error)?;
310 active_indexes = indexes
311 .iter()
312 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
313 .cloned()
314 .collect();
315 }
316 }
317
318 let stale = active_indexes
319 .iter()
320 .any(|status| status.is_stale_for(graph_version));
321 metadata = metadata_for_indexes(&context, graph_version, &active_indexes);
322 if plan.freshness == FreshnessPolicy::AllowStale && stale {
323 degraded_reasons
324 .push("one or more indexes are behind the graph version".to_owned());
325 }
326 read_model_backend_statuses(&plan, graph_version, &indexes, &self.runtime.retrieval)
327 };
328 if backend_statuses
329 .iter()
330 .any(|status| status.state == crate::domain::RetrievalBackendState::Unavailable)
331 {
332 degraded_reasons.push(
333 "semantic/vector retrieval backends unavailable; using bm25, graph evidence, and code graph fallback"
334 .to_owned(),
335 );
336 }
337 let mut disabled_retriever_sources = self.runtime.retrieval.disabled_retriever_sources();
338 if plan.freshness == FreshnessPolicy::GraphOnly {
339 for source in [RetrieverSource::Semantic, RetrieverSource::Vector] {
340 if !disabled_retriever_sources.contains(&source) {
341 disabled_retriever_sources.push(source);
342 }
343 }
344 }
345 let candidate_limit = self.runtime.retrieval.rerank.candidate_limit(plan.limit);
346 let results = store
347 .search(GraphSearchRequest {
348 query: plan.query.clone(),
349 source_scope: plan.source_scope.clone(),
350 graph_version,
351 limit: candidate_limit,
352 disabled_retriever_sources,
353 })
354 .await
355 .map_err(storage_api_error)?;
356 let (mut results, mut rerank) = self.runtime.retrieval.rerank.rerank(&plan.query, results);
357 let truncated = results.len() > plan.limit;
358 results.truncate(plan.limit);
359 rerank.returned_count = results.len();
360 if rerank.degraded {
361 if let Some(reason) = &rerank.reason {
362 degraded_reasons.push(reason.clone());
363 }
364 }
365 let context_pack = RetrievedContextPack {
366 graph_version,
367 source_scope: plan.source_scope.clone(),
368 freshness: plan.freshness,
369 truncated,
370 backend_statuses: backend_statuses.clone(),
371 items: results
372 .iter()
373 .map(|hit| ContextPackItem {
374 result_id: hit.evidence_id.clone(),
375 source_scope: hit.source_scope.clone(),
376 source_path: hit.source_path.clone(),
377 source_span: hit.source_span,
378 entities: hit.entities.clone(),
379 graph_facts: hit.graph_facts.clone(),
380 graph_paths: hit
381 .graph_facts
382 .iter()
383 .map(ContextGraphPath::from_fact)
384 .collect(),
385 code_artifact: hit.code_artifact.clone(),
386 retriever_sources: hit.retriever_sources.clone(),
387 ranking: hit.ranking.clone(),
388 rerank: hit.rerank.clone(),
389 })
390 .collect(),
391 };
392 let budget_used = RetrievalBudgetUsed {
393 limit: plan.limit,
394 candidate_count: rerank.candidate_count,
395 returned_count: results.len(),
396 context_bytes: retrieval_context_bytes(&results, &context_pack, &backend_statuses),
397 };
398 let fusion = FusionDiagnostics {
399 algorithm: "reciprocal_rank_fusion".to_owned(),
400 k: RECIPROCAL_RANK_FUSION_K,
401 candidate_count: budget_used.candidate_count,
402 };
403 let degraded_reason = (!degraded_reasons.is_empty()).then(|| degraded_reasons.join("; "));
404
405 Ok(HybridRetrievalResponse {
406 metadata,
407 context_pack,
408 retrieval_mode,
409 source_scope: plan.source_scope,
410 freshness: plan.freshness,
411 results,
412 fusion,
413 rerank,
414 backend_statuses,
415 truncated,
416 budget_used,
417 degraded_reason,
418 indexes,
419 })
420 }
421
422 pub async fn inspect_graph(
424 &self,
425 _request: GraphInspectionRequest,
426 context: RequestContext,
427 ) -> Result<GraphInspectionResponse, ApiError> {
428 let store = self.storage.get().await.map_err(storage_api_error)?;
429 let repository_code_totals = store
430 .code_repository_totals()
431 .await
432 .map_err(storage_api_error)?;
433 let graph = graph_with_repository_code_totals(
434 store.inspect_graph().await.map_err(storage_api_error)?,
435 &repository_code_totals,
436 );
437
438 Ok(GraphInspectionResponse {
439 metadata: ApiMetadata::graph_only(&context, graph.graph_version),
440 graph,
441 repository_code_totals,
442 })
443 }
444
445 pub async fn graph_canvas(
447 &self,
448 request: GraphCanvasRequest,
449 context: RequestContext,
450 ) -> Result<GraphCanvasResponse, ApiError> {
451 if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
452 return Err(ApiError::invalid_argument(format!(
453 "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
454 )));
455 }
456 let store = self.storage.get().await.map_err(storage_api_error)?;
457 let graph_version = store
458 .current_graph_version()
459 .await
460 .map_err(storage_api_error)?;
461 let snapshot = store
462 .graph_canvas(GraphCanvasStorageRequest {
463 selection: canvas_selection(request.kind),
464 source_scope: request.source_scope,
465 query: request.query,
466 graph_version,
467 limit: request.limit,
468 })
469 .await
470 .map_err(storage_api_error)?;
471 let node_count = snapshot.nodes.len();
472 let edge_count = snapshot.edges.len();
473
474 Ok(GraphCanvasResponse {
475 metadata: ApiMetadata::graph_only(&context, graph_version),
476 nodes: snapshot
477 .nodes
478 .into_iter()
479 .map(|node| GraphCanvasNode {
480 id: node.id,
481 kind: node.kind,
482 label: node.label,
483 subtitle: node.subtitle,
484 source_scope: node.source_scope,
485 graph_version: node.graph_version.get(),
486 weight: node.weight,
487 status: node.status,
488 details: node.details,
489 })
490 .collect(),
491 edges: snapshot
492 .edges
493 .into_iter()
494 .map(|edge| GraphCanvasEdge {
495 id: edge.id,
496 kind: edge.kind,
497 source: edge.source,
498 target: edge.target,
499 label: edge.label,
500 graph_version: edge.graph_version.get(),
501 confidence_basis_points: edge.confidence_basis_points,
502 evidence_count: edge.evidence_count,
503 details: edge.details,
504 })
505 .collect(),
506 summary: GraphCanvasSummary {
507 kind: request.kind,
508 node_count,
509 edge_count,
510 truncated: snapshot.truncated,
511 available_kinds: snapshot.available_kinds,
512 },
513 })
514 }
515
516 pub async fn refresh_indexes(
518 &self,
519 request: IndexRefreshRequest,
520 context: RequestContext,
521 ) -> Result<IndexRefreshResponse, ApiError> {
522 let store = self.storage.get().await.map_err(storage_api_error)?;
523 let graph_version = store
524 .current_graph_version()
525 .await
526 .map_err(storage_api_error)?;
527 let outcome = refresh_index_kinds(
528 &store,
529 request.kinds,
530 graph_version,
531 &self.runtime.retrieval,
532 )
533 .await?;
534 let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
535
536 Ok(IndexRefreshResponse {
537 metadata,
538 indexes: outcome.indexes,
539 index_cursors: outcome.cursors,
540 diagnostics: outcome.diagnostics,
541 })
542 }
543
544 pub async fn health(&self, context: RequestContext) -> Result<HealthResponse, ApiError> {
546 let store = self.storage.get().await.map_err(storage_api_error)?;
547 let repository_code_totals = store
548 .code_repository_totals()
549 .await
550 .map_err(storage_api_error)?;
551 let graph = graph_with_repository_code_totals(
552 store.inspect_graph().await.map_err(storage_api_error)?,
553 &repository_code_totals,
554 );
555 reconcile_index_refreshes(&store, graph.graph_version, &self.runtime.retrieval).await?;
556 let outcome = filter_outcome_to_read_models(
557 index_refresh_outcome(&store).await?,
558 &self.runtime.retrieval,
559 );
560 let healthy = outcome
561 .indexes
562 .iter()
563 .all(|status| !status.is_stale_for(graph.graph_version));
564 let file_index = file_index_diagnostics_or_default(&store).await?;
565
566 Ok(HealthResponse {
567 metadata: metadata_for_indexes(&context, graph.graph_version, &outcome.indexes),
568 healthy,
569 graph,
570 repository_code_totals,
571 indexes: outcome.indexes,
572 index_cursors: outcome.cursors,
573 index_refresh: outcome.diagnostics,
574 file_index,
575 runtime: runtime_status_with_model_profiles(
576 &self.runtime,
577 self.model_provider_config()
578 .profile_summary(&self.runtime.retrieval)
579 .await,
580 ),
581 })
582 }
583
584 pub async fn probe_embedding_provider(
586 &self,
587 context: RequestContext,
588 ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
589 let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
590 return Ok(EmbeddingProviderProbeResponse {
591 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
592 ok: false,
593 provider: None,
594 model: self.runtime.retrieval.vector_model.name.clone(),
595 dimension: self.runtime.retrieval.vector_model.dimension,
596 latency_ms: None,
597 error_code: Some("remote_embedding_not_configured".to_owned()),
598 error_message: Some("remote embedding provider is not configured".to_owned()),
599 retryable: Some(false),
600 });
601 };
602 let network = self.runtime.network.current();
603 let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
604 ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
605 })?;
606 let provider_name = remote.provider.as_str().to_owned();
607 let provider = embedding_provider(remote, client);
608 let started = Instant::now();
609 let result = provider
610 .embed(EmbeddingRequest {
611 inputs: vec!["relay-knowledge provider probe".to_owned()],
612 model: self.runtime.retrieval.vector_model.name.clone(),
613 dimension: self.runtime.retrieval.vector_model.dimension,
614 })
615 .await;
616
617 match result {
618 Ok(_) => Ok(EmbeddingProviderProbeResponse {
619 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
620 ok: true,
621 provider: Some(provider_name),
622 model: self.runtime.retrieval.vector_model.name.clone(),
623 dimension: self.runtime.retrieval.vector_model.dimension,
624 latency_ms: Some(duration_millis(started.elapsed())),
625 error_code: None,
626 error_message: None,
627 retryable: None,
628 }),
629 Err(error) => Ok(EmbeddingProviderProbeResponse {
630 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
631 ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
632 provider: Some(provider_name),
633 model: self.runtime.retrieval.vector_model.name.clone(),
634 dimension: self.runtime.retrieval.vector_model.dimension,
635 latency_ms: Some(duration_millis(started.elapsed())),
636 error_code: Some(error.code),
637 error_message: Some(error.message),
638 retryable: Some(error.retry == ProviderRetryClass::Retryable),
639 }),
640 }
641 }
642
643 pub async fn reconcile_startup_indexes(
645 &self,
646 context: RequestContext,
647 ) -> Result<ServiceRecoveryReport, ApiError> {
648 let store = self.storage.get().await.map_err(storage_api_error)?;
649 let graph_version = store
650 .current_graph_version()
651 .await
652 .map_err(storage_api_error)?;
653 let before = store.index_statuses().await.map_err(storage_api_error)?;
654 let active_before = before
655 .iter()
656 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
657 .cloned()
658 .collect::<Vec<_>>();
659 let stale_index_kinds = active_before
660 .iter()
661 .filter(|status| status.is_stale_for(graph_version))
662 .map(|status| status.kind)
663 .collect::<Vec<_>>();
664 let index_lag_max = active_before
665 .iter()
666 .map(|status| {
667 graph_version
668 .get()
669 .saturating_sub(status.indexed_graph_version.get())
670 })
671 .max()
672 .unwrap_or(0);
673 let outcome = if stale_index_kinds.is_empty() {
674 index_refresh_outcome(&store).await?
675 } else {
676 recover_index_kinds(
677 &store,
678 stale_index_kinds.clone(),
679 graph_version,
680 &self.runtime.retrieval,
681 )
682 .await?
683 };
684 let refreshed = outcome
685 .indexes
686 .iter()
687 .filter(|status| {
688 stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
689 })
690 .map(|status| status.kind)
691 .collect::<Vec<_>>();
692 let after = outcome.indexes;
693 let active_after = after
694 .iter()
695 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
696 .cloned()
697 .collect::<Vec<_>>();
698 let metadata = metadata_for_indexes(&context, graph_version, &active_after);
699
700 Ok(ServiceRecoveryReport {
701 metadata,
702 graph_version: graph_version.get(),
703 stale_index_kinds,
704 refreshed_index_kinds: refreshed,
705 index_lag_max,
706 task_queue_depth: outcome.diagnostics.queue_depth,
707 dead_letter_count: outcome.diagnostics.dead_letter_count,
708 heartbeat_state: "ready".to_owned(),
709 })
710 }
711
712 pub async fn service_status(
714 &self,
715 context: RequestContext,
716 ) -> Result<ServiceStatusResponse, ApiError> {
717 let store = self.storage.get().await.map_err(storage_api_error)?;
718 let graph_version = store
719 .current_graph_version()
720 .await
721 .map_err(storage_api_error)?;
722 let index_refresh =
723 reconcile_index_refreshes(&store, graph_version, &self.runtime.retrieval).await?;
724 let file_index = file_index_diagnostics_or_default(&store).await?;
725 let service_definition_path = self
726 .runtime
727 .paths
728 .service_dir
729 .join(service_definition_filename())
730 .display()
731 .to_string();
732 let operator = store
733 .service_operator_status()
734 .await
735 .map_err(storage_api_error)?;
736 let workers = super::operations::overlay_worker_runtime(
737 store.worker_statuses().await.map_err(storage_api_error)?,
738 &self.runtime.workers,
739 );
740 let proposal_backlog = store
741 .list_proposals(ProposalListRequest {
742 state: Some(ProposalState::Proposed),
743 limit: usize::MAX,
744 })
745 .await
746 .map_err(storage_api_error)?
747 .len();
748 let audit_event_count = store.audit_event_count().await.map_err(storage_api_error)?;
749
750 Ok(ServiceStatusResponse {
751 metadata: ApiMetadata::graph_only(&context, graph_version),
752 service_name: PROJECT_NAME.to_owned(),
753 mode: operator.state.as_str().to_owned(),
754 background_enabled: operator.state != crate::domain::ServiceOperatorState::Disabled,
755 silent_updates_enabled: operator.silent_updates_enabled,
756 service_definition_path,
757 index_refresh,
758 file_index,
759 agent_protocols: agent_protocol_status(&self.runtime),
760 operator,
761 workers,
762 proposal_backlog,
763 audit_sink: crate::api::AuditSinkStatus {
764 durable: true,
765 event_count: audit_event_count,
766 last_error: None,
767 },
768 })
769 }
770
771 pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
772 self.storage.get().await
773 }
774
775 pub fn agent_audit_log_path(&self) -> PathBuf {
777 self.runtime.paths.agent_audit_log_file()
778 }
779}
780
781#[derive(Debug, Clone)]
783pub struct AgentDurableAuditInput {
784 pub operation: String,
785 pub interface: String,
786 pub request_id: String,
787 pub trace_id: String,
788 pub status: AuditStatus,
789 pub actor: Option<String>,
790 pub source_scope: Option<String>,
791 pub graph_version: u64,
792 pub detail_json: String,
793 pub message: Option<String>,
794}
795
796fn current_time_millis() -> u64 {
797 std::time::SystemTime::now()
798 .duration_since(std::time::UNIX_EPOCH)
799 .map_or(0, |duration| {
800 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
801 })
802}
803
804#[derive(Clone)]
805pub(super) struct StorageProvider {
806 path: Option<PathBuf>,
807 ready: Arc<OnceLock<Arc<dyn KnowledgeStore>>>,
808 init_lock: Arc<tokio::sync::Mutex<()>>,
809}
810
811impl StorageProvider {
812 fn sqlite(path: PathBuf) -> Self {
813 Self {
814 path: Some(path),
815 ready: Arc::new(OnceLock::new()),
816 init_lock: Arc::new(tokio::sync::Mutex::new(())),
817 }
818 }
819
820 fn ready(store: Arc<dyn KnowledgeStore>) -> Self {
821 let ready = OnceLock::new();
822 let _ = ready.set(store);
823
824 Self {
825 path: None,
826 ready: Arc::new(ready),
827 init_lock: Arc::new(tokio::sync::Mutex::new(())),
828 }
829 }
830
831 pub(super) async fn get(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
832 if let Some(store) = self.ready.get() {
833 return Ok(Arc::clone(store));
834 }
835 let _guard = self.init_lock.lock().await;
836 if let Some(store) = self.ready.get() {
837 return Ok(Arc::clone(store));
838 }
839
840 let Some(path) = self.path.clone() else {
841 return Err(StorageError::InvalidInput(
842 "storage provider was not initialized".to_owned(),
843 ));
844 };
845 let ready = Arc::clone(&self.ready);
846 tokio::task::spawn_blocking(move || {
847 if let Some(store) = ready.get() {
848 return Ok(Arc::clone(store));
849 }
850 let store = Arc::new(SqliteGraphStore::open(path)?) as Arc<dyn KnowledgeStore>;
851 let _ = ready.set(Arc::clone(&store));
852 Ok(store)
853 })
854 .await?
855 }
856}
857
858fn storage_api_error(error: StorageError) -> ApiError {
859 ApiError::storage_unavailable(error.to_string())
860}
861
862async fn file_index_diagnostics_or_default(
863 store: &Arc<dyn KnowledgeStore>,
864) -> Result<FileIndexDiagnostics, ApiError> {
865 match store.file_index_diagnostics().await {
866 Ok(diagnostics) => Ok(diagnostics),
867 Err(StorageError::InvalidInput(message))
868 if message == "file index storage is unavailable" =>
869 {
870 Ok(FileIndexDiagnostics::default())
871 }
872 Err(error) => Err(storage_api_error(error)),
873 }
874}
875
876fn normalize_optional_source_scope(value: Option<String>) -> Result<Option<String>, String> {
877 value
878 .map(|scope| {
879 SourceScope::parse(scope)
880 .map(String::from)
881 .map_err(|error| error.to_string())
882 })
883 .transpose()
884}
885
886fn retrieval_context_bytes(
887 results: &[RetrievalHit],
888 context_pack: &RetrievedContextPack,
889 backend_statuses: &[RetrievalBackendStatus],
890) -> usize {
891 serialized_context_bytes(&context_pack.backend_statuses)
892 .saturating_add(serialized_context_bytes(backend_statuses))
893 .saturating_add(results.iter().map(serialized_context_bytes).sum::<usize>())
894 .saturating_add(
895 context_pack
896 .items
897 .iter()
898 .map(serialized_context_bytes)
899 .sum::<usize>(),
900 )
901}
902
903fn graph_with_repository_code_totals(
904 mut graph: GraphInspection,
905 repository_totals: &CodeRepositoryTotals,
906) -> GraphInspection {
907 graph.code_file_count = graph
908 .code_file_count
909 .saturating_add(repository_totals.indexed_file_count);
910 graph.code_symbol_count = graph
911 .code_symbol_count
912 .saturating_add(repository_totals.symbol_count);
913 graph.code_reference_count = graph
914 .code_reference_count
915 .saturating_add(repository_totals.reference_count);
916 graph.code_chunk_count = graph
917 .code_chunk_count
918 .saturating_add(repository_totals.chunk_count);
919 graph.code_parse_status_counts = add_parse_status_counts(
920 graph.code_parse_status_counts,
921 repository_totals.parse_status_counts,
922 );
923
924 graph
925}
926
927fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
928 match kind {
929 GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
930 GraphCanvasKind::Code => GraphCanvasSelection::Code,
931 GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
932 }
933}
934
935fn add_parse_status_counts(
936 left: CodeParseStatusCounts,
937 right: CodeParseStatusCounts,
938) -> CodeParseStatusCounts {
939 CodeParseStatusCounts {
940 parsed: left.parsed.saturating_add(right.parsed),
941 partial: left.partial.saturating_add(right.partial),
942 text_only: left.text_only.saturating_add(right.text_only),
943 failed: left.failed.saturating_add(right.failed),
944 }
945}
946
947fn serialized_context_bytes<T: Serialize + ?Sized>(value: &T) -> usize {
948 serde_json::to_vec(value)
949 .map(|bytes| bytes.len())
950 .unwrap_or(usize::MAX / 4)
951}
952
953fn duration_millis(duration: std::time::Duration) -> u64 {
954 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
955}
956
957fn service_definition_filename() -> &'static str {
958 if cfg!(target_os = "windows") {
959 WINDOWS_SERVICE_DEFINITION_FILE_NAME
960 } else if cfg!(target_os = "macos") {
961 MACOS_SERVICE_DEFINITION_FILE_NAME
962 } else {
963 LINUX_SERVICE_DEFINITION_FILE_NAME
964 }
965}
966
967#[cfg(test)]
968mod id_tests;
969
970#[cfg(test)]
971mod graph_only_tests;
972
973#[cfg(test)]
974mod recovery_tests;
975
976#[cfg(test)]
977mod refresh_tests;
978
979#[cfg(test)]
980mod storage_tests;
981
982#[cfg(test)]
983mod operations_tests;
984
985#[cfg(test)]
986mod tests;