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 wait = self
1654                            .create_backoff
1655                            .map(|(_, w)| (w * 2).min(Duration::from_secs(300)))
1656                            .unwrap_or(Duration::from_secs(10));
1657                        self.create_backoff = Some((Instant::now(), wait));
1658                        self.state = "failing".into();
1659                        self.message = Some(format!("{msg}; retrying in {wait:?}"));
1660                        self.publish_status(def);
1661                        return Ok(false);
1662                    }
1663                    match uc.failure_action.unwrap_or_default() {
1664                        FailureAction::Continue => continue,
1665                        FailureAction::Pause => {
1666                            self.event("warn", None, &format!("rollout of rev {rev} paused"));
1667                            self.paused = Some((rev.clone(), format!("rollout paused: {msg}")));
1668                            return Ok(false);
1669                        }
1670                        FailureAction::Rollback => {
1671                            self.paused = Some((rev.clone(), format!("rolled back: {msg}")));
1672                            let ctl = Controller {
1673                                inner: self.inner.clone(),
1674                            };
1675                            self.event(
1676                                "warn",
1677                                None,
1678                                &format!("rollout of rev {rev} failed; rolling back"),
1679                            );
1680                            if let Err(e) = ctl.rollback(&self.q) {
1681                                self.event("error", None, &format!("rollback failed: {e}"));
1682                            }
1683                            return Ok(false);
1684                        }
1685                    }
1686                }
1687            }
1688        }
1689        Ok(true)
1690    }
1691
1692    /// Update one slot of the rollout display, and publish it.
1693    fn slot_state(
1694        &mut self,
1695        def: &StackDef,
1696        slot: u32,
1697        old_state: Option<&str>,
1698        new: Option<&str>,
1699        new_state: Option<&str>,
1700    ) {
1701        if let Some(ro) = &mut self.rollout {
1702            if let Some(s) = ro.slots.iter_mut().find(|s| s.slot == slot) {
1703                if let Some(o) = old_state {
1704                    s.old_state = o.into();
1705                }
1706                if let Some(n) = new {
1707                    s.new = Some(n.into());
1708                }
1709                if let Some(n) = new_state {
1710                    s.new_state = n.into();
1711                    if n == "serving" {
1712                        ro.done += 1;
1713                    }
1714                }
1715            }
1716        }
1717        self.publish_status(def);
1718    }
1719
1720    /// Unmet `depends_on`, as a message.
1721    fn waiting_for(&self, spec: &SandboxSpec) -> Option<String> {
1722        let st = self.inner.status.lock().unwrap();
1723        for (dep, d) in &spec.depends_on {
1724            let s = st.get(&(self.q.clone(), dep.clone()));
1725            let ok = match d.condition {
1726                DependCondition::ServiceStarted => s.is_some_and(|s| s.running > 0),
1727                DependCondition::ServiceHealthy => s.is_some_and(|s| s.healthy > 0),
1728            };
1729            if !ok {
1730                return Some(format!("waiting for {dep} ({:?})", d.condition));
1731            }
1732        }
1733        None
1734    }
1735
1736    fn handle(&self, name: &str) -> Sandbox {
1737        let (_, d) = self.template.as_ref().expect("resolved in pass");
1738        Sandbox::like(self.client(), name, d)
1739    }
1740
1741    /// Keep one instance running, set up after every boot, health-checked,
1742    /// and in or out of rotation.
1743    #[expect(
1744        clippy::too_many_lines,
1745        reason = "predates the lint ratchet; split it when next changed"
1746    )]
1747    fn maintain(
1748        &mut self,
1749        def: &StackDef,
1750        i: &Inst,
1751        spec: &SandboxSpec,
1752        oci: bool,
1753        probe: Option<&HealthProbe>,
1754    ) -> Result<()> {
1755        let policy = spec.deploy.as_ref().and_then(|d| d.restart_policy.clone());
1756        let condition = policy
1757            .as_ref()
1758            .and_then(|p| p.condition)
1759            .unwrap_or_default();
1760        let sb = self.handle(&i.name);
1761        if !i.running() {
1762            self.set_rotation(&i.name, false);
1763            if condition == RestartCondition::None {
1764                return Ok(());
1765            }
1766            if self.restart_budget_spent(&i.name, policy.as_ref()) {
1767                self.message = Some(format!("{}: restart limit reached", i.name));
1768                return Ok(());
1769            }
1770            self.event(
1771                "warn",
1772                Some(&i.name),
1773                &format!("{} is {}; starting it", i.name, i.status),
1774            );
1775            self.count_restart(&i.name);
1776            if let Err(e) = sb.start() {
1777                self.log(&format!("cannot start {}: {e}", i.name));
1778                return Ok(());
1779            }
1780        }
1781        let (pid, ip) = match instance_state(self.client(), &i.name) {
1782            Ok(s) => s,
1783            Err(e) if e.is_not_found() => return Ok(()),
1784            Err(e) => return Err(e),
1785        };
1786        let rt = self.rt.entry(i.name.clone()).or_default();
1787        rt.ip = ip;
1788        if rt.pid != pid {
1789            // A (re)boot: /run/secrets is a fresh tmpfs, and the unit may be
1790            // from an older isb. Probing starts over.
1791            let restarted = rt.pid != 0;
1792            rt.pid = pid;
1793            rt.since = Some(Instant::now());
1794            rt.failures = 0;
1795            rt.healthy = None;
1796            rt.next_probe = None;
1797            let r = self.setup(def, &sb, spec, oci);
1798            if let Err(e) = r {
1799                // Leave pid unset so the next pass tries again.
1800                if let Some(rt) = self.rt.get_mut(&i.name) {
1801                    rt.pid = 0;
1802                }
1803                self.set_rotation(&i.name, false);
1804                return Err(Error::OperationFailed {
1805                    step: format!("set up {}", i.name),
1806                    message: e.to_string(),
1807                });
1808            }
1809            // Seen restarting: its app now runs what setup delivered.
1810            if restarted {
1811                self.mark_started(def, &i.name)?;
1812            }
1813        }
1814        let alive = self.alive(&sb, spec, oci);
1815        let healthy = match probe {
1816            None => alive,
1817            Some(p) => {
1818                let rt = self.rt.get_mut(&i.name).unwrap();
1819                let since = rt.since.unwrap_or_else(Instant::now);
1820                let in_start = since.elapsed() < p.start_period;
1821                if rt.next_probe.is_none_or(|t| Instant::now() >= t) {
1822                    let r = supervise::probe(&sb, p);
1823                    let rt = self.rt.get_mut(&i.name).unwrap();
1824                    rt.last_probe = r.output.clone();
1825                    if r.ok {
1826                        rt.failures = 0;
1827                        rt.healthy = Some(true);
1828                        rt.unhealthy_restarts = 0;
1829                    } else if !in_start {
1830                        rt.failures += 1;
1831                        if rt.failures >= p.retries {
1832                            rt.healthy = Some(false);
1833                        }
1834                    }
1835                    let wait = if rt.healthy.is_none() {
1836                        p.start_interval
1837                    } else {
1838                        p.interval
1839                    };
1840                    rt.next_probe = Some(Instant::now() + wait);
1841                }
1842                self.rt[&i.name].healthy == Some(true) && alive
1843            }
1844        };
1845        self.set_rotation(&i.name, healthy);
1846        let unhealthy = self.rt[&i.name].healthy == Some(false);
1847        if unhealthy && condition != RestartCondition::None {
1848            if self.restart_budget_spent(&i.name, policy.as_ref()) {
1849                self.message = Some(format!("{}: unhealthy, restart limit reached", i.name));
1850                return Ok(());
1851            }
1852            let n = self.rt[&i.name].unhealthy_restarts;
1853            if n >= RESTARTS_BEFORE_REPLACE {
1854                self.log(&format!(
1855                    "{} stayed unhealthy through {n} restarts; replacing it",
1856                    i.name
1857                ));
1858                let rt = self.rt.get_mut(&i.name).unwrap();
1859                rt.unhealthy_restarts = 0;
1860                drop(sb);
1861                self.retire(&i.name)?;
1862                return Ok(());
1863            }
1864            self.event(
1865                "warn",
1866                Some(&i.name),
1867                &format!(
1868                    "{} is unhealthy ({}); restarting its app",
1869                    i.name,
1870                    self.rt[&i.name]
1871                        .last_probe
1872                        .lines()
1873                        .last()
1874                        .unwrap_or("probe failed")
1875                ),
1876            );
1877            self.count_restart(&i.name);
1878            let rt = self.rt.get_mut(&i.name).unwrap();
1879            rt.unhealthy_restarts += 1;
1880            rt.failures = 0;
1881            rt.healthy = None;
1882            rt.since = Some(Instant::now());
1883            supervise::restart_app(&sb, &self.service, oci)?;
1884            if !oci {
1885                // The unit restarts within the same instance; it reads the
1886                // files and variables delivered last.
1887                self.mark_started(def, &i.name)?;
1888            }
1889        }
1890        Ok(())
1891    }
1892
1893    /// Secrets and the app's unit, after a boot or on creation.
1894    fn setup(&self, def: &StackDef, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> Result<()> {
1895        // Read from the store now, never from the definition.
1896        let keys = spec.secret_keys();
1897        let values = if keys.is_empty() {
1898            BTreeMap::new()
1899        } else {
1900            super::secrets::values(&self.inner.secrets, &def.org, &def.secrets, keys)?
1901        };
1902        // A new value for a secret the service takes with `on_change: none`
1903        // is delivered without restarting the app.
1904        let none = |k: &str| def.on_change(&self.service, k) == OnChange::None;
1905        // An OCI app started before its files arrived: restart it once so
1906        // it reads them (the next pass finds them in place).
1907        let pushed = supervise::push_secrets_detailed(sb, spec, &values)?;
1908        if oci && (pushed.missing || pushed.changed.iter().any(|k| !none(k))) {
1909            supervise::restart_app(sb, &self.service, oci)?;
1910        }
1911        if spec.command.is_some() && !oci {
1912            let mut s = spec.clone();
1913            s.restart = Some(RestartMode::Always);
1914            let env = supervise::secret_env(spec, &values)?;
1915            // Only a change all of whose variables are `none` secrets is
1916            // left for the next start.
1917            let env_restarts =
1918                spec.env.secrets.is_empty() || spec.env.secrets.values().any(|k| !none(k));
1919            supervise::install_with(
1920                sb,
1921                &self.service,
1922                &s,
1923                !spec.secrets.is_empty(),
1924                &env,
1925                env_restarts,
1926            )?;
1927        }
1928        Ok(())
1929    }
1930
1931    /// An OCI instance's secret variables, for its config (`environment.KEY`).
1932    fn oci_secret_env(
1933        &self,
1934        def: &StackDef,
1935        spec: &SandboxSpec,
1936    ) -> Result<BTreeMap<String, String>> {
1937        if spec.env.secrets.is_empty() {
1938            return Ok(BTreeMap::new());
1939        }
1940        let values = super::secrets::values(
1941            &self.inner.secrets,
1942            &def.org,
1943            &def.secrets,
1944            spec.env.secrets.values().map(String::as_str),
1945        )?;
1946        supervise::secret_env(spec, &values)
1947    }
1948
1949    /// The app's process is up: its unit is active, or (OCI, or no command)
1950    /// the instance runs.
1951    fn alive(&self, sb: &Sandbox, spec: &SandboxSpec, oci: bool) -> bool {
1952        if spec.command.is_some() && !oci {
1953            return supervise::unit_state(sb, &self.service).is_ok_and(|s| s == "active");
1954        }
1955        sb.info()
1956            .is_ok_and(|i| i.status.eq_ignore_ascii_case("running"))
1957    }
1958
1959    fn count_restart(&mut self, name: &str) {
1960        self.rt
1961            .entry(name.to_string())
1962            .or_default()
1963            .restarts
1964            .push_back(Instant::now());
1965    }
1966
1967    fn restart_budget_spent(
1968        &mut self,
1969        name: &str,
1970        policy: Option<&crate::spec::RestartPolicy>,
1971    ) -> bool {
1972        let Some(p) = policy else { return false };
1973        let Some(max) = p.max_attempts else {
1974            return false;
1975        };
1976        let window = p
1977            .window
1978            .as_deref()
1979            .and_then(|w| crate::flex::parse_duration(w).ok());
1980        let rt = self.rt.entry(name.to_string()).or_default();
1981        if let Some(w) = window {
1982            while rt.restarts.front().is_some_and(|t| t.elapsed() > w) {
1983                rt.restarts.pop_front();
1984            }
1985        }
1986        rt.restarts.len() as u32 >= max
1987    }
1988
1989    /// Replace (or fill) one slot with an instance of the current revision.
1990    #[expect(clippy::too_many_arguments)]
1991    fn replace(
1992        &mut self,
1993        def: &StackDef,
1994        spec: &SandboxSpec,
1995        rev: &str,
1996        slot: u32,
1997        old: Option<&str>,
1998        order: UpdateOrder,
1999        oci: bool,
2000        probe: Option<&HealthProbe>,
2001        monitor: Duration,
2002    ) -> Result<()> {
2003        // Before anything is stopped: a secret that cannot be read leaves
2004        // the old instance serving.
2005        let secret_env = if oci {
2006            self.oci_secret_env(def, spec)?
2007        } else {
2008            BTreeMap::new()
2009        };
2010        if let (Some(o), UpdateOrder::StopFirst) = (old, order) {
2011            self.log(&format!("slot {slot}: replacing {o} (stop-first)"));
2012            self.slot_state(def, slot, Some("draining"), None, None);
2013            self.retire(o)?;
2014            self.slot_state(def, slot, Some("retired"), None, None);
2015        }
2016        let name = instance_name(&self.stack, &self.service, slot, &new_id())?;
2017        self.log(&format!("slot {slot}: creating {name} (rev {rev})"));
2018        self.slot_state(def, slot, None, Some(&name), Some("creating"));
2019        let mut s = instance_spec(def, &self.service, spec, slot, rev)?;
2020        s.name = Some(name.clone());
2021        // `env.secrets` stays set, so the values are redacted in reports.
2022        s.env.vars.extend(secret_env);
2023        let d = crate::sandbox::resolve(self.client(), &s, &def.file.volumes, &def.base_dir)?;
2024        let stack = self.q.clone();
2025        let mut report = |m: &str| eprintln!("isb serve: {stack}: {m}");
2026        let created =
2027            crate::sandbox::ensure(self.client(), &d, EnsureOptions::default(), &mut report);
2028        let result = created.and_then(|_| {
2029            let inst = Inst {
2030                name: name.clone(),
2031                slot,
2032                rev: rev.to_string(),
2033                status: "Running".into(),
2034                secrets: Some(live_versions(def, &self.service)),
2035            };
2036            self.insts.push(inst.clone());
2037            self.slot_state(def, slot, None, None, Some("probing"));
2038            self.wait_serving(def, &inst, spec, oci, probe, monitor)
2039        });
2040        if let Err(e) = result {
2041            // Read its output before it is deleted: it explains the failure.
2042            let (e, attempt) =
2043                super::failure::explain(self.client(), &name, &self.service, oci, e, now_ms());
2044            self.inner.failures.record(&self.q, &self.service, attempt);
2045            self.event(
2046                "error",
2047                Some(&name),
2048                &format!("{name} did not come up: {e}"),
2049            );
2050            self.slot_state(def, slot, None, None, Some("failed"));
2051            let _ = self.retire(&name);
2052            return Err(e);
2053        }
2054        if let (Some(o), UpdateOrder::StartFirst) = (old, order) {
2055            self.log(&format!("slot {slot}: {name} is serving; retiring {o}"));
2056            self.slot_state(def, slot, Some("draining"), None, None);
2057            self.retire(o)?;
2058            self.slot_state(def, slot, Some("retired"), None, None);
2059        }
2060        self.slot_state(def, slot, None, None, Some("serving"));
2061        Ok(())
2062    }
2063
2064    /// Wait for a new instance to pass its health (or, without a
2065    /// healthcheck, for its app to run), put it in rotation, then watch it
2066    /// for `monitor`.
2067    fn wait_serving(
2068        &mut self,
2069        def: &StackDef,
2070        i: &Inst,
2071        spec: &SandboxSpec,
2072        oci: bool,
2073        probe: Option<&HealthProbe>,
2074        monitor: Duration,
2075    ) -> Result<()> {
2076        let deadline = match probe {
2077            Some(p) => p.start_period + p.interval * p.retries + Duration::from_secs(30),
2078            None => Duration::from_secs(60),
2079        }
2080        .max(Duration::from_secs(60));
2081        let started = Instant::now();
2082        loop {
2083            self.maintain(def, i, spec, oci, probe)?;
2084            if self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2085                break;
2086            }
2087            if self
2088                .rt
2089                .get(&i.name)
2090                .is_some_and(|r| r.healthy == Some(false))
2091            {
2092                return Err(Error::invalid(format!(
2093                    "unhealthy: {}",
2094                    self.rt[&i.name].last_probe
2095                )));
2096            }
2097            if started.elapsed() > deadline {
2098                let why = self
2099                    .rt
2100                    .get(&i.name)
2101                    .map(|r| r.last_probe.clone())
2102                    .filter(|s| !s.is_empty())
2103                    .unwrap_or_else(|| "its app is not running".into());
2104                return Err(Error::invalid(format!(
2105                    "not serving after {:?}: {why}",
2106                    started.elapsed()
2107                )));
2108            }
2109            // Probe on the start interval rather than waiting for the next one.
2110            if let Some(rt) = self.rt.get_mut(&i.name) {
2111                rt.next_probe = None;
2112            }
2113            std::thread::sleep(
2114                probe
2115                    .map(|p| p.start_interval)
2116                    .unwrap_or(Duration::from_secs(1))
2117                    .min(Duration::from_secs(2)),
2118            );
2119        }
2120        self.sync_routes();
2121        self.slot_state(def, i.slot, None, None, Some("monitoring"));
2122        let watch = Instant::now();
2123        while watch.elapsed() < monitor {
2124            std::thread::sleep(Duration::from_secs(1).min(monitor));
2125            if let Some(rt) = self.rt.get_mut(&i.name) {
2126                rt.next_probe = None;
2127            }
2128            self.maintain(def, i, spec, oci, probe)?;
2129            if !self.rt.get(&i.name).is_some_and(|r| r.in_rotation) {
2130                return Err(Error::invalid(format!(
2131                    "failed within the {monitor:?} monitor period"
2132                )));
2133            }
2134        }
2135        Ok(())
2136    }
2137
2138    /// Take an instance out of rotation, let its connections drain, then
2139    /// delete it.
2140    fn retire(&mut self, name: &str) -> Result<()> {
2141        self.drain(name);
2142        if let Ok(sb) = Sandbox::get(self.client(), name) {
2143            let _ = sb.stop(false, Duration::from_secs(10));
2144        }
2145        self.rt.remove(name);
2146        self.insts.retain(|i| i.name != name);
2147        match Sandbox::remove(self.client(), name, true) {
2148            Err(e) if !e.is_not_found() => Err(e),
2149            _ => Ok(()),
2150        }
2151    }
2152
2153    /// Take an instance out of rotation and let its connections drain (up
2154    /// to [`DRAIN`]).
2155    fn drain(&mut self, name: &str) {
2156        let ip = self.rt.get(name).and_then(|r| r.ip);
2157        self.set_rotation(name, false);
2158        self.sync_routes();
2159        if let Some(ip) = ip {
2160            let started = Instant::now();
2161            if let Some(o) = &self.inner.observer {
2162                o.drain(&self.q, &self.service, ip, DRAIN);
2163            }
2164            let left = DRAIN.saturating_sub(started.elapsed());
2165            for (k, p) in &self.routes {
2166                self.inner
2167                    .balancer
2168                    .wait_drained(k, SocketAddr::new(ip, p.target), left);
2169            }
2170        }
2171    }
2172
2173    fn set_rotation(&mut self, name: &str, on: bool) {
2174        let rt = self.rt.entry(name.to_string()).or_default();
2175        if rt.in_rotation != on {
2176            rt.in_rotation = on;
2177            self.sync_routes();
2178        }
2179    }
2180
2181    /// Bring the balancer's routes in line with the spec's published ports.
2182    fn set_routes(&mut self, spec: &SandboxSpec) {
2183        let want: BTreeMap<String, Published> = match published(spec) {
2184            Ok(ps) => ps
2185                .into_iter()
2186                .map(|p| (format!("{}/{}/{}", self.q, self.service, p.display()), p))
2187                .collect(),
2188            Err(e) => {
2189                self.message = Some(e.to_string());
2190                BTreeMap::new()
2191            }
2192        };
2193        let stale: Vec<String> = self
2194            .routes
2195            .keys()
2196            .filter(|k| !want.contains_key(*k))
2197            .cloned()
2198            .collect();
2199        for k in stale {
2200            self.inner.balancer.remove_route(&k);
2201            self.routes.remove(&k);
2202            self.route_errors.remove(&k);
2203        }
2204        for (k, p) in want {
2205            self.routes.entry(k).or_insert(p);
2206        }
2207        self.sync_routes();
2208    }
2209
2210    fn sync_routes(&mut self) {
2211        let mut errors = BTreeMap::new();
2212        // UDP ports are proxy devices on the replica (see super::ports).
2213        for (k, p) in self.routes.iter().filter(|(_, p)| !p.udp) {
2214            let backends: Vec<SocketAddr> = self
2215                .rt
2216                .values()
2217                .filter(|r| r.in_rotation)
2218                .filter_map(|r| r.ip)
2219                .map(|ip| SocketAddr::new(ip, p.target))
2220                .collect();
2221            if let Err(e) = self.inner.balancer.set_route(k, p.listen, backends) {
2222                errors.insert(k.clone(), e.to_string());
2223            }
2224        }
2225        for (k, e) in &errors {
2226            if self.route_errors.get(k) != Some(e) {
2227                self.log(&format!("cannot publish {k}: {e}"));
2228            }
2229        }
2230        self.route_errors = errors;
2231        self.sync_observer();
2232        self.sync_dns();
2233    }
2234
2235    /// Tell the observer (the ingress) when the in-rotation set changed.
2236    fn sync_observer(&mut self) {
2237        let Some(o) = self.inner.observer.clone() else {
2238            return;
2239        };
2240        let mut ips: Vec<IpAddr> = self
2241            .rt
2242            .values()
2243            .filter(|r| r.in_rotation)
2244            .filter_map(|r| r.ip)
2245            .collect();
2246        ips.sort();
2247        if self.observed.as_ref() != Some(&ips) {
2248            o.rotation(&self.q, &self.service, &ips);
2249            self.observed = Some(ips);
2250        }
2251    }
2252
2253    #[expect(
2254        clippy::too_many_lines,
2255        reason = "predates the lint ratchet; split it when next changed"
2256    )]
2257    fn publish_status(&mut self, def: &StackDef) {
2258        // An address can change without a rotation change (a restart).
2259        self.sync_observer();
2260        self.sync_dns();
2261        let Ok(spec) = def.service(&self.service) else {
2262            return;
2263        };
2264        let rev = def.revision(&self.service).unwrap_or_default();
2265        let probe = matches!(spec.health_probe(), Ok(Some(_)));
2266        let live = live_versions(def, &self.service);
2267        let snap = self.inner.snapshot.lock().unwrap().instances.clone();
2268        let instances: Vec<InstanceStatus> = self
2269            .insts
2270            .iter()
2271            .map(|i| {
2272                let rt = self.rt.get(&i.name);
2273                let m = snap.get(&format!("{}/{}", self.oclient.project_name(), i.name));
2274                let health = match (probe, rt.and_then(|r| r.healthy)) {
2275                    (false, _) => "none",
2276                    (true, Some(true)) => "healthy",
2277                    (true, Some(false)) => "unhealthy",
2278                    (true, None) => "starting",
2279                };
2280                InstanceStatus {
2281                    name: i.name.clone(),
2282                    slot: i.slot,
2283                    rev: i.rev.clone(),
2284                    status: m
2285                        .map(|m| m.status.clone())
2286                        .unwrap_or_else(|| i.status.clone()),
2287                    health: health.into(),
2288                    ip: rt.and_then(|r| r.ip).map(|ip| ip.to_string()),
2289                    in_rotation: rt.is_some_and(|r| r.in_rotation),
2290                    restarts: rt.map(|r| r.restarts.len() as u32).unwrap_or(0),
2291                    last_probe: rt.map(|r| r.last_probe.clone()).unwrap_or_default(),
2292                    cpu_pct: m.and_then(|m| m.cpu_pct),
2293                    cpu_history: m.map(|m| m.cpu_history.clone()).unwrap_or_default(),
2294                    mem_bytes: m.and_then(|m| m.mem_bytes),
2295                    disk_bytes: m.and_then(|m| m.disk_bytes),
2296                    stale_secrets: i
2297                        .secrets
2298                        .as_ref()
2299                        .map(|have| super::secrets::stale(have, &live))
2300                        .unwrap_or_default(),
2301                }
2302            })
2303            .collect();
2304        let mut instances = instances;
2305        instances.sort_by(|a, b| (a.slot, &a.name).cmp(&(b.slot, &b.name)));
2306        let routes = self.inner.balancer.routes();
2307        let now = Instant::now();
2308        let mut ports = Vec::new();
2309        for (k, p) in &self.routes {
2310            let r = routes.iter().find(|r| r.key == *k);
2311            let accepted = r.map(|r| r.accepted).unwrap_or(0);
2312            let e = self
2313                .rates
2314                .entry(k.clone())
2315                .or_insert_with(|| (accepted, now, VecDeque::new()));
2316            let dt = now.duration_since(e.1).as_secs_f32();
2317            // Only sample on a reconcile-sized step, so a burst of status
2318            // updates during a rollout does not flatten the curve.
2319            if dt >= 1.0 {
2320                let rate = accepted.saturating_sub(e.0) as f32 / dt;
2321                if e.2.len() == crate::metrics::HISTORY {
2322                    e.2.pop_front();
2323                }
2324                e.2.push_back(rate);
2325                e.0 = accepted;
2326                e.1 = now;
2327            }
2328            ports.push(PortStatus {
2329                listen: p.display(),
2330                target: p.target,
2331                backends: match r {
2332                    _ if p.udp => {
2333                        let ips = self.rt.values().filter_map(|r| r.ip);
2334                        ips.map(|ip| SocketAddr::new(ip, p.target).to_string())
2335                            .collect()
2336                    }
2337                    Some(r) => r.backends.iter().map(|b| b.addr.to_string()).collect(),
2338                    None => Vec::new(),
2339                },
2340                error: self.route_errors.get(k).cloned(),
2341                accepted,
2342                active: r
2343                    .map(|r| {
2344                        r.backends
2345                            .iter()
2346                            .chain(r.draining.iter())
2347                            .map(|b| b.active)
2348                            .sum()
2349                    })
2350                    .unwrap_or(0),
2351                rate_history: e.2.iter().copied().collect(),
2352            });
2353        }
2354        self.rates.retain(|k, _| self.routes.contains_key(k));
2355        let running = instances
2356            .iter()
2357            .filter(|i| i.status.eq_ignore_ascii_case("running"))
2358            .count() as u32;
2359        let healthy = instances.iter().filter(|i| i.in_rotation).count() as u32;
2360        self.watch_health(healthy, spec.replicas());
2361        let st = ServiceStatus {
2362            service: self.service.clone(),
2363            image: spec.image.clone(),
2364            rev,
2365            replicas: spec.replicas(),
2366            running,
2367            healthy,
2368            state: self.state.clone(),
2369            message: self.message.clone(),
2370            instances,
2371            ports,
2372            rollout: self.rollout.clone(),
2373            checked_at: now_secs(),
2374            domains: Vec::new(),
2375        };
2376        self.inner.status.lock().unwrap().insert(self.key(), st);
2377    }
2378}
2379
2380/// Names of the services that are not converged, for a deploy waiting on
2381/// its rollout.
2382pub fn unsettled(st: &StackStatus) -> BTreeSet<String> {
2383    st.services
2384        .iter()
2385        .filter(|s| s.state != "converged")
2386        .map(|s| s.service.clone())
2387        .collect()
2388}
2389
2390#[path = "rotation.rs"]
2391mod rotation;
2392
2393#[cfg(test)]
2394#[path = "controller_tests.rs"]
2395mod tests;