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, PartialEq, Eq)]
43#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
44#[serde(rename_all = "camelCase")]
45pub struct ReportedOperation {
46 pub plugin: String,
48 pub plugin_version: String,
50 pub name: String,
52}
53
54#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
57#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
58#[serde(rename_all = "camelCase")]
59pub struct OperationsReport {
60 #[serde(default, skip_serializing_if = "Option::is_none")]
63 pub loaded_bundle_hash: Option<String>,
64 #[serde(default, skip_serializing_if = "Vec::is_empty")]
66 pub operations: Vec<ReportedOperation>,
67}
68
69#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
74#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
75#[serde(rename_all = "camelCase")]
76pub struct OperationsBundleDownload {
77 pub plugin: String,
79 pub plugin_version: String,
81 pub url: String,
83}
84
85#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
90#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
91#[serde(rename_all = "camelCase")]
92pub struct TargetOperationsBundleSet {
93 pub hash: String,
96 #[serde(default, skip_serializing_if = "Vec::is_empty")]
100 pub bundles: Vec<OperationsBundleDownload>,
101}
102
103#[derive(Debug, Clone, Serialize, Deserialize)]
105#[serde(rename_all = "camelCase")]
106pub struct SyncRequest {
107 pub deployment_id: String,
109 #[serde(default)]
111 pub session: String,
112 #[serde(default)]
114 pub supports_execution_claims: bool,
115 #[serde(default, skip_serializing_if = "Option::is_none")]
117 pub execution_claim: Option<SyncExecutionClaim>,
118 #[serde(skip_serializing_if = "Option::is_none")]
120 pub current_state: Option<DeploymentState>,
121 #[serde(
123 default,
124 rename = "resourceHeartbeats",
125 skip_serializing_if = "Vec::is_empty"
126 )]
127 pub heartbeats: Vec<ResourceHeartbeat>,
128 #[serde(
130 default,
131 rename = "observedInventoryBatches",
132 skip_serializing_if = "Vec::is_empty"
133 )]
134 pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
135 #[serde(default, skip_serializing_if = "Vec::is_empty")]
137 pub capabilities: Vec<OperatorCapabilityReport>,
138 #[serde(default, skip_serializing_if = "Option::is_none")]
140 pub operator_version: Option<String>,
141 #[serde(default, skip_serializing_if = "Option::is_none")]
145 pub operations_report: Option<OperationsReport>,
146}
147
148#[derive(Debug, Clone, Serialize, Deserialize)]
150#[serde(rename_all = "camelCase")]
151pub struct SyncResponse {
152 #[serde(default, skip_serializing_if = "Option::is_none")]
154 pub execution_claim: Option<SyncExecutionClaim>,
155 #[serde(default, skip_serializing_if = "Option::is_none")]
161 pub current_state: Option<DeploymentState>,
162 #[serde(skip_serializing_if = "Option::is_none")]
165 pub target: Option<TargetDeployment>,
166 #[serde(default, skip_serializing_if = "Option::is_none")]
171 pub commands_url: Option<String>,
172 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub target_operations_bundle_set: Option<TargetOperationsBundleSet>,
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
182#[serde(rename_all = "camelCase")]
183pub struct SyncExecutionClaim {
184 pub operation_id: String,
185 pub attempt_id: String,
186}
187
188#[derive(Debug, Clone, Serialize, Deserialize)]
190#[serde(rename_all = "camelCase")]
191pub struct TargetDeployment {
192 pub release_info: ReleaseInfo,
194 pub config: DeploymentConfig,
196}
197
198#[cfg(test)]
199mod tests {
200 use super::*;
201
202 #[test]
203 fn test_sync_request_serialization() {
204 let req = SyncRequest {
205 deployment_id: "dep_abc123".to_string(),
206 session: "operator-test".to_string(),
207 supports_execution_claims: true,
208 execution_claim: None,
209 current_state: None,
210 heartbeats: Vec::new(),
211 observed_inventory_batches: Vec::new(),
212 capabilities: Vec::new(),
213 operator_version: None,
214 operations_report: None,
215 };
216
217 let json = serde_json::to_value(&req).unwrap();
218 assert_eq!(json["deploymentId"], "dep_abc123");
219 assert_eq!(json["supportsExecutionClaims"], true);
220 assert!(json.get("currentState").is_none());
222 assert!(json.get("resourceHeartbeats").is_none());
223 assert!(json.get("capabilities").is_none());
224 assert!(json.get("operatorVersion").is_none());
225 assert!(json.get("operationsReport").is_none());
226 }
227
228 #[test]
229 fn test_sync_request_deserialization() {
230 let json = r#"{"deploymentId": "dep_xyz"}"#;
231 let req: SyncRequest = serde_json::from_str(json).unwrap();
232 assert_eq!(req.deployment_id, "dep_xyz");
233 assert!(req.session.is_empty());
234 assert!(!req.supports_execution_claims);
235 assert!(req.current_state.is_none());
236 assert!(req.heartbeats.is_empty());
237 assert!(req.observed_inventory_batches.is_empty());
238 assert!(req.capabilities.is_empty());
239 assert!(req.operator_version.is_none());
240 assert!(req.operations_report.is_none());
241 }
242
243 #[test]
244 fn test_sync_response_empty() {
245 let resp = SyncResponse {
246 execution_claim: None,
247 current_state: None,
248 target: None,
249 commands_url: None,
250 target_operations_bundle_set: None,
251 };
252 let json = serde_json::to_value(&resp).unwrap();
253 assert!(json.get("target").is_none());
255 assert!(json.get("currentState").is_none());
256 assert!(json.get("targetOperationsBundleSet").is_none());
257 }
258
259 #[test]
260 fn test_sync_response_roundtrip_no_target() {
261 let resp = SyncResponse {
262 execution_claim: None,
263 current_state: None,
264 target: None,
265 commands_url: None,
266 target_operations_bundle_set: None,
267 };
268 let serialized = serde_json::to_string(&resp).unwrap();
269 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
270 assert!(deserialized.target.is_none());
271 assert!(deserialized.current_state.is_none());
272 }
273
274 #[test]
275 fn test_sync_request_with_camel_case() {
276 let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
278 let req: SyncRequest = serde_json::from_str(json).unwrap();
279 assert_eq!(req.deployment_id, "dep_1");
280 assert!(req.current_state.is_none());
281 assert!(req.heartbeats.is_empty());
282 assert!(req.capabilities.is_empty());
283
284 let json = r#"{"deployment_id": "dep_1"}"#;
286 assert!(serde_json::from_str::<SyncRequest>(json).is_err());
287 }
288
289 #[test]
290 fn test_sync_request_heartbeats_roundtrip() {
291 let json = serde_json::json!({
292 "deploymentId": "dep_1",
293 "resourceHeartbeats": [{
294 "deploymentId": "dep_1",
295 "resourceId": "api",
296 "resourceType": "container",
297 "controllerPlatform": "kubernetes",
298 "backend": "kubernetes",
299 "observedAt": "2026-01-01T00:00:00Z",
300 "data": {
301 "resourceType": "container",
302 "data": {
303 "backend": "kubernetes",
304 "status": {
305 "health": "healthy",
306 "lifecycle": "running",
307 "message": null,
308 "stale": false,
309 "partial": false,
310 "collectionIssues": []
311 },
312 "namespace": "default",
313 "name": "api",
314 "workloadKind": "deployment",
315 "replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
316 "restarts": 0,
317 "cpu": null,
318 "memory": null,
319 "workload": null,
320 "pods": [],
321 "instances": [],
322 "events": []
323 }
324 },
325 "raw": []
326 }]
327 });
328
329 let req: SyncRequest = serde_json::from_value(json).unwrap();
330 assert_eq!(req.heartbeats.len(), 1);
331 assert_eq!(req.heartbeats[0].resource_id, "api");
332 assert!(req.capabilities.is_empty());
333
334 let serialized = serde_json::to_value(&req).unwrap();
335 assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
336 }
337
338 #[test]
339 fn test_sync_response_observe_only_state_roundtrip() {
340 let state = DeploymentState {
341 status: crate::DeploymentStatus::Running,
342 platform: crate::Platform::Kubernetes,
343 current_release: None,
344 target_release: None,
345 stack_state: None,
346 error: None,
347 environment_info: None,
348 runtime_metadata: None,
349 retry_requested: false,
350 protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
351 };
352 assert!(!state.has_desired());
353
354 let resp = SyncResponse {
355 execution_claim: None,
356 current_state: Some(state),
357 target: None,
358 commands_url: None,
359 target_operations_bundle_set: None,
360 };
361
362 let serialized = serde_json::to_string(&resp).unwrap();
363 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
364 let current_state = deserialized.current_state.unwrap();
365
366 assert_eq!(current_state.status, crate::DeploymentStatus::Running);
367 assert!(!current_state.has_desired());
368 assert!(deserialized.target.is_none());
369 }
370
371 #[test]
372 fn test_sync_request_operations_report_absent_by_default() {
373 let json = r#"{"deploymentId": "dep_old_operator", "operatorVersion": "0.9.0"}"#;
376 let req: SyncRequest = serde_json::from_str(json).unwrap();
377 assert_eq!(req.operator_version.as_deref(), Some("0.9.0"));
378 assert!(req.operations_report.is_none());
379 }
380
381 #[test]
382 fn test_sync_request_operations_report_roundtrip() {
383 let req = SyncRequest {
384 deployment_id: "dep_1".to_string(),
385 session: String::new(),
386 supports_execution_claims: false,
387 execution_claim: None,
388 current_state: None,
389 heartbeats: Vec::new(),
390 observed_inventory_batches: Vec::new(),
391 capabilities: Vec::new(),
392 operator_version: Some("1.2.3".to_string()),
393 operations_report: Some(OperationsReport {
394 loaded_bundle_hash: Some("builtin:s3@1.0.0:key|".to_string()),
395 operations: vec![ReportedOperation {
396 plugin: "s3".to_string(),
397 plugin_version: "1.0.0".to_string(),
398 name: "list-buckets".to_string(),
399 }],
400 }),
401 };
402
403 let json = serde_json::to_value(&req).unwrap();
404 assert_eq!(
405 json["operationsReport"]["loadedBundleHash"],
406 "builtin:s3@1.0.0:key|"
407 );
408 assert_eq!(json["operationsReport"]["operations"][0]["plugin"], "s3");
409 assert_eq!(
410 json["operationsReport"]["operations"][0]["pluginVersion"],
411 "1.0.0"
412 );
413 assert_eq!(
414 json["operationsReport"]["operations"][0]["name"],
415 "list-buckets"
416 );
417
418 let deserialized: SyncRequest = serde_json::from_value(json).unwrap();
419 assert_eq!(
420 deserialized.operations_report.as_ref().unwrap().loaded_bundle_hash,
421 req.operations_report.as_ref().unwrap().loaded_bundle_hash
422 );
423 assert_eq!(
424 deserialized.operations_report.unwrap().operations,
425 req.operations_report.unwrap().operations
426 );
427 }
428
429 #[test]
430 fn test_sync_response_target_operations_bundle_set_roundtrip() {
431 let resp = SyncResponse {
432 execution_claim: None,
433 current_state: None,
434 target: None,
435 commands_url: None,
436 target_operations_bundle_set: Some(TargetOperationsBundleSet {
437 hash: "builtin:s3@1.0.0:key|".to_string(),
438 bundles: vec![OperationsBundleDownload {
439 plugin: "s3".to_string(),
440 plugin_version: "1.0.0".to_string(),
441 url: "https://storage.example.com/bundle.zip?sig=abc".to_string(),
442 }],
443 }),
444 };
445
446 let json = serde_json::to_value(&resp).unwrap();
447 assert_eq!(json["targetOperationsBundleSet"]["hash"], "builtin:s3@1.0.0:key|");
448 assert_eq!(
449 json["targetOperationsBundleSet"]["bundles"][0]["plugin"],
450 "s3"
451 );
452
453 let deserialized: SyncResponse = serde_json::from_value(json).unwrap();
454 assert_eq!(
455 deserialized.target_operations_bundle_set,
456 resp.target_operations_bundle_set
457 );
458 }
459
460 #[test]
461 fn test_sync_response_target_operations_bundle_set_absent_by_default() {
462 let json = serde_json::json!({});
463 let resp: SyncResponse = serde_json::from_value(json).unwrap();
464 assert!(resp.target_operations_bundle_set.is_none());
465 }
466}