Skip to main content

isb_core/stack/
controller.rs

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