Skip to main content

relay_knowledge/application/service/
service_status.rs

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