google_cloud_bigquery/write/
append_builder.rs1use super::append_future::AppendFuture;
16use super::append_response::to_result;
17use super::error::AppendError;
18use super::runner::WriteRequest;
19use crate::Error;
20use crate::model::AppendRowsRequest;
21use gaxi::prost::{FromProto, ToProto};
22use tokio::sync::{mpsc, oneshot};
23
24#[derive(Clone, Debug)]
27pub struct AppendWithOffset {
28 req_tx: mpsc::UnboundedSender<WriteRequest>,
29 pub(crate) req: AppendRowsRequest,
30}
31
32impl AppendWithOffset {
33 #[allow(dead_code)]
34 pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
35 Self { req_tx, req }
36 }
37
38 pub fn set_offset(mut self, offset: i64) -> Self {
56 self.req.offset = Some(offset);
57 self
58 }
59
60 pub fn send(self) -> AppendFuture {
86 let (tx, rx) = oneshot::channel();
87 let (resp_tx, resp_rx) = oneshot::channel();
88 let req = match self.req.to_proto().map_err(Error::deser) {
89 Ok(req) => req,
90 Err(e) => {
91 let _ = tx.send(Err(e.into()));
92 return AppendFuture::new(rx);
93 }
94 };
95 let write = WriteRequest { req, resp_tx };
96 let _ = self.req_tx.send(write);
97 tokio::spawn(async move {
98 let res = async {
99 let resp = resp_rx
100 .await
101 .map_err(|_| AppendError::UnexpectedEndOfStream)??;
102 let resp = resp.cnv().map_err(Error::ser)?;
103 to_result(resp)
104 }
105 .await;
106 let _ = tx.send(res);
107 });
108 AppendFuture::new(rx)
109 }
110}
111
112#[derive(Clone, Debug)]
114pub struct Append {
115 req_tx: mpsc::UnboundedSender<WriteRequest>,
116 pub(crate) req: AppendRowsRequest,
117}
118
119impl Append {
120 pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
121 Self { req_tx, req }
122 }
123
124 pub fn send(self) -> AppendFuture {
147 let (tx, rx) = oneshot::channel();
148 tokio::spawn(async move {
149 let (resp_tx, resp_rx) = oneshot::channel();
150 let res = async move {
151 let req = self.req.to_proto().map_err(Error::deser)?;
152 let write = WriteRequest { req, resp_tx };
153 let _ = self.req_tx.send(write);
154 let resp = resp_rx
155 .await
156 .map_err(|_| AppendError::UnexpectedEndOfStream)??;
157 let resp = resp.cnv().map_err(Error::ser)?;
158 to_result(resp)
159 }
160 .await;
161 let _ = tx.send(res);
162 });
163 AppendFuture::new(rx)
164 }
165}
166
167#[cfg(test)]
168mod tests {
169 use super::*;
170 use crate::google::cloud::bigquery::storage::v1;
171 use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
172 AppendResult, Response,
173 };
174 use crate::model::TableSchema;
175
176 #[tokio::test]
177 async fn success() -> anyhow::Result<()> {
178 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
179 let req = AppendRowsRequest::new().set_write_stream(write_stream());
180
181 let builder = Append::new(req_tx, req);
182 let handle = tokio::spawn(async move { builder.send().await });
183
184 let write = req_rx.recv().await.expect("should receive request");
186 assert_eq!(write.req.write_stream, write_stream());
187
188 let resp = v1::AppendRowsResponse {
190 response: Some(Response::AppendResult(AppendResult::default())),
191 write_stream: write_stream(),
192 updated_schema: Some(v1::TableSchema::default()),
193 ..Default::default()
194 };
195 write
196 .resp_tx
197 .send(Ok(resp))
198 .expect("sending on channel always succeeds");
199
200 let resp = handle.await??;
201 assert_eq!(resp.offset, None);
202 assert_eq!(resp.updated_schema, Some(TableSchema::default()));
203 Ok(())
204 }
205
206 #[tokio::test]
207 async fn stream_closed() -> anyhow::Result<()> {
208 let (req_tx, req_rx) = mpsc::unbounded_channel();
209 let req = AppendRowsRequest::new().set_write_stream(write_stream());
210
211 let builder = Append::new(req_tx, req);
212 let handle = tokio::spawn(async move { builder.send().await });
213
214 drop(req_rx);
216
217 let err = handle.await?.expect_err("should return an error");
218 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
219 Ok(())
220 }
221
222 #[tokio::test]
223 async fn rpc_error() -> anyhow::Result<()> {
224 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
225 let req = AppendRowsRequest::new().set_write_stream(write_stream());
226
227 let builder = Append::new(req_tx, req);
228 let handle = tokio::spawn(async move { builder.send().await });
229
230 let write = req_rx.recv().await.expect("should receive request");
232 let append_err: AppendError = Error::io("fail").into();
233 write
234 .resp_tx
235 .send(Err(append_err))
236 .expect("sending on channel always succeeds");
237
238 let err = handle.await?.expect_err("should return an error");
239 assert!(matches!(err, AppendError::Rpc { source: _ }));
240 Ok(())
241 }
242
243 #[tokio::test]
244 async fn row_errors() -> anyhow::Result<()> {
245 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
246 let req = AppendRowsRequest::new().set_write_stream(write_stream());
247
248 let builder = Append::new(req_tx, req);
249 let handle = tokio::spawn(async move { builder.send().await });
250
251 let write = req_rx.recv().await.expect("should receive request");
252
253 let row_error = v1::RowError {
254 index: 42,
255 code: v1::row_error::RowErrorCode::FieldsError as i32,
256 message: "fail".to_string(),
257 };
258 let resp = v1::AppendRowsResponse {
259 row_errors: vec![row_error],
260 write_stream: write_stream(),
261 ..Default::default()
262 };
263 write
264 .resp_tx
265 .send(Ok(resp))
266 .expect("sending on channel always succeeds");
267
268 let err = handle.await?.expect_err("should return an error");
269 assert!(matches!(err, AppendError::RowErrors(_)));
270 Ok(())
271 }
272
273 #[tokio::test]
274 async fn offset_success() -> anyhow::Result<()> {
275 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
276 let req = AppendRowsRequest::new().set_write_stream(write_stream());
277
278 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
279 let future = builder.send();
280
281 let write = req_rx.recv().await.expect("should receive request");
282 assert_eq!(write.req.offset, Some(100));
283
284 let resp = v1::AppendRowsResponse {
285 response: Some(Response::AppendResult(AppendResult::default())),
286 write_stream: write_stream(),
287 ..Default::default()
289 };
290 write
291 .resp_tx
292 .send(Ok(resp))
293 .expect("sending on channel always succeeds");
294
295 let resp = future.await?;
296 assert_eq!(resp.offset, None);
297 assert_eq!(resp.updated_schema, None);
298 Ok(())
299 }
300
301 #[tokio::test]
302 async fn offset_stream_closed() -> anyhow::Result<()> {
303 let (req_tx, req_rx) = mpsc::unbounded_channel();
304 let req = AppendRowsRequest::new().set_write_stream(write_stream());
305
306 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
307 let future = builder.send();
308
309 drop(req_rx);
310
311 let err = future.await.expect_err("should return an error");
312 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
313 Ok(())
314 }
315
316 #[tokio::test]
317 async fn offset_rpc_error() -> anyhow::Result<()> {
318 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
319 let req = AppendRowsRequest::new().set_write_stream(write_stream());
320
321 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
322 let future = builder.send();
323
324 let write = req_rx.recv().await.expect("should receive request");
326 let append_err: AppendError = Error::io("fail").into();
327 write
328 .resp_tx
329 .send(Err(append_err))
330 .expect("sending on channel always succeeds");
331
332 let err = future.await.expect_err("should return an error");
333 assert!(matches!(err, AppendError::Rpc { source: _ }));
334 Ok(())
335 }
336
337 #[tokio::test]
338 async fn offset_row_errors() -> anyhow::Result<()> {
339 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
340 let req = AppendRowsRequest::new().set_write_stream(write_stream());
341
342 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
343 let future = builder.send();
344
345 let write = req_rx.recv().await.expect("should receive request");
346
347 let row_error = v1::RowError {
348 index: 42,
349 code: v1::row_error::RowErrorCode::FieldsError as i32,
350 message: "fail".to_string(),
351 };
352 let resp = v1::AppendRowsResponse {
353 row_errors: vec![row_error],
354 write_stream: write_stream(),
355 ..Default::default()
356 };
357 write
358 .resp_tx
359 .send(Ok(resp))
360 .expect("sending on channel always succeeds");
361
362 let err = future.await.expect_err("should return an error");
363 assert!(matches!(err, AppendError::RowErrors(_)));
364 Ok(())
365 }
366
367 #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
368 async fn synchronous_queueing() -> anyhow::Result<()> {
369 const NUM_WRITES: i64 = 1000;
370 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
371 let write_handle = tokio::spawn(async move {
372 let mut writes = tokio::task::JoinSet::new();
373 for i in 0..NUM_WRITES {
374 writes.spawn(
375 AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
376 .set_offset(i)
377 .send(),
378 );
379 }
380 let _ = writes.join_all().await;
381 });
382
383 for i in 0..NUM_WRITES {
384 let write = req_rx.recv().await.expect("should receive request");
385 assert_eq!(write.req.offset, Some(i), "received out of order write");
386 }
387 write_handle.await?;
388 Ok(())
389 }
390
391 fn write_stream() -> String {
392 "projects/p/datasets/d/tables/t/streams/_default".to_string()
393 }
394}