Coverage Report

Created: 2026-10-01 05:28

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
    /// Whether the worker is draining.
131
    #[metric(help = "If the worker is draining.")]
132
    pub is_draining: bool,
133
134
    /// Maximum inflight tasks for this worker (or 0 for unlimited)
135
    #[metric(help = "Maximum inflight tasks for this worker (or 0 for unlimited)")]
136
    pub max_inflight_tasks: u64,
137
138
    /// What the worker last reported having to spare, from its keepalive;
139
    /// `None` until it reports, and for workers that never do.
140
    pub last_load: Option<WorkerLoad>,
141
142
    /// Stats about the worker.
143
    #[metric]
144
    metrics: Arc<Metrics>,
145
}
146
147
/// Messages the channel holds beyond the worker's own concurrency: kills,
148
/// keepalives and the disconnect, which arrive alongside the dispatches.
149
const CHANNEL_HEADROOM: usize = 16;
150
/// A worker that set no `max_inflight_tasks` gets this much queue.
151
const DEFAULT_CHANNEL_DISPATCHES: usize = 256;
152
153
/// How many messages the worker's channel holds: what its concurrency
154
/// allows in flight plus the headroom. Bounded so that a worker that has
155
/// stopped reading cannot be handed work without limit; a full channel is
156
/// `ResourceExhausted` on the send, and the caller requeues.
157
0
pub fn channel_capacity(max_inflight_tasks: u64) -> usize {
158
0
    let ceiling = DEFAULT_CHANNEL_DISPATCHES * 4;
159
0
    let dispatches = match usize::try_from(max_inflight_tasks) {
160
0
        Ok(0) => DEFAULT_CHANNEL_DISPATCHES,
161
0
        Ok(n) if n <= ceiling => n,
162
0
        _ => ceiling,
163
    };
164
0
    dispatches + CHANNEL_HEADROOM
165
0
}
166
167
294
fn send_msg_to_worker(
168
294
    tx: &Sender<UpdateForWorker>,
169
294
    msg: update_for_worker::Update,
170
294
) -> Result<(), Error> {
171
294
    match tx.try_send(UpdateForWorker { update: Some(msg) }) {
172
284
        Ok(()) => Ok(()),
173
2
        Err(TrySendError::Full(_)) => Err(Error::new(
174
2
            Code::ResourceExhausted,
175
2
            "Worker channel full".to_string(),
176
2
        )),
177
8
        Err(err @ TrySendError::Closed(_)) => {
178
8
            Err(Error::from_std_err(Code::Internal, &err).append("Worker disconnected"))
179
        }
180
    }
181
294
}
182
183
/// Reduces the platform properties available on the worker based on the platform properties provided.
184
/// This is used because we allow more than 1 job to run on a worker at a time, and this is how the
185
/// scheduler knows if more jobs can run on a given worker.
186
139
fn reduce_platform_properties(
187
139
    parent_props: &mut PlatformProperties,
188
139
    reduction_props: &PlatformProperties,
189
139
) {
190
139
    debug_assert!(
reduction_props0
.
is_satisfied_by0
(
parent_props0
, false));
191
139
    for (
property94
,
prop_value94
) in &reduction_props.properties {
192
94
        if let PlatformPropertyValue::Minimum(
value72
) = prop_value {
193
72
            let worker_props = &mut parent_props.properties;
194
72
            if let PlatformPropertyValue::Minimum(worker_value) =
195
72
                worker_props.get_mut(property).unwrap()
196
72
            {
197
72
                *worker_value -= value;
198
72
            
}0
199
22
        }
200
    }
201
139
}
202
203
impl Worker {
204
125
    pub fn new(
205
125
        id: WorkerId,
206
125
        platform_properties: PlatformProperties,
207
125
        tx: Sender<UpdateForWorker>,
208
125
        timestamp: WorkerTimestamp,
209
125
        max_inflight_tasks: u64,
210
125
    ) -> Self {
211
125
        Self {
212
125
            id,
213
125
            total_platform_properties: platform_properties.clone(),
214
125
            platform_properties,
215
125
            tx,
216
125
            running_action_infos: HashMap::new(),
217
125
            last_update_timestamp: timestamp,
218
125
            is_paused: false,
219
125
            pause_needs_kb: None,
220
125
            is_draining: false,
221
125
            max_inflight_tasks,
222
125
            last_load: None,
223
125
            metrics: Arc::new(Metrics {
224
125
                connected_timestamp: SystemTime::now()
225
125
                    .duration_since(UNIX_EPOCH)
226
125
                    .unwrap()
227
125
                    .as_secs(),
228
125
                actions_completed: CounterWithTime::default(),
229
125
                run_action: AsyncCounterWrapper::default(),
230
125
                keep_alive: FuncCounterWrapper::default(),
231
125
                notify_disconnect: CounterWithTime::default(),
232
125
                kill_operation: CounterWithTime::default(),
233
125
                keep_alives_dropped: CounterWithTime::default(),
234
125
            }),
235
125
        }
236
125
    }
237
238
    /// Sends the initial connection information to the worker. This generally is just meta info.
239
    /// This should only be sent once and should always be the first item in the stream.
240
    /// `memory_property` is the platform property the scheduler vetoes
241
    /// placement on; the worker declines for load against the same one, so
242
    /// the two sides cannot disagree. None means the worker never declines
243
    /// for load.
244
125
    pub fn send_initial_connection_result(
245
125
        &mut self,
246
125
        memory_property: Option<&str>,
247
125
    ) -> Result<(), Error> {
248
125
        send_msg_to_worker(
249
125
            &self.tx,
250
125
            update_for_worker::Update::ConnectionResult(ConnectionResult {
251
125
                worker_id: self.id.clone().into(),
252
125
                dispatch_ack: true,
253
125
                memory_property: memory_property.unwrap_or_default().to_string(),
254
125
            }),
255
        )
256
125
        .err_tip(|| 
format!0
("Failed to send ConnectionResult to worker : {}", self.id))
257
125
    }
258
259
    /// Notifies the worker of a requested state change.
260
167
    pub async fn notify_update(&mut self, worker_update: WorkerUpdate) -> Result<(), Error> {
261
167
        match worker_update {
262
141
            WorkerUpdate::RunAction(action) => {
263
141
                let (operation_id, action_info, dispatched_at) = *action;
264
141
                self.run_action(operation_id, action_info, dispatched_at)
265
141
                    .await
266
            }
267
            WorkerUpdate::Disconnect => {
268
21
                self.metrics.notify_disconnect.inc();
269
21
                send_msg_to_worker(&self.tx, update_for_worker::Update::Disconnect(()))
270
            }
271
5
            WorkerUpdate::KillOperation(operation_id) => {
272
5
                let last_seen = self.last_update_timestamp;
273
5
                let pending_action_info = self
274
5
                    .running_action_infos
275
5
                    .get_mut(&operation_id)
276
5
                    .err_tip(|| 
{0
277
0
                        format!(
278
                            "Worker {} asked to kill operation {operation_id} that is not running on it",
279
                            self.id
280
                        )
281
0
                    })?;
282
                // Set before the send so a racing update_action cannot slip
283
                // through in between.
284
5
                pending_action_info.kill_requested_at = Some(last_seen);
285
5
                self.metrics.kill_operation.inc();
286
5
                send_msg_to_worker(
287
5
                    &self.tx,
288
5
                    update_for_worker::Update::KillOperationRequest(KillOperationRequest {
289
5
                        operation_id: operation_id.to_string(),
290
5
                    }),
291
                )
292
            }
293
        }
294
167
    }
295
296
    /// Tells the worker to stop an operation this scheduler no longer holds
297
    /// for it: a late acknowledgement of a dispatch the sweep already took
298
    /// back means the worker is about to run an action that now belongs
299
    /// elsewhere. Nothing to book, since the operation is not on the ledger.
300
1
    pub(crate) fn kill_unknown_operation(&self, operation_id: &OperationId) -> Result<(), Error> {
301
1
        self.metrics.kill_operation.inc();
302
1
        send_msg_to_worker(
303
1
            &self.tx,
304
1
            update_for_worker::Update::KillOperationRequest(KillOperationRequest {
305
1
                operation_id: operation_id.to_string(),
306
1
            }),
307
        )
308
1
    }
309
310
    /// Forgets a kill request whose message never reached the worker, so
311
    /// the next revoked-operation sweep sends it again.
312
1
    pub(crate) fn clear_kill_request(&mut self, operation_id: &OperationId) {
313
1
        if let Some(pending_action_info) = self.running_action_infos.get_mut(operation_id) {
314
1
            pending_action_info.kill_requested_at = None;
315
1
        
}0
316
1
    }
317
318
    /// Whether the worker has been told to kill this operation.
319
46
    pub(crate) fn is_kill_requested(&self, operation_id: &OperationId) -> bool {
320
46
        self.running_action_infos
321
46
            .get(operation_id)
322
46
            .is_some_and(|pending_action_info| pending_action_info.kill_requested_at.is_some())
323
46
    }
324
325
1
    pub fn keep_alive(&mut self) -> Result<(), Error> {
326
1
        let tx = &mut self.tx;
327
1
        let id = &self.id;
328
1
        let dropped = &self.metrics.keep_alives_dropped;
329
1
        self.metrics.keep_alive.wrap(move || {
330
1
            match send_msg_to_worker(tx, update_for_worker::Update::KeepAlive(())) {
331
                // A worker that is not reading has a full queue of work
332
                // ahead of the keepalive; skipping one costs nothing.
333
0
                Err(err) if err.code == Code::ResourceExhausted => {
334
0
                    dropped.inc();
335
0
                    Ok(())
336
                }
337
1
                other => other.err_tip(|| 
format!0
("Failed to send KeepAlive to worker : {id}")),
338
            }
339
1
        })
340
1
    }
341
342
    /// The worker said it took this operation. False when the operation is
343
    /// not on this worker, which a late acknowledgement can cause.
344
4
    pub fn mark_accepted(&mut self, operation_id: &OperationId) -> bool {
345
4
        match self.running_action_infos.get_mut(operation_id) {
346
3
            Some(pending) => {
347
3
                pending.accepted = true;
348
3
                true
349
            }
350
1
            None => false,
351
        }
352
4
    }
353
354
    /// Operations dispatched at or before `cutoff` that the worker never
355
    /// acknowledged.
356
8
    pub fn unacknowledged(&self, cutoff: WorkerTimestamp) -> Vec<OperationId> {
357
8
        self.running_action_infos
358
8
            .iter()
359
8
            .filter(|(_, pending)| 
!pending.accepted7
&&
pending.dispatched_at <= cutoff4
)
360
8
            .map(|(operation_id, _)| 
operation_id2
.
clone2
())
361
8
            .collect()
362
8
    }
363
364
141
    async fn run_action(
365
141
        &mut self,
366
141
        operation_id: OperationId,
367
141
        action_info: ActionInfoWithProps,
368
141
        dispatched_at: WorkerTimestamp,
369
141
    ) -> Result<(), Error> {
370
141
        let tx = &mut self.tx;
371
141
        let worker_platform_properties = &mut self.platform_properties;
372
141
        let running_action_infos = &mut self.running_action_infos;
373
141
        let worker_id = self.id.clone().into();
374
141
        self.metrics
375
141
            .run_action
376
141
            .wrap(async move {
377
141
                let action_info_clone = action_info.clone();
378
141
                let operation_id_string = operation_id.to_string();
379
141
                let start_execute = StartExecute {
380
141
                    execute_request: Some(action_info_clone.inner.as_ref().into()),
381
141
                    operation_id: operation_id_string,
382
141
                    queued_timestamp: Some(action_info.inner.insert_timestamp.into()),
383
141
                    platform: Some((&action_info.platform_properties).into()),
384
141
                    worker_id,
385
141
                    request_metadata: action_info.origin_metadata.bazel_metadata.clone(),
386
141
                };
387
                // Send first: a channel that will not take the message
388
                // leaves the ledger untouched, and the caller requeues.
389
141
                send_msg_to_worker(tx, update_for_worker::Update::StartAction(start_execute))
?2
;
390
139
                reduce_platform_properties(
391
139
                    worker_platform_properties,
392
139
                    &action_info.platform_properties,
393
                );
394
139
                running_action_infos.insert(
395
139
                    operation_id,
396
139
                    PendingActionInfoData {
397
139
                        action_info,
398
139
                        kill_requested_at: None,
399
139
                        dispatched_at,
400
139
                        accepted: false,
401
139
                        last_usage: None,
402
139
                    },
403
                );
404
139
                Ok(())
405
141
            })
406
141
            .await
407
141
    }
408
409
    /// Releases everything the action reserved. This is the only place the
410
    /// `Minimum` budget comes back: the worker's `ExecuteComplete` used to
411
    /// restore it when the process exited, but the action is still resident
412
    /// through output upload, which is where its memory peaks (output buffers,
413
    /// the upload fan-out), so handing the budget back then admitted new work
414
    /// onto a worker at its fullest.
415
40
    pub(crate) fn complete_action(&mut self, operation_id: &OperationId) -> Result<(), Error> {
416
40
        let pending_action_info = self.running_action_infos.remove(operation_id).err_tip(|| 
{0
417
0
            format!(
418
                "Worker {} tried to complete operation {} that was not running",
419
                self.id, operation_id
420
            )
421
0
        })?;
422
40
        self.restore_platform_properties(&pending_action_info.action_info.platform_properties);
423
40
        self.is_paused = false;
424
40
        self.pause_needs_kb = None;
425
40
        self.metrics.actions_completed.inc();
426
40
        Ok(())
427
40
    }
428
429
0
    pub fn has_actions(&self) -> bool {
430
0
        !self.running_action_infos.is_empty()
431
0
    }
432
433
40
    fn restore_platform_properties(&mut self, props: &PlatformProperties) {
434
40
        for (
property29
,
prop_value29
) in &props.properties {
435
29
            if let PlatformPropertyValue::Minimum(value) = prop_value {
436
29
                let worker_props = &mut self.platform_properties.properties;
437
29
                if let PlatformPropertyValue::Minimum(worker_value) =
438
29
                    worker_props.get_mut(property).unwrap()
439
29
                {
440
29
                    *worker_value += value;
441
29
                
}0
442
0
            }
443
        }
444
40
    }
445
446
618
    pub fn can_accept_work(&self) -> bool {
447
618
        !self.is_paused
448
605
            && !self.is_draining
449
603
            && (self.max_inflight_tasks == 0
450
210
                || u64::try_from(self.running_action_infos.len()).unwrap_or(u64::MAX)
451
210
                    < self.max_inflight_tasks)
452
618
    }
453
}
454
455
impl PartialEq for Worker {
456
0
    fn eq(&self, other: &Self) -> bool {
457
0
        self.id == other.id
458
0
    }
459
}
460
461
impl Eq for Worker {}
462
463
impl Hash for Worker {
464
0
    fn hash<H: Hasher>(&self, state: &mut H) {
465
0
        self.id.hash(state);
466
0
    }
467
}
468
469
#[derive(Debug, Default, MetricsComponent)]
470
struct Metrics {
471
    #[metric(help = "The timestamp of when this worker connected.")]
472
    connected_timestamp: u64,
473
    #[metric(help = "The number of actions completed for this worker.")]
474
    actions_completed: CounterWithTime,
475
    #[metric(help = "The number of actions started for this worker.")]
476
    run_action: AsyncCounterWrapper,
477
    #[metric(help = "The number of keep_alive sent to this worker.")]
478
    keep_alive: FuncCounterWrapper,
479
    #[metric(help = "The number of notify_disconnect sent to this worker.")]
480
    notify_disconnect: CounterWithTime,
481
    #[metric(help = "The number of kill_operation sent to this worker.")]
482
    kill_operation: CounterWithTime,
483
    #[metric(help = "Keepalives not sent because the worker's channel was full.")]
484
    keep_alives_dropped: CounterWithTime,
485
}