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)]
33pub struct BufferedWriter<F> {
34 pub(crate) inner: BaseWriter<F>,
35}
36
37impl<F> BufferedWriter<F>
38where
39 F: DataFormat,
40{
41 pub(crate) fn new(inner: Arc<Transport>, write_stream: String, format: F) -> Self {
42 Self {
43 inner: BaseWriter::new(inner, write_stream, format),
44 }
45 }
46
47 pub fn write_stream(&self) -> &str {
49 &self.inner.write_stream
50 }
51
52 pub fn append(&self, rows: F::Rows) -> AppendWithOffset {
54 AppendWithOffset::new(
55 self.inner.runner.req_tx.clone(),
56 self.inner.append_request(rows),
57 )
58 }
59
60 pub async fn flush(&self, offset: i64) -> Result<FlushRowsResponse> {
86 self.inner
87 .client
88 .flush_rows()
89 .set_write_stream(&self.inner.write_stream)
90 .set_offset(offset)
91 .send()
92 .await
93 }
94
95 pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
97 self.inner.finalize().await
98 }
99}
100
101#[cfg(test)]
102mod tests {
103 use super::*;
104 use crate::error::AppendError;
105 use crate::write::test::*;
106 use bigquery_grpc_mock::{MockBigQueryWrite, start};
107 use gaxi::grpc::tonic::Response as TonicResponse;
108 use tokio::sync::mpsc;
109
110 #[tokio::test]
111 async fn basic_success() -> anyhow::Result<()> {
112 let (response_tx, response_rx) = mpsc::channel(10);
113
114 let mut mock = MockBigQueryWrite::new();
115 mock.expect_append_rows()
116 .return_once(|_| Ok(TonicResponse::from(response_rx)));
117
118 mock.expect_flush_rows()
119 .return_once(|req| {
120 assert_eq!(req.get_ref().offset, Some(3));
121 assert_eq!(req.get_ref().write_stream, write_stream());
122 Ok(TonicResponse::new(
123 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse::default()
124 ))
125 });
126
127 mock.expect_finalize_write_stream()
128 .return_once(|req| {
129 assert_eq!(req.get_ref().name, write_stream());
130 Ok(TonicResponse::new(
131 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default()
132 ))
133 });
134
135 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
136 let transport = Arc::new(test_transport(endpoint).await?);
137
138 let writer = BufferedWriter::new(transport, write_stream(), format());
139 assert_eq!(writer.write_stream(), write_stream());
140
141 response_tx.send(Ok(convert(&test_response(1)))).await?;
142 let resp = writer.append(rows(1)).send().await?;
143 assert_eq!(resp.offset, Some(1));
144
145 response_tx.send(Ok(convert(&test_response(2)))).await?;
146 let resp = writer.append(rows(2)).send().await?;
147 assert_eq!(resp.offset, Some(2));
148
149 response_tx.send(Ok(convert(&test_response(3)))).await?;
150 let resp = writer.append(rows(3)).send().await?;
151 assert_eq!(resp.offset, Some(3));
152
153 drop(response_tx);
154 let err = writer.append(rows(4)).send().await.expect_err("channel");
155 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
156
157 writer.flush(3).await?;
158 writer.finalize().await?;
159
160 Ok(())
161 }
162
163 #[tokio::test]
164 async fn multiple_flushes() -> anyhow::Result<()> {
165 let (response_tx, response_rx) = mpsc::channel(10);
166 let mut mock = MockBigQueryWrite::new();
167 mock.expect_append_rows()
168 .return_once(|_| Ok(TonicResponse::from(response_rx)));
169
170 mock.expect_flush_rows().times(2).returning(|req| {
171 Ok(TonicResponse::new(
172 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse {
173 offset: req.get_ref().offset.unwrap_or(0),
174 },
175 ))
176 });
177
178 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
179 let transport = Arc::new(test_transport(endpoint).await?);
180 let writer = BufferedWriter::new(transport, write_stream(), format());
181 assert_eq!(writer.write_stream(), write_stream());
182
183 response_tx.send(Ok(convert(&test_response(1)))).await?;
184 let _ = writer.append(rows(1)).send().await?;
185 let flush1 = writer.flush(1).await?;
186 assert_eq!(flush1.offset, 1);
187
188 response_tx.send(Ok(convert(&test_response(2)))).await?;
189 let _ = writer.append(rows(2)).send().await?;
190 let flush2 = writer.flush(2).await?;
191 assert_eq!(flush2.offset, 2);
192
193 Ok(())
194 }
195}