/build/source/nativelink-scheduler/src/match_outcome.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 | | //! The result of trying to place one queued action on the worker fleet. |
16 | | |
17 | | use core::fmt; |
18 | | use std::collections::{BTreeSet, HashMap}; |
19 | | use std::sync::Arc; |
20 | | |
21 | | use nativelink_util::action_messages::WorkerId; |
22 | | use nativelink_util::platform_properties::{PlatformProperties, PlatformPropertyValue}; |
23 | | |
24 | | /// How many distinct worker values are named for one unsatisfied property. |
25 | | const MAX_OFFERED_VALUES: usize = 8; |
26 | | |
27 | | /// What happened when the scheduler looked for a worker for an action. |
28 | | #[derive(Debug, Clone, PartialEq, Eq)] |
29 | | pub enum MatchOutcome { |
30 | | /// This worker can run the action now. |
31 | | Matched(WorkerId), |
32 | | /// A worker could run the action when idle, but none has room right now. |
33 | | WaitingForCapacity, |
34 | | /// No workers are connected at all. |
35 | | NoWorkersConnected, |
36 | | /// No worker could run the action even when fully idle. |
37 | | Unsatisfiable(Arc<UnsatisfiableReason>), |
38 | | } |
39 | | |
40 | | /// One action property that the fleet cannot satisfy. |
41 | | #[derive(Debug, Clone, PartialEq, Eq)] |
42 | | pub struct UnsatisfiedProperty { |
43 | | pub name: String, |
44 | | pub requested: PlatformPropertyValue, |
45 | | /// For a `Minimum` property, the largest total any worker registered |
46 | | /// with. Otherwise the distinct values workers registered with, capped at |
47 | | /// `MAX_OFFERED_VALUES`. Empty when no worker declares the property. |
48 | | pub offered: Vec<String>, |
49 | | /// How many distinct values were left out of `offered`. |
50 | | pub offered_omitted: usize, |
51 | | } |
52 | | |
53 | | /// Why no worker can run an action. |
54 | | #[derive(Debug, Clone, PartialEq, Eq)] |
55 | | pub struct UnsatisfiableReason { |
56 | | /// Sorted by property name. |
57 | | pub properties: Vec<UnsatisfiedProperty>, |
58 | | /// True when every property is satisfied by some worker, but no single |
59 | | /// worker satisfies all of them. `properties` then lists every property |
60 | | /// that restricts matching. |
61 | | pub combination_only: bool, |
62 | | } |
63 | | |
64 | | impl UnsatisfiableReason { |
65 | | /// No worker is connected at all, here or on any peer: nothing to judge |
66 | | /// the action's properties against, so no property is named. |
67 | | #[must_use] |
68 | 525 | pub const fn no_workers() -> Self { |
69 | 525 | Self { |
70 | 525 | properties: Vec::new(), |
71 | 525 | combination_only: false, |
72 | 525 | } |
73 | 525 | } |
74 | | |
75 | | /// Whether this is `no_workers`. |
76 | | #[must_use] |
77 | 168 | pub const fn is_no_workers(&self) -> bool { |
78 | 168 | self.properties.is_empty() && !self.combination_only5 |
79 | 168 | } |
80 | | |
81 | | /// The unsatisfied property names, comma separated, for metric labels. |
82 | | #[must_use] |
83 | 29 | pub fn property_names(&self) -> String { |
84 | 29 | if self.is_no_workers() { |
85 | 1 | return "no_workers".to_string(); |
86 | 28 | } |
87 | 28 | self.properties |
88 | 28 | .iter() |
89 | 53 | .map28 (|property| property.name.as_str()) |
90 | 28 | .collect::<Vec<_>>() |
91 | 28 | .join(",") |
92 | 29 | } |
93 | | |
94 | | /// The reason as told to the client that submitted the action: the |
95 | | /// properties it asked for, without the values workers were registered |
96 | | /// with, which can name internal images or pools. The largest total of |
97 | | /// a `minimum` property is kept, as it is only a number. |
98 | | #[must_use] |
99 | 28 | pub const fn for_client(&self) -> ClientReason<'_> { |
100 | 28 | ClientReason(self) |
101 | 28 | } |
102 | | |
103 | 139 | fn write(&self, f: &mut fmt::Formatter<'_>, with_offered_values: bool) -> fmt::Result { |
104 | 139 | if self.is_no_workers() { |
105 | 4 | return f.write_str("no worker is connected to the scheduler"); |
106 | 135 | } |
107 | 135 | if self.combination_only { |
108 | 1 | f.write_str("no single worker satisfies the combination of: ")?0 ; |
109 | | } else { |
110 | 134 | f.write_str("no worker can satisfy: ")?0 ; |
111 | | } |
112 | 250 | for (i, property) in self.properties.iter()135 .enumerate135 () { |
113 | 250 | if i > 0 { |
114 | 115 | f.write_str("; ")?0 ; |
115 | 135 | } |
116 | | // A priority property only needs the key on the worker, so its |
117 | | // value is not what went unsatisfied. |
118 | 250 | if matches!249 (property.requested, PlatformPropertyValue::Priority(_)) { |
119 | 1 | write!(f, "'{}' is required", property.name)?0 ; |
120 | | } else { |
121 | 249 | write!( |
122 | 249 | f, |
123 | | "'{}' requested {}", |
124 | | property.name, |
125 | 249 | property.requested.as_str() |
126 | 0 | )?; |
127 | | } |
128 | 250 | if property.offered.is_empty() { |
129 | 2 | f.write_str(", no worker declares this property")?0 ; |
130 | 248 | } else if matches!7 (property.requested, PlatformPropertyValue::Minimum(_)) { |
131 | 241 | write!(f, ", largest worker total {}", property.offered[0])?0 ; |
132 | 7 | } else if with_offered_values { |
133 | 6 | write!(f, ", workers offer [{}]", property.offered.join(", "))?0 ; |
134 | 6 | if property.offered_omitted > 0 { |
135 | 1 | write!(f, " and {} more", property.offered_omitted)?0 ; |
136 | 5 | } |
137 | | } else { |
138 | 1 | f.write_str(", no worker offers it")?0 ; |
139 | | } |
140 | | } |
141 | 135 | Ok(()) |
142 | 139 | } |
143 | | } |
144 | | |
145 | | impl fmt::Display for UnsatisfiableReason { |
146 | 111 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
147 | 111 | self.write(f, true) |
148 | 111 | } |
149 | | } |
150 | | |
151 | | /// See `UnsatisfiableReason::for_client`. |
152 | | #[derive(Debug, Clone, Copy)] |
153 | | pub struct ClientReason<'a>(&'a UnsatisfiableReason); |
154 | | |
155 | | impl fmt::Display for ClientReason<'_> { |
156 | 28 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
157 | 28 | self.0.write(f, false) |
158 | 28 | } |
159 | | } |
160 | | |
161 | | /// The parts of an action's platform properties that decide which workers can |
162 | | /// run it, in a form that can key a map. `Ignore` properties never restrict |
163 | | /// matching and a `Priority` property only requires the key, so actions that |
164 | | /// differ only in those share a shape. |
165 | | #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] |
166 | | pub struct PropertyShape(Vec<(String, PlatformPropertyValue)>); |
167 | | |
168 | | impl From<&PlatformProperties> for PropertyShape { |
169 | 620 | fn from(platform_properties: &PlatformProperties) -> Self { |
170 | 620 | let mut shape: Vec<_> = platform_properties |
171 | 620 | .properties |
172 | 620 | .iter() |
173 | 1.71k | .filter_map620 (|(name, value)| match value { |
174 | 2 | PlatformPropertyValue::Ignore(_) => None, |
175 | | PlatformPropertyValue::Priority(_) => { |
176 | 311 | Some((name.clone(), PlatformPropertyValue::Priority(String::new()))) |
177 | | } |
178 | 1.40k | _ => Some((name.clone(), value.clone())), |
179 | 1.71k | }) |
180 | 620 | .collect(); |
181 | 620 | shape.sort_unstable(); |
182 | 620 | Self(shape) |
183 | 620 | } |
184 | | } |
185 | | |
186 | 163 | fn offered_values( |
187 | 163 | name: &str, |
188 | 163 | requested: &PlatformPropertyValue, |
189 | 163 | fleet: &[&PlatformProperties], |
190 | 163 | ) -> (Vec<String>, usize) { |
191 | 163 | let declared = fleet |
192 | 163 | .iter() |
193 | 176 | .filter_map163 (|worker_totals| worker_totals.properties.get(name)); |
194 | 163 | if matches!6 (requested, PlatformPropertyValue::Minimum(_)) { |
195 | 157 | let largest = declared |
196 | 160 | .filter_map157 (|value| match value { |
197 | 160 | PlatformPropertyValue::Minimum(total) => Some(*total), |
198 | 0 | _ => None, |
199 | 160 | }) |
200 | 157 | .max(); |
201 | 157 | return (largest.map(|v| v156 .to_string156 ()).into_iter().collect(), 0); |
202 | 6 | } |
203 | 14 | let distinct6 : BTreeSet<String>6 = declared6 .map6 (|value| value.as_str().into_owned()).collect6 (); |
204 | 6 | let omitted = distinct.len().saturating_sub(MAX_OFFERED_VALUES); |
205 | 6 | ( |
206 | 6 | distinct.into_iter().take(MAX_OFFERED_VALUES).collect(), |
207 | 6 | omitted, |
208 | 6 | ) |
209 | 163 | } |
210 | | |
211 | | /// Works out which properties stop every worker in `fleet` from running the |
212 | | /// action. `fleet` holds the properties each worker registered with. Only |
213 | | /// call this once it is known that no worker in `fleet` satisfies the action. |
214 | | #[must_use] |
215 | 89 | pub fn explain_unsatisfiable( |
216 | 89 | action_properties: &PlatformProperties, |
217 | 89 | fleet: &[&PlatformProperties], |
218 | 89 | ) -> UnsatisfiableReason { |
219 | 89 | let mut restricting: Vec<_> = action_properties |
220 | 89 | .properties |
221 | 89 | .iter() |
222 | 386 | .filter89 (|(_, value)| !matches!(value, PlatformPropertyValue::Ignore(_))) |
223 | 89 | .collect(); |
224 | 607 | restricting89 .sort_unstable_by89 (|a, b| a.0.cmp(b.0)); |
225 | | |
226 | 163 | let to_unsatisfied89 = |(name, requested): &(&String, &PlatformPropertyValue)| { |
227 | 163 | let (offered, offered_omitted) = offered_values(name, requested, fleet); |
228 | 163 | UnsatisfiedProperty { |
229 | 163 | name: (*name).clone(), |
230 | 163 | requested: (*requested).clone(), |
231 | 163 | offered, |
232 | 163 | offered_omitted, |
233 | 163 | } |
234 | 163 | }; |
235 | | |
236 | | // How many workers each property rules out on its own. |
237 | 393 | let rules_out89 = |(name, requested): &(&String, &PlatformPropertyValue)| -> usize { |
238 | 393 | let alone = |
239 | 393 | PlatformProperties::new(HashMap::from([((*name).clone(), (*requested).clone())])); |
240 | 393 | fleet |
241 | 393 | .iter() |
242 | 416 | .filter393 (|worker_totals| !alone.is_satisfied_by(worker_totals, false)) |
243 | 393 | .count() |
244 | 393 | }; |
245 | | |
246 | 89 | let properties: Vec<_> = restricting |
247 | 89 | .iter() |
248 | 386 | .filter89 (|property| rules_out(property) == fleet.len()) |
249 | 89 | .map(to_unsatisfied) |
250 | 89 | .collect(); |
251 | | |
252 | 89 | if properties.is_empty() { |
253 | | // Only the properties that rule out some worker take part in the |
254 | | // combination; one every worker satisfies is not a cause. |
255 | | return UnsatisfiableReason { |
256 | 2 | properties: restricting |
257 | 2 | .iter() |
258 | 7 | .filter2 (|property| rules_out(property) > 0) |
259 | 2 | .map(to_unsatisfied) |
260 | 2 | .collect(), |
261 | | combination_only: true, |
262 | | }; |
263 | 87 | } |
264 | 87 | UnsatisfiableReason { |
265 | 87 | properties, |
266 | 87 | combination_only: false, |
267 | 87 | } |
268 | 89 | } |