google_cloud_bigquery/write/
append_future.rs1use super::append_response::AppendResponse;
16use super::error::{AppendError, AppendResult};
17use std::future::Future;
18use std::pin::Pin;
19use std::task::{Context, Poll};
20use tokio::sync::oneshot;
21
22#[derive(Debug)]
31pub struct AppendFuture {
32 inner: Inner,
33}
34
35enum Inner {
36 Rx(oneshot::Receiver<AppendResult<AppendResponse>>),
37 Boxed(Pin<Box<dyn Future<Output = AppendResult<AppendResponse>> + Send>>),
38}
39
40impl std::fmt::Debug for Inner {
41 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
42 match self {
43 Self::Rx(rx) => f.debug_tuple("Rx").field(rx).finish(),
44 Self::Boxed(_) => f.debug_tuple("Boxed").finish(),
45 }
46 }
47}
48
49impl AppendFuture {
50 pub(crate) fn new(rx: oneshot::Receiver<AppendResult<AppendResponse>>) -> Self {
51 Self {
52 inner: Inner::Rx(rx),
53 }
54 }
55
56 pub(crate) fn from_future<F>(fut: F) -> Self
57 where
58 F: Future<Output = AppendResult<AppendResponse>> + Send + 'static,
59 {
60 Self {
61 inner: Inner::Boxed(Box::pin(fut)),
62 }
63 }
64}
65
66impl Future for AppendFuture {
67 type Output = AppendResult<AppendResponse>;
68
69 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
70 match &mut self.inner {
71 Inner::Rx(rx) => {
72 let result = std::task::ready!(Pin::new(rx).poll(cx));
73 match result {
74 Ok(res) => Poll::Ready(res),
75 Err(_) => Poll::Ready(Err(AppendError::UnexpectedEndOfStream)),
76 }
77 }
78 Inner::Boxed(fut) => fut.as_mut().poll(cx),
79 }
80 }
81}
82
83#[cfg(test)]
84mod tests {
85 use super::*;
86 use crate::model::TableSchema;
87
88 #[tokio::test]
89 async fn happy_path() {
90 let (tx, rx) = oneshot::channel();
91 let _ = tx.send(Ok(AppendResponse {
92 offset: None,
93 updated_schema: Some(TableSchema::default()),
94 }));
95 let future = AppendFuture::new(rx);
96 let resp = future.await.expect("should succeed");
97 assert_eq!(resp.offset, None);
98 assert_eq!(resp.updated_schema, Some(TableSchema::default()));
99 }
100
101 #[tokio::test]
102 async fn dropped_sender() {
103 let (tx, rx) = oneshot::channel::<AppendResult<AppendResponse>>();
104 drop(tx);
106
107 let future = AppendFuture::new(rx);
108 let err = future
109 .await
110 .expect_err("should return unexpected end of stream");
111 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
112 }
113
114 #[tokio::test]
115 async fn channel_returns_error() {
116 let (tx, rx) = oneshot::channel();
117 let _ = tx.send(Err(AppendError::UnexpectedEndOfStream));
118 let future = AppendFuture::new(rx);
119 let err = future.await.expect_err("should return error from task");
120 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
121 }
122
123 #[tokio::test]
124 async fn from_future_success() {
125 let future = AppendFuture::from_future(async {
126 Ok(AppendResponse {
127 offset: Some(7),
128 updated_schema: None,
129 })
130 });
131 let resp = future.await.expect("should succeed");
132 assert_eq!(resp.offset, Some(7));
133 }
134
135 #[tokio::test]
136 async fn from_future_error() {
137 let future = AppendFuture::from_future(async { Err(AppendError::UnexpectedEndOfStream) });
138 let err = future.await.expect_err("should return error");
139 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
140 }
141
142 #[test]
143 fn debug_format() {
144 let (_, rx) = oneshot::channel();
145 let future_rx = AppendFuture::new(rx);
146 assert!(format!("{future_rx:?}").contains("Rx"));
147
148 let future_boxed = AppendFuture::from_future(async { Ok(AppendResponse::default()) });
149 assert!(format!("{future_boxed:?}").contains("Boxed"));
150 }
151}