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