use super::arrow::WriterBuilder as ArrowWriterBuilder;
use super::client_builder::ClientBuilder;
use super::transport::Transport;
use crate::ClientBuilderResult as BuilderResult;
use crate::model::ArrowSchema;
use std::sync::Arc;
#[derive(Debug)]
pub struct Write {
#[allow(unused)]
inner: Arc<Transport>,
}
impl Write {
pub fn builder() -> ClientBuilder {
ClientBuilder::new()
}
pub(crate) async fn new(builder: ClientBuilder) -> BuilderResult<Self> {
let transport = Transport::new(builder.config).await?;
Ok(Self {
inner: Arc::new(transport),
})
}
pub fn arrow(&self, schema: ArrowSchema) -> ArrowWriterBuilder {
ArrowWriterBuilder::new(self.inner.clone(), schema)
}
}
#[cfg(test)]
mod tests {
use super::super::error::AppendError;
use super::*;
use crate::model::{ArrowRecordBatch, ArrowSchema};
use bigquery_grpc_mock::{MockBigQueryWrite, start};
use gaxi::grpc::tonic::Status as TonicStatus;
use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
#[tokio::test]
async fn arrow() -> anyhow::Result<()> {
let mut mock = MockBigQueryWrite::new();
mock.expect_append_rows()
.return_once(|_| Err(TonicStatus::failed_precondition("fail")));
let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
let client = Write::builder()
.with_endpoint(endpoint)
.with_credentials(Anonymous::new().build())
.build()
.await?;
let writer = client
.arrow(ArrowSchema::new())
.default("projects/p/datasets/d/tables/t")?;
let err = writer
.append(ArrowRecordBatch::new())
.send()
.await
.expect_err("write should fail");
assert!(matches!(err, AppendError::Rpc { source: _ }));
Ok(())
}
}