relay_knowledge/application/service/
service_status.rs1use crate::{
2 api::{
3 ApiError, ApiMetadata, AuditSinkStatus, CodeIndexWorkerStatus, RequestContext,
4 ServiceStatusResponse, WatcherDiagnostics,
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 watcher: self.watcher_diagnostics().await,
112 }))
113 }
114
115 async fn storage_free_service_status(&self, context: RequestContext) -> ServiceStatusResponse {
116 let operator = crate::domain::ServiceOperatorStatus {
117 state: ServiceOperatorState::Disabled,
118 silent_updates_enabled: self.runtime.workers.silent_updates_enabled,
119 allowed_scopes: Vec::new(),
120 last_run_at_ms: None,
121 next_retry_at_ms: None,
122 last_error: None,
123 updated_at_ms: 0,
124 };
125
126 self.service_status_response(ServiceStatusParts {
127 context,
128 graph_version: GraphVersion::ZERO,
129 storage: self.storage_topology_diagnostics().await,
130 index_refresh: IndexRefreshDiagnostics::default(),
131 file_index: FileIndexDiagnostics::default(),
132 operator,
133 workers: overlay_worker_runtime(Vec::new(), &self.runtime.workers),
134 code_index_workers: CodeIndexWorkerStatus::from_queue(
135 self.runtime.workers.code_index_max_in_flight,
136 CodeIndexTaskQueueStatus::default(),
137 ),
138 proposal_backlog: 0,
139 audit_sink: AuditSinkStatus {
140 durable: true,
141 event_count: 0,
142 last_error: Some(
143 "audit event count not sampled because storage is not open".to_owned(),
144 ),
145 },
146 watcher: self.watcher_diagnostics().await,
147 })
148 }
149
150 fn service_status_response(&self, parts: ServiceStatusParts) -> ServiceStatusResponse {
151 let ServiceStatusParts {
152 context,
153 graph_version,
154 storage,
155 index_refresh,
156 file_index,
157 operator,
158 workers,
159 code_index_workers,
160 proposal_backlog,
161 audit_sink,
162 watcher,
163 } = parts;
164
165 ServiceStatusResponse {
166 metadata: ApiMetadata::graph_only(&context, graph_version),
167 service_name: PROJECT_NAME.to_owned(),
168 mode: operator.state.as_str().to_owned(),
169 background_enabled: operator.state != ServiceOperatorState::Disabled,
170 silent_updates_enabled: operator.silent_updates_enabled,
171 service_definition_path: self
172 .runtime
173 .paths
174 .service_dir
175 .join(service_definition_filename())
176 .display()
177 .to_string(),
178 storage,
179 index_refresh,
180 file_index,
181 agent_protocols: agent_protocol_status(&self.runtime),
182 operator,
183 workers,
184 code_index_workers,
185 proposal_backlog,
186 audit_sink,
187 watcher: watcher
188 .as_ref()
189 .map(WatcherDiagnostics::from_watcher_state)
190 .unwrap_or_else(WatcherDiagnostics::default_disabled),
191 }
192 }
193}
194
195struct ServiceStatusParts {
196 context: RequestContext,
197 graph_version: GraphVersion,
198 storage: crate::api::StorageTopologyDiagnostics,
199 index_refresh: IndexRefreshDiagnostics,
200 file_index: FileIndexDiagnostics,
201 operator: crate::domain::ServiceOperatorStatus,
202 workers: Vec<crate::domain::WorkerStatus>,
203 code_index_workers: CodeIndexWorkerStatus,
204 proposal_backlog: usize,
205 audit_sink: AuditSinkStatus,
206 watcher: Option<crate::watcher::WatcherDiagnostics>,
207}