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},
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}
65
66impl RelayKnowledgeService {
67    /// Creates a service from already validated foundational configuration.
68    pub fn new(runtime: RuntimeConfiguration) -> Self {
69        Self {
70            storage: StorageProvider::configured(&runtime),
71            runtime,
72            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
73        }
74    }
75
76    /// Creates a service backed by an explicit store for deterministic tests.
77    pub fn with_store(runtime: RuntimeConfiguration, store: Arc<dyn KnowledgeStore>) -> Self {
78        Self {
79            runtime,
80            storage: StorageProvider::ready(store),
81            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
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 index_cursors = Vec::new();
295        let mut index_refresh = IndexRefreshDiagnostics::default();
296        let mut metadata = ApiMetadata::graph_only(&context, graph_version);
297        let mut degraded_reasons = Vec::new();
298        let backend_statuses = if plan.freshness == FreshnessPolicy::GraphOnly {
299            retrieval_mode = RetrievalMode::GraphOnly;
300            degraded_reasons.push("graph_only freshness policy selected".to_owned());
301            Vec::new()
302        } else {
303            let mut index_outcome = retrieval_index_freshness_snapshot(&store).await?;
304            indexes = index_outcome.indexes;
305            index_cursors = index_outcome.cursors;
306            index_refresh = index_outcome.diagnostics;
307            let mut active_indexes = indexes
308                .iter()
309                .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
310                .cloned()
311                .collect::<Vec<_>>();
312            if plan.freshness == FreshnessPolicy::WaitUntilFresh {
313                let stale_kinds = active_indexes
314                    .iter()
315                    .filter(|status| status.is_stale_for(graph_version))
316                    .map(|status| status.kind)
317                    .collect::<Vec<_>>();
318                if !stale_kinds.is_empty() {
319                    refresh_index_kinds(
320                        &store,
321                        stale_kinds,
322                        graph_version,
323                        &self.runtime.retrieval,
324                    )
325                    .await?;
326                    index_outcome = retrieval_index_freshness_snapshot(&store).await?;
327                    indexes = index_outcome.indexes;
328                    index_cursors = index_outcome.cursors;
329                    index_refresh = index_outcome.diagnostics;
330                    active_indexes = indexes
331                        .iter()
332                        .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
333                        .cloned()
334                        .collect();
335                }
336            }
337
338            let stale = active_indexes
339                .iter()
340                .any(|status| status.is_stale_for(graph_version));
341            metadata = metadata_for_indexes(&context, graph_version, &active_indexes);
342            if plan.freshness == FreshnessPolicy::AllowStale && stale {
343                degraded_reasons
344                    .push("one or more indexes are behind the graph version".to_owned());
345            }
346            read_model_backend_statuses(&plan, graph_version, &indexes, &self.runtime.retrieval)
347        };
348        if backend_statuses
349            .iter()
350            .any(|status| status.state == crate::domain::RetrievalBackendState::Unavailable)
351        {
352            degraded_reasons.push(
353                "semantic/vector retrieval backends unavailable; using bm25, graph evidence, and code graph fallback"
354                    .to_owned(),
355            );
356        }
357        let mut disabled_retriever_sources = self.runtime.retrieval.disabled_retriever_sources();
358        if plan.freshness == FreshnessPolicy::GraphOnly {
359            for source in [RetrieverSource::Semantic, RetrieverSource::Vector] {
360                if !disabled_retriever_sources.contains(&source) {
361                    disabled_retriever_sources.push(source);
362                }
363            }
364        }
365        let candidate_limit = self.runtime.retrieval.rerank.candidate_limit(plan.limit);
366        let results = store
367            .search(GraphSearchRequest {
368                query: plan.query.clone(),
369                source_scope: plan.source_scope.clone(),
370                graph_version,
371                limit: candidate_limit,
372                disabled_retriever_sources,
373            })
374            .await
375            .map_err(storage_api_error)?;
376        let (mut results, mut rerank) = self.runtime.retrieval.rerank.rerank(&plan.query, results);
377        let truncated = results.len() > plan.limit;
378        results.truncate(plan.limit);
379        rerank.returned_count = results.len();
380        if rerank.degraded {
381            if let Some(reason) = &rerank.reason {
382                degraded_reasons.push(reason.clone());
383            }
384        }
385        let context_pack = RetrievedContextPack {
386            graph_version,
387            source_scope: plan.source_scope.clone(),
388            freshness: plan.freshness,
389            truncated,
390            backend_statuses: backend_statuses.clone(),
391            items: results
392                .iter()
393                .map(|hit| ContextPackItem {
394                    result_id: hit.evidence_id.clone(),
395                    source_scope: hit.source_scope.clone(),
396                    source_path: hit.source_path.clone(),
397                    source_span: hit.source_span,
398                    entities: hit.entities.clone(),
399                    graph_facts: hit.graph_facts.clone(),
400                    graph_paths: hit
401                        .graph_facts
402                        .iter()
403                        .map(ContextGraphPath::from_fact)
404                        .collect(),
405                    code_artifact: hit.code_artifact.clone(),
406                    retriever_sources: hit.retriever_sources.clone(),
407                    ranking: hit.ranking.clone(),
408                    rerank: hit.rerank.clone(),
409                })
410                .collect(),
411        };
412        let budget_used = RetrievalBudgetUsed {
413            limit: plan.limit,
414            candidate_count: rerank.candidate_count,
415            returned_count: results.len(),
416            context_bytes: retrieval_context_bytes(&results, &context_pack, &backend_statuses),
417        };
418        let fusion = FusionDiagnostics {
419            algorithm: "reciprocal_rank_fusion".to_owned(),
420            k: RECIPROCAL_RANK_FUSION_K,
421            candidate_count: budget_used.candidate_count,
422        };
423        let degraded_reason = (!degraded_reasons.is_empty()).then(|| degraded_reasons.join("; "));
424
425        Ok(HybridRetrievalResponse {
426            metadata,
427            context_pack,
428            retrieval_mode,
429            source_scope: plan.source_scope,
430            freshness: plan.freshness,
431            results,
432            fusion,
433            rerank,
434            backend_statuses,
435            truncated,
436            budget_used,
437            degraded_reason,
438            indexes,
439            index_cursors,
440            index_refresh,
441        })
442    }
443
444    /// Returns graph inspection information without exposing storage internals.
445    pub async fn inspect_graph(
446        &self,
447        _request: GraphInspectionRequest,
448        context: RequestContext,
449    ) -> Result<GraphInspectionResponse, ApiError> {
450        let store = self.storage.get().await.map_err(storage_api_error)?;
451        let repository_code_totals = store
452            .code_repository_totals()
453            .await
454            .map_err(storage_api_error)?;
455        let graph = graph_with_repository_code_totals(
456            store.inspect_graph().await.map_err(storage_api_error)?,
457            &repository_code_totals,
458        );
459
460        Ok(GraphInspectionResponse {
461            metadata: ApiMetadata::graph_only(&context, graph.graph_version),
462            graph,
463            repository_code_totals,
464        })
465    }
466
467    /// Returns a bounded read-only graph canvas snapshot for the Web workspace.
468    pub async fn graph_canvas(
469        &self,
470        request: GraphCanvasRequest,
471        context: RequestContext,
472    ) -> Result<GraphCanvasResponse, ApiError> {
473        if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
474            return Err(ApiError::invalid_argument(format!(
475                "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
476            )));
477        }
478        let store = self.storage.get().await.map_err(storage_api_error)?;
479        let graph_version = store
480            .current_graph_version()
481            .await
482            .map_err(storage_api_error)?;
483        let snapshot = store
484            .graph_canvas(GraphCanvasStorageRequest {
485                selection: canvas_selection(request.kind),
486                source_scope: request.source_scope,
487                query: request.query,
488                graph_version,
489                limit: request.limit,
490            })
491            .await
492            .map_err(storage_api_error)?;
493        let node_count = snapshot.nodes.len();
494        let edge_count = snapshot.edges.len();
495
496        Ok(GraphCanvasResponse {
497            metadata: ApiMetadata::graph_only(&context, graph_version),
498            nodes: snapshot
499                .nodes
500                .into_iter()
501                .map(|node| GraphCanvasNode {
502                    id: node.id,
503                    kind: node.kind,
504                    label: node.label,
505                    subtitle: node.subtitle,
506                    source_scope: node.source_scope,
507                    graph_version: node.graph_version.get(),
508                    weight: node.weight,
509                    status: node.status,
510                    details: node.details,
511                })
512                .collect(),
513            edges: snapshot
514                .edges
515                .into_iter()
516                .map(|edge| GraphCanvasEdge {
517                    id: edge.id,
518                    kind: edge.kind,
519                    source: edge.source,
520                    target: edge.target,
521                    label: edge.label,
522                    graph_version: edge.graph_version.get(),
523                    confidence_basis_points: edge.confidence_basis_points,
524                    evidence_count: edge.evidence_count,
525                    details: edge.details,
526                })
527                .collect(),
528            summary: GraphCanvasSummary {
529                kind: request.kind,
530                node_count,
531                edge_count,
532                truncated: snapshot.truncated,
533                available_kinds: snapshot.available_kinds,
534            },
535        })
536    }
537
538    /// Refreshes derived index metadata up to the current graph version.
539    pub async fn refresh_indexes(
540        &self,
541        request: IndexRefreshRequest,
542        context: RequestContext,
543    ) -> Result<IndexRefreshResponse, ApiError> {
544        let store = self.storage.get().await.map_err(storage_api_error)?;
545        let graph_version = store
546            .current_graph_version()
547            .await
548            .map_err(storage_api_error)?;
549        let outcome = refresh_index_kinds(
550            &store,
551            request.kinds,
552            graph_version,
553            &self.runtime.retrieval,
554        )
555        .await?;
556        let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
557
558        Ok(IndexRefreshResponse {
559            metadata,
560            indexes: outcome.indexes,
561            index_cursors: outcome.cursors,
562            diagnostics: outcome.diagnostics,
563        })
564    }
565
566    /// Probes the configured remote embedding provider without exposing secrets.
567    pub async fn probe_embedding_provider(
568        &self,
569        context: RequestContext,
570    ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
571        let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
572            return Ok(EmbeddingProviderProbeResponse {
573                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
574                ok: false,
575                provider: None,
576                model: self.runtime.retrieval.vector_model.name.clone(),
577                dimension: self.runtime.retrieval.vector_model.dimension,
578                latency_ms: None,
579                error_code: Some("remote_embedding_not_configured".to_owned()),
580                error_message: Some("remote embedding provider is not configured".to_owned()),
581                retryable: Some(false),
582            });
583        };
584        let network = self.runtime.network.current();
585        let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
586            ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
587        })?;
588        let provider_name = remote.provider.as_str().to_owned();
589        let provider = embedding_provider(remote, client);
590        let started = Instant::now();
591        let result = provider
592            .embed(EmbeddingRequest {
593                inputs: vec!["relay-knowledge provider probe".to_owned()],
594                model: self.runtime.retrieval.vector_model.name.clone(),
595                dimension: self.runtime.retrieval.vector_model.dimension,
596            })
597            .await;
598
599        match result {
600            Ok(_) => Ok(EmbeddingProviderProbeResponse {
601                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
602                ok: true,
603                provider: Some(provider_name),
604                model: self.runtime.retrieval.vector_model.name.clone(),
605                dimension: self.runtime.retrieval.vector_model.dimension,
606                latency_ms: Some(duration_millis(started.elapsed())),
607                error_code: None,
608                error_message: None,
609                retryable: None,
610            }),
611            Err(error) => Ok(EmbeddingProviderProbeResponse {
612                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
613                ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
614                provider: Some(provider_name),
615                model: self.runtime.retrieval.vector_model.name.clone(),
616                dimension: self.runtime.retrieval.vector_model.dimension,
617                latency_ms: Some(duration_millis(started.elapsed())),
618                error_code: Some(error.code),
619                error_message: Some(error.message),
620                retryable: Some(error.retry == ProviderRetryClass::Retryable),
621            }),
622        }
623    }
624
625    /// Reconciles derived index cursors before resident service work starts.
626    pub async fn reconcile_startup_indexes(
627        &self,
628        context: RequestContext,
629    ) -> Result<ServiceRecoveryReport, ApiError> {
630        let store = self.storage.get().await.map_err(storage_api_error)?;
631        let graph_version = store
632            .current_graph_version()
633            .await
634            .map_err(storage_api_error)?;
635        let before = store.index_statuses().await.map_err(storage_api_error)?;
636        let active_before = before
637            .iter()
638            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
639            .cloned()
640            .collect::<Vec<_>>();
641        let stale_index_kinds = active_before
642            .iter()
643            .filter(|status| status.is_stale_for(graph_version))
644            .map(|status| status.kind)
645            .collect::<Vec<_>>();
646        let index_lag_max = active_before
647            .iter()
648            .map(|status| {
649                graph_version
650                    .get()
651                    .saturating_sub(status.indexed_graph_version.get())
652            })
653            .max()
654            .unwrap_or(0);
655        let outcome = if stale_index_kinds.is_empty() {
656            index_refresh_outcome(&store).await?
657        } else {
658            recover_index_kinds(
659                &store,
660                stale_index_kinds.clone(),
661                graph_version,
662                &self.runtime.retrieval,
663            )
664            .await?
665        };
666        let refreshed = outcome
667            .indexes
668            .iter()
669            .filter(|status| {
670                stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
671            })
672            .map(|status| status.kind)
673            .collect::<Vec<_>>();
674        let after = outcome.indexes;
675        let active_after = after
676            .iter()
677            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
678            .cloned()
679            .collect::<Vec<_>>();
680        let metadata = metadata_for_indexes(&context, graph_version, &active_after);
681
682        Ok(ServiceRecoveryReport {
683            metadata,
684            graph_version: graph_version.get(),
685            stale_index_kinds,
686            refreshed_index_kinds: refreshed,
687            index_lag_max,
688            task_queue_depth: outcome.diagnostics.queue_depth,
689            dead_letter_count: outcome.diagnostics.dead_letter_count,
690            heartbeat_state: "ready".to_owned(),
691        })
692    }
693
694    pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
695        self.storage.get().await
696    }
697
698    /// Runs one split-worker preview code-index attempt through durable task leases.
699    pub async fn run_code_index_worker_preview(
700        &self,
701        request: CodeIndexWorkerRunRequest,
702        context: RequestContext,
703    ) -> Result<CodeIndexWorkerRunResponse, ApiError> {
704        let store = self.storage.get().await.map_err(storage_api_error)?;
705        let task = self
706            .run_code_index_task_once(request.task_id, context.clone())
707            .await?;
708        let graph_version = store
709            .current_graph_version()
710            .await
711            .map_err(storage_api_error)?;
712
713        Ok(CodeIndexWorkerRunResponse {
714            metadata: ApiMetadata::graph_only(&context, graph_version),
715            worker_kind: "code_index".to_owned(),
716            claimed: task.is_some(),
717            task,
718        })
719    }
720
721    /// Returns the persistent agent audit log path resolved by the path boundary.
722    pub fn agent_audit_log_path(&self) -> PathBuf {
723        self.runtime.paths.agent_audit_log_file()
724    }
725}
726
727/// Durable audit event input accepted from resident agent adapters.
728#[derive(Debug, Clone)]
729pub struct AgentDurableAuditInput {
730    pub operation: String,
731    pub interface: String,
732    pub request_id: String,
733    pub trace_id: String,
734    pub status: AuditStatus,
735    pub actor: Option<String>,
736    pub source_scope: Option<String>,
737    pub graph_version: u64,
738    pub detail_json: String,
739    pub message: Option<String>,
740}
741
742pub(super) fn current_time_millis() -> u64 {
743    std::time::SystemTime::now()
744        .duration_since(std::time::UNIX_EPOCH)
745        .map_or(0, |duration| {
746            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
747        })
748}
749
750pub(super) fn storage_api_error(error: StorageError) -> ApiError {
751    ApiError::storage_unavailable(error.to_string())
752}
753
754pub(super) async fn file_index_diagnostics_or_default(
755    store: &Arc<dyn KnowledgeStore>,
756) -> Result<FileIndexDiagnostics, ApiError> {
757    match store.file_index_diagnostics().await {
758        Ok(diagnostics) => Ok(diagnostics),
759        Err(StorageError::InvalidInput(message))
760            if message == "file index storage is unavailable" =>
761        {
762            Ok(FileIndexDiagnostics::default())
763        }
764        Err(error) => Err(storage_api_error(error)),
765    }
766}
767
768async fn retrieval_index_freshness_snapshot(
769    store: &Arc<dyn KnowledgeStore>,
770) -> Result<IndexRefreshOutcome, ApiError> {
771    let indexes = store.index_statuses().await.map_err(storage_api_error)?;
772    let cursors = match store.index_cursors().await {
773        Ok(cursors) => cursors,
774        Err(StorageError::InvalidInput(message))
775            if message == "index cursor storage is unavailable" =>
776        {
777            Vec::new()
778        }
779        Err(error) => return Err(storage_api_error(error)),
780    };
781    let diagnostics = match store.index_refresh_diagnostics(current_time_millis()).await {
782        Ok(diagnostics) => diagnostics,
783        Err(StorageError::InvalidInput(message))
784            if message == "index refresh diagnostics are unavailable" =>
785        {
786            IndexRefreshDiagnostics::default()
787        }
788        Err(error) => return Err(storage_api_error(error)),
789    };
790
791    Ok(IndexRefreshOutcome {
792        indexes,
793        cursors,
794        diagnostics,
795    })
796}
797
798fn normalize_optional_source_scope(value: Option<String>) -> Result<Option<String>, String> {
799    value
800        .map(|scope| {
801            SourceScope::parse(scope)
802                .map(String::from)
803                .map_err(|error| error.to_string())
804        })
805        .transpose()
806}
807
808fn retrieval_context_bytes(
809    results: &[RetrievalHit],
810    context_pack: &RetrievedContextPack,
811    backend_statuses: &[RetrievalBackendStatus],
812) -> usize {
813    serialized_context_bytes(&context_pack.backend_statuses)
814        .saturating_add(serialized_context_bytes(backend_statuses))
815        .saturating_add(results.iter().map(serialized_context_bytes).sum::<usize>())
816        .saturating_add(
817            context_pack
818                .items
819                .iter()
820                .map(serialized_context_bytes)
821                .sum::<usize>(),
822        )
823}
824
825pub(super) fn graph_with_repository_code_totals(
826    mut graph: GraphInspection,
827    repository_totals: &CodeRepositoryTotals,
828) -> GraphInspection {
829    graph.code_file_count = graph
830        .code_file_count
831        .saturating_add(repository_totals.indexed_file_count);
832    graph.code_symbol_count = graph
833        .code_symbol_count
834        .saturating_add(repository_totals.symbol_count);
835    graph.code_reference_count = graph
836        .code_reference_count
837        .saturating_add(repository_totals.reference_count);
838    graph.code_chunk_count = graph
839        .code_chunk_count
840        .saturating_add(repository_totals.chunk_count);
841    graph.code_parse_status_counts = add_parse_status_counts(
842        graph.code_parse_status_counts,
843        repository_totals.parse_status_counts,
844    );
845
846    graph
847}
848
849fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
850    match kind {
851        GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
852        GraphCanvasKind::Code => GraphCanvasSelection::Code,
853        GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
854    }
855}
856
857fn add_parse_status_counts(
858    left: CodeParseStatusCounts,
859    right: CodeParseStatusCounts,
860) -> CodeParseStatusCounts {
861    CodeParseStatusCounts {
862        parsed: left.parsed.saturating_add(right.parsed),
863        partial: left.partial.saturating_add(right.partial),
864        text_only: left.text_only.saturating_add(right.text_only),
865        failed: left.failed.saturating_add(right.failed),
866    }
867}
868
869fn serialized_context_bytes<T: Serialize + ?Sized>(value: &T) -> usize {
870    serde_json::to_vec(value)
871        .map(|bytes| bytes.len())
872        .unwrap_or(usize::MAX / 4)
873}
874
875fn duration_millis(duration: std::time::Duration) -> u64 {
876    u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
877}
878
879fn service_definition_filename() -> &'static str {
880    if cfg!(target_os = "windows") {
881        WINDOWS_SERVICE_DEFINITION_FILE_NAME
882    } else if cfg!(target_os = "macos") {
883        MACOS_SERVICE_DEFINITION_FILE_NAME
884    } else {
885        LINUX_SERVICE_DEFINITION_FILE_NAME
886    }
887}
888
889mod health;
890pub(crate) mod knowledge_map;
891mod service_status;
892mod storage_diagnostics;
893mod storage_provider;
894
895#[cfg(test)]
896mod id_tests;
897
898#[cfg(test)]
899mod graph_only_tests;
900
901#[cfg(test)]
902mod recovery_tests;
903
904#[cfg(test)]
905mod refresh_tests;
906
907#[cfg(test)]
908mod storage_tests;
909
910#[cfg(test)]
911mod operations_tests;
912
913#[cfg(test)]
914mod tests;