1use std::{
2 path::PathBuf,
3 sync::{Arc, OnceLock},
4 time::Instant,
5};
6
7use serde::Serialize;
8
9use crate::{
10 api::{
11 AgentProtocolStatus, ApiError, ApiMetadata, EmbeddingProviderProbeResponse,
12 GRAPH_CANVAS_MAX_LIMIT, GraphCanvasEdge, GraphCanvasKind, GraphCanvasNode,
13 GraphCanvasRequest, GraphCanvasResponse, GraphCanvasSummary, GraphInspectionRequest,
14 GraphInspectionResponse, HealthResponse, HybridRetrievalRequest, HybridRetrievalResponse,
15 IndexRefreshRequest, IndexRefreshResponse, IngestRequest, IngestResponse,
16 MultimodalExtractionRequest, MultimodalExtractionResponse, ProjectStatusResponse,
17 RequestContext, ServiceRecoveryReport, ServiceStatusResponse,
18 },
19 domain::{
20 AuditStatus, CodeParseStatusCounts, CodeRepositoryTotals, ContextGraphPath,
21 ContextPackItem, FreshnessPolicy, FusionDiagnostics, IndexKind, ProposalState,
22 RECIPROCAL_RANK_FUSION_K, RetrievalBackendStatus, RetrievalBudgetUsed, RetrievalHit,
23 RetrievalMode, RetrievedContextPack, RetrieverSource, SourceScope,
24 },
25 env::EnvironmentConfig,
26 model_provider::ModelProviderConfigService,
27 observability::ObservabilityRuntime,
28 project::{
29 DATABASE_FILE_NAME, LINUX_SERVICE_DEFINITION_FILE_NAME, MACOS_SERVICE_DEFINITION_FILE_NAME,
30 PROJECT_NAME, WINDOWS_SERVICE_DEFINITION_FILE_NAME,
31 },
32 retrieval::{
33 RetrievalPlan,
34 provider::{EmbeddingRequest, ProviderRetryClass, embedding_provider},
35 read_model_backend_statuses,
36 },
37 storage::{
38 FileIndexDiagnostics, GraphCanvasSelection, GraphCanvasStorageRequest, GraphInspection,
39 GraphSearchRequest, KnowledgeStore, NewAuditEvent, ProposalListRequest, SqliteGraphStore,
40 StorageError,
41 },
42};
43
44use super::{
45 RuntimeConfiguration, RuntimeConfigurationError,
46 index_refresh::{
47 filter_outcome_to_read_models, index_refresh_outcome, metadata_for_indexes,
48 reconcile_index_refreshes, recover_index_kinds, refresh_index_kinds,
49 },
50 ingest::mutation_batch_from_request,
51 multimodal::extraction_ingest_request,
52 status::{agent_protocol_status, runtime_status, runtime_status_with_model_profiles},
53 update::{VersionCheckResponse, check_for_updates},
54};
55
56#[cfg(test)]
57use super::ingest::generated_evidence_id;
58
59#[derive(Clone)]
61pub struct RelayKnowledgeService {
62 pub(super) runtime: RuntimeConfiguration,
63 pub(super) storage: StorageProvider,
64}
65
66impl RelayKnowledgeService {
67 pub fn new(runtime: RuntimeConfiguration) -> Self {
69 let database_path = runtime.paths.data_dir.join(DATABASE_FILE_NAME);
70
71 Self {
72 runtime,
73 storage: StorageProvider::sqlite(database_path),
74 }
75 }
76
77 pub fn with_store(runtime: RuntimeConfiguration, store: Arc<dyn KnowledgeStore>) -> Self {
79 Self {
80 runtime,
81 storage: StorageProvider::ready(store),
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 metadata = ApiMetadata::graph_only(&context, graph_version);
295 let mut degraded_reasons = Vec::new();
296 let backend_statuses = if plan.freshness == FreshnessPolicy::GraphOnly {
297 retrieval_mode = RetrievalMode::GraphOnly;
298 degraded_reasons.push("graph_only freshness policy selected".to_owned());
299 Vec::new()
300 } else {
301 indexes = store.index_statuses().await.map_err(storage_api_error)?;
302 let mut active_indexes = indexes
303 .iter()
304 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
305 .cloned()
306 .collect::<Vec<_>>();
307 if plan.freshness == FreshnessPolicy::WaitUntilFresh {
308 let stale_kinds = active_indexes
309 .iter()
310 .filter(|status| status.is_stale_for(graph_version))
311 .map(|status| status.kind)
312 .collect::<Vec<_>>();
313 if !stale_kinds.is_empty() {
314 refresh_index_kinds(
315 &store,
316 stale_kinds,
317 graph_version,
318 &self.runtime.retrieval,
319 )
320 .await?;
321 indexes = store.index_statuses().await.map_err(storage_api_error)?;
322 active_indexes = indexes
323 .iter()
324 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
325 .cloned()
326 .collect();
327 }
328 }
329
330 let stale = active_indexes
331 .iter()
332 .any(|status| status.is_stale_for(graph_version));
333 metadata = metadata_for_indexes(&context, graph_version, &active_indexes);
334 if plan.freshness == FreshnessPolicy::AllowStale && stale {
335 degraded_reasons
336 .push("one or more indexes are behind the graph version".to_owned());
337 }
338 read_model_backend_statuses(&plan, graph_version, &indexes, &self.runtime.retrieval)
339 };
340 if backend_statuses
341 .iter()
342 .any(|status| status.state == crate::domain::RetrievalBackendState::Unavailable)
343 {
344 degraded_reasons.push(
345 "semantic/vector retrieval backends unavailable; using bm25, graph evidence, and code graph fallback"
346 .to_owned(),
347 );
348 }
349 let mut disabled_retriever_sources = self.runtime.retrieval.disabled_retriever_sources();
350 if plan.freshness == FreshnessPolicy::GraphOnly {
351 for source in [RetrieverSource::Semantic, RetrieverSource::Vector] {
352 if !disabled_retriever_sources.contains(&source) {
353 disabled_retriever_sources.push(source);
354 }
355 }
356 }
357 let candidate_limit = self.runtime.retrieval.rerank.candidate_limit(plan.limit);
358 let results = store
359 .search(GraphSearchRequest {
360 query: plan.query.clone(),
361 source_scope: plan.source_scope.clone(),
362 graph_version,
363 limit: candidate_limit,
364 disabled_retriever_sources,
365 })
366 .await
367 .map_err(storage_api_error)?;
368 let (mut results, mut rerank) = self.runtime.retrieval.rerank.rerank(&plan.query, results);
369 let truncated = results.len() > plan.limit;
370 results.truncate(plan.limit);
371 rerank.returned_count = results.len();
372 if rerank.degraded {
373 if let Some(reason) = &rerank.reason {
374 degraded_reasons.push(reason.clone());
375 }
376 }
377 let context_pack = RetrievedContextPack {
378 graph_version,
379 source_scope: plan.source_scope.clone(),
380 freshness: plan.freshness,
381 truncated,
382 backend_statuses: backend_statuses.clone(),
383 items: results
384 .iter()
385 .map(|hit| ContextPackItem {
386 result_id: hit.evidence_id.clone(),
387 source_scope: hit.source_scope.clone(),
388 source_path: hit.source_path.clone(),
389 source_span: hit.source_span,
390 entities: hit.entities.clone(),
391 graph_facts: hit.graph_facts.clone(),
392 graph_paths: hit
393 .graph_facts
394 .iter()
395 .map(ContextGraphPath::from_fact)
396 .collect(),
397 code_artifact: hit.code_artifact.clone(),
398 retriever_sources: hit.retriever_sources.clone(),
399 ranking: hit.ranking.clone(),
400 rerank: hit.rerank.clone(),
401 })
402 .collect(),
403 };
404 let budget_used = RetrievalBudgetUsed {
405 limit: plan.limit,
406 candidate_count: rerank.candidate_count,
407 returned_count: results.len(),
408 context_bytes: retrieval_context_bytes(&results, &context_pack, &backend_statuses),
409 };
410 let fusion = FusionDiagnostics {
411 algorithm: "reciprocal_rank_fusion".to_owned(),
412 k: RECIPROCAL_RANK_FUSION_K,
413 candidate_count: budget_used.candidate_count,
414 };
415 let degraded_reason = (!degraded_reasons.is_empty()).then(|| degraded_reasons.join("; "));
416
417 Ok(HybridRetrievalResponse {
418 metadata,
419 context_pack,
420 retrieval_mode,
421 source_scope: plan.source_scope,
422 freshness: plan.freshness,
423 results,
424 fusion,
425 rerank,
426 backend_statuses,
427 truncated,
428 budget_used,
429 degraded_reason,
430 indexes,
431 })
432 }
433
434 pub async fn inspect_graph(
436 &self,
437 _request: GraphInspectionRequest,
438 context: RequestContext,
439 ) -> Result<GraphInspectionResponse, ApiError> {
440 let store = self.storage.get().await.map_err(storage_api_error)?;
441 let repository_code_totals = store
442 .code_repository_totals()
443 .await
444 .map_err(storage_api_error)?;
445 let graph = graph_with_repository_code_totals(
446 store.inspect_graph().await.map_err(storage_api_error)?,
447 &repository_code_totals,
448 );
449
450 Ok(GraphInspectionResponse {
451 metadata: ApiMetadata::graph_only(&context, graph.graph_version),
452 graph,
453 repository_code_totals,
454 })
455 }
456
457 pub async fn graph_canvas(
459 &self,
460 request: GraphCanvasRequest,
461 context: RequestContext,
462 ) -> Result<GraphCanvasResponse, ApiError> {
463 if request.limit == 0 || request.limit > GRAPH_CANVAS_MAX_LIMIT {
464 return Err(ApiError::invalid_argument(format!(
465 "graph canvas limit must be between 1 and {GRAPH_CANVAS_MAX_LIMIT}"
466 )));
467 }
468 let store = self.storage.get().await.map_err(storage_api_error)?;
469 let graph_version = store
470 .current_graph_version()
471 .await
472 .map_err(storage_api_error)?;
473 let snapshot = store
474 .graph_canvas(GraphCanvasStorageRequest {
475 selection: canvas_selection(request.kind),
476 source_scope: request.source_scope,
477 query: request.query,
478 graph_version,
479 limit: request.limit,
480 })
481 .await
482 .map_err(storage_api_error)?;
483 let node_count = snapshot.nodes.len();
484 let edge_count = snapshot.edges.len();
485
486 Ok(GraphCanvasResponse {
487 metadata: ApiMetadata::graph_only(&context, graph_version),
488 nodes: snapshot
489 .nodes
490 .into_iter()
491 .map(|node| GraphCanvasNode {
492 id: node.id,
493 kind: node.kind,
494 label: node.label,
495 subtitle: node.subtitle,
496 source_scope: node.source_scope,
497 graph_version: node.graph_version.get(),
498 weight: node.weight,
499 status: node.status,
500 details: node.details,
501 })
502 .collect(),
503 edges: snapshot
504 .edges
505 .into_iter()
506 .map(|edge| GraphCanvasEdge {
507 id: edge.id,
508 kind: edge.kind,
509 source: edge.source,
510 target: edge.target,
511 label: edge.label,
512 graph_version: edge.graph_version.get(),
513 confidence_basis_points: edge.confidence_basis_points,
514 evidence_count: edge.evidence_count,
515 details: edge.details,
516 })
517 .collect(),
518 summary: GraphCanvasSummary {
519 kind: request.kind,
520 node_count,
521 edge_count,
522 truncated: snapshot.truncated,
523 available_kinds: snapshot.available_kinds,
524 },
525 })
526 }
527
528 pub async fn refresh_indexes(
530 &self,
531 request: IndexRefreshRequest,
532 context: RequestContext,
533 ) -> Result<IndexRefreshResponse, ApiError> {
534 let store = self.storage.get().await.map_err(storage_api_error)?;
535 let graph_version = store
536 .current_graph_version()
537 .await
538 .map_err(storage_api_error)?;
539 let outcome = refresh_index_kinds(
540 &store,
541 request.kinds,
542 graph_version,
543 &self.runtime.retrieval,
544 )
545 .await?;
546 let metadata = metadata_for_indexes(&context, graph_version, &outcome.indexes);
547
548 Ok(IndexRefreshResponse {
549 metadata,
550 indexes: outcome.indexes,
551 index_cursors: outcome.cursors,
552 diagnostics: outcome.diagnostics,
553 })
554 }
555
556 pub async fn health(&self, context: RequestContext) -> Result<HealthResponse, ApiError> {
558 let store = self.storage.get().await.map_err(storage_api_error)?;
559 let repository_code_totals = store
560 .code_repository_totals()
561 .await
562 .map_err(storage_api_error)?;
563 let graph = graph_with_repository_code_totals(
564 store.inspect_graph().await.map_err(storage_api_error)?,
565 &repository_code_totals,
566 );
567 reconcile_index_refreshes(&store, graph.graph_version, &self.runtime.retrieval).await?;
568 let outcome = filter_outcome_to_read_models(
569 index_refresh_outcome(&store).await?,
570 &self.runtime.retrieval,
571 );
572 let healthy = outcome
573 .indexes
574 .iter()
575 .all(|status| !status.is_stale_for(graph.graph_version));
576 let file_index = file_index_diagnostics_or_default(&store).await?;
577
578 Ok(HealthResponse {
579 metadata: metadata_for_indexes(&context, graph.graph_version, &outcome.indexes),
580 healthy,
581 graph,
582 repository_code_totals,
583 indexes: outcome.indexes,
584 index_cursors: outcome.cursors,
585 index_refresh: outcome.diagnostics,
586 file_index,
587 runtime: runtime_status_with_model_profiles(
588 &self.runtime,
589 self.model_provider_config()
590 .profile_summary(&self.runtime.retrieval)
591 .await,
592 ),
593 })
594 }
595
596 pub async fn probe_embedding_provider(
598 &self,
599 context: RequestContext,
600 ) -> Result<EmbeddingProviderProbeResponse, ApiError> {
601 let Some(remote) = self.runtime.retrieval.remote_embedding.clone() else {
602 return Ok(EmbeddingProviderProbeResponse {
603 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
604 ok: false,
605 provider: None,
606 model: self.runtime.retrieval.vector_model.name.clone(),
607 dimension: self.runtime.retrieval.vector_model.dimension,
608 latency_ms: None,
609 error_code: Some("remote_embedding_not_configured".to_owned()),
610 error_message: Some("remote embedding provider is not configured".to_owned()),
611 retryable: Some(false),
612 });
613 };
614 let network = self.runtime.network.current();
615 let client = crate::net::http::outbound_json_client(&network.http).map_err(|error| {
616 ApiError::invalid_argument(format!("failed to build HTTP client: {error}"))
617 })?;
618 let provider_name = remote.provider.as_str().to_owned();
619 let provider = embedding_provider(remote, client);
620 let started = Instant::now();
621 let result = provider
622 .embed(EmbeddingRequest {
623 inputs: vec!["relay-knowledge provider probe".to_owned()],
624 model: self.runtime.retrieval.vector_model.name.clone(),
625 dimension: self.runtime.retrieval.vector_model.dimension,
626 })
627 .await;
628
629 match result {
630 Ok(_) => Ok(EmbeddingProviderProbeResponse {
631 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
632 ok: true,
633 provider: Some(provider_name),
634 model: self.runtime.retrieval.vector_model.name.clone(),
635 dimension: self.runtime.retrieval.vector_model.dimension,
636 latency_ms: Some(duration_millis(started.elapsed())),
637 error_code: None,
638 error_message: None,
639 retryable: None,
640 }),
641 Err(error) => Ok(EmbeddingProviderProbeResponse {
642 metadata: ApiMetadata::graph_only(&context, crate::domain::GraphVersion::ZERO),
643 ok: error.code == "rate_limited" && error.retry == ProviderRetryClass::Retryable,
644 provider: Some(provider_name),
645 model: self.runtime.retrieval.vector_model.name.clone(),
646 dimension: self.runtime.retrieval.vector_model.dimension,
647 latency_ms: Some(duration_millis(started.elapsed())),
648 error_code: Some(error.code),
649 error_message: Some(error.message),
650 retryable: Some(error.retry == ProviderRetryClass::Retryable),
651 }),
652 }
653 }
654
655 pub async fn reconcile_startup_indexes(
657 &self,
658 context: RequestContext,
659 ) -> Result<ServiceRecoveryReport, ApiError> {
660 let store = self.storage.get().await.map_err(storage_api_error)?;
661 let graph_version = store
662 .current_graph_version()
663 .await
664 .map_err(storage_api_error)?;
665 let before = store.index_statuses().await.map_err(storage_api_error)?;
666 let active_before = before
667 .iter()
668 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
669 .cloned()
670 .collect::<Vec<_>>();
671 let stale_index_kinds = active_before
672 .iter()
673 .filter(|status| status.is_stale_for(graph_version))
674 .map(|status| status.kind)
675 .collect::<Vec<_>>();
676 let index_lag_max = active_before
677 .iter()
678 .map(|status| {
679 graph_version
680 .get()
681 .saturating_sub(status.indexed_graph_version.get())
682 })
683 .max()
684 .unwrap_or(0);
685 let outcome = if stale_index_kinds.is_empty() {
686 index_refresh_outcome(&store).await?
687 } else {
688 recover_index_kinds(
689 &store,
690 stale_index_kinds.clone(),
691 graph_version,
692 &self.runtime.retrieval,
693 )
694 .await?
695 };
696 let refreshed = outcome
697 .indexes
698 .iter()
699 .filter(|status| {
700 stale_index_kinds.contains(&status.kind) && !status.is_stale_for(graph_version)
701 })
702 .map(|status| status.kind)
703 .collect::<Vec<_>>();
704 let after = outcome.indexes;
705 let active_after = after
706 .iter()
707 .filter(|status| self.runtime.retrieval.refreshes_index(status.kind))
708 .cloned()
709 .collect::<Vec<_>>();
710 let metadata = metadata_for_indexes(&context, graph_version, &active_after);
711
712 Ok(ServiceRecoveryReport {
713 metadata,
714 graph_version: graph_version.get(),
715 stale_index_kinds,
716 refreshed_index_kinds: refreshed,
717 index_lag_max,
718 task_queue_depth: outcome.diagnostics.queue_depth,
719 dead_letter_count: outcome.diagnostics.dead_letter_count,
720 heartbeat_state: "ready".to_owned(),
721 })
722 }
723
724 pub async fn service_status(
726 &self,
727 context: RequestContext,
728 ) -> Result<ServiceStatusResponse, ApiError> {
729 let store = self.storage.get().await.map_err(storage_api_error)?;
730 let graph_version = store
731 .current_graph_version()
732 .await
733 .map_err(storage_api_error)?;
734 let index_refresh =
735 reconcile_index_refreshes(&store, graph_version, &self.runtime.retrieval).await?;
736 let file_index = file_index_diagnostics_or_default(&store).await?;
737 let service_definition_path = self
738 .runtime
739 .paths
740 .service_dir
741 .join(service_definition_filename())
742 .display()
743 .to_string();
744 let operator = store
745 .service_operator_status()
746 .await
747 .map_err(storage_api_error)?;
748 let workers = super::operations::overlay_worker_runtime(
749 store.worker_statuses().await.map_err(storage_api_error)?,
750 &self.runtime.workers,
751 );
752 let proposal_backlog = store
753 .list_proposals(ProposalListRequest {
754 state: Some(ProposalState::Proposed),
755 limit: usize::MAX,
756 })
757 .await
758 .map_err(storage_api_error)?
759 .len();
760 let audit_event_count = store.audit_event_count().await.map_err(storage_api_error)?;
761
762 Ok(ServiceStatusResponse {
763 metadata: ApiMetadata::graph_only(&context, graph_version),
764 service_name: PROJECT_NAME.to_owned(),
765 mode: operator.state.as_str().to_owned(),
766 background_enabled: operator.state != crate::domain::ServiceOperatorState::Disabled,
767 silent_updates_enabled: operator.silent_updates_enabled,
768 service_definition_path,
769 index_refresh,
770 file_index,
771 agent_protocols: agent_protocol_status(&self.runtime),
772 operator,
773 workers,
774 proposal_backlog,
775 audit_sink: crate::api::AuditSinkStatus {
776 durable: true,
777 event_count: audit_event_count,
778 last_error: None,
779 },
780 })
781 }
782
783 pub(super) async fn store(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
784 self.storage.get().await
785 }
786
787 pub fn agent_audit_log_path(&self) -> PathBuf {
789 self.runtime.paths.agent_audit_log_file()
790 }
791}
792
793#[derive(Debug, Clone)]
795pub struct AgentDurableAuditInput {
796 pub operation: String,
797 pub interface: String,
798 pub request_id: String,
799 pub trace_id: String,
800 pub status: AuditStatus,
801 pub actor: Option<String>,
802 pub source_scope: Option<String>,
803 pub graph_version: u64,
804 pub detail_json: String,
805 pub message: Option<String>,
806}
807
808fn current_time_millis() -> u64 {
809 std::time::SystemTime::now()
810 .duration_since(std::time::UNIX_EPOCH)
811 .map_or(0, |duration| {
812 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
813 })
814}
815
816#[derive(Clone)]
817pub(super) struct StorageProvider {
818 path: Option<PathBuf>,
819 ready: Arc<OnceLock<Arc<dyn KnowledgeStore>>>,
820 init_lock: Arc<tokio::sync::Mutex<()>>,
821}
822
823impl StorageProvider {
824 fn sqlite(path: PathBuf) -> Self {
825 Self {
826 path: Some(path),
827 ready: Arc::new(OnceLock::new()),
828 init_lock: Arc::new(tokio::sync::Mutex::new(())),
829 }
830 }
831
832 fn ready(store: Arc<dyn KnowledgeStore>) -> Self {
833 let ready = OnceLock::new();
834 let _ = ready.set(store);
835
836 Self {
837 path: None,
838 ready: Arc::new(ready),
839 init_lock: Arc::new(tokio::sync::Mutex::new(())),
840 }
841 }
842
843 pub(super) async fn get(&self) -> Result<Arc<dyn KnowledgeStore>, StorageError> {
844 if let Some(store) = self.ready.get() {
845 return Ok(Arc::clone(store));
846 }
847 let _guard = self.init_lock.lock().await;
848 if let Some(store) = self.ready.get() {
849 return Ok(Arc::clone(store));
850 }
851
852 let Some(path) = self.path.clone() else {
853 return Err(StorageError::InvalidInput(
854 "storage provider was not initialized".to_owned(),
855 ));
856 };
857 let ready = Arc::clone(&self.ready);
858 tokio::task::spawn_blocking(move || {
859 if let Some(store) = ready.get() {
860 return Ok(Arc::clone(store));
861 }
862 let store = Arc::new(SqliteGraphStore::open(path)?) as Arc<dyn KnowledgeStore>;
863 let _ = ready.set(Arc::clone(&store));
864 Ok(store)
865 })
866 .await?
867 }
868}
869
870fn storage_api_error(error: StorageError) -> ApiError {
871 ApiError::storage_unavailable(error.to_string())
872}
873
874async fn file_index_diagnostics_or_default(
875 store: &Arc<dyn KnowledgeStore>,
876) -> Result<FileIndexDiagnostics, ApiError> {
877 match store.file_index_diagnostics().await {
878 Ok(diagnostics) => Ok(diagnostics),
879 Err(StorageError::InvalidInput(message))
880 if message == "file index storage is unavailable" =>
881 {
882 Ok(FileIndexDiagnostics::default())
883 }
884 Err(error) => Err(storage_api_error(error)),
885 }
886}
887
888fn normalize_optional_source_scope(value: Option<String>) -> Result<Option<String>, String> {
889 value
890 .map(|scope| {
891 SourceScope::parse(scope)
892 .map(String::from)
893 .map_err(|error| error.to_string())
894 })
895 .transpose()
896}
897
898fn retrieval_context_bytes(
899 results: &[RetrievalHit],
900 context_pack: &RetrievedContextPack,
901 backend_statuses: &[RetrievalBackendStatus],
902) -> usize {
903 serialized_context_bytes(&context_pack.backend_statuses)
904 .saturating_add(serialized_context_bytes(backend_statuses))
905 .saturating_add(results.iter().map(serialized_context_bytes).sum::<usize>())
906 .saturating_add(
907 context_pack
908 .items
909 .iter()
910 .map(serialized_context_bytes)
911 .sum::<usize>(),
912 )
913}
914
915fn graph_with_repository_code_totals(
916 mut graph: GraphInspection,
917 repository_totals: &CodeRepositoryTotals,
918) -> GraphInspection {
919 graph.code_file_count = graph
920 .code_file_count
921 .saturating_add(repository_totals.indexed_file_count);
922 graph.code_symbol_count = graph
923 .code_symbol_count
924 .saturating_add(repository_totals.symbol_count);
925 graph.code_reference_count = graph
926 .code_reference_count
927 .saturating_add(repository_totals.reference_count);
928 graph.code_chunk_count = graph
929 .code_chunk_count
930 .saturating_add(repository_totals.chunk_count);
931 graph.code_parse_status_counts = add_parse_status_counts(
932 graph.code_parse_status_counts,
933 repository_totals.parse_status_counts,
934 );
935
936 graph
937}
938
939fn canvas_selection(kind: GraphCanvasKind) -> GraphCanvasSelection {
940 match kind {
941 GraphCanvasKind::Knowledge => GraphCanvasSelection::Knowledge,
942 GraphCanvasKind::Code => GraphCanvasSelection::Code,
943 GraphCanvasKind::Mixed => GraphCanvasSelection::Mixed,
944 }
945}
946
947fn add_parse_status_counts(
948 left: CodeParseStatusCounts,
949 right: CodeParseStatusCounts,
950) -> CodeParseStatusCounts {
951 CodeParseStatusCounts {
952 parsed: left.parsed.saturating_add(right.parsed),
953 partial: left.partial.saturating_add(right.partial),
954 text_only: left.text_only.saturating_add(right.text_only),
955 failed: left.failed.saturating_add(right.failed),
956 }
957}
958
959fn serialized_context_bytes<T: Serialize + ?Sized>(value: &T) -> usize {
960 serde_json::to_vec(value)
961 .map(|bytes| bytes.len())
962 .unwrap_or(usize::MAX / 4)
963}
964
965fn duration_millis(duration: std::time::Duration) -> u64 {
966 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
967}
968
969fn service_definition_filename() -> &'static str {
970 if cfg!(target_os = "windows") {
971 WINDOWS_SERVICE_DEFINITION_FILE_NAME
972 } else if cfg!(target_os = "macos") {
973 MACOS_SERVICE_DEFINITION_FILE_NAME
974 } else {
975 LINUX_SERVICE_DEFINITION_FILE_NAME
976 }
977}
978
979#[cfg(test)]
980mod id_tests;
981
982#[cfg(test)]
983mod graph_only_tests;
984
985#[cfg(test)]
986mod recovery_tests;
987
988#[cfg(test)]
989mod refresh_tests;
990
991#[cfg(test)]
992mod storage_tests;
993
994#[cfg(test)]
995mod operations_tests;
996
997#[cfg(test)]
998mod tests;