Skip to main content

google_cloud_bigquery/write/builder/
append_with_offset.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::super::append_future::AppendFuture;
16use super::super::append_response::to_result;
17use super::super::error::AppendError;
18use super::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 to an [application-created stream],
25/// with an optional stream offset to ensure [exactly-once] semantics.
26///
27/// [application-created stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#application-created_streams
28/// [exactly-once]: https://docs.cloud.google.com/bigquery/docs/write-api-best-practices#manage_stream_offsets_to_achieve_exactly-once_semantics
29#[derive(Clone, Debug)]
30pub struct AppendWithOffset {
31    req_tx: mpsc::UnboundedSender<WriteRequest>,
32    pub(crate) req: AppendRowsRequest,
33}
34
35impl AppendWithOffset {
36    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
37        Self { req_tx, req }
38    }
39
40    /// Sets the target stream offset to guarantee [exactly-once] writes.
41    ///
42    /// The offset is a 0-based cumulative row index in the stream. For example,
43    /// after appending a batch of 10 rows at offset `0`, the next batch should
44    /// be appended at offset `10`.
45    ///
46    /// # Example
47    ///
48    /// ```
49    /// # use google_cloud_bigquery::write::format::Arrow;
50    /// # use google_cloud_bigquery::write::PendingWriter;
51    /// # async fn sample(writer: PendingWriter<Arrow>) -> anyhow::Result<()> {
52    /// let resp = writer.append(rows()).set_offset(0).send().await?;
53    /// # Ok(()) }
54    ///
55    /// use google_cloud_bigquery::model::ArrowRecordBatch;
56    /// fn rows() -> ArrowRecordBatch {
57    ///   todo!("Define your rows...")
58    /// }
59    /// ```
60    ///
61    /// [exactly-once]: https://docs.cloud.google.com/bigquery/docs/write-api-best-practices#manage_stream_offsets_to_achieve_exactly-once_semantics
62    pub fn set_offset(mut self, offset: i64) -> Self {
63        self.req.offset = Some(offset);
64        self
65    }
66
67    /// Append rows to the stream.
68    ///
69    /// Applications are encouraged to queue up requests and await their
70    /// responses independently.
71    ///
72    /// Note that the service will reject requests with a mismatched offset, so
73    /// requests must be queued in order.
74    ///
75    /// # Example
76    ///
77    /// ```
78    /// # use google_cloud_bigquery::write::format::Arrow;
79    /// # use google_cloud_bigquery::write::PendingWriter;
80    /// # async fn sample(writer: PendingWriter<Arrow>) -> anyhow::Result<()> {
81    /// let f1 = writer.append(ten_rows()).set_offset(0).send();
82    /// let f2 = writer.append(ten_rows()).set_offset(10).send();
83    ///
84    /// let resp1 = f1.await?;
85    /// let resp2 = f2.await?;
86    /// # Ok(()) }
87    ///
88    /// use google_cloud_bigquery::model::ArrowRecordBatch;
89    /// fn ten_rows() -> ArrowRecordBatch {
90    ///   todo!("Define 10 rows...")
91    /// }
92    /// ```
93    pub fn send(self) -> AppendFuture {
94        let (resp_tx, resp_rx) = oneshot::channel();
95        let req = match self.req.to_proto().map_err(Error::ser) {
96            Ok(req) => req,
97            Err(e) => return AppendFuture::from_future(async move { Err(e.into()) }),
98        };
99        let write = WriteRequest { req, resp_tx };
100        if self.req_tx.send(write).is_err() {
101            return AppendFuture::from_future(
102                async move { Err(AppendError::UnexpectedEndOfStream) },
103            );
104        }
105        AppendFuture::from_future(async move {
106            let resp = resp_rx
107                .await
108                .map_err(|_| AppendError::UnexpectedEndOfStream)??;
109            let resp = resp.cnv().map_err(Error::deser)?;
110            to_result(resp)
111        })
112    }
113}
114
115#[cfg(test)]
116mod tests {
117    use super::*;
118    use crate::google::cloud::bigquery::storage::v1;
119    use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
120        AppendResult, Response,
121    };
122    use crate::write::test::*;
123
124    #[tokio::test]
125    async fn success() -> anyhow::Result<()> {
126        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
127        let req = AppendRowsRequest::new().set_write_stream(write_stream());
128
129        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
130        let future = builder.send();
131
132        let write = req_rx.recv().await.expect("should receive request");
133        assert_eq!(write.req.offset, Some(100));
134
135        let resp = v1::AppendRowsResponse {
136            response: Some(Response::AppendResult(AppendResult::default())),
137            write_stream: write_stream(),
138            // Ensure schema matches none
139            ..Default::default()
140        };
141        write
142            .resp_tx
143            .send(Ok(resp))
144            .expect("sending on channel always succeeds");
145
146        let resp = future.await?;
147        assert_eq!(resp.offset, None);
148        assert_eq!(resp.updated_schema, None);
149        Ok(())
150    }
151
152    #[tokio::test]
153    async fn stream_closed() -> anyhow::Result<()> {
154        let (req_tx, req_rx) = mpsc::unbounded_channel();
155        let req = AppendRowsRequest::new().set_write_stream(write_stream());
156
157        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
158        let future = builder.send();
159
160        drop(req_rx);
161
162        let err = future.await.expect_err("should return an error");
163        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
164        Ok(())
165    }
166
167    #[tokio::test]
168    async fn rpc_error() -> anyhow::Result<()> {
169        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
170        let req = AppendRowsRequest::new().set_write_stream(write_stream());
171
172        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
173        let future = builder.send();
174
175        // Simulate a stream ending in a known error
176        let write = req_rx.recv().await.expect("should receive request");
177        let append_err: AppendError = Error::io("fail").into();
178        write
179            .resp_tx
180            .send(Err(append_err))
181            .expect("sending on channel always succeeds");
182
183        let err = future.await.expect_err("should return an error");
184        assert!(matches!(err, AppendError::Rpc { source: _ }));
185        Ok(())
186    }
187
188    #[tokio::test]
189    async fn row_errors() -> anyhow::Result<()> {
190        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
191        let req = AppendRowsRequest::new().set_write_stream(write_stream());
192
193        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
194        let future = builder.send();
195
196        let write = req_rx.recv().await.expect("should receive request");
197
198        let row_error = v1::RowError {
199            index: 42,
200            code: v1::row_error::RowErrorCode::FieldsError as i32,
201            message: "fail".to_string(),
202        };
203        let resp = v1::AppendRowsResponse {
204            row_errors: vec![row_error],
205            write_stream: write_stream(),
206            ..Default::default()
207        };
208        write
209            .resp_tx
210            .send(Ok(resp))
211            .expect("sending on channel always succeeds");
212
213        let err = future.await.expect_err("should return an error");
214        assert!(matches!(err, AppendError::RowErrors { .. }));
215        Ok(())
216    }
217
218    #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
219    async fn synchronous_queueing() -> anyhow::Result<()> {
220        const NUM_WRITES: i64 = 1000;
221        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
222        let write_handle = tokio::spawn(async move {
223            let mut writes = tokio::task::JoinSet::new();
224            for i in 0..NUM_WRITES {
225                writes.spawn(
226                    AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
227                        .set_offset(i)
228                        .send(),
229                );
230            }
231            let _ = writes.join_all().await;
232        });
233
234        for i in 0..NUM_WRITES {
235            let write = req_rx.recv().await.expect("should receive request");
236            assert_eq!(write.req.offset, Some(i), "received out of order write");
237        }
238        write_handle.await?;
239        Ok(())
240    }
241
242    #[tokio::test]
243    async fn send_when_req_tx_closed_returns_unexpected_end_of_stream() {
244        let (req_tx, req_rx) = mpsc::unbounded_channel();
245        drop(req_rx);
246
247        let req = AppendRowsRequest::new().set_write_stream(write_stream());
248        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
249        let future = builder.send();
250
251        let err = future
252            .await
253            .expect_err("should return unexpected end of stream");
254        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
255    }
256
257    #[tokio::test]
258    async fn send_serialization_error() {
259        use crate::model::append_rows_request::MissingValueInterpretation;
260
261        let (req_tx, _req_rx) = mpsc::unbounded_channel();
262        let invalid: MissingValueInterpretation = serde_json::from_str("\"INVALID\"").unwrap();
263        let req = AppendRowsRequest::new().set_default_missing_value_interpretation(invalid);
264        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
265        let future = builder.send();
266
267        let err = future.await.expect_err("should return serialization error");
268        assert!(matches!(err, AppendError::Rpc { .. }));
269    }
270}