1use serde::{Deserialize, Serialize};
7
8use crate::{
9 DeploymentConfig, DeploymentState, ObservedInventoryBatch, ReleaseInfo, ResourceHeartbeat,
10};
11
12#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
14#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
15#[serde(rename_all = "kebab-case")]
16pub enum OperatorCapabilityState {
17 Granted,
19 Denied,
21 Unavailable,
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
27#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
28#[serde(rename_all = "camelCase")]
29pub struct OperatorCapabilityReport {
30 pub key: String,
32 pub state: OperatorCapabilityState,
34 #[serde(default, skip_serializing_if = "Option::is_none")]
36 pub detail: Option<String>,
37}
38
39#[derive(Debug, Clone, Serialize, Deserialize)]
41#[serde(rename_all = "camelCase")]
42pub struct SyncRequest {
43 pub deployment_id: String,
45 #[serde(default)]
47 pub session: String,
48 #[serde(default)]
50 pub supports_execution_claims: bool,
51 #[serde(default, skip_serializing_if = "Option::is_none")]
53 pub execution_claim: Option<SyncExecutionClaim>,
54 #[serde(skip_serializing_if = "Option::is_none")]
56 pub current_state: Option<DeploymentState>,
57 #[serde(
59 default,
60 rename = "resourceHeartbeats",
61 skip_serializing_if = "Vec::is_empty"
62 )]
63 pub heartbeats: Vec<ResourceHeartbeat>,
64 #[serde(
66 default,
67 rename = "observedInventoryBatches",
68 skip_serializing_if = "Vec::is_empty"
69 )]
70 pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
71 #[serde(default, skip_serializing_if = "Vec::is_empty")]
73 pub capabilities: Vec<OperatorCapabilityReport>,
74 #[serde(default, skip_serializing_if = "Option::is_none")]
76 pub operator_version: Option<String>,
77}
78
79#[derive(Debug, Clone, Serialize, Deserialize)]
81#[serde(rename_all = "camelCase")]
82pub struct SyncResponse {
83 #[serde(default, skip_serializing_if = "Option::is_none")]
85 pub execution_claim: Option<SyncExecutionClaim>,
86 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub current_state: Option<DeploymentState>,
93 #[serde(skip_serializing_if = "Option::is_none")]
96 pub target: Option<TargetDeployment>,
97 #[serde(default, skip_serializing_if = "Option::is_none")]
102 pub commands_url: Option<String>,
103}
104
105#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
106#[serde(rename_all = "camelCase")]
107pub struct SyncExecutionClaim {
108 pub operation_id: String,
109 pub attempt_id: String,
110}
111
112#[derive(Debug, Clone, Serialize, Deserialize)]
114#[serde(rename_all = "camelCase")]
115pub struct TargetDeployment {
116 pub release_info: ReleaseInfo,
118 pub config: DeploymentConfig,
120}
121
122#[cfg(test)]
123mod tests {
124 use super::*;
125
126 #[test]
127 fn test_sync_request_serialization() {
128 let req = SyncRequest {
129 deployment_id: "dep_abc123".to_string(),
130 session: "operator-test".to_string(),
131 supports_execution_claims: true,
132 execution_claim: None,
133 current_state: None,
134 heartbeats: Vec::new(),
135 observed_inventory_batches: Vec::new(),
136 capabilities: Vec::new(),
137 operator_version: None,
138 };
139
140 let json = serde_json::to_value(&req).unwrap();
141 assert_eq!(json["deploymentId"], "dep_abc123");
142 assert_eq!(json["supportsExecutionClaims"], true);
143 assert!(json.get("currentState").is_none());
145 assert!(json.get("resourceHeartbeats").is_none());
146 assert!(json.get("capabilities").is_none());
147 assert!(json.get("operatorVersion").is_none());
148 }
149
150 #[test]
151 fn test_sync_request_deserialization() {
152 let json = r#"{"deploymentId": "dep_xyz"}"#;
153 let req: SyncRequest = serde_json::from_str(json).unwrap();
154 assert_eq!(req.deployment_id, "dep_xyz");
155 assert!(req.session.is_empty());
156 assert!(!req.supports_execution_claims);
157 assert!(req.current_state.is_none());
158 assert!(req.heartbeats.is_empty());
159 assert!(req.observed_inventory_batches.is_empty());
160 assert!(req.capabilities.is_empty());
161 assert!(req.operator_version.is_none());
162 }
163
164 #[test]
165 fn test_sync_response_empty() {
166 let resp = SyncResponse {
167 execution_claim: None,
168 current_state: None,
169 target: None,
170 commands_url: None,
171 };
172 let json = serde_json::to_value(&resp).unwrap();
173 assert!(json.get("target").is_none());
175 assert!(json.get("currentState").is_none());
176 }
177
178 #[test]
179 fn test_sync_response_roundtrip_no_target() {
180 let resp = SyncResponse {
181 execution_claim: None,
182 current_state: None,
183 target: None,
184 commands_url: None,
185 };
186 let serialized = serde_json::to_string(&resp).unwrap();
187 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
188 assert!(deserialized.target.is_none());
189 assert!(deserialized.current_state.is_none());
190 }
191
192 #[test]
193 fn test_sync_request_with_camel_case() {
194 let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
196 let req: SyncRequest = serde_json::from_str(json).unwrap();
197 assert_eq!(req.deployment_id, "dep_1");
198 assert!(req.current_state.is_none());
199 assert!(req.heartbeats.is_empty());
200 assert!(req.capabilities.is_empty());
201
202 let json = r#"{"deployment_id": "dep_1"}"#;
204 assert!(serde_json::from_str::<SyncRequest>(json).is_err());
205 }
206
207 #[test]
208 fn test_sync_request_heartbeats_roundtrip() {
209 let json = serde_json::json!({
210 "deploymentId": "dep_1",
211 "resourceHeartbeats": [{
212 "deploymentId": "dep_1",
213 "resourceId": "api",
214 "resourceType": "container",
215 "controllerPlatform": "kubernetes",
216 "backend": "kubernetes",
217 "observedAt": "2026-01-01T00:00:00Z",
218 "data": {
219 "resourceType": "container",
220 "data": {
221 "backend": "kubernetes",
222 "status": {
223 "health": "healthy",
224 "lifecycle": "running",
225 "message": null,
226 "stale": false,
227 "partial": false,
228 "collectionIssues": []
229 },
230 "namespace": "default",
231 "name": "api",
232 "workloadKind": "deployment",
233 "replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
234 "restarts": 0,
235 "cpu": null,
236 "memory": null,
237 "workload": null,
238 "pods": [],
239 "instances": [],
240 "events": []
241 }
242 },
243 "raw": []
244 }]
245 });
246
247 let req: SyncRequest = serde_json::from_value(json).unwrap();
248 assert_eq!(req.heartbeats.len(), 1);
249 assert_eq!(req.heartbeats[0].resource_id, "api");
250 assert!(req.capabilities.is_empty());
251
252 let serialized = serde_json::to_value(&req).unwrap();
253 assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
254 }
255
256 #[test]
257 fn test_sync_response_observe_only_state_roundtrip() {
258 let state = DeploymentState {
259 status: crate::DeploymentStatus::Running,
260 platform: crate::Platform::Kubernetes,
261 current_release: None,
262 target_release: None,
263 stack_state: None,
264 error: None,
265 environment_info: None,
266 runtime_metadata: None,
267 retry_requested: false,
268 protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
269 };
270 assert!(!state.has_desired());
271
272 let resp = SyncResponse {
273 execution_claim: None,
274 current_state: Some(state),
275 target: None,
276 commands_url: None,
277 };
278
279 let serialized = serde_json::to_string(&resp).unwrap();
280 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
281 let current_state = deserialized.current_state.unwrap();
282
283 assert_eq!(current_state.status, crate::DeploymentStatus::Running);
284 assert!(!current_state.has_desired());
285 assert!(deserialized.target.is_none());
286 }
287}