/build/source/nativelink-scheduler/src/unsatisfiable_tracker.rs
Line | Count | Source |
1 | | // Copyright 2026 The NativeLink Authors. All rights reserved. |
2 | | // |
3 | | // Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License"); |
4 | | // you may not use this file except in compliance with the License. |
5 | | // You may obtain a copy of the License at |
6 | | // |
7 | | // See LICENSE file for details |
8 | | // |
9 | | // Unless required by applicable law or agreed to in writing, software |
10 | | // distributed under the License is distributed on an "AS IS" BASIS, |
11 | | // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
12 | | // See the License for the specific language governing permissions and |
13 | | // limitations under the License. |
14 | | |
15 | | //! Tracks how long each action property shape has gone without any worker |
16 | | //! able to run it. |
17 | | //! |
18 | | //! Time is kept per shape rather than per operation, so a build whose actions |
19 | | //! all queue together waits once rather than once per action. An action is |
20 | | //! only failed once its shape is due *and* the action has itself been queued |
21 | | //! for the timeout, so one that arrives while its shape is already due is not |
22 | | //! failed on sight. The tracker therefore also knows when the youngest such |
23 | | //! action becomes eligible, so the matching loop can wake for it, and it only |
24 | | //! backs off for a shape whose eligible actions could not be retired. |
25 | | |
26 | | use core::time::Duration; |
27 | | use std::collections::HashMap; |
28 | | use std::time::SystemTime; |
29 | | |
30 | | use crate::match_outcome::PropertyShape; |
31 | | |
32 | | /// How often the same shape is logged while it stays unsatisfiable. |
33 | | const WARN_INTERVAL: Duration = Duration::from_mins(1); |
34 | | |
35 | | /// Shortest wait before re-running a matching pass on behalf of a shape. |
36 | | const MIN_RECHECK: Duration = Duration::from_secs(1); |
37 | | |
38 | | /// Most doublings of `MIN_RECHECK` while a due shape keeps failing to be |
39 | | /// retired, so a stuck action wakes the loop at most every 32 seconds. |
40 | | const MAX_DUE_BACKOFF_DOUBLINGS: u32 = 5; |
41 | | |
42 | | #[derive(Debug)] |
43 | | struct ShapeState { |
44 | | first_unsatisfiable: SystemTime, |
45 | | last_seen: SystemTime, |
46 | | /// Fleet generation when the shape was last seen unsatisfiable. |
47 | | fleet_generation: u64, |
48 | | seen_this_pass: bool, |
49 | | seen_last_pass: bool, |
50 | | last_warned: Option<SystemTime>, |
51 | | /// Passes in which the shape was due and an eligible action of it was |
52 | | /// seen but could not be failed (the per-pass cap, or a version |
53 | | /// conflict). Passes where every action was merely too young do not |
54 | | /// count: nothing was stuck, it just was not time yet. |
55 | | due_passes: u32, |
56 | | /// Whether the shape was due in the pass being recorded. |
57 | | due_this_pass: bool, |
58 | | /// Actions of this shape seen this pass that had waited the timeout. |
59 | | eligible_this_pass: u32, |
60 | | /// Actions of this shape failed this pass (see `note_failed`). |
61 | | failed_this_pass: u32, |
62 | | /// When the youngest action of this shape seen this pass will have |
63 | | /// waited the timeout, so the loop can wake for it. Kept until the |
64 | | /// shape is next seen, because `next_deadline` runs after `end_pass`. |
65 | | earliest_eligible: Option<SystemTime>, |
66 | | } |
67 | | |
68 | | /// What the matching pass should do with one unsatisfiable action. |
69 | | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
70 | | pub struct Observation { |
71 | | /// How long the shape has been unsatisfiable. |
72 | | pub waited: Duration, |
73 | | /// The shape has been unsatisfiable for the configured timeout. |
74 | | pub is_due: bool, |
75 | | /// The action has itself been queued for the configured timeout. It is |
76 | | /// failed only when this and `is_due` both hold. |
77 | | pub action_eligible: bool, |
78 | | /// The shape has not been logged recently. |
79 | | pub should_warn: bool, |
80 | | } |
81 | | |
82 | | #[derive(Debug)] |
83 | | pub struct UnsatisfiableTracker { |
84 | | /// `None` never fails an action. |
85 | | timeout: Option<Duration>, |
86 | | shapes: HashMap<PropertyShape, ShapeState>, |
87 | | seen_this_pass: u64, |
88 | | } |
89 | | |
90 | | impl UnsatisfiableTracker { |
91 | | #[must_use] |
92 | 102 | pub fn new(timeout: Option<Duration>) -> Self { |
93 | 102 | Self { |
94 | 102 | timeout, |
95 | 102 | shapes: HashMap::new(), |
96 | 102 | seen_this_pass: 0, |
97 | 102 | } |
98 | 102 | } |
99 | | |
100 | | /// Records that an action with this shape, first submitted at |
101 | | /// `insert_timestamp`, was found unsatisfiable. |
102 | 153 | pub fn observe( |
103 | 153 | &mut self, |
104 | 153 | shape: PropertyShape, |
105 | 153 | now: SystemTime, |
106 | 153 | fleet_generation: u64, |
107 | 153 | insert_timestamp: SystemTime, |
108 | 153 | ) -> Observation { |
109 | 153 | self.seen_this_pass += 1; |
110 | 153 | let timeout = self.timeout; |
111 | 153 | let state = self.shapes.entry(shape).or_insert(ShapeState { |
112 | 153 | first_unsatisfiable: now, |
113 | 153 | last_seen: now, |
114 | 153 | fleet_generation, |
115 | 153 | seen_this_pass: false, |
116 | 153 | seen_last_pass: false, |
117 | 153 | last_warned: None, |
118 | 153 | due_passes: 0, |
119 | 153 | due_this_pass: false, |
120 | 153 | eligible_this_pass: 0, |
121 | 153 | failed_this_pass: 0, |
122 | 153 | earliest_eligible: None, |
123 | 153 | }); |
124 | 153 | let first_this_pass = !state.seen_this_pass; |
125 | 153 | if first_this_pass { |
126 | | // A shape kept from an earlier build only carries over the time |
127 | | // it was actually queued, not the idle gap since. |
128 | 121 | if !state.seen_last_pass { |
129 | 36 | let gap = now.duration_since(state.last_seen).unwrap_or_default(); |
130 | 36 | state.first_unsatisfiable += gap; |
131 | 85 | } |
132 | 121 | state.due_this_pass = false; |
133 | 121 | state.eligible_this_pass = 0; |
134 | 121 | state.failed_this_pass = 0; |
135 | 121 | state.earliest_eligible = None; |
136 | 32 | } |
137 | 153 | state.last_seen = now; |
138 | 153 | state.fleet_generation = fleet_generation; |
139 | 153 | state.seen_this_pass = true; |
140 | | |
141 | 153 | let should_warn = state.last_warned.is_none_or(|last_warned| {118 |
142 | 118 | now.duration_since(last_warned).unwrap_or_default() >= WARN_INTERVAL |
143 | 118 | }); |
144 | 153 | if should_warn { |
145 | 60 | state.last_warned = Some(now); |
146 | 93 | } |
147 | | |
148 | 153 | let waited = now |
149 | 153 | .duration_since(state.first_unsatisfiable) |
150 | 153 | .unwrap_or_default(); |
151 | 153 | let is_due = timeout.is_some_and(|timeout| waited139 >= timeout139 ); |
152 | 153 | state.due_this_pass |= is_due; |
153 | | |
154 | | // A clock that went backwards reads as not yet waited, the safe |
155 | | // direction; `now` never being reached reads the same way. |
156 | 153 | let action_eligible = timeout.is_some_and(|timeout| {139 |
157 | 139 | now.duration_since(insert_timestamp) |
158 | 139 | .is_ok_and(|action_waited| action_waited >= timeout) |
159 | 139 | }); |
160 | 153 | if action_eligible { |
161 | 49 | state.eligible_this_pass += 1; |
162 | 90 | } else if let Some(eligible_at) = |
163 | 104 | timeout.and_then(|timeout| insert_timestamp90 .checked_add90 (timeout90 )) |
164 | | { |
165 | | state.earliest_eligible = Some( |
166 | 90 | state |
167 | 90 | .earliest_eligible |
168 | 90 | .map_or(eligible_at, |earliest| earliest20 .min20 (eligible_at20 )), |
169 | | ); |
170 | 14 | } |
171 | 153 | Observation { |
172 | 153 | waited, |
173 | 153 | is_due, |
174 | 153 | action_eligible, |
175 | 153 | should_warn, |
176 | 153 | } |
177 | 153 | } |
178 | | |
179 | | /// Records that an action with this shape was failed this pass, so a |
180 | | /// pass that retired every eligible action does not read as stuck. |
181 | 29 | pub fn note_failed(&mut self, shape: &PropertyShape) { |
182 | 29 | if let Some(state) = self.shapes.get_mut(shape) { |
183 | 29 | state.failed_this_pass += 1; |
184 | 29 | }0 |
185 | 29 | } |
186 | | |
187 | | /// Whether a shape that stays unsatisfiable is ever failed. |
188 | | #[must_use] |
189 | 122 | pub const fn fails_actions(&self) -> bool { |
190 | 122 | self.timeout.is_some() |
191 | 122 | } |
192 | | |
193 | | /// Closes a matching pass and returns how many unsatisfiable actions it |
194 | | /// saw. |
195 | | /// |
196 | | /// A shape that was not seen is kept for one timeout period so its clock |
197 | | /// keeps running and it stays due for actions that arrive one after |
198 | | /// another; each such action still waits its own minimum before it is |
199 | | /// failed. The shape is dropped at once if the fleet changed, because a |
200 | | /// new worker may be able to run it. |
201 | | /// |
202 | | /// A due shape backs off (see `next_deadline`) only when an eligible |
203 | | /// action of it was seen but not failed this pass. |
204 | 558 | pub fn end_pass(&mut self, now: SystemTime, fleet_generation: u64) -> u64 { |
205 | 558 | let linger = self.timeout.unwrap_or_default(); |
206 | 558 | self.shapes.retain(|_, state| {125 |
207 | 125 | if state.seen_this_pass |
208 | 115 | && state.due_this_pass |
209 | 46 | && state.eligible_this_pass > state.failed_this_pass |
210 | 8 | { |
211 | 8 | state.due_passes = state.due_passes.saturating_add(1); |
212 | 117 | } |
213 | 125 | state.seen_last_pass = state.seen_this_pass; |
214 | 125 | state.seen_this_pass = false; |
215 | 125 | state.seen_last_pass |
216 | 10 | || (state.fleet_generation == fleet_generation |
217 | 5 | && now.duration_since(state.last_seen).unwrap_or_default() < linger) |
218 | 125 | }); |
219 | 558 | core::mem::take(&mut self.seen_this_pass) |
220 | 558 | } |
221 | | |
222 | | /// How long until the loop should run again for a shape seen in the |
223 | | /// last pass, if any: when the shape comes due, when its youngest action |
224 | | /// becomes eligible, or, for a due shape whose eligible actions could |
225 | | /// not be retired, a backoff that doubles up to 32 seconds. |
226 | | #[must_use] |
227 | 525 | pub fn next_deadline(&self, now: SystemTime) -> Option<Duration> { |
228 | 525 | let timeout127 = self.timeout?398 ; |
229 | 127 | self.shapes |
230 | 127 | .values() |
231 | 127 | .filter(|state| state.seen_last_pass) |
232 | 127 | .filter_map(|state| {44 |
233 | 44 | let waited = now |
234 | 44 | .duration_since(state.first_unsatisfiable) |
235 | 44 | .unwrap_or_default(); |
236 | 44 | let shape_term = if waited < timeout { |
237 | 26 | Some(timeout.saturating_sub(waited)) |
238 | 18 | } else if state.eligible_this_pass > state.failed_this_pass { |
239 | | // Still queued though eligible, so wake less and less |
240 | | // often for it. |
241 | 7 | Some(MIN_RECHECK * 2u32.pow(state.due_passes.min(MAX_DUE_BACKOFF_DOUBLINGS))) |
242 | | } else { |
243 | | // Due, but nothing of it is stuck: only a young action can |
244 | | // change anything, and that is covered below. |
245 | 11 | None |
246 | | }; |
247 | 44 | let eligible_term = state |
248 | 44 | .earliest_eligible |
249 | 44 | .map(|eligible_at| eligible_at34 .duration_since34 (now34 ).unwrap_or_default34 ()); |
250 | 44 | match (shape_term, eligible_term) { |
251 | 26 | (Some(a), Some(b)) => Some(a.min(b)), |
252 | 18 | (a, b) => a.or(b), |
253 | | } |
254 | 44 | .map(|deadline| deadline41 .max41 (MIN_RECHECK)) |
255 | 44 | }) |
256 | 127 | .min() |
257 | 525 | } |
258 | | } |