/build/source/nativelink-worker/src/local_worker.rs
Line | Count | Source |
1 | | // Copyright 2024 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 | | use core::hash::BuildHasher; |
16 | | use core::pin::Pin; |
17 | | use core::str; |
18 | | use core::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}; |
19 | | use core::time::Duration; |
20 | | use std::borrow::Cow; |
21 | | use std::collections::HashMap; |
22 | | use std::env; |
23 | | use std::process::Stdio; |
24 | | use std::sync::{Arc, Weak}; |
25 | | use std::time::{SystemTime, UNIX_EPOCH}; |
26 | | |
27 | | use futures::future::BoxFuture; |
28 | | use futures::stream::FuturesUnordered; |
29 | | use futures::{Future, FutureExt, StreamExt, TryFutureExt, select}; |
30 | | use nativelink_config::cas_server::{EnvironmentSource, LocalWorkerConfig}; |
31 | | use nativelink_error::{Code, Error, ResultExt, make_err, make_input_err}; |
32 | | use nativelink_metric::{MetricsComponent, RootMetricsComponent}; |
33 | | use nativelink_proto::com::github::trace_machina::nativelink::remote_execution::update_for_worker::Update; |
34 | | use nativelink_proto::com::github::trace_machina::nativelink::remote_execution::worker_api_client::WorkerApiClient; |
35 | | use nativelink_proto::com::github::trace_machina::nativelink::remote_execution::{ |
36 | | ActionResourceUsage, ExecuteAccepted, ExecuteComplete, ExecuteDeclined, ExecuteResult, |
37 | | GoingAwayRequest, KeepAliveRequest, StartExecute, UpdateForWorker, WorkerLoad, |
38 | | execute_declined, execute_result, |
39 | | }; |
40 | | use nativelink_store::fast_slow_store::FastSlowStore; |
41 | | use nativelink_util::action_messages::{ActionResult, ActionStage, OperationId}; |
42 | | use nativelink_util::common::fs; |
43 | | use nativelink_util::digest_hasher::DigestHasherFunc; |
44 | | use nativelink_util::health_utils::{HealthStatus, HealthStatusIndicator}; |
45 | | use nativelink_util::metrics_utils::{AsyncCounterWrapper, CounterWithTime}; |
46 | | use nativelink_util::shutdown_guard::ShutdownGuard; |
47 | | use nativelink_util::store_trait::Store; |
48 | | use nativelink_util::{background_spawn, spawn, spawn_blocking, tls_utils}; |
49 | | use opentelemetry::context::Context; |
50 | | use tokio::sync::{broadcast, mpsc}; |
51 | | use tokio::{process, time}; |
52 | | use tokio_stream::wrappers::UnboundedReceiverStream; |
53 | | use tonic::{Streaming, async_trait}; |
54 | | use tracing::{Level, debug, error, event, info, info_span, instrument, trace, warn}; |
55 | | |
56 | | use crate::running_actions_manager::{ |
57 | | ExecutionConfiguration, MISSING_INPUT_ERROR_TIP, Metrics as RunningActionManagerMetrics, |
58 | | ResourceEnforcement, RunningAction, RunningActionsManager, RunningActionsManagerArgs, |
59 | | RunningActionsManagerImpl, |
60 | | }; |
61 | | use crate::worker_api_client_wrapper::{WorkerApiClientTrait, WorkerApiClientWrapper}; |
62 | | use crate::worker_utils::make_connect_worker_request; |
63 | | |
64 | | /// Amount of time to wait if we have actions in transit before we try to |
65 | | /// consider an error to have occurred. |
66 | | const ACTIONS_IN_TRANSIT_TIMEOUT_S: f32 = 10.; |
67 | | |
68 | | /// Increments `actions_in_transit` on creation and decrements it on drop, so |
69 | | /// the count stays accurate even when the owning action future is aborted by |
70 | | /// a disconnect. A stranded count makes the disconnect handler conclude the |
71 | | /// in-transit actions never drained, turning every disconnect-with-work into |
72 | | /// a fatal error instead of a reconnect. |
73 | | struct ActionsInTransitGuard { |
74 | | actions_in_transit: Arc<AtomicU64>, |
75 | | } |
76 | | |
77 | | impl ActionsInTransitGuard { |
78 | 0 | fn new(actions_in_transit: Arc<AtomicU64>) -> Self { |
79 | 0 | actions_in_transit.fetch_add(1, Ordering::Release); |
80 | 0 | Self { actions_in_transit } |
81 | 0 | } |
82 | | } |
83 | | |
84 | | impl Drop for ActionsInTransitGuard { |
85 | 0 | fn drop(&mut self) { |
86 | 0 | self.actions_in_transit.fetch_sub(1, Ordering::Release); |
87 | 0 | } |
88 | | } |
89 | | |
90 | | /// If we lose connection to the worker api server we will wait this many seconds |
91 | | /// before trying to connect, doubling on every failed attempt up to |
92 | | /// `CONNECTION_RETRY_MAX_DELAY_S`, with jitter so a fleet that lost its |
93 | | /// scheduler together does not redial together. |
94 | | const CONNECTION_RETRY_DELAY_S: f32 = 0.5; |
95 | | const CONNECTION_RETRY_MAX_DELAY_S: f32 = 30.0; |
96 | | |
97 | | /// Delay before the `attempt`th consecutive reconnect (0 = first retry), |
98 | | /// between half and one and a half times the exponential figure. |
99 | 6 | fn reconnect_delay(attempt: u32) -> Duration { |
100 | 6 | let exponent = i32::try_from(attempt.min(16)).unwrap_or(16); |
101 | 6 | let base = (CONNECTION_RETRY_DELAY_S * 2f32.powi(exponent)).min(CONNECTION_RETRY_MAX_DELAY_S); |
102 | 6 | let nanos = SystemTime::now() |
103 | 6 | .duration_since(UNIX_EPOCH) |
104 | 6 | .map_or(0, |d| d.subsec_nanos()); |
105 | 6 | let jitter = 0.5 + (nanos % 1_000) as f32 / 1_000.0; |
106 | 6 | Duration::from_secs_f32(base * jitter) |
107 | 6 | } |
108 | | |
109 | | /// Default endpoint timeout. If this value gets modified the documentation in |
110 | | /// `cas_server.rs` must also be updated. |
111 | | const DEFAULT_ENDPOINT_TIMEOUT_S: f32 = 5.; |
112 | | |
113 | | /// Default maximum amount of time a task is allowed to run for. |
114 | | /// If this value gets modified the documentation in `cas_server.rs` must also be updated. |
115 | | const DEFAULT_MAX_ACTION_TIMEOUT: Duration = Duration::from_mins(20); |
116 | | const DEFAULT_MAX_UPLOAD_TIMEOUT: Duration = Duration::from_mins(10); |
117 | | /// Default for `max_download_timeout_s`: an input fetch is quick or wedged. |
118 | | const DEFAULT_MAX_DOWNLOAD_TIMEOUT: Duration = Duration::from_mins(10); |
119 | | const DEFAULT_MAX_CLEANUP_WAIT: Duration = Duration::from_secs(30); |
120 | | const DEFAULT_MAX_CLEANUP_BACKOFF: Duration = Duration::from_millis(500); |
121 | | /// If this value gets modified the documentation in `cas_server.rs` must also be updated. |
122 | | const DEFAULT_PRECONDITION_TIMEOUT: Duration = Duration::from_secs(30); |
123 | | /// If this value gets modified the documentation in `cas_server.rs` must also be updated. |
124 | | const DEFAULT_KILL_GRACE: Duration = Duration::from_secs(5); |
125 | | |
126 | | struct FinishedActionResult { |
127 | | action_result: ActionResult, |
128 | | resource_usage: Option<ActionResourceUsage>, |
129 | | } |
130 | | |
131 | | struct LocalWorkerImpl<'a, T: WorkerApiClientTrait + 'static, U: RunningActionsManager> { |
132 | | config: &'a LocalWorkerConfig, |
133 | | // According to the tonic documentation it is a cheap operation to clone this. |
134 | | grpc_client: T, |
135 | | worker_id: String, |
136 | | running_actions_manager: Arc<U>, |
137 | | // Number of actions that have been received in `Update::StartAction`, but |
138 | | // not yet processed by running_actions_manager's spawn. This number should |
139 | | // always be zero if there are no actions running and no actions being waited |
140 | | // on by the scheduler. |
141 | | actions_in_transit: Arc<AtomicU64>, |
142 | | accepted_action: AtomicBool, |
143 | | /// The scheduler said it understands `ExecuteAccepted` and |
144 | | /// `ExecuteDeclined`. Without it the worker runs whatever it is sent, |
145 | | /// as every earlier release did. |
146 | | dispatch_ack: bool, |
147 | | /// The platform property the scheduler reads an action's memory |
148 | | /// reservation from, as it told us on connection; empty when it does |
149 | | /// not veto on memory. Read from the same property, a refusal for load |
150 | | /// here agrees with what the scheduler would have vetoed. |
151 | | memory_property: String, |
152 | | metrics: Arc<Metrics>, |
153 | | } |
154 | | |
155 | | /// Why this worker will not run an action it was just sent. |
156 | | enum Refusal { |
157 | | AtCapacity { in_flight: u64, max: u64 }, |
158 | | Load { needed_kb: u64, free_kb: u64 }, |
159 | | ShuttingDown, |
160 | | } |
161 | | |
162 | | impl Refusal { |
163 | 2 | fn into_declined(self, operation_id: String) -> ExecuteDeclined { |
164 | 2 | let (reason, detail, needed_kb, free_kb) = match self { |
165 | 1 | Self::AtCapacity { in_flight, max } => ( |
166 | 1 | execute_declined::Reason::AtCapacity, |
167 | 1 | format!("{in_flight} of {max} in flight"), |
168 | 1 | 0, |
169 | 1 | 0, |
170 | 1 | ), |
171 | 1 | Self::Load { needed_kb, free_kb } => ( |
172 | 1 | execute_declined::Reason::Load, |
173 | 1 | format!("needs {needed_kb} KiB, {free_kb} KiB free"), |
174 | 1 | needed_kb, |
175 | 1 | free_kb, |
176 | 1 | ), |
177 | 0 | Self::ShuttingDown => ( |
178 | 0 | execute_declined::Reason::ShuttingDown, |
179 | 0 | "worker shutting down".to_string(), |
180 | 0 | 0, |
181 | 0 | 0, |
182 | 0 | ), |
183 | | }; |
184 | 2 | ExecuteDeclined { |
185 | 2 | operation_id, |
186 | 2 | reason: reason as i32, |
187 | 2 | detail, |
188 | 2 | needed_kb, |
189 | 2 | free_kb, |
190 | 2 | } |
191 | 2 | } |
192 | | } |
193 | | |
194 | | /// The memory reservation an action carries under `property`, in KiB, if |
195 | | /// it declares one. An empty property name is a scheduler that does not |
196 | | /// veto on memory, so nothing is ever read. |
197 | 3 | fn memory_reservation_kb(start_execute: &StartExecute, property: &str) -> Option<u64> { |
198 | 3 | if property.is_empty() { |
199 | 1 | return None; |
200 | 2 | } |
201 | 2 | start_execute |
202 | 2 | .platform |
203 | 2 | .as_ref()?0 |
204 | | .properties |
205 | 2 | .iter() |
206 | 2 | .find(|p| p.name1 == property1 ) |
207 | 2 | .and_then(|p| p.value.parse::<u64>()1 .ok1 ()) |
208 | 2 | .filter(|kb| *kb1 > 0) |
209 | 3 | } |
210 | | |
211 | 13 | pub async fn preconditions_met<H: BuildHasher + Sync>( |
212 | 13 | precondition_script: Option<String>, |
213 | 13 | extra_envs: &HashMap<String, String, H>, |
214 | 13 | timeout: Duration, |
215 | 13 | ) -> Result<(), Error> { |
216 | 13 | let Some(precondition_script3 ) = &precondition_script else { |
217 | | // No script means we are always ok to proceed. |
218 | 10 | return Ok(()); |
219 | | }; |
220 | | // TODO: Might want to pass some information about the command to the |
221 | | // script, but at this point it's not even been downloaded yet, |
222 | | // so that's not currently possible. Perhaps we'll move this in |
223 | | // future to pass useful information through? Or perhaps we'll |
224 | | // have a pre-condition and a pre-execute script instead, although |
225 | | // arguably entrypoint already gives us that. |
226 | | |
227 | 3 | let maybe_split_cmd = shlex::split(precondition_script); |
228 | 3 | let (command, args) = match &maybe_split_cmd { |
229 | 3 | Some(split_cmd) => (&split_cmd[0], &split_cmd[1..]), |
230 | | None => { |
231 | 0 | return Err(make_input_err!( |
232 | 0 | "Could not parse the value of precondition_script: '{}'", |
233 | 0 | precondition_script, |
234 | 0 | )); |
235 | | } |
236 | | }; |
237 | | |
238 | 3 | let precondition_process = process::Command::new(command) |
239 | 3 | .args(args) |
240 | 3 | .kill_on_drop(true) |
241 | 3 | .stdin(Stdio::null()) |
242 | 3 | .stdout(Stdio::piped()) |
243 | 3 | .stderr(Stdio::null()) |
244 | 3 | .env_clear() |
245 | 3 | .envs(extra_envs) |
246 | 3 | .spawn() |
247 | 3 | .err_tip(|| format!0 ("Could not execute precondition command {precondition_script:?}"))?0 ; |
248 | 3 | let _owned = crate::reaper::OwnedChild::new(precondition_process.id()); |
249 | | // Bounded: a script that hangs held the action forever, and |
250 | | // `kill_on_drop` ends the script when the timeout drops it. |
251 | 3 | let output2 = match time::timeout(timeout, precondition_process.wait_with_output()).await { |
252 | 2 | Ok(output) => output?0 , |
253 | | Err(_) => { |
254 | 1 | return Err(make_err!( |
255 | 1 | Code::ResourceExhausted, |
256 | 1 | "Preconditions script {precondition_script:?} did not finish within {} ms", |
257 | 1 | timeout.as_millis() |
258 | 1 | )); |
259 | | } |
260 | | }; |
261 | 2 | let stdout = str::from_utf8(&output.stdout).unwrap_or(""); |
262 | 2 | trace!(status = %output.status, %stdout, "Preconditions script returned"); |
263 | 2 | if output.status.code() == Some(0) { |
264 | 1 | Ok(()) |
265 | | } else { |
266 | 1 | Err(make_err!( |
267 | 1 | Code::ResourceExhausted, |
268 | 1 | "Preconditions script returned status {} - {}", |
269 | 1 | output.status, |
270 | 1 | stdout |
271 | 1 | )) |
272 | | } |
273 | 13 | } |
274 | | |
275 | | impl<'a, T: WorkerApiClientTrait + 'static, U: RunningActionsManager> LocalWorkerImpl<'a, T, U> { |
276 | 15 | fn new( |
277 | 15 | config: &'a LocalWorkerConfig, |
278 | 15 | grpc_client: T, |
279 | 15 | worker_id: String, |
280 | 15 | dispatch_ack: bool, |
281 | 15 | memory_property: String, |
282 | 15 | running_actions_manager: Arc<U>, |
283 | 15 | metrics: Arc<Metrics>, |
284 | 15 | ) -> Self { |
285 | 15 | Self { |
286 | 15 | config, |
287 | 15 | grpc_client, |
288 | 15 | worker_id, |
289 | 15 | dispatch_ack, |
290 | 15 | memory_property, |
291 | 15 | running_actions_manager, |
292 | 15 | // Number of actions that have been received in `Update::StartAction`, but |
293 | 15 | // not yet processed by running_actions_manager's spawn. This number should |
294 | 15 | // always be zero if there are no actions running and no actions being waited |
295 | 15 | // on by the scheduler. |
296 | 15 | actions_in_transit: Arc::new(AtomicU64::new(0)), |
297 | 15 | accepted_action: AtomicBool::new(false), |
298 | 15 | metrics, |
299 | 15 | } |
300 | 15 | } |
301 | | |
302 | | /// Local admission: the worker's own word on whether it can take this |
303 | | /// action now. The scheduler's ledger says what it believes the worker |
304 | | /// has; this is what the worker has. |
305 | 4 | fn admission(&self, start_execute: &StartExecute, in_flight: u64) -> Option<Refusal> { |
306 | 4 | let max = self.config.max_inflight_tasks; |
307 | 4 | if max > 0 && in_flight >= max2 { |
308 | 1 | return Some(Refusal::AtCapacity { in_flight, max }); |
309 | 3 | } |
310 | 1 | if let (Some(needed_kb), Some(free_kb)) = ( |
311 | 3 | memory_reservation_kb(start_execute, &self.memory_property), |
312 | 3 | crate::capacity::free_memory_kb(), |
313 | 1 | ) && free_kb < needed_kb |
314 | | { |
315 | 1 | return Some(Refusal::Load { needed_kb, free_kb }); |
316 | 2 | } |
317 | 2 | None |
318 | 4 | } |
319 | | |
320 | | /// Tells the scheduler the action is refused; only meaningful when it |
321 | | /// understands the message. |
322 | 2 | async fn decline(&self, operation_id: String, refusal: Refusal) -> Result<(), Error> { |
323 | 2 | self.metrics.actions_declined.inc(); |
324 | 2 | self.grpc_client |
325 | 2 | .clone() |
326 | 2 | .execute_declined(refusal.into_declined(operation_id)) |
327 | 2 | .await |
328 | 1 | .err_tip(|| "Could not send ExecuteDeclined") |
329 | 1 | } |
330 | | |
331 | | /// Starts a background spawn/thread that will send a message to the server every `timeout / 2`. |
332 | 15 | async fn start_keep_alive(&self) -> Result<(), Error> { |
333 | | // According to tonic's documentation this call should be cheap and is the same stream. |
334 | 15 | let mut grpc_client = self.grpc_client.clone(); |
335 | 15 | let timeout = self |
336 | 15 | .config |
337 | 15 | .worker_api_endpoint |
338 | 15 | .timeout |
339 | 15 | .unwrap_or(DEFAULT_ENDPOINT_TIMEOUT_S); |
340 | | |
341 | 15 | info!(timeout, "Started KeepAlive"); |
342 | | |
343 | | // We always send 2 keep alive requests per timeout. Http2 should manage most of our |
344 | | // timeout issues, this is a secondary check to ensure we can still send data. |
345 | 15 | let mut interval = time::interval(Duration::from_secs_f32(timeout) / 2); |
346 | 15 | interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip); |
347 | | |
348 | | // Skip the first interval as it happens immediately and we don't need a keep alive until timeout/2 has passed |
349 | 15 | interval.tick().await; |
350 | | |
351 | | // Explicitly spawn the keep alive loop so it goes onto a different thread from the execute commands. |
352 | | // Its failure is this worker's failure: a worker whose keepalives |
353 | | // stop reaching the scheduler is evicted `worker_timeout_s` later |
354 | | // with every action it holds requeued, so ending the stream now and |
355 | | // reconnecting is the cheaper outcome. Before, the task died quietly |
356 | | // and the worker ran on without keepalives. |
357 | 4 | spawn!("keep alives", async move { |
358 | | loop { |
359 | 5 | interval.tick().await; |
360 | | // What the worker has to spare rides on the keepalive, so a |
361 | | // scheduler with `live_memory_veto` can skip a worker whose |
362 | | // actions declared less than they use. |
363 | 2 | let load = crate::capacity::free_memory_kb() |
364 | 2 | .map(|free_memory_kb| WorkerLoad { free_memory_kb }); |
365 | 2 | if let Err(e1 ) = grpc_client.keep_alive(KeepAliveRequest { load }).await { |
366 | 1 | error!(?e, "Failed to send KeepAlive in LocalWorker"); |
367 | 1 | return Err(e.append("KeepAlive failed; reconnecting to the scheduler")); |
368 | 1 | } |
369 | 1 | debug!("Sent KeepAlive"); |
370 | | } |
371 | 1 | }) |
372 | 4 | .await |
373 | 1 | .map_err(|e| make_err!0 (Code::Internal0 , "KeepAlive task ended: {e:?}"))?0 |
374 | 1 | } |
375 | | |
376 | 15 | async fn run( |
377 | 15 | &self, |
378 | 15 | update_for_worker_stream: Streaming<UpdateForWorker>, |
379 | 15 | shutdown_rx: &mut broadcast::Receiver<ShutdownGuard>, |
380 | 15 | ) -> Result<(), Error> { |
381 | | // This big block of logic is designed to help simplify upstream components. Upstream |
382 | | // components can write standard futures that return a `Result<(), Error>` and this block |
383 | | // will forward the error up to the client and disconnect from the scheduler. |
384 | | // It is a common use case that an item sent through update_for_worker_stream will always |
385 | | // have a response but the response will be triggered through a callback to the scheduler. |
386 | | // This can be quite tricky to manage, so what we have done here is given access to a |
387 | | // `futures` variable which because this is in a single thread as well as a channel that you |
388 | | // send a future into that makes it into the `futures` variable. |
389 | | // This means that if you want to perform an action based on the result of the future |
390 | | // you use the `.map()` method and the new action will always come to live in this spawn, |
391 | | // giving mutable access to stuff in this struct. |
392 | | // NOTE: If you ever return from this function it will disconnect from the scheduler. |
393 | 15 | let mut futures = FuturesUnordered::new(); |
394 | 15 | futures.push(self.start_keep_alive().boxed()); |
395 | | |
396 | 15 | let (add_future_channel, add_future_rx) = mpsc::unbounded_channel(); |
397 | 15 | let mut add_future_rx = UnboundedReceiverStream::new(add_future_rx).fuse(); |
398 | | |
399 | 15 | let mut update_for_worker_stream = update_for_worker_stream.fuse(); |
400 | | // A notify which is triggered every time actions_in_flight is subtracted. |
401 | 15 | let actions_notify = Arc::new(tokio::sync::Notify::new()); |
402 | | // A counter of actions that are in-flight, this is similar to actions_in_transit but |
403 | | // includes the AC upload and notification to the scheduler. |
404 | 15 | let actions_in_flight = Arc::new(AtomicU64::new(0)); |
405 | | // Set to true when shutting down, this stops any new StartAction. |
406 | 15 | let mut shutting_down = false; |
407 | | |
408 | | loop { |
409 | 43 | if self.config.single_use |
410 | 11 | && self.accepted_action.load(Ordering::Acquire) |
411 | 7 | && actions_in_flight.load(Ordering::Acquire) == 0 |
412 | | { |
413 | | // The action, CAS/AC uploads, cleanup, and execution_response |
414 | | // acknowledgment have all completed. This container is spent. |
415 | 1 | if let Err(err0 ) = self |
416 | 1 | .grpc_client |
417 | 1 | .clone() |
418 | 1 | .going_away(GoingAwayRequest { drain: false }) |
419 | 1 | .await |
420 | | { |
421 | 0 | warn!(?err, "Could not unregister completed single-use worker"); |
422 | 1 | } |
423 | 1 | return Ok(()); |
424 | 42 | } |
425 | 42 | select! { |
426 | 42 | maybe_update21 = update_for_worker_stream.next() => if !shutting_down21 || maybe_update0 .is_some0 () { |
427 | 21 | match maybe_update |
428 | 21 | .err_tip(|| "UpdateForWorker stream closed early")?5 |
429 | 16 | .err_tip(|| "Got error in UpdateForWorker stream")?0 |
430 | | .update |
431 | 16 | .err_tip(|| "Expected update to exist in UpdateForWorker")?0 |
432 | | { |
433 | | Update::ConnectionResult(_) => { |
434 | 0 | return Err(make_input_err!( |
435 | 0 | "Got ConnectionResult in LocalWorker::run which should never happen" |
436 | 0 | )); |
437 | | } |
438 | | // TODO(palfrey) We should possibly do something with this notification. |
439 | 0 | Update::Disconnect(()) => { |
440 | 0 | self.metrics.disconnects_received.inc(); |
441 | 0 | } |
442 | 0 | Update::KeepAlive(()) => { |
443 | 0 | self.metrics.keep_alives_received.inc(); |
444 | 0 | } |
445 | 1 | Update::KillOperationRequest(kill_operation_request) => { |
446 | 1 | let operation_id = OperationId::from(kill_operation_request.operation_id); |
447 | 1 | if let Err(err0 ) = self.running_actions_manager.kill_operation(&operation_id).await { |
448 | 0 | error!( |
449 | | %operation_id, |
450 | | ?err, |
451 | | "Failed to send kill request for operation" |
452 | | ); |
453 | 1 | } |
454 | | } |
455 | 15 | Update::StartAction(start_execute) => { |
456 | | // Don't accept any new requests if we're shutting down. |
457 | 15 | if shutting_down || (self.config.single_use |
458 | 5 | && self.accepted_action.load(Ordering::Acquire)) { |
459 | 1 | if self.dispatch_ack { |
460 | 0 | self.decline(start_execute.operation_id, Refusal::ShuttingDown).await?; |
461 | 1 | } else if let Some(instance_name) = start_execute.execute_request.map(|request| request.instance_name) { |
462 | 1 | self.grpc_client.clone().execution_response( |
463 | 1 | ExecuteResult{ |
464 | 1 | instance_name, |
465 | 1 | operation_id: start_execute.operation_id, |
466 | 1 | result: Some(execute_result::Result::InternalError(make_err!(Code::ResourceExhausted, "Worker shutting down").into())), |
467 | 1 | resource_usage: None, |
468 | 1 | } |
469 | 1 | ).await?0 ; |
470 | 0 | } |
471 | 1 | continue; |
472 | 14 | } |
473 | | |
474 | | // Admission, then the acknowledgement: what the |
475 | | // scheduler charged on the send is confirmed or |
476 | | // handed back before anything runs. A scheduler |
477 | | // that does not speak the acknowledgement gets |
478 | | // the old behaviour, run whatever arrives. |
479 | 14 | if self.dispatch_ack { |
480 | 4 | if let Some(refusal2 ) = self.admission( |
481 | 4 | &start_execute, |
482 | 4 | actions_in_flight.load(Ordering::Acquire), |
483 | 4 | ) { |
484 | 2 | self.decline(start_execute.operation_id, refusal).await?0 ; |
485 | 1 | continue; |
486 | 2 | } |
487 | 2 | self.grpc_client |
488 | 2 | .clone() |
489 | 2 | .execute_accepted(ExecuteAccepted { |
490 | 2 | operation_id: start_execute.operation_id.clone(), |
491 | 2 | }) |
492 | 2 | .await |
493 | 2 | .err_tip(|| "Could not send ExecuteAccepted")?0 ; |
494 | 10 | } |
495 | | // Admitted: a single-use worker is spent from |
496 | | // here. A decline above must not spend it, or |
497 | | // every dispatch it turns away costs a pod. |
498 | 12 | if self.config.single_use { |
499 | 3 | self.accepted_action.store(true, Ordering::Release); |
500 | 9 | } |
501 | | |
502 | 12 | self.metrics.start_actions_received.inc(); |
503 | | |
504 | 12 | let execute_request = start_execute.execute_request.as_ref(); |
505 | 12 | let operation_id = start_execute.operation_id.clone(); |
506 | 12 | let operation_id_to_log = operation_id.clone(); |
507 | 12 | let maybe_instance_name = execute_request.map(|v| v.instance_name.clone()); |
508 | 12 | let action_digest = execute_request.and_then(|v| v.action_digest.clone()); |
509 | 12 | let digest_hasher = execute_request |
510 | 12 | .ok_or_else(|| make_input_err!0 ("Expected execute_request to be set")) |
511 | 12 | .and_then(|v| DigestHasherFunc::try_from(v.digest_function)) |
512 | 12 | .err_tip(|| "In LocalWorkerImpl::new()")?0 ; |
513 | | |
514 | 12 | let start_action_fut = { |
515 | 12 | let precondition_script_cfg = self.config.experimental_precondition_script.clone(); |
516 | 12 | let precondition_timeout = if self.config.precondition_timeout_ms == 0 { |
517 | 12 | DEFAULT_PRECONDITION_TIMEOUT |
518 | | } else { |
519 | 0 | Duration::from_millis(self.config.precondition_timeout_ms) |
520 | | }; |
521 | 12 | let mut extra_envs: HashMap<String, String> = HashMap::new(); |
522 | 12 | if let Some(ref additional_environment0 ) = self.config.additional_environment { |
523 | 0 | for (name, source) in additional_environment { |
524 | 0 | let value = match source { |
525 | 0 | EnvironmentSource::Property(property) => start_execute |
526 | 0 | .platform.as_ref().and_then(|p|p.properties.iter().find(|pr| &pr.name == property)) |
527 | 0 | .map_or_else(|| Cow::Borrowed(""), |v| Cow::Borrowed(v.value.as_str())), |
528 | 0 | EnvironmentSource::Value(value) => Cow::Borrowed(value.as_str()), |
529 | 0 | EnvironmentSource::FromEnvironment => Cow::Owned(env::var(name).unwrap_or_default()), |
530 | 0 | other => { |
531 | 0 | debug!(?other, "Worker doesn't support this type of additional environment"); |
532 | 0 | continue; |
533 | | } |
534 | | }; |
535 | 0 | extra_envs.insert(name.clone(), value.into_owned()); |
536 | | } |
537 | 12 | } |
538 | 12 | let actions_in_transit_guard = |
539 | 12 | ActionsInTransitGuard::new(self.actions_in_transit.clone()); |
540 | 12 | let worker_id = self.worker_id.clone(); |
541 | 12 | let running_actions_manager = self.running_actions_manager.clone(); |
542 | 12 | let mut grpc_client = self.grpc_client.clone(); |
543 | 12 | let complete = ExecuteComplete { |
544 | 12 | operation_id: operation_id.clone(), |
545 | 12 | }; |
546 | 12 | let single_use = self.config.single_use; |
547 | 12 | self.metrics.clone().wrap(move |metrics| async move {11 |
548 | 11 | metrics.preconditions.wrap(preconditions_met(precondition_script_cfg, &extra_envs, precondition_timeout)) |
549 | 11 | .and_then(|()| running_actions_manager10 .create_and_add_action10 (worker_id10 , start_execute10 )) |
550 | 11 | .map(move |r| {8 |
551 | | // Now that we either failed or registered our action, we can |
552 | | // consider the action to no longer be in transit. |
553 | 8 | drop(actions_in_transit_guard); |
554 | 8 | r |
555 | 8 | }) |
556 | 11 | .and_then(|action| {7 |
557 | 7 | debug!( |
558 | 7 | operation_id = %action.get_operation_id(), |
559 | | "Received request to run action" |
560 | | ); |
561 | 7 | action |
562 | 7 | .clone() |
563 | 7 | .prepare_action() |
564 | 7 | .and_then(RunningAction::execute) |
565 | 7 | .and_then(|result| async move {3 |
566 | | // Reusable workers release their slot during upload. |
567 | | // A single-use worker must never advertise another slot. |
568 | 3 | if !single_use { |
569 | 2 | drop(grpc_client.execution_complete(complete).await); |
570 | 1 | } |
571 | 3 | Ok(result) |
572 | 6 | }) |
573 | 7 | .and_then(RunningAction::upload_results) |
574 | 7 | .and_then(|action| async move {3 |
575 | 3 | let resource_usage = action.resource_usage(); |
576 | 3 | let action_result = action.get_finished_result().await?0 ; |
577 | 3 | Ok(FinishedActionResult { |
578 | 3 | action_result, |
579 | 3 | resource_usage, |
580 | 3 | }) |
581 | 6 | }) |
582 | | // Note: We need ensure we run cleanup even if one of the other steps fail. |
583 | 7 | .then(|result| async move {5 |
584 | 5 | if let Err(e0 ) = action.cleanup().await { |
585 | 0 | return Result::<FinishedActionResult, Error>::Err(e).merge(result); |
586 | 5 | } |
587 | 5 | result |
588 | 10 | }) |
589 | 11 | }7 ).await |
590 | 17 | }) |
591 | | }; |
592 | | |
593 | 12 | let make_publish_future = { |
594 | 12 | let mut grpc_client = self.grpc_client.clone(); |
595 | | |
596 | 12 | let running_actions_manager = self.running_actions_manager.clone(); |
597 | 12 | let worker_id = self.worker_id.clone(); |
598 | 6 | move |res: Result<FinishedActionResult, Error>| async move { |
599 | 6 | let instance_name = maybe_instance_name |
600 | 6 | .err_tip(|| "`instance_name` could not be resolved; this is likely an internal error in local_worker.")?0 ; |
601 | 6 | match res { |
602 | 3 | Ok(FinishedActionResult { mut action_result, resource_usage }) => { |
603 | | // Save in the action cache before notifying the scheduler that we've completed. |
604 | 3 | if let Some(digest_info) = action_digest.clone().and_then(|action_digest| action_digest.try_into().ok()) && |
605 | 3 | let Err(err0 ) = running_actions_manager.cache_action_result(digest_info, &mut action_result, digest_hasher).await { |
606 | 0 | error!( |
607 | | ?err, |
608 | | ?action_digest, |
609 | | "Error saving action in store", |
610 | | ); |
611 | 3 | } |
612 | 3 | let action_stage = ActionStage::Completed(action_result); |
613 | 3 | let resource_usage = resource_usage.map(|mut resource_usage| {0 |
614 | 0 | resource_usage.operation_id.clone_from(&operation_id); |
615 | 0 | resource_usage.worker_id.clone_from(&worker_id); |
616 | 0 | resource_usage |
617 | 0 | }); |
618 | 3 | grpc_client.execution_response( |
619 | 3 | ExecuteResult{ |
620 | 3 | instance_name, |
621 | 3 | operation_id, |
622 | 3 | result: Some(execute_result::Result::ExecuteResponse(action_stage.into())), |
623 | 3 | resource_usage, |
624 | 3 | } |
625 | 3 | ) |
626 | 3 | .await |
627 | 1 | .err_tip(|| "Error while calling execution_response")?0 ; |
628 | | }, |
629 | 3 | Err(e) => { |
630 | | // The manager tags a NotFound on anything the client |
631 | | // referenced by digest; no store's wording involved. |
632 | 3 | let is_cas_blob_missing = e.code == Code::NotFound |
633 | 3 | && e.messages.iter()2 .any2 (|m| m == MISSING_INPUT_ERROR_TIP); |
634 | 3 | if is_cas_blob_missing { |
635 | 1 | warn!( |
636 | | ?e, |
637 | | "Missing CAS inputs during prepare_action, returning FAILED_PRECONDITION" |
638 | | ); |
639 | | // The context names the digest; the |
640 | | // execute response turns it into the |
641 | | // PreconditionFailure detail Bazel |
642 | | // re-uploads and retries on. |
643 | 1 | let action_result = ActionResult { |
644 | 1 | error: Some( |
645 | 1 | make_err!( |
646 | 1 | Code::FailedPrecondition, |
647 | 1 | "{}", |
648 | 1 | e.message_string() |
649 | 1 | ) |
650 | 1 | .with_context(e.context.clone()), |
651 | 1 | ), |
652 | 1 | ..ActionResult::default() |
653 | 1 | }; |
654 | 1 | let action_stage = ActionStage::Completed(action_result); |
655 | 1 | grpc_client.execution_response(ExecuteResult{ |
656 | 1 | instance_name, |
657 | 1 | operation_id, |
658 | 1 | result: Some(execute_result::Result::ExecuteResponse(action_stage.into())), |
659 | 1 | resource_usage: None, |
660 | 1 | }).await.err_tip0 (|| "Error calling execution_response with missing inputs")?0 ; |
661 | | } else { |
662 | 2 | grpc_client.execution_response(ExecuteResult{ |
663 | 2 | instance_name, |
664 | 2 | operation_id, |
665 | 2 | result: Some(execute_result::Result::InternalError(e.into())), |
666 | 2 | resource_usage: None, |
667 | 2 | }).await.err_tip0 (|| "Error calling execution_response with error")?0 ; |
668 | | } |
669 | | }, |
670 | | } |
671 | 1 | Ok(()) |
672 | 7 | } |
673 | | }; |
674 | | |
675 | 12 | let add_future_channel = add_future_channel.clone(); |
676 | | |
677 | 12 | info_span!( |
678 | | "worker_start_action_ctx", |
679 | | operation_id = operation_id_to_log, |
680 | 12 | digest_function = %digest_hasher.to_string(), |
681 | 12 | ).in_scope(|| { |
682 | 12 | let _guard = Context::current_with_value(digest_hasher) |
683 | 12 | .attach(); |
684 | | |
685 | 12 | let actions_in_flight = actions_in_flight.clone(); |
686 | 12 | let actions_notify = actions_notify.clone(); |
687 | 12 | let actions_in_flight_fail = actions_in_flight.clone(); |
688 | 12 | let actions_notify_fail = actions_notify.clone(); |
689 | 12 | actions_in_flight.fetch_add(1, Ordering::Release); |
690 | | |
691 | 12 | futures.push( |
692 | 12 | spawn!("worker_start_action", start_action_fut).map(move |res| {6 |
693 | 6 | let res = res.err_tip(|| "Failed to launch spawn")?0 ; |
694 | 6 | if let Err(err3 ) = &res { |
695 | 3 | error!(?err, "Error executing action"); |
696 | 3 | } |
697 | 6 | add_future_channel |
698 | 6 | .send(make_publish_future(res).then(move |res| {1 |
699 | 1 | actions_in_flight.fetch_sub(1, Ordering::Release); |
700 | 1 | actions_notify.notify_one(); |
701 | 1 | core::future::ready(res) |
702 | 6 | }1 ).boxed()) |
703 | 6 | .map_err(|err| |
704 | 0 | Error::from_std_err(Code::Internal, &err).append("LocalWorker could not send future") |
705 | 0 | )?; |
706 | 6 | Ok(()) |
707 | 6 | }) |
708 | 12 | .or_else(move |err| {0 |
709 | | // If the make_publish_future is not run we still need to notify. |
710 | 0 | actions_in_flight_fail.fetch_sub(1, Ordering::Release); |
711 | 0 | actions_notify_fail.notify_one(); |
712 | 0 | core::future::ready(Err(err)) |
713 | 0 | }) |
714 | 12 | .boxed() |
715 | | ); |
716 | 12 | }); |
717 | | } |
718 | | } |
719 | 0 | }, |
720 | 42 | res6 = add_future_rx.next() => { |
721 | 6 | let fut = res.err_tip(|| "New future stream receives should never be closed")?0 ; |
722 | 6 | futures.push(fut); |
723 | | }, |
724 | 42 | res8 = futures.next() => res8 .err_tip8 (|| "Keep-alive should always pending. Likely unable to send data to scheduler")?0 ?1 , |
725 | 42 | complete_msg0 = shutdown_rx.recv().fuse() => { |
726 | 0 | warn!("Worker loop received shutdown signal. Shutting down worker...",); |
727 | 0 | let mut grpc_client = self.grpc_client.clone(); |
728 | 0 | let shutdown_guard = complete_msg.map_err(|e| |
729 | 0 | Error::from_std_err(Code::Internal, &e).append("Failed to receive shutdown message"))?; |
730 | 0 | let actions_in_flight = actions_in_flight.clone(); |
731 | 0 | let actions_notify = actions_notify.clone(); |
732 | 0 | let drain_on_shutdown = self.config.drain_on_shutdown; |
733 | 0 | let drain_deadline = if self.config.max_action_timeout_s == 0 { |
734 | 0 | DEFAULT_MAX_ACTION_TIMEOUT |
735 | | } else { |
736 | 0 | Duration::from_secs(self.config.max_action_timeout_s as u64) |
737 | | }; |
738 | 0 | let shutdown_future = async move { |
739 | 0 | if drain_on_shutdown { |
740 | | // Say so first, so nothing new is dispatched here |
741 | | // while the running actions finish; the scheduler |
742 | | // removes this worker when the stream closes. |
743 | 0 | if let Err(e) = grpc_client.going_away(GoingAwayRequest { drain: true }).await { |
744 | 0 | error!("Failed to send GoingAwayRequest: {e}",); |
745 | 0 | return Err(e); |
746 | 0 | } |
747 | 0 | } |
748 | | // Wait for in-flight operations to be fully completed, |
749 | | // for as long as one action is allowed to run. |
750 | 0 | let wait = async { |
751 | 0 | while actions_in_flight.load(Ordering::Acquire) > 0 { |
752 | 0 | actions_notify.notified().await; |
753 | | } |
754 | 0 | }; |
755 | 0 | if time::timeout(drain_deadline, wait).await.is_err() { |
756 | 0 | error!( |
757 | 0 | actions_in_flight = actions_in_flight.load(Ordering::Acquire), |
758 | | "Drain deadline passed with actions still in flight; shutting down anyway" |
759 | | ); |
760 | 0 | } |
761 | 0 | if !drain_on_shutdown { |
762 | | // Sending this message immediately evicts all jobs from |
763 | | // this worker, of which there should be none. |
764 | 0 | if let Err(e) = grpc_client.going_away(GoingAwayRequest { drain: false }).await { |
765 | 0 | error!("Failed to send GoingAwayRequest: {e}",); |
766 | 0 | return Err(e); |
767 | 0 | } |
768 | 0 | } |
769 | | // Allow shutdown to occur now. |
770 | 0 | drop(shutdown_guard); |
771 | 0 | Ok::<(), Error>(()) |
772 | 0 | }; |
773 | 0 | futures.push(shutdown_future.boxed()); |
774 | 0 | shutting_down = true; |
775 | | }, |
776 | | }; |
777 | | } |
778 | | // Unreachable. |
779 | 7 | } |
780 | | } |
781 | | |
782 | | type ConnectionFactory<T> = Box<dyn Fn() -> BoxFuture<'static, Result<T, Error>> + Send + Sync>; |
783 | | |
784 | | pub struct LocalWorker<T: WorkerApiClientTrait + 'static, U: RunningActionsManager> { |
785 | | config: Arc<LocalWorkerConfig>, |
786 | | running_actions_manager: Arc<U>, |
787 | | connection_factory: ConnectionFactory<T>, |
788 | | sleep_fn: Option<Box<dyn Fn(Duration) -> BoxFuture<'static, ()> + Send + Sync>>, |
789 | | metrics: Arc<Metrics>, |
790 | | registration: Arc<WorkerRegistration>, |
791 | | } |
792 | | |
793 | | /// Whether the worker currently holds a registration with its scheduler. |
794 | | /// A worker that has not registered, or lost its connection and is |
795 | | /// reconnecting, cannot take work; as a health indicator this is what a |
796 | | /// readiness probe reads. |
797 | | #[derive(Debug)] |
798 | | pub struct WorkerRegistration { |
799 | | name: String, |
800 | | registered: AtomicBool, |
801 | | ever_registered: AtomicBool, |
802 | | } |
803 | | |
804 | | impl WorkerRegistration { |
805 | 18 | pub fn new(name: &str) -> Arc<Self> { |
806 | 18 | Arc::new(Self { |
807 | 18 | name: name.to_string(), |
808 | 18 | registered: AtomicBool::new(false), |
809 | 18 | ever_registered: AtomicBool::new(false), |
810 | 18 | }) |
811 | 18 | } |
812 | | |
813 | 6 | pub fn is_registered(&self) -> bool { |
814 | 6 | self.registered.load(Ordering::Acquire) |
815 | 6 | } |
816 | | |
817 | 22 | fn set_registered(&self, registered: bool) { |
818 | 22 | self.registered.store(registered, Ordering::Release); |
819 | 22 | if registered { |
820 | 15 | self.ever_registered.store(true, Ordering::Release); |
821 | 15 | }7 |
822 | 22 | } |
823 | | } |
824 | | |
825 | | #[async_trait] |
826 | | impl HealthStatusIndicator for WorkerRegistration { |
827 | 0 | fn get_name(&self) -> &'static str { |
828 | 0 | "WorkerRegistration" |
829 | 0 | } |
830 | | |
831 | 1 | async fn check_health(&self, _namespace: Cow<'static, str>) -> HealthStatus { |
832 | | if self.is_registered() { |
833 | | HealthStatus::Ok { |
834 | | struct_name: "WorkerRegistration", |
835 | | message: Cow::Owned(format!( |
836 | | "worker '{}' registered with the scheduler", |
837 | | self.name |
838 | | )), |
839 | | } |
840 | | } else if self.ever_registered.load(Ordering::Acquire) { |
841 | | // Lost after it was there: the worker is reconnecting, and until |
842 | | // it does it holds no work. Initializing, not Failed: Failed |
843 | | // turns the plain status check red too, and a liveness probe on |
844 | | // it would restart every worker during a scheduler roll longer |
845 | | // than its threshold, and redden a co-hosted CAS on one worker's |
846 | | // blip. This belongs to readiness alone. |
847 | | HealthStatus::Initializing { |
848 | | struct_name: "WorkerRegistration", |
849 | | message: Cow::Owned(format!( |
850 | | "worker '{}' lost its scheduler connection, reconnecting", |
851 | | self.name |
852 | | )), |
853 | | } |
854 | | } else { |
855 | | HealthStatus::Initializing { |
856 | | struct_name: "WorkerRegistration", |
857 | | message: Cow::Owned(format!( |
858 | | "worker '{}' not yet registered with the scheduler", |
859 | | self.name |
860 | | )), |
861 | | } |
862 | | } |
863 | 1 | } |
864 | | } |
865 | | |
866 | | impl< |
867 | | T: WorkerApiClientTrait + core::fmt::Debug + 'static, |
868 | | U: RunningActionsManager + core::fmt::Debug, |
869 | | > core::fmt::Debug for LocalWorker<T, U> |
870 | | { |
871 | 0 | fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { |
872 | 0 | f.debug_struct("LocalWorker") |
873 | 0 | .field("config", &self.config) |
874 | 0 | .field("running_actions_manager", &self.running_actions_manager) |
875 | 0 | .field("metrics", &self.metrics) |
876 | 0 | .finish_non_exhaustive() |
877 | 0 | } |
878 | | } |
879 | | |
880 | | /// Creates a new `LocalWorker`. The `cas_store` must be an instance of |
881 | | /// `FastSlowStore` and will be checked at runtime. |
882 | 2 | pub async fn new_local_worker( |
883 | 2 | config: Arc<LocalWorkerConfig>, |
884 | 2 | cas_store: Store, |
885 | 2 | ac_store: Option<Store>, |
886 | 2 | historical_store: Store, |
887 | 2 | ) -> Result<LocalWorker<WorkerApiClientWrapper, RunningActionsManagerImpl>, Error> { |
888 | | #[cfg(not(target_os = "linux"))] |
889 | | if config.experimental_buck2_file_capture.is_some() { |
890 | | return Err(make_input_err!( |
891 | | "Buck2 container file capture requires a Linux execution container" |
892 | | )); |
893 | | } |
894 | 2 | let fast_slow_store = cas_store |
895 | 2 | .downcast_ref::<FastSlowStore>(None) |
896 | 2 | .err_tip(|| "Expected store for LocalWorker's store to be a FastSlowStore")?0 |
897 | 2 | .get_arc() |
898 | 2 | .err_tip(|| "FastSlowStore's Arc doesn't exist")?0 ; |
899 | | |
900 | | // Log warning about CAS configuration for multi-worker setups |
901 | 2 | event!( |
902 | 2 | Level::INFO, |
903 | 2 | worker_name = %config.name, |
904 | | "Starting worker '{}'. IMPORTANT: If running multiple workers, all workers \ |
905 | | must share the same CAS storage path to avoid 'Object not found' errors.", |
906 | 2 | config.name |
907 | | ); |
908 | | |
909 | 2 | if let Ok(path1 ) = fs::canonicalize(&config.work_directory).await { |
910 | 1 | fs::remove_dir_all(&path).await.err_tip(|| {0 |
911 | 0 | format!( |
912 | | "Could not remove work_directory '{}' in LocalWorker", |
913 | 0 | path.as_path().to_str().unwrap_or("bad path") |
914 | | ) |
915 | 0 | })?; |
916 | 1 | } |
917 | | |
918 | 2 | fs::create_dir_all(&config.work_directory) |
919 | 2 | .await |
920 | 2 | .err_tip(|| format!0 ("Could not make work_directory : {}", config.work_directory0 ))?0 ; |
921 | 2 | let entrypoint = if config.entrypoint.is_empty() { |
922 | 2 | None |
923 | | } else { |
924 | 0 | Some(config.entrypoint.clone()) |
925 | | }; |
926 | 2 | let max_action_timeout = if config.max_action_timeout_s == 0 { |
927 | 2 | DEFAULT_MAX_ACTION_TIMEOUT |
928 | | } else { |
929 | 0 | Duration::from_secs(config.max_action_timeout_s as u64) |
930 | | }; |
931 | 2 | let max_upload_timeout = if config.max_upload_timeout_s == 0 { |
932 | 2 | DEFAULT_MAX_UPLOAD_TIMEOUT |
933 | | } else { |
934 | 0 | Duration::from_secs(config.max_upload_timeout_s as u64) |
935 | | }; |
936 | 2 | let max_download_timeout = if config.max_download_timeout_s == 0 { |
937 | 2 | DEFAULT_MAX_DOWNLOAD_TIMEOUT |
938 | | } else { |
939 | 0 | Duration::from_secs(config.max_download_timeout_s as u64) |
940 | | }; |
941 | 2 | let max_cleanup_wait = if config.max_cleanup_wait_s == 0 { |
942 | 2 | DEFAULT_MAX_CLEANUP_WAIT |
943 | | } else { |
944 | 0 | Duration::from_secs(config.max_cleanup_wait_s as u64) |
945 | | }; |
946 | 2 | let max_cleanup_backoff = if config.max_cleanup_backoff_ms == 0 { |
947 | 2 | DEFAULT_MAX_CLEANUP_BACKOFF |
948 | | } else { |
949 | 0 | Duration::from_millis(config.max_cleanup_backoff_ms as u64) |
950 | | }; |
951 | | |
952 | | // Initialize directory cache if configured |
953 | 2 | let directory_cache = if let Some(cache_config0 ) = &config.directory_cache { |
954 | | use std::path::PathBuf; |
955 | | |
956 | | use crate::directory_cache::{ |
957 | | DirectoryCache, DirectoryCacheConfig as WorkerDirCacheConfig, |
958 | | }; |
959 | | |
960 | 0 | let cache_root = if cache_config.cache_root.is_empty() { |
961 | 0 | PathBuf::from(&config.work_directory).parent().map_or_else( |
962 | 0 | || PathBuf::from("/tmp/nativelink_directory_cache"), |
963 | 0 | |p| p.join("directory_cache"), |
964 | | ) |
965 | | } else { |
966 | 0 | PathBuf::from(&cache_config.cache_root) |
967 | | }; |
968 | | |
969 | 0 | let worker_cache_config = WorkerDirCacheConfig { |
970 | 0 | max_entries: cache_config.max_entries, |
971 | 0 | max_size_bytes: cache_config.max_size_bytes, |
972 | 0 | cache_root, |
973 | 0 | experimental_subtree_caching: cache_config.experimental_subtree_caching, |
974 | 0 | max_concurrent_fetches: cache_config.max_concurrent_fetches, |
975 | 0 | experimental_get_tree_prefetch: cache_config.experimental_get_tree_prefetch, |
976 | 0 | }; |
977 | | |
978 | 0 | match DirectoryCache::new(worker_cache_config, fast_slow_store.clone()).await { |
979 | 0 | Ok(cache) => { |
980 | 0 | tracing::info!("Directory cache initialized successfully"); |
981 | 0 | Some(Arc::new(cache)) |
982 | | } |
983 | 0 | Err(e) => { |
984 | 0 | tracing::warn!("Failed to initialize directory cache: {:?}", e); |
985 | 0 | None |
986 | | } |
987 | | } |
988 | | } else { |
989 | 2 | None |
990 | | }; |
991 | | |
992 | | #[cfg(target_os = "linux")] |
993 | 2 | let use_namespaces = { |
994 | 2 | let use_mount_namespace = config.use_mount_namespace.unwrap_or_default(); |
995 | | // The mount namespace only exists when `use_namespaces` is on too; |
996 | | // `isolate_tmp` is judged against that, not against the flag alone, |
997 | | // or an explicit `true` would be dropped without a word. |
998 | 2 | let mount_namespace_on = config.use_namespaces == Some(true) && use_mount_namespace0 ; |
999 | | // A private /tmp is part of the mount isolation unless explicitly |
1000 | | // turned off. A worker without a /tmp, such as a minimal container |
1001 | | // image, has nothing for actions to collide on, so the default is |
1002 | | // off there instead of failing to start. |
1003 | 2 | let isolate_tmp = config.isolate_tmp.unwrap_or_else(|| { |
1004 | 2 | let has_tmp = std::path::Path::new("/tmp").is_dir(); |
1005 | 2 | if mount_namespace_on && !has_tmp0 { |
1006 | 0 | warn!("/tmp does not exist on this worker, so actions will not get a private /tmp"); |
1007 | 2 | } |
1008 | 2 | mount_namespace_on && has_tmp0 |
1009 | 2 | }); |
1010 | 2 | if isolate_tmp && !mount_namespace_on0 { |
1011 | 0 | return Err(make_err!( |
1012 | 0 | Code::InvalidArgument, |
1013 | 0 | "isolate_tmp requires use_namespaces and use_mount_namespace to be true" |
1014 | 0 | )); |
1015 | 2 | } |
1016 | 2 | if let Some(use_namespaces0 ) = &config.use_namespaces { |
1017 | 0 | if *use_namespaces |
1018 | 0 | && !crate::namespace_utils::namespaces_supported(use_mount_namespace, isolate_tmp) |
1019 | | { |
1020 | 0 | return Err(make_err!(Code::Unavailable, "Namespaces not supported")); |
1021 | 0 | } |
1022 | 0 | if !*use_namespaces { |
1023 | 0 | crate::running_actions_manager::UseNamespaces::No |
1024 | 0 | } else if use_mount_namespace { |
1025 | 0 | crate::running_actions_manager::UseNamespaces::YesAndMount { isolate_tmp } |
1026 | | } else { |
1027 | 0 | crate::running_actions_manager::UseNamespaces::Yes |
1028 | | } |
1029 | 2 | } else if use_mount_namespace { |
1030 | 0 | return Err(make_err!( |
1031 | 0 | Code::Unavailable, |
1032 | 0 | "Mount namespaces not supported" |
1033 | 0 | )); |
1034 | | } else { |
1035 | 2 | crate::running_actions_manager::UseNamespaces::No |
1036 | | } |
1037 | | }; |
1038 | | |
1039 | | #[cfg(not(target_os = "linux"))] |
1040 | | if config.use_namespaces.is_some_and(core::convert::identity) { |
1041 | | return Err(make_err!( |
1042 | | Code::Unavailable, |
1043 | | "Namespaces not supported on non-Linux OSes" |
1044 | | )); |
1045 | | } |
1046 | | #[cfg(not(target_os = "linux"))] |
1047 | | if config |
1048 | | .use_mount_namespace |
1049 | | .is_some_and(core::convert::identity) |
1050 | | { |
1051 | | return Err(make_err!( |
1052 | | Code::Unavailable, |
1053 | | "Mount namespaces not supported on non-Linux OSes" |
1054 | | )); |
1055 | | } |
1056 | | #[cfg(not(target_os = "linux"))] |
1057 | | if config.isolate_tmp.is_some_and(core::convert::identity) { |
1058 | | return Err(make_err!( |
1059 | | Code::Unavailable, |
1060 | | "isolate_tmp is not supported on non-Linux OSes" |
1061 | | )); |
1062 | | } |
1063 | | |
1064 | | // A pooled worker process gets the same PID, user, UTS and IPC |
1065 | | // namespaces a one-shot action gets, never the mount namespace. |
1066 | | #[cfg(target_os = "linux")] |
1067 | 2 | let persistent_workers_namespaced = !matches!0 ( |
1068 | 2 | use_namespaces, |
1069 | | crate::running_actions_manager::UseNamespaces::No |
1070 | | ); |
1071 | | #[cfg(not(target_os = "linux"))] |
1072 | | let persistent_workers_namespaced = false; |
1073 | | |
1074 | 2 | let running_actions_manager = |
1075 | 2 | Arc::new(RunningActionsManagerImpl::new(RunningActionsManagerArgs { |
1076 | 2 | root_action_directory: config.work_directory.clone(), |
1077 | | execution_configuration: ExecutionConfiguration { |
1078 | 2 | max_captured_output_bytes: config.max_captured_output_bytes, |
1079 | 2 | kill_grace: if config.kill_grace_ms == 0 { |
1080 | 2 | DEFAULT_KILL_GRACE |
1081 | | } else { |
1082 | 0 | Duration::from_millis(config.kill_grace_ms) |
1083 | | }, |
1084 | 2 | set_tmpdir: config.set_tmpdir, |
1085 | 2 | persistent_workers: persistent_worker_settings( |
1086 | 2 | config.persistent_workers.as_ref(), |
1087 | 2 | persistent_workers_namespaced, |
1088 | | ), |
1089 | 2 | buck2_file_capture: config.experimental_buck2_file_capture.clone(), |
1090 | 2 | entrypoint, |
1091 | 2 | additional_environment: config.additional_environment.clone(), |
1092 | 2 | resource_enforcement: config |
1093 | 2 | .resource_enforcement |
1094 | 2 | .as_ref() |
1095 | 2 | .and_then(ResourceEnforcement::from_config), |
1096 | | }, |
1097 | 2 | cas_store: fast_slow_store, |
1098 | 2 | ac_store, |
1099 | 2 | historical_store, |
1100 | 2 | upload_action_result_config: &config.upload_action_result, |
1101 | 2 | max_action_timeout, |
1102 | 2 | max_upload_timeout, |
1103 | 2 | max_download_timeout, |
1104 | 2 | max_cleanup_wait, |
1105 | 2 | max_cleanup_backoff, |
1106 | 2 | timeout_handled_externally: config.timeout_handled_externally, |
1107 | 2 | directory_cache, |
1108 | 2 | active_input_leases: config.experimental_active_input_leases, |
1109 | | #[cfg(target_os = "linux")] |
1110 | 2 | use_namespaces, |
1111 | 0 | })?); |
1112 | | // What actions leave behind is reaped here; see `crate::reaper`. |
1113 | 2 | if let Err(err0 ) = crate::reaper::become_subreaper() { |
1114 | 0 | warn!( |
1115 | | ?err, |
1116 | | "Could not become the subreaper; orphans of actions will not be reaped here" |
1117 | | ); |
1118 | 2 | } |
1119 | 2 | drop(background_spawn!("orphan_reaper", async move {1 |
1120 | 1 | let mut seen_last = std::collections::BTreeSet::new(); |
1121 | | loop { |
1122 | 1 | time::sleep(crate::reaper::REAP_INTERVAL).await; |
1123 | | // The sweep reads /proc for every process; off the runtime. |
1124 | 0 | let (reaped, seen) = spawn_blocking!("orphan_reaper_sweep", move || { |
1125 | 0 | crate::reaper::reap_orphaned_zombies(&seen_last) |
1126 | 0 | }) |
1127 | 0 | .await |
1128 | 0 | .unwrap_or_default(); |
1129 | 0 | seen_last = seen; |
1130 | 0 | if reaped > 0 { |
1131 | 0 | info!(reaped, "Reaped zombies that actions left behind"); |
1132 | 0 | } |
1133 | | } |
1134 | | })); |
1135 | 2 | if config.orphan_sweep_interval_s > 0 { |
1136 | 0 | let interval = Duration::from_secs(config.orphan_sweep_interval_s); |
1137 | 0 | let manager = running_actions_manager.clone(); |
1138 | | // Detached on purpose: the guarded spawn aborts its task when the |
1139 | | // handle drops, and this one runs for the worker's whole life. |
1140 | 0 | drop(background_spawn!("orphan_sweep", async move { |
1141 | 0 | info!(interval_s = interval.as_secs(), "Orphan sweep scheduled"); |
1142 | | loop { |
1143 | 0 | time::sleep(interval).await; |
1144 | 0 | match manager.sweep_orphaned_directories().await { |
1145 | 0 | Ok(0) => debug!("Orphan sweep found nothing"), |
1146 | 0 | Ok(removed) => info!(removed, "Orphan sweep finished"), |
1147 | 0 | Err(err) => warn!(?err, "Orphan sweep failed"), |
1148 | | } |
1149 | | } |
1150 | | })); |
1151 | 2 | } |
1152 | 2 | let local_worker = LocalWorker::new_with_connection_factory_and_actions_manager( |
1153 | 2 | config.clone(), |
1154 | 2 | running_actions_manager, |
1155 | 2 | Box::new(move || {0 |
1156 | 0 | let config = config.clone(); |
1157 | 0 | Box::pin(async move { |
1158 | 0 | let timeout = config |
1159 | 0 | .worker_api_endpoint |
1160 | 0 | .timeout |
1161 | 0 | .unwrap_or(DEFAULT_ENDPOINT_TIMEOUT_S); |
1162 | 0 | let timeout_duration = Duration::from_secs_f32(timeout); |
1163 | 0 | let tls_config = |
1164 | 0 | tls_utils::load_client_config(&config.worker_api_endpoint.tls_config) |
1165 | 0 | .err_tip(|| "Parsing local worker TLS configuration")?; |
1166 | 0 | let endpoint = |
1167 | 0 | tls_utils::endpoint_from(&config.worker_api_endpoint.uri, tls_config) |
1168 | 0 | .map_err(|e| { |
1169 | 0 | Error::from_std_err(Code::InvalidArgument, &e) |
1170 | 0 | .append("Invalid URI for worker endpoint") |
1171 | 0 | })? |
1172 | 0 | .connect_timeout(timeout_duration) |
1173 | 0 | .timeout(timeout_duration); |
1174 | | |
1175 | 0 | let transport = endpoint.connect().await.map_err(|e| { |
1176 | 0 | Error::from_std_err(Code::Internal, &e).append(format!( |
1177 | | "Could not connect to endpoint {}", |
1178 | 0 | config.worker_api_endpoint.uri |
1179 | | )) |
1180 | 0 | })?; |
1181 | 0 | Ok(WorkerApiClient::new(transport).into()) |
1182 | 0 | }) |
1183 | 0 | }), |
1184 | 2 | Box::new(move |d| Box::pin0 (time::sleep0 (d0 ))), |
1185 | | ); |
1186 | 2 | Ok(local_worker) |
1187 | 2 | } |
1188 | | |
1189 | | impl<T: WorkerApiClientTrait + 'static, U: RunningActionsManager> LocalWorker<T, U> { |
1190 | 18 | pub fn new_with_connection_factory_and_actions_manager( |
1191 | 18 | config: Arc<LocalWorkerConfig>, |
1192 | 18 | running_actions_manager: Arc<U>, |
1193 | 18 | connection_factory: ConnectionFactory<T>, |
1194 | 18 | sleep_fn: Box<dyn Fn(Duration) -> BoxFuture<'static, ()> + Send + Sync>, |
1195 | 18 | ) -> Self { |
1196 | 18 | let metrics = Arc::new(Metrics::new(Arc::downgrade( |
1197 | 18 | running_actions_manager.metrics(), |
1198 | | ))); |
1199 | 18 | let registration = WorkerRegistration::new(&config.name); |
1200 | 18 | Self { |
1201 | 18 | config, |
1202 | 18 | running_actions_manager, |
1203 | 18 | connection_factory, |
1204 | 18 | sleep_fn: Some(sleep_fn), |
1205 | 18 | metrics, |
1206 | 18 | registration, |
1207 | 18 | } |
1208 | 18 | } |
1209 | | |
1210 | | /// The registration flag this worker flips, for a health registry. |
1211 | 16 | pub fn registration(&self) -> Arc<WorkerRegistration> { |
1212 | 16 | self.registration.clone() |
1213 | 16 | } |
1214 | | |
1215 | | /// Flip a flag that was registered with a health registry before the |
1216 | | /// worker existed, as the server binary has to. |
1217 | | #[must_use] |
1218 | 0 | pub fn with_registration(mut self, registration: Arc<WorkerRegistration>) -> Self { |
1219 | 0 | self.registration = registration; |
1220 | 0 | self |
1221 | 0 | } |
1222 | | |
1223 | | #[allow( |
1224 | | clippy::missing_const_for_fn, |
1225 | | reason = "False positive on stable, but not on nightly" |
1226 | | )] |
1227 | 0 | pub fn name(&self) -> &String { |
1228 | 0 | &self.config.name |
1229 | 0 | } |
1230 | | |
1231 | 22 | async fn register_worker( |
1232 | 22 | &self, |
1233 | 22 | client: &mut T, |
1234 | 22 | ) -> Result<(String, bool, String, Streaming<UpdateForWorker>), Error> { |
1235 | 22 | let mut extra_envs: HashMap<String, String> = HashMap::new(); |
1236 | 22 | if let Some(ref additional_environment0 ) = self.config.additional_environment { |
1237 | 0 | for (name, source) in additional_environment { |
1238 | 0 | let value = match source { |
1239 | 0 | EnvironmentSource::Value(value) => Cow::Borrowed(value.as_str()), |
1240 | | EnvironmentSource::FromEnvironment => { |
1241 | 0 | Cow::Owned(env::var(name).unwrap_or_default()) |
1242 | | } |
1243 | 0 | other => { |
1244 | 0 | debug!( |
1245 | | ?other, |
1246 | | "Worker registration doesn't support this type of additional environment" |
1247 | | ); |
1248 | 0 | continue; |
1249 | | } |
1250 | | }; |
1251 | 0 | extra_envs.insert(name.clone(), value.into_owned()); |
1252 | | } |
1253 | 22 | } |
1254 | | |
1255 | 22 | let mut platform_properties = self.config.platform_properties.clone(); |
1256 | 22 | if let Some(capacity0 ) = &self.config.capacity { |
1257 | 0 | let memory_headroom_percent = crate::capacity::memory_headroom_percent( |
1258 | 0 | capacity, |
1259 | 0 | self.config.resource_enforcement.as_ref(), |
1260 | | ); |
1261 | 0 | crate::capacity::apply(capacity, memory_headroom_percent, &mut platform_properties) |
1262 | 0 | .err_tip(|| "Advertising capacity from the cgroup")?; |
1263 | 22 | } |
1264 | 22 | let connect_worker_request = make_connect_worker_request( |
1265 | 22 | self.config.name.clone(), |
1266 | 22 | &platform_properties, |
1267 | 22 | &extra_envs, |
1268 | 22 | if self.config.single_use { |
1269 | 3 | 1 |
1270 | | } else { |
1271 | 19 | self.config.max_inflight_tasks |
1272 | | }, |
1273 | | ) |
1274 | 22 | .await?0 ; |
1275 | 22 | let mut update_for_worker_stream16 = client |
1276 | 22 | .connect_worker(connect_worker_request) |
1277 | 22 | .await |
1278 | 16 | .err_tip(|| "Could not call connect_worker() in worker")?0 |
1279 | 16 | .into_inner(); |
1280 | | |
1281 | 16 | let first_msg_update15 = update_for_worker_stream |
1282 | 16 | .next() |
1283 | 16 | .await |
1284 | 16 | .err_tip(|| "Got EOF expected UpdateForWorker")?1 |
1285 | 15 | .err_tip(|| "Got error when receiving UpdateForWorker")?0 |
1286 | | .update; |
1287 | | |
1288 | 15 | let (worker_id, dispatch_ack, memory_property) = match first_msg_update { |
1289 | 15 | Some(Update::ConnectionResult(connection_result)) => ( |
1290 | 15 | connection_result.worker_id, |
1291 | 15 | connection_result.dispatch_ack, |
1292 | 15 | connection_result.memory_property, |
1293 | 15 | ), |
1294 | 0 | other => { |
1295 | 0 | return Err(make_input_err!( |
1296 | 0 | "Expected first response from scheduler to be a ConnectionResult got : {:?}", |
1297 | 0 | other |
1298 | 0 | )); |
1299 | | } |
1300 | | }; |
1301 | 15 | Ok(( |
1302 | 15 | worker_id, |
1303 | 15 | dispatch_ack, |
1304 | 15 | memory_property, |
1305 | 15 | update_for_worker_stream, |
1306 | 15 | )) |
1307 | 16 | } |
1308 | | |
1309 | | #[instrument(skip(self), level = Level::INFO)] |
1310 | 16 | pub async fn run( |
1311 | 16 | mut self, |
1312 | 16 | mut shutdown_rx: broadcast::Receiver<ShutdownGuard>, |
1313 | 16 | ) -> Result<(), Error> { |
1314 | | // Belt-and-suspenders QoS bump: the main binary already calls |
1315 | | // this before runtime creation so the tokio worker threads |
1316 | | // inherit P-core preference via pthread QoS inheritance, but |
1317 | | // any thread that reaches this point should also be tagged in |
1318 | | // case it was spawned by a path that bypassed `on_thread_start`. |
1319 | | // No-op on non-macOS. |
1320 | | let _ = crate::qos::set_user_initiated(); |
1321 | | |
1322 | | let sleep_fn = self |
1323 | | .sleep_fn |
1324 | | .take() |
1325 | | .err_tip(|| "Could not unwrap sleep_fn in LocalWorker::run")?; |
1326 | | let sleep_fn_pin = Pin::new(&sleep_fn); |
1327 | | let attempts = AtomicU32::new(0); |
1328 | | let attempts_ref = &attempts; |
1329 | 6 | let error_handler = Box::pin(move |err| async move { |
1330 | 6 | let attempt = attempts_ref.fetch_add(1, Ordering::AcqRel); |
1331 | 6 | let delay = reconnect_delay(attempt); |
1332 | 6 | error!(?err, attempt, delay_ms = delay.as_millis(), "Error"); |
1333 | 6 | (sleep_fn_pin)(delay).await; |
1334 | 12 | }); |
1335 | | |
1336 | | loop { |
1337 | | // First connect to our endpoint. |
1338 | | let mut client = match (self.connection_factory)().await { |
1339 | | Ok(client) => client, |
1340 | | Err(e) => { |
1341 | | (error_handler)(e).await; |
1342 | | continue; // Try to connect again. |
1343 | | } |
1344 | | }; |
1345 | | |
1346 | | debug!("Connected to endpoint"); |
1347 | | |
1348 | | // Next register our worker with the scheduler. |
1349 | | let (inner, update_for_worker_stream) = match self.register_worker(&mut client).await { |
1350 | | Err(e) => { |
1351 | | (error_handler)(e).await; |
1352 | | continue; // Try to connect again. |
1353 | | } |
1354 | | Ok((worker_id, dispatch_ack, memory_property, update_for_worker_stream)) => ( |
1355 | | LocalWorkerImpl::new( |
1356 | | &self.config, |
1357 | | client, |
1358 | | worker_id, |
1359 | | dispatch_ack, |
1360 | | memory_property, |
1361 | | self.running_actions_manager.clone(), |
1362 | | self.metrics.clone(), |
1363 | | ), |
1364 | | update_for_worker_stream, |
1365 | | ), |
1366 | | }; |
1367 | | info!( |
1368 | | worker_id = %inner.worker_id, |
1369 | | "Worker registered with scheduler" |
1370 | | ); |
1371 | | attempts.store(0, Ordering::Release); |
1372 | | self.registration.set_registered(true); |
1373 | | |
1374 | | // Now listen for connections and run all other services. |
1375 | | let run_result = inner.run(update_for_worker_stream, &mut shutdown_rx).await; |
1376 | | self.registration.set_registered(false); |
1377 | | if let Err(err) = run_result { |
1378 | | // Give in-transit actions a chance to settle before we kill |
1379 | | // them, so their results still reach the scheduler. |
1380 | | const ITERATIONS: usize = 1_000; |
1381 | | |
1382 | | let sleep_duration = ACTIONS_IN_TRANSIT_TIMEOUT_S / ITERATIONS as f32; |
1383 | | let mut drained = false; |
1384 | | for _ in 0..ITERATIONS { |
1385 | | if inner.actions_in_transit.load(Ordering::Acquire) == 0 { |
1386 | | drained = true; |
1387 | | break; |
1388 | | } |
1389 | | (sleep_fn_pin)(Duration::from_secs_f32(sleep_duration)).await; |
1390 | | } |
1391 | | if !drained { |
1392 | | // Deliberately not fatal. Returning here propagates out of |
1393 | | // the worker's main loop and aborts the process, so a |
1394 | | // scheduler blip that happened to catch an action in |
1395 | | // transit took the whole worker down — and every action it |
1396 | | // held then had to run again elsewhere. At fleet scale that |
1397 | | // is a restart storm. kill_all() below discards these |
1398 | | // actions anyway, so the wait is a courtesy and overrunning |
1399 | | // it costs nothing beyond the actions we were already |
1400 | | // giving up on. |
1401 | | error!( |
1402 | | actions_in_transit = inner.actions_in_transit.load(Ordering::Acquire), |
1403 | | "Actions in transit did not reach zero before we disconnected from the scheduler" |
1404 | | ); |
1405 | | } |
1406 | | |
1407 | | // Kill off any existing actions because if we re-connect, we'll |
1408 | | // get some more and it might resource lock us. |
1409 | | self.running_actions_manager.kill_all().await; |
1410 | | |
1411 | | if self.config.single_use && inner.accepted_action.load(Ordering::Acquire) { |
1412 | | return Err( |
1413 | | err.append("Single-use worker disconnected after accepting its action") |
1414 | | ); |
1415 | | } |
1416 | | error!(?err, "Worker disconnected from scheduler, reconnecting"); |
1417 | | (error_handler)(err).await; // Try to connect again. |
1418 | | } else if self.config.single_use { |
1419 | | return Ok(()); |
1420 | | } |
1421 | | } |
1422 | | // Unreachable. |
1423 | 2 | } |
1424 | | } |
1425 | | |
1426 | | #[derive(Debug, MetricsComponent)] |
1427 | | pub struct Metrics { |
1428 | | #[metric( |
1429 | | help = "Total number of actions sent to this worker to process. This does not mean it started them, it just means it received a request to execute it." |
1430 | | )] |
1431 | | start_actions_received: CounterWithTime, |
1432 | | #[metric(help = "Total number of disconnects received from the scheduler.")] |
1433 | | disconnects_received: CounterWithTime, |
1434 | | #[metric( |
1435 | | help = "Dispatches this worker declined: at capacity, short of memory, or shutting down." |
1436 | | )] |
1437 | | actions_declined: CounterWithTime, |
1438 | | #[metric(help = "Total number of keep-alives received from the scheduler.")] |
1439 | | keep_alives_received: CounterWithTime, |
1440 | | #[metric( |
1441 | | help = "Stats about the calls to check if an action satisfies the config supplied script." |
1442 | | )] |
1443 | | preconditions: AsyncCounterWrapper, |
1444 | | #[metric] |
1445 | | #[allow( |
1446 | | clippy::struct_field_names, |
1447 | | reason = "TODO Fix this. Triggers on nightly" |
1448 | | )] |
1449 | | running_actions_manager_metrics: Weak<RunningActionManagerMetrics>, |
1450 | | } |
1451 | | |
1452 | | impl RootMetricsComponent for Metrics {} |
1453 | | |
1454 | | impl Metrics { |
1455 | 18 | fn new(running_actions_manager_metrics: Weak<RunningActionManagerMetrics>) -> Self { |
1456 | 18 | Self { |
1457 | 18 | start_actions_received: CounterWithTime::default(), |
1458 | 18 | disconnects_received: CounterWithTime::default(), |
1459 | 18 | actions_declined: CounterWithTime::default(), |
1460 | 18 | keep_alives_received: CounterWithTime::default(), |
1461 | 18 | preconditions: AsyncCounterWrapper::default(), |
1462 | 18 | running_actions_manager_metrics, |
1463 | 18 | } |
1464 | 18 | } |
1465 | | } |
1466 | | |
1467 | | impl Metrics { |
1468 | 12 | async fn wrap<U, T: Future<Output = U>, F: FnOnce(Arc<Self>) -> T>( |
1469 | 12 | self: Arc<Self>, |
1470 | 12 | fut: F, |
1471 | 12 | ) -> U {11 |
1472 | 11 | fut(self).await |
1473 | 6 | } |
1474 | | } |
1475 | | |
1476 | | /// The pool settings from the worker config, defaults for what it leaves at |
1477 | | /// zero or unset. |
1478 | 2 | fn persistent_worker_settings( |
1479 | 2 | config: Option<&nativelink_config::cas_server::PersistentWorkersConfig>, |
1480 | 2 | namespaced: bool, |
1481 | 2 | ) -> crate::running_actions_manager::PersistentWorkersSettings { |
1482 | | use crate::persistent_worker::PoolConfig; |
1483 | 2 | let defaults = PoolConfig::default(); |
1484 | 2 | let pool = config.map_or(defaults, |c| {0 |
1485 | 0 | let or_default = |value: u64, default: u64| if value == 0 { default } else { value }; |
1486 | | PoolConfig { |
1487 | 0 | max_workers_per_key: if c.max_workers_per_key == 0 { |
1488 | 0 | defaults.max_workers_per_key |
1489 | | } else { |
1490 | 0 | c.max_workers_per_key |
1491 | | }, |
1492 | 0 | idle_timeout: Duration::from_secs(or_default( |
1493 | 0 | c.idle_timeout_s, |
1494 | 0 | defaults.idle_timeout.as_secs(), |
1495 | 0 | )), |
1496 | 0 | max_requests_per_worker: or_default( |
1497 | 0 | c.max_requests_per_worker, |
1498 | 0 | defaults.max_requests_per_worker, |
1499 | 0 | ), |
1500 | 0 | shutdown_grace: Duration::from_millis(or_default( |
1501 | 0 | c.shutdown_grace_ms, |
1502 | 0 | u64::try_from(defaults.shutdown_grace.as_millis()).unwrap_or(5000), |
1503 | 0 | )), |
1504 | 0 | acquire_timeout: Duration::from_secs(or_default( |
1505 | 0 | c.acquire_timeout_s, |
1506 | 0 | defaults.acquire_timeout.as_secs(), |
1507 | 0 | )), |
1508 | 0 | namespaced, |
1509 | | } |
1510 | 0 | }); |
1511 | | crate::running_actions_manager::PersistentWorkersSettings { |
1512 | 2 | enabled: config.is_none_or(|c| c.enabled), |
1513 | 2 | pool: PoolConfig { namespaced, ..pool }, |
1514 | | } |
1515 | 2 | } |