google_cloud_bigquery/write/builder/
append.rs1use super::super::append_future::AppendFuture;
16use super::super::dispatcher::Dispatcher;
17use crate::model::AppendRowsRequest;
18use std::sync::Arc;
19use tokio::sync::oneshot;
20
21#[derive(Clone, Debug)]
23pub struct Append {
24 inner: Arc<Dispatcher>,
25 pub(crate) req: AppendRowsRequest,
26}
27
28impl Append {
29 pub(crate) fn new(inner: Arc<Dispatcher>, req: AppendRowsRequest) -> Self {
30 Self { inner, req }
31 }
32
33 pub fn send(self) -> AppendFuture {
57 let (tx, rx) = oneshot::channel();
58 tokio::spawn(async move {
59 let res = self.inner.send(self.req).await;
60 let _ = tx.send(res);
61 });
62 AppendFuture::new(rx)
63 }
64}
65
66#[cfg(test)]
67mod tests {
68 use super::*;
69 use crate::Error;
70 use crate::error::AppendError;
71 use crate::google::cloud::bigquery::storage::v1;
72 use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
73 AppendResult, Response,
74 };
75 use crate::model::TableSchema;
76 use crate::write::test::*;
77 use tokio::sync::mpsc;
78
79 #[tokio::test]
80 async fn success() -> anyhow::Result<()> {
81 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
82 let dispatcher = test_dispatcher(req_tx).await?;
83 let req = AppendRowsRequest::new().set_write_stream(write_stream());
84
85 let builder = Append::new(dispatcher, req);
86 let handle = tokio::spawn(async move { builder.send().await });
87
88 let write = req_rx.recv().await.expect("should receive request");
90 assert_eq!(write.req.write_stream, write_stream());
91
92 let resp = v1::AppendRowsResponse {
94 response: Some(Response::AppendResult(AppendResult::default())),
95 write_stream: write_stream(),
96 updated_schema: Some(v1::TableSchema::default()),
97 ..Default::default()
98 };
99 write
100 .resp_tx
101 .send(Ok(resp))
102 .expect("sending on channel always succeeds");
103
104 let resp = handle.await??;
105 assert_eq!(resp.offset, None);
106 assert_eq!(resp.updated_schema, Some(TableSchema::default()));
107 Ok(())
108 }
109
110 #[tokio::test]
111 async fn stream_closed() -> anyhow::Result<()> {
112 let (req_tx, req_rx) = mpsc::unbounded_channel();
113 let dispatcher = test_dispatcher(req_tx).await?;
114 let req = AppendRowsRequest::new().set_write_stream(write_stream());
115
116 let builder = Append::new(dispatcher, req);
117 let handle = tokio::spawn(async move { builder.send().await });
118
119 drop(req_rx);
121
122 let err = handle.await?.expect_err("should return an error");
123 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
124 Ok(())
125 }
126
127 #[tokio::test]
128 async fn rpc_error() -> anyhow::Result<()> {
129 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
130 let dispatcher = test_dispatcher(req_tx).await?;
131 let req = AppendRowsRequest::new().set_write_stream(write_stream());
132
133 let builder = Append::new(dispatcher, req);
134 let handle = tokio::spawn(async move { builder.send().await });
135
136 let write = req_rx.recv().await.expect("should receive request");
138 let append_err: AppendError = Error::io("fail").into();
139 write
140 .resp_tx
141 .send(Err(append_err))
142 .expect("sending on channel always succeeds");
143
144 let err = handle.await?.expect_err("should return an error");
145 assert!(matches!(err, AppendError::Rpc { source: _ }));
146 Ok(())
147 }
148
149 #[tokio::test]
150 async fn row_errors() -> anyhow::Result<()> {
151 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
152 let dispatcher = test_dispatcher(req_tx).await?;
153 let req = AppendRowsRequest::new().set_write_stream(write_stream());
154
155 let builder = Append::new(dispatcher, req);
156 let handle = tokio::spawn(async move { builder.send().await });
157
158 let write = req_rx.recv().await.expect("should receive request");
159
160 let row_error = v1::RowError {
161 index: 42,
162 code: v1::row_error::RowErrorCode::FieldsError as i32,
163 message: "fail".to_string(),
164 };
165 let resp = v1::AppendRowsResponse {
166 row_errors: vec![row_error],
167 write_stream: write_stream(),
168 ..Default::default()
169 };
170 write
171 .resp_tx
172 .send(Ok(resp))
173 .expect("sending on channel always succeeds");
174
175 let err = handle.await?.expect_err("should return an error");
176 assert!(matches!(err, AppendError::RowErrors { .. }));
177 Ok(())
178 }
179}