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 chrono::{DateTime, Utc};
7use serde::{Deserialize, Serialize};
8use std::collections::BTreeMap;
9
10use crate::{
11    DeploymentConfig, DeploymentState, ObservedInventoryBatch, ReleaseInfo, ResourceHeartbeat,
12};
13
14/// State of an Operator capability as observed inside the environment.
15#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
16#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
17#[serde(rename_all = "kebab-case")]
18pub enum OperatorCapabilityState {
19    /// The Operator has the permission or local facility needed for the capability.
20    Granted,
21    /// The environment explicitly denied the capability.
22    Denied,
23    /// The capability does not apply in this environment.
24    Unavailable,
25}
26
27/// Report-only Operator capability status.
28#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
29#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
30#[serde(rename_all = "camelCase")]
31pub struct OperatorCapabilityReport {
32    /// Stable capability key, such as `k8s-workloads` or `logs`.
33    pub key: String,
34    /// Whether the capability is currently usable.
35    pub state: OperatorCapabilityState,
36    /// Optional human-readable detail from the Operator.
37    #[serde(default, skip_serializing_if = "Option::is_none")]
38    pub detail: Option<String>,
39}
40
41/// A single operation the Operator currently has loaded. Opaque identifiers
42/// only — no tier, description, or other plugin business logic crosses into
43/// this public crate.
44#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
45#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
46#[serde(rename_all = "camelCase")]
47pub struct ReportedOperation {
48    /// Name of the plugin that owns this operation.
49    pub plugin: String,
50    /// Version of the plugin that owns this operation.
51    pub plugin_version: String,
52    /// Name of the operation within the plugin.
53    pub name: String,
54}
55
56/// Report-only summary of the operations the Operator has loaded, used to
57/// confirm a bundle sync actually took effect.
58#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
59#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
60#[serde(rename_all = "camelCase")]
61pub struct OperationsReport {
62    /// Hash of the enabled-plugin bundle set the Operator currently has
63    /// loaded. Compared against the platform's target hash to detect drift.
64    #[serde(default, skip_serializing_if = "Option::is_none")]
65    pub loaded_bundle_hash: Option<String>,
66    /// Operations currently loaded and executable by the Operator.
67    #[serde(default, skip_serializing_if = "Vec::is_empty")]
68    pub operations: Vec<ReportedOperation>,
69}
70
71/// Origin of the exact Operator image running this process.
72#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
73#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
74#[serde(rename_all = "kebab-case")]
75pub enum OperatorImageSource {
76    /// The image came from an immutable generated package.
77    Package,
78    /// The image was configured directly by the installer.
79    Configured,
80}
81
82/// Exact immutable Operator image identity observed by the running process.
83#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
84#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
85#[serde(rename_all = "camelCase")]
86pub struct OperatorImageReport {
87    /// How the installer selected this image.
88    pub source: OperatorImageSource,
89    /// Package ID when `source` is `package`; otherwise `null`.
90    pub package_id: Option<String>,
91    /// Package version when `source` is `package`; otherwise `null`.
92    pub package_version: Option<String>,
93    /// Exact OCI image reference in `repository@sha256:digest` form.
94    pub image: String,
95    /// Exact lowercase OCI digest in `sha256:digest` form.
96    pub digest: String,
97}
98
99impl OperatorImageReport {
100    /// Validate that this report identifies one immutable image and that its
101    /// package fields agree with `source`.
102    pub fn validate(&self) -> Result<(), &'static str> {
103        let Some(encoded_digest) = self.digest.strip_prefix("sha256:") else {
104            return Err("operator image digest must start with sha256:");
105        };
106        if encoded_digest.len() != 64
107            || !encoded_digest
108                .bytes()
109                .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
110        {
111            return Err("operator image digest must contain 64 lowercase hex characters");
112        }
113
114        let Some((repository, image_digest)) = self.image.rsplit_once('@') else {
115            return Err("operator image must use repository@sha256:digest form");
116        };
117        if repository.is_empty()
118            || repository.chars().any(char::is_whitespace)
119            || repository.contains('@')
120            || image_digest != self.digest
121        {
122            return Err("operator image must contain one repository and the reported digest");
123        }
124
125        let package_id = self.package_id.as_deref().map(str::trim);
126        let package_version = self.package_version.as_deref().map(str::trim);
127        match self.source {
128            OperatorImageSource::Package
129                if package_id.is_some_and(|value| !value.is_empty())
130                    && package_version.is_some_and(|value| !value.is_empty()) =>
131            {
132                Ok(())
133            }
134            OperatorImageSource::Package => {
135                Err("package operator images require package ID and version")
136            }
137            OperatorImageSource::Configured
138                if self.package_id.is_none() && self.package_version.is_none() =>
139            {
140                Ok(())
141            }
142            OperatorImageSource::Configured => {
143                Err("configured operator images must not include package identity")
144            }
145        }
146    }
147}
148
149/// One bundle the Operator needs to download to reach `targetBundleHash`.
150/// The manager mints a short-lived presigned GET URL per bundle — the
151/// Operator never holds real cloud storage credentials, mirroring the OCI
152/// registry proxy's credential-injection pattern.
153#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
154#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
155#[serde(rename_all = "camelCase")]
156pub struct OperationsBundleDownload {
157    /// Plugin name this bundle provides.
158    pub plugin: String,
159    /// Plugin version this bundle provides.
160    pub plugin_version: String,
161    /// Presigned URL to GET the bundle ZIP from. Short-lived.
162    pub url: String,
163}
164
165/// Target operations-bundle set for the Operator to converge its loaded
166/// plugin registry toward, independent of any release/config target — a
167/// plugin can be enabled with no release change, so this is not nested under
168/// `TargetDeployment`.
169#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
170#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
171#[serde(rename_all = "camelCase")]
172pub struct TargetOperationsBundleSet {
173    /// Hash identifying this exact enabled-plugin set. Compare against
174    /// `OperationsReport.loadedBundleHash` to detect drift.
175    pub hash: String,
176    /// Presigned downloads for every bundle in the target set. The Operator
177    /// fetches only the ones it doesn't already have loaded at the right
178    /// version; already-loaded bundles are harmless to re-download.
179    #[serde(default, skip_serializing_if = "Vec::is_empty")]
180    pub bundles: Vec<OperationsBundleDownload>,
181}
182
183/// One release-independent container the Operator should run in its namespace.
184/// The manager sends the complete set on every sync. An empty set removes
185/// containers previously owned by this deployment.
186#[derive(Clone, Serialize, Deserialize, PartialEq, Eq)]
187#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
188#[serde(rename_all = "camelCase")]
189pub struct TargetDynamicContainer {
190    pub name: String,
191    pub generation: u64,
192    pub image: String,
193    pub cpu: String,
194    pub memory: String,
195    pub replicas: u32,
196    pub ports: Vec<u16>,
197    pub deleted: bool,
198    #[serde(default)]
199    pub env: BTreeMap<String, String>,
200    #[serde(default)]
201    pub secret_env: BTreeMap<String, String>,
202    #[serde(default, skip_serializing_if = "Option::is_none")]
203    pub health_check: Option<DynamicContainerHealthCheck>,
204    /// Stop an installed workload when its release no longer admits the image.
205    #[serde(default, skip_serializing_if = "Option::is_none")]
206    pub suspended_reason: Option<String>,
207}
208
209// Never include secret values in sync diagnostics.
210impl std::fmt::Debug for TargetDynamicContainer {
211    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
212        f.debug_struct("TargetDynamicContainer")
213            .field("name", &self.name)
214            .field("generation", &self.generation)
215            .field("image", &self.image)
216            .field("replicas", &self.replicas)
217            .finish_non_exhaustive()
218    }
219}
220
221#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
222#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
223#[serde(rename_all = "camelCase")]
224pub struct DynamicContainerHealthCheck {
225    pub path: String,
226    pub port: u16,
227}
228
229/// What the Operator observed after applying one target generation.
230#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
231#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
232#[serde(rename_all = "camelCase")]
233pub struct DynamicContainerReport {
234    pub name: String,
235    pub generation: u64,
236    pub status: DynamicContainerStatus,
237    #[serde(default, skip_serializing_if = "Option::is_none")]
238    pub message: Option<String>,
239}
240
241#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
242#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
243#[serde(rename_all = "lowercase")]
244pub enum DynamicContainerStatus {
245    Pending,
246    Running,
247    Failing,
248    Stopped,
249}
250
251/// Request sent by the agent to the manager during periodic sync.
252#[derive(Debug, Clone, Serialize, Deserialize)]
253#[serde(rename_all = "camelCase")]
254pub struct SyncRequest {
255    /// The deployment ID this agent is managing.
256    pub deployment_id: String,
257    /// Stable identity for this Operator process. Used to fence update work.
258    #[serde(default)]
259    pub session: String,
260    /// Signals that this Operator persists and echoes execution claims.
261    #[serde(default)]
262    pub supports_execution_claims: bool,
263    /// Signals that this Operator understands container tunnels. Older
264    /// Operators reject stacks that declare one, so the manager leaves
265    /// tunnels out of their targets.
266    #[serde(default)]
267    pub supports_tunnels: bool,
268    /// Exact update claim returned by the previous sync response.
269    #[serde(default, skip_serializing_if = "Option::is_none")]
270    pub execution_claim: Option<SyncExecutionClaim>,
271    /// Current deployment state as seen by the agent.
272    #[serde(skip_serializing_if = "Option::is_none")]
273    pub current_state: Option<DeploymentState>,
274    /// Managed Alien resource status samples emitted by the Operator's deployment step.
275    #[serde(
276        default,
277        rename = "resourceHeartbeats",
278        skip_serializing_if = "Vec::is_empty"
279    )]
280    pub heartbeats: Vec<ResourceHeartbeat>,
281    /// Observed raw-resource inventory batches successfully read by the Operator.
282    #[serde(
283        default,
284        rename = "observedInventoryBatches",
285        skip_serializing_if = "Vec::is_empty"
286    )]
287    pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
288    /// Report-only capabilities observed by the Operator.
289    #[serde(default, skip_serializing_if = "Vec::is_empty")]
290    pub capabilities: Vec<OperatorCapabilityReport>,
291    /// Version of the Operator binary reporting this sync.
292    #[serde(default, skip_serializing_if = "Option::is_none")]
293    pub operator_version: Option<String>,
294    /// Report-only summary of the operations the Operator currently has
295    /// loaded. Absent means the Operator does not yet support reporting
296    /// this (older Operator versions), not that it has no operations.
297    #[serde(default, skip_serializing_if = "Option::is_none")]
298    pub operations_report: Option<OperationsReport>,
299}
300
301/// Where the Operator read the application identity from.
302#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
303#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
304#[serde(rename_all = "camelCase")]
305pub enum ObservedApplicationSource {
306    /// Labels and pod statuses of the observed Kubernetes workloads.
307    Kubernetes,
308}
309
310/// A container image running in one observed application workload.
311#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
312#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
313#[serde(rename_all = "camelCase")]
314pub struct ObservedApplicationImage {
315    /// Inventory identity of the workload, matching the `rawIdentity` of its
316    /// observed resource sample (for example `apps/v1:Deployment:shop:api`).
317    pub workload: String,
318    /// Container name within the workload.
319    pub container: String,
320    /// Image reference reported by the container runtime.
321    pub image: String,
322    /// Registry manifest digest in `sha256:<hex>` form, when the runtime
323    /// reports one.
324    #[serde(default, skip_serializing_if = "Option::is_none")]
325    pub digest: Option<String>,
326}
327
328/// Application release the Operator observes running in its environment.
329///
330/// This identifies the customer's application, not the Operator: the
331/// Operator's own image is reported separately as [`OperatorImageReport`].
332#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
333#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
334#[serde(rename_all = "camelCase")]
335pub struct ObservedApplicationReport {
336    /// Where this identity was read from.
337    pub source: ObservedApplicationSource,
338    /// Helm chart name from the workloads' `helm.sh/chart` label. Present only
339    /// when every labelled workload names the same chart.
340    #[serde(default, skip_serializing_if = "Option::is_none")]
341    pub chart_name: Option<String>,
342    /// Helm chart version from the same label.
343    #[serde(default, skip_serializing_if = "Option::is_none")]
344    pub chart_version: Option<String>,
345    /// Distinct container images running in the observed workloads.
346    #[serde(default, skip_serializing_if = "Vec::is_empty")]
347    pub images: Vec<ObservedApplicationImage>,
348    /// Whether every workload kind could be listed. When `false`, the chart
349    /// and images describe only the workloads the Operator could read.
350    pub complete: bool,
351    /// When the workloads were read.
352    pub observed_at: DateTime<Utc>,
353}
354
355/// Extensible wire input for sync metadata that is not part of the
356/// long-standing [`SyncRequest`] struct-literal contract.
357#[derive(Debug, Clone, Serialize)]
358#[serde(rename_all = "camelCase")]
359pub struct SyncInput {
360    #[serde(flatten)]
361    request: SyncRequest,
362    /// Exact immutable image identity injected by the installer. Absent for
363    /// older installations that do not carry an image receipt.
364    #[serde(skip_serializing_if = "Option::is_none")]
365    operator_image: Option<OperatorImageReport>,
366    /// Application release observed in the environment. Absent when the
367    /// Operator observed none or predates this report.
368    #[serde(skip_serializing_if = "Option::is_none")]
369    application: Option<ObservedApplicationReport>,
370    /// Absent for older Operators. An empty report means no containers remain.
371    #[serde(skip_serializing_if = "Option::is_none")]
372    dynamic_containers: Option<Vec<DynamicContainerReport>>,
373}
374
375impl SyncInput {
376    /// Start building a sync payload while preserving the stable
377    /// [`SyncRequest`] literal surface for downstream callers.
378    pub fn builder(request: SyncRequest) -> SyncInputBuilder {
379        SyncInputBuilder {
380            request,
381            operator_image: None,
382            application: None,
383            dynamic_containers: None,
384        }
385    }
386}
387
388/// Builder for optional sync receipts. New optional wire metadata belongs
389/// here so adding it does not break downstream [`SyncRequest`] literals.
390pub struct SyncInputBuilder {
391    request: SyncRequest,
392    operator_image: Option<OperatorImageReport>,
393    application: Option<ObservedApplicationReport>,
394    dynamic_containers: Option<Vec<DynamicContainerReport>>,
395}
396
397impl SyncInputBuilder {
398    /// Attach the immutable image receipt for the running Operator.
399    pub fn operator_image(mut self, operator_image: OperatorImageReport) -> Self {
400        self.operator_image = Some(operator_image);
401        self
402    }
403
404    /// Attach the application release observed in the environment.
405    pub fn application(mut self, application: ObservedApplicationReport) -> Self {
406        self.application = Some(application);
407        self
408    }
409
410    pub fn dynamic_containers(mut self, reports: Vec<DynamicContainerReport>) -> Self {
411        self.dynamic_containers = Some(reports);
412        self
413    }
414
415    /// Finish the serializable sync payload.
416    pub fn build(self) -> SyncInput {
417        SyncInput {
418            request: self.request,
419            operator_image: self.operator_image,
420            application: self.application,
421            dynamic_containers: self.dynamic_containers,
422        }
423    }
424}
425
426/// Response from the manager to the agent sync request.
427#[derive(Debug, Clone, Serialize, Deserialize)]
428#[serde(rename_all = "camelCase")]
429pub struct SyncResponse {
430    /// Exact operation attempt associated with `target`.
431    #[serde(default, skip_serializing_if = "Option::is_none")]
432    pub execution_claim: Option<SyncExecutionClaim>,
433    /// Authoritative deployment state from the manager.
434    ///
435    /// Pull agents use this to hydrate local state when attaching to an
436    /// already-imported deployment. Absent means the agent's local state is
437    /// already authoritative or no state has been established yet.
438    #[serde(default, skip_serializing_if = "Option::is_none")]
439    pub current_state: Option<DeploymentState>,
440    /// Target deployment the agent should converge toward.
441    /// None means no changes needed or this is an observe-only deployment.
442    #[serde(skip_serializing_if = "Option::is_none")]
443    pub target: Option<TargetDeployment>,
444    /// Public URL for the commands API (e.g. `https://manager.example.com/v1`).
445    /// Operators and app-owned receivers use this to lease pending commands;
446    /// Workers themselves receive pushes.
447    /// When absent, the agent falls back to its sync URL.
448    #[serde(default, skip_serializing_if = "Option::is_none")]
449    pub commands_url: Option<String>,
450    /// Target operations-bundle set the Operator should converge its loaded
451    /// plugin registry toward. None means no enabled plugin set has ever
452    /// been established for this project (nothing to sync), not that the
453    /// Operator is up to date — compare `hash` against the Operator's own
454    /// loaded hash to decide whether to download anything.
455    #[serde(default, skip_serializing_if = "Option::is_none")]
456    pub target_operations_bundle_set: Option<TargetOperationsBundleSet>,
457    /// Complete release-independent target set. None means the manager does
458    /// not support this protocol; Some(empty) means remove owned containers.
459    #[serde(default, skip_serializing_if = "Option::is_none")]
460    pub target_dynamic_containers: Option<Vec<TargetDynamicContainer>>,
461    /// Base URL the Operator opens tunnel connections to. None means the
462    /// manager does not accept tunnels, so the Operator never dials.
463    #[serde(default, skip_serializing_if = "Option::is_none")]
464    pub tunnel_url: Option<String>,
465    /// Operator image the manager wants this Operator to run. Operators that
466    /// manage their own workload update to it; None means no opinion.
467    #[serde(default, skip_serializing_if = "Option::is_none")]
468    pub target_operator_image: Option<String>,
469}
470
471#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
472#[serde(rename_all = "camelCase")]
473pub struct SyncExecutionClaim {
474    pub operation_id: String,
475    pub attempt_id: String,
476}
477
478/// Target deployment state for the agent to converge toward.
479#[derive(Debug, Clone, Serialize, Deserialize)]
480#[serde(rename_all = "camelCase")]
481pub struct TargetDeployment {
482    /// Release information (ID, version, stack definition).
483    pub release_info: ReleaseInfo,
484    /// Full deployment configuration (settings, env vars, etc.).
485    pub config: DeploymentConfig,
486}
487
488#[cfg(test)]
489mod tests {
490    use super::*;
491
492    #[test]
493    fn test_sync_request_serialization() {
494        let req = SyncRequest {
495            deployment_id: "dep_abc123".to_string(),
496            session: "operator-test".to_string(),
497            supports_execution_claims: true,
498            supports_tunnels: true,
499            execution_claim: None,
500            current_state: None,
501            heartbeats: Vec::new(),
502            observed_inventory_batches: Vec::new(),
503            capabilities: Vec::new(),
504            operator_version: None,
505            operations_report: None,
506        };
507
508        let json = serde_json::to_value(&req).unwrap();
509        assert_eq!(json["deploymentId"], "dep_abc123");
510        assert_eq!(json["supportsExecutionClaims"], true);
511        // current_state is None → should be omitted
512        assert!(json.get("currentState").is_none());
513        assert!(json.get("resourceHeartbeats").is_none());
514        assert!(json.get("capabilities").is_none());
515        assert!(json.get("operatorVersion").is_none());
516        assert!(json.get("operatorImage").is_none());
517        assert!(json.get("operationsReport").is_none());
518    }
519
520    #[test]
521    fn test_sync_request_deserialization() {
522        let json = r#"{"deploymentId": "dep_xyz"}"#;
523        let req: SyncRequest = serde_json::from_str(json).unwrap();
524        assert_eq!(req.deployment_id, "dep_xyz");
525        assert!(req.session.is_empty());
526        assert!(!req.supports_execution_claims);
527        assert!(req.current_state.is_none());
528        assert!(req.heartbeats.is_empty());
529        assert!(req.observed_inventory_batches.is_empty());
530        assert!(req.capabilities.is_empty());
531        assert!(req.operator_version.is_none());
532        assert!(req.operations_report.is_none());
533    }
534
535    #[test]
536    fn test_sync_response_empty() {
537        let resp = SyncResponse {
538            execution_claim: None,
539            current_state: None,
540            target: None,
541            commands_url: None,
542            target_operations_bundle_set: None,
543            target_dynamic_containers: None,
544            tunnel_url: None,
545            target_operator_image: None,
546        };
547        let json = serde_json::to_value(&resp).unwrap();
548        // target is None → should be omitted
549        assert!(json.get("target").is_none());
550        assert!(json.get("currentState").is_none());
551        assert!(json.get("targetOperationsBundleSet").is_none());
552    }
553
554    #[test]
555    fn test_sync_response_roundtrip_no_target() {
556        let resp = SyncResponse {
557            execution_claim: None,
558            current_state: None,
559            target: None,
560            commands_url: None,
561            target_operations_bundle_set: None,
562            target_dynamic_containers: None,
563            tunnel_url: None,
564            target_operator_image: None,
565        };
566        let serialized = serde_json::to_string(&resp).unwrap();
567        let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
568        assert!(deserialized.target.is_none());
569        assert!(deserialized.current_state.is_none());
570    }
571
572    #[test]
573    fn test_sync_request_with_camel_case() {
574        // Verify camelCase renaming works correctly
575        let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
576        let req: SyncRequest = serde_json::from_str(json).unwrap();
577        assert_eq!(req.deployment_id, "dep_1");
578        assert!(req.current_state.is_none());
579        assert!(req.heartbeats.is_empty());
580        assert!(req.capabilities.is_empty());
581
582        // snake_case should NOT work
583        let json = r#"{"deployment_id": "dep_1"}"#;
584        assert!(serde_json::from_str::<SyncRequest>(json).is_err());
585    }
586
587    #[test]
588    fn test_sync_request_heartbeats_roundtrip() {
589        let json = serde_json::json!({
590            "deploymentId": "dep_1",
591            "resourceHeartbeats": [{
592                "deploymentId": "dep_1",
593                "resourceId": "api",
594                "resourceType": "container",
595                "controllerPlatform": "kubernetes",
596                "backend": "kubernetes",
597                "observedAt": "2026-01-01T00:00:00Z",
598                "data": {
599                    "resourceType": "container",
600                    "data": {
601                        "backend": "kubernetes",
602                        "status": {
603                            "health": "healthy",
604                            "lifecycle": "running",
605                            "message": null,
606                            "stale": false,
607                            "partial": false,
608                            "collectionIssues": []
609                        },
610                        "namespace": "default",
611                        "name": "api",
612                        "workloadKind": "deployment",
613                        "replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
614                        "restarts": 0,
615                        "cpu": null,
616                        "memory": null,
617                        "workload": null,
618                        "pods": [],
619                        "instances": [],
620                        "events": []
621                    }
622                },
623                "raw": []
624            }]
625        });
626
627        let req: SyncRequest = serde_json::from_value(json).unwrap();
628        assert_eq!(req.heartbeats.len(), 1);
629        assert_eq!(req.heartbeats[0].resource_id, "api");
630        assert!(req.capabilities.is_empty());
631
632        let serialized = serde_json::to_value(&req).unwrap();
633        assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
634    }
635
636    #[test]
637    fn test_sync_response_observe_only_state_roundtrip() {
638        let state = DeploymentState {
639            status: crate::DeploymentStatus::Running,
640            platform: crate::Platform::Kubernetes,
641            current_release: None,
642            target_release: None,
643            stack_state: None,
644            error: None,
645            environment_info: None,
646            runtime_metadata: None,
647            retry_requested: false,
648            protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
649        };
650        assert!(!state.has_desired());
651
652        let resp = SyncResponse {
653            execution_claim: None,
654            current_state: Some(state),
655            target: None,
656            commands_url: None,
657            target_operations_bundle_set: None,
658            target_dynamic_containers: None,
659            tunnel_url: None,
660            target_operator_image: None,
661        };
662
663        let serialized = serde_json::to_string(&resp).unwrap();
664        let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
665        let current_state = deserialized.current_state.unwrap();
666
667        assert_eq!(current_state.status, crate::DeploymentStatus::Running);
668        assert!(!current_state.has_desired());
669        assert!(deserialized.target.is_none());
670    }
671
672    #[test]
673    fn test_sync_request_operations_report_absent_by_default() {
674        // Simulates an old Operator that doesn't know about operationsReport:
675        // the field must default to None, never be required.
676        let json = r#"{"deploymentId": "dep_old_operator", "operatorVersion": "0.9.0"}"#;
677        let req: SyncRequest = serde_json::from_str(json).unwrap();
678        assert_eq!(req.operator_version.as_deref(), Some("0.9.0"));
679        assert!(req.operations_report.is_none());
680    }
681
682    #[test]
683    fn test_sync_request_operator_image_roundtrip() {
684        let digest = format!("sha256:{}", "a".repeat(64));
685        let report = OperatorImageReport {
686            source: OperatorImageSource::Package,
687            package_id: Some("pkg_operator".to_string()),
688            package_version: Some("1.2.3".to_string()),
689            image: format!("registry.example.com/operator@{digest}"),
690            digest,
691        };
692        report.validate().expect("valid package image report");
693
694        let value = serde_json::to_value(&report).unwrap();
695        assert_eq!(value["source"], "package");
696        assert_eq!(value["packageId"], "pkg_operator");
697        assert_eq!(value["packageVersion"], "1.2.3");
698
699        let decoded: OperatorImageReport = serde_json::from_value(value).unwrap();
700        assert_eq!(decoded, report);
701    }
702
703    #[test]
704    fn sync_input_adds_operator_image_without_expanding_sync_request() {
705        let digest = format!("sha256:{}", "a".repeat(64));
706        let request = SyncRequest {
707            deployment_id: "dep_1".to_string(),
708            session: String::new(),
709            supports_execution_claims: false,
710            supports_tunnels: false,
711            execution_claim: None,
712            current_state: None,
713            heartbeats: Vec::new(),
714            observed_inventory_batches: Vec::new(),
715            capabilities: Vec::new(),
716            operator_version: None,
717            operations_report: None,
718        };
719        let input = SyncInput::builder(request)
720            .operator_image(OperatorImageReport {
721                source: OperatorImageSource::Configured,
722                package_id: None,
723                package_version: None,
724                image: format!("registry.example.com/operator@{digest}"),
725                digest,
726            })
727            .build();
728
729        let value = serde_json::to_value(input).unwrap();
730        assert_eq!(value["deploymentId"], "dep_1");
731        assert_eq!(value["operatorImage"]["source"], "configured");
732    }
733
734    #[test]
735    fn sync_input_omits_absent_operator_image() {
736        let request = SyncRequest {
737            deployment_id: "dep_1".to_string(),
738            session: String::new(),
739            supports_execution_claims: false,
740            supports_tunnels: false,
741            execution_claim: None,
742            current_state: None,
743            heartbeats: Vec::new(),
744            observed_inventory_batches: Vec::new(),
745            capabilities: Vec::new(),
746            operator_version: None,
747            operations_report: None,
748        };
749
750        let value = serde_json::to_value(SyncInput::builder(request).build()).unwrap();
751        assert!(value.get("operatorImage").is_none());
752    }
753
754    #[test]
755    fn operator_image_report_rejects_mutable_or_inconsistent_identity() {
756        let digest = format!("sha256:{}", "a".repeat(64));
757        let mutable = OperatorImageReport {
758            source: OperatorImageSource::Configured,
759            package_id: None,
760            package_version: None,
761            image: "registry.example.com/operator:latest".to_string(),
762            digest: digest.clone(),
763        };
764        assert_eq!(
765            mutable.validate(),
766            Err("operator image must use repository@sha256:digest form")
767        );
768
769        let mismatched = OperatorImageReport {
770            source: OperatorImageSource::Configured,
771            package_id: None,
772            package_version: None,
773            image: format!("registry.example.com/operator@sha256:{}", "b".repeat(64)),
774            digest,
775        };
776        assert_eq!(
777            mismatched.validate(),
778            Err("operator image must contain one repository and the reported digest")
779        );
780    }
781
782    #[test]
783    fn operator_image_report_enforces_source_specific_package_fields() {
784        let digest = format!("sha256:{}", "a".repeat(64));
785        let image = format!("registry.example.com/operator@{digest}");
786        let missing_package = OperatorImageReport {
787            source: OperatorImageSource::Package,
788            package_id: None,
789            package_version: None,
790            image: image.clone(),
791            digest: digest.clone(),
792        };
793        assert_eq!(
794            missing_package.validate(),
795            Err("package operator images require package ID and version")
796        );
797
798        let configured_with_package = OperatorImageReport {
799            source: OperatorImageSource::Configured,
800            package_id: Some("pkg_operator".to_string()),
801            package_version: Some("1.2.3".to_string()),
802            image,
803            digest,
804        };
805        assert_eq!(
806            configured_with_package.validate(),
807            Err("configured operator images must not include package identity")
808        );
809    }
810
811    #[test]
812    fn test_sync_request_operations_report_roundtrip() {
813        let req = SyncRequest {
814            deployment_id: "dep_1".to_string(),
815            session: String::new(),
816            supports_execution_claims: false,
817            supports_tunnels: false,
818            execution_claim: None,
819            current_state: None,
820            heartbeats: Vec::new(),
821            observed_inventory_batches: Vec::new(),
822            capabilities: Vec::new(),
823            operator_version: Some("1.2.3".to_string()),
824            operations_report: Some(OperationsReport {
825                loaded_bundle_hash: Some("builtin:s3@1.0.0:key|".to_string()),
826                operations: vec![ReportedOperation {
827                    plugin: "s3".to_string(),
828                    plugin_version: "1.0.0".to_string(),
829                    name: "list-buckets".to_string(),
830                }],
831            }),
832        };
833
834        let json = serde_json::to_value(&req).unwrap();
835        assert_eq!(
836            json["operationsReport"]["loadedBundleHash"],
837            "builtin:s3@1.0.0:key|"
838        );
839        assert_eq!(json["operationsReport"]["operations"][0]["plugin"], "s3");
840        assert_eq!(
841            json["operationsReport"]["operations"][0]["pluginVersion"],
842            "1.0.0"
843        );
844        assert_eq!(
845            json["operationsReport"]["operations"][0]["name"],
846            "list-buckets"
847        );
848
849        let deserialized: SyncRequest = serde_json::from_value(json).unwrap();
850        assert_eq!(
851            deserialized
852                .operations_report
853                .as_ref()
854                .unwrap()
855                .loaded_bundle_hash,
856            req.operations_report.as_ref().unwrap().loaded_bundle_hash
857        );
858        assert_eq!(
859            deserialized.operations_report.unwrap().operations,
860            req.operations_report.unwrap().operations
861        );
862    }
863
864    #[test]
865    fn test_sync_response_target_operations_bundle_set_roundtrip() {
866        let resp = SyncResponse {
867            execution_claim: None,
868            current_state: None,
869            target: None,
870            commands_url: None,
871            target_operations_bundle_set: Some(TargetOperationsBundleSet {
872                hash: "builtin:s3@1.0.0:key|".to_string(),
873                bundles: vec![OperationsBundleDownload {
874                    plugin: "s3".to_string(),
875                    plugin_version: "1.0.0".to_string(),
876                    url: "https://storage.example.com/bundle.zip?sig=abc".to_string(),
877                }],
878            }),
879            target_dynamic_containers: None,
880            tunnel_url: None,
881            target_operator_image: None,
882        };
883
884        let json = serde_json::to_value(&resp).unwrap();
885        assert_eq!(
886            json["targetOperationsBundleSet"]["hash"],
887            "builtin:s3@1.0.0:key|"
888        );
889        assert_eq!(
890            json["targetOperationsBundleSet"]["bundles"][0]["plugin"],
891            "s3"
892        );
893
894        let deserialized: SyncResponse = serde_json::from_value(json).unwrap();
895        assert_eq!(
896            deserialized.target_operations_bundle_set,
897            resp.target_operations_bundle_set
898        );
899    }
900
901    #[test]
902    fn test_sync_response_target_operations_bundle_set_absent_by_default() {
903        let json = serde_json::json!({});
904        let resp: SyncResponse = serde_json::from_value(json).unwrap();
905        assert!(resp.target_operations_bundle_set.is_none());
906    }
907
908    #[test]
909    fn dynamic_container_sync_preserves_empty_target_and_hides_secrets_in_debug() {
910        let old_manager_response: SyncResponse = serde_json::from_str("{}").unwrap();
911        assert!(old_manager_response.target_dynamic_containers.is_none());
912
913        let empty_target = SyncResponse {
914            execution_claim: None,
915            current_state: None,
916            target: None,
917            commands_url: None,
918            target_operations_bundle_set: None,
919            target_dynamic_containers: Some(vec![]),
920            tunnel_url: None,
921            target_operator_image: None,
922        };
923        let json = serde_json::to_value(&empty_target).unwrap();
924        assert_eq!(json["targetDynamicContainers"], serde_json::json!([]));
925
926        let target = TargetDynamicContainer {
927            name: "api".to_string(),
928            generation: 2,
929            image: "example.com/api@sha256:abc".to_string(),
930            cpu: "0.5".to_string(),
931            memory: "512Mi".to_string(),
932            replicas: 1,
933            ports: vec![8080],
934            deleted: false,
935            env: BTreeMap::new(),
936            secret_env: BTreeMap::from([("TOKEN".to_string(), "private-value".to_string())]),
937            health_check: None,
938            suspended_reason: None,
939        };
940        assert!(!format!("{target:?}").contains("private-value"));
941        let roundtrip: TargetDynamicContainer =
942            serde_json::from_value(serde_json::to_value(&target).unwrap()).unwrap();
943        assert_eq!(roundtrip, target);
944    }
945}