Skip to main content

relay_knowledge/application/service/
mod.rs

1use std::{
2    path::PathBuf,
3    sync::{Arc, atomic::AtomicUsize},
4    time::Instant,
5};
6
7use crate::{
8    api::{
9        AgentProtocolStatus, ApiError, ApiMetadata, CodeIndexWorkerRunRequest,
10        CodeIndexWorkerRunResponse, EmbeddingProviderProbeResponse, GRAPH_CANVAS_MAX_LIMIT,
11        GraphCanvasEdge, GraphCanvasKind, GraphCanvasNode, GraphCanvasRequest, GraphCanvasResponse,
12        GraphCanvasSummary, GraphInspectionRequest, GraphInspectionResponse, HealthResponse,
13        IndexRefreshRequest, IndexRefreshResponse, IngestRequest, IngestResponse,
14        MultimodalExtractionRequest, MultimodalExtractionResponse, ProjectStatusResponse,
15        RequestContext, ServiceRecoveryReport,
16    },
17    clock::system_now_millis_or_zero as current_time_millis,
18    domain::{AuditStatus, CodeParseStatusCounts, CodeRepositoryTotals, IndexKind},
19    env::EnvironmentConfig,
20    model_provider::ModelProviderConfigService,
21    observability::ObservabilityRuntime,
22    ports::{
23        embedding::{EmbeddingProvider, EmbeddingRequest, ProviderRetryClass},
24        worker_outbound::WorkerOutboundPort,
25    },
26    project::{
27        LINUX_SERVICE_DEFINITION_FILE_NAME, MACOS_SERVICE_DEFINITION_FILE_NAME, PROJECT_NAME,
28        WINDOWS_SERVICE_DEFINITION_FILE_NAME,
29    },
30    storage::{
31        FileIndexDiagnostics, GraphCanvasSelection, GraphCanvasStorageRequest, GraphInspection,
32        KnowledgeStore, KnowledgeStoreFactory, NewAuditEvent, StorageError,
33    },
34};
35
36use storage_provider::StorageProvider;
37
38use super::{
39    RuntimeConfiguration, RuntimeConfigurationError,
40    knowledge::{
41        index_refresh::{
42            index_refresh_outcome, metadata_for_indexes, recover_index_kinds, refresh_index_kinds,
43        },
44        ingest::mutation_batch_from_request,
45        multimodal::extraction_ingest_request,
46    },
47    runtime::{agent_protocol_status, runtime_status, runtime_status_with_model_profiles},
48    update::{VersionCheckResponse, check_for_updates},
49};
50
51/// Shared application service used by CLI, Web, and future API adapters.
52#[derive(Clone)]
53pub struct RelayKnowledgeService {
54    pub(super) runtime: RuntimeConfiguration,
55    pub(super) storage: StorageProvider,
56    pub(super) health_cache: Arc<tokio::sync::RwLock<Option<HealthResponse>>>,
57    pub(super) watcher: Arc<tokio::sync::RwLock<Option<crate::watcher::WatcherHandle>>>,
58    pub(super) code_retention_cursor: Arc<AtomicUsize>,
59    pub(super) embedding_provider: Option<Arc<dyn EmbeddingProvider>>,
60    pub(super) worker_outbound: Option<Arc<dyn WorkerOutboundPort>>,
61}
62
63impl RelayKnowledgeService {
64    /// Creates a service from validated configuration and injected runtime adapters.
65    pub fn with_runtime_adapters(
66        runtime: RuntimeConfiguration,
67        factory: Arc<dyn KnowledgeStoreFactory>,
68        embedding_provider: Option<Arc<dyn EmbeddingProvider>>,
69        worker_outbound: Option<Arc<dyn WorkerOutboundPort>>,
70    ) -> Self {
71        Self {
72            storage: StorageProvider::configured(factory),
73            runtime,
74            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
75            watcher: Arc::new(tokio::sync::RwLock::new(None)),
76            code_retention_cursor: Arc::new(AtomicUsize::new(0)),
77            embedding_provider,
78            worker_outbound,
79        }
80    }
81
82    /// Creates a service backed by an explicit store and injected runtime adapters.
83    pub fn with_store_and_runtime_adapters(
84        runtime: RuntimeConfiguration,
85        store: Arc<dyn KnowledgeStore>,
86        embedding_provider: Option<Arc<dyn EmbeddingProvider>>,
87        worker_outbound: Option<Arc<dyn WorkerOutboundPort>>,
88    ) -> Self {
89        Self {
90            runtime,
91            storage: StorageProvider::ready(store),
92            health_cache: Arc::new(tokio::sync::RwLock::new(None)),
93            watcher: Arc::new(tokio::sync::RwLock::new(None)),
94            code_retention_cursor: Arc::new(AtomicUsize::new(0)),
95            embedding_provider,
96            worker_outbound,
97        }
98    }
99
100    /// Applies network-related settings from a typed environment snapshot.
101    pub async fn refresh_network_from_environment(
102        &self,
103        environment: &EnvironmentConfig,
104    ) -> Result<(), RuntimeConfigurationError> {
105        self.runtime
106            .network
107            .refresh_from_environment(environment)
108            .map(|_| ())
109            .map_err(RuntimeConfigurationError::Network)
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 provider_name = remote.provider.as_str().to_owned();
400        let provider = self.embedding_provider.as_ref().ok_or_else(|| {
401            ApiError::invalid_argument(
402                "remote embedding provider is configured without an embedding adapter",
403            )
404        })?;
405        let started = Instant::now();
406        let result = provider
407            .embed(EmbeddingRequest {
408                inputs: vec!["relay-knowledge provider probe".to_owned()],
409                model: self.runtime.retrieval.vector_model.name.clone(),
410                dimension: self.runtime.retrieval.vector_model.dimension,
411            })
412            .await;
413
414        match result {
415            Ok(_) => Ok(EmbeddingProviderProbeResponse {
416                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
417                ok: true,
418                provider: Some(provider_name),
419                model: self.runtime.retrieval.vector_model.name.clone(),
420                dimension: self.runtime.retrieval.vector_model.dimension,
421                latency_ms: Some(duration_millis(started.elapsed())),
422                error_code: None,
423                error_message: None,
424                retryable: None,
425            }),
426            Err(error) => Ok(EmbeddingProviderProbeResponse {
427                metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
428                ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
429                provider: Some(provider_name),
430                model: self.runtime.retrieval.vector_model.name.clone(),
431                dimension: self.runtime.retrieval.vector_model.dimension,
432                latency_ms: Some(duration_millis(started.elapsed())),
433                error_code: Some(error.code),
434                error_message: Some(error.message),
435                retryable: Some(error.retry == ProviderRetryClass::Retryable),
436            }),
437        }
438    }
439
440    /// Reconciles derived index cursors before resident service work starts.
441    pub async fn reconcile_startup_indexes(
442        &self,
443        context: RequestContext,
444    ) -> Result<ServiceRecoveryReport, ApiError> {
445        let store = self.storage.get().await.map_err(storage_api_error)?;
446        let graph_version = store
447            .current_graph_version()
448            .await
449            .map_err(storage_api_error)?;
450        let before = store.index_statuses().await.map_err(storage_api_error)?;
451        let active_before = before
452            .iter()
453            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
454            .cloned()
455            .collect::<Vec<_>>();
456        let stale_index_kinds = active_before
457            .iter()
458            .filter(|status| status.is_stale_for(graph_version))
459            .map(|status| status.kind)
460            .collect::<Vec<_>>();
461        let index_lag_max = active_before
462            .iter()
463            .map(|status| {
464                graph_version
465                    .get()
466                    .saturating_sub(status.indexed_graph_version.get())
467            })
468            .max()
469            .unwrap_or(0);
470        let outcome = if stale_index_kinds.is_empty() {
471            index_refresh_outcome(&store).await?
472        } else {
473            recover_index_kinds(
474                &store,
475                stale_index_kinds.clone(),
476                graph_version,
477                &self.runtime.retrieval,
478            )
479            .await?
480        };
481        let refreshed = outcome
482            .indexes
483            .iter()
484            .filter(|status| {
485                stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
486            })
487            .map(|status| status.kind)
488            .collect::<Vec<_>>();
489        let after = outcome.indexes;
490        let active_after = after
491            .iter()
492            .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
493            .cloned()
494            .collect::<Vec<_>>();
495        let metadata = metadata_for_indexes(&context, graph_version, &active_after);
496
497        Ok(ServiceRecoveryReport {
498            metadata,
499            graph_version: graph_version.get(),
500            stale_index_kinds,
501            refreshed_index_kinds: refreshed,
502            index_lag_max,
503            task_queue_depth: outcome.diagnostics.queue_depth,
504            dead_letter_count: outcome.diagnostics.dead_letter_count,
505            heartbeat_state: "ready".to_owned(),
506        })
507    }
508
509    pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
510        self.storage.get().await
511    }
512
513    /// Returns whether graph storage has already been opened by this service.
514    pub fn storage_is_ready(&self) -> bool {
515        self.storage.ready_store().is_some()
516    }
517
518    /// Runs one split-worker preview code-index attempt through durable task leases.
519    pub async fn run_code_index_worker_preview(
520        &self,
521        request: CodeIndexWorkerRunRequest,
522        context: RequestContext,
523    ) -> Result<CodeIndexWorkerRunResponse, ApiError> {
524        let store = self.storage.get().await.map_err(storage_api_error)?;
525        let task = self
526            .run_code_index_task_once(request.task_id, context.clone())
527            .await?;
528        let graph_version = store
529            .current_graph_version()
530            .await
531            .map_err(storage_api_error)?;
532
533        Ok(CodeIndexWorkerRunResponse {
534            metadata: ApiMetadata::graph_only(&context, graph_version),
535            worker_kind: "code_index".to_owned(),
536            claimed: task.is_some(),
537            task,
538        })
539    }
540
541    /// Returns the persistent agent audit log path resolved by the path boundary.
542    pub fn agent_audit_log_path(&self) -> PathBuf {
543        self.runtime.paths.agent_audit_log_file()
544    }
545}
546
547/// Durable audit event input accepted from resident agent adapters.
548#[derive(Debug, Clone)]
549pub struct AgentDurableAuditInput {
550    pub operation: String,
551    pub interface: String,
552    pub request_id: String,
553    pub trace_id: String,
554    pub status: AuditStatus,
555    pub actor: Option<String>,
556    pub source_scope: Option<String>,
557    pub graph_version: u64,
558    pub detail_json: String,
559    pub message: Option<String>,
560}
561
562pub(super) fn storage_api_error(error: StorageError) -> ApiError {
563    match error {
564        StorageError::CapacityExceeded(message) => ApiError::qos_rejected(message),
565        other => ApiError::storage_unavailable(other.to_string()),
566    }
567}
568
569pub(super) async fn file_index_diagnostics_or_default(
570    store: &Arc<dyn KnowledgeStore>,
571) -> Result<FileIndexDiagnostics, ApiError> {
572    match store.file_index_diagnostics().await {
573        Ok(diagnostics) => Ok(diagnostics),
574        Err(StorageError::InvalidInput(message))
575            if message == "file index storage is unavailable" =>
576        {
577            Ok(FileIndexDiagnostics::default())
578        }
579        Err(error) => Err(storage_api_error(error)),
580    }
581}
582
583pub(super) fn graph_with_repository_code_totals(
584    mut graph: GraphInspection,
585    repository_totals: &CodeRepositoryTotals,
586) -> GraphInspection {
587    graph.code_file_count = graph
588        .code_file_count
589        .saturating_add(repository_totals.indexed_file_count);
590    graph.code_symbol_count = graph
591        .code_symbol_count
592        .saturating_add(repository_totals.symbol_count);
593    graph.code_reference_count = graph
594        .code_reference_count
595        .saturating_add(repository_totals.reference_count);
596    graph.code_chunk_count = graph
597        .code_chunk_count
598        .saturating_add(repository_totals.chunk_count);
599    graph.code_parse_status_counts = add_parse_status_counts(
600        graph.code_parse_status_counts,
601        repository_totals.parse_status_counts,
602    );
603
604    graph
605}
606
607fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
608    match kind {
609        GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
610        GraphCanvasKind::Code => GraphCanvasSelection::Code,
611        GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
612    }
613}
614
615fn add_parse_status_counts(
616    left: CodeParseStatusCounts,
617    right: CodeParseStatusCounts,
618) -> CodeParseStatusCounts {
619    CodeParseStatusCounts {
620        parsed: left.parsed.saturating_add(right.parsed),
621        partial: left.partial.saturating_add(right.partial),
622        text_only: left.text_only.saturating_add(right.text_only),
623        failed: left.failed.saturating_add(right.failed),
624    }
625}
626
627fn duration_millis(duration: std::time::Duration) -> u64 {
628    u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
629}
630
631fn service_definition_filename() -> &'static str {
632    if cfg!(target_os = "windows") {
633        WINDOWS_SERVICE_DEFINITION_FILE_NAME
634    } else if cfg!(target_os = "macos") {
635        MACOS_SERVICE_DEFINITION_FILE_NAME
636    } else {
637        LINUX_SERVICE_DEFINITION_FILE_NAME
638    }
639}
640
641mod health;
642mod lifecycle_plan;
643mod retrieval;
644mod service_status;
645mod storage_diagnostics;
646mod storage_provider;
647mod watcher;
648
649#[cfg(test)]
650#[path = "graph_only_tests.rs"]
651mod graph_only_tests;
652
653#[cfg(test)]
654#[path = "recovery_tests.rs"]
655mod recovery_tests;
656
657#[cfg(test)]
658#[path = "refresh_tests.rs"]
659mod refresh_tests;
660
661#[cfg(test)]
662#[path = "storage_tests.rs"]
663mod storage_tests;
664
665#[cfg(test)]
666#[path = "operations_tests.rs"]
667mod operations_tests;
668
669#[cfg(test)]
670#[path = "mod_tests.rs"]
671mod mod_tests;