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 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    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
34        Self { req_tx, req }
35    }
36
37    /// Sets the target stream offset to guarantee [exactly-once] writes.
38    ///
39    /// # Example
40    ///
41    /// ```
42    /// # use google_cloud_bigquery::write::format::Arrow;
43    /// # use google_cloud_bigquery::write::PendingWriter;
44    /// # async fn sample(writer: PendingWriter<Arrow>) -> 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::format::Arrow;
72    /// # use google_cloud_bigquery::write::PendingWriter;
73    /// # async fn sample(writer: PendingWriter<Arrow>) -> anyhow::Result<()> {
74    /// let f1 = writer.append(rows()).set_offset(0).send();
75    /// let f2 = writer.append(rows()).set_offset(1).send();
76    ///
77    /// let resp1 = f1.await?;
78    /// let resp2 = f2.await?;
79    /// # Ok(()) }
80    ///
81    /// use google_cloud_bigquery::model::ArrowRecordBatch;
82    /// fn rows() -> ArrowRecordBatch {
83    ///   todo!("Define your rows...")
84    /// }
85    /// ```
86    pub fn send(self) -> AppendFuture {
87        let (resp_tx, resp_rx) = oneshot::channel();
88        let req = match self.req.to_proto().map_err(Error::ser) {
89            Ok(req) => req,
90            Err(e) => return AppendFuture::from_future(async move { Err(e.into()) }),
91        };
92        let write = WriteRequest { req, resp_tx };
93        if self.req_tx.send(write).is_err() {
94            return AppendFuture::from_future(
95                async move { Err(AppendError::UnexpectedEndOfStream) },
96            );
97        }
98        AppendFuture::from_future(async move {
99            let resp = resp_rx
100                .await
101                .map_err(|_| AppendError::UnexpectedEndOfStream)??;
102            let resp = resp.cnv().map_err(Error::deser)?;
103            to_result(resp)
104        })
105    }
106}
107
108#[cfg(test)]
109mod tests {
110    use super::*;
111    use crate::google::cloud::bigquery::storage::v1;
112    use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
113        AppendResult, Response,
114    };
115    use crate::write::test::*;
116
117    #[tokio::test]
118    async fn success() -> anyhow::Result<()> {
119        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
120        let req = AppendRowsRequest::new().set_write_stream(write_stream());
121
122        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
123        let future = builder.send();
124
125        let write = req_rx.recv().await.expect("should receive request");
126        assert_eq!(write.req.offset, Some(100));
127
128        let resp = v1::AppendRowsResponse {
129            response: Some(Response::AppendResult(AppendResult::default())),
130            write_stream: write_stream(),
131            // Ensure schema matches none
132            ..Default::default()
133        };
134        write
135            .resp_tx
136            .send(Ok(resp))
137            .expect("sending on channel always succeeds");
138
139        let resp = future.await?;
140        assert_eq!(resp.offset, None);
141        assert_eq!(resp.updated_schema, None);
142        Ok(())
143    }
144
145    #[tokio::test]
146    async fn stream_closed() -> anyhow::Result<()> {
147        let (req_tx, req_rx) = mpsc::unbounded_channel();
148        let req = AppendRowsRequest::new().set_write_stream(write_stream());
149
150        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
151        let future = builder.send();
152
153        drop(req_rx);
154
155        let err = future.await.expect_err("should return an error");
156        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
157        Ok(())
158    }
159
160    #[tokio::test]
161    async fn rpc_error() -> anyhow::Result<()> {
162        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
163        let req = AppendRowsRequest::new().set_write_stream(write_stream());
164
165        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
166        let future = builder.send();
167
168        // Simulate a stream ending in a known error
169        let write = req_rx.recv().await.expect("should receive request");
170        let append_err: AppendError = Error::io("fail").into();
171        write
172            .resp_tx
173            .send(Err(append_err))
174            .expect("sending on channel always succeeds");
175
176        let err = future.await.expect_err("should return an error");
177        assert!(matches!(err, AppendError::Rpc { source: _ }));
178        Ok(())
179    }
180
181    #[tokio::test]
182    async fn row_errors() -> anyhow::Result<()> {
183        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
184        let req = AppendRowsRequest::new().set_write_stream(write_stream());
185
186        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
187        let future = builder.send();
188
189        let write = req_rx.recv().await.expect("should receive request");
190
191        let row_error = v1::RowError {
192            index: 42,
193            code: v1::row_error::RowErrorCode::FieldsError as i32,
194            message: "fail".to_string(),
195        };
196        let resp = v1::AppendRowsResponse {
197            row_errors: vec![row_error],
198            write_stream: write_stream(),
199            ..Default::default()
200        };
201        write
202            .resp_tx
203            .send(Ok(resp))
204            .expect("sending on channel always succeeds");
205
206        let err = future.await.expect_err("should return an error");
207        assert!(matches!(err, AppendError::RowErrors { .. }));
208        Ok(())
209    }
210
211    #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
212    async fn synchronous_queueing() -> anyhow::Result<()> {
213        const NUM_WRITES: i64 = 1000;
214        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
215        let write_handle = tokio::spawn(async move {
216            let mut writes = tokio::task::JoinSet::new();
217            for i in 0..NUM_WRITES {
218                writes.spawn(
219                    AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
220                        .set_offset(i)
221                        .send(),
222                );
223            }
224            let _ = writes.join_all().await;
225        });
226
227        for i in 0..NUM_WRITES {
228            let write = req_rx.recv().await.expect("should receive request");
229            assert_eq!(write.req.offset, Some(i), "received out of order write");
230        }
231        write_handle.await?;
232        Ok(())
233    }
234
235    #[tokio::test]
236    async fn send_when_req_tx_closed_returns_unexpected_end_of_stream() {
237        let (req_tx, req_rx) = mpsc::unbounded_channel();
238        drop(req_rx);
239
240        let req = AppendRowsRequest::new().set_write_stream(write_stream());
241        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
242        let future = builder.send();
243
244        let err = future
245            .await
246            .expect_err("should return unexpected end of stream");
247        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
248    }
249
250    #[tokio::test]
251    async fn send_serialization_error() {
252        use crate::model::append_rows_request::MissingValueInterpretation;
253
254        let (req_tx, _req_rx) = mpsc::unbounded_channel();
255        let invalid: MissingValueInterpretation = serde_json::from_str("\"INVALID\"").unwrap();
256        let req = AppendRowsRequest::new().set_default_missing_value_interpretation(invalid);
257        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
258        let future = builder.send();
259
260        let err = future.await.expect_err("should return serialization error");
261        assert!(matches!(err, AppendError::Rpc { .. }));
262    }
263}