pitchfork-cli 2.27.0

Daemons with DX
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
//! Idle shutdown of daemons the proxy started.
//!
//! A daemon is eligible only while its record carries a
//! `proxy_idle_timeout_ms`, which only a proxy auto-start sets. Anything
//! started another way — and a proxy-started daemon that has since been
//! started explicitly, which [`Supervisor::claim_daemons`] records — is never
//! stopped here.
//!
//! The interval watcher calls [`Supervisor::check_idle_daemons`]. It stops an
//! eligible daemon once all of these hold:
//!
//! - the proxy has carried nothing for it (see [`crate::proxy::activity`]) for
//!   its grace period;
//! - nothing running or starting depends on it, other than daemons being
//!   stopped in the same pass;
//! - no tracked shell or project session is inside its directory.
//!
//! Daemons are stopped dependents first, so a dependency goes only after the
//! last daemon that needs it, and never while anything else still does.

use super::autostop::is_within;
use super::{SUPERVISOR, Supervisor};
use crate::daemon::Daemon;
use crate::daemon_id::DaemonId;
use crate::daemon_status::DaemonStatus;
use crate::ipc::IpcResponse;
use crate::proxy::activity::ACTIVITY;
use log::LevelFilter::Info;
use std::collections::{BTreeMap, HashMap, HashSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;

/// Held by an idle stop while it makes its final check and marks the daemon
/// stopping, and briefly by a shell or project session entering a directory.
///
/// A shell entering a daemon's directory keeps the daemon running, so the
/// two must not interleave: the shell is either seen by the final check, or
/// registered once the daemon is already marked stopping — the same as
/// entering a moment later, which the shell hook handles by starting what it
/// manages, after the stop. The lock is released before the stop itself, so
/// entering a directory never waits for a daemon to exit.
static SHELL_ADMISSION: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());

/// Hold off idle stops while a shell or project session enters a directory.
pub(crate) async fn admit_shell() -> tokio::sync::MutexGuard<'static, ()> {
    SHELL_ADMISSION.lock().await
}

/// Set while a sweep's stops are under way, so a slow stop does not let the
/// next tick start a second, overlapping sweep.
static SWEEPING: AtomicBool = AtomicBool::new(false);

/// Holds [`SWEEPING`] for one sweep and clears it when dropped, so a sweep
/// that ends early — or panics — does not turn idle shutdown off for good.
struct Sweep;

impl Sweep {
    fn begin() -> Option<Self> {
        (!SWEEPING.swap(true, Ordering::SeqCst)).then_some(Sweep)
    }
}

impl Drop for Sweep {
    fn drop(&mut self) {
        SWEEPING.store(false, Ordering::SeqCst);
    }
}

/// A daemon's idle grace period, if it may be stopped for inactivity.
fn grace(daemon: &Daemon) -> Option<Duration> {
    daemon.proxy_idle_timeout_ms.map(Duration::from_millis)
}

/// Whether `daemon` still needs what it depends on: it is up, on its way up,
/// or on its way back. A restart (file watch, `--force`) stops first, and an
/// errored daemon with retries left is about to start again, so both keep
/// their dependencies as a running daemon does.
fn is_live(daemon: &Daemon) -> bool {
    let status = &daemon.status;
    status.is_running()
        || status.is_waiting()
        || status.is_stopping()
        || (status.is_errored() && daemon.retry_count < daemon.retry.count())
}

/// Whether a tracked shell or project session is inside `daemon`'s directory.
fn shell_inside(daemon: &Daemon, active_dirs: &[PathBuf]) -> bool {
    daemon
        .dir
        .as_deref()
        .is_some_and(|dir| active_dirs.iter().any(|d| is_within(dir, d)))
}

/// Live daemons that depend on `id`, other than those in `excluding`.
fn live_dependents<'a>(
    id: &'a DaemonId,
    daemons: &'a BTreeMap<DaemonId, Daemon>,
    excluding: &'a HashSet<DaemonId>,
) -> impl Iterator<Item = &'a DaemonId> + 'a {
    daemons
        .values()
        .filter(move |d| is_live(d) && d.depends.contains(id) && !excluding.contains(&d.id))
        .map(|d| &d.id)
}

