use crate::defer;
use crate::platform::{memcpy, memset};
use crate::types::{DataType, RecordHeader, RecordStatus, RemDbError, Result, TableDef, Value};
use core::ptr::NonNull;
extern crate alloc;
use alloc::vec::Vec;
pub struct MemoryTable {
pub def: alloc::sync::Arc<TableDef>,
pub data_start: NonNull<u8>,
pub status_array: NonNull<RecordHeader>,
pub record_count: usize,
pub lock: u32,
pub record_size: usize,
pub free_slots: NonNull<usize>,
pub free_slot_count: usize,
pub low_power_mode: bool,
pub low_power_max_records: Option<usize>,
pub snapshot_version: u32,
pub max_pk: u64,
}
#[derive(Copy, Clone)]
pub struct RecordRef<'a> {
table: &'a MemoryTable,
id: usize,
record_ptr: *const u8,
}
impl<'a> RecordRef<'a> {
pub fn id(&self) -> usize {
self.id
}
pub fn table_def(&self) -> &'a TableDef {
self.table.def.as_ref()
}
fn field_def(&self, col: usize) -> Result<&'a crate::types::FieldDef> {
self.table
.def
.fields
.get(col)
.ok_or(RemDbError::FieldNotFound)
}
fn field_ptr(&self, field: &crate::types::FieldDef) -> *const u8 {
unsafe { self.record_ptr.add(field.offset) }
}
pub fn get_u8(&self, col: usize) -> Result<u8> {
let field = self.field_def(col)?;
if field.data_type != DataType::UInt8 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const u8) })
}
pub fn get_u16(&self, col: usize) -> Result<u16> {
let field = self.field_def(col)?;
if field.data_type != DataType::UInt16 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const u16) })
}
pub fn get_u32(&self, col: usize) -> Result<u32> {
let field = self.field_def(col)?;
if field.data_type != DataType::UInt32 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const u32) })
}
pub fn get_u64(&self, col: usize) -> Result<u64> {
let field = self.field_def(col)?;
if field.data_type != DataType::UInt64 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const u64) })
}
pub fn get_i8(&self, col: usize) -> Result<i8> {
let field = self.field_def(col)?;
if field.data_type != DataType::Int8 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const i8) })
}
pub fn get_i16(&self, col: usize) -> Result<i16> {
let field = self.field_def(col)?;
if field.data_type != DataType::Int16 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const i16) })
}
pub fn get_i32(&self, col: usize) -> Result<i32> {
let field = self.field_def(col)?;
if field.data_type != DataType::Int32 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const i32) })
}
pub fn get_i64(&self, col: usize) -> Result<i64> {
let field = self.field_def(col)?;
if field.data_type != DataType::Int64 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const i64) })
}
pub fn get_f32(&self, col: usize) -> Result<f32> {
let field = self.field_def(col)?;
if field.data_type != DataType::Float32 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const f32) })
}
pub fn get_f64(&self, col: usize) -> Result<f64> {
let field = self.field_def(col)?;
if field.data_type != DataType::Float64 {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const f64) })
}
pub fn get_bool(&self, col: usize) -> Result<bool> {
let field = self.field_def(col)?;
if field.data_type != DataType::Bool {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe { core::ptr::read_unaligned(self.field_ptr(field) as *const u8) != 0 })
}
pub fn get_timestamp(&self, col: usize) -> Result<crate::types::db_timestamp> {
let field = self.field_def(col)?;
if field.data_type != DataType::Timestamp && field.data_type != DataType::TimestampTZ {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe {
core::ptr::read_unaligned(self.field_ptr(field) as *const crate::types::db_timestamp)
})
}
pub fn get_interval(&self, col: usize) -> Result<crate::types::db_interval> {
let field = self.field_def(col)?;
if field.data_type != DataType::Interval {
return Err(RemDbError::TypeMismatch);
}
Ok(unsafe {
core::ptr::read_unaligned(self.field_ptr(field) as *const crate::types::db_interval)
})
}
pub fn get_str(&self, col: usize) -> Result<&'a str> {
let field = self.field_def(col)?;
if field.data_type != DataType::VarChar
&& field.data_type != DataType::Char
&& field.data_type != DataType::Text
{
return Err(RemDbError::TypeMismatch);
}
let bytes = unsafe { core::slice::from_raw_parts(self.field_ptr(field), field.size) };
let end = bytes.iter().position(|b| *b == 0).unwrap_or(bytes.len());
core::str::from_utf8(&bytes[..end]).map_err(|_| RemDbError::TypeMismatch)
}
pub fn get_bytes(&self, col: usize) -> Result<&'a [u8]> {
let field = self.field_def(col)?;
Ok(unsafe { core::slice::from_raw_parts(self.field_ptr(field), field.size) })
}
pub fn get_json(&self, col: usize) -> Result<crate::json::JsonDocument> {
let field = self.field_def(col)?;
if field.data_type != DataType::Json {
return Err(RemDbError::TypeMismatch);
}
let json_storage = unsafe {
core::ptr::read_unaligned(self.field_ptr(field) as *const crate::types::JsonStorage)
};
match json_storage {
crate::types::JsonStorage::Inline(data) => {
let size = field.size;
let data_slice = &data[..size];
crate::json::JsonDocument::from_binary(data_slice, size)
.map_err(|_| RemDbError::TypeMismatch)
}
crate::types::JsonStorage::External {
pool_id,
offset,
length,
} => {
let pool_manager = crate::json::memory_pool::get_global_json_pool_manager()
.ok_or(RemDbError::UnsupportedOperation)?;
let pool = pool_manager
.get_pool(pool_id)
.ok_or(RemDbError::UnsupportedOperation)?;
if let Some(data_ptr) = pool.get_block_data(offset as usize, 0) {
let data_slice =
unsafe { core::slice::from_raw_parts(data_ptr, length as usize) };
crate::json::JsonDocument::from_binary(data_slice, length as usize)
.map_err(|_| RemDbError::TypeMismatch)
} else {
Err(RemDbError::UnsupportedOperation)
}
}
crate::types::JsonStorage::Null => Ok(crate::json::JsonDocument::from_binary(&[], 0)
.map_err(|_| RemDbError::TypeMismatch)?),
}
}
}
pub struct RecordCursor<'a> {
table: &'a MemoryTable,
next_id: usize,
}
impl<'a> RecordCursor<'a> {
fn new(table: &'a MemoryTable) -> Self {
Self { table, next_id: 0 }
}
}
impl<'a> Iterator for RecordCursor<'a> {
type Item = RecordRef<'a>;
fn next(&mut self) -> Option<Self::Item> {
while self.next_id < self.table.def.max_records {
let current = self.next_id;
self.next_id += 1;
unsafe {
let status_ptr = self.table.status_array.as_ptr().add(current);
if (*status_ptr).status == RecordStatus::Used {
let record_ptr = self.table.get_record_ptr(current);
return Some(RecordRef {
table: self.table,
id: current,
record_ptr,
});
}
}
}
None
}
}
pub struct RecordIdCursor<'a> {
table: &'a MemoryTable,
ids: Vec<usize>,
pos: usize,
}
impl<'a> RecordIdCursor<'a> {
fn new(table: &'a MemoryTable, ids: Vec<usize>) -> Self {
Self { table, ids, pos: 0 }
}
}
impl<'a> Iterator for RecordIdCursor<'a> {
type Item = RecordRef<'a>;
fn next(&mut self) -> Option<Self::Item> {
while self.pos < self.ids.len() {
let id = self.ids[self.pos];
self.pos += 1;
if id >= self.table.def.max_records {
continue;
}
unsafe {
let status_ptr = self.table.status_array.as_ptr().add(id);
if (*status_ptr).status == RecordStatus::Used {
let record_ptr = self.table.get_record_ptr(id);
return Some(RecordRef {
table: self.table,
id,
record_ptr,
});
}
}
}
None
}
}
impl Drop for MemoryTable {
fn drop(&mut self) {
unsafe {
crate::memory::allocator::free(self.data_start);
crate::memory::allocator::free(self.status_array.cast());
crate::memory::allocator::free(self.free_slots.cast());
}
}
}
impl MemoryTable {
pub fn new(def: alloc::sync::Arc<TableDef>) -> Result<Self> {
Self::new_with_options(def, false)
}
pub fn new_with_options(
def: alloc::sync::Arc<TableDef>,
skip_status_init: bool,
) -> Result<Self> {
if def.max_records == 0 {
return Err(RemDbError::ConfigError);
}
let data_size = def.record_size * def.max_records;
let status_size = core::mem::size_of::<RecordHeader>() * def.max_records;
let free_slots_size = core::mem::size_of::<usize>() * def.max_records;
let data_start = crate::memory::allocator::alloc(data_size)?;
let status_start = crate::memory::allocator::alloc(status_size)?;
let free_slots_start = crate::memory::allocator::alloc(free_slots_size)?;
let mut free_slot_count = def.max_records;
unsafe {
if !skip_status_init {
let status_array = status_start.cast::<RecordHeader>();
for i in 0..def.max_records {
let status_ptr = status_array.as_ptr().add(i);
(*status_ptr).status = RecordStatus::Free;
(*status_ptr).version = 0;
(*status_ptr).lock_type = crate::types::LockType::None;
(*status_ptr).lock_owner = 0;
(*status_ptr).lock_count = 0;
}
let free_slots = free_slots_start.cast::<usize>();
for i in 0..def.max_records {
*free_slots.as_ptr().add(i) = (def.max_records - 1 - i) as usize;
}
} else {
free_slot_count = 0;
let free_slots = free_slots_start.cast::<usize>();
for i in 0..def.max_records {
*free_slots.as_ptr().add(i) = 0;
}
}
}
Ok(MemoryTable {
def: def.clone(),
data_start,
status_array: status_start.cast(),
record_count: 0,
lock: 0,
record_size: def.record_size, free_slots: free_slots_start.cast(),
free_slot_count: free_slot_count,
low_power_mode: false, low_power_max_records: None, snapshot_version: 0, max_pk: 0, })
}
pub const fn calculate_memory_size(def: &TableDef) -> usize {
let data_size = def.record_size * def.max_records;
let status_size = core::mem::size_of::<RecordHeader>() * def.max_records;
let free_slots_size = core::mem::size_of::<usize>() * def.max_records;
data_size + status_size + free_slots_size
}
pub unsafe fn validate_constraints(
&self,
record_data: *const u8,
exclude_slot: Option<usize>,
) -> Result<()> {
for field in self.def.fields.iter() {
if field.not_null {
let is_null = match field.data_type {
DataType::VarChar | DataType::Char | DataType::Text => {
let str_ptr = record_data.add(field.offset) as *const u8;
let mut all_zero = true;
for i in 0..field.size {
if *str_ptr.add(i) != 0 {
all_zero = false;
break;
}
}
all_zero
}
_ => false,
};
if is_null {
return Err(RemDbError::NotNullViolation);
}
}
match field.data_type {
DataType::Float32 => {
let value =
core::ptr::read_unaligned(record_data.add(field.offset) as *const f32);
if value.is_nan() || value.is_infinite() {
return Err(RemDbError::TypeMismatch);
}
}
DataType::Float64 => {
let value =
core::ptr::read_unaligned(record_data.add(field.offset) as *const f64);
if value.is_nan() || value.is_infinite() {
return Err(RemDbError::TypeMismatch);
}
}
_ => {}
}
}
if !self.def.primary_key.is_empty() {
for slot_id in 0..self.def.max_records {
if exclude_slot == Some(slot_id) {
continue;
}
let status_ptr = self.status_array.as_ptr().add(slot_id);
let status = &*status_ptr;
if status.status == RecordStatus::Used {
let current_tx_id = crate::transaction::tx_id_counter();
let is_visible = crate::transaction::is_visible(
status.create_tx_id,
status.delete_tx_id,
current_tx_id,
);
if !is_visible {
continue;
}
let record_ptr = self.data_start.as_ptr().add(slot_id * self.record_size);
let mut is_duplicate = true;
for &pk_col_idx in &self.def.primary_key {
let pk_field = &self.def.fields[pk_col_idx];
let fields_equal = match pk_field.data_type {
DataType::UInt8 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const u8,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u8,
);
current == existing
}
DataType::UInt16 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const u16,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u16,
);
current == existing
}
DataType::UInt32 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const u32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u32,
);
current == existing
}
DataType::UInt64 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const u64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u64,
);
current == existing
}
DataType::Int8 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const i8,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i8,
);
current == existing
}
DataType::Int16 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const i16,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i16,
);
current == existing
}
DataType::Int32 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const i32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i32,
);
current == existing
}
DataType::Int64 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const i64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i64,
);
current == existing
}
DataType::Float32 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const f32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const f32,
);
current == existing
}
DataType::Float64 => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const f64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const f64,
);
current == existing
}
DataType::Bool => {
let current = core::ptr::read_unaligned(
record_data.add(pk_field.offset) as *const bool,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const bool,
);
current == existing
}
DataType::Timestamp => {
let current =
core::ptr::read_unaligned(record_data.add(pk_field.offset)
as *const crate::types::db_timestamp);
let existing =
core::ptr::read_unaligned(record_ptr.add(pk_field.offset)
as *const crate::types::db_timestamp);
current.value == existing.value
}
DataType::TimestampTZ => {
let current =
core::ptr::read_unaligned(record_data.add(pk_field.offset)
as *const crate::types::db_timestamp);
let existing =
core::ptr::read_unaligned(record_ptr.add(pk_field.offset)
as *const crate::types::db_timestamp);
current.value == existing.value
&& current.tz_offset == existing.tz_offset
}
DataType::VarChar | DataType::Char | DataType::Text => {
let current_str = record_data.add(pk_field.offset) as *const u8;
let existing_str = record_ptr.add(pk_field.offset) as *const u8;
let mut is_equal = true;
for i in 0..pk_field.size {
if *current_str.add(i) != *existing_str.add(i) {
is_equal = false;
break;
}
}
is_equal
}
DataType::Interval => {
let current =
core::ptr::read_unaligned(record_data.add(pk_field.offset)
as *const crate::types::db_interval);
let existing =
core::ptr::read_unaligned(record_ptr.add(pk_field.offset)
as *const crate::types::db_interval);
current.value == existing.value
}
DataType::Vector => false, DataType::Json => false, };
if !fields_equal {
is_duplicate = false;
break;
}
}
if is_duplicate {
return Err(RemDbError::DuplicateKey);
}
}
}
}
for unique_field in self.def.fields.iter().filter(|f| f.unique) {
if unique_field.primary_key {
continue;
}
for slot_id in 0..self.def.max_records {
if exclude_slot == Some(slot_id) {
continue;
}
let status_ptr = self.status_array.as_ptr().add(slot_id);
let status = &*status_ptr;
if status.status == RecordStatus::Used {
let current_tx_id = crate::transaction::tx_id_counter();
let is_visible = crate::transaction::is_visible(
status.create_tx_id,
status.delete_tx_id,
current_tx_id,
);
if !is_visible {
continue;
}
let record_ptr = self.data_start.as_ptr().add(slot_id * self.record_size);
let is_duplicate = match unique_field.data_type {
DataType::UInt8 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const u8,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const u8,
);
current == existing
}
DataType::UInt16 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const u16,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const u16,
);
current == existing
}
DataType::UInt32 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const u32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const u32,
);
current == existing
}
DataType::UInt64 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const u64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const u64,
);
current == existing
}
DataType::Int8 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const i8,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const i8,
);
current == existing
}
DataType::Int16 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const i16,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const i16,
);
current == existing
}
DataType::Int32 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const i32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const i32,
);
current == existing
}
DataType::Int64 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const i64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const i64,
);
current == existing
}
DataType::Float32 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const f32,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const f32,
);
current == existing
}
DataType::Float64 => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const f64,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const f64,
);
current == existing
}
DataType::Bool => {
let current = core::ptr::read_unaligned(
record_data.add(unique_field.offset) as *const bool,
);
let existing = core::ptr::read_unaligned(
record_ptr.add(unique_field.offset) as *const bool,
);
current == existing
}
DataType::Timestamp => {
let current =
core::ptr::read_unaligned(record_data.add(unique_field.offset)
as *const crate::types::db_timestamp);
let existing =
core::ptr::read_unaligned(record_ptr.add(unique_field.offset)
as *const crate::types::db_timestamp);
current.value == existing.value
}
DataType::TimestampTZ => {
let current =
core::ptr::read_unaligned(record_data.add(unique_field.offset)
as *const crate::types::db_timestamp);
let existing =
core::ptr::read_unaligned(record_ptr.add(unique_field.offset)
as *const crate::types::db_timestamp);
current.value == existing.value
&& current.tz_offset == existing.tz_offset
}
DataType::VarChar | DataType::Char | DataType::Text => {
let current_str = record_data.add(unique_field.offset) as *const u8;
let existing_str = record_ptr.add(unique_field.offset) as *const u8;
let mut is_equal = true;
for i in 0..unique_field.size {
if *current_str.add(i) != *existing_str.add(i) {
is_equal = false;
break;
}
}
is_equal
}
DataType::Interval => {
let current =
core::ptr::read_unaligned(record_data.add(unique_field.offset)
as *const crate::types::db_interval);
let existing =
core::ptr::read_unaligned(record_ptr.add(unique_field.offset)
as *const crate::types::db_interval);
current.value == existing.value
}
DataType::Vector => false, DataType::Json => false, };
if is_duplicate {
return Err(RemDbError::DuplicateKey);
}
}
}
}
Ok(())
}
unsafe fn get_field_by_offset(
&self,
record_data: *const u8,
offset: usize,
data_type: DataType,
size: usize,
) -> Result<Value> {
let field_ptr = record_data.add(offset);
let value = match data_type {
DataType::UInt8 => Value {
u8: *field_ptr as u8,
},
DataType::UInt16 => Value {
u16: core::ptr::read_unaligned(field_ptr as *const u16),
},
DataType::UInt32 => Value {
u32: core::ptr::read_unaligned(field_ptr as *const u32),
},
DataType::UInt64 => Value {
u64: core::ptr::read_unaligned(field_ptr as *const u64),
},
DataType::Int8 => Value {
i8: core::ptr::read_unaligned(field_ptr as *const i8),
},
DataType::Int16 => Value {
i16: core::ptr::read_unaligned(field_ptr as *const i16),
},
DataType::Int32 => Value {
i32: core::ptr::read_unaligned(field_ptr as *const i32),
},
DataType::Int64 => Value {
i64: core::ptr::read_unaligned(field_ptr as *const i64),
},
DataType::Float32 => Value {
float32: core::ptr::read_unaligned(field_ptr as *const f32),
},
DataType::Float64 => Value {
float64: core::ptr::read_unaligned(field_ptr as *const f64),
},
DataType::Bool => Value {
bool: *field_ptr != 0,
},
DataType::Timestamp => Value {
time: core::ptr::read_unaligned(field_ptr as *const crate::types::db_timestamp),
},
DataType::TimestampTZ => Value {
time: core::ptr::read_unaligned(field_ptr as *const crate::types::db_timestamp),
},
DataType::Interval => Value {
interval: core::ptr::read_unaligned(field_ptr as *const crate::types::db_interval),
},
DataType::VarChar | DataType::Char | DataType::Text => {
let mut str_value = [0u8; crate::types::MAX_STRING_LEN];
let copy_size = core::cmp::min(size, crate::types::MAX_STRING_LEN);
memcpy(str_value.as_mut_ptr(), field_ptr, copy_size);
Value { string: str_value }
}
DataType::Vector => Value {
vector: field_ptr as *const f32,
},
DataType::Json => {
let json_storage =
core::ptr::read_unaligned(field_ptr as *const crate::types::JsonStorage);
Value { json_storage }
}
};
Ok(value)
}
pub fn insert(&mut self, record_data: *const u8) -> Result<usize> {
crate::get_global_db().map(|db| db.metrics.inc_write_ops());
unsafe {
self.validate_constraints(record_data, None)?;
}
crate::platform::spin_lock(&mut self.lock);
let lock_ptr = &mut self.lock;
defer! { crate::platform::spin_unlock(lock_ptr); }
let max_records = if self.low_power_mode {
self.low_power_max_records.unwrap_or(self.def.max_records)
} else {
self.def.max_records
};
let mut slot_id = 0;
let mut is_overwrite = false;
if self.record_count >= max_records {
if self.low_power_mode {
let mut oldest_id = None;
let mut oldest_version = u16::MAX;
for i in 0..self.def.max_records {
unsafe {
let status_ptr = self.status_array.as_ptr().add(i);
let status = &*status_ptr;
if status.status == RecordStatus::Used && status.version < oldest_version {
oldest_id = Some(i);
oldest_version = status.version;
}
}
}
let slot_id_val = match oldest_id {
Some(id) => id,
None => return Err(RemDbError::NoRecordsToOverwrite),
};
slot_id = slot_id_val;
is_overwrite = true;
} else {
return Err(RemDbError::OutOfMemory);
}
} else {
if self.free_slot_count == 0 {
return Err(RemDbError::OutOfMemory);
}
slot_id = unsafe {
self.free_slot_count -= 1;
*self.free_slots.as_ptr().add(self.free_slot_count)
};
}
let record_ptr = unsafe { self.data_start.as_ptr().add(slot_id * self.record_size) };
if crate::transaction::has_active_tx() {
let mut new_data = Vec::with_capacity(self.record_size);
new_data.resize(self.record_size, 0);
memcpy(new_data.as_mut_ptr(), record_data, self.record_size);
unsafe {
if let Some(mut tx) = crate::transaction::get_current_tx() {
let tx_id = tx.as_mut().id;
tx.as_mut().begin_log_item(
tx_id,
crate::transaction::LogOperation::Insert,
self.def.id,
slot_id as u16,
self.record_size as u16,
None,
Some(&new_data),
);
}
}
}
memcpy(record_ptr, record_data, self.record_size);
let status_ptr = unsafe { self.status_array.as_ptr().add(slot_id) };
unsafe {
(*status_ptr).status = RecordStatus::Used;
(*status_ptr).version += 1;
}
if self.def.primary_key.len() == 1 {
let pk_col_idx = self.def.primary_key[0];
if let Some(pk_field) = self.def.fields.get(pk_col_idx) {
let pk_value =
unsafe {
match pk_field.data_type {
DataType::UInt8 => *record_ptr.add(pk_field.offset) as u64,
DataType::UInt16 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u16,
) as u64,
DataType::UInt32 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u32,
) as u64,
DataType::UInt64 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const u64,
),
DataType::Int8 => *record_ptr.add(pk_field.offset) as i8 as u64,
DataType::Int16 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i16,
) as u64,
DataType::Int32 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i32,
) as u64,
DataType::Int64 => core::ptr::read_unaligned(
record_ptr.add(pk_field.offset) as *const i64,
) as u64,
_ => 0, }
};
if pk_value > self.max_pk {
self.max_pk = pk_value;
}
}
}
if !is_overwrite {
self.record_count += 1;
crate::get_global_db().map(|db| db.metrics.add_used_memory(self.record_size));
}
let inserted_slot_id = slot_id;
Ok(inserted_slot_id)
}
#[cfg(feature = "pubsub")]
unsafe fn publish_to_pubsub_inline(
table_name: &str,
record_size: usize,
id: usize,
record_data: *const u8,
is_insert: bool,
) {
let table_topic = crate::pubsub::topics::get_table_content_topic(table_name);
if let Some(topic_id) = crate::pubsub::get_topic_id(&table_topic) {
let op_type = if is_insert { "INSERT" } else { "UPDATE" };
let mut msg = alloc::format!("{}:table={},id={},data=", op_type, table_name, id);
for i in 0..record_size {
let byte = *record_data.add(i);
msg.push_str(&format!("{:02x}", byte));
}
let _ = crate::pubsub::publish(topic_id, msg.as_bytes());
}
}
pub unsafe fn update(&mut self, id: usize, record_data: *const u8) -> Result<()> {
crate::get_global_db().map(|db| db.metrics.inc_update_ops());
if id >= self.def.max_records {
return Err(RemDbError::RecordNotFound);
}
let status_ptr = self.status_array.as_ptr().add(id);
if (*status_ptr).status != RecordStatus::Used {
return Err(RemDbError::RecordNotFound);
}
self.validate_constraints(record_data, Some(id))?;
crate::platform::spin_lock(&mut self.lock);
let lock_ptr = &mut self.lock;
defer! { crate::platform::spin_unlock(lock_ptr); }
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
if crate::transaction::has_active_tx() {
let mut old_data = Vec::with_capacity(self.record_size);
old_data.resize(self.record_size, 0);
memcpy(old_data.as_mut_ptr(), record_ptr, self.record_size);
let mut new_data = Vec::with_capacity(self.record_size);
new_data.resize(self.record_size, 0);
memcpy(new_data.as_mut_ptr(), record_data, self.record_size);
if let Some(mut tx) = crate::transaction::get_current_tx() {
unsafe {
let tx_id = tx.as_mut().id;
tx.as_mut().begin_log_item(
tx_id,
crate::transaction::LogOperation::Update,
self.def.id,
id as u16,
self.record_size as u16,
Some(&old_data),
Some(&new_data),
);
}
}
}
memcpy(record_ptr, record_data, self.record_size);
(*status_ptr).version += 1;
Ok(())
}
pub unsafe fn delete(&mut self, id: usize) -> Result<()> {
crate::get_global_db().map(|db| db.metrics.inc_delete_ops());
crate::platform::spin_lock(&mut self.lock);
defer! { crate::platform::spin_unlock(&mut self.lock); }
if id >= self.def.max_records {
return Err(RemDbError::RecordNotFound);
}
let status_ptr = self.status_array.as_ptr().add(id);
if (*status_ptr).status != RecordStatus::Used {
return Err(RemDbError::RecordNotFound);
}
if crate::transaction::has_active_tx() {
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
let mut old_data = Vec::with_capacity(self.record_size);
old_data.resize(self.record_size, 0);
memcpy(old_data.as_mut_ptr(), record_ptr, self.record_size);
if let Some(mut tx) = crate::transaction::get_current_tx() {
unsafe {
let tx_id = tx.as_mut().id;
tx.as_mut().begin_log_item(
tx_id,
crate::transaction::LogOperation::Delete,
self.def.id,
id as u16,
self.record_size as u16,
Some(&old_data),
None,
);
}
}
}
(*status_ptr).status = RecordStatus::Free;
(*status_ptr).version += 1;
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
memset(record_ptr, 0, self.record_size);
if self.free_slot_count < self.def.max_records {
*self.free_slots.as_ptr().add(self.free_slot_count) = id;
self.free_slot_count += 1;
}
self.record_count -= 1;
crate::get_global_db().map(|db| db.metrics.sub_used_memory(self.record_size));
Ok(())
}
pub unsafe fn get_by_id(&self, id: usize, dest: *mut u8) -> Result<()> {
crate::get_global_db().map(|db| db.metrics.inc_read_ops());
if id >= self.def.max_records {
return Err(RemDbError::RecordNotFound);
}
let status_ptr = self.status_array.as_ptr().add(id);
if (*status_ptr).status != RecordStatus::Used {
return Err(RemDbError::RecordNotFound);
}
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
memcpy(dest, record_ptr, self.record_size);
Ok(())
}
pub fn get_by_id_ref(&self, id: usize) -> Option<RecordRef<'_>> {
crate::get_global_db().map(|db| db.metrics.inc_read_ops());
if id >= self.def.max_records {
return None;
}
unsafe {
let status_ptr = self.status_array.as_ptr().add(id);
if (*status_ptr).status != RecordStatus::Used {
return None;
}
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
Some(RecordRef {
table: self,
id,
record_ptr,
})
}
}
pub fn scan_ref(&self) -> RecordCursor<'_> {
RecordCursor::new(self)
}
pub fn scan_ids_ref(&self, ids: Vec<usize>) -> RecordIdCursor<'_> {
RecordIdCursor::new(self, ids)
}
#[cfg(feature = "pubsub")]
unsafe fn publish_to_pubsub(&self, id: usize, record_data: *const u8, is_insert: bool) {
let table_name = &self.def.name;
let table_topic = crate::pubsub::topics::get_table_content_topic(table_name);
if let Some(topic_id) = crate::pubsub::get_topic_id(&table_topic) {
let op_type = if is_insert { "INSERT" } else { "UPDATE" };
let mut msg = alloc::format!("{}:table={},id={},data=", op_type, table_name, id);
for i in 0..self.record_size {
let byte = *record_data.add(i);
msg.push_str(&format!("{:02x}", byte));
}
let _ = crate::pubsub::publish(topic_id, msg.as_bytes());
}
}
pub unsafe fn get_field(&self, record_data: *const u8, field_index: usize) -> Result<Value> {
if field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let field = &self.def.fields[field_index];
let field_ptr = record_data.add(field.offset);
let value = match field.data_type {
crate::types::DataType::UInt8 => Value {
u8: *field_ptr as u8,
},
crate::types::DataType::UInt16 => Value {
u16: core::ptr::read_unaligned(field_ptr as *const u16),
},
crate::types::DataType::UInt32 => Value {
u32: core::ptr::read_unaligned(field_ptr as *const u32),
},
crate::types::DataType::UInt64 => Value {
u64: core::ptr::read_unaligned(field_ptr as *const u64),
},
crate::types::DataType::Int8 => Value {
i8: core::ptr::read_unaligned(field_ptr as *const i8),
},
crate::types::DataType::Int16 => Value {
i16: core::ptr::read_unaligned(field_ptr as *const i16),
},
crate::types::DataType::Int32 => Value {
i32: core::ptr::read_unaligned(field_ptr as *const i32),
},
crate::types::DataType::Int64 => Value {
i64: core::ptr::read_unaligned(field_ptr as *const i64),
},
crate::types::DataType::Float32 => Value {
float32: core::ptr::read_unaligned(field_ptr as *const f32),
},
crate::types::DataType::Float64 => Value {
float64: core::ptr::read_unaligned(field_ptr as *const f64),
},
crate::types::DataType::Bool => Value {
bool: *field_ptr != 0,
},
crate::types::DataType::Timestamp => Value {
time: core::ptr::read_unaligned(field_ptr as *const crate::types::db_timestamp),
},
crate::types::DataType::TimestampTZ => Value {
time: core::ptr::read_unaligned(field_ptr as *const crate::types::db_timestamp),
},
crate::types::DataType::VarChar
| crate::types::DataType::Char
| crate::types::DataType::Text => {
let mut str_value = [0u8; crate::types::MAX_STRING_LEN];
let copy_size = core::cmp::min(field.size, crate::types::MAX_STRING_LEN);
memcpy(str_value.as_mut_ptr(), field_ptr, copy_size);
Value { string: str_value }
}
crate::types::DataType::Interval => Value {
interval: core::ptr::read_unaligned(field_ptr as *const crate::types::db_interval),
},
crate::types::DataType::Vector => Value {
vector: field_ptr as *const f32,
},
crate::types::DataType::Json => {
let json_storage =
core::ptr::read_unaligned(field_ptr as *const crate::types::JsonStorage);
Value { json_storage }
}
};
Ok(value)
}
pub unsafe fn set_field(
&self,
record_data: *mut u8,
field_index: usize,
value: &Value,
) -> Result<()> {
if field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let field = &self.def.fields[field_index];
let field_ptr = record_data.add(field.offset);
match field.data_type {
crate::types::DataType::UInt8 => {
*(field_ptr as *mut u8) = value.u8;
}
crate::types::DataType::UInt16 => {
*(field_ptr as *mut u16) = value.u16;
}
crate::types::DataType::UInt32 => {
*(field_ptr as *mut u32) = value.u32;
}
crate::types::DataType::UInt64 => {
*(field_ptr as *mut u64) = value.u64;
}
crate::types::DataType::Int8 => {
*(field_ptr as *mut i8) = value.i8;
}
crate::types::DataType::Int16 => {
*(field_ptr as *mut i16) = value.i16;
}
crate::types::DataType::Int32 => {
*(field_ptr as *mut i32) = value.i32;
}
crate::types::DataType::Int64 => {
*(field_ptr as *mut i64) = value.i64;
}
crate::types::DataType::Float32 => {
*(field_ptr as *mut f32) = value.float32;
}
crate::types::DataType::Float64 => {
*(field_ptr as *mut f64) = value.float64;
}
crate::types::DataType::Bool => {
*field_ptr = if value.bool { 1 } else { 0 };
}
crate::types::DataType::Timestamp => {
*(field_ptr as *mut crate::types::db_timestamp) = value.time;
}
crate::types::DataType::TimestampTZ => {
*(field_ptr as *mut crate::types::db_timestamp) = value.time;
}
crate::types::DataType::VarChar
| crate::types::DataType::Char
| crate::types::DataType::Text => {
memcpy(field_ptr, value.string.as_ptr(), field.size);
}
crate::types::DataType::Interval => {
*(field_ptr as *mut crate::types::db_interval) = value.interval;
}
crate::types::DataType::Vector => {
let vector_metadata = field
.vector_metadata
.as_ref()
.ok_or(RemDbError::InvalidData("vector_metadata not set"))?;
let dimension = vector_metadata.dimension as usize;
crate::compression::compress_vector(value.vector, dimension, field_ptr);
}
crate::types::DataType::Json => {
*(field_ptr as *mut crate::types::JsonStorage) = value.json_storage;
}
}
Ok(())
}
pub fn record_count(&self) -> usize {
self.record_count
}
pub fn max_records(&self) -> usize {
self.def.max_records
}
pub fn is_full(&self) -> bool {
self.record_count >= self.def.max_records
}
pub fn set_low_power_mode(&mut self, enabled: bool, max_records: Option<usize>) {
self.low_power_mode = enabled;
self.low_power_max_records = max_records;
}
pub fn is_low_power_mode(&self) -> bool {
self.low_power_mode
}
pub unsafe fn iterate<F>(&self, mut f: F) -> Result<()>
where
F: FnMut(usize, *const u8) -> bool,
{
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status == RecordStatus::Used {
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
if !f(i, record_ptr) {
break;
}
}
}
Ok(())
}
pub unsafe fn get_status_ptr(&self, index: usize) -> *mut RecordHeader {
debug_assert!(
index < self.def.max_records,
"Record index out of bounds: {} (max: {})",
index,
self.def.max_records
);
self.status_array.as_ptr().add(index)
}
pub unsafe fn get_record_ptr(&self, index: usize) -> *const u8 {
debug_assert!(
index < self.def.max_records,
"Record index out of bounds: {} (max: {})",
index,
self.def.max_records
);
self.data_start.as_ptr().add(index * self.record_size)
}
pub unsafe fn get_record_ptr_mut(&mut self, index: usize) -> *mut u8 {
debug_assert!(
index < self.def.max_records,
"Record index out of bounds: {} (max: {})",
index,
self.def.max_records
);
self.data_start.as_ptr().add(index * self.record_size) as *mut u8
}
pub fn set_json(
&mut self,
record_data: *mut u8,
col: usize,
json_doc: &crate::json::JsonDocument,
) -> Result<()> {
let field = self.def.fields.get(col).ok_or(RemDbError::FieldNotFound)?;
if field.data_type != DataType::Json {
return Err(RemDbError::TypeMismatch);
}
let field_ptr = unsafe { record_data.add(field.offset) };
match json_doc.storage() {
crate::types::JsonStorage::Inline(data) => {
unsafe {
let json_storage = crate::types::JsonStorage::Inline(*data);
core::ptr::write_unaligned(
field_ptr as *mut crate::types::JsonStorage,
json_storage,
);
}
}
crate::types::JsonStorage::External {
pool_id,
offset,
length,
} => {
unsafe {
let json_storage = crate::types::JsonStorage::External {
pool_id: *pool_id,
offset: *offset,
length: *length,
};
core::ptr::write_unaligned(
field_ptr as *mut crate::types::JsonStorage,
json_storage,
);
}
}
crate::types::JsonStorage::Null => {
unsafe {
let json_storage = crate::types::JsonStorage::Null;
core::ptr::write_unaligned(
field_ptr as *mut crate::types::JsonStorage,
json_storage,
);
}
}
}
Ok(())
}
pub unsafe fn set_record_count(&mut self, count: usize) {
self.record_count = count;
}
pub unsafe fn inc_record_count(&mut self) {
self.record_count += 1;
}
pub unsafe fn batch_insert(
&mut self,
records: *const u8,
count: usize,
out_ids: *mut usize,
) -> Result<usize> {
if records.is_null() {
return Err(RemDbError::UnsupportedOperation);
}
crate::platform::spin_lock(&mut self.lock);
let max_records = if self.low_power_mode {
self.low_power_max_records.unwrap_or(self.def.max_records)
} else {
self.def.max_records
};
let available = max_records - self.record_count;
let mut actual_count = count;
if self.low_power_mode && self.record_count >= max_records {
actual_count = count;
} else if available < count {
actual_count = available;
}
if actual_count == 0 {
crate::platform::spin_unlock(&mut self.lock);
return Err(RemDbError::OutOfMemory);
}
let mut slot_ids = [0usize; 256]; assert!(
actual_count <= slot_ids.len(),
"Batch insert count exceeds maximum"
);
let mut inserted_count = 0;
let mut i = 0;
while i < actual_count && self.free_slot_count > 0 {
slot_ids[i] = *self.free_slots.as_ptr().add(self.free_slot_count - 1);
self.free_slot_count -= 1;
inserted_count += 1;
i += 1;
}
if i < actual_count && self.low_power_mode {
let mut oldest_ids = [0usize; 256];
let mut oldest_versions = [u16::MAX; 256];
for record_id in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(record_id);
let status = &*status_ptr;
if status.status == crate::types::RecordStatus::Used {
for j in 0..(actual_count - i) {
if status.version < oldest_versions[j] {
for k in (j + 1)..(actual_count - i) {
if oldest_versions[k] > oldest_versions[k - 1] {
break;
}
oldest_ids[k] = oldest_ids[k - 1];
oldest_versions[k] = oldest_versions[k - 1];
}
oldest_ids[j] = record_id;
oldest_versions[j] = status.version;
break;
}
}
}
}
for j in 0..(actual_count - i) {
slot_ids[i + j] = oldest_ids[j];
}
inserted_count = actual_count;
}
let has_active_tx = crate::transaction::has_active_tx();
let current_tx = if has_active_tx {
crate::transaction::get_current_tx()
} else {
None
};
crate::platform::spin_unlock(&mut self.lock);
for j in 0..inserted_count {
let slot_id = slot_ids[j];
if !out_ids.is_null() {
*out_ids.add(j) = slot_id;
}
let record_ptr = self.data_start.as_ptr().add(slot_id * self.record_size);
let src_ptr = records.add(j * self.record_size);
if let Some(mut tx) = current_tx {
let tx_mut = tx.as_mut();
if tx_mut.is_active() && !tx_mut.is_read_only() {
let mut new_data = Vec::with_capacity(self.record_size);
new_data.resize(self.record_size, 0);
memcpy(new_data.as_mut_ptr(), src_ptr, self.record_size);
let var_log_item = tx_mut.begin_variable_size_log_item(
tx_mut.id,
crate::transaction::LogOperation::Insert,
self.def.id,
slot_id as u16,
None,
Some(&new_data),
);
if let Some(log_manager) = crate::transaction::get_log_manager() {
log_manager
.write_variable_size_log_item(&var_log_item)
.unwrap_or(());
}
}
}
memcpy(record_ptr, src_ptr, self.record_size);
let status_ptr = self.status_array.as_ptr().add(slot_id);
(*status_ptr).status = crate::types::RecordStatus::Used;
(*status_ptr).version += 1;
}
crate::platform::spin_lock(&mut self.lock);
let new_records_count = if self.low_power_mode && self.record_count >= max_records {
0 } else {
inserted_count };
self.record_count += new_records_count;
crate::platform::spin_unlock(&mut self.lock);
Ok(inserted_count)
}
pub unsafe fn time_series_batch_insert(
&mut self,
records: *const u8,
count: usize,
out_ids: *mut usize,
) -> Result<usize> {
if records.is_null() {
return Err(RemDbError::UnsupportedOperation);
}
crate::platform::spin_lock(&mut self.lock);
let available = self.def.max_records - self.record_count;
let actual_count = core::cmp::min(count, available);
let actual_count = core::cmp::min(actual_count, self.free_slot_count);
if actual_count == 0 {
crate::platform::spin_unlock(&mut self.lock);
return Err(RemDbError::OutOfMemory);
}
let original_free_slot_count = self.free_slot_count;
let end_free_slot = self.free_slot_count - actual_count;
self.free_slot_count = end_free_slot;
let has_active_tx = crate::transaction::has_active_tx();
let current_tx = if has_active_tx {
crate::transaction::get_current_tx()
} else {
None
};
crate::platform::spin_unlock(&mut self.lock);
let mut inserted_count = 0;
for i in 0..actual_count {
let free_slot_index = original_free_slot_count - 1 - i;
let slot_id = *self.free_slots.as_ptr().add(free_slot_index);
if !out_ids.is_null() {
*out_ids.add(i) = slot_id;
}
let record_ptr = self.data_start.as_ptr().add(slot_id * self.record_size);
let src_ptr = records.add(i * self.record_size);
if let Some(mut tx) = current_tx {
let tx_mut = tx.as_mut();
if tx_mut.is_active() && !tx_mut.is_read_only() {
let mut new_data = Vec::with_capacity(self.record_size);
new_data.resize(self.record_size, 0);
memcpy(new_data.as_mut_ptr(), src_ptr, self.record_size);
let var_log_item = tx_mut.begin_variable_size_log_item(
tx_mut.id,
crate::transaction::LogOperation::TimeSeriesInsert,
self.def.id,
slot_id as u16,
None,
Some(&new_data),
);
if let Some(log_manager) = crate::transaction::get_log_manager() {
log_manager
.write_variable_size_log_item(&var_log_item)
.unwrap_or(());
}
}
}
memcpy(record_ptr, src_ptr, self.record_size);
let status_ptr = self.status_array.as_ptr().add(slot_id);
(*status_ptr).status = RecordStatus::Used;
inserted_count += 1;
}
crate::platform::spin_lock(&mut self.lock);
self.record_count += inserted_count;
crate::platform::spin_unlock(&mut self.lock);
Ok(inserted_count)
}
pub unsafe fn batch_get(&self, ids: &[usize], dest: *mut u8) -> Result<usize> {
let mut success_count = 0;
for (i, &id) in ids.iter().enumerate() {
if id >= self.def.max_records {
continue;
}
let status_ptr = self.status_array.as_ptr().add(id);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(id * self.record_size);
let dest_ptr = dest.add(i * self.record_size);
memcpy(dest_ptr, record_ptr, self.record_size);
success_count += 1;
}
Ok(success_count)
}
pub unsafe fn aggregate_count(
&self,
time_field_index: usize,
start_time: u64,
end_time: u64,
) -> Result<usize> {
if time_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let mut count = 0;
let time_field = &self.def.fields[time_field_index];
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let timestamp = match time_field.data_type {
crate::types::DataType::UInt64 => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
crate::types::DataType::Timestamp => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
_ => {
let field_ptr = record_ptr.add(time_field.offset);
match time_field.data_type {
crate::types::DataType::UInt8 => {
core::ptr::read_unaligned(field_ptr as *const u8) as u64
}
crate::types::DataType::UInt16 => {
core::ptr::read_unaligned(field_ptr as *const u16) as u64
}
crate::types::DataType::UInt32 => {
core::ptr::read_unaligned(field_ptr as *const u32) as u64
}
crate::types::DataType::Int8 => {
core::ptr::read_unaligned(field_ptr as *const i8) as u64
}
crate::types::DataType::Int16 => {
core::ptr::read_unaligned(field_ptr as *const i16) as u64
}
crate::types::DataType::Int32 => {
core::ptr::read_unaligned(field_ptr as *const i32) as u64
}
crate::types::DataType::Int64 => {
core::ptr::read_unaligned(field_ptr as *const i64) as u64
}
_ => continue, }
}
};
if timestamp >= start_time && timestamp <= end_time {
count += 1;
}
}
Ok(count)
}
pub unsafe fn aggregate_sum(
&self,
time_field_index: usize,
value_field_index: usize,
start_time: u64,
end_time: u64,
) -> Result<f64> {
if time_field_index >= self.def.fields.len() || value_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let mut sum = 0.0;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let time_field = &self.def.fields[time_field_index];
let timestamp = match time_field.data_type {
crate::types::DataType::UInt64 => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
crate::types::DataType::Timestamp => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
_ => {
let field_ptr = record_ptr.add(time_field.offset);
match time_field.data_type {
crate::types::DataType::UInt8 => {
core::ptr::read_unaligned(field_ptr as *const u8) as u64
}
crate::types::DataType::UInt16 => {
core::ptr::read_unaligned(field_ptr as *const u16) as u64
}
crate::types::DataType::UInt32 => {
core::ptr::read_unaligned(field_ptr as *const u32) as u64
}
crate::types::DataType::Int8 => {
core::ptr::read_unaligned(field_ptr as *const i8) as u64
}
crate::types::DataType::Int16 => {
core::ptr::read_unaligned(field_ptr as *const i16) as u64
}
crate::types::DataType::Int32 => {
core::ptr::read_unaligned(field_ptr as *const i32) as u64
}
crate::types::DataType::Int64 => {
core::ptr::read_unaligned(field_ptr as *const i64) as u64
}
_ => continue, }
}
};
if timestamp >= start_time && timestamp <= end_time {
let value = self.get_field(record_ptr, value_field_index)?;
let numeric_value = match self.def.fields[value_field_index].data_type {
crate::types::DataType::UInt8 => value.u8 as f64,
crate::types::DataType::UInt16 => value.u16 as f64,
crate::types::DataType::UInt32 => value.u32 as f64,
crate::types::DataType::UInt64 => value.u64 as f64,
crate::types::DataType::Float32 => value.float32 as f64,
crate::types::DataType::Float64 => value.float64,
_ => return Err(RemDbError::TypeMismatch),
};
sum += numeric_value;
}
}
Ok(sum)
}
pub unsafe fn aggregate_avg(
&self,
time_field_index: usize,
value_field_index: usize,
start_time: u64,
end_time: u64,
) -> Result<f64> {
if time_field_index >= self.def.fields.len() || value_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let mut sum = 0.0;
let mut count = 0;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let time_field = &self.def.fields[time_field_index];
let timestamp = match time_field.data_type {
crate::types::DataType::UInt64 => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
crate::types::DataType::Timestamp => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
_ => {
let field_ptr = record_ptr.add(time_field.offset);
match time_field.data_type {
crate::types::DataType::UInt8 => {
core::ptr::read_unaligned(field_ptr as *const u8) as u64
}
crate::types::DataType::UInt16 => {
core::ptr::read_unaligned(field_ptr as *const u16) as u64
}
crate::types::DataType::UInt32 => {
core::ptr::read_unaligned(field_ptr as *const u32) as u64
}
crate::types::DataType::Int8 => {
core::ptr::read_unaligned(field_ptr as *const i8) as u64
}
crate::types::DataType::Int16 => {
core::ptr::read_unaligned(field_ptr as *const i16) as u64
}
crate::types::DataType::Int32 => {
core::ptr::read_unaligned(field_ptr as *const i32) as u64
}
crate::types::DataType::Int64 => {
core::ptr::read_unaligned(field_ptr as *const i64) as u64
}
_ => continue, }
}
};
if timestamp >= start_time && timestamp <= end_time {
let value = self.get_field(record_ptr, value_field_index)?;
let numeric_value = match self.def.fields[value_field_index].data_type {
crate::types::DataType::UInt8 => value.u8 as f64,
crate::types::DataType::UInt16 => value.u16 as f64,
crate::types::DataType::UInt32 => value.u32 as f64,
crate::types::DataType::UInt64 => value.u64 as f64,
crate::types::DataType::Float32 => value.float32 as f64,
crate::types::DataType::Float64 => value.float64,
_ => return Err(RemDbError::TypeMismatch),
};
sum += numeric_value;
count += 1;
}
}
if count == 0 {
Ok(0.0)
} else {
Ok(sum / count as f64)
}
}
pub unsafe fn aggregate_min(
&self,
time_field_index: usize,
value_field_index: usize,
start_time: u64,
end_time: u64,
) -> Result<f64> {
if time_field_index >= self.def.fields.len() || value_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let mut min_value: Option<f64> = None;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let time_field = &self.def.fields[time_field_index];
let timestamp = match time_field.data_type {
crate::types::DataType::UInt64 => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
crate::types::DataType::Timestamp => {
core::ptr::read_unaligned(record_ptr.add(time_field.offset) as *const u64)
}
_ => {
let field_ptr = record_ptr.add(time_field.offset);
match time_field.data_type {
crate::types::DataType::UInt8 => {
core::ptr::read_unaligned(field_ptr as *const u8) as u64
}
crate::types::DataType::UInt16 => {
core::ptr::read_unaligned(field_ptr as *const u16) as u64
}
crate::types::DataType::UInt32 => {
core::ptr::read_unaligned(field_ptr as *const u32) as u64
}
crate::types::DataType::Int8 => {
core::ptr::read_unaligned(field_ptr as *const i8) as u64
}
crate::types::DataType::Int16 => {
core::ptr::read_unaligned(field_ptr as *const i16) as u64
}
crate::types::DataType::Int32 => {
core::ptr::read_unaligned(field_ptr as *const i32) as u64
}
crate::types::DataType::Int64 => {
core::ptr::read_unaligned(field_ptr as *const i64) as u64
}
_ => continue, }
}
};
if timestamp >= start_time && timestamp <= end_time {
let value = self.get_field(record_ptr, value_field_index)?;
let numeric_value = match self.def.fields[value_field_index].data_type {
crate::types::DataType::UInt8 => value.u8 as f64,
crate::types::DataType::UInt16 => value.u16 as f64,
crate::types::DataType::UInt32 => value.u32 as f64,
crate::types::DataType::UInt64 => value.u64 as f64,
crate::types::DataType::Float32 => value.float32 as f64,
crate::types::DataType::Float64 => value.float64,
_ => return Err(RemDbError::TypeMismatch),
};
if let Some(current_min) = min_value {
if numeric_value < current_min {
min_value = Some(numeric_value);
}
} else {
min_value = Some(numeric_value);
}
}
}
min_value.ok_or(RemDbError::RecordNotFound)
}
unsafe fn read_timestamp_value(
&self,
record_ptr: *const u8,
time_field_index: usize,
) -> Option<u64> {
let time_field = &self.def.fields[time_field_index];
match time_field.data_type {
crate::types::DataType::UInt64 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const u64,
)),
crate::types::DataType::Timestamp => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const u64,
)),
crate::types::DataType::UInt8 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const u8,
) as u64),
crate::types::DataType::UInt16 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const u16,
) as u64),
crate::types::DataType::UInt32 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const u32,
) as u64),
crate::types::DataType::Int8 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const i8,
) as u64),
crate::types::DataType::Int16 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const i16,
) as u64),
crate::types::DataType::Int32 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const i32,
) as u64),
crate::types::DataType::Int64 => Some(core::ptr::read_unaligned(
record_ptr.add(time_field.offset) as *const i64,
) as u64),
_ => None, }
}
pub unsafe fn aggregate_max(
&self,
time_field_index: usize,
value_field_index: usize,
start_time: u64,
end_time: u64,
) -> Result<f64> {
if time_field_index >= self.def.fields.len() || value_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let mut max_value: Option<f64> = None;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let Some(timestamp) = self.read_timestamp_value(record_ptr, time_field_index) else {
continue;
};
if timestamp >= start_time && timestamp <= end_time {
let value = self.get_field(record_ptr, value_field_index)?;
let numeric_value = match self.def.fields[value_field_index].data_type {
crate::types::DataType::UInt8 => value.u8 as f64,
crate::types::DataType::UInt16 => value.u16 as f64,
crate::types::DataType::UInt32 => value.u32 as f64,
crate::types::DataType::UInt64 => value.u64 as f64,
crate::types::DataType::Float32 => value.float32 as f64,
crate::types::DataType::Float64 => value.float64,
_ => return Err(RemDbError::TypeMismatch),
};
if let Some(current_max) = max_value {
if numeric_value > current_max {
max_value = Some(numeric_value);
}
} else {
max_value = Some(numeric_value);
}
}
}
max_value.ok_or(RemDbError::RecordNotFound)
}
pub unsafe fn get_latest_records(
&self,
time_field_index: usize,
count: usize,
dest: *mut u8,
) -> Result<usize> {
if time_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
if dest.is_null() {
return Err(RemDbError::UnsupportedOperation);
}
if self.record_count == 0 {
return Ok(0);
}
let mut record_times = [(0usize, 0u64); 1024]; let mut record_count = 0;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
if let Some(timestamp) = self.read_timestamp_value(record_ptr, time_field_index) {
record_times[record_count] = (i, timestamp);
record_count += 1;
};
}
for i in 0..record_count {
for j in i + 1..record_count {
if record_times[i].1 < record_times[j].1 {
record_times.swap(i, j);
}
}
}
let actual_count = core::cmp::min(count, record_count);
for i in 0..actual_count {
let (record_id, _) = record_times[i];
let src_ptr = self.data_start.as_ptr().add(record_id * self.record_size);
let dest_ptr = dest.add(i * self.record_size);
memcpy(dest_ptr, src_ptr, self.record_size);
}
Ok(actual_count)
}
pub unsafe fn get_records_in_time_window(
&self,
time_field_index: usize,
start_time: u64,
end_time: u64,
dest: *mut u8,
max_records: usize,
) -> Result<usize> {
if time_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
if dest.is_null() {
return Err(RemDbError::UnsupportedOperation);
}
if self.record_count == 0 {
return Ok(0);
}
let mut matched_records = [(0usize, 0u64); 1024]; let mut match_count = 0;
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
if let Some(timestamp) = self.read_timestamp_value(record_ptr, time_field_index) {
if timestamp >= start_time && timestamp <= end_time {
matched_records[match_count] = (i, timestamp);
match_count += 1;
}
};
}
for i in 0..match_count {
for j in i + 1..match_count {
if matched_records[i].1 > matched_records[j].1 {
matched_records.swap(i, j);
}
}
}
let actual_count = core::cmp::min(max_records, match_count);
for i in 0..actual_count {
let (record_id, _) = matched_records[i];
let src_ptr = self.data_start.as_ptr().add(record_id * self.record_size);
let dest_ptr = dest.add(i * self.record_size);
memcpy(dest_ptr, src_ptr, self.record_size);
}
Ok(actual_count)
}
#[cfg(feature = "std")]
pub unsafe fn get_aggregate_in_time_window(
&self,
time_field_index: usize,
value_field_index: usize,
start_time: u64,
end_time: u64,
window_size: u64,
) -> Result<Vec<(u64, f64, f64, f64, f64, usize)>> {
if time_field_index >= self.def.fields.len() || value_field_index >= self.def.fields.len() {
return Err(RemDbError::FieldNotFound);
}
let time_field = &self.def.fields[time_field_index];
let value_field = &self.def.fields[value_field_index];
if time_field.data_type != crate::types::DataType::Timestamp
&& time_field.data_type != crate::types::DataType::UInt64
{
return Err(RemDbError::TypeMismatch);
}
use alloc::collections::BTreeMap;
let mut window_aggregates: BTreeMap<u64, (f64, f64, f64, f64, usize)> = BTreeMap::new();
for i in 0..self.def.max_records {
let status_ptr = self.status_array.as_ptr().add(i);
if (*status_ptr).status != RecordStatus::Used {
continue;
}
let record_ptr = self.data_start.as_ptr().add(i * self.record_size);
let Some(time_value) = self.read_timestamp_value(record_ptr, time_field_index) else {
continue;
};
if time_value >= start_time && time_value <= end_time {
let value = self.get_field(record_ptr, value_field_index)?;
let numeric_value = match value_field.data_type {
crate::types::DataType::UInt8 => value.u8 as f64,
crate::types::DataType::UInt16 => value.u16 as f64,
crate::types::DataType::UInt32 => value.u32 as f64,
crate::types::DataType::UInt64 => value.u64 as f64,
crate::types::DataType::Float32 => value.float32 as f64,
crate::types::DataType::Float64 => value.float64,
_ => return Err(RemDbError::TypeMismatch),
};
let window_key = time_value - (time_value % window_size);
let entry = window_aggregates.entry(window_key).or_insert((
0.0,
numeric_value,
numeric_value,
0.0,
0,
));
entry.0 += numeric_value; if numeric_value < entry.1 {
entry.1 = numeric_value;
} if numeric_value > entry.2 {
entry.2 = numeric_value;
} entry.3 = numeric_value; entry.4 += 1; }
}
let mut result = Vec::with_capacity(window_aggregates.len());
for (window_start, (sum, min, max, _last, count)) in window_aggregates {
let avg = if count > 0 { sum / count as f64 } else { 0.0 };
result.push((window_start, sum, avg, min, max, count));
}
Ok(result)
}
#[cfg(not(feature = "std"))]
pub unsafe fn get_aggregate_in_time_window(
&self,
_time_field_index: usize,
_value_field_index: usize,
_start_time: u64,
_end_time: u64,
_window_size: u64,
) -> Result<()> {
Err(RemDbError::UnsupportedOperation)
}
}
#[macro_export]
macro_rules! defer {
($($code:tt)*) => {
let _defer = $crate::table::Defer::new(|| { $($code)* });
};
}
pub struct Defer<F: FnMut()>(Option<F>);
impl<F: FnMut()> Defer<F> {
pub fn new(f: F) -> Self {
Defer(Some(f))
}
}
impl<F: FnMut()> Drop for Defer<F> {
fn drop(&mut self) {
if let Some(mut f) = self.0.take() {
f();
}
}
}