Coverage Report

Created: 2026-07-21 15:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}