src/Store.bend source
src/Store.bend on the hub · documented module
import Baseimport mylsm-lsm-store@0.3.2.0/mylsm.bend as MyLsmStoreimport mylsm-lsm-store@0.3.2.0/src/Db.bend as DbTypesimport mylsm-lsm-store@0.3.2.0/src/Wal.bend as WalTypes# Session effects for ber-core. Op is a newtype over Db transitions and# Handle over MyLSM Db; do-blocks desugar through Op.bind/Op.pure exactly# like Sess. Pairs sit only in return and parameter positions (never as a# Sess answer), so Op stays kind-polymorphic; computed pairs thread through# helpers (match on parameters, never on computed values). Defs are ordered# leaf-first: no def calls one written below it.## Every op returns its answer with a journal of physical mutations# (execution order). Pure run_op drops it; run_op_tracked exposes it and# run_op_durable appends it through one WAL frame plus fsync. oput/odel# journal exactly what they apply, so the WAL mirrors the memtable.# Journal rope: O(1) per-bind combination (no calls, just nesting); flattened# once at flush with an explicit work stack (constant native stack) into# chronological order. Reads build Empty spines that flatten to Nil.type Journal is Data: JEmpty{} JOne{mut: WalTypes.Mut} JMore{prior: Journal, post: Journal}# Counts rope nodes structurally; the work loop below visits each once.def rope_nodes(+journal: Journal) -> Nat: match journal: case JEmpty{}: 1n case JOne{unlogged_mut}: 1n case JMore{prior, post}: Nat.add(1n, Nat.add(rope_nodes(prior), rope_nodes(post)))# Flattens a rope oldest-first with an explicit work stack; the accumulator# holds the reversed prefix. Fuel covers nodes plus the final empty check.def flatten_go(fuel: Nat, +work: List<&2, Journal>, +acc: List<&2, WalTypes.Mut>) -> List<&2, WalTypes.Mut>: match fuel: case 0n: acc case 1n+fuel_left: match work: case Nil{}: List.reverse(&2, WalTypes.Mut, acc) case Con{head, tail}: match head: case JEmpty{}: flatten_go(fuel_left, tail, acc) case JOne{mut}: flatten_go(fuel_left, tail, Con{mut, acc}) case JMore{prior, post}: flatten_go(fuel_left, Con{prior, Con{post, tail}}, acc)# Materializes a journal rope into its chronological mutation list.def flatten_journal(+journal: Journal) -> List<&2, WalTypes.Mut>: flatten_go(Nat.add(rope_nodes(journal), 1n), Con{journal, Nil{}}, Nil{})# One session step chain: a Db transition answering with its journal.type Op<a, -A: Kind(a)> is Kind(a <&> &1): Op{run: DbTypes.Db -> (DbTypes.Db & A) & Journal}# An opened store handle; affine, threaded through run_op.type Handle is Type: Handle{db: DbTypes.Db}# Unwraps the inner transition for translators like bind and run_op.def Op.inner(b, -B: Kind(b), op: Op<b, B>) -> DbTypes.Db -> (DbTypes.Db & B) & Journal: match op: case Op{run_fn}: run_fn# Lifts a pure value with an empty journal.def Op.pure(a, -A: Kind(a), val: A) -> Op<a, A>: Op{st => ((st, val), JEmpty{})}# Concatenates a prior journal before a later one, finishing the second step.def Op.seq2(a, -B: Kind(a), +prior: Journal, second: (DbTypes.Db & B) & Journal) -> (DbTypes.Db & B) & Journal: match second: case ((db_next, answer), journal): ((db_next, answer), JMore{prior, journal})# Runs the next op on the first answer, threading Db and journal.def Op.seq(a, -A: Kind(a), -B: Kind(a), first: (DbTypes.Db & A) & Journal, fun: A -> Op<a, B>) -> (DbTypes.Db & B) & Journal: match first: case ((db_mid, answer), prior): Op.seq2(a, B, prior, Op.inner(a, B, fun(answer))(db_mid))# Sequences two ops, concatenating journals in execution order.def Op.bind(a, -A: Kind(a), -B: Kind(a), op: Op<a, A>, fun: A -> Op<a, B>) -> Op<a, B>: match op: case Op{run1}: Op{st => Op.seq(a, A, B, run1(st), fun)}# Pairs a finished point write with its journaled mutation.def pair_put(+key: String, +val: String, db_next: DbTypes.Db) -> (DbTypes.Db & Unit) & Journal: ((db_next, Unit{}), JOne{WalTypes.Put{key, val}})# Writes one key and journals the mutation.def oput(+key: String, +val: String) -> Op<&2, Unit>: Op{st => pair_put(key, val, MyLsmStore.put_go(st, key, val))}# Pairs a finished point read with an empty journal.def pair_get(finished: DbTypes.Db & Maybe<&2, String>) -> (DbTypes.Db & Maybe<&2, String>) & Journal: match finished: case (db_next, found_val): ((db_next, found_val), JEmpty{})# Reads one key; missing reads as None. Reads journal nothing.def oget(+key: String) -> Op<&2, Maybe<&2, String>>: Op{st => pair_get(MyLsmStore.sget_go(st, key))}# Pairs a finished point delete with its journaled mutation.def pair_del(+key: String, db_next: DbTypes.Db) -> (DbTypes.Db & Unit) & Journal: ((db_next, Unit{}), JOne{WalTypes.Del{key}})# Deletes one key and journals the mutation.def odel(+key: String) -> Op<&2, Unit>: Op{st => pair_del(key, MyLsmStore.del_go(st, key))}# Opens a store handle over a directory label (pure; no IO touched).def open_store(+dir: String) -> Handle: Handle{MyLsmStore.open(dir)}# Default store label for CLI use (pure; no IO touched).def open_default_store() -> Handle: open_store("./.ber")# Wraps a reopened database; IO errors propagate as Fail.def wrap_reopened(reopened: Result<&1, &1, U32 & String, DbTypes.Db>) -> Result<&1, &1, U32 & String, Handle>: match reopened: case Fail{error}: Fail{error} case Done{db}: Done{Handle{db}}# Opens a durable handle, replaying wal.log when present; missing file# answers empty. The caller ensures the directory exists.def open_durable(+dir: String) -> IO(Result<&1, &1, U32 & String, Handle>): do IO<Result<&1, &1, U32 & String, Handle>>: reopened : Result<&1, &1, U32 & String, DbTypes.Db> <- MyLsmStore.reopen(dir) return wrap_reopened(reopened)# Projects the materialized journal out of a tracked run (laws observe this).def tracked_journal(a, -A: Kind(a), tracked: (Handle & A) & Journal) -> List<&2, WalTypes.Mut>: match tracked: case ((unused_handle, unused_answer), journal): flatten_journal(journal)# Rewraps a computed runner pair into Handle form, dropping the journal.def rewrap_pair(a, -A: Kind(a), computed: (DbTypes.Db & A) & Journal) -> Handle & A: match computed: case ((db_next, answer), unused_journal): (Handle{db_next}, answer)# Executes a whole op chain against a handle (pure; journal dropped).def run_op(a, -A: Kind(a), store: Handle, op: Op<a, A>) -> Handle & A: match store: case Handle{db}: rewrap_pair(a, A, Op.inner(a, A, op)(db))# Rewraps a tracked runner pair, keeping the journal.def rewrap_tracked(a, -A: Kind(a), computed: (DbTypes.Db & A) & Journal) -> (Handle & A) & Journal: match computed: case ((db_next, answer), journal): ((Handle{db_next}, answer), journal)# Executes a whole op chain, exposing the journal for the WAL runner.def run_op_tracked(a, -A: Kind(a), store: Handle, op: Op<a, A>) -> (Handle & A) & Journal: match store: case Handle{db}: rewrap_tracked(a, A, Op.inner(a, A, op)(db))# Wraps a written database with its answer; IO errors propagate as Fail.def wrap_written(a, -A: Kind(a), written: Result<&1, &1, U32 & String, DbTypes.Db>, answer: A) -> Result<&1, &1, U32 & String, Handle & A>: match written: case Fail{error}: Fail{error} case Done{db_next}: Done{(Handle{db_next}, answer)}# Skips the WAL when the journal is empty (pure reads never fsync).def flush_clean(a, -A: Kind(a), store_next: Handle, answer: A) -> IO(Result<&1, &1, U32 & String, Handle & A>): match store_next: case Handle{db}: IO.pure(Result<&1, &1, U32 & String, Handle & A>, Done{(Handle{db}, answer)})# Flushes a nonempty journal through one WAL frame plus fsync.def flush_dirty(a, -A: Kind(a), store_next: Handle, answer: A, +head_mut: WalTypes.Mut, +tail_muts: List<&2, WalTypes.Mut>) -> IO(Result<&1, &1, U32 & String, Handle & A>): match store_next: case Handle{db}: do IO<Result<&1, &1, U32 & String, Handle & A>>: written : Result<&1, &1, U32 & String, DbTypes.Db> <- DbTypes.db_write(db, WalTypes.Batch{Con{head_mut, tail_muts}}) return wrap_written(a, A, written, answer)# Dispatches on journal emptiness without matching computed values.def flush_list(a, -A: Kind(a), store_next: Handle, answer: A, +mutations: List<&2, WalTypes.Mut>) -> IO(Result<&1, &1, U32 & String, Handle & A>): match mutations: case Nil{}: flush_clean(a, A, store_next, answer) case Con{head_mut, tail_muts}: flush_dirty(a, A, store_next, answer, head_mut, tail_muts)# Splits a tracked runner pair for the flush dispatch.def flush_tracked(a, -A: Kind(a), tracked: (Handle & A) & Journal) -> IO(Result<&1, &1, U32 & String, Handle & A>): match tracked: case ((store_next, answer), journal): flush_list(a, A, store_next, answer, flatten_journal(journal))# Executes a whole op chain durably: pure run, then one WAL append plus# fsync of exactly the journaled mutations.def run_op_durable(a, -A: Kind(a), store: Handle, op: Op<a, A>) -> IO(Result<&1, &1, U32 & String, Handle & A>): flush_tracked(a, A, run_op_tracked(a, A, store, op))