Skip to main content

moirai_executor/registry/
retention.rs

1//! Completed-task retention: which settled blocks the registry releases, and
2//! the bounded sweep that releases them.
3
4use std::{sync::Arc, time::Duration, time::Instant};
5
6use moirai_core::executor::CleanupConfig;
7
8use super::directory::SweepWindow;
9use super::registry::TaskRegistry;
10use super::state::{Retirement, TASK_STATE_BLOCK_SIZE, TaskStateBlock};
11
12/// Queued blocks one sweep step examines.
13///
14/// A step runs once per block the registry creates, i.e. once per
15/// [`TASK_STATE_BLOCK_SIZE`] registrations, so the examination cost per spawn is
16/// this budget over that block size. The budget exceeds one so the sweep laps
17/// the resident blocks faster than allocation adds to them.
18pub(super) const SWEEP_WINDOW: usize = 8;
19
20/// How long, and how many, completed tasks the registry keeps observable.
21///
22/// Retention is block-granular: a block of 1,024 tasks is released as one unit
23/// once **every** task in it has completed and released its lifecycle token,
24/// and either the block's newest completion is older than `max_age` or more
25/// than `max_completed_tasks` worth of blocks are resident. The cap holds
26/// regardless of age, and it counts resident blocks, so a block pinned by a
27/// running task and the block being filled count toward it.
28///
29/// A single long-running task keeps its own block resident (about 73 KiB)
30/// however many later tasks complete, and no other. The sweep visits resident
31/// blocks in a queue, eight per created block, so a settled block waits at
32/// most one lap of that queue. With `k` the larger of the cap in blocks and the
33/// pinned blocks plus the one being filled, the resident count stays at most
34/// `M + ceil(M / 8)` for the least `M` with `M >= k + ceil(M / 8)`, about 9/7
35/// of `k`, however many blocks are created.
36///
37/// After release, [`TaskRegistry::is_completed`] stays `true` and a cancel
38/// request still reports "already completed"; only per-task metadata and
39/// statistics are gone.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub struct RetentionPolicy {
42    /// Completed-task metadata older than this may be released.
43    pub max_age: Duration,
44    /// Most completed tasks kept regardless of age, rounded up to whole blocks.
45    pub max_completed_tasks: usize,
46}
47
48impl RetentionPolicy {
49    /// The policy an executor derives from its cleanup configuration, or
50    /// `None` when automatic cleanup is disabled.
51    #[must_use]
52    pub fn from_cleanup(config: &CleanupConfig) -> Option<Self> {
53        config.enable_automatic_cleanup.then_some(Self {
54            max_age: config.task_retention_duration,
55            max_completed_tasks: config.max_retained_tasks,
56        })
57    }
58
59    /// Resident blocks above which settled blocks retire regardless of age: the
60    /// retained completed tasks in whole blocks, plus the block being filled.
61    pub(super) fn max_resident_blocks(&self) -> usize {
62        self.max_completed_tasks
63            .div_ceil(TASK_STATE_BLOCK_SIZE)
64            .saturating_add(1)
65    }
66}
67
68/// What a sweep decided about one examined block, before the directory lock is
69/// retaken.
70#[derive(Clone, Copy)]
71enum Verdict {
72    /// The block is pinned, partly unregistered, or younger than the window
73    /// while the resident cap is not exceeded.
74    Keep,
75    /// Every task completed before the age cutoff.
76    Expired,
77    /// The block is settled but young: it retires only while more blocks than
78    /// the cap are resident when the sweep commits.
79    OverCap,
80}
81
82impl Verdict {
83    /// `forceable` is whether the resident cap can retire a settled block.
84    fn of(
85        block: &TaskStateBlock,
86        first_slot_unissued: bool,
87        cutoff: Option<Instant>,
88        forceable: bool,
89    ) -> Self {
90        let expired = || {
91            cutoff.is_some_and(|cutoff| {
92                block.is_settled(first_slot_unissued, Retirement::CompletedBefore(cutoff))
93            })
94        };
95        if forceable && !block.is_settled(first_slot_unissued, Retirement::Forced) {
96            Self::Keep
97        } else if expired() {
98            Self::Expired
99        } else if forceable {
100            Self::OverCap
101        } else {
102            Self::Keep
103        }
104    }
105}
106
107impl TaskRegistry {
108    /// Advance the sweep by one bounded window.
109    ///
110    /// Called from the registration slow path, once per created block, so the
111    /// reclamation work is proportional to allocation: memory can grow only as
112    /// fast as tasks register, and each new block pays for examining a few
113    /// queued ones. That is at most [`SWEEP_WINDOW`] blocks, each up to two
114    /// settledness scans of at most [`TASK_STATE_BLOCK_SIZE`] slots (the cap
115    /// test, then the age test), plus two brief directory write locks, per
116    /// [`TASK_STATE_BLOCK_SIZE`] registrations; the other registrations do none
117    /// of it. A registry without a retention policy never sweeps.
118    pub(super) fn sweep_step(&self) {
119        let Some(policy) = self.retention else {
120            return;
121        };
122        let cutoff = Instant::now().checked_sub(policy.max_age);
123        self.sweep_window(cutoff, policy.max_resident_blocks());
124    }
125
126    /// Release every block whose tasks all completed at least `older_than` ago,
127    /// returning how many blocks were released.
128    ///
129    /// Every resident block is examined, including any that a concurrent
130    /// automatic sweep has checked out; that sweep finds such a block already
131    /// retired and moves on. The directory lock is held only to list the
132    /// resident blocks and to retire the settled ones, never across a scan, so
133    /// registrations proceed throughout. Automatic retention makes calling this
134    /// unnecessary for an executor; it serves callers that manage a registry
135    /// directly.
136    pub fn cleanup_completed(&self, older_than: Duration) -> usize {
137        // `Instant - Duration` panics when the result predates the platform's
138        // clock origin, which a caller-supplied retention window longer than the
139        // process uptime reaches. No recorded completion can be older than a
140        // cutoff before the clock started, so that case is an empty sweep.
141        let Some(cutoff) = Instant::now().checked_sub(older_than) else {
142            return 0;
143        };
144        let resident: Vec<_> = self
145            .blocks
146            .read()
147            .expect("task registry block directory is never poisoned")
148            .resident_indexed()
149            .map(|(index, block)| (index, Arc::clone(block)))
150            .collect();
151        let settled: Vec<usize> = resident
152            .iter()
153            .filter(|(index, block)| {
154                block.is_settled(*index == 0, Retirement::CompletedBefore(cutoff))
155            })
156            .map(|&(index, _)| index)
157            .collect();
158        // The retired blocks drop with `resident`, after the directory lock.
159        let released = self
160            .blocks
161            .write()
162            .expect("task registry block directory is never poisoned")
163            .cleanup(&settled);
164        released.len()
165    }
166
167    /// Check out the next window of queued blocks, retire the settled ones, and
168    /// requeue the rest.
169    ///
170    /// `cutoff` bounds completion age (`None`: no block is old enough), and
171    /// settled blocks retire regardless of age while more than `resident_cap`
172    /// blocks are resident. Whether the cap is exceeded is decided under the
173    /// directory lock against the live count, so concurrent sweeps cannot each
174    /// retire down from a stale copy of it; only the settledness scans, which
175    /// cost up to a block's slots each, run outside the lock.
176    fn sweep_window(&self, cutoff: Option<Instant>, resident_cap: usize) {
177        let SweepWindow { blocks, resident } = self
178            .blocks
179            .write()
180            .expect("task registry block directory is never poisoned")
181            .check_out::<SWEEP_WINDOW>();
182        let forceable = resident > resident_cap;
183        let judged = blocks.map(|slot| {
184            slot.map(|(index, block)| {
185                let verdict = Verdict::of(&block, index == 0, cutoff, forceable);
186                (index, block, verdict)
187            })
188        });
189
190        // Released blocks drop with `judged`, after the directory lock.
191        let mut released = [const { None }; SWEEP_WINDOW];
192        let mut directory = self
193            .blocks
194            .write()
195            .expect("task registry block directory is never poisoned");
196        for (slot, (index, _, verdict)) in released.iter_mut().zip(judged.iter().flatten()) {
197            let retire = match verdict {
198                Verdict::Expired => true,
199                Verdict::OverCap => directory.resident() > resident_cap,
200                Verdict::Keep => false,
201            };
202            if retire {
203                *slot = directory.retire(*index);
204            } else {
205                directory.requeue(*index);
206            }
207        }
208        drop(directory);
209    }
210}