#![no_std]
#![cfg_attr(target_arch = "aarch64", allow(unsafe_code))]
extern crate alloc;
pub mod bignum;
pub mod bloom;
mod codec;
pub mod fts_simple;
pub mod halfvec;
pub mod jsonb_gin;
mod nsw;
pub mod persistent;
pub mod persistent_btree;
pub mod quantize;
pub mod row_header;
pub mod row_locator;
pub mod segment;
pub mod snapshot;
mod table;
pub mod trgm;
pub mod vacuum;
pub use self::bloom::{BloomError, BloomFilter};
pub(crate) use self::codec::*;
pub use self::codec::{
decode_row_body_dense, decode_row_body_dense_pruned, encode_row_body_dense,
encode_row_body_dense_into, encode_row_body_dense_masked_into, row_body_encoded_len,
};
pub(crate) use self::nsw::nsw_insert_at;
pub use self::nsw::{NswMetric, cosine_dot_norms_f32, inner_product_f32, nsw_index_on, nsw_query};
pub use self::row_locator::{RowLocator, RowLocatorError};
pub use self::segment::{
BRIN_SIDECAR_MAGIC, BrinSummary, OwnedSegment, SEGMENT_COMPRESS_ALGO_LZSS,
SEGMENT_COMPRESS_ALGO_NONE, SEGMENT_MAGIC, SEGMENT_MAGIC_V2, SEGMENT_PAGE_BYTES, SegmentError,
SegmentMeta, SegmentReader, derive_brin_summaries, encode_segment, wrap_v2_envelope,
wrap_v2_envelope_with_brin,
};
use alloc::borrow::Cow;
use alloc::boxed::Box;
use alloc::collections::{BTreeMap, BTreeSet};
use alloc::format;
use alloc::string::{String, ToString};
use alloc::sync::Arc;
use alloc::vec::Vec;
use core::fmt;
use self::persistent::PersistentVec;
use self::persistent_btree::PersistentBTreeMap;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum VecEncoding {
#[default]
F32,
Sq8,
F16,
}
impl fmt::Display for VecEncoding {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::F32 => f.write_str("F32"),
Self::Sq8 => f.write_str("SQ8"),
Self::F16 => f.write_str("HALF"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DataType {
SmallInt,
Int, BigInt, Float, Real,
Text,
Varchar(u32),
Char(u32),
Bool,
Vector {
dim: u32,
encoding: VecEncoding,
},
Numeric {
precision: u16,
scale: i16,
},
Date,
Timestamp,
Timestamptz,
Name,
Xid,
Xid8,
Oid,
Interval,
Json,
Jsonb,
Bytes,
TextArray,
IntArray,
BigIntArray,
OidArray,
IntervalArray,
BoolArray, SmallIntArray, FloatArray, NumericArray, DateArray, TimestampArray, TimestamptzArray, UuidArray, JsonArray, JsonbArray, BytesArray, VarcharArray, CharArray, Multirange(RangeKind),
Point,
Lseg,
Path,
PgBox,
Polygon,
Line,
Circle,
Inet,
Cidr,
Macaddr,
Macaddr8,
PgLsn,
Bit(u32),
BitVarying(u32),
Xml,
Char1,
MoneyArray,
TsVector,
TsQuery,
Uuid,
Time,
Year,
TimeTz,
Money,
Range(RangeKind),
Hstore,
IntArray2D,
BigIntArray2D,
TextArray2D,
BoolArray2D,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum RangeKind {
Int4,
Int8,
Num,
Ts,
TsTz,
Date,
}
impl RangeKind {
pub const fn tag(self) -> u8 {
match self {
Self::Int4 => 0,
Self::Int8 => 1,
Self::Num => 2,
Self::Ts => 3,
Self::TsTz => 4,
Self::Date => 5,
}
}
pub const fn from_tag(t: u8) -> Option<Self> {
Some(match t {
0 => Self::Int4,
1 => Self::Int8,
2 => Self::Num,
3 => Self::Ts,
4 => Self::TsTz,
5 => Self::Date,
_ => return None,
})
}
pub const fn keyword(self) -> &'static str {
match self {
Self::Int4 => "INT4RANGE",
Self::Int8 => "INT8RANGE",
Self::Num => "NUMRANGE",
Self::Ts => "TSRANGE",
Self::TsTz => "TSTZRANGE",
Self::Date => "DATERANGE",
}
}
}
impl fmt::Display for DataType {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::SmallInt => f.write_str("SMALLINT"),
Self::Int => f.write_str("INT"),
Self::BigInt => f.write_str("BIGINT"),
Self::Xid => f.write_str("XID"),
Self::Xid8 => f.write_str("XID8"),
Self::Oid => f.write_str("OID"),
Self::OidArray => f.write_str("OID[]"),
Self::Float => f.write_str("FLOAT"),
Self::Real => f.write_str("REAL"),
Self::Text => f.write_str("TEXT"),
Self::Varchar(n) => write!(f, "VARCHAR({n})"),
Self::Char(n) => write!(f, "CHAR({n})"),
Self::Bool => f.write_str("BOOL"),
Self::Vector { dim, encoding } => match encoding {
VecEncoding::F32 => write!(f, "VECTOR({dim})"),
VecEncoding::Sq8 => write!(f, "VECTOR({dim}) USING SQ8"),
VecEncoding::F16 => write!(f, "VECTOR({dim}) USING HALF"),
},
Self::Numeric { precision, scale } => {
if *scale == 0 {
write!(f, "NUMERIC({precision})")
} else {
write!(f, "NUMERIC({precision}, {scale})")
}
}
Self::Date => f.write_str("DATE"),
Self::Timestamp => f.write_str("TIMESTAMP"),
Self::Timestamptz => f.write_str("TIMESTAMPTZ"),
Self::Name => f.write_str("NAME"),
Self::Interval => f.write_str("INTERVAL"),
Self::Json => f.write_str("JSON"),
Self::Jsonb => f.write_str("JSONB"),
Self::Bytes => f.write_str("BYTEA"),
Self::TextArray => f.write_str("TEXT[]"),
Self::IntArray => f.write_str("INT[]"),
Self::BigIntArray => f.write_str("BIGINT[]"),
Self::IntervalArray => f.write_str("INTERVAL[]"),
Self::BoolArray => f.write_str("BOOL[]"),
Self::SmallIntArray => f.write_str("SMALLINT[]"),
Self::FloatArray => f.write_str("FLOAT[]"),
Self::NumericArray => f.write_str("NUMERIC[]"),
Self::DateArray => f.write_str("DATE[]"),
Self::TimestampArray => f.write_str("TIMESTAMP[]"),
Self::TimestamptzArray => f.write_str("TIMESTAMPTZ[]"),
Self::UuidArray => f.write_str("UUID[]"),
Self::JsonArray => f.write_str("JSON[]"),
Self::JsonbArray => f.write_str("JSONB[]"),
Self::BytesArray => f.write_str("BYTEA[]"),
Self::VarcharArray => f.write_str("VARCHAR[]"),
Self::CharArray => f.write_str("CHAR[]"),
Self::Multirange(k) => f.write_str(match k {
RangeKind::Int4 => "INT4MULTIRANGE",
RangeKind::Int8 => "INT8MULTIRANGE",
RangeKind::Num => "NUMMULTIRANGE",
RangeKind::Ts => "TSMULTIRANGE",
RangeKind::TsTz => "TSTZMULTIRANGE",
RangeKind::Date => "DATEMULTIRANGE",
}),
Self::Point => f.write_str("POINT"),
Self::Lseg => f.write_str("LSEG"),
Self::Path => f.write_str("PATH"),
Self::PgBox => f.write_str("BOX"),
Self::Polygon => f.write_str("POLYGON"),
Self::Line => f.write_str("LINE"),
Self::Circle => f.write_str("CIRCLE"),
Self::Inet => f.write_str("INET"),
Self::Cidr => f.write_str("CIDR"),
Self::Macaddr => f.write_str("MACADDR"),
Self::Macaddr8 => f.write_str("MACADDR8"),
Self::PgLsn => f.write_str("PG_LSN"),
Self::Bit(0) => f.write_str("BIT"),
Self::Bit(n) => write!(f, "BIT({n})"),
Self::BitVarying(0) => f.write_str("VARBIT"),
Self::BitVarying(n) => write!(f, "VARBIT({n})"),
Self::Xml => f.write_str("XML"),
Self::Char1 => f.write_str("\"char\""),
Self::MoneyArray => f.write_str("MONEY[]"),
Self::TsVector => f.write_str("TSVECTOR"),
Self::TsQuery => f.write_str("TSQUERY"),
Self::Uuid => f.write_str("UUID"),
Self::Time => f.write_str("TIME"),
Self::Year => f.write_str("YEAR"),
Self::TimeTz => f.write_str("TIMETZ"),
Self::Money => f.write_str("MONEY"),
Self::Range(k) => f.write_str(k.keyword()),
Self::Hstore => f.write_str("HSTORE"),
Self::IntArray2D => f.write_str("INT[][]"),
Self::BigIntArray2D => f.write_str("BIGINT[][]"),
Self::TextArray2D => f.write_str("TEXT[][]"),
Self::BoolArray2D => f.write_str("BOOL[][]"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TsLexeme {
pub word: String,
pub positions: Vec<u16>,
pub weight: u8,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TsQueryAst {
Term {
word: String,
weight_mask: u8,
},
And(Box<TsQueryAst>, Box<TsQueryAst>),
Or(Box<TsQueryAst>, Box<TsQueryAst>),
Not(Box<TsQueryAst>),
Phrase {
left: Box<TsQueryAst>,
right: Box<TsQueryAst>,
distance: u16,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Hash)]
pub enum NumericKind {
#[default]
Finite,
NaN,
PosInf,
NegInf,
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum Value<'arena> {
SmallInt(i16),
Int(i32),
BigInt(i64),
Float(f64),
Real(f32),
Text(Cow<'arena, str>),
Bool(bool),
Vector(Cow<'arena, [f32]>),
Sq8Vector(crate::quantize::Sq8Vector),
HalfVector(crate::halfvec::HalfVector),
Numeric {
scaled: i128,
scale: u16,
kind: NumericKind,
},
NumericBig(alloc::boxed::Box<crate::bignum::BigNumeric>),
Date(i32),
Timestamp(i64),
Interval {
months: i32,
days: i32,
micros: i64,
},
Json(Cow<'arena, str>),
Bytes(Cow<'arena, [u8]>),
TextArray(Vec<Option<String>>),
IntArray(Vec<Option<i32>>),
BigIntArray(Vec<Option<i64>>),
IntervalArray(Vec<Option<IntervalSpan>>),
BoolArray(Vec<Option<bool>>),
SmallIntArray(Vec<Option<i16>>),
FloatArray(Vec<Option<f64>>),
NumericArray(Vec<Option<(i128, u16)>>),
DateArray(Vec<Option<i32>>),
TimestampArray(Vec<Option<i64>>),
TimestamptzArray(Vec<Option<i64>>),
UuidArray(Vec<Option<[u8; 16]>>),
JsonArray(Vec<Option<String>>),
JsonbArray(Vec<Option<String>>),
BytesArray(Vec<Option<Vec<u8>>>),
VarcharArray(Vec<Option<String>>),
CharArray(Vec<Option<String>>),
Multirange {
kind: RangeKind,
ranges: Vec<RangeSpan>,
},
Point(Point2D),
Lseg(Point2D, Point2D),
Path {
points: Vec<Point2D>,
closed: bool,
},
PgBox(Point2D, Point2D),
Polygon(Vec<Point2D>),
Line {
a: f64,
b: f64,
c: f64,
},
Circle {
center: Point2D,
radius: f64,
},
Inet {
family: u8,
bits: u8,
addr: [u8; 16],
},
Cidr {
family: u8,
bits: u8,
addr: [u8; 16],
},
Macaddr([u8; 6]),
Macaddr8([u8; 8]),
PgLsn(u64),
RegClass(i64, alloc::boxed::Box<str>),
RegProc(i64, alloc::boxed::Box<str>),
RegType(i64, alloc::boxed::Box<str>),
Xid(u32),
Cid(u32),
Tid(u32, u32),
BitString {
nbits: u32,
bytes: Cow<'arena, [u8]>,
},
Xml(Cow<'arena, str>),
Char1(u8),
BpChar(Cow<'arena, str>),
MoneyArray(Vec<Option<i64>>),
TsVector(Vec<TsLexeme>),
TsQuery(TsQueryAst),
Uuid([u8; 16]),
Time(i64),
Year(u16),
TimeTz {
us: i64,
offset_secs: i32,
},
Money(i64),
Hstore(Vec<(String, Option<String>)>),
IntArray2D(Vec<Vec<Option<i32>>>),
BigIntArray2D(Vec<Vec<Option<i64>>>),
TextArray2D(Vec<Vec<Option<String>>>),
BoolArray2D(Vec<Vec<Option<bool>>>),
Range {
kind: RangeKind,
lower: Option<alloc::boxed::Box<Value<'static>>>,
upper: Option<alloc::boxed::Box<Value<'static>>>,
lower_inc: bool,
upper_inc: bool,
empty: bool,
},
Composite(alloc::vec::Vec<(alloc::string::String, Value<'static>)>),
Null,
}
pub type ValueOwned = Value<'static>;
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct Point2D {
pub x: f64,
pub y: f64,
}
#[derive(Debug, Clone, PartialEq)]
pub struct RangeSpan {
pub lower: Option<alloc::boxed::Box<Value<'static>>>,
pub upper: Option<alloc::boxed::Box<Value<'static>>>,
pub lower_inc: bool,
pub upper_inc: bool,
pub empty: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct IntervalSpan {
pub months: i32,
pub days: i32,
pub micros: i64,
}
impl<'arena> Value<'arena> {
pub fn data_type(&self) -> Option<DataType> {
match self {
Self::SmallInt(_) => Some(DataType::SmallInt),
Self::Int(_) => Some(DataType::Int),
Self::BigInt(_) => Some(DataType::BigInt),
Self::Float(_) => Some(DataType::Float),
Self::Real(_) => Some(DataType::Real),
Self::Text(_) => Some(DataType::Text),
Self::Bool(_) => Some(DataType::Bool),
Self::Vector(v) => Some(DataType::Vector {
dim: u32::try_from(v.len()).expect("vector dim ≤ u32"),
encoding: VecEncoding::F32,
}),
Self::Sq8Vector(q) => Some(DataType::Vector {
dim: u32::try_from(q.bytes.len()).expect("vector dim ≤ u32"),
encoding: VecEncoding::Sq8,
}),
Self::HalfVector(h) => Some(DataType::Vector {
dim: u32::try_from(h.dim()).expect("vector dim ≤ u32"),
encoding: VecEncoding::F16,
}),
Self::Numeric { scale, .. } => Some(DataType::Numeric {
precision: 0,
scale: i16::try_from(*scale).unwrap_or(i16::MAX),
}),
Self::NumericBig(b) => Some(DataType::Numeric {
precision: 0,
scale: i16::try_from(b.scale()).unwrap_or(i16::MAX),
}),
Self::Date(_) => Some(DataType::Date),
Self::Timestamp(_) => Some(DataType::Timestamp),
Self::Interval { .. } => Some(DataType::Interval),
Self::Json(_) => Some(DataType::Json),
Self::Bytes(_) => Some(DataType::Bytes),
Self::TextArray(_) => Some(DataType::TextArray),
Self::IntArray(_) => Some(DataType::IntArray),
Self::BigIntArray(_) => Some(DataType::BigIntArray),
Self::IntervalArray(_) => Some(DataType::IntervalArray),
Self::BoolArray(_) => Some(DataType::BoolArray),
Self::SmallIntArray(_) => Some(DataType::SmallIntArray),
Self::FloatArray(_) => Some(DataType::FloatArray),
Self::NumericArray(_) => Some(DataType::NumericArray),
Self::DateArray(_) => Some(DataType::DateArray),
Self::TimestampArray(_) => Some(DataType::TimestampArray),
Self::TimestamptzArray(_) => Some(DataType::TimestamptzArray),
Self::UuidArray(_) => Some(DataType::UuidArray),
Self::JsonArray(_) => Some(DataType::JsonArray),
Self::JsonbArray(_) => Some(DataType::JsonbArray),
Self::BytesArray(_) => Some(DataType::BytesArray),
Self::VarcharArray(_) => Some(DataType::VarcharArray),
Self::CharArray(_) => Some(DataType::CharArray),
Self::Multirange { kind, .. } => Some(DataType::Multirange(*kind)),
Self::Point(_) => Some(DataType::Point),
Self::Lseg(_, _) => Some(DataType::Lseg),
Self::Path { .. } => Some(DataType::Path),
Self::PgBox(_, _) => Some(DataType::PgBox),
Self::Polygon(_) => Some(DataType::Polygon),
Self::Line { .. } => Some(DataType::Line),
Self::Circle { .. } => Some(DataType::Circle),
Self::Inet { .. } => Some(DataType::Inet),
Self::Cidr { .. } => Some(DataType::Cidr),
Self::Macaddr(_) => Some(DataType::Macaddr),
Self::Macaddr8(_) => Some(DataType::Macaddr8),
Self::PgLsn(_) => Some(DataType::PgLsn),
Self::BitString { .. } => Some(DataType::BitVarying(0)),
Self::Xml(_) => Some(DataType::Xml),
Self::Char1(_) => Some(DataType::Char1),
Self::BpChar(s) => Some(DataType::Char(
u32::try_from(s.chars().count()).unwrap_or(0),
)),
Self::MoneyArray(_) => Some(DataType::MoneyArray),
Self::TsVector(_) => Some(DataType::TsVector),
Self::TsQuery(_) => Some(DataType::TsQuery),
Self::Uuid(_) => Some(DataType::Uuid),
Self::Time(_) => Some(DataType::Time),
Self::Year(_) => Some(DataType::Year),
Self::TimeTz { .. } => Some(DataType::TimeTz),
Self::Money(_) => Some(DataType::Money),
Self::Range { kind, .. } => Some(DataType::Range(*kind)),
Self::Hstore(_) => Some(DataType::Hstore),
Self::IntArray2D(_) => Some(DataType::IntArray2D),
Self::BigIntArray2D(_) => Some(DataType::BigIntArray2D),
Self::TextArray2D(_) => Some(DataType::TextArray2D),
Self::BoolArray2D(_) => Some(DataType::BoolArray2D),
Self::Composite(_) => None,
Self::Xid(_) => Some(DataType::Xid),
Self::RegClass(..)
| Self::RegProc(..)
| Self::RegType(..)
| Self::Tid(..)
| Self::Cid(_) => None,
Self::Null => None,
}
}
pub const fn is_null(&self) -> bool {
matches!(self, Self::Null)
}
pub fn into_owned(self) -> Value<'static> {
match self {
Value::SmallInt(n) => Value::SmallInt(n),
Value::Int(n) => Value::Int(n),
Value::BigInt(n) => Value::BigInt(n),
Value::Float(f) => Value::Float(f),
Value::Real(f) => Value::Real(f),
Value::Text(s) => Value::Text(Cow::Owned(s.into_owned())),
Value::Bool(b) => Value::Bool(b),
Value::Vector(v) => Value::Vector(Cow::Owned(v.into_owned())),
Value::Sq8Vector(q) => Value::Sq8Vector(q),
Value::HalfVector(h) => Value::HalfVector(h),
Value::Numeric {
scaled,
scale,
kind,
} => Value::Numeric {
scaled,
scale,
kind,
},
Value::NumericBig(b) => Value::NumericBig(b),
Value::Date(d) => Value::Date(d),
Value::Timestamp(t) => Value::Timestamp(t),
Value::Interval {
months,
days,
micros,
} => Value::Interval {
months,
days,
micros,
},
Value::Json(s) => Value::Json(Cow::Owned(s.into_owned())),
Value::Bytes(b) => Value::Bytes(Cow::Owned(b.into_owned())),
Value::TextArray(v) => Value::TextArray(v),
Value::IntArray(v) => Value::IntArray(v),
Value::BigIntArray(v) => Value::BigIntArray(v),
Value::IntervalArray(v) => Value::IntervalArray(v),
Value::BoolArray(v) => Value::BoolArray(v),
Value::SmallIntArray(v) => Value::SmallIntArray(v),
Value::FloatArray(v) => Value::FloatArray(v),
Value::NumericArray(v) => Value::NumericArray(v),
Value::DateArray(v) => Value::DateArray(v),
Value::TimestampArray(v) => Value::TimestampArray(v),
Value::TimestamptzArray(v) => Value::TimestamptzArray(v),
Value::UuidArray(v) => Value::UuidArray(v),
Value::JsonArray(v) => Value::JsonArray(v),
Value::JsonbArray(v) => Value::JsonbArray(v),
Value::BytesArray(v) => Value::BytesArray(v),
Value::VarcharArray(v) => Value::VarcharArray(v),
Value::CharArray(v) => Value::CharArray(v),
Value::Multirange { kind, ranges } => Value::Multirange { kind, ranges },
Value::Composite(fields) => Value::Composite(fields),
Value::RegClass(oid, name) => Value::RegClass(oid, name),
Value::Tid(b, o) => Value::Tid(b, o),
Value::Xid(x) => Value::Xid(x),
Value::Cid(c) => Value::Cid(c),
Value::RegProc(oid, name) => Value::RegProc(oid, name),
Value::RegType(oid, name) => Value::RegType(oid, name),
Value::Point(p) => Value::Point(p),
Value::Lseg(a, b) => Value::Lseg(a, b),
Value::Path { points, closed } => Value::Path { points, closed },
Value::PgBox(a, b) => Value::PgBox(a, b),
Value::Polygon(p) => Value::Polygon(p),
Value::Line { a, b, c } => Value::Line { a, b, c },
Value::Circle { center, radius } => Value::Circle { center, radius },
Value::Inet { family, bits, addr } => Value::Inet { family, bits, addr },
Value::Cidr { family, bits, addr } => Value::Cidr { family, bits, addr },
Value::Macaddr(m) => Value::Macaddr(m),
Value::Macaddr8(m) => Value::Macaddr8(m),
Value::PgLsn(l) => Value::PgLsn(l),
Value::BitString { nbits, bytes } => Value::BitString {
nbits,
bytes: Cow::Owned(bytes.into_owned()),
},
Value::Xml(s) => Value::Xml(Cow::Owned(s.into_owned())),
Value::Char1(c) => Value::Char1(c),
Value::BpChar(s) => Value::BpChar(Cow::Owned(s.into_owned())),
Value::MoneyArray(v) => Value::MoneyArray(v),
Value::TsVector(v) => Value::TsVector(v),
Value::TsQuery(q) => Value::TsQuery(q),
Value::Uuid(u) => Value::Uuid(u),
Value::Time(t) => Value::Time(t),
Value::Year(y) => Value::Year(y),
Value::TimeTz { us, offset_secs } => Value::TimeTz { us, offset_secs },
Value::Money(m) => Value::Money(m),
Value::Range {
kind,
lower,
upper,
lower_inc,
upper_inc,
empty,
} => Value::Range {
kind,
lower,
upper,
lower_inc,
upper_inc,
empty,
},
Value::Hstore(h) => Value::Hstore(h),
Value::IntArray2D(a) => Value::IntArray2D(a),
Value::BigIntArray2D(a) => Value::BigIntArray2D(a),
Value::TextArray2D(a) => Value::TextArray2D(a),
Value::BoolArray2D(a) => Value::BoolArray2D(a),
Value::Null => Value::Null,
}
}
pub fn clone_into<'a>(&self, arena: &'a bumpalo::Bump) -> Value<'a> {
match self {
Value::Text(s) => Value::Text(Cow::Borrowed(arena.alloc_str(s))),
Value::Json(s) => Value::Json(Cow::Borrowed(arena.alloc_str(s))),
Value::Xml(s) => Value::Xml(Cow::Borrowed(arena.alloc_str(s))),
Value::BpChar(s) => Value::BpChar(Cow::Borrowed(arena.alloc_str(s))),
Value::Bytes(b) => {
let slot = arena.alloc_slice_copy::<u8>(b);
Value::Bytes(Cow::Borrowed(slot))
}
Value::Vector(v) => {
let slot = arena.alloc_slice_copy::<f32>(v);
Value::Vector(Cow::Borrowed(slot))
}
Value::BitString { nbits, bytes } => {
let slot = arena.alloc_slice_copy::<u8>(bytes);
Value::BitString {
nbits: *nbits,
bytes: Cow::Borrowed(slot),
}
}
other => other.clone().into_owned(),
}
}
}
impl Value<'static> {
pub fn text<S: Into<String>>(s: S) -> Self {
Value::Text(Cow::Owned(s.into()))
}
pub const fn numeric(scaled: i128, scale: u16) -> Self {
Value::Numeric {
scaled,
scale,
kind: NumericKind::Finite,
}
}
pub const fn numeric_special(kind: NumericKind) -> Self {
Value::Numeric {
scaled: 0,
scale: 0,
kind,
}
}
pub fn json<S: Into<String>>(s: S) -> Self {
Value::Json(Cow::Owned(s.into()))
}
pub fn xml<S: Into<String>>(s: S) -> Self {
Value::Xml(Cow::Owned(s.into()))
}
pub fn bytes<B: Into<Vec<u8>>>(b: B) -> Self {
Value::Bytes(Cow::Owned(b.into()))
}
pub fn vector<V: Into<Vec<f32>>>(v: V) -> Self {
Value::Vector(Cow::Owned(v.into()))
}
pub fn bit_string<B: Into<Vec<u8>>>(nbits: u32, bytes: B) -> Self {
Value::BitString {
nbits,
bytes: Cow::Owned(bytes.into()),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct Row<'arena> {
pub values: Vec<Value<'arena>>,
}
pub type RowOwned = Row<'static>;
impl<'arena> Row<'arena> {
pub const fn new(values: Vec<Value<'arena>>) -> Self {
Self { values }
}
pub fn len(&self) -> usize {
self.values.len()
}
pub fn is_empty(&self) -> bool {
self.values.is_empty()
}
}
impl<'arena> Row<'arena> {
pub fn clone_into<'a>(&self, arena: &'a bumpalo::Bump) -> Row<'a> {
Row {
values: self.values.iter().map(|v| v.clone_into(arena)).collect(),
}
}
pub fn into_owned(self) -> Row<'static> {
Row {
values: self.values.into_iter().map(Value::into_owned).collect(),
}
}
}
impl Row<'static> {
pub fn from_arena(row: Row<'_>) -> Self {
Self {
values: row.values.into_iter().map(Value::into_owned).collect(),
}
}
}
#[allow(clippy::struct_excessive_bools)]
#[derive(Debug, Clone, PartialEq)]
pub struct ColumnSchema {
pub name: String,
pub ty: DataType,
pub nullable: bool,
pub default: Option<Value<'static>>,
pub runtime_default: Option<String>,
pub collation_name: Option<String>,
pub auto_increment: bool,
pub user_enum_type: Option<String>,
pub user_domain_type: Option<String>,
pub user_composite_type: Option<String>,
pub acl: Vec<AclItem>,
pub on_update_runtime: Option<String>,
pub collation: Collation,
pub is_unsigned: bool,
pub inline_enum_variants: Option<Vec<String>>,
pub inline_set_variants: Option<Vec<String>>,
pub generated_stored_expr: Option<String>,
pub identity_always: bool,
pub default_text: Option<String>,
pub auto_restart: Option<i64>,
pub scalar_row_source: bool,
pub mysql_int_width: Option<MysqlIntWidth>,
pub mysql_fsp: Option<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Collation {
Binary,
CaseInsensitive,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MysqlIntWidth {
Tiny,
Small,
Medium,
Int,
Big,
}
#[must_use]
pub fn mysql_ci_fold(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for ch in s.chars() {
for lc in ch.to_lowercase() {
match fold_latin_base(lc) {
Some(base) => out.push_str(base),
None => out.push(lc),
}
}
}
out
}
pub fn mysql_compare_fold(s: &str) -> String {
mysql_ci_fold(s.trim_end_matches(' '))
}
fn fold_latin_base(c: char) -> Option<&'static str> {
Some(match c {
'à' | 'á' | 'â' | 'ã' | 'ä' | 'å' | 'ā' | 'ă' | 'ą' => "a",
'æ' => "ae",
'ç' | 'ć' | 'č' | 'ĉ' | 'ċ' => "c",
'ð' | 'ď' | 'đ' => "d",
'è' | 'é' | 'ê' | 'ë' | 'ē' | 'ĕ' | 'ė' | 'ę' | 'ě' => "e",
'ĝ' | 'ğ' | 'ġ' | 'ģ' => "g",
'ì' | 'í' | 'î' | 'ï' | 'ĩ' | 'ī' | 'ĭ' | 'į' => "i",
'ĵ' => "j",
'ķ' => "k",
'ł' | 'ĺ' | 'ļ' | 'ľ' => "l",
'ñ' | 'ń' | 'ņ' | 'ň' => "n",
'ò' | 'ó' | 'ô' | 'õ' | 'ö' | 'ø' | 'ō' | 'ŏ' | 'ő' => "o",
'œ' => "oe",
'ŕ' | 'ŗ' | 'ř' => "r",
'ś' | 'š' | 'ŝ' | 'ş' => "s",
'ß' => "ss",
'ţ' | 'ť' | 'ŧ' => "t",
'ù' | 'ú' | 'û' | 'ü' | 'ũ' | 'ū' | 'ŭ' | 'ů' | 'ű' | 'ų' => "u",
'ý' | 'ÿ' => "y",
'ź' | 'ž' | 'ż' => "z",
_ => return None,
})
}
#[allow(clippy::derivable_impls)]
impl Default for Collation {
fn default() -> Self {
Self::Binary
}
}
impl Collation {
pub const TAG_BINARY: u8 = 0;
pub const TAG_CASE_INSENSITIVE: u8 = 1;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PolicyCmd {
All,
Select,
Insert,
Update,
Delete,
}
impl PolicyCmd {
#[must_use]
pub const fn as_pg_char(self) -> char {
match self {
Self::All => '*',
Self::Select => 'r',
Self::Insert => 'a',
Self::Update => 'w',
Self::Delete => 'd',
}
}
#[must_use]
pub const fn as_pg_word(self) -> &'static str {
match self {
Self::All => "ALL",
Self::Select => "SELECT",
Self::Insert => "INSERT",
Self::Update => "UPDATE",
Self::Delete => "DELETE",
}
}
#[must_use]
pub const fn to_wire_byte(self) -> u8 {
match self {
Self::All => 0,
Self::Select => 1,
Self::Insert => 2,
Self::Update => 3,
Self::Delete => 4,
}
}
#[must_use]
pub const fn from_wire_byte(b: u8) -> Option<Self> {
match b {
0 => Some(Self::All),
1 => Some(Self::Select),
2 => Some(Self::Insert),
3 => Some(Self::Update),
4 => Some(Self::Delete),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct PolicyDef {
pub name: String,
pub cmd: PolicyCmd,
pub permissive: bool,
pub roles: Vec<String>,
pub using_expr: Option<String>,
pub with_check_expr: Option<String>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct TableSchema {
pub name: String,
pub columns: Vec<ColumnSchema>,
pub hot_tier_bytes: Option<u64>,
pub foreign_keys: Vec<ForeignKeyConstraint>,
pub uniqueness_constraints: Vec<UniquenessConstraint>,
pub exclusion_constraints: Vec<ExclusionConstraint>,
pub checks: Vec<CheckConstraint>,
pub partition_role: Option<PartitionRole>,
pub policies: Vec<PolicyDef>,
pub row_security: bool,
pub force_row_security: bool,
pub owner: Option<String>,
pub acl: Vec<AclItem>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AclItem {
pub grantee: String,
pub privs: u16,
pub grantable: u16,
pub grantor: String,
}
pub mod priv_bits {
pub const INSERT: u16 = 1 << 0; pub const SELECT: u16 = 1 << 1; pub const UPDATE: u16 = 1 << 2; pub const DELETE: u16 = 1 << 3; pub const TRUNCATE: u16 = 1 << 4; pub const REFERENCES: u16 = 1 << 5; pub const TRIGGER: u16 = 1 << 6; pub const MAINTAIN: u16 = 1 << 7; pub const USAGE: u16 = 1 << 8; pub const CREATE: u16 = 1 << 9; pub const CONNECT: u16 = 1 << 10; pub const TEMPORARY: u16 = 1 << 11; pub const EXECUTE: u16 = 1 << 12; pub const ALL: u16 =
INSERT | SELECT | UPDATE | DELETE | TRUNCATE | REFERENCES | TRIGGER | MAINTAIN;
pub const ALL_SEQUENCE: u16 = SELECT | UPDATE | USAGE;
pub const ALL_SCHEMA: u16 = USAGE | CREATE;
pub const ALL_DATABASE: u16 = CREATE | CONNECT | TEMPORARY;
pub const ALL_FUNCTION: u16 = EXECUTE;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PartitionRole {
Parent {
kind: PartitionKind,
key_column_positions: Vec<usize>,
index_template_sources: Vec<String>,
},
Range {
parent_name: String,
lower: PartitionBound,
upper: PartitionBound,
},
List {
parent_name: String,
values: Vec<PartitionBound>,
},
Inherits {
parent_names: Vec<String>,
},
Hash {
parent_name: String,
modulus: u32,
remainder: u32,
},
Default {
parent_name: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PartitionKind {
Range,
List,
Hash,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PartitionBound {
MinValue,
MaxValue,
TimestampTz(i64),
BigInt(i64),
Int(i32),
SmallInt(i16),
Date(i32),
Text(alloc::string::String),
}
impl PartitionBound {
#[must_use]
pub fn equals_value(&self, other: &Value<'_>) -> bool {
match (self, other) {
(PartitionBound::TimestampTz(a), Value::Timestamp(b)) => a == b,
(PartitionBound::BigInt(a), Value::BigInt(b)) => a == b,
(PartitionBound::Int(a), Value::Int(b)) => a == b,
(PartitionBound::SmallInt(a), Value::SmallInt(b)) => a == b,
(PartitionBound::Date(a), Value::Date(b)) => a == b,
(PartitionBound::Text(a), Value::Text(b)) => a.as_str() == b.as_ref(),
_ => false,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CheckConstraint {
pub name: Option<String>,
pub expr: String,
pub validated: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UniquenessConstraint {
pub is_primary_key: bool,
pub columns: Vec<usize>,
pub nulls_not_distinct: bool,
pub name: Option<String>,
pub deferrable: bool,
pub initially_deferred: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ExclusionConstraint {
pub name: String,
pub method: Option<String>,
pub elements: Vec<(usize, String)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ForeignKeyConstraint {
pub name: Option<String>,
pub local_columns: Vec<usize>,
pub parent_table: String,
pub parent_columns: Vec<usize>,
pub on_delete: FkAction,
pub on_update: FkAction,
pub match_type: MatchType,
pub deferrable: bool,
pub initially_deferred: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum MatchType {
#[default]
Simple,
Full,
}
impl MatchType {
pub const fn tag(self) -> u8 {
match self {
Self::Simple => 0,
Self::Full => 1,
}
}
pub const fn from_tag(b: u8) -> Option<Self> {
Some(match b {
0 => Self::Simple,
1 => Self::Full,
_ => return None,
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FkAction {
Restrict,
Cascade,
SetNull,
SetDefault,
NoAction,
}
impl FkAction {
pub const fn tag(self) -> u8 {
match self {
Self::Restrict => 0,
Self::Cascade => 1,
Self::SetNull => 2,
Self::SetDefault => 3,
Self::NoAction => 4,
}
}
pub const fn from_tag(b: u8) -> Option<Self> {
Some(match b {
0 => Self::Restrict,
1 => Self::Cascade,
2 => Self::SetNull,
3 => Self::SetDefault,
4 => Self::NoAction,
_ => return None,
})
}
}
impl TableSchema {
pub fn column_position(&self, name: &str) -> Option<usize> {
self.columns.iter().position(|c| c.name == name)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub enum IndexKey {
Int(i64),
Text(String),
Bool(bool),
Uuid([u8; 16]),
}
impl IndexKey {
#[inline]
pub fn from_i64(n: i64) -> Self {
Self::Int(n)
}
pub fn from_value(v: &Value<'_>) -> Option<Self> {
match v {
Value::BigInt(n) => Some(Self::Int(*n)),
Value::SmallInt(n) => Some(Self::Int(i64::from(*n))),
Value::Int(n) => Some(Self::Int(i64::from(*n))),
Value::Text(s) => Some(Self::Text(s.clone().into_owned())),
Value::BpChar(s) => Some(Self::Text(s.trim_end_matches(' ').to_string())),
Value::Bool(b) => Some(Self::Bool(*b)),
Value::Date(d) => Some(Self::Int(i64::from(*d))),
Value::Timestamp(t) => Some(Self::Int(*t)),
Value::Uuid(b) => Some(Self::Uuid(*b)),
Value::Time(us) => Some(Self::Int(*us)),
Value::Year(y) => Some(Self::Int(i64::from(*y))),
Value::TimeTz { us, offset_secs } => {
Some(Self::Int(us - i64::from(*offset_secs) * 1_000_000))
}
Value::Money(c) => Some(Self::Int(*c)),
Value::Range { .. } => None,
Value::Hstore(_) => None,
Value::NumericBig(_) => None,
Value::IntArray2D(_)
| Value::BigIntArray2D(_)
| Value::TextArray2D(_)
| Value::BoolArray2D(_) => None,
Value::IntervalArray(_) => None,
Value::BoolArray(_)
| Value::SmallIntArray(_)
| Value::FloatArray(_)
| Value::NumericArray(_)
| Value::DateArray(_)
| Value::TimestampArray(_)
| Value::TimestamptzArray(_)
| Value::UuidArray(_)
| Value::JsonArray(_)
| Value::JsonbArray(_)
| Value::BytesArray(_)
| Value::VarcharArray(_)
| Value::CharArray(_)
| Value::Multirange { .. }
| Value::Point(_)
| Value::Lseg(_, _)
| Value::Path { .. }
| Value::PgBox(_, _)
| Value::Polygon(_)
| Value::Line { .. }
| Value::Circle { .. }
| Value::Inet { .. }
| Value::Cidr { .. }
| Value::Macaddr(_)
| Value::Macaddr8(_)
| Value::PgLsn(_)
| Value::BitString { .. }
| Value::Xml(_)
| Value::Char1(_)
| Value::MoneyArray(_)
| Value::Composite(_)
| Value::Tid(..)
| Value::Xid(_)
| Value::Cid(_)
| Value::RegClass(..)
| Value::RegProc(..)
| Value::RegType(..) => None,
Value::Null
| Value::Float(_)
| Value::Vector(_)
| Value::Sq8Vector(_)
| Value::HalfVector(_)
| Value::Numeric { .. }
| Value::Interval { .. }
| Value::Json(_)
| Value::Bytes(_)
| Value::TextArray(_)
| Value::IntArray(_)
| Value::BigIntArray(_)
| Value::TsVector(_)
| Value::TsQuery(_)
| Value::Real(_) => None,
}
}
}
#[derive(Debug, Clone)]
pub struct Index {
pub name: String,
pub column_position: usize,
pub kind: IndexKind,
pub included_columns: Vec<usize>,
pub partial_predicate: Option<String>,
pub expression: Option<String>,
pub nulls_not_distinct: bool,
pub descending: bool,
pub nulls_first: Option<bool>,
pub collation: Option<String>,
pub is_unique: bool,
pub extra_column_positions: Vec<usize>,
}
pub const NSW_DEFAULT_M: usize = 16;
#[derive(Debug, Clone)]
pub struct FreezeReport {
pub segment_id: u32,
pub frozen_rows: usize,
pub bytes_freed: u64,
pub segment_bytes: Vec<u8>,
}
#[derive(Debug, Clone)]
pub struct FreezeSlice {
pub row_range: core::ops::Range<usize>,
pub rows: Vec<(u64, Vec<u8>, IndexKey)>,
}
#[derive(Debug, Clone)]
pub struct CompactReport {
pub sources: Vec<u32>,
pub merged_segment_id: Option<u32>,
pub merged_segment_bytes: Vec<u8>,
pub merged_rows: usize,
pub deleted_rows_pruned: usize,
pub bytes_reclaimed_estimate: u64,
}
#[derive(Debug, Clone)]
pub enum IndexKind {
BTree(PersistentBTreeMap<IndexKey, Vec<RowLocator>>),
Nsw(NswGraph),
Brin {
column_type: DataType,
},
Gin(PersistentBTreeMap<alloc::string::String, Vec<RowLocator>>),
GinTrgm(PersistentBTreeMap<alloc::string::String, Vec<RowLocator>>),
GinFulltext(PersistentBTreeMap<alloc::string::String, Vec<RowLocator>>),
GinJsonb(PersistentBTreeMap<alloc::string::String, Vec<RowLocator>>),
}
impl IndexKind {
#[must_use]
pub fn approx_resident_bytes(&self) -> u64 {
const HEADER: usize = 24; let loc = core::mem::size_of::<RowLocator>();
match self {
IndexKind::BTree(map) => {
let key = core::mem::size_of::<IndexKey>();
map.iter()
.map(|(_, locs)| (key + HEADER + locs.len() * loc) as u64)
.sum()
}
IndexKind::Nsw(g) => {
let mut b = g.levels.len() as u64;
for layer in &g.layers {
for nbrs in layer.iter() {
b += (HEADER + nbrs.len() * core::mem::size_of::<u32>()) as u64;
}
}
b
}
IndexKind::Brin { .. } => core::mem::size_of::<DataType>() as u64,
IndexKind::Gin(map)
| IndexKind::GinTrgm(map)
| IndexKind::GinFulltext(map)
| IndexKind::GinJsonb(map) => map
.iter()
.map(|(word, postings)| {
(word.len() + HEADER + HEADER + postings.len() * loc) as u64
})
.sum(),
}
}
}
#[derive(Debug, Clone)]
pub struct NswGraph {
pub m: usize,
pub m_max_0: usize,
pub entry: Option<usize>,
pub entry_level: u8,
pub levels: PersistentVec<u8>,
pub layers: Vec<PersistentVec<Vec<u32>>>,
}
impl NswGraph {
fn new(m: usize) -> Self {
Self {
m,
m_max_0: m.saturating_mul(2),
entry: None,
entry_level: 0,
levels: PersistentVec::new(),
layers: alloc::vec![PersistentVec::new()],
}
}
pub const fn cap_for_layer(&self, layer: u8) -> usize {
if layer == 0 { self.m_max_0 } else { self.m }
}
}
#[allow(clippy::verbose_bit_mask)] pub fn nsw_assign_level(row_idx: usize) -> u8 {
const MAX_LEVEL: u8 = 7; let mut x = (row_idx as u64).wrapping_mul(0x9E37_79B9_7F4A_7C15);
x ^= x >> 30;
x = x.wrapping_mul(0xBF58_476D_1CE4_E5B9);
x ^= x >> 27;
x = x.wrapping_mul(0x94D0_49BB_1331_11EB);
x ^= x >> 31;
let mut level: u8 = 0;
while x & 0xF == 0 && level < MAX_LEVEL {
level += 1;
x >>= 4;
}
level
}
impl Index {
fn new_btree(name: String, column_position: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::BTree(PersistentBTreeMap::new()),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_nsw(name: String, column_position: usize, m: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::Nsw(NswGraph::new(m)),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_brin(name: String, column_position: usize, column_type: DataType) -> Self {
Self {
name,
column_position,
kind: IndexKind::Brin { column_type },
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_gin(name: String, column_position: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::Gin(PersistentBTreeMap::new()),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_gin_trgm(name: String, column_position: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::GinTrgm(PersistentBTreeMap::new()),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_gin_fulltext(name: String, column_position: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::GinFulltext(PersistentBTreeMap::new()),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
fn new_gin_jsonb(name: String, column_position: usize) -> Self {
Self {
name,
column_position,
kind: IndexKind::GinJsonb(PersistentBTreeMap::new()),
included_columns: Vec::new(),
partial_predicate: None,
expression: None,
is_unique: false,
nulls_not_distinct: false,
descending: false,
nulls_first: None,
collation: None,
extra_column_positions: Vec::new(),
}
}
pub fn iter_desc(
&self,
) -> alloc::boxed::Box<dyn Iterator<Item = (&IndexKey, &alloc::vec::Vec<RowLocator>)> + '_>
{
match &self.kind {
IndexKind::BTree(m) => alloc::boxed::Box::new(m.iter_rev()),
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => alloc::boxed::Box::new(core::iter::empty()),
}
}
pub fn iter_asc(
&self,
) -> alloc::boxed::Box<dyn Iterator<Item = (&IndexKey, &alloc::vec::Vec<RowLocator>)> + '_>
{
match &self.kind {
IndexKind::BTree(m) => alloc::boxed::Box::new(m.iter()),
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => alloc::boxed::Box::new(core::iter::empty()),
}
}
pub fn lookup_eq(&self, key: &IndexKey) -> &[RowLocator] {
match &self.kind {
IndexKind::BTree(m) => m.get(key).map_or(&[][..], Vec::as_slice),
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => &[][..],
}
}
#[inline]
pub fn lookup_eq_i64(&self, n: i64) -> &[RowLocator] {
match &self.kind {
IndexKind::BTree(m) => m.get(&IndexKey::Int(n)).map_or(&[][..], Vec::as_slice),
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => &[][..],
}
}
pub fn lookup_range_capped(
&self,
lo: core::ops::Bound<&IndexKey>,
hi: core::ops::Bound<&IndexKey>,
cap: usize,
) -> Option<Vec<RowLocator>> {
self.lookup_range_capped_by(lo, hi, cap, |_| true)
}
pub fn lookup_range_capped_by(
&self,
lo: core::ops::Bound<&IndexKey>,
hi: core::ops::Bound<&IndexKey>,
cap: usize,
keep: impl Fn(RowLocator) -> bool,
) -> Option<Vec<RowLocator>> {
match &self.kind {
IndexKind::BTree(m) => {
let mut out: Vec<RowLocator> = Vec::new();
for (_, locs) in m.range(lo, hi) {
out.extend(locs.iter().copied().filter(|l| keep(*l)));
if out.len() > cap {
return None;
}
}
Some(out)
}
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => None,
}
}
pub fn range_keyed(
&self,
lo: core::ops::Bound<&IndexKey>,
hi: core::ops::Bound<&IndexKey>,
) -> Option<impl Iterator<Item = (&IndexKey, RowLocator)> + '_> {
match &self.kind {
IndexKind::BTree(m) => Some(
m.range(lo, hi)
.flat_map(|(k, locs)| locs.iter().map(move |l| (k, *l))),
),
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => None,
}
}
pub fn gin_lookup_word(&self, word: &str) -> &[RowLocator] {
match &self.kind {
IndexKind::Gin(m) | IndexKind::GinFulltext(m) => {
m.get(&String::from(word)).map_or(&[][..], Vec::as_slice)
}
IndexKind::BTree(_)
| IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::GinTrgm(_)
| IndexKind::GinJsonb(_) => &[][..],
}
}
pub fn gin_trgm_lookup(&self, tri: &str) -> &[RowLocator] {
match &self.kind {
IndexKind::GinTrgm(m) => m.get(&String::from(tri)).map_or(&[][..], Vec::as_slice),
IndexKind::BTree(_)
| IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => &[][..],
}
}
pub fn gin_jsonb_lookup(&self, token: &str) -> &[RowLocator] {
match &self.kind {
IndexKind::GinJsonb(m) => m.get(&String::from(token)).map_or(&[][..], Vec::as_slice),
IndexKind::BTree(_)
| IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_) => &[][..],
}
}
pub const fn nsw(&self) -> Option<&NswGraph> {
match &self.kind {
IndexKind::Nsw(g) => Some(g),
IndexKind::BTree(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => None,
}
}
pub const fn is_brin(&self) -> bool {
matches!(self.kind, IndexKind::Brin { .. })
}
pub const fn is_gin_trgm(&self) -> bool {
matches!(self.kind, IndexKind::GinTrgm(_))
}
pub const fn is_gin(&self) -> bool {
matches!(self.kind, IndexKind::Gin(_))
}
pub const fn is_gin_fulltext(&self) -> bool {
matches!(self.kind, IndexKind::GinFulltext(_))
}
pub const fn is_gin_jsonb(&self) -> bool {
matches!(self.kind, IndexKind::GinJsonb(_))
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum RowChange {
Insert {
table: String,
row: Row<'static>,
rowid: row_header::RowId,
writer_version: u64,
},
Update {
table: String,
pos: usize,
new_row: Vec<Value<'static>>,
rowid: row_header::RowId,
writer_version: u64,
},
Delete {
table: String,
positions: Vec<usize>,
rowids: Vec<row_header::RowId>,
writer_version: u64,
},
Tombstone {
table: String,
rowids: Vec<row_header::RowId>,
xmax: u64,
},
}
impl RowChange {
#[must_use]
pub fn table_name(&self) -> &str {
match self {
Self::Insert { table, .. }
| Self::Update { table, .. }
| Self::Delete { table, .. }
| Self::Tombstone { table, .. } => table,
}
}
pub fn set_writer_version(&mut self, v: u64) {
match self {
RowChange::Insert { writer_version, .. }
| RowChange::Update { writer_version, .. }
| RowChange::Delete { writer_version, .. } => *writer_version = v,
RowChange::Tombstone { xmax, .. } => {
debug_assert_eq!(
*xmax, v,
"tombstone xmax must match the statement writer version"
);
*xmax = v;
}
}
}
}
const REDO_META_MARKER: u8 = 0xFF;
const REDO_META_VERSION: u8 = 1;
static UNRESOLVED_TOMBSTONES: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(0);
#[must_use]
pub fn unresolved_tombstones() -> u64 {
UNRESOLVED_TOMBSTONES.load(core::sync::atomic::Ordering::Relaxed)
}
#[must_use]
pub fn unresolved_tombstone_count() -> u64 {
UNRESOLVED_TOMBSTONES.load(core::sync::atomic::Ordering::Relaxed)
}
const _: () = assert!(FILE_VERSION < REDO_META_MARKER);
#[must_use]
pub fn encode_redo_log(changes: &[RowChange]) -> Vec<u8> {
let mut out = Vec::new();
out.push(REDO_META_MARKER);
out.push(REDO_META_VERSION);
out.push(FILE_VERSION);
codec::write_u32(&mut out, changes.len() as u32);
let write_values = |out: &mut Vec<u8>, vals: &[Value<'static>]| {
codec::write_u32(out, vals.len() as u32);
for v in vals {
codec::write_value(out, v);
}
};
for change in changes {
match change {
RowChange::Insert {
table,
row,
rowid,
writer_version,
} => {
out.push(0);
codec::write_str(&mut out, table);
write_values(&mut out, &row.values);
codec::write_u64(&mut out, rowid.0);
codec::write_u64(&mut out, *writer_version);
}
RowChange::Update {
table,
pos,
new_row,
rowid,
writer_version,
} => {
out.push(1);
codec::write_str(&mut out, table);
codec::write_u32(&mut out, *pos as u32);
write_values(&mut out, new_row);
codec::write_u64(&mut out, rowid.0);
codec::write_u64(&mut out, *writer_version);
}
RowChange::Delete {
table,
positions,
rowids,
writer_version,
} => {
out.push(2);
codec::write_str(&mut out, table);
codec::write_u32(&mut out, positions.len() as u32);
for p in positions {
codec::write_u32(&mut out, *p as u32);
}
debug_assert_eq!(
rowids.len(),
positions.len(),
"redo Delete: rowids must be parallel to positions"
);
for rid in rowids {
codec::write_u64(&mut out, rid.0);
}
codec::write_u64(&mut out, *writer_version);
}
RowChange::Tombstone {
table,
rowids,
xmax,
} => {
out.push(3);
codec::write_str(&mut out, table);
codec::write_u32(&mut out, rowids.len() as u32);
for rid in rowids {
codec::write_u64(&mut out, rid.0);
}
codec::write_u64(&mut out, *xmax);
}
}
}
out
}
pub fn decode_redo_log(bytes: &[u8]) -> Result<Vec<RowChange>, StorageError> {
let first = *bytes
.first()
.ok_or_else(|| StorageError::Corrupt("redo log: empty".into()))?;
let has_meta = first == REDO_META_MARKER;
let (codec_version, header_len) = if has_meta {
let meta_version = *bytes
.get(1)
.ok_or_else(|| StorageError::Corrupt("redo log: short header".into()))?;
if meta_version != REDO_META_VERSION {
return Err(StorageError::Corrupt(alloc::format!(
"redo log: unknown metadata version {meta_version}"
)));
}
let file_version = *bytes
.get(2)
.ok_or_else(|| StorageError::Corrupt("redo log: short header".into()))?;
(file_version, 3usize)
} else {
(first, 1usize)
};
let mut cur = codec::Cursor::new(bytes).with_codec_version(codec_version);
for _ in 0..header_len {
cur.read_u8()?;
}
let count = cur.read_u32()? as usize;
let mut read_values =
|cur: &mut codec::Cursor<'_>| -> Result<Vec<Value<'static>>, StorageError> {
let n = cur.read_u32()? as usize;
let mut vals = Vec::with_capacity(n);
for _ in 0..n {
vals.push(cur.read_value()?);
}
Ok(vals)
};
let mut changes = Vec::with_capacity(count);
for _ in 0..count {
let op = cur.read_u8()?;
let table = cur.read_str()?;
let change = match op {
0 => {
let row = Row::new(read_values(&mut cur)?);
let (rowid, writer_version) = if has_meta {
(row_header::RowId(cur.read_u64()?), cur.read_u64()?)
} else {
(row_header::RowId::UNASSIGNED, 0)
};
RowChange::Insert {
table,
row,
rowid,
writer_version,
}
}
1 => {
let pos = cur.read_u32()? as usize;
let new_row = read_values(&mut cur)?;
let (rowid, writer_version) = if has_meta {
(row_header::RowId(cur.read_u64()?), cur.read_u64()?)
} else {
(row_header::RowId::UNASSIGNED, 0)
};
RowChange::Update {
table,
pos,
new_row,
rowid,
writer_version,
}
}
2 => {
let n = cur.read_u32()? as usize;
let mut positions = Vec::with_capacity(n);
for _ in 0..n {
positions.push(cur.read_u32()? as usize);
}
let (rowids, writer_version) = if has_meta {
let mut rowids = Vec::with_capacity(n);
for _ in 0..n {
rowids.push(row_header::RowId(cur.read_u64()?));
}
(rowids, cur.read_u64()?)
} else {
(Vec::new(), 0)
};
RowChange::Delete {
table,
positions,
rowids,
writer_version,
}
}
3 if has_meta => {
let n = cur.read_u32()? as usize;
let mut rowids = Vec::with_capacity(n);
for _ in 0..n {
rowids.push(row_header::RowId(cur.read_u64()?));
}
let xmax = cur.read_u64()?;
RowChange::Tombstone {
table,
rowids,
xmax,
}
}
other => {
return Err(StorageError::Corrupt(alloc::format!(
"redo log: unknown op {other}"
)));
}
};
changes.push(change);
}
Ok(changes)
}
#[derive(Debug, Default)]
pub struct ScanStats {
pub seq_scan: core::sync::atomic::AtomicU64,
pub seq_tup_read: core::sync::atomic::AtomicU64,
pub idx_scan: core::sync::atomic::AtomicU64,
pub idx_tup_fetch: core::sync::atomic::AtomicU64,
}
impl Clone for ScanStats {
fn clone(&self) -> Self {
use core::sync::atomic::{AtomicU64, Ordering};
Self {
seq_scan: AtomicU64::new(self.seq_scan.load(Ordering::Relaxed)),
seq_tup_read: AtomicU64::new(self.seq_tup_read.load(Ordering::Relaxed)),
idx_scan: AtomicU64::new(self.idx_scan.load(Ordering::Relaxed)),
idx_tup_fetch: AtomicU64::new(self.idx_tup_fetch.load(Ordering::Relaxed)),
}
}
}
#[must_use]
pub fn range_excl_index_key(v: &Value<'_>) -> Option<(i128, u8)> {
let Value::Range {
lower,
lower_inc,
empty,
..
} = v
else {
return None;
};
if *empty {
return None;
}
let key = match lower {
None => i128::MIN,
Some(b) => match b.as_ref() {
Value::SmallInt(n) => i128::from(*n),
Value::Int(n) => i128::from(*n),
Value::BigInt(n) => i128::from(*n),
Value::Date(n) => i128::from(*n),
Value::Timestamp(n) => i128::from(*n),
_ => return None,
},
};
Some((key, u8::from(!*lower_inc)))
}
#[derive(Debug, Clone)]
pub struct ExclRangeIndex {
pub column_position: usize,
pub map: PersistentBTreeMap<(i128, u8), Vec<RowLocator>>,
}
#[derive(Debug, Clone)]
pub struct Table {
schema: TableSchema,
rel_id: row_header::RelId,
rows: PersistentVec<Row<'static>>,
headers: PersistentVec<row_header::RowHeader>,
rowids: PersistentVec<row_header::RowId>,
next_rowid: u64,
dead_rows: u64,
stat_tup_ins: u64,
stat_tup_upd: u64,
stat_tup_del: u64,
scan_stats: ScanStats,
last_autovacuum_us: Option<i64>,
last_analyze_us: Option<i64>,
indices: Vec<Index>,
hot_bytes: u64,
cold_row_count: u64,
cold_row_count_stale: bool,
redo_log: Option<Vec<RowChange>>,
excl_indexes: Vec<ExclRangeIndex>,
prune_horizon: u64,
}
#[derive(Debug, Default)]
pub struct ColdReadStats {
pub cold_reads: core::sync::atomic::AtomicU64,
}
impl Clone for ColdReadStats {
fn clone(&self) -> Self {
Self {
cold_reads: core::sync::atomic::AtomicU64::new(
self.cold_reads.load(core::sync::atomic::Ordering::Relaxed),
),
}
}
}
#[derive(Debug, Clone, Default)]
pub struct Catalog {
pub cold_read_stats: ColdReadStats,
tables: Vec<Table>,
by_name: BTreeMap<String, usize>,
temp_prefix: Option<String>,
dirty_tables: alloc::collections::BTreeSet<String>,
next_rel_id: u64,
cold_segments: Vec<Option<Arc<OwnedSegment>>>,
functions: BTreeMap<String, FunctionDef>,
triggers: Vec<TriggerDef>,
rules: Vec<RuleDef>,
statistics_ext: Vec<StatisticsExtDef>,
large_objects: alloc::collections::BTreeMap<u32, Vec<u8>>,
sequences: BTreeMap<String, SequenceDef>,
schema_acl: Vec<AclItem>,
database_acl: Vec<AclItem>,
views: BTreeMap<String, ViewDef>,
materialized_views: BTreeMap<String, String>,
enum_types: BTreeMap<String, EnumDef>,
domain_types: BTreeMap<String, DomainDef>,
comments: BTreeMap<String, String>,
db_role_settings: BTreeMap<(String, String), BTreeMap<String, String>>,
replication_slots: BTreeMap<String, (String, String)>,
composite_types: BTreeMap<String, CompositeDef>,
schemas: alloc::collections::BTreeSet<String>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct FunctionDef {
pub name: String,
pub args_repr: String,
pub returns: String,
pub language: String,
pub body: String,
pub owner: Option<String>,
pub acl: Vec<AclItem>,
pub volatility: u8,
pub strict: bool,
pub security_definer: bool,
pub leakproof: bool,
pub parallel: u8,
pub cost: Option<f64>,
pub rows: Option<f64>,
}
pub const FN_VOLATILE: u8 = b'v';
pub const FN_IMMUTABLE: u8 = b'i';
pub const FN_STABLE: u8 = b's';
pub const FN_PARALLEL_UNSAFE: u8 = b'u';
pub const FN_PARALLEL_RESTRICTED: u8 = b'r';
pub const FN_PARALLEL_SAFE: u8 = b's';
#[must_use]
pub fn resolve_stored_function_key(
functions: &BTreeMap<String, FunctionDef>,
stored: &str,
) -> Option<String> {
if functions.contains_key(stored) {
return Some(stored.to_string());
}
functions
.values()
.find(|f| function_signature_key_legacy(&f.name, &f.args_repr) == stored)
.map(|f| function_signature_key(&f.name, &f.args_repr))
}
pub use spg_sql::parser::is_multiword_type_phrase;
#[must_use]
pub fn function_signature_key_legacy(name: &str, args_repr: &str) -> String {
let inner = args_repr
.trim()
.trim_start_matches('(')
.trim_end_matches(')');
let types: Vec<String> = if inner.trim().is_empty() {
Vec::new()
} else {
inner
.split(',')
.map(|part| {
let mut words: Vec<&str> = part.split_whitespace().collect();
if !words.is_empty()
&& (words[0].eq_ignore_ascii_case("OUT")
|| words[0].eq_ignore_ascii_case("INOUT"))
{
words.remove(0);
}
let ty = if words.len() >= 2 {
words[1..].join(" ")
} else {
words.first().map_or(String::new(), |w| (*w).to_string())
};
normalize_type_name(&ty)
})
.collect()
};
format!("{}({})", name.to_ascii_lowercase(), types.join(","))
}
pub fn function_signature_key(name: &str, args_repr: &str) -> String {
let types = function_arg_types(args_repr);
format!("{}({})", name.to_ascii_lowercase(), types.join(","))
}
#[must_use]
pub fn function_arg_types(args_repr: &str) -> Vec<String> {
let inner = args_repr
.trim()
.trim_start_matches('(')
.trim_end_matches(')');
if inner.trim().is_empty() {
return Vec::new();
}
inner
.split(',')
.map(|part| {
let mut words: Vec<&str> = part.split_whitespace().collect();
if !words.is_empty()
&& (words[0].eq_ignore_ascii_case("OUT") || words[0].eq_ignore_ascii_case("INOUT"))
{
words.remove(0);
}
let whole = words.join(" ");
let ty = if words.len() >= 2 && !is_multiword_type_phrase(&whole) {
words[1..].join(" ")
} else {
whole
};
normalize_type_name(&ty)
})
.collect()
}
#[must_use]
pub fn function_arg_names(args_repr: &str) -> Vec<String> {
let inner = args_repr
.trim()
.trim_start_matches('(')
.trim_end_matches(')');
if inner.trim().is_empty() {
return Vec::new();
}
inner
.split(',')
.map(|part| {
let mut words: Vec<&str> = part.split_whitespace().collect();
if !words.is_empty()
&& (words[0].eq_ignore_ascii_case("OUT") || words[0].eq_ignore_ascii_case("INOUT"))
{
words.remove(0);
}
if words.len() >= 2 {
words[0].to_string()
} else {
String::new()
}
})
.collect()
}
#[must_use]
pub fn normalize_type_name(ty: &str) -> String {
let t = ty.trim().to_ascii_lowercase();
let base = t.split_once('(').map_or(t.as_str(), |(h, _)| h).trim();
match base {
"int" | "int4" | "integer" => "int",
"bigint" | "int8" => "bigint",
"smallint" | "int2" => "smallint",
"text" | "varchar" | "character varying" | "char" | "character" | "bpchar" => "text",
"bool" | "boolean" => "bool",
"float" | "float8" | "double precision" => "float",
"real" | "float4" => "real",
"numeric" | "decimal" => "numeric",
"timestamptz" | "timestamp with time zone" => "timestamptz",
"timestamp" | "timestamp without time zone" => "timestamp",
other => other,
}
.to_string()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TriggerDef {
pub name: String,
pub table: String,
pub timing: String,
pub events: Vec<String>,
pub for_each: String,
pub function: String,
pub update_columns: Vec<String>,
pub enabled: bool,
pub when_condition: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StatisticsExtDef {
pub name: String,
pub table: String,
pub kinds: Vec<String>,
pub columns: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RuleDef {
pub name: String,
pub table: String,
pub event: String,
pub instead: bool,
pub when_condition: String,
pub commands: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SequenceDef {
pub name: String,
pub data_type: SequenceDataType,
pub start: i64,
pub increment: i64,
pub min_value: i64,
pub max_value: i64,
pub cache: i64,
pub cycle: bool,
pub owned_by: Option<(String, String)>,
pub last_value: i64,
pub is_called: bool,
pub owner: Option<String>,
pub acl: Vec<AclItem>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SequenceDataType {
SmallInt,
Int,
BigInt,
}
#[must_use]
pub fn is_builtin_schema(name: &str) -> bool {
name.eq_ignore_ascii_case("public")
|| name.eq_ignore_ascii_case("pg_catalog")
|| name.eq_ignore_ascii_case("information_schema")
}
#[must_use]
pub fn parse_uuid_str(input: &str) -> Option<[u8; 16]> {
let s = input.trim();
let s = if let Some(inner) = s.strip_prefix('{').and_then(|x| x.strip_suffix('}')) {
inner
} else {
s
};
let hex: String = match s.len() {
32 => s.to_ascii_lowercase(),
36 => {
let b = s.as_bytes();
if b[8] != b'-' || b[13] != b'-' || b[18] != b'-' || b[23] != b'-' {
return None;
}
let mut out = String::with_capacity(32);
out.push_str(&s[0..8]);
out.push_str(&s[9..13]);
out.push_str(&s[14..18]);
out.push_str(&s[19..23]);
out.push_str(&s[24..36]);
out.make_ascii_lowercase();
out
}
_ => return None,
};
let bytes = hex.as_bytes();
let mut out = [0u8; 16];
for i in 0..16 {
let hi = hex_nibble(bytes[i * 2])?;
let lo = hex_nibble(bytes[i * 2 + 1])?;
out[i] = (hi << 4) | lo;
}
Some(out)
}
fn hex_nibble(b: u8) -> Option<u8> {
match b {
b'0'..=b'9' => Some(b - b'0'),
b'a'..=b'f' => Some(10 + b - b'a'),
b'A'..=b'F' => Some(10 + b - b'A'),
_ => None,
}
}
#[must_use]
pub fn format_uuid(b: &[u8; 16]) -> String {
const HEX: &[u8; 16] = b"0123456789abcdef";
let mut out = String::with_capacity(36);
for (i, byte) in b.iter().enumerate() {
if matches!(i, 4 | 6 | 8 | 10) {
out.push('-');
}
out.push(HEX[(byte >> 4) as usize] as char);
out.push(HEX[(byte & 0x0f) as usize] as char);
}
out
}
#[derive(Debug, Clone, Default)]
pub struct TxWriteSet {
pub inserted: Vec<(row_header::RowId, Row<'static>)>,
pub tombstoned: Vec<row_header::RowId>,
}
impl TxWriteSet {
#[must_use]
pub fn is_empty(&self) -> bool {
self.inserted.is_empty() && self.tombstoned.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DomainCheck {
pub name: String,
pub expr: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DomainDef {
pub name: String,
pub base_type: DataType,
pub nullable: bool,
pub default: Option<String>,
pub checks: Vec<DomainCheck>,
pub base_domain: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EnumDef {
pub name: String,
pub labels: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompositeDef {
pub name: String,
pub fields: Vec<(String, DataType)>,
pub field_user_types: Vec<Option<String>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ViewDef {
pub name: String,
pub columns: Vec<String>,
pub body: String,
pub check_option: u8,
}
impl SequenceDataType {
pub fn default_bounds(self, increment_positive: bool) -> (i64, i64) {
match self {
Self::SmallInt => {
if increment_positive {
(1, i64::from(i16::MAX))
} else {
(i64::from(i16::MIN), -1)
}
}
Self::Int => {
if increment_positive {
(1, i64::from(i32::MAX))
} else {
(i64::from(i32::MIN), -1)
}
}
Self::BigInt => {
if increment_positive {
(1, i64::MAX)
} else {
(i64::MIN, -1)
}
}
}
}
}
impl Catalog {
pub fn vacuum_all(
&mut self,
oldest_active_snapshot: u64,
dry_run: bool,
) -> vacuum::VacuumReport {
let mut total = vacuum::VacuumReport::default();
let names: Vec<String> = self
.tables
.iter()
.map(|t| t.schema().name.clone())
.collect();
for name in names {
let Some(t) = self.get_mut(&name) else {
continue;
};
let r = t.vacuum(oldest_active_snapshot, dry_run);
if r.rows_reclaimed > 0 {
total.per_table.push((name, r.rows_reclaimed));
}
total.rows_reclaimed += r.rows_reclaimed;
total.rows_examined += r.rows_examined;
}
total
}
pub const fn new() -> Self {
Self {
cold_read_stats: ColdReadStats {
cold_reads: core::sync::atomic::AtomicU64::new(0),
},
tables: Vec::new(),
by_name: BTreeMap::new(),
temp_prefix: None,
dirty_tables: alloc::collections::BTreeSet::new(),
next_rel_id: 0,
cold_segments: Vec::new(),
functions: BTreeMap::new(),
triggers: Vec::new(),
rules: Vec::new(),
statistics_ext: Vec::new(),
large_objects: alloc::collections::BTreeMap::new(),
sequences: BTreeMap::new(),
schema_acl: Vec::new(),
database_acl: Vec::new(),
views: BTreeMap::new(),
materialized_views: BTreeMap::new(),
enum_types: BTreeMap::new(),
domain_types: BTreeMap::new(),
comments: BTreeMap::new(),
db_role_settings: BTreeMap::new(),
replication_slots: BTreeMap::new(),
composite_types: BTreeMap::new(),
schemas: alloc::collections::BTreeSet::new(),
}
}
pub const fn functions(&self) -> &BTreeMap<String, FunctionDef> {
&self.functions
}
pub fn create_function(
&mut self,
def: FunctionDef,
or_replace: bool,
) -> Result<(), StorageError> {
let key = function_signature_key(&def.name, &def.args_repr);
if !or_replace && self.functions.contains_key(&key) {
return Err(StorageError::Corrupt(format!(
"function {:?} already exists (drop or use CREATE OR REPLACE)",
def.name
)));
}
self.functions.insert(key, def);
Ok(())
}
#[must_use]
pub fn functions_named(&self, name: &str) -> Vec<&FunctionDef> {
self.functions
.values()
.filter(|f| f.name.eq_ignore_ascii_case(name))
.collect()
}
#[must_use]
pub fn function_by_key(&self, key: &str) -> Option<&FunctionDef> {
self.functions.get(key)
}
pub fn drop_function_by_key(&mut self, key: &str) -> bool {
self.functions.remove(key).is_some()
}
pub fn drop_function(&mut self, name: &str) -> bool {
let keys: Vec<String> = self
.functions
.iter()
.filter(|(_, f)| f.name.eq_ignore_ascii_case(name))
.map(|(k, _)| k.clone())
.collect();
let hit = !keys.is_empty();
for k in keys {
self.functions.remove(&k);
}
hit
}
#[must_use]
pub fn schema_acl(&self) -> &[AclItem] {
&self.schema_acl
}
pub fn schema_acl_mut(&mut self) -> &mut Vec<AclItem> {
&mut self.schema_acl
}
#[must_use]
pub fn database_acl(&self) -> &[AclItem] {
&self.database_acl
}
pub fn database_acl_mut(&mut self) -> &mut Vec<AclItem> {
&mut self.database_acl
}
pub fn sequence_mut(&mut self, name: &str) -> Option<&mut SequenceDef> {
let key = self.sequence_key(name);
self.sequences.get_mut(&key)
}
pub fn function_mut(&mut self, name: &str) -> Option<&mut FunctionDef> {
self.functions.get_mut(name)
}
pub const fn sequences_all(&self) -> &BTreeMap<String, SequenceDef> {
&self.sequences
}
#[must_use]
pub fn sequence(&self, name: &str) -> Option<&SequenceDef> {
if let Some(mangled) = self.temp_name_for(name)
&& let Some(def) = self.sequences.get(&mangled)
{
return Some(def);
}
self.sequences.get(name)
}
#[must_use]
pub fn has_sequence(&self, name: &str) -> bool {
self.sequence(name).is_some()
}
#[must_use]
pub fn sequence_key(&self, name: &str) -> String {
if let Some(mangled) = self.temp_name_for(name)
&& self.sequences.contains_key(&mangled)
{
return mangled;
}
name.into()
}
pub fn create_sequence(
&mut self,
def: SequenceDef,
if_not_exists: bool,
) -> Result<(), StorageError> {
if self.sequences.contains_key(&def.name) {
if if_not_exists {
return Ok(());
}
return Err(StorageError::Corrupt(format!(
"relation {:?} already exists",
def.name
)));
}
self.sequences.insert(def.name.clone(), def);
Ok(())
}
pub fn rename_sequence(&mut self, old: &str, new: &str) -> Result<(), StorageError> {
if !self.sequences.contains_key(old) {
return Err(StorageError::Corrupt(format!(
"relation {old:?} does not exist"
)));
}
if self.sequences.contains_key(new) {
return Err(StorageError::Corrupt(format!(
"relation {new:?} already exists"
)));
}
if let Some(mut def) = self.sequences.remove(old) {
def.name = new.to_string();
self.sequences.insert(new.to_string(), def);
}
Ok(())
}
pub fn drop_sequence(&mut self, name: &str) -> bool {
self.sequences.remove(name).is_some()
}
#[must_use]
pub fn sequence_counters(&self) -> Vec<(String, i64, bool)> {
self.sequences
.iter()
.map(|(k, d)| (k.clone(), d.last_value, d.is_called))
.collect()
}
pub fn restore_sequence_counters(&mut self, saved: &[(String, i64, bool)]) {
for (k, last, called) in saved {
if let Some(d) = self.sequences.get_mut(k) {
d.last_value = *last;
d.is_called = *called;
}
}
}
pub fn sequence_next_value(&mut self, name: &str) -> Result<i64, StorageError> {
let key = self.sequence_key(name);
let Some(seq) = self.sequences.get_mut(&key) else {
return Err(StorageError::TableNotFound { name: name.into() });
};
let candidate = if seq.is_called {
let next = seq.last_value.checked_add(seq.increment).ok_or_else(|| {
StorageError::Corrupt(format!("sequence {name:?} arithmetic overflow"))
})?;
if seq.increment > 0 {
if next > seq.max_value {
if seq.cycle {
seq.min_value
} else {
return Err(StorageError::SequenceExhausted {
name: name.into(),
limit: seq.max_value,
is_max: true,
});
}
} else {
next
}
} else if next < seq.min_value {
if seq.cycle {
seq.max_value
} else {
return Err(StorageError::SequenceExhausted {
name: name.into(),
limit: seq.min_value,
is_max: false,
});
}
} else {
next
}
} else {
seq.last_value
};
seq.last_value = candidate;
seq.is_called = true;
Ok(candidate)
}
pub fn sequence_current_value(&self, name: &str) -> Result<i64, StorageError> {
let Some(seq) = self.sequences.get(name) else {
return Err(StorageError::TableNotFound { name: name.into() });
};
if !seq.is_called {
return Err(StorageError::Corrupt(format!(
"currval of sequence {name:?} is not yet defined in this session"
)));
}
Ok(seq.last_value)
}
pub fn sequence_set_value(
&mut self,
name: &str,
value: i64,
is_called: bool,
) -> Result<i64, StorageError> {
let key = self.sequence_key(name);
let Some(seq) = self.sequences.get_mut(&key) else {
return Err(StorageError::TableNotFound { name: name.into() });
};
if value < seq.min_value || value > seq.max_value {
return Err(StorageError::Unsupported(format!(
"setval: value {value} is out of bounds for sequence \"{name}\" ({}..{})",
seq.min_value, seq.max_value
)));
}
seq.last_value = value;
seq.is_called = is_called;
Ok(value)
}
pub const fn views_all(&self) -> &BTreeMap<String, ViewDef> {
&self.views
}
#[must_use]
pub fn view(&self, name: &str) -> Option<&ViewDef> {
if let Some(mangled) = self.temp_name_for(name)
&& let Some(def) = self.views.get(&mangled)
{
return Some(def);
}
self.views.get(name)
}
#[must_use]
pub fn has_view(&self, name: &str) -> bool {
self.view(name).is_some()
}
#[must_use]
pub fn view_key(&self, name: &str) -> String {
if let Some(mangled) = self.temp_name_for(name)
&& self.views.contains_key(&mangled)
{
return mangled;
}
name.into()
}
pub fn create_view(
&mut self,
def: ViewDef,
or_replace: bool,
if_not_exists: bool,
) -> Result<(), StorageError> {
if self.views.contains_key(&def.name) {
if or_replace {
self.views.insert(def.name.clone(), def);
return Ok(());
}
if if_not_exists {
return Ok(());
}
return Err(StorageError::Corrupt(format!(
"relation {:?} already exists",
def.name
)));
}
if self.by_name.contains_key(&def.name) {
return Err(StorageError::Corrupt(format!(
"view {:?} would shadow an existing table",
def.name
)));
}
if self.sequences.contains_key(&def.name) {
return Err(StorageError::Corrupt(format!(
"view {:?} would shadow an existing sequence",
def.name
)));
}
self.views.insert(def.name.clone(), def);
Ok(())
}
pub fn drop_view(&mut self, name: &str) -> bool {
self.views.remove(name).is_some()
}
pub const fn materialized_views(&self) -> &BTreeMap<String, String> {
&self.materialized_views
}
pub fn register_materialized_view(&mut self, name: String, body: String) {
self.materialized_views.insert(name, body);
}
pub fn drop_materialized_view_source(&mut self, name: &str) -> bool {
self.materialized_views.remove(name).is_some()
}
pub const fn enum_types(&self) -> &BTreeMap<String, EnumDef> {
&self.enum_types
}
pub fn create_enum_type(&mut self, def: EnumDef) -> Result<(), StorageError> {
if self.enum_types.contains_key(&def.name) {
return Err(StorageError::Corrupt(format!(
"type {:?} already exists",
def.name
)));
}
self.enum_types.insert(def.name.clone(), def);
Ok(())
}
pub fn rename_enum_value(
&mut self,
type_name: &str,
old: &str,
new: &str,
) -> Result<(), StorageError> {
let def = self
.enum_types
.get_mut(type_name)
.ok_or_else(|| StorageError::Corrupt(format!("type {type_name:?} does not exist")))?;
if def.labels.iter().any(|l| l == new) {
return Err(StorageError::Corrupt(format!(
"enum label {new:?} already exists"
)));
}
let at = def.labels.iter().position(|l| l == old).ok_or_else(|| {
StorageError::Corrupt(format!("{old:?} is not an existing enum label"))
})?;
def.labels[at] = new.to_string();
Ok(())
}
pub fn set_comment(&mut self, key: &str, text: Option<&str>) {
match text {
Some(t) => {
self.comments.insert(key.to_string(), t.to_string());
}
None => {
self.comments.remove(key);
}
}
}
#[must_use]
pub fn comment(&self, key: &str) -> Option<&str> {
self.comments.get(key).map(String::as_str)
}
pub fn set_db_role_setting(
&mut self,
database: &str,
role: &str,
param: &str,
value: Option<&str>,
) {
let key = (database.to_string(), role.to_string());
match value {
Some(v) => {
self.db_role_settings
.entry(key)
.or_default()
.insert(param.to_ascii_lowercase(), v.to_string());
}
None => {
if let Some(m) = self.db_role_settings.get_mut(&key) {
m.remove(¶m.to_ascii_lowercase());
if m.is_empty() {
self.db_role_settings.remove(&key);
}
}
}
}
}
pub fn create_replication_slot(
&mut self,
name: &str,
plugin: &str,
slot_type: &str,
) -> Result<(), String> {
if self.replication_slots.contains_key(name) {
return Err(alloc::format!("replication slot \"{name}\" already exists"));
}
self.replication_slots.insert(
name.to_string(),
(plugin.to_string(), slot_type.to_string()),
);
Ok(())
}
pub fn drop_replication_slot(&mut self, name: &str) -> Result<(), String> {
if self.replication_slots.remove(name).is_none() {
return Err(alloc::format!("replication slot \"{name}\" does not exist"));
}
Ok(())
}
#[must_use]
pub const fn replication_slots(&self) -> &BTreeMap<String, (String, String)> {
&self.replication_slots
}
pub fn reset_db_role_settings(&mut self, database: &str, role: &str) {
self.db_role_settings
.remove(&(database.to_string(), role.to_string()));
}
#[must_use]
pub const fn db_role_settings(&self) -> &BTreeMap<(String, String), BTreeMap<String, String>> {
&self.db_role_settings
}
#[must_use]
pub const fn comments(&self) -> &BTreeMap<String, String> {
&self.comments
}
pub fn drop_comments_for(&mut self, kind: &str, name: &str) {
let exact = alloc::format!("{kind}:{name}");
let col_prefix = alloc::format!("column:{name}.");
self.comments
.retain(|k, _| *k != exact && !k.starts_with(&col_prefix));
}
pub fn add_enum_value(
&mut self,
type_name: &str,
label: &str,
if_not_exists: bool,
position: Option<(bool, String)>,
) -> Result<bool, StorageError> {
let def = self
.enum_types
.get_mut(type_name)
.ok_or_else(|| StorageError::Corrupt(format!("type {type_name:?} does not exist")))?;
if def.labels.iter().any(|l| l == label) {
if if_not_exists {
return Ok(false);
}
return Err(StorageError::Corrupt(format!(
"enum label {label:?} already exists"
)));
}
match position {
None => def.labels.push(label.to_string()),
Some((is_before, anchor)) => {
let at = def
.labels
.iter()
.position(|l| l == &anchor)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"enum label {anchor:?} does not exist in type {type_name:?}"
))
})?;
let idx = if is_before { at } else { at + 1 };
def.labels.insert(idx, label.to_string());
}
}
Ok(true)
}
pub fn drop_enum_type(&mut self, name: &str) -> bool {
self.enum_types.remove(name).is_some()
}
pub const fn domain_types(&self) -> &BTreeMap<String, DomainDef> {
&self.domain_types
}
pub fn create_domain_type(&mut self, def: DomainDef) -> Result<(), StorageError> {
if self.domain_types.contains_key(&def.name) {
return Err(StorageError::Corrupt(format!(
"domain {:?} already exists",
def.name
)));
}
self.domain_types.insert(def.name.clone(), def);
Ok(())
}
pub fn drop_domain_type(&mut self, name: &str) -> bool {
self.domain_types.remove(name).is_some()
}
pub const fn composite_types(&self) -> &BTreeMap<String, CompositeDef> {
&self.composite_types
}
pub fn create_composite_type(&mut self, def: CompositeDef) -> Result<(), StorageError> {
if self.composite_types.contains_key(&def.name) {
return Err(StorageError::Corrupt(format!(
"type {:?} already exists",
def.name
)));
}
self.composite_types.insert(def.name.clone(), def);
Ok(())
}
pub fn drop_composite_type(&mut self, name: &str) -> bool {
self.composite_types.remove(name).is_some()
}
pub const fn user_schemas(&self) -> &alloc::collections::BTreeSet<String> {
&self.schemas
}
pub fn schema_exists(&self, name: &str) -> bool {
is_builtin_schema(name) || self.schemas.contains(name)
}
pub fn create_schema(&mut self, name: String, if_not_exists: bool) -> Result<(), StorageError> {
if is_builtin_schema(&name) {
if if_not_exists {
return Ok(());
}
return Err(StorageError::Corrupt(format!(
"schema {name:?} is built-in and cannot be redeclared"
)));
}
if self.schemas.contains(&name) {
if if_not_exists {
return Ok(());
}
return Err(StorageError::Corrupt(format!(
"schema {name:?} already exists"
)));
}
self.schemas.insert(name);
Ok(())
}
pub fn drop_schema(&mut self, name: &str) -> Result<bool, StorageError> {
if is_builtin_schema(name) {
return Err(StorageError::Corrupt(format!(
"schema {name:?} is built-in and cannot be dropped"
)));
}
Ok(self.schemas.remove(name))
}
#[allow(clippy::too_many_arguments)]
pub fn alter_sequence(
&mut self,
name: &str,
increment: Option<i64>,
min_value: Option<i64>,
max_value: Option<i64>,
start: Option<i64>,
restart: Option<Option<i64>>,
cache: Option<i64>,
cycle: Option<bool>,
owned_by: Option<Option<(String, String)>>,
) -> Result<(), StorageError> {
let Some(seq) = self.sequences.get_mut(name) else {
return Err(StorageError::TableNotFound { name: name.into() });
};
if let Some(v) = increment {
seq.increment = v;
}
if let Some(v) = min_value {
seq.min_value = v;
}
if let Some(v) = max_value {
seq.max_value = v;
}
if let Some(v) = start {
seq.start = v;
}
if let Some(restart_value) = restart {
seq.last_value = restart_value.unwrap_or(seq.start);
seq.is_called = false;
}
if let Some(v) = cache {
seq.cache = v;
}
if let Some(v) = cycle {
seq.cycle = v;
}
if let Some(v) = owned_by {
seq.owned_by = v;
}
Ok(())
}
pub fn triggers(&self) -> &[TriggerDef] {
&self.triggers
}
pub fn triggers_mut(&mut self) -> &mut Vec<TriggerDef> {
&mut self.triggers
}
pub fn create_trigger(
&mut self,
def: TriggerDef,
or_replace: bool,
) -> Result<(), StorageError> {
if !self.by_name.contains_key(&def.table) && !self.views.contains_key(&def.table) {
return Err(StorageError::TableNotFound {
name: def.table.clone(),
});
}
if self.functions_named(&def.function).is_empty() {
return Err(StorageError::Corrupt(format!(
"function {}() does not exist",
def.function
)));
}
let dup = self
.triggers
.iter()
.position(|t| t.name == def.name && t.table == def.table);
match (dup, or_replace) {
(Some(_), false) => Err(StorageError::Corrupt(format!(
"trigger {:?} already exists on table {:?}",
def.name, def.table
))),
(Some(i), true) => {
self.triggers[i] = def;
Ok(())
}
(None, _) => {
self.triggers.push(def);
Ok(())
}
}
}
pub fn drop_trigger(&mut self, name: &str, table: &str) -> bool {
let before = self.triggers.len();
self.triggers
.retain(|t| !(t.name == name && t.table == table));
before != self.triggers.len()
}
pub fn rules(&self) -> &[RuleDef] {
&self.rules
}
#[must_use]
pub fn statistics_ext(&self) -> &[StatisticsExtDef] {
&self.statistics_ext
}
#[must_use]
pub fn large_objects(&self) -> &alloc::collections::BTreeMap<u32, Vec<u8>> {
&self.large_objects
}
#[must_use]
pub fn large_object(&self, oid: u32) -> Option<&[u8]> {
self.large_objects.get(&oid).map(Vec::as_slice)
}
pub fn create_large_object(&mut self, oid: u32, bytes: Vec<u8>) -> Result<u32, String> {
let id = if oid == 0 {
self.next_large_object_oid()
} else {
oid
};
if self.large_objects.contains_key(&id) {
return Err(format!("large object {id} already exists"));
}
self.large_objects.insert(id, bytes);
Ok(id)
}
pub fn put_large_object(&mut self, oid: u32, offset: usize, data: &[u8]) -> Result<(), String> {
let Some(buf) = self.large_objects.get_mut(&oid) else {
return Err(format!("large object {oid} does not exist"));
};
let end = offset.saturating_add(data.len());
if buf.len() < end {
buf.resize(end, 0);
}
buf[offset..end].copy_from_slice(data);
Ok(())
}
pub fn truncate_large_object(&mut self, oid: u32, len: usize) -> Result<(), String> {
let Some(buf) = self.large_objects.get_mut(&oid) else {
return Err(format!("large object {oid} does not exist"));
};
buf.resize(len, 0);
Ok(())
}
pub fn unlink_large_object(&mut self, oid: u32) -> bool {
self.large_objects.remove(&oid).is_some()
}
fn next_large_object_oid(&self) -> u32 {
self.large_objects
.keys()
.next_back()
.map_or(500_000, |m| m.saturating_add(1))
}
pub fn create_statistics_ext(&mut self, def: StatisticsExtDef) -> Result<(), String> {
if self.statistics_ext.iter().any(|s| s.name == def.name) {
return Err(def.name);
}
self.statistics_ext.push(def);
Ok(())
}
pub fn drop_statistics_ext(&mut self, name: &str) -> bool {
let before = self.statistics_ext.len();
self.statistics_ext.retain(|s| s.name != name);
before != self.statistics_ext.len()
}
pub fn create_rule(&mut self, def: RuleDef, or_replace: bool) -> Result<(), StorageError> {
if !self.by_name.contains_key(&def.table) && !self.views.contains_key(&def.table) {
return Err(StorageError::TableNotFound {
name: def.table.clone(),
});
}
let dup = self
.rules
.iter()
.position(|r| r.name == def.name && r.table == def.table);
match (dup, or_replace) {
(Some(_), false) => Err(StorageError::Corrupt(format!(
"rule {:?} for relation {:?} already exists",
def.name, def.table
))),
(Some(i), true) => {
self.rules[i] = def;
Ok(())
}
(None, _) => {
self.rules.push(def);
Ok(())
}
}
}
pub fn drop_rule(&mut self, name: &str, table: &str) -> bool {
let before = self.rules.len();
self.rules.retain(|r| !(r.name == name && r.table == table));
before != self.rules.len()
}
pub fn create_table(&mut self, schema: TableSchema) -> Result<(), StorageError> {
if self.by_name.contains_key(&schema.name) {
return Err(StorageError::DuplicateTable {
name: schema.name.clone(),
});
}
let idx = self.tables.len();
let name = schema.name.clone();
self.tables.push(Table::new(schema));
self.by_name.insert(name.clone(), idx);
self.dirty_tables.insert(name);
self.next_rel_id += 1;
let rid = row_header::RelId(self.next_rel_id);
self.tables[idx].set_rel_id(rid);
Ok(())
}
fn resolve_index(&self, name: &str) -> Option<usize> {
if let Some(prefix) = &self.temp_prefix {
let mut mangled = String::with_capacity(prefix.len() + name.len());
mangled.push_str(prefix);
mangled.push_str(name);
if let Some(idx) = self.by_name.get(&mangled) {
return Some(*idx);
}
}
self.by_name.get(name).copied()
}
pub fn set_temp_prefix(&mut self, prefix: Option<String>) {
self.temp_prefix = prefix;
}
#[must_use]
pub fn temp_name_for(&self, name: &str) -> Option<String> {
self.temp_prefix
.as_ref()
.map(|p| alloc::format!("{p}{name}"))
}
pub fn get(&self, name: &str) -> Option<&Table> {
let idx = self.resolve_index(name)?;
self.tables.get(idx)
}
pub fn get_mut(&mut self, name: &str) -> Option<&mut Table> {
let idx = self.resolve_index(name)?;
let recorded = self.tables.get(idx).map(|t| t.schema().name.clone());
if let Some(n) = recorded {
self.dirty_tables.insert(n);
}
self.tables.get_mut(idx)
}
#[must_use]
pub fn dirty_tables(&self) -> &alloc::collections::BTreeSet<String> {
&self.dirty_tables
}
pub fn clear_dirty_tables(&mut self) {
self.dirty_tables.clear();
}
pub fn install_table(&mut self, name: &str, table: Table) {
match self.by_name.get(name).copied() {
Some(idx) => self.tables[idx] = table,
None => {
let idx = self.tables.len();
self.tables.push(table);
self.by_name.insert(name.into(), idx);
}
}
self.dirty_tables.insert(name.into());
}
pub fn tables_position_of(&self, name: &str) -> Option<usize> {
self.resolve_index(name)
}
pub fn tables_at(&self, idx: usize) -> Option<&Table> {
self.tables.get(idx)
}
pub fn apply_redo(&mut self, changes: &[RowChange]) -> Result<(), StorageError> {
let mut runs: alloc::vec::Vec<(String, alloc::vec::Vec<&RowChange>)> =
alloc::vec::Vec::new();
for change in changes {
if let RowChange::Tombstone { xmax, .. } = change {
row_header::observe_persisted_version(*xmax);
}
let table = match change {
RowChange::Insert { table, .. }
| RowChange::Update { table, .. }
| RowChange::Delete { table, .. }
| RowChange::Tombstone { table, .. } => table.clone(),
};
if runs.last().map(|(t, _)| t.as_str()) != Some(table.as_str()) {
runs.push((table, alloc::vec::Vec::new()));
}
runs.last_mut().unwrap().1.push(change);
}
for (table_name, run) in runs {
self.apply_redo_run_on_table(&table_name, &run)?;
}
Ok(())
}
fn apply_redo_run_on_table(
&mut self,
table_name: &str,
run: &[&RowChange],
) -> Result<(), StorageError> {
let table = self.get_mut(table_name).ok_or_else(|| {
StorageError::Corrupt(alloc::format!("redo: unknown table {table_name:?}"))
})?;
let original_rows: alloc::vec::Vec<Row<'static>> = table.rows().iter().cloned().collect();
let mut live: alloc::vec::Vec<bool> = alloc::vec![true; original_rows.len()];
let mut tail: alloc::vec::Vec<Row<'static>> = alloc::vec::Vec::new();
let mut overlay: alloc::collections::BTreeMap<usize, alloc::vec::Vec<Value<'static>>> =
alloc::collections::BTreeMap::new();
let has_tomb = run.iter().any(|c| matches!(c, RowChange::Tombstone { .. }));
let orig_rowids: alloc::vec::Vec<row_header::RowId> =
table.rowids().iter().copied().collect();
let orig_headers: alloc::vec::Vec<row_header::RowHeader> =
table.headers().iter().copied().collect();
let mut tail_rowids: alloc::vec::Vec<row_header::RowId> = alloc::vec::Vec::new();
let mut tomb_targets: alloc::vec::Vec<(row_header::RowId, u64)> = alloc::vec::Vec::new();
fn translate(live: &[bool], tail_len: usize, current_pos: usize) -> Option<usize> {
let mut seen = 0usize;
for (i, &alive) in live.iter().enumerate() {
if alive {
if seen == current_pos {
return Some(i);
}
seen += 1;
}
}
let off = current_pos - seen;
if off < tail_len {
Some(live.len() + off)
} else {
None
}
}
for change in run {
match *change {
RowChange::Insert { row, rowid, .. } => {
if row.len() != table.schema().columns.len() {
return Err(StorageError::ArityMismatch {
expected: table.schema().columns.len(),
actual: row.len(),
});
}
tail.push(row.clone());
tail_rowids.push(*rowid);
}
RowChange::Update { pos, new_row, .. } => {
if new_row.len() != table.schema().columns.len() {
return Err(StorageError::ArityMismatch {
expected: table.schema().columns.len(),
actual: new_row.len(),
});
}
let abs = translate(&live, tail.len(), *pos).ok_or_else(|| {
StorageError::Corrupt(alloc::format!(
"redo: update_row position {pos} out of bounds in table {table_name:?}",
))
})?;
if abs < live.len() {
overlay.insert(abs, new_row.clone());
} else {
tail[abs - live.len()] = Row::new(new_row.clone());
}
}
RowChange::Delete { positions, .. } => {
let mut sorted: alloc::vec::Vec<usize> = positions.clone();
sorted.sort_unstable();
sorted.dedup();
let mut to_flip_live: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
let mut to_flip_tail: alloc::vec::Vec<usize> = alloc::vec::Vec::new();
let mut seen = 0usize;
let mut sp = sorted.iter().peekable();
for (i, &alive) in live.iter().enumerate() {
if !alive {
continue;
}
while let Some(&&p) = sp.peek() {
if seen == p {
to_flip_live.push(i);
sp.next();
} else {
break;
}
}
if sp.peek().is_none() {
break;
}
seen += 1;
}
for &p in sp {
let off = p - seen;
if off < tail.len() {
to_flip_tail.push(off);
}
}
for i in to_flip_live {
live[i] = false;
overlay.remove(&i);
}
to_flip_tail.sort_unstable();
to_flip_tail.dedup();
for off in to_flip_tail.into_iter().rev() {
tail.remove(off);
{
tail_rowids.remove(off);
}
}
}
RowChange::Tombstone { rowids, xmax, .. } => {
for rid in rowids {
tomb_targets.push((*rid, *xmax));
}
}
}
}
let mut new_rows: PersistentVec<Row> = PersistentVec::new();
let mut new_hot_bytes: u64 = 0;
let schema_snapshot = table.schema().clone();
let mut final_rowids: alloc::vec::Vec<row_header::RowId> = alloc::vec::Vec::new();
let mut final_headers: alloc::vec::Vec<row_header::RowHeader> = alloc::vec::Vec::new();
for (i, row) in original_rows.into_iter().enumerate() {
if !live[i] {
continue;
}
let final_row = if let Some(new_values) = overlay.remove(&i) {
Row::new(new_values)
} else {
row
};
new_hot_bytes = new_hot_bytes
.saturating_add(row_body_encoded_len(&final_row, &schema_snapshot) as u64);
new_rows.push_mut(final_row);
final_rowids.push(
orig_rowids
.get(i)
.copied()
.unwrap_or(row_header::RowId::UNASSIGNED),
);
final_headers.push(
orig_headers
.get(i)
.copied()
.unwrap_or_else(row_header::RowHeader::frozen),
);
}
for (off, row) in tail.into_iter().enumerate() {
new_hot_bytes =
new_hot_bytes.saturating_add(row_body_encoded_len(&row, &schema_snapshot) as u64);
new_rows.push_mut(row);
final_rowids.push(
tail_rowids
.get(off)
.copied()
.unwrap_or(row_header::RowId::UNASSIGNED),
);
final_headers.push(row_header::RowHeader::frozen());
}
table.set_rows_and_rebuild_indices_with_rowids(
new_rows,
new_hot_bytes,
&final_rowids,
&final_headers,
);
if has_tomb && !tomb_targets.is_empty() {
let mut id_to_slot: alloc::collections::BTreeMap<row_header::RowId, usize> =
alloc::collections::BTreeMap::new();
for (slot, rid) in final_rowids.iter().enumerate() {
if *rid != row_header::RowId::UNASSIGNED {
id_to_slot.insert(*rid, slot);
}
}
let table = self.get_mut(table_name).ok_or_else(|| {
StorageError::Corrupt(alloc::format!("redo: unknown table {table_name:?}"))
})?;
for (rid, xmax) in &tomb_targets {
match id_to_slot.get(rid) {
Some(&slot) => {
let _ = table.mark_row_deleted(slot, *xmax);
}
None => {
UNRESOLVED_TOMBSTONES.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
}
}
}
}
Ok(())
}
fn table_for_redo(&mut self, name: &str) -> Result<&mut Table, StorageError> {
self.get_mut(name)
.ok_or_else(|| StorageError::Corrupt(alloc::format!("redo: unknown table {name:?}")))
}
pub fn enable_redo_all(&mut self) {
for t in &mut self.tables {
t.enable_redo();
}
}
pub fn drain_redo(&mut self) -> Vec<RowChange> {
let mut all = Vec::new();
for t in &mut self.tables {
all.extend(t.take_redo());
}
all
}
pub fn table_count(&self) -> usize {
self.tables.len()
}
pub fn drop_table(&mut self, name: &str) -> bool {
let key = match self.temp_prefix.as_ref() {
Some(p) => {
let mangled = alloc::format!("{p}{name}");
if self.by_name.contains_key(&mangled) {
mangled
} else {
name.into()
}
}
None => name.into(),
};
let Some(idx) = self.by_name.remove(&key) else {
return false;
};
self.dirty_tables.insert(key.clone());
self.tables.swap_remove(idx);
if idx < self.tables.len() {
let moved_name = self.tables[idx].schema.name.clone();
self.by_name.insert(moved_name, idx);
}
true
}
pub fn rename_table(&mut self, old: &str, new: &str) -> Result<(), StorageError> {
if old == new {
return Ok(());
}
if self.by_name.contains_key(new) {
return Err(StorageError::Corrupt(format!(
"rename_table: target name {new:?} already exists"
)));
}
let idx = self
.by_name
.remove(old)
.ok_or_else(|| StorageError::TableNotFound { name: old.into() })?;
self.tables[idx].schema.name = new.to_string();
self.by_name.insert(new.to_string(), idx);
for t in &mut self.tables {
for fk in &mut t.schema.foreign_keys {
if fk.parent_table == old {
fk.parent_table = new.to_string();
}
}
}
for trig in &mut self.triggers {
if trig.table == old {
trig.table = new.to_string();
}
}
Ok(())
}
pub fn rename_index(&mut self, old: &str, new: &str) -> Result<(), StorageError> {
if old == new {
return Ok(());
}
for t in &self.tables {
if t.indices.iter().any(|i| i.name == new) {
return Err(StorageError::Corrupt(format!(
"rename_index: target name {new:?} already exists"
)));
}
}
for t in &mut self.tables {
for i in &mut t.indices {
if i.name == old {
i.name = new.to_string();
return Ok(());
}
}
}
Err(StorageError::IndexNotFound { name: old.into() })
}
pub fn drop_named_index(&mut self, name: &str) -> bool {
for t in &mut self.tables {
let before = t.indices.len();
t.indices.retain(|i| i.name != name);
if t.indices.len() != before {
return true;
}
}
false
}
pub fn table_names(&self) -> Vec<String> {
self.tables.iter().map(|t| t.schema.name.clone()).collect()
}
pub const TEMP_NAME_MARKER: &'static str = "__spg_temp_";
#[must_use]
pub fn listed_name<'a>(&self, stored: &'a str) -> Option<&'a str> {
if !stored.starts_with(Self::TEMP_NAME_MARKER) {
return Some(stored);
}
let prefix = self.temp_prefix.as_ref()?;
stored.strip_prefix(prefix.as_str())
}
#[must_use]
pub fn visible_table_names(&self) -> Vec<String> {
self.tables
.iter()
.filter_map(|t| self.listed_name(&t.schema.name).map(String::from))
.collect()
}
pub fn load_segment_bytes(&mut self, bytes: Vec<u8>) -> Result<u32, StorageError> {
let id = u32::try_from(self.cold_segments.len()).map_err(|_| {
StorageError::Corrupt("cold segment count would exceed u32::MAX".into())
})?;
let seg = OwnedSegment::from_bytes(bytes)
.map_err(|e| StorageError::Corrupt(format!("cold segment parse failed: {e}")))?;
self.cold_segments.push(Some(Arc::new(seg)));
Ok(id)
}
pub fn load_segment_bytes_at(
&mut self,
target_id: u32,
bytes: Vec<u8>,
) -> Result<(), StorageError> {
let seg = OwnedSegment::from_bytes(bytes)
.map_err(|e| StorageError::Corrupt(format!("cold segment parse failed: {e}")))?;
let idx = target_id as usize;
while self.cold_segments.len() <= idx {
self.cold_segments.push(None);
}
if self.cold_segments[idx].is_some() {
return Err(StorageError::Corrupt(format!(
"load_segment_bytes_at: segment_id {target_id} already occupied"
)));
}
self.cold_segments[idx] = Some(Arc::new(seg));
Ok(())
}
pub fn tombstone_segment(&mut self, segment_id: u32) -> Result<(), StorageError> {
let idx = segment_id as usize;
if idx >= self.cold_segments.len() {
return Err(StorageError::Corrupt(format!(
"tombstone_segment: segment_id {segment_id} out of bounds (len={})",
self.cold_segments.len()
)));
}
self.cold_segments[idx] = None;
Ok(())
}
#[must_use]
pub fn cold_segment_count(&self) -> usize {
self.cold_segments.iter().filter(|s| s.is_some()).count()
}
#[must_use]
pub fn has_any_cold_segments(&self) -> bool {
self.cold_segments.iter().any(Option::is_some)
}
#[must_use]
pub fn cold_segment_slot_count(&self) -> usize {
self.cold_segments.len()
}
#[must_use]
pub fn cold_segment_ids_global(&self) -> Vec<u32> {
self.cold_segments
.iter()
.enumerate()
.filter_map(|(i, s)| s.as_ref().map(|_| i as u32))
.collect()
}
#[must_use]
pub fn hot_tier_bytes(&self) -> u64 {
self.tables
.iter()
.map(Table::hot_bytes)
.fold(0u64, u64::saturating_add)
}
pub fn freeze_oldest_to_cold(
&mut self,
table_name: &str,
index_name: &str,
max_rows: usize,
) -> Result<FreezeReport, StorageError> {
if max_rows == 0 {
return Err(StorageError::Corrupt(
"freeze_oldest_to_cold: max_rows must be > 0".into(),
));
}
let table = self.get(table_name).ok_or_else(|| {
StorageError::Corrupt(format!(
"freeze_oldest_to_cold: table {table_name:?} not found"
))
})?;
if max_rows > table.rows.len() {
return Err(StorageError::Corrupt(format!(
"freeze_oldest_to_cold: max_rows {max_rows} > row_count {}",
table.rows.len()
)));
}
let idx = table
.indices
.iter()
.find(|i| i.name == index_name)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"freeze_oldest_to_cold: index {index_name:?} not found on {table_name:?}"
))
})?;
if !matches!(idx.kind, IndexKind::BTree(_)) {
return Err(StorageError::Corrupt(format!(
"freeze_oldest_to_cold: index {index_name:?} is NSW; only BTree indices may freeze"
)));
}
let column_position = idx.column_position;
let schema = table.schema.clone();
let mut to_freeze: Vec<(u64, Vec<u8>, IndexKey)> = Vec::with_capacity(max_rows);
for row_idx in 0..max_rows {
let row = table.rows.get(row_idx).expect("bounds-checked above");
let key = IndexKey::from_value(&row.values[column_position]).ok_or_else(|| {
StorageError::Corrupt(format!(
"freeze_oldest_to_cold: row {row_idx} has NULL / non-key value in index column"
))
})?;
let pk_u64 = index_key_as_u64(&key).ok_or_else(|| {
StorageError::Corrupt(format!(
"freeze_oldest_to_cold: index {index_name:?} column type is non-integer; \
v5.2.2 cold tier requires IndexKey::Int (Text PK lands in v5.5+)"
))
})?;
to_freeze.push((pk_u64, encode_row_body_dense(row, &schema), key));
}
to_freeze.sort_by_key(|(k, _, _)| *k);
for w in to_freeze.windows(2) {
if w[0].0 == w[1].0 {
return Err(StorageError::Corrupt(format!(
"freeze_oldest_to_cold: duplicate PK {} in freeze batch",
w[0].0
)));
}
}
let post_swap_keys: Vec<IndexKey> = to_freeze.iter().map(|(_, _, k)| k.clone()).collect();
let seg_rows: Vec<(u64, Vec<u8>)> = to_freeze
.into_iter()
.map(|(k, body, _)| (k, body))
.collect();
let frozen_rows = seg_rows.len();
let (seg_bytes, _meta) = encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES)
.map_err(|e| StorageError::Corrupt(format!("freeze_oldest_to_cold: encode: {e}")))?;
let bytes_before = self.get(table_name).expect("just validated").hot_bytes();
let positions: Vec<usize> = (0..max_rows).collect();
let t_mut = self
.get_mut(table_name)
.expect("just validated; still present");
let removed = t_mut.delete_rows(&positions);
debug_assert_eq!(removed, max_rows, "delete_rows count matches request");
let bytes_after = t_mut.hot_bytes();
let bytes_freed = bytes_before.saturating_sub(bytes_after);
let segment_id = self
.load_segment_bytes(seg_bytes.clone())
.map_err(|e| StorageError::Corrupt(format!("freeze_oldest_to_cold: load: {e}")))?;
let new_cold = post_swap_keys.into_iter().map(|k| {
(
k,
RowLocator::Cold {
segment_id,
page_offset: 0,
},
)
});
let t_mut = self.get_mut(table_name).expect("still present");
t_mut.register_cold_locators(index_name, new_cold)?;
t_mut.mark_cold_row_count_stale();
Ok(FreezeReport {
segment_id,
frozen_rows,
bytes_freed,
segment_bytes: seg_bytes,
})
}
#[must_use]
pub fn cold_segment(&self, segment_id: u32) -> Option<&OwnedSegment> {
self.cold_segments
.get(segment_id as usize)
.and_then(|s| s.as_deref())
}
pub fn resolve_cold_locator(
&self,
table_name: &str,
segment_id: u32,
key: &IndexKey,
) -> Option<Row<'static>> {
let t = self.get(table_name)?;
let u64_key = index_key_as_u64(key)?;
let seg = self.cold_segments.get(segment_id as usize)?.as_ref()?;
let payload = seg.lookup(u64_key)?;
let (row, _) = decode_row_body_dense(&payload, &t.schema, seg.codec_version()).ok()?;
self.cold_read_stats
.cold_reads
.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
Some(row)
}
pub fn lookup_by_pk(&self, table: &str, index_name: &str, key: &IndexKey) -> Option<Row<'_>> {
let t = self.get(table)?;
let idx = t.indices.iter().find(|i| i.name == index_name)?;
let locators = idx.lookup_eq(key);
let cold_u64_key = index_key_as_u64(key);
for loc in locators {
match *loc {
RowLocator::Hot(i) => {
if let Some(row) = t.rows.get(i) {
return Some(row.clone());
}
}
RowLocator::Cold {
segment_id,
page_offset: _,
} => {
let Some(u64_key) = cold_u64_key else {
continue;
};
let Some(seg) = self
.cold_segments
.get(segment_id as usize)
.and_then(|s| s.as_deref())
else {
continue;
};
let Some(payload) = seg.lookup(u64_key) else {
continue;
};
let (row, _) =
decode_row_body_dense(&payload, &t.schema, seg.codec_version()).ok()?;
return Some(row);
}
}
}
None
}
pub fn promote_cold_row(
&mut self,
table_name: &str,
index_name: &str,
key: &IndexKey,
) -> Result<Option<usize>, StorageError> {
let cold_loc = self.find_cold_locator(table_name, index_name, key)?;
let Some((segment_id, _page_offset)) = cold_loc else {
return Ok(None);
};
let u64_key = index_key_as_u64(key).ok_or_else(|| {
StorageError::Corrupt(
"promote_cold_row: key type not coercible to u64 (cold tier requires integer PK)"
.into(),
)
})?;
let schema = self
.get(table_name)
.ok_or_else(|| {
StorageError::Corrupt(format!("promote_cold_row: table {table_name:?} not found"))
})?
.schema
.clone();
let seg = self
.cold_segments
.get(segment_id as usize)
.and_then(|s| s.as_ref())
.ok_or_else(|| {
StorageError::Corrupt(format!(
"promote_cold_row: segment {segment_id} not registered on catalog"
))
})?;
let payload = seg.lookup(u64_key).ok_or_else(|| {
StorageError::Corrupt(format!(
"promote_cold_row: key {u64_key} resolves to segment {segment_id} \
but the segment's bloom/page lookup didn't return a row"
))
})?;
let (row, _consumed) = decode_row_body_dense(&payload, &schema, seg.codec_version())?;
let t = self
.get_mut(table_name)
.expect("table existed at lookup time");
t.insert(row)?;
let new_hot_idx =
t.rows.len().checked_sub(1).ok_or_else(|| {
StorageError::Corrupt("promote_cold_row: empty after insert".into())
})?;
t.remove_cold_locators_for_key(index_name, key)?;
Ok(Some(new_hot_idx))
}
pub fn shadow_cold_row(
&mut self,
table_name: &str,
index_name: &str,
key: &IndexKey,
) -> Result<usize, StorageError> {
let t = self.get_mut(table_name).ok_or_else(|| {
StorageError::Corrupt(format!("shadow_cold_row: table {table_name:?} not found"))
})?;
t.remove_cold_locators_for_key(index_name, key)
}
pub fn prepare_freeze_slice(
&self,
table_name: &str,
index_name: &str,
row_range: core::ops::Range<usize>,
) -> Result<FreezeSlice, StorageError> {
let table = self.get(table_name).ok_or_else(|| {
StorageError::Corrupt(format!(
"prepare_freeze_slice: table {table_name:?} not found"
))
})?;
let idx = table
.indices
.iter()
.find(|i| i.name == index_name)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"prepare_freeze_slice: index {index_name:?} not found on {table_name:?}"
))
})?;
if !matches!(idx.kind, IndexKind::BTree(_)) {
return Err(StorageError::Corrupt(format!(
"prepare_freeze_slice: index {index_name:?} is NSW; only BTree indices may freeze"
)));
}
if row_range.end > table.rows.len() {
return Err(StorageError::Corrupt(format!(
"prepare_freeze_slice: row_range end {} > row_count {}",
row_range.end,
table.rows.len()
)));
}
let column_position = idx.column_position;
let schema = table.schema.clone();
let mut rows: Vec<(u64, Vec<u8>, IndexKey)> = Vec::with_capacity(row_range.len());
for row_idx in row_range.clone() {
let row = table.rows.get(row_idx).expect("bounds-checked above");
let key = IndexKey::from_value(&row.values[column_position]).ok_or_else(|| {
StorageError::Corrupt(format!(
"prepare_freeze_slice: row {row_idx} has NULL / non-key value in index column"
))
})?;
let pk_u64 = index_key_as_u64(&key).ok_or_else(|| {
StorageError::Corrupt(format!(
"prepare_freeze_slice: index {index_name:?} column type is non-integer; \
v5.2.2 cold tier requires IndexKey::Int (Text PK lands in v5.5+)"
))
})?;
rows.push((pk_u64, encode_row_body_dense(row, &schema), key));
}
rows.sort_by_key(|(k, _, _)| *k);
Ok(FreezeSlice { row_range, rows })
}
pub fn commit_freeze_slices(
&mut self,
table_name: &str,
index_name: &str,
slices: Vec<FreezeSlice>,
) -> Result<FreezeReport, StorageError> {
let table = self.get(table_name).ok_or_else(|| {
StorageError::Corrupt(format!(
"commit_freeze_slices: table {table_name:?} not found"
))
})?;
let idx = table
.indices
.iter()
.find(|i| i.name == index_name)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"commit_freeze_slices: index {index_name:?} not found on {table_name:?}"
))
})?;
if !matches!(idx.kind, IndexKind::BTree(_)) {
return Err(StorageError::Corrupt(format!(
"commit_freeze_slices: index {index_name:?} is NSW; only BTree indices may freeze"
)));
}
let mut ordered = slices;
ordered.sort_by_key(|s| s.row_range.start);
let mut expected_start = 0usize;
for s in &ordered {
if s.row_range.start != expected_start {
return Err(StorageError::Corrupt(format!(
"commit_freeze_slices: gap/overlap at row {}; expected start {}",
s.row_range.start, expected_start
)));
}
expected_start = s.row_range.end;
}
let max_rows = expected_start;
if max_rows > table.rows.len() {
return Err(StorageError::Corrupt(format!(
"commit_freeze_slices: total row range {} exceeds row_count {}",
max_rows,
table.rows.len()
)));
}
if max_rows == 0 {
return Ok(FreezeReport {
segment_id: u32::MAX,
frozen_rows: 0,
bytes_freed: 0,
segment_bytes: Vec::new(),
});
}
let total_rows: usize = ordered.iter().map(|s| s.rows.len()).sum();
if total_rows != max_rows {
return Err(StorageError::Corrupt(format!(
"commit_freeze_slices: total slice rows {total_rows} ≠ row_range coverage {max_rows}"
)));
}
let mut cursors: Vec<usize> = alloc::vec![0; ordered.len()];
let mut merged: Vec<(u64, Vec<u8>, IndexKey)> = Vec::with_capacity(total_rows);
loop {
let mut pick: Option<usize> = None;
for (i, c) in cursors.iter().enumerate() {
let slice = &ordered[i];
if *c >= slice.rows.len() {
continue;
}
match pick {
None => pick = Some(i),
Some(j) => {
if slice.rows[*c].0 < ordered[j].rows[cursors[j]].0 {
pick = Some(i);
}
}
}
}
let Some(i) = pick else { break };
let row = ordered[i].rows[cursors[i]].clone();
cursors[i] += 1;
merged.push(row);
}
for w in merged.windows(2) {
if w[0].0 == w[1].0 {
return Err(StorageError::Corrupt(format!(
"commit_freeze_slices: duplicate PK {} across slices",
w[0].0
)));
}
}
let post_swap_keys: Vec<IndexKey> = merged.iter().map(|(_, _, k)| k.clone()).collect();
let seg_rows: Vec<(u64, Vec<u8>)> =
merged.into_iter().map(|(k, body, _)| (k, body)).collect();
let frozen_rows = seg_rows.len();
let (seg_bytes, _meta) = encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES)
.map_err(|e| StorageError::Corrupt(format!("commit_freeze_slices: encode: {e}")))?;
let bytes_before = self.get(table_name).expect("just validated").hot_bytes();
let positions: Vec<usize> = (0..max_rows).collect();
let t_mut = self
.get_mut(table_name)
.expect("just validated; still present");
let removed = t_mut.delete_rows(&positions);
debug_assert_eq!(removed, max_rows, "delete_rows count matches request");
let bytes_after = t_mut.hot_bytes();
let bytes_freed = bytes_before.saturating_sub(bytes_after);
let segment_id = self
.load_segment_bytes(seg_bytes.clone())
.map_err(|e| StorageError::Corrupt(format!("commit_freeze_slices: load: {e}")))?;
let new_cold = post_swap_keys.into_iter().map(|k| {
(
k,
RowLocator::Cold {
segment_id,
page_offset: 0,
},
)
});
let t_mut = self.get_mut(table_name).expect("still present");
t_mut.register_cold_locators(index_name, new_cold)?;
t_mut.mark_cold_row_count_stale();
Ok(FreezeReport {
segment_id,
frozen_rows,
bytes_freed,
segment_bytes: seg_bytes,
})
}
pub fn compact_cold_segments(
&mut self,
table_name: &str,
index_name: &str,
target_segment_bytes: u64,
) -> Result<CompactReport, StorageError> {
let t = self.get(table_name).ok_or_else(|| {
StorageError::Corrupt(format!(
"compact_cold_segments: table {table_name:?} not found"
))
})?;
let idx = t
.indices
.iter()
.find(|i| i.name == index_name)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"compact_cold_segments: index {index_name:?} not found on {table_name:?}"
))
})?;
let map = match &idx.kind {
IndexKind::BTree(m) => m,
IndexKind::Nsw(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => {
return Err(StorageError::Corrupt(format!(
"compact_cold_segments: index {index_name:?} is not BTree; \
compaction applies only to BTree cold-tier indices"
)));
}
};
let mut referenced_ids: BTreeSet<u32> = BTreeSet::new();
for (_key, locators) in map.iter() {
for loc in locators {
if let RowLocator::Cold { segment_id, .. } = loc {
referenced_ids.insert(*segment_id);
}
}
}
let candidate_set: BTreeSet<u32> = referenced_ids
.into_iter()
.filter(|id| {
self.cold_segments
.get(*id as usize)
.and_then(|s| s.as_deref())
.is_some_and(|s| (s.bytes().len() as u64) < target_segment_bytes)
})
.collect();
if candidate_set.len() < 2 {
return Ok(CompactReport {
sources: Vec::new(),
merged_segment_id: None,
merged_segment_bytes: Vec::new(),
merged_rows: 0,
deleted_rows_pruned: 0,
bytes_reclaimed_estimate: 0,
});
}
let mut source_row_count: usize = 0;
let mut source_byte_total: u64 = 0;
for &id in &candidate_set {
let seg = self.cold_segments[id as usize]
.as_ref()
.expect("candidate selected only when slot is Some");
source_row_count = source_row_count.saturating_add(seg.meta().num_rows as usize);
source_byte_total = source_byte_total.saturating_add(seg.bytes().len() as u64);
}
let mut collected: BTreeMap<u64, (Vec<u8>, IndexKey)> = BTreeMap::new();
for (key, locators) in map.iter() {
for loc in locators {
let RowLocator::Cold { segment_id, .. } = loc else {
continue;
};
if !candidate_set.contains(segment_id) {
continue;
}
let u64_key = index_key_as_u64(key).ok_or_else(|| {
StorageError::Corrupt(format!(
"compact_cold_segments: index {index_name:?} has non-integer Cold key; \
cold tier requires IndexKey::Int (Text PK lands in v5.5+)"
))
})?;
let seg = self.cold_segments[*segment_id as usize]
.as_ref()
.expect("candidate slot guaranteed Some above");
let payload = seg.lookup(u64_key).ok_or_else(|| {
StorageError::Corrupt(format!(
"compact_cold_segments: BTree {index_name:?} points key={u64_key} \
at segment {segment_id} but the segment lookup missed"
))
})?;
collected.insert(u64_key, (payload, key.clone()));
break;
}
}
let merged_rows = collected.len();
let deleted_rows_pruned = source_row_count.saturating_sub(merged_rows);
let seg_rows: Vec<(u64, Vec<u8>)> = collected
.iter()
.map(|(k, (body, _))| (*k, body.clone()))
.collect();
let (seg_bytes, _meta) = encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES)
.map_err(|e| StorageError::Corrupt(format!("compact_cold_segments: encode: {e}")))?;
let merged_bytes_len = seg_bytes.len() as u64;
let merged_segment_id = self
.load_segment_bytes(seg_bytes.clone())
.map_err(|e| StorageError::Corrupt(format!("compact_cold_segments: load: {e}")))?;
let entries: Vec<(IndexKey, Vec<RowLocator>)> = {
let t = self
.get(table_name)
.expect("table existed at the start of this fn");
let idx = t
.indices
.iter()
.find(|i| i.name == index_name)
.expect("index existed at the start of this fn");
let IndexKind::BTree(map) = &idx.kind else {
unreachable!("validated above");
};
map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
};
let t_mut = self
.get_mut(table_name)
.expect("table existed at the start of this fn");
let idx_mut = t_mut
.indices
.iter_mut()
.find(|i| i.name == index_name)
.expect("index existed at the start of this fn");
let IndexKind::BTree(map_mut) = &mut idx_mut.kind else {
unreachable!("validated above");
};
for (key, locators) in entries {
let mut new_locs: Vec<RowLocator> = Vec::with_capacity(locators.len());
let mut changed = false;
for loc in &locators {
match *loc {
RowLocator::Cold {
segment_id,
page_offset: _,
} if candidate_set.contains(&segment_id) => {
let replacement = RowLocator::Cold {
segment_id: merged_segment_id,
page_offset: 0,
};
if !new_locs.contains(&replacement) {
new_locs.push(replacement);
}
changed = true;
}
other => new_locs.push(other),
}
}
if changed {
map_mut.insert_mut(key, new_locs);
}
}
for &id in &candidate_set {
self.tombstone_segment(id)?;
}
let bytes_reclaimed_estimate = source_byte_total.saturating_sub(merged_bytes_len);
Ok(CompactReport {
sources: candidate_set.into_iter().collect(),
merged_segment_id: Some(merged_segment_id),
merged_segment_bytes: seg_bytes,
merged_rows,
deleted_rows_pruned,
bytes_reclaimed_estimate,
})
}
fn find_cold_locator(
&self,
table_name: &str,
index_name: &str,
key: &IndexKey,
) -> Result<Option<(u32, u32)>, StorageError> {
let t = self.get(table_name).ok_or_else(|| {
StorageError::Corrupt(format!("find_cold_locator: table {table_name:?} not found"))
})?;
let idx = t
.indices
.iter()
.find(|i| i.name == index_name)
.ok_or_else(|| {
StorageError::Corrupt(format!(
"find_cold_locator: index {index_name:?} not found on {table_name:?}"
))
})?;
if !matches!(idx.kind, IndexKind::BTree(_)) {
return Err(StorageError::Corrupt(format!(
"find_cold_locator: index {index_name:?} is NSW; promote-on-write only applies to BTree indices"
)));
}
for loc in idx.lookup_eq(key) {
if let RowLocator::Cold {
segment_id,
page_offset,
} = *loc
{
return Ok(Some((segment_id, page_offset)));
}
}
Ok(None)
}
}
fn index_key_as_u64(key: &IndexKey) -> Option<u64> {
match key {
IndexKey::Int(n) => Some(n.cast_unsigned()),
IndexKey::Text(_) | IndexKey::Bool(_) | IndexKey::Uuid(_) => None,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum StorageError {
DuplicateTable {
name: String,
},
TableNotFound {
name: String,
},
ArityMismatch {
expected: usize,
actual: usize,
},
TypeMismatch {
column: String,
expected: DataType,
actual: DataType,
position: usize,
},
NullInNotNull {
column: String,
},
DuplicateIndex {
name: String,
},
ColumnNotFound {
column: String,
},
Corrupt(String),
IndexNotFound {
name: String,
},
Unsupported(String),
SequenceExhausted {
name: String,
limit: i64,
is_max: bool,
},
}
impl fmt::Display for StorageError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::DuplicateTable { name } => write!(f, "relation \"{name}\" already exists"),
Self::TableNotFound { name } => write!(f, "relation \"{name}\" does not exist"),
Self::ArityMismatch { expected, actual } => write!(
f,
"row arity mismatch: expected {expected} columns, got {actual}"
),
Self::TypeMismatch {
column,
expected,
actual,
position,
} => write!(
f,
"type mismatch in column {column:?} (position {position}): expected {expected}, got {actual}"
),
Self::NullInNotNull { column } => {
write!(
f,
"null value in column \"{column}\" violates not-null constraint"
)
}
Self::DuplicateIndex { name } => write!(f, "relation \"{name}\" already exists"),
Self::ColumnNotFound { column } => write!(f, "column \"{column}\" does not exist"),
Self::Corrupt(detail) => write!(f, "corrupt on-disk format: {detail}"),
Self::IndexNotFound { name } => write!(f, "index \"{name}\" does not exist"),
Self::Unsupported(detail) => write!(f, "unsupported: {detail}"),
Self::SequenceExhausted {
name,
limit,
is_max,
} => write!(
f,
"nextval: reached {} value of sequence \"{name}\" ({limit})",
if *is_max { "maximum" } else { "minimum" }
),
}
}
}
impl ColumnSchema {
pub fn new(name: impl Into<String>, ty: DataType, nullable: bool) -> Self {
Self {
name: name.into(),
ty,
nullable,
collation_name: None,
default: None,
runtime_default: None,
auto_increment: false,
user_enum_type: None,
user_domain_type: None,
user_composite_type: None,
acl: Vec::new(),
on_update_runtime: None,
collation: Collation::Binary,
is_unsigned: false,
inline_enum_variants: None,
inline_set_variants: None,
generated_stored_expr: None,
identity_always: false,
default_text: None,
auto_restart: None,
scalar_row_source: false,
mysql_int_width: None,
mysql_fsp: None,
}
}
#[must_use]
pub fn with_default(mut self, default: Value<'static>) -> Self {
self.default = Some(default);
self
}
#[must_use]
pub fn with_runtime_default(mut self, expr: impl Into<String>) -> Self {
self.runtime_default = Some(expr.into());
self
}
#[must_use]
pub const fn with_auto_increment(mut self) -> Self {
self.auto_increment = true;
self
}
}
impl TableSchema {
pub fn new(name: impl Into<String>, columns: Vec<ColumnSchema>) -> Self {
Self {
name: name.into(),
columns,
hot_tier_bytes: None,
foreign_keys: Vec::new(),
uniqueness_constraints: Vec::new(),
exclusion_constraints: Vec::new(),
checks: Vec::new(),
partition_role: None,
policies: Vec::new(),
row_security: false,
force_row_security: false,
owner: None,
acl: Vec::new(),
}
}
}
const FILE_MAGIC: &[u8; 8] = b"SPGDB001";
const FILE_VERSION: u8 = 89;
pub const CURRENT_ROW_CODEC_VERSION: u8 = FILE_VERSION;
const FILE_VERSION_CRC_TRAILER: u8 = 54;
const MIN_SUPPORTED_FILE_VERSION: u8 = 8;
const INDEX_KEY_TAG_INT: u8 = 0;
const INDEX_KEY_TAG_TEXT: u8 = 1;
const INDEX_KEY_TAG_BOOL: u8 = 2;
const INDEX_KEY_TAG_UUID: u8 = 3;
impl Catalog {
pub fn serialize(&self) -> Vec<u8> {
let mut out = Vec::with_capacity(64);
out.extend_from_slice(FILE_MAGIC);
out.push(FILE_VERSION);
write_u32(
&mut out,
u32::try_from(self.tables.len()).expect("≤ 4G tables"),
);
for t in &self.tables {
write_str(&mut out, &t.schema.name);
write_u16(
&mut out,
u16::try_from(t.schema.columns.len()).expect("≤ 65k columns/table"),
);
for c in &t.schema.columns {
write_str(&mut out, &c.name);
write_data_type(&mut out, c.ty);
out.push(u8::from(c.nullable));
match &c.default {
None => out.push(0),
Some(v) => {
out.push(1);
write_value(&mut out, v);
}
}
out.push(u8::from(c.auto_increment));
}
write_u32(
&mut out,
u32::try_from(t.rows.len()).expect("≤ 4G rows/table"),
);
for row in &t.rows {
out.extend_from_slice(&encode_row_body_dense(row, &t.schema));
}
write_u16(
&mut out,
u16::try_from(t.indices.len()).expect("≤ 65k indices/table"),
);
for idx in &t.indices {
write_str(&mut out, &idx.name);
write_u16(
&mut out,
u16::try_from(idx.column_position).expect("≤ 65k columns/table"),
);
match &idx.kind {
IndexKind::BTree(map) => {
out.push(0);
write_u32(
&mut out,
u32::try_from(map.len()).expect("≤ 4G index entries/index"),
);
for (key, locators) in map {
write_index_key(&mut out, key);
write_u32(
&mut out,
u32::try_from(locators.len()).expect("≤ 4G locators/key"),
);
for loc in locators {
loc.write_le(&mut out);
}
}
}
IndexKind::Nsw(g) => {
out.push(1);
write_u16(&mut out, u16::try_from(g.m).expect("≤ 65k NSW neighbours"));
write_nsw_graph(&mut out, g);
}
IndexKind::Brin { column_type } => {
out.push(2);
write_data_type(&mut out, *column_type);
}
IndexKind::Gin(map) => {
out.push(3);
write_u32(
&mut out,
u32::try_from(map.len()).expect("≤ 4G GIN posting lists"),
);
for (word, locators) in map {
write_str(&mut out, word);
write_u32(
&mut out,
u32::try_from(locators.len()).expect("≤ 4G locators/posting list"),
);
for loc in locators {
loc.write_le(&mut out);
}
}
}
IndexKind::GinTrgm(map) => {
out.push(4);
write_u32(
&mut out,
u32::try_from(map.len()).expect("≤ 4G trigram-GIN posting lists"),
);
for (tri, locators) in map {
write_str(&mut out, tri);
write_u32(
&mut out,
u32::try_from(locators.len()).expect("≤ 4G locators/posting list"),
);
for loc in locators {
loc.write_le(&mut out);
}
}
}
IndexKind::GinFulltext(map) => {
out.push(5);
write_u32(
&mut out,
u32::try_from(map.len()).expect("≤ 4G fulltext-GIN posting lists"),
);
for (lex, locators) in map {
write_str(&mut out, lex);
write_u32(
&mut out,
u32::try_from(locators.len()).expect("≤ 4G locators/posting list"),
);
for loc in locators {
loc.write_le(&mut out);
}
}
}
IndexKind::GinJsonb(map) => {
out.push(6);
write_u32(
&mut out,
u32::try_from(map.len()).expect("≤ 4G JSONB-GIN posting lists"),
);
for (token, locators) in map {
write_str(&mut out, token);
write_u32(
&mut out,
u32::try_from(locators.len()).expect("≤ 4G locators/posting list"),
);
for loc in locators {
loc.write_le(&mut out);
}
}
}
}
write_u16(
&mut out,
u16::try_from(idx.included_columns.len()).expect("≤ 65k INCLUDE columns/index"),
);
for col_pos in &idx.included_columns {
write_u16(
&mut out,
u16::try_from(*col_pos).expect("≤ 65k columns/table"),
);
}
match &idx.partial_predicate {
None => out.push(0),
Some(pred) => {
out.push(1);
write_str(&mut out, pred);
}
}
match &idx.expression {
None => out.push(0),
Some(expr) => {
out.push(1);
write_str(&mut out, expr);
}
}
out.push(u8::from(idx.is_unique));
write_u16(
&mut out,
u16::try_from(idx.extra_column_positions.len())
.expect("≤ 65k extra cols / index"),
);
for cp in &idx.extra_column_positions {
write_u16(&mut out, u16::try_from(*cp).expect("≤ 65k columns/table"));
}
out.push(u8::from(idx.nulls_not_distinct));
out.push(u8::from(idx.descending));
out.push(match idx.nulls_first {
None => 0,
Some(true) => 1,
Some(false) => 2,
});
match &idx.collation {
Some(c) => {
out.push(1);
write_str(&mut out, c);
}
None => out.push(0),
}
}
match t.schema.hot_tier_bytes {
None => out.push(0),
Some(n) => {
out.push(1);
out.extend_from_slice(&n.to_le_bytes());
}
}
write_u16(
&mut out,
u16::try_from(t.schema.foreign_keys.len()).expect("≤ 65k FKs/table"),
);
for fk in &t.schema.foreign_keys {
match &fk.name {
None => out.push(0),
Some(n) => {
out.push(1);
write_str(&mut out, n);
}
}
write_u16(
&mut out,
u16::try_from(fk.local_columns.len()).expect("≤ 65k FK columns"),
);
for &p in &fk.local_columns {
write_u16(&mut out, u16::try_from(p).expect("≤ 65k columns/table"));
}
write_str(&mut out, &fk.parent_table);
write_u16(
&mut out,
u16::try_from(fk.parent_columns.len()).expect("≤ 65k FK parent columns"),
);
for &p in &fk.parent_columns {
write_u16(&mut out, u16::try_from(p).expect("≤ 65k columns/table"));
}
out.push(fk.on_delete.tag());
out.push(fk.on_update.tag());
out.push(fk.match_type.tag());
out.push(u8::from(fk.deferrable) | (u8::from(fk.initially_deferred) << 1));
}
write_u16(
&mut out,
u16::try_from(t.schema.uniqueness_constraints.len())
.expect("≤ 65k uniqueness constraints/table"),
);
for uc in &t.schema.uniqueness_constraints {
out.push(u8::from(uc.is_primary_key));
write_u16(
&mut out,
u16::try_from(uc.columns.len()).expect("≤ 65k cols in uniqueness constraint"),
);
for &p in &uc.columns {
write_u16(&mut out, u16::try_from(p).expect("≤ 65k columns/table"));
}
out.push(u8::from(uc.nulls_not_distinct));
}
let mut rt_defaults: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(e) = &c.runtime_default {
rt_defaults.push((i, e.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(rt_defaults.len()).expect("≤ 65k runtime defaults/table"),
);
for (pos, expr) in rt_defaults {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, expr);
}
write_u16(
&mut out,
u16::try_from(t.schema.checks.len()).expect("≤ 65k CHECK constraints/table"),
);
for c in &t.schema.checks {
write_str(&mut out, c.expr.as_str());
}
let mut enum_bindings: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(e) = &c.user_enum_type {
enum_bindings.push((i, e.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(enum_bindings.len()).expect("≤ 65k enum-typed columns/table"),
);
for (pos, ename) in enum_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, ename);
}
let mut domain_bindings: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(d) = &c.user_domain_type {
domain_bindings.push((i, d.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(domain_bindings.len()).expect("≤ 65k domain-typed columns/table"),
);
for (pos, dname) in domain_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, dname);
}
let mut on_update_bindings: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(e) = &c.on_update_runtime {
on_update_bindings.push((i, e.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(on_update_bindings.len()).expect("≤ 65k ON UPDATE columns/table"),
);
for (pos, expr_src) in on_update_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, expr_src);
}
let mut coll_bindings: Vec<(usize, u8)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
let tag = match c.collation {
Collation::Binary => continue,
Collation::CaseInsensitive => Collation::TAG_CASE_INSENSITIVE,
};
coll_bindings.push((i, tag));
}
write_u16(
&mut out,
u16::try_from(coll_bindings.len()).expect("≤ 65k collation bindings/table"),
);
for (pos, tag) in coll_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
out.push(tag);
}
let mut unsigned_bindings: Vec<usize> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if c.is_unsigned {
unsigned_bindings.push(i);
}
}
write_u16(
&mut out,
u16::try_from(unsigned_bindings.len()).expect("≤ 65k UNSIGNED columns/table"),
);
for pos in unsigned_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
}
let mut enum_inline_bindings: Vec<(usize, &[String])> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(vs) = &c.inline_enum_variants {
enum_inline_bindings.push((i, vs.as_slice()));
}
}
write_u16(
&mut out,
u16::try_from(enum_inline_bindings.len()).expect("≤ 65k inline-ENUM columns/table"),
);
for (pos, variants) in enum_inline_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_u16(
&mut out,
u16::try_from(variants.len()).expect("≤ 65k variants/ENUM"),
);
for v in variants {
write_str(&mut out, v.as_str());
}
}
let mut set_inline_bindings: Vec<(usize, &[String])> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(vs) = &c.inline_set_variants {
set_inline_bindings.push((i, vs.as_slice()));
}
}
write_u16(
&mut out,
u16::try_from(set_inline_bindings.len()).expect("≤ 65k inline-SET columns/table"),
);
for (pos, variants) in set_inline_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_u16(
&mut out,
u16::try_from(variants.len()).expect("≤ 65k variants/SET"),
);
for v in variants {
write_str(&mut out, v.as_str());
}
}
write_partition_role(&mut out, t.schema.partition_role.as_ref());
let mut gen_bindings: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(src) = &c.generated_stored_expr {
gen_bindings.push((i, src.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(gen_bindings.len()).expect("≤ 65k GENERATED STORED columns/table"),
);
for (pos, src) in gen_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, src);
}
let mut default_texts: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(src) = &c.default_text {
default_texts.push((i, src.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(default_texts.len()).expect("≤ 65k defaulted columns/table"),
);
for (pos, src) in default_texts {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, src);
}
out.push(u8::from(t.schema.row_security));
out.push(u8::from(t.schema.force_row_security));
write_u16(
&mut out,
u16::try_from(t.schema.policies.len()).expect("≤ 65k policies/table"),
);
for p in &t.schema.policies {
write_str(&mut out, &p.name);
out.push(p.cmd.to_wire_byte());
out.push(u8::from(p.permissive));
write_u16(
&mut out,
u16::try_from(p.roles.len()).expect("≤ 65k roles/policy"),
);
for r in &p.roles {
write_str(&mut out, r);
}
match &p.using_expr {
Some(s) => {
out.push(1);
write_str(&mut out, s);
}
None => out.push(0),
}
match &p.with_check_expr {
Some(s) => {
out.push(1);
write_str(&mut out, s);
}
None => out.push(0),
}
}
debug_assert_eq!(
t.rows.len(),
t.headers.len(),
"headers must be lock-step with rows at serialize"
);
debug_assert_eq!(
t.rows.len(),
t.rowids.len(),
"rowids must be lock-step with rows at serialize"
);
write_u32(
&mut out,
u32::try_from(t.rows.len()).expect("≤ 4G rows/table"),
);
for (h, rid) in t.headers.iter().zip(t.rowids.iter()) {
out.extend_from_slice(&h.xmin.to_le_bytes());
out.extend_from_slice(&h.xmax.to_le_bytes());
out.push(h.flags);
out.extend_from_slice(&rid.0.to_le_bytes());
}
out.extend_from_slice(&t.next_rowid.to_le_bytes());
write_u16(
&mut out,
u16::try_from(t.schema.checks.len()).expect("≤ 65k CHECK constraints/table"),
);
for c in &t.schema.checks {
match &c.name {
Some(n) => {
out.push(1);
write_str(&mut out, n);
}
None => out.push(0),
}
}
write_u16(
&mut out,
u16::try_from(t.schema.uniqueness_constraints.len())
.expect("≤ 65k uniqueness constraints/table"),
);
for uc in &t.schema.uniqueness_constraints {
match &uc.name {
Some(n) => {
out.push(1);
write_str(&mut out, n);
}
None => out.push(0),
}
}
let mut comp_bindings: Vec<(usize, &str)> = Vec::new();
for (i, c) in t.schema.columns.iter().enumerate() {
if let Some(n) = &c.user_composite_type {
comp_bindings.push((i, n.as_str()));
}
}
write_u16(
&mut out,
u16::try_from(comp_bindings.len()).expect("≤ 65k composite-typed columns/table"),
);
for (pos, n) in comp_bindings {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_str(&mut out, n);
}
match &t.schema.owner {
Some(o) => {
out.push(1);
write_str(&mut out, o);
}
None => out.push(0),
}
write_u16(
&mut out,
u16::try_from(t.schema.acl.len()).expect("≤ 65k aclitems/table"),
);
for a in &t.schema.acl {
write_str(&mut out, &a.grantee);
write_u16(&mut out, a.privs);
write_u16(&mut out, a.grantable);
write_str(&mut out, &a.grantor);
}
let granted: Vec<(usize, &ColumnSchema)> = t
.schema
.columns
.iter()
.enumerate()
.filter(|(_, c)| !c.acl.is_empty())
.collect();
write_u16(
&mut out,
u16::try_from(granted.len()).expect("≤ 65k granted columns/table"),
);
for (pos, c) in granted {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
write_u16(
&mut out,
u16::try_from(c.acl.len()).expect("≤ 65k aclitems/column"),
);
for a in &c.acl {
write_str(&mut out, &a.grantee);
write_u16(&mut out, a.privs);
write_u16(&mut out, a.grantable);
write_str(&mut out, &a.grantor);
}
}
write_u16(
&mut out,
u16::try_from(t.schema.exclusion_constraints.len())
.expect("≤ 65k exclusion constraints/table"),
);
for ex in &t.schema.exclusion_constraints {
write_str(&mut out, &ex.name);
match &ex.method {
Some(m) => {
out.push(1);
write_str(&mut out, m);
}
None => out.push(0),
}
write_u16(
&mut out,
u16::try_from(ex.elements.len()).expect("≤ 65k elements/exclusion"),
);
for (pos, op) in &ex.elements {
write_u16(&mut out, u16::try_from(*pos).expect("≤ 65k columns/table"));
write_str(&mut out, op);
}
}
let restarts: Vec<(usize, i64)> = t
.schema
.columns
.iter()
.enumerate()
.filter_map(|(i, c)| c.auto_restart.map(|n| (i, n)))
.collect();
write_u16(
&mut out,
u16::try_from(restarts.len()).expect("≤ 65k restart columns/table"),
);
for (pos, n) in restarts {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
out.extend_from_slice(&n.to_le_bytes());
}
let int_widths: Vec<(usize, u8)> = t
.schema
.columns
.iter()
.enumerate()
.filter_map(|(i, c)| {
c.mysql_int_width.map(|w| {
let tag = match w {
MysqlIntWidth::Tiny => 0u8,
MysqlIntWidth::Medium => 1u8,
MysqlIntWidth::Small => 2u8,
MysqlIntWidth::Int => 3u8,
MysqlIntWidth::Big => 4u8,
};
(i, tag)
})
})
.collect();
write_u16(
&mut out,
u16::try_from(int_widths.len()).expect("≤ 65k narrow-int columns/table"),
);
for (pos, tag) in int_widths {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
out.push(tag);
}
let fsps: Vec<(usize, u8)> = t
.schema
.columns
.iter()
.enumerate()
.filter_map(|(i, c)| c.mysql_fsp.map(|p| (i, p)))
.collect();
write_u16(
&mut out,
u16::try_from(fsps.len()).expect("≤ 65k temporal columns/table"),
);
for (pos, fsp) in fsps {
write_u16(&mut out, u16::try_from(pos).expect("≤ 65k columns/table"));
out.push(fsp);
}
let unvalidated: Vec<usize> = t
.schema
.checks
.iter()
.enumerate()
.filter_map(|(i, c)| (!c.validated).then_some(i))
.collect();
write_u16(
&mut out,
u16::try_from(unvalidated.len()).expect("≤ 65k CHECK constraints/table"),
);
for idx in unvalidated {
write_u16(&mut out, u16::try_from(idx).expect("≤ 65k CHECK/table"));
}
let collated: Vec<(usize, &str)> = t
.schema
.columns
.iter()
.enumerate()
.filter_map(|(i, c)| c.collation_name.as_deref().map(|n| (i, n)))
.collect();
write_u16(
&mut out,
u16::try_from(collated.len()).expect("≤ 65k columns/table"),
);
for (idx, name) in collated {
write_u16(&mut out, u16::try_from(idx).expect("≤ 65k columns/table"));
write_str(&mut out, name);
}
write_u16(
&mut out,
u16::try_from(t.schema.uniqueness_constraints.len())
.expect("≤ 65k uniqueness constraints/table"),
);
for uc in &t.schema.uniqueness_constraints {
out.push(u8::from(uc.deferrable) | (u8::from(uc.initially_deferred) << 1));
}
}
write_u32(
&mut out,
u32::try_from(self.functions.len()).expect("≤ 4G functions"),
);
for fd in self.functions.values() {
write_str(&mut out, &fd.name);
write_str(&mut out, &fd.args_repr);
write_str(&mut out, &fd.returns);
write_str(&mut out, &fd.language);
write_str_long(&mut out, &fd.body);
}
write_u32(
&mut out,
u32::try_from(self.triggers.len()).expect("≤ 4G triggers"),
);
for td in &self.triggers {
write_str(&mut out, &td.name);
write_str(&mut out, &td.table);
write_str(&mut out, &td.timing);
write_u16(
&mut out,
u16::try_from(td.events.len()).expect("≤ 65k events / trigger"),
);
for ev in &td.events {
write_str(&mut out, ev);
}
write_str(&mut out, &td.for_each);
write_str(&mut out, &td.function);
write_u16(
&mut out,
u16::try_from(td.update_columns.len()).expect("≤ 65k cols / trigger"),
);
for c in &td.update_columns {
write_str(&mut out, c);
}
out.push(u8::from(td.enabled));
write_str(&mut out, &td.when_condition);
}
write_u32(
&mut out,
u32::try_from(self.sequences.len()).expect("≤ 4G sequences"),
);
for seq in self.sequences.values() {
write_str(&mut out, &seq.name);
out.push(match seq.data_type {
SequenceDataType::SmallInt => 0,
SequenceDataType::Int => 1,
SequenceDataType::BigInt => 2,
});
out.extend_from_slice(&seq.start.to_le_bytes());
out.extend_from_slice(&seq.increment.to_le_bytes());
out.extend_from_slice(&seq.min_value.to_le_bytes());
out.extend_from_slice(&seq.max_value.to_le_bytes());
out.extend_from_slice(&seq.cache.to_le_bytes());
out.push(u8::from(seq.cycle));
match &seq.owned_by {
None => out.push(0),
Some((table, column)) => {
out.push(1);
write_str(&mut out, table);
write_str(&mut out, column);
}
}
out.extend_from_slice(&seq.last_value.to_le_bytes());
out.push(u8::from(seq.is_called));
}
write_u32(
&mut out,
u32::try_from(self.views.len()).expect("≤ 4G views"),
);
for view in self.views.values() {
write_str(&mut out, &view.name);
write_u16(
&mut out,
u16::try_from(view.columns.len()).expect("≤ 65k cols / view"),
);
for c in &view.columns {
write_str(&mut out, c);
}
write_str_long(&mut out, &view.body);
out.push(view.check_option);
}
write_u32(
&mut out,
u32::try_from(self.materialized_views.len()).expect("≤ 4G materialized views"),
);
for (name, body) in &self.materialized_views {
write_str(&mut out, name);
write_str_long(&mut out, body);
}
write_u32(
&mut out,
u32::try_from(self.enum_types.len()).expect("≤ 4G enum types"),
);
for e in self.enum_types.values() {
write_str(&mut out, &e.name);
write_u16(
&mut out,
u16::try_from(e.labels.len()).expect("≤ 65k labels / enum"),
);
for l in &e.labels {
write_str(&mut out, l);
}
}
write_u32(
&mut out,
u32::try_from(self.domain_types.len()).expect("≤ 4G domain types"),
);
for d in self.domain_types.values() {
write_str(&mut out, &d.name);
write_data_type(&mut out, d.base_type);
out.push(u8::from(d.nullable));
match &d.default {
None => out.push(0),
Some(s) => {
out.push(1);
write_str(&mut out, s);
}
}
write_u16(
&mut out,
u16::try_from(d.checks.len()).expect("≤ 65k CHECKs / domain"),
);
for c in &d.checks {
write_str(&mut out, &c.expr);
write_str(&mut out, &c.name);
}
match &d.base_domain {
None => out.push(0),
Some(s) => {
out.push(1);
write_str(&mut out, s);
}
}
}
write_u32(
&mut out,
u32::try_from(self.schemas.len()).expect("≤ 4G schemas"),
);
for name in &self.schemas {
write_str(&mut out, name);
}
write_u32(
&mut out,
u32::try_from(self.composite_types.len()).expect("≤ 4G composite types"),
);
for c in self.composite_types.values() {
write_str(&mut out, &c.name);
write_u16(
&mut out,
u16::try_from(c.fields.len()).expect("≤ 65k fields / composite"),
);
for (i, (fname, fty)) in c.fields.iter().enumerate() {
write_str(&mut out, fname);
write_data_type(&mut out, *fty);
match c.field_user_types.get(i).and_then(Option::as_ref) {
None => out.push(0),
Some(n) => {
out.push(1);
write_str(&mut out, n);
}
}
}
}
write_u32(
&mut out,
u32::try_from(self.comments.len()).expect("≤ 4G comments"),
);
for (k, v) in &self.comments {
write_str(&mut out, k);
write_str_long(&mut out, v);
}
let acl_out = |out: &mut Vec<u8>, acl: &[AclItem]| {
write_u16(out, u16::try_from(acl.len()).expect("≤ 65k aclitems"));
for a in acl {
write_str(out, &a.grantee);
write_u16(out, a.privs);
write_u16(out, a.grantable);
write_str(out, &a.grantor);
}
};
let owned: Vec<&SequenceDef> = self
.sequences
.values()
.filter(|s| s.owner.is_some() || !s.acl.is_empty())
.collect();
write_u32(
&mut out,
u32::try_from(owned.len()).expect("≤ 4G sequences"),
);
for seq in owned {
write_str(&mut out, &seq.name);
match &seq.owner {
Some(o) => {
out.push(1);
write_str(&mut out, o);
}
None => out.push(0),
}
acl_out(&mut out, &seq.acl);
}
acl_out(&mut out, &self.schema_acl);
acl_out(&mut out, &self.database_acl);
let fns: Vec<&FunctionDef> = self
.functions
.values()
.filter(|f| f.owner.is_some() || !f.acl.is_empty())
.collect();
write_u32(&mut out, u32::try_from(fns.len()).expect("≤ 4G functions"));
for f in fns {
write_str(&mut out, &function_signature_key(&f.name, &f.args_repr));
match &f.owner {
Some(o) => {
out.push(1);
write_str(&mut out, o);
}
None => out.push(0),
}
acl_out(&mut out, &f.acl);
}
write_u32(
&mut out,
u32::try_from(self.rules.len()).expect("≤ 4G rules"),
);
for r in &self.rules {
write_str(&mut out, &r.name);
write_str(&mut out, &r.table);
write_str(&mut out, &r.event);
out.push(u8::from(r.instead));
write_str(&mut out, &r.when_condition);
write_u16(
&mut out,
u16::try_from(r.commands.len()).expect("≤ 65k commands / rule"),
);
for c in &r.commands {
write_str(&mut out, c);
}
}
write_u32(
&mut out,
u32::try_from(self.statistics_ext.len()).expect("≤ 4G statistics objects"),
);
for st in &self.statistics_ext {
write_str(&mut out, &st.name);
write_str(&mut out, &st.table);
write_u16(
&mut out,
u16::try_from(st.kinds.len()).expect("≤ 65k kinds"),
);
for k in &st.kinds {
write_str(&mut out, k);
}
write_u16(
&mut out,
u16::try_from(st.columns.len()).expect("≤ 65k columns"),
);
for c in &st.columns {
write_str(&mut out, c);
}
}
write_u32(
&mut out,
u32::try_from(self.large_objects.len()).expect("≤ 4G large objects"),
);
for (oid, bytes) in &self.large_objects {
write_u32(&mut out, *oid);
write_u32(
&mut out,
u32::try_from(bytes.len()).expect("≤ 4G per object"),
);
out.extend_from_slice(bytes);
}
let attr_fns: Vec<(&String, &FunctionDef)> = self
.functions
.iter()
.filter(|(_, f)| {
f.volatility != FN_VOLATILE
|| f.strict
|| f.security_definer
|| f.leakproof
|| f.parallel != FN_PARALLEL_UNSAFE
|| f.cost.is_some()
|| f.rows.is_some()
})
.collect();
write_u32(
&mut out,
u32::try_from(attr_fns.len()).expect("≤ 4G functions"),
);
for (key, f) in attr_fns {
write_str(&mut out, key);
out.push(f.volatility);
let flags = u8::from(f.strict)
| (u8::from(f.security_definer) << 1)
| (u8::from(f.leakproof) << 2);
out.push(flags);
out.push(f.parallel);
out.extend_from_slice(&f.cost.unwrap_or(f64::NAN).to_le_bytes());
out.extend_from_slice(&f.rows.unwrap_or(f64::NAN).to_le_bytes());
}
write_u32(
&mut out,
u32::try_from(self.db_role_settings.len()).expect("≤ 4G scopes"),
);
for ((db, role), params) in &self.db_role_settings {
write_str(&mut out, db);
write_str(&mut out, role);
write_u32(&mut out, u32::try_from(params.len()).expect("≤ 4G params"));
for (name, value) in params {
write_str(&mut out, name);
write_str(&mut out, value);
}
}
write_u32(
&mut out,
u32::try_from(self.replication_slots.len()).expect("≤ 4G slots"),
);
for (name, (plugin, slot_type)) in &self.replication_slots {
write_str(&mut out, name);
write_str(&mut out, plugin);
write_str(&mut out, slot_type);
}
let crc = spg_crypto::crc32c::crc32c(&out);
write_u32(&mut out, crc);
out
}
pub fn deserialize(buf: &[u8]) -> Result<Self, StorageError> {
let mut cur = Cursor::new(buf);
let magic = cur.take(8)?;
if magic != FILE_MAGIC {
return Err(StorageError::Corrupt(format!(
"bad magic: expected SPGDB001, got {magic:?}"
)));
}
let version = cur.read_u8()?;
if !(MIN_SUPPORTED_FILE_VERSION..=FILE_VERSION).contains(&version) {
return Err(StorageError::Corrupt(format!(
"unsupported file version: {version} (supported: {MIN_SUPPORTED_FILE_VERSION}..={FILE_VERSION})"
)));
}
cur.codec_version = version;
let table_count = cur.read_u32()? as usize;
let mut cat = Self::new();
for _ in 0..table_count {
deserialize_table(&mut cur, &mut cat, version)?;
}
for (i, t) in cat.tables.iter_mut().enumerate() {
t.set_rel_id(row_header::RelId((i as u64) + 1));
}
cat.next_rel_id = cat.tables.len() as u64;
if version >= 22 {
let fn_count = cur.read_u32()? as usize;
for _ in 0..fn_count {
let name = cur.read_str()?;
let args_repr = cur.read_str()?;
let returns = cur.read_str()?;
let language = cur.read_str()?;
let body = cur.read_str_long()?;
let key = function_signature_key(&name, &args_repr);
cat.functions.insert(
key,
FunctionDef {
name,
args_repr,
returns,
language,
body,
owner: None,
acl: Vec::new(),
volatility: FN_VOLATILE,
strict: false,
security_definer: false,
leakproof: false,
parallel: FN_PARALLEL_UNSAFE,
cost: None,
rows: None,
},
);
}
let trg_count = cur.read_u32()? as usize;
for _ in 0..trg_count {
let name = cur.read_str()?;
let table = cur.read_str()?;
let timing = cur.read_str()?;
let ev_count = cur.read_u16()? as usize;
let mut events = Vec::with_capacity(ev_count);
for _ in 0..ev_count {
events.push(cur.read_str()?);
}
let for_each = cur.read_str()?;
let function = cur.read_str()?;
let update_columns = if version >= 23 {
let n = cur.read_u16()? as usize;
let mut cols = Vec::with_capacity(n);
for _ in 0..n {
cols.push(cur.read_str()?);
}
cols
} else {
Vec::new()
};
let enabled = if version >= 25 {
cur.read_u8()? != 0
} else {
true
};
let when_condition = if version >= 70 {
cur.read_str()?
} else {
String::new()
};
cat.triggers.push(TriggerDef {
name,
table,
timing,
events,
for_each,
function,
update_columns,
enabled,
when_condition,
});
}
}
if version >= 26 {
let seq_count = cur.read_u32()? as usize;
for _ in 0..seq_count {
let name = cur.read_str()?;
let data_type = match cur.read_u8()? {
0 => SequenceDataType::SmallInt,
1 => SequenceDataType::Int,
2 => SequenceDataType::BigInt,
other => {
return Err(StorageError::Corrupt(format!(
"unknown SEQUENCE data-type tag {other}"
)));
}
};
let start = cur.read_i64()?;
let increment = cur.read_i64()?;
let min_value = cur.read_i64()?;
let max_value = cur.read_i64()?;
let cache = cur.read_i64()?;
let cycle = cur.read_u8()? != 0;
let owned_by = match cur.read_u8()? {
0 => None,
1 => {
let t = cur.read_str()?;
let c = cur.read_str()?;
Some((t, c))
}
other => {
return Err(StorageError::Corrupt(format!(
"unknown SEQUENCE owned-by tag {other}"
)));
}
};
let last_value = cur.read_i64()?;
let is_called = cur.read_u8()? != 0;
cat.sequences.insert(
name.clone(),
SequenceDef {
name,
data_type,
start,
increment,
min_value,
max_value,
cache,
cycle,
owned_by,
last_value,
is_called,
owner: None,
acl: Vec::new(),
},
);
}
}
if version >= 27 {
let view_count = cur.read_u32()? as usize;
for _ in 0..view_count {
let name = cur.read_str()?;
let col_count = cur.read_u16()? as usize;
let mut columns = Vec::with_capacity(col_count);
for _ in 0..col_count {
columns.push(cur.read_str()?);
}
let body = cur.read_str_long()?;
let check_option = if version >= 69 { cur.read_u8()? } else { 0 };
cat.views.insert(
name.clone(),
ViewDef {
name,
columns,
body,
check_option,
},
);
}
}
if version >= 28 {
let mv_count = cur.read_u32()? as usize;
for _ in 0..mv_count {
let name = cur.read_str()?;
let body = cur.read_str_long()?;
cat.materialized_views.insert(name, body);
}
}
if version >= 29 {
let etype_count = cur.read_u32()? as usize;
for _ in 0..etype_count {
let name = cur.read_str()?;
let label_count = cur.read_u16()? as usize;
let mut labels = Vec::with_capacity(label_count);
for _ in 0..label_count {
labels.push(cur.read_str()?);
}
cat.enum_types
.insert(name.clone(), EnumDef { name, labels });
}
}
if version >= 30 {
let dtype_count = cur.read_u32()? as usize;
for _ in 0..dtype_count {
let name = cur.read_str()?;
let base_type = cur.read_data_type()?;
let nullable = cur.read_u8()? != 0;
let default = match cur.read_u8()? {
0 => None,
1 => Some(cur.read_str()?),
other => {
return Err(StorageError::Corrupt(format!(
"unknown DOMAIN default tag {other}"
)));
}
};
let check_count = cur.read_u16()? as usize;
let mut checks: Vec<DomainCheck> = Vec::with_capacity(check_count);
for i in 0..check_count {
let expr = cur.read_str()?;
let cname = if version >= 75 {
cur.read_str()?
} else if i == 0 {
alloc::format!("{name}_check")
} else {
alloc::format!("{name}_check{i}")
};
checks.push(DomainCheck { name: cname, expr });
}
let base_domain = if version >= 74 {
match cur.read_u8()? {
0 => None,
1 => Some(cur.read_str()?),
other => {
return Err(StorageError::Corrupt(alloc::format!(
"domain base_domain tag {other}"
)));
}
}
} else {
None
};
cat.domain_types.insert(
name.clone(),
DomainDef {
name,
base_type,
nullable,
default,
checks,
base_domain,
},
);
}
}
if version >= 31 {
let sch_count = cur.read_u32()? as usize;
for _ in 0..sch_count {
let name = cur.read_str()?;
cat.schemas.insert(name);
}
}
if version >= 52 {
let ctype_count = cur.read_u32()? as usize;
for _ in 0..ctype_count {
let name = cur.read_str()?;
let field_count = cur.read_u16()? as usize;
let mut fields = Vec::with_capacity(field_count);
let mut field_user_types: Vec<Option<String>> = Vec::with_capacity(field_count);
for _ in 0..field_count {
let fname = cur.read_str()?;
let fty = cur.read_data_type()?;
let ut = if version >= 76 {
match cur.read_u8()? {
0 => None,
1 => Some(cur.read_str()?),
other => {
return Err(StorageError::Corrupt(alloc::format!(
"composite field user-type tag {other}"
)));
}
}
} else {
None
};
fields.push((fname, fty));
field_user_types.push(ut);
}
cat.composite_types.insert(
name.clone(),
CompositeDef {
name,
fields,
field_user_types,
},
);
}
}
if version >= 61 {
let comment_count = cur.read_u32()? as usize;
for _ in 0..comment_count {
let key = cur.read_str()?;
let text = cur.read_str_long()?;
cat.comments.insert(key, text);
}
}
if version >= 66 {
let read_acl = |cur: &mut Cursor| -> Result<Vec<AclItem>, StorageError> {
let n = cur.read_u16()? as usize;
let mut acl = Vec::with_capacity(n);
for _ in 0..n {
let grantee = cur.read_str()?;
let privs = cur.read_u16()?;
let grantable = cur.read_u16()?;
let grantor = cur.read_str()?;
acl.push(AclItem {
grantee,
privs,
grantable,
grantor,
});
}
Ok(acl)
};
let seq_count = cur.read_u32()? as usize;
for _ in 0..seq_count {
let name = cur.read_str()?;
let owner = if cur.read_u8()? == 1 {
Some(cur.read_str()?)
} else {
None
};
let acl = read_acl(&mut cur)?;
if let Some(seq) = cat.sequences.get_mut(&name) {
seq.owner = owner;
seq.acl = acl;
}
}
cat.schema_acl = read_acl(&mut cur)?;
cat.database_acl = read_acl(&mut cur)?;
if version >= 67 {
let fn_count = cur.read_u32()? as usize;
for _ in 0..fn_count {
let name = cur.read_str()?;
let owner = if cur.read_u8()? == 1 {
Some(cur.read_str()?)
} else {
None
};
let acl = read_acl(&mut cur)?;
let target = resolve_stored_function_key(&cat.functions, &name);
if let Some(k) = target
&& let Some(f) = cat.functions.get_mut(&k)
{
f.owner = owner;
f.acl = acl;
}
}
}
}
if version >= 71 {
let rule_count = cur.read_u32()? as usize;
for _ in 0..rule_count {
let name = cur.read_str()?;
let table = cur.read_str()?;
let event = cur.read_str()?;
let instead = cur.read_u8()? != 0;
let when_condition = cur.read_str()?;
let cmd_count = cur.read_u16()? as usize;
let mut commands = Vec::with_capacity(cmd_count);
for _ in 0..cmd_count {
commands.push(cur.read_str()?);
}
cat.rules.push(RuleDef {
name,
table,
event,
instead,
when_condition,
commands,
});
}
}
if version >= 77 {
let count = cur.read_u32()? as usize;
for _ in 0..count {
let name = cur.read_str()?;
let table = cur.read_str()?;
let nk = cur.read_u16()? as usize;
let mut kinds = Vec::with_capacity(nk);
for _ in 0..nk {
kinds.push(cur.read_str()?);
}
let nc = cur.read_u16()? as usize;
let mut columns = Vec::with_capacity(nc);
for _ in 0..nc {
columns.push(cur.read_str()?);
}
cat.statistics_ext.push(StatisticsExtDef {
name,
table,
kinds,
columns,
});
}
}
if version >= 78 {
let count = cur.read_u32()? as usize;
for _ in 0..count {
let oid = cur.read_u32()?;
let len = cur.read_u32()? as usize;
let bytes = cur.read_bytes(len)?;
cat.large_objects.insert(oid, bytes);
}
}
if version >= 80 {
let count = cur.read_u32()? as usize;
for _ in 0..count {
let key = cur.read_str()?;
let volatility = cur.read_u8()?;
let flags = cur.read_u8()?;
let parallel = cur.read_u8()?;
let cost = f64::from_le_bytes(cur.read_bytes(8)?.try_into().unwrap_or([0; 8]));
let rows = f64::from_le_bytes(cur.read_bytes(8)?.try_into().unwrap_or([0; 8]));
if let Some(f) = cat.functions.get_mut(&key) {
f.volatility = volatility;
f.strict = flags & 1 != 0;
f.security_definer = flags & 2 != 0;
f.leakproof = flags & 4 != 0;
f.parallel = parallel;
f.cost = (!cost.is_nan()).then_some(cost);
f.rows = (!rows.is_nan()).then_some(rows);
}
}
}
if version >= 85 {
let scopes = cur.read_u32()? as usize;
for _ in 0..scopes {
let db = cur.read_str()?;
let role = cur.read_str()?;
let params = cur.read_u32()? as usize;
let mut m: BTreeMap<String, String> = BTreeMap::new();
for _ in 0..params {
let name = cur.read_str()?;
let value = cur.read_str()?;
m.insert(name, value);
}
if !m.is_empty() {
cat.db_role_settings.insert((db, role), m);
}
}
}
if version >= 86 {
let count = cur.read_u32()? as usize;
for _ in 0..count {
let name = cur.read_str()?;
let plugin = cur.read_str()?;
let slot_type = cur.read_str()?;
cat.replication_slots.insert(name, (plugin, slot_type));
}
}
if version >= FILE_VERSION_CRC_TRAILER {
let crc_start = cur.pos;
let stored = cur.read_u32()?;
let computed = spg_crypto::crc32c::crc32c(&buf[..crc_start]);
if computed != stored {
return Err(StorageError::Corrupt(format!(
"base snapshot CRC mismatch: computed {computed:#010x}, stored {stored:#010x}"
)));
}
}
if cur.pos < buf.len() {
return Err(StorageError::Corrupt(format!(
"trailing bytes: {} unread",
buf.len() - cur.pos
)));
}
Ok(cat)
}
}
#[cfg(test)]
mod tests;