Skip to main content

google_cloud_bigquery/write/
default.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::builder::Append;
16use super::dispatcher::Dispatcher;
17use super::format::DataFormat;
18use super::pool::StreamPool;
19use super::retry_policy::RetryOptions;
20use std::sync::Arc;
21
22/// A writer for the [default stream].
23///
24/// The default stream is designed for streaming scenarios where you have
25/// continuously arriving data. It has the following characteristics:
26///
27/// - Data written to the default stream is available immediately for query.
28/// - The default stream supports at-least-once semantics.
29/// - You don't need to explicitly create the default stream.
30///
31/// [default stream]: https://docs.cloud.google.com/bigquery/docs/write-api#default_stream
32#[derive(Debug)]
33pub struct DefaultWriter<F> {
34    pub(crate) inner: Arc<Dispatcher>,
35    pub(crate) write_stream: String,
36    pub(crate) format: F,
37}
38
39impl<F> DefaultWriter<F>
40where
41    F: DataFormat,
42{
43    pub(crate) fn new(
44        pool: Arc<StreamPool>,
45        retry_options: RetryOptions,
46        write_stream: String,
47        format: F,
48    ) -> Self {
49        let inner = Arc::new(Dispatcher::new(pool, retry_options));
50        Self {
51            inner,
52            write_stream,
53            format,
54        }
55    }
56
57    /// Appends rows to the stream.
58    pub fn append(&self, rows: F::Rows) -> Append {
59        let req = self.format.make_request(&self.write_stream, rows);
60        Append::new(self.inner.clone(), req)
61    }
62}
63
64#[cfg(test)]
65mod tests {
66    use super::super::pool::StreamPoolOptions;
67    use super::*;
68    use crate::error::AppendError;
69    use crate::write::test::*;
70    use bigquery_grpc_mock::{MockBigQueryWrite, start};
71    use gaxi::grpc::tonic::{Response as TonicResponse, Status as TonicStatus};
72    use tokio::sync::mpsc;
73
74    #[tokio::test]
75    async fn basic_success() -> anyhow::Result<()> {
76        let (response_tx, response_rx) = mpsc::channel(10);
77
78        let mut mock = MockBigQueryWrite::new();
79        mock.expect_append_rows()
80            .return_once(|_| Ok(TonicResponse::from(response_rx)));
81        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
82        let transport = Arc::new(test_transport(endpoint).await?);
83        let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default()));
84
85        let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), format());
86
87        response_tx.send(Ok(convert(&test_response(1)))).await?;
88        let resp = writer.append(rows(1)).send().await?;
89        assert_eq!(resp.offset, Some(1));
90
91        response_tx.send(Ok(convert(&test_response(2)))).await?;
92        let resp = writer.append(rows(2)).send().await?;
93        assert_eq!(resp.offset, Some(2));
94
95        response_tx.send(Ok(convert(&test_response(3)))).await?;
96        let resp = writer.append(rows(3)).send().await?;
97        assert_eq!(resp.offset, Some(3));
98
99        response_tx
100            .send(Err(TonicStatus::failed_precondition("fail")))
101            .await?;
102        let err = writer.append(rows(4)).send().await.expect_err("fail");
103        assert!(matches!(err, AppendError::Rpc { source: _ }), "{err:?}");
104
105        Ok(())
106    }
107}