Skip to main content

google_cloud_bigquery/write/
buffered.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::base::BaseWriter;
16use super::format::DataFormat;
17use crate::Result;
18use crate::model::{FinalizeWriteStreamResponse, FlushRowsResponse};
19use crate::write::builder::AppendWithOffset;
20use crate::write::transport::Transport;
21use std::sync::Arc;
22
23/// A writer for a [buffered stream].
24///
25/// [buffered stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#buffered_type
26#[derive(Debug)]
27pub struct BufferedWriter<F> {
28    pub(crate) inner: BaseWriter<F>,
29}
30
31impl<F> BufferedWriter<F>
32where
33    F: DataFormat,
34{
35    pub(crate) fn new(inner: Arc<Transport>, write_stream: String, format: F) -> Self {
36        Self {
37            inner: BaseWriter::new(inner, write_stream, format),
38        }
39    }
40
41    /// Return the full resource name of the underlying write stream.
42    pub fn write_stream(&self) -> &str {
43        &self.inner.write_stream
44    }
45
46    /// Append rows to the buffered stream.
47    pub fn append(&self, rows: F::Rows) -> AppendWithOffset {
48        AppendWithOffset::new(
49            self.inner.runner.req_tx.clone(),
50            self.inner.append_request(rows),
51        )
52    }
53
54    /// Flush the buffered stream, making rows up to the specified offset available for reading.
55    pub async fn flush(&self, offset: i64) -> Result<FlushRowsResponse> {
56        self.inner
57            .client
58            .flush_rows()
59            .set_write_stream(&self.inner.write_stream)
60            .set_offset(offset)
61            .send()
62            .await
63    }
64
65    /// Finalize the buffered stream, preventing further writes.
66    pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
67        self.inner.finalize().await
68    }
69}
70
71#[cfg(test)]
72mod tests {
73    use super::*;
74    use crate::error::AppendError;
75    use crate::write::test::*;
76    use bigquery_grpc_mock::{MockBigQueryWrite, start};
77    use gaxi::grpc::tonic::Response as TonicResponse;
78    use tokio::sync::mpsc;
79
80    #[tokio::test]
81    async fn basic_success() -> anyhow::Result<()> {
82        let (response_tx, response_rx) = mpsc::channel(10);
83
84        let mut mock = MockBigQueryWrite::new();
85        mock.expect_append_rows()
86            .return_once(|_| Ok(TonicResponse::from(response_rx)));
87
88        mock.expect_flush_rows()
89            .return_once(|req| {
90                assert_eq!(req.get_ref().offset, Some(3));
91                assert_eq!(req.get_ref().write_stream, write_stream());
92                Ok(TonicResponse::new(
93                    bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse::default()
94                ))
95            });
96
97        mock.expect_finalize_write_stream()
98            .return_once(|req| {
99                assert_eq!(req.get_ref().name, write_stream());
100                Ok(TonicResponse::new(
101                    bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default()
102                ))
103            });
104
105        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
106        let transport = Arc::new(test_transport(endpoint).await?);
107
108        let writer = BufferedWriter::new(transport, write_stream(), format());
109        assert_eq!(writer.write_stream(), write_stream());
110
111        response_tx.send(Ok(convert(&test_response(1)))).await?;
112        let resp = writer.append(rows(1)).send().await?;
113        assert_eq!(resp.offset, Some(1));
114
115        response_tx.send(Ok(convert(&test_response(2)))).await?;
116        let resp = writer.append(rows(2)).send().await?;
117        assert_eq!(resp.offset, Some(2));
118
119        response_tx.send(Ok(convert(&test_response(3)))).await?;
120        let resp = writer.append(rows(3)).send().await?;
121        assert_eq!(resp.offset, Some(3));
122
123        drop(response_tx);
124        let err = writer.append(rows(4)).send().await.expect_err("channel");
125        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
126
127        writer.flush(3).await?;
128        writer.finalize().await?;
129
130        Ok(())
131    }
132
133    #[tokio::test]
134    async fn multiple_flushes() -> anyhow::Result<()> {
135        let (response_tx, response_rx) = mpsc::channel(10);
136        let mut mock = MockBigQueryWrite::new();
137        mock.expect_append_rows()
138            .return_once(|_| Ok(TonicResponse::from(response_rx)));
139
140        mock.expect_flush_rows().times(2).returning(|req| {
141            Ok(TonicResponse::new(
142                bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse {
143                    offset: req.get_ref().offset.unwrap_or(0),
144                },
145            ))
146        });
147
148        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
149        let transport = Arc::new(test_transport(endpoint).await?);
150        let writer = BufferedWriter::new(transport, write_stream(), format());
151        assert_eq!(writer.write_stream(), write_stream());
152
153        response_tx.send(Ok(convert(&test_response(1)))).await?;
154        let _ = writer.append(rows(1)).send().await?;
155        let flush1 = writer.flush(1).await?;
156        assert_eq!(flush1.offset, 1);
157
158        response_tx.send(Ok(convert(&test_response(2)))).await?;
159        let _ = writer.append(rows(2)).send().await?;
160        let flush2 = writer.flush(2).await?;
161        assert_eq!(flush2.offset, 2);
162
163        Ok(())
164    }
165}