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