google-cloud-bigquery 0.17.0-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::super::append_future::AppendFuture;
use super::super::append_response::to_result;
use super::super::error::AppendError;
use super::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 {
    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::format::Arrow;
    /// # use google_cloud_bigquery::write::PendingWriter;
    /// # async fn sample(writer: PendingWriter<Arrow>) -> 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::format::Arrow;
    /// # use google_cloud_bigquery::write::PendingWriter;
    /// # async fn sample(writer: PendingWriter<Arrow>) -> 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 (resp_tx, resp_rx) = oneshot::channel();
        let req = match self.req.to_proto().map_err(Error::ser) {
            Ok(req) => req,
            Err(e) => return AppendFuture::from_future(async move { Err(e.into()) }),
        };
        let write = WriteRequest { req, resp_tx };
        if self.req_tx.send(write).is_err() {
            return AppendFuture::from_future(
                async move { Err(AppendError::UnexpectedEndOfStream) },
            );
        }
        AppendFuture::from_future(async move {
            let resp = resp_rx
                .await
                .map_err(|_| AppendError::UnexpectedEndOfStream)??;
            let resp = resp.cnv().map_err(Error::deser)?;
            to_result(resp)
        })
    }
}

#[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::write::test::*;

    #[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 = 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 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 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 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(())
    }

    #[tokio::test]
    async fn send_when_req_tx_closed_returns_unexpected_end_of_stream() {
        let (req_tx, req_rx) = mpsc::unbounded_channel();
        drop(req_rx);

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

        let err = future
            .await
            .expect_err("should return unexpected end of stream");
        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
    }

    #[tokio::test]
    async fn send_serialization_error() {
        use crate::model::append_rows_request::MissingValueInterpretation;

        let (req_tx, _req_rx) = mpsc::unbounded_channel();
        let invalid: MissingValueInterpretation = serde_json::from_str("\"INVALID\"").unwrap();
        let req = AppendRowsRequest::new().set_default_missing_value_interpretation(invalid);
        let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
        let future = builder.send();

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