use std::ptr;
use reifydb_core::interface::{
catalog::object::ObjectId,
change::{Change, ChangeOrigin, Diff},
};
use reifydb_value::value::diff_type::DiffType;
use tracing::instrument;
use crate::{
common::extern_c::wire::columns::ExternCColumns,
flow::{
extern_c::wire::change::{ExternCChange, ExternCDiff, ExternCOrigin},
operator::extern_c::binding::arena::Arena,
},
};
impl Arena {
#[instrument(name = "flow::marshal::change", level = "trace", skip_all, fields(diff_count = change.diffs.len()))]
pub fn marshal_change(&mut self, change: &Change) -> ExternCChange {
let diffs_count = change.diffs.len();
let diffs_ptr = if diffs_count > 0 {
let diffs_array = self.alloc(diffs_count * size_of::<ExternCDiff>()) as *mut ExternCDiff;
unsafe {
for (i, diff) in change.diffs.iter().enumerate() {
let marshalled = self.marshal_diff(diff);
*diffs_array.add(i) = marshalled;
}
}
diffs_array
} else {
ptr::null_mut()
};
ExternCChange {
origin: Self::marshal_origin(&change.origin),
diff_count: diffs_count,
diffs: diffs_ptr,
version: change.version.source.0,
changed_at: change.changed_at.to_nanos(),
}
}
fn marshal_origin(origin: &ChangeOrigin) -> ExternCOrigin {
match origin {
ChangeOrigin::Flow(operator_id) => ExternCOrigin {
origin: 0,
id: operator_id.0,
},
ChangeOrigin::Object(object_id) => match object_id {
ObjectId::Table(id) => ExternCOrigin {
origin: 1,
id: id.0,
},
ObjectId::View(id) => ExternCOrigin {
origin: 2,
id: id.0,
},
ObjectId::TableVirtual(id) => ExternCOrigin {
origin: 3,
id: id.0,
},
ObjectId::RingBuffer(id) => ExternCOrigin {
origin: 4,
id: id.0,
},
ObjectId::Dictionary(id) => ExternCOrigin {
origin: 6,
id: id.0,
},
ObjectId::Series(id) => ExternCOrigin {
origin: 7,
id: id.0,
},
ObjectId::Queue(id) => ExternCOrigin {
origin: 8,
id: id.0,
},
},
}
}
#[instrument(name = "flow::marshal::diff", level = "trace", skip_all, fields(diff_type = ?diff.kind()))]
fn marshal_diff(&mut self, diff: &Diff) -> ExternCDiff {
match diff {
Diff::Insert {
post,
..
} => ExternCDiff {
diff_type: DiffType::Insert,
pre: ExternCColumns::empty(),
post: self.marshal_columns(post),
},
Diff::Update {
pre,
post,
..
} => ExternCDiff {
diff_type: DiffType::Update,
pre: self.marshal_columns(pre),
post: self.marshal_columns(post),
},
Diff::Remove {
pre,
..
} => ExternCDiff {
diff_type: DiffType::Remove,
pre: self.marshal_columns(pre),
post: ExternCColumns::empty(),
},
}
}
}