google-cloud-bigquery 0.16.1-preview

Google Cloud Client Libraries for Rust - BigQuery
Documentation
// Copyright 2026 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
//     https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

use super::append_future::AppendFuture;
use super::append_response::to_result;
use super::error::AppendError;
use super::runner::WriteRequest;
use crate::Error;
use crate::model::AppendRowsRequest;
use gaxi::prost::{FromProto, ToProto};
use tokio::sync::{mpsc, oneshot};

/// A request builder for appending rows with a specific stream offset,
/// ensuring exactly-once semantics.
#[derive(Clone, Debug)]
pub struct AppendWithOffset {
    req_tx: mpsc::UnboundedSender<WriteRequest>,
    pub(crate) req: AppendRowsRequest,
}

impl AppendWithOffset {
    #[allow(dead_code)]
    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
        Self { req_tx, req }
    }

    /// Sets the target stream offset to guarantee [exactly-once] writes.
    ///
    /// # Example
    ///
    /// ```
    /// # use google_cloud_bigquery::write::arrow::PendingWriter;
    /// # async fn sample(writer: PendingWriter) -> anyhow::Result<()> {
    /// let resp = writer.append(rows()).set_offset(0).send().await?;
    /// # Ok(()) }
    ///
    /// use google_cloud_bigquery::model::ArrowRecordBatch;
    /// fn rows() -> ArrowRecordBatch {
    ///   todo!("Define your rows...")
    /// }
    /// ```
    ///
    /// [exactly-once]: https://docs.cloud.google.com/bigquery/docs/write-api-best-practices#manage_stream_offsets_to_achieve_exactly-once_semantics
    pub fn set_offset(mut self, offset: i64) -> Self {
        self.req.offset = Some(offset);
        self
    }

    /// Append rows to the stream.
    ///
    /// Applications are encouraged to queue up requests and await their
    /// responses independently.
    ///
    /// Note that the service will reject requests with a mismatched offset, so
    /// requests must be queued in order.
    ///
    /// # Example
    ///
    /// ```
    /// # use google_cloud_bigquery::write::arrow::PendingWriter;
    /// # async fn sample(writer: PendingWriter) -> anyhow::Result<()> {
    /// let f1 = writer.append(rows()).set_offset(0).send();
    /// let f2 = writer.append(rows()).set_offset(1).send();
    ///
    /// let resp1 = f1.await?;
    /// let resp2 = f2.await?;
    /// # Ok(()) }
    ///
    /// use google_cloud_bigquery::model::ArrowRecordBatch;
    /// fn rows() -> ArrowRecordBatch {
    ///   todo!("Define your rows...")
    /// }
    /// ```
    pub fn send(self) -> AppendFuture {
        let (tx, rx) = oneshot::channel();
        let (resp_tx, resp_rx) = oneshot::channel();
        let req = match self.req.to_proto().map_err(Error::deser) {
            Ok(req) => req,
            Err(e) => {
                let _ = tx.send(Err(e.into()));
                return AppendFuture::new(rx);
            }
        };
        let write = WriteRequest { req, resp_tx };
        let _ = self.req_tx.send(write);
        tokio::spawn(async move {
            let res = async {
                let resp = resp_rx
                    .await
                    .map_err(|_| AppendError::UnexpectedEndOfStream)??;
                let resp = resp.cnv().map_err(Error::ser)?;
                to_result(resp)
            }
            .await;
            let _ = tx.send(res);
        });
        AppendFuture::new(rx)
    }
}

/// A request builder for appending rows on the default stream.
#[derive(Clone, Debug)]
pub struct Append {
    req_tx: mpsc::UnboundedSender<WriteRequest>,
    pub(crate) req: AppendRowsRequest,
}

