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#[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 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 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 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 pub fn observability(&self) -> ObservabilityRuntime {
114 self.runtime.observability.clone()
115 }
116
117 pub fn model_provider_config(&self) -> ModelProviderConfigService {
119 ModelProviderConfigService::new(self.runtime.paths.clone())
120 }
121
122 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 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 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 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 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 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 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 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 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 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 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 pub fn storage_is_ready(&self) -> bool {
515 self.storage.ready_store().is_some()
516 }
517
518 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 pub fn agent_audit_log_path(&self) -> PathBuf {
543 self.runtime.paths.agent_audit_log_file()
544 }
545}
546
547#[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;