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.
//
// Code generated by sidekick. DO NOT EDIT.
use crate::Result;

/// Implements a [Read](super::stub::Read) decorator for logging and tracing.
#[derive(Clone, Debug)]
pub struct Read<T>
where
    T: super::stub::Read + std::fmt::Debug + Send + Sync,
{
    inner: T,
    #[allow(dead_code)]
    duration: gaxi::observability::DurationMetric,
}

impl<T> Read<T>
where
    T: super::stub::Read + std::fmt::Debug + Send + Sync,
{
    pub fn new(inner: T) -> Self {
        Self {
            inner,
            duration: gaxi::observability::DurationMetric::new(&info::INSTRUMENTATION_CLIENT_INFO),
        }
    }
}

impl<T> super::stub::Read for Read<T>
where
    T: super::stub::Read + std::fmt::Debug + Send + Sync,
{
    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn create_read_session(
        &self,
        req: crate::write::generated::gapic_storage::model::CreateReadSessionRequest,
        options: crate::RequestOptions,
    ) -> Result<crate::Response<crate::write::generated::gapic_storage::model::ReadSession>> {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::Read::create_read_session",
            self.inner.create_read_session(req, options));
        pending.await
    }

    async fn read_rows(
        &self,
        req: crate::write::generated::gapic_storage::model::ReadRowsRequest,
        options: crate::RequestOptions,
    ) -> Result<
        google_cloud_gax::streaming::ResponseStream<
            crate::write::generated::gapic_storage::model::ReadRowsResponse,
        >,
    > {
        self.inner.read_rows(req, options).await
    }

    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn split_read_stream(
        &self,
        req: crate::write::generated::gapic_storage::model::SplitReadStreamRequest,
        options: crate::RequestOptions,
    ) -> Result<
        crate::Response<crate::write::generated::gapic_storage::model::SplitReadStreamResponse>,
    > {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::Read::split_read_stream",
            self.inner.split_read_stream(req, options));
        pending.await
    }
}

/// Implements a [BigQueryWrite](super::stub::BigQueryWrite) decorator for logging and tracing.
#[derive(Clone, Debug)]
pub struct BigQueryWrite<T>
where
    T: super::stub::BigQueryWrite + std::fmt::Debug + Send + Sync,
{
    inner: T,
    #[allow(dead_code)]
    duration: gaxi::observability::DurationMetric,
}

impl<T> BigQueryWrite<T>
where
    T: super::stub::BigQueryWrite + std::fmt::Debug + Send + Sync,
{
    pub fn new(inner: T) -> Self {
        Self {
            inner,
            duration: gaxi::observability::DurationMetric::new(&info::INSTRUMENTATION_CLIENT_INFO),
        }
    }
}

impl<T> super::stub::BigQueryWrite for BigQueryWrite<T>
where
    T: super::stub::BigQueryWrite + std::fmt::Debug + Send + Sync,
{
    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn create_write_stream(
        &self,
        req: crate::write::generated::gapic_storage::model::CreateWriteStreamRequest,
        options: crate::RequestOptions,
    ) -> Result<crate::Response<crate::write::generated::gapic_storage::model::WriteStream>> {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::BigQueryWrite::create_write_stream",
            self.inner.create_write_stream(req, options));
        pending.await
    }

    fn append_rows(
        &self,
        options: crate::RequestOptions,
    ) -> (
        google_cloud_gax::streaming::RequestSender<
            crate::write::generated::gapic_storage::model::AppendRowsRequest,
        >,
        google_cloud_gax::streaming::ResponseStream<
            crate::write::generated::gapic_storage::model::AppendRowsResponse,
        >,
    ) {
        self.inner.append_rows(options)
    }

    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn get_write_stream(
        &self,
        req: crate::write::generated::gapic_storage::model::GetWriteStreamRequest,
        options: crate::RequestOptions,
    ) -> Result<crate::Response<crate::write::generated::gapic_storage::model::WriteStream>> {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::BigQueryWrite::get_write_stream",
            self.inner.get_write_stream(req, options));
        pending.await
    }

    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn finalize_write_stream(
        &self,
        req: crate::write::generated::gapic_storage::model::FinalizeWriteStreamRequest,
        options: crate::RequestOptions,
    ) -> Result<
        crate::Response<crate::write::generated::gapic_storage::model::FinalizeWriteStreamResponse>,
    > {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::BigQueryWrite::finalize_write_stream",
            self.inner.finalize_write_stream(req, options));
        pending.await
    }

    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn batch_commit_write_streams(
        &self,
        req: crate::write::generated::gapic_storage::model::BatchCommitWriteStreamsRequest,
        options: crate::RequestOptions,
    ) -> Result<
        crate::Response<
            crate::write::generated::gapic_storage::model::BatchCommitWriteStreamsResponse,
        >,
    > {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::BigQueryWrite::batch_commit_write_streams",
            self.inner.batch_commit_write_streams(req, options));
        pending.await
    }

    #[tracing::instrument(level = tracing::Level::DEBUG, ret)]
    async fn flush_rows(
        &self,
        req: crate::write::generated::gapic_storage::model::FlushRowsRequest,
        options: crate::RequestOptions,
    ) -> Result<crate::Response<crate::write::generated::gapic_storage::model::FlushRowsResponse>>
    {
        let (_span, pending) = gaxi::client_request_signals!(
            metric: self.duration.clone(),
            info: *info::INSTRUMENTATION_CLIENT_INFO,
            method: "client::BigQueryWrite::flush_rows",
            self.inner.flush_rows(req, options));
        pending.await
    }
}

pub(crate) mod info {
    const NAME: &str = env!("CARGO_PKG_NAME");
    const VERSION: &str = env!("CARGO_PKG_VERSION");
    pub(crate) static INSTRUMENTATION_CLIENT_INFO: std::sync::LazyLock<
        gaxi::options::InstrumentationClientInfo,
    > = std::sync::LazyLock::new(|| {
        let mut info = gaxi::options::InstrumentationClientInfo::default();
        info.service_name = "bigquerystorage";
        info.client_version = VERSION;
        info.client_artifact = NAME;
        info.default_host = "bigquerystorage";
        info
    });
}