src/DbIo.bend source
src/DbIo.bend on the hub · documented module
import Baseimport ./Db.bend as Dbimport ./Wal.bend as Walimport ./Fs.bend as Fsimport ./CrashPoint.bend as CrashPointimport ./DurableCommit.bend as Commitimport bend-kit-bytes@0.3.2.0/bytes.bend as Bytes# Host boundary for durable Db operations. Pure transitions and their laws live# in Db; filesystem ordering remains an explicitly unproved host interaction.# Handle wal tail in the database filesystem effects.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(Db.wal_path(dir))) _cp2 : Unit <- IO.try(Unit, CrashPoint.hit("wal.synced")) _prm : Unit <- IO.try(Unit, Fs.chmod(Db.wal_path(dir), U32.from_nat(384n))) return Done{res}def wal.initialize.fsync( result: Result<&1, &1, U32 & String, Unit>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match result: case Fail{error}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{error}) case Done{Unit{}}: Fs.chmod(path, U32.from_nat(384n))def wal.initialize.synced( result: Result<&1, &1, U32 & String, Unit>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match result: case Fail{error}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{error}) case Done{Unit{}}: do IO<Result<&1, &1, U32 & String, Unit>>: synced : Result<&1, &1, U32 & String, Unit> <- Fs.fsync(path) wal.initialize.fsync(synced, path)def wal.initialize.written( pair: File & Result<&1, &1, U32 & String, Unit>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match pair: case (file, result): do IO<Result<&1, &1, U32 & String, Unit>>: _closed : Unit <- File.close(file) wal.initialize.synced(result, path)def wal.initialize.opened( opened: Result<&1, &1, U32 & String, File>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match opened: case Fail{error}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{error}) case Done{file}: do IO<Result<&1, &1, U32 & String, Unit>>: written : File & Result<&1, &1, U32 & String, Unit> <- Fs.write_bytes(file, Wal.log_header()) wal.initialize.written(written, path)def wal.initialize.header(+path: String) -> IO(Result<&1, &1, U32 & String, Unit>): do IO<Result<&1, &1, U32 & String, Unit>>: opened : Result<&1, &1, U32 & String, File> <- File.open(path, "a") wal.initialize.opened(opened, path)def wal.initialize.header_if_empty( empty: Bool, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match empty: case True{}: wal.initialize.header(path) case False{}: IO.pure(Result<&1, &1, U32 & String, Unit>, Done{Unit{}})def wal.initialize.existing( result: Result<&1, &1, U32 & String, Nat>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match result: case Fail{error}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{error}) case Done{size}: wal.initialize.header_if_empty(Nat.is_eq(size, 0n), path)def wal.initialize.present( present: Bool, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match present: case False{}: wal.initialize.header(path) case True{}: do IO<Result<&1, &1, U32 & String, Unit>>: size : Result<&1, &1, U32 & String, Nat> <- Fs.file_size(path) wal.initialize.existing(size, path)def wal.initialize.exists( result: Result<&1, &1, U32 & String, Bool>, +path: String) -> IO(Result<&1, &1, U32 & String, Unit>): match result: case Fail{error}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{error}) case Done{present}: wal.initialize.present(present, path)# Creates the WAL only when it is absent.def wal_initialize(+dir: String) -> IO(Result<&1, &1, U32 & String, Unit>): do IO<Result<&1, &1, U32 & String, Unit>>: present : Result<&1, &1, U32 & String, Bool> <- Fs.exists(Db.wal_path(dir)) wal.initialize.exists(present, Db.wal_path(dir))# Handle wal encoded in the database filesystem effects.def wal_encoded( +dir: String, result: Result<&1, &1, Wal.Error, Bytes.Bytes>) -> IO(Result<&1, &1, U32 & String, Unit>): match result: case Fail{_}: IO.pure(Result<&1, &1, U32 & String, Unit>, Fail{(U32.from_nat(3n), "WAL batch cannot be encoded")}) case Done{frame}: do IO<Result<&1, &1, U32 & String, Unit>>: f : File <- IO.try(File, File.open(Db.wal_path(dir), "a")) fr : File & Result<&1, &1, U32 & String, Unit> <- Fs.write_bytes(f, frame) wal_tail(fr, dir)# Handle wal append in the database filesystem effects.def wal_append(+dir: String, batch: Wal.Batch) -> IO(Result<&1, &1, U32 & String, Unit>): wal_encoded(dir, Wal.encode_frame(batch))# Turns the fsync result into a durable or uncertain commit stage.def append_synced(result: Result<&1, &1, U32 & String, Unit>, file: File) -> IO(Commit.Stage): match result: case Fail{(code, message)}: do IO<Commit.Stage>: _closed : Unit <- File.close(file) return Commit.AppendUnknown{code, message} case Done{Unit{}}: do IO<Commit.Stage>: _checkpoint : Unit <- IO.try(Unit, CrashPoint.hit("wal.synced")) _closed : Unit <- File.close(file) return Commit.DurableAppend{}# Syncs the WAL after a successful append and closes the file.def append_written(pair: File & Result<&1, &1, U32 & String, Unit>, +path: String) -> IO(Commit.Stage): match pair: case (file, Fail{(code, message)}): do IO<Commit.Stage>: _closed : Unit <- File.close(file) return Commit.AppendUnknown{code, message} case (file, Done{Unit{}}): do IO<Commit.Stage>: _checkpoint : Unit <- IO.try(Unit, CrashPoint.hit("wal.appended")) synced : Result<&1, &1, U32 & String, Unit> <- Fs.fsync(path) append_synced(synced, file)# Rejects an open failure before append or writes the encoded frame.def append_opened(opened: Result<&1, &1, U32 & String, File>, frame: Bytes.Bytes, +path: String) -> IO(Commit.Stage): match opened: case Fail{(code, message)}: IO.pure(Commit.Stage, Commit.AppendRejected{code, message}) case Done{file}: do IO<Commit.Stage>: _checkpoint : Unit <- IO.try(Unit, CrashPoint.hit("wal.before_append")) written : File & Result<&1, &1, U32 & String, Unit> <- Fs.write_bytes(file, frame) append_written(written, path)# Opens the WAL and appends one encoded frame.def append_frame(+path: String, frame: Bytes.Bytes) -> IO(Commit.Stage): do IO<Commit.Stage>: opened : Result<&1, &1, U32 & String, File> <- File.open(path, "a") append_opened(opened, frame, path)# Rejects encoding errors before opening the WAL.def encoded_append( result: Result<&1, &1, Wal.Error, Bytes.Bytes>, +dir: String, muts: List<&2, Wal.Mut>,) -> IO(Commit.Stage & List<&2, Wal.Mut>): match result: case Fail{error}: IO.pure(Commit.Stage & List<&2, Wal.Mut>, (Commit.EncodingRejected{error}, muts)) case Done{frame}: do IO<Commit.Stage & List<&2, Wal.Mut>>: stage : Commit.Stage <- append_frame(Db.wal_path(dir), frame) return (stage, muts)# Encodes a batch and appends it to the WAL.def wal_commit(+dir: String, +batch: Wal.Batch) -> IO(Commit.Stage & List<&2, Wal.Mut>): match batch: case Wal.Batch{muts}: encoded_append(Wal.encode_frame(Wal.Batch{muts}), dir, muts)# Persists a batch before applying its pure commit transition.def db_write_durable(+db: Db.Db, +batch: Wal.Batch) -> IO(Db.Db & Commit.Stage): match db: case Db.Db{dir, _, _, _, _, _, _, _, _, _}: do IO<Db.Db & Commit.Stage>: pair : Commit.Stage & List<&2, Wal.Mut> <- wal_commit(dir, batch) return Commit.apply(pair, db)# Persists a batch and rotates the MemTable when needed.def db_write(+db: Db.Db, +batch: Wal.Batch) -> 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}: match batch: case Wal.Batch{muts}: do IO<Result<&1, &1, U32 & String, Db.Db>>: _res : Unit <- IO.try(Unit, wal_append(dir, Wal.Batch{muts})) return Db.rot_done(Db.apply_and_rotate(muts, mem, frozen, mem_count, frozen_count), dir, batch_cap, bcache, levels, flushed, manifest_token)# Writes one key/value pair through the legacy core path.def db_put(+db: Db.Db, +key: String, +val: String) -> IO(Result<&1, &1, U32 & String, Db.Db>): db_write(db, Wal.Batch{Con{Wal.Put{key, val}, Nil{}}})# Restores staged mutation order before writing the batch.def db_write_staged(+db: Db.Db, +staged: List<&2, Wal.Batch>) -> 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}: db_write(Db.Db{dir, mem, frozen, batch_cap, bcache, levels, flushed, manifest_token, mem_count, frozen_count}, Wal.Batch{Db.staged_muts(List.reverse(&2, Wal.Batch, staged))})# Writes one tombstone through the legacy core path.def db_del(+db: Db.Db, +key: String) -> IO(Result<&1, &1, U32 & String, Db.Db>): db_write(db, Wal.Batch{Con{Wal.Del{key}, Nil{}}})