import Base
import ./src/Keys.bend as Keys
import ./src/MemTable.bend as MemTable
import ./src/SortedRun.bend as SortedRun
import ./src/Sstable.bend as Sstable
import ./src/SstFile.bend as SstFile
import ./src/Wal.bend as Wal
import ./src/Manifest.bend as Manifest
import ./src/MergeIter.bend as MergeIter
import ./src/Db.bend as Db
import ./src/Recover.bend as Recover
import ./src/DurableDb.bend as DurableDb
import ./src/DurableError.bend as DurableError
import ./src/DurableDbPolicy.bend as DurableDbPolicy
import bend-kit-bytes@0.3.2.0/bytes.bend as Bytes
type InMemory.Session is Kind(a <&> &1):
InMemorySession{run: Db.Db -> (Db.Db & A)}
# --- Level 2: part control (1:1 delegates) ---
def cmp(+s1: String, +s2: String) -> Cmp:
Keys.cmp(s1, s2)
# Key equality: true exactly for identical strings.
def eq(+s1: String, +s2: String) -> Bool:
Keys.eq(s1, s2)
# Empty MemTable: no entries, count zero.
def mem_empty() -> MemTable.MemTable:
MemTable.empty()
# MemTable with one live version prepended (newest-first log).
def mem_put(tab: MemTable.MemTable, key: String, val: String) -> MemTable.MemTable:
MemTable.put(tab, key, val)
# MemTable with one tombstone prepended; hides older versions on read.
def mem_del(tab: MemTable.MemTable, key: String) -> MemTable.MemTable:
MemTable.del(tab, key)
# Newest-first read; a tombstone answers None and hides older versions.
def mem_get(tab: MemTable.MemTable, +key: String) -> Maybe<&2, String>:
MemTable.get(tab, key)
# Number of log entries, live and tombstoned.
def mem_count(tab: MemTable.MemTable) -> Nat:
MemTable.count(tab)
# Canonicalize newest-first entries into a strict ordered run.
def sort_newest(entries: List<&2, MemTable.Entry>) -> List<&2, MemTable.Entry>:
SortedRun.sort_newest(entries)
# Range scan over a merged run: lo <= key < hi, tombstones dropped.
def range_scan(+merged: List<&2, MemTable.Entry>, lo: String, hi: String) -> List<&2, MemTable.Entry>:
MergeIter.scan(merged, lo, hi)
# Build an SSTable, canonicalizing entries newest-first.
def sst_build(+entries: List<&2, MemTable.Entry>, level: Nat, est_keys: Nat) -> Sstable.Table:
Sstable.build(entries, level, est_keys)
# Build an SSTable trusting an already strict, unique run (no sorting).
def sst_from_sorted_unique(+entries: List<&2, MemTable.Entry>, level: Nat) -> Sstable.Table:
Sstable.from_sorted_unique(entries, level)
# Build an SSTable from a sorted run with an explicit key estimate.
def sst_build_sorted(+entries: List<&2, MemTable.Entry>, level: Nat, est_keys: Nat) -> Sstable.Table:
Sstable.build_sorted(entries, level, est_keys)
# Handle wal encode in the public module exports.
def wal_encode(batch: Wal.Batch) -> Result<&1, &1, Wal.Error, Bytes.Bytes>:
Wal.encode_frame(batch)
# Handle wal decode in the public module exports.
def wal_decode(encoded: Bytes.Bytes) -> Result<&1, &1, Wal.Error, Wal.Batch>:
Wal.decode_frame(encoded)
# Encode an SST as packed v3 bytes.
def sst_encode(
+entries: List<&2, MemTable.Entry>,
+level: U32
) -> Result<&1, &1, SstFile.Error, Bytes.Bytes>:
SstFile.encode_file(entries, level)
# Parse packed v3 bytes with strict bounds and checksum verification.
def sst_parse(
encoded: Bytes.Bytes
) -> Result<&1, &1, SstFile.Error, Sstable.Table>:
SstFile.parse(encoded)
# Serialize a Manifest to its exact on-disk bytes.
def mfst_serialize(mfst: Manifest.Manifest) -> Result<&1, &1, Manifest.Error, Bytes.Bytes>:
Manifest.serialize(mfst)
# Handle mfst parse in the public module exports.
def mfst_parse(encoded: Bytes.Bytes) -> Maybe<&2, Manifest.Manifest>:
Manifest.parse(encoded)
def InMemory.Session.pure(a, -A: Kind(a), value: A) -> InMemory.Session:
InMemorySession{database => (database, value)}
def InMemory.Session.run(a, -A: Kind(a), session: InMemory.Session, database: Db.Db) -> Db.Db & A:
match session:
case InMemorySession{run}: run(database)
def InMemory.Session.apply(a, -A: Kind(a), b, -B: Kind(b), pair: Db.Db & A, next: A -> InMemory.Session) -> Db.Db & B:
match pair:
case (database, value): InMemory.Session.run(b, B, next(value), database)
def InMemory.Session.bind(a, -A: Kind(a), -B: Kind(a), session: InMemory.Session, next: A -> InMemory.Session) -> InMemory.Session:
InMemorySession{database => InMemory.Session.apply(a, A, a, B, InMemory.Session.run(a, A, session, database), next)}
def InMemory.open(+path: String) -> Db.Db:
Db.open_db(path)
def InMemory.put_database(+database: Db.Db, +key: String, +value: String) -> Db.Db:
match database:
case Db.Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:
Db.apply_done(Db.apply_and_rotate(Con{Wal.Put{key, value}, Nil{}}, mem, frozen, mem_count, frozen_count), dir, batch_cap, bcache, levels, flushed, manifest_token)
def InMemory.delete_database(+database: Db.Db, +key: String) -> Db.Db:
match database:
case Db.Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:
Db.apply_done(Db.apply_and_rotate(Con{Wal.Del{key}, Nil{}}, mem, frozen, mem_count, frozen_count), dir, batch_cap, bcache, levels, flushed, manifest_token)
def InMemory.write_database_batch(+database: Db.Db, +mutations: List<&2, Wal.Mut>) -> Db.Db:
match database:
case Db.Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:
Db.apply_done(Db.apply_and_rotate(mutations, mem, frozen, mem_count, frozen_count), dir, batch_cap, bcache, levels, flushed, manifest_token)
def InMemory.read_database(+database: Db.Db, +key: String) -> Db.Db & Maybe<&2, String>:
Db.db_get_cached(database, key)
def InMemory.Session.put(+key: String, +value: String) -> InMemory.Session<&2, Unit>:
InMemorySession{database => (InMemory.put_database(database, key, value), Unit{})}
def InMemory.Session.delete(+key: String) -> InMemory.Session<&2, Unit>:
InMemorySession{database => (InMemory.delete_database(database, key), Unit{})}
def InMemory.Session.write_batch(+mutations: List<&2, Wal.Mut>) -> InMemory.Session<&2, Unit>:
InMemorySession{database => (InMemory.write_database_batch(database, mutations), Unit{})}
def InMemory.Session.get(+key: String) -> InMemory.Session<&2, Maybe<&2, String>>:
InMemorySession{database => InMemory.read_database(database, key)}
def InMemory.run_session(a, -A: Kind(a), database: Db.Db, session: InMemory.Session) -> Db.Db & A:
InMemory.Session.run(a, A, session, database)
def InMemory.value_of(a, -A: Kind(a), pair: Db.Db & A) -> A:
match pair:
case (database, value): value
# Recover a persisted database from its directory.
def open_recovering(+dir: String) -> IO(Result<&1, &1, U32 & String, Db.Db>):
Recover.open_db(dir)
# Lists the error categories exposed by the durable API.
type ErrorKind is Data:
NotFoundError{}
AlreadyExistsError{}
BusyError{}
InvalidArgumentError{}
UnsupportedFormatError{}
CorruptionError{}
ResourceLimitError{}
IoError{}
CommitUnknownError{}
# Carries structured error details for API callers.
type Error is Data:
Error{kind: ErrorKind, operation: String, path: String, host_code: U32, host_message: String}
# Owns the database and its process lock.
type Handle is Type:
Handle{inner: DurableDb.Handle}
# Carries failures from scoped session execution.
type SessionError is Data:
OpenFailure{error: Error}
OperationFailure{error: Error}
CloseFailure{error: Error}
OperationAndCloseFailure{operation_error: Error, close_error: Error}
# Holds durable operations that thread one owned handle.
type Session is Kind(a <&> &1):
Session{run: Handle -> IO(Handle & Result<&1, &1, Error, A>)}
# Reports database size, cache use, errors, and maintenance state.
type DurableStats is Data:
DurableStats{active_bytes: Nat, wal_bytes: Nat, memtable_entries: Nat, memtable_payload_bytes: Nat, cache_entries: Nat, operation_errors: Nat, maintenance_blocked: Bool}
# Reports whether post-commit maintenance completed.
type MaintenanceStatus is Data:
MaintenanceComplete{}
MaintenancePending{error: Error}
# Reports maintenance status for a confirmed write.
type WriteOutcome is Data:
WriteOutcome{maintenance: MaintenanceStatus}
# Describes a public key/value update or deletion.
type Mutation is Data:
SetValue{key: String, value: String}
DeleteKey{key: String}
# Converts a public mutation to the WAL representation.
def to_storage_mut(mutation: Mutation) -> Wal.Mut:
match mutation:
case SetValue{key, value}: Wal.Put{key, value}
case DeleteKey{key}: Wal.Del{key}
# Converts each public mutation to its WAL representation.
def to_storage_muts(+mutations: List<&2, Mutation>) -> List<&2, Wal.Mut>:
match mutations:
case Nil{}: Nil{}
case Con{mutation, tail}: Con{to_storage_mut(mutation), to_storage_muts(tail)}
# Maps internal error categories to the public error kind.
def public_error_kind(kind: DurableError.ErrorKind) -> ErrorKind:
match kind:
case DurableError.NotFound{}: NotFoundError{}
case DurableError.AlreadyExists{}: AlreadyExistsError{}
case DurableError.Busy{}: BusyError{}
case DurableError.InvalidArgument{}: InvalidArgumentError{}
case DurableError.UnsupportedFormat{}: UnsupportedFormatError{}
case DurableError.Corruption{}: CorruptionError{}
case DurableError.ResourceLimit{}: ResourceLimitError{}
case DurableError.Io{}: IoError{}
case DurableError.CommitUnknown{}: CommitUnknownError{}
# Builds the public error while preserving its structured details.
def public_error(error: DurableError.Error) -> Error:
match error:
case DurableError.Error{kind, operation, path, host_code, host_message}:
Error{public_error_kind(kind), operation, path, host_code, host_message}
# Converts an internal result to the public error type.
def public_error_result(result: Result<&1, &1, DurableError.Error, Unit>) -> Result<&1, &1, Error, Unit>:
match result:
case Fail{error}: Fail{public_error(error)}
case Done{value}: Done{value}
# Converts a durable open result to the public handle.
def public_open_result(result: Result<&1, &1, DurableError.Error, DurableDb.Handle>) -> Result<&1, &1, Error, Handle>:
match result:
case Fail{error}: Fail{public_error(error)}
case Done{handle}: Done{Handle{handle}}
# Creates a database and returns its owned durable handle.
def create(+dir: String) -> IO(Result<&1, &1, Error, Handle>):
do IO>:
result : Result<&1, &1, DurableError.Error, DurableDb.Handle> <- DurableDb.create(dir)
return public_open_result(result)
# Opens and recovers an existing database under its lock.
def open_existing(+dir: String) -> IO(Result<&1, &1, Error, Handle>):
do IO>:
result : Result<&1, &1, DurableError.Error, DurableDb.Handle> <- DurableDb.open_existing(dir)
return public_open_result(result)
# Converts an optional stored value and its error.
def public_maybe_result(
result: Result<&1, &1, DurableError.Error, Maybe<&2, String>>,
) -> Result<&1, &1, Error, Maybe<&2, String>>:
match result:
case Fail{error}: Fail{public_error(error)}
case Done{value}: Done{value}
# Converts a durable read result to the public value type.
def durable_get_result(
pair: DurableDb.Handle & Result<&1, &1, DurableError.Error, Maybe<&2, String>>,
) -> Handle & Result<&1, &1, Error, Maybe<&2, String>>:
match pair:
case (handle, result): (Handle{handle}, public_maybe_result(result))
# Reads a key from the owned durable database handle.
def get(handle: Handle, +key: String) -> IO(Handle & Result<&1, &1, Error, Maybe<&2, String>>):
match handle:
case Handle{inner}:
do IO>>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, Maybe<&2, String>> <- DurableDb.get(inner, key)
return durable_get_result(pair)
# Maps internal maintenance state to its public status.
def public_maintenance(maintenance: DurableDb.Maintenance) -> MaintenanceStatus:
match maintenance:
case DurableDb.Maintained{}: MaintenanceComplete{}
case DurableDb.Pending{error}: MaintenancePending{public_error(error)}
# Converts a confirmed write outcome to the public type.
def public_write_result(
result: Result<&1, &1, DurableError.Error, DurableDb.WriteOutcome>,
) -> Result<&1, &1, Error, WriteOutcome>:
match result:
case Fail{error}: Fail{public_error(error)}
case Done{DurableDb.WriteOutcome{maintenance}}: Done{WriteOutcome{public_maintenance(maintenance)}}
# Preserves maintenance status while converting write errors.
def durable_write_result(
pair: DurableDb.Handle & Result<&1, &1, DurableError.Error, Unit>,
) -> Handle & Result<&1, &1, Error, Unit>:
match pair:
case (handle, result): (Handle{handle}, public_error_result(result))
# Converts a batch outcome while preserving handle ownership.
def durable_batch_result(
pair: DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.WriteOutcome>,
) -> Handle & Result<&1, &1, Error, WriteOutcome>:
match pair:
case (handle, result): (Handle{handle}, public_write_result(result))
# Validates and commits a public mutation batch atomically.
def write_batch(handle: Handle, mutations: List<&2, Mutation>) -> IO(Handle & Result<&1, &1, Error, WriteOutcome>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.WriteOutcome> <- DurableDb.write_batch(inner, to_storage_muts(mutations))
return durable_batch_result(pair)
# Writes one key and value through the durable handle.
def put(handle: Handle, +key: String, +value: String) -> IO(Handle & Result<&1, &1, Error, WriteOutcome>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.WriteOutcome> <- DurableDb.put(inner, key, value)
return durable_batch_result(pair)
# Deletes one key through the durable handle.
def delete(handle: Handle, +key: String) -> IO(Handle & Result<&1, &1, Error, WriteOutcome>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.WriteOutcome> <- DurableDb.delete(inner, key)
return durable_batch_result(pair)
# Flushes the durable handle and returns maintenance status.
def flush(handle: Handle) -> IO(Handle & Result<&1, &1, Error, Unit>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, Unit> <- DurableDb.flush(inner)
return durable_write_result(pair)
# Compacts the durable handle and returns maintenance status.
def compact(handle: Handle) -> IO(Handle & Result<&1, &1, Error, Unit>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, Unit> <- DurableDb.compact(inner)
return durable_write_result(pair)
# Maps internal counters to the public stats record.
def public_stats(result: Result<&1, &1, DurableError.Error, DurableDb.Stats>) -> Result<&1, &1, Error, DurableStats>:
match result:
case Fail{error}: Fail{public_error(error)}
case Done{DurableDb.Stats{active_bytes, wal_bytes, memtable_entries, memtable_payload_bytes, cache_entries, operation_errors, state}}:
match state:
case DurableDbPolicy.MaintenanceBlocked{}: Done{DurableStats{active_bytes, wal_bytes, memtable_entries, memtable_payload_bytes, cache_entries, operation_errors, True{}}}
case _: Done{DurableStats{active_bytes, wal_bytes, memtable_entries, memtable_payload_bytes, cache_entries, operation_errors, False{}}}
# Converts durable stats errors to the public error type.
def durable_stats_result(
pair: DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.Stats>,
) -> Handle & Result<&1, &1, Error, DurableStats>:
match pair:
case (handle, result): (Handle{handle}, public_stats(result))
# Reads bounded operational stats from the durable handle.
def stats(handle: Handle) -> IO(Handle & Result<&1, &1, Error, DurableStats>):
match handle:
case Handle{inner}:
do IO>:
pair : DurableDb.Handle & Result<&1, &1, DurableError.Error, DurableDb.Stats> <- DurableDb.stats(inner)
return durable_stats_result(pair)
# Closes the durable handle and releases its database lock.
def close(handle: Handle) -> IO(Result<&1, &1, Error, Unit>):
match handle:
case Handle{inner}:
do IO>:
result : Result<&1, &1, DurableError.Error, Unit> <- DurableDb.close(inner)
return public_error_result(result)
def Session.pure(a, -A: Kind(a), value: A) -> Session:
Session{handle => IO.pure(Handle & Result<&1, &1, Error, A>, (handle, Done{value}))}
def Session.run(a, -A: Kind(a), session: Session, handle: Handle) -> IO(Handle & Result<&1, &1, Error, A>):
match session:
case Session{run}: run(handle)
def Session.bind_result(a, -A: Kind(a), -B: Kind(a), pair: Handle & Result<&1, &1, Error, A>, next: A -> Session) -> IO(Handle & Result<&1, &1, Error, B>):
match pair:
case (handle, Fail{error}):
IO.pure(Handle & Result<&1, &1, Error, B>, (handle, Fail{error}))
case (handle, Done{value}):
Session.run(a, B, next(value), handle)
def Session.bind(a, -A: Kind(a), -B: Kind(a), session: Session, next: A -> Session) -> Session:
Session{handle =>
do IO>:
pair : Handle & Result<&1, &1, Error, A> <- Session.run(a, A, session, handle)
Session.bind_result(a, A, B, pair, next)
}
def Session.put(+key: String, +value: String) -> Session<&2, WriteOutcome>:
Session{handle => put(handle, key, value)}
def Session.delete(+key: String) -> Session<&2, WriteOutcome>:
Session{handle => delete(handle, key)}
def Session.write_batch(mutations: List<&2, Mutation>) -> Session<&2, WriteOutcome>:
Session{handle => write_batch(handle, mutations)}
def Session.get(+key: String) -> Session<&2, Maybe<&2, String>>:
Session{handle => get(handle, key)}
def Session.flush() -> Session<&2, Unit>:
Session{handle => flush(handle)}
def Session.compact() -> Session<&2, Unit>:
Session{handle => compact(handle)}
def Session.stats() -> Session<&2, DurableStats>:
Session{handle => stats(handle)}
def Session.finish_result(a, -A: Kind(a), result: Result<&1, &1, Error, A>, close_result: Result<&1, &1, Error, Unit>) -> Result<&1, &1, SessionError, A>:
match result close_result:
case Done{value} Done{Unit{}}: Done{value}
case Fail{error} Done{Unit{}}: Fail{OperationFailure{error}}
case Done{_} Fail{error}: Fail{CloseFailure{error}}
case Fail{operation_error} Fail{close_error}:
Fail{OperationAndCloseFailure{operation_error, close_error}}
def Session.finish(a, -A: Kind(a), handle: Handle, result: Result<&1, &1, Error, A>) -> IO(Result<&1, &1, SessionError, A>):
do IO>:
close_result : Result<&1, &1, Error, Unit> <- close(handle)
return Session.finish_result(a, A, result, close_result)
def Session.finish_pair(a, -A: Kind(a), pair: Handle & Result<&1, &1, Error, A>) -> IO(Result<&1, &1, SessionError, A>):
match pair:
case (handle, result): Session.finish(a, A, handle, result)
def Session.run_opened(a, -A: Kind(a), opened: Result<&1, &1, Error, Handle>, session: Session) -> IO(Result<&1, &1, SessionError, A>):
match opened:
case Fail{error}:
IO.pure(Result<&1, &1, SessionError, A>, Fail{OpenFailure{error}})
case Done{handle}:
do IO>:
pair : Handle & Result<&1, &1, Error, A> <- Session.run(a, A, session, handle)
Session.finish_pair(a, A, pair)
# Runs a program and returns its still-owned handle.
def run_session(a, -A: Kind(a), handle: Handle, session: Session) -> IO(Handle & Result<&1, &1, Error, A>):
Session.run(a, A, session, handle)
# Creates, runs, and closes a scoped durable session.
def create_database(a, -A: Kind(a), +path: String, session: Session) -> IO(Result<&1, &1, SessionError, A>):
do IO>:
opened : Result<&1, &1, Error, Handle> <- create(path)
Session.run_opened(a, A, opened, session)
# Opens, runs, and closes a scoped durable session.
def open_database(a, -A: Kind(a), +path: String, session: Session) -> IO(Result<&1, &1, SessionError, A>):
do IO>:
opened : Result<&1, &1, Error, Handle> <- open_existing(path)
Session.run_opened(a, A, opened, session)
# Checks whether the public error means the path exists.
def error_is_already_exists(error: Error) -> Bool:
match error:
case Error{AlreadyExistsError{}, _, _, _, _}: True{}
case _: False{}
# Checks whether the public error means the database is locked.
def error_is_busy(error: Error) -> Bool:
match error:
case Error{BusyError{}, _, _, _, _}: True{}
case _: False{}