Skip to main content

relay_knowledge/application/service/
mod.rs

1use std::{path::PathBuf, sync::Arc, time::Instant};
2
3use crate::{
4    api::{
5        AgentProtocolStatus, ApiError, ApiMetadata, CodeIndexWorkerRunRequest,
6        CodeIndexWorkerRunResponse, EmbeddingProviderProbeResponse, GRAPH_CANVAS_MAX_LIMIT,
7        GraphCanvasEdge, GraphCanvasKind, GraphCanvasNode, GraphCanvasRequest, GraphCanvasResponse,
8        GraphCanvasSummary, GraphInspectionRequest, GraphInspectionResponse, HealthResponse,
9        IndexRefreshRequest, IndexRefreshResponse, IngestRequest, IngestResponse,
10        MultimodalExtractionRequest, MultimodalExtractionResponse, ProjectStatusResponse,
11        RequestContext, ServiceRecoveryReport,
12    },
13    domain::{AuditStatus, CodeParseStatusCounts, CodeRepositoryTotals, IndexKind},
14    env::EnvironmentConfig,
15    model_provider::ModelProviderConfigService,
16    observability::ObservabilityRuntime,
17    project::{
18        LINUX_SERVICE_DEFINITION_FILE_NAME, MACOS_SERVICE_DEFINITION_FILE_NAME, PROJECT_NAME,
19        WINDOWS_SERVICE_DEFINITION_FILE_NAME,
20    },
21    retrieval::provider::{EmbeddingRequest, ProviderRetryClass, embedding_provider_with_qos},
22    storage::{
23        FileIndexDiagnostics, GraphCanvasSelection, GraphCanvasStorageRequest, GraphInspection,
24        KnowledgeStore, NewAuditEvent, StorageError,
25    },
26};
27
28use storage_provider::StorageProvider;
29
30use super::{
31    RuntimeConfiguration, RuntimeConfigurationError,
32    knowledge::{
33        index_refresh::{
34            index_refresh_outcome, metadata_for_indexes, recover_index_kinds, refresh_index_kinds,
35        },
36        ingest::mutation_batch_from_request,
37        multimodal::extraction_ingest_request,
38    },
39    runtime::{agent_protocol_status, runtime_status, runtime_status_with_model_profiles},
40    update::{VersionCheckResponse, check_for_updates},
41};
42
43/// Shared application service used by CLI, Web, and future API adapters.
44#[derive(Clone)]
45pub struct RelayKnowledgeService {
46    pub(super) runtime: RuntimeConfiguration,
47    pub(super) storage: StorageProvider,
48    pub(super) health_cache: Arc<tokio::sync::RwLock<Option<HealthResponse>>>,
49    pub(super) watcher: Arc<tokio::sync::RwLock<Option<crate::watcher::WatcherHandle>>>,
50}
51
52impl RelayKnowledgeService {
53    /// Creates a service from already validated foundational configuration.
54    pub fn new(runtime: RuntimeConfiguration) -> Self {
55        Self {
56            storage: StorageProvider::configured(&runtime),
57            runtime,
58            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
59            watcher: Arc::new(tokio::sync::RwLock::new(None)),
60        }
61    }
62
63    /// Creates a service backed by an explicit store for deterministic tests.
64    pub fn with_store(runtime: RuntimeConfiguration, store: Arc<dyn KnowledgeStore>) -> Self {
65        Self {
66            runtime,
67            storage: StorageProvider::ready(store),
68            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
69            watcher: Arc::new(tokio::sync::RwLock::new(None)),
70        }
71    }
72
73    /// Creates a service by reading the current process environment once.
74    pub async fn from_process_environment() -> Result<Self, RuntimeConfigurationError> {
75        RuntimeConfiguration::from_process_environment()
76            .await
77            .map(Self::new)
78    }
79
80    /// Creates a service from a deterministic environment snapshot.
81    pub async fn from_environment(
82        environment: &EnvironmentConfig,
83    ) -> Result<Self, RuntimeConfigurationError> {
84        RuntimeConfiguration::from_environment(environment)
85            .await
86            .map(Self::new)
87    }
88
89    /// Applies network-related settings from a typed environment snapshot.
90    pub async fn refresh_network_from_environment(
91        &self,
92        environment: &EnvironmentConfig,
93    ) -> Result<(), RuntimeConfigurationError> {
94        self.runtime
95            .network
96            .refresh_from_environment(environment)
97            .map(|_| ())
98            .map_err(RuntimeConfigurationError::Network)
99    }
100
101    /// Re-reads process environment variables and applies network changes.
102    pub async fn refresh_network_from_process_environment(
103        &self,
104    ) -> Result<(), RuntimeConfigurationError> {
105        self.runtime
106            .network
107            .refresh_from_process_environment()
108            .map(|_| ())
109            .map_err(RuntimeConfigurationError::NetworkRuntime)
110    }
111
112    /// Returns the shared observability runtime for interface adapters.
113    pub fn observability(&self) -> ObservabilityRuntime {
114        self.runtime.observability.clone()
115    }
116
117    /// Returns the model provider configuration service rooted in runtime paths.
118    pub fn model_provider_config(&self) -> ModelProviderConfigService {
119        ModelProviderConfigService::new(self.runtime.paths.clone())
120    }
121
122    /// Checks configured release sources without opening graph storage.
123    pub async fn check_for_updates(&self, force_refresh: bool) -> VersionCheckResponse {
124        check_for_updates(
125            &self.runtime.paths,
126            &self.runtime.network,
127            &self.runtime.updates,
128            force_refresh,
129        )
130        .await
131    }
132
133    /// Persists a redacted agent protocol audit event through the durable sink.
134    pub async fn record_agent_audit(&self, event: AgentDurableAuditInput) -> Result<(), ApiError> {
135        let store = self.storage.get().await.map_err(storage_api_error)?;
136        store
137            .insert_audit_event(NewAuditEvent {
138                operation: event.operation,
139                interface: event.interface,
140                request_id: event.request_id,
141                trace_id: event.trace_id,
142                status: event.status,
143                actor: event.actor,
144                source_scope: event.source_scope,
145                graph_version: event.graph_version,
146                detail_json: event.detail_json,
147                message: event.message,
148                now_ms: current_time_millis(),
149            })
150            .await
151            .map(|_| ())
152            .map_err(storage_api_error)
153    }
154
155    /// Returns the current project status through the unified API contract.
156    pub async fn project_status(
157        &self,
158        context: RequestContext,
159    ) -> Result<ProjectStatusResponse, ApiError> {
160        let store = self.storage.get().await.map_err(storage_api_error)?;
161        let graph_version = store
162            .current_graph_version()
163            .await
164            .map_err(storage_api_error)?;
165
166        let model_profiles = self
167            .model_provider_config()
168            .profile_summary(&self.runtime.retrieval)
169            .await;
170
171        Ok(ProjectStatusResponse {
172            project_name: PROJECT_NAME.to_owned(),
173            metadata: ApiMetadata::graph_only(&context, graph_version),
174            runtime: runtime_status_with_model_profiles(&self.runtime, model_profiles),
175        })
176    }
177
178    /// Returns runtime diagnostics without opening or migrating graph storage.
179    pub fn runtime_diagnostics(
180        &self,
181        context: RequestContext,
182    ) -> (ProjectStatusResponse, AgentProtocolStatus) {
183        (
184            ProjectStatusResponse {
185                project_name: PROJECT_NAME.to_owned(),
186                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
187                runtime: runtime_status(&self.runtime),
188            },
189            agent_protocol_status(&self.runtime),
190        )
191    }
192
193    /// Commits evidence into graph storage and refreshes all v1 index metadata.
194    pub async fn ingest(
195        &self,
196        request: IngestRequest,
197        context: RequestContext,
198    ) -> Result<IngestResponse, ApiError> {
199        let batch = mutation_batch_from_request(request)
200            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
201        let worker_evidence = batch.evidence.clone();
202        let store = self.storage.get().await.map_err(storage_api_error)?;
203        let receipt = store
204            .commit_mutation_batch(batch)
205            .await
206            .map_err(storage_api_error)?;
207        self.queue_worker_tasks_for_evidence(&store, &worker_evidence, receipt.graph_version)
208            .await?;
209        let (indexes, metadata, index_refresh_error) = match refresh_index_kinds(
210            &store,
211            IndexKind::ALL,
212            receipt.graph_version,
213            &self.runtime.retrieval,
214        )
215        .await
216        {
217            Ok(outcome) => {
218                let metadata =
219                    metadata_for_indexes(&context, receipt.graph_version, &outcome.indexes);
220
221                (outcome.indexes, metadata, None)
222            }
223            Err(error) => (
224                Vec::new(),
225                ApiMetadata::indexed(&context, receipt.graph_version, None, None, true),
226                Some(error.message),
227            ),
228        };
229
230        Ok(IngestResponse {
231            metadata,
232            receipt,
233            indexes,
234            index_refresh_error,
235        })
236    }
237
238    /// Commits derived multimodal worker output through the same bounded ingest path.
239    pub async fn commit_multimodal_extraction(
240        &self,
241        request: MultimodalExtractionRequest,
242        context: RequestContext,
243    ) -> Result<MultimodalExtractionResponse, ApiError> {
244        let converted = extraction_ingest_request(request).map_err(ApiError::invalid_argument)?;
245        let parent_evidence_id = converted.parent_evidence_id;
246        let derived_evidence_count = converted.derived_evidence_count;
247        let response = self.ingest(converted.ingest, context).await?;
248
249        Ok(MultimodalExtractionResponse {
250            metadata: response.metadata,
251            parent_evidence_id,
252            derived_evidence_count,
253            receipt: response.receipt,
254            indexes: response.indexes,
255            index_refresh_error: response.index_refresh_error,
256        })
257    }
258
259    /// Returns graph inspection information without exposing storage internals.
260    pub async fn inspect_graph(
261        &self,
262        _request: GraphInspectionRequest,
263        context: RequestContext,
264    ) -> Result<GraphInspectionResponse, ApiError> {
265        let store = self.storage.get().await.map_err(storage_api_error)?;
266        let repository_code_totals = store
267            .code_repository_totals()
268            .await
269            .map_err(storage_api_error)?;
270        let graph = graph_with_repository_code_totals(
271            store.inspect_graph().await.map_err(storage_api_error)?,
272            &repository_code_totals,
273        );
274
275        Ok(GraphInspectionResponse {
276            metadata: ApiMetadata::graph_only(&context, graph.graph_version),
277            graph,
278            repository_code_totals,
279        })
280    }
281
282    /// Returns a bounded read-only graph canvas snapshot for the Web workspace.
283    pub async fn graph_canvas(
284        &self,
285        request: GraphCanvasRequest,
286        context: RequestContext,
287    ) -> Result<GraphCanvasResponse, ApiError> {
288        if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
289            return Err(ApiError::invalid_argument(format!(
290                "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
291            )));
292        }
293        let store = self.storage.get().await.map_err(storage_api_error)?;
294        let graph_version = store
295            .current_graph_version()
296            .await
297            .map_err(storage_api_error)?;
298        let snapshot = store
299            .graph_canvas(GraphCanvasStorageRequest {
300                selection: canvas_selection(request.kind),
301                source_scope: request.source_scope,
302                query: request.query,
303                graph_version,
304                limit: request.limit,
305            })
306            .await
307            .map_err(storage_api_error)?;
308        let node_count = snapshot.nodes.len();
309        let edge_count = snapshot.edges.len();
310
311        Ok(GraphCanvasResponse {
312            metadata: ApiMetadata::graph_only(&context, graph_version),
313            nodes: snapshot
314                .nodes
315                .into_iter()
316                .map(|node| GraphCanvasNode {
317                    id: node.id,
318                    kind: node.kind,
319                    label: node.label,
320                    subtitle: node.subtitle,
321                    source_scope: node.source_scope,
322                    graph_version: node.graph_version.get(),
323                    weight: node.weight,
324                    status: node.status,
325                    details: node.details,
326                })
327                .collect(),
328            edges: snapshot
329                .edges
330                .into_iter()
331                .map(|edge| GraphCanvasEdge {
332                    id: edge.id,
333                    kind: edge.kind,
334                    source: edge.source,
335                    target: edge.target,
336                    label: edge.label,
337                    graph_version: edge.graph_version.get(),
338                    confidence_basis_points: edge.confidence_basis_points,
339                    evidence_count: edge.evidence_count,
340                    details: edge.details,
341                })
342                .collect(),
343            summary: GraphCanvasSummary {
344                kind: request.kind,
345                node_count,
346                edge_count,
347                truncated: snapshot.truncated,
348                available_kinds: snapshot.available_kinds,
349            },
350        })
351    }
352
353    /// Refreshes derived index metadata up to the current graph version.
354    pub async fn refresh_indexes(
355        &self,
356        request: IndexRefreshRequest,
357        context: RequestContext,
358    ) -> Result<IndexRefreshResponse, ApiError> {
359        let store = self.storage.get().await.map_err(storage_api_error)?;
360        let graph_version = store
361            .current_graph_version()
362            .await
363            .map_err(storage_api_error)?;
364        let outcome = refresh_index_kinds(
365            &store,
366            request.kinds,
367            graph_version,
368            &self.runtime.retrieval,
369        )
370        .await?;
371        let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
372
373        Ok(IndexRefreshResponse {
374            metadata,
375            indexes: outcome.indexes,
376            index_cursors: outcome.cursors,
377            diagnostics: outcome.diagnostics,
378        })
379    }
380
381    /// Probes the configured remote embedding provider without exposing secrets.
382    pub async fn probe_embedding_provider(
383        &self,
384        context: RequestContext,
385    ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
386        let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
387            return Ok(EmbeddingProviderProbeResponse {
388                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
389                ok: false,
390                provider: None,
391                model: self.runtime.retrieval.vector_model.name.clone(),
392                dimension: self.runtime.retrieval.vector_model.dimension,
393                latency_ms: None,
394                error_code: Some("remote_embedding_not_configured".to_owned()),
395                error_message: Some("remote embedding provider is not configured".to_owned()),
396                retryable: Some(false),
397            });
398        };
399        let network = self.runtime.network.current();
400        let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
401            ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
402        })?;
403        let provider_name = remote.provider.as_str().to_owned();
404        let provider = embedding_provider_with_qos(
405            remote,
406            client,
407            self.runtime.network.qos_runtime(),
408            network.qos,
409        );
410        let started = Instant::now();
411        let result = provider
412            .embed(EmbeddingRequest {
413                inputs: vec!["relay-knowledge provider probe".to_owned()],
414                model: self.runtime.retrieval.vector_model.name.clone(),
415                dimension: self.runtime.retrieval.vector_model.dimension,
416            })
417            .await;
418
419        match result {
420            Ok(_) => Ok(EmbeddingProviderProbeResponse {
421                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
422                ok: true,
423                provider: Some(provider_name),
424                model: self.runtime.retrieval.vector_model.name.clone(),
425                dimension: self.runtime.retrieval.vector_model.dimension,
426                latency_ms: Some(duration_millis(started.elapsed())),
427                error_code: None,
428                error_message: None,
429                retryable: None,
430            }),
431            Err(error) => Ok(EmbeddingProviderProbeResponse {
432                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
433                ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
434                provider: Some(provider_name),
435                model: self.runtime.retrieval.vector_model.name.clone(),
436                dimension: self.runtime.retrieval.vector_model.dimension,
437                latency_ms: Some(duration_millis(started.elapsed())),
438                error_code: Some(error.code),
439                error_message: Some(error.message),
440                retryable: Some(error.retry == ProviderRetryClass::Retryable),
441            }),
442        }
443    }
444
445    /// Reconciles derived index cursors before resident service work starts.
446    pub async fn reconcile_startup_indexes(
447        &self,
448        context: RequestContext,
449    ) -> Result<ServiceRecoveryReport, ApiError> {
450        let store = self.storage.get().await.map_err(storage_api_error)?;
451        let graph_version = store
452            .current_graph_version()
453            .await
454            .map_err(storage_api_error)?;
455        let before = store.index_statuses().await.map_err(storage_api_error)?;
456        let active_before = before
457            .iter()
458            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
459            .cloned()
460            .collect::<Vec<_>>();
461        let stale_index_kinds = active_before
462            .iter()
463            .filter(|status| status.is_stale_for(graph_version))
464            .map(|status| status.kind)
465            .collect::<Vec<_>>();
466        let index_lag_max = active_before
467            .iter()
468            .map(|status| {
469                graph_version
470                    .get()
471                    .saturating_sub(status.indexed_graph_version.get())
472            })
473            .max()
474            .unwrap_or(0);
475        let outcome = if stale_index_kinds.is_empty() {
476            index_refresh_outcome(&store).await?
477        } else {
478            recover_index_kinds(
479                &store,
480                stale_index_kinds.clone(),
481                graph_version,
482                &self.runtime.retrieval,
483            )
484            .await?
485        };
486        let refreshed = outcome
487            .indexes
488            .iter()
489            .filter(|status| {
490                stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
491            })
492            .map(|status| status.kind)
493            .collect::<Vec<_>>();
494        let after = outcome.indexes;
495        let active_after = after
496            .iter()
497            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
498            .cloned()
499            .collect::<Vec<_>>();
500        let metadata = metadata_for_indexes(&context, graph_version, &active_after);
501
502        Ok(ServiceRecoveryReport {
503            metadata,
504            graph_version: graph_version.get(),
505            stale_index_kinds,
506            refreshed_index_kinds: refreshed,
507            index_lag_max,
508            task_queue_depth: outcome.diagnostics.queue_depth,
509            dead_letter_count: outcome.diagnostics.dead_letter_count,
510            heartbeat_state: "ready".to_owned(),
511        })
512    }
513
514    pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
515        self.storage.get().await
516    }
517
518    /// Returns whether graph storage has already been opened by this service.
519    pub fn storage_is_ready(&self) -> bool {
520        self.storage.ready_store().is_some()
521    }
522
523    /// Runs one split-worker preview code-index attempt through durable task leases.
524    pub async fn run_code_index_worker_preview(
525        &self,
526        request: CodeIndexWorkerRunRequest,
527        context: RequestContext,
528    ) -> Result<CodeIndexWorkerRunResponse, ApiError> {
529        let store = self.storage.get().await.map_err(storage_api_error)?;
530        let task = self
531            .run_code_index_task_once(request.task_id, context.clone())
532            .await?;
533        let graph_version = store
534            .current_graph_version()
535            .await
536            .map_err(storage_api_error)?;
537
538        Ok(CodeIndexWorkerRunResponse {
539            metadata: ApiMetadata::graph_only(&context, graph_version),
540            worker_kind: "code_index".to_owned(),
541            claimed: task.is_some(),
542            task,
543        })
544    }
545
546    /// Returns the persistent agent audit log path resolved by the path boundary.
547    pub fn agent_audit_log_path(&self) -> PathBuf {
548        self.runtime.paths.agent_audit_log_file()
549    }
550}
551
552/// Durable audit event input accepted from resident agent adapters.
553#[derive(Debug, Clone)]
554pub struct AgentDurableAuditInput {
555    pub operation: String,
556    pub interface: String,
557    pub request_id: String,
558    pub trace_id: String,
559    pub status: AuditStatus,
560    pub actor: Option<String>,
561    pub source_scope: Option<String>,
562    pub graph_version: u64,
563    pub detail_json: String,
564    pub message: Option<String>,
565}
566
567pub(super) fn current_time_millis() -> u64 {
568    std::time::SystemTime::now()
569        .duration_since(std::time::UNIX_EPOCH)
570        .map_or(0, |duration| {
571            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
572        })
573}
574
575pub(super) fn storage_api_error(error: StorageError) -> ApiError {
576    ApiError::storage_unavailable(error.to_string())
577}
578
579pub(super) async fn file_index_diagnostics_or_default(
580    store: &Arc<dyn KnowledgeStore>,
581) -> Result<FileIndexDiagnostics, ApiError> {
582    match store.file_index_diagnostics().await {
583        Ok(diagnostics) => Ok(diagnostics),
584        Err(StorageError::InvalidInput(message))
585            if message == "file index storage is unavailable" =>
586        {
587            Ok(FileIndexDiagnostics::default())
588        }
589        Err(error) => Err(storage_api_error(error)),
590    }
591}
592
593pub(super) fn graph_with_repository_code_totals(
594    mut graph: GraphInspection,
595    repository_totals: &CodeRepositoryTotals,
596) -> GraphInspection {
597    graph.code_file_count = graph
598        .code_file_count
599        .saturating_add(repository_totals.indexed_file_count);
600    graph.code_symbol_count = graph
601        .code_symbol_count
602        .saturating_add(repository_totals.symbol_count);
603    graph.code_reference_count = graph
604        .code_reference_count
605        .saturating_add(repository_totals.reference_count);
606    graph.code_chunk_count = graph
607        .code_chunk_count
608        .saturating_add(repository_totals.chunk_count);
609    graph.code_parse_status_counts = add_parse_status_counts(
610        graph.code_parse_status_counts,
611        repository_totals.parse_status_counts,
612    );
613
614    graph
615}
616
617fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
618    match kind {
619        GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
620        GraphCanvasKind::Code => GraphCanvasSelection::Code,
621        GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
622    }
623}
624
625fn add_parse_status_counts(
626    left: CodeParseStatusCounts,
627    right: CodeParseStatusCounts,
628) -> CodeParseStatusCounts {
629    CodeParseStatusCounts {
630        parsed: left.parsed.saturating_add(right.parsed),
631        partial: left.partial.saturating_add(right.partial),
632        text_only: left.text_only.saturating_add(right.text_only),
633        failed: left.failed.saturating_add(right.failed),
634    }
635}
636
637fn duration_millis(duration: std::time::Duration) -> u64 {
638    u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
639}
640
641fn service_definition_filename() -> &'static str {
642    if cfg!(target_os = "windows") {
643        WINDOWS_SERVICE_DEFINITION_FILE_NAME
644    } else if cfg!(target_os = "macos") {
645        MACOS_SERVICE_DEFINITION_FILE_NAME
646    } else {
647        LINUX_SERVICE_DEFINITION_FILE_NAME
648    }
649}
650
651mod health;
652mod lifecycle_plan;
653mod retrieval;
654mod service_status;
655mod storage_diagnostics;
656mod storage_provider;
657mod watcher;
658
659#[cfg(test)]
660#[path = "graph_only_tests.rs"]
661mod graph_only_tests;
662
663#[cfg(test)]
664#[path = "recovery_tests.rs"]
665mod recovery_tests;
666
667#[cfg(test)]
668#[path = "refresh_tests.rs"]
669mod refresh_tests;
670
671#[cfg(test)]
672#[path = "storage_tests.rs"]
673mod storage_tests;
674
675#[cfg(test)]
676#[path = "operations_tests.rs"]
677mod operations_tests;
678
679#[cfg(test)]
680#[path = "mod_tests.rs"]
681mod mod_tests;