/build/source/nativelink-worker/src/worker_utils.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::BuildHasher; |
16 | | use std::collections::HashMap; |
17 | | use std::io::{BufRead, BufReader, Cursor}; |
18 | | use std::process::Stdio; |
19 | | |
20 | | use futures::future::try_join_all; |
21 | | use nativelink_config::cas_server::WorkerProperty; |
22 | | use nativelink_error::{Code, Error, ResultExt, make_err, make_input_err}; |
23 | | use nativelink_proto::build::bazel::remote::execution::v2::platform::Property; |
24 | | use nativelink_proto::com::github::trace_machina::nativelink::remote_execution::ConnectWorkerRequest; |
25 | | use tokio::process; |
26 | | use tracing::{info, warn}; |
27 | | |
28 | | #[expect(clippy::future_not_send)] // TODO(jhpratt) remove this |
29 | 36 | pub async fn make_connect_worker_request<S: BuildHasher>( |
30 | 36 | worker_id_prefix: String, |
31 | 36 | worker_properties: &HashMap<String, WorkerProperty, S>, |
32 | 36 | extra_envs: &HashMap<String, String, S>, |
33 | 36 | max_inflight_tasks: u64, |
34 | 36 | admits_when_idle: bool, |
35 | 36 | ) -> Result<ConnectWorkerRequest, Error> { |
36 | 36 | let mut futures = vec![]; |
37 | 36 | for (property_name8 , worker_property8 ) in worker_properties { |
38 | 8 | futures.push(async move { |
39 | 8 | match worker_property { |
40 | 1 | WorkerProperty::Values(values) => { |
41 | 1 | let mut props = Vec::with_capacity(values.len()); |
42 | 1 | for value in values { |
43 | 1 | props.push(Property { |
44 | 1 | name: property_name.clone(), |
45 | 1 | value: value.clone(), |
46 | 1 | }); |
47 | 1 | } |
48 | 1 | Ok(props) |
49 | | } |
50 | 7 | WorkerProperty::QueryCmd(cmd) => { |
51 | 7 | let maybe_split_cmd = shlex::split(cmd); |
52 | 7 | let (command6 , args6 ) = match &maybe_split_cmd { |
53 | 6 | Some(split_cmd) => (&split_cmd[0], &split_cmd[1..]), |
54 | | None => { |
55 | 1 | return Err(make_input_err!( |
56 | 1 | "Could not parse the value of worker property: {}: '{}'", |
57 | 1 | property_name, |
58 | 1 | cmd |
59 | 1 | )); |
60 | | } |
61 | | }; |
62 | 6 | let mut process = process::Command::new(command); |
63 | 6 | process.env_clear(); |
64 | 6 | process.envs(extra_envs); |
65 | 6 | process.args(args); |
66 | 6 | process.stdin(Stdio::null()); |
67 | 6 | let err_fn = || {2 |
68 | 2 | format!("Error executing property_name {property_name} command: '{cmd}'") |
69 | 2 | }; |
70 | 6 | info!(cmd, property_name, "Spawning process",); |
71 | 6 | let process_output = process.output().await.err_tip(err_fn)?0 ; |
72 | 6 | if !process_output.stderr.is_empty() { |
73 | 2 | warn!( |
74 | 2 | stderr = ?String::from_utf8_lossy(&process_output.stderr), |
75 | | cmd = cmd, |
76 | | property_name = property_name, |
77 | | "Got stderr when running query cmd" |
78 | | ); |
79 | 4 | } |
80 | 6 | if !process_output.status.success() { |
81 | 2 | let Some(exit_code1 ) = process_output.status.code() else { |
82 | 1 | return Err(make_err!( |
83 | 1 | Code::Internal, |
84 | 1 | "{}: {}", |
85 | 1 | err_fn(), |
86 | 1 | process_output.status |
87 | 1 | )); |
88 | | }; |
89 | 1 | return Err(make_err!(exit_code.into(), "{}", err_fn())); |
90 | 4 | } |
91 | 4 | let reader = BufReader::new(Cursor::new(process_output.stdout)); |
92 | | |
93 | 4 | let mut props = vec![]; |
94 | 4 | for value3 in reader.lines() { |
95 | 3 | props.push(Property { |
96 | 3 | name: property_name.clone(), |
97 | 3 | value: value |
98 | 3 | .err_tip(|| "Could split input by lines")?0 |
99 | 3 | .trim() |
100 | 3 | .to_string(), |
101 | | }); |
102 | | } |
103 | 4 | Ok(props) |
104 | | } |
105 | | } |
106 | 8 | }); |
107 | | } |
108 | | |
109 | | Ok(ConnectWorkerRequest { |
110 | 36 | worker_id_prefix, |
111 | 36 | properties: try_join_all(futures).await?3 .into_iter33 ().flatten33 ().collect33 (), |
112 | 33 | max_inflight_tasks, |
113 | 33 | admits_when_idle, |
114 | | }) |
115 | 36 | } |