impl Append {
    pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
        Self { req_tx, req }
    }

    /// Append rows to the stream.
    ///
    /// Applications are encouraged to queue up requests and await their
    /// responses independently.
    ///
    /// # Example
    ///
    /// ```
    /// # use google_cloud_bigquery::write::arrow::DefaultWriter;
    /// # async fn sample(writer: DefaultWriter) -> anyhow::Result<()> {
    /// let f1 = writer.append(rows()).send();
    /// let f2 = writer.append(rows()).send();
    ///
    /// let resp1 = f1.await?;
    /// let resp2 = f2.await?;
    /// # Ok(()) }
    ///
    /// use google_cloud_bigquery::model::ArrowRecordBatch;
    /// fn rows() -> ArrowRecordBatch {
    ///   todo!("Define your rows...")
    /// }
    /// ```
    pub fn send(self) -> AppendFuture {
        let (tx, rx) = oneshot::channel();
        tokio::spawn(async move {
            let (resp_tx, resp_rx) = oneshot::channel();
            let res = async move {
                let req = self.req.to_proto().map_err(Error::deser)?;
                let write = WriteRequest { req, resp_tx };
                let _ = self.req_tx.send(write);
                let resp = resp_rx
                    .await
                    .map_err(|_| AppendError::UnexpectedEndOfStream)??;
                let resp = resp.cnv().map_err(Error::ser)?;
                to_result(resp)
            }
            .await;
            let _ = tx.send(res);
        });
        AppendFuture::new(rx)
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::google::cloud::bigquery::storage::v1;
    use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
        AppendResult, Response,
    };
    use crate::model::TableSchema;

    #[tokio::test]
    async fn success() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = Append::new(req_tx, req);
        let handle = tokio::spawn(async move { builder.send().await });

        // Receive and verify the request
        let write = req_rx.recv().await.expect("should receive request");
        assert_eq!(write.req.write_stream, write_stream());

        // Provide a successful response
        let resp = v1::AppendRowsResponse {
            response: Some(Response::AppendResult(AppendResult::default())),
            write_stream: write_stream(),
            updated_schema: Some(v1::TableSchema::default()),
            ..Default::default()
        };
        write
            .resp_tx
            .send(Ok(resp))
            .expect("sending on channel always succeeds");

        let resp = handle.await??;
        assert_eq!(resp.offset, None);
        assert_eq!(resp.updated_schema, Some(TableSchema::default()));
        Ok(())
    }

    #[tokio::test]
    async fn stream_closed() -> anyhow::Result<()> {
        let (req_tx, req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = Append::new(req_tx, req);
        let handle = tokio::spawn(async move { builder.send().await });

        // Simulate a stream closure
        drop(req_rx);

        let err = handle.await?.expect_err("should return an error");
        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
        Ok(())
    }

    #[tokio::test]
    async fn rpc_error() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = Append::new(req_tx, req);
        let handle = tokio::spawn(async move { builder.send().await });

        // Simulate a stream ending in a known error
        let write = req_rx.recv().await.expect("should receive request");
        let append_err: AppendError = Error::io("fail").into();
        write
            .resp_tx
            .send(Err(append_err))
            .expect("sending on channel always succeeds");

        let err = handle.await?.expect_err("should return an error");
        assert!(matches!(err, AppendError::Rpc { source: _ }));
        Ok(())
    }

    #[tokio::test]
    async fn row_errors() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = Append::new(req_tx, req);
        let handle = tokio::spawn(async move { builder.send().await });

        let write = req_rx.recv().await.expect("should receive request");

        let row_error = v1::RowError {
            index: 42,
            code: v1::row_error::RowErrorCode::FieldsError as i32,
            message: "fail".to_string(),
        };
        let resp = v1::AppendRowsResponse {
            row_errors: vec![row_error],
            write_stream: write_stream(),
            ..Default::default()
        };
        write
            .resp_tx
            .send(Ok(resp))
            .expect("sending on channel always succeeds");

        let err = handle.await?.expect_err("should return an error");
        assert!(matches!(err, AppendError::RowErrors(_)));
        Ok(())
    }

    #[tokio::test]
    async fn offset_success() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
        let future = builder.send();

        let write = req_rx.recv().await.expect("should receive request");
        assert_eq!(write.req.offset, Some(100));

        let resp = v1::AppendRowsResponse {
            response: Some(Response::AppendResult(AppendResult::default())),
            write_stream: write_stream(),
            // Ensure schema matches none
            ..Default::default()
        };
        write
            .resp_tx
            .send(Ok(resp))
            .expect("sending on channel always succeeds");

        let resp = future.await?;
        assert_eq!(resp.offset, None);
        assert_eq!(resp.updated_schema, None);
        Ok(())
    }

    #[tokio::test]
    async fn offset_stream_closed() -> anyhow::Result<()> {
        let (req_tx, req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
        let future = builder.send();

        drop(req_rx);

        let err = future.await.expect_err("should return an error");
        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
        Ok(())
    }

    #[tokio::test]
    async fn offset_rpc_error() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
        let future = builder.send();

        // Simulate a stream ending in a known error
        let write = req_rx.recv().await.expect("should receive request");
        let append_err: AppendError = Error::io("fail").into();
        write
            .resp_tx
            .send(Err(append_err))
            .expect("sending on channel always succeeds");

        let err = future.await.expect_err("should return an error");
        assert!(matches!(err, AppendError::Rpc { source: _ }));
        Ok(())
    }

    #[tokio::test]
    async fn offset_row_errors() -> anyhow::Result<()> {
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let req = AppendRowsRequest::new().set_write_stream(write_stream());

        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
        let future = builder.send();

        let write = req_rx.recv().await.expect("should receive request");

        let row_error = v1::RowError {
            index: 42,
            code: v1::row_error::RowErrorCode::FieldsError as i32,
            message: "fail".to_string(),
        };
        let resp = v1::AppendRowsResponse {
            row_errors: vec![row_error],
            write_stream: write_stream(),
            ..Default::default()
        };
        write
            .resp_tx
            .send(Ok(resp))
            .expect("sending on channel always succeeds");

        let err = future.await.expect_err("should return an error");
        assert!(matches!(err, AppendError::RowErrors(_)));
        Ok(())
    }

    #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
    async fn synchronous_queueing() -> anyhow::Result<()> {
        const NUM_WRITES: i64 = 1000;
        let (req_tx, mut req_rx) = mpsc::unbounded_channel();
        let write_handle = tokio::spawn(async move {
            let mut writes = tokio::task::JoinSet::new();
            for i in 0..NUM_WRITES {
                writes.spawn(
                    AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
                        .set_offset(i)
                        .send(),
                );
            }
            let _ = writes.join_all().await;
        });

        for i in 0..NUM_WRITES {
            let write = req_rx.recv().await.expect("should receive request");
            assert_eq!(write.req.offset, Some(i), "received out of order write");
        }
        write_handle.await?;
        Ok(())
    }

    fn write_stream() -> String {
        "projects/p/datasets/d/tables/t/streams/_default".to_string()
    }
}