use std::collections::{BTreeMap, BTreeSet, VecDeque};
use std::net::{IpAddr, SocketAddr};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use serde::Serialize;
use serde_json::Value;
use super::changes::diff;
use super::{
LABEL_REV, LABEL_SERVICE, LABEL_SLOT, LABEL_STACK, StackDef, Store, instance_name, new_id,
now_secs, validate_stack_name,
};
use crate::balance::Balancer;
use crate::client::{Client, encode_query, encode_segment};
use crate::error::{Error, Result};
use crate::org::OrgId;
use crate::plan::{Desired, split_addr};
use crate::sandbox::{EnsureOptions, Sandbox};
use crate::secrets::Secrets;
use crate::spec::{
DependCondition, FailureAction, HealthProbe, PortBind, RestartCondition, RestartMode,
SandboxSpec, UpdateConfig, UpdateOrder,
};
use crate::supervise;
const DRAIN: Duration = Duration::from_secs(10);
const RESTARTS_BEFORE_REPLACE: u32 = 3;
const HEALTH_DEBOUNCE: Duration = Duration::from_secs(20);
#[derive(Debug, Clone, Serialize)]
pub struct InstanceStatus {
pub name: String,
pub slot: u32,
pub rev: String,
pub status: String,
pub health: String,
pub ip: Option<String>,
pub in_rotation: bool,
pub restarts: u32,
#[serde(skip_serializing_if = "String::is_empty")]
pub last_probe: String,
pub cpu_pct: Option<f32>,
pub cpu_history: Vec<f32>,
pub mem_bytes: Option<u64>,
pub disk_bytes: Option<u64>,
}
#[derive(Debug, Clone, Serialize)]
pub struct PortStatus {
pub listen: String,
pub target: u16,
pub backends: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
pub accepted: u64,
pub active: usize,
pub rate_history: Vec<f32>,
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct ServiceStatus {
pub service: String,
pub image: String,
pub rev: String,
pub replicas: u32,
pub running: u32,
pub healthy: u32,
pub state: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub message: Option<String>,
pub instances: Vec<InstanceStatus>,
pub ports: Vec<PortStatus>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rollout: Option<RolloutStatus>,
pub checked_at: u64,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub domains: Vec<crate::ingress::DomainStatus>,
}
pub trait Observer: Send + Sync {
fn rotation(&self, stack: &str, service: &str, ips: &[IpAddr]);
fn drain(&self, stack: &str, service: &str, ip: IpAddr, timeout: Duration);
fn stacks_changed(&self, defs: Vec<Arc<StackDef>>);
fn domains(&self, stack: &str, service: &str) -> Vec<crate::ingress::DomainStatus>;
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct RolloutStatus {
pub to_rev: String,
pub order: String,
pub parallelism: usize,
pub done: usize,
pub total: usize,
pub started_at: u64,
pub slots: Vec<SlotRollout>,
}
#[derive(Debug, Clone, Serialize, Default)]
pub struct SlotRollout {
pub slot: u32,
pub old: Option<String>,
pub old_rev: Option<String>,
pub old_state: String,
pub new: Option<String>,
pub new_state: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct Event {
pub seq: u64,
pub at: u64,
pub level: String,
pub stack: String,
#[serde(skip_serializing_if = "String::is_empty")]
pub service: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub instance: Option<String>,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub kind: Option<String>,
}
pub type EventSink = Arc<dyn Fn(&Event) + Send + Sync>;
const EVENTS_KEPT: usize = 1000;
#[derive(Debug, Clone, Default)]
pub struct Snapshot {
pub host: crate::metrics::HostSample,
pub instances: BTreeMap<String, crate::metrics::InstanceSample>,
pub at: u64,
}
pub fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[derive(Debug, Clone, Serialize)]
pub struct StackStatus {
pub name: String,
pub org: String,
pub deployed_at: u64,
pub deployed_by: String,
pub has_previous: bool,
pub converged: bool,
pub services: Vec<ServiceStatus>,
}
#[derive(Debug, Clone, Serialize)]
pub struct DeployChange {
pub service: String,
pub change: String,
pub rev: String,
pub replicas: u32,
}
struct Slot {
def: Arc<StackDef>,
remove: bool,
remove_volumes: bool,
}
struct WorkerShared {
slot: Mutex<Slot>,
wake: Condvar,
stop: AtomicBool,
kick: AtomicBool,
}
struct Inner {
client: Client,
store: Store,
balancer: Balancer,
interval: Duration,
workers: Mutex<BTreeMap<(String, String), Arc<WorkerShared>>>,
stacks: Mutex<BTreeMap<String, Arc<StackDef>>>,
status: Mutex<BTreeMap<(String, String), ServiceStatus>>,
failures: super::failure::Failures,
events: Mutex<(u64, VecDeque<Event>)>,
snapshot: Mutex<Snapshot>,
secrets: Arc<Secrets>,
edit: Mutex<()>,
refresh: Mutex<super::secrets::RefreshSchedule>,
observer: Option<Arc<dyn Observer>>,
event_sink: Mutex<Option<EventSink>>,
metrics_sink: Mutex<Option<std::sync::mpsc::SyncSender<crate::metrics_history::Sample>>>,
}
impl Inner {
fn emit(
&self,
level: &str,
stack: &str,
service: &str,
instance: Option<&str>,
message: String,
) {
self.emit_kind(None, level, stack, service, instance, message);
}
fn emit_kind(
&self,
kind: Option<&str>,
level: &str,
stack: &str,
service: &str,
instance: Option<&str>,
message: String,
) {
let mut e = self.events.lock().unwrap();
e.0 += 1;
let ev = Event {
seq: e.0,
at: now_ms(),
level: level.into(),
stack: stack.into(),
service: service.into(),
instance: instance.map(String::from),
message,
kind: kind.map(String::from),
};
if e.1.len() == EVENTS_KEPT {
e.1.pop_front();
}
let sink = self.event_sink.lock().unwrap().clone();
if let Some(s) = sink {
s(&ev);
}
e.1.push_back(ev);
}
}
const SAMPLE_EVERY: Duration = Duration::from_secs(2);
const SECRET_TICK: Duration = Duration::from_secs(10);
#[derive(Clone)]
pub struct Controller {
inner: Arc<Inner>,
}
impl Controller {
pub fn start(
client: Client,
store: Store,
interval: Duration,
secrets: Arc<Secrets>,
) -> Result<Controller> {
Controller::start_with(client, store, interval, secrets, None)
}
pub fn start_with(
client: Client,
store: Store,
interval: Duration,
secrets: Arc<Secrets>,
observer: Option<Arc<dyn Observer>>,
) -> Result<Controller> {
let c = Controller {
inner: Arc::new(Inner {
client,
store,
balancer: Balancer::new(),
interval,
workers: Mutex::new(BTreeMap::new()),
stacks: Mutex::new(BTreeMap::new()),
status: Mutex::new(BTreeMap::new()),
failures: Default::default(),
events: Mutex::new((0, VecDeque::new())),
snapshot: Mutex::new(Snapshot::default()),
secrets,
edit: Mutex::new(()),
refresh: Mutex::new(Default::default()),
observer,
event_sink: Mutex::new(None),
metrics_sink: Mutex::new(None),
}),
};
let weak = Arc::downgrade(&c.inner);
let _ = std::thread::Builder::new()
.name("isb-secrets".into())
.spawn(move || {
while let Some(inner) = weak.upgrade() {
let tick = inner.interval.clamp(Duration::from_secs(1), SECRET_TICK);
Controller { inner }.check_due_secrets();
std::thread::sleep(tick);
}
});
let weak = Arc::downgrade(&c.inner);
let _ = std::thread::Builder::new()
.name("isb-metrics".into())
.spawn(move || {
let mut sampler = crate::metrics::Sampler::new();
while let Some(inner) = weak.upgrade() {
match sampler.sample(&inner.client) {
Ok((host, insts)) => {
if let Some(tx) = &*inner.metrics_sink.lock().unwrap() {
crate::metrics_history::offer(tx, (now_ms(), insts.clone()));
}
*inner.snapshot.lock().unwrap() = Snapshot {
host,
instances: insts
.into_iter()
.map(|i| (format!("{}/{}", i.project, i.name), i))
.collect(),
at: now_ms(),
};
}
Err(e) => eprintln!("isb serve: metrics: {e}"),
}
drop(inner);
std::thread::sleep(SAMPLE_EVERY);
}
});
let defs = c.inner.store.load_all()?;
for def in &defs {
let dir = crate::discovery::org_dir(&def.org);
let keep: Vec<(String, String)> = defs
.iter()
.filter(|d| d.org == def.org)
.flat_map(|d| d.file.services.keys().map(|s| (d.name.clone(), s.clone())))
.collect();
crate::discovery::prune(&dir, &keep);
}
for def in defs {
eprintln!("isb serve: resuming stack {}", def.name);
c.apply(Arc::new(def));
}
for def in c.definitions() {
let keys: Vec<String> = def.secrets.keys().cloned().collect();
c.check_bindings(&def.qualified(), &keys, false);
}
Ok(c)
}
pub fn secret_changed(&self, org: &OrgId, name: &str) -> Vec<String> {
let mut rolled = Vec::new();
for def in self.definitions() {
if def.org != *org {
continue;
}
let keys: Vec<String> = def
.secrets
.iter()
.filter(|(_, b)| b.name == name)
.map(|(k, _)| k.clone())
.collect();
if !keys.is_empty() && self.check_bindings(&def.qualified(), &keys, false) {
rolled.push(def.qualified());
}
}
rolled
}
pub fn refresh_secret(
&self,
org: &OrgId,
name: &str,
) -> Result<crate::stack::secrets::Refreshed> {
let mut found = Vec::new();
let mut rolled = Vec::new();
for def in self.definitions() {
if def.org != *org {
continue;
}
let q = def.qualified();
let keys: Vec<String> = def
.secrets
.iter()
.filter(|(_, b)| b.name == name)
.map(|(k, _)| k.clone())
.collect();
for k in &keys {
let b = &def.secrets[k];
let v = self.inner.secrets.refresh_in(&b.driver, org, name)?;
found.push((b.driver.clone(), v));
let every = def
.file
.secrets
.get(k)
.map(crate::spec::SecretDef::refresh_interval)
.unwrap_or(crate::spec::DEFAULT_SECRET_REFRESH);
self.inner
.refresh
.lock()
.unwrap()
.reset(&q, k, every, Instant::now());
}
if !keys.is_empty() && self.check_bindings(&q, &keys, true) {
rolled.push(q);
}
}
Ok((found, rolled))
}
fn check_due_secrets(&self) {
let defs: Vec<(String, Arc<StackDef>)> = self
.inner
.stacks
.lock()
.unwrap()
.iter()
.map(|(q, d)| (q.clone(), d.clone()))
.collect();
let due = self
.inner
.refresh
.lock()
.unwrap()
.due(defs.iter().map(|(q, d)| (q.as_str(), &**d)), Instant::now());
let mut by_stack: BTreeMap<String, Vec<String>> = BTreeMap::new();
for (q, k) in due {
by_stack.entry(q).or_default().push(k);
}
for (q, keys) in by_stack {
self.check_bindings(&q, &keys, false);
}
}
fn check_bindings(&self, q: &str, keys: &[String], quiet: bool) -> bool {
let _g = self.inner.edit.lock().unwrap();
let Ok(cur) = self.get_def(q) else {
return false;
};
let mut def = (*cur).clone();
let mut moved = Vec::new();
for k in keys {
let Some(b) = def.secrets.get_mut(k) else {
continue;
};
match self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
Ok(v) if v != b.version => {
moved.push(format!("{} v{} -> v{v}", b.name, b.version));
b.version = v;
}
Ok(_) => {}
Err(e) if !quiet => self.note(
"warn",
q,
format!("secret {}: cannot check its version: {e}", b.name),
),
Err(_) => {}
}
}
if moved.is_empty() {
return false;
}
if let Err(e) = self.inner.store.save(&def) {
self.note("error", q, format!("cannot save new secret versions: {e}"));
return false;
}
self.apply(Arc::new(def));
self.note(
"info",
q,
format!("new secret version ({}): rolling", moved.join(", ")),
);
true
}
pub fn balancer(&self) -> &Balancer {
&self.inner.balancer
}
pub fn events(&self, since: u64, limit: usize) -> (u64, Vec<Event>) {
let e = self.inner.events.lock().unwrap();
let out: Vec<Event> = e.1.iter().filter(|x| x.seq > since).cloned().collect();
let skip = out.len().saturating_sub(limit);
(e.0, out.into_iter().skip(skip).collect())
}
pub fn wait_events(&self, since: u64, limit: usize, timeout: Duration) -> (u64, Vec<Event>) {
let started = Instant::now();
loop {
let r = self.events(since, limit);
if !r.1.is_empty() || started.elapsed() >= timeout {
return r;
}
std::thread::sleep(Duration::from_millis(250));
}
}
pub fn set_event_sink(&self, sink: EventSink) {
let e = self.inner.events.lock().unwrap();
for ev in &e.1 {
sink(ev);
}
*self.inner.event_sink.lock().unwrap() = Some(sink);
}
pub fn set_metrics_sink(
&self,
tx: std::sync::mpsc::SyncSender<crate::metrics_history::Sample>,
) {
*self.inner.metrics_sink.lock().unwrap() = Some(tx);
}
pub fn snapshot(&self) -> Snapshot {
self.inner.snapshot.lock().unwrap().clone()
}
pub fn event(&self, kind: &str, level: &str, stack: &str, service: &str, message: String) {
eprintln!("isb serve: {stack}: {message}");
self.inner
.emit_kind(Some(kind), level, stack, service, None, message);
}
pub fn relay(
&self,
kind: Option<&str>,
level: &str,
stack: &str,
service: &str,
instance: Option<&str>,
message: String,
) {
self.inner
.emit_kind(kind, level, stack, service, instance, message);
}
pub fn note(&self, level: &str, stack: &str, message: String) {
eprintln!("isb serve: {stack}: {message}");
self.inner.emit(level, stack, "", None, message);
}
pub fn note_service(&self, level: &str, stack: &str, service: &str, message: String) {
self.inner.emit(level, stack, service, None, message);
}
pub fn service_event(&self, level: &str, stack: &str, service: &str, message: String) {
eprintln!("isb serve: {stack}/{service}: {message}");
self.inner.emit(level, stack, service, None, message);
}
fn notify_stacks(&self) {
if let Some(o) = &self.inner.observer {
o.stacks_changed(self.definitions());
}
}
pub fn plan(&self, def: &StackDef) -> Result<Vec<DeployChange>> {
self.validate(def)?;
let mut def = def.clone();
self.pin_images(&mut def)?;
let def = &def;
let old = self
.inner
.stacks
.lock()
.unwrap()
.get(&def.qualified())
.cloned();
diff(old.as_deref(), def)
}
pub fn client(&self) -> &Client {
&self.inner.client
}
pub fn validate(&self, def: &StackDef) -> Result<()> {
validate_stack_name(&def.name)?;
let host = crate::sandbox::host_facts(&crate::org::client(&self.inner.client, &def.org))?;
for (svc, spec) in &def.file.services {
let mut s = instance_spec(def, svc, spec, 1, "0000")?;
s.name = Some(instance_name(&def.name, svc, 1, "0000")?);
crate::plan::resolve(&s, &def.file.volumes, &host, &def.base_dir)?;
published(spec)?;
crate::ingress::domain::validate(svc, &spec.domains)?;
}
Ok(())
}
pub fn deploy(&self, mut def: StackDef) -> Result<Vec<DeployChange>> {
self.validate(&def)?;
self.pin_images(&mut def)?;
let _g = self.inner.edit.lock().unwrap();
let old = self
.inner
.stacks
.lock()
.unwrap()
.get(&def.qualified())
.cloned();
if let Some(old) = &old {
let mut prev = (**old).clone();
prev.previous = None;
def.previous = Some(Box::new(prev));
for (k, v) in &old.force {
def.force.entry(k.clone()).or_insert(*v);
}
}
let changes = diff(old.as_deref(), &def)?;
self.inner.store.save(&def)?;
self.apply(Arc::new(def));
Ok(changes)
}
fn apply(&self, def: Arc<StackDef>) {
let name = def.qualified();
self.inner
.stacks
.lock()
.unwrap()
.insert(name.clone(), def.clone());
let mut workers = self.inner.workers.lock().unwrap();
for svc in def.file.services.keys() {
let key = (name.clone(), svc.clone());
match workers.get(&key) {
Some(w) => {
let mut slot = w.slot.lock().unwrap();
slot.def = def.clone();
slot.remove = false;
if let Some(st) = self.inner.status.lock().unwrap().get_mut(&key) {
if st.state == "failing" {
st.state = "updating".into();
st.message = None;
}
}
w.wake.notify_all();
}
None => {
let shared = Arc::new(WorkerShared {
slot: Mutex::new(Slot {
def: def.clone(),
remove: false,
remove_volumes: false,
}),
wake: Condvar::new(),
stop: AtomicBool::new(false),
kick: AtomicBool::new(false),
});
workers.insert(key, shared.clone());
spawn_worker(self.inner.clone(), &def, svc.clone(), shared);
}
}
}
for ((stack, svc), w) in workers.iter() {
if *stack == name && !def.file.services.contains_key(svc) {
let mut slot = w.slot.lock().unwrap();
slot.remove = true;
w.wake.notify_all();
}
}
drop(workers);
self.notify_stacks();
}
pub fn org_limits_changed(&self, org: &OrgId) -> usize {
let stacks = self.inner.stacks.lock().unwrap();
let workers = self.inner.workers.lock().unwrap();
let status = self.inner.status.lock().unwrap();
let mut n = 0;
for (key, w) in workers.iter() {
let ours = stacks.get(&key.0).is_some_and(|d| d.org == *org);
let limited = status.get(key).is_some_and(|s| {
s.state == "failing" && s.message.as_deref().is_some_and(limit_error)
});
if ours && limited {
w.kick.store(true, Ordering::SeqCst);
let _slot = w.slot.lock().unwrap();
w.wake.notify_all();
n += 1;
}
}
n
}
pub fn remove(&self, name: &str, volumes: bool, timeout: Duration) -> Result<()> {
{
let _g = self.inner.edit.lock().unwrap();
let Some(def) = self.inner.stacks.lock().unwrap().remove(name) else {
return Err(Error::NotFound(format!("stack {name}")));
};
self.inner.store.remove(&def.org, &def.name)?;
}
self.notify_stacks();
let ws: Vec<Arc<WorkerShared>> = self
.inner
.workers
.lock()
.unwrap()
.iter()
.filter(|((s, _), _)| s == name)
.map(|(_, w)| w.clone())
.collect();
for w in &ws {
let mut slot = w.slot.lock().unwrap();
slot.remove = true;
slot.remove_volumes = volumes;
w.wake.notify_all();
}
let started = Instant::now();
while started.elapsed() < timeout {
let left = self
.inner
.workers
.lock()
.unwrap()
.keys()
.any(|(s, _)| s == name);
if !left {
return Ok(());
}
std::thread::sleep(Duration::from_millis(200));
}
Err(Error::invalid(format!(
"stack {name}: still removing after {timeout:?}; it carries on in the background"
)))
}
pub fn rollback(&self, name: &str) -> Result<Vec<DeployChange>> {
let _g = self.inner.edit.lock().unwrap();
let cur = self.get_def(name)?;
let prev = cur
.previous
.clone()
.ok_or_else(|| Error::invalid(format!("stack {name} has no previous deployment")))?;
let mut def = *prev;
def.deployed_at = now_secs();
for b in def.secrets.values_mut() {
if let Ok(v) = self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
b.version = v;
}
}
let mut cur2 = (*cur).clone();
cur2.previous = None;
let changes = diff(Some(&cur), &def)?;
def.previous = Some(Box::new(cur2));
self.inner.store.save(&def)?;
self.apply(Arc::new(def));
Ok(changes)
}
pub fn scale(&self, name: &str, service: &str, replicas: u32) -> Result<()> {
let _g = self.inner.edit.lock().unwrap();
let cur = self.get_def(name)?;
let mut def = (*cur).clone();
let spec = def
.file
.services
.get_mut(service)
.ok_or_else(|| Error::NotFound(format!("service {service} in stack {name}")))?;
spec.deploy.get_or_insert_with(Default::default).replicas = Some(replicas);
self.inner.store.save(&def)?;
self.apply(Arc::new(def));
Ok(())
}
pub fn redeploy(&self, name: &str, service: &str) -> Result<()> {
let _g = self.inner.edit.lock().unwrap();
let cur = self.get_def(name)?;
cur.service(service)?;
let mut def = (*cur).clone();
*def.force.entry(service.to_string()).or_insert(0) += 1;
self.pin_images(&mut def)?;
self.inner.store.save(&def)?;
self.apply(Arc::new(def));
Ok(())
}
fn pin_images(&self, def: &mut StackDef) -> Result<()> {
def.images.clear();
for (svc, spec) in &def.file.services {
let Some(r) = spec.image.strip_prefix("registry:") else {
continue;
};
let r = crate::registry::ImageRef::parse(r)?;
let reg = crate::registry::Registry::shared(&self.inner.client)?;
let d = reg.resolve(&def.org, &r)?;
if r.digest.is_none() {
def.images.insert(svc.clone(), d);
}
}
Ok(())
}
fn get_def(&self, name: &str) -> Result<Arc<StackDef>> {
self.inner
.stacks
.lock()
.unwrap()
.get(name)
.cloned()
.ok_or_else(|| Error::NotFound(format!("stack {name}")))
}
pub fn definitions(&self) -> Vec<Arc<StackDef>> {
self.inner
.stacks
.lock()
.unwrap()
.values()
.cloned()
.collect()
}
pub fn definition(&self, name: &str) -> Result<StackDef> {
self.get_def(name).map(|d| (*d).clone())
}
pub fn list(&self) -> Vec<StackStatus> {
let names: Vec<String> = self.inner.stacks.lock().unwrap().keys().cloned().collect();
names.iter().filter_map(|n| self.status(n).ok()).collect()
}
pub fn status(&self, name: &str) -> Result<StackStatus> {
let def = self.get_def(name)?;
let st = self.inner.status.lock().unwrap();
let services: Vec<ServiceStatus> = def
.file
.services
.keys()
.map(|svc| {
st.get(&(name.to_string(), svc.clone()))
.cloned()
.unwrap_or_else(|| ServiceStatus {
service: svc.clone(),
state: "starting".into(),
..Default::default()
})
})
.collect();
drop(st);
let mut services = services;
if let Some(o) = &self.inner.observer {
for s in &mut services {
s.domains = o.domains(name, &s.service);
}
}
let converged = services.iter().all(|s| s.state == "converged");
Ok(StackStatus {
name: def.name.clone(),
org: def.org.to_string(),
deployed_at: def.deployed_at,
deployed_by: def.deployed_by.clone(),
has_previous: def.previous.is_some(),
converged,
services,
})
}
pub fn logs(
&self,
name: &str,
service: &str,
slot: Option<u32>,
lines: usize,
) -> Result<BTreeMap<String, String>> {
let def = self.get_def(name)?;
let oci = crate::plan::ImageSource::parse(&def.service(service)?.image)?.is_oci();
let oc = crate::org::client(&self.inner.client, &def.org);
super::failure::replica_logs(&oc, &def.name, service, oci, slot, lines)
}
pub fn last_failure(&self, name: &str, service: &str) -> Option<super::failure::FailedAttempt> {
let state = self.inner.status.lock().unwrap();
let converged = state
.get(&(name.to_string(), service.to_string()))
.is_some_and(|s| s.state == "converged");
drop(state);
if converged {
return None;
}
self.inner.failures.last(name, service)
}
pub fn shutdown(&self) {
for w in self.inner.workers.lock().unwrap().values() {
w.stop.store(true, Ordering::SeqCst);
w.wake.notify_all();
}
self.inner.balancer.clear();
}
}
fn limit_error(msg: &str) -> bool {
msg.contains(" quota (") || msg.contains(" limit (")
}
fn instance_spec(
def: &StackDef,
service: &str,
spec: &SandboxSpec,
slot: u32,
rev: &str,
) -> Result<SandboxSpec> {
let mut s = spec.clone();
s.image = def.instance_image(service, &spec.image);
s.restart = Some(RestartMode::Always);
s.ports.retain(|p| p.bind == PortBind::Guest);
s.domains.clear();
if let Some(d) = &s.deploy {
s.labels.extend(d.labels.clone());
}
s.labels.insert(LABEL_STACK.into(), def.name.clone());
s.labels.insert(LABEL_SERVICE.into(), service.into());
s.labels.insert(LABEL_SLOT.into(), slot.to_string());
s.labels.insert(LABEL_REV.into(), rev.into());
Ok(s)
}
#[derive(Debug, Clone, PartialEq)]
struct Published {
listen: SocketAddr,
target: u16,
}
fn published(spec: &SandboxSpec) -> Result<Vec<Published>> {
let mut out = Vec::new();
for p in &spec.ports {
if p.bind == PortBind::Guest {
continue;
}
let listen = crate::plan::normalize_addr(&p.listen, "127.0.0.1").map_err(Error::invalid)?;
let connect =
crate::plan::normalize_addr(&p.connect, "127.0.0.1").map_err(Error::invalid)?;
let (lp, lh, lport) = split_addr(&listen).ok_or_else(|| {
Error::invalid(format!(
"port {listen}: a stack publishes single tcp ports (no ranges)"
))
})?;
let (_, _, cport) = split_addr(&connect).ok_or_else(|| {
Error::invalid(format!(
"port {connect}: a stack publishes single tcp ports (no ranges)"
))
})?;
if lp != "tcp" {
return Err(Error::invalid(format!(
"port {listen}: the stack balancer is tcp only"
)));
}
if p.search.is_some() {
return Err(Error::invalid(
"a stack's published ports are fixed; port search is for isb up",
));
}
let host: IpAddr = lh
.trim_start_matches('[')
.trim_end_matches(']')
.parse()
.map_err(|_| {
Error::invalid(format!("port {listen}: the host must be an IP address"))
})?;
out.push(Published {
listen: SocketAddr::new(host, lport),
target: cport,
});
}
Ok(out)
}
#[derive(Debug, Clone)]
#[doc(hidden)]
pub struct Inst {
pub name: String,
pub(super) slot: u32,
pub rev: String,
status: String,
}
impl Inst {
#[doc(hidden)]
pub fn is_running(&self) -> bool {
self.running()
}
fn running(&self) -> bool {
self.status.eq_ignore_ascii_case("running")
}
}
#[doc(hidden)]
pub fn list_instances(client: &Client, stack: &str, service: Option<&str>) -> Result<Vec<Inst>> {
let mut filter = format!("config.user.{LABEL_STACK} eq {stack}");
if let Some(s) = service {
filter.push_str(&format!(" and config.user.{LABEL_SERVICE} eq {s}"));
}
let v = client.get(&format!(
"/1.0/instances?recursion=1&filter={}",
encode_query(&filter)
))?;
let mut out = Vec::new();
for i in v.as_array().into_iter().flatten() {
let info = crate::sandbox::SandboxInfo::from_api(i);
let c = &info.config;
if c.get(&format!("user.{LABEL_STACK}")).map(String::as_str) != Some(stack) {
continue;
}
let svc = c
.get(&format!("user.{LABEL_SERVICE}"))
.cloned()
.unwrap_or_default();
if service.is_some_and(|s| s != svc) {
continue;
}
out.push(Inst {
name: info.name.clone(),
slot: c
.get(&format!("user.{LABEL_SLOT}"))
.and_then(|s| s.parse().ok())
.unwrap_or(0),
rev: c
.get(&format!("user.{LABEL_REV}"))
.cloned()
.unwrap_or_default(),
status: info.status.clone(),
});
}
out.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
Ok(out)
}
fn instance_state(client: &Client, name: &str) -> Result<(i64, Option<IpAddr>)> {
let v = client.get(&format!("/1.0/instances/{}/state", encode_segment(name)))?;
let pid = v.get("pid").and_then(Value::as_i64).unwrap_or(0);
let mut v4 = None;
let mut v6 = None;
if let Some(nets) = v.get("network").and_then(Value::as_object) {
for (ifname, n) in nets {
if ifname == "lo" {
continue;
}
for a in n
.get("addresses")
.and_then(Value::as_array)
.into_iter()
.flatten()
{
if a.get("scope").and_then(Value::as_str) != Some("global") {
continue;
}
let Some(ip) = a
.get("address")
.and_then(Value::as_str)
.and_then(|s| s.parse::<IpAddr>().ok())
else {
continue;
};
match ip {
IpAddr::V4(_) if v4.is_none() => v4 = Some(ip),
IpAddr::V6(_) if v6.is_none() => v6 = Some(ip),
_ => {}
}
}
}
}
Ok((pid, v4.or(v6)))
}
#[derive(Debug, Default)]
struct InstRt {
pid: i64,
since: Option<Instant>,
ip: Option<IpAddr>,
failures: u32,
healthy: Option<bool>,
next_probe: Option<Instant>,
last_probe: String,
unhealthy_restarts: u32,
restarts: VecDeque<Instant>,
in_rotation: bool,
}
struct Worker {
inner: Arc<Inner>,
stack: String,
q: String,
oclient: Client,
service: String,
shared: Arc<WorkerShared>,
rt: BTreeMap<String, InstRt>,
paused: Option<(String, String)>,
create_backoff: Option<(Instant, Duration)>,
routes: BTreeMap<String, Published>,
route_errors: BTreeMap<String, String>,
template: Option<(String, Desired)>,
state: String,
message: Option<String>,
last_error: Option<String>,
insts: Vec<Inst>,
rollout: Option<RolloutStatus>,
rates: BTreeMap<String, (u64, Instant, VecDeque<f32>)>,
deps_met: bool,
org: OrgId,
dns_last: Option<Vec<IpAddr>>,
dns_error: Option<String>,
observed: Option<Vec<IpAddr>>,
ever_healthy: bool,
health_down: Option<Instant>,
health_alarm: bool,
seen: Option<Arc<StackDef>>,
}
fn spawn_worker(inner: Arc<Inner>, def: &StackDef, service: String, shared: Arc<WorkerShared>) {
let (stack, q, org) = (def.name.clone(), def.qualified(), def.org.clone());
let oclient = crate::org::client(&inner.client, &def.org);
let name = format!("isb-{q}-{service}");
let r = std::thread::Builder::new().name(name).spawn(move || {
let mut w = Worker::new(inner, stack, q, oclient, service, shared, org);
w.run();
});
if let Err(e) = r {
eprintln!("isb serve: cannot start a worker thread: {e}");
}
}
impl Worker {
fn new(
inner: Arc<Inner>,
stack: String,
q: String,
oclient: Client,
service: String,
shared: Arc<WorkerShared>,
org: OrgId,
) -> Worker {
Worker {
inner,
stack,
q,
oclient,
service,
shared,
rt: BTreeMap::new(),
paused: None,
create_backoff: None,
routes: BTreeMap::new(),
route_errors: BTreeMap::new(),
template: None,
state: "starting".into(),
message: None,
last_error: None,
insts: Vec::new(),
rollout: None,
rates: BTreeMap::new(),
deps_met: false,
org,
dns_last: None,
dns_error: None,
observed: None,
ever_healthy: false,
health_down: None,
health_alarm: false,
seen: None,
}
}
fn begin_pass(&mut self, def: &Arc<StackDef>) {
let new_def = self.seen.as_ref().is_none_or(|d| !Arc::ptr_eq(d, def));
let kicked = self.shared.kick.swap(false, Ordering::SeqCst);
if new_def {
self.seen = Some(def.clone());
self.last_error = None;
if self.state == "failing" {
self.state = "updating".into();
self.message = None;
}
}
if new_def || kicked {
self.create_backoff = None;
}
}
}
impl Worker {
fn log(&self, msg: &str) {
self.event("info", None, msg);
}
fn event(&self, level: &str, instance: Option<&str>, msg: &str) {
eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
self.inner
.emit(level, &self.q, &self.service, instance, msg.to_string());
}
fn client(&self) -> &Client {
&self.oclient
}
fn watch_health(&mut self, healthy: u32, replicas: u32) {
if healthy > 0 {
self.ever_healthy = true;
self.health_down = None;
if self.health_alarm {
self.health_alarm = false;
let msg = format!("{healthy} of {replicas} replicas healthy again");
eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
self.inner.emit_kind(
Some("health.recovered"),
"info",
&self.q,
&self.service,
None,
msg,
);
}
return;
}
if replicas == 0 || !self.ever_healthy || self.rollout.is_some() {
self.health_down = None;
return;
}
let since = *self.health_down.get_or_insert_with(Instant::now);
if !self.health_alarm && since.elapsed() >= HEALTH_DEBOUNCE {
self.health_alarm = true;
let msg = format!(
"no healthy replica (of {replicas}) for {}s",
since.elapsed().as_secs()
);
eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
self.inner.emit_kind(
Some("health.unhealthy"),
"error",
&self.q,
&self.service,
None,
msg,
);
}
}
fn key(&self) -> (String, String) {
(self.q.clone(), self.service.clone())
}
fn run(&mut self) {
loop {
if self.shared.stop.load(Ordering::SeqCst) {
return;
}
let (def, remove, remove_volumes) = {
let s = self.shared.slot.lock().unwrap();
(s.def.clone(), s.remove, s.remove_volumes)
};
if remove {
self.teardown(&def, remove_volumes);
let mut ws = self.inner.workers.lock().unwrap();
let slot = self.shared.slot.lock().unwrap();
if slot.remove || self.shared.stop.load(Ordering::SeqCst) {
ws.remove(&self.key());
self.inner.status.lock().unwrap().remove(&self.key());
return;
}
continue;
}
self.begin_pass(&def);
if let Err(e) = self.pass(&def) {
self.state = "failing".into();
self.message = Some(e.to_string());
if self.last_error.as_deref() != Some(&e.to_string()) {
self.event("error", None, &e.to_string());
self.last_error = Some(e.to_string());
}
self.publish_status(&def);
let slot = self.shared.slot.lock().unwrap();
if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
}
continue;
}
self.last_error = None;
let slot = self.shared.slot.lock().unwrap();
if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
}
}
}
fn superseded(&self, def: &Arc<StackDef>) -> bool {
let s = self.shared.slot.lock().unwrap();
s.remove || !Arc::ptr_eq(&s.def, def) || self.shared.stop.load(Ordering::SeqCst)
}
fn teardown(&mut self, def: &StackDef, volumes: bool) {
if let Some(dir) = Some(crate::discovery::org_dir(&self.org)).filter(|d| d.is_dir()) {
if let Err(e) =
crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &[])
{
self.log(&format!("cannot remove the service name: {e}"));
}
}
self.dns_last = Some(Vec::new());
if let Some(o) = &self.inner.observer {
o.rotation(&self.q, &self.service, &[]);
}
self.observed = Some(Vec::new());
for (k, _) in std::mem::take(&mut self.routes) {
self.inner.balancer.remove_route(&k);
}
match list_instances(self.client(), &self.stack, Some(&self.service)) {
Ok(insts) => {
for i in insts {
self.log(&format!("removing {}", i.name));
if let Err(e) = Sandbox::remove(self.client(), &i.name, true) {
if !e.is_not_found() {
self.log(&format!("cannot remove {}: {e}", i.name));
}
}
}
}
Err(e) => self.log(&format!("cannot list instances to remove: {e}")),
}
if volumes {
self.remove_volumes(def);
}
self.rt.clear();
self.template = None;
}
fn remove_volumes(&self, def: &StackDef) {
let Ok(spec) = def.service(&self.service) else {
return;
};
let Ok(host) = crate::sandbox::host_facts(self.client()) else {
return;
};
let Ok(pool) = host.pick_pool(spec.storage.as_deref()) else {
return;
};
for v in &spec.volumes {
if v.mount_type != crate::spec::MountType::Volume {
continue;
}
let d = def.file.volumes.get(&v.source);
if v.external || d.is_some_and(|d| d.external) {
continue;
}
let name = d
.and_then(|d| d.name.clone())
.unwrap_or_else(|| v.source.clone());
let vpool = match v.pool.as_deref().or(d.and_then(|d| d.pool.as_deref())) {
Some(p) if p != "auto" => p.to_string(),
_ => pool.clone(),
};
match crate::volume::remove(self.client(), &vpool, &name) {
Ok(()) => self.log(&format!("volume {name}: deleted")),
Err(e) if e.is_not_found() => {}
Err(e) => self.log(&format!("volume {name}: kept ({e})")),
}
}
}
#[expect(
clippy::too_many_lines,
reason = "predates the lint ratchet; split it when next changed"
)]
fn pass(&mut self, def: &Arc<StackDef>) -> Result<()> {
let spec = def.service(&self.service)?.clone();
let rev = def.revision(&self.service)?;
let replicas = spec.replicas();
let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
let probe = spec.health_probe().map_err(Error::invalid)?;
if !self.deps_met {
if let Some(msg) = self.waiting_for(&spec) {
self.state = "waiting".into();
self.message = Some(msg);
self.publish_status(def);
return Ok(());
}
self.deps_met = true;
}
if self.template.as_ref().is_none_or(|(r, _)| *r != rev) {
let mut s = instance_spec(def, &self.service, &spec, 1, &rev)?;
s.name = Some(instance_name(&self.stack, &self.service, 1, "0000")?);
let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
self.template = Some((rev.clone(), d));
}
self.set_routes(&spec);
let mut insts = list_instances(self.client(), &self.stack, Some(&self.service))?;
self.rt.retain(|n, _| insts.iter().any(|i| i.name == *n));
let extra: Vec<Inst> = insts
.iter()
.filter(|i| i.slot > replicas || i.slot == 0)
.cloned()
.collect();
for i in extra.iter().rev() {
self.log(&format!("scaling down: removing {}", i.name));
self.retire(&i.name)?;
}
insts.retain(|i| i.slot >= 1 && i.slot <= replicas);
self.insts = insts.clone();
for i in &insts {
self.maintain(def, i, &spec, oci, probe.as_ref())?;
}
for slot in 1..=replicas {
let current_ok = insts.iter().any(|i| {
i.slot == slot
&& i.rev == rev
&& self.rt.get(&i.name).is_some_and(|r| r.in_rotation)
});
if current_ok {
for i in insts.iter().filter(|i| i.slot == slot && i.rev != rev) {
self.log(&format!("removing leftover {}", i.name));
self.retire(&i.name)?;
}
}
let mut current: Vec<&Inst> = insts
.iter()
.filter(|i| i.slot == slot && i.rev == rev)
.collect();
while current.len() > 1 {
let i = current.pop().unwrap();
self.log(&format!("removing duplicate {}", i.name));
self.retire(&i.name)?;
}
}
self.sync_routes();
self.publish_status(def);
let mut pending: Vec<(u32, Option<String>)> = Vec::new();
for slot in 1..=replicas {
if insts.iter().any(|i| i.slot == slot && i.rev == rev) {
continue;
}
let old = insts
.iter()
.find(|i| i.slot == slot)
.map(|i| i.name.clone());
pending.push((slot, old));
}
if pending.is_empty() {
self.create_backoff = None;
let all_ok = insts
.iter()
.all(|i| i.rev == rev && self.rt.get(&i.name).is_some_and(|r| r.in_rotation));
self.state = if all_ok { "converged" } else { "failing" }.into();
if all_ok {
self.message = None;
} else if self.message.is_none() {
self.message = Some("some replicas are not healthy".into());
}
self.publish_status(def);
return Ok(());
}
if self.paused.as_ref().is_some_and(|(r, _)| *r == rev) {
self.state = "paused".into();
self.message = self.paused.as_ref().map(|(_, m)| m.clone());
self.publish_status(def);
return Ok(());
}
if let Some((at, wait)) = self.create_backoff {
if at.elapsed() < wait {
return Ok(());
}
}
self.state = "updating".into();
self.message = None;
self.publish_status(def);
let uc: UpdateConfig = spec
.deploy
.as_ref()
.and_then(|d| d.update_config.clone())
.unwrap_or_default();
let parallel = match uc.parallelism.unwrap_or(1) {
0 => pending.len(),
n => n as usize,
};
let delay = uc
.delay
.as_deref()
.map(crate::flex::parse_duration)
.transpose()
.map_err(Error::invalid)?
.unwrap_or_default();
let monitor = uc
.monitor
.as_deref()
.map(crate::flex::parse_duration)
.transpose()
.map_err(Error::invalid)?
.unwrap_or(Duration::from_secs(5));
let order = uc.order.unwrap_or_default();
let order_name = match order {
UpdateOrder::StopFirst => "stop-first",
UpdateOrder::StartFirst => "start-first",
};
self.rollout = Some(RolloutStatus {
to_rev: rev.clone(),
order: order_name.into(),
parallelism: parallel,
done: 0,
total: pending.len(),
started_at: now_secs(),
slots: pending
.iter()
.map(|(slot, old)| SlotRollout {
slot: *slot,
old: old.clone(),
old_rev: old
.as_ref()
.and_then(|o| insts.iter().find(|i| i.name == *o))
.map(|i| i.rev.clone()),
old_state: if old.is_some() { "serving" } else { "none" }.into(),
new: None,
new_state: "waiting".into(),
})
.collect(),
});
let rollout_started = Instant::now();
self.log(&format!(
"rolling out rev {rev} to {} slot(s), {order_name}",
pending.len()
));
self.publish_status(def);
let r = self.roll(
def,
&spec,
&rev,
&pending,
parallel,
delay,
monitor,
order,
oci,
probe.as_ref(),
&uc,
);
let rollout = self.rollout.take();
if let (Ok(true), Some(ro)) = (&r, rollout) {
self.log(&format!(
"rollout of rev {rev} complete: {}/{} slot(s) in {:.0?}",
ro.done,
ro.total,
rollout_started.elapsed()
));
}
self.publish_status(def);
r.map(|_| ())
}
#[expect(clippy::too_many_arguments)]
#[expect(
clippy::excessive_nesting,
reason = "predates the lint ratchet; split it when next changed"
)]
fn roll(
&mut self,
def: &Arc<StackDef>,
spec: &SandboxSpec,
rev: &String,
pending: &[(u32, Option<String>)],
parallel: usize,
delay: Duration,
monitor: Duration,
order: UpdateOrder,
oci: bool,
probe: Option<&HealthProbe>,
uc: &UpdateConfig,
) -> Result<bool> {
for (n, batch) in pending.chunks(parallel).enumerate() {
if n > 0 && !delay.is_zero() {
std::thread::sleep(delay);
}
if self.superseded(def) {
return Ok(false);
}
for (slot, old) in batch {
let r = self.replace(
def,
spec,
rev,
*slot,
old.as_deref(),
order,
oci,
probe,
monitor,
);
if let Err(e) = r {
let msg = format!("slot {slot}: {e}");
self.event("error", None, &msg);
if old.is_none() {
let wait = self
.create_backoff
.map(|(_, w)| (w * 2).min(Duration::from_secs(300)))
.unwrap_or(Duration::from_secs(10));
self.create_backoff = Some((Instant::now(), wait));
self.state = "failing".into();
self.message = Some(format!("{msg}; retrying in {wait:?}"));
self.publish_status(def);
return Ok(false);
}
match uc.failure_action.unwrap_or_default() {
FailureAction::Continue => continue,
FailureAction::Pause => {
self.event("warn", None, &format!("rollout of rev {rev} paused"));
self.paused = Some((rev.clone(), format!("rollout paused: {msg}")));
return Ok(false);
}
FailureAction::Rollback => {
self.paused = Some((rev.clone(), format!("rolled back: {msg}")));
let ctl = Controller {
inner: self.inner.clone(),
};
self.event(
"warn",
None,
&format!("rollout of rev {rev} failed; rolling back"),
);
if let Err(e) = ctl.rollback(&self.q) {
self.event("error", None, &format!("rollback failed: {e}"));
}
return Ok(false);
}
}
}
}
}
Ok(true)
}
fn slot_state(
&mut self,
def: &StackDef,
slot: u32,
old_state: Option<&str>,
new: Option<&str>,
new_state: Option<&str>,
) {
if let Some(ro) = &mut self.rollout {
if let Some(s) = ro.slots.iter_mut().find(|s| s.slot == slot) {
if let Some(o) = old_state {
s.old_state = o.into();
}
if let Some(n) = new {
s.new = Some(n.into());
}
if let Some(n) = new_state {
s.new_state = n.into();
if n == "serving" {
ro.done += 1;
}
}
}
}
self.publish_status(def);
}
fn waiting_for(&self, spec: &SandboxSpec) -> Option<String> {
let st = self.inner.status.lock().unwrap();
for (dep, d) in &spec.depends_on {
let s = st.get(&(self.q.clone(), dep.clone()));
let ok = match d.condition {
DependCondition::ServiceStarted => s.is_some_and(|s| s.running > 0),
DependCondition::ServiceHealthy => s.is_some_and(|s| s.healthy > 0),
};
if !ok {
return Some(format!("waiting for {dep} ({:?})", d.condition));
}
}
None
}
fn handle(&self, name: &str) -> Sandbox {
let (_, d) = self.template.as_ref().expect("resolved in pass");
Sandbox::like(self.client(), name, d)
}
#[expect(
clippy::too_many_lines,
reason = "predates the lint ratchet; split it when next changed"
)]
fn maintain(
&mut self,
def: &StackDef,
i: &Inst,
spec: &SandboxSpec,
oci: bool,
probe: Option<&HealthProbe>,
) -> Result<()> {
let policy = spec.deploy.as_ref().and_then(|d| d.restart_policy.clone());
let condition = policy
.as_ref()
.and_then(|p| p.condition)
.unwrap_or_default();
let sb = self.handle(&i.name);
if !i.running() {
self.set_rotation(&i.name, false);
if condition == RestartCondition::None {
return Ok(());
}
if self.restart_budget_spent(&i.name, policy.as_ref()) {
self.message = Some(format!("{}: restart limit reached", i.name));
return Ok(());
}
self.event(
"warn",
Some(&i.name),
&format!("{} is {}; starting it", i.name, i.status),
);
self.count_restart(&i.name);
if let Err(e) = sb.start() {
self.log(&format!("cannot start {}: {e}", i.name));
return Ok(());
}
}
let (pid, ip) = match instance_state(self.client(), &i.name) {
Ok(s) => s,
Err(e) if e.is_not_found() => return Ok(()),
Err(e) => return Err(e),
};
let rt = self.rt.entry(i.name.clone()).or_default();
rt.ip = ip;
if rt.pid != pid {
rt.pid = pid;
rt.since = Some(Instant::now());
rt.failures = 0;
rt.healthy = None;
rt.next_probe = None;
let r = self.setup(def, &sb, spec, oci);
if let Err(e) = r {
if let Some(rt) = self.rt.get_mut(&i.name) {
rt.pid = 0;
}
self.set_rotation(&i.name, false);
return Err(Error::OperationFailed {
step: format!("set up {}", i.name),
message: e.to_string(),
});
}
}
let alive = self.alive(&sb, spec, oci);
let healthy = match probe {
None => alive,
Some(p) => {
let rt = self.rt.get_mut(&i.name).unwrap();
let since = rt.since.unwrap_or_else(Instant::now);
let in_start = since.elapsed() < p.start_period;
if rt.next_probe.is_none_or(|t| Instant::now() >= t) {
let r = supervise::probe(&sb, p);
let rt = self.rt.get_mut(&i.name).unwrap();
rt.last_probe = r.output.clone();
if r.ok {
rt.failures = 0;
rt.healthy = Some(true);
rt.unhealthy_restarts = 0;
} else if !in_start {
rt.failures += 1;
if rt.failures >= p.retries {
rt.healthy = Some(false);
}
}
let wait = if rt.healthy.is_none() {
p.start_interval
} else {
p.interval
};
rt.next_probe = Some(Instant::now() + wait);
}
self.rt[&i.name].healthy == Some(true) && alive
}
};
self.set_rotation(&i.name, healthy);
let unhealthy = self.rt[&i.name].healthy == Some(false);
if unhealthy && condition != RestartCondition::None {
if self.restart_budget_spent(&i.name, policy.as_ref()) {
self.message = Some(format!("{}: unhealthy, restart limit reached", i.name));
return Ok(());
}
let n = self.rt[&i.name].unhealthy_restarts;
if n >= RESTARTS_BEFORE_REPLACE {
self.log(&format!(
"{} stayed unhealthy through {n} restarts; replacing it",
i.name
));
let rt = self.rt.get_mut(&i.name).unwrap();
rt.unhealthy_restarts = 0;
drop(sb);
self.retire(&i.name)?;
return Ok(());
}
self.event(
"warn",
Some(&i.name),
&format!(
"{} is unhealthy ({}); restarting its app",
i.name,
self.rt[&i.name]
.last_probe
.lines()
.last()
.unwrap_or("probe failed")
),
);
self.count_restart(&i.name);
let rt = self.rt.get_mut(&i.name).unwrap();
rt.unhealthy_restarts += 1;
rt.failures = 0;
rt.healthy = None;
rt.since = Some(Instant::now());
supervise::restart_app(&sb, &self.service, oci)?;
}
Ok(())
}
fn setup(&self, def: &StackDef, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> Result<()> {
let keys = spec.secret_keys();
let values = if keys.is_empty() {
BTreeMap::new()
} else {
super::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?
};
if supervise::push_secrets(sb, spec, &values)? && oci {
supervise::restart_app(sb, &self.service, oci)?;
}
if spec.command.is_some() && !oci {
let mut s = spec.clone();
s.restart = Some(RestartMode::Always);
let env = supervise::secret_env(spec, &values)?;
supervise::install(sb, &self.service, &s, !spec.secrets.is_empty(), &env)?;
}
Ok(())
}
fn oci_secret_env(
&self,
def: &StackDef,
spec: &SandboxSpec,
) -> Result<BTreeMap<String, String>> {
if spec.env.secrets.is_empty() {
return Ok(BTreeMap::new());
}
let values = super::secrets::values(
&self.inner.secrets,
&def.org,
&def.secrets,
spec.env.secrets.values().map(String::as_str),
)?;
supervise::secret_env(spec, &values)
}
fn alive(&self, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> bool {
if spec.command.is_some() && !oci {
return supervise::unit_state(sb, &self.service).is_ok_and(|s| s == "active");
}
sb.info()
.is_ok_and(|i| i.status.eq_ignore_ascii_case("running"))
}
fn count_restart(&mut self, name: &str) {
self.rt
.entry(name.to_string())
.or_default()
.restarts
.push_back(Instant::now());
}
fn restart_budget_spent(
&mut self,
name: &str,
policy: Option<&crate::spec::RestartPolicy>,
) -> bool {
let Some(p) = policy else { return false };
let Some(max) = p.max_attempts else {
return false;
};
let window = p
.window
.as_deref()
.and_then(|w| crate::flex::parse_duration(w).ok());
let rt = self.rt.entry(name.to_string()).or_default();
if let Some(w) = window {
while rt.restarts.front().is_some_and(|t| t.elapsed() > w) {
rt.restarts.pop_front();
}
}
rt.restarts.len() as u32 >= max
}
#[expect(clippy::too_many_arguments)]
fn replace(
&mut self,
def: &StackDef,
spec: &SandboxSpec,
rev: &str,
slot: u32,
old: Option<&str>,
order: UpdateOrder,
oci: bool,
probe: Option<&HealthProbe>,
monitor: Duration,
) -> Result<()> {
let secret_env = if oci {
self.oci_secret_env(def, spec)?
} else {
BTreeMap::new()
};
if let (Some(o), UpdateOrder::StopFirst) = (old, order) {
self.log(&format!("slot {slot}: replacing {o} (stop-first)"));
self.slot_state(def, slot, Some("draining"), None, None);
self.retire(o)?;
self.slot_state(def, slot, Some("retired"), None, None);
}
let name = instance_name(&self.stack, &self.service, slot, &new_id())?;
self.log(&format!("slot {slot}: creating {name} (rev {rev})"));
self.slot_state(def, slot, None, Some(&name), Some("creating"));
let mut s = instance_spec(def, &self.service, spec, slot, rev)?;
s.name = Some(name.clone());
s.env.vars.extend(secret_env);
let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
let stack = self.q.clone();
let mut report = |m: &str| eprintln!("isb serve: {stack}: {m}");
let created =
crate::sandbox::ensure(self.client(), &d, EnsureOptions::default(), &mut report);
let result = created.and_then(|_| {
let inst = Inst {
name: name.clone(),
slot,
rev: rev.to_string(),
status: "Running".into(),
};
self.insts.push(inst.clone());
self.slot_state(def, slot, None, None, Some("probing"));
self.wait_serving(def, &inst, spec, oci, probe, monitor)
});
if let Err(e) = result {
let (e, attempt) =
super::failure::explain(self.client(), &name, &self.service, oci, e, now_ms());
self.inner.failures.record(&self.q, &self.service, attempt);
self.event(
"error",
Some(&name),
&format!("{name} did not come up: {e}"),
);
self.slot_state(def, slot, None, None, Some("failed"));
let _ = self.retire(&name);
return Err(e);
}
if let (Some(o), UpdateOrder::StartFirst) = (old, order) {
self.log(&format!("slot {slot}: {name} is serving; retiring {o}"));
self.slot_state(def, slot, Some("draining"), None, None);
self.retire(o)?;
self.slot_state(def, slot, Some("retired"), None, None);
}
self.slot_state(def, slot, None, None, Some("serving"));
Ok(())
}
fn wait_serving(
&mut self,
def: &StackDef,
i: &Inst,
spec: &SandboxSpec,
oci: bool,
probe: Option<&HealthProbe>,
monitor: Duration,
) -> Result<()> {
let deadline = match probe {
Some(p) => p.start_period + p.interval * p.retries + Duration::from_secs(30),
None => Duration::from_secs(60),
}
.max(Duration::from_secs(60));
let started = Instant::now();
loop {
self.maintain(def, i, spec, oci, probe)?;
if self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
break;
}
if self
.rt
.get(&i.name)
.is_some_and(|r| r.healthy == Some(false))
{
return Err(Error::invalid(format!(
"unhealthy: {}",
self.rt[&i.name].last_probe
)));
}
if started.elapsed() > deadline {
let why = self
.rt
.get(&i.name)
.map(|r| r.last_probe.clone())
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "its app is not running".into());
return Err(Error::invalid(format!(
"not serving after {:?}: {why}",
started.elapsed()
)));
}
if let Some(rt) = self.rt.get_mut(&i.name) {
rt.next_probe = None;
}
std::thread::sleep(
probe
.map(|p| p.start_interval)
.unwrap_or(Duration::from_secs(1))
.min(Duration::from_secs(2)),
);
}
self.sync_routes();
self.slot_state(def, i.slot, None, None, Some("monitoring"));
let watch = Instant::now();
while watch.elapsed() < monitor {
std::thread::sleep(Duration::from_secs(1).min(monitor));
if let Some(rt) = self.rt.get_mut(&i.name) {
rt.next_probe = None;
}
self.maintain(def, i, spec, oci, probe)?;
if !self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
return Err(Error::invalid(format!(
"failed within the {monitor:?} monitor period"
)));
}
}
Ok(())
}
fn retire(&mut self, name: &str) -> Result<()> {
let ip = self.rt.get(name).and_then(|r| r.ip);
self.set_rotation(name, false);
self.sync_routes();
if let Some(ip) = ip {
let started = Instant::now();
if let Some(o) = &self.inner.observer {
o.drain(&self.q, &self.service, ip, DRAIN);
}
let left = DRAIN.saturating_sub(started.elapsed());
for (k, p) in &self.routes {
self.inner
.balancer
.wait_drained(k, SocketAddr::new(ip, p.target), left);
}
}
if let Ok(sb) = Sandbox::get(self.client(), name) {
let _ = sb.stop(false, Duration::from_secs(10));
}
self.rt.remove(name);
self.insts.retain(|i| i.name != name);
match Sandbox::remove(self.client(), name, true) {
Err(e) if !e.is_not_found() => Err(e),
_ => Ok(()),
}
}
fn set_rotation(&mut self, name: &str, on: bool) {
let rt = self.rt.entry(name.to_string()).or_default();
if rt.in_rotation != on {
rt.in_rotation = on;
self.sync_routes();
}
}
fn set_routes(&mut self, spec: &SandboxSpec) {
let want: BTreeMap<String, Published> = match published(spec) {
Ok(ps) => ps
.into_iter()
.map(|p| (format!("{}/{}/{}", self.q, self.service, p.listen), p))
.collect(),
Err(e) => {
self.message = Some(e.to_string());
BTreeMap::new()
}
};
let stale: Vec<String> = self
.routes
.keys()
.filter(|k| !want.contains_key(*k))
.cloned()
.collect();
for k in stale {
self.inner.balancer.remove_route(&k);
self.routes.remove(&k);
self.route_errors.remove(&k);
}
for (k, p) in want {
self.routes.entry(k).or_insert(p);
}
self.sync_routes();
}
fn sync_routes(&mut self) {
let mut errors = BTreeMap::new();
for (k, p) in &self.routes {
let backends: Vec<SocketAddr> = self
.rt
.values()
.filter(|r| r.in_rotation)
.filter_map(|r| r.ip)
.map(|ip| SocketAddr::new(ip, p.target))
.collect();
if let Err(e) = self.inner.balancer.set_route(k, p.listen, backends) {
errors.insert(k.clone(), e.to_string());
}
}
for (k, e) in &errors {
if self.route_errors.get(k) != Some(e) {
self.log(&format!("cannot publish {k}: {e}"));
}
}
self.route_errors = errors;
self.sync_observer();
self.sync_dns();
}
fn sync_observer(&mut self) {
let Some(o) = self.inner.observer.clone() else {
return;
};
let mut ips: Vec<IpAddr> = self
.rt
.values()
.filter(|r| r.in_rotation)
.filter_map(|r| r.ip)
.collect();
ips.sort();
if self.observed.as_ref() != Some(&ips) {
o.rotation(&self.q, &self.service, &ips);
self.observed = Some(ips);
}
}
fn sync_dns(&mut self) {
let dir = crate::discovery::org_dir(&self.org);
let mut ips: Vec<IpAddr> = self
.rt
.values()
.filter(|r| r.in_rotation)
.filter_map(|r| r.ip)
.collect();
ips.sort();
let fresh = self.dns_last.is_none() && ips.is_empty();
if fresh || self.dns_last.as_ref() == Some(&ips) || !dir.is_dir() {
return;
}
match crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &ips) {
Ok(()) => {
self.dns_last = Some(ips);
self.dns_error = None;
}
Err(e) => {
let e = e.to_string();
if self.dns_error.as_deref() != Some(&e) {
self.event(
"warn",
None,
&format!("cannot publish the service name: {e}"),
);
self.dns_error = Some(e);
}
}
}
}
#[expect(
clippy::too_many_lines,
reason = "predates the lint ratchet; split it when next changed"
)]
fn publish_status(&mut self, def: &StackDef) {
self.sync_observer();
self.sync_dns();
let Ok(spec) = def.service(&self.service) else {
return;
};
let rev = def.revision(&self.service).unwrap_or_default();
let probe = matches!(spec.health_probe(), Ok(Some(_)));
let snap = self.inner.snapshot.lock().unwrap().instances.clone();
let instances: Vec<InstanceStatus> = self
.insts
.iter()
.map(|i| {
let rt = self.rt.get(&i.name);
let m = snap.get(&format!("{}/{}", self.oclient.project_name(), i.name));
let health = match (probe, rt.and_then(|r| r.healthy)) {
(false, _) => "none",
(true, Some(true)) => "healthy",
(true, Some(false)) => "unhealthy",
(true, None) => "starting",
};
InstanceStatus {
name: i.name.clone(),
slot: i.slot,
rev: i.rev.clone(),
status: m
.map(|m| m.status.clone())
.unwrap_or_else(|| i.status.clone()),
health: health.into(),
ip: rt.and_then(|r| r.ip).map(|ip| ip.to_string()),
in_rotation: rt.is_some_and(|r| r.in_rotation),
restarts: rt.map(|r| r.restarts.len() as u32).unwrap_or(0),
last_probe: rt.map(|r| r.last_probe.clone()).unwrap_or_default(),
cpu_pct: m.and_then(|m| m.cpu_pct),
cpu_history: m.map(|m| m.cpu_history.clone()).unwrap_or_default(),
mem_bytes: m.and_then(|m| m.mem_bytes),
disk_bytes: m.and_then(|m| m.disk_bytes),
}
})
.collect();
let mut instances = instances;
instances.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
let routes = self.inner.balancer.routes();
let now = Instant::now();
let mut ports = Vec::new();
for (k, p) in &self.routes {
let r = routes.iter().find(|r| r.key == *k);
let accepted = r.map(|r| r.accepted).unwrap_or(0);
let e = self
.rates
.entry(k.clone())
.or_insert_with(|| (accepted, now, VecDeque::new()));
let dt = now.duration_since(e.1).as_secs_f32();
if dt >= 1.0 {
let rate = accepted.saturating_sub(e.0) as f32 / dt;
if e.2.len() == crate::metrics::HISTORY {
e.2.pop_front();
}
e.2.push_back(rate);
e.0 = accepted;
e.1 = now;
}
ports.push(PortStatus {
listen: p.listen.to_string(),
target: p.target,
backends: r
.map(|r| r.backends.iter().map(|b| b.addr.to_string()).collect())
.unwrap_or_default(),
error: self.route_errors.get(k).cloned(),
accepted,
active: r
.map(|r| {
r.backends
.iter()
.chain(r.draining.iter())
.map(|b| b.active)
.sum()
})
.unwrap_or(0),
rate_history: e.2.iter().copied().collect(),
});
}
self.rates.retain(|k, _| self.routes.contains_key(k));
let running = instances
.iter()
.filter(|i| i.status.eq_ignore_ascii_case("running"))
.count() as u32;
let healthy = instances.iter().filter(|i| i.in_rotation).count() as u32;
self.watch_health(healthy, spec.replicas());
let st = ServiceStatus {
service: self.service.clone(),
image: spec.image.clone(),
rev,
replicas: spec.replicas(),
running,
healthy,
state: self.state.clone(),
message: self.message.clone(),
instances,
ports,
rollout: self.rollout.clone(),
checked_at: now_secs(),
domains: Vec::new(),
};
self.inner.status.lock().unwrap().insert(self.key(), st);
}
}
pub fn unsettled(st: &StackStatus) -> BTreeSet<String> {
st.services
.iter()
.filter(|s| s.state != "converged")
.map(|s| s.service.clone())
.collect()
}
#[cfg(test)]
#[path = "controller_tests.rs"]
mod tests;