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