Skip to main content

stasis/application/
dto.rs

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}