Skip to main content

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