use std::panic::{catch_unwind, AssertUnwindSafe};
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
pub use powdb_query::ast::ParamValue;
pub use powdb_query::executor::{Engine, WalSyncMode};
pub use powdb_query::result::{QueryError, QueryResult};
pub use powdb_storage::types::{TypeId, Value};
pub use powdb_sync::RETAINED_SEGMENT_FORMAT_VERSION;
pub use powdb_storage::pj1::pj1_to_text;
fn archive_wal_records_if_sync_enabled(
data_dir: &Path,
records: &[powdb_storage::wal::WalRecord],
) -> std::io::Result<()> {
match powdb_sync::read_identity(data_dir) {
Ok(identity) => powdb_sync::archive_wal_records_for_identity(data_dir, identity, records),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(err) => Err(err),
}
}
#[derive(Debug)]
pub enum Error {
Open(std::io::Error),
Query(QueryError),
Poisoned,
OpenPanicked,
InvalidArgument(String),
Sync(std::io::Error),
}
impl std::fmt::Display for Error {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Error::Open(e) => write!(f, "failed to open database: {e}"),
Error::Query(e) => write!(f, "query failed: {e}"),
Error::Poisoned => write!(
f,
"database handle is poisoned (a previous call panicked); reopen the database"
),
Error::OpenPanicked => write!(
f,
"opening the database panicked (data directory may be corrupt); restore from a backup"
),
Error::InvalidArgument(msg) => write!(f, "{msg}"),
Error::Sync(e) => write!(f, "sync apply failed: {e}"),
}
}
}
pub fn parse_sync_mode(mode: &str) -> Option<WalSyncMode> {
match mode.to_ascii_lowercase().as_str() {
"full" => Some(WalSyncMode::Full),
"normal" => Some(WalSyncMode::Normal),
"off" => Some(WalSyncMode::Off),
_ => None,
}
}
fn values_to_params(params: &[Value]) -> Result<Vec<ParamValue>, Error> {
params
.iter()
.map(|value| match value {
Value::Int(v) => Ok(ParamValue::Int(*v)),
Value::Float(v) => Ok(ParamValue::Float(*v)),
Value::Bool(v) => Ok(ParamValue::Bool(*v)),
Value::Str(v) => Ok(ParamValue::Str(v.clone())),
Value::Empty => Ok(ParamValue::Null),
other => Err(Error::InvalidArgument(format!(
"cannot bind a {:?} value as a query parameter; supported parameter types are int, float, bool, str, and null",
other.type_id()
))),
})
.collect()
}
impl std::error::Error for Error {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Error::Open(e) => Some(e),
Error::Sync(e) => Some(e),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub struct SyncApplyIdentity {
pub database_id: [u8; 16],
pub primary_generation: u64,
pub wal_format_version: u16,
pub catalog_version: u16,
pub segment_format_version: u16,
}
#[derive(Debug, Clone)]
pub struct RetainedUnitInput {
pub tx_id: u64,
pub record_type: u8,
pub lsn: u64,
pub data: Vec<u8>,
}
#[derive(Debug, Clone)]
pub struct RetainedApplyRequest {
pub since_lsn: u64,
pub identity: SyncApplyIdentity,
pub units: Vec<RetainedUnitInput>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RetainedApplyResult {
pub through_lsn: u64,
pub units_applied: usize,
}
pub struct Database {
engine: Option<Engine>,
poisoned: AtomicBool,
}
impl Database {
pub fn open(dir: impl AsRef<Path>) -> Result<Self, Error> {
let dir = dir.as_ref();
Self::wrap_open(catch_unwind(AssertUnwindSafe(|| {
Engine::new_with_wal_archive(dir, archive_wal_records_if_sync_enabled)
})))
}
pub fn open_read_only(dir: impl AsRef<Path>) -> Result<Self, Error> {
let dir = dir.as_ref();
Self::wrap_open(catch_unwind(AssertUnwindSafe(|| {
Engine::open_read_only(dir)
})))
}
pub fn open_read_only_with_memory_limit(
dir: impl AsRef<Path>,
limit_bytes: usize,
) -> Result<Self, Error> {
let dir = dir.as_ref();
Self::wrap_open(catch_unwind(AssertUnwindSafe(|| {
Engine::open_read_only_with_memory_limit(dir, limit_bytes)
})))
}
pub fn open_with_memory_limit(
dir: impl AsRef<Path>,
limit_bytes: usize,
) -> Result<Self, Error> {
let dir = dir.as_ref();
Self::wrap_open(catch_unwind(AssertUnwindSafe(|| {
Engine::with_memory_limit_and_wal_archive(
dir,
limit_bytes,
archive_wal_records_if_sync_enabled,
)
})))
}
fn wrap_open(result: std::thread::Result<std::io::Result<Engine>>) -> Result<Self, Error> {
let engine = result
.map_err(|_| Error::OpenPanicked)?
.map_err(Error::Open)?;
Ok(Self {
engine: Some(engine),
poisoned: AtomicBool::new(false),
})
}
pub fn query(&mut self, powql: &str) -> Result<QueryResult, Error> {
self.run_mut(|e| e.execute_powql(powql))
}
pub fn query_sql(&mut self, sql: &str) -> Result<QueryResult, Error> {
self.run_mut(|e| e.execute_sql(sql))
}
pub fn query_readonly(&self, powql: &str) -> Result<QueryResult, Error> {
match self.run_ref(|e| e.execute_powql_readonly(powql)) {
Err(Error::Query(QueryError::ReadonlyNeedsWrite)) => Err(Error::InvalidArgument(
"statement would mutate the database; use query() instead of query_readonly()"
.to_string(),
)),
other => other,
}
}
pub fn query_with_params(
&mut self,
powql: &str,
params: &[Value],
) -> Result<QueryResult, Error> {
let bound = values_to_params(params)?;
self.run_mut(|e| e.execute_powql_with_params(powql, &bound))
}
pub fn query_readonly_with_params(
&self,
powql: &str,
params: &[Value],
) -> Result<QueryResult, Error> {
let bound = values_to_params(params)?;
match self.run_ref(|e| e.execute_powql_readonly_with_params(powql, &bound)) {
Err(Error::Query(QueryError::ReadonlyNeedsWrite)) => Err(Error::InvalidArgument(
"statement would mutate the database; use query_with_params() instead of query_readonly_with_params()"
.to_string(),
)),
other => other,
}
}
pub fn set_sync_mode(&mut self, mode: WalSyncMode) {
if let Some(engine) = self.engine.as_mut() {
engine.set_wal_sync_mode(mode);
}
}
pub fn set_sync_mode_str(&mut self, mode: &str) -> Result<(), Error> {
let parsed = parse_sync_mode(mode).ok_or_else(|| {
Error::InvalidArgument(format!(
"unknown sync mode {mode:?}; expected \"full\", \"normal\", or \"off\""
))
})?;
self.set_sync_mode(parsed);
Ok(())
}
pub fn is_poisoned(&self) -> bool {
self.poisoned.load(Ordering::Acquire)
}
pub fn apply_retained_units(
&mut self,
request: RetainedApplyRequest,
) -> Result<RetainedApplyResult, Error> {
if request.identity.segment_format_version != RETAINED_SEGMENT_FORMAT_VERSION {
return Err(Error::InvalidArgument(format!(
"unsupported retained segment format {}; expected {}",
request.identity.segment_format_version, RETAINED_SEGMENT_FORMAT_VERSION
)));
}
let identity = powdb_sync::SegmentIdentity {
database_id: request.identity.database_id,
primary_generation: request.identity.primary_generation,
wal_format_version: request.identity.wal_format_version,
catalog_version: request.identity.catalog_version,
};
let units = request
.units
.into_iter()
.map(|unit| powdb_sync::RetainedUnit {
tx_id: unit.tx_id,
record_type: unit.record_type,
lsn: unit.lsn,
data: unit.data,
})
.collect::<Vec<_>>();
self.run_sync_mut(|engine| {
let summary = powdb_sync::apply_retained_units_chunk(
engine.catalog_mut(),
identity,
request.since_lsn,
&units,
)?;
Ok(RetainedApplyResult {
through_lsn: summary.through_lsn,
units_applied: summary.units_applied,
})
})
}
pub fn close(self) {
}
fn run_mut(
&mut self,
f: impl FnOnce(&mut Engine) -> Result<QueryResult, QueryError>,
) -> Result<QueryResult, Error> {
if self.poisoned.load(Ordering::Acquire) {
return Err(Error::Poisoned);
}
let engine = self.engine.as_mut().expect("engine present until drop");
match catch_unwind(AssertUnwindSafe(|| f(engine))) {
Ok(inner) => inner.map_err(Error::Query),
Err(_) => {
self.poisoned.store(true, Ordering::Release);
Err(Error::Poisoned)
}
}
}
fn run_ref(
&self,
f: impl FnOnce(&Engine) -> Result<QueryResult, QueryError>,
) -> Result<QueryResult, Error> {
if self.poisoned.load(Ordering::Acquire) {
return Err(Error::Poisoned);
}
let engine = self.engine.as_ref().expect("engine present until drop");
match catch_unwind(AssertUnwindSafe(|| f(engine))) {
Ok(inner) => inner.map_err(Error::Query),
Err(_) => {
self.poisoned.store(true, Ordering::Release);
Err(Error::Poisoned)
}
}
}
fn run_sync_mut<T>(
&mut self,
f: impl FnOnce(&mut Engine) -> std::io::Result<T>,
) -> Result<T, Error> {
if self.poisoned.load(Ordering::Acquire) {
return Err(Error::Poisoned);
}
let engine = self.engine.as_mut().expect("engine present until drop");
match catch_unwind(AssertUnwindSafe(|| f(engine))) {
Ok(inner) => inner.map_err(Error::Sync),
Err(_) => {
self.poisoned.store(true, Ordering::Release);
Err(Error::Poisoned)
}
}
}
}
impl Drop for Database {
fn drop(&mut self) {
if let Some(engine) = self.engine.take() {
if self.poisoned.load(Ordering::Acquire) {
std::mem::forget(engine);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn caught_panic_poisons_and_returns_error() {
let poisoned = AtomicBool::new(false);
let result: Result<(), Error> = (|| {
if poisoned.load(Ordering::Acquire) {
return Err(Error::Poisoned);
}
match catch_unwind(AssertUnwindSafe(|| panic!("boom"))) {
Ok(()) => Ok(()),
Err(_) => {
poisoned.store(true, Ordering::Release);
Err(Error::Poisoned)
}
}
})();
assert!(matches!(result, Err(Error::Poisoned)));
assert!(poisoned.load(Ordering::Acquire));
}
#[test]
fn poisoned_handle_rejects_then_reopen_recovers() {
let dir = std::env::temp_dir().join(format!("powdb_facade_poison_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
{
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
db.query("insert T { id := 1 }").unwrap(); db.poisoned.store(true, Ordering::Release);
assert!(db.is_poisoned());
assert!(matches!(db.query("count(T)"), Err(Error::Poisoned)));
}
let mut db = Database::open(&dir).unwrap();
match db.query("count(T)").unwrap() {
QueryResult::Scalar(Value::Int(n)) => assert_eq!(n, 1),
other => panic!("expected 1 row after poisoned-drop reopen, got {other:?}"),
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_error_displays_human_message_not_debug() {
let dir = std::env::temp_dir().join(format!("powdb_facade_disp_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
let err = db.query("Missing filter .x > 1").unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("table 'Missing' not found"),
"expected human Display, got {msg:?}"
);
assert!(!msg.contains("TableNotFound"), "leaked Debug: {msg:?}");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_readonly_rejects_mutation_with_actionable_message() {
let dir = std::env::temp_dir().join(format!("powdb_facade_ro_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
let err = db.query_readonly("insert T { id := 1 }").unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("use query()"),
"expected actionable hint, got {msg:?}"
);
assert!(
!msg.contains("__POWDB_READONLY_NEEDS_WRITE__"),
"internal sentinel leaked to caller: {msg:?}"
);
assert!(matches!(err, Error::InvalidArgument(_)));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn parse_sync_mode_is_case_insensitive_and_rejects_unknown() {
assert!(matches!(parse_sync_mode("full"), Some(WalSyncMode::Full)));
assert!(matches!(
parse_sync_mode("Normal"),
Some(WalSyncMode::Normal)
));
assert!(matches!(parse_sync_mode("OFF"), Some(WalSyncMode::Off)));
assert!(parse_sync_mode("bogus").is_none());
}
#[test]
fn set_sync_mode_str_sets_known_and_errors_on_unknown() {
let dir = std::env::temp_dir().join(format!("powdb_syncstr_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
db.set_sync_mode_str("normal").unwrap();
db.query("insert T { id := 1 }").unwrap();
assert!(db.set_sync_mode_str("bogus").is_err());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn apply_retained_units_noop_uses_seeded_sync_boundary() {
let dir = std::env::temp_dir().join(format!("powdb_facade_apply_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let database_id = [7u8; 16];
seed_apply_boundary(&dir, database_id, 1, 0);
let mut db = Database::open(&dir).unwrap();
let result = db
.apply_retained_units(RetainedApplyRequest {
since_lsn: 0,
identity: SyncApplyIdentity {
database_id,
primary_generation: 1,
wal_format_version: powdb_storage::wal::WAL_FORMAT_VERSION,
catalog_version: powdb_storage::catalog::CATALOG_VERSION,
segment_format_version: RETAINED_SEGMENT_FORMAT_VERSION,
},
units: Vec::new(),
})
.unwrap();
assert_eq!(
result,
RetainedApplyResult {
through_lsn: 0,
units_applied: 0
}
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn apply_retained_units_rejects_wrong_segment_format() {
let dir =
std::env::temp_dir().join(format!("powdb_facade_apply_badfmt_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
let err = db
.apply_retained_units(RetainedApplyRequest {
since_lsn: 0,
identity: SyncApplyIdentity {
database_id: [7u8; 16],
primary_generation: 1,
wal_format_version: powdb_storage::wal::WAL_FORMAT_VERSION,
catalog_version: powdb_storage::catalog::CATALOG_VERSION,
segment_format_version: RETAINED_SEGMENT_FORMAT_VERSION + 1,
},
units: Vec::new(),
})
.unwrap_err();
assert!(matches!(err, Error::InvalidArgument(_)));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn open_archives_pending_sync_wal_before_recovery_truncates() {
let dir =
std::env::temp_dir().join(format!("powdb_facade_sync_recovery_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let identity = {
let mut catalog = powdb_storage::catalog::Catalog::create(&dir).unwrap();
catalog
.create_table(powdb_storage::types::Schema {
table_name: "T".into(),
columns: vec![powdb_storage::types::ColumnDef {
name: "id".into(),
type_id: powdb_storage::types::TypeId::Int,
required: true,
position: 0,
}],
})
.unwrap();
catalog.insert("T", &vec![Value::Int(1)]).unwrap();
catalog.sync_wal().unwrap();
powdb_sync::open_or_create_identity(&dir).unwrap()
};
let db = Database::open(&dir).unwrap();
drop(db);
let units = powdb_sync::read_units_since(
&powdb_sync::retained_segments_dir(&dir),
identity.segment_identity(),
0,
100,
)
.unwrap();
assert!(
!units.is_empty(),
"embedded facade open must preserve replayed sync WAL"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn open_with_corrupt_heap_returns_error_not_panic() {
let dir = std::env::temp_dir().join(format!("powdb_corrupt_heap_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
{
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
db.query("insert T { id := 1 }").unwrap();
}
let heap = dir.join("T.heap");
let mut bytes = std::fs::read(&heap).unwrap();
for b in bytes.iter_mut().take(20) {
*b = 0xFF;
}
std::fs::write(&heap, &bytes).unwrap();
let result = Database::open(&dir);
let _ = std::fs::remove_dir_all(&dir);
assert!(
result.is_err(),
"corrupt heap must return Err, not panic/abort the host"
);
}
#[test]
fn query_returns_lossless_typed_values() {
let dir = std::env::temp_dir().join(format!("powdb_facade_typed_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type Doc { required id: int, big: int, body: json }")
.unwrap();
db.query("insert Doc { id := 1, big := 9007199254740993, body := \"null\" }")
.unwrap();
db.query("insert Doc { id := 2, big := 9223372036854775807 }")
.unwrap();
match db.query("Doc { id, big, body }").unwrap() {
QueryResult::Rows { rows, .. } => {
assert_eq!(rows.len(), 2);
assert_eq!(rows[0][1], Value::Int(9_007_199_254_740_993));
assert_eq!(rows[1][1], Value::Int(9_223_372_036_854_775_807));
assert!(matches!(rows[0][2], Value::Json(_)));
assert_eq!(rows[1][2], Value::Empty);
assert_ne!(rows[0][2], rows[1][2]);
}
other => panic!("expected rows, got {other:?}"),
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_with_params_binds_positional_values() {
let dir = std::env::temp_dir().join(format!("powdb_facade_params_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type User { required name: str, age: int }")
.unwrap();
db.query(r#"insert User { name := "Ada", age := 36 }"#)
.unwrap();
db.query(r#"insert User { name := "Bo", age := 20 }"#)
.unwrap();
match db
.query_with_params(
"User filter .name = $1 and .age > $2 { .name, .age }",
&[Value::Str("Ada".into()), Value::Int(30)],
)
.unwrap()
{
QueryResult::Rows { rows, .. } => {
assert_eq!(rows.len(), 1);
assert_eq!(rows[0][0], Value::Str("Ada".into()));
assert_eq!(rows[0][1], Value::Int(36));
}
other => panic!("expected rows, got {other:?}"),
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_with_params_rejects_unbindable_shapes() {
let dir =
std::env::temp_dir().join(format!("powdb_facade_badparam_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
let err = db
.query_with_params("T filter .id = $1 { .id }", &[Value::Bytes(vec![1, 2, 3])])
.unwrap_err();
assert!(matches!(err, Error::InvalidArgument(_)));
assert!(
err.to_string().contains("cannot bind"),
"expected actionable message, got {err}"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_readonly_with_params_rejects_mutation() {
let dir = std::env::temp_dir().join(format!("powdb_facade_roparam_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
let err = db
.query_readonly_with_params("insert T { id := $1 }", &[Value::Int(1)])
.unwrap_err();
assert!(matches!(err, Error::InvalidArgument(_)));
assert!(
!err.to_string().contains("__POWDB_READONLY_NEEDS_WRITE__"),
"internal sentinel leaked: {err}"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn query_readonly_with_params_reads_under_shared_borrow() {
let dir = std::env::temp_dir().join(format!("powdb_facade_rords_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let mut db = Database::open(&dir).unwrap();
db.query("type T { required id: int }").unwrap();
db.query("insert T { id := 7 }").unwrap();
match db
.query_readonly_with_params("T filter .id = $1 { .id }", &[Value::Int(7)])
.unwrap()
{
QueryResult::Rows { rows, .. } => {
assert_eq!(rows, vec![vec![Value::Int(7)]]);
}
other => panic!("expected rows, got {other:?}"),
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn open_read_only_serves_reads_and_rejects_mutations() {
let dir = std::env::temp_dir().join(format!("powdb_facade_ro_open_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
{
let mut db = Database::open(&dir).unwrap();
db.query("type User { required name: str, age: int }")
.unwrap();
db.query(r#"insert User { name := "Ada", age := 36 }"#)
.unwrap();
}
let mut db = Database::open_read_only(&dir).unwrap();
match db.query("count(User)").unwrap() {
QueryResult::Scalar(Value::Int(n)) => assert_eq!(n, 1),
other => panic!("expected scalar 1, got {other:?}"),
}
match db
.query_readonly("User filter .age = 36 { .name }")
.unwrap()
{
QueryResult::Rows { rows, .. } => {
assert_eq!(rows[0][0], Value::Str("Ada".into()))
}
other => panic!("expected rows, got {other:?}"),
}
let err = db
.query(r#"insert User { name := "Bo", age := 20 }"#)
.unwrap_err();
assert!(
err.to_string().contains("readonly mode"),
"expected a read-only error, got {err}"
);
assert!(!err.to_string().contains("__POWDB_READONLY_NEEDS_WRITE__"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn open_read_only_refuses_non_empty_wal() {
let dir =
std::env::temp_dir().join(format!("powdb_facade_ro_dirtywal_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
{
let mut catalog = powdb_storage::catalog::Catalog::create(&dir).unwrap();
catalog
.create_table(powdb_storage::types::Schema {
table_name: "T".into(),
columns: vec![powdb_storage::types::ColumnDef {
name: "id".into(),
type_id: powdb_storage::types::TypeId::Int,
required: true,
position: 0,
}],
})
.unwrap();
catalog.insert("T", &vec![Value::Int(1)]).unwrap();
catalog.sync_wal().unwrap();
std::mem::forget(catalog);
}
let err = match Database::open_read_only(&dir) {
Ok(_) => panic!("read-only open must refuse a non-empty WAL"),
Err(err) => err,
};
assert!(
err.to_string().contains("WAL is not empty"),
"expected WAL-not-empty refusal, got {err}"
);
let _ = std::fs::remove_dir_all(&dir);
}
fn seed_apply_boundary(
dir: &std::path::Path,
database_id: [u8; 16],
generation: u64,
lsn: u64,
) {
let sync_dir = dir.join(".powdb-sync");
std::fs::create_dir_all(&sync_dir).unwrap();
let database_id_hex = encode_hex_16(database_id);
std::fs::write(
sync_dir.join("identity.json"),
format!(
r#"{{"format_version":1,"database_id":"{database_id_hex}","primary_generation":{generation},"created_unix_secs":1}}"#
),
)
.unwrap();
std::fs::write(
sync_dir.join("apply-state.json"),
format!(
r#"{{"format_version":1,"database_id":{},"primary_generation":{generation},"wal_format_version":{},"catalog_version":{},"from_lsn":{lsn},"through_lsn":{lsn},"applied_lsn":{lsn},"status":"complete","started_unix_secs":1,"updated_unix_secs":1}}"#,
json_u8_array(database_id),
powdb_storage::wal::WAL_FORMAT_VERSION,
powdb_storage::catalog::CATALOG_VERSION,
),
)
.unwrap();
}
fn encode_hex_16(bytes: [u8; 16]) -> String {
bytes.iter().map(|byte| format!("{byte:02x}")).collect()
}
fn json_u8_array(bytes: [u8; 16]) -> String {
let body = bytes
.iter()
.map(u8::to_string)
.collect::<Vec<_>>()
.join(",");
format!("[{body}]")
}
}