/// Which daemons to stop for inactivity, and in what order.
///
/// Returns levels: every daemon in a level is stopped before any in the next,
/// and nothing in a later level depends on anything in an earlier one — so
/// dependents come first. A daemon is included only when it is eligible, idle
/// by `is_idle`, has no shell inside its directory, and every live daemon
/// depending on it is itself included.
pub(crate) fn plan_idle_stops(
    daemons: &BTreeMap<DaemonId, Daemon>,
    active_dirs: &[PathBuf],
    is_idle: impl Fn(&DaemonId, Duration) -> bool,
) -> Vec<Vec<DaemonId>> {
    let mut chosen: HashSet<DaemonId> = daemons
        .values()
        .filter(|d| d.status.is_running() && !shell_inside(d, active_dirs))
        .filter(|d| grace(d).is_some_and(|g| is_idle(&d.id, g)))
        .map(|d| d.id.clone())
        .collect();

    // Drop anything something outside the set still needs. Dropping one can
    // leave its own dependencies needed in turn, so repeat until nothing
    // changes.
    loop {
        let needed: Vec<DaemonId> = chosen
            .iter()
            .filter(|id| live_dependents(id, daemons, &chosen).next().is_some())
            .cloned()
            .collect();
        if needed.is_empty() {
            break;
        }
        for id in needed {
            chosen.remove(&id);
        }
    }

    // Peel off, level by level, the daemons nothing left in the set depends on.
    let mut levels = Vec::new();
    while !chosen.is_empty() {
        let mut level: Vec<DaemonId> = chosen
            .iter()
            .filter(|id| {
                !chosen
                    .iter()
                    .any(|other| daemons.get(other).is_some_and(|d| d.depends.contains(id)))
            })
            .cloned()
            .collect();
        if level.is_empty() {
            // What is left is a dependency cycle and what it depends on.
            // Nothing in a cycle can go first, and stopping its members
            // together cannot be undone if one of them then fails to stop,
            // which would leave the rest running without it. `depends` cycles
            // are refused at start, so one only appears through a later
            // config edit; leave it running.
            break;
        }
        level.sort();
        for id in &level {
            chosen.remove(id);
        }
        levels.push(level);
    }
    levels
}

impl Supervisor {
    /// Record that `ids` were started explicitly, so none of them is stopped
    /// for inactivity from now on.
    ///
    /// Sent by every start that is not the proxy's, for the daemons it names
    /// and everything they depend on, before anything is started: the result
    /// is the same as if the explicit start had come first. Waits out an idle
    /// stop already under way for any of them, so the caller then sees the
    /// daemon stopped and starts it, rather than skipping it as running while
    /// it goes away.
    pub(crate) async fn claim_daemons(&self, ids: &[DaemonId]) {
        let claimed: Vec<DaemonId> = {
            let mut state_file = self.state_file.lock().await;
            ids.iter()
                .filter(|id| state_file.clear_proxy_idle_timeout(id))
                .cloned()
                .collect()
        };
        for id in &claimed {
            info!("{id} was started explicitly; it will no longer be stopped when idle");
        }
        // An idle stop revalidates ownership under the daemon's stop lock, so
        // one that has not taken the lock yet will now call itself off; one
        // that holds it is waited for here.
        for id in ids {
            if ACTIVITY.is_idle_stopping(id) {
                drop(self.stop_lock(id).await.lock().await);
            }
        }
    }

