Skip to main content

relay_knowledge/application/
service.rs

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/// Shared application service used by CLI, Web, and future API adapters.
60#[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    /// Creates a service from already validated foundational configuration.
69    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    /// Creates a service backed by an explicit store for deterministic tests.
80    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    /// Creates a service by reading the current process environment once.
89    pub async fn from_process_environment() -> Result<Self, RuntimeConfigurationError> {
90        RuntimeConfiguration::from_process_environment()
91            .await
92            .map(Self::new)
93    }
94
95    /// Creates a service from a deterministic environment snapshot.
96    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    /// Applies network-related settings from a typed environment snapshot.
105    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    /// Re-reads process environment variables and applies network changes.
117    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    /// Returns the shared observability runtime for interface adapters.
128    pub fn observability(&self) -> ObservabilityRuntime {
129        self.runtime.observability.clone()
130    }
131
132    /// Returns the model provider configuration service rooted in runtime paths.
133    pub fn model_provider_config(&self) -> ModelProviderConfigService {
134        ModelProviderConfigService::new(self.runtime.paths.clone())
135    }
136
137    /// Checks configured release sources without opening graph storage.
138    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    /// Persists a redacted agent protocol audit event through the durable sink.
149    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    /// Returns the current project status through the unified API contract.
171    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    /// Returns runtime diagnostics without opening or migrating graph storage.
194    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    /// Commits evidence into graph storage and refreshes all v1 index metadata.
209    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    /// Commits derived multimodal worker output through the same bounded ingest path.
254    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    /// Retrieves context through the unified hybrid retrieval contract.
275    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    /// Returns graph inspection information without exposing storage internals.
438    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    /// Returns a bounded read-only graph canvas snapshot for the Web workspace.
461    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    /// Refreshes derived index metadata up to the current graph version.
532    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    /// Probes the configured remote embedding provider without exposing secrets.
560    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    /// Reconciles derived index cursors before resident service work starts.
619    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    /// Returns the managed background service definition location and defaults.
688    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    /// Returns the persistent agent audit log path resolved by the path boundary.
751    pub fn agent_audit_log_path(&self) -> PathBuf {
752        self.runtime.paths.agent_audit_log_file()
753    }
754}
755
756/// Durable audit event input accepted from resident agent adapters.
757#[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;