~/bend-docscommunity

src/Flush.bend source

src/Flush.bend on the hub · documented module

import Baseimport ./Keys.bend as Keysimport ./MemTable.bend as MemTableimport ./Sstable.bend as Sstableimport ./SortedRun.bend as SortedRunimport ./SstStreamIo.bend as SstStreamIoimport ./Wal.bend as Walimport ./Manifest.bend as Manifestimport ./StorageBytes.bend as StorageBytesimport ./Fs.bend as Fsimport ./DbIo.bend as DbIoimport bend-kit-bytes@0.3.2.0/bytes.bend as Bytesimport ./Db.bend as Dbimport ./CrashPoint.bend as CrashPointimport ./FlushPolicy.bend as Policy# Represent ManifestReadState data used by the memtable flush.type ManifestReadState is Type:  ManifestNeed{file: File, remaining: Nat, chunks: List<&1, Bytes.Bytes>}  ManifestGot{pair: File & Result<&1, &1, U32 & String, Bytes.Bytes>, remaining: Nat, chunks: List<&1, Bytes.Bytes>}# Flush (Task 8): sealed memtable -> L0 SSTable file + manifest + WAL reset.## Crash protocol (each step durable before the next, tmp + rename for# atomic visibility):#   1. ensure <dir>/l0; write <name>.tmp; fsync; rename to <name>; fsync l0.#   2. load manifest (absent file = fresh M{Nil}; corrupt = ABORT Fail).#   3. write MANIFEST.tmp; fsync; rename; fsync dir.#   4. remove wal.log (failure ignored: a stale WAL replays idempotently —#      same versions, mem-first reads return equal values; merge dedups).# A crash before step 3 leaves the table file orphaned (recovery ignores# files absent from the manifest; Task 12 sweeps orphans). A crash after# step 3 with a live WAL replays duplicates harmlessly (see above).# Empty memtable flush is a no-op (no files touched).# CrashPoint calls are test-only host effects. Unset or unequal checkpoints are# successful no-ops; signal delivery and filesystem ordering are not Bend proofs.# Handle l0 add in the memtable flush.def l0_add(+levels: List<&2, List<&2, Sstable.Table>>, +tbl: Sstable.Table) -> List<&2, List<&2, Sstable.Table>>:  match levels:    case Nil{}:      Con{Con{tbl, Nil{}}, Nil{}}    case Con{l0, rest}:      Con{Con{tbl, l0}, rest}# Handle mfst add in the memtable flush.def mfst_add(+mfst: Manifest.Manifest, +name: String) -> Manifest.Manifest:  match mfst:    case Manifest.M{lvs}:      match lvs:        case Nil{}:          Manifest.M{Con{Con{name, Nil{}}, Nil{}}}        case Con{l0, rest}:          Manifest.M{Con{Con{name, l0}, rest}}# --- File helpers (pair-splits contained; straight-line callers above) ---# Write tail for the memtable flush.def write_tail(  fr: (File & Result<&1, &1, U32 & String, Unit>),  +path: 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))        _cls : Unit <- File.close(f2)        _syn : Unit <- IO.try(Unit, Fs.fsync(path))        _p_rm : Unit <- IO.try(Unit, Fs.chmod(path, U32.from_nat(384n)))        return Done{res}# Write file for the memtable flush.def write_file(+path: 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(path, "w"))    fr : (File & Result<&1, &1, U32 & String, Unit>) <- File.write(f, data)    write_tail(fr, path)def write_manifest.bytes(+path: String, bytes: Bytes.Bytes) -> IO(Result<&1, &1, U32 & String, Unit>):  do IO<Result<&1, &1, U32 & String, Unit>>:    f : File <- IO.try(File, File.open(path, "w"))    fr : File & Result<&1, &1, U32 & String, Unit> <- Fs.write_bytes(f, bytes)    write_tail(fr, path)def write_manifest.encoded(  path: String,  encoded: Result<&1, &1, Manifest.Error, Bytes.Bytes>) -> IO(Result<&1, &1, U32 & String, Unit>):  match encoded:    case Fail{_}:      IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{(U32.from_nat(1n), "cannot encode manifest")})    case Done{bytes}:      write_manifest.bytes(path, bytes)# Write manifest for the memtable flush.def write_manifest(path: String, manifest: Manifest.Manifest) -> IO(Result<&1, &1, U32 & String, Unit>):  write_manifest.encoded(path, Manifest.serialize(manifest))# Read tail3 for the memtable flush.def read_tail3(fr: (File & Result<&1, &1, U32 & String, String>)) -> IO(Result<&1, &1, U32 & String, String>):  match fr:    case (f3, r1):      do IO<Result<&1, &1, U32 & String, String>>:        _cls : Unit <- File.close(f3)        IO.pure(Result<&1, &1, U32 & String, String>, r1)# Base File.read decodes each returned chunk independently. A chunk boundary can# split a multi-byte UTF-8 character, so table files are read in one bounded call# until a byte-oriented streaming effect with decoder carry is available.def read_file(+path: String) -> IO(Result<&1, &1, U32 & String, String>):  do IO<Result<&1, &1, U32 & String, String>>:    f : File <- IO.try(File, File.open(path, "r"))    fr : (File & Result<&1, &1, U32 & String, String>) <- File.read(f, U32.from_nat(4294967295n))    read_tail3(fr)# Handle mfst parsed in the memtable flush.def mfst_parsed(+opt: Maybe<&2, Manifest.Manifest>) -> Result<&1, &1, U32 & String, Manifest.Manifest>:  match opt:    case None{}:      Fail{(U32.from_nat(1n), "bad manifest")}    case Some{mm}:      Done{mm}def manifest.read.result(  pair: File & Result<&1, &1, U32 & String, Bytes.Bytes>) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  match pair:    case (file, Fail{error}):      do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:        _closed : Unit <- File.close(file)        return Fail{error}    case (file, Done{bytes}):      do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:        _closed : Unit <- File.close(file)        return mfst_parsed(Manifest.parse(bytes))def manifest.read.loop(fuel: Nat, state: ManifestReadState) -> IO(File & Result<&1, &1, U32 & String, Bytes.Bytes>):  match fuel:    case 0n:      match state:        case ManifestNeed{file, +remaining, chunks}:          match remaining:            case 0n:              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Done{Bytes.concat(List.reverse(&1, Bytes.Bytes, chunks))}))            case 1n+_:              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Fail{(U32.from_nat(7n), "bounded manifest read exhausted")}))        case ManifestGot{pair, _, _}:          match pair:            case (file, Fail{error}):              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Fail{error}))            case (file, Done{_}):              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Fail{(U32.from_nat(7n), "bounded manifest read exhausted")}))    case 1n+rest:      match state:        case ManifestNeed{file, +remaining, chunks}:          match remaining:            case 0n:              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Done{Bytes.concat(List.reverse(&1, Bytes.Bytes, chunks))}))            case 1n+_:              do IO<File & Result<&1, &1, U32 & String, Bytes.Bytes>>:                pair : File & Result<&1, &1, U32 & String, Bytes.Bytes> <- Fs.read_bytes(file, U32.from_nat(remaining))                manifest.read.loop(rest, ManifestGot{pair, remaining, chunks})        case ManifestGot{pair, +remaining, chunks}:          match pair:            case (file, Fail{error}):              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Fail{error}))            case (file, Done{Bytes.Bytes{0, _}}):              IO.pure(File & Result<&1, &1, U32 & String, Bytes.Bytes>, (file, Fail{(U32.from_nat(8n), "truncated manifest file")}))            case (file, Done{Bytes.Bytes{+len, buf}}):              manifest.read.loop(rest, ManifestNeed{file, Nat.sub(remaining, U32.to_nat(len)), Con{Bytes.Bytes{len, buf}, chunks}})def manifest.read.size.run(file: File, +size: Nat) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:    pair : File & Result<&1, &1, U32 & String, Bytes.Bytes> <- manifest.read.loop(Nat.add(size, 1n), ManifestNeed{file, size, Nil{}})    manifest.read.result(pair)def manifest.read.size.valid(  file: File,  size: Nat,  within: Bool) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  match within:    case False{}:      do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:        _closed : Unit <- File.close(file)        return Fail{(U32.from_nat(9n), "manifest too large")}    case True{}:      manifest.read.size.run(file, size)def manifest.read.size(  file: File,  result: Result<&1, &1, U32 & String, Nat>) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  match result:    case Fail{error}:      do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:        _closed : Unit <- File.close(file)        return Fail{error}    case Done{+size}:      manifest.read.size.valid(file, size, Nat.is_le(size, U32.to_nat(StorageBytes.MAX_MANIFEST_BYTES())))# Handle mfst open error in the memtable flush.def mfst_open_error(  missing: Bool,  +code: U32) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  match missing:    case True{}:      IO.pure(Result<&1, &1, U32 & String, Manifest.Manifest>, Done{Manifest.M{Nil{}}})    case False{}:      IO.pure(Result<&1, &1, U32 & String, Manifest.Manifest>, Fail{(code, "cannot open manifest")})# Handle mfst opened in the memtable flush.def mfst_opened(  path: String,  result: Result<&1, &1, U32 & String, File>) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  match result:    case Fail{(+code, _)}:      mfst_open_error(U32.is_eq(code, 2), code)    case Done{file}:      do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:        size : Result<&1, &1, U32 & String, Nat> <- Fs.file_size(path)        manifest.read.size(file, size)# Handle load mfst in the memtable flush.def load_mfst(+path: String) -> IO(Result<&1, &1, U32 & String, Manifest.Manifest>):  do IO<Result<&1, &1, U32 & String, Manifest.Manifest>>:    opened : Result<&1, &1, U32 & String, File> <- File.open(path, "r")    mfst_opened(path, opened)# Handle dir tail in the memtable flush.def dir_tail(res: Result<&1, &1, U32 & String, Nat>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>):  match res:    case Done{n}:      IO.pure(Result<&1, &1, U32 & String, Unit>, Done{Unit{}})    case Fail{e}:      Fs.make_dir(path)# Handle ensure dir in the memtable flush.def ensure_dir(+path: String) -> IO(Result<&1, &1, U32 & String, Unit>):  do IO<Result<&1, &1, U32 & String, Unit>>:    r : Result<&1, &1, U32 & String, Nat> <- Fs.read_dir_count(path)    dir_tail(r, path)# --- Flush chain (defined bottom-up: tails first, entry last) ---# Flush goB for the memtable flush.def flush_goB(  +dir: String,  +mpath: String,  +mtmp: String,  +m2: Manifest.Manifest,  +tbl: Sstable.Table,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +mem: MemTable.MemTable,  +batch_cap: Nat,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  do IO<Result<&1, &1, U32 & String, Db.Db>>:    _v1 : Unit <- IO.try(Unit, write_manifest(mtmp, m2))    _cp1 : Unit <- IO.try(Unit, CrashPoint.hit("flush.manifest_synced"))    _v2 : Unit <- IO.try(Unit, Fs.rename(mtmp, mpath))    _v3 : Unit <- IO.try(Unit, Fs.fsync(dir))    _cp2 : Unit <- IO.try(Unit, CrashPoint.hit("flush.manifest_published"))    _rm : Result<&1, &1, U32 & String, Unit> <- Fs.remove(Db.wal_path(dir))    _wal : Unit <- IO.try(Unit, DbIo.wal_initialize(dir))    return Done{Db.Db{dir, mem, MemTable.MT{Nil{}}, batch_cap, Nil{}, l0_add(levels, tbl), Nat.add(flushed, 1n), Manifest.token(m2), 0n, 0n}}# Flush drift for the memtable flush.def flush_drift(  ok: Bool,  +dir: String,  +mpath: String,  +mtmp: String,  +mfst: Manifest.Manifest,  +tbl: Sstable.Table,  +name: String,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +mem: MemTable.MemTable,  +batch_cap: Nat,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  match ok:    case False{}:      IO.pure(Result<&1, &1, U32 & String, Db.Db>, Fail{(U32.from_nat(2n), "manifest drift")})    case True{}:      flush_goB(dir, mpath, mtmp, mfst_add(mfst, name), tbl, levels, flushed, mem, batch_cap)# Flush goA for the memtable flush.def flush_goA(  +dir: String,  +l0dir: String,  +tmp: String,  +final: String,  +mpath: String,  +mtmp: String,  +entries: List<&2, MemTable.Entry>,  +level: U32,  +tbl: Sstable.Table,  +name: String,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +manifest_token: String,  +mem: MemTable.MemTable,  +batch_cap: Nat,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  do IO<Result<&1, &1, U32 & String, Db.Db>>:    _u1 : Unit <- IO.try(Unit, ensure_dir(l0dir))    _u2 : Unit <- IO.try(Unit, SstStreamIo.write_table(tmp, entries, level))    _cp1 : Unit <- IO.try(Unit, CrashPoint.hit("flush.table_synced"))    _u3 : Unit <- IO.try(Unit, Fs.rename(tmp, final))    _u4 : Unit <- IO.try(Unit, Fs.fsync(l0dir))    _cp2 : Unit <- IO.try(Unit, CrashPoint.hit("flush.table_published"))    +m : Manifest.Manifest <- IO.try(Manifest.Manifest, load_mfst(mpath))    flush_drift(String.eq(Manifest.token(m), manifest_token), dir, mpath, mtmp, m, tbl, name, levels, flushed, mem, batch_cap)# Flush pre for the memtable flush.def flush_pre(  +dir: String,  +entries: List<&2, MemTable.Entry>,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +manifest_token: String,  +batch_cap: Nat,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  +tbl = Sstable.build(entries, 0n, List.length(&2, MemTable.Entry, entries))  +name = Policy.table_name(flushed)  +l0dir = dir ++ "/l0"  flush_goA(dir, l0dir, l0dir ++ "/" ++ name ++ ".tmp", l0dir ++ "/" ++ name, dir ++ "/MANIFEST", dir ++ "/MANIFEST.tmp", Db.table_entries(tbl), 0, tbl, name, levels, flushed,  manifest_token, MemTable.MT{Nil{}}, batch_cap)# Frozen-path flush MERGES mem (newer) over frozen (older) into one L0 table# and drains both: the WAL is truncated on publish, so anything left in mem# would otherwise exist nowhere on disk (loss window caught by 20485 smoke).def flush_frozen(  +dir: String,  +mem: MemTable.MemTable,  +entries: List<&2, MemTable.Entry>,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +manifest_token: String,  +batch_cap: Nat,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  match mem:    case MemTable.MT{live}:      +merged = SortedRun.merge_many_newest(Con{live, Con{entries, Nil{}}})      +tbl = Sstable.build(merged, 0n, List.length(&2, MemTable.Entry, merged))      +name = Policy.table_name(flushed)      +l0dir = dir ++ "/l0"      flush_goA(dir, l0dir, l0dir ++ "/" ++ name ++ ".tmp", l0dir ++ "/" ++ name, dir ++ "/MANIFEST", dir ++ "/MANIFEST.tmp", Db.table_entries(tbl), 0, tbl, name, levels, flushed,      manifest_token, MemTable.MT{Nil{}}, batch_cap)# Flush pick for the memtable flush.def flush_pick(  +frozen: MemTable.MemTable,  +mem: MemTable.MemTable,  +dir: String,  +levels: List<&2, List<&2, Sstable.Table>>,  +flushed: Nat,  +manifest_token: String,  +batch_cap: Nat,  +bcache: List<&2, Db.BEntry>,) -> IO(Result<&1, &1, U32 & String, Db.Db>):  match frozen mem:    case MemTable.MT{Nil{}} MemTable.MT{Nil{}}:      IO.pure(Result<&1, &1, U32 & String, Db.Db>, Done{Db.Db{dir, MemTable.MT{Nil{}}, MemTable.MT{Nil{}}, batch_cap, bcache, levels, flushed, manifest_token, 0n, 0n}})    case MemTable.MT{Nil{}} MemTable.MT{Con{e, t}}:      flush_pre(dir, Con{e, t}, levels, flushed, manifest_token, batch_cap)    case MemTable.MT{Con{fe, ft}} _:      flush_frozen(dir, mem, Con{fe, ft}, levels, flushed, manifest_token, batch_cap)# Handle flush in the memtable flush.def flush(+db: Db.Db) -> IO(Result<&1, &1, U32 & String, Db.Db>):  match db:    case Db.Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}:      flush_pick(frozen, mem, dir, levels, flushed, manifest_token, batch_cap, bcache)# --- Closed-vector laws live in laws/Flush.bend (spec law 6) ---