Skip to main content

relay_knowledge/application/service/service_status/
mod.rs

1//! Managed-service status assembly and operator diagnostics.
2
3use 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    /// Returns the managed background service definition location and defaults.
33    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    /// Returns control-plane service diagnostics without opening cold storage.
42    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}