Skip to main content

isb_core/stack/
controller.rs

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