Skip to main content

google_cloud_bigquery/write/
append_response.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::error::{AppendError, AppendResult};
16use crate::Error;
17use crate::model::append_rows_response::Response;
18use crate::model::{AppendRowsResponse, TableSchema};
19
20/// The return type of an `append()` operation.
21#[derive(Clone, Debug, Default, PartialEq)]
22#[non_exhaustive]
23pub struct AppendResponse {
24    /// The row offset at which the last append occurred. The offset will not be
25    /// set if appending using the default stream.
26    pub offset: Option<i64>,
27
28    /// If set, the table schema has changed.
29    ///
30    /// Note that this notification is best effort. Changing a table schema can
31    /// take several minutes to propagate.
32    ///
33    /// The client library does not use this information to modify any internal
34    /// state. It only forwards the notification to the application, which
35    /// should react accordingly (if necessary).
36    pub updated_schema: Option<TableSchema>,
37}
38
39pub(crate) fn to_result(resp: AppendRowsResponse) -> AppendResult<AppendResponse> {
40    let mut status = None;
41    let offset = match resp.response {
42        None => None,
43        Some(Response::AppendResult(r)) => r.offset,
44        Some(Response::Error(s)) => {
45            status = Some((*s).into());
46            None
47        }
48    };
49
50    if !resp.row_errors.is_empty() {
51        return Err(AppendError::RowErrors {
52            status: status.unwrap_or_default(),
53            row_errors: resp.row_errors,
54        });
55    }
56    if let Some(status) = status {
57        return Err(Error::service(status).into());
58    }
59
60    Ok(AppendResponse {
61        offset,
62        updated_schema: resp.updated_schema,
63    })
64}
65
66#[cfg(test)]
67mod tests {
68    use super::*;
69    use crate::model::RowError;
70    use crate::model::append_rows_response::AppendResult;
71    use crate::model::row_error::RowErrorCode;
72    use google_cloud_gax::error::rpc::Code;
73    use google_cloud_rpc::model::Status as RpcStatus;
74
75    fn schema() -> TableSchema {
76        TableSchema::new()
77    }
78
79    fn row_error(index: i64) -> RowError {
80        RowError::new()
81            .set_index(index)
82            .set_code(RowErrorCode::FieldsError)
83            .set_message("fail")
84    }
85
86    #[test]
87    fn success() -> anyhow::Result<()> {
88        let resp = AppendRowsResponse::new()
89            .set_append_result(AppendResult::new().set_offset(42))
90            .set_updated_schema(schema());
91
92        let res = to_result(resp)?;
93        assert_eq!(res.offset, Some(42));
94        assert_eq!(res.updated_schema, Some(schema()));
95        Ok(())
96    }
97
98    #[test]
99    fn rpc_error() {
100        let resp = AppendRowsResponse::new().set_error(
101            RpcStatus::new()
102                .set_code(Code::InvalidArgument as i32)
103                .set_message("fail"),
104        );
105
106        let err = to_result(resp).expect_err("should error");
107        let AppendError::Rpc { source } = err else {
108            panic!("Expected AppendError::Rpc, got {:?}", err);
109        };
110        let status = source.status().expect("status should be set");
111        assert_eq!(status.code, Code::InvalidArgument);
112        assert_eq!(status.message, "fail");
113    }
114
115    #[test]
116    fn row_errors() {
117        let resp = AppendRowsResponse::new()
118            .set_error(
119                RpcStatus::new()
120                    .set_code(Code::InvalidArgument as i32)
121                    .set_message("Errors found while processing rows."),
122            )
123            .set_row_errors(vec![row_error(1), row_error(2)]);
124
125        let err = to_result(resp).expect_err("should error");
126        let AppendError::RowErrors { status, row_errors } = err else {
127            panic!("Expected AppendError::RowErrors, got {:?}", err);
128        };
129        assert_eq!(status.code, Code::InvalidArgument);
130        assert_eq!(status.message, "Errors found while processing rows.");
131        assert_eq!(row_errors, vec![row_error(1), row_error(2)]);
132    }
133
134    #[test]
135    fn row_errors_without_status() {
136        let resp = AppendRowsResponse::new().set_row_errors(vec![row_error(1), row_error(2)]);
137
138        let err = to_result(resp).expect_err("should error");
139        let AppendError::RowErrors { status, row_errors } = err else {
140            panic!("Expected AppendError::RowErrors, got {:?}", err);
141        };
142        assert_eq!(status.code, Code::Unknown);
143        assert_eq!(row_errors, vec![row_error(1), row_error(2)]);
144    }
145
146    #[test]
147    fn unset_response() -> anyhow::Result<()> {
148        let resp = AppendRowsResponse::new().set_updated_schema(schema());
149        let res = to_result(resp)?;
150        assert_eq!(res.offset, None);
151        assert_eq!(res.updated_schema, Some(schema()));
152        Ok(())
153    }
154}