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