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