Skip to main content

google_cloud_bigquery/write/builder/
append.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::dispatcher::Dispatcher;
17use crate::model::AppendRowsRequest;
18use std::sync::Arc;
19use tokio::sync::oneshot;
20
21/// A request builder for appending rows to the default stream.
22#[derive(Clone, Debug)]
23pub struct Append {
24    inner: Arc<Dispatcher>,
25    pub(crate) req: AppendRowsRequest,
26}
27
28impl Append {
29    pub(crate) fn new(inner: Arc<Dispatcher>, req: AppendRowsRequest) -> Self {
30        Self { inner, req }
31    }
32
33    /// Append rows to the stream.
34    ///
35    /// Applications are encouraged to queue up requests and await their
36    /// responses independently.
37    ///
38    /// # Example
39    ///
40    /// ```
41    /// # use google_cloud_bigquery::write::format::Arrow;
42    /// # use google_cloud_bigquery::write::DefaultWriter;
43    /// # async fn sample(writer: DefaultWriter<Arrow>) -> anyhow::Result<()> {
44    /// let f1 = writer.append(rows()).send();
45    /// let f2 = writer.append(rows()).send();
46    ///
47    /// let resp1 = f1.await?;
48    /// let resp2 = f2.await?;
49    /// # Ok(()) }
50    ///
51    /// use google_cloud_bigquery::model::ArrowRecordBatch;
52    /// fn rows() -> ArrowRecordBatch {
53    ///   todo!("Define your rows...")
54    /// }
55    /// ```
56    pub fn send(self) -> AppendFuture {
57        let (tx, rx) = oneshot::channel();
58        tokio::spawn(async move {
59            let res = self.inner.send(self.req).await;
60            let _ = tx.send(res);
61        });
62        AppendFuture::new(rx)
63    }
64}
65
66#[cfg(test)]
67mod tests {
68    use super::*;
69    use crate::Error;
70    use crate::error::AppendError;
71    use crate::google::cloud::bigquery::storage::v1;
72    use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
73        AppendResult, Response,
74    };
75    use crate::model::TableSchema;
76    use crate::write::test::*;
77    use tokio::sync::mpsc;
78
79    #[tokio::test]
80    async fn success() -> anyhow::Result<()> {
81        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
82        let dispatcher = test_dispatcher(req_tx).await?;
83        let req = AppendRowsRequest::new().set_write_stream(write_stream());
84
85        let builder = Append::new(dispatcher, req);
86        let handle = tokio::spawn(async move { builder.send().await });
87
88        // Receive and verify the request
89        let write = req_rx.recv().await.expect("should receive request");
90        assert_eq!(write.req.write_stream, write_stream());
91
92        // Provide a successful response
93        let resp = v1::AppendRowsResponse {
94            response: Some(Response::AppendResult(AppendResult::default())),
95            write_stream: write_stream(),
96            updated_schema: Some(v1::TableSchema::default()),
97            ..Default::default()
98        };
99        write
100            .resp_tx
101            .send(Ok(resp))
102            .expect("sending on channel always succeeds");
103
104        let resp = handle.await??;
105        assert_eq!(resp.offset, None);
106        assert_eq!(resp.updated_schema, Some(TableSchema::default()));
107        Ok(())
108    }
109
110    #[tokio::test]
111    async fn stream_closed() -> anyhow::Result<()> {
112        let (req_tx, req_rx) = mpsc::unbounded_channel();
113        let dispatcher = test_dispatcher(req_tx).await?;
114        let req = AppendRowsRequest::new().set_write_stream(write_stream());
115
116        let builder = Append::new(dispatcher, req);
117        let handle = tokio::spawn(async move { builder.send().await });
118
119        // Simulate a stream closure
120        drop(req_rx);
121
122        let err = handle.await?.expect_err("should return an error");
123        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
124        Ok(())
125    }
126
127    #[tokio::test]
128    async fn rpc_error() -> anyhow::Result<()> {
129        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
130        let dispatcher = test_dispatcher(req_tx).await?;
131        let req = AppendRowsRequest::new().set_write_stream(write_stream());
132
133        let builder = Append::new(dispatcher, req);
134        let handle = tokio::spawn(async move { builder.send().await });
135
136        // Simulate a stream ending in a known error
137        let write = req_rx.recv().await.expect("should receive request");
138        let append_err: AppendError = Error::io("fail").into();
139        write
140            .resp_tx
141            .send(Err(append_err))
142            .expect("sending on channel always succeeds");
143
144        let err = handle.await?.expect_err("should return an error");
145        assert!(matches!(err, AppendError::Rpc { source: _ }));
146        Ok(())
147    }
148
149    #[tokio::test]
150    async fn row_errors() -> anyhow::Result<()> {
151        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
152        let dispatcher = test_dispatcher(req_tx).await?;
153        let req = AppendRowsRequest::new().set_write_stream(write_stream());
154
155        let builder = Append::new(dispatcher, req);
156        let handle = tokio::spawn(async move { builder.send().await });
157
158        let write = req_rx.recv().await.expect("should receive request");
159
160        let row_error = v1::RowError {
161            index: 42,
162            code: v1::row_error::RowErrorCode::FieldsError as i32,
163            message: "fail".to_string(),
164        };
165        let resp = v1::AppendRowsResponse {
166            row_errors: vec![row_error],
167            write_stream: write_stream(),
168            ..Default::default()
169        };
170        write
171            .resp_tx
172            .send(Ok(resp))
173            .expect("sending on channel always succeeds");
174
175        let err = handle.await?.expect_err("should return an error");
176        assert!(matches!(err, AppendError::RowErrors { .. }));
177        Ok(())
178    }
179}