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