Skip to main content

google_cloud_bigquery/write/
committed.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;
19use crate::write::builder::AppendWithOffset;
20use crate::write::transport::Transport;
21use std::sync::Arc;
22
23/// A writer for a [committed stream].
24///
25/// [committed stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#committed_type
26#[derive(Debug)]
27pub struct CommittedWriter<F> {
28    pub(crate) inner: BaseWriter<F>,
29}
30
31impl<F> CommittedWriter<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 committed 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    /// Finalize the stream, preventing further writes.
55    pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
56        self.inner.finalize().await
57    }
58}
59
60#[cfg(test)]
61mod tests {
62    use super::*;
63    use crate::error::AppendError;
64    use crate::write::test::*;
65    use bigquery_grpc_mock::{MockBigQueryWrite, start};
66    use gaxi::grpc::tonic::Response as TonicResponse;
67    use tokio::sync::mpsc;
68
69    #[tokio::test]
70    async fn basic_success() -> anyhow::Result<()> {
71        let (response_tx, response_rx) = mpsc::channel(10);
72
73        let mut mock = MockBigQueryWrite::new();
74        mock.expect_append_rows()
75            .return_once(|_| Ok(TonicResponse::from(response_rx)));
76
77        mock.expect_finalize_write_stream().return_once(|_| {
78            Ok(TonicResponse::new(
79                bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default(),
80            ))
81        });
82
83        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
84        let transport = Arc::new(test_transport(endpoint).await?);
85
86        let writer = CommittedWriter::new(transport, write_stream(), format());
87        assert_eq!(writer.write_stream(), write_stream());
88
89        response_tx.send(Ok(convert(&test_response(1)))).await?;
90        let resp = writer.append(rows(1)).send().await?;
91        assert_eq!(resp.offset, Some(1));
92
93        response_tx.send(Ok(convert(&test_response(2)))).await?;
94        let resp = writer.append(rows(2)).send().await?;
95        assert_eq!(resp.offset, Some(2));
96
97        response_tx.send(Ok(convert(&test_response(3)))).await?;
98        let resp = writer.append(rows(3)).send().await?;
99        assert_eq!(resp.offset, Some(3));
100
101        drop(response_tx);
102        let err = writer.append(rows(4)).send().await.expect_err("channel");
103        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
104
105        // We can still finalize the stream even if row appends hit a closed bidirectional stream
106        writer.finalize().await?;
107
108        Ok(())
109    }
110}