use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use crate::config::ServiceWithMirrors;
use crate::{MirrorConfig, MirrorProviderSlot, MirrorShape, Provider};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum SyncState {
Synced,
OutOfSync,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum HealthState {
Healthy,
Progressing,
Degraded,
Missing,
Idle,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "kebab-case")]
pub enum WireContainerStatus {
Running,
Restarting {
restart_count: u32,
last_exit_code: i32,
last_finished_at_unix_ms: u64,
},
Stopped,
Failed {
reason: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum Runtime {
Process,
Containers,
}
impl Runtime {
pub fn from_mirror(mirror: &MirrorConfig) -> Self {
let any_container = mirror.providers.values().any(|slot| match slot {
MirrorProviderSlot::Inline { kind, .. } => matches!(
kind,
Provider::MiniflareContainer | Provider::MinioContainer | Provider::LocalContainer
),
MirrorProviderSlot::Reference { .. } => true,
});
if any_container {
Runtime::Containers
} else {
Runtime::Process
}
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct MirrorObservation {
pub running: bool,
pub ready: bool,
pub errored: bool,
pub live_revision: Option<String>,
pub live_fields: BTreeMap<String, BTreeMap<String, String>>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DriftEntry {
pub path: String,
pub desired: String,
pub live: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CellStatus {
pub env: String,
pub sync: SyncState,
pub health: HealthState,
pub runtime: Runtime,
pub shape: MirrorShape,
pub declared_revision: Option<String>,
pub live_revision: Option<String>,
pub provider_label: String,
pub drift: Vec<DriftEntry>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workload_status: Option<WireContainerStatus>,
}
impl CellStatus {
pub fn drift_count(&self) -> usize {
self.drift.len()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ServiceStatus {
pub name: String,
pub address: String,
pub cells: BTreeMap<String, CellStatus>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct StatusSummary {
pub synced: usize,
pub out_of_sync: usize,
pub unknown: usize,
pub healthy: usize,
pub progressing: usize,
pub degraded: usize,
pub missing: usize,
pub idle: usize,
}
pub fn compute_cell(
env: &str,
mirror: &MirrorConfig,
obs: Option<&MirrorObservation>,
) -> CellStatus {
let runtime = Runtime::from_mirror(mirror);
let declared_revision = declared_revision(mirror);
let provider_label = provider_label(mirror);
let (sync, health, live_revision, drift) = match obs {
None => (SyncState::Unknown, HealthState::Missing, None, Vec::new()),
Some(o) => {
let drift = compute_drift(mirror, o);
let rev_diverges = matches!(
(&declared_revision, &o.live_revision),
(Some(d), Some(l)) if d != l
);
let sync = if !drift.is_empty() || rev_diverges {
SyncState::OutOfSync
} else {
SyncState::Synced
};
let health = if o.errored {
HealthState::Degraded
} else if o.running && o.ready {
HealthState::Healthy
} else if o.running {
HealthState::Progressing
} else if mirror.shape == MirrorShape::Local {
HealthState::Idle
} else {
HealthState::Missing
};
(sync, health, o.live_revision.clone(), drift)
}
};
CellStatus {
env: env.to_string(),
sync,
health,
runtime,
shape: mirror.shape,
declared_revision,
live_revision,
provider_label,
drift,
workload_status: None,
}
}
pub fn compute_service(
svc: &ServiceWithMirrors,
observations: &BTreeMap<String, MirrorObservation>,
) -> ServiceStatus {
let cells = svc
.mirrors
.iter()
.map(|(env, mirror)| {
let cell = compute_cell(env, mirror, observations.get(env));
(env.clone(), cell)
})
.collect();
ServiceStatus {
name: svc.service.name.clone(),
address: svc.service.address.label(),
cells,
}
}
pub fn summarize(services: &[ServiceStatus]) -> StatusSummary {
let mut s = StatusSummary::default();
for svc in services {
for cell in svc.cells.values() {
match cell.sync {
SyncState::Synced => s.synced += 1,
SyncState::OutOfSync => s.out_of_sync += 1,
SyncState::Unknown => s.unknown += 1,
}
match cell.health {
HealthState::Healthy => s.healthy += 1,
HealthState::Progressing => s.progressing += 1,
HealthState::Degraded => s.degraded += 1,
HealthState::Missing => s.missing += 1,
HealthState::Idle => s.idle += 1,
}
}
}
s
}
const REVISION_KEYS: [&str; 3] = ["image", "version", "tag"];
fn declared_revision(mirror: &MirrorConfig) -> Option<String> {
for slot in mirror.providers.values() {
let fields = slot_fields(slot);
for key in REVISION_KEYS {
if let Some(v) = fields.get(key).and_then(toml_value_to_string) {
return Some(v);
}
}
}
None
}
fn provider_label(mirror: &MirrorConfig) -> String {
let mut seen: Vec<String> = Vec::new();
for slot in mirror.providers.values() {
let label = match slot {
MirrorProviderSlot::Reference { provider_id, .. } => provider_id.clone(),
MirrorProviderSlot::Inline { kind, .. } => provider_kind_label(*kind),
};
if !seen.contains(&label) {
seen.push(label);
}
}
seen.join(" + ")
}
fn compute_drift(mirror: &MirrorConfig, obs: &MirrorObservation) -> Vec<DriftEntry> {
let mut out = Vec::new();
for (role, slot) in &mirror.providers {
let Some(live_slot) = obs.live_fields.get(role) else {
continue;
};
let declared = slot_fields(slot);
for (key, live_val) in live_slot {
let desired = declared.get(key).and_then(toml_value_to_string);
if let Some(desired) = desired {
if &desired != live_val {
out.push(DriftEntry {
path: format!("providers.{role}.{key}"),
desired,
live: live_val.clone(),
});
}
}
}
}
out.sort_by(|a, b| a.path.cmp(&b.path));
out
}
fn slot_fields(slot: &MirrorProviderSlot) -> &BTreeMap<String, toml::Value> {
match slot {
MirrorProviderSlot::Reference { fields, .. } => fields,
MirrorProviderSlot::Inline { fields, .. } => fields,
}
}
fn toml_value_to_string(v: &toml::Value) -> Option<String> {
match v {
toml::Value::String(s) => Some(s.clone()),
toml::Value::Integer(n) => Some(n.to_string()),
toml::Value::Float(f) => Some(f.to_string()),
toml::Value::Boolean(b) => Some(b.to_string()),
_ => None,
}
}
fn provider_kind_label(kind: Provider) -> String {
match kind {
Provider::Cloudflare => "cloudflare",
Provider::Hetzner => "hetzner",
Provider::Vultr => "vultr",
Provider::Static => "static",
Provider::MiniflareNative => "miniflare-native",
Provider::LocalContainer => "local-container",
Provider::LocalProcess => "local-process",
Provider::Device => "device",
Provider::MiniflareContainer => "miniflare-container",
Provider::MinioContainer => "minio-container",
Provider::LocalPgDev => "local-pg-dev",
Provider::LocalMailcrab => "local-mailcrab",
Provider::LocalS3Fs => "local-s3-fs",
}
.to_string()
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SyncHistoryEntry {
pub id: String,
pub service: String,
pub env: String,
pub status: SyncOutcome,
pub started_at: chrono::DateTime<chrono::Utc>,
pub completed_at: chrono::DateTime<chrono::Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub triggered_by: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rev: Option<String>,
pub workload_count: u32,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum SyncOutcome {
Success,
Failed,
Cancelled,
}
pub fn new_sync_id() -> String {
let mut bytes = [0u8; 8];
getrandom::getrandom(&mut bytes).unwrap_or(());
hex::encode(bytes)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::{ServiceConfig, ServiceWithMirrors};
fn mirror(src: &str) -> MirrorConfig {
toml::from_str(src).expect("mirror toml")
}
fn local_static_mirror() -> MirrorConfig {
mirror(
"schema_version = 1\nshape = \"local\"\n\n[providers.static]\nkind = \"miniflare-native\"\nport = 4321\n",
)
}
fn cloudflare_mirror() -> MirrorConfig {
mirror(
"schema_version = 1\nshape = \"single-machine\"\n\n[providers.static]\nuse = \"cloudflare\"\nimage = \"caddy:2.8.1\"\n",
)
}
fn sim_miniflare_mirror() -> MirrorConfig {
mirror(
"schema_version = 1\nshape = \"local\"\n\n[providers.static]\nkind = \"miniflare-container\"\nimage = \"caddy:2.8.1\"\nport = 8080\n\n[providers.object_store]\nkind = \"minio-container\"\n",
)
}
#[test]
fn runtime_process_for_local_static_only() {
assert_eq!(
Runtime::from_mirror(&local_static_mirror()),
Runtime::Process
);
}
#[test]
fn runtime_containers_for_inline_container_kinds() {
assert_eq!(
Runtime::from_mirror(&sim_miniflare_mirror()),
Runtime::Containers
);
}
#[test]
fn runtime_containers_for_referenced_provider() {
assert_eq!(
Runtime::from_mirror(&cloudflare_mirror()),
Runtime::Containers
);
}
#[test]
fn no_observation_is_unknown_missing() {
let cell = compute_cell("ha", &cloudflare_mirror(), None);
assert_eq!(cell.sync, SyncState::Unknown);
assert_eq!(cell.health, HealthState::Missing);
assert!(cell.drift.is_empty());
assert_eq!(cell.live_revision, None);
}
#[test]
fn running_ready_local_is_synced_healthy() {
let obs = MirrorObservation {
running: true,
ready: true,
..Default::default()
};
let cell = compute_cell("dev", &local_static_mirror(), Some(&obs));
assert_eq!(cell.sync, SyncState::Synced);
assert_eq!(cell.health, HealthState::Healthy);
assert_eq!(cell.runtime, Runtime::Process);
}
#[test]
fn declared_but_down_local_is_idle_not_missing() {
let obs = MirrorObservation {
running: false,
..Default::default()
};
let cell = compute_cell("sim", &sim_miniflare_mirror(), Some(&obs));
assert_eq!(cell.health, HealthState::Idle);
assert_eq!(cell.sync, SyncState::Synced);
}
#[test]
fn down_continuous_tier_is_missing() {
let obs = MirrorObservation {
running: false,
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.health, HealthState::Missing);
}
#[test]
fn running_not_ready_is_progressing() {
let obs = MirrorObservation {
running: true,
ready: false,
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.health, HealthState::Progressing);
}
#[test]
fn errored_is_degraded() {
let obs = MirrorObservation {
running: true,
ready: true,
errored: true,
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.health, HealthState::Degraded);
}
#[test]
fn diverging_live_revision_is_out_of_sync() {
let obs = MirrorObservation {
running: true,
ready: true,
live_revision: Some("caddy:2.7.6".into()),
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.declared_revision.as_deref(), Some("caddy:2.8.1"));
assert_eq!(cell.live_revision.as_deref(), Some("caddy:2.7.6"));
assert_eq!(cell.sync, SyncState::OutOfSync);
}
#[test]
fn matching_live_revision_is_synced() {
let obs = MirrorObservation {
running: true,
ready: true,
live_revision: Some("caddy:2.8.1".into()),
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.sync, SyncState::Synced);
}
#[test]
fn unknown_live_revision_does_not_force_out_of_sync() {
let obs = MirrorObservation {
running: true,
ready: true,
live_revision: None,
..Default::default()
};
let cell = compute_cell("prod", &cloudflare_mirror(), Some(&obs));
assert_eq!(cell.sync, SyncState::Synced);
}
#[test]
fn field_drift_is_detected_and_makes_out_of_sync() {
let mut live_fields = BTreeMap::new();
let mut static_slot = BTreeMap::new();
static_slot.insert("image".to_string(), "caddy:2.7.6".to_string());
static_slot.insert("port".to_string(), "8080".to_string()); live_fields.insert("static".to_string(), static_slot);
let obs = MirrorObservation {
running: true,
ready: true,
live_fields,
..Default::default()
};
let cell = compute_cell("sim", &sim_miniflare_mirror(), Some(&obs));
assert_eq!(cell.sync, SyncState::OutOfSync);
assert_eq!(cell.drift_count(), 1, "only the image field drifts");
assert_eq!(cell.drift[0].path, "providers.static.image");
assert_eq!(cell.drift[0].desired, "caddy:2.8.1");
assert_eq!(cell.drift[0].live, "caddy:2.7.6");
}
#[test]
fn provider_label_joins_inline_kinds() {
assert_eq!(
provider_label(&sim_miniflare_mirror()),
"minio-container + miniflare-container"
);
assert_eq!(provider_label(&cloudflare_mirror()), "cloudflare");
}
#[test]
fn declared_revision_none_when_no_version_field() {
assert_eq!(declared_revision(&local_static_mirror()), None);
}
#[test]
fn compute_service_and_summary_roll_up() {
let svc = ServiceWithMirrors {
service: ServiceConfig {
schema_version: 1,
name: "yah-dev".into(),
address: crate::config::ServiceAddress::front_door("yah.dev"),
description: None,
components: vec![],
},
mirrors: BTreeMap::from([
("dev".to_string(), local_static_mirror()),
("prod".to_string(), cloudflare_mirror()),
]),
component_transform_recipes: BTreeMap::new(),
passway_machines: BTreeMap::new(),
};
let mut obs = BTreeMap::new();
obs.insert(
"dev".to_string(),
MirrorObservation {
running: true,
ready: true,
..Default::default()
},
);
let status = compute_service(&svc, &obs);
assert_eq!(status.name, "yah-dev");
assert_eq!(status.cells.len(), 2);
assert_eq!(status.cells["dev"].sync, SyncState::Synced);
assert_eq!(status.cells["dev"].health, HealthState::Healthy);
assert_eq!(status.cells["prod"].sync, SyncState::Unknown);
assert_eq!(status.cells["prod"].health, HealthState::Missing);
let summary = summarize(&[status]);
assert_eq!(summary.synced, 1);
assert_eq!(summary.unknown, 1);
assert_eq!(summary.healthy, 1);
assert_eq!(summary.missing, 1);
}
#[test]
fn states_serialize_in_kebab_case_for_the_wire() {
assert_eq!(
serde_json::to_string(&SyncState::OutOfSync).unwrap(),
"\"out-of-sync\""
);
assert_eq!(
serde_json::to_string(&HealthState::Idle).unwrap(),
"\"idle\""
);
assert_eq!(
serde_json::to_string(&Runtime::Containers).unwrap(),
"\"containers\""
);
}
}