Skip to main content

relay_knowledge/application/service/
service_status.rs

1use 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    /// Returns the managed background service definition location and defaults.
31    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    /// Returns control-plane service diagnostics without opening cold storage.
40    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}