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