1use chrono::{DateTime, Utc};
2
3use crate::domain::runtime::cluster_node::{
4 ClusterNode, ClusterNodeHealth, ClusterNodeHealthSnapshot, ClusterNodeRole, QueueOwnershipMode,
5};
6use crate::domain::runtime::delivery_endpoint::DeliveryProtocol;
7
8#[derive(Clone, Debug)]
9pub struct RegisterAgentRequest {
10 pub id: String,
11 pub name: String,
12 pub system_prompt: String,
13}
14
15#[derive(Clone, Debug)]
16pub struct InvokeAgentRequest {
17 pub agent_id: String,
18 pub user_prompt: String,
19}
20
21#[derive(Clone, Debug)]
22pub struct InvokeAgentResponse {
23 pub agent_id: String,
24 pub completion: String,
25}
26
27#[derive(Clone, Debug)]
28pub struct RegisterDeliveryEndpointRequest {
29 pub endpoint_id: String,
30 pub name: String,
31 pub protocol: DeliveryProtocol,
32 pub target: String,
33 pub metadata: Option<String>,
34}
35
36#[derive(Clone, Debug)]
37pub struct RegisterDeliveryEndpointResponse {
38 pub endpoint_id: String,
39 pub enabled: bool,
40}
41
42#[derive(Clone, Debug)]
43pub struct SetDeliveryEndpointEnabledRequest {
44 pub endpoint_id: String,
45 pub enabled: bool,
46}
47
48#[derive(Clone, Debug, Default)]
49pub struct ListEndpointDiagnosticsReadModelRequest {
50 pub endpoint_ids: Option<Vec<String>>,
51 pub protocol: Option<DeliveryProtocol>,
52 pub min_failure_count: Option<u64>,
53 pub stale_after_seconds: Option<i64>,
54 pub unhealthy_only: bool,
55 pub include_disabled: bool,
56 pub offset: usize,
57 pub limit: Option<usize>,
58}
59
60#[derive(Clone, Debug)]
61pub struct EndpointDiagnosticsReadModelRow {
62 pub endpoint_id: String,
63 pub endpoint_name: String,
64 pub protocol: DeliveryProtocol,
65 pub target: String,
66 pub enabled: bool,
67 pub success_count: u64,
68 pub failure_count: u64,
69 pub last_event_id: Option<String>,
70 pub last_error: Option<String>,
71 pub last_success_at: Option<DateTime<Utc>>,
72 pub last_failure_at: Option<DateTime<Utc>>,
73 pub updated_at: DateTime<Utc>,
74 pub unhealthy: bool,
75}
76
77#[derive(Clone, Debug, Eq, PartialEq)]
78pub enum EndpointFailureTrendDirection {
79 Improving,
80 Stable,
81 Worsening,
82}
83
84#[derive(Clone, Debug)]
85pub struct EndpointFailureRateTrendRow {
86 pub endpoint_id: String,
87 pub endpoint_name: String,
88 pub protocol: DeliveryProtocol,
89 pub enabled: bool,
90 pub success_count: u64,
91 pub failure_count: u64,
92 pub total_attempts: u64,
93 pub failure_rate: f64,
94 pub trend: EndpointFailureTrendDirection,
95 pub last_success_at: Option<DateTime<Utc>>,
96 pub last_failure_at: Option<DateTime<Utc>>,
97 pub updated_at: DateTime<Utc>,
98}
99
100#[derive(Clone, Debug, Default)]
101pub struct ListTopUnhealthyEndpointsRequest {
102 pub protocol: Option<DeliveryProtocol>,
103 pub include_disabled: bool,
104 pub limit: usize,
105}
106
107#[derive(Clone, Debug, Default)]
108pub struct ListEndpointFailureRateTrendsRequest {
109 pub protocol: Option<DeliveryProtocol>,
110 pub include_disabled: bool,
111 pub min_total_attempts: Option<u64>,
112 pub limit: usize,
113}
114
115#[derive(Clone, Debug)]
116pub struct PruneEndpointDeliveryStatusesRequest {
117 pub updated_before: DateTime<Utc>,
118}
119
120#[derive(Clone, Debug)]
121pub struct RegisterClusterNodeRequest {
122 pub node_id: String,
123 pub role: ClusterNodeRole,
124 pub region: String,
125 pub queue_ownership: Vec<String>,
126 pub capability_tags: Vec<String>,
127 pub heartbeat_at: DateTime<Utc>,
128 pub lease_ttl_seconds: i64,
129 pub queue_ownership_mode: Option<QueueOwnershipMode>,
130 pub metadata: Option<String>,
131}
132
133#[derive(Clone, Debug)]
134pub struct HeartbeatClusterNodeRequest {
135 pub node_id: String,
136 pub heartbeat_at: DateTime<Utc>,
137 pub lease_ttl_seconds: i64,
138 pub queue_ownership_mode: Option<QueueOwnershipMode>,
139 pub queue_ownership: Option<Vec<String>>,
140 pub capability_tags: Option<Vec<String>>,
141 pub metadata: Option<String>,
142}
143
144#[derive(Clone, Debug, Default)]
145pub struct ListClusterNodeHealthRequest {
146 pub role: Option<ClusterNodeRole>,
147 pub region: Option<String>,
148 pub capability_tag: Option<String>,
149 pub queue: Option<String>,
150 pub health: Option<ClusterNodeHealth>,
151 pub offset: usize,
152 pub limit: Option<usize>,
153}
154
155#[derive(Clone, Debug)]
156pub struct ListQueueOwnershipHealthRequest {
157 pub queue_prefix: Option<String>,
158}
159
160#[derive(Clone, Debug)]
161pub struct QueueOwnershipHealthRow {
162 pub queue: String,
163 pub owners: Vec<String>,
164 pub healthy_owners: usize,
165 pub degraded_owners: usize,
166 pub offline_owners: usize,
167}
168
169#[derive(Clone, Debug)]
170pub struct PruneExpiredClusterNodesRequest {
171 pub now: DateTime<Utc>,
172}
173
174#[derive(Clone, Debug)]
175pub struct RunClusterHeartbeatSweepRequest {
176 pub now: DateTime<Utc>,
177}
178
179#[derive(Clone, Debug)]
180pub struct RunClusterHeartbeatSweepResponse {
181 pub pruned_nodes: u64,
182 pub emitted_events: u64,
183}
184
185#[derive(Clone, Debug)]
186pub struct ForwardClusterCommandRequest {
187 pub target_region: String,
188 pub command_name: String,
189 pub payload: String,
190 pub correlation_id: Option<String>,
191 pub issued_at: DateTime<Utc>,
192}
193
194#[derive(Clone, Debug)]
195pub struct ForwardClusterCommandResponse {
196 pub accepted: bool,
197}
198
199#[derive(Clone, Debug)]
200pub struct InitiateCoordinatorHandoffRequest {
201 pub target_region: String,
202 pub coordinator_node_id: String,
203 pub queue_scope: Option<Vec<String>>,
204 pub reason: Option<String>,
205 pub correlation_id: Option<String>,
206 pub issued_at: DateTime<Utc>,
207}
208
209#[derive(Clone, Debug)]
210pub struct InitiateCoordinatorHandoffResponse {
211 pub accepted: bool,
212 pub command_name: String,
213}
214
215#[derive(Clone, Debug)]
216pub struct InitiateCoordinatorFailoverRequest {
217 pub target_region: String,
218 pub coordinator_node_id: String,
219 pub failover_to_node_id: Option<String>,
220 pub queue_scope: Option<Vec<String>>,
221 pub reason: Option<String>,
222 pub correlation_id: Option<String>,
223 pub issued_at: DateTime<Utc>,
224}
225
226#[derive(Clone, Debug)]
227pub struct InitiateCoordinatorFailoverResponse {
228 pub accepted: bool,
229 pub command_name: String,
230}
231
232#[derive(Clone, Debug)]
233pub struct RebalanceQueueOwnershipRequest {
234 pub target_region: String,
235 pub queue: String,
236 pub desired_owners: Vec<String>,
237 pub strategy: Option<String>,
238 pub reason: Option<String>,
239 pub correlation_id: Option<String>,
240 pub issued_at: DateTime<Utc>,
241}
242
243#[derive(Clone, Debug)]
244pub struct RebalanceQueueOwnershipResponse {
245 pub accepted: bool,
246 pub command_name: String,
247}
248
249#[derive(Clone, Debug, Default)]
250pub struct ListClusterForwardOutcomesRequest {
251 pub target_region: Option<String>,
252 pub command_name: Option<String>,
253 pub accepted: Option<bool>,
254 pub limit: usize,
255}
256
257#[derive(Clone, Debug)]
258pub struct ClusterForwardOutcomeRow {
259 pub target_region: String,
260 pub command_name: String,
261 pub correlation_id: Option<String>,
262 pub accepted: bool,
263 pub attempts: u32,
264 pub error: Option<String>,
265 pub completed_at: DateTime<Utc>,
266}
267
268#[derive(Clone, Debug)]
269pub struct ClusterNodeHealthRow {
270 pub snapshot: ClusterNodeHealthSnapshot,
271}
272
273impl ClusterNodeHealthRow {
274 pub fn node(&self) -> &ClusterNode {
275 &self.snapshot.node
276 }
277}