relay_knowledge/api/operations/
worker.rs1use serde::{Deserialize, Serialize};
2
3use crate::{
4 api::ApiMetadata,
5 domain::{
6 CodeIndexTaskQueueStatus, CodeIndexTaskRecord, ProposalRecord, WorkerKind, WorkerStatus,
7 WorkerTaskRecord,
8 },
9};
10
11#[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 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#[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
56pub struct WorkerStatusResponse {
57 pub metadata: ApiMetadata,
58 pub workers: Vec<WorkerStatus>,
59}
60
61#[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#[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#[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#[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;