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, Copy, Serialize, Deserialize, PartialEq, Eq)]
71#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
72#[serde(rename_all = "kebab-case")]
73pub enum OperatorImageSource {
74 Package,
76 Configured,
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
82#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
83#[serde(rename_all = "camelCase")]
84pub struct OperatorImageReport {
85 pub source: OperatorImageSource,
87 pub package_id: Option<String>,
89 pub package_version: Option<String>,
91 pub image: String,
93 pub digest: String,
95}
96
97impl OperatorImageReport {
98 pub fn validate(&self) -> Result<(), &'static str> {
101 let Some(encoded_digest) = self.digest.strip_prefix("sha256:") else {
102 return Err("operator image digest must start with sha256:");
103 };
104 if encoded_digest.len() != 64
105 || !encoded_digest
106 .bytes()
107 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
108 {
109 return Err("operator image digest must contain 64 lowercase hex characters");
110 }
111
112 let Some((repository, image_digest)) = self.image.rsplit_once('@') else {
113 return Err("operator image must use repository@sha256:digest form");
114 };
115 if repository.is_empty()
116 || repository.chars().any(char::is_whitespace)
117 || repository.contains('@')
118 || image_digest != self.digest
119 {
120 return Err("operator image must contain one repository and the reported digest");
121 }
122
123 let package_id = self.package_id.as_deref().map(str::trim);
124 let package_version = self.package_version.as_deref().map(str::trim);
125 match self.source {
126 OperatorImageSource::Package
127 if package_id.is_some_and(|value| !value.is_empty())
128 && package_version.is_some_and(|value| !value.is_empty()) =>
129 {
130 Ok(())
131 }
132 OperatorImageSource::Package => {
133 Err("package operator images require package ID and version")
134 }
135 OperatorImageSource::Configured
136 if self.package_id.is_none() && self.package_version.is_none() =>
137 {
138 Ok(())
139 }
140 OperatorImageSource::Configured => {
141 Err("configured operator images must not include package identity")
142 }
143 }
144 }
145}
146
147#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
152#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
153#[serde(rename_all = "camelCase")]
154pub struct OperationsBundleDownload {
155 pub plugin: String,
157 pub plugin_version: String,
159 pub url: String,
161}
162
163#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
168#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
169#[serde(rename_all = "camelCase")]
170pub struct TargetOperationsBundleSet {
171 pub hash: String,
174 #[serde(default, skip_serializing_if = "Vec::is_empty")]
178 pub bundles: Vec<OperationsBundleDownload>,
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize)]
183#[serde(rename_all = "camelCase")]
184pub struct SyncRequest {
185 pub deployment_id: String,
187 #[serde(default)]
189 pub session: String,
190 #[serde(default)]
192 pub supports_execution_claims: bool,
193 #[serde(default, skip_serializing_if = "Option::is_none")]
195 pub execution_claim: Option<SyncExecutionClaim>,
196 #[serde(skip_serializing_if = "Option::is_none")]
198 pub current_state: Option<DeploymentState>,
199 #[serde(
201 default,
202 rename = "resourceHeartbeats",
203 skip_serializing_if = "Vec::is_empty"
204 )]
205 pub heartbeats: Vec<ResourceHeartbeat>,
206 #[serde(
208 default,
209 rename = "observedInventoryBatches",
210 skip_serializing_if = "Vec::is_empty"
211 )]
212 pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
213 #[serde(default, skip_serializing_if = "Vec::is_empty")]
215 pub capabilities: Vec<OperatorCapabilityReport>,
216 #[serde(default, skip_serializing_if = "Option::is_none")]
218 pub operator_version: Option<String>,
219 #[serde(default, skip_serializing_if = "Option::is_none")]
223 pub operations_report: Option<OperationsReport>,
224}
225
226#[derive(Debug, Clone, Serialize)]
229#[serde(rename_all = "camelCase")]
230pub struct SyncInput {
231 #[serde(flatten)]
232 request: SyncRequest,
233 #[serde(skip_serializing_if = "Option::is_none")]
236 operator_image: Option<OperatorImageReport>,
237}
238
239impl SyncInput {
240 pub fn builder(request: SyncRequest) -> SyncInputBuilder {
243 SyncInputBuilder {
244 request,
245 operator_image: None,
246 }
247 }
248}
249
250pub struct SyncInputBuilder {
253 request: SyncRequest,
254 operator_image: Option<OperatorImageReport>,
255}
256
257impl SyncInputBuilder {
258 pub fn operator_image(mut self, operator_image: OperatorImageReport) -> Self {
260 self.operator_image = Some(operator_image);
261 self
262 }
263
264 pub fn build(self) -> SyncInput {
266 SyncInput {
267 request: self.request,
268 operator_image: self.operator_image,
269 }
270 }
271}
272
273#[derive(Debug, Clone, Serialize, Deserialize)]
275#[serde(rename_all = "camelCase")]
276pub struct SyncResponse {
277 #[serde(default, skip_serializing_if = "Option::is_none")]
279 pub execution_claim: Option<SyncExecutionClaim>,
280 #[serde(default, skip_serializing_if = "Option::is_none")]
286 pub current_state: Option<DeploymentState>,
287 #[serde(skip_serializing_if = "Option::is_none")]
290 pub target: Option<TargetDeployment>,
291 #[serde(default, skip_serializing_if = "Option::is_none")]
296 pub commands_url: Option<String>,
297 #[serde(default, skip_serializing_if = "Option::is_none")]
303 pub target_operations_bundle_set: Option<TargetOperationsBundleSet>,
304}
305
306#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
307#[serde(rename_all = "camelCase")]
308pub struct SyncExecutionClaim {
309 pub operation_id: String,
310 pub attempt_id: String,
311}
312
313#[derive(Debug, Clone, Serialize, Deserialize)]
315#[serde(rename_all = "camelCase")]
316pub struct TargetDeployment {
317 pub release_info: ReleaseInfo,
319 pub config: DeploymentConfig,
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326
327 #[test]
328 fn test_sync_request_serialization() {
329 let req = SyncRequest {
330 deployment_id: "dep_abc123".to_string(),
331 session: "operator-test".to_string(),
332 supports_execution_claims: true,
333 execution_claim: None,
334 current_state: None,
335 heartbeats: Vec::new(),
336 observed_inventory_batches: Vec::new(),
337 capabilities: Vec::new(),
338 operator_version: None,
339 operations_report: None,
340 };
341
342 let json = serde_json::to_value(&req).unwrap();
343 assert_eq!(json["deploymentId"], "dep_abc123");
344 assert_eq!(json["supportsExecutionClaims"], true);
345 assert!(json.get("currentState").is_none());
347 assert!(json.get("resourceHeartbeats").is_none());
348 assert!(json.get("capabilities").is_none());
349 assert!(json.get("operatorVersion").is_none());
350 assert!(json.get("operatorImage").is_none());
351 assert!(json.get("operationsReport").is_none());
352 }
353
354 #[test]
355 fn test_sync_request_deserialization() {
356 let json = r#"{"deploymentId": "dep_xyz"}"#;
357 let req: SyncRequest = serde_json::from_str(json).unwrap();
358 assert_eq!(req.deployment_id, "dep_xyz");
359 assert!(req.session.is_empty());
360 assert!(!req.supports_execution_claims);
361 assert!(req.current_state.is_none());
362 assert!(req.heartbeats.is_empty());
363 assert!(req.observed_inventory_batches.is_empty());
364 assert!(req.capabilities.is_empty());
365 assert!(req.operator_version.is_none());
366 assert!(req.operations_report.is_none());
367 }
368
369 #[test]
370 fn test_sync_response_empty() {
371 let resp = SyncResponse {
372 execution_claim: None,
373 current_state: None,
374 target: None,
375 commands_url: None,
376 target_operations_bundle_set: None,
377 };
378 let json = serde_json::to_value(&resp).unwrap();
379 assert!(json.get("target").is_none());
381 assert!(json.get("currentState").is_none());
382 assert!(json.get("targetOperationsBundleSet").is_none());
383 }
384
385 #[test]
386 fn test_sync_response_roundtrip_no_target() {
387 let resp = SyncResponse {
388 execution_claim: None,
389 current_state: None,
390 target: None,
391 commands_url: None,
392 target_operations_bundle_set: None,
393 };
394 let serialized = serde_json::to_string(&resp).unwrap();
395 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
396 assert!(deserialized.target.is_none());
397 assert!(deserialized.current_state.is_none());
398 }
399
400 #[test]
401 fn test_sync_request_with_camel_case() {
402 let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
404 let req: SyncRequest = serde_json::from_str(json).unwrap();
405 assert_eq!(req.deployment_id, "dep_1");
406 assert!(req.current_state.is_none());
407 assert!(req.heartbeats.is_empty());
408 assert!(req.capabilities.is_empty());
409
410 let json = r#"{"deployment_id": "dep_1"}"#;
412 assert!(serde_json::from_str::<SyncRequest>(json).is_err());
413 }
414
415 #[test]
416 fn test_sync_request_heartbeats_roundtrip() {
417 let json = serde_json::json!({
418 "deploymentId": "dep_1",
419 "resourceHeartbeats": [{
420 "deploymentId": "dep_1",
421 "resourceId": "api",
422 "resourceType": "container",
423 "controllerPlatform": "kubernetes",
424 "backend": "kubernetes",
425 "observedAt": "2026-01-01T00:00:00Z",
426 "data": {
427 "resourceType": "container",
428 "data": {
429 "backend": "kubernetes",
430 "status": {
431 "health": "healthy",
432 "lifecycle": "running",
433 "message": null,
434 "stale": false,
435 "partial": false,
436 "collectionIssues": []
437 },
438 "namespace": "default",
439 "name": "api",
440 "workloadKind": "deployment",
441 "replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
442 "restarts": 0,
443 "cpu": null,
444 "memory": null,
445 "workload": null,
446 "pods": [],
447 "instances": [],
448 "events": []
449 }
450 },
451 "raw": []
452 }]
453 });
454
455 let req: SyncRequest = serde_json::from_value(json).unwrap();
456 assert_eq!(req.heartbeats.len(), 1);
457 assert_eq!(req.heartbeats[0].resource_id, "api");
458 assert!(req.capabilities.is_empty());
459
460 let serialized = serde_json::to_value(&req).unwrap();
461 assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
462 }
463
464 #[test]
465 fn test_sync_response_observe_only_state_roundtrip() {
466 let state = DeploymentState {
467 status: crate::DeploymentStatus::Running,
468 platform: crate::Platform::Kubernetes,
469 current_release: None,
470 target_release: None,
471 stack_state: None,
472 error: None,
473 environment_info: None,
474 runtime_metadata: None,
475 retry_requested: false,
476 protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
477 };
478 assert!(!state.has_desired());
479
480 let resp = SyncResponse {
481 execution_claim: None,
482 current_state: Some(state),
483 target: None,
484 commands_url: None,
485 target_operations_bundle_set: None,
486 };
487
488 let serialized = serde_json::to_string(&resp).unwrap();
489 let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
490 let current_state = deserialized.current_state.unwrap();
491
492 assert_eq!(current_state.status, crate::DeploymentStatus::Running);
493 assert!(!current_state.has_desired());
494 assert!(deserialized.target.is_none());
495 }
496
497 #[test]
498 fn test_sync_request_operations_report_absent_by_default() {
499 let json = r#"{"deploymentId": "dep_old_operator", "operatorVersion": "0.9.0"}"#;
502 let req: SyncRequest = serde_json::from_str(json).unwrap();
503 assert_eq!(req.operator_version.as_deref(), Some("0.9.0"));
504 assert!(req.operations_report.is_none());
505 }
506
507 #[test]
508 fn test_sync_request_operator_image_roundtrip() {
509 let digest = format!("sha256:{}", "a".repeat(64));
510 let report = OperatorImageReport {
511 source: OperatorImageSource::Package,
512 package_id: Some("pkg_operator".to_string()),
513 package_version: Some("1.2.3".to_string()),
514 image: format!("registry.example.com/operator@{digest}"),
515 digest,
516 };
517 report.validate().expect("valid package image report");
518
519 let value = serde_json::to_value(&report).unwrap();
520 assert_eq!(value["source"], "package");
521 assert_eq!(value["packageId"], "pkg_operator");
522 assert_eq!(value["packageVersion"], "1.2.3");
523
524 let decoded: OperatorImageReport = serde_json::from_value(value).unwrap();
525 assert_eq!(decoded, report);
526 }
527
528 #[test]
529 fn sync_input_adds_operator_image_without_expanding_sync_request() {
530 let digest = format!("sha256:{}", "a".repeat(64));
531 let request = SyncRequest {
532 deployment_id: "dep_1".to_string(),
533 session: String::new(),
534 supports_execution_claims: false,
535 execution_claim: None,
536 current_state: None,
537 heartbeats: Vec::new(),
538 observed_inventory_batches: Vec::new(),
539 capabilities: Vec::new(),
540 operator_version: None,
541 operations_report: None,
542 };
543 let input = SyncInput::builder(request)
544 .operator_image(OperatorImageReport {
545 source: OperatorImageSource::Configured,
546 package_id: None,
547 package_version: None,
548 image: format!("registry.example.com/operator@{digest}"),
549 digest,
550 })
551 .build();
552
553 let value = serde_json::to_value(input).unwrap();
554 assert_eq!(value["deploymentId"], "dep_1");
555 assert_eq!(value["operatorImage"]["source"], "configured");
556 }
557
558 #[test]
559 fn sync_input_omits_absent_operator_image() {
560 let request = SyncRequest {
561 deployment_id: "dep_1".to_string(),
562 session: String::new(),
563 supports_execution_claims: false,
564 execution_claim: None,
565 current_state: None,
566 heartbeats: Vec::new(),
567 observed_inventory_batches: Vec::new(),
568 capabilities: Vec::new(),
569 operator_version: None,
570 operations_report: None,
571 };
572
573 let value = serde_json::to_value(SyncInput::builder(request).build()).unwrap();
574 assert!(value.get("operatorImage").is_none());
575 }
576
577 #[test]
578 fn operator_image_report_rejects_mutable_or_inconsistent_identity() {
579 let digest = format!("sha256:{}", "a".repeat(64));
580 let mutable = OperatorImageReport {
581 source: OperatorImageSource::Configured,
582 package_id: None,
583 package_version: None,
584 image: "registry.example.com/operator:latest".to_string(),
585 digest: digest.clone(),
586 };
587 assert_eq!(
588 mutable.validate(),
589 Err("operator image must use repository@sha256:digest form")
590 );
591
592 let mismatched = OperatorImageReport {
593 source: OperatorImageSource::Configured,
594 package_id: None,
595 package_version: None,
596 image: format!("registry.example.com/operator@sha256:{}", "b".repeat(64)),
597 digest,
598 };
599 assert_eq!(
600 mismatched.validate(),
601 Err("operator image must contain one repository and the reported digest")
602 );
603 }
604
605 #[test]
606 fn operator_image_report_enforces_source_specific_package_fields() {
607 let digest = format!("sha256:{}", "a".repeat(64));
608 let image = format!("registry.example.com/operator@{digest}");
609 let missing_package = OperatorImageReport {
610 source: OperatorImageSource::Package,
611 package_id: None,
612 package_version: None,
613 image: image.clone(),
614 digest: digest.clone(),
615 };
616 assert_eq!(
617 missing_package.validate(),
618 Err("package operator images require package ID and version")
619 );
620
621 let configured_with_package = OperatorImageReport {
622 source: OperatorImageSource::Configured,
623 package_id: Some("pkg_operator".to_string()),
624 package_version: Some("1.2.3".to_string()),
625 image,
626 digest,
627 };
628 assert_eq!(
629 configured_with_package.validate(),
630 Err("configured operator images must not include package identity")
631 );
632 }
633
634 #[test]
635 fn test_sync_request_operations_report_roundtrip() {
636 let req = SyncRequest {
637 deployment_id: "dep_1".to_string(),
638 session: String::new(),
639 supports_execution_claims: false,
640 execution_claim: None,
641 current_state: None,
642 heartbeats: Vec::new(),
643 observed_inventory_batches: Vec::new(),
644 capabilities: Vec::new(),
645 operator_version: Some("1.2.3".to_string()),
646 operations_report: Some(OperationsReport {
647 loaded_bundle_hash: Some("builtin:s3@1.0.0:key|".to_string()),
648 operations: vec![ReportedOperation {
649 plugin: "s3".to_string(),
650 plugin_version: "1.0.0".to_string(),
651 name: "list-buckets".to_string(),
652 }],
653 }),
654 };
655
656 let json = serde_json::to_value(&req).unwrap();
657 assert_eq!(
658 json["operationsReport"]["loadedBundleHash"],
659 "builtin:s3@1.0.0:key|"
660 );
661 assert_eq!(json["operationsReport"]["operations"][0]["plugin"], "s3");
662 assert_eq!(
663 json["operationsReport"]["operations"][0]["pluginVersion"],
664 "1.0.0"
665 );
666 assert_eq!(
667 json["operationsReport"]["operations"][0]["name"],
668 "list-buckets"
669 );
670
671 let deserialized: SyncRequest = serde_json::from_value(json).unwrap();
672 assert_eq!(
673 deserialized
674 .operations_report
675 .as_ref()
676 .unwrap()
677 .loaded_bundle_hash,
678 req.operations_report.as_ref().unwrap().loaded_bundle_hash
679 );
680 assert_eq!(
681 deserialized.operations_report.unwrap().operations,
682 req.operations_report.unwrap().operations
683 );
684 }
685
686 #[test]
687 fn test_sync_response_target_operations_bundle_set_roundtrip() {
688 let resp = SyncResponse {
689 execution_claim: None,
690 current_state: None,
691 target: None,
692 commands_url: None,
693 target_operations_bundle_set: Some(TargetOperationsBundleSet {
694 hash: "builtin:s3@1.0.0:key|".to_string(),
695 bundles: vec![OperationsBundleDownload {
696 plugin: "s3".to_string(),
697 plugin_version: "1.0.0".to_string(),
698 url: "https://storage.example.com/bundle.zip?sig=abc".to_string(),
699 }],
700 }),
701 };
702
703 let json = serde_json::to_value(&resp).unwrap();
704 assert_eq!(
705 json["targetOperationsBundleSet"]["hash"],
706 "builtin:s3@1.0.0:key|"
707 );
708 assert_eq!(
709 json["targetOperationsBundleSet"]["bundles"][0]["plugin"],
710 "s3"
711 );
712
713 let deserialized: SyncResponse = serde_json::from_value(json).unwrap();
714 assert_eq!(
715 deserialized.target_operations_bundle_set,
716 resp.target_operations_bundle_set
717 );
718 }
719
720 #[test]
721 fn test_sync_response_target_operations_bundle_set_absent_by_default() {
722 let json = serde_json::json!({});
723 let resp: SyncResponse = serde_json::from_value(json).unwrap();
724 assert!(resp.target_operations_bundle_set.is_none());
725 }
726}