Coverage Report

Created: 2026-10-08 11:27

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/build/source/nativelink-scheduler/src/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::{Hash, Hasher};
16
use std::collections::HashMap;
17
use std::sync::Arc;
18
use std::time::{SystemTime, UNIX_EPOCH};
19
20
use nativelink_error::{Code, Error, ResultExt};
21
use nativelink_metric::MetricsComponent;
22
use nativelink_proto::com::github::trace_machina::nativelink::remote_execution::{
23
    ActionResourceUsage, ConnectionResult, KillOperationRequest, StartExecute, UpdateForWorker,
24
    WorkerLoad, update_for_worker,
25
};
26
use nativelink_util::action_messages::{ActionInfo, OperationId, WorkerId};
27
use nativelink_util::metrics_utils::{AsyncCounterWrapper, CounterWithTime, FuncCounterWrapper};
28
use nativelink_util::origin_event::OriginMetadata;
29
use nativelink_util::platform_properties::{PlatformProperties, PlatformPropertyValue};
30
use tokio::sync::mpsc::Sender;
31
use tokio::sync::mpsc::error::TrySendError;
32
33
pub type WorkerTimestamp = u64;
34
35
/// Represents the action info and the platform properties of the action.
36
/// These platform properties have the type of the properties as well as
37
/// the value of the properties, unlike `ActionInfo`, which only has the
38
/// string value of the properties.
39
#[derive(Clone, Debug, MetricsComponent)]
40
pub struct ActionInfoWithProps {
41
    /// The action info of the action.
42
    #[metric(group = "action_info")]
43
    pub inner: Arc<ActionInfo>,
44
    /// The platform properties of the action.
45
    #[metric(group = "platform_properties")]
46
    pub platform_properties: PlatformProperties,
47
    /// Origin metadata used when publishing scheduler-side telemetry for this action.
48
    pub origin_metadata: OriginMetadata,
49
    /// `OriginEvent` id for the `scheduler_start_execute` request.
50
    pub scheduler_start_execute_event_id: Option<String>,
51
}
52
53
/// Notifications to send worker about a requested state change.
54
#[derive(Debug)]
55
pub enum WorkerUpdate {
56
    /// Requests that the worker begin executing this action, dispatched at
57
    /// this scheduler time.
58
    RunAction(Box<(OperationId, ActionInfoWithProps, WorkerTimestamp)>),
59
60
    /// Request that the worker is no longer in the pool and may discard any jobs.
61
    Disconnect,
62
63
    /// Requests that the worker stop executing this operation.
64
    KillOperation(OperationId),
65
}
66
67
#[derive(Debug, MetricsComponent)]
68
pub struct PendingActionInfoData {
69
    #[metric]
70
    pub action_info: ActionInfoWithProps,
71
    /// When the worker was told to kill this operation, recorded as the
72
    /// worker's last-seen timestamp at send time; keepalives keep that
73
    /// within the worker timeout of now, so a deadline computed from it
74
    /// fires at most one worker timeout late, never early. `None` until a
75
    /// kill is sent. The operation's later report then only settles the
76
    /// worker's own bookkeeping.
77
    #[metric(help = "When the worker was asked to kill this operation.")]
78
    pub kill_requested_at: Option<WorkerTimestamp>,
79
    /// The worker's last-seen timestamp when the dispatch was sent.
80
    #[metric(help = "When this operation was dispatched to the worker.")]
81
    pub dispatched_at: WorkerTimestamp,
82
    /// The worker said it took the action. A worker that does not speak
83
    /// the acknowledgement never sets this, which is why the
84
    /// unacknowledged sweep is opt-in.
85
    #[metric(help = "Whether the worker acknowledged the dispatch.")]
86
    pub accepted: bool,
87
    /// The last resource usage the worker reported for this operation; it
88
    /// arrives just before the result and says whether the result is a kill.
89
    pub last_usage: Option<ActionResourceUsage>,
90
}
91
92
/// Represents a connection to a worker and used as the medium to
93
/// interact with the worker from the client/scheduler.
94
#[derive(Debug, MetricsComponent)]
95
pub struct Worker {
96
    /// Unique identifier of the worker.
97
    #[metric(help = "The unique identifier of the worker.")]
98
    pub id: WorkerId,
99
100
    /// Properties that describe the capabilities of this worker.
101
    #[metric(group = "platform_properties")]
102
    pub platform_properties: PlatformProperties,
103
104
    /// The properties the worker registered with. `platform_properties` is
105
    /// reduced while actions run; this is what the worker offers when idle.
106
    pub total_platform_properties: PlatformProperties,
107
108
    /// Channel to send commands from scheduler to worker.
109
    pub tx: Sender<UpdateForWorker>,
110
111
    /// The action info of the running actions on the worker.
112
    #[metric(group = "running_action_infos")]
113
    pub running_action_infos: HashMap<OperationId, PendingActionInfoData>,
114
115
    /// Timestamp of last time this worker had been communicated with.
116
    // Warning: Do not update this timestamp without updating the placement of the worker in
117
    // the LRUCache in the Workers struct.
118
    #[metric(help = "Last time this worker was communicated with.")]
119
    pub last_update_timestamp: WorkerTimestamp,
120
121
    /// Whether the worker rejected the last action due to back pressure.
122
    #[metric(help = "If the worker is paused.")]
123
    pub is_paused: bool,
124
125
    /// Set when the pause came from a decline for load: it lifts only on a
126
    /// keepalive whose free memory covers this much, not on any keepalive.
127
    #[metric(help = "Free memory in KiB the worker must report before it is unpaused.")]
128
    pub pause_needs_kb: Option<u64>,
129
130
    /// The worker said on connection that it admits any action while it
131
    /// holds nothing else, so a decline for load from it comes from a busy
132
    /// worker even when this scheduler thinks it idle.
133
    pub admits_when_idle: bool,
134
135
    /// The smallest reservation this worker declined for load while it held
136
    /// nothing else, from a worker that does not admit when idle (an older
137
    /// one, or one with no cgroup limit). It is not offered that much or
138
    /// more until a keepalive reports that much free.
139
    #[metric(help = "Smallest reservation in KiB the worker declined while idle.")]
140
    pub idle_declined_kb: Option<u64>,
141
142
    /// Whether the worker is draining.
143
    #[metric(help = "If the worker is draining.")]
144
    pub is_draining: bool,
145
146
    /// Maximum inflight tasks for this worker (or 0 for unlimited)
147
    #[metric(help = "Maximum inflight tasks for this worker (or 0 for unlimited)")]
148
    pub max_inflight_tasks: u64,
149
150
    /// What the worker last reported having to spare, from its keepalive;
151
    /// `None` until it reports, and for workers that never do.
152
    pub last_load: Option<WorkerLoad>,
153
154
    /// Stats about the worker.
155
    #[metric]
156
    metrics: Arc<Metrics>,
157
}
158
159
/// Messages the channel holds beyond the worker's own concurrency: kills,
160
/// keepalives and the disconnect, which arrive alongside the dispatches.
161
const CHANNEL_HEADROOM: usize = 16;
162
/// A worker that set no `max_inflight_tasks` gets this much queue.
163
const DEFAULT_CHANNEL_DISPATCHES: usize = 256;
164
165
/// How many messages the worker's channel holds: what its concurrency
166
/// allows in flight plus the headroom. Bounded so that a worker that has
167
/// stopped reading cannot be handed work without limit; a full channel is
168
/// `ResourceExhausted` on the send, and the caller requeues.
169
0
pub fn channel_capacity(max_inflight_tasks: u64) -> usize {
170
0
    let ceiling = DEFAULT_CHANNEL_DISPATCHES * 4;
171
0
    let dispatches = match usize::try_from(max_inflight_tasks) {
172
0
        Ok(0) => DEFAULT_CHANNEL_DISPATCHES,
173
0
        Ok(n) if n <= ceiling => n,
174
0
        _ => ceiling,
175
    };
176
0
    dispatches + CHANNEL_HEADROOM
177
0
}
178
179
324
fn send_msg_to_worker(
180
324
    tx: &Sender<UpdateForWorker>,
181
324
    msg: update_for_worker::Update,
182
324
) -> Result<(), Error> {
183
324
    match tx.try_send(UpdateForWorker { update: Some(msg) }) {
184
314
        Ok(()) => Ok(()),
185
2
        Err(TrySendError::Full(_)) => Err(Error::new(
186
2
            Code::ResourceExhausted,
187
2
            "Worker channel full".to_string(),
188
2
        )),
189
8
        Err(err @ TrySendError::Closed(_)) => {
190
8
            Err(Error::from_std_err(Code::Internal, &err).append("Worker disconnected"))
191
        }
192
    }
193
324
}
194
195
/// Reduces the platform properties available on the worker based on the platform properties provided.
196
/// This is used because we allow more than 1 job to run on a worker at a time, and this is how the
197
/// scheduler knows if more jobs can run on a given worker.
198
156
fn reduce_platform_properties(
199
156
    parent_props: &mut PlatformProperties,
200
156
    reduction_props: &PlatformProperties,
201
156
) {
202
156
    debug_assert!(
reduction_props0
.
is_satisfied_by0
(
parent_props0
, false));
203
156
    for (
property107
,
prop_value107
) in &reduction_props.properties {
204
107
        if let PlatformPropertyValue::Minimum(
value85
) = prop_value {
205
85
            let worker_props = &mut parent_props.properties;
206
85
            if let PlatformPropertyValue::Minimum(worker_value) =
207
85
                worker_props.get_mut(property).unwrap()
208
85
            {
209
85
                *worker_value -= value;
210
85
            
}0
211
22
        }
212
    }
213
156
}
214
215
impl Worker {
216
136
    pub fn new(
217
136
        id: WorkerId,
218
136
        platform_properties: PlatformProperties,
219
136
        tx: Sender<UpdateForWorker>,
220
136
        timestamp: WorkerTimestamp,
221
136
        max_inflight_tasks: u64,
222
136
    ) -> Self {
223
136
        Self {
224
136
            id,
225
136
            total_platform_properties: platform_properties.clone(),
226
136
            platform_properties,
227
136
            tx,
228
136
            running_action_infos: HashMap::new(),
229
136
            last_update_timestamp: timestamp,
230
136
            is_paused: false,
231
136
            pause_needs_kb: None,
232
136
            admits_when_idle: false,
233
136
            idle_declined_kb: None,
234
136
            is_draining: false,
235
136
            max_inflight_tasks,
236
136
            last_load: None,
237
136
            metrics: Arc::new(Metrics {
238
136
                connected_timestamp: SystemTime::now()
239
136
                    .duration_since(UNIX_EPOCH)
240
136
                    .unwrap()
241
136
                    .as_secs(),
242
136
                actions_completed: CounterWithTime::default(),
243
136
                run_action: AsyncCounterWrapper::default(),
244
136
                keep_alive: FuncCounterWrapper::default(),
245
136
                notify_disconnect: CounterWithTime::default(),
246
136
                kill_operation: CounterWithTime::default(),
247
136
                keep_alives_dropped: CounterWithTime::default(),
248
136
            }),
249
136
        }
250
136
    }
251
252
    /// Sends the initial connection information to the worker. This generally is just meta info.
253
    /// This should only be sent once and should always be the first item in the stream.
254
    /// `memory_property` is the platform property the scheduler vetoes
255
    /// placement on; the worker declines for load against the same one, so
256
    /// the two sides cannot disagree. None means the worker never declines
257
    /// for load.
258
136
    pub fn send_initial_connection_result(
259
136
        &mut self,
260
136
        memory_property: Option<&str>,
261
136
    ) -> Result<(), Error> {
262
136
        send_msg_to_worker(
263
136
            &self.tx,
264
136
            update_for_worker::Update::ConnectionResult(ConnectionResult {
265
136
                worker_id: self.id.clone().into(),
266
136
                dispatch_ack: true,
267
136
                memory_property: memory_property.unwrap_or_default().to_string(),
268
136
            }),
269
        )
270
136
        .err_tip(|| 
format!0
("Failed to send ConnectionResult to worker : {}", self.id))
271
136
    }
272
273
    /// Notifies the worker of a requested state change.
274
186
    pub async fn notify_update(&mut self, worker_update: WorkerUpdate) -> Result<(), Error> {
275
186
        match worker_update {
276
158
            WorkerUpdate::RunAction(action) => {
277
158
                let (operation_id, action_info, dispatched_at) = *action;
278
158
                self.run_action(operation_id, action_info, dispatched_at)
279
158
                    .await
280
            }
281
            WorkerUpdate::Disconnect => {
282
23
                self.metrics.notify_disconnect.inc();
283
23
                send_msg_to_worker(&self.tx, update_for_worker::Update::Disconnect(()))
284
            }
285
5
            WorkerUpdate::KillOperation(operation_id) => {
286
5
                let last_seen = self.last_update_timestamp;
287
5
                let pending_action_info = self
288
5
                    .running_action_infos
289
5
                    .get_mut(&operation_id)
290
5
                    .err_tip(|| 
{0
291
0
                        format!(
292
                            "Worker {} asked to kill operation {operation_id} that is not running on it",
293
                            self.id
294
                        )
295
0
                    })?;
296
                // Set before the send so a racing update_action cannot slip
297
                // through in between.
298
5
                pending_action_info.kill_requested_at = Some(last_seen);
299
5
                self.metrics.kill_operation.inc();
300
5
                send_msg_to_worker(
301
5
                    &self.tx,
302
5
                    update_for_worker::Update::KillOperationRequest(KillOperationRequest {
303
5
                        operation_id: operation_id.to_string(),
304
5
                    }),
305
                )
306
            }
307
        }
308
186
    }
309
310
    /// Tells the worker to stop an operation this scheduler no longer holds
311
    /// for it: a late acknowledgement of a dispatch the sweep already took
312
    /// back means the worker is about to run an action that now belongs
313
    /// elsewhere. Nothing to book, since the operation is not on the ledger.
314
1
    pub(crate) fn kill_unknown_operation(&self, operation_id: &OperationId) -> Result<(), Error> {
315
1
        self.metrics.kill_operation.inc();
316
1
        send_msg_to_worker(
317
1
            &self.tx,
318
1
            update_for_worker::Update::KillOperationRequest(KillOperationRequest {
319
1
                operation_id: operation_id.to_string(),
320
1
            }),
321
        )
322
1
    }
323
324
    /// Forgets a kill request whose message never reached the worker, so
325
    /// the next revoked-operation sweep sends it again.
326
1
    pub(crate) fn clear_kill_request(&mut self, operation_id: &OperationId) {
327
1
        if let Some(pending_action_info) = self.running_action_infos.get_mut(operation_id) {
328
1
            pending_action_info.kill_requested_at = None;
329
1
        
}0
330
1
    }
331
332
    /// Whether the worker has been told to kill this operation.
333
54
    pub(crate) fn is_kill_requested(&self, operation_id: &OperationId) -> bool {
334
54
        self.running_action_infos
335
54
            .get(operation_id)
336
54
            .is_some_and(|pending_action_info| pending_action_info.kill_requested_at.is_some())
337
54
    }
338
339
1
    pub fn keep_alive(&mut self) -> Result<(), Error> {
340
1
        let tx = &mut self.tx;
341
1
        let id = &self.id;
342
1
        let dropped = &self.metrics.keep_alives_dropped;
343
1
        self.metrics.keep_alive.wrap(move || {
344
1
            match send_msg_to_worker(tx, update_for_worker::Update::KeepAlive(())) {
345
                // A worker that is not reading has a full queue of work
346
                // ahead of the keepalive; skipping one costs nothing.
347
0
                Err(err) if err.code == Code::ResourceExhausted => {
348
0
                    dropped.inc();
349
0
                    Ok(())
350
                }
351
1
                other => other.err_tip(|| 
format!0
("Failed to send KeepAlive to worker : {id}")),
352
            }
353
1
        })
354
1
    }
355
356
    /// The worker said it took this operation. False when the operation is
357
    /// not on this worker, which a late acknowledgement can cause.
358
4
    pub fn mark_accepted(&mut self, operation_id: &OperationId) -> bool {
359
4
        match self.running_action_infos.get_mut(operation_id) {
360
3
            Some(pending) => {
361
3
                pending.accepted = true;
362
3
                true
363
            }
364
1
            None => false,
365
        }
366
4
    }
367
368
    /// Operations dispatched at or before `cutoff` that the worker never
369
    /// acknowledged.
370
8
    pub fn unacknowledged(&self, cutoff: WorkerTimestamp) -> Vec<OperationId> {
371
8
        self.running_action_infos
372
8
            .iter()
373
8
            .filter(|(_, pending)| 
!pending.accepted7
&&
pending.dispatched_at <= cutoff4
)
374
8
            .map(|(operation_id, _)| 
operation_id2
.
clone2
())
375
8
            .collect()
376
8
    }
377
378
158
    async fn run_action(
379
158
        &mut self,
380
158
        operation_id: OperationId,
381
158
        action_info: ActionInfoWithProps,
382
158
        dispatched_at: WorkerTimestamp,
383
158
    ) -> Result<(), Error> {
384
158
        let tx = &mut self.tx;
385
158
        let worker_platform_properties = &mut self.platform_properties;
386
158
        let running_action_infos = &mut self.running_action_infos;
387
158
        let worker_id = self.id.clone().into();
388
158
        self.metrics
389
158
            .run_action
390
158
            .wrap(async move {
391
158
                let action_info_clone = action_info.clone();
392
158
                let operation_id_string = operation_id.to_string();
393
158
                let start_execute = StartExecute {
394
158
                    execute_request: Some(action_info_clone.inner.as_ref().into()),
395
158
                    operation_id: operation_id_string,
396
158
                    queued_timestamp: Some(action_info.inner.insert_timestamp.into()),
397
158
                    platform: Some((&action_info.platform_properties).into()),
398
158
                    worker_id,
399
158
                    request_metadata: action_info.origin_metadata.bazel_metadata.clone(),
400
158
                };
401
                // Send first: a channel that will not take the message
402
                // leaves the ledger untouched, and the caller requeues.
403
158
                send_msg_to_worker(tx, update_for_worker::Update::StartAction(start_execute))
?2
;
404
156
                reduce_platform_properties(
405
156
                    worker_platform_properties,
406
156
                    &action_info.platform_properties,
407
                );
408
156
                running_action_infos.insert(
409
156
                    operation_id,
410
156
                    PendingActionInfoData {
411
156
                        action_info,
412
156
                        kill_requested_at: None,
413
156
                        dispatched_at,
414
156
                        accepted: false,
415
156
                        last_usage: None,
416
156
                    },
417
                );
418
156
                Ok(())
419
158
            })
420
158
            .await
421
158
    }
422
423
    /// Releases everything the action reserved. This is the only place the
424
    /// `Minimum` budget comes back: the worker's `ExecuteComplete` used to
425
    /// restore it when the process exited, but the action is still resident
426
    /// through output upload, which is where its memory peaks (output buffers,
427
    /// the upload fan-out), so handing the budget back then admitted new work
428
    /// onto a worker at its fullest.
429
48
    pub(crate) fn complete_action(&mut self, operation_id: &OperationId) -> Result<(), Error> {
430
48
        let pending_action_info = self.running_action_infos.remove(operation_id).err_tip(|| 
{0
431
0
            format!(
432
                "Worker {} tried to complete operation {} that was not running",
433
                self.id, operation_id
434
            )
435
0
        })?;
436
48
        self.restore_platform_properties(&pending_action_info.action_info.platform_properties);
437
48
        self.is_paused = false;
438
48
        self.pause_needs_kb = None;
439
48
        self.metrics.actions_completed.inc();
440
48
        Ok(())
441
48
    }
442
443
0
    pub fn has_actions(&self) -> bool {
444
0
        !self.running_action_infos.is_empty()
445
0
    }
446
447
48
    fn restore_platform_properties(&mut self, props: &PlatformProperties) {
448
48
        for (
property35
,
prop_value35
) in &props.properties {
449
35
            if let PlatformPropertyValue::Minimum(value) = prop_value {
450
35
                let worker_props = &mut self.platform_properties.properties;
451
35
                if let PlatformPropertyValue::Minimum(worker_value) =
452
35
                    worker_props.get_mut(property).unwrap()
453
35
                {
454
35
                    *worker_value += value;
455
35
                
}0
456
0
            }
457
        }
458
48
    }
459
460
668
    pub fn can_accept_work(&self) -> bool {
461
668
        !self.is_paused
462
648
            && !self.is_draining
463
646
            && (self.max_inflight_tasks == 0
464
251
                || u64::try_from(self.running_action_infos.len()).unwrap_or(u64::MAX)
465
251
                    < self.max_inflight_tasks)
466
668
    }
467
}
468
469
impl PartialEq for Worker {
470
0
    fn eq(&self, other: &Self) -> bool {
471
0
        self.id == other.id
472
0
    }
473
}
474
475
impl Eq for Worker {}
476
477
impl Hash for Worker {
478
0
    fn hash<H: Hasher>(&self, state: &mut H) {
479
0
        self.id.hash(state);
480
0
    }
481
}
482
483
#[derive(Debug, Default, MetricsComponent)]
484
struct Metrics {
485
    #[metric(help = "The timestamp of when this worker connected.")]
486
    connected_timestamp: u64,
487
    #[metric(help = "The number of actions completed for this worker.")]
488
    actions_completed: CounterWithTime,
489
    #[metric(help = "The number of actions started for this worker.")]
490
    run_action: AsyncCounterWrapper,
491
    #[metric(help = "The number of keep_alive sent to this worker.")]
492
    keep_alive: FuncCounterWrapper,
493
    #[metric(help = "The number of notify_disconnect sent to this worker.")]
494
    notify_disconnect: CounterWithTime,
495
    #[metric(help = "The number of kill_operation sent to this worker.")]
496
    kill_operation: CounterWithTime,
497
    #[metric(help = "Keepalives not sent because the worker's channel was full.")]
498
    keep_alives_dropped: CounterWithTime,
499
}