use std::collections::{BTreeMap, HashMap};
use arrow::array::RecordBatch;
use arrow::buffer::NullBuffer;
use arrow::datatypes::Field;
use arrow::{
array::{BooleanArray, UInt64Array},
error::ArrowError,
};
use itertools::Itertools as _;
use re_chunk::external::nohash_hasher::IntMap;
use re_chunk::external::re_byte_size;
use re_chunk::{ArchetypeName, ChunkError, ChunkId, ComponentIdentifier, ComponentType, Timeline};
use re_log_types::{AbsoluteTimeRange, EntityPath, StoreId, TimeType, TimelineName};
use re_types_core::{
ComponentDescriptor, FIELD_METADATA_KEY_ARCHETYPE, FIELD_METADATA_KEY_COMPONENT,
FIELD_METADATA_KEY_COMPONENT_TYPE,
};
use crate::{CodecError, CodecResult, Decodable as _, StreamFooterEntry, ToApplication as _};
#[derive(Clone, Debug, re_byte_size::SizeBytes)]
pub struct RawRrdManifest {
pub store_id: StoreId,
pub sorbet_schema: arrow::datatypes::Schema,
pub sorbet_schema_sha256: [u8; 32],
pub data: arrow::array::RecordBatch,
}
pub type RrdManifestStaticMap = IntMap<EntityPath, IntMap<ComponentIdentifier, ChunkId>>;
#[derive(Debug, Clone, Copy, PartialEq, Eq, re_byte_size::SizeBytes)]
pub struct RrdManifestTemporalMapEntry {
pub time_range: AbsoluteTimeRange,
pub num_rows: u64,
}
pub type RrdManifestTemporalMap = IntMap<
EntityPath,
IntMap<Timeline, IntMap<ComponentIdentifier, BTreeMap<ChunkId, RrdManifestTemporalMapEntry>>>,
>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RrdManifestSha256(pub [u8; 32]);
impl std::fmt::Display for RrdManifestSha256 {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "RrdManifest#{}", sha256_to_hex(&self.0))
}
}
pub fn sha256_to_hex(sha256: &[u8; 32]) -> String {
use std::fmt::Write as _;
sha256
.iter()
.fold(String::with_capacity(2 * sha256.len()), |mut hex, byte| {
write!(hex, "{byte:02x}").ok();
hex
})
}
impl RawRrdManifest {
pub fn concat(manifests: &[&Self]) -> Result<Self, ArrowError> {
re_tracing::profile_function!();
let first = manifests.first().ok_or_else(|| {
ArrowError::InvalidArgumentError("No manifests to concatenate".to_owned())
})?;
for other in &manifests[1..] {
if first.store_id != other.store_id {
return Err(ArrowError::SchemaError(
"Mismatching store_id in RawRrdManifest::concat".to_owned(),
));
}
if first.sorbet_schema_sha256 != other.sorbet_schema_sha256 {
return Err(ArrowError::SchemaError(
"Mismatching sorbet recording schemas in RawRrdManifest::concat".to_owned(),
));
}
if first.data.schema() != other.data.schema() {
re_log::debug!(
"Different schemas in the RrdManifest ({} columns in existing, {} in the new part)",
first.data.num_columns(),
other.data.num_columns(),
);
}
}
let batches: Vec<&RecordBatch> = manifests.iter().map(|m| &m.data).collect();
let data = arrow::compute::concat_batches(&first.data.schema(), batches)?;
Ok(Self {
store_id: first.store_id.clone(),
sorbet_schema: first.sorbet_schema.clone(),
sorbet_schema_sha256: first.sorbet_schema_sha256,
data,
})
}
pub fn without_recording_properties(self) -> CodecResult<Self> {
re_tracing::profile_function!();
let keep: BooleanArray = self
.col_chunk_entity_path()?
.iter()
.map(|entity_path| Some(!EntityPath::parse_forgiving(entity_path).is_property()))
.collect();
if keep.true_count() == keep.len() {
return Ok(self);
}
let data = arrow::compute::filter_record_batch(&self.data, &keep)
.map_err(CodecError::ArrowDeserialization)?;
Ok(Self { data, ..self })
}
pub fn merge(store_id: StoreId, manifests: Vec<Self>) -> CodecResult<Self> {
re_tracing::profile_function!();
if manifests.is_empty() {
return Err(CodecError::ArrowDeserialization(
ArrowError::InvalidArgumentError("cannot merge 0 manifests".to_owned()),
));
}
let parts: Vec<_> = manifests
.into_iter()
.map(|m| (m.sorbet_schema, m.data))
.collect();
let (sorbet_schema, sorbet_schema_sha256, data) = Self::merge_polymorphic_parts(parts)?;
let data = strip_null_mask_on_default_columns(data)?;
Ok(Self {
store_id,
sorbet_schema,
sorbet_schema_sha256,
data,
})
}
pub fn merge_polymorphic_parts(
parts: Vec<(arrow::datatypes::Schema, RecordBatch)>,
) -> CodecResult<(arrow::datatypes::Schema, [u8; 32], RecordBatch)> {
use re_arrow_util::{RecordBatchExt as _, concat_polymorphic_batches};
let (sorbet_schemas, data_batches): (Vec<_>, Vec<_>) = parts.into_iter().unzip();
let sorbet_schema = arrow::datatypes::Schema::try_merge(sorbet_schemas)
.map_err(CodecError::ArrowDeserialization)?;
let sorbet_schema_sha256 = Self::compute_sorbet_schema_sha256(&sorbet_schema)
.map_err(CodecError::ArrowSerialization)?;
let nullable_batches: Vec<RecordBatch> = data_batches
.into_iter()
.map(|b| b.make_nullable())
.collect();
let data = concat_polymorphic_batches(&nullable_batches)
.map_err(CodecError::ArrowDeserialization)?;
Ok((sorbet_schema, sorbet_schema_sha256, data))
}
pub fn from_rrd_bytes(rrd_bytes: &[u8]) -> CodecResult<Vec<Self>> {
let stream_footer = match crate::StreamFooter::from_rrd_bytes(rrd_bytes) {
Ok(footer) => footer,
Err(CodecError::FrameDecoding(_)) => return Ok(vec![]),
Err(err) => Err(err)?,
};
let mut manifests = Vec::new();
for entry in stream_footer.entries {
let StreamFooterEntry {
rrd_footer_byte_span_from_start_excluding_header,
crc_excluding_header,
} = entry;
let rrd_footer_byte_span = rrd_footer_byte_span_from_start_excluding_header;
let rrd_footer_byte_span = rrd_footer_byte_span
.try_cast::<usize>()
.ok_or_else(|| {
CodecError::FrameDecoding(
"RRD footer too large for native bit width".to_owned(),
)
})?
.range();
let rrd_footer_bytes = &rrd_bytes[rrd_footer_byte_span];
let crc = crate::StreamFooter::compute_crc(rrd_footer_bytes);
if crc != crc_excluding_header {
return Err(CodecError::CrcMismatch {
expected: crc_excluding_header,
got: crc,
});
}
let rrd_footer =
re_protos::log_msg::v1alpha1::RrdFooter::from_rrd_bytes(rrd_footer_bytes)?;
let new_manifests: Vec<_> = rrd_footer
.manifests
.iter()
.map(|manifest| manifest.to_application(()))
.try_collect()?;
manifests.extend(new_manifests);
}
Ok(manifests)
}
pub fn build_in_memory_from_chunks<'a>(
store_id: StoreId,
chunks: impl Iterator<Item = &'a re_chunk::Chunk>,
) -> CodecResult<Self> {
let mut rrd_manifest_builder = crate::RrdManifestBuilder::default();
let mut offset = 0;
for chunk in chunks {
let chunk_batch = chunk.to_chunk_batch()?;
use re_byte_size::SizeBytes as _;
let byte_size_uncompressed = chunk.heap_size_bytes();
let uncompressed_byte_span =
re_span::Span::from_start_len(offset, byte_size_uncompressed);
offset += byte_size_uncompressed;
rrd_manifest_builder.append(
&chunk_batch,
uncompressed_byte_span,
byte_size_uncompressed,
)?;
}
let rrd_manifest = rrd_manifest_builder.build(store_id)?;
rrd_manifest.sanity_check_cheap()?;
rrd_manifest.sanity_check_heavy()?;
Ok(rrd_manifest)
}
pub fn compute_sha256(&self) -> Result<RrdManifestSha256, ArrowError> {
re_tracing::profile_function!();
let data_ipc = {
let mut data_ipc = Vec::new();
let mut w =
arrow::ipc::writer::StreamWriter::try_new(&mut data_ipc, self.data.schema_ref())?;
w.write(&self.data)?;
data_ipc
};
use sha2::Digest as _;
let mut hash = [0u8; 32];
let mut hasher = sha2::Sha256::new();
hasher.update(&data_ipc);
hasher.finalize_into(sha2::digest::generic_array::GenericArray::from_mut_slice(
&mut hash,
));
Ok(RrdManifestSha256(hash))
}
pub fn calc_static_map(&self) -> CodecResult<RrdManifestStaticMap> {
re_tracing::profile_function!();
use re_arrow_util::ArrowArrayDowncastRef as _;
let mut per_entity: RrdManifestStaticMap = IntMap::default();
let chunk_ids = self.col_chunk_id_iter()?;
let chunk_entity_paths = self.col_chunk_entity_path_iter()?;
let chunk_is_static = self.col_chunk_is_static_iter()?;
let has_static_component_data: Vec<_> =
itertools::izip!(self.data.schema_ref().fields(), self.data.columns(),)
.filter(|(f, _c)| Self::is_index_has_static_data(f))
.map(|(f, c)| {
c.try_downcast_array_ref::<arrow::array::BooleanArray>()
.map_err(|err| {
CodecError::ArrowDeserialization(ArrowError::SchemaError(format!(
"cannot downcast column '{}': {err}",
f.name(),
)))
})
.map(|c| (f, c))
})
.try_collect()?;
for (i, (chunk_id, is_static, entity_path)) in
itertools::izip!(chunk_ids, chunk_is_static, chunk_entity_paths).enumerate()
{
if !is_static {
continue;
}
for (f, has_static_component_data) in &has_static_component_data {
let has_static_component_data = has_static_component_data.value(i);
if !has_static_component_data {
continue;
}
let Some(component) = f.metadata().get(FIELD_METADATA_KEY_COMPONENT) else {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"column '{}' is missing rerun:component metadata",
f.name()
),
}));
};
let component = ComponentIdentifier::try_new(component).map_err(|err| {
CodecError::from(ChunkError::Malformed {
reason: err.to_string(),
})
})?;
let per_component = per_entity.entry(entity_path.clone()).or_default();
per_component
.entry(component)
.and_modify(|id| *id = chunk_id)
.or_insert(chunk_id);
}
}
Ok(per_entity)
}
pub fn calc_temporal_map(&self) -> CodecResult<RrdManifestTemporalMap> {
re_tracing::profile_function!();
use re_arrow_util::ArrowArrayDowncastRef as _;
let fields = self.data.schema_ref().fields();
let columns = self.data.columns();
let indexes = fields
.iter()
.filter_map(|f| {
f.metadata()
.get(Self::FIELD_METADATA_KEY_INDEX)
.and_then(|index| {
f.metadata()
.get(FIELD_METADATA_KEY_COMPONENT)
.map(|c| (index, c, f))
})
})
.filter(|(_index, _component, field)| Self::is_index_start(field))
.collect_vec();
let mut per_entity: RrdManifestTemporalMap = Default::default();
let chunk_ids = self.col_chunk_id_iter()?;
let chunk_entity_paths = self.col_chunk_entity_path_iter()?;
let chunk_is_static = self.col_chunk_is_static_iter()?;
struct IndexColumns<'a> {
index: &'a str,
component: &'a String,
time_type: TimeType,
col_start_nulls: NullBuffer,
col_start_raw: &'a [i64],
col_end_nulls: NullBuffer,
col_end_raw: &'a [i64],
col_num_rows_raw: &'a [u64],
}
let sibling = |field: &Field, marker: &str| {
itertools::izip!(fields, columns)
.find(|(f, _col)| {
Self::has_index_marker(f, marker) && f.metadata() == field.metadata()
})
.ok_or_else(|| {
CodecError::from(ChunkError::Malformed {
reason: format!("{marker} index is missing for {}", field.name()),
})
})
};
let mut columns_per_index = Vec::<IndexColumns<'_>>::new();
for (index, component, field) in indexes {
let index = index.as_str();
if index == Self::INDEX_NAME_STATIC {
continue;
}
let (_, col_start) = sibling(field, Self::INDEX_MARKER_START)?;
let (_, col_end) = sibling(field, Self::INDEX_MARKER_END)?;
let (field_num_rows, col_num_rows) = sibling(field, Self::INDEX_MARKER_NUM_ROWS)?;
let (time_type, col_start_raw) =
TimeType::from_arrow_array(col_start).map_err(CodecError::ArrowDeserialization)?;
let (_, col_end_raw) =
TimeType::from_arrow_array(col_end).map_err(CodecError::ArrowDeserialization)?;
let col_num_rows_raw: &[u64] = col_num_rows
.try_downcast_array_ref::<UInt64Array>()
.map_err(|err| {
CodecError::ArrowDeserialization(ArrowError::SchemaError(format!(
"cannot downcast column '{}': {err}",
field_num_rows.name(),
)))
})?
.values();
let col_start_nulls = col_start
.nulls()
.cloned()
.unwrap_or_else(|| NullBuffer::new_valid(col_start.len()));
let col_end_nulls = col_end
.nulls()
.cloned()
.unwrap_or_else(|| NullBuffer::new_valid(col_end.len()));
columns_per_index.push(IndexColumns {
index,
component,
time_type,
col_start_nulls,
col_start_raw,
col_end_nulls,
col_end_raw,
col_num_rows_raw,
});
}
for (i, (chunk_id, is_static, entity_path)) in
itertools::izip!(chunk_ids, chunk_is_static, chunk_entity_paths).enumerate()
{
if is_static {
continue;
}
for columns in &columns_per_index {
let IndexColumns {
index,
component,
time_type,
col_start_nulls,
col_start_raw,
col_end_nulls,
col_end_raw,
col_num_rows_raw,
} = columns;
if !col_start_nulls.is_valid(i) || !col_end_nulls.is_valid(i) {
continue;
}
let component = ComponentIdentifier::try_new(component).map_err(|err| {
CodecError::from(ChunkError::Malformed {
reason: err.to_string(),
})
})?;
let timeline = Timeline::new(TimelineName::try_new(*index)?, *time_type);
let per_timeline = per_entity.entry(entity_path.clone()).or_default();
let per_component = per_timeline.entry(timeline).or_default();
let per_chunk = per_component.entry(component).or_default();
let start = col_start_raw[i];
let end = col_end_raw[i];
let num_rows = col_num_rows_raw[i];
let entry = RrdManifestTemporalMapEntry {
time_range: AbsoluteTimeRange::new(start, end),
num_rows,
};
per_chunk
.entry(chunk_id)
.and_modify(|tr| *tr = entry)
.or_insert(entry);
}
}
Ok(per_entity)
}
}
impl PartialEq for RawRrdManifest {
fn eq(&self, other: &Self) -> bool {
let Self {
store_id,
sorbet_schema,
sorbet_schema_sha256,
data,
} = self;
*store_id == other.store_id
&& *data == other.data
&& *sorbet_schema_sha256 == other.sorbet_schema_sha256
&& sorbet_schema.metadata() == other.sorbet_schema.metadata()
&& sorbet_schema.fields.len() == other.sorbet_schema.fields.len()
&& {
let sorted_fields = sorbet_schema.fields.iter().sorted_by_key(|f| f.name());
let other_sorted_fields = other
.sorbet_schema
.fields
.iter()
.sorted_by_key(|f| f.name());
itertools::izip!(sorted_fields, other_sorted_fields).all(|(f1, f2)| f1 == f2)
}
}
}
impl RawRrdManifest {
const FIELD_METADATA_KEY_INDEX: &str = "rerun:index";
const FIELD_METADATA_KEY_INDEX_KIND: &str = "rerun:index_kind";
const FIELD_METADATA_KEY_INDEX_MARKER: &str = "rerun:index_marker";
const INDEX_NAME_STATIC: &str = "rerun:static";
const INDEX_NAME_STATIC_SEGMENT_MANIFEST: &str = "static";
const INDEX_MARKER_START: &str = "start";
const INDEX_MARKER_END: &str = "end";
const INDEX_MARKER_NUM_ROWS: &str = "num_rows";
const INDEX_MARKER_HAS_DATA: &str = "has_data";
const INDEX_MARKER_HAS_STATIC_DATA: &str = "has_static_data";
const INDEX_MARKER_LEN: &str = "len";
pub fn is_index(field: &Field) -> bool {
field
.metadata()
.contains_key(Self::FIELD_METADATA_KEY_INDEX)
}
pub fn get_component(field: &Field) -> Option<&str> {
field
.metadata()
.get(FIELD_METADATA_KEY_COMPONENT)
.map(|s| s.as_str())
}
pub fn get_component_type(field: &Field) -> Option<&str> {
field
.metadata()
.get(FIELD_METADATA_KEY_COMPONENT_TYPE)
.map(|s| s.as_str())
}
pub fn is_index_per_component(field: &Field) -> bool {
field.metadata().contains_key(FIELD_METADATA_KEY_COMPONENT)
}
pub fn get_index_name(field: &Field) -> Option<&str> {
field
.metadata()
.get(Self::FIELD_METADATA_KEY_INDEX)
.map(|s| s.as_str())
}
pub fn get_index_kind(field: &Field) -> Option<&str> {
field
.metadata()
.get(Self::FIELD_METADATA_KEY_INDEX_KIND)
.map(|s| s.as_str())
}
pub fn is_specific_index(field: &Field, index_name: &str) -> bool {
Self::get_index_name(field) == Some(index_name)
}
pub fn is_index_static(field: &Field) -> bool {
Self::get_index_name(field).is_some_and(|name| {
name == Self::INDEX_NAME_STATIC || name == Self::INDEX_NAME_STATIC_SEGMENT_MANIFEST
})
}
pub fn has_index_marker(field: &Field, marker: &str) -> bool {
field
.metadata()
.get(Self::FIELD_METADATA_KEY_INDEX_MARKER)
.is_some_and(|m| m == marker)
|| field
.name()
.rsplit_once(':')
.is_some_and(|(_, suffix)| suffix == marker)
}
pub fn is_index_start(field: &Field) -> bool {
Self::has_index_marker(field, Self::INDEX_MARKER_START)
}
pub fn is_index_end(field: &Field) -> bool {
Self::has_index_marker(field, Self::INDEX_MARKER_END)
}
pub fn is_index_num_rows(field: &Field) -> bool {
Self::has_index_marker(field, Self::INDEX_MARKER_NUM_ROWS)
}
pub fn is_index_length(field: &Field) -> bool {
Self::has_index_marker(field, Self::INDEX_MARKER_LEN)
}
pub fn is_index_has_static_data(field: &Field) -> bool {
Self::has_index_marker(field, Self::INDEX_MARKER_HAS_STATIC_DATA)
}
pub fn is_index_global_temporal(field: &Field) -> bool {
Self::is_index(field)
&& !Self::is_index_static(field)
&& !Self::is_index_per_component(field)
}
}
impl RawRrdManifest {
pub fn compute_sorbet_schema_sha256(
schema: &arrow::datatypes::Schema,
) -> Result<[u8; 32], ArrowError> {
let schema = {
let mut fields = schema.fields().to_vec();
fields.sort();
arrow::datatypes::Schema::new_with_metadata(fields, Default::default()) };
let partition_schema_ipc = {
let mut schema_ipc = Vec::new();
arrow::ipc::writer::StreamWriter::try_new(&mut schema_ipc, &schema)?;
schema_ipc
};
use sha2::Digest as _;
let mut hash = [0u8; 32];
let mut hasher = sha2::Sha256::new();
hasher.update(&partition_schema_ipc);
hasher.finalize_into(sha2::digest::generic_array::GenericArray::from_mut_slice(
&mut hash,
));
Ok(hash)
}
pub fn compute_column_name(
entity_path: Option<&EntityPath>,
strip_entity_prefix: Option<&str>,
component_desc: Option<&ComponentDescriptor>,
prefix: Option<&str>,
suffix: Option<&str>,
) -> String {
use re_types_core::reflection::ComponentDescriptorExt as _;
let full_name = [
prefix.map(ToOwned::to_owned),
entity_path.map(|p| {
let path = p.to_string();
let path = strip_entity_prefix
.and_then(|prefix| path.strip_prefix(prefix))
.unwrap_or(&path);
path.strip_suffix("/").unwrap_or(path).to_owned()
}),
component_desc
.and_then(|descr| descr.archetype)
.map(|archetype| archetype.short_name().to_owned()),
component_desc.map(|descr| descr.archetype_field_name().to_owned()),
suffix.map(ToOwned::to_owned),
]
.into_iter()
.flatten()
.filter(|s| !s.is_empty())
.collect::<Vec<_>>()
.join(":");
let sanitized = full_name.replace([',', ' ', '-', '.', '\\'], "_");
sanitized.trim_start_matches('_').to_owned()
}
pub(super) fn chunk_fetcher_record_batch(&self) -> RecordBatch {
let columns_to_keep = super::RrdManifest::CHUNK_FETCHER_COLUMNS;
let schema = self.data.schema_ref();
let indices: Vec<usize> = schema
.fields()
.iter()
.enumerate()
.filter(|(_, field)| columns_to_keep.contains(&field.name().as_str()))
.map(|(i, _)| i)
.collect();
self.data
.project(&indices)
.unwrap_or_else(|_| self.data.clone())
}
}
impl RawRrdManifest {
pub fn sanity_check_cheap(&self) -> CodecResult<()> {
re_tracing::profile_function!();
self.check_global_columns_are_correct()?;
self.check_index_columns_are_correct()?;
self.check_manifest_schema_matches_sorbet_schema()?;
Ok(())
}
pub fn sanity_check_heavy(&self) -> CodecResult<()> {
re_tracing::profile_function!();
self.check_sorbet_schema_sha256_is_correct()?;
Ok(())
}
fn check_global_columns_are_correct(&self) -> CodecResult<()> {
_ = self.col_chunk_id()?;
_ = self.col_chunk_is_static()?;
_ = self.col_chunk_num_rows()?;
_ = self.col_chunk_entity_path()?;
_ = self.col_chunk_byte_size_uncompressed()?;
_ = self.col_chunk_byte_offset()?;
_ = self.col_chunk_byte_size()?;
if self
.data
.schema_ref()
.column_with_name(Self::COLUMN_CHUNK_KEY.name)
.is_some()
{
_ = self.col_chunk_key()?;
}
Ok(())
}
const SORBET_KIND_INDEX: &str = "index";
const SORBET_KIND_DATA: &str = "data";
fn check_index_columns_are_correct(&self) -> CodecResult<()> {
{
for field in self.data.schema().fields() {
if let Some((_, suffix)) = field.name().rsplit_once(':') {
match suffix {
Self::INDEX_MARKER_START | Self::INDEX_MARKER_END => {
}
Self::INDEX_MARKER_HAS_STATIC_DATA => {
if *field.data_type() != Self::COLUMN_CHUNK_IS_STATIC.data_type() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' should be {} but is actually {}",
field.name(),
Self::COLUMN_CHUNK_IS_STATIC.data_type(),
field.data_type(),
),
}));
}
}
Self::INDEX_MARKER_NUM_ROWS => {
if *field.data_type() != Self::COLUMN_CHUNK_NUM_ROWS.data_type() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' should be {} but is actually {}",
field.name(),
Self::COLUMN_CHUNK_NUM_ROWS.data_type(),
field.data_type(),
),
}));
}
}
suffix => {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' has invalid suffix '{suffix}'",
field.name(),
),
}));
}
}
} else {
match field.name().as_str() {
name if Self::GLOBAL_COLUMN_NAMES.contains(&name) => {}
name if Self::COMMON_IMPL_SPECIFIC_FIELDS.contains(&name) => {}
name => {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"unexpected field '{name}' should not be present in an RRD manifest",
),
}));
}
}
}
}
}
{
for field in self.data.schema().fields() {
if let Some((prefix, suffix)) = field.name().rsplit_once(':') {
let counterpart = match suffix {
"start" => "end",
"end" => "start",
_ => continue,
};
let field_counterpart = self
.data
.schema_ref()
.field_with_name(&format!("{prefix}:{counterpart}"))
.map_err(|_err| {
CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' does not have matching `:{counterpart}` field",
field.name()
),
})
})?;
match field.data_type() {
arrow::datatypes::DataType::Int64
| arrow::datatypes::DataType::Timestamp(_, _)
| arrow::datatypes::DataType::Duration(_) => {}
datatype => {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' is {datatype} which is not a supported index datatype",
field.name(),
),
}));
}
}
if field.data_type() != field_counterpart.data_type() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' is {} but field '{}' is {}",
field.name(),
field.data_type(),
field_counterpart.name(),
field_counterpart.data_type()
),
}));
}
}
}
}
{
for field in self.data.schema().fields() {
if let Some((prefix, "num_rows")) = field.name().rsplit_once(':') {
let field_num_rows = self
.data
.schema_ref()
.field_with_name(&format!("{prefix}:num_rows"))
.map_err(|_err| {
CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' does not have matching `:num_rows` field",
field.name()
),
})
})?;
match field_num_rows.data_type() {
arrow::datatypes::DataType::UInt64 => {}
datatype => {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"field '{}' is {datatype} while it should be UInt64Array",
field_num_rows.name(),
),
}));
}
}
}
}
}
Ok(())
}
fn check_manifest_schema_matches_sorbet_schema(&self) -> CodecResult<()> {
let any_static_chunks = self.col_chunk_is_static_iter()?.any(|b| b);
let sorbet_indexes = self
.sorbet_schema
.fields()
.iter()
.filter_map(|f| {
let md = f.metadata();
(md.get(re_sorbet::RERUN_KIND).map(|s| s.as_str()) == Some(Self::SORBET_KIND_INDEX))
.then(|| md.contains_key(re_sorbet::SORBET_INDEX_NAME).then_some(f))
.flatten()
})
.unique()
.collect_vec();
let sorbet_columns = self
.sorbet_schema
.fields()
.iter()
.filter(|f| {
f.metadata().get(re_sorbet::RERUN_KIND).map(|s| s.as_str())
== Some(Self::SORBET_KIND_DATA)
})
.unique()
.collect_vec();
if any_static_chunks {
for column in &sorbet_columns {
let md = column.metadata();
let Some(component) = md.get(FIELD_METADATA_KEY_COMPONENT) else {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"column '{}' is missing rerun:component metadata",
column.name()
),
}));
};
let descr = ComponentDescriptor {
archetype: md
.get(FIELD_METADATA_KEY_ARCHETYPE)
.and_then(|s| ArchetypeName::try_new(s).ok()),
component: ComponentIdentifier::try_new(component).map_err(|err| {
CodecError::from(ChunkError::Malformed {
reason: err.to_string(),
})
})?,
component_type: md
.get(FIELD_METADATA_KEY_COMPONENT_TYPE)
.and_then(|s| ComponentType::try_new(s).ok()),
};
let column_name = Self::compute_column_name(
None,
None,
Some(&descr),
None,
Some(Self::INDEX_MARKER_HAS_STATIC_DATA),
);
self.data
.schema_ref()
.field_with_name(&column_name)
.map_err(|_err| {
CodecError::from(ChunkError::Malformed {
reason: format!("static index '{column_name}' is missing"),
})
})?;
}
}
let mut rrd_manifest_fields: HashMap<_, _> = self
.data
.schema_ref()
.fields()
.iter()
.filter(|f| Self::is_index_start(f) || Self::is_index_end(f))
.map(|f| (f.name(), f))
.collect();
for sorbet_index in &sorbet_indexes {
let sorbet_index_name_normalized =
Self::compute_column_name(None, None, None, Some(sorbet_index.name()), None);
for suffix in [Self::INDEX_MARKER_START, Self::INDEX_MARKER_END] {
let field = rrd_manifest_fields.remove(&format!("{sorbet_index_name_normalized}:{suffix}"))
.ok_or_else(|| {
CodecError::from(ChunkError::Malformed {
reason: format!(
"global index '{sorbet_index}' does not have matching `:{suffix}` field"
),
})
})?;
if sorbet_index.data_type() != field.data_type() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"global index '{}' is {} but '{}' is {}",
sorbet_index.name(),
sorbet_index.data_type(),
field.name(),
field.data_type()
),
}));
}
}
for sorbet_column in &sorbet_columns {
let md = sorbet_column.metadata();
let Some(component) = md.get(FIELD_METADATA_KEY_COMPONENT) else {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"column '{}' is missing rerun:component metadata",
sorbet_column.name()
),
}));
};
let descr = ComponentDescriptor {
archetype: md
.get(FIELD_METADATA_KEY_ARCHETYPE)
.and_then(|s| ArchetypeName::try_new(s).ok()),
component: ComponentIdentifier::try_new(component).map_err(|err| {
CodecError::from(ChunkError::Malformed {
reason: err.to_string(),
})
})?,
component_type: md
.get(FIELD_METADATA_KEY_COMPONENT_TYPE)
.and_then(|s| ComponentType::try_new(s).ok()),
};
for suffix in [Self::INDEX_MARKER_START, Self::INDEX_MARKER_END] {
let column_name = Self::compute_column_name(
None,
None,
Some(&descr),
Some(sorbet_index.name()),
Some(suffix),
);
if md.get(re_sorbet::SORBET_IS_STATIC).map(|s| s.as_str()) == Some("true") {
_ = rrd_manifest_fields.remove(&column_name);
continue;
}
let Some(field) = rrd_manifest_fields.remove(&column_name) else {
continue;
};
if sorbet_index.data_type() != field.data_type() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"local index '{}' is {} but '{}' is {}",
sorbet_index.name(),
sorbet_index.data_type(),
field.name(),
field.data_type()
),
}));
}
}
}
}
if !rrd_manifest_fields.is_empty() {
return Err(CodecError::from(ChunkError::Malformed {
reason: format!(
"detected dangling indexes (present in manifest but not in Sorbet schema): {:?}",
rrd_manifest_fields.keys()
),
}));
}
Ok(())
}
fn check_sorbet_schema_sha256_is_correct(&self) -> CodecResult<()> {
let expected_sorbet_schema_sha256 = Self::compute_sorbet_schema_sha256(&self.sorbet_schema)
.map_err(CodecError::ArrowDeserialization)?;
if self.sorbet_schema_sha256 != expected_sorbet_schema_sha256 {
return Err(CodecError::ArrowDeserialization(ArrowError::SchemaError(
format!(
"invalid schema hash: expected {} but got {}",
sha256_to_hex(&expected_sorbet_schema_sha256),
sha256_to_hex(&self.sorbet_schema_sha256),
),
)));
}
Ok(())
}
}
impl RawRrdManifest {
pub const COLUMN_CHUNK_ID: quiver::ColumnDesc<ChunkId> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_id");
pub const COLUMN_CHUNK_IS_STATIC: quiver::ColumnDesc<bool> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_is_static");
pub const COLUMN_CHUNK_NUM_ROWS: quiver::ColumnDesc<u64> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_num_rows");
pub const COLUMN_CHUNK_ENTITY_PATH: quiver::ColumnDesc<EntityPath> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_entity_path");
pub const COLUMN_CHUNK_BYTE_OFFSET: quiver::ColumnDesc<u64> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_byte_offset");
pub const COLUMN_CHUNK_BYTE_SIZE: quiver::ColumnDesc<u64> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_byte_size");
pub const COLUMN_CHUNK_BYTE_SIZE_UNCOMPRESSED: quiver::ColumnDesc<u64> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_byte_size_uncompressed");
pub const COLUMN_CHUNK_KEY: quiver::ColumnDesc<quiver::Binary> =
quiver::ColumnDesc::new("RawRrdManifest", "chunk_key");
pub const GLOBAL_COLUMN_NAMES: &[&str] = &[
Self::COLUMN_CHUNK_ID.name,
Self::COLUMN_CHUNK_IS_STATIC.name,
Self::COLUMN_CHUNK_NUM_ROWS.name,
Self::COLUMN_CHUNK_ENTITY_PATH.name,
Self::COLUMN_CHUNK_BYTE_OFFSET.name,
Self::COLUMN_CHUNK_BYTE_SIZE.name,
Self::COLUMN_CHUNK_BYTE_SIZE_UNCOMPRESSED.name,
Self::COLUMN_CHUNK_KEY.name,
];
pub const COMMON_IMPL_SPECIFIC_FIELDS: &[&str] = &[
"chunk_partition_id",
"chunk_partition_layer",
"rerun_partition_id",
"rerun_partition_layer",
"chunk_segment_id",
"chunk_segment_layer",
"rerun_segment_id",
"rerun_segment_layer",
];
pub fn field_index_start(timeline: &Timeline, desc: Option<&ComponentDescriptor>) -> Field {
Self::any_index_field(
timeline,
timeline.datatype(),
desc,
Self::INDEX_MARKER_START,
)
}
pub fn field_index_end(timeline: &Timeline, desc: Option<&ComponentDescriptor>) -> Field {
Self::any_index_field(timeline, timeline.datatype(), desc, Self::INDEX_MARKER_END)
}
pub fn field_index_num_rows(timeline: &Timeline, desc: Option<&ComponentDescriptor>) -> Field {
Self::any_index_field(
timeline,
arrow::datatypes::DataType::UInt64,
desc,
Self::INDEX_MARKER_NUM_ROWS,
)
}
pub fn field_index_has_data(timeline: &Timeline, desc: &ComponentDescriptor) -> Field {
Self::any_index_field(
timeline,
arrow::datatypes::DataType::Boolean,
Some(desc),
Self::INDEX_MARKER_HAS_DATA,
)
}
pub fn field_has_static_data(desc: &ComponentDescriptor) -> Field {
let field_name = Self::compute_column_name(
None,
None,
Some(desc),
None,
Some(Self::INDEX_MARKER_HAS_STATIC_DATA),
);
let mut metadata = std::collections::HashMap::default();
metadata.extend(
[
Some((
Self::FIELD_METADATA_KEY_INDEX.to_owned(),
Self::INDEX_NAME_STATIC.to_owned(),
)),
desc.component_type.map(|component_type| {
(
FIELD_METADATA_KEY_COMPONENT_TYPE.to_owned(),
component_type.full_name().to_owned(),
)
}),
desc.archetype.as_ref().map(|name| {
(
FIELD_METADATA_KEY_ARCHETYPE.to_owned(),
name.full_name().to_owned(),
)
}),
Some((
FIELD_METADATA_KEY_COMPONENT.to_owned(),
desc.component.to_string(),
)),
]
.into_iter()
.flatten(),
);
let nullable = false;
Field::new(field_name, arrow::datatypes::DataType::Boolean, nullable)
.with_metadata(metadata)
}
fn any_index_field(
timeline: &Timeline,
datatype: arrow::datatypes::DataType,
desc: Option<&ComponentDescriptor>,
marker: &str,
) -> Field {
let index_name = timeline.name();
let field_name =
Self::compute_column_name(None, None, desc, Some(index_name), Some(marker));
let mut metadata = std::collections::HashMap::default();
metadata.extend([(
Self::FIELD_METADATA_KEY_INDEX.to_owned(),
timeline.name().to_string(),
)]);
if let Some(desc) = desc {
metadata.extend(
[
desc.component_type.map(|component_type| {
(
FIELD_METADATA_KEY_COMPONENT_TYPE.to_owned(),
component_type.full_name().to_owned(),
)
}),
desc.archetype.as_ref().map(|name| {
(
FIELD_METADATA_KEY_ARCHETYPE.to_owned(),
name.full_name().to_owned(),
)
}),
Some((
FIELD_METADATA_KEY_COMPONENT.to_owned(),
desc.component.to_string(),
)),
]
.into_iter()
.flatten(),
);
}
let nullable = true; Field::new(field_name, datatype, nullable).with_metadata(metadata)
}
}
impl RawRrdManifest {
pub fn col_chunk_entity_path(&self) -> CodecResult<quiver::Column<EntityPath>> {
Ok(Self::COLUMN_CHUNK_ENTITY_PATH.extract(&self.data)?)
}
pub fn col_chunk_entity_path_iter(&self) -> CodecResult<impl Iterator<Item = EntityPath>> {
Ok(self.col_chunk_entity_path()?.into_iter_owned())
}
pub fn col_chunk_id(&self) -> CodecResult<quiver::Column<ChunkId>> {
Ok(Self::COLUMN_CHUNK_ID.extract(&self.data)?)
}
pub fn col_chunk_id_iter(&self) -> CodecResult<impl Iterator<Item = ChunkId>> {
Ok(self.col_chunk_id()?.into_iter_owned())
}
pub fn col_chunk_is_static(&self) -> CodecResult<quiver::Column<bool>> {
Ok(Self::COLUMN_CHUNK_IS_STATIC.extract(&self.data)?)
}
pub fn col_chunk_is_static_iter(&self) -> CodecResult<impl Iterator<Item = bool>> {
Ok(self.col_chunk_is_static()?.into_iter_owned())
}
pub fn col_chunk_num_rows(&self) -> CodecResult<quiver::Column<u64>> {
Ok(Self::COLUMN_CHUNK_NUM_ROWS.extract(&self.data)?)
}
pub fn col_chunk_num_rows_iter(&self) -> CodecResult<impl Iterator<Item = u64>> {
Ok(self.col_chunk_num_rows()?.into_iter_owned())
}
pub fn col_chunk_byte_offset(&self) -> CodecResult<quiver::Column<u64>> {
Ok(Self::COLUMN_CHUNK_BYTE_OFFSET.extract(&self.data)?)
}
pub fn col_chunk_byte_offset_iter(&self) -> CodecResult<impl Iterator<Item = u64>> {
Ok(self.col_chunk_byte_offset()?.into_iter_owned())
}
pub fn col_chunk_byte_size(&self) -> CodecResult<quiver::Column<u64>> {
Ok(Self::COLUMN_CHUNK_BYTE_SIZE.extract(&self.data)?)
}
pub fn col_chunk_byte_size_iter(&self) -> CodecResult<impl Iterator<Item = u64>> {
Ok(self.col_chunk_byte_size()?.into_iter_owned())
}
pub fn col_chunk_byte_size_uncompressed(&self) -> CodecResult<quiver::Column<u64>> {
Ok(Self::COLUMN_CHUNK_BYTE_SIZE_UNCOMPRESSED.extract(&self.data)?)
}
pub fn col_chunk_byte_size_uncompressed_iter(&self) -> CodecResult<impl Iterator<Item = u64>> {
Ok(self.col_chunk_byte_size_uncompressed()?.into_iter_owned())
}
pub fn col_chunk_key(&self) -> CodecResult<quiver::Column<Option<quiver::Binary>>> {
Ok(Self::COLUMN_CHUNK_KEY.optional().extract(&self.data)?)
}
}
fn strip_null_mask_on_default_columns(data: RecordBatch) -> CodecResult<RecordBatch> {
use re_arrow_util::ArrowArrayDowncastRef as _;
let (schema, mut columns, num_rows) = data.into_parts();
let mut new_fields = schema.fields.to_vec();
for (field, column) in itertools::izip!(&mut new_fields, &mut columns) {
if column.null_count() == 0 {
continue;
}
let name = field.name().as_str();
if RawRrdManifest::is_index_has_static_data(field) {
let Some(c) = column.downcast_array_ref::<BooleanArray>() else {
return Err(CodecError::ArrowDeserialization(ArrowError::SchemaError(
format!(
"'{name}' should be a BooleanArray, got {}",
column.data_type()
),
)));
};
let (bools, _nulls) = c.clone().into_parts();
*column = std::sync::Arc::new(BooleanArray::new(bools, None));
*field = std::sync::Arc::new((**field).clone().with_nullable(false));
} else if RawRrdManifest::is_index_num_rows(field) {
let Some(c) = column.downcast_array_ref::<UInt64Array>() else {
return Err(CodecError::ArrowDeserialization(ArrowError::SchemaError(
format!(
"'{name}' should be a UInt64Array, got {}",
column.data_type()
),
)));
};
let (_dt, ints, _nulls) = c.clone().into_parts();
*column = std::sync::Arc::new(UInt64Array::new(ints, None));
*field = std::sync::Arc::new((**field).clone().with_nullable(false));
}
}
let schema = arrow::datatypes::Schema::new_with_metadata(new_fields, schema.metadata.clone());
RecordBatch::try_new_with_options(
std::sync::Arc::new(schema),
columns,
&arrow::array::RecordBatchOptions::new().with_row_count(Some(num_rows)),
)
.map_err(CodecError::ArrowSerialization)
}