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        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}