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 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 results = 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.runtime.retrieval.rerank.rerank(&plan.query, results);
380 let truncated = results.len() > plan.limit;
381 results.truncate(plan.limit);
382 rerank.returned_count = results.len();
383 if rerank.degraded {
384 if let Some(reason) = &rerank.reason {
385 degraded_reasons.push(reason.clone());
386 }
387 }
388 let context_pack = RetrievedContextPack {
389 graph_version,
390 source_scope: plan.source_scope.clone(),
391 freshness: plan.freshness,
392 truncated,
393 backend_statuses: backend_statuses.clone(),
394 items: results
395 .iter()
396 .map(|hit| ContextPackItem {
397 result_id: hit.evidence_id.clone(),
398 source_scope: hit.source_scope.clone(),
399 source_path: hit.source_path.clone(),
400 source_span: hit.source_span,
401 entities: hit.entities.clone(),
402 graph_facts: hit.graph_facts.clone(),
403 graph_paths: hit
404 .graph_facts
405 .iter()
406 .map(ContextGraphPath::from_fact)
407 .collect(),
408 code_artifact: hit.code_artifact.clone(),
409 retriever_sources: hit.retriever_sources.clone(),
410 ranking: hit.ranking.clone(),
411 rerank: hit.rerank.clone(),
412 })
413 .collect(),
414 };
415 let budget_used = RetrievalBudgetUsed {
416 limit: plan.limit,
417 candidate_count: rerank.candidate_count,
418 returned_count: results.len(),
419 context_bytes: retrieval_context_bytes(&results, &context_pack, &backend_statuses),
420 };
421 let fusion = FusionDiagnostics {
422 algorithm: "reciprocal_rank_fusion".to_owned(),
423 k: RECIPROCAL_RANK_FUSION_K,
424 candidate_count: budget_used.candidate_count,
425 };
426 let degraded_reason = (!degraded_reasons.is_empty()).then(|| degraded_reasons.join("; "));
427
428 Ok(HybridRetrievalResponse {
429 metadata,
430 context_pack,
431 retrieval_mode,
432 source_scope: plan.source_scope,
433 freshness: plan.freshness,
434 results,
435 fusion,
436 rerank,
437 backend_statuses,
438 truncated,
439 budget_used,
440 degraded_reason,
441 indexes,
442 index_cursors,
443 index_refresh,
444 })
445 }
446
447 pub async fn inspect_graph(
449 &self,
450 _request: GraphInspectionRequest,
451 context: RequestContext,
452 ) -> Result<GraphInspectionResponse, ApiError> {
453 let store = self.storage.get().await.map_err(storage_api_error)?;
454 let repository_code_totals = store
455 .code_repository_totals()
456 .await
457 .map_err(storage_api_error)?;
458 let graph = graph_with_repository_code_totals(
459 store.inspect_graph().await.map_err(storage_api_error)?,
460 &repository_code_totals,
461 );
462
463 Ok(GraphInspectionResponse {
464 metadata: ApiMetadata::graph_only(&context, graph.graph_version),
465 graph,
466 repository_code_totals,
467 })
468 }
469
470 pub async fn graph_canvas(
472 &self,
473 request: GraphCanvasRequest,
474 context: RequestContext,
475 ) -> Result<GraphCanvasResponse, ApiError> {
476 if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
477 return Err(ApiError::invalid_argument(format!(
478 "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
479 )));
480 }
481 let store = self.storage.get().await.map_err(storage_api_error)?;
482 let graph_version = store
483 .current_graph_version()
484 .await
485 .map_err(storage_api_error)?;
486 let snapshot = store
487 .graph_canvas(GraphCanvasStorageRequest {
488 selection: canvas_selection(request.kind),
489 source_scope: request.source_scope,
490 query: request.query,
491 graph_version,
492 limit: request.limit,
493 })
494 .await
495 .map_err(storage_api_error)?;
496 let node_count = snapshot.nodes.len();
497 let edge_count = snapshot.edges.len();
498
499 Ok(GraphCanvasResponse {
500 metadata: ApiMetadata::graph_only(&context, graph_version),
501 nodes: snapshot
502 .nodes
503 .into_iter()
504 .map(|node| GraphCanvasNode {
505 id: node.id,
506 kind: node.kind,
507 label: node.label,
508 subtitle: node.subtitle,
509 source_scope: node.source_scope,
510 graph_version: node.graph_version.get(),
511 weight: node.weight,
512 status: node.status,
513 details: node.details,
514 })
515 .collect(),
516 edges: snapshot
517 .edges
518 .into_iter()
519 .map(|edge| GraphCanvasEdge {
520 id: edge.id,
521 kind: edge.kind,
522 source: edge.source,
523 target: edge.target,
524 label: edge.label,
525 graph_version: edge.graph_version.get(),
526 confidence_basis_points: edge.confidence_basis_points,
527 evidence_count: edge.evidence_count,
528 details: edge.details,
529 })
530 .collect(),
531 summary: GraphCanvasSummary {
532 kind: request.kind,
533 node_count,
534 edge_count,
535 truncated: snapshot.truncated,
536 available_kinds: snapshot.available_kinds,
537 },
538 })
539 }
540
541 pub async fn refresh_indexes(
543 &self,
544 request: IndexRefreshRequest,
545 context: RequestContext,
546 ) -> Result<IndexRefreshResponse, ApiError> {
547 let store = self.storage.get().await.map_err(storage_api_error)?;
548 let graph_version = store
549 .current_graph_version()
550 .await
551 .map_err(storage_api_error)?;
552 let outcome = refresh_index_kinds(
553 &store,
554 request.kinds,
555 graph_version,
556 &self.runtime.retrieval,
557 )
558 .await?;
559 let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
560
561 Ok(IndexRefreshResponse {
562 metadata,
563 indexes: outcome.indexes,
564 index_cursors: outcome.cursors,
565 diagnostics: outcome.diagnostics,
566 })
567 }
568
569 pub async fn probe_embedding_provider(
571 &self,
572 context: RequestContext,
573 ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
574 let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
575 return Ok(EmbeddingProviderProbeResponse {
576 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
577 ok: false,
578 provider: None,
579 model: self.runtime.retrieval.vector_model.name.clone(),
580 dimension: self.runtime.retrieval.vector_model.dimension,
581 latency_ms: None,
582 error_code: Some("remote_embedding_not_configured".to_owned()),
583 error_message: Some("remote embedding provider is not configured".to_owned()),
584 retryable: Some(false),
585 });
586 };
587 let network = self.runtime.network.current();
588 let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
589 ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
590 })?;
591 let provider_name = remote.provider.as_str().to_owned();
592 let provider = embedding_provider(remote, client);
593 let started = Instant::now();
594 let result = provider
595 .embed(EmbeddingRequest {
596 inputs: vec!["relay-knowledge provider probe".to_owned()],
597 model: self.runtime.retrieval.vector_model.name.clone(),
598 dimension: self.runtime.retrieval.vector_model.dimension,
599 })
600 .await;
601
602 match result {
603 Ok(_) => Ok(EmbeddingProviderProbeResponse {
604 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
605 ok: true,
606 provider: Some(provider_name),
607 model: self.runtime.retrieval.vector_model.name.clone(),
608 dimension: self.runtime.retrieval.vector_model.dimension,
609 latency_ms: Some(duration_millis(started.elapsed())),
610 error_code: None,
611 error_message: None,
612 retryable: None,
613 }),
614 Err(error) => Ok(EmbeddingProviderProbeResponse {
615 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
616 ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
617 provider: Some(provider_name),
618 model: self.runtime.retrieval.vector_model.name.clone(),
619 dimension: self.runtime.retrieval.vector_model.dimension,
620 latency_ms: Some(duration_millis(started.elapsed())),
621 error_code: Some(error.code),
622 error_message: Some(error.message),
623 retryable: Some(error.retry == ProviderRetryClass::Retryable),
624 }),
625 }
626 }
627
628 pub async fn reconcile_startup_indexes(
630 &self,
631 context: RequestContext,
632 ) -> Result<ServiceRecoveryReport, ApiError> {
633 let store = self.storage.get().await.map_err(storage_api_error)?;
634 let graph_version = store
635 .current_graph_version()
636 .await
637 .map_err(storage_api_error)?;
638 let before = store.index_statuses().await.map_err(storage_api_error)?;
639 let active_before = before
640 .iter()
641 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
642 .cloned()
643 .collect::<Vec<_>>();
644 let stale_index_kinds = active_before
645 .iter()
646 .filter(|status| status.is_stale_for(graph_version))
647 .map(|status| status.kind)
648 .collect::<Vec<_>>();
649 let index_lag_max = active_before
650 .iter()
651 .map(|status| {
652 graph_version
653 .get()
654 .saturating_sub(status.indexed_graph_version.get())
655 })
656 .max()
657 .unwrap_or(0);
658 let outcome = if stale_index_kinds.is_empty() {
659 index_refresh_outcome(&store).await?
660 } else {
661 recover_index_kinds(
662 &store,
663 stale_index_kinds.clone(),
664 graph_version,
665 &self.runtime.retrieval,
666 )
667 .await?
668 };
669 let refreshed = outcome
670 .indexes
671 .iter()
672 .filter(|status| {
673 stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
674 })
675 .map(|status| status.kind)
676 .collect::<Vec<_>>();
677 let after = outcome.indexes;
678 let active_after = after
679 .iter()
680 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
681 .cloned()
682 .collect::<Vec<_>>();
683 let metadata = metadata_for_indexes(&context, graph_version, &active_after);
684
685 Ok(ServiceRecoveryReport {
686 metadata,
687 graph_version: graph_version.get(),
688 stale_index_kinds,
689 refreshed_index_kinds: refreshed,
690 index_lag_max,
691 task_queue_depth: outcome.diagnostics.queue_depth,
692 dead_letter_count: outcome.diagnostics.dead_letter_count,
693 heartbeat_state: "ready".to_owned(),
694 })
695 }
696
697 pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
698 self.storage.get().await
699 }
700
701 pub async fn run_code_index_worker_preview(
703 &self,
704 request: CodeIndexWorkerRunRequest,
705 context: RequestContext,
706 ) -> Result<CodeIndexWorkerRunResponse, ApiError> {
707 let store = self.storage.get().await.map_err(storage_api_error)?;
708 let task = self
709 .run_code_index_task_once(request.task_id, context.clone())
710 .await?;
711 let graph_version = store
712 .current_graph_version()
713 .await
714 .map_err(storage_api_error)?;
715
716 Ok(CodeIndexWorkerRunResponse {
717 metadata: ApiMetadata::graph_only(&context, graph_version),
718 worker_kind: "code_index".to_owned(),
719 claimed: task.is_some(),
720 task,
721 })
722 }
723
724 pub fn agent_audit_log_path(&self) -> PathBuf {
726 self.runtime.paths.agent_audit_log_file()
727 }
728}
729
730#[derive(Debug, Clone)]
732pub struct AgentDurableAuditInput {
733 pub operation: String,
734 pub interface: String,
735 pub request_id: String,
736 pub trace_id: String,
737 pub status: AuditStatus,
738 pub actor: Option<String>,
739 pub source_scope: Option<String>,
740 pub graph_version: u64,
741 pub detail_json: String,
742 pub message: Option<String>,
743}
744
745pub(super) fn current_time_millis() -> u64 {
746 std::time::SystemTime::now()
747 .duration_since(std::time::UNIX_EPOCH)
748 .map_or(0, |duration| {
749 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
750 })
751}
752
753pub(super) fn storage_api_error(error: StorageError) -> ApiError {
754 ApiError::storage_unavailable(error.to_string())
755}
756
757pub(super) async fn file_index_diagnostics_or_default(
758 store: &Arc<dyn KnowledgeStore>,
759) -> Result<FileIndexDiagnostics, ApiError> {
760 match store.file_index_diagnostics().await {
761 Ok(diagnostics) => Ok(diagnostics),
762 Err(StorageError::InvalidInput(message))
763 if message == "file index storage is unavailable" =>
764 {
765 Ok(FileIndexDiagnostics::default())
766 }
767 Err(error) => Err(storage_api_error(error)),
768 }
769}
770
771async fn retrieval_index_freshness_snapshot(
772 store: &Arc<dyn KnowledgeStore>,
773) -> Result<IndexRefreshOutcome, ApiError> {
774 let indexes = store.index_statuses().await.map_err(storage_api_error)?;
775 let cursors = match store.index_cursors().await {
776 Ok(cursors) => cursors,
777 Err(StorageError::InvalidInput(message))
778 if message == "index cursor storage is unavailable" =>
779 {
780 Vec::new()
781 }
782 Err(error) => return Err(storage_api_error(error)),
783 };
784 let diagnostics = match store.index_refresh_diagnostics(current_time_millis()).await {
785 Ok(diagnostics) => diagnostics,
786 Err(StorageError::InvalidInput(message))
787 if message == "index refresh diagnostics are unavailable" =>
788 {
789 IndexRefreshDiagnostics::default()
790 }
791 Err(error) => return Err(storage_api_error(error)),
792 };
793
794 Ok(IndexRefreshOutcome {
795 indexes,
796 cursors,
797 diagnostics,
798 })
799}
800
801fn normalize_optional_source_scope(value: Option<String>) -> Result<Option<String>, String> {
802 value
803 .map(|scope| {
804 SourceScope::parse(scope)
805 .map(String::from)
806 .map_err(|error| error.to_string())
807 })
808 .transpose()
809}
810
811fn retrieval_context_bytes(
812 results: &[RetrievalHit],
813 context_pack: &RetrievedContextPack,
814 backend_statuses: &[RetrievalBackendStatus],
815) -> usize {
816 serialized_context_bytes(&context_pack.backend_statuses)
817 .saturating_add(serialized_context_bytes(backend_statuses))
818 .saturating_add(results.iter().map(serialized_context_bytes).sum::<usize>())
819 .saturating_add(
820 context_pack
821 .items
822 .iter()
823 .map(serialized_context_bytes)
824 .sum::<usize>(),
825 )
826}
827
828pub(super) fn graph_with_repository_code_totals(
829 mut graph: GraphInspection,
830 repository_totals: &CodeRepositoryTotals,
831) -> GraphInspection {
832 graph.code_file_count = graph
833 .code_file_count
834 .saturating_add(repository_totals.indexed_file_count);
835 graph.code_symbol_count = graph
836 .code_symbol_count
837 .saturating_add(repository_totals.symbol_count);
838 graph.code_reference_count = graph
839 .code_reference_count
840 .saturating_add(repository_totals.reference_count);
841 graph.code_chunk_count = graph
842 .code_chunk_count
843 .saturating_add(repository_totals.chunk_count);
844 graph.code_parse_status_counts = add_parse_status_counts(
845 graph.code_parse_status_counts,
846 repository_totals.parse_status_counts,
847 );
848
849 graph
850}
851
852fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
853 match kind {
854 GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
855 GraphCanvasKind::Code => GraphCanvasSelection::Code,
856 GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
857 }
858}
859
860fn add_parse_status_counts(
861 left: CodeParseStatusCounts,
862 right: CodeParseStatusCounts,
863) -> CodeParseStatusCounts {
864 CodeParseStatusCounts {
865 parsed: left.parsed.saturating_add(right.parsed),
866 partial: left.partial.saturating_add(right.partial),
867 text_only: left.text_only.saturating_add(right.text_only),
868 failed: left.failed.saturating_add(right.failed),
869 }
870}
871
872fn serialized_context_bytes<T: Serialize + ?Sized>(value: &T) -> usize {
873 serde_json::to_vec(value)
874 .map(|bytes| bytes.len())
875 .unwrap_or(usize::MAX / 4)
876}
877
878fn duration_millis(duration: std::time::Duration) -> u64 {
879 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
880}
881
882fn service_definition_filename() -> &'static str {
883 if cfg!(target_os = "windows") {
884 WINDOWS_SERVICE_DEFINITION_FILE_NAME
885 } else if cfg!(target_os = "macos") {
886 MACOS_SERVICE_DEFINITION_FILE_NAME
887 } else {
888 LINUX_SERVICE_DEFINITION_FILE_NAME
889 }
890}
891
892mod health;
893pub(crate) mod knowledge_map;
894mod lifecycle_plan;
895mod service_status;
896mod storage_diagnostics;
897mod storage_provider;
898mod watcher;
899
900#[cfg(test)]
901mod id_tests;
902
903#[cfg(test)]
904mod graph_only_tests;
905
906#[cfg(test)]
907mod recovery_tests;
908
909#[cfg(test)]
910mod refresh_tests;
911
912#[cfg(test)]
913mod storage_tests;
914
915#[cfg(test)]
916mod operations_tests;
917
918#[cfg(test)]
919mod tests;