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