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