reifydb-sdk 0.9.3

SDK for building ReifyDB operators, procedures, transforms and more
Documentation
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2026 ReifyDB

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;

			// SAFETY: `diffs_count > 0` makes the arena block non-null, and it reserved
			// `diffs_count * size_of::<ExternCDiff>()` bytes at alignment 8, so every `add(i)` with
			// `i < diffs_count` is in bounds; ExternCDiff is Copy, so the stores drop nothing. The
			// writes go through the raw pointer because a reference to the block would be invalid
			// until every `diff_type` discriminant is written.
			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(),
			},
		}
	}
}