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#[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 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 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 pub async fn from_process_environment() -> Result<Self, RuntimeConfigurationError> {
75 RuntimeConfiguration::from_process_environment()
76 .await
77 .map(Self::new)
78 }
79
80 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 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 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 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 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 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 pub fn storage_is_ready(&self) -> bool {
520 self.storage.ready_store().is_some()
521 }
522
523 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 pub fn agent_audit_log_path(&self) -> PathBuf {
548 self.runtime.paths.agent_audit_log_file()
549 }
550}
551
552#[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;