use std::{io::Read, marker::PhantomData};
use bon::bon;
use serde::de::DeserializeOwned;
use crate::{
AvroResult, AvroSchema, Schema,
decode::decode_internal,
schema::{ResolvedOwnedSchema, ResolvedSchema},
serde::deser_schema::{Config, SchemaAwareDeserializer},
types::Value,
util::is_human_readable,
};
pub struct GenericDatumReader<'s> {
writer: &'s Schema,
resolved: ResolvedSchema<'s>,
reader: Option<(&'s Schema, ResolvedSchema<'s>)>,
human_readable: bool,
}
#[bon]
impl<'s> GenericDatumReader<'s> {
#[builder]
pub fn new(
#[builder(start_fn)]
writer_schema: &'s Schema,
resolved_writer_schemata: Option<ResolvedSchema<'s>>,
reader_schema: Option<&'s Schema>,
resolved_reader_schemata: Option<ResolvedSchema<'s>>,
#[builder(default = is_human_readable())]
human_readable: bool,
) -> AvroResult<Self> {
let resolved_writer_schemata = if let Some(resolved) = resolved_writer_schemata {
resolved
} else {
ResolvedSchema::try_from(writer_schema)?
};
let reader = if let Some(reader) = reader_schema {
if let Some(resolved) = resolved_reader_schemata {
Some((reader, resolved))
} else {
Some((reader, ResolvedSchema::try_from(reader)?))
}
} else {
None
};
Ok(Self {
writer: writer_schema,
resolved: resolved_writer_schemata,
reader,
human_readable,
})
}
}
impl<'s, S: generic_datum_reader_builder::State> GenericDatumReaderBuilder<'s, S> {
pub fn writer_schemata(
self,
schemata: Vec<&'s Schema>,
) -> AvroResult<
GenericDatumReaderBuilder<'s, generic_datum_reader_builder::SetResolvedWriterSchemata<S>>,
>
where
S::ResolvedWriterSchemata: generic_datum_reader_builder::IsUnset,
{
let resolved = ResolvedSchema::new_with_schemata(schemata)?;
Ok(self.resolved_writer_schemata(resolved))
}
pub fn reader_schemata(
self,
schemata: Vec<&'s Schema>,
) -> AvroResult<
GenericDatumReaderBuilder<'s, generic_datum_reader_builder::SetResolvedReaderSchemata<S>>,
>
where
S::ResolvedReaderSchemata: generic_datum_reader_builder::IsUnset,
S::ReaderSchema: generic_datum_reader_builder::IsSet,
{
let resolved = ResolvedSchema::new_with_schemata(schemata)?;
Ok(self.resolved_reader_schemata(resolved))
}
}
impl<'s> GenericDatumReader<'s> {
pub fn read_value<R: Read>(&self, reader: &mut R) -> AvroResult<Value> {
let value = decode_internal(self.writer, self.resolved.get_names(), None, reader)?;
if let Some((reader, resolved)) = &self.reader {
value.resolve_internal(reader, resolved.get_names(), None, None)
} else {
Ok(value)
}
}
pub fn read_deser<T: DeserializeOwned>(&self, reader: &mut impl Read) -> AvroResult<T> {
if let Some((_, _)) = &self.reader {
panic!("Schema aware deserialisation does not resolve schemas yet");
} else {
T::deserialize(SchemaAwareDeserializer::new(
reader,
self.writer,
Config {
names: self.resolved.get_names(),
human_readable: self.human_readable,
},
)?)
}
}
}
pub struct SpecificDatumReader<T: AvroSchema> {
resolved: ResolvedOwnedSchema,
human_readable: bool,
phantom: PhantomData<T>,
}
#[bon]
impl<T: AvroSchema> SpecificDatumReader<T> {
#[builder]
pub fn new(#[builder(default = is_human_readable())] human_readable: bool) -> AvroResult<Self> {
Ok(Self {
resolved: T::get_schema().try_into()?,
human_readable,
phantom: PhantomData,
})
}
}
impl<T: AvroSchema + DeserializeOwned> SpecificDatumReader<T> {
pub fn read<R: Read>(&self, reader: &mut R) -> AvroResult<T> {
T::deserialize(SchemaAwareDeserializer::new(
reader,
self.resolved.get_root_schema(),
Config {
names: self.resolved.get_names(),
human_readable: self.human_readable,
},
)?)
}
}
#[deprecated(since = "0.22.0", note = "Use `GenericDatumReader` instead")]
pub fn from_avro_datum<R: Read>(
writer_schema: &Schema,
reader: &mut R,
reader_schema: Option<&Schema>,
) -> AvroResult<Value> {
GenericDatumReader::builder(writer_schema)
.maybe_reader_schema(reader_schema)
.build()?
.read_value(reader)
}
#[deprecated(since = "0.22.0", note = "Use `GenericDatumReader` instead")]
pub fn from_avro_datum_schemata<R: Read>(
writer_schema: &Schema,
writer_schemata: Vec<&Schema>,
reader: &mut R,
reader_schema: Option<&Schema>,
) -> AvroResult<Value> {
GenericDatumReader::builder(writer_schema)
.writer_schemata(writer_schemata)?
.maybe_reader_schema(reader_schema)
.build()?
.read_value(reader)
}
#[deprecated(since = "0.22.0", note = "Use `GenericDatumReader` instead")]
pub fn from_avro_datum_reader_schemata<R: Read>(
writer_schema: &Schema,
writer_schemata: Vec<&Schema>,
reader: &mut R,
reader_schema: Option<&Schema>,
reader_schemata: Vec<&Schema>,
) -> AvroResult<Value> {
GenericDatumReader::builder(writer_schema)
.writer_schemata(writer_schemata)?
.maybe_reader_schema(reader_schema)
.reader_schemata(reader_schemata)?
.build()?
.read_value(reader)
}
#[cfg(test)]
mod tests {
use apache_avro_test_helper::TestResult;
use serde::Deserialize;
use crate::{
Schema,
reader::datum::GenericDatumReader,
types::{Record, Value},
};
#[test]
fn test_from_avro_datum() -> TestResult {
let schema = Schema::parse_str(
r#"{
"type": "record",
"name": "test",
"fields": [
{
"name": "a",
"type": "long",
"default": 42
},
{
"name": "b",
"type": "string"
}
]
}"#,
)?;
let mut encoded: &'static [u8] = &[54, 6, 102, 111, 111];
let mut record = Record::new(&schema).unwrap();
record.put("a", 27i64);
record.put("b", "foo");
let expected = record.into();
let avro_datum = GenericDatumReader::builder(&schema)
.build()?
.read_value(&mut encoded)?;
assert_eq!(avro_datum, expected);
Ok(())
}
#[test]
fn test_from_avro_datum_with_union_to_struct() -> TestResult {
const TEST_RECORD_SCHEMA_3240: &str = r#"
{
"type": "record",
"name": "TestRecord3240",
"fields": [
{
"name": "a",
"type": "long",
"default": 42
},
{
"name": "b",
"type": "string"
},
{
"name": "a_nullable_array",
"type": ["null", {"type": "array", "items": {"type": "string"}}],
"default": null
},
{
"name": "a_nullable_boolean",
"type": ["null", {"type": "boolean"}],
"default": null
},
{
"name": "a_nullable_string",
"type": ["null", {"type": "string"}],
"default": null
}
]
}
"#;
#[derive(Default, Debug, Deserialize, PartialEq, Eq)]
struct TestRecord3240 {
a: i64,
b: String,
a_nullable_array: Option<Vec<String>>,
a_nullable_string: Option<String>,
}
let schema = Schema::parse_str(TEST_RECORD_SCHEMA_3240)?;
let mut encoded: &[u8] = &[54, 6, 102, 111, 111];
let error = GenericDatumReader::builder(&schema)
.build()?
.read_deser::<TestRecord3240>(&mut encoded)
.unwrap_err();
assert_eq!(
error.to_string(),
"Failed to read bytes for decoding variable length integer: failed to fill whole buffer"
);
Ok(())
}
#[test]
fn test_null_union() -> TestResult {
let schema = Schema::parse_str(r#"["null", "long"]"#)?;
let mut encoded: &'static [u8] = &[2, 0];
let avro_datum = GenericDatumReader::builder(&schema)
.build()?
.read_value(&mut encoded)?;
assert_eq!(avro_datum, Value::Union(1, Box::new(Value::Long(0))));
Ok(())
}
}