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