google_cloud_bigquery/write/
buffered.rs1use super::base::BaseWriter;
16use super::format::DataFormat;
17use crate::Result;
18use crate::model::{FinalizeWriteStreamResponse, FlushRowsResponse};
19use crate::write::builder::AppendWithOffset;
20use crate::write::transport::Transport;
21use std::sync::Arc;
22
23#[derive(Debug)]
27pub struct BufferedWriter<F> {
28 pub(crate) inner: BaseWriter<F>,
29}
30
31impl<F> BufferedWriter<F>
32where
33 F: DataFormat,
34{
35 pub(crate) fn new(inner: Arc<Transport>, write_stream: String, format: F) -> Self {
36 Self {
37 inner: BaseWriter::new(inner, write_stream, format),
38 }
39 }
40
41 pub fn write_stream(&self) -> &str {
43 &self.inner.write_stream
44 }
45
46 pub fn append(&self, rows: F::Rows) -> AppendWithOffset {
48 AppendWithOffset::new(
49 self.inner.runner.req_tx.clone(),
50 self.inner.append_request(rows),
51 )
52 }
53
54 pub async fn flush(&self, offset: i64) -> Result<FlushRowsResponse> {
56 self.inner
57 .client
58 .flush_rows()
59 .set_write_stream(&self.inner.write_stream)
60 .set_offset(offset)
61 .send()
62 .await
63 }
64
65 pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
67 self.inner.finalize().await
68 }
69}
70
71#[cfg(test)]
72mod tests {
73 use super::*;
74 use crate::error::AppendError;
75 use crate::write::test::*;
76 use bigquery_grpc_mock::{MockBigQueryWrite, start};
77 use gaxi::grpc::tonic::Response as TonicResponse;
78 use tokio::sync::mpsc;
79
80 #[tokio::test]
81 async fn basic_success() -> anyhow::Result<()> {
82 let (response_tx, response_rx) = mpsc::channel(10);
83
84 let mut mock = MockBigQueryWrite::new();
85 mock.expect_append_rows()
86 .return_once(|_| Ok(TonicResponse::from(response_rx)));
87
88 mock.expect_flush_rows()
89 .return_once(|req| {
90 assert_eq!(req.get_ref().offset, Some(3));
91 assert_eq!(req.get_ref().write_stream, write_stream());
92 Ok(TonicResponse::new(
93 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse::default()
94 ))
95 });
96
97 mock.expect_finalize_write_stream()
98 .return_once(|req| {
99 assert_eq!(req.get_ref().name, write_stream());
100 Ok(TonicResponse::new(
101 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default()
102 ))
103 });
104
105 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
106 let transport = Arc::new(test_transport(endpoint).await?);
107
108 let writer = BufferedWriter::new(transport, write_stream(), format());
109 assert_eq!(writer.write_stream(), write_stream());
110
111 response_tx.send(Ok(convert(&test_response(1)))).await?;
112 let resp = writer.append(rows(1)).send().await?;
113 assert_eq!(resp.offset, Some(1));
114
115 response_tx.send(Ok(convert(&test_response(2)))).await?;
116 let resp = writer.append(rows(2)).send().await?;
117 assert_eq!(resp.offset, Some(2));
118
119 response_tx.send(Ok(convert(&test_response(3)))).await?;
120 let resp = writer.append(rows(3)).send().await?;
121 assert_eq!(resp.offset, Some(3));
122
123 drop(response_tx);
124 let err = writer.append(rows(4)).send().await.expect_err("channel");
125 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
126
127 writer.flush(3).await?;
128 writer.finalize().await?;
129
130 Ok(())
131 }
132
133 #[tokio::test]
134 async fn multiple_flushes() -> anyhow::Result<()> {
135 let (response_tx, response_rx) = mpsc::channel(10);
136 let mut mock = MockBigQueryWrite::new();
137 mock.expect_append_rows()
138 .return_once(|_| Ok(TonicResponse::from(response_rx)));
139
140 mock.expect_flush_rows().times(2).returning(|req| {
141 Ok(TonicResponse::new(
142 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse {
143 offset: req.get_ref().offset.unwrap_or(0),
144 },
145 ))
146 });
147
148 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
149 let transport = Arc::new(test_transport(endpoint).await?);
150 let writer = BufferedWriter::new(transport, write_stream(), format());
151 assert_eq!(writer.write_stream(), write_stream());
152
153 response_tx.send(Ok(convert(&test_response(1)))).await?;
154 let _ = writer.append(rows(1)).send().await?;
155 let flush1 = writer.flush(1).await?;
156 assert_eq!(flush1.offset, 1);
157
158 response_tx.send(Ok(convert(&test_response(2)))).await?;
159 let _ = writer.append(rows(2)).send().await?;
160 let flush2 = writer.flush(2).await?;
161 assert_eq!(flush2.offset, 2);
162
163 Ok(())
164 }
165}