import Base
import ./Keys.bend as Keys
import ./MemTable.bend as MemTable
import ./Sstable.bend as Sstable
import ./SortedRun.bend as SortedRun
import ./SstStreamIo.bend as SstStreamIo
import ./Wal.bend as Wal
import ./Manifest.bend as Manifest
import ./StorageBytes.bend as StorageBytes
import ./Fs.bend as Fs
import ./DbIo.bend as DbIo
import bend-kit-bytes@0.3.2.0/bytes.bend as Bytes
import ./Db.bend as Db
import ./CrashPoint.bend as CrashPoint
import ./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
/l0; write .tmp; fsync; rename to ; 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>:
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>:
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>:
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>:
_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>:
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>:
_closed : Unit <- File.close(file)
return Fail{error}
case (file, Done{bytes}):
do IO>:
_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>:
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>:
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>:
_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>:
_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>:
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>:
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>:
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>:
_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>:
_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) ---