Skip to main content

google_cloud_bigquery/write/
append_builder.rs

1// Copyright 2026 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (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//     https://www.apache.org/licenses/LICENSE-2.0
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
15use super::append_future::AppendFuture;
16use super::append_response::to_result;
17use super::error::AppendError;
18use super::runner::WriteRequest;
19use crate::Error;
20use crate::model::AppendRowsRequest;
21use gaxi::prost::{FromProto, ToProto};
22use tokio::sync::{mpsc, oneshot};
23
24/// A request builder for appending rows with a specific stream offset,
25/// ensuring exactly-once semantics.
26#[derive(Clone, Debug)]
27pub struct AppendWithOffset {
28    req_tx: mpsc::UnboundedSender<WriteRequest>,
29    pub(crate) req: AppendRowsRequest,
30}
31
32impl AppendWithOffset {
33    #[allow(dead_code)]
34    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
35        Self { req_tx, req }
36    }
37
38    /// Sets the target stream offset to guarantee [exactly-once] writes.
39    ///
40    /// # Example
41    ///
42    /// ```
43    /// # use google_cloud_bigquery::write::arrow::PendingWriter;
44    /// # async fn sample(writer: PendingWriter) -> anyhow::Result<()> {
45    /// let resp = writer.append(rows()).set_offset(0).send().await?;
46    /// # Ok(()) }
47    ///
48    /// use google_cloud_bigquery::model::ArrowRecordBatch;
49    /// fn rows() -> ArrowRecordBatch {
50    ///   todo!("Define your rows...")
51    /// }
52    /// ```
53    ///
54    /// [exactly-once]: https://docs.cloud.google.com/bigquery/docs/write-api-best-practices#manage_stream_offsets_to_achieve_exactly-once_semantics
55    pub fn set_offset(mut self, offset: i64) -> Self {
56        self.req.offset = Some(offset);
57        self
58    }
59
60    /// Append rows to the stream.
61    ///
62    /// Applications are encouraged to queue up requests and await their
63    /// responses independently.
64    ///
65    /// Note that the service will reject requests with a mismatched offset, so
66    /// requests must be queued in order.
67    ///
68    /// # Example
69    ///
70    /// ```
71    /// # use google_cloud_bigquery::write::arrow::PendingWriter;
72    /// # async fn sample(writer: PendingWriter) -> anyhow::Result<()> {
73    /// let f1 = writer.append(rows()).set_offset(0).send();
74    /// let f2 = writer.append(rows()).set_offset(1).send();
75    ///
76    /// let resp1 = f1.await?;
77    /// let resp2 = f2.await?;
78    /// # Ok(()) }
79    ///
80    /// use google_cloud_bigquery::model::ArrowRecordBatch;
81    /// fn rows() -> ArrowRecordBatch {
82    ///   todo!("Define your rows...")
83    /// }
84    /// ```
85    pub fn send(self) -> AppendFuture {
86        let (tx, rx) = oneshot::channel();
87        let (resp_tx, resp_rx) = oneshot::channel();
88        let req = match self.req.to_proto().map_err(Error::deser) {
89            Ok(req) => req,
90            Err(e) => {
91                let _ = tx.send(Err(e.into()));
92                return AppendFuture::new(rx);
93            }
94        };
95        let write = WriteRequest { req, resp_tx };
96        let _ = self.req_tx.send(write);
97        tokio::spawn(async move {
98            let res = async {
99                let resp = resp_rx
100                    .await
101                    .map_err(|_| AppendError::UnexpectedEndOfStream)??;
102                let resp = resp.cnv().map_err(Error::ser)?;
103                to_result(resp)
104            }
105            .await;
106            let _ = tx.send(res);
107        });
108        AppendFuture::new(rx)
109    }
110}
111
112/// A request builder for appending rows on the default stream.
113#[derive(Clone, Debug)]
114pub struct Append {
115    req_tx: mpsc::UnboundedSender<WriteRequest>,
116    pub(crate) req: AppendRowsRequest,
117}
118
119impl Append {
120    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
121        Self { req_tx, req }
122    }
123
124    /// Append rows to the stream.
125    ///
126    /// Applications are encouraged to queue up requests and await their
127    /// responses independently.
128    ///
129    /// # Example
130    ///
131    /// ```
132    /// # use google_cloud_bigquery::write::arrow::DefaultWriter;
133    /// # async fn sample(writer: DefaultWriter) -> anyhow::Result<()> {
134    /// let f1 = writer.append(rows()).send();
135    /// let f2 = writer.append(rows()).send();
136    ///
137    /// let resp1 = f1.await?;
138    /// let resp2 = f2.await?;
139    /// # Ok(()) }
140    ///
141    /// use google_cloud_bigquery::model::ArrowRecordBatch;
142    /// fn rows() -> ArrowRecordBatch {
143    ///   todo!("Define your rows...")
144    /// }
145    /// ```
146    pub fn send(self) -> AppendFuture {
147        let (tx, rx) = oneshot::channel();
148        tokio::spawn(async move {
149            let (resp_tx, resp_rx) = oneshot::channel();
150            let res = async move {
151                let req = self.req.to_proto().map_err(Error::deser)?;
152                let write = WriteRequest { req, resp_tx };
153                let _ = self.req_tx.send(write);
154                let resp = resp_rx
155                    .await
156                    .map_err(|_| AppendError::UnexpectedEndOfStream)??;
157                let resp = resp.cnv().map_err(Error::ser)?;
158                to_result(resp)
159            }
160            .await;
161            let _ = tx.send(res);
162        });
163        AppendFuture::new(rx)
164    }
165}
166
167#[cfg(test)]
168mod tests {
169    use super::*;
170    use crate::google::cloud::bigquery::storage::v1;
171    use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
172        AppendResult, Response,
173    };
174    use crate::model::TableSchema;
175
176    #[tokio::test]
177    async fn success() -> anyhow::Result<()> {
178        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
179        let req = AppendRowsRequest::new().set_write_stream(write_stream());
180
181        let builder = Append::new(req_tx, req);
182        let handle = tokio::spawn(async move { builder.send().await });
183
184        // Receive and verify the request
185        let write = req_rx.recv().await.expect("should receive request");
186        assert_eq!(write.req.write_stream, write_stream());
187
188        // Provide a successful response
189        let resp = v1::AppendRowsResponse {
190            response: Some(Response::AppendResult(AppendResult::default())),
191            write_stream: write_stream(),
192            updated_schema: Some(v1::TableSchema::default()),
193            ..Default::default()
194        };
195        write
196            .resp_tx
197            .send(Ok(resp))
198            .expect("sending on channel always succeeds");
199
200        let resp = handle.await??;
201        assert_eq!(resp.offset, None);
202        assert_eq!(resp.updated_schema, Some(TableSchema::default()));
203        Ok(())
204    }
205
206    #[tokio::test]
207    async fn stream_closed() -> anyhow::Result<()> {
208        let (req_tx, req_rx) = mpsc::unbounded_channel();
209        let req = AppendRowsRequest::new().set_write_stream(write_stream());
210
211        let builder = Append::new(req_tx, req);
212        let handle = tokio::spawn(async move { builder.send().await });
213
214        // Simulate a stream closure
215        drop(req_rx);
216
217        let err = handle.await?.expect_err("should return an error");
218        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
219        Ok(())
220    }
221
222    #[tokio::test]
223    async fn rpc_error() -> anyhow::Result<()> {
224        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
225        let req = AppendRowsRequest::new().set_write_stream(write_stream());
226
227        let builder = Append::new(req_tx, req);
228        let handle = tokio::spawn(async move { builder.send().await });
229
230        // Simulate a stream ending in a known error
231        let write = req_rx.recv().await.expect("should receive request");
232        let append_err: AppendError = Error::io("fail").into();
233        write
234            .resp_tx
235            .send(Err(append_err))
236            .expect("sending on channel always succeeds");
237
238        let err = handle.await?.expect_err("should return an error");
239        assert!(matches!(err, AppendError::Rpc { source: _ }));
240        Ok(())
241    }
242
243    #[tokio::test]
244    async fn row_errors() -> anyhow::Result<()> {
245        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
246        let req = AppendRowsRequest::new().set_write_stream(write_stream());
247
248        let builder = Append::new(req_tx, req);
249        let handle = tokio::spawn(async move { builder.send().await });
250
251        let write = req_rx.recv().await.expect("should receive request");
252
253        let row_error = v1::RowError {
254            index: 42,
255            code: v1::row_error::RowErrorCode::FieldsError as i32,
256            message: "fail".to_string(),
257        };
258        let resp = v1::AppendRowsResponse {
259            row_errors: vec![row_error],
260            write_stream: write_stream(),
261            ..Default::default()
262        };
263        write
264            .resp_tx
265            .send(Ok(resp))
266            .expect("sending on channel always succeeds");
267
268        let err = handle.await?.expect_err("should return an error");
269        assert!(matches!(err, AppendError::RowErrors(_)));
270        Ok(())
271    }
272
273    #[tokio::test]
274    async fn offset_success() -> anyhow::Result<()> {
275        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
276        let req = AppendRowsRequest::new().set_write_stream(write_stream());
277
278        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
279        let future = builder.send();
280
281        let write = req_rx.recv().await.expect("should receive request");
282        assert_eq!(write.req.offset, Some(100));
283
284        let resp = v1::AppendRowsResponse {
285            response: Some(Response::AppendResult(AppendResult::default())),
286            write_stream: write_stream(),
287            // Ensure schema matches none
288            ..Default::default()
289        };
290        write
291            .resp_tx
292            .send(Ok(resp))
293            .expect("sending on channel always succeeds");
294
295        let resp = future.await?;
296        assert_eq!(resp.offset, None);
297        assert_eq!(resp.updated_schema, None);
298        Ok(())
299    }
300
301    #[tokio::test]
302    async fn offset_stream_closed() -> anyhow::Result<()> {
303        let (req_tx, req_rx) = mpsc::unbounded_channel();
304        let req = AppendRowsRequest::new().set_write_stream(write_stream());
305
306        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
307        let future = builder.send();
308
309        drop(req_rx);
310
311        let err = future.await.expect_err("should return an error");
312        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
313        Ok(())
314    }
315
316    #[tokio::test]
317    async fn offset_rpc_error() -> anyhow::Result<()> {
318        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
319        let req = AppendRowsRequest::new().set_write_stream(write_stream());
320
321        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
322        let future = builder.send();
323
324        // Simulate a stream ending in a known error
325        let write = req_rx.recv().await.expect("should receive request");
326        let append_err: AppendError = Error::io("fail").into();
327        write
328            .resp_tx
329            .send(Err(append_err))
330            .expect("sending on channel always succeeds");
331
332        let err = future.await.expect_err("should return an error");
333        assert!(matches!(err, AppendError::Rpc { source: _ }));
334        Ok(())
335    }
336
337    #[tokio::test]
338    async fn offset_row_errors() -> anyhow::Result<()> {
339        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
340        let req = AppendRowsRequest::new().set_write_stream(write_stream());
341
342        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
343        let future = builder.send();
344
345        let write = req_rx.recv().await.expect("should receive request");
346
347        let row_error = v1::RowError {
348            index: 42,
349            code: v1::row_error::RowErrorCode::FieldsError as i32,
350            message: "fail".to_string(),
351        };
352        let resp = v1::AppendRowsResponse {
353            row_errors: vec![row_error],
354            write_stream: write_stream(),
355            ..Default::default()
356        };
357        write
358            .resp_tx
359            .send(Ok(resp))
360            .expect("sending on channel always succeeds");
361
362        let err = future.await.expect_err("should return an error");
363        assert!(matches!(err, AppendError::RowErrors(_)));
364        Ok(())
365    }
366
367    #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
368    async fn synchronous_queueing() -> anyhow::Result<()> {
369        const NUM_WRITES: i64 = 1000;
370        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
371        let write_handle = tokio::spawn(async move {
372            let mut writes = tokio::task::JoinSet::new();
373            for i in 0..NUM_WRITES {
374                writes.spawn(
375                    AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
376                        .set_offset(i)
377                        .send(),
378                );
379            }
380            let _ = writes.join_all().await;
381        });
382
383        for i in 0..NUM_WRITES {
384            let write = req_rx.recv().await.expect("should receive request");
385            assert_eq!(write.req.offset, Some(i), "received out of order write");
386        }
387        write_handle.await?;
388        Ok(())
389    }
390
391    fn write_stream() -> String {
392        "projects/p/datasets/d/tables/t/streams/_default".to_string()
393    }
394}