1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
// Copyright Materialize, Inc. and contributors. All rights reserved.
//
// Use of this software is governed by the Business Source License
// included in the LICENSE file.
//
// As of the Change Date specified in that file, in accordance with
// the Business Source License, use of this software will be governed
// by the Apache License, Version 2.0.

use mz_interchange::protobuf::{DecodedDescriptors, Decoder};
use mz_ore::error::ErrorExt;
use mz_repr::Row;
use mz_storage_types::errors::DecodeErrorKind;
use mz_storage_types::sources::encoding::ProtobufEncoding;

#[derive(Debug)]
pub struct ProtobufDecoderState {
    decoder: Decoder,
    events_success: i64,
    events_error: i64,
}

impl ProtobufDecoderState {
    pub fn new(
        ProtobufEncoding {
            descriptors,
            message_name,
            confluent_wire_format,
        }: ProtobufEncoding,
    ) -> Result<Self, anyhow::Error> {
        let descriptors = DecodedDescriptors::from_bytes(&descriptors, message_name)
            .expect("descriptors provided to protobuf source are pre-validated");
        Ok(ProtobufDecoderState {
            decoder: Decoder::new(descriptors, confluent_wire_format)?,
            events_success: 0,
            events_error: 0,
        })
    }
    pub fn get_value(&mut self, bytes: &[u8]) -> Option<Result<Row, DecodeErrorKind>> {
        match self.decoder.decode(bytes) {
            Ok(row) => {
                if let Some(row) = row {
                    self.events_success += 1;
                    Some(Ok(row))
                } else {
                    self.events_error += 1;
                    Some(Err(DecodeErrorKind::Text(
                        "protobuf deserialization returned None".into(),
                    )))
                }
            }
            Err(err) => {
                self.events_error += 1;
                Some(Err(DecodeErrorKind::Text(
                    format!(
                        "protobuf deserialization error: {}",
                        err.display_with_causes()
                    )
                    .into(),
                )))
            }
        }
    }
}