/build/source/nativelink-service/src/bep_server.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::pin::Pin; |
16 | | use core::time::Duration; |
17 | | use std::borrow::Cow; |
18 | | |
19 | | use bytes::BytesMut; |
20 | | use futures::Stream; |
21 | | use futures::stream::unfold; |
22 | | use nativelink_error::{Error, ResultExt}; |
23 | | use nativelink_proto::com::github::trace_machina::nativelink::events::{BepEvent, bep_event}; |
24 | | use nativelink_proto::google::devtools::build::v1::publish_build_event_server::{ |
25 | | PublishBuildEvent, PublishBuildEventServer, |
26 | | }; |
27 | | use nativelink_proto::google::devtools::build::v1::{ |
28 | | PublishBuildToolEventStreamRequest, PublishBuildToolEventStreamResponse, |
29 | | PublishLifecycleEventRequest, |
30 | | }; |
31 | | use nativelink_store::store_manager::StoreManager; |
32 | | use nativelink_util::store_trait::{Store, StoreDriver, StoreKey, StoreLike}; |
33 | | use opentelemetry::baggage::BaggageExt; |
34 | | use opentelemetry::context::Context; |
35 | | use opentelemetry_semantic_conventions::attribute::ENDUSER_ID; |
36 | | use prost::Message; |
37 | | use tokio::time::sleep; |
38 | | use tonic::{Request, Response, Result, Status, Streaming}; |
39 | | use tracing::{Level, instrument, warn}; |
40 | | |
41 | | /// Bounded retries for persisting a BEP event so a transient store failure |
42 | | /// (e.g. a Redis Sentinel failover) doesn't drop it — BEP events are the |
43 | | /// authoritative build record, the same durability class as origin events. |
44 | | const MAX_BEP_UPLOAD_ATTEMPTS: u32 = 5; |
45 | | |
46 | | /// Current version of the BEP event. This might be used in the future if |
47 | | /// there is a breaking change in the BEP event format. |
48 | | const BEP_EVENT_VERSION: u32 = 0; |
49 | | |
50 | | #[allow(clippy::result_large_err, reason = "TODO Fix this. Breaks on nightly")] |
51 | 4 | fn get_identity() -> Option<String> { |
52 | 4 | Context::current() |
53 | 4 | .baggage() |
54 | 4 | .get(ENDUSER_ID) |
55 | 4 | .map(|value| value.as_str()0 .to_string0 ()) |
56 | 4 | } |
57 | | |
58 | | #[derive(Debug)] |
59 | | pub struct BepServer { |
60 | | store: Store, |
61 | | } |
62 | | |
63 | | impl BepServer { |
64 | 4 | pub fn new( |
65 | 4 | config: &nativelink_config::cas_server::BepConfig, |
66 | 4 | store_manager: &StoreManager, |
67 | 4 | ) -> Result<Self, Error> { |
68 | 4 | let store = store_manager |
69 | 4 | .get_store(&config.store) |
70 | 4 | .err_tip(|| format!0 ("Expected store {} to exist in store manager", &config.store0 ))?0 ; |
71 | | |
72 | 4 | Ok(Self { store }) |
73 | 4 | } |
74 | | |
75 | 0 | pub fn into_service(self) -> PublishBuildEventServer<Self> { |
76 | 0 | PublishBuildEventServer::new(self) |
77 | 0 | } |
78 | | |
79 | 2 | async fn inner_publish_lifecycle_event( |
80 | 2 | &self, |
81 | 2 | request: PublishLifecycleEventRequest, |
82 | 2 | identity: Option<String>, |
83 | 2 | ) -> Result<Response<()>, Error> { |
84 | 2 | let build_event = request |
85 | 2 | .build_event |
86 | 2 | .as_ref() |
87 | 2 | .err_tip(|| "Expected build_event to be set")?0 ; |
88 | 2 | let stream_id = build_event |
89 | 2 | .stream_id |
90 | 2 | .as_ref() |
91 | 2 | .err_tip(|| "Expected stream_id to be set")?0 ; |
92 | | |
93 | 2 | let sequence_number = build_event.sequence_number; |
94 | | |
95 | 2 | let store_key = StoreKey::Str(Cow::Owned(format!( |
96 | 2 | "BepEvent:le:{}:{}:{}", |
97 | 2 | &stream_id.build_id, &stream_id.invocation_id, sequence_number, |
98 | 2 | ))); |
99 | | |
100 | 2 | let bep_event = BepEvent { |
101 | 2 | version: BEP_EVENT_VERSION, |
102 | 2 | identity: identity.unwrap_or_default(), |
103 | 2 | event: Some(bep_event::Event::LifecycleEvent(request)), |
104 | 2 | }; |
105 | 2 | let mut buf = BytesMut::new(); |
106 | 2 | bep_event |
107 | 2 | .encode(&mut buf) |
108 | 2 | .err_tip(|| "Could not encode PublishLifecycleEventRequest proto")?0 ; |
109 | | |
110 | 2 | let data = buf.freeze(); |
111 | 4 | for attempt in 1..=MAX_BEP_UPLOAD_ATTEMPTS2 { |
112 | 4 | match self |
113 | 4 | .store |
114 | 4 | .update_oneshot(store_key.borrow(), data.clone()) |
115 | 4 | .await |
116 | | { |
117 | 2 | Ok(()) => break, |
118 | 2 | Err(err) if attempt < MAX_BEP_UPLOAD_ATTEMPTS => { |
119 | 2 | warn!( |
120 | | attempt, |
121 | | max = MAX_BEP_UPLOAD_ATTEMPTS, |
122 | | ?err, |
123 | | %store_key, |
124 | | "Failed to store BEP lifecycle event, retrying" |
125 | | ); |
126 | 2 | sleep(Duration::from_secs_f32(0.1 * attempt as f32)).await; |
127 | | } |
128 | 0 | Err(err) => { |
129 | 0 | return Err(err).err_tip(|| { |
130 | 0 | format!( |
131 | | "Failed to store PublishLifecycleEventRequest for {store_key} after retries" |
132 | | ) |
133 | 0 | }); |
134 | | } |
135 | | } |
136 | | } |
137 | | |
138 | 2 | Ok(Response::new(())) |
139 | 2 | } |
140 | | |
141 | 2 | async fn inner_publish_build_tool_event_stream( |
142 | 2 | &self, |
143 | 2 | stream: Streaming<PublishBuildToolEventStreamRequest>, |
144 | 2 | identity: Option<String>, |
145 | 2 | ) -> Result<Response<PublishBuildToolEventStreamStream>, Error> { |
146 | 6 | async fn process_request( |
147 | 6 | store: Pin<&dyn StoreDriver>, |
148 | 6 | request: PublishBuildToolEventStreamRequest, |
149 | 6 | identity: String, |
150 | 6 | ) -> Result<PublishBuildToolEventStreamResponse, Status> { |
151 | 6 | let ordered_build_event = request |
152 | 6 | .ordered_build_event |
153 | 6 | .as_ref() |
154 | 6 | .err_tip(|| "Expected ordered_build_event to be set")?0 ; |
155 | 6 | let stream_id = ordered_build_event |
156 | 6 | .stream_id |
157 | 6 | .as_ref() |
158 | 6 | .err_tip(|| "Expected stream_id to be set")?0 |
159 | 6 | .clone(); |
160 | | |
161 | 6 | let sequence_number = ordered_build_event.sequence_number; |
162 | | |
163 | 6 | let bep_event = BepEvent { |
164 | 6 | version: BEP_EVENT_VERSION, |
165 | 6 | identity, |
166 | 6 | event: Some(bep_event::Event::BuildToolEvent(request)), |
167 | 6 | }; |
168 | 6 | let mut buf = BytesMut::new(); |
169 | | |
170 | 6 | bep_event |
171 | 6 | .encode(&mut buf) |
172 | 6 | .err_tip(|| "Could not encode PublishBuildToolEventStreamRequest proto")?0 ; |
173 | | |
174 | 6 | let store_key = StoreKey::Str(Cow::Owned(format!( |
175 | 6 | "BepEvent:be:{}:{}:{}", |
176 | 6 | &stream_id.build_id, &stream_id.invocation_id, sequence_number, |
177 | 6 | ))); |
178 | 6 | let data = buf.freeze(); |
179 | 6 | for attempt in 1..=MAX_BEP_UPLOAD_ATTEMPTS { |
180 | 6 | match store.update_oneshot(store_key.borrow(), data.clone()).await { |
181 | 6 | Ok(()) => break, |
182 | 0 | Err(err) if attempt < MAX_BEP_UPLOAD_ATTEMPTS => { |
183 | 0 | warn!( |
184 | | attempt, |
185 | | max = MAX_BEP_UPLOAD_ATTEMPTS, |
186 | | ?err, |
187 | | %store_key, |
188 | | "Failed to store BEP build-tool event, retrying" |
189 | | ); |
190 | 0 | sleep(Duration::from_secs_f32(0.1 * attempt as f32)).await; |
191 | | } |
192 | 0 | Err(err) => { |
193 | 0 | return Err(err).err_tip( |
194 | | || "Failed to store PublishBuildToolEventStreamRequest after retries", |
195 | 0 | )?; |
196 | | } |
197 | | } |
198 | | } |
199 | | |
200 | 6 | Ok(PublishBuildToolEventStreamResponse { |
201 | 6 | stream_id: Some(stream_id.clone()), |
202 | 6 | sequence_number, |
203 | 6 | }) |
204 | 6 | } |
205 | | |
206 | | struct State { |
207 | | store: Store, |
208 | | stream: Streaming<PublishBuildToolEventStreamRequest>, |
209 | | identity: String, |
210 | | } |
211 | | |
212 | 2 | let response_stream = |
213 | 2 | unfold( |
214 | 2 | Some(State { |
215 | 2 | store: self.store.clone(), |
216 | 2 | stream, |
217 | 2 | identity: identity.unwrap_or_default(), |
218 | 2 | }), |
219 | 7 | move |maybe_state| async move { |
220 | 7 | let mut state = maybe_state?0 ; |
221 | 6 | let request = |
222 | 7 | match state.stream.message().await.err_tip( |
223 | | || "While receiving message in publish_build_tool_event_stream", |
224 | | ) { |
225 | 6 | Ok(Some(request)) => request, |
226 | 1 | Ok(None) => return None, |
227 | 0 | Err(e) => return Some((Err(e.into()), None)), |
228 | | }; |
229 | 6 | process_request( |
230 | 6 | state.store.as_store_driver_pin(), |
231 | 6 | request, |
232 | 6 | state.identity.clone(), |
233 | 6 | ) |
234 | 6 | .await |
235 | 6 | .map_or_else( |
236 | 0 | |e| Some((Err(e), None)), |
237 | 6 | |response| Some((Ok(response), Some(state))), |
238 | | ) |
239 | 14 | }, |
240 | | ); |
241 | | |
242 | 2 | Ok(Response::new(Box::pin(response_stream))) |
243 | 2 | } |
244 | | } |
245 | | |
246 | | type PublishBuildToolEventStreamStream = Pin< |
247 | | Box<dyn Stream<Item = Result<PublishBuildToolEventStreamResponse, Status>> + Send + 'static>, |
248 | | >; |
249 | | |
250 | | #[tonic::async_trait] |
251 | | impl PublishBuildEvent for BepServer { |
252 | | type PublishBuildToolEventStreamStream = PublishBuildToolEventStreamStream; |
253 | | |
254 | | #[instrument( |
255 | | err, |
256 | | ret(level = Level::INFO), |
257 | | level = Level::ERROR, |
258 | | skip_all, |
259 | | fields(request = ?grpc_request.get_ref()) |
260 | | )] |
261 | | async fn publish_lifecycle_event( |
262 | | &self, |
263 | | grpc_request: Request<PublishLifecycleEventRequest>, |
264 | | ) -> Result<Response<()>, Status> { |
265 | | self.inner_publish_lifecycle_event(grpc_request.into_inner(), get_identity()) |
266 | | .await |
267 | | .map_err(Error::into) |
268 | | } |
269 | | |
270 | | #[instrument( |
271 | | err, |
272 | | level = Level::ERROR, |
273 | | skip_all, |
274 | | fields(request = ?grpc_request.get_ref()) |
275 | | )] |
276 | | async fn publish_build_tool_event_stream( |
277 | | &self, |
278 | | grpc_request: Request<Streaming<PublishBuildToolEventStreamRequest>>, |
279 | | ) -> Result<Response<Self::PublishBuildToolEventStreamStream>, Status> { |
280 | | self.inner_publish_build_tool_event_stream(grpc_request.into_inner(), get_identity()) |
281 | | .await |
282 | | .map_err(Error::into) |
283 | | } |
284 | | } |