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    /// Exact update claim returned by the previous sync response.
264    #[serde(default, skip_serializing_if = "Option::is_none")]
265    pub execution_claim: Option<SyncExecutionClaim>,
266    /// Current deployment state as seen by the agent.
267    #[serde(skip_serializing_if = "Option::is_none")]
268    pub current_state: Option<DeploymentState>,
269    /// Managed Alien resource status samples emitted by the Operator's deployment step.
270    #[serde(
271        default,
272        rename = "resourceHeartbeats",
273        skip_serializing_if = "Vec::is_empty"
274    )]
275    pub heartbeats: Vec<ResourceHeartbeat>,
276    /// Observed raw-resource inventory batches successfully read by the Operator.
277    #[serde(
278        default,
279        rename = "observedInventoryBatches",
280        skip_serializing_if = "Vec::is_empty"
281    )]
282    pub observed_inventory_batches: Vec<ObservedInventoryBatch>,
283    /// Report-only capabilities observed by the Operator.
284    #[serde(default, skip_serializing_if = "Vec::is_empty")]
285    pub capabilities: Vec<OperatorCapabilityReport>,
286    /// Version of the Operator binary reporting this sync.
287    #[serde(default, skip_serializing_if = "Option::is_none")]
288    pub operator_version: Option<String>,
289    /// Report-only summary of the operations the Operator currently has
290    /// loaded. Absent means the Operator does not yet support reporting
291    /// this (older Operator versions), not that it has no operations.
292    #[serde(default, skip_serializing_if = "Option::is_none")]
293    pub operations_report: Option<OperationsReport>,
294}
295
296/// Where the Operator read the application identity from.
297#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
298#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
299#[serde(rename_all = "camelCase")]
300pub enum ObservedApplicationSource {
301    /// Labels and pod statuses of the observed Kubernetes workloads.
302    Kubernetes,
303}
304
305/// A container image running in one observed application workload.
306#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
307#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
308#[serde(rename_all = "camelCase")]
309pub struct ObservedApplicationImage {
310    /// Inventory identity of the workload, matching the `rawIdentity` of its
311    /// observed resource sample (for example `apps/v1:Deployment:shop:api`).
312    pub workload: String,
313    /// Container name within the workload.
314    pub container: String,
315    /// Image reference reported by the container runtime.
316    pub image: String,
317    /// Registry manifest digest in `sha256:<hex>` form, when the runtime
318    /// reports one.
319    #[serde(default, skip_serializing_if = "Option::is_none")]
320    pub digest: Option<String>,
321}
322
323/// Application release the Operator observes running in its environment.
324///
325/// This identifies the customer's application, not the Operator: the
326/// Operator's own image is reported separately as [`OperatorImageReport`].
327#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
328#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
329#[serde(rename_all = "camelCase")]
330pub struct ObservedApplicationReport {
331    /// Where this identity was read from.
332    pub source: ObservedApplicationSource,
333    /// Helm chart name from the workloads' `helm.sh/chart` label. Present only
334    /// when every labelled workload names the same chart.
335    #[serde(default, skip_serializing_if = "Option::is_none")]
336    pub chart_name: Option<String>,
337    /// Helm chart version from the same label.
338    #[serde(default, skip_serializing_if = "Option::is_none")]
339    pub chart_version: Option<String>,
340    /// Distinct container images running in the observed workloads.
341    #[serde(default, skip_serializing_if = "Vec::is_empty")]
342    pub images: Vec<ObservedApplicationImage>,
343    /// Whether every workload kind could be listed. When `false`, the chart
344    /// and images describe only the workloads the Operator could read.
345    pub complete: bool,
346    /// When the workloads were read.
347    pub observed_at: DateTime<Utc>,
348}
349
350/// Extensible wire input for sync metadata that is not part of the
351/// long-standing [`SyncRequest`] struct-literal contract.
352#[derive(Debug, Clone, Serialize)]
353#[serde(rename_all = "camelCase")]
354pub struct SyncInput {
355    #[serde(flatten)]
356    request: SyncRequest,
357    /// Exact immutable image identity injected by the installer. Absent for
358    /// older installations that do not carry an image receipt.
359    #[serde(skip_serializing_if = "Option::is_none")]
360    operator_image: Option<OperatorImageReport>,
361    /// Application release observed in the environment. Absent when the
362    /// Operator observed none or predates this report.
363    #[serde(skip_serializing_if = "Option::is_none")]
364    application: Option<ObservedApplicationReport>,
365    /// Absent for older Operators. An empty report means no containers remain.
366    #[serde(skip_serializing_if = "Option::is_none")]
367    dynamic_containers: Option<Vec<DynamicContainerReport>>,
368}
369
370impl SyncInput {
371    /// Start building a sync payload while preserving the stable
372    /// [`SyncRequest`] literal surface for downstream callers.
373    pub fn builder(request: SyncRequest) -> SyncInputBuilder {
374        SyncInputBuilder {
375            request,
376            operator_image: None,
377            application: None,
378            dynamic_containers: None,
379        }
380    }
381}
382
383/// Builder for optional sync receipts. New optional wire metadata belongs
384/// here so adding it does not break downstream [`SyncRequest`] literals.
385pub struct SyncInputBuilder {
386    request: SyncRequest,
387    operator_image: Option<OperatorImageReport>,
388    application: Option<ObservedApplicationReport>,
389    dynamic_containers: Option<Vec<DynamicContainerReport>>,
390}
391
392impl SyncInputBuilder {
393    /// Attach the immutable image receipt for the running Operator.
394    pub fn operator_image(mut self, operator_image: OperatorImageReport) -> Self {
395        self.operator_image = Some(operator_image);
396        self
397    }
398
399    /// Attach the application release observed in the environment.
400    pub fn application(mut self, application: ObservedApplicationReport) -> Self {
401        self.application = Some(application);
402        self
403    }
404
405    pub fn dynamic_containers(mut self, reports: Vec<DynamicContainerReport>) -> Self {
406        self.dynamic_containers = Some(reports);
407        self
408    }
409
410    /// Finish the serializable sync payload.
411    pub fn build(self) -> SyncInput {
412        SyncInput {
413            request: self.request,
414            operator_image: self.operator_image,
415            application: self.application,
416            dynamic_containers: self.dynamic_containers,
417        }
418    }
419}
420
421/// Response from the manager to the agent sync request.
422#[derive(Debug, Clone, Serialize, Deserialize)]
423#[serde(rename_all = "camelCase")]
424pub struct SyncResponse {
425    /// Exact operation attempt associated with `target`.
426    #[serde(default, skip_serializing_if = "Option::is_none")]
427    pub execution_claim: Option<SyncExecutionClaim>,
428    /// Authoritative deployment state from the manager.
429    ///
430    /// Pull agents use this to hydrate local state when attaching to an
431    /// already-imported deployment. Absent means the agent's local state is
432    /// already authoritative or no state has been established yet.
433    #[serde(default, skip_serializing_if = "Option::is_none")]
434    pub current_state: Option<DeploymentState>,
435    /// Target deployment the agent should converge toward.
436    /// None means no changes needed or this is an observe-only deployment.
437    #[serde(skip_serializing_if = "Option::is_none")]
438    pub target: Option<TargetDeployment>,
439    /// Public URL for the commands API (e.g. `https://manager.example.com/v1`).
440    /// Operators and app-owned receivers use this to lease pending commands;
441    /// Workers themselves receive pushes.
442    /// When absent, the agent falls back to its sync URL.
443    #[serde(default, skip_serializing_if = "Option::is_none")]
444    pub commands_url: Option<String>,
445    /// Target operations-bundle set the Operator should converge its loaded
446    /// plugin registry toward. None means no enabled plugin set has ever
447    /// been established for this project (nothing to sync), not that the
448    /// Operator is up to date — compare `hash` against the Operator's own
449    /// loaded hash to decide whether to download anything.
450    #[serde(default, skip_serializing_if = "Option::is_none")]
451    pub target_operations_bundle_set: Option<TargetOperationsBundleSet>,
452    /// Complete release-independent target set. None means the manager does
453    /// not support this protocol; Some(empty) means remove owned containers.
454    #[serde(default, skip_serializing_if = "Option::is_none")]
455    pub target_dynamic_containers: Option<Vec<TargetDynamicContainer>>,
456}
457
458#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
459#[serde(rename_all = "camelCase")]
460pub struct SyncExecutionClaim {
461    pub operation_id: String,
462    pub attempt_id: String,
463}
464
465/// Target deployment state for the agent to converge toward.
466#[derive(Debug, Clone, Serialize, Deserialize)]
467#[serde(rename_all = "camelCase")]
468pub struct TargetDeployment {
469    /// Release information (ID, version, stack definition).
470    pub release_info: ReleaseInfo,
471    /// Full deployment configuration (settings, env vars, etc.).
472    pub config: DeploymentConfig,
473}
474
475#[cfg(test)]
476mod tests {
477    use super::*;
478
479    #[test]
480    fn test_sync_request_serialization() {
481        let req = SyncRequest {
482            deployment_id: "dep_abc123".to_string(),
483            session: "operator-test".to_string(),
484            supports_execution_claims: true,
485            execution_claim: None,
486            current_state: None,
487            heartbeats: Vec::new(),
488            observed_inventory_batches: Vec::new(),
489            capabilities: Vec::new(),
490            operator_version: None,
491            operations_report: None,
492        };
493
494        let json = serde_json::to_value(&req).unwrap();
495        assert_eq!(json["deploymentId"], "dep_abc123");
496        assert_eq!(json["supportsExecutionClaims"], true);
497        // current_state is None → should be omitted
498        assert!(json.get("currentState").is_none());
499        assert!(json.get("resourceHeartbeats").is_none());
500        assert!(json.get("capabilities").is_none());
501        assert!(json.get("operatorVersion").is_none());
502        assert!(json.get("operatorImage").is_none());
503        assert!(json.get("operationsReport").is_none());
504    }
505
506    #[test]
507    fn test_sync_request_deserialization() {
508        let json = r#"{"deploymentId": "dep_xyz"}"#;
509        let req: SyncRequest = serde_json::from_str(json).unwrap();
510        assert_eq!(req.deployment_id, "dep_xyz");
511        assert!(req.session.is_empty());
512        assert!(!req.supports_execution_claims);
513        assert!(req.current_state.is_none());
514        assert!(req.heartbeats.is_empty());
515        assert!(req.observed_inventory_batches.is_empty());
516        assert!(req.capabilities.is_empty());
517        assert!(req.operator_version.is_none());
518        assert!(req.operations_report.is_none());
519    }
520
521    #[test]
522    fn test_sync_response_empty() {
523        let resp = SyncResponse {
524            execution_claim: None,
525            current_state: None,
526            target: None,
527            commands_url: None,
528            target_operations_bundle_set: None,
529            target_dynamic_containers: None,
530        };
531        let json = serde_json::to_value(&resp).unwrap();
532        // target is None → should be omitted
533        assert!(json.get("target").is_none());
534        assert!(json.get("currentState").is_none());
535        assert!(json.get("targetOperationsBundleSet").is_none());
536    }
537
538    #[test]
539    fn test_sync_response_roundtrip_no_target() {
540        let resp = SyncResponse {
541            execution_claim: None,
542            current_state: None,
543            target: None,
544            commands_url: None,
545            target_operations_bundle_set: None,
546            target_dynamic_containers: None,
547        };
548        let serialized = serde_json::to_string(&resp).unwrap();
549        let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
550        assert!(deserialized.target.is_none());
551        assert!(deserialized.current_state.is_none());
552    }
553
554    #[test]
555    fn test_sync_request_with_camel_case() {
556        // Verify camelCase renaming works correctly
557        let json = r#"{"deploymentId": "dep_1", "currentState": null}"#;
558        let req: SyncRequest = serde_json::from_str(json).unwrap();
559        assert_eq!(req.deployment_id, "dep_1");
560        assert!(req.current_state.is_none());
561        assert!(req.heartbeats.is_empty());
562        assert!(req.capabilities.is_empty());
563
564        // snake_case should NOT work
565        let json = r#"{"deployment_id": "dep_1"}"#;
566        assert!(serde_json::from_str::<SyncRequest>(json).is_err());
567    }
568
569    #[test]
570    fn test_sync_request_heartbeats_roundtrip() {
571        let json = serde_json::json!({
572            "deploymentId": "dep_1",
573            "resourceHeartbeats": [{
574                "deploymentId": "dep_1",
575                "resourceId": "api",
576                "resourceType": "container",
577                "controllerPlatform": "kubernetes",
578                "backend": "kubernetes",
579                "observedAt": "2026-01-01T00:00:00Z",
580                "data": {
581                    "resourceType": "container",
582                    "data": {
583                        "backend": "kubernetes",
584                        "status": {
585                            "health": "healthy",
586                            "lifecycle": "running",
587                            "message": null,
588                            "stale": false,
589                            "partial": false,
590                            "collectionIssues": []
591                        },
592                        "namespace": "default",
593                        "name": "api",
594                        "workloadKind": "deployment",
595                        "replicas": { "desired": 1, "current": 1, "ready": 1, "available": 1, "updated": null, "misscheduled": null },
596                        "restarts": 0,
597                        "cpu": null,
598                        "memory": null,
599                        "workload": null,
600                        "pods": [],
601                        "instances": [],
602                        "events": []
603                    }
604                },
605                "raw": []
606            }]
607        });
608
609        let req: SyncRequest = serde_json::from_value(json).unwrap();
610        assert_eq!(req.heartbeats.len(), 1);
611        assert_eq!(req.heartbeats[0].resource_id, "api");
612        assert!(req.capabilities.is_empty());
613
614        let serialized = serde_json::to_value(&req).unwrap();
615        assert_eq!(serialized["resourceHeartbeats"][0]["resourceId"], "api");
616    }
617
618    #[test]
619    fn test_sync_response_observe_only_state_roundtrip() {
620        let state = DeploymentState {
621            status: crate::DeploymentStatus::Running,
622            platform: crate::Platform::Kubernetes,
623            current_release: None,
624            target_release: None,
625            stack_state: None,
626            error: None,
627            environment_info: None,
628            runtime_metadata: None,
629            retry_requested: false,
630            protocol_version: crate::DEPLOYMENT_PROTOCOL_VERSION,
631        };
632        assert!(!state.has_desired());
633
634        let resp = SyncResponse {
635            execution_claim: None,
636            current_state: Some(state),
637            target: None,
638            commands_url: None,
639            target_operations_bundle_set: None,
640            target_dynamic_containers: None,
641        };
642
643        let serialized = serde_json::to_string(&resp).unwrap();
644        let deserialized: SyncResponse = serde_json::from_str(&serialized).unwrap();
645        let current_state = deserialized.current_state.unwrap();
646
647        assert_eq!(current_state.status, crate::DeploymentStatus::Running);
648        assert!(!current_state.has_desired());
649        assert!(deserialized.target.is_none());
650    }
651
652    #[test]
653    fn test_sync_request_operations_report_absent_by_default() {
654        // Simulates an old Operator that doesn't know about operationsReport:
655        // the field must default to None, never be required.
656        let json = r#"{"deploymentId": "dep_old_operator", "operatorVersion": "0.9.0"}"#;
657        let req: SyncRequest = serde_json::from_str(json).unwrap();
658        assert_eq!(req.operator_version.as_deref(), Some("0.9.0"));
659        assert!(req.operations_report.is_none());
660    }
661
662    #[test]
663    fn test_sync_request_operator_image_roundtrip() {
664        let digest = format!("sha256:{}", "a".repeat(64));
665        let report = OperatorImageReport {
666            source: OperatorImageSource::Package,
667            package_id: Some("pkg_operator".to_string()),
668            package_version: Some("1.2.3".to_string()),
669            image: format!("registry.example.com/operator@{digest}"),
670            digest,
671        };
672        report.validate().expect("valid package image report");
673
674        let value = serde_json::to_value(&report).unwrap();
675        assert_eq!(value["source"], "package");
676        assert_eq!(value["packageId"], "pkg_operator");
677        assert_eq!(value["packageVersion"], "1.2.3");
678
679        let decoded: OperatorImageReport = serde_json::from_value(value).unwrap();
680        assert_eq!(decoded, report);
681    }
682
683    #[test]
684    fn sync_input_adds_operator_image_without_expanding_sync_request() {
685        let digest = format!("sha256:{}", "a".repeat(64));
686        let request = SyncRequest {
687            deployment_id: "dep_1".to_string(),
688            session: String::new(),
689            supports_execution_claims: false,
690            execution_claim: None,
691            current_state: None,
692            heartbeats: Vec::new(),
693            observed_inventory_batches: Vec::new(),
694            capabilities: Vec::new(),
695            operator_version: None,
696            operations_report: None,
697        };
698        let input = SyncInput::builder(request)
699            .operator_image(OperatorImageReport {
700                source: OperatorImageSource::Configured,
701                package_id: None,
702                package_version: None,
703                image: format!("registry.example.com/operator@{digest}"),
704                digest,
705            })
706            .build();
707
708        let value = serde_json::to_value(input).unwrap();
709        assert_eq!(value["deploymentId"], "dep_1");
710        assert_eq!(value["operatorImage"]["source"], "configured");
711    }
712
713    #[test]
714    fn sync_input_omits_absent_operator_image() {
715        let request = SyncRequest {
716            deployment_id: "dep_1".to_string(),
717            session: String::new(),
718            supports_execution_claims: false,
719            execution_claim: None,
720            current_state: None,
721            heartbeats: Vec::new(),
722            observed_inventory_batches: Vec::new(),
723            capabilities: Vec::new(),
724            operator_version: None,
725            operations_report: None,
726        };
727
728        let value = serde_json::to_value(SyncInput::builder(request).build()).unwrap();
729        assert!(value.get("operatorImage").is_none());
730    }
731
732    #[test]
733    fn operator_image_report_rejects_mutable_or_inconsistent_identity() {
734        let digest = format!("sha256:{}", "a".repeat(64));
735        let mutable = OperatorImageReport {
736            source: OperatorImageSource::Configured,
737            package_id: None,
738            package_version: None,
739            image: "registry.example.com/operator:latest".to_string(),
740            digest: digest.clone(),
741        };
742        assert_eq!(
743            mutable.validate(),
744            Err("operator image must use repository@sha256:digest form")
745        );
746
747        let mismatched = OperatorImageReport {
748            source: OperatorImageSource::Configured,
749            package_id: None,
750            package_version: None,
751            image: format!("registry.example.com/operator@sha256:{}", "b".repeat(64)),
752            digest,
753        };
754        assert_eq!(
755            mismatched.validate(),
756            Err("operator image must contain one repository and the reported digest")
757        );
758    }
759
760    #[test]
761    fn operator_image_report_enforces_source_specific_package_fields() {
762        let digest = format!("sha256:{}", "a".repeat(64));
763        let image = format!("registry.example.com/operator@{digest}");
764        let missing_package = OperatorImageReport {
765            source: OperatorImageSource::Package,
766            package_id: None,
767            package_version: None,
768            image: image.clone(),
769            digest: digest.clone(),
770        };
771        assert_eq!(
772            missing_package.validate(),
773            Err("package operator images require package ID and version")
774        );
775
776        let configured_with_package = OperatorImageReport {
777            source: OperatorImageSource::Configured,
778            package_id: Some("pkg_operator".to_string()),
779            package_version: Some("1.2.3".to_string()),
780            image,
781            digest,
782        };
783        assert_eq!(
784            configured_with_package.validate(),
785            Err("configured operator images must not include package identity")
786        );
787    }
788
789    #[test]
790    fn test_sync_request_operations_report_roundtrip() {
791        let req = SyncRequest {
792            deployment_id: "dep_1".to_string(),
793            session: String::new(),
794            supports_execution_claims: false,
795            execution_claim: None,
796            current_state: None,
797            heartbeats: Vec::new(),
798            observed_inventory_batches: Vec::new(),
799            capabilities: Vec::new(),
800            operator_version: Some("1.2.3".to_string()),
801            operations_report: Some(OperationsReport {
802                loaded_bundle_hash: Some("builtin:s3@1.0.0:key|".to_string()),
803                operations: vec![ReportedOperation {
804                    plugin: "s3".to_string(),
805                    plugin_version: "1.0.0".to_string(),
806                    name: "list-buckets".to_string(),
807                }],
808            }),
809        };
810
811        let json = serde_json::to_value(&req).unwrap();
812        assert_eq!(
813            json["operationsReport"]["loadedBundleHash"],
814            "builtin:s3@1.0.0:key|"
815        );
816        assert_eq!(json["operationsReport"]["operations"][0]["plugin"], "s3");
817        assert_eq!(
818            json["operationsReport"]["operations"][0]["pluginVersion"],
819            "1.0.0"
820        );
821        assert_eq!(
822            json["operationsReport"]["operations"][0]["name"],
823            "list-buckets"
824        );
825
826        let deserialized: SyncRequest = serde_json::from_value(json).unwrap();
827        assert_eq!(
828            deserialized
829                .operations_report
830                .as_ref()
831                .unwrap()
832                .loaded_bundle_hash,
833            req.operations_report.as_ref().unwrap().loaded_bundle_hash
834        );
835        assert_eq!(
836            deserialized.operations_report.unwrap().operations,
837            req.operations_report.unwrap().operations
838        );
839    }
840
841    #[test]
842    fn test_sync_response_target_operations_bundle_set_roundtrip() {
843        let resp = SyncResponse {
844            execution_claim: None,
845            current_state: None,
846            target: None,
847            commands_url: None,
848            target_operations_bundle_set: Some(TargetOperationsBundleSet {
849                hash: "builtin:s3@1.0.0:key|".to_string(),
850                bundles: vec![OperationsBundleDownload {
851                    plugin: "s3".to_string(),
852                    plugin_version: "1.0.0".to_string(),
853                    url: "https://storage.example.com/bundle.zip?sig=abc".to_string(),
854                }],
855            }),
856            target_dynamic_containers: None,
857        };
858
859        let json = serde_json::to_value(&resp).unwrap();
860        assert_eq!(
861            json["targetOperationsBundleSet"]["hash"],
862            "builtin:s3@1.0.0:key|"
863        );
864        assert_eq!(
865            json["targetOperationsBundleSet"]["bundles"][0]["plugin"],
866            "s3"
867        );
868
869        let deserialized: SyncResponse = serde_json::from_value(json).unwrap();
870        assert_eq!(
871            deserialized.target_operations_bundle_set,
872            resp.target_operations_bundle_set
873        );
874    }
875
876    #[test]
877    fn test_sync_response_target_operations_bundle_set_absent_by_default() {
878        let json = serde_json::json!({});
879        let resp: SyncResponse = serde_json::from_value(json).unwrap();
880        assert!(resp.target_operations_bundle_set.is_none());
881    }
882
883    #[test]
884    fn dynamic_container_sync_preserves_empty_target_and_hides_secrets_in_debug() {
885        let old_manager_response: SyncResponse = serde_json::from_str("{}").unwrap();
886        assert!(old_manager_response.target_dynamic_containers.is_none());
887
888        let empty_target = SyncResponse {
889            execution_claim: None,
890            current_state: None,
891            target: None,
892            commands_url: None,
893            target_operations_bundle_set: None,
894            target_dynamic_containers: Some(vec![]),
895        };
896        let json = serde_json::to_value(&empty_target).unwrap();
897        assert_eq!(json["targetDynamicContainers"], serde_json::json!([]));
898
899        let target = TargetDynamicContainer {
900            name: "api".to_string(),
901            generation: 2,
902            image: "example.com/api@sha256:abc".to_string(),
903            cpu: "0.5".to_string(),
904            memory: "512Mi".to_string(),
905            replicas: 1,
906            ports: vec![8080],
907            deleted: false,
908            env: BTreeMap::new(),
909            secret_env: BTreeMap::from([("TOKEN".to_string(), "private-value".to_string())]),
910            health_check: None,
911            suspended_reason: None,
912        };
913        assert!(!format!("{target:?}").contains("private-value"));
914        let roundtrip: TargetDynamicContainer =
915            serde_json::from_value(serde_json::to_value(&target).unwrap()).unwrap();
916        assert_eq!(roundtrip, target);
917    }
918}