use serde::{Deserialize, Serialize};
use crate::{
DeploymentConfig, DeploymentState, ObservedInventoryBatch, ReleaseInfo, ResourceHeartbeat,
};
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "kebab-case")]
pub enum OperatorCapabilityState {
Granted,
Denied,
Unavailable,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
pub struct OperatorCapabilityReport {
pub key: String,
pub state: OperatorCapabilityState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
pub struct ReportedOperation {
pub plugin: String,
pub plugin_version: String,
pub name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
pub struct OperationsReport {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub loaded_bundle_hash: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub operations: Vec<ReportedOperation>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
pub struct OperationsBundleDownload {
pub plugin: String,
pub plugin_version: String,
pub url: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
pub struct TargetOperationsBundleSet {
pub hash: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub bundles: Vec<OperationsBundleDownload>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SyncRequest {
pub deployment_id: String,
#[serde(default)]
pub session: String,
#[serde(default)]
pub supports_execution_claims: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution_claim: Option<SyncExecutionClaim>,
#[serde(skip_serializing_if = "Option::is_none")]
pub current_state: Option<DeploymentState>,
#[serde(
default,
rename = "resourceHeartbeats",
skip_serializing_if = "Vec::is_empty"
)]
pub heartbeats: Vec<ResourceHeartbeat>,
#[serde(
default,
rename = "observedInventoryBatches",
skip_serializing_if = "Vec::is_empty"
)]
pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub capabilities: Vec<OperatorCapabilityReport>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operator_version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operations_report: Option<OperationsReport>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SyncResponse {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution_claim: Option<SyncExecutionClaim>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_state: Option<DeploymentState>,
#[serde(skip_serializing_if = "Option::is_none")]
pub target: Option<TargetDeployment>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub commands_url: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target_operations_bundle_set: Option<TargetOperationsBundleSet>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct SyncExecutionClaim {
pub operation_id: String,
pub attempt_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct TargetDeployment {
pub release_info: ReleaseInfo,
pub config: DeploymentConfig,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_sync_request_serialization() {
let req = SyncRequest {
deployment_id: "dep_abc123".to_string(),
session: "operator-test".to_string(),
supports_execution_claims: true,
execution_claim: None,
current_state: None,
heartbeats: Vec::new(),
observed_inventory_batches: Vec::new(),
capabilities: Vec::new(),
operator_version: None,
operations_report: None,
};
let json = serde_json::to_value(&req).unwrap();
assert_eq!(json["deploymentId"], "dep_abc123");
assert_eq!(json["supportsExecutionClaims"], true);
assert!(json.get("currentState").is_none());
assert!(json.get("resourceHeartbeats").is_none());
assert!(json.get("capabilities").is_none());
assert!(json.get("operatorVersion").is_none());
assert!(json.get("operationsReport").is_none());
}
#[test]
fn test_sync_request_deserialization() {
let json = r#"{"deploymentId": "dep_xyz"}"#;
let req: SyncRequest = serde_json::from_str(json).unwrap();
assert_eq!(req.deployment_id, "dep_xyz");
assert!(req.session.is_empty());
assert!(!req.supports_execution_claims);
assert!(req.current_state.is_none());
assert!(req.heartbeats.is_empty());
assert!(req.observed_inventory_batches.is_empty());
assert!(req.capabilities.is_empty());
assert!(req.operator_version.is_none());
assert!(req.operations_report.is_none());
}
#[test]
fn test_sync_response_empty() {
let resp = SyncResponse {
execution_claim: None,
current_state: None,
target: None,
commands_url: None,
target_operations_bundle_set: None,
};
let json = serde_json::to_value(&resp).unwrap();
assert!(json.get("target").is_none());
assert!(json.get("currentState").is_none());
assert!(json.get("targetOperationsBundleSet").is_none());
}
#[test]
fn test_sync_response_roundtrip_no_target() {
let resp = SyncResponse {
execution_claim: None,
current_state: None,
target: None,
commands_url: None,
target_operations_bundle_set: None,
};
let serialized = serde_json::to_string(&resp).unwrap();
let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
assert!(deserialized.target.is_none());
assert!(deserialized.current_state.is_none());
}
#[test]
fn test_sync_request_with_camel_case() {
let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
let req: SyncRequest = serde_json::from_str(json).unwrap();
assert_eq!(req.deployment_id, "dep_1");
assert!(req.current_state.is_none());
assert!(req.heartbeats.is_empty());
assert!(req.capabilities.is_empty());
let json = r#"{"deployment_id": "dep_1"}"#;
assert!(serde_json::from_str::<SyncRequest>(json).is_err());
}
#[test]
fn test_sync_request_heartbeats_roundtrip() {
let json = serde_json::json!({
"deploymentId": "dep_1",
"resourceHeartbeats": [{
"deploymentId": "dep_1",
"resourceId": "api",
"resourceType": "container",
"controllerPlatform": "kubernetes",
"backend": "kubernetes",
"observedAt": "2026-01-01T00:00:00Z",
"data": {
"resourceType": "container",
"data": {
"backend": "kubernetes",
"status": {
"health": "healthy",
"lifecycle": "running",
"message": null,
"stale": false,
"partial": false,
"collectionIssues": []
},
"namespace": "default",
"name": "api",
"workloadKind": "deployment",
"replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
"restarts": 0,
"cpu": null,
"memory": null,
"workload": null,
"pods": [],
"instances": [],
"events": []
}
},
"raw": []
}]
});
let req: SyncRequest = serde_json::from_value(json).unwrap();
assert_eq!(req.heartbeats.len(), 1);
assert_eq!(req.heartbeats[0].resource_id, "api");
assert!(req.capabilities.is_empty());
let serialized = serde_json::to_value(&req).unwrap();
assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
}
#[test]
fn test_sync_response_observe_only_state_roundtrip() {
let state = DeploymentState {
status: crate::DeploymentStatus::Running,
platform: crate::Platform::Kubernetes,
current_release: None,
target_release: None,
stack_state: None,
error: None,
environment_info: None,
runtime_metadata: None,
retry_requested: false,
protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
};
assert!(!state.has_desired());
let resp = SyncResponse {
execution_claim: None,
current_state: Some(state),
target: None,
commands_url: None,
target_operations_bundle_set: None,
};
let serialized = serde_json::to_string(&resp).unwrap();
let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
let current_state = deserialized.current_state.unwrap();
assert_eq!(current_state.status, crate::DeploymentStatus::Running);
assert!(!current_state.has_desired());
assert!(deserialized.target.is_none());
}
#[test]
fn test_sync_request_operations_report_absent_by_default() {
let json = r#"{"deploymentId": "dep_old_operator", "operatorVersion": "0.9.0"}"#;
let req: SyncRequest = serde_json::from_str(json).unwrap();
assert_eq!(req.operator_version.as_deref(), Some("0.9.0"));
assert!(req.operations_report.is_none());
}
#[test]
fn test_sync_request_operations_report_roundtrip() {
let req = SyncRequest {
deployment_id: "dep_1".to_string(),
session: String::new(),
supports_execution_claims: false,
execution_claim: None,
current_state: None,
heartbeats: Vec::new(),
observed_inventory_batches: Vec::new(),
capabilities: Vec::new(),
operator_version: Some("1.2.3".to_string()),
operations_report: Some(OperationsReport {
loaded_bundle_hash: Some("builtin:s3@1.0.0:key|".to_string()),
operations: vec![ReportedOperation {
plugin: "s3".to_string(),
plugin_version: "1.0.0".to_string(),
name: "list-buckets".to_string(),
}],
}),
};
let json = serde_json::to_value(&req).unwrap();
assert_eq!(
json["operationsReport"]["loadedBundleHash"],
"builtin:s3@1.0.0:key|"
);
assert_eq!(json["operationsReport"]["operations"][0]["plugin"], "s3");
assert_eq!(
json["operationsReport"]["operations"][0]["pluginVersion"],
"1.0.0"
);
assert_eq!(
json["operationsReport"]["operations"][0]["name"],
"list-buckets"
);
let deserialized: SyncRequest = serde_json::from_value(json).unwrap();
assert_eq!(
deserialized.operations_report.as_ref().unwrap().loaded_bundle_hash,
req.operations_report.as_ref().unwrap().loaded_bundle_hash
);
assert_eq!(
deserialized.operations_report.unwrap().operations,
req.operations_report.unwrap().operations
);
}
#[test]
fn test_sync_response_target_operations_bundle_set_roundtrip() {
let resp = SyncResponse {
execution_claim: None,
current_state: None,
target: None,
commands_url: None,
target_operations_bundle_set: Some(TargetOperationsBundleSet {
hash: "builtin:s3@1.0.0:key|".to_string(),
bundles: vec![OperationsBundleDownload {
plugin: "s3".to_string(),
plugin_version: "1.0.0".to_string(),
url: "https://storage.example.com/bundle.zip?sig=abc".to_string(),
}],
}),
};
let json = serde_json::to_value(&resp).unwrap();
assert_eq!(json["targetOperationsBundleSet"]["hash"], "builtin:s3@1.0.0:key|");
assert_eq!(
json["targetOperationsBundleSet"]["bundles"][0]["plugin"],
"s3"
);
let deserialized: SyncResponse = serde_json::from_value(json).unwrap();
assert_eq!(
deserialized.target_operations_bundle_set,
resp.target_operations_bundle_set
);
}
#[test]
fn test_sync_response_target_operations_bundle_set_absent_by_default() {
let json = serde_json::json!({});
let resp: SyncResponse = serde_json::from_value(json).unwrap();
assert!(resp.target_operations_bundle_set.is_none());
}
}