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