Skip to main content

isb_core/stack/
controller.rs

1//! The reconciler behind `isb serve`: one worker thread per service keeps its
2//! replicas created, current, running, healthy and in the load balancer.
3//!
4//! A worker never holds a lock across incus calls, and a slow service (an
5//! image pull, a long rollout) never delays another's health checks. Workers
6//! are told about a new deployment through their shared slot and woken early;
7//! a worker whose service went away deletes its own instances and exits, so
8//! nothing ever has to join a worker from inside another.
9
10use std::collections::{BTreeMap, BTreeSet, VecDeque};
11use std::net::{IpAddr, SocketAddr};
12use std::sync::atomic::{AtomicBool, Ordering};
13use std::sync::{Arc, Condvar, Mutex};
14use std::time::{Duration, Instant};
15
16use serde::Serialize;
17use serde_json::Value;
18
19use super::changes::diff;
20use super::ports::{Published, published};
21use super::{
22    LABEL_REV, LABEL_SERVICE, LABEL_SLOT, LABEL_STACK, StackDef, Store, instance_name, new_id,
23    now_secs, validate_stack_name,
24};
25use crate::balance::Balancer;
26use crate::client::{Client, encode_query, encode_segment};
27use crate::error::{Error, Result};
28use crate::org::OrgId;
29use crate::plan::Desired;
30use crate::sandbox::{EnsureOptions, Sandbox};
31use crate::secrets::Secrets;
32use crate::spec::{
33    DependCondition, FailureAction, HealthProbe, RestartCondition, RestartMode, SandboxSpec,
34    UpdateConfig, UpdateOrder,
35};
36use crate::supervise;
37
38/// How long a replaced instance's connections may drain before it is stopped.
39const DRAIN: Duration = Duration::from_secs(10);
40/// How often an unhealthy app is restarted before its instance is replaced.
41const RESTARTS_BEFORE_REPLACE: u32 = 3;
42/// How long a service that was healthy must have no healthy replica before
43/// `health.unhealthy` is raised (restarts and short blips stay quiet).
44const HEALTH_DEBOUNCE: Duration = Duration::from_secs(20);
45
46/// One replica, as `stack_status` reports it.
47#[derive(Debug, Clone, Serialize)]
48pub struct InstanceStatus {
49    pub name: String,
50    pub slot: u32,
51    pub rev: String,
52    /// incus status: Running, Stopped, ...
53    pub status: String,
54    /// `healthy`, `unhealthy`, `starting`, or `none` (no healthcheck: the app
55    /// is judged by its process alone).
56    pub health: String,
57    pub ip: Option<String>,
58    /// Receiving traffic from the balancer.
59    pub in_rotation: bool,
60    pub restarts: u32,
61    #[serde(skip_serializing_if = "String::is_empty")]
62    pub last_probe: String,
63    /// Percent of one core; none until two samples were taken.
64    pub cpu_pct: Option<f32>,
65    pub cpu_history: Vec<f32>,
66    pub mem_bytes: Option<u64>,
67    /// Root disk usage, where the storage driver reports it.
68    pub disk_bytes: Option<u64>,
69}
70
71/// A published port: TCP served by the balancer, UDP by a NAT proxy on the
72/// replica (`listen` ends in `/udp`, `backends` is the replica).
73#[derive(Debug, Clone, Serialize)]
74pub struct PortStatus {
75    pub listen: String,
76    pub target: u16,
77    pub backends: Vec<String>,
78    #[serde(skip_serializing_if = "Option::is_none")]
79    pub error: Option<String>,
80    /// Connections accepted since the daemon started, and per second lately.
81    pub accepted: u64,
82    pub active: usize,
83    pub rate_history: Vec<f32>,
84}
85
86/// One service of a stack, as `stack_status` reports it.
87#[derive(Debug, Clone, Serialize, Default)]
88pub struct ServiceStatus {
89    pub service: String,
90    pub image: String,
91    pub rev: String,
92    pub replicas: u32,
93    pub running: u32,
94    pub healthy: u32,
95    /// `starting`, `converged`, `updating`, `paused`, `waiting`, `failing`.
96    pub state: String,
97    #[serde(skip_serializing_if = "Option::is_none")]
98    pub message: Option<String>,
99    pub instances: Vec<InstanceStatus>,
100    pub ports: Vec<PortStatus>,
101    /// The rollout in progress, slot by slot.
102    #[serde(skip_serializing_if = "Option::is_none")]
103    pub rollout: Option<RolloutStatus>,
104    /// Unix seconds of the last completed reconcile pass.
105    pub checked_at: u64,
106    /// The service's domains as the ingress serves them.
107    #[serde(skip_serializing_if = "Vec::is_empty")]
108    pub domains: Vec<crate::ingress::DomainStatus>,
109}
110
111/// Something that follows which replicas receive traffic: the ingress.
112/// Called from worker threads, never with a controller lock held.
113pub trait Observer: Send + Sync {
114    /// The in-rotation replicas of a service changed. `stack` is qualified.
115    fn rotation(&self, stack: &str, service: &str, ips: &[IpAddr]);
116    /// Wait (at most `timeout`) until requests to `ip` have drained from the
117    /// observer's proxies, before the replica is stopped.
118    fn drain(&self, stack: &str, service: &str, ip: IpAddr, timeout: Duration);
119    /// The deployed definitions changed (deploy, scale, rollback, removal).
120    fn stacks_changed(&self, defs: Vec<Arc<StackDef>>);
121    /// A service's domains for its status.
122    fn domains(&self, stack: &str, service: &str) -> Vec<crate::ingress::DomainStatus>;
123}
124
125/// A rollout in progress.
126#[derive(Debug, Clone, Serialize, Default)]
127pub struct RolloutStatus {
128    pub to_rev: String,
129    /// `stop-first` or `start-first`.
130    pub order: String,
131    pub parallelism: usize,
132    pub done: usize,
133    pub total: usize,
134    /// Unix seconds.
135    pub started_at: u64,
136    pub slots: Vec<SlotRollout>,
137}
138
139/// One slot's handover: the instance going away and the one replacing it.
140#[derive(Debug, Clone, Serialize, Default)]
141pub struct SlotRollout {
142    pub slot: u32,
143    pub old: Option<String>,
144    pub old_rev: Option<String>,
145    /// `serving`, `draining`, `retired`, or `none`.
146    pub old_state: String,
147    pub new: Option<String>,
148    /// `waiting`, `creating`, `probing`, `monitoring`, `serving`, `failed`.
149    pub new_state: String,
150}
151
152/// Something that happened, for the event feed.
153#[derive(Debug, Clone, Serialize)]
154pub struct Event {
155    /// Increases by one per event; pass the last one seen as `since`.
156    pub seq: u64,
157    /// Unix milliseconds.
158    pub at: u64,
159    /// `info`, `warn` or `error`.
160    pub level: String,
161    pub stack: String,
162    #[serde(skip_serializing_if = "String::is_empty")]
163    pub service: String,
164    #[serde(skip_serializing_if = "Option::is_none")]
165    pub instance: Option<String>,
166    pub message: String,
167    /// What happened, for consumers that act on events (notifications):
168    /// dotted, `<subject>.<outcome>`. In use: `deploy.succeeded`,
169    /// `deploy.failed`, `health.unhealthy`, `health.recovered`,
170    /// `backup.succeeded`, `backup.failed`, `restore.succeeded`,
171    /// `restore.failed`, `job.succeeded`, `job.failed`, `cert.issued`,
172    /// `cert.failed`, `preview.created`, `preview.removed` (a preview's
173    /// deploys are `deploy.*` under its own stack), `server.unreachable`,
174    /// `server.recovered` (on a control plane, stack `<org>/@servers` for
175    /// each org on the server and `system/@servers`, service = the server).
176    /// Most events have none.
177    #[serde(default, skip_serializing_if = "Option::is_none")]
178    pub kind: Option<String>,
179}
180
181/// Hears every event as it is emitted. It must not block (the history
182/// queues and returns).
183pub type EventSink = Arc<dyn Fn(&Event) + Send + Sync>;
184
185/// Events kept for late readers.
186const EVENTS_KEPT: usize = 1000;
187
188/// The latest host and instance sample.
189#[derive(Debug, Clone, Default)]
190pub struct Snapshot {
191    pub host: crate::metrics::HostSample,
192    /// Keyed by `<project>/<name>`: names are unique per project only.
193    pub instances: BTreeMap<String, crate::metrics::InstanceSample>,
194    /// Unix milliseconds; 0 before the first sample.
195    pub at: u64,
196}
197
198pub fn now_ms() -> u64 {
199    std::time::SystemTime::now()
200        .duration_since(std::time::UNIX_EPOCH)
201        .map(|d| d.as_millis() as u64)
202        .unwrap_or(0)
203}
204
205/// A stack, as `stack_status` and `stack_list` report it.
206#[derive(Debug, Clone, Serialize)]
207pub struct StackStatus {
208    pub name: String,
209    pub org: String,
210    pub deployed_at: u64,
211    pub deployed_by: String,
212    pub has_previous: bool,
213    /// True when every service has its replicas, all current and healthy.
214    pub converged: bool,
215    pub services: Vec<ServiceStatus>,
216}
217
218/// What a deploy is about to do, per service.
219#[derive(Debug, Clone, Serialize)]
220pub struct DeployChange {
221    pub service: String,
222    /// `create`, `update` (new revision: rolling replace), `scale`, `remove`,
223    /// or `unchanged`.
224    pub change: String,
225    pub rev: String,
226    pub replicas: u32,
227}
228
229/// The shared slot a worker reads its instructions from.
230struct Slot {
231    def: Arc<StackDef>,
232    /// Set when the service left its stack: delete its instances and exit.
233    remove: bool,
234    /// Delete the service's named volumes too (with `remove`).
235    remove_volumes: bool,
236}
237
238struct WorkerShared {
239    slot: Mutex<Slot>,
240    wake: Condvar,
241    stop: AtomicBool,
242    /// Set when something the service was failing on may have changed (an
243    /// org's limits): the worker drops its retry backoff on its next pass.
244    kick: AtomicBool,
245}
246
247struct Inner {
248    client: Client,
249    store: Store,
250    balancer: Balancer,
251    interval: Duration,
252    /// Live workers, by (stack, service).
253    workers: Mutex<BTreeMap<(String, String), Arc<WorkerShared>>>,
254    /// Deployed stacks, by name.
255    stacks: Mutex<BTreeMap<String, Arc<StackDef>>>,
256    status: Mutex<BTreeMap<(String, String), ServiceStatus>>,
257    /// The output of each service's last replica that failed to come up.
258    failures: super::failure::Failures,
259    events: Mutex<(u64, VecDeque<Event>)>,
260    snapshot: Mutex<Snapshot>,
261    /// The org stores secret values are read from at delivery.
262    secrets: Arc<Secrets>,
263    /// Held across read-modify-save of a definition, so a version bump and
264    /// a deploy never overwrite each other.
265    edit: Mutex<()>,
266    /// When driver-backed secrets are next checked for a new version.
267    refresh: Mutex<super::secrets::RefreshSchedule>,
268    observer: Option<Arc<dyn Observer>>,
269    /// Where each event also goes: the persistent history.
270    event_sink: Mutex<Option<EventSink>>,
271    /// Where each sample also goes: the metrics history.
272    metrics_sink: Mutex<Option<std::sync::mpsc::SyncSender<crate::metrics_history::Sample>>>,
273}
274
275impl Inner {
276    fn emit(
277        &self,
278        level: &str,
279        stack: &str,
280        service: &str,
281        instance: Option<&str>,
282        message: String,
283    ) {
284        self.emit_kind(None, level, stack, service, instance, message);
285    }
286
287    fn emit_kind(
288        &self,
289        kind: Option<&str>,
290        level: &str,
291        stack: &str,
292        service: &str,
293        instance: Option<&str>,
294        message: String,
295    ) {
296        let mut e = self.events.lock().unwrap();
297        e.0 += 1;
298        let ev = Event {
299            seq: e.0,
300            at: now_ms(),
301            level: level.into(),
302            stack: stack.into(),
303            service: service.into(),
304            instance: instance.map(String::from),
305            message,
306            kind: kind.map(String::from),
307        };
308        if e.1.len() == EVENTS_KEPT {
309            e.1.pop_front();
310        }
311        let sink = self.event_sink.lock().unwrap().clone();
312        if let Some(s) = sink {
313            s(&ev);
314        }
315        e.1.push_back(ev);
316    }
317}
318
319/// How often the metrics sampler runs.
320const SAMPLE_EVERY: Duration = Duration::from_secs(2);
321/// The longest wait between looking for driver-backed secrets that are due.
322const SECRET_TICK: Duration = Duration::from_secs(10);
323
324/// The daemon's stack controller.
325#[derive(Clone)]
326pub struct Controller {
327    inner: Arc<Inner>,
328}
329
330impl Controller {
331    /// Load every stored stack and start reconciling it. Secret values are
332    /// read from `secrets` whenever they are delivered.
333    pub fn start(
334        client: Client,
335        store: Store,
336        interval: Duration,
337        secrets: Arc<Secrets>,
338    ) -> Result<Controller> {
339        Controller::start_with(client, store, interval, secrets, None)
340    }
341
342    /// [`Controller::start`], telling `observer` about rotation changes from
343    /// the first one on.
344    pub fn start_with(
345        client: Client,
346        store: Store,
347        interval: Duration,
348        secrets: Arc<Secrets>,
349        observer: Option<Arc<dyn Observer>>,
350    ) -> Result<Controller> {
351        let c = Controller {
352            inner: Arc::new(Inner {
353                client,
354                store,
355                balancer: Balancer::new(),
356                interval,
357                workers: Mutex::new(BTreeMap::new()),
358                stacks: Mutex::new(BTreeMap::new()),
359                status: Mutex::new(BTreeMap::new()),
360                failures: Default::default(),
361                events: Mutex::new((0, VecDeque::new())),
362                snapshot: Mutex::new(Snapshot::default()),
363                secrets,
364                edit: Mutex::new(()),
365                refresh: Mutex::new(Default::default()),
366                observer,
367                event_sink: Mutex::new(None),
368                metrics_sink: Mutex::new(None),
369            }),
370        };
371        // Driver-backed secrets are polled on their refresh intervals; the
372        // tick only decides which are due.
373        let weak = Arc::downgrade(&c.inner);
374        let _ = std::thread::Builder::new()
375            .name("isb-secrets".into())
376            .spawn(move || {
377                while let Some(inner) = weak.upgrade() {
378                    let tick = inner.interval.clamp(Duration::from_secs(1), SECRET_TICK);
379                    Controller { inner }.check_due_secrets();
380                    std::thread::sleep(tick);
381                }
382            });
383        // A weak handle, so the sampler ends with the controller.
384        let weak = Arc::downgrade(&c.inner);
385        let _ = std::thread::Builder::new()
386            .name("isb-metrics".into())
387            .spawn(move || {
388                let mut sampler = crate::metrics::Sampler::new();
389                while let Some(inner) = weak.upgrade() {
390                    match sampler.sample(&inner.client) {
391                        Ok((host, insts)) => {
392                            if let Some(tx) = &*inner.metrics_sink.lock().unwrap() {
393                                crate::metrics_history::offer(tx, (now_ms(), insts.clone()));
394                            }
395                            *inner.snapshot.lock().unwrap() = Snapshot {
396                                host,
397                                instances: insts
398                                    .into_iter()
399                                    .map(|i| (format!("{}/{}", i.project, i.name), i))
400                                    .collect(),
401                                at: now_ms(),
402                            };
403                        }
404                        Err(e) => eprintln!("isb serve: metrics: {e}"),
405                    }
406                    drop(inner);
407                    std::thread::sleep(SAMPLE_EVERY);
408                }
409            });
410        let defs = c.inner.store.load_all()?;
411        // Service names of stacks removed while no daemon ran.
412        for def in &defs {
413            let dir = crate::discovery::org_dir(&def.org);
414            let keep: Vec<(String, String)> = defs
415                .iter()
416                .filter(|d| d.org == def.org)
417                .flat_map(|d| d.file.services.keys().map(|s| (d.name.clone(), s.clone())))
418                .collect();
419            crate::discovery::prune(&dir, &keep);
420        }
421        for def in defs {
422            eprintln!("isb serve: resuming stack {}", def.name);
423            c.apply(Arc::new(def));
424        }
425        // A vault may have moved on while the daemon was down.
426        for def in c.definitions() {
427            let keys: Vec<String> = def.secrets.keys().cloned().collect();
428            c.check_bindings(&def.qualified(), &keys, false);
429        }
430        Ok(c)
431    }
432
433    /// A stored secret got a new value (`isb secret set`): every stack in
434    /// the org bound to it moves to the new version and rolls. Returns the
435    /// stacks rolled.
436    pub fn secret_changed(&self, org: &OrgId, name: &str) -> Vec<String> {
437        let mut rolled = Vec::new();
438        for def in self.definitions() {
439            if def.org != *org {
440                continue;
441            }
442            let keys: Vec<String> = def
443                .secrets
444                .iter()
445                .filter(|(_, b)| b.name == name)
446                .map(|(k, _)| k.clone())
447                .collect();
448            if !keys.is_empty() && self.check_bindings(&def.qualified(), &keys, false) {
449                rolled.push(def.qualified());
450            }
451        }
452        rolled
453    }
454
455    /// Re-read every binding to `name` in the org from its driver now
456    /// (`isb secret refresh`). Returns each driver and version found, and
457    /// the stacks rolled.
458    pub fn refresh_secret(
459        &self,
460        org: &OrgId,
461        name: &str,
462    ) -> Result<crate::stack::secrets::Refreshed> {
463        let mut found = Vec::new();
464        let mut rolled = Vec::new();
465        for def in self.definitions() {
466            if def.org != *org {
467                continue;
468            }
469            let q = def.qualified();
470            let keys: Vec<String> = def
471                .secrets
472                .iter()
473                .filter(|(_, b)| b.name == name)
474                .map(|(k, _)| k.clone())
475                .collect();
476            for k in &keys {
477                let b = &def.secrets[k];
478                let v = self.inner.secrets.refresh_in(&b.driver, org, name)?;
479                found.push((b.driver.clone(), v));
480                let every = def
481                    .file
482                    .secrets
483                    .get(k)
484                    .map(crate::spec::SecretDef::refresh_interval)
485                    .unwrap_or(crate::spec::DEFAULT_SECRET_REFRESH);
486                self.inner
487                    .refresh
488                    .lock()
489                    .unwrap()
490                    .reset(&q, k, every, Instant::now());
491            }
492            if !keys.is_empty() && self.check_bindings(&q, &keys, true) {
493                rolled.push(q);
494            }
495        }
496        Ok((found, rolled))
497    }
498
499    /// The driver-backed bindings whose refresh interval is up.
500    fn check_due_secrets(&self) {
501        let defs: Vec<(String, Arc<StackDef>)> = self
502            .inner
503            .stacks
504            .lock()
505            .unwrap()
506            .iter()
507            .map(|(q, d)| (q.clone(), d.clone()))
508            .collect();
509        let due = self
510            .inner
511            .refresh
512            .lock()
513            .unwrap()
514            .due(defs.iter().map(|(q, d)| (q.as_str(), &**d)), Instant::now());
515        let mut by_stack: BTreeMap<String, Vec<String>> = BTreeMap::new();
516        for (q, k) in due {
517            by_stack.entry(q).or_default().push(k);
518        }
519        for (q, keys) in by_stack {
520            self.check_bindings(&q, &keys, false);
521        }
522    }
523
524    /// Compare the given bindings of a stack with the store's current
525    /// versions; on any change, save the new versions and roll. `quiet`
526    /// skips logging lookups that fail (the caller reports them).
527    fn check_bindings(&self, q: &str, keys: &[String], quiet: bool) -> bool {
528        let _g = self.inner.edit.lock().unwrap();
529        let Ok(cur) = self.get_def(q) else {
530            return false;
531        };
532        let mut def = (*cur).clone();
533        let mut moved = Vec::new();
534        for k in keys {
535            let Some(b) = def.secrets.get_mut(k) else {
536                continue;
537            };
538            match self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
539                Ok(v) if v != b.version => {
540                    moved.push(format!("{} v{} -> v{v}", b.name, b.version));
541                    b.version = v;
542                }
543                Ok(_) => {}
544                // The services keep the value they have.
545                Err(e) if !quiet => self.note(
546                    "warn",
547                    q,
548                    format!("secret {}: cannot check its version: {e}", b.name),
549                ),
550                Err(_) => {}
551            }
552        }
553        if moved.is_empty() {
554            return false;
555        }
556        if let Err(e) = self.inner.store.save(&def) {
557            self.note("error", q, format!("cannot save new secret versions: {e}"));
558            return false;
559        }
560        self.apply(Arc::new(def));
561        self.note(
562            "info",
563            q,
564            format!("new secret version ({}): rolling", moved.join(", ")),
565        );
566        true
567    }
568
569    pub fn balancer(&self) -> &Balancer {
570        &self.inner.balancer
571    }
572
573    /// Events after `since` (0: all kept), oldest first, at most `limit`.
574    pub fn events(&self, since: u64, limit: usize) -> (u64, Vec<Event>) {
575        let e = self.inner.events.lock().unwrap();
576        let out: Vec<Event> = e.1.iter().filter(|x| x.seq > since).cloned().collect();
577        let skip = out.len().saturating_sub(limit);
578        (e.0, out.into_iter().skip(skip).collect())
579    }
580
581    /// Wait up to `timeout` for an event after `since`.
582    pub fn wait_events(&self, since: u64, limit: usize, timeout: Duration) -> (u64, Vec<Event>) {
583        let started = Instant::now();
584        loop {
585            let r = self.events(since, limit);
586            if !r.1.is_empty() || started.elapsed() >= timeout {
587                return r;
588            }
589            std::thread::sleep(Duration::from_millis(250));
590        }
591    }
592
593    /// Hand every event to `sink` too (the persistent history), starting
594    /// with the ones already kept, so none emitted before it was set is
595    /// missed and none is handed over twice.
596    pub fn set_event_sink(&self, sink: EventSink) {
597        let e = self.inner.events.lock().unwrap();
598        for ev in &e.1 {
599            sink(ev);
600        }
601        *self.inner.event_sink.lock().unwrap() = Some(sink);
602    }
603
604    /// Send every metrics sample to `tx` too (the metrics history).
605    pub fn set_metrics_sink(
606        &self,
607        tx: std::sync::mpsc::SyncSender<crate::metrics_history::Sample>,
608    ) {
609        *self.inner.metrics_sink.lock().unwrap() = Some(tx);
610    }
611
612    /// The latest metrics sample.
613    pub fn snapshot(&self) -> Snapshot {
614        self.inner.snapshot.lock().unwrap().clone()
615    }
616
617    /// Record an event from outside a worker (a deploy, a removal).
618    /// Record an event of a known kind ([`Event::kind`]) about a stack, or
619    /// one of its services when `service` is not empty.
620    pub fn event(&self, kind: &str, level: &str, stack: &str, service: &str, message: String) {
621        eprintln!("isb serve: {stack}: {message}");
622        self.inner
623            .emit_kind(Some(kind), level, stack, service, None, message);
624    }
625
626    /// Re-emit an event another daemon recorded (a control plane mirroring
627    /// a server's feed), under this feed's numbering.
628    pub fn relay(
629        &self,
630        kind: Option<&str>,
631        level: &str,
632        stack: &str,
633        service: &str,
634        instance: Option<&str>,
635        message: String,
636    ) {
637        self.inner
638            .emit_kind(kind, level, stack, service, instance, message);
639    }
640
641    pub fn note(&self, level: &str, stack: &str, message: String) {
642        eprintln!("isb serve: {stack}: {message}");
643        self.inner.emit(level, stack, "", None, message);
644    }
645
646    /// Record an event about one service of a stack (an app's deployment
647    /// log, say). Not echoed to stderr: it can be chatty.
648    pub fn note_service(&self, level: &str, stack: &str, service: &str, message: String) {
649        self.inner.emit(level, stack, service, None, message);
650    }
651
652    /// Record an event about one service (`stack` qualified).
653    pub fn service_event(&self, level: &str, stack: &str, service: &str, message: String) {
654        eprintln!("isb serve: {stack}/{service}: {message}");
655        self.inner.emit(level, stack, service, None, message);
656    }
657
658    fn notify_stacks(&self) {
659        if let Some(o) = &self.inner.observer {
660            o.stacks_changed(self.definitions());
661        }
662    }
663
664    /// What deploying `def` would change, without deploying it.
665    pub fn plan(&self, def: &StackDef) -> Result<Vec<DeployChange>> {
666        self.validate(def)?;
667        let mut def = def.clone();
668        self.pin_images(&mut def)?;
669        let def = &def;
670        let old = self
671            .inner
672            .stacks
673            .lock()
674            .unwrap()
675            .get(&def.qualified())
676            .cloned();
677        diff(old.as_deref(), def)
678    }
679
680    pub fn client(&self) -> &Client {
681        &self.inner.client
682    }
683
684    /// Check a stack definition against this host without deploying it:
685    /// every service must resolve (image source, paths, ports).
686    pub fn validate(&self, def: &StackDef) -> Result<()> {
687        validate_stack_name(&def.name)?;
688        // In the org's project, which is what `registry:` images resolve in.
689        let host = crate::sandbox::host_facts(&crate::org::client(&self.inner.client, &def.org))?;
690        for (svc, spec) in &def.file.services {
691            let mut s = instance_spec(def, svc, spec, 1, "0000")?;
692            s.name = Some(instance_name(&def.name, svc, 1, "0000")?);
693            crate::plan::resolve(&s, &def.file.volumes, &host, &def.base_dir)?;
694            published(spec)?;
695            crate::ingress::domain::validate(svc, &spec.domains)?;
696        }
697        super::ports::validate_udp(&self.inner.client, def, &self.definitions())
698    }
699
700    /// Deploy (or update) a stack. Returns what will change; the rollout
701    /// itself happens in the background.
702    pub fn deploy(&self, mut def: StackDef) -> Result<Vec<DeployChange>> {
703        self.validate(&def)?;
704        self.pin_images(&mut def)?;
705        let _g = self.inner.edit.lock().unwrap();
706        let old = self
707            .inner
708            .stacks
709            .lock()
710            .unwrap()
711            .get(&def.qualified())
712            .cloned();
713        if let Some(old) = &old {
714            let mut prev = (**old).clone();
715            prev.previous = None;
716            def.previous = Some(Box::new(prev));
717            // A forced update survives a redeploy that does not ask for one.
718            for (k, v) in &old.force {
719                def.force.entry(k.clone()).or_insert(*v);
720            }
721        }
722        let changes = diff(old.as_deref(), &def)?;
723        self.inner.store.save(&def)?;
724        self.apply(Arc::new(def));
725        Ok(changes)
726    }
727
728    /// Hand a definition to the workers: update existing ones, start new
729    /// ones, and tell those whose service is gone to clean up.
730    fn apply(&self, def: Arc<StackDef>) {
731        let name = def.qualified();
732        self.inner
733            .stacks
734            .lock()
735            .unwrap()
736            .insert(name.clone(), def.clone());
737        let mut workers = self.inner.workers.lock().unwrap();
738        for svc in def.file.services.keys() {
739            let key = (name.clone(), svc.clone());
740            match workers.get(&key) {
741                Some(w) => {
742                    let mut slot = w.slot.lock().unwrap();
743                    slot.def = def.clone();
744                    slot.remove = false;
745                    // A new deployment starts from a clean slate: the
746                    // failure the service reported before it is not the
747                    // deployment's, and a waiter must not read it as such.
748                    if let Some(st) = self.inner.status.lock().unwrap().get_mut(&key) {
749                        if st.state == "failing" {
750                            st.state = "updating".into();
751                            st.message = None;
752                        }
753                    }
754                    w.wake.notify_all();
755                }
756                None => {
757                    let shared = Arc::new(WorkerShared {
758                        slot: Mutex::new(Slot {
759                            def: def.clone(),
760                            remove: false,
761                            remove_volumes: false,
762                        }),
763                        wake: Condvar::new(),
764                        stop: AtomicBool::new(false),
765                        kick: AtomicBool::new(false),
766                    });
767                    workers.insert(key, shared.clone());
768                    spawn_worker(self.inner.clone(), &def, svc.clone(), shared);
769                }
770            }
771        }
772        for ((stack, svc), w) in workers.iter() {
773            if *stack == name && !def.file.services.contains_key(svc) {
774                let mut slot = w.slot.lock().unwrap();
775                slot.remove = true;
776                w.wake.notify_all();
777            }
778        }
779        drop(workers);
780        self.notify_stacks();
781    }
782
783    /// An org's limits changed: services of the org that are failing on a
784    /// limit (a quota refusal from incus) retry now instead of after their
785    /// backoff. Returns how many were woken.
786    pub fn org_limits_changed(&self, org: &OrgId) -> usize {
787        let stacks = self.inner.stacks.lock().unwrap();
788        let workers = self.inner.workers.lock().unwrap();
789        let status = self.inner.status.lock().unwrap();
790        let mut n = 0;
791        for (key, w) in workers.iter() {
792            let ours = stacks.get(&key.0).is_some_and(|d| d.org == *org);
793            let limited = status.get(key).is_some_and(|s| {
794                s.state == "failing" && s.message.as_deref().is_some_and(limit_error)
795            });
796            if ours && limited {
797                w.kick.store(true, Ordering::SeqCst);
798                // Under the slot lock, the worker is either waiting (and
799                // wakes) or yet to look at the flag.
800                let _slot = w.slot.lock().unwrap();
801                w.wake.notify_all();
802                n += 1;
803            }
804        }
805        n
806    }
807
808    /// Remove a stack: every instance and published port; with `volumes`,
809    /// its named volumes too. Returns once the workers have cleaned up (or
810    /// after `timeout`).
811    pub fn remove(&self, name: &str, volumes: bool, timeout: Duration) -> Result<()> {
812        {
813            let _g = self.inner.edit.lock().unwrap();
814            let Some(def) = self.inner.stacks.lock().unwrap().remove(name) else {
815                return Err(Error::NotFound(format!("stack {name}")));
816            };
817            self.inner.store.remove(&def.org, &def.name)?;
818        }
819        self.notify_stacks();
820        let ws: Vec<Arc<WorkerShared>> = self
821            .inner
822            .workers
823            .lock()
824            .unwrap()
825            .iter()
826            .filter(|((s, _), _)| s == name)
827            .map(|(_, w)| w.clone())
828            .collect();
829        for w in &ws {
830            let mut slot = w.slot.lock().unwrap();
831            slot.remove = true;
832            slot.remove_volumes = volumes;
833            w.wake.notify_all();
834        }
835        let started = Instant::now();
836        while started.elapsed() < timeout {
837            let left = self
838                .inner
839                .workers
840                .lock()
841                .unwrap()
842                .keys()
843                .any(|(s, _)| s == name);
844            if !left {
845                return Ok(());
846            }
847            std::thread::sleep(Duration::from_millis(200));
848        }
849        Err(Error::invalid(format!(
850            "stack {name}: still removing after {timeout:?}; it carries on in the background"
851        )))
852    }
853
854    /// Go back to the previous deployment (the current one becomes the
855    /// previous, so a second rollback undoes the first).
856    pub fn rollback(&self, name: &str) -> Result<Vec<DeployChange>> {
857        let _g = self.inner.edit.lock().unwrap();
858        let cur = self.get_def(name)?;
859        let prev = cur
860            .previous
861            .clone()
862            .ok_or_else(|| Error::invalid(format!("stack {name} has no previous deployment")))?;
863        let mut def = *prev;
864        def.deployed_at = now_secs();
865        // The store keeps only each secret's current value: that is what a
866        // rollback delivers, under its current version.
867        for b in def.secrets.values_mut() {
868            if let Ok(v) = self.inner.secrets.version_in(&b.driver, &def.org, &b.name) {
869                b.version = v;
870            }
871        }
872        let mut cur2 = (*cur).clone();
873        cur2.previous = None;
874        let changes = diff(Some(&cur), &def)?;
875        def.previous = Some(Box::new(cur2));
876        self.inner.store.save(&def)?;
877        self.apply(Arc::new(def));
878        Ok(changes)
879    }
880
881    /// Change one service's replica count.
882    pub fn scale(&self, name: &str, service: &str, replicas: u32) -> Result<()> {
883        let _g = self.inner.edit.lock().unwrap();
884        let cur = self.get_def(name)?;
885        let mut def = (*cur).clone();
886        let spec = def
887            .file
888            .services
889            .get_mut(service)
890            .ok_or_else(|| Error::NotFound(format!("service {service} in stack {name}")))?;
891        spec.deploy.get_or_insert_with(Default::default).replicas = Some(replicas);
892        super::ports::check_replicas(service, spec)?;
893        self.inner.store.save(&def)?;
894        self.apply(Arc::new(def));
895        Ok(())
896    }
897
898    /// Replace every instance of a service even though its spec is the same
899    /// (`docker service update --force`): picks up a moved image tag or a
900    /// changed bind-mounted file.
901    pub fn redeploy(&self, name: &str, service: &str) -> Result<()> {
902        let _g = self.inner.edit.lock().unwrap();
903        let cur = self.get_def(name)?;
904        cur.service(service)?;
905        let mut def = (*cur).clone();
906        *def.force.entry(service.to_string()).or_insert(0) += 1;
907        // A moved tag is what a redeploy is usually for.
908        self.pin_images(&mut def)?;
909        self.inner.store.save(&def)?;
910        self.apply(Arc::new(def));
911        Ok(())
912    }
913
914    /// Resolve every `registry:` image given by tag to the digest the tag
915    /// names now, in the stack's org (a digest in the file only has to
916    /// exist there).
917    fn pin_images(&self, def: &mut StackDef) -> Result<()> {
918        def.images.clear();
919        for (svc, spec) in &def.file.services {
920            let Some(r) = spec.image.strip_prefix("registry:") else {
921                continue;
922            };
923            let r = crate::registry::ImageRef::parse(r)?;
924            let reg = crate::registry::Registry::shared(&self.inner.client)?;
925            // A digest is checked too, so a typo (or another org's digest)
926            // fails the deploy rather than the rollout.
927            let d = reg.resolve(&def.org, &r)?;
928            if r.digest.is_none() {
929                def.images.insert(svc.clone(), d);
930            }
931        }
932        Ok(())
933    }
934
935    fn get_def(&self, name: &str) -> Result<Arc<StackDef>> {
936        self.inner
937            .stacks
938            .lock()
939            .unwrap()
940            .get(name)
941            .cloned()
942            .ok_or_else(|| Error::NotFound(format!("stack {name}")))
943    }
944
945    /// Every stored definition, without computing status.
946    pub fn definitions(&self) -> Vec<Arc<StackDef>> {
947        self.inner
948            .stacks
949            .lock()
950            .unwrap()
951            .values()
952            .cloned()
953            .collect()
954    }
955
956    /// The stored definition of a stack.
957    pub fn definition(&self, name: &str) -> Result<StackDef> {
958        self.get_def(name).map(|d| (*d).clone())
959    }
960
961    pub fn list(&self) -> Vec<StackStatus> {
962        let names: Vec<String> = self.inner.stacks.lock().unwrap().keys().cloned().collect();
963        names.iter().filter_map(|n| self.status(n).ok()).collect()
964    }
965
966    pub fn status(&self, name: &str) -> Result<StackStatus> {
967        let def = self.get_def(name)?;
968        let st = self.inner.status.lock().unwrap();
969        let services: Vec<ServiceStatus> = def
970            .file
971            .services
972            .keys()
973            .map(|svc| {
974                st.get(&(name.to_string(), svc.clone()))
975                    .cloned()
976                    .unwrap_or_else(|| ServiceStatus {
977                        service: svc.clone(),
978                        state: "starting".into(),
979                        ..Default::default()
980                    })
981            })
982            .collect();
983        drop(st);
984        let mut services = services;
985        if let Some(o) = &self.inner.observer {
986            for s in &mut services {
987                s.domains = o.domains(name, &s.service);
988            }
989        }
990        let converged = services.iter().all(|s| s.state == "converged");
991        Ok(StackStatus {
992            name: def.name.clone(),
993            org: def.org.to_string(),
994            deployed_at: def.deployed_at,
995            deployed_by: def.deployed_by.clone(),
996            has_previous: def.previous.is_some(),
997            converged,
998            services,
999        })
1000    }
1001
1002    /// Recent output of a service's replicas (or one slot's).
1003    pub fn logs(
1004        &self,
1005        name: &str,
1006        service: &str,
1007        slot: Option<u32>,
1008        lines: usize,
1009    ) -> Result<BTreeMap<String, String>> {
1010        let def = self.get_def(name)?;
1011        let oci = crate::plan::ImageSource::parse(&def.service(service)?.image)?.is_oci();
1012        let oc = crate::org::client(&self.inner.client, &def.org);
1013        super::failure::replica_logs(&oc, &def.name, service, oci, slot, lines)
1014    }
1015
1016    /// The last replica of a service that failed to come up, with its
1017    /// output, while the service is not converged: after the instance is
1018    /// deleted, this is all that is left to read.
1019    pub fn last_failure(&self, name: &str, service: &str) -> Option<super::failure::FailedAttempt> {
1020        let state = self.inner.status.lock().unwrap();
1021        let converged = state
1022            .get(&(name.to_string(), service.to_string()))
1023            .is_some_and(|s| s.state == "converged");
1024        drop(state);
1025        if converged {
1026            return None;
1027        }
1028        self.inner.failures.last(name, service)
1029    }
1030
1031    /// Stop every worker and the balancer. Apps keep running in their
1032    /// instances; published ports stop until the next start.
1033    pub fn shutdown(&self) {
1034        for w in self.inner.workers.lock().unwrap().values() {
1035            w.stop.store(true, Ordering::SeqCst);
1036            w.wake.notify_all();
1037        }
1038        self.inner.balancer.clear();
1039    }
1040}
1041
1042/// Whether a failure message is an incus project-limit refusal (see
1043/// `org::limits`): one that raising the org's limits can fix.
1044fn limit_error(msg: &str) -> bool {
1045    msg.contains(" quota (") || msg.contains(" limit (")
1046}
1047
1048/// The spec an instance of `service` is created from: labelled, with its
1049/// published TCP ports removed (the balancer serves them) and its UDP ports
1050/// as NAT proxies (see [`super::ports`]), and always long-running, as swarm
1051/// ignores `restart` in favour of `restart_policy`.
1052fn instance_spec(
1053    def: &StackDef,
1054    service: &str,
1055    spec: &SandboxSpec,
1056    slot: u32,
1057    rev: &str,
1058) -> Result<SandboxSpec> {
1059    let mut s = spec.clone();
1060    s.image = def.instance_image(service, &spec.image);
1061    s.restart = Some(RestartMode::Always);
1062    (s.ports, s.stack_udp) = super::ports::instance_ports(spec)?;
1063    s.domains.clear();
1064    if let Some(d) = &s.deploy {
1065        s.labels.extend(d.labels.clone());
1066    }
1067    s.labels.insert(LABEL_STACK.into(), def.name.clone());
1068    s.labels.insert(LABEL_SERVICE.into(), service.into());
1069    s.labels.insert(LABEL_SLOT.into(), slot.to_string());
1070    s.labels.insert(LABEL_REV.into(), rev.into());
1071    Ok(s)
1072}
1073
1074/// A stack's instance as listed.
1075#[derive(Debug, Clone)]
1076#[doc(hidden)]
1077pub struct Inst {
1078    pub name: String,
1079    pub(super) slot: u32,
1080    pub rev: String,
1081    status: String,
1082}
1083
1084impl Inst {
1085    #[doc(hidden)]
1086    pub fn is_running(&self) -> bool {
1087        self.running()
1088    }
1089
1090    fn running(&self) -> bool {
1091        self.status.eq_ignore_ascii_case("running")
1092    }
1093}
1094
1095/// A stack's instances (of one service), using incus' server-side filter.
1096#[doc(hidden)]
1097pub fn list_instances(client: &Client, stack: &str, service: Option<&str>) -> Result<Vec<Inst>> {
1098    let mut filter = format!("config.user.{LABEL_STACK} eq {stack}");
1099    if let Some(s) = service {
1100        filter.push_str(&format!(" and config.user.{LABEL_SERVICE} eq {s}"));
1101    }
1102    let v = client.get(&format!(
1103        "/1.0/instances?recursion=1&filter={}",
1104        encode_query(&filter)
1105    ))?;
1106    let mut out = Vec::new();
1107    for i in v.as_array().into_iter().flatten() {
1108        let info = crate::sandbox::SandboxInfo::from_api(i);
1109        let c = &info.config;
1110        // Filter again: an incus without filter support returns everything.
1111        if c.get(&format!("user.{LABEL_STACK}")).map(String::as_str) != Some(stack) {
1112            continue;
1113        }
1114        let svc = c
1115            .get(&format!("user.{LABEL_SERVICE}"))
1116            .cloned()
1117            .unwrap_or_default();
1118        if service.is_some_and(|s| s != svc) {
1119            continue;
1120        }
1121        out.push(Inst {
1122            name: info.name.clone(),
1123            slot: c
1124                .get(&format!("user.{LABEL_SLOT}"))
1125                .and_then(|s| s.parse().ok())
1126                .unwrap_or(0),
1127            rev: c
1128                .get(&format!("user.{LABEL_REV}"))
1129                .cloned()
1130                .unwrap_or_default(),
1131            status: info.status.clone(),
1132        });
1133    }
1134    out.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
1135    Ok(out)
1136}
1137
1138/// An instance's init pid (changes on every start) and its first global
1139/// address on any interface but loopback, IPv4 preferred.
1140fn instance_state(client: &Client, name: &str) -> Result<(i64, Option<IpAddr>)> {
1141    let v = client.get(&format!("/1.0/instances/{}/state", encode_segment(name)))?;
1142    let pid = v.get("pid").and_then(Value::as_i64).unwrap_or(0);
1143    let mut v4 = None;
1144    let mut v6 = None;
1145    if let Some(nets) = v.get("network").and_then(Value::as_object) {
1146        for (ifname, n) in nets {
1147            if ifname == "lo" {
1148                continue;
1149            }
1150            for a in n
1151                .get("addresses")
1152                .and_then(Value::as_array)
1153                .into_iter()
1154                .flatten()
1155            {
1156                if a.get("scope").and_then(Value::as_str) != Some("global") {
1157                    continue;
1158                }
1159                let Some(ip) = a
1160                    .get("address")
1161                    .and_then(Value::as_str)
1162                    .and_then(|s| s.parse::<IpAddr>().ok())
1163                else {
1164                    continue;
1165                };
1166                match ip {
1167                    IpAddr::V4(_) if v4.is_none() => v4 = Some(ip),
1168                    IpAddr::V6(_) if v6.is_none() => v6 = Some(ip),
1169                    _ => {}
1170                }
1171            }
1172        }
1173    }
1174    Ok((pid, v4.or(v6)))
1175}
1176
1177/// Per-instance memory of a worker.
1178#[derive(Debug, Default)]
1179struct InstRt {
1180    /// The init pid secrets and the unit were last set up for.
1181    pid: i64,
1182    /// When that pid was first seen: the start of `start_period`.
1183    since: Option<Instant>,
1184    ip: Option<IpAddr>,
1185    failures: u32,
1186    healthy: Option<bool>,
1187    next_probe: Option<Instant>,
1188    last_probe: String,
1189    /// App restarts for failing health, since it was last healthy.
1190    unhealthy_restarts: u32,
1191    /// Restarts counted against `restart_policy.max_attempts`.
1192    restarts: VecDeque<Instant>,
1193    in_rotation: bool,
1194}
1195
1196/// A service's worker.
1197struct Worker {
1198    inner: Arc<Inner>,
1199    /// The stack's own name: labels and instance names use it.
1200    stack: String,
1201    /// `org/stack`: the controller's key, and how events name it.
1202    q: String,
1203    /// A client on the stack's org (its incus project).
1204    oclient: Client,
1205    service: String,
1206    shared: Arc<WorkerShared>,
1207    rt: BTreeMap<String, InstRt>,
1208    /// The revision whose rollout failed and paused; not retried until the
1209    /// revision changes.
1210    paused: Option<(String, String)>,
1211    /// Backoff for creating into an empty slot that keeps failing.
1212    create_backoff: Option<(Instant, Duration)>,
1213    routes: BTreeMap<String, Published>,
1214    route_errors: BTreeMap<String, String>,
1215    /// The resolved spec of the current revision, for exec defaults.
1216    template: Option<(String, Desired)>,
1217    state: String,
1218    message: Option<String>,
1219    last_error: Option<String>,
1220    /// The instances as last listed, kept current through a rollout.
1221    insts: Vec<Inst>,
1222    rollout: Option<RolloutStatus>,
1223    /// Per route: accepted count at the last status, when, and conn/s history.
1224    rates: BTreeMap<String, (u64, Instant, VecDeque<f32>)>,
1225    /// `depends_on` was met once: from then on the service is reconciled
1226    /// whatever its dependencies do.
1227    deps_met: bool,
1228    org: OrgId,
1229    /// The addresses last published as the service's name.
1230    dns_last: Option<Vec<IpAddr>>,
1231    dns_error: Option<String>,
1232    /// The addresses last told to the observer.
1233    observed: Option<Vec<IpAddr>>,
1234    /// The service had a healthy replica at some point in this run.
1235    ever_healthy: bool,
1236    /// Since when no replica has been healthy.
1237    health_down: Option<Instant>,
1238    /// `health.unhealthy` was raised and not yet answered by `recovered`.
1239    health_alarm: bool,
1240    /// The definition the last pass ran on.
1241    seen: Option<Arc<StackDef>>,
1242}
1243
1244fn spawn_worker(inner: Arc<Inner>, def: &StackDef, service: String, shared: Arc<WorkerShared>) {
1245    let (stack, q, org) = (def.name.clone(), def.qualified(), def.org.clone());
1246    let oclient = crate::org::client(&inner.client, &def.org);
1247    let name = format!("isb-{q}-{service}");
1248    let r = std::thread::Builder::new().name(name).spawn(move || {
1249        let mut w = Worker::new(inner, stack, q, oclient, service, shared, org);
1250        w.run();
1251    });
1252    if let Err(e) = r {
1253        eprintln!("isb serve: cannot start a worker thread: {e}");
1254    }
1255}
1256
1257impl Worker {
1258    fn new(
1259        inner: Arc<Inner>,
1260        stack: String,
1261        q: String,
1262        oclient: Client,
1263        service: String,
1264        shared: Arc<WorkerShared>,
1265        org: OrgId,
1266    ) -> Worker {
1267        Worker {
1268            inner,
1269            stack,
1270            q,
1271            oclient,
1272            service,
1273            shared,
1274            rt: BTreeMap::new(),
1275            paused: None,
1276            create_backoff: None,
1277            routes: BTreeMap::new(),
1278            route_errors: BTreeMap::new(),
1279            template: None,
1280            state: "starting".into(),
1281            message: None,
1282            last_error: None,
1283            insts: Vec::new(),
1284            rollout: None,
1285            rates: BTreeMap::new(),
1286            deps_met: false,
1287            org,
1288            dns_last: None,
1289            dns_error: None,
1290            observed: None,
1291            ever_healthy: false,
1292            health_down: None,
1293            health_alarm: false,
1294            seen: None,
1295        }
1296    }
1297
1298    /// Drop the retry backoff when the instructions changed (a new
1299    /// deployment) or were kicked (an org's limits changed), so the next
1300    /// pass makes a fresh attempt, and forget the failure the old
1301    /// instructions ended in.
1302    fn begin_pass(&mut self, def: &Arc<StackDef>) {
1303        let new_def = self.seen.as_ref().is_none_or(|d| !Arc::ptr_eq(d, def));
1304        let kicked = self.shared.kick.swap(false, Ordering::SeqCst);
1305        if new_def {
1306            self.seen = Some(def.clone());
1307            self.last_error = None;
1308            if self.state == "failing" {
1309                self.state = "updating".into();
1310                self.message = None;
1311            }
1312        }
1313        if new_def || kicked {
1314            self.create_backoff = None;
1315        }
1316    }
1317}
1318
1319impl Worker {
1320    fn log(&self, msg: &str) {
1321        self.event("info", None, msg);
1322    }
1323
1324    fn event(&self, level: &str, instance: Option<&str>, msg: &str) {
1325        eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1326        self.inner
1327            .emit(level, &self.q, &self.service, instance, msg.to_string());
1328    }
1329
1330    fn client(&self) -> &Client {
1331        &self.oclient
1332    }
1333
1334    /// Raise `health.unhealthy` once a service that was healthy has had no
1335    /// healthy replica for [`HEALTH_DEBOUNCE`] (not while a rollout runs:
1336    /// it reports its own failure), and `health.recovered` when one is back.
1337    fn watch_health(&mut self, healthy: u32, replicas: u32) {
1338        if healthy > 0 {
1339            self.ever_healthy = true;
1340            self.health_down = None;
1341            if self.health_alarm {
1342                self.health_alarm = false;
1343                let msg = format!("{healthy} of {replicas} replicas healthy again");
1344                eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1345                self.inner.emit_kind(
1346                    Some("health.recovered"),
1347                    "info",
1348                    &self.q,
1349                    &self.service,
1350                    None,
1351                    msg,
1352                );
1353            }
1354            return;
1355        }
1356        if replicas == 0 || !self.ever_healthy || self.rollout.is_some() {
1357            self.health_down = None;
1358            return;
1359        }
1360        let since = *self.health_down.get_or_insert_with(Instant::now);
1361        if !self.health_alarm && since.elapsed() >= HEALTH_DEBOUNCE {
1362            self.health_alarm = true;
1363            let msg = format!(
1364                "no healthy replica (of {replicas}) for {}s",
1365                since.elapsed().as_secs()
1366            );
1367            eprintln!("isb serve: {}/{}: {msg}", self.q, self.service);
1368            self.inner.emit_kind(
1369                Some("health.unhealthy"),
1370                "error",
1371                &self.q,
1372                &self.service,
1373                None,
1374                msg,
1375            );
1376        }
1377    }
1378
1379    fn key(&self) -> (String, String) {
1380        (self.q.clone(), self.service.clone())
1381    }
1382
1383    fn run(&mut self) {
1384        loop {
1385            if self.shared.stop.load(Ordering::SeqCst) {
1386                return;
1387            }
1388            let (def, remove, remove_volumes) = {
1389                let s = self.shared.slot.lock().unwrap();
1390                (s.def.clone(), s.remove, s.remove_volumes)
1391            };
1392            if remove {
1393                self.teardown(&def, remove_volumes);
1394                // The service may have been deployed again meanwhile: then
1395                // this worker carries on instead of leaving it unattended.
1396                let mut ws = self.inner.workers.lock().unwrap();
1397                let slot = self.shared.slot.lock().unwrap();
1398                if slot.remove || self.shared.stop.load(Ordering::SeqCst) {
1399                    ws.remove(&self.key());
1400                    self.inner.status.lock().unwrap().remove(&self.key());
1401                    return;
1402                }
1403                continue;
1404            }
1405            self.begin_pass(&def);
1406            if let Err(e) = self.pass(&def) {
1407                self.state = "failing".into();
1408                self.message = Some(e.to_string());
1409                // Once per distinct error, not once per pass.
1410                if self.last_error.as_deref() != Some(&e.to_string()) {
1411                    self.event("error", None, &e.to_string());
1412                    self.last_error = Some(e.to_string());
1413                }
1414                self.publish_status(&def);
1415                let slot = self.shared.slot.lock().unwrap();
1416                if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
1417                    let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
1418                }
1419                continue;
1420            }
1421            self.last_error = None;
1422            let slot = self.shared.slot.lock().unwrap();
1423            if Arc::ptr_eq(&slot.def, &def) && !slot.remove {
1424                let _ = self.shared.wake.wait_timeout(slot, self.inner.interval);
1425            }
1426        }
1427    }
1428
1429    /// A deployment arrived (or removal was asked) since `def` was read.
1430    fn superseded(&self, def: &Arc<StackDef>) -> bool {
1431        let s = self.shared.slot.lock().unwrap();
1432        s.remove || !Arc::ptr_eq(&s.def, def) || self.shared.stop.load(Ordering::SeqCst)
1433    }
1434
1435    /// Delete this service's instances and routes, then leave.
1436    fn teardown(&mut self, def: &StackDef, volumes: bool) {
1437        // The name goes first, whoever published it.
1438        if let Some(dir) = Some(crate::discovery::org_dir(&self.org)).filter(|d| d.is_dir()) {
1439            if let Err(e) =
1440                crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &[])
1441            {
1442                self.log(&format!("cannot remove the service name: {e}"));
1443            }
1444        }
1445        self.dns_last = Some(Vec::new());
1446        if let Some(o) = &self.inner.observer {
1447            o.rotation(&self.q, &self.service, &[]);
1448        }
1449        self.observed = Some(Vec::new());
1450        for (k, _) in std::mem::take(&mut self.routes) {
1451            self.inner.balancer.remove_route(&k);
1452        }
1453        match list_instances(self.client(), &self.stack, Some(&self.service)) {
1454            Ok(insts) => {
1455                for i in insts {
1456                    self.log(&format!("removing {}", i.name));
1457                    if let Err(e) = Sandbox::remove(self.client(), &i.name, true) {
1458                        if !e.is_not_found() {
1459                            self.log(&format!("cannot remove {}: {e}", i.name));
1460                        }
1461                    }
1462                }
1463            }
1464            Err(e) => self.log(&format!("cannot list instances to remove: {e}")),
1465        }
1466        if volumes {
1467            self.remove_volumes(def);
1468        }
1469        self.rt.clear();
1470        self.template = None;
1471    }
1472
1473    fn remove_volumes(&self, def: &StackDef) {
1474        let Ok(spec) = def.service(&self.service) else {
1475            return;
1476        };
1477        let Ok(host) = crate::sandbox::host_facts(self.client()) else {
1478            return;
1479        };
1480        let Ok(pool) = host.pick_pool(spec.storage.as_deref()) else {
1481            return;
1482        };
1483        for v in &spec.volumes {
1484            if v.mount_type != crate::spec::MountType::Volume {
1485                continue;
1486            }
1487            let d = def.file.volumes.get(&v.source);
1488            if v.external || d.is_some_and(|d| d.external) {
1489                continue;
1490            }
1491            let name = d
1492                .and_then(|d| d.name.clone())
1493                .unwrap_or_else(|| v.source.clone());
1494            let vpool = match v.pool.as_deref().or(d.and_then(|d| d.pool.as_deref())) {
1495                Some(p) if p != "auto" => p.to_string(),
1496                _ => pool.clone(),
1497            };
1498            match crate::volume::remove(self.client(), &vpool, &name) {
1499                Ok(()) => self.log(&format!("volume {name}: deleted")),
1500                Err(e) if e.is_not_found() => {}
1501                // Another service of the stack may still be using it.
1502                Err(e) => self.log(&format!("volume {name}: kept ({e})")),
1503            }
1504        }
1505    }
1506
1507    /// One reconcile pass.
1508    #[expect(
1509        clippy::too_many_lines,
1510        reason = "predates the lint ratchet; split it when next changed"
1511    )]
1512    fn pass(&mut self, def: &Arc<StackDef>) -> Result<()> {
1513        let spec = def.service(&self.service)?.clone();
1514        let rev = def.revision(&self.service)?;
1515        let replicas = spec.replicas();
1516        let oci = crate::plan::ImageSource::parse(&spec.image)?.is_oci();
1517        let probe = spec.health_probe().map_err(Error::invalid)?;
1518
1519        if !self.deps_met {
1520            if let Some(msg) = self.waiting_for(&spec) {
1521                self.state = "waiting".into();
1522                self.message = Some(msg);
1523                self.publish_status(def);
1524                return Ok(());
1525            }
1526            self.deps_met = true;
1527        }
1528        if self.template.as_ref().is_none_or(|(r, _)| *r != rev) {
1529            let mut s = instance_spec(def, &self.service, &spec, 1, &rev)?;
1530            s.name = Some(instance_name(&self.stack, &self.service, 1, "0000")?);
1531            let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
1532            self.template = Some((rev.clone(), d));
1533        }
1534        self.set_routes(&spec);
1535
1536        let mut insts = list_instances(self.client(), &self.stack, Some(&self.service))?;
1537        self.rt.retain(|n, _| insts.iter().any(|i| i.name == *n));
1538
1539        // Scale down, highest slots first.
1540        let extra: Vec<Inst> = insts
1541            .iter()
1542            .filter(|i| i.slot > replicas || i.slot == 0)
1543            .cloned()
1544            .collect();
1545        for i in extra.iter().rev() {
1546            self.log(&format!("scaling down: removing {}", i.name));
1547            self.retire(&i.name)?;
1548        }
1549        insts.retain(|i| i.slot >= 1 && i.slot <= replicas);
1550        self.insts = insts.clone();
1551
1552        // Keep what exists running, set up and healthy.
1553        for i in &insts {
1554            self.maintain(def, i, &spec, oci, probe.as_ref())?;
1555        }
1556        // An interrupted rollout can leave an old instance next to a current
1557        // one in the same slot: once the current one serves, drop the old.
1558        for slot in 1..=replicas {
1559            let current_ok = insts.iter().any(|i| {
1560                i.slot == slot
1561                    && i.rev == rev
1562                    && self.rt.get(&i.name).is_some_and(|r| r.in_rotation)
1563            });
1564            if current_ok {
1565                for i in insts.iter().filter(|i| i.slot == slot && i.rev != rev) {
1566                    self.log(&format!("removing leftover {}", i.name));
1567                    self.retire(&i.name)?;
1568                }
1569            }
1570            let mut current: Vec<&Inst> = insts
1571                .iter()
1572                .filter(|i| i.slot == slot && i.rev == rev)
1573                .collect();
1574            // Two current instances in one slot (a crash mid-create): keep one.
1575            while current.len() > 1 {
1576                let i = current.pop().unwrap();
1577                self.log(&format!("removing duplicate {}", i.name));
1578                self.retire(&i.name)?;
1579            }
1580        }
1581        self.sync_routes();
1582        self.publish_status(def);
1583
1584        // Roll out: slots without a current instance.
1585        let mut pending: Vec<(u32, Option<String>)> = Vec::new();
1586        for slot in 1..=replicas {
1587            if insts.iter().any(|i| i.slot == slot && i.rev == rev) {
1588                continue;
1589            }
1590            let old = insts
1591                .iter()
1592                .find(|i| i.slot == slot)
1593                .map(|i| i.name.clone());
1594            pending.push((slot, old));
1595        }
1596        if pending.is_empty() {
1597            self.create_backoff = None;
1598            let all_ok = insts
1599                .iter()
1600                .all(|i| i.rev == rev && self.rt.get(&i.name).is_some_and(|r| r.in_rotation));
1601            self.state = if all_ok { "converged" } else { "failing" }.into();
1602            if all_ok {
1603                self.message = None;
1604            } else if self.message.is_none() {
1605                self.message = Some("some replicas are not healthy".into());
1606            }
1607            self.publish_status(def);
1608            return Ok(());
1609        }
1610        if self.paused.as_ref().is_some_and(|(r, _)| *r == rev) {
1611            self.state = "paused".into();
1612            self.message = self.paused.as_ref().map(|(_, m)| m.clone());
1613            self.publish_status(def);
1614            return Ok(());
1615        }
1616        if let Some((at, wait)) = self.create_backoff {
1617            if at.elapsed() < wait {
1618                return Ok(());
1619            }
1620        }
1621        self.state = "updating".into();
1622        self.message = None;
1623        self.publish_status(def);
1624
1625        let uc: UpdateConfig = spec
1626            .deploy
1627            .as_ref()
1628            .and_then(|d| d.update_config.clone())
1629            .unwrap_or_default();
1630        let parallel = match uc.parallelism.unwrap_or(1) {
1631            0 => pending.len(),
1632            n => n as usize,
1633        };
1634        let delay = uc
1635            .delay
1636            .as_deref()
1637            .map(crate::flex::parse_duration)
1638            .transpose()
1639            .map_err(Error::invalid)?
1640            .unwrap_or_default();
1641        let monitor = uc
1642            .monitor
1643            .as_deref()
1644            .map(crate::flex::parse_duration)
1645            .transpose()
1646            .map_err(Error::invalid)?
1647            .unwrap_or(Duration::from_secs(5));
1648        let order = uc.order.unwrap_or_default();
1649        let order_name = match order {
1650            UpdateOrder::StopFirst => "stop-first",
1651            UpdateOrder::StartFirst => "start-first",
1652        };
1653        self.rollout = Some(RolloutStatus {
1654            to_rev: rev.clone(),
1655            order: order_name.into(),
1656            parallelism: parallel,
1657            done: 0,
1658            total: pending.len(),
1659            started_at: now_secs(),
1660            slots: pending
1661                .iter()
1662                .map(|(slot, old)| SlotRollout {
1663                    slot: *slot,
1664                    old: old.clone(),
1665                    old_rev: old
1666                        .as_ref()
1667                        .and_then(|o| insts.iter().find(|i| i.name == *o))
1668                        .map(|i| i.rev.clone()),
1669                    old_state: if old.is_some() { "serving" } else { "none" }.into(),
1670                    new: None,
1671                    new_state: "waiting".into(),
1672                })
1673                .collect(),
1674        });
1675        let rollout_started = Instant::now();
1676        self.log(&format!(
1677            "rolling out rev {rev} to {} slot(s), {order_name}",
1678            pending.len()
1679        ));
1680        self.publish_status(def);
1681        let r = self.roll(
1682            def,
1683            &spec,
1684            &rev,
1685            &pending,
1686            parallel,
1687            delay,
1688            monitor,
1689            order,
1690            oci,
1691            probe.as_ref(),
1692            &uc,
1693        );
1694        let rollout = self.rollout.take();
1695        if let (Ok(true), Some(ro)) = (&r, rollout) {
1696            self.log(&format!(
1697                "rollout of rev {rev} complete: {}/{} slot(s) in {:.0?}",
1698                ro.done,
1699                ro.total,
1700                rollout_started.elapsed()
1701            ));
1702        }
1703        self.publish_status(def);
1704        r.map(|_| ())
1705    }
1706
1707    /// Replace the pending slots in batches. Ok(true) when every slot made
1708    /// it, Ok(false) when the rollout stopped (paused, rolled back, retrying).
1709    #[expect(clippy::too_many_arguments)]
1710    #[expect(
1711        clippy::excessive_nesting,
1712        reason = "predates the lint ratchet; split it when next changed"
1713    )]
1714    fn roll(
1715        &mut self,
1716        def: &Arc<StackDef>,
1717        spec: &SandboxSpec,
1718        rev: &String,
1719        pending: &[(u32, Option<String>)],
1720        parallel: usize,
1721        delay: Duration,
1722        monitor: Duration,
1723        order: UpdateOrder,
1724        oci: bool,
1725        probe: Option<&HealthProbe>,
1726        uc: &UpdateConfig,
1727    ) -> Result<bool> {
1728        for (n, batch) in pending.chunks(parallel).enumerate() {
1729            if n > 0 && !delay.is_zero() {
1730                std::thread::sleep(delay);
1731            }
1732            if self.superseded(def) {
1733                return Ok(false);
1734            }
1735            for (slot, old) in batch {
1736                let r = self.replace(
1737                    def,
1738                    spec,
1739                    rev,
1740                    *slot,
1741                    old.as_deref(),
1742                    order,
1743                    oci,
1744                    probe,
1745                    monitor,
1746                );
1747                if let Err(e) = r {
1748                    let msg = format!("slot {slot}: {e}");
1749                    self.event("error", None, &msg);
1750                    if old.is_none() {
1751                        // Nothing to protect: keep trying, slower each time.
1752                        let wait = self
1753                            .create_backoff
1754                            .map(|(_, w)| (w * 2).min(Duration::from_secs(300)))
1755                            .unwrap_or(Duration::from_secs(10));
1756                        self.create_backoff = Some((Instant::now(), wait));
1757                        self.state = "failing".into();
1758                        self.message = Some(format!("{msg}; retrying in {wait:?}"));
1759                        self.publish_status(def);
1760                        return Ok(false);
1761                    }
1762                    match uc.failure_action.unwrap_or_default() {
1763                        FailureAction::Continue => continue,
1764                        FailureAction::Pause => {
1765                            self.event("warn", None, &format!("rollout of rev {rev} paused"));
1766                            self.paused = Some((rev.clone(), format!("rollout paused: {msg}")));
1767                            return Ok(false);
1768                        }
1769                        FailureAction::Rollback => {
1770                            self.paused = Some((rev.clone(), format!("rolled back: {msg}")));
1771                            let ctl = Controller {
1772                                inner: self.inner.clone(),
1773                            };
1774                            self.event(
1775                                "warn",
1776                                None,
1777                                &format!("rollout of rev {rev} failed; rolling back"),
1778                            );
1779                            if let Err(e) = ctl.rollback(&self.q) {
1780                                self.event("error", None, &format!("rollback failed: {e}"));
1781                            }
1782                            return Ok(false);
1783                        }
1784                    }
1785                }
1786            }
1787        }
1788        Ok(true)
1789    }
1790
1791    /// Update one slot of the rollout display, and publish it.
1792    fn slot_state(
1793        &mut self,
1794        def: &StackDef,
1795        slot: u32,
1796        old_state: Option<&str>,
1797        new: Option<&str>,
1798        new_state: Option<&str>,
1799    ) {
1800        if let Some(ro) = &mut self.rollout {
1801            if let Some(s) = ro.slots.iter_mut().find(|s| s.slot == slot) {
1802                if let Some(o) = old_state {
1803                    s.old_state = o.into();
1804                }
1805                if let Some(n) = new {
1806                    s.new = Some(n.into());
1807                }
1808                if let Some(n) = new_state {
1809                    s.new_state = n.into();
1810                    if n == "serving" {
1811                        ro.done += 1;
1812                    }
1813                }
1814            }
1815        }
1816        self.publish_status(def);
1817    }
1818
1819    /// Unmet `depends_on`, as a message.
1820    fn waiting_for(&self, spec: &SandboxSpec) -> Option<String> {
1821        let st = self.inner.status.lock().unwrap();
1822        for (dep, d) in &spec.depends_on {
1823            let s = st.get(&(self.q.clone(), dep.clone()));
1824            let ok = match d.condition {
1825                DependCondition::ServiceStarted => s.is_some_and(|s| s.running > 0),
1826                DependCondition::ServiceHealthy => s.is_some_and(|s| s.healthy > 0),
1827            };
1828            if !ok {
1829                return Some(format!("waiting for {dep} ({:?})", d.condition));
1830            }
1831        }
1832        None
1833    }
1834
1835    fn handle(&self, name: &str) -> Sandbox {
1836        let (_, d) = self.template.as_ref().expect("resolved in pass");
1837        Sandbox::like(self.client(), name, d)
1838    }
1839
1840    /// Keep one instance running, set up after every boot, health-checked,
1841    /// and in or out of rotation.
1842    #[expect(
1843        clippy::too_many_lines,
1844        reason = "predates the lint ratchet; split it when next changed"
1845    )]
1846    fn maintain(
1847        &mut self,
1848        def: &StackDef,
1849        i: &Inst,
1850        spec: &SandboxSpec,
1851        oci: bool,
1852        probe: Option<&HealthProbe>,
1853    ) -> Result<()> {
1854        let policy = spec.deploy.as_ref().and_then(|d| d.restart_policy.clone());
1855        let condition = policy
1856            .as_ref()
1857            .and_then(|p| p.condition)
1858            .unwrap_or_default();
1859        let sb = self.handle(&i.name);
1860        if !i.running() {
1861            self.set_rotation(&i.name, false);
1862            if condition == RestartCondition::None {
1863                return Ok(());
1864            }
1865            if self.restart_budget_spent(&i.name, policy.as_ref()) {
1866                self.message = Some(format!("{}: restart limit reached", i.name));
1867                return Ok(());
1868            }
1869            self.event(
1870                "warn",
1871                Some(&i.name),
1872                &format!("{} is {}; starting it", i.name, i.status),
1873            );
1874            self.count_restart(&i.name);
1875            if let Err(e) = sb.start() {
1876                self.log(&format!("cannot start {}: {e}", i.name));
1877                return Ok(());
1878            }
1879        }
1880        let (pid, ip) = match instance_state(self.client(), &i.name) {
1881            Ok(s) => s,
1882            Err(e) if e.is_not_found() => return Ok(()),
1883            Err(e) => return Err(e),
1884        };
1885        let rt = self.rt.entry(i.name.clone()).or_default();
1886        rt.ip = ip;
1887        if rt.pid != pid {
1888            // A (re)boot: /run/secrets is a fresh tmpfs, and the unit may be
1889            // from an older isb. Probing starts over.
1890            rt.pid = pid;
1891            rt.since = Some(Instant::now());
1892            rt.failures = 0;
1893            rt.healthy = None;
1894            rt.next_probe = None;
1895            let r = self.setup(def, &sb, spec, oci);
1896            if let Err(e) = r {
1897                // Leave pid unset so the next pass tries again.
1898                if let Some(rt) = self.rt.get_mut(&i.name) {
1899                    rt.pid = 0;
1900                }
1901                self.set_rotation(&i.name, false);
1902                return Err(Error::OperationFailed {
1903                    step: format!("set up {}", i.name),
1904                    message: e.to_string(),
1905                });
1906            }
1907        }
1908        let alive = self.alive(&sb, spec, oci);
1909        let healthy = match probe {
1910            None => alive,
1911            Some(p) => {
1912                let rt = self.rt.get_mut(&i.name).unwrap();
1913                let since = rt.since.unwrap_or_else(Instant::now);
1914                let in_start = since.elapsed() < p.start_period;
1915                if rt.next_probe.is_none_or(|t| Instant::now() >= t) {
1916                    let r = supervise::probe(&sb, p);
1917                    let rt = self.rt.get_mut(&i.name).unwrap();
1918                    rt.last_probe = r.output.clone();
1919                    if r.ok {
1920                        rt.failures = 0;
1921                        rt.healthy = Some(true);
1922                        rt.unhealthy_restarts = 0;
1923                    } else if !in_start {
1924                        rt.failures += 1;
1925                        if rt.failures >= p.retries {
1926                            rt.healthy = Some(false);
1927                        }
1928                    }
1929                    let wait = if rt.healthy.is_none() {
1930                        p.start_interval
1931                    } else {
1932                        p.interval
1933                    };
1934                    rt.next_probe = Some(Instant::now() + wait);
1935                }
1936                self.rt[&i.name].healthy == Some(true) && alive
1937            }
1938        };
1939        self.set_rotation(&i.name, healthy);
1940        let unhealthy = self.rt[&i.name].healthy == Some(false);
1941        if unhealthy && condition != RestartCondition::None {
1942            if self.restart_budget_spent(&i.name, policy.as_ref()) {
1943                self.message = Some(format!("{}: unhealthy, restart limit reached", i.name));
1944                return Ok(());
1945            }
1946            let n = self.rt[&i.name].unhealthy_restarts;
1947            if n >= RESTARTS_BEFORE_REPLACE {
1948                self.log(&format!(
1949                    "{} stayed unhealthy through {n} restarts; replacing it",
1950                    i.name
1951                ));
1952                let rt = self.rt.get_mut(&i.name).unwrap();
1953                rt.unhealthy_restarts = 0;
1954                drop(sb);
1955                self.retire(&i.name)?;
1956                return Ok(());
1957            }
1958            self.event(
1959                "warn",
1960                Some(&i.name),
1961                &format!(
1962                    "{} is unhealthy ({}); restarting its app",
1963                    i.name,
1964                    self.rt[&i.name]
1965                        .last_probe
1966                        .lines()
1967                        .last()
1968                        .unwrap_or("probe failed")
1969                ),
1970            );
1971            self.count_restart(&i.name);
1972            let rt = self.rt.get_mut(&i.name).unwrap();
1973            rt.unhealthy_restarts += 1;
1974            rt.failures = 0;
1975            rt.healthy = None;
1976            rt.since = Some(Instant::now());
1977            supervise::restart_app(&sb, &self.service, oci)?;
1978        }
1979        Ok(())
1980    }
1981
1982    /// Secrets and the app's unit, after a boot or on creation.
1983    fn setup(&self, def: &StackDef, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> Result<()> {
1984        // Read from the store now, never from the definition.
1985        let keys = spec.secret_keys();
1986        let values = if keys.is_empty() {
1987            BTreeMap::new()
1988        } else {
1989            super::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?
1990        };
1991        // An OCI app started before its files arrived: restart it once so
1992        // it reads them (the next pass finds them in place).
1993        if supervise::push_secrets(sb, spec, &values)? && oci {
1994            supervise::restart_app(sb, &self.service, oci)?;
1995        }
1996        if spec.command.is_some() && !oci {
1997            let mut s = spec.clone();
1998            s.restart = Some(RestartMode::Always);
1999            let env = supervise::secret_env(spec, &values)?;
2000            supervise::install(sb, &self.service, &s, !spec.secrets.is_empty(), &env)?;
2001        }
2002        Ok(())
2003    }
2004
2005    /// An OCI instance's secret variables, for its config (`environment.KEY`).
2006    fn oci_secret_env(
2007        &self,
2008        def: &StackDef,
2009        spec: &SandboxSpec,
2010    ) -> Result<BTreeMap<String, String>> {
2011        if spec.env.secrets.is_empty() {
2012            return Ok(BTreeMap::new());
2013        }
2014        let values = super::secrets::values(
2015            &self.inner.secrets,
2016            &def.org,
2017            &def.secrets,
2018            spec.env.secrets.values().map(String::as_str),
2019        )?;
2020        supervise::secret_env(spec, &values)
2021    }
2022
2023    /// The app's process is up: its unit is active, or (OCI, or no command)
2024    /// the instance runs.
2025    fn alive(&self, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> bool {
2026        if spec.command.is_some() && !oci {
2027            return supervise::unit_state(sb, &self.service).is_ok_and(|s| s == "active");
2028        }
2029        sb.info()
2030            .is_ok_and(|i| i.status.eq_ignore_ascii_case("running"))
2031    }
2032
2033    fn count_restart(&mut self, name: &str) {
2034        self.rt
2035            .entry(name.to_string())
2036            .or_default()
2037            .restarts
2038            .push_back(Instant::now());
2039    }
2040
2041    fn restart_budget_spent(
2042        &mut self,
2043        name: &str,
2044        policy: Option<&crate::spec::RestartPolicy>,
2045    ) -> bool {
2046        let Some(p) = policy else { return false };
2047        let Some(max) = p.max_attempts else {
2048            return false;
2049        };
2050        let window = p
2051            .window
2052            .as_deref()
2053            .and_then(|w| crate::flex::parse_duration(w).ok());
2054        let rt = self.rt.entry(name.to_string()).or_default();
2055        if let Some(w) = window {
2056            while rt.restarts.front().is_some_and(|t| t.elapsed() > w) {
2057                rt.restarts.pop_front();
2058            }
2059        }
2060        rt.restarts.len() as u32 >= max
2061    }
2062
2063    /// Replace (or fill) one slot with an instance of the current revision.
2064    #[expect(clippy::too_many_arguments)]
2065    fn replace(
2066        &mut self,
2067        def: &StackDef,
2068        spec: &SandboxSpec,
2069        rev: &str,
2070        slot: u32,
2071        old: Option<&str>,
2072        order: UpdateOrder,
2073        oci: bool,
2074        probe: Option<&HealthProbe>,
2075        monitor: Duration,
2076    ) -> Result<()> {
2077        // Before anything is stopped: a secret that cannot be read leaves
2078        // the old instance serving.
2079        let secret_env = if oci {
2080            self.oci_secret_env(def, spec)?
2081        } else {
2082            BTreeMap::new()
2083        };
2084        if let (Some(o), UpdateOrder::StopFirst) = (old, order) {
2085            self.log(&format!("slot {slot}: replacing {o} (stop-first)"));
2086            self.slot_state(def, slot, Some("draining"), None, None);
2087            self.retire(o)?;
2088            self.slot_state(def, slot, Some("retired"), None, None);
2089        }
2090        let name = instance_name(&self.stack, &self.service, slot, &new_id())?;
2091        self.log(&format!("slot {slot}: creating {name} (rev {rev})"));
2092        self.slot_state(def, slot, None, Some(&name), Some("creating"));
2093        let mut s = instance_spec(def, &self.service, spec, slot, rev)?;
2094        s.name = Some(name.clone());
2095        // `env.secrets` stays set, so the values are redacted in reports.
2096        s.env.vars.extend(secret_env);
2097        let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
2098        let stack = self.q.clone();
2099        let mut report = |m: &str| eprintln!("isb serve: {stack}: {m}");
2100        let created =
2101            crate::sandbox::ensure(self.client(), &d, EnsureOptions::default(), &mut report);
2102        let result = created.and_then(|_| {
2103            let inst = Inst {
2104                name: name.clone(),
2105                slot,
2106                rev: rev.to_string(),
2107                status: "Running".into(),
2108            };
2109            self.insts.push(inst.clone());
2110            self.slot_state(def, slot, None, None, Some("probing"));
2111            self.wait_serving(def, &inst, spec, oci, probe, monitor)
2112        });
2113        if let Err(e) = result {
2114            // Read its output before it is deleted: it explains the failure.
2115            let (e, attempt) =
2116                super::failure::explain(self.client(), &name, &self.service, oci, e, now_ms());
2117            self.inner.failures.record(&self.q, &self.service, attempt);
2118            self.event(
2119                "error",
2120                Some(&name),
2121                &format!("{name} did not come up: {e}"),
2122            );
2123            self.slot_state(def, slot, None, None, Some("failed"));
2124            let _ = self.retire(&name);
2125            return Err(e);
2126        }
2127        if let (Some(o), UpdateOrder::StartFirst) = (old, order) {
2128            self.log(&format!("slot {slot}: {name} is serving; retiring {o}"));
2129            self.slot_state(def, slot, Some("draining"), None, None);
2130            self.retire(o)?;
2131            self.slot_state(def, slot, Some("retired"), None, None);
2132        }
2133        self.slot_state(def, slot, None, None, Some("serving"));
2134        Ok(())
2135    }
2136
2137    /// Wait for a new instance to pass its health (or, without a
2138    /// healthcheck, for its app to run), put it in rotation, then watch it
2139    /// for `monitor`.
2140    fn wait_serving(
2141        &mut self,
2142        def: &StackDef,
2143        i: &Inst,
2144        spec: &SandboxSpec,
2145        oci: bool,
2146        probe: Option<&HealthProbe>,
2147        monitor: Duration,
2148    ) -> Result<()> {
2149        let deadline = match probe {
2150            Some(p) => p.start_period + p.interval * p.retries + Duration::from_secs(30),
2151            None => Duration::from_secs(60),
2152        }
2153        .max(Duration::from_secs(60));
2154        let started = Instant::now();
2155        loop {
2156            self.maintain(def, i, spec, oci, probe)?;
2157            if self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2158                break;
2159            }
2160            if self
2161                .rt
2162                .get(&i.name)
2163                .is_some_and(|r| r.healthy == Some(false))
2164            {
2165                return Err(Error::invalid(format!(
2166                    "unhealthy: {}",
2167                    self.rt[&i.name].last_probe
2168                )));
2169            }
2170            if started.elapsed() > deadline {
2171                let why = self
2172                    .rt
2173                    .get(&i.name)
2174                    .map(|r| r.last_probe.clone())
2175                    .filter(|s| !s.is_empty())
2176                    .unwrap_or_else(|| "its app is not running".into());
2177                return Err(Error::invalid(format!(
2178                    "not serving after {:?}: {why}",
2179                    started.elapsed()
2180                )));
2181            }
2182            // Probe on the start interval rather than waiting for the next one.
2183            if let Some(rt) = self.rt.get_mut(&i.name) {
2184                rt.next_probe = None;
2185            }
2186            std::thread::sleep(
2187                probe
2188                    .map(|p| p.start_interval)
2189                    .unwrap_or(Duration::from_secs(1))
2190                    .min(Duration::from_secs(2)),
2191            );
2192        }
2193        self.sync_routes();
2194        self.slot_state(def, i.slot, None, None, Some("monitoring"));
2195        let watch = Instant::now();
2196        while watch.elapsed() < monitor {
2197            std::thread::sleep(Duration::from_secs(1).min(monitor));
2198            if let Some(rt) = self.rt.get_mut(&i.name) {
2199                rt.next_probe = None;
2200            }
2201            self.maintain(def, i, spec, oci, probe)?;
2202            if !self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2203                return Err(Error::invalid(format!(
2204                    "failed within the {monitor:?} monitor period"
2205                )));
2206            }
2207        }
2208        Ok(())
2209    }
2210
2211    /// Take an instance out of rotation, let its connections drain, then
2212    /// delete it.
2213    fn retire(&mut self, name: &str) -> Result<()> {
2214        let ip = self.rt.get(name).and_then(|r| r.ip);
2215        self.set_rotation(name, false);
2216        self.sync_routes();
2217        if let Some(ip) = ip {
2218            let started = Instant::now();
2219            if let Some(o) = &self.inner.observer {
2220                o.drain(&self.q, &self.service, ip, DRAIN);
2221            }
2222            let left = DRAIN.saturating_sub(started.elapsed());
2223            for (k, p) in &self.routes {
2224                self.inner
2225                    .balancer
2226                    .wait_drained(k, SocketAddr::new(ip, p.target), left);
2227            }
2228        }
2229        if let Ok(sb) = Sandbox::get(self.client(), name) {
2230            let _ = sb.stop(false, Duration::from_secs(10));
2231        }
2232        self.rt.remove(name);
2233        self.insts.retain(|i| i.name != name);
2234        match Sandbox::remove(self.client(), name, true) {
2235            Err(e) if !e.is_not_found() => Err(e),
2236            _ => Ok(()),
2237        }
2238    }
2239
2240    fn set_rotation(&mut self, name: &str, on: bool) {
2241        let rt = self.rt.entry(name.to_string()).or_default();
2242        if rt.in_rotation != on {
2243            rt.in_rotation = on;
2244            self.sync_routes();
2245        }
2246    }
2247
2248    /// Bring the balancer's routes in line with the spec's published ports.
2249    fn set_routes(&mut self, spec: &SandboxSpec) {
2250        let want: BTreeMap<String, Published> = match published(spec) {
2251            Ok(ps) => ps
2252                .into_iter()
2253                .map(|p| (format!("{}/{}/{}", self.q, self.service, p.display()), p))
2254                .collect(),
2255            Err(e) => {
2256                self.message = Some(e.to_string());
2257                BTreeMap::new()
2258            }
2259        };
2260        let stale: Vec<String> = self
2261            .routes
2262            .keys()
2263            .filter(|k| !want.contains_key(*k))
2264            .cloned()
2265            .collect();
2266        for k in stale {
2267            self.inner.balancer.remove_route(&k);
2268            self.routes.remove(&k);
2269            self.route_errors.remove(&k);
2270        }
2271        for (k, p) in want {
2272            self.routes.entry(k).or_insert(p);
2273        }
2274        self.sync_routes();
2275    }
2276
2277    fn sync_routes(&mut self) {
2278        let mut errors = BTreeMap::new();
2279        // UDP ports are proxy devices on the replica (see super::ports).
2280        for (k, p) in self.routes.iter().filter(|(_, p)| !p.udp) {
2281            let backends: Vec<SocketAddr> = self
2282                .rt
2283                .values()
2284                .filter(|r| r.in_rotation)
2285                .filter_map(|r| r.ip)
2286                .map(|ip| SocketAddr::new(ip, p.target))
2287                .collect();
2288            if let Err(e) = self.inner.balancer.set_route(k, p.listen, backends) {
2289                errors.insert(k.clone(), e.to_string());
2290            }
2291        }
2292        for (k, e) in &errors {
2293            if self.route_errors.get(k) != Some(e) {
2294                self.log(&format!("cannot publish {k}: {e}"));
2295            }
2296        }
2297        self.route_errors = errors;
2298        self.sync_observer();
2299        self.sync_dns();
2300    }
2301
2302    /// Tell the observer (the ingress) when the in-rotation set changed.
2303    fn sync_observer(&mut self) {
2304        let Some(o) = self.inner.observer.clone() else {
2305            return;
2306        };
2307        let mut ips: Vec<IpAddr> = self
2308            .rt
2309            .values()
2310            .filter(|r| r.in_rotation)
2311            .filter_map(|r| r.ip)
2312            .collect();
2313        ips.sort();
2314        if self.observed.as_ref() != Some(&ips) {
2315            o.rotation(&self.q, &self.service, &ips);
2316            self.observed = Some(ips);
2317        }
2318    }
2319
2320    /// Publish the in-rotation replicas' addresses as the service's name
2321    /// (see [`crate::discovery`]). Nothing to do in an org created without
2322    /// service names.
2323    fn sync_dns(&mut self) {
2324        let dir = crate::discovery::org_dir(&self.org);
2325        let mut ips: Vec<IpAddr> = self
2326            .rt
2327            .values()
2328            .filter(|r| r.in_rotation)
2329            .filter_map(|r| r.ip)
2330            .collect();
2331        ips.sort();
2332        // A worker that has published nothing yet (a daemon restart) leaves
2333        // the last records alone until a replica is back in rotation, rather
2334        // than blanking the name while health is being re-established.
2335        let fresh = self.dns_last.is_none() && ips.is_empty();
2336        if fresh || self.dns_last.as_ref() == Some(&ips) || !dir.is_dir() {
2337            return;
2338        }
2339        match crate::discovery::publish(&dir, &self.org, &self.stack, &self.service, &ips) {
2340            Ok(()) => {
2341                self.dns_last = Some(ips);
2342                self.dns_error = None;
2343            }
2344            Err(e) => {
2345                let e = e.to_string();
2346                if self.dns_error.as_deref() != Some(&e) {
2347                    self.event(
2348                        "warn",
2349                        None,
2350                        &format!("cannot publish the service name: {e}"),
2351                    );
2352                    self.dns_error = Some(e);
2353                }
2354            }
2355        }
2356    }
2357
2358    #[expect(
2359        clippy::too_many_lines,
2360        reason = "predates the lint ratchet; split it when next changed"
2361    )]
2362    fn publish_status(&mut self, def: &StackDef) {
2363        // An address can change without a rotation change (a restart).
2364        self.sync_observer();
2365        self.sync_dns();
2366        let Ok(spec) = def.service(&self.service) else {
2367            return;
2368        };
2369        let rev = def.revision(&self.service).unwrap_or_default();
2370        let probe = matches!(spec.health_probe(), Ok(Some(_)));
2371        let snap = self.inner.snapshot.lock().unwrap().instances.clone();
2372        let instances: Vec<InstanceStatus> = self
2373            .insts
2374            .iter()
2375            .map(|i| {
2376                let rt = self.rt.get(&i.name);
2377                let m = snap.get(&format!("{}/{}", self.oclient.project_name(), i.name));
2378                let health = match (probe, rt.and_then(|r| r.healthy)) {
2379                    (false, _) => "none",
2380                    (true, Some(true)) => "healthy",
2381                    (true, Some(false)) => "unhealthy",
2382                    (true, None) => "starting",
2383                };
2384                InstanceStatus {
2385                    name: i.name.clone(),
2386                    slot: i.slot,
2387                    rev: i.rev.clone(),
2388                    status: m
2389                        .map(|m| m.status.clone())
2390                        .unwrap_or_else(|| i.status.clone()),
2391                    health: health.into(),
2392                    ip: rt.and_then(|r| r.ip).map(|ip| ip.to_string()),
2393                    in_rotation: rt.is_some_and(|r| r.in_rotation),
2394                    restarts: rt.map(|r| r.restarts.len() as u32).unwrap_or(0),
2395                    last_probe: rt.map(|r| r.last_probe.clone()).unwrap_or_default(),
2396                    cpu_pct: m.and_then(|m| m.cpu_pct),
2397                    cpu_history: m.map(|m| m.cpu_history.clone()).unwrap_or_default(),
2398                    mem_bytes: m.and_then(|m| m.mem_bytes),
2399                    disk_bytes: m.and_then(|m| m.disk_bytes),
2400                }
2401            })
2402            .collect();
2403        let mut instances = instances;
2404        instances.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
2405        let routes = self.inner.balancer.routes();
2406        let now = Instant::now();
2407        let mut ports = Vec::new();
2408        for (k, p) in &self.routes {
2409            let r = routes.iter().find(|r| r.key == *k);
2410            let accepted = r.map(|r| r.accepted).unwrap_or(0);
2411            let e = self
2412                .rates
2413                .entry(k.clone())
2414                .or_insert_with(|| (accepted, now, VecDeque::new()));
2415            let dt = now.duration_since(e.1).as_secs_f32();
2416            // Only sample on a reconcile-sized step, so a burst of status
2417            // updates during a rollout does not flatten the curve.
2418            if dt >= 1.0 {
2419                let rate = accepted.saturating_sub(e.0) as f32 / dt;
2420                if e.2.len() == crate::metrics::HISTORY {
2421                    e.2.pop_front();
2422                }
2423                e.2.push_back(rate);
2424                e.0 = accepted;
2425                e.1 = now;
2426            }
2427            ports.push(PortStatus {
2428                listen: p.display(),
2429                target: p.target,
2430                backends: match r {
2431                    _ if p.udp => {
2432                        let ips = self.rt.values().filter_map(|r| r.ip);
2433                        ips.map(|ip| SocketAddr::new(ip, p.target).to_string())
2434                            .collect()
2435                    }
2436                    Some(r) => r.backends.iter().map(|b| b.addr.to_string()).collect(),
2437                    None => Vec::new(),
2438                },
2439                error: self.route_errors.get(k).cloned(),
2440                accepted,
2441                active: r
2442                    .map(|r| {
2443                        r.backends
2444                            .iter()
2445                            .chain(r.draining.iter())
2446                            .map(|b| b.active)
2447                            .sum()
2448                    })
2449                    .unwrap_or(0),
2450                rate_history: e.2.iter().copied().collect(),
2451            });
2452        }
2453        self.rates.retain(|k, _| self.routes.contains_key(k));
2454        let running = instances
2455            .iter()
2456            .filter(|i| i.status.eq_ignore_ascii_case("running"))
2457            .count() as u32;
2458        let healthy = instances.iter().filter(|i| i.in_rotation).count() as u32;
2459        self.watch_health(healthy, spec.replicas());
2460        let st = ServiceStatus {
2461            service: self.service.clone(),
2462            image: spec.image.clone(),
2463            rev,
2464            replicas: spec.replicas(),
2465            running,
2466            healthy,
2467            state: self.state.clone(),
2468            message: self.message.clone(),
2469            instances,
2470            ports,
2471            rollout: self.rollout.clone(),
2472            checked_at: now_secs(),
2473            domains: Vec::new(),
2474        };
2475        self.inner.status.lock().unwrap().insert(self.key(), st);
2476    }
2477}
2478
2479/// Names of the services that are not converged, for a deploy waiting on
2480/// its rollout.
2481pub fn unsettled(st: &StackStatus) -> BTreeSet<String> {
2482    st.services
2483        .iter()
2484        .filter(|s| s.state != "converged")
2485        .map(|s| s.service.clone())
2486        .collect()
2487}
2488
2489#[cfg(test)]
2490#[path = "controller_tests.rs"]
2491mod tests;