/build/source/nativelink-worker/src/buck2_file_capture.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 | | //! Buck2-only execution-container capture, held through action cleanup. |
16 | | //! The helper initiates authenticated uploads; `NativeLink` never opens a public |
17 | | //! file server or grants a browser access to the worker API. |
18 | | |
19 | | use core::time::Duration; |
20 | | use std::path::Path; |
21 | | use std::process::Stdio; |
22 | | |
23 | | use nativelink_config::cas_server::Buck2FileCaptureConfig; |
24 | | use nativelink_error::{Code, Error, ResultExt, make_err}; |
25 | | use nativelink_proto::build::bazel::remote::execution::v2::RequestMetadata; |
26 | | use serde_json::json; |
27 | | use tokio::io::{AsyncReadExt, AsyncWriteExt}; |
28 | | use tokio::process::{Child, Command}; |
29 | | use uuid::Uuid; |
30 | | |
31 | | #[derive(Debug)] |
32 | | pub struct Buck2FileCapture { |
33 | | child: Child, |
34 | | /// Registered with the reaper for as long as the helper is ours to wait on. |
35 | | _owned: crate::reaper::OwnedChild, |
36 | | finish_timeout: Duration, |
37 | | } |
38 | | |
39 | | impl Buck2FileCapture { |
40 | 19 | pub fn eligible(metadata: Option<&RequestMetadata>) -> bool { |
41 | 19 | metadata.is_some_and(|metadata| {18 |
42 | 18 | !metadata.tool_invocation_id.is_empty() |
43 | 16 | && metadata |
44 | 16 | .tool_details |
45 | 16 | .as_ref() |
46 | 16 | .is_some_and(|tool| tool.tool_name.eq_ignore_ascii_case("buck2")) |
47 | 18 | }) |
48 | 19 | } |
49 | | /// Only authenticated scheduler metadata identifying Buck2 may start a |
50 | | /// capture. A missing invocation or another tool leaves the old flow intact. |
51 | 3 | pub async fn start( |
52 | 3 | config: &Buck2FileCaptureConfig, |
53 | 3 | metadata: Option<&RequestMetadata>, |
54 | 3 | worker: &str, |
55 | 3 | attempt: &str, |
56 | 3 | action_digest: &str, |
57 | 7 | ) -> Result<Option<Self>, Error> { |
58 | 7 | let Some(metadata4 ) = metadata.filter(|metadata| Self::eligible6 (Some(metadata)6 )) else { |
59 | 3 | return Ok(None); |
60 | | }; |
61 | 4 | let session = Uuid::new_v4().to_string(); |
62 | 4 | let finalize_seconds = if config.finalize_timeout_s == 0 { |
63 | 0 | 120 |
64 | | } else { |
65 | 4 | config.finalize_timeout_s |
66 | | }; |
67 | 4 | let configuration = json!({ |
68 | 4 | "root": "/", |
69 | 4 | "stateDirectory": Path::new(&config.state_directory).join(&session), |
70 | 4 | "protectedPaths": config.protected_paths, |
71 | 4 | "internalPaths": config.internal_paths, |
72 | 4 | "gateway": config.gateway, |
73 | 4 | "tokenFile": config.token_file, |
74 | 4 | "finalizeSeconds": finalize_seconds, |
75 | 4 | "session": { |
76 | 4 | "id": session, |
77 | 4 | "tool": "buck2", |
78 | 4 | "build": metadata.tool_invocation_id, |
79 | 4 | "action": action_digest, |
80 | 4 | "attempt": attempt, |
81 | 4 | "worker": worker, |
82 | 4 | "container": config.container_id, |
83 | | }, |
84 | | }); |
85 | 4 | let mut body = serde_json::to_vec(&configuration) |
86 | 4 | .map_err(|err| Error::from_std_err0 (Code::Internal0 , &err0 ))?0 ; |
87 | 4 | body.push(b'\n'); |
88 | 4 | let mut command = Command::new(&config.executable); |
89 | 4 | command |
90 | 4 | .stdin(Stdio::piped()) |
91 | 4 | .stdout(Stdio::piped()) |
92 | 4 | .stderr(Stdio::inherit()) |
93 | 4 | .env_clear() |
94 | 4 | .kill_on_drop(true); |
95 | | // Customer-provided trust roots are configuration, not action env. |
96 | 8 | for name in ["SSL_CERT_FILE", "SSL_CERT_DIR"]4 { |
97 | 8 | if let Some(value4 ) = std::env::var_os(name) { |
98 | 4 | command.env(name, value); |
99 | 4 | } |
100 | | } |
101 | 4 | let startup = async move { |
102 | 4 | let mut child = command |
103 | 4 | .spawn() |
104 | 4 | .err_tip(|| "Starting Buck2 file capture helper")?0 ; |
105 | 4 | child |
106 | 4 | .stdin |
107 | 4 | .as_mut() |
108 | 4 | .err_tip(|| "Capture helper has no control pipe")?0 |
109 | 4 | .write_all(&body) |
110 | 4 | .await |
111 | 4 | .err_tip(|| "Configuring Buck2 file capture")?0 ; |
112 | 4 | let mut stdout = child |
113 | 4 | .stdout |
114 | 4 | .take() |
115 | 4 | .err_tip(|| "Capture helper has no readiness pipe")?0 ; |
116 | 4 | let mut ready = [0_u8; 6]; |
117 | 4 | stdout |
118 | 4 | .read_exact(&mut ready) |
119 | 4 | .await |
120 | 4 | .err_tip(|| "Waiting for Buck2 capture readiness")?0 ; |
121 | 4 | if &ready != b"ready\n" { |
122 | 0 | return Err(make_err!( |
123 | 0 | Code::FailedPrecondition, |
124 | 0 | "Invalid Buck2 capture readiness acknowledgement" |
125 | 0 | )); |
126 | 4 | } |
127 | 4 | let owned = crate::reaper::OwnedChild::new(child.id()); |
128 | 4 | Ok(Some(Self { |
129 | 4 | child, |
130 | 4 | _owned: owned, |
131 | 4 | // The helper reserves five seconds to acknowledge a partial |
132 | 4 | // final scan, plus five seconds for shutdown/reaping here. |
133 | 4 | finish_timeout: Duration::from_secs(finalize_seconds as u64 + 10), |
134 | 4 | })) |
135 | 4 | }; |
136 | 4 | tokio::time::timeout(Duration::from_secs(30), startup) |
137 | 4 | .await |
138 | 4 | .err_tip(|| "Buck2 capture registration timed out")?0 |
139 | 7 | } |
140 | | |
141 | | /// Closing the control pipe asks the helper to stop live work, capture a |
142 | | /// final revision and wait for durable acknowledgement. This must run before |
143 | | /// deleting action files or reporting the single-use worker as completed. |
144 | 5 | pub async fn finish(mut self) -> Result<(), Error> { |
145 | 5 | drop(self.child.stdin.take()); |
146 | 5 | if let Ok(status4 ) = tokio::time::timeout(self.finish_timeout, self.child.wait()).await { |
147 | 4 | let status = status.err_tip(|| "Reaping Buck2 capture helper")?0 ; |
148 | 4 | if status.success() { |
149 | 4 | Ok(()) |
150 | | } else { |
151 | 0 | Err(make_err!( |
152 | 0 | Code::Unavailable, |
153 | 0 | "Buck2 file capture ended with {status}; archive may be partial" |
154 | 0 | )) |
155 | | } |
156 | | } else { |
157 | 1 | self.child |
158 | 1 | .kill() |
159 | 1 | .await |
160 | 1 | .err_tip(|| "Stopping timed-out Buck2 capture helper")?0 ; |
161 | 1 | Err(make_err!( |
162 | 1 | Code::DeadlineExceeded, |
163 | 1 | "Buck2 capture finalization deadline exceeded; archive is partial" |
164 | 1 | )) |
165 | | } |
166 | 5 | } |
167 | | } |
168 | | |
169 | | #[cfg(test)] |
170 | | mod tests { |
171 | | use nativelink_proto::build::bazel::remote::execution::v2::ToolDetails; |
172 | | |
173 | | use super::*; |
174 | | |
175 | 9 | fn metadata(tool: &str, invocation: &str) -> RequestMetadata { |
176 | 9 | RequestMetadata { |
177 | 9 | tool_details: Some(ToolDetails { |
178 | 9 | tool_name: tool.to_string(), |
179 | 9 | ..Default::default() |
180 | 9 | }), |
181 | 9 | tool_invocation_id: invocation.to_string(), |
182 | 9 | ..Default::default() |
183 | 9 | } |
184 | 9 | } |
185 | | |
186 | | #[test] |
187 | 1 | fn only_buck2_with_an_invocation_is_eligible() { |
188 | 1 | assert!(!Buck2FileCapture::eligible(None)); |
189 | 4 | for tool in ["", 1 "bazel"1 , "pants", "recc"] { |
190 | 4 | assert!(!Buck2FileCapture::eligible(Some(&metadata( |
191 | 4 | tool, |
192 | 4 | "invocation" |
193 | 4 | )))); |
194 | | } |
195 | 1 | assert!(!Buck2FileCapture::eligible(Some(&metadata("buck2", "")))); |
196 | 1 | assert!(Buck2FileCapture::eligible(Some(&metadata( |
197 | 1 | "buck2", |
198 | 1 | "invocation" |
199 | 1 | )))); |
200 | 1 | assert!(Buck2FileCapture::eligible(Some(&metadata( |
201 | 1 | "Buck2", |
202 | 1 | "invocation" |
203 | 1 | )))); |
204 | 1 | } |
205 | | |
206 | | #[cfg(target_family = "unix")] |
207 | | #[nativelink_macro::nativelink_test] |
208 | | async fn non_buck2_does_not_launch_the_helper() -> Result<(), Error> { |
209 | | let config = Buck2FileCaptureConfig { |
210 | | executable: "/nonexistent-capture-helper".to_string(), |
211 | | gateway: String::new(), |
212 | | token_file: String::new(), |
213 | | state_directory: String::new(), |
214 | | container_id: String::new(), |
215 | | protected_paths: vec![], |
216 | | internal_paths: vec![], |
217 | | finalize_timeout_s: 0, |
218 | | }; |
219 | | for request in [ |
220 | | None, |
221 | | Some(metadata("bazel", "invocation")), |
222 | | Some(metadata("buck2", "")), |
223 | | ] { |
224 | | assert!( |
225 | | Buck2FileCapture::start(&config, request.as_ref(), "worker", "attempt", "digest") |
226 | | .await? |
227 | | .is_none() |
228 | | ); |
229 | | } |
230 | | Ok(()) |
231 | | } |
232 | | |
233 | | #[cfg(target_family = "unix")] |
234 | | #[nativelink_macro::nativelink_test] |
235 | | async fn unresponsive_capture_has_a_bounded_shutdown() -> Result<(), Error> { |
236 | | // Nix supplies coreutils through PATH and has no /bin/sleep. |
237 | | let child = Command::new("sleep") |
238 | | .arg("30") |
239 | | .kill_on_drop(true) |
240 | | .spawn() |
241 | | .err_tip(|| "Starting unresponsive test helper")?; |
242 | | let owned = crate::reaper::OwnedChild::new(child.id()); |
243 | | let capture = Buck2FileCapture { |
244 | | child, |
245 | | _owned: owned, |
246 | | finish_timeout: Duration::from_millis(20), |
247 | | }; |
248 | | let result = tokio::time::timeout(Duration::from_secs(2), capture.finish()) |
249 | | .await |
250 | | .err_tip(|| "Capture deadline did not bound shutdown")?; |
251 | | assert_eq!(result.unwrap_err().code, Code::DeadlineExceeded); |
252 | | Ok(()) |
253 | | } |
254 | | } |