#[allow(unused_imports)]
use crate::Error;
use crate::Result;
const DEFAULT_HOST: &str = "https://bigquerystorage.googleapis.com";
mod info {
const NAME: &str = env!("CARGO_PKG_NAME");
const VERSION: &str = env!("CARGO_PKG_VERSION");
pub(crate) static X_GOOG_API_CLIENT_HEADER: std::sync::LazyLock<String> =
std::sync::LazyLock::new(|| {
let ac = gaxi::api_header::XGoogApiClient {
name: NAME,
version: VERSION,
library_type: gaxi::api_header::GAPIC,
};
ac.grpc_header_value()
});
}
#[derive(Clone)]
pub struct Read {
pub(crate) inner: gaxi::grpc::Client,
}
impl std::fmt::Debug for Read {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::result::Result<(), std::fmt::Error> {
f.debug_struct("Read").field("inner", &self.inner).finish()
}
}
impl Read {
pub async fn new(config: gaxi::options::ClientConfig) -> crate::ClientBuilderResult<Self> {
let inner = if gaxi::options::tracing_enabled(&config) {
gaxi::grpc::Client::new_with_instrumentation(
config,
DEFAULT_HOST,
&super::tracing::info::INSTRUMENTATION_CLIENT_INFO,
)
.await?
} else {
gaxi::grpc::Client::new(config, DEFAULT_HOST).await?
};
Ok(Self { inner })
}
}
impl super::stub::Read for Read {
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>> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, false);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryRead",
"CreateReadSession",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryRead/CreateReadSession",
);
let x_goog_request_params = [Some(&req)
.and_then(|m| m.read_session.as_ref())
.map(|m| &m.table)
.map(|s| s.as_str())
.map(|v| format!("read_session.table={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::ReadSession;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method("google.cloud.bigquery.storage.v1.BigQueryRead/CreateReadSession");
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req)
.and_then(|m| m.read_session.as_ref())
.map(|m| &m.table)
.map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::ReadSession,
>,
)
}
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,
>,
> {
let extensions = {
let mut e = gaxi::grpc::tonic::Extensions::new();
e.insert(gaxi::grpc::tonic::GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryRead",
"ReadRows",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryRead/ReadRows",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.read_stream)
.map(|s| s.as_str())
.map(|v| format!("read_stream={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
self.inner
.execute_server_streaming::<
crate::write::generated::gapic_storage::model::ReadRowsRequest,
crate::write::generated::gapic_storage::model::ReadRowsResponse,
crate::google::cloud::bigquery::storage::v1::ReadRowsRequest,
crate::google::cloud::bigquery::storage::v1::ReadRowsResponse,
>(
extensions,
path,
req,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
}
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>,
> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, true);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryRead",
"SplitReadStream",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryRead/SplitReadStream",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.name)
.map(|s| s.as_str())
.map(|v| format!("name={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::SplitReadStreamResponse;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method("google.cloud.bigquery.storage.v1.BigQueryRead/SplitReadStream");
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.name).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::SplitReadStreamResponse,
>,
)
}
}
#[derive(Clone)]
pub struct BigQueryWrite {
pub(crate) inner: gaxi::grpc::Client,
}
impl std::fmt::Debug for BigQueryWrite {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::result::Result<(), std::fmt::Error> {
f.debug_struct("BigQueryWrite")
.field("inner", &self.inner)
.finish()
}
}
impl BigQueryWrite {
pub async fn new(config: gaxi::options::ClientConfig) -> crate::ClientBuilderResult<Self> {
let inner = if gaxi::options::tracing_enabled(&config) {
gaxi::grpc::Client::new_with_instrumentation(
config,
DEFAULT_HOST,
&super::tracing::info::INSTRUMENTATION_CLIENT_INFO,
)
.await?
} else {
gaxi::grpc::Client::new(config, DEFAULT_HOST).await?
};
Ok(Self { inner })
}
}
impl super::stub::BigQueryWrite for BigQueryWrite {
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>> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, false);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"CreateWriteStream",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/CreateWriteStream",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.parent)
.map(|s| s.as_str())
.map(|v| format!("parent={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::WriteStream;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method("google.cloud.bigquery.storage.v1.BigQueryWrite/CreateWriteStream");
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.parent).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::WriteStream,
>,
)
}
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,
>,
) {
let x_goog_request_params = "";
let extensions = {
let mut e = gaxi::grpc::tonic::Extensions::new();
e.insert(gaxi::grpc::tonic::GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"AppendRows",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/AppendRows",
);
self.inner
.execute_bidi_streaming::<
crate::write::generated::gapic_storage::model::AppendRowsRequest,
crate::write::generated::gapic_storage::model::AppendRowsResponse,
crate::google::cloud::bigquery::storage::v1::AppendRowsRequest,
crate::google::cloud::bigquery::storage::v1::AppendRowsResponse,
>(
extensions,
path,
options,
&info::X_GOOG_API_CLIENT_HEADER,
x_goog_request_params,
)
}
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>> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, false);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"GetWriteStream",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/GetWriteStream",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.name)
.map(|s| s.as_str())
.map(|v| format!("name={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::WriteStream;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method("google.cloud.bigquery.storage.v1.BigQueryWrite/GetWriteStream");
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.name).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::WriteStream,
>,
)
}
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>,
> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, false);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"FinalizeWriteStream",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/FinalizeWriteStream",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.name)
.map(|s| s.as_str())
.map(|v| format!("name={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method(
"google.cloud.bigquery.storage.v1.BigQueryWrite/FinalizeWriteStream",
);
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.name).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::FinalizeWriteStreamResponse,
>,
)
}
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,
>,
> {
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, true);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"BatchCommitWriteStreams",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/BatchCommitWriteStreams",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.parent)
.map(|s| s.as_str())
.map(|v| format!("parent={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::BatchCommitWriteStreamsResponse;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method(
"google.cloud.bigquery.storage.v1.BigQueryWrite/BatchCommitWriteStreams",
);
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.parent).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::BatchCommitWriteStreamsResponse,
>,
)
}
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>>
{
use gaxi::{
grpc::tonic::{Extensions, GrpcMethod},
prost::ToProto,
};
let options = google_cloud_gax::options::internal::set_default_idempotency(options, false);
let extensions = {
let mut e = Extensions::new();
e.insert(GrpcMethod::new(
"google.cloud.bigquery.storage.v1.BigQueryWrite",
"FlushRows",
));
e
};
let path = http::uri::PathAndQuery::from_static(
"/google.cloud.bigquery.storage.v1.BigQueryWrite/FlushRows",
);
let x_goog_request_params = [Some(&req)
.map(|m| &m.write_stream)
.map(|s| s.as_str())
.map(|v| format!("write_stream={v}"))]
.into_iter()
.flatten()
.fold(String::new(), |b, p| b + "&" + &p);
type TR = crate::google::cloud::bigquery::storage::v1::FlushRowsResponse;
if let Some(recorder) = gaxi::observability::RequestRecorder::current() {
let attributes = gaxi::observability::ClientRequestAttributes::default()
.set_rpc_method("google.cloud.bigquery.storage.v1.BigQueryWrite/FlushRows");
let resource_name = (|| {
Some(format!(
"//bigquerystorage.googleapis.com/{}",
Some(&req).map(|m| &m.write_stream).map(|s| s.as_str())?,
))
})();
let attributes = if let Some(rn) = resource_name.filter(|s| !s.is_empty()) {
attributes.set_resource_name(rn)
} else {
attributes
};
recorder.on_client_request(attributes);
}
self.inner
.execute(
extensions,
path,
req.to_proto().map_err(Error::deser)?,
options,
&info::X_GOOG_API_CLIENT_HEADER,
&x_goog_request_params,
)
.await
.and_then(
gaxi::grpc::to_gax_response::<
TR,
crate::write::generated::gapic_storage::model::FlushRowsResponse,
>,
)
}
}