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 ServiceStatusResponse {
168 metadata: ApiMetadata::graph_only(&context, graph_version),
169 service_name: PROJECT_NAME.to_owned(),
170 mode: operator.state.as_str().to_owned(),
171 background_enabled: operator.state != ServiceOperatorState::Disabled,
172 silent_updates_enabled: operator.silent_updates_enabled,
173 service_definition_path: self
174 .runtime
175 .paths
176 .service_dir
177 .join(service_definition_filename())
178 .display()
179 .to_string(),
180 storage,
181 index_refresh,
182 file_index,
183 agent_protocols: agent_protocol_status(&self.runtime),
184 operator,
185 workers,
186 code_index_workers,
187 proposal_backlog,
188 audit_sink,
189 watcher: watcher
190 .as_ref()
191 .map(WatcherDiagnostics::from_watcher_state)
192 .unwrap_or_else(WatcherDiagnostics::default_disabled),
193 }
194 }
195}
196
197struct ServiceStatusParts {
198 context: RequestContext,
199 graph_version: GraphVersion,
200 storage: crate::api::StorageTopologyDiagnostics,
201 index_refresh: IndexRefreshDiagnostics,
202 file_index: FileIndexDiagnostics,
203 operator: crate::domain::ServiceOperatorStatus,
204 workers: Vec<crate::domain::WorkerStatus>,
205 code_index_workers: CodeIndexWorkerStatus,
206 proposal_backlog: usize,
207 audit_sink: AuditSinkStatus,
208 watcher: Option<crate::watcher::WatcherDiagnostics>,
209}