relay_knowledge/application/service/
service_status.rs1use crate::{
2 api::{
3 ApiError, ApiMetadata, AuditSinkStatus, CodeIndexWorkerStatus, RequestContext,
4 ServiceStatusResponse,
5 },
6 domain::{CodeIndexTaskQueueStatus, GraphVersion, ProposalState, ServiceOperatorState},
7 project::PROJECT_NAME,
8 storage::{FileIndexDiagnostics, IndexRefreshDiagnostics},
9};
10
11use super::{
12 RelayKnowledgeService, file_index_diagnostics_or_default, service_definition_filename,
13 storage_api_error,
14};
15use crate::application::{
16 knowledge::index_refresh::{
17 filter_outcome_to_read_models, index_refresh_outcome, reconcile_index_refreshes,
18 },
19 status::agent_protocol_status,
20 worker::operations::overlay_worker_runtime,
21};
22
23#[derive(Clone, Copy)]
24enum ServiceStatusRefreshMode {
25 Reconcile,
26 Observe,
27}
28
29impl RelayKnowledgeService {
30 pub async fn service_status(
32 &self,
33 context: RequestContext,
34 ) -> Result<ServiceStatusResponse, ApiError> {
35 self.service_status_with_refresh_mode(context, ServiceStatusRefreshMode::Reconcile)
36 .await
37 }
38
39 pub async fn read_only_service_status(
41 &self,
42 context: RequestContext,
43 ) -> Result<ServiceStatusResponse, ApiError> {
44 if self.storage.ready_store().is_none() {
45 return Ok(self.storage_free_service_status(context).await);
46 }
47
48 self.service_status_with_refresh_mode(context, ServiceStatusRefreshMode::Observe)
49 .await
50 }
51
52 async fn service_status_with_refresh_mode(
53 &self,
54 context: RequestContext,
55 refresh_mode: ServiceStatusRefreshMode,
56 ) -> Result<ServiceStatusResponse, ApiError> {
57 let store = self.storage.get().await.map_err(storage_api_error)?;
58 let graph_version = store
59 .current_graph_version()
60 .await
61 .map_err(storage_api_error)?;
62 let index_refresh = match refresh_mode {
63 ServiceStatusRefreshMode::Reconcile => {
64 reconcile_index_refreshes(&store, graph_version, &self.runtime.retrieval).await?
65 }
66 ServiceStatusRefreshMode::Observe => {
67 filter_outcome_to_read_models(
68 index_refresh_outcome(&store).await?,
69 &self.runtime.retrieval,
70 )
71 .diagnostics
72 }
73 };
74 let file_index = file_index_diagnostics_or_default(&store).await?;
75 let storage = self.storage_topology_diagnostics().await;
76 let operator = store
77 .service_operator_status()
78 .await
79 .map_err(storage_api_error)?;
80 let workers = overlay_worker_runtime(
81 store.worker_statuses().await.map_err(storage_api_error)?,
82 &self.runtime.workers,
83 );
84 let code_index_workers = match refresh_mode {
85 ServiceStatusRefreshMode::Reconcile => self.code_index_worker_status(&store).await?,
86 ServiceStatusRefreshMode::Observe => {
87 self.read_only_code_index_worker_status(&store).await?
88 }
89 };
90 let proposal_backlog = store
91 .proposal_count(Some(ProposalState::Proposed))
92 .await
93 .map_err(storage_api_error)?;
94 let audit_event_count = store.audit_event_count().await.map_err(storage_api_error)?;
95
96 Ok(self.service_status_response(ServiceStatusParts {
97 context,
98 graph_version,
99 storage,
100 index_refresh,
101 file_index,
102 operator,
103 workers,
104 code_index_workers,
105 proposal_backlog,
106 audit_sink: AuditSinkStatus {
107 durable: true,
108 event_count: audit_event_count,
109 last_error: None,
110 },
111 }))
112 }
113
114 async fn storage_free_service_status(&self, context: RequestContext) -> ServiceStatusResponse {
115 let operator = crate::domain::ServiceOperatorStatus {
116 state: ServiceOperatorState::Disabled,
117 silent_updates_enabled: self.runtime.workers.silent_updates_enabled,
118 allowed_scopes: Vec::new(),
119 last_run_at_ms: None,
120 next_retry_at_ms: None,
121 last_error: None,
122 updated_at_ms: 0,
123 };
124
125 self.service_status_response(ServiceStatusParts {
126 context,
127 graph_version: GraphVersion::ZERO,
128 storage: self.storage_topology_diagnostics().await,
129 index_refresh: IndexRefreshDiagnostics::default(),
130 file_index: FileIndexDiagnostics::default(),
131 operator,
132 workers: overlay_worker_runtime(Vec::new(), &self.runtime.workers),
133 code_index_workers: CodeIndexWorkerStatus::from_queue(
134 self.runtime.workers.code_index_max_in_flight,
135 CodeIndexTaskQueueStatus::default(),
136 ),
137 proposal_backlog: 0,
138 audit_sink: AuditSinkStatus {
139 durable: true,
140 event_count: 0,
141 last_error: Some(
142 "audit event count not sampled because storage is not open".to_owned(),
143 ),
144 },
145 })
146 }
147
148 fn service_status_response(&self, parts: ServiceStatusParts) -> ServiceStatusResponse {
149 let ServiceStatusParts {
150 context,
151 graph_version,
152 storage,
153 index_refresh,
154 file_index,
155 operator,
156 workers,
157 code_index_workers,
158 proposal_backlog,
159 audit_sink,
160 } = parts;
161
162 ServiceStatusResponse {
163 metadata: ApiMetadata::graph_only(&context, graph_version),
164 service_name: PROJECT_NAME.to_owned(),
165 mode: operator.state.as_str().to_owned(),
166 background_enabled: operator.state != ServiceOperatorState::Disabled,
167 silent_updates_enabled: operator.silent_updates_enabled,
168 service_definition_path: self
169 .runtime
170 .paths
171 .service_dir
172 .join(service_definition_filename())
173 .display()
174 .to_string(),
175 storage,
176 index_refresh,
177 file_index,
178 agent_protocols: agent_protocol_status(&self.runtime),
179 operator,
180 workers,
181 code_index_workers,
182 proposal_backlog,
183 audit_sink,
184 }
185 }
186}
187
188struct ServiceStatusParts {
189 context: RequestContext,
190 graph_version: GraphVersion,
191 storage: crate::api::StorageTopologyDiagnostics,
192 index_refresh: IndexRefreshDiagnostics,
193 file_index: FileIndexDiagnostics,
194 operator: crate::domain::ServiceOperatorStatus,
195 workers: Vec<crate::domain::WorkerStatus>,
196 code_index_workers: CodeIndexWorkerStatus,
197 proposal_backlog: usize,
198 audit_sink: AuditSinkStatus,
199}