    /// Stop proxy-started daemons that have been idle for their grace period.
    ///
    /// Cheap when nothing is eligible, which is the default. The stops run in
    /// a detached task so a slow one does not hold up the interval watcher.
    pub(crate) async fn check_idle_daemons(&self) {
        let daemons = {
            let state_file = self.state_file.lock().await;
            if !state_file
                .daemons
                .values()
                .any(|d| d.proxy_idle_timeout_ms.is_some() && d.status.is_running())
            {
                return;
            }
            state_file.daemons.clone()
        };
        let Some(sweep) = Sweep::begin() else {
            return;
        };
        let active_dirs = self.get_active_directories().await;
        let plan = plan_idle_stops(&daemons, &active_dirs, |id, grace| {
            let activity = ACTIVITY.snapshot(id);
            activity.in_flight == 0 && !activity.idle_stopping && activity.idle_for >= grace
        });
        if plan.is_empty() {
            return;
        }
        debug!("idle shutdown plan: {plan:?}");
        let graces: HashMap<DaemonId, Duration> = daemons
            .values()
            .filter_map(|d| grace(d).map(|g| (d.id.clone(), g)))
            .collect();
        tokio::spawn(async move {
            let _sweep = sweep;
            for level in plan {
                SUPERVISOR.idle_stop_level(level, &graces).await;
            }
        });
    }

    /// Stop one level of the plan for inactivity, as far as it still
    /// qualifies.
    ///
    /// Every member is claimed first. A claim keeps new proxy work from
    /// starting while the level is handled: a request arriving meanwhile
    /// waits and starts the daemon again once it has stopped.
    ///
    /// The planner leaves dependency cycles out, so members of a level do not
    /// depend on each other; each is still checked with the other claimed
    /// members set aside, and one that no longer qualifies leaves that set
    /// and has its claim released at once, so requests for it do not wait on
    /// the rest of the level and nothing goes while a member that stays still
    /// needs it.
    /// The last check for each member runs under its stop lock, since a
    /// request, an explicit start, a shell or a new dependent may have
    /// arrived since the plan was made.
    async fn idle_stop_level(&self, level: Vec<DaemonId>, graces: &HashMap<DaemonId, Duration>) {
        let mut group: HashSet<DaemonId> = HashSet::new();
        for id in level {
            match graces.get(&id) {
                Some(&grace) if ACTIVITY.claim_idle_stop(&id, grace) => {
                    group.insert(id);
                }
                _ => debug!("idle stop of {id} called off: it was active again"),
            }
        }

        loop {
            let mut blocked = Vec::new();
            for id in &group {
                if let Some(reason) = self.idle_stop_blocker(id, &group, false).await {
                    debug!("idle stop of {id} called off: {reason}");
                    blocked.push(id.clone());
                }
            }
            if blocked.is_empty() {
                break;
            }
            for id in blocked {
                // Staying up, so requests for it must not wait on this level.
                group.remove(&id);
                ACTIVITY.release_idle_stop(&id);
            }
        }

        let mut ordered: Vec<DaemonId> = group.iter().cloned().collect();
        ordered.sort();
        for id in ordered {
            let lock = self.stop_lock(&id).await;
            let stopped = {
                let _guard = lock.lock().await;
                // The decision and marking the daemon stopping are one step as
                // far as a shell entering its directory is concerned; the stop
                // itself, which can take the whole stop timeout, runs after.
                let blocker = {
                    let _admission = SHELL_ADMISSION.lock().await;
                    self.idle_stop_blocker(&id, &group, true).await
                };
                match blocker {
                    Some(reason) => {
                        debug!("idle stop of {id} called off: {reason}");
                        false
                    }
                    None => {
                        let grace = graces.get(&id).copied().unwrap_or_default();
                        info!("stopping {id}: no proxy activity for {grace:?}");
                        match self.stop_locked(&id).await {
                            Ok(IpcResponse::Ok | IpcResponse::DaemonWasNotRunning) => true,
                            // No process to stop, so nothing recorded an
                            // outcome over the stopping mark.
                            Ok(IpcResponse::DaemonNotRunning) => {
                                self.settle_stopping(&id, DaemonStatus::Stopped).await;
                                true
                            }
                            Ok(rsp) => {
                                error!("failed to stop idle daemon {id}: {rsp:?}");
                                self.settle_stopping(&id, DaemonStatus::Running).await;
                                false
                            }
                            Err(e) => {
                                error!("failed to stop idle daemon {id}: {e}");
                                self.settle_stopping(&id, DaemonStatus::Running).await;
                                false
                            }
                        }
                    }
                }
            };
            // Handled either way: requests for it no longer wait on this level.
            ACTIVITY.release_idle_stop(&id);
            if stopped {
                self.add_notification(Info, format!("stopped idle {id}"))
                    .await;
            } else {
                // Still running, so the members checked after it must keep
                // what it needs.
                group.remove(&id);
            }
        }
    }

