google_cloud_bigquery/write/
committed.rs1use super::base::BaseWriter;
16use super::format::DataFormat;
17use crate::Result;
18use crate::model::FinalizeWriteStreamResponse;
19use crate::write::builder::AppendWithOffset;
20use crate::write::transport::Transport;
21use std::sync::Arc;
22
23#[derive(Debug)]
32pub struct CommittedWriter<F> {
33 pub(crate) inner: BaseWriter<F>,
34}
35
36impl<F> CommittedWriter<F>
37where
38 F: DataFormat,
39{
40 pub(crate) fn new(inner: Arc<Transport>, write_stream: String, format: F) -> Self {
41 Self {
42 inner: BaseWriter::new(inner, write_stream, format),
43 }
44 }
45
46 pub fn write_stream(&self) -> &str {
48 &self.inner.write_stream
49 }
50
51 pub fn append(&self, rows: F::Rows) -> AppendWithOffset {
53 AppendWithOffset::new(
54 self.inner.runner.req_tx.clone(),
55 self.inner.append_request(rows),
56 )
57 }
58
59 pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
61 self.inner.finalize().await
62 }
63}
64
65#[cfg(test)]
66mod tests {
67 use super::*;
68 use crate::error::AppendError;
69 use crate::write::test::*;
70 use bigquery_grpc_mock::{MockBigQueryWrite, start};
71 use gaxi::grpc::tonic::Response as TonicResponse;
72 use tokio::sync::mpsc;
73
74 #[tokio::test]
75 async fn basic_success() -> anyhow::Result<()> {
76 let (response_tx, response_rx) = mpsc::channel(10);
77
78 let mut mock = MockBigQueryWrite::new();
79 mock.expect_append_rows()
80 .return_once(|_| Ok(TonicResponse::from(response_rx)));
81
82 mock.expect_finalize_write_stream().return_once(|_| {
83 Ok(TonicResponse::new(
84 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default(),
85 ))
86 });
87
88 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
89 let transport = Arc::new(test_transport(endpoint).await?);
90
91 let writer = CommittedWriter::new(transport, write_stream(), format());
92 assert_eq!(writer.write_stream(), write_stream());
93
94 response_tx.send(Ok(convert(&test_response(1)))).await?;
95 let resp = writer.append(rows(1)).send().await?;
96 assert_eq!(resp.offset, Some(1));
97
98 response_tx.send(Ok(convert(&test_response(2)))).await?;
99 let resp = writer.append(rows(2)).send().await?;
100 assert_eq!(resp.offset, Some(2));
101
102 response_tx.send(Ok(convert(&test_response(3)))).await?;
103 let resp = writer.append(rows(3)).send().await?;
104 assert_eq!(resp.offset, Some(3));
105
106 drop(response_tx);
107 let err = writer.append(rows(4)).send().await.expect_err("channel");
108 assert!(matches!(err, AppendError::UnexpectedEndOfStream));
109
110 writer.finalize().await?;
112
113 Ok(())
114 }
115}