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