    /// Replace the stopping mark an idle stop left on `id` with `status`, if
    /// nothing has recorded an outcome over it.
    async fn settle_stopping(&self, id: &DaemonId, status: DaemonStatus) {
        let mut state_file = self.state_file.lock().await;
        if state_file
            .daemons
            .get(id)
            .is_some_and(|d| d.status.is_stopping())
        {
            state_file.set_status(id, status);
        }
    }

    /// Why `id` must not be stopped for inactivity right now, if anything.
    ///
    /// Dependents in `stopping_with` are being stopped along with it and do
    /// not count. With `commit`, a daemon that may be stopped is marked
    /// stopping under the same lock, so anything that looks at it afterwards —
    /// a shell hook deciding what to start — sees it on its way down rather
    /// than running.
    async fn idle_stop_blocker(
        &self,
        id: &DaemonId,
        stopping_with: &HashSet<DaemonId>,
        commit: bool,
    ) -> Option<&'static str> {
        // Shells and sessions are read under the same lock as the rest, so a
        // shell entering the directory cannot slip in between the two.
        let mut state_file = self.state_file.lock().await;
        let active_dirs = state_file.active_directories();
        let Some(daemon) = state_file.daemons.get(id) else {
            return Some("it is no longer known");
        };
        if !daemon.status.is_running() {
            return Some("it is no longer running");
        }
        if daemon.proxy_idle_timeout_ms.is_none() {
            return Some("it was started explicitly");
        }
        if shell_inside(daemon, &active_dirs) {
            return Some("a shell is inside its directory");
        }
        if live_dependents(id, &state_file.daemons, stopping_with)
            .next()
            .is_some()
        {
            return Some("a running daemon depends on it");
        }
        if commit {
            state_file.set_status(id, DaemonStatus::Stopping);
        }
        None
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::daemon_status::DaemonStatus;

    fn id(name: &str) -> DaemonId {
        DaemonId::new("proj", name)
    }

    struct Fixture(BTreeMap<DaemonId, Daemon>);

    impl Fixture {
        fn new() -> Self {
            Self(BTreeMap::new())
        }

        fn add(mut self, name: &str, idle_ms: Option<u64>, depends: &[&str]) -> Self {
            self.0.insert(
                id(name),
                Daemon {
                    id: id(name),
                    status: DaemonStatus::Running,
                    dir: Some(PathBuf::from("/work/proj")),
                    depends: depends.iter().map(|d| id(d)).collect(),
                    proxy_idle_timeout_ms: idle_ms,
                    ..Default::default()
                },
            );
            self
        }

        fn status(mut self, name: &str, status: DaemonStatus) -> Self {
            self.0.get_mut(&id(name)).unwrap().status = status;
            self
        }

        fn plan(&self, active_dirs: &[&str], idle: &[&str]) -> Vec<Vec<String>> {
            let dirs: Vec<PathBuf> = active_dirs.iter().map(PathBuf::from).collect();
            let idle: HashSet<DaemonId> = idle.iter().map(|n| id(n)).collect();
            plan_idle_stops(&self.0, &dirs, |d, _| idle.contains(d))
                .into_iter()
                .map(|l| l.into_iter().map(|d| d.name().to_string()).collect())
                .collect()
        }
    }

    const G: Option<u64> = Some(60_000);

    #[test]
    fn only_proxy_started_daemons_are_eligible() {
        let f = Fixture::new().add("web", G, &[]).add("manual", None, &[]);
        assert_eq!(f.plan(&[], &["web", "manual"]), vec![vec!["web"]]);
    }

