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}