Skip to main content

alien_core/
sync.rs

1//! Sync protocol types for agent ↔ manager communication.
2//!
3//! The agent periodically calls `POST /v1/sync` with a `SyncRequest` and
4//! receives a `SyncResponse` containing the target deployment state.
5
6use serde::{Deserialize, Serialize};
7
8use crate::{
9    DeploymentConfig, DeploymentState, ObservedInventoryBatch, ReleaseInfo, ResourceHeartbeat,
10};
11
12/// State of an Operator capability as observed inside the environment.
13#[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    /// The Operator has the permission or local facility needed for the capability.
18    Granted,
19    /// The environment explicitly denied the capability.
20    Denied,
21    /// The capability does not apply in this environment.
22    Unavailable,
23}
24
25/// Report-only Operator capability status.
26#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
27#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
28#[serde(rename_all = "camelCase")]
29pub struct OperatorCapabilityReport {
30    /// Stable capability key, such as `k8s-workloads` or `logs`.
31    pub key: String,
32    /// Whether the capability is currently usable.
33    pub state: OperatorCapabilityState,
34    /// Optional human-readable detail from the Operator.
35    #[serde(default, skip_serializing_if = "Option::is_none")]
36    pub detail: Option<String>,
37}
38
39/// Request sent by the agent to the manager during periodic sync.
40#[derive(Debug, Clone, Serialize, Deserialize)]
41#[serde(rename_all = "camelCase")]
42pub struct SyncRequest {
43    /// The deployment ID this agent is managing.
44    pub deployment_id: String,
45    /// Stable identity for this Operator process. Used to fence update work.
46    #[serde(default)]
47    pub session: String,
48    /// Signals that this Operator persists and echoes execution claims.
49    #[serde(default)]
50    pub supports_execution_claims: bool,
51    /// Exact update claim returned by the previous sync response.
52    #[serde(default, skip_serializing_if = "Option::is_none")]
53    pub execution_claim: Option<SyncExecutionClaim>,
54    /// Current deployment state as seen by the agent.
55    #[serde(skip_serializing_if = "Option::is_none")]
56    pub current_state: Option<DeploymentState>,
57    /// Managed Alien resource status samples emitted by the Operator's deployment step.
58    #[serde(
59        default,
60        rename = "resourceHeartbeats",
61        skip_serializing_if = "Vec::is_empty"
62    )]
63    pub heartbeats: Vec<ResourceHeartbeat>,
64    /// Observed raw-resource inventory batches successfully read by the Operator.
65    #[serde(
66        default,
67        rename = "observedInventoryBatches",
68        skip_serializing_if = "Vec::is_empty"
69    )]
70    pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
71    /// Report-only capabilities observed by the Operator.
72    #[serde(default, skip_serializing_if = "Vec::is_empty")]
73    pub capabilities: Vec<OperatorCapabilityReport>,
74    /// Version of the Operator binary reporting this sync.
75    #[serde(default, skip_serializing_if = "Option::is_none")]
76    pub operator_version: Option<String>,
77}
78
79/// Response from the manager to the agent sync request.
80#[derive(Debug, Clone, Serialize, Deserialize)]
81#[serde(rename_all = "camelCase")]
82pub struct SyncResponse {
83    /// Exact operation attempt associated with `target`.
84    #[serde(default, skip_serializing_if = "Option::is_none")]
85    pub execution_claim: Option<SyncExecutionClaim>,
86    /// Authoritative deployment state from the manager.
87    ///
88    /// Pull agents use this to hydrate local state when attaching to an
89    /// already-imported deployment. Absent means the agent's local state is
90    /// already authoritative or no state has been established yet.
91    #[serde(default, skip_serializing_if = "Option::is_none")]
92    pub current_state: Option<DeploymentState>,
93    /// Target deployment the agent should converge toward.
94    /// None means no changes needed or this is an observe-only deployment.
95    #[serde(skip_serializing_if = "Option::is_none")]
96    pub target: Option<TargetDeployment>,
97    /// Public URL for the commands API (e.g. `https://manager.example.com/v1`).
98    /// Operators and app-owned receivers use this to lease pending commands;
99    /// Workers themselves receive pushes.
100    /// When absent, the agent falls back to its sync URL.
101    #[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/// Target deployment state for the agent to converge toward.
113#[derive(Debug, Clone, Serialize, Deserialize)]
114#[serde(rename_all = "camelCase")]
115pub struct TargetDeployment {
116    /// Release information (ID, version, stack definition).
117    pub release_info: ReleaseInfo,
118    /// Full deployment configuration (settings, env vars, etc.).
119    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        // current_state is None → should be omitted
144        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        // target is None → should be omitted
174        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        // Verify camelCase renaming works correctly
195        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        // snake_case should NOT work
203        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}