google_cloud_bigquery/write/builder/
append_with_offset.rs1use super::super::append_future::AppendFuture;
16use super::super::append_response::to_result;
17use super::super::error::AppendError;
18use super::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)]
30pub struct AppendWithOffset {
31 req_tx: mpsc::UnboundedSender<WriteRequest>,
32 pub(crate) req: AppendRowsRequest,
33}
34
35impl AppendWithOffset {
36 pub(crate) fn new(req_tx: mpsc::UnboundedSender<WriteRequest>, req: AppendRowsRequest) -> Self {
37 Self { req_tx, req }
38 }
39
40 pub fn set_offset(mut self, offset: i64) -> Self {
63 self.req.offset = Some(offset);
64 self
65 }
66
67 pub fn send(self) -> AppendFuture {
94 let (resp_tx, resp_rx) = oneshot::channel();
95 let req = match self.req.to_proto().map_err(Error::ser) {
96 Ok(req) => req,
97 Err(e) => return AppendFuture::from_future(async move { Err(e.into()) }),
98 };
99 let write = WriteRequest { req, resp_tx };
100 if self.req_tx.send(write).is_err() {
101 return AppendFuture::from_future(
102 async move { Err(AppendError::UnexpectedEndOfStream) },
103 );
104 }
105 AppendFuture::from_future(async move {
106 let resp = resp_rx
107 .await
108 .map_err(|_| AppendError::UnexpectedEndOfStream)??;
109 let resp = resp.cnv().map_err(Error::deser)?;
110 to_result(resp)
111 })
112 }
113}
114
115#[cfg(test)]
116mod tests {
117 use super::*;
118 use crate::google::cloud::bigquery::storage::v1;
119 use crate::google::cloud::bigquery::storage::v1::append_rows_response::{
120 AppendResult, Response,
121 };
122 use crate::write::test::*;
123
124 #[tokio::test]
125 async fn success() -> anyhow::Result<()> {
126 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
127 let req = AppendRowsRequest::new().set_write_stream(write_stream());
128
129 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
130 let future = builder.send();
131
132 let write = req_rx.recv().await.expect("should receive request");
133 assert_eq!(write.req.offset, Some(100));
134
135 let resp = v1::AppendRowsResponse {
136 response: Some(Response::AppendResult(AppendResult::default())),
137 write_stream: write_stream(),
138 ..Default::default()
140 };
141 write
142 .resp_tx
143 .send(Ok(resp))
144 .expect("sending on channel always succeeds");
145
146 let resp = future.await?;
147 assert_eq!(resp.offset, None);
148 assert_eq!(resp.updated_schema, None);
149 Ok(())
150 }
151
152 #[tokio::test]
153 async fn stream_closed() -> anyhow::Result<()> {
154 let (req_tx, req_rx) = mpsc::unbounded_channel();
155 let req = AppendRowsRequest::new().set_write_stream(write_stream());
156
157 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
158 let future = builder.send();
159
160 drop(req_rx);
161
162 let err = future.await.expect_err("should return an error");
163 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
164 Ok(())
165 }
166
167 #[tokio::test]
168 async fn rpc_error() -> anyhow::Result<()> {
169 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
170 let req = AppendRowsRequest::new().set_write_stream(write_stream());
171
172 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
173 let future = builder.send();
174
175 let write = req_rx.recv().await.expect("should receive request");
177 let append_err: AppendError = Error::io("fail").into();
178 write
179 .resp_tx
180 .send(Err(append_err))
181 .expect("sending on channel always succeeds");
182
183 let err = future.await.expect_err("should return an error");
184 assert!(matches!(err, AppendError::Rpc { source: _ }));
185 Ok(())
186 }
187
188 #[tokio::test]
189 async fn row_errors() -> anyhow::Result<()> {
190 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
191 let req = AppendRowsRequest::new().set_write_stream(write_stream());
192
193 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
194 let future = builder.send();
195
196 let write = req_rx.recv().await.expect("should receive request");
197
198 let row_error = v1::RowError {
199 index: 42,
200 code: v1::row_error::RowErrorCode::FieldsError as i32,
201 message: "fail".to_string(),
202 };
203 let resp = v1::AppendRowsResponse {
204 row_errors: vec![row_error],
205 write_stream: write_stream(),
206 ..Default::default()
207 };
208 write
209 .resp_tx
210 .send(Ok(resp))
211 .expect("sending on channel always succeeds");
212
213 let err = future.await.expect_err("should return an error");
214 assert!(matches!(err, AppendError::RowErrors { .. }));
215 Ok(())
216 }
217
218 #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
219 async fn synchronous_queueing() -> anyhow::Result<()> {
220 const NUM_WRITES: i64 = 1000;
221 let (req_tx, mut req_rx) = mpsc::unbounded_channel();
222 let write_handle = tokio::spawn(async move {
223 let mut writes = tokio::task::JoinSet::new();
224 for i in 0..NUM_WRITES {
225 writes.spawn(
226 AppendWithOffset::new(req_tx.clone(), AppendRowsRequest::new())
227 .set_offset(i)
228 .send(),
229 );
230 }
231 let _ = writes.join_all().await;
232 });
233
234 for i in 0..NUM_WRITES {
235 let write = req_rx.recv().await.expect("should receive request");
236 assert_eq!(write.req.offset, Some(i), "received out of order write");
237 }
238 write_handle.await?;
239 Ok(())
240 }
241
242 #[tokio::test]
243 async fn send_when_req_tx_closed_returns_unexpected_end_of_stream() {
244 let (req_tx, req_rx) = mpsc::unbounded_channel();
245 drop(req_rx);
246
247 let req = AppendRowsRequest::new().set_write_stream(write_stream());
248 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
249 let future = builder.send();
250
251 let err = future
252 .await
253 .expect_err("should return unexpected end of stream");
254 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
255 }
256
257 #[tokio::test]
258 async fn send_serialization_error() {
259 use crate::model::append_rows_request::MissingValueInterpretation;
260
261 let (req_tx, _req_rx) = mpsc::unbounded_channel();
262 let invalid: MissingValueInterpretation = serde_json::from_str("\"INVALID\"").unwrap();
263 let req = AppendRowsRequest::new().set_default_missing_value_interpretation(invalid);
264 let builder = AppendWithOffset::new(req_tx, req).set_offset(100);
265 let future = builder.send();
266
267 let err = future.await.expect_err("should return serialization error");
268 assert!(matches!(err, AppendError::Rpc { .. }));
269 }
270}