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/admin.rs
Line
Count
Source
1
// Copyright 2026 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
//! Read-only views for the admin API: what the scheduler has queued and what
16
//! it has connected. A provisioner sizes pods from the first and drains
17
//! through the second, instead of reading the scheduler's store behind its
18
//! back.
19
20
use std::collections::HashMap;
21
use std::time::{SystemTime, UNIX_EPOCH};
22
23
use futures::StreamExt;
24
use nativelink_error::{Error, ResultExt};
25
use nativelink_util::action_messages::ActionStage;
26
use nativelink_util::metrics::record_awaited_action_orphan;
27
use nativelink_util::operation_state_manager::{
28
    ClientStateManager, OperationFilter, OperationStageFlags, OrderDirection,
29
};
30
use serde::{Deserialize, Serialize};
31
use tracing::warn;
32
33
use crate::simple_scheduler_state_manager::is_lost_record;
34
use crate::worker_scheduler::WorkerSummary;
35
36
/// A queued action as the admin API lists it: enough to size a pod for it
37
/// and to tell how long it has waited.
38
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
39
pub struct QueuedDemand {
40
    pub operation_id: String,
41
    /// Milliseconds since the epoch when the action was queued.
42
    pub queued_since_ms: u64,
43
    pub priority: i32,
44
    /// The action's properties as the scheduler matches them, reservations
45
    /// included.
46
    pub platform_properties: HashMap<String, String>,
47
    /// Milliseconds since the epoch when a client last checked in on the
48
    /// operation. A queue entry no client has touched for longer than the
49
    /// client timeout is one the scheduler will not dispatch.
50
    pub client_last_seen_ms: Option<u64>,
51
}
52
53
/// Every queued action of the scheduler, oldest first as the store yields
54
/// them.
55
6
pub async fn queued_demand(scheduler: &dyn ClientStateManager) -> Result<Vec<QueuedDemand>, Error> {
56
6
    let mut stream = scheduler
57
6
        .filter_operations(OperationFilter {
58
6
            stages: OperationStageFlags::Queued,
59
6
            // The Redis backend serves the queue in one direction only, the
60
6
            // matcher's: highest priority first, then oldest.
61
6
            order_by_priority_direction: Some(OrderDirection::Desc),
62
6
            ..Default::default()
63
6
        })
64
6
        .await
65
6
        .err_tip(|| "Listing queued operations for the admin API")
?0
;
66
6
    let mut out = Vec::new();
67
16
    while let Some(
result11
) = stream.next().await {
68
        // The filter read the record once; each read below is a fresh one,
69
        // and the record can have moved on or gone in between. One that is
70
        // gone or cannot be decoded is left out and counted, one that has
71
        // since been dispatched is no longer demand, and a store failure is
72
        // the listing's error.
73
11
        let Some((
state8
, _)) = still_listed(result.as_state().await)
?1
else {
74
2
            continue;
75
        };
76
8
        if !
matches!1
(state.stage, ActionStage::Queued) {
77
1
            continue;
78
7
        }
79
7
        let Some((action_info, _)) = still_listed(result.as_action_info().await)
?0
else {
80
0
            continue;
81
        };
82
7
        let Some(client_last_seen) = still_listed(result.client_last_seen().await)
?0
else {
83
0
            continue;
84
        };
85
7
        let client_last_seen_ms = client_last_seen.map(epoch_ms);
86
7
        out.push(QueuedDemand {
87
7
            operation_id: state.client_operation_id.to_string(),
88
7
            queued_since_ms: epoch_ms(action_info.insert_timestamp),
89
7
            priority: action_info.priority,
90
7
            platform_properties: action_info.platform_properties.clone(),
91
7
            client_last_seen_ms,
92
7
        });
93
    }
94
5
    Ok(out)
95
6
}
96
97
/// A read of a listed record: the value, `None` for a record that is gone
98
/// or cannot be decoded (skipped and counted), or the store's error.
99
25
fn still_listed<T>(read: Result<T, Error>) -> Result<Option<T>, Error> {
100
3
    match read {
101
22
        Ok(value) => Ok(Some(value)),
102
3
        Err(
err2
) if is_lost_record(&err
)2
=> {
103
2
            warn!(
104
                ?err,
105
                "Queued operation listed but its record cannot be read; left out of the demand"
106
            );
107
2
            record_awaited_action_orphan("admin");
108
2
            Ok(None)
109
        }
110
1
        Err(err) => Err(err).err_tip(|| "Reading a queued operation for the admin API"),
111
    }
112
25
}
113
114
11
fn epoch_ms(at: SystemTime) -> u64 {
115
11
    u64::try_from(
116
11
        at.duration_since(UNIX_EPOCH)
117
11
            .unwrap_or_default()
118
11
            .as_millis(),
119
    )
120
11
    .unwrap_or(u64::MAX)
121
11
}
122
123
/// The demand listing as the JSON body the admin API serves.
124
1
pub async fn queued_demand_json(scheduler: &dyn ClientStateManager) -> Result<String, Error>0
{
125
1
    let demand = queued_demand(scheduler).await
?0
;
126
1
    serde_json::to_string(&demand).map_err(|e| 
{0
127
0
        Error::new(
128
0
            nativelink_error::Code::Internal,
129
0
            format!("Serializing the demand listing: {e}"),
130
        )
131
0
    })
132
1
}
133
134
/// The worker listing as the JSON body the admin API serves.
135
1
pub fn workers_json(workers: &[WorkerSummary]) -> Result<String, Error> {
136
1
    serde_json::to_string(workers).map_err(|e| 
{0
137
0
        Error::new(
138
0
            nativelink_error::Code::Internal,
139
0
            format!("Serializing the worker listing: {e}"),
140
        )
141
0
    })
142
1
}