~/bend-docscommunity

src/Db.bend source

src/Db.bend on the hub · documented module

import Baseimport ./Keys.bend as Keysimport ./MemTable.bend as MemTableimport ./Sstable.bend as Sstableimport ./Wal.bend as Walimport ./Fs.bend as Fsimport ./Manifest.bend as Manifestimport ./CrashPoint.bend as CrashPoint# Db API (Task 7): durable write path + reads.## Write path (durability ordering, the load-bearing property): encode the# batch -> append the WAL -> fsync -> THEN update mem. A crash before fsync# loses at most unacked writes (recovery replays the WAL prefix, Task 10).# This ordering is a code-review property (no law over IO exists); Task 12# fault-injection crashes between each step to verify it empirically.# CrashPoint calls are test-only host effects. Unset or unequal checkpoints are# successful no-ops; signal delivery and filesystem ordering are not Bend proofs.## Read path: mem (newest-first) ++ L0 newest-first ++ L1 ... — first match# wins INCLUDING tombstones (MemTable.scan_go freezes), so deletes can never# resurrect older versions. Cross-level newest-first order is CORRECT iff# Task 9 maintains the L0-drain invariant: compaction always drains ALL of# L0 (inputs deleted), so every L0 table is strictly newer than every L1+# table; within a level, tables are stored newest-first. Flush (Task 8)# Cons'es new tables at the head; compaction (Task 9) preserves the order.type BEntry is Data:  BEntry{tab: String, key: String, val: Maybe<&2, String>}def default_batch_cap() -> Nat:  1ndef bcache_bound() -> Nat:  256n# manifest_token is the exact serialized Manifest observed by this handle. Flush# and compaction compare it before publication to reject stale-handle drift.type Db is Data:  Db{dir: String, mem: MemTable.MemTable, frozen: MemTable.MemTable, batch_cap: Nat, bcache: List<&2, BEntry>, levels: List<&2, List<&2, Sstable.Table>>, flushed: Nat, manifest_token: String, mem_count: Nat, frozen_count: Nat}# Rotation result: both tables plus their exact counts, so no pass ever# walks a memtable to learn its size. Counts stay exact by construction# (open_db starts 0/0; every transition does arithmetic, never List.length).type RotRes is Data:  Rot{mem: MemTable.MemTable, frozen: MemTable.MemTable, mem_count: Nat, frozen_count: Nat}def wal_path(+dir: String) -> String:  dir ++ "/wal.log"def open_db(+dir: String) -> Db:  Db{dir, MemTable.empty(), MemTable.empty(), default_batch_cap(), Nil{}, Nil{}, 0n, Manifest.serialize(Manifest.M{Nil{}}), 0n, 0n}# --- Pure write core (laws below pin batch == sequential) ---def apply_mut(+mem: MemTable.MemTable, +mut: Wal.Mut) -> MemTable.MemTable:  match mut:    case Wal.Put{key, val}:      MemTable.put(mem, key, val)    case Wal.Del{key}:      MemTable.del(mem, key)def apply_batch(+muts: List<&2, Wal.Mut>, +mem: MemTable.MemTable) -> MemTable.MemTable:  match muts:    case Nil{}:      mem    case Con{+h, t}:      apply_batch(t, apply_mut(mem, h))# --- Pure read core (concat newest-first, first match wins) ---def table_entries(+tab: Sstable.Table) -> List<&2, MemTable.Entry>:  match tab:    case Sstable.Tbl{entries, filter, nbits, smallest, largest, count, chunks, blocks}:      entriesdef level_entries(+tabs: List<&2, Sstable.Table>) -> List<&2, MemTable.Entry>:  match tabs:    case Nil{}:      Nil{}    case Con{+h, t}:      List.append(&2, MemTable.Entry, table_entries(h), level_entries(t))def all_level_entries(+lvls: List<&2, List<&2, Sstable.Table>>) -> List<&2, MemTable.Entry>:  match lvls:    case Nil{}:      Nil{}    case Con{+h, t}:      List.append(&2, MemTable.Entry, level_entries(h), all_level_entries(t))def all_entries(+db: Db) -> List<&2, MemTable.Entry>:  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      match mem:        case MemTable.MT{entries}:          List.append(&2, MemTable.Entry, entries, all_level_entries(levels))def table_hit_if(present: Bool, tab: Sstable.Table, +key: String) -> Maybe<&2, Maybe<&2, String>>:  match present:    case False{}:      None{}    case True{}:      Sstable.block_get_hit(tab, key)def table_hit(+tab: Sstable.Table, +key: String) -> Maybe<&2, Maybe<&2, String>>:  table_hit_if(Sstable.maybe_present(tab, key), tab, key)def tables_hits(+tabs: List<&2, Sstable.Table>, +key: String) -> List<&2, Maybe<&2, Maybe<&2, String>>>:  match tabs:    case Nil{}:      Nil{}    case Con{+h, t}:      Con{table_hit(h, key), tables_hits(t, key)}def levels_hits(+lvls: List<&2, List<&2, Sstable.Table>>, +key: String) -> List<&2, Maybe<&2, Maybe<&2, String>>>:  match lvls:    case Nil{}:      Nil{}    case Con{h, t}:      List.append(&2, Maybe<&2, Maybe<&2, String>>, tables_hits(h, key), levels_hits(t, key))def hits_first(+hits: List<&2, Maybe<&2, Maybe<&2, String>>>) -> Maybe<&2, Maybe<&2, String>>:  match hits:    case Nil{}:      None{}    case Con{None{}, t}:      hits_first(t)    case Con{Some{found}, t}:      Some{found}def flatten_opt(+opt: Maybe<&2, Maybe<&2, String>>) -> Maybe<&2, String>:  match opt:    case None{}:      None{}    case Some{inner}:      innerdef db_get(+db: Db, +key: String) -> Maybe<&2, String>:  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      flatten_opt(hits_first(Con{MemTable.get_hit(mem, key), Con{MemTable.get_hit(frozen, key), levels_hits(levels, key)}}))# --- Phase A3: table-identity read cache (assoc-list, bounded FIFO) ---# LRU verdict (spike-by-inspection, 2026-09-25): hub lru has no simple public# put, reads hand the cache back, and every op needs a W.U64 clock — adopting# it threads time through the pure read path and changes db_get's shape for a# marginal gain over this list. Custom assoc-list wins; see the Task 3 commit.def pick_hit(hit: Bool, ans: Maybe<&2, String>, best: Maybe<&2, Maybe<&2, String>>) -> Maybe<&2, Maybe<&2, String>>:  match hit:    case True{}:      Some{ans}    case False{}:      bestdef bcache_go(+cache: List<&2, BEntry>, +tab: String, +key: String, +best: Maybe<&2, Maybe<&2, String>>) -> Maybe<&2, Maybe<&2, String>>:  match cache:    case Nil{}:      best    case Con{BEntry{t, ky, val}, rest}:      bcache_go(rest, tab, key, pick_hit(Bool.and(String.eq(t, tab), String.eq(ky, key)), val, best))def bcache_lookup(+cache: List<&2, BEntry>, +tab: String, +key: String) -> Maybe<&2, Maybe<&2, String>>:  bcache_go(cache, tab, key, None{})def bcache_push_all(+fresh: List<&2, BEntry>, +old: List<&2, BEntry>) -> List<&2, BEntry>:  List.take(&2, BEntry, List.append(&2, BEntry, fresh, old), bcache_bound())def probe_fill(hit: Maybe<&2, Maybe<&2, String>>, +tab: String, +key: String) -> (Maybe<&2, Maybe<&2, String>> & List<&2, BEntry>):  match hit:    case None{}:      (None{}, Nil{})    case Some{+inner}:      (Some{inner}, Con{BEntry{tab, key, inner}, Nil{}})def probe_cached(cached: Maybe<&2, Maybe<&2, String>>, +tab: Sstable.Table, +key: String) -> (Maybe<&2, Maybe<&2, String>> & List<&2, BEntry>):  match cached:    case None{}:      probe_fill(table_hit(tab, key), Sstable.table_id(tab), key)    case Some{ans}:      (Some{ans}, Nil{})def table_probe(+tab: Sstable.Table, +key: String, +cache: List<&2, BEntry>) -> (Maybe<&2, Maybe<&2, String>> & List<&2, BEntry>):  probe_cached(bcache_lookup(cache, Sstable.table_id(tab), key), tab, key)def tables_cons(cur: (Maybe<&2, Maybe<&2, String>> & List<&2, BEntry>), rest: (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>)) -> (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>):  match cur:    case (hit, delta):      match rest:        case (rhits, rdelta):          (Con{hit, rhits}, List.append(&2, BEntry, delta, rdelta))def tables_walk(+tabs: List<&2, Sstable.Table>, +key: String, +cache: List<&2, BEntry>) -> (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>):  match tabs:    case Nil{}:      (Nil{}, Nil{})    case Con{+h, t}:      tables_cons(table_probe(h, key, cache), tables_walk(t, key, cache))def levels_cons(cur: (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>), rest: (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>)) -> (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>):  match cur:    case (chits, cdelta):      match rest:        case (rhits, rdelta):          (List.append(&2, Maybe<&2, Maybe<&2, String>>, chits, rhits), List.append(&2, BEntry, cdelta, rdelta))def levels_cached(+lvls: List<&2, List<&2, Sstable.Table>>, +key: String, +cache: List<&2, BEntry>) -> (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>):  match lvls:    case Nil{}:      (Nil{}, Nil{})    case Con{+h, t}:      levels_cons(tables_walk(h, key, cache), levels_cached(t, key, cache))def cached_levels(lr: (List<&2, Maybe<&2, Maybe<&2, String>>> & List<&2, BEntry>), +dir: String, +mem: MemTable.MemTable, +frozen: MemTable.MemTable, +batch_cap: Nat, +bcache: List<&2, BEntry>, +levels: List<&2, List<&2, Sstable.Table>>, +flushed: Nat, +manifest_token: String, +mem_count: Nat, +frozen_count: Nat) -> Db & Maybe<&2, String>:  match lr:    case (hits, delta):      (Db{dir, mem, frozen, batch_cap, bcache_push_all(delta, bcache), levels, flushed, manifest_token, mem_count, frozen_count}, flatten_opt(hits_first(hits)))def cached_frozen(fhit: Maybe<&2, Maybe<&2, String>>, +dir: String, +mem: MemTable.MemTable, +frozen: MemTable.MemTable, +batch_cap: Nat, +bcache: List<&2, BEntry>, +levels: List<&2, List<&2, Sstable.Table>>, +flushed: Nat, +manifest_token: String, +key: String, +mem_count: Nat, +frozen_count: Nat) -> Db & Maybe<&2, String>:  match fhit:    case Some{found}:      (Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}, flatten_opt(Some{found}))    case None{}:      cached_levels(levels_cached(levels, key, bcache), dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count)def cached_mem(mhit: Maybe<&2, Maybe<&2, String>>, +dir: String, +mem: MemTable.MemTable, +frozen: MemTable.MemTable, +batch_cap: Nat, +bcache: List<&2, BEntry>, +levels: List<&2, List<&2, Sstable.Table>>, +flushed: Nat, +manifest_token: String, +key: String, +mem_count: Nat, +frozen_count: Nat) -> Db & Maybe<&2, String>:  match mhit:    case Some{found}:      (Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}, flatten_opt(Some{found}))    case None{}:      cached_frozen(MemTable.get_hit(frozen, key), dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, key, mem_count, frozen_count)def db_get_cached(+db: Db, +key: String) -> Db & Maybe<&2, String>:  match db:    case Db{dir, +mem, +frozen, +batch_cap, +bcache, +levels, +flushed, +manifest_token, +mem_count, +frozen_count}:      cached_mem(MemTable.get_hit(mem, key), dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, key, mem_count, frozen_count)# WAL framing v1 is applied here (see Recover): each batch is stored as# dashes(len) ++ ";" ++ encode(batch) so replay can stream frames and# truncate a torn tail. Files are permissioned 0600 at creation.def wal_frame(+data: String) -> String:  Wal.dashes(String.length(data)) ++ ";" ++ data# Tail of the append: match heads the def body (params always# destructurable), so the write's pair splits with single uses.def wal_tail(fr: (File & Result<&1, &1, U32 & String, Unit>), +dir: String) -> IO(Result<&1, &1, U32 & String, Unit>):  match fr:    case (f2, r):      do IO<Result<&1, &1, U32 & String, Unit>>:        res : Unit <- IO.try(Unit, IO.pure(Result<&1, &1, U32 & String, Unit>, r))        _cp1 : Unit <- IO.try(Unit, CrashPoint.hit("wal.appended"))        _cls : Unit <- File.close(f2)        _syn : Unit <- IO.try(Unit, Fs.fsync(wal_path(dir)))        _cp2 : Unit <- IO.try(Unit, CrashPoint.hit("wal.synced"))        _prm : Unit <- IO.try(Unit, Fs.chmod(wal_path(dir), U32.from_nat(384n)))        return Done{res}def wal_append(+dir: String, +data: String) -> IO(Result<&1, &1, U32 & String, Unit>):  do IO<Result<&1, &1, U32 & String, Unit>>:    f : File <- IO.try(File, File.open(wal_path(dir), "a"))    fr : (File & Result<&1, &1, U32 & String, Unit>) <- File.write(f, wal_frame(data))    wal_tail(fr, dir)def batch_len(+muts: List<&2, Wal.Mut>) -> Nat:  List.length(&2, Wal.Mut, muts)def rotate_cnt(full: Bool, +grown: MemTable.MemTable, +frozen: MemTable.MemTable, +count: Nat, +fcount: Nat) -> RotRes:  match full:    case True{}:      Rot{MemTable.empty(), grown, 0n, count}    case False{}:      Rot{grown, frozen, count, fcount}def frozen_empty(+frozen: MemTable.MemTable) -> Bool:  match frozen:    case MemTable.MT{entries}:      match entries:        case Nil{}:          True{}        case Con{_, _}:          False{}def rotate_after(+grown: MemTable.MemTable, +frozen: MemTable.MemTable, +count: Nat, +fcount: Nat) -> RotRes:  rotate_cnt(Bool.and(Nat.is_le(4096n, count), frozen_empty(frozen)), grown, frozen, count, fcount)def apply_and_rotate(+muts: List<&2, Wal.Mut>, +mem: MemTable.MemTable, +frozen: MemTable.MemTable, +mem_count: Nat, +frozen_count: Nat) -> RotRes:  rotate_after(apply_batch(muts, mem), frozen, Nat.add(mem_count, batch_len(muts)), frozen_count)def apply_done(res: RotRes, +dir: String, +batch_cap: Nat, +bcache: List<&2, BEntry>, +levels: List<&2, List<&2, Sstable.Table>>, +flushed: Nat, +manifest_token: String) -> Db:  match res:    case Rot{mem2, frozen2, mc, fc}:      Db{dir, mem2, frozen2, batch_cap, bcache, levels, flushed, manifest_token, mc, fc}def rot_done(res: RotRes, +dir: String, +batch_cap: Nat, +bcache: List<&2, BEntry>, +levels: List<&2, List<&2, Sstable.Table>>, +flushed: Nat, +manifest_token: String) -> Result<&1, &1, U32 & String, Db>:  Done{apply_done(res, dir, batch_cap, bcache, levels, flushed, manifest_token)}def db_write(+db: Db, +batch: Wal.Batch) -> IO(Result<&1, &1, U32 & String, Db>):  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      match batch:        case Wal.Batch{muts}:          do IO<Result<&1, &1, U32 & String, Db>>:            _res : Unit <- IO.try(Unit, wal_append(dir, Wal.encode(Wal.Batch{muts})))            return rot_done(apply_and_rotate(muts, mem, frozen, mem_count, frozen_count), dir, batch_cap, bcache, levels, flushed, manifest_token)def db_put(+db: Db, +key: String, +val: String) -> IO(Result<&1, &1, U32 & String, Db>):  db_write(db, Wal.Batch{Con{Wal.Put{key, val}, Nil{}}})def batch_cap_of(+db: Db) -> Nat:  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      batch_capdef with_batch_cap(+db: Db, +cap: Nat) -> Db:  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      Db{dir, mem, frozen, cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}def staged_muts(+staged: List<&2, Wal.Batch>) -> List<&2, Wal.Mut>:  match staged:    case Nil{}:      Nil{}    case Con{Wal.Batch{muts}, t}:      List.append(&2, Wal.Mut, muts, staged_muts(t))def db_write_staged(+db: Db, +staged: List<&2, Wal.Batch>) -> IO(Result<&1, &1, U32 & String, Db>):  match db:    case Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      db_write(Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}, Wal.Batch{staged_muts(List.reverse(&2, Wal.Batch, staged))})def db_del(+db: Db, +key: String) -> IO(Result<&1, &1, U32 & String, Db>):  db_write(db, Wal.Batch{Con{Wal.Del{key}, Nil{}}})# --- Closed-vector laws live in laws/Db.bend (spec laws 1, 2, 3, 7) ---