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