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