1use std::{path::PathBuf, sync::Arc, time::Instant};
2
3use serde::Serialize;
4
5use crate::{
6 api::{
7 AgentProtocolStatus, ApiError, ApiMetadata, CodeIndexWorkerRunRequest,
8 CodeIndexWorkerRunResponse, EmbeddingProviderProbeResponse, GRAPH_CANVAS_MAX_LIMIT,
9 GraphCanvasEdge, GraphCanvasKind, GraphCanvasNode, GraphCanvasRequest, GraphCanvasResponse,
10 GraphCanvasSummary, GraphInspectionRequest, GraphInspectionResponse, HealthResponse,
11 HybridRetrievalRequest, HybridRetrievalResponse, IndexRefreshRequest, IndexRefreshResponse,
12 IngestRequest, IngestResponse, MultimodalExtractionRequest, MultimodalExtractionResponse,
13 ProjectStatusResponse, RequestContext, ServiceRecoveryReport,
14 },
15 domain::{
16 AuditStatus, CodeParseStatusCounts, CodeRepositoryTotals, ContextGraphPath,
17 ContextPackItem, FreshnessPolicy, FusionDiagnostics, IndexKind, RECIPROCAL_RANK_FUSION_K,
18 RetrievalBackendStatus, RetrievalBudgetUsed, RetrievalHit, RetrievalMode,
19 RetrievedContextPack, RetrieverSource, SourceScope,
20 },
21 env::EnvironmentConfig,
22 model_provider::ModelProviderConfigService,
23 observability::ObservabilityRuntime,
24 project::{
25 LINUX_SERVICE_DEFINITION_FILE_NAME, MACOS_SERVICE_DEFINITION_FILE_NAME, PROJECT_NAME,
26 WINDOWS_SERVICE_DEFINITION_FILE_NAME,
27 },
28 retrieval::{
29 RetrievalPlan,
30 provider::{EmbeddingRequest, ProviderRetryClass, embedding_provider_with_qos},
31 read_model_backend_statuses,
32 },
33 storage::{
34 FileIndexDiagnostics, GraphCanvasSelection, GraphCanvasStorageRequest, GraphInspection,
35 GraphSearchRequest, IndexRefreshDiagnostics, KnowledgeStore, NewAuditEvent, StorageError,
36 },
37};
38
39use storage_provider::StorageProvider;
40
41use super::{
42 RuntimeConfiguration, RuntimeConfigurationError,
43 knowledge::{
44 index_refresh::{
45 IndexRefreshOutcome, index_refresh_outcome, metadata_for_indexes, recover_index_kinds,
46 refresh_index_kinds,
47 },
48 ingest::mutation_batch_from_request,
49 multimodal::extraction_ingest_request,
50 },
51 status::{agent_protocol_status, runtime_status, runtime_status_with_model_profiles},
52 update::{VersionCheckResponse, check_for_updates},
53};
54
55#[cfg(test)]
56use super::knowledge::ingest::generated_evidence_id;
57
58#[derive(Clone)]
60pub struct RelayKnowledgeService {
61 pub(super) runtime: RuntimeConfiguration,
62 pub(super) storage: StorageProvider,
63 pub(super) health_cache: Arc<tokio::sync::RwLock<Option<HealthResponse>>>,
64 pub(super) watcher: Arc<tokio::sync::RwLock<Option<crate::watcher::WatcherHandle>>>,
65}
66
67impl RelayKnowledgeService {
68 pub fn new(runtime: RuntimeConfiguration) -> Self {
70 Self {
71 storage: StorageProvider::configured(&runtime),
72 runtime,
73 health_cache: Arc::new(tokio::sync::RwLock::new(None)),
74 watcher: Arc::new(tokio::sync::RwLock::new(None)),
75 }
76 }
77
78 pub fn with_store(runtime: RuntimeConfiguration, store: Arc<dyn KnowledgeStore>) -> Self {
80 Self {
81 runtime,
82 storage: StorageProvider::ready(store),
83 health_cache: Arc::new(tokio::sync::RwLock::new(None)),
84 watcher: Arc::new(tokio::sync::RwLock::new(None)),
85 }
86 }
87
88 pub async fn from_process_environment() -> Result<Self, RuntimeConfigurationError> {
90 RuntimeConfiguration::from_process_environment()
91 .await
92 .map(Self::new)
93 }
94
95 pub async fn from_environment(
97 environment: &EnvironmentConfig,
98 ) -> Result<Self, RuntimeConfigurationError> {
99 RuntimeConfiguration::from_environment(environment)
100 .await
101 .map(Self::new)
102 }
103
104 pub async fn refresh_network_from_environment(
106 &self,
107 environment: &EnvironmentConfig,
108 ) -> Result<(), RuntimeConfigurationError> {
109 self.runtime
110 .network
111 .refresh_from_environment(environment)
112 .map(|_| ())
113 .map_err(RuntimeConfigurationError::Network)
114 }
115
116 pub async fn refresh_network_from_process_environment(
118 &self,
119 ) -> Result<(), RuntimeConfigurationError> {
120 self.runtime
121 .network
122 .refresh_from_process_environment()
123 .map(|_| ())
124 .map_err(RuntimeConfigurationError::NetworkRuntime)
125 }
126
127 pub fn observability(&self) -> ObservabilityRuntime {
129 self.runtime.observability.clone()
130 }
131
132 pub fn model_provider_config(&self) -> ModelProviderConfigService {
134 ModelProviderConfigService::new(self.runtime.paths.clone())
135 }
136
137 pub async fn check_for_updates(&self, force_refresh: bool) -> VersionCheckResponse {
139 check_for_updates(
140 &self.runtime.paths,
141 &self.runtime.network,
142 &self.runtime.updates,
143 force_refresh,
144 )
145 .await
146 }
147
148 pub async fn record_agent_audit(&self, event: AgentDurableAuditInput) -> Result<(), ApiError> {
150 let store = self.storage.get().await.map_err(storage_api_error)?;
151 store
152 .insert_audit_event(NewAuditEvent {
153 operation: event.operation,
154 interface: event.interface,
155 request_id: event.request_id,
156 trace_id: event.trace_id,
157 status: event.status,
158 actor: event.actor,
159 source_scope: event.source_scope,
160 graph_version: event.graph_version,
161 detail_json: event.detail_json,
162 message: event.message,
163 now_ms: current_time_millis(),
164 })
165 .await
166 .map(|_| ())
167 .map_err(storage_api_error)
168 }
169
170 pub async fn project_status(
172 &self,
173 context: RequestContext,
174 ) -> Result<ProjectStatusResponse, ApiError> {
175 let store = self.storage.get().await.map_err(storage_api_error)?;
176 let graph_version = store
177 .current_graph_version()
178 .await
179 .map_err(storage_api_error)?;
180
181 let model_profiles = self
182 .model_provider_config()
183 .profile_summary(&self.runtime.retrieval)
184 .await;
185
186 Ok(ProjectStatusResponse {
187 project_name: PROJECT_NAME.to_owned(),
188 metadata: ApiMetadata::graph_only(&context, graph_version),
189 runtime: runtime_status_with_model_profiles(&self.runtime, model_profiles),
190 })
191 }
192
193 pub fn runtime_diagnostics(
195 &self,
196 context: RequestContext,
197 ) -> (ProjectStatusResponse, AgentProtocolStatus) {
198 (
199 ProjectStatusResponse {
200 project_name: PROJECT_NAME.to_owned(),
201 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
202 runtime: runtime_status(&self.runtime),
203 },
204 agent_protocol_status(&self.runtime),
205 )
206 }
207
208 pub async fn ingest(
210 &self,
211 request: IngestRequest,
212 context: RequestContext,
213 ) -> Result<IngestResponse, ApiError> {
214 let batch = mutation_batch_from_request(request)
215 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
216 let worker_evidence = batch.evidence.clone();
217 let store = self.storage.get().await.map_err(storage_api_error)?;
218 let receipt = store
219 .commit_mutation_batch(batch)
220 .await
221 .map_err(storage_api_error)?;
222 self.queue_worker_tasks_for_evidence(&store, &worker_evidence, receipt.graph_version)
223 .await?;
224 let (indexes, metadata, index_refresh_error) = match refresh_index_kinds(
225 &store,
226 IndexKind::ALL,
227 receipt.graph_version,
228 &self.runtime.retrieval,
229 )
230 .await
231 {
232 Ok(outcome) => {
233 let metadata =
234 metadata_for_indexes(&context, receipt.graph_version, &outcome.indexes);
235
236 (outcome.indexes, metadata, None)
237 }
238 Err(error) => (
239 Vec::new(),
240 ApiMetadata::indexed(&context, receipt.graph_version, None, None, true),
241 Some(error.message),
242 ),
243 };
244
245 Ok(IngestResponse {
246 metadata,
247 receipt,
248 indexes,
249 index_refresh_error,
250 })
251 }
252
253 pub async fn commit_multimodal_extraction(
255 &self,
256 request: MultimodalExtractionRequest,
257 context: RequestContext,
258 ) -> Result<MultimodalExtractionResponse, ApiError> {
259 let converted = extraction_ingest_request(request).map_err(ApiError::invalid_argument)?;
260 let parent_evidence_id = converted.parent_evidence_id;
261 let derived_evidence_count = converted.derived_evidence_count;
262 let response = self.ingest(converted.ingest, context).await?;
263
264 Ok(MultimodalExtractionResponse {
265 metadata: response.metadata,
266 parent_evidence_id,
267 derived_evidence_count,
268 receipt: response.receipt,
269 indexes: response.indexes,
270 index_refresh_error: response.index_refresh_error,
271 })
272 }
273
274 pub async fn retrieve_context(
276 &self,
277 request: HybridRetrievalRequest,
278 context: RequestContext,
279 ) -> Result<HybridRetrievalResponse, ApiError> {
280 let source_scope = normalize_optional_source_scope(request.source_scope)
281 .map_err(ApiError::invalid_argument)?;
282 let plan = RetrievalPlan::new(
283 request.query,
284 source_scope,
285 request.limit,
286 request.freshness,
287 )
288 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
289 let store = self.storage.get().await.map_err(storage_api_error)?;
290 let graph_version = store
291 .current_graph_version()
292 .await
293 .map_err(storage_api_error)?;
294
295 let mut retrieval_mode = RetrievalMode::Hybrid;
296 let mut indexes = Vec::new();
297 let mut index_cursors = Vec::new();
298 let mut index_refresh = IndexRefreshDiagnostics::default();
299 let mut metadata = ApiMetadata::graph_only(&context, graph_version);
300 let mut degraded_reasons = Vec::new();
301 let backend_statuses = if plan.freshness == FreshnessPolicy::GraphOnly {
302 retrieval_mode = RetrievalMode::GraphOnly;
303 degraded_reasons.push("graph_only freshness policy selected".to_owned());
304 Vec::new()
305 } else {
306 let mut index_outcome = retrieval_index_freshness_snapshot(&store).await?;
307 indexes = index_outcome.indexes;
308 index_cursors = index_outcome.cursors;
309 index_refresh = index_outcome.diagnostics;
310 let mut active_indexes = indexes
311 .iter()
312 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
313 .cloned()
314 .collect::<Vec<_>>();
315 if plan.freshness == FreshnessPolicy::WaitUntilFresh {
316 let stale_kinds = active_indexes
317 .iter()
318 .filter(|status| status.is_stale_for(graph_version))
319 .map(|status| status.kind)
320 .collect::<Vec<_>>();
321 if !stale_kinds.is_empty() {
322 refresh_index_kinds(
323 &store,
324 stale_kinds,
325 graph_version,
326 &self.runtime.retrieval,
327 )
328 .await?;
329 index_outcome = retrieval_index_freshness_snapshot(&store).await?;
330 indexes = index_outcome.indexes;
331 index_cursors = index_outcome.cursors;
332 index_refresh = index_outcome.diagnostics;
333 active_indexes = indexes
334 .iter()
335 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
336 .cloned()
337 .collect();
338 }
339 }
340
341 let stale = active_indexes
342 .iter()
343 .any(|status| status.is_stale_for(graph_version));
344 metadata = metadata_for_indexes(&context, graph_version, &active_indexes);
345 if plan.freshness == FreshnessPolicy::AllowStale && stale {
346 degraded_reasons
347 .push("one or more indexes are behind the graph version".to_owned());
348 }
349 read_model_backend_statuses(&plan, graph_version, &indexes, &self.runtime.retrieval)
350 };
351 if backend_statuses
352 .iter()
353 .any(|status| status.state == crate::domain::RetrievalBackendState::Unavailable)
354 {
355 degraded_reasons.push(
356 "semantic/vector retrieval backends unavailable; using bm25, graph evidence, and code graph fallback"
357 .to_owned(),
358 );
359 }
360 let mut disabled_retriever_sources = self.runtime.retrieval.disabled_retriever_sources();
361 if plan.freshness == FreshnessPolicy::GraphOnly {
362 for source in [RetrieverSource::Semantic, RetrieverSource::Vector] {
363 if !disabled_retriever_sources.contains(&source) {
364 disabled_retriever_sources.push(source);
365 }
366 }
367 }
368 let candidate_limit = self.runtime.retrieval.rerank.candidate_limit(plan.limit);
369 let search_outcome = store
370 .search(GraphSearchRequest {
371 query: plan.query.clone(),
372 source_scope: plan.source_scope.clone(),
373 graph_version,
374 limit: candidate_limit,
375 disabled_retriever_sources,
376 })
377 .await
378 .map_err(storage_api_error)?;
379 let (mut results, mut rerank) = self
380 .runtime
381 .retrieval
382 .rerank
383 .rerank(&plan.query, search_outcome.hits);
384 let result_truncated = results.len() > plan.limit;
385 results.truncate(plan.limit);
386 rerank.returned_count = results.len();
387 if rerank.degraded {
388 if let Some(reason) = &rerank.reason {
389 degraded_reasons.push(reason.clone());
390 }
391 }
392 let degraded_reason = (!degraded_reasons.is_empty()).then(|| degraded_reasons.join("; "));
393 let mut provenance_trace = search_outcome.trace;
394 provenance_trace.mark_citations_for_hits(results.iter());
395 provenance_trace.stale = degraded_reasons
396 .iter()
397 .any(|reason| reason.contains("behind the graph version"));
398 provenance_trace.degraded_reason = degraded_reason.clone();
399 provenance_trace.truncated |= result_truncated;
400 provenance_trace.apply_budget(plan.limit.saturating_mul(4).max(plan.limit + 8).min(64));
401 let truncated = result_truncated || provenance_trace.truncated;
402
403 let context_pack = RetrievedContextPack {
404 graph_version,
405 source_scope: plan.source_scope.clone(),
406 freshness: plan.freshness,
407 truncated,
408 backend_statuses: backend_statuses.clone(),
409 provenance_trace: Some(provenance_trace),
410 items: results
411 .iter()
412 .map(|hit| ContextPackItem {
413 result_id: hit.evidence_id.clone(),
414 source_scope: hit.source_scope.clone(),
415 source_path: hit.source_path.clone(),
416 source_span: hit.source_span,
417 entities: hit.entities.clone(),
418 graph_facts: hit.graph_facts.clone(),
419 graph_paths: hit
420 .graph_facts
421 .iter()
422 .map(ContextGraphPath::from_fact)
423 .collect(),
424 code_artifact: hit.code_artifact.clone(),
425 retriever_sources: hit.retriever_sources.clone(),
426 ranking: hit.ranking.clone(),
427 rerank: hit.rerank.clone(),
428 })
429 .collect(),
430 };
431 let budget_used = RetrievalBudgetUsed {
432 limit: plan.limit,
433 candidate_count: rerank.candidate_count,
434 returned_count: results.len(),
435 context_bytes: retrieval_context_bytes(&results, &context_pack, &backend_statuses),
436 };
437 let fusion = FusionDiagnostics {
438 algorithm: "reciprocal_rank_fusion".to_owned(),
439 k: RECIPROCAL_RANK_FUSION_K,
440 candidate_count: budget_used.candidate_count,
441 };
442
443 Ok(HybridRetrievalResponse {
444 metadata,
445 context_pack,
446 retrieval_mode,
447 source_scope: plan.source_scope,
448 freshness: plan.freshness,
449 results,
450 fusion,
451 rerank,
452 backend_statuses,
453 truncated,
454 budget_used,
455 degraded_reason,
456 indexes,
457 index_cursors,
458 index_refresh,
459 })
460 }
461
462 pub async fn inspect_graph(
464 &self,
465 _request: GraphInspectionRequest,
466 context: RequestContext,
467 ) -> Result<GraphInspectionResponse, ApiError> {
468 let store = self.storage.get().await.map_err(storage_api_error)?;
469 let repository_code_totals = store
470 .code_repository_totals()
471 .await
472 .map_err(storage_api_error)?;
473 let graph = graph_with_repository_code_totals(
474 store.inspect_graph().await.map_err(storage_api_error)?,
475 &repository_code_totals,
476 );
477
478 Ok(GraphInspectionResponse {
479 metadata: ApiMetadata::graph_only(&context, graph.graph_version),
480 graph,
481 repository_code_totals,
482 })
483 }
484
485 pub async fn graph_canvas(
487 &self,
488 request: GraphCanvasRequest,
489 context: RequestContext,
490 ) -> Result<GraphCanvasResponse, ApiError> {
491 if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
492 return Err(ApiError::invalid_argument(format!(
493 "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
494 )));
495 }
496 let store = self.storage.get().await.map_err(storage_api_error)?;
497 let graph_version = store
498 .current_graph_version()
499 .await
500 .map_err(storage_api_error)?;
501 let snapshot = store
502 .graph_canvas(GraphCanvasStorageRequest {
503 selection: canvas_selection(request.kind),
504 source_scope: request.source_scope,
505 query: request.query,
506 graph_version,
507 limit: request.limit,
508 })
509 .await
510 .map_err(storage_api_error)?;
511 let node_count = snapshot.nodes.len();
512 let edge_count = snapshot.edges.len();
513
514 Ok(GraphCanvasResponse {
515 metadata: ApiMetadata::graph_only(&context, graph_version),
516 nodes: snapshot
517 .nodes
518 .into_iter()
519 .map(|node| GraphCanvasNode {
520 id: node.id,
521 kind: node.kind,
522 label: node.label,
523 subtitle: node.subtitle,
524 source_scope: node.source_scope,
525 graph_version: node.graph_version.get(),
526 weight: node.weight,
527 status: node.status,
528 details: node.details,
529 })
530 .collect(),
531 edges: snapshot
532 .edges
533 .into_iter()
534 .map(|edge| GraphCanvasEdge {
535 id: edge.id,
536 kind: edge.kind,
537 source: edge.source,
538 target: edge.target,
539 label: edge.label,
540 graph_version: edge.graph_version.get(),
541 confidence_basis_points: edge.confidence_basis_points,
542 evidence_count: edge.evidence_count,
543 details: edge.details,
544 })
545 .collect(),
546 summary: GraphCanvasSummary {
547 kind: request.kind,
548 node_count,
549 edge_count,
550 truncated: snapshot.truncated,
551 available_kinds: snapshot.available_kinds,
552 },
553 })
554 }
555
556 pub async fn refresh_indexes(
558 &self,
559 request: IndexRefreshRequest,
560 context: RequestContext,
561 ) -> Result<IndexRefreshResponse, ApiError> {
562 let store = self.storage.get().await.map_err(storage_api_error)?;
563 let graph_version = store
564 .current_graph_version()
565 .await
566 .map_err(storage_api_error)?;
567 let outcome = refresh_index_kinds(
568 &store,
569 request.kinds,
570 graph_version,
571 &self.runtime.retrieval,
572 )
573 .await?;
574 let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
575
576 Ok(IndexRefreshResponse {
577 metadata,
578 indexes: outcome.indexes,
579 index_cursors: outcome.cursors,
580 diagnostics: outcome.diagnostics,
581 })
582 }
583
584 pub async fn probe_embedding_provider(
586 &self,
587 context: RequestContext,
588 ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
589 let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
590 return Ok(EmbeddingProviderProbeResponse {
591 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
592 ok: false,
593 provider: None,
594 model: self.runtime.retrieval.vector_model.name.clone(),
595 dimension: self.runtime.retrieval.vector_model.dimension,
596 latency_ms: None,
597 error_code: Some("remote_embedding_not_configured".to_owned()),
598 error_message: Some("remote embedding provider is not configured".to_owned()),
599 retryable: Some(false),
600 });
601 };
602 let network = self.runtime.network.current();
603 let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
604 ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
605 })?;
606 let provider_name = remote.provider.as_str().to_owned();
607 let provider = embedding_provider_with_qos(
608 remote,
609 client,
610 self.runtime.network.qos_runtime(),
611 network.qos,
612 );
613 let started = Instant::now();
614 let result = provider
615 .embed(EmbeddingRequest {
616 inputs: vec!["relay-knowledge provider probe".to_owned()],
617 model: self.runtime.retrieval.vector_model.name.clone(),
618 dimension: self.runtime.retrieval.vector_model.dimension,
619 })
620 .await;
621
622 match result {
623 Ok(_) => Ok(EmbeddingProviderProbeResponse {
624 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
625 ok: true,
626 provider: Some(provider_name),
627 model: self.runtime.retrieval.vector_model.name.clone(),
628 dimension: self.runtime.retrieval.vector_model.dimension,
629 latency_ms: Some(duration_millis(started.elapsed())),
630 error_code: None,
631 error_message: None,
632 retryable: None,
633 }),
634 Err(error) => Ok(EmbeddingProviderProbeResponse {
635 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
636 ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
637 provider: Some(provider_name),
638 model: self.runtime.retrieval.vector_model.name.clone(),
639 dimension: self.runtime.retrieval.vector_model.dimension,
640 latency_ms: Some(duration_millis(started.elapsed())),
641 error_code: Some(error.code),
642 error_message: Some(error.message),
643 retryable: Some(error.retry == ProviderRetryClass::Retryable),
644 }),
645 }
646 }
647
648 pub async fn reconcile_startup_indexes(
650 &self,
651 context: RequestContext,
652 ) -> Result<ServiceRecoveryReport, ApiError> {
653 let store = self.storage.get().await.map_err(storage_api_error)?;
654 let graph_version = store
655 .current_graph_version()
656 .await
657 .map_err(storage_api_error)?;
658 let before = store.index_statuses().await.map_err(storage_api_error)?;
659 let active_before = before
660 .iter()
661 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
662 .cloned()
663 .collect::<Vec<_>>();
664 let stale_index_kinds = active_before
665 .iter()
666 .filter(|status| status.is_stale_for(graph_version))
667 .map(|status| status.kind)
668 .collect::<Vec<_>>();
669 let index_lag_max = active_before
670 .iter()
671 .map(|status| {
672 graph_version
673 .get()
674 .saturating_sub(status.indexed_graph_version.get())
675 })
676 .max()
677 .unwrap_or(0);
678 let outcome = if stale_index_kinds.is_empty() {
679 index_refresh_outcome(&store).await?
680 } else {
681 recover_index_kinds(
682 &store,
683 stale_index_kinds.clone(),
684 graph_version,
685 &self.runtime.retrieval,
686 )
687 .await?
688 };
689 let refreshed = outcome
690 .indexes
691 .iter()
692 .filter(|status| {
693 stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
694 })
695 .map(|status| status.kind)
696 .collect::<Vec<_>>();
697 let after = outcome.indexes;
698 let active_after = after
699 .iter()
700 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
701 .cloned()
702 .collect::<Vec<_>>();
703 let metadata = metadata_for_indexes(&context, graph_version, &active_after);
704
705 Ok(ServiceRecoveryReport {
706 metadata,
707 graph_version: graph_version.get(),
708 stale_index_kinds,
709 refreshed_index_kinds: refreshed,
710 index_lag_max,
711 task_queue_depth: outcome.diagnostics.queue_depth,
712 dead_letter_count: outcome.diagnostics.dead_letter_count,
713 heartbeat_state: "ready".to_owned(),
714 })
715 }
716
717 pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
718 self.storage.get().await
719 }
720
721 pub fn storage_is_ready(&self) -> bool {
723 self.storage.ready_store().is_some()
724 }
725
726 pub async fn run_code_index_worker_preview(
728 &self,
729 request: CodeIndexWorkerRunRequest,
730 context: RequestContext,
731 ) -> Result<CodeIndexWorkerRunResponse, ApiError> {
732 let store = self.storage.get().await.map_err(storage_api_error)?;
733 let task = self
734 .run_code_index_task_once(request.task_id, context.clone())
735 .await?;
736 let graph_version = store
737 .current_graph_version()
738 .await
739 .map_err(storage_api_error)?;
740
741 Ok(CodeIndexWorkerRunResponse {
742 metadata: ApiMetadata::graph_only(&context, graph_version),
743 worker_kind: "code_index".to_owned(),
744 claimed: task.is_some(),
745 task,
746 })
747 }
748
749 pub fn agent_audit_log_path(&self) -> PathBuf {
751 self.runtime.paths.agent_audit_log_file()
752 }
753}
754
755#[derive(Debug, Clone)]
757pub struct AgentDurableAuditInput {
758 pub operation: String,
759 pub interface: String,
760 pub request_id: String,
761 pub trace_id: String,
762 pub status: AuditStatus,
763 pub actor: Option<String>,
764 pub source_scope: Option<String>,
765 pub graph_version: u64,
766 pub detail_json: String,
767 pub message: Option<String>,
768}
769
770pub(super) fn current_time_millis() -> u64 {
771 std::time::SystemTime::now()
772 .duration_since(std::time::UNIX_EPOCH)
773 .map_or(0, |duration| {
774 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
775 })
776}
777
778pub(super) fn storage_api_error(error: StorageError) -> ApiError {
779 ApiError::storage_unavailable(error.to_string())
780}
781
782pub(super) async fn file_index_diagnostics_or_default(
783 store: &Arc<dyn KnowledgeStore>,
784) -> Result<FileIndexDiagnostics, ApiError> {
785 match store.file_index_diagnostics().await {
786 Ok(diagnostics) => Ok(diagnostics),
787 Err(StorageError::InvalidInput(message))
788 if message == "file index storage is unavailable" =>
789 {
790 Ok(FileIndexDiagnostics::default())
791 }
792 Err(error) => Err(storage_api_error(error)),
793 }
794}
795
796async fn retrieval_index_freshness_snapshot(
797 store: &Arc<dyn KnowledgeStore>,
798) -> Result<IndexRefreshOutcome, ApiError> {
799 let indexes = store.index_statuses().await.map_err(storage_api_error)?;
800 let cursors = match store.index_cursors().await {
801 Ok(cursors) => cursors,
802 Err(StorageError::InvalidInput(message))
803 if message == "index cursor storage is unavailable" =>
804 {
805 Vec::new()
806 }
807 Err(error) => return Err(storage_api_error(error)),
808 };
809 let diagnostics = match store.index_refresh_diagnostics(current_time_millis()).await {
810 Ok(diagnostics) => diagnostics,
811 Err(StorageError::InvalidInput(message))
812 if message == "index refresh diagnostics are unavailable" =>
813 {
814 IndexRefreshDiagnostics::default()
815 }
816 Err(error) => return Err(storage_api_error(error)),
817 };
818
819 Ok(IndexRefreshOutcome {
820 indexes,
821 cursors,
822 diagnostics,
823 })
824}
825
826fn normalize_optional_source_scope(value: Option<String>) -> Result<Option<String>, String> {
827 value
828 .map(|scope| {
829 SourceScope::parse(scope)
830 .map(String::from)
831 .map_err(|error| error.to_string())
832 })
833 .transpose()
834}
835
836fn retrieval_context_bytes(
837 results: &[RetrievalHit],
838 context_pack: &RetrievedContextPack,
839 backend_statuses: &[RetrievalBackendStatus],
840) -> usize {
841 serialized_context_bytes(&context_pack.backend_statuses)
842 .saturating_add(serialized_context_bytes(backend_statuses))
843 .saturating_add(
844 context_pack
845 .provenance_trace
846 .as_ref()
847 .map(serialized_context_bytes)
848 .unwrap_or_default(),
849 )
850 .saturating_add(results.iter().map(serialized_context_bytes).sum::<usize>())
851 .saturating_add(
852 context_pack
853 .items
854 .iter()
855 .map(serialized_context_bytes)
856 .sum::<usize>(),
857 )
858}
859
860pub(super) fn graph_with_repository_code_totals(
861 mut graph: GraphInspection,
862 repository_totals: &CodeRepositoryTotals,
863) -> GraphInspection {
864 graph.code_file_count = graph
865 .code_file_count
866 .saturating_add(repository_totals.indexed_file_count);
867 graph.code_symbol_count = graph
868 .code_symbol_count
869 .saturating_add(repository_totals.symbol_count);
870 graph.code_reference_count = graph
871 .code_reference_count
872 .saturating_add(repository_totals.reference_count);
873 graph.code_chunk_count = graph
874 .code_chunk_count
875 .saturating_add(repository_totals.chunk_count);
876 graph.code_parse_status_counts = add_parse_status_counts(
877 graph.code_parse_status_counts,
878 repository_totals.parse_status_counts,
879 );
880
881 graph
882}
883
884fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
885 match kind {
886 GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
887 GraphCanvasKind::Code => GraphCanvasSelection::Code,
888 GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
889 }
890}
891
892fn add_parse_status_counts(
893 left: CodeParseStatusCounts,
894 right: CodeParseStatusCounts,
895) -> CodeParseStatusCounts {
896 CodeParseStatusCounts {
897 parsed: left.parsed.saturating_add(right.parsed),
898 partial: left.partial.saturating_add(right.partial),
899 text_only: left.text_only.saturating_add(right.text_only),
900 failed: left.failed.saturating_add(right.failed),
901 }
902}
903
904fn serialized_context_bytes<T: Serialize + ?Sized>(value: &T) -> usize {
905 serde_json::to_vec(value)
906 .map(|bytes| bytes.len())
907 .unwrap_or(usize::MAX / 4)
908}
909
910fn duration_millis(duration: std::time::Duration) -> u64 {
911 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
912}
913
914fn service_definition_filename() -> &'static str {
915 if cfg!(target_os = "windows") {
916 WINDOWS_SERVICE_DEFINITION_FILE_NAME
917 } else if cfg!(target_os = "macos") {
918 MACOS_SERVICE_DEFINITION_FILE_NAME
919 } else {
920 LINUX_SERVICE_DEFINITION_FILE_NAME
921 }
922}
923
924mod health;
925mod lifecycle_plan;
926mod service_status;
927mod storage_diagnostics;
928mod storage_provider;
929mod watcher;
930
931#[cfg(test)]
932mod id_tests;
933
934#[cfg(test)]
935mod graph_only_tests;
936
937#[cfg(test)]
938mod recovery_tests;
939
940#[cfg(test)]
941mod refresh_tests;
942
943#[cfg(test)]
944mod storage_tests;
945
946#[cfg(test)]
947mod operations_tests;
948
949#[cfg(test)]
950mod tests;
951
952#[cfg(test)]
953mod trace_tests;