Skip to main content

relay_knowledge/api/operations/
worker.rs

1use serde::{Deserialize, Serialize};
2
3use crate::{
4    api::ApiMetadata,
5    domain::{
6        CodeIndexTaskQueueStatus, CodeIndexTaskRecord, ProposalRecord, WorkerKind, WorkerStatus,
7        WorkerTaskRecord,
8    },
9};
10
11/// Master-worker diagnostics for repository code indexing.
12#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
13pub struct CodeIndexWorkerStatus {
14    pub configured_worker_count: usize,
15    pub active_worker_slots: usize,
16    pub queue_depth: usize,
17    pub queued_task_count: usize,
18    pub running_task_count: usize,
19    pub retrying_task_count: usize,
20    pub dead_letter_task_count: usize,
21    pub running_lease_count: usize,
22    #[serde(skip_serializing_if = "Option::is_none")]
23    pub last_error: Option<String>,
24}
25
26impl CodeIndexWorkerStatus {
27    /// Overlays resident master runtime configuration onto durable task queue state.
28    pub fn from_queue(configured_worker_count: usize, queue: CodeIndexTaskQueueStatus) -> Self {
29        let queue_depth = queue
30            .queued_task_count
31            .saturating_add(queue.retrying_task_count);
32
33        Self {
34            configured_worker_count,
35            active_worker_slots: configured_worker_count.saturating_sub(queue.running_task_count),
36            queue_depth,
37            queued_task_count: queue.queued_task_count,
38            running_task_count: queue.running_task_count,
39            retrying_task_count: queue.retrying_task_count,
40            dead_letter_task_count: queue.dead_letter_task_count,
41            running_lease_count: queue.running_lease_count,
42            last_error: queue.last_error,
43        }
44    }
45}
46
47/// Worker status filter. Missing kind means all worker families.
48#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
49pub struct WorkerStatusRequest {
50    #[serde(skip_serializing_if = "Option::is_none")]
51    pub kind: Option<WorkerKind>,
52}
53
54/// Worker status response.
55#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
56pub struct WorkerStatusResponse {
57    pub metadata: ApiMetadata,
58    pub workers: Vec<WorkerStatus>,
59}
60
61/// Bounded foreground worker run request.
62#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
63pub struct WorkerRunRequest {
64    #[serde(skip_serializing_if = "Option::is_none")]
65    pub kind: Option<WorkerKind>,
66}
67
68/// Bounded foreground worker run response.
69#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
70pub struct WorkerRunResponse {
71    pub metadata: ApiMetadata,
72    #[serde(skip_serializing_if = "Option::is_none")]
73    pub task: Option<WorkerTaskRecord>,
74    #[serde(default)]
75    pub proposals: Vec<ProposalRecord>,
76    pub workers: Vec<WorkerStatus>,
77    #[serde(skip_serializing_if = "Option::is_none")]
78    pub degraded_reason: Option<String>,
79}
80
81/// Preview split-worker request for one durable code-index task.
82#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
83pub struct CodeIndexWorkerRunRequest {
84    #[serde(skip_serializing_if = "Option::is_none")]
85    pub task_id: Option<String>,
86}
87
88/// Preview split-worker result for one durable code-index task.
89#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
90pub struct CodeIndexWorkerRunResponse {
91    pub metadata: ApiMetadata,
92    pub worker_kind: String,
93    pub claimed: bool,
94    #[serde(skip_serializing_if = "Option::is_none")]
95    pub task: Option<CodeIndexTaskRecord>,
96}
97
98#[cfg(test)]
99#[path = "worker_tests.rs"]
100mod tests;