Coverage Report

Created: 2026-10-06 13:02

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
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
}