use std::{cell::RefCell, collections::BTreeMap, mem, ops::Deref};
use postcard::{from_bytes, to_extend};
use reifydb_value::{
byte_size::ByteSize,
encoding::LeBytes,
error::{Error as ValueError, TypeError},
util::cowvec::CowVec,
value::datetime::DateTime,
};
use serde::{Serialize, de::DeserializeOwned};
use thiserror::Error;
use crate::row::bytes::{EncodedBytes, EncodedRowBuilder, RowBuilder, read_defined_at, sealed::Sealed};
const CREATED_AT_OFFSET: usize = 0;
const UPDATED_AT_OFFSET: usize = CREATED_AT_OFFSET + DateTime::ENCODED_SIZE;
const TIME_OFFSET: usize = UPDATED_AT_OFFSET + DateTime::ENCODED_SIZE;
pub const OPERATOR_HEADER_SIZE: usize = TIME_OFFSET + DateTime::ENCODED_SIZE;
impl From<OperatorError> for ValueError {
fn from(err: OperatorError) -> Self {
match err {
OperatorError::Serialization(_) => TypeError::SerdeSerialize {
message: err.to_string(),
}
.into(),
_ => TypeError::SerdeDeserialize {
message: err.to_string(),
}
.into(),
}
}
}
#[derive(Debug, Error, PartialEq)]
pub enum OperatorError {
#[error("operator state serialization failed: {0}")]
Serialization(String),
#[error("operator state deserialization failed: {0}")]
Deserialization(String),
#[error("operator row is {len} bytes, too short to carry the time header")]
Truncated {
len: usize,
},
}
#[inline]
pub fn read_time(buf: &[u8]) -> Option<DateTime> {
let time = DateTime::from_le_bytes(
buf[TIME_OFFSET..OPERATOR_HEADER_SIZE].try_into().expect("the operator header is length-checked"),
);
(time != DateTime::MAX).then_some(time)
}
#[inline]
pub fn write_time(buf: &mut [u8], time: DateTime) {
buf[TIME_OFFSET..OPERATOR_HEADER_SIZE].copy_from_slice(&time.to_le_bytes());
}
#[inline]
pub fn read_created_at(buf: &[u8]) -> DateTime {
DateTime::from_le_bytes(
buf[CREATED_AT_OFFSET..UPDATED_AT_OFFSET].try_into().expect("the operator header is length-checked"),
)
}
#[inline]
pub fn write_created_at(buf: &mut [u8], created_at: DateTime) {
buf[CREATED_AT_OFFSET..UPDATED_AT_OFFSET].copy_from_slice(&created_at.to_le_bytes());
}
#[inline]
pub fn read_updated_at(buf: &[u8]) -> DateTime {
DateTime::from_le_bytes(
buf[UPDATED_AT_OFFSET..TIME_OFFSET].try_into().expect("the operator header is length-checked"),
)
}
#[inline]
pub fn write_updated_at(buf: &mut [u8], updated_at: DateTime) {
buf[UPDATED_AT_OFFSET..TIME_OFFSET].copy_from_slice(&updated_at.to_le_bytes());
}
#[repr(transparent)]
#[derive(Debug, Clone, PartialEq)]
pub struct EncodedOperatorRow(EncodedBytes);
impl EncodedOperatorRow {
pub fn new(body: &[u8], time: DateTime) -> Self {
let mut buffer = Vec::with_capacity(OPERATOR_HEADER_SIZE + body.len());
buffer.extend_from_slice(&DateTime::EPOCH.to_le_bytes());
buffer.extend_from_slice(&DateTime::EPOCH.to_le_bytes());
buffer.extend_from_slice(&time.to_le_bytes());
buffer.extend_from_slice(body);
Self(EncodedBytes(CowVec::new(buffer)))
}
pub fn timeless(body: &[u8]) -> Self {
Self::new(body, DateTime::MAX)
}
pub fn into_bytes(self) -> EncodedBytes {
self.0
}
pub fn bytes(&self) -> &EncodedBytes {
&self.0
}
pub fn view(bytes: &EncodedBytes) -> &Self {
unsafe { &*(bytes as *const EncodedBytes as *const Self) }
}
#[inline]
pub fn row_time(&self) -> Option<DateTime> {
read_time(&self.0)
}
#[inline]
pub fn time(&self) -> DateTime {
DateTime::from_le_bytes(
self.0[TIME_OFFSET..OPERATOR_HEADER_SIZE].try_into().expect("the header is length-checked"),
)
}
pub fn set_time(&mut self, time: DateTime) {
self.0.make_mut()[TIME_OFFSET..OPERATOR_HEADER_SIZE].copy_from_slice(&time.to_le_bytes());
}
pub fn body(&self) -> &[u8] {
&self.0[OPERATOR_HEADER_SIZE..]
}
pub fn body_mut(&mut self) -> &mut [u8] {
&mut self.0.make_mut()[OPERATOR_HEADER_SIZE..]
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.body().is_empty()
}
pub fn byte_size(&self) -> ByteSize {
ByteSize::from(self.0.len() as u64)
}
}
impl TryFrom<EncodedBytes> for EncodedOperatorRow {
type Error = OperatorError;
fn try_from(bytes: EncodedBytes) -> Result<Self, Self::Error> {
if bytes.len() < OPERATOR_HEADER_SIZE {
return Err(OperatorError::Truncated {
len: bytes.len(),
});
}
Ok(Self(bytes))
}
}
impl From<EncodedOperatorRow> for EncodedBytes {
fn from(row: EncodedOperatorRow) -> Self {
row.0
}
}
#[repr(transparent)]
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EncodedOperatorRowBuilder(EncodedRowBuilder);
impl EncodedOperatorRowBuilder {
pub(crate) fn wrap(builder: EncodedRowBuilder) -> Self {
Self(builder)
}
#[inline]
pub fn row_time(&self) -> Option<DateTime> {
read_time(self.as_slice())
}
pub fn set_time(&mut self, time: DateTime) {
write_time(self.as_mut_slice(), time);
}
#[inline]
pub fn created_at(&self) -> DateTime {
read_created_at(self.as_slice())
}
#[inline]
pub fn updated_at(&self) -> DateTime {
read_updated_at(self.as_slice())
}
pub fn set_timestamps(&mut self, created_at: DateTime, updated_at: DateTime) {
write_created_at(self.as_mut_slice(), created_at);
write_updated_at(self.as_mut_slice(), updated_at);
}
#[inline]
pub fn is_defined(&self, index: usize) -> bool {
read_defined_at(self.as_slice(), OPERATOR_HEADER_SIZE, index)
}
pub fn body(&self) -> &[u8] {
&self.as_slice()[OPERATOR_HEADER_SIZE..]
}
pub fn freeze(self) -> EncodedOperatorRow {
EncodedOperatorRow(self.0.freeze())
}
}
impl Sealed for EncodedOperatorRowBuilder {
fn buffer(&self) -> &Vec<u8> {
self.0.buffer()
}
fn buffer_mut(&mut self) -> &mut Vec<u8> {
self.0.buffer_mut()
}
fn take_buffer(self) -> Vec<u8> {
self.0.take_buffer()
}
}
impl EncodedOperatorRow {
pub fn thaw(self) -> EncodedOperatorRowBuilder {
EncodedOperatorRowBuilder(self.0.thaw())
}
}
pub trait OperatorState: Sized + Send + 'static {
fn encode_state(&self, now: DateTime) -> Result<EncodedOperatorRow, OperatorError>;
fn decode_state(row: &EncodedOperatorRow) -> Result<Self, OperatorError>;
}
thread_local! {
static ENCODE_BUFFER: RefCell<Vec<u8>> = const { RefCell::new(Vec::new()) };
}
pub fn encode<T>(value: &T, now: DateTime) -> Result<EncodedOperatorRow, OperatorError>
where
T: Serialize,
{
let mut buffer = ENCODE_BUFFER.with(|cell| mem::take(&mut *cell.borrow_mut()));
buffer.clear();
let mut filled = to_extend(value, buffer).map_err(|e| OperatorError::Serialization(e.to_string()))?;
let result = EncodedOperatorRow::new(&filled, now);
filled.clear();
ENCODE_BUFFER.with(|cell| *cell.borrow_mut() = filled);
Ok(result)
}
pub fn decode_body<T>(row: &EncodedOperatorRow) -> Result<T, OperatorError>
where
T: DeserializeOwned,
{
from_bytes(row.body()).map_err(|e| OperatorError::Deserialization(e.to_string()))
}
pub fn decode<T: OperatorState>(row: &EncodedOperatorRow) -> Result<T, OperatorError> {
T::decode_state(row)
}
pub mod derive {
pub use serde::{self, Deserialize, Serialize};
}
pub trait StateCodec: Sized + Send + 'static + Serialize + DeserializeOwned {}
impl<T> StateCodec for T where T: Sized + Send + 'static + Serialize + DeserializeOwned {}
macro_rules! leaf_operator_state {
($($ty:ty),* $(,)?) => {
$(impl OperatorState for $ty {
fn encode_state(&self, now: DateTime) -> Result<EncodedOperatorRow, OperatorError> {
encode(self, now)
}
fn decode_state(row: &EncodedOperatorRow) -> Result<Self, OperatorError> {
decode_body::<Self>(row)
}
})*
};
}
leaf_operator_state!(u64, i64, Vec<u8>, (i64, i64, i64), DateTime);
impl<K, V> OperatorState for BTreeMap<K, V>
where
K: Send + 'static,
V: Send + 'static,
Self: Serialize + DeserializeOwned,
{
fn encode_state(&self, now: DateTime) -> Result<EncodedOperatorRow, OperatorError> {
encode(self, now)
}
fn decode_state(row: &EncodedOperatorRow) -> Result<Self, OperatorError> {
decode_body::<Self>(row)
}
}
impl Deref for EncodedOperatorRowBuilder {
type Target = [u8];
fn deref(&self) -> &Self::Target {
self.as_slice()
}
}