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