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/// A single operation the Operator currently has loaded. Opaque identifiers
40/// only — no tier, description, or other plugin business logic crosses into
41/// this public crate.
42#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
43#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
44#[serde(rename_all = "camelCase")]
45pub struct ReportedOperation {
46    /// Name of the plugin that owns this operation.
47    pub plugin: String,
48    /// Version of the plugin that owns this operation.
49    pub plugin_version: String,
50    /// Name of the operation within the plugin.
51    pub name: String,
52}
53
54/// Report-only summary of the operations the Operator has loaded, used to
55/// confirm a bundle sync actually took effect.
56#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
57#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
58#[serde(rename_all = "camelCase")]
59pub struct OperationsReport {
60    /// Hash of the enabled-plugin bundle set the Operator currently has
61    /// loaded. Compared against the platform's target hash to detect drift.
62    #[serde(default, skip_serializing_if = "Option::is_none")]
63    pub loaded_bundle_hash: Option<String>,
64    /// Operations currently loaded and executable by the Operator.
65    #[serde(default, skip_serializing_if = "Vec::is_empty")]
66    pub operations: Vec<ReportedOperation>,
67}
68
69/// One bundle the Operator needs to download to reach `targetBundleHash`.
70/// The manager mints a short-lived presigned GET URL per bundle — the
71/// Operator never holds real cloud storage credentials, mirroring the OCI
72/// registry proxy's credential-injection pattern.
73#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
74#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
75#[serde(rename_all = "camelCase")]
76pub struct OperationsBundleDownload {
77    /// Plugin name this bundle provides.
78    pub plugin: String,
79    /// Plugin version this bundle provides.
80    pub plugin_version: String,
81    /// Presigned URL to GET the bundle ZIP from. Short-lived.
82    pub url: String,
83}
84
85/// Target operations-bundle set for the Operator to converge its loaded
86/// plugin registry toward, independent of any release/config target — a
87/// plugin can be enabled with no release change, so this is not nested under
88/// `TargetDeployment`.
89#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
90#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
91#[serde(rename_all = "camelCase")]
92pub struct TargetOperationsBundleSet {
93    /// Hash identifying this exact enabled-plugin set. Compare against
94    /// `OperationsReport.loadedBundleHash` to detect drift.
95    pub hash: String,
96    /// Presigned downloads for every bundle in the target set. The Operator
97    /// fetches only the ones it doesn't already have loaded at the right
98    /// version; already-loaded bundles are harmless to re-download.
99    #[serde(default, skip_serializing_if = "Vec::is_empty")]
100    pub bundles: Vec<OperationsBundleDownload>,
101}
102
103/// Request sent by the agent to the manager during periodic sync.
104#[derive(Debug, Clone, Serialize, Deserialize)]
105#[serde(rename_all = "camelCase")]
106pub struct SyncRequest {
107    /// The deployment ID this agent is managing.
108    pub deployment_id: String,
109    /// Stable identity for this Operator process. Used to fence update work.
110    #[serde(default)]
111    pub session: String,
112    /// Signals that this Operator persists and echoes execution claims.
113    #[serde(default)]
114    pub supports_execution_claims: bool,
115    /// Exact update claim returned by the previous sync response.
116    #[serde(default, skip_serializing_if = "Option::is_none")]
117    pub execution_claim: Option<SyncExecutionClaim>,
118    /// Current deployment state as seen by the agent.
119    #[serde(skip_serializing_if = "Option::is_none")]
120    pub current_state: Option<DeploymentState>,
121    /// Managed Alien resource status samples emitted by the Operator's deployment step.
122    #[serde(
123        default,
124        rename = "resourceHeartbeats",
125        skip_serializing_if = "Vec::is_empty"
126    )]
127    pub heartbeats: Vec<ResourceHeartbeat>,
128    /// Observed raw-resource inventory batches successfully read by the Operator.
129    #[serde(
130        default,
131        rename = "observedInventoryBatches",
132        skip_serializing_if = "Vec::is_empty"
133    )]
134    pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
135    /// Report-only capabilities observed by the Operator.
136    #[serde(default, skip_serializing_if = "Vec::is_empty")]
137    pub capabilities: Vec<OperatorCapabilityReport>,
138    /// Version of the Operator binary reporting this sync.
139    #[serde(default, skip_serializing_if = "Option::is_none")]
140    pub operator_version: Option<String>,
141    /// Report-only summary of the operations the Operator currently has
142    /// loaded. Absent means the Operator does not yet support reporting
143    /// this (older Operator versions), not that it has no operations.
144    #[serde(default, skip_serializing_if = "Option::is_none")]
145    pub operations_report: Option<OperationsReport>,
146}
147
148/// Response from the manager to the agent sync request.
149#[derive(Debug, Clone, Serialize, Deserialize)]
150#[serde(rename_all = "camelCase")]
151pub struct SyncResponse {
152    /// Exact operation attempt associated with `target`.
153    #[serde(default, skip_serializing_if = "Option::is_none")]
154    pub execution_claim: Option<SyncExecutionClaim>,
155    /// Authoritative deployment state from the manager.
156    ///
157    /// Pull agents use this to hydrate local state when attaching to an
158    /// already-imported deployment. Absent means the agent's local state is
159    /// already authoritative or no state has been established yet.
160    #[serde(default, skip_serializing_if = "Option::is_none")]
161    pub current_state: Option<DeploymentState>,
162    /// Target deployment the agent should converge toward.
163    /// None means no changes needed or this is an observe-only deployment.
164    #[serde(skip_serializing_if = "Option::is_none")]
165    pub target: Option<TargetDeployment>,
166    /// Public URL for the commands API (e.g. `https://manager.example.com/v1`).
167    /// Operators and app-owned receivers use this to lease pending commands;
168    /// Workers themselves receive pushes.
169    /// When absent, the agent falls back to its sync URL.
170    #[serde(default, skip_serializing_if = "Option::is_none")]
171    pub commands_url: Option<String>,
172    /// Target operations-bundle set the Operator should converge its loaded
173    /// plugin registry toward. None means no enabled plugin set has ever
174    /// been established for this project (nothing to sync), not that the
175    /// Operator is up to date — compare `hash` against the Operator's own
176    /// loaded hash to decide whether to download anything.
177    #[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/// Target deployment state for the agent to converge toward.
189#[derive(Debug, Clone, Serialize, Deserialize)]
190#[serde(rename_all = "camelCase")]
191pub struct TargetDeployment {
192    /// Release information (ID, version, stack definition).
193    pub release_info: ReleaseInfo,
194    /// Full deployment configuration (settings, env vars, etc.).
195    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        // current_state is None → should be omitted
221        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        // target is None → should be omitted
254        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        // Verify camelCase renaming works correctly
277        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        // snake_case should NOT work
285        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        // Simulates an old Operator that doesn't know about operationsReport:
374        // the field must default to None, never be required.
375        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}