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