Skip to main content

google_cloud_bigquery/write/
append_future.rs

1// Copyright 2026 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     https://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// A future that resolves to the result of an async append operation.
23///
24/// This future represents a write request that has already been queued by the
25/// client library to send over the network. Awaiting this future yields the server's acknowledgment
26/// or an error if the write fails.
27///
28/// The underlying operation begins immediately and runs independently in the
29/// background, even if this future is dropped or never awaited.
30#[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 the sender immediately
105        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}