    #[test]
    fn active_daemons_are_kept() {
        let f = Fixture::new().add("web", G, &[]);
        assert!(f.plan(&[], &[]).is_empty());
    }

    #[test]
    fn dependencies_stop_after_their_dependents() {
        let f = Fixture::new()
            .add("db", G, &[])
            .add("cache", G, &[])
            .add("api", G, &["db", "cache"])
            .add("web", G, &["api"]);
        assert_eq!(
            f.plan(&[], &["db", "cache", "api", "web"]),
            vec![vec!["web"], vec!["api"], vec!["cache", "db"]]
        );
    }

    #[test]
    fn a_dependency_stays_while_a_busy_dependent_needs_it() {
        let f = Fixture::new().add("db", G, &[]).add("api", G, &["db"]);
        // db has no traffic of its own, but api is still busy.
        assert!(f.plan(&[], &["db"]).is_empty());
    }

    #[test]
    fn a_shared_dependency_stays_for_an_explicitly_started_consumer() {
        let f =
            Fixture::new()
                .add("db", G, &[])
                .add("api", G, &["db"])
                .add("worker", None, &["db"]);
        assert_eq!(f.plan(&[], &["db", "api"]), vec![vec!["api"]]);
    }

    #[test]
    fn a_starting_dependent_keeps_its_dependency() {
        let f = Fixture::new()
            .add("db", G, &[])
            .add("worker", None, &["db"])
            .status("worker", DaemonStatus::Waiting);
        assert!(f.plan(&[], &["db"]).is_empty());
    }

    #[test]
    fn a_dependent_on_its_way_back_keeps_its_dependency() {
        // Stopping for a restart.
        let f = Fixture::new()
            .add("db", G, &[])
            .add("api", None, &["db"])
            .status("api", DaemonStatus::Stopping);
        assert!(f.plan(&[], &["db"]).is_empty());

        // Crashed, with a retry still to come.
        let mut f = Fixture::new()
            .add("db", G, &[])
            .add("api", None, &["db"])
            .status("api", DaemonStatus::Errored(1));
        f.0.get_mut(&id("api")).unwrap().retry = crate::config_types::Retry(3);
        assert!(f.plan(&[], &["db"]).is_empty());

        // Out of retries: it is not coming back.
        f.0.get_mut(&id("api")).unwrap().retry_count = 3;
        assert_eq!(f.plan(&[], &["db"]), vec![vec!["db"]]);
    }

    #[test]
    fn a_stopped_dependent_does_not_keep_its_dependency() {
        let f = Fixture::new()
            .add("db", G, &[])
            .add("worker", None, &["db"])
            .status("worker", DaemonStatus::Stopped);
        assert_eq!(f.plan(&[], &["db"]), vec![vec!["db"]]);
    }

    #[test]
    fn a_proxied_daemon_depending_on_a_busy_one_keeps_it() {
        // `admin` depends on `api`; `api` is idle, but `admin` is serving
        // traffic, so `api` must stay even though it saw none itself.
        let f = Fixture::new().add("api", G, &[]).add("admin", G, &["api"]);
        assert!(f.plan(&[], &["api"]).is_empty());
    }

    #[test]
    fn a_shell_inside_the_directory_keeps_the_stack() {
        let f = Fixture::new().add("db", G, &[]).add("api", G, &["db"]);
        assert!(f.plan(&["/work/proj/src"], &["db", "api"]).is_empty());
        // A shell elsewhere does not.
        assert_eq!(
            f.plan(&["/work/other"], &["db", "api"]),
            vec![vec!["api"], vec!["db"]]
        );
    }

    #[test]
    fn a_dependency_cycle_is_left_running() {
        // `web` depends on the cycle and goes; the cycle and `db`, which it
        // needs, stay.
        let f = Fixture::new()
            .add("db", G, &[])
            .add("a", G, &["b", "db"])
            .add("b", G, &["a"])
            .add("web", G, &["a"]);
        assert_eq!(f.plan(&[], &["db", "a", "b", "web"]), vec![vec!["web"]]);
    }
}