Coverage Report

Created: 2026-10-01 05:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}