use serde::{Deserialize, Serialize};
use crate::transport::{BulkTransportReady, LocalTransportReady, RelayLeaseReady};
pub const WORKLOAD_TRANSPORT_BARRIER_VERSION: u8 = 2;
pub const WORKLOAD_TRANSPORT_CONTROL_BYTES: u64 = 8 * 1024 * 1024;
pub const WORKLOAD_TRANSPORT_CONTROL_FRAMES: u64 = 256;
pub const WORKLOAD_TRANSPORT_BULK_BYTES: u64 = 32 * 1024 * 1024;
pub const WORKLOAD_TRANSPORT_BULK_FRAMES: u64 = 256;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Ready {
pub boot_time_ns: u64,
pub init_time_ns: u64,
pub ready_time_ns: u64,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub agent_version: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub bulk_transport: Option<BulkTransportReady>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub relay_lease: Option<RelayLeaseReady>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub local_transport: Option<LocalTransportReady>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workload_transport_barrier_version: Option<u8>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ClockSync {
pub unix_time_nanos: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Ping {}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Pong {}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Touch {}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Touched {
pub activity_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadFreeze {
#[serde(default)]
pub external_mount_tags: Vec<String>,
pub attempt_id: String,
pub host_input: WorkloadTransportPosition,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadFrozen {
pub attempt_id: String,
pub guest_bulk_bytes_target: u64,
pub input_credit: WorkloadTransportCredit,
#[serde(default)]
pub external_mounts_synced: bool,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadTransportPosition {
pub control_bytes: u64,
pub control_frames: u64,
pub bulk_bytes: u64,
pub bulk_frames: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadTransportCredit {
pub control_bytes: u64,
pub control_frames: u64,
pub bulk_bytes: u64,
pub bulk_frames: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadThaw {
pub attempt_id: String,
pub mode: WorkloadThawMode,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WorkloadThawMode {
Continue,
Restore,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkloadThawed {
pub attempt_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RootDiskGrow {
pub size_bytes: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RootDiskState {
pub filesystem_bytes: u64,
pub device_bytes: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CoreError {
pub kind: CoreErrorKind,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub offending_type: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workload_failure: Option<WorkloadFailure>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkloadFailure {
pub attempt_id: String,
pub disposition: WorkloadFailureDisposition,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WorkloadFailureDisposition {
Unavailable,
RecoveryRequired,
#[serde(other)]
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CoreErrorKind {
MalformedMessage,
UnsupportedMessageType,
UnsupportedProtocolGeneration,
InvalidFlags,
InvalidPayload,
InvalidSession,
CapabilityUnavailable,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InitResolved {
pub default_user: ResolvedUser,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct ResolvedUser {
pub uid: u32,
pub gid: u32,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct InitAck {}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RelayClientDisconnected {
pub id_start: u32,
pub id_end_exclusive: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub incarnation: Option<[u8; 16]>,
}
#[cfg(test)]
mod tests {
use serde::{Deserialize, Serialize};
use super::{Ready, RelayClientDisconnected};
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
struct LegacyReady {
boot_time_ns: u64,
init_time_ns: u64,
ready_time_ns: u64,
#[serde(default, skip_serializing_if = "String::is_empty")]
agent_version: String,
}
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
struct LegacyRelayClientDisconnected {
id_start: u32,
id_end_exclusive: u32,
}
#[test]
fn ready_without_transport_capabilities_is_byte_compatible() {
let legacy = LegacyReady {
boot_time_ns: 11,
init_time_ns: 22,
ready_time_ns: 33,
agent_version: "0.6.8".into(),
};
let current = Ready {
boot_time_ns: legacy.boot_time_ns,
init_time_ns: legacy.init_time_ns,
ready_time_ns: legacy.ready_time_ns,
agent_version: legacy.agent_version.clone(),
bulk_transport: None,
relay_lease: None,
local_transport: None,
workload_transport_barrier_version: None,
};
let mut legacy_bytes = Vec::new();
ciborium::into_writer(&legacy, &mut legacy_bytes).unwrap();
let mut current_bytes = Vec::new();
ciborium::into_writer(¤t, &mut current_bytes).unwrap();
assert_eq!(current_bytes, legacy_bytes);
let decoded: Ready = ciborium::from_reader(legacy_bytes.as_slice()).unwrap();
assert!(decoded.bulk_transport.is_none());
assert!(decoded.relay_lease.is_none());
assert!(decoded.local_transport.is_none());
assert!(decoded.workload_transport_barrier_version.is_none());
}
#[test]
fn relay_disconnect_without_incarnation_is_byte_compatible() {
let legacy = LegacyRelayClientDisconnected {
id_start: 1,
id_end_exclusive: 1024,
};
let current = RelayClientDisconnected {
id_start: legacy.id_start,
id_end_exclusive: legacy.id_end_exclusive,
incarnation: None,
};
let mut legacy_bytes = Vec::new();
ciborium::into_writer(&legacy, &mut legacy_bytes).unwrap();
let mut current_bytes = Vec::new();
ciborium::into_writer(¤t, &mut current_bytes).unwrap();
assert_eq!(current_bytes, legacy_bytes);
let decoded: RelayClientDisconnected =
ciborium::from_reader(legacy_bytes.as_slice()).unwrap();
assert_eq!(decoded.id_start, current.id_start);
assert_eq!(decoded.id_end_exclusive, current.id_end_exclusive);
assert_eq!(decoded.incarnation, None);
}
}
#[cfg(test)]
mod workload_tests {
use super::*;
#[test]
fn workload_barrier_payloads_roundtrip_without_resetting_counters() {
let position = WorkloadTransportPosition {
control_bytes: 73 * WORKLOAD_TRANSPORT_CONTROL_BYTES,
control_frames: 20_000,
bulk_bytes: 91 * WORKLOAD_TRANSPORT_BULK_BYTES,
bulk_frames: 30_000,
};
let freeze = WorkloadFreeze {
external_mount_tags: Vec::new(),
attempt_id: "captured-generation".into(),
host_input: position,
};
let mut bytes = Vec::new();
ciborium::into_writer(&freeze, &mut bytes).unwrap();
let decoded: WorkloadFreeze = ciborium::from_reader(bytes.as_slice()).unwrap();
assert_eq!(decoded, freeze);
let frozen = WorkloadFrozen {
external_mounts_synced: false,
attempt_id: freeze.attempt_id,
guest_bulk_bytes_target: 987_654_321,
input_credit: WorkloadTransportCredit {
control_bytes: position.control_bytes + 100,
control_frames: position.control_frames + 2,
bulk_bytes: position.bulk_bytes + 200,
bulk_frames: position.bulk_frames + 3,
},
};
bytes.clear();
ciborium::into_writer(&frozen, &mut bytes).unwrap();
let decoded: WorkloadFrozen = ciborium::from_reader(bytes.as_slice()).unwrap();
assert_eq!(decoded, frozen);
}
#[test]
fn superseded_development_freeze_payloads_do_not_imply_safe_boundaries() {
let old = serde_json::json!({"attempt_id":"old-development-capture"});
assert!(serde_json::from_value::<WorkloadFreeze>(old.clone()).is_err());
assert!(serde_json::from_value::<WorkloadFrozen>(old).is_err());
}
#[test]
fn unknown_barrier_capability_is_preserved_for_explicit_negotiation() {
let ready = Ready {
workload_transport_barrier_version: Some(99),
..Ready::default()
};
let mut bytes = Vec::new();
ciborium::into_writer(&ready, &mut bytes).unwrap();
let decoded: Ready = ciborium::from_reader(bytes.as_slice()).unwrap();
assert_eq!(decoded.workload_transport_barrier_version, Some(99));
assert_ne!(
decoded.workload_transport_barrier_version,
Some(WORKLOAD_TRANSPORT_BARRIER_VERSION)
);
}
#[test]
fn input_window_fits_a_maximum_primary_frame() {
let credit = WorkloadTransportCredit {
control_bytes: WORKLOAD_TRANSPORT_CONTROL_BYTES,
control_frames: WORKLOAD_TRANSPORT_CONTROL_FRAMES,
bulk_bytes: WORKLOAD_TRANSPORT_BULK_BYTES,
bulk_frames: WORKLOAD_TRANSPORT_BULK_FRAMES,
};
assert!(credit.control_bytes >= crate::codec::MAX_FRAME_SIZE as u64 + 4);
assert!(credit.control_frames > 0);
assert!(credit.bulk_frames > 0);
}
#[test]
fn freezer_error_details_are_additive_and_unknown_details_are_not_unavailable() {
let old = serde_json::json!({"kind":"capability_unavailable", "message":"freezer failed"});
let decoded: CoreError = serde_json::from_value(old.clone()).unwrap();
assert!(decoded.workload_failure.is_none());
let mut new = old;
new["workload_failure"] =
serde_json::json!({"attempt_id":"a", "disposition":"future_state"});
let decoded: CoreError = serde_json::from_value(new.clone()).unwrap();
assert_eq!(
decoded.workload_failure.unwrap().disposition,
WorkloadFailureDisposition::Unknown
);
#[derive(Deserialize)]
struct OldCoreError {
kind: CoreErrorKind,
message: String,
}
let old_reader: OldCoreError = serde_json::from_value(new).unwrap();
assert_eq!(old_reader.kind, CoreErrorKind::CapabilityUnavailable);
assert_eq!(old_reader.message, "freezer failed");
}
}