Skip to main content

relay_knowledge/application/service/
mod.rs

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