use crate::Result;
#[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
}
}
#[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
});
}