google_cloud_bigquery/write/
default.rs1use super::builder::Append;
16use super::dispatcher::Dispatcher;
17use super::format::DataFormat;
18use super::pool::StreamPool;
19use super::retry_policy::RetryOptions;
20use std::sync::Arc;
21
22#[derive(Debug)]
33pub struct DefaultWriter<F> {
34 pub(crate) inner: Arc<Dispatcher>,
35 pub(crate) write_stream: String,
36 pub(crate) format: F,
37}
38
39impl<F> DefaultWriter<F>
40where
41 F: DataFormat,
42{
43 pub(crate) fn new(
44 pool: Arc<StreamPool>,
45 retry_options: RetryOptions,
46 write_stream: String,
47 format: F,
48 ) -> Self {
49 let inner = Arc::new(Dispatcher::new(pool, retry_options));
50 Self {
51 inner,
52 write_stream,
53 format,
54 }
55 }
56
57 pub fn append(&self, rows: F::Rows) -> Append {
59 let req = self.format.make_request(&self.write_stream, rows);
60 Append::new(self.inner.clone(), req)
61 }
62}
63
64#[cfg(test)]
65mod tests {
66 use super::super::pool::StreamPoolOptions;
67 use super::*;
68 use crate::error::AppendError;
69 use crate::write::test::*;
70 use bigquery_grpc_mock::{MockBigQueryWrite, start};
71 use gaxi::grpc::tonic::{Response as TonicResponse, Status as TonicStatus};
72 use tokio::sync::mpsc;
73
74 #[tokio::test]
75 async fn basic_success() -> anyhow::Result<()> {
76 let (response_tx, response_rx) = mpsc::channel(10);
77
78 let mut mock = MockBigQueryWrite::new();
79 mock.expect_append_rows()
80 .return_once(|_| Ok(TonicResponse::from(response_rx)));
81 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
82 let transport = Arc::new(test_transport(endpoint).await?);
83 let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default()));
84
85 let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), format());
86
87 response_tx.send(Ok(convert(&test_response(1)))).await?;
88 let resp = writer.append(rows(1)).send().await?;
89 assert_eq!(resp.offset, Some(1));
90
91 response_tx.send(Ok(convert(&test_response(2)))).await?;
92 let resp = writer.append(rows(2)).send().await?;
93 assert_eq!(resp.offset, Some(2));
94
95 response_tx.send(Ok(convert(&test_response(3)))).await?;
96 let resp = writer.append(rows(3)).send().await?;
97 assert_eq!(resp.offset, Some(3));
98
99 response_tx
100 .send(Err(TonicStatus::failed_precondition("fail")))
101 .await?;
102 let err = writer.append(rows(4)).send().await.expect_err("fail");
103 assert!(matches!(err, AppendError::Rpc { source: _ }), "{err:?}");
104
105 Ok(())
106 }
107}