◐ Off-By-One · answer catalog

go-wal-torn-tail-recovery-manifest-fencing

2 answer(s)gogo1.26gogo1.26

Problem class: go-wal-torn-tail-recovery-manifest-fencing

📦 Source in repository (JSON)

Answer 1

The complete, verified solution is written to ~/walstore/SOLUTION.md (32 KB). Here it is in full.


Crash Recovery for a Single-Writer WAL + SSTable Store

Problem class: go-wal-torn-tail-recovery-manifest-fencing Language / toolchain: Go 1.26, stdlib only (hash/crc32 Castagnoli, encoding/binary, os).

This is a complete, self-contained diagnosis and fix. Drop the files into a directory, run go test ./..., and the property test plus fencing test pass. All code was compiled and executed with go1.26.0 linux/amd64.


1. Symptoms

A single-writer store keeps an append-only WAL and periodically compacts into SSTables referenced by a manifest. After a crash (or after kill -9 mid-write) the store exhibits one or more of:

2. Root-Cause Analysis

Every symptom is a durability-boundary bug: the code conflates "bytes were handed to the kernel" with "bytes are durable and complete".

  1. No frame integrity boundary. A naive WAL writes records with no length/checksum, so a torn tail cannot be distinguished from a complete record. Fix: length-prefixed frames with a CRC32C over the payload, and a scanner that stops at the first frame that is short, mis-magiced, or fails CRC. The trailing region is the torn tail and is discarded.

  2. Batch spread across multiple frames. If a batch writes one frame per key, a crash can leave a subset applied. Fix: encode the whole batch (seq, count, then all records) as one frame. The frame either survives CRC validation (whole batch applied) or it does not (whole batch dropped).

  3. No contiguous sequence check on replay. Even with checksums, a stale or duplicated frame can be replayed. Fix: give every batch a global monotonically increasing seq; replay only accepts seq == prev+1, so recovered state is necessarily a prefix of committed batches.

  4. Manifest pointer treated as ground truth. The manifest is a hint, not a durability proof: it can be updated before the referenced SSTable body is on stable storage. Fix: every SSTable carries a CRC32C + magic footer; recovery scans all generations, validates each, and selects the newest checksum-clean generation.

  5. No writer fencing. "Single writer" must be enforced, not assumed. Fix: a durable epoch in the manifest is incremented (and fsynced) on every Open. Each mutating operation re-reads the manifest and returns a typed *StaleEpochError if the on-disk epoch moved past the handle's epoch.

Invariant achieved: recovered state = some prefix of committed batches; each batch atomic; seq strictly increasing and contiguous; newest durable generation chosen; stale writers rejected with a typed error.


3. Layout

<dir>/
  MANIFEST              magic | epoch | baseGen | baseSeq | crc32c
  sst/<20-digit gen>.sst   body | crc32c(body) | magic
  wal/000001.wal        repeated: magic | len | crc32c(payload) | payload

4. Exact Fix (source)

go.mod

module walstore

go 1.26

wal.go

package walstore

import (
    "encoding/binary"
    "errors"
    "hash/crc32"
)

// CRC32C (Castagnoli) is used, as required.
var crcTable = crc32.MakeTable(crc32.Castagnoli)

const (
    frameMagic     = uint32(0x57414c31) // "WAL1"
    frameHeaderLen = 12                 // magic(4) + length(4) + crc(4)
)

// ErrTornTail indicates the WAL ended in an incomplete or corrupt frame.
var ErrTornTail = errors.New("walstore: torn or corrupt WAL frame")

// Record is a single key/value mutation inside a batch.
type Record struct {
    Key string
    Val string
}

// Batch is an atomically applied group of records carrying a monotonic seq.
type Batch struct {
    Seq     uint64
    Records []Record
}

// encodeBatch serialises a batch into a single WAL frame payload.
func encodeBatch(b Batch) []byte {
    // seq(8) + count(4) + sum(4 + len(k) + 4 + len(v))
    n := 12
    for _, r := range b.Records {
        n += 8 + len(r.Key) + len(r.Val)
    }
    buf := make([]byte, n)
    binary.LittleEndian.PutUint64(buf[0:8], b.Seq)
    binary.LittleEndian.PutUint32(buf[8:12], uint32(len(b.Records)))
    off := 12
    for _, r := range b.Records {
        binary.LittleEndian.PutUint32(buf[off:off+4], uint32(len(r.Key)))
        off += 4
        copy(buf[off:], r.Key)
        off += len(r.Key)
        binary.LittleEndian.PutUint32(buf[off:off+4], uint32(len(r.Val)))
        off += 4
        copy(buf[off:], r.Val)
        off += len(r.Val)
    }
    return buf
}

// decodeBatch parses a frame payload into a Batch.
func decodeBatch(p []byte) (Batch, error) {
    if len(p) < 12 {
        return Batch{}, ErrTornTail
    }
    var b Batch
    b.Seq = binary.LittleEndian.Uint64(p[0:8])
    count := binary.LittleEndian.Uint32(p[8:12])
    if count > uint32(len(p))/8+1 { // cheap sanity guard against absurd counts
        return Batch{}, ErrTornTail
    }
    off := 12
    for i := uint32(0); i < count; i++ {
        if off+4 > len(p) {
            return Batch{}, ErrTornTail
        }
        kl := binary.LittleEndian.Uint32(p[off : off+4])
        off += 4
        if off+int(kl) > len(p) {
            return Batch{}, ErrTornTail
        }
        k := string(p[off : off+int(kl)])
        off += int(kl)
        if off+4 > len(p) {
            return Batch{}, ErrTornTail
        }
        vl := binary.LittleEndian.Uint32(p[off : off+4])
        off += 4
        if off+int(vl) > len(p) {
            return Batch{}, ErrTornTail
        }
        v := string(p[off : off+int(vl)])
        off += int(vl)
        b.Records = append(b.Records, Record{Key: k, Val: v})
    }
    if off != len(p) {
        return Batch{}, ErrTornTail
    }
    return b, nil
}

// appendFrame writes a length-prefixed, CRC32C-protected frame.
func appendFrame(dst []byte, payload []byte) []byte {
    var hdr [frameHeaderLen]byte
    binary.LittleEndian.PutUint32(hdr[0:4], frameMagic)
    binary.LittleEndian.PutUint32(hdr[4:8], uint32(len(payload)))
    binary.LittleEndian.PutUint32(hdr[8:12], crc32.Checksum(payload, crcTable))
    dst = append(dst, hdr[:]...)
    dst = append(dst, payload...)
    return dst
}

// scanFrames returns every fully-valid batch encoded in data, in order, and the
// number of leading bytes consumed. The scan stops at the first frame that is
// incomplete or fails its CRC. The trailing region is the torn tail and is
// discarded on recovery.
func scanFrames(data []byte) (batches []Batch, ends []int, consumed int, torn bool) {
    off := 0
    for {
        if off == len(data) {
            return batches, ends, off, false
        }
        if off+frameHeaderLen > len(data) {
            return batches, ends, off, true
        }
        magic := binary.LittleEndian.Uint32(data[off : off+4])
        if magic != frameMagic {
            return batches, ends, off, true
        }
        l := binary.LittleEndian.Uint32(data[off+4 : off+8])
        want := binary.LittleEndian.Uint32(data[off+8 : off+12])
        start := off + frameHeaderLen
        end := start + int(l)
        if l > uint32(len(data)) || end > len(data) {
            return batches, ends, off, true
        }
        payload := data[start:end]
        if crc32.Checksum(payload, crcTable) != want {
            return batches, ends, off, true
        }
        b, err := decodeBatch(payload)
        if err != nil {
            return batches, ends, off, true
        }
        batches = append(batches, b)
        ends = append(ends, end)
        off = end
    }
}

manifest.go

package walstore

import (
    "encoding/binary"
    "errors"
    "hash/crc32"
    "os"
    "path/filepath"
)

// Manifest is the durable fence record. Epoch is bumped by every writer
// acquisition; any writer holding a lower epoch is stale and must stop.
type Manifest struct {
    Epoch   uint64
    BaseGen uint64
    BaseSeq uint64
}

const manifestMagic = uint32(0x4d414e31) // "MAN1"

func manifestPath(dir string) string { return filepath.Join(dir, "MANIFEST") }

func encodeManifest(m Manifest) []byte {
    var buf [32]byte
    binary.LittleEndian.PutUint32(buf[0:4], manifestMagic)
    binary.LittleEndian.PutUint64(buf[4:12], m.Epoch)
    binary.LittleEndian.PutUint64(buf[12:20], m.BaseGen)
    binary.LittleEndian.PutUint64(buf[20:28], m.BaseSeq)
    binary.LittleEndian.PutUint32(buf[28:32], crc32.Checksum(buf[0:28], crcTable))
    out := make([]byte, 32)
    copy(out, buf[:])
    return out
}

func decodeManifest(data []byte) (Manifest, error) {
    if len(data) != 32 {
        return Manifest{}, errors.New("walstore: bad manifest length")
    }
    if binary.LittleEndian.Uint32(data[0:4]) != manifestMagic {
        return Manifest{}, errors.New("walstore: bad manifest magic")
    }
    if crc32.Checksum(data[0:28], crcTable) != binary.LittleEndian.Uint32(data[28:32]) {
        return Manifest{}, errors.New("walstore: manifest checksum mismatch")
    }
    return Manifest{
        Epoch:   binary.LittleEndian.Uint64(data[4:12]),
        BaseGen: binary.LittleEndian.Uint64(data[12:20]),
        BaseSeq: binary.LittleEndian.Uint64(data[20:28]),
    }, nil
}

func readManifest(dir string) (Manifest, error) {
    raw, err := os.ReadFile(manifestPath(dir))
    if err != nil {
        return Manifest{}, err
    }
    return decodeManifest(raw)
}

// writeManifestFsync atomically replaces the manifest and fsyncs both the file
// and its parent directory so the epoch bump is durable before use.
func writeManifestFsync(dir string, m Manifest) error {
    tmp := manifestPath(dir) + ".tmp"
    f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
    if err != nil {
        return err
    }
    if _, err := f.Write(encodeManifest(m)); err != nil {
        f.Close()
        return err
    }
    if err := f.Sync(); err != nil {
        f.Close()
        return err
    }
    if err := f.Close(); err != nil {
        return err
    }
    if err := os.Rename(tmp, manifestPath(dir)); err != nil {
        return err
    }
    return fsyncDir(dir)
}

func fsyncDir(dir string) error {
    d, err := os.Open(dir)
    if err != nil {
        return err
    }
    defer d.Close()
    return d.Sync()
}

sst.go

package walstore

import (
    "encoding/binary"
    "errors"
    "fmt"
    "hash/crc32"
    "os"
    "path/filepath"
    "sort"
    "strconv"
    "strings"
)

const sstMagic = uint32(0x53535431) // "SST1"

// SSTable is an immutable, checksummed sorted snapshot at a generation.
type SSTable struct {
    Gen uint64
    Seq uint64 // highest sequence number included
    KV  map[string]string
}

// encodeSST serialises the table plus a footer protected by CRC32C.
func encodeSST(t SSTable) []byte {
    keys := make([]string, 0, len(t.KV))
    for k := range t.KV {
        keys = append(keys, k)
    }
    sort.Strings(keys)
    buf := make([]byte, 0, 16+len(keys)*16)
    var tmp [8]byte
    binary.LittleEndian.PutUint64(tmp[:], t.Seq)
    buf = append(buf, tmp[:]...)
    binary.LittleEndian.PutUint32(tmp[:4], uint32(len(keys)))
    buf = append(buf, tmp[:4]...)
    for _, k := range keys {
        v := t.KV[k]
        binary.LittleEndian.PutUint32(tmp[:4], uint32(len(k)))
        buf = append(buf, tmp[:4]...)
        buf = append(buf, k...)
        binary.LittleEndian.PutUint32(tmp[:4], uint32(len(v)))
        buf = append(buf, tmp[:4]...)
        buf = append(buf, v...)
    }
    // footer: crc(4) over all preceding bytes, then magic(4)
    var crcBuf [4]byte
    binary.LittleEndian.PutUint32(crcBuf[:], crc32.Checksum(buf, crcTable))
    buf = append(buf, crcBuf[:]...)
    binary.LittleEndian.PutUint32(crcBuf[:], sstMagic)
    buf = append(buf, crcBuf[:]...)
    return buf
}

// decodeSST validates and parses an SSTable. A truncated or corrupt table
// (e.g. an fsync that never completed) returns an error.
func decodeSST(data []byte) (SSTable, error) {
    if len(data) < 16 {
        return SSTable{}, errors.New("walstore: short SSTable")
    }
    body := data[:len(data)-8]
    footer := data[len(data)-8:]
    if binary.LittleEndian.Uint32(footer[4:8]) != sstMagic {
        return SSTable{}, errors.New("walstore: bad SSTable magic")
    }
    if crc32.Checksum(body, crcTable) != binary.LittleEndian.Uint32(footer[0:4]) {
        return SSTable{}, errors.New("walstore: SSTable checksum mismatch")
    }
    var t SSTable
    t.Seq = binary.LittleEndian.Uint64(body[0:8])
    count := binary.LittleEndian.Uint32(body[8:12])
    off := 12
    t.KV = make(map[string]string, count)
    for i := uint32(0); i < count; i++ {
        if off+4 > len(body) {
            return SSTable{}, errors.New("walstore: truncated SSTable body")
        }
        kl := binary.LittleEndian.Uint32(body[off : off+4])
        off += 4
        if off+int(kl) > len(body) {
            return SSTable{}, errors.New("walstore: truncated key")
        }
        k := string(body[off : off+int(kl)])
        off += int(kl)
        if off+4 > len(body) {
            return SSTable{}, errors.New("walstore: truncated SSTable body")
        }
        vl := binary.LittleEndian.Uint32(body[off : off+4])
        off += 4
        if off+int(vl) > len(body) {
            return SSTable{}, errors.New("walstore: truncated value")
        }
        v := string(body[off : off+int(vl)])
        off += int(vl)
        t.KV[k] = v
    }
    if off != len(body) {
        return SSTable{}, errors.New("walstore: trailing SSTable bytes")
    }
    return t, nil
}

// maxSSTGen returns the highest generation number present on disk, valid or
// not, so fresh writes never reuse a generation slot.
func maxSSTGen(dir string) uint64 {
    entries, err := os.ReadDir(filepath.Join(dir, "sst"))
    if err != nil {
        return 0
    }
    var max uint64
    for _, e := range entries {
        name := e.Name()
        if !strings.HasSuffix(name, ".sst") {
            continue
        }
        gen, err := strconv.ParseUint(strings.TrimSuffix(name, ".sst"), 10, 64)
        if err != nil {
            continue
        }
        if gen > max {
            max = gen
        }
    }
    return max
}

// sstPath returns the canonical path of a generation's SSTable.
func sstPath(dir string, gen uint64) string {
    return filepath.Join(dir, "sst", fmt.Sprintf("%020d.sst", gen))
}

// scanSSTGenerations validates every generation present on disk and returns the
// newest one whose contents are durable and checksum-clean. The manifest's
// pointer is only a hint; this is the authoritative durability check.
func scanSSTGenerations(dir string) (SSTable, bool, error) {
    entries, err := os.ReadDir(filepath.Join(dir, "sst"))
    if err != nil {
        if os.IsNotExist(err) {
            return SSTable{}, false, nil
        }
        return SSTable{}, false, err
    }
    var gens []uint64
    for _, e := range entries {
        name := e.Name()
        if !strings.HasSuffix(name, ".sst") {
            continue
        }
        gen, err := strconv.ParseUint(strings.TrimSuffix(name, ".sst"), 10, 64)
        if err != nil {
            continue
        }
        gens = append(gens, gen)
    }
    sort.Slice(gens, func(i, j int) bool { return gens[i] > gens[j] })
    for _, gen := range gens {
        raw, err := os.ReadFile(sstPath(dir, gen))
        if err != nil {
            continue
        }
        t, err := decodeSST(raw)
        if err != nil {
            continue // incomplete fsync / torn generation: skip
        }
        t.Gen = gen
        return t, true, nil
    }
    return SSTable{}, false, nil
}

store.go

package walstore

import (
    "errors"
    "fmt"
    "os"
    "path/filepath"
    "sync"
)

// StaleEpochError is returned when a writer tries to mutate the store while a
// newer writer has fenced it by advancing the manifest epoch.
type StaleEpochError struct {
    Held   uint64
    OnDisk uint64
}

func (e *StaleEpochError) Error() string {
    return fmt.Sprintf("walstore: stale manifest epoch: writer holds %d, on-disk is %d", e.Held, e.OnDisk)
}

// Store is the single-writer handle over one directory.
type Store struct {
    mu       sync.Mutex
    dir      string
    epoch    uint64
    seq      uint64
    nextGen  uint64
    mem      map[string]string
    wal      *os.File
    done     bool
}

// Recovered is the result of replaying a WAL image.
type Recovered struct {
    Mem      map[string]string
    Seq      uint64
    Applied  []Batch
    Consumed int
    Torn     bool
}

// recoverState replays WAL frames on top of a base snapshot. Only batches with
// contiguous, strictly-increasing sequence numbers are applied; the first
// incomplete, corrupt, or out-of-order frame terminates replay (torn tail).
func recoverState(raw []byte, base map[string]string, baseSeq uint64) (Recovered, error) {
    mem := make(map[string]string, len(base))
    for k, v := range base {
        mem[k] = v
    }
    batches, ends, consumed, torn := scanFrames(raw)
    res := Recovered{Mem: mem, Seq: baseSeq, Consumed: consumed, Torn: torn}
    prev := baseSeq
    applied := 0
    for _, b := range batches {
        if b.Seq != prev+1 {
            // Stop before the out-of-order frame; only the validated
            // prefix may be physically retained.
            res.Torn = true
            if applied == 0 {
                res.Consumed = 0
            } else {
                res.Consumed = ends[applied-1]
            }
            return res, nil
        }
        for _, r := range b.Records {
            mem[r.Key] = r.Val
        }
        prev = b.Seq
        res.Seq = b.Seq
        res.Applied = append(res.Applied, b)
        applied++
    }
    return res, nil
}

// Open recovers the store and acquires a fresh writer epoch. Any previously
// open writer becomes stale.
func Open(dir string) (*Store, error) {
    if err := os.MkdirAll(filepath.Join(dir, "sst"), 0o755); err != nil {
        return nil, err
    }
    if err := os.MkdirAll(filepath.Join(dir, "wal"), 0o755); err != nil {
        return nil, err
    }

    m, err := readManifest(dir)
    if err != nil {
        if !os.IsNotExist(err) {
            return nil, err
        }
        m = Manifest{}
    }

    sst, ok, err := scanSSTGenerations(dir)
    if err != nil {
        return nil, err
    }
    var base map[string]string
    var baseSeq uint64
    var baseGen uint64
    if ok {
        base = sst.KV
        baseSeq = sst.Seq
        baseGen = sst.Gen
    } else {
        base = map[string]string{}
    }

    walFile := filepath.Join(dir, "wal", "000001.wal")
    raw, err := os.ReadFile(walFile)
    if err != nil && !os.IsNotExist(err) {
        return nil, err
    }
    rec, err := recoverState(raw, base, baseSeq)
    if err != nil {
        return nil, err
    }

    // Physically drop the torn tail so subsequent appends are clean.
    f, err := os.OpenFile(walFile, os.O_CREATE|os.O_RDWR, 0o644)
    if err != nil {
        return nil, err
    }
    if err := f.Truncate(int64(rec.Consumed)); err != nil {
        f.Close()
        return nil, err
    }
    if _, err := f.Seek(0, os.SEEK_END); err != nil {
        f.Close()
        return nil, err
    }

    s := &Store{
        dir:     dir,
        epoch:   m.Epoch + 1,
        seq:     rec.Seq,
        nextGen: maxSSTGen(dir) + 1,
        mem:     rec.Mem,
        wal:     f,
    }
    // Persist the epoch fence before returning the writer.
    if err := writeManifestFsync(dir, Manifest{Epoch: s.epoch, BaseGen: baseGen, BaseSeq: rec.Seq}); err != nil {
        f.Close()
        return nil, err
    }
    return s, nil
}

// Epoch returns the writer epoch held by this handle.
func (s *Store) Epoch() uint64 {
    s.mu.Lock()
    defer s.mu.Unlock()
    return s.epoch
}

// Seq returns the highest applied sequence number.
func (s *Store) Seq() uint64 {
    s.mu.Lock()
    defer s.mu.Unlock()
    return s.seq
}

// Get returns the current value for a key.
func (s *Store) Get(key string) (string, bool) {
    s.mu.Lock()
    defer s.mu.Unlock()
    v, ok := s.mem[key]
    return v, ok
}

// checkEpochLocked fails if another writer advanced the manifest epoch.
func (s *Store) checkEpochLocked() error {
    m, err := readManifest(s.dir)
    if err != nil {
        return err
    }
    if m.Epoch != s.epoch {
        return &StaleEpochError{Held: s.epoch, OnDisk: m.Epoch}
    }
    return nil
}

// Append atomically applies one batch, fsyncing its single WAL frame first.
func (s *Store) Append(records []Record) error {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return errors.New("walstore: store is closed")
    }
    if err := s.checkEpochLocked(); err != nil {
        return err
    }
    batch := Batch{Seq: s.seq + 1, Records: records}
    frame := appendFrame(nil, encodeBatch(batch))
    if _, err := s.wal.Write(frame); err != nil {
        return err
    }
    if err := s.wal.Sync(); err != nil {
        return err
    }
    for _, r := range records {
        s.mem[r.Key] = r.Val
    }
    s.seq = batch.Seq
    return nil
}

// Flush writes the current in-memory state as a new checksummed SSTable
// generation, fsyncs it, then publishes it in the manifest. The manifest
// pointer is only a hint; recovery re-validates every generation.
func (s *Store) Flush() (uint64, error) {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return 0, errors.New("walstore: store is closed")
    }
    if err := s.checkEpochLocked(); err != nil {
        return 0, err
    }
    gen := s.nextGen
    table := SSTable{Gen: gen, Seq: s.seq, KV: s.mem}
    raw := encodeSST(table)

    tmp := sstPath(s.dir, gen) + ".tmp"
    f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
    if err != nil {
        return 0, err
    }
    if _, err := f.Write(raw); err != nil {
        f.Close()
        return 0, err
    }
    if err := f.Sync(); err != nil {
        f.Close()
        return 0, err
    }
    if err := f.Close(); err != nil {
        return 0, err
    }
    if err := os.Rename(tmp, sstPath(s.dir, gen)); err != nil {
        return 0, err
    }
    if err := fsyncDir(filepath.Join(s.dir, "sst")); err != nil {
        return 0, err
    }
    if err := writeManifestFsync(s.dir, Manifest{Epoch: s.epoch, BaseGen: gen, BaseSeq: s.seq}); err != nil {
        return 0, err
    }
    s.nextGen = gen + 1
    // The WAL prefix covered by this SSTable is no longer needed.
    if err := s.wal.Truncate(0); err != nil {
        return 0, err
    }
    if _, err := s.wal.Seek(0, os.SEEK_END); err != nil {
        return 0, err
    }
    if err := s.wal.Sync(); err != nil {
        return 0, err
    }
    return gen, nil
}

// Close releases the writer handle.
func (s *Store) Close() error {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return nil
    }
    s.done = true
    return s.wal.Close()
}

store_test.go

package walstore

import (
    "errors"
    "fmt"
    "math/rand"
    "os"
    "path/filepath"
    "testing"
)

// frameLen returns the on-disk size of one encoded batch frame.
func frameLen(b Batch) int {
    return frameHeaderLen + len(encodeBatch(b))
}

// buildWAL encodes batches into a raw WAL image and returns per-frame end
// offsets.
func buildWAL(batches []Batch) ([]byte, []int) {
    var raw []byte
    ends := make([]int, len(batches))
    for i, b := range batches {
        raw = appendFrame(raw, encodeBatch(b))
        ends[i] = len(raw)
    }
    return raw, ends
}

func modelStates(batches []Batch) []map[string]string {
    states := make([]map[string]string, len(batches)+1)
    states[0] = map[string]string{}
    for i, b := range batches {
        m := make(map[string]string, len(states[i]))
        for k, v := range states[i] {
            m[k] = v
        }
        for _, r := range b.Records {
            m[r.Key] = r.Val
        }
        states[i+1] = m
    }
    return states
}

func randomBatches(rng *rand.Rand, n int) []Batch {
    batches := make([]Batch, n)
    for i := range batches {
        recs := make([]Record, 1+rng.Intn(4))
        for j := range recs {
            recs[j] = Record{
                Key: fmt.Sprintf("k%d_%d", i, rng.Intn(3)),
                Val: fmt.Sprintf("v%d_%d", i, rng.Intn(1000)),
            }
        }
        batches[i] = Batch{Seq: uint64(i + 1), Records: recs}
    }
    return batches
}

// TestTornTailPrefixProperty is the required property test: at least 200 random
// truncation offsets crossed with random batch boundaries. For every image the
// recovered state must equal some prefix of committed batches, never a partial
// batch and never a non-prefix.
func TestTornTailPrefixProperty(t *testing.T) {
    rng := rand.New(rand.NewSource(0xC0FFEE))
    const seeds = 8
    checks := 0

    for seed := 0; seed < seeds; seed++ {
        batches := randomBatches(rng, 1+rng.Intn(40))
        raw, ends := buildWAL(batches)
        states := modelStates(batches)

        // 200+ offsets per seed, plus every frame boundary.
        offsets := make([]int, 0, 260)
        for i := 0; i < 256; i++ {
            offsets = append(offsets, rng.Intn(len(raw)+1))
        }
        for _, e := range ends {
            offsets = append(offsets, e, e-1)
        }

        for _, off := range offsets {
            if off < 0 {
                continue
            }
            if off > len(raw) {
                continue
            }
            checks++
            rec, err := recoverState(raw[:off], map[string]string{}, 0)
            if err != nil {
                t.Fatalf("recover error: %v", err)
            }

            // Highest fully-contained frame index.
            wantApplied := 0
            for i, e := range ends {
                if e <= off {
                    wantApplied = i + 1
                }
            }
            if len(rec.Applied) != wantApplied {
                t.Fatalf("seed=%d off=%d: applied=%d want=%d", seed, off, len(rec.Applied), wantApplied)
            }
            // Sequence numbers must be strictly monotonic & contiguous.
            for i, b := range rec.Applied {
                if b.Seq != uint64(i+1) {
                    t.Fatalf("seed=%d off=%d: non-monotonic seq at %d: %d", seed, off, i, b.Seq)
                }
            }
            // State must equal the exact committed prefix.
            if !equalKV(rec.Mem, states[wantApplied]) {
                t.Fatalf("seed=%d off=%d: state mismatch\n got=%v\nwant=%v", seed, off, rec.Mem, states[wantApplied])
            }
            if rec.Seq != uint64(wantApplied) {
                t.Fatalf("seed=%d off=%d: seq=%d want=%d", seed, off, rec.Seq, wantApplied)
            }
        }
    }
    if checks < 200 {
        t.Fatalf("only %d checks, need >= 200", checks)
    }
    t.Logf("verified %d truncation offsets across %d random batch streams", checks, seeds)
}

// TestAtomicBatchDiscard proves a torn multi-key batch is fully discarded.
func TestAtomicBatchDiscard(t *testing.T) {
    recs := []Record{{Key: "a", Val: "1"}, {Key: "b", Val: "2"}, {Key: "c", Val: "3"}}
    b := Batch{Seq: 1, Records: recs}
    full := appendFrame(nil, encodeBatch(b))
    // Cut one byte before the CRC-valid end: the entire batch must vanish.
    rec, err := recoverState(full[:len(full)-1], map[string]string{}, 0)
    if err != nil {
        t.Fatal(err)
    }
    if len(rec.Applied) != 0 || len(rec.Mem) != 0 {
        t.Fatalf("partial batch applied: %+v", rec)
    }
    // With the last byte present, all keys appear together.
    rec, err = recoverState(full, map[string]string{}, 0)
    if err != nil {
        t.Fatal(err)
    }
    if len(rec.Mem) != 3 {
        t.Fatalf("atomic batch not fully applied: %+v", rec)
    }
}

// TestStaleEpochWriterRejected is the required fencing test.
func TestStaleEpochWriterRejected(t *testing.T) {
    dir := t.TempDir()
    a, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer a.Close()
    if err := a.Append([]Record{{Key: "x", Val: "1"}}); err != nil {
        t.Fatal(err)
    }

    // A second writer acquires a newer epoch.
    b, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer b.Close()
    if b.Epoch() <= a.Epoch() {
        t.Fatalf("epoch did not advance: a=%d b=%d", a.Epoch(), b.Epoch())
    }

    // The old writer must be fenced with a typed error.
    err = a.Append([]Record{{Key: "x", Val: "2"}})
    var se *StaleEpochError
    if !errors.As(err, &se) {
        t.Fatalf("want *StaleEpochError, got %T: %v", err, err)
    }
    // The new writer still works.
    if err := b.Append([]Record{{Key: "x", Val: "3"}}); err != nil {
        t.Fatal(err)
    }
}

// TestNewestDurableGenerationSelected proves recovery skips an SSTable whose
// fsync never completed (truncated/corrupt) even though the manifest points at
// it, and falls back to the newest checksum-clean generation.
func TestNewestDurableGenerationSelected(t *testing.T) {
    dir := t.TempDir()
    s, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    if err := s.Append([]Record{{Key: "k", Val: "durable"}}); err != nil {
        t.Fatal(err)
    }
    if _, err := s.Flush(); err != nil { // generation 1
        t.Fatal(err)
    }
    if err := s.Close(); err != nil {
        t.Fatal(err)
    }

    // Simulate a generation 2 whose write/fsync never completed: truncated.
    if err := os.WriteFile(sstPath(dir, 2), []byte("half-written"), 0o644); err != nil {
        t.Fatal(err)
    }
    // Manifest dishonestly points at the torn generation 2.
    if err := writeManifestFsync(dir, Manifest{Epoch: 1, BaseGen: 2, BaseSeq: 1}); err != nil {
        t.Fatal(err)
    }

    s2, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer s2.Close()
    v, ok := s2.Get("k")
    if !ok || v != "durable" {
        t.Fatalf("recovery picked torn generation: got %q ok=%v", v, ok)
    }
    // seq must not jump to the phantom generation.
    if s2.Seq() != 1 {
        t.Fatalf("seq=%d want 1", s2.Seq())
    }
    // New flush must not reuse generation 2 (it exists on disk).
    gen, err := s2.Flush()
    if err != nil {
        t.Fatal(err)
    }
    if gen <= 2 {
        t.Fatalf("reused generation %d", gen)
    }
}

// TestReopenRecoversAndTruncatesTornTail exercises the full Open path.
func TestReopenRecoversAndTruncatesTornTail(t *testing.T) {
    dir := t.TempDir()
    s, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    for i := 0; i < 5; i++ {
        if err := s.Append([]Record{{Key: "n", Val: fmt.Sprintf("%d", i)}}); err != nil {
            t.Fatal(err)
        }
    }
    if err := s.Close(); err != nil {
        t.Fatal(err)
    }

    // Append junk to the live WAL, then corrupt the last valid frame's CRC.
    walFile := filepath.Join(dir, "wal", "000001.wal")
    f, err := os.OpenFile(walFile, os.O_APPEND|os.O_WRONLY, 0o644)
    if err != nil {
        t.Fatal(err)
    }
    if _, err := f.Write([]byte("torn-garbage")); err != nil {
        t.Fatal(err)
    }
    f.Close()

    s2, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer s2.Close()
    if got, _ := s2.Get("n"); got != "4" {
        t.Fatalf("recovered n=%q want 4", got)
    }
    if s2.Seq() != 5 {
        t.Fatalf("seq=%d want 5", s2.Seq())
    }
    raw, err := os.ReadFile(walFile)
    if err != nil {
        t.Fatal(err)
    }
    // The torn garbage must be gone, and the 5 valid frames retained so they
    // can be replayed on the next open.
    if string(raw) == "torn-garbage" || len(raw) == 0 {
        t.Fatalf("torn tail not cleaned: %q", raw)
    }
    batches, _, _, torn := scanFrames(raw)
    if torn || len(batches) != 5 {
        t.Fatalf("expected 5 clean frames, got %d torn=%v", len(batches), torn)
    }
}

func equalKV(a, b map[string]string) bool {
    if len(a) != len(b) {
        return false
    }
    for k, v := range a {
        if b[k] != v {
            return false
        }
    }
    return true
}

5. Verification

go vet ./...
go test -race -count=3 ./...

Observed on go1.26.0 linux/amd64:

$ go test -v ./...
=== RUN   TestTornTailPrefixProperty
    store_test.go:124: verified 2370 truncation offsets across 8 random batch streams
--- PASS: TestTornTailPrefixProperty (0.02s)
=== RUN   TestAtomicBatchDiscard
--- PASS: TestAtomicBatchDiscard (0.00s)
=== RUN   TestStaleEpochWriterRejected
--- PASS: TestStaleEpochWriterRejected (0.00s)
=== RUN   TestNewestDurableGenerationSelected
--- PASS: TestNewestDurableGenerationSelected (0.00s)
=== RUN   TestReopenRecoversAndTruncatesTornTail
--- PASS: TestReopenRecoversAndTruncatesTornTail (0.00s)
PASS
ok      walstore        0.031s

Evidence map

Requirement Test
≥200 random truncation offsets × random batch boundaries; recovered state is a committed prefix TestTornTailPrefixProperty (2,370 offsets over 8 random streams)
Multi-key batch fully applied or fully discarded TestAtomicBatchDiscard
Stale-epoch writer rejected with a typed error TestStaleEpochWriterRejected (errors.As → *StaleEpochError)
Newest durable generation selected when manifest points at an un-fsynced SSTable TestNewestDurableGenerationSelected
End-to-end reopen drops the torn tail and replays the valid prefix TestReopenRecoversAndTruncatesTornTail

How the property test works

  1. For each random seed, generate a stream of batches with random sizes and random keys/values; encode them into one raw WAL image.
  2. Record each frame's end offset and the exact model state after each prefix.
  3. Cross 256 uniformly random truncation offsets with every frame boundary (±1 byte), yielding 2,370 evaluated images.
  4. recoverState must (a) apply exactly the frames fully contained before the cut, (b) never apply a partial batch, (c) keep seq contiguous and strictly increasing, and (d) produce a map equal to the corresponding model prefix. Any deviation is a hard t.Fatalf.

Why the fencing test is meaningful

Open fsyncs Manifest{Epoch: old+1} before returning. Writer A then holds epoch e. Opening writer B publishes e+1. A.Append re-reads the manifest, sees e+1, and returns *StaleEpochError before touching the WAL; B.Append succeeds. Fencing is therefore durable and checked on every mutation, not just at startup.


The artifact is saved at ~/walstore/ (wal.go, manifest.go, sst.go, store.go, store_test.go, go.mod, and SOLUTION.md). Run cd ~/walstore && go test -race ./... to reproduce.

Evidence & signatures

# Evidence
- Problem class: go-wal-torn-tail-recovery-manifest-fencing
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-27T16:32:27.219Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement crash recovery for a single-writer WAL plus SSTable store in Go: length-prefixed CRC32C frames, replay after the WAL is truncated at an arbitrary byte offset, multi-key batches that are either fully applied or fully discarded, and strictly monotonic sequence numbers. The manifest may point at an SSTable whose fsync never completed, so recovery must select the newest durable generation and reject a writer still holding a stale manifest epoch with a typed error. Evidence: a property test over at least 200 random truncation offsets crossed with random batch boundaries proving the recovered state equals some prefix of committed batches, plus a stale-epoch writer rejection test.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-wal-torn-tail-recovery-manifest-fencing", "provider": "openrouter", "solved_at": "2026-09-27T16:32:27.221Z", "version": "1.26"}

Answer 2

The complete, verified solution is written to ~/walstore/SOLUTION.md (32 KB). Here it is in full.


Crash Recovery for a Single-Writer WAL + SSTable Store

Problem class: go-wal-torn-tail-recovery-manifest-fencing Language / toolchain: Go 1.26, stdlib only (hash/crc32 Castagnoli, encoding/binary, os).

This is a complete, self-contained diagnosis and fix. Drop the files into a directory, run go test ./..., and the property test plus fencing test pass. All code was compiled and executed with go1.26.0 linux/amd64.


1. Symptoms

A single-writer store keeps an append-only WAL and periodically compacts into SSTables referenced by a manifest. After a crash (or after kill -9 mid-write) the store exhibits one or more of:

2. Root-Cause Analysis

Every symptom is a durability-boundary bug: the code conflates "bytes were handed to the kernel" with "bytes are durable and complete".

  1. No frame integrity boundary. A naive WAL writes records with no length/checksum, so a torn tail cannot be distinguished from a complete record. Fix: length-prefixed frames with a CRC32C over the payload, and a scanner that stops at the first frame that is short, mis-magiced, or fails CRC. The trailing region is the torn tail and is discarded.

  2. Batch spread across multiple frames. If a batch writes one frame per key, a crash can leave a subset applied. Fix: encode the whole batch (seq, count, then all records) as one frame. The frame either survives CRC validation (whole batch applied) or it does not (whole batch dropped).

  3. No contiguous sequence check on replay. Even with checksums, a stale or duplicated frame can be replayed. Fix: give every batch a global monotonically increasing seq; replay only accepts seq == prev+1, so recovered state is necessarily a prefix of committed batches.

  4. Manifest pointer treated as ground truth. The manifest is a hint, not a durability proof: it can be updated before the referenced SSTable body is on stable storage. Fix: every SSTable carries a CRC32C + magic footer; recovery scans all generations, validates each, and selects the newest checksum-clean generation.

  5. No writer fencing. "Single writer" must be enforced, not assumed. Fix: a durable epoch in the manifest is incremented (and fsynced) on every Open. Each mutating operation re-reads the manifest and returns a typed *StaleEpochError if the on-disk epoch moved past the handle's epoch.

Invariant achieved: recovered state = some prefix of committed batches; each batch atomic; seq strictly increasing and contiguous; newest durable generation chosen; stale writers rejected with a typed error.


3. Layout

<dir>/
  MANIFEST              magic | epoch | baseGen | baseSeq | crc32c
  sst/<20-digit gen>.sst   body | crc32c(body) | magic
  wal/000001.wal        repeated: magic | len | crc32c(payload) | payload

4. Exact Fix (source)

go.mod

module walstore

go 1.26

wal.go

package walstore

import (
    "encoding/binary"
    "errors"
    "hash/crc32"
)

// CRC32C (Castagnoli) is used, as required.
var crcTable = crc32.MakeTable(crc32.Castagnoli)

const (
    frameMagic     = uint32(0x57414c31) // "WAL1"
    frameHeaderLen = 12                 // magic(4) + length(4) + crc(4)
)

// ErrTornTail indicates the WAL ended in an incomplete or corrupt frame.
var ErrTornTail = errors.New("walstore: torn or corrupt WAL frame")

// Record is a single key/value mutation inside a batch.
type Record struct {
    Key string
    Val string
}

// Batch is an atomically applied group of records carrying a monotonic seq.
type Batch struct {
    Seq     uint64
    Records []Record
}

// encodeBatch serialises a batch into a single WAL frame payload.
func encodeBatch(b Batch) []byte {
    // seq(8) + count(4) + sum(4 + len(k) + 4 + len(v))
    n := 12
    for _, r := range b.Records {
        n += 8 + len(r.Key) + len(r.Val)
    }
    buf := make([]byte, n)
    binary.LittleEndian.PutUint64(buf[0:8], b.Seq)
    binary.LittleEndian.PutUint32(buf[8:12], uint32(len(b.Records)))
    off := 12
    for _, r := range b.Records {
        binary.LittleEndian.PutUint32(buf[off:off+4], uint32(len(r.Key)))
        off += 4
        copy(buf[off:], r.Key)
        off += len(r.Key)
        binary.LittleEndian.PutUint32(buf[off:off+4], uint32(len(r.Val)))
        off += 4
        copy(buf[off:], r.Val)
        off += len(r.Val)
    }
    return buf
}

// decodeBatch parses a frame payload into a Batch.
func decodeBatch(p []byte) (Batch, error) {
    if len(p) < 12 {
        return Batch{}, ErrTornTail
    }
    var b Batch
    b.Seq = binary.LittleEndian.Uint64(p[0:8])
    count := binary.LittleEndian.Uint32(p[8:12])
    if count > uint32(len(p))/8+1 { // cheap sanity guard against absurd counts
        return Batch{}, ErrTornTail
    }
    off := 12
    for i := uint32(0); i < count; i++ {
        if off+4 > len(p) {
            return Batch{}, ErrTornTail
        }
        kl := binary.LittleEndian.Uint32(p[off : off+4])
        off += 4
        if off+int(kl) > len(p) {
            return Batch{}, ErrTornTail
        }
        k := string(p[off : off+int(kl)])
        off += int(kl)
        if off+4 > len(p) {
            return Batch{}, ErrTornTail
        }
        vl := binary.LittleEndian.Uint32(p[off : off+4])
        off += 4
        if off+int(vl) > len(p) {
            return Batch{}, ErrTornTail
        }
        v := string(p[off : off+int(vl)])
        off += int(vl)
        b.Records = append(b.Records, Record{Key: k, Val: v})
    }
    if off != len(p) {
        return Batch{}, ErrTornTail
    }
    return b, nil
}

// appendFrame writes a length-prefixed, CRC32C-protected frame.
func appendFrame(dst []byte, payload []byte) []byte {
    var hdr [frameHeaderLen]byte
    binary.LittleEndian.PutUint32(hdr[0:4], frameMagic)
    binary.LittleEndian.PutUint32(hdr[4:8], uint32(len(payload)))
    binary.LittleEndian.PutUint32(hdr[8:12], crc32.Checksum(payload, crcTable))
    dst = append(dst, hdr[:]...)
    dst = append(dst, payload...)
    return dst
}

// scanFrames returns every fully-valid batch encoded in data, in order, and the
// number of leading bytes consumed. The scan stops at the first frame that is
// incomplete or fails its CRC. The trailing region is the torn tail and is
// discarded on recovery.
func scanFrames(data []byte) (batches []Batch, ends []int, consumed int, torn bool) {
    off := 0
    for {
        if off == len(data) {
            return batches, ends, off, false
        }
        if off+frameHeaderLen > len(data) {
            return batches, ends, off, true
        }
        magic := binary.LittleEndian.Uint32(data[off : off+4])
        if magic != frameMagic {
            return batches, ends, off, true
        }
        l := binary.LittleEndian.Uint32(data[off+4 : off+8])
        want := binary.LittleEndian.Uint32(data[off+8 : off+12])
        start := off + frameHeaderLen
        end := start + int(l)
        if l > uint32(len(data)) || end > len(data) {
            return batches, ends, off, true
        }
        payload := data[start:end]
        if crc32.Checksum(payload, crcTable) != want {
            return batches, ends, off, true
        }
        b, err := decodeBatch(payload)
        if err != nil {
            return batches, ends, off, true
        }
        batches = append(batches, b)
        ends = append(ends, end)
        off = end
    }
}

manifest.go

package walstore

import (
    "encoding/binary"
    "errors"
    "hash/crc32"
    "os"
    "path/filepath"
)

// Manifest is the durable fence record. Epoch is bumped by every writer
// acquisition; any writer holding a lower epoch is stale and must stop.
type Manifest struct {
    Epoch   uint64
    BaseGen uint64
    BaseSeq uint64
}

const manifestMagic = uint32(0x4d414e31) // "MAN1"

func manifestPath(dir string) string { return filepath.Join(dir, "MANIFEST") }

func encodeManifest(m Manifest) []byte {
    var buf [32]byte
    binary.LittleEndian.PutUint32(buf[0:4], manifestMagic)
    binary.LittleEndian.PutUint64(buf[4:12], m.Epoch)
    binary.LittleEndian.PutUint64(buf[12:20], m.BaseGen)
    binary.LittleEndian.PutUint64(buf[20:28], m.BaseSeq)
    binary.LittleEndian.PutUint32(buf[28:32], crc32.Checksum(buf[0:28], crcTable))
    out := make([]byte, 32)
    copy(out, buf[:])
    return out
}

func decodeManifest(data []byte) (Manifest, error) {
    if len(data) != 32 {
        return Manifest{}, errors.New("walstore: bad manifest length")
    }
    if binary.LittleEndian.Uint32(data[0:4]) != manifestMagic {
        return Manifest{}, errors.New("walstore: bad manifest magic")
    }
    if crc32.Checksum(data[0:28], crcTable) != binary.LittleEndian.Uint32(data[28:32]) {
        return Manifest{}, errors.New("walstore: manifest checksum mismatch")
    }
    return Manifest{
        Epoch:   binary.LittleEndian.Uint64(data[4:12]),
        BaseGen: binary.LittleEndian.Uint64(data[12:20]),
        BaseSeq: binary.LittleEndian.Uint64(data[20:28]),
    }, nil
}

func readManifest(dir string) (Manifest, error) {
    raw, err := os.ReadFile(manifestPath(dir))
    if err != nil {
        return Manifest{}, err
    }
    return decodeManifest(raw)
}

// writeManifestFsync atomically replaces the manifest and fsyncs both the file
// and its parent directory so the epoch bump is durable before use.
func writeManifestFsync(dir string, m Manifest) error {
    tmp := manifestPath(dir) + ".tmp"
    f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
    if err != nil {
        return err
    }
    if _, err := f.Write(encodeManifest(m)); err != nil {
        f.Close()
        return err
    }
    if err := f.Sync(); err != nil {
        f.Close()
        return err
    }
    if err := f.Close(); err != nil {
        return err
    }
    if err := os.Rename(tmp, manifestPath(dir)); err != nil {
        return err
    }
    return fsyncDir(dir)
}

func fsyncDir(dir string) error {
    d, err := os.Open(dir)
    if err != nil {
        return err
    }
    defer d.Close()
    return d.Sync()
}

sst.go

package walstore

import (
    "encoding/binary"
    "errors"
    "fmt"
    "hash/crc32"
    "os"
    "path/filepath"
    "sort"
    "strconv"
    "strings"
)

const sstMagic = uint32(0x53535431) // "SST1"

// SSTable is an immutable, checksummed sorted snapshot at a generation.
type SSTable struct {
    Gen uint64
    Seq uint64 // highest sequence number included
    KV  map[string]string
}

// encodeSST serialises the table plus a footer protected by CRC32C.
func encodeSST(t SSTable) []byte {
    keys := make([]string, 0, len(t.KV))
    for k := range t.KV {
        keys = append(keys, k)
    }
    sort.Strings(keys)
    buf := make([]byte, 0, 16+len(keys)*16)
    var tmp [8]byte
    binary.LittleEndian.PutUint64(tmp[:], t.Seq)
    buf = append(buf, tmp[:]...)
    binary.LittleEndian.PutUint32(tmp[:4], uint32(len(keys)))
    buf = append(buf, tmp[:4]...)
    for _, k := range keys {
        v := t.KV[k]
        binary.LittleEndian.PutUint32(tmp[:4], uint32(len(k)))
        buf = append(buf, tmp[:4]...)
        buf = append(buf, k...)
        binary.LittleEndian.PutUint32(tmp[:4], uint32(len(v)))
        buf = append(buf, tmp[:4]...)
        buf = append(buf, v...)
    }
    // footer: crc(4) over all preceding bytes, then magic(4)
    var crcBuf [4]byte
    binary.LittleEndian.PutUint32(crcBuf[:], crc32.Checksum(buf, crcTable))
    buf = append(buf, crcBuf[:]...)
    binary.LittleEndian.PutUint32(crcBuf[:], sstMagic)
    buf = append(buf, crcBuf[:]...)
    return buf
}

// decodeSST validates and parses an SSTable. A truncated or corrupt table
// (e.g. an fsync that never completed) returns an error.
func decodeSST(data []byte) (SSTable, error) {
    if len(data) < 16 {
        return SSTable{}, errors.New("walstore: short SSTable")
    }
    body := data[:len(data)-8]
    footer := data[len(data)-8:]
    if binary.LittleEndian.Uint32(footer[4:8]) != sstMagic {
        return SSTable{}, errors.New("walstore: bad SSTable magic")
    }
    if crc32.Checksum(body, crcTable) != binary.LittleEndian.Uint32(footer[0:4]) {
        return SSTable{}, errors.New("walstore: SSTable checksum mismatch")
    }
    var t SSTable
    t.Seq = binary.LittleEndian.Uint64(body[0:8])
    count := binary.LittleEndian.Uint32(body[8:12])
    off := 12
    t.KV = make(map[string]string, count)
    for i := uint32(0); i < count; i++ {
        if off+4 > len(body) {
            return SSTable{}, errors.New("walstore: truncated SSTable body")
        }
        kl := binary.LittleEndian.Uint32(body[off : off+4])
        off += 4
        if off+int(kl) > len(body) {
            return SSTable{}, errors.New("walstore: truncated key")
        }
        k := string(body[off : off+int(kl)])
        off += int(kl)
        if off+4 > len(body) {
            return SSTable{}, errors.New("walstore: truncated SSTable body")
        }
        vl := binary.LittleEndian.Uint32(body[off : off+4])
        off += 4
        if off+int(vl) > len(body) {
            return SSTable{}, errors.New("walstore: truncated value")
        }
        v := string(body[off : off+int(vl)])
        off += int(vl)
        t.KV[k] = v
    }
    if off != len(body) {
        return SSTable{}, errors.New("walstore: trailing SSTable bytes")
    }
    return t, nil
}

// maxSSTGen returns the highest generation number present on disk, valid or
// not, so fresh writes never reuse a generation slot.
func maxSSTGen(dir string) uint64 {
    entries, err := os.ReadDir(filepath.Join(dir, "sst"))
    if err != nil {
        return 0
    }
    var max uint64
    for _, e := range entries {
        name := e.Name()
        if !strings.HasSuffix(name, ".sst") {
            continue
        }
        gen, err := strconv.ParseUint(strings.TrimSuffix(name, ".sst"), 10, 64)
        if err != nil {
            continue
        }
        if gen > max {
            max = gen
        }
    }
    return max
}

// sstPath returns the canonical path of a generation's SSTable.
func sstPath(dir string, gen uint64) string {
    return filepath.Join(dir, "sst", fmt.Sprintf("%020d.sst", gen))
}

// scanSSTGenerations validates every generation present on disk and returns the
// newest one whose contents are durable and checksum-clean. The manifest's
// pointer is only a hint; this is the authoritative durability check.
func scanSSTGenerations(dir string) (SSTable, bool, error) {
    entries, err := os.ReadDir(filepath.Join(dir, "sst"))
    if err != nil {
        if os.IsNotExist(err) {
            return SSTable{}, false, nil
        }
        return SSTable{}, false, err
    }
    var gens []uint64
    for _, e := range entries {
        name := e.Name()
        if !strings.HasSuffix(name, ".sst") {
            continue
        }
        gen, err := strconv.ParseUint(strings.TrimSuffix(name, ".sst"), 10, 64)
        if err != nil {
            continue
        }
        gens = append(gens, gen)
    }
    sort.Slice(gens, func(i, j int) bool { return gens[i] > gens[j] })
    for _, gen := range gens {
        raw, err := os.ReadFile(sstPath(dir, gen))
        if err != nil {
            continue
        }
        t, err := decodeSST(raw)
        if err != nil {
            continue // incomplete fsync / torn generation: skip
        }
        t.Gen = gen
        return t, true, nil
    }
    return SSTable{}, false, nil
}

store.go

package walstore

import (
    "errors"
    "fmt"
    "os"
    "path/filepath"
    "sync"
)

// StaleEpochError is returned when a writer tries to mutate the store while a
// newer writer has fenced it by advancing the manifest epoch.
type StaleEpochError struct {
    Held   uint64
    OnDisk uint64
}

func (e *StaleEpochError) Error() string {
    return fmt.Sprintf("walstore: stale manifest epoch: writer holds %d, on-disk is %d", e.Held, e.OnDisk)
}

// Store is the single-writer handle over one directory.
type Store struct {
    mu       sync.Mutex
    dir      string
    epoch    uint64
    seq      uint64
    nextGen  uint64
    mem      map[string]string
    wal      *os.File
    done     bool
}

// Recovered is the result of replaying a WAL image.
type Recovered struct {
    Mem      map[string]string
    Seq      uint64
    Applied  []Batch
    Consumed int
    Torn     bool
}

// recoverState replays WAL frames on top of a base snapshot. Only batches with
// contiguous, strictly-increasing sequence numbers are applied; the first
// incomplete, corrupt, or out-of-order frame terminates replay (torn tail).
func recoverState(raw []byte, base map[string]string, baseSeq uint64) (Recovered, error) {
    mem := make(map[string]string, len(base))
    for k, v := range base {
        mem[k] = v
    }
    batches, ends, consumed, torn := scanFrames(raw)
    res := Recovered{Mem: mem, Seq: baseSeq, Consumed: consumed, Torn: torn}
    prev := baseSeq
    applied := 0
    for _, b := range batches {
        if b.Seq != prev+1 {
            // Stop before the out-of-order frame; only the validated
            // prefix may be physically retained.
            res.Torn = true
            if applied == 0 {
                res.Consumed = 0
            } else {
                res.Consumed = ends[applied-1]
            }
            return res, nil
        }
        for _, r := range b.Records {
            mem[r.Key] = r.Val
        }
        prev = b.Seq
        res.Seq = b.Seq
        res.Applied = append(res.Applied, b)
        applied++
    }
    return res, nil
}

// Open recovers the store and acquires a fresh writer epoch. Any previously
// open writer becomes stale.
func Open(dir string) (*Store, error) {
    if err := os.MkdirAll(filepath.Join(dir, "sst"), 0o755); err != nil {
        return nil, err
    }
    if err := os.MkdirAll(filepath.Join(dir, "wal"), 0o755); err != nil {
        return nil, err
    }

    m, err := readManifest(dir)
    if err != nil {
        if !os.IsNotExist(err) {
            return nil, err
        }
        m = Manifest{}
    }

    sst, ok, err := scanSSTGenerations(dir)
    if err != nil {
        return nil, err
    }
    var base map[string]string
    var baseSeq uint64
    var baseGen uint64
    if ok {
        base = sst.KV
        baseSeq = sst.Seq
        baseGen = sst.Gen
    } else {
        base = map[string]string{}
    }

    walFile := filepath.Join(dir, "wal", "000001.wal")
    raw, err := os.ReadFile(walFile)
    if err != nil && !os.IsNotExist(err) {
        return nil, err
    }
    rec, err := recoverState(raw, base, baseSeq)
    if err != nil {
        return nil, err
    }

    // Physically drop the torn tail so subsequent appends are clean.
    f, err := os.OpenFile(walFile, os.O_CREATE|os.O_RDWR, 0o644)
    if err != nil {
        return nil, err
    }
    if err := f.Truncate(int64(rec.Consumed)); err != nil {
        f.Close()
        return nil, err
    }
    if _, err := f.Seek(0, os.SEEK_END); err != nil {
        f.Close()
        return nil, err
    }

    s := &Store{
        dir:     dir,
        epoch:   m.Epoch + 1,
        seq:     rec.Seq,
        nextGen: maxSSTGen(dir) + 1,
        mem:     rec.Mem,
        wal:     f,
    }
    // Persist the epoch fence before returning the writer.
    if err := writeManifestFsync(dir, Manifest{Epoch: s.epoch, BaseGen: baseGen, BaseSeq: rec.Seq}); err != nil {
        f.Close()
        return nil, err
    }
    return s, nil
}

// Epoch returns the writer epoch held by this handle.
func (s *Store) Epoch() uint64 {
    s.mu.Lock()
    defer s.mu.Unlock()
    return s.epoch
}

// Seq returns the highest applied sequence number.
func (s *Store) Seq() uint64 {
    s.mu.Lock()
    defer s.mu.Unlock()
    return s.seq
}

// Get returns the current value for a key.
func (s *Store) Get(key string) (string, bool) {
    s.mu.Lock()
    defer s.mu.Unlock()
    v, ok := s.mem[key]
    return v, ok
}

// checkEpochLocked fails if another writer advanced the manifest epoch.
func (s *Store) checkEpochLocked() error {
    m, err := readManifest(s.dir)
    if err != nil {
        return err
    }
    if m.Epoch != s.epoch {
        return &StaleEpochError{Held: s.epoch, OnDisk: m.Epoch}
    }
    return nil
}

// Append atomically applies one batch, fsyncing its single WAL frame first.
func (s *Store) Append(records []Record) error {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return errors.New("walstore: store is closed")
    }
    if err := s.checkEpochLocked(); err != nil {
        return err
    }
    batch := Batch{Seq: s.seq + 1, Records: records}
    frame := appendFrame(nil, encodeBatch(batch))
    if _, err := s.wal.Write(frame); err != nil {
        return err
    }
    if err := s.wal.Sync(); err != nil {
        return err
    }
    for _, r := range records {
        s.mem[r.Key] = r.Val
    }
    s.seq = batch.Seq
    return nil
}

// Flush writes the current in-memory state as a new checksummed SSTable
// generation, fsyncs it, then publishes it in the manifest. The manifest
// pointer is only a hint; recovery re-validates every generation.
func (s *Store) Flush() (uint64, error) {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return 0, errors.New("walstore: store is closed")
    }
    if err := s.checkEpochLocked(); err != nil {
        return 0, err
    }
    gen := s.nextGen
    table := SSTable{Gen: gen, Seq: s.seq, KV: s.mem}
    raw := encodeSST(table)

    tmp := sstPath(s.dir, gen) + ".tmp"
    f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644)
    if err != nil {
        return 0, err
    }
    if _, err := f.Write(raw); err != nil {
        f.Close()
        return 0, err
    }
    if err := f.Sync(); err != nil {
        f.Close()
        return 0, err
    }
    if err := f.Close(); err != nil {
        return 0, err
    }
    if err := os.Rename(tmp, sstPath(s.dir, gen)); err != nil {
        return 0, err
    }
    if err := fsyncDir(filepath.Join(s.dir, "sst")); err != nil {
        return 0, err
    }
    if err := writeManifestFsync(s.dir, Manifest{Epoch: s.epoch, BaseGen: gen, BaseSeq: s.seq}); err != nil {
        return 0, err
    }
    s.nextGen = gen + 1
    // The WAL prefix covered by this SSTable is no longer needed.
    if err := s.wal.Truncate(0); err != nil {
        return 0, err
    }
    if _, err := s.wal.Seek(0, os.SEEK_END); err != nil {
        return 0, err
    }
    if err := s.wal.Sync(); err != nil {
        return 0, err
    }
    return gen, nil
}

// Close releases the writer handle.
func (s *Store) Close() error {
    s.mu.Lock()
    defer s.mu.Unlock()
    if s.done {
        return nil
    }
    s.done = true
    return s.wal.Close()
}

store_test.go

package walstore

import (
    "errors"
    "fmt"
    "math/rand"
    "os"
    "path/filepath"
    "testing"
)

// frameLen returns the on-disk size of one encoded batch frame.
func frameLen(b Batch) int {
    return frameHeaderLen + len(encodeBatch(b))
}

// buildWAL encodes batches into a raw WAL image and returns per-frame end
// offsets.
func buildWAL(batches []Batch) ([]byte, []int) {
    var raw []byte
    ends := make([]int, len(batches))
    for i, b := range batches {
        raw = appendFrame(raw, encodeBatch(b))
        ends[i] = len(raw)
    }
    return raw, ends
}

func modelStates(batches []Batch) []map[string]string {
    states := make([]map[string]string, len(batches)+1)
    states[0] = map[string]string{}
    for i, b := range batches {
        m := make(map[string]string, len(states[i]))
        for k, v := range states[i] {
            m[k] = v
        }
        for _, r := range b.Records {
            m[r.Key] = r.Val
        }
        states[i+1] = m
    }
    return states
}

func randomBatches(rng *rand.Rand, n int) []Batch {
    batches := make([]Batch, n)
    for i := range batches {
        recs := make([]Record, 1+rng.Intn(4))
        for j := range recs {
            recs[j] = Record{
                Key: fmt.Sprintf("k%d_%d", i, rng.Intn(3)),
                Val: fmt.Sprintf("v%d_%d", i, rng.Intn(1000)),
            }
        }
        batches[i] = Batch{Seq: uint64(i + 1), Records: recs}
    }
    return batches
}

// TestTornTailPrefixProperty is the required property test: at least 200 random
// truncation offsets crossed with random batch boundaries. For every image the
// recovered state must equal some prefix of committed batches, never a partial
// batch and never a non-prefix.
func TestTornTailPrefixProperty(t *testing.T) {
    rng := rand.New(rand.NewSource(0xC0FFEE))
    const seeds = 8
    checks := 0

    for seed := 0; seed < seeds; seed++ {
        batches := randomBatches(rng, 1+rng.Intn(40))
        raw, ends := buildWAL(batches)
        states := modelStates(batches)

        // 200+ offsets per seed, plus every frame boundary.
        offsets := make([]int, 0, 260)
        for i := 0; i < 256; i++ {
            offsets = append(offsets, rng.Intn(len(raw)+1))
        }
        for _, e := range ends {
            offsets = append(offsets, e, e-1)
        }

        for _, off := range offsets {
            if off < 0 {
                continue
            }
            if off > len(raw) {
                continue
            }
            checks++
            rec, err := recoverState(raw[:off], map[string]string{}, 0)
            if err != nil {
                t.Fatalf("recover error: %v", err)
            }

            // Highest fully-contained frame index.
            wantApplied := 0
            for i, e := range ends {
                if e <= off {
                    wantApplied = i + 1
                }
            }
            if len(rec.Applied) != wantApplied {
                t.Fatalf("seed=%d off=%d: applied=%d want=%d", seed, off, len(rec.Applied), wantApplied)
            }
            // Sequence numbers must be strictly monotonic & contiguous.
            for i, b := range rec.Applied {
                if b.Seq != uint64(i+1) {
                    t.Fatalf("seed=%d off=%d: non-monotonic seq at %d: %d", seed, off, i, b.Seq)
                }
            }
            // State must equal the exact committed prefix.
            if !equalKV(rec.Mem, states[wantApplied]) {
                t.Fatalf("seed=%d off=%d: state mismatch\n got=%v\nwant=%v", seed, off, rec.Mem, states[wantApplied])
            }
            if rec.Seq != uint64(wantApplied) {
                t.Fatalf("seed=%d off=%d: seq=%d want=%d", seed, off, rec.Seq, wantApplied)
            }
        }
    }
    if checks < 200 {
        t.Fatalf("only %d checks, need >= 200", checks)
    }
    t.Logf("verified %d truncation offsets across %d random batch streams", checks, seeds)
}

// TestAtomicBatchDiscard proves a torn multi-key batch is fully discarded.
func TestAtomicBatchDiscard(t *testing.T) {
    recs := []Record{{Key: "a", Val: "1"}, {Key: "b", Val: "2"}, {Key: "c", Val: "3"}}
    b := Batch{Seq: 1, Records: recs}
    full := appendFrame(nil, encodeBatch(b))
    // Cut one byte before the CRC-valid end: the entire batch must vanish.
    rec, err := recoverState(full[:len(full)-1], map[string]string{}, 0)
    if err != nil {
        t.Fatal(err)
    }
    if len(rec.Applied) != 0 || len(rec.Mem) != 0 {
        t.Fatalf("partial batch applied: %+v", rec)
    }
    // With the last byte present, all keys appear together.
    rec, err = recoverState(full, map[string]string{}, 0)
    if err != nil {
        t.Fatal(err)
    }
    if len(rec.Mem) != 3 {
        t.Fatalf("atomic batch not fully applied: %+v", rec)
    }
}

// TestStaleEpochWriterRejected is the required fencing test.
func TestStaleEpochWriterRejected(t *testing.T) {
    dir := t.TempDir()
    a, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer a.Close()
    if err := a.Append([]Record{{Key: "x", Val: "1"}}); err != nil {
        t.Fatal(err)
    }

    // A second writer acquires a newer epoch.
    b, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer b.Close()
    if b.Epoch() <= a.Epoch() {
        t.Fatalf("epoch did not advance: a=%d b=%d", a.Epoch(), b.Epoch())
    }

    // The old writer must be fenced with a typed error.
    err = a.Append([]Record{{Key: "x", Val: "2"}})
    var se *StaleEpochError
    if !errors.As(err, &se) {
        t.Fatalf("want *StaleEpochError, got %T: %v", err, err)
    }
    // The new writer still works.
    if err := b.Append([]Record{{Key: "x", Val: "3"}}); err != nil {
        t.Fatal(err)
    }
}

// TestNewestDurableGenerationSelected proves recovery skips an SSTable whose
// fsync never completed (truncated/corrupt) even though the manifest points at
// it, and falls back to the newest checksum-clean generation.
func TestNewestDurableGenerationSelected(t *testing.T) {
    dir := t.TempDir()
    s, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    if err := s.Append([]Record{{Key: "k", Val: "durable"}}); err != nil {
        t.Fatal(err)
    }
    if _, err := s.Flush(); err != nil { // generation 1
        t.Fatal(err)
    }
    if err := s.Close(); err != nil {
        t.Fatal(err)
    }

    // Simulate a generation 2 whose write/fsync never completed: truncated.
    if err := os.WriteFile(sstPath(dir, 2), []byte("half-written"), 0o644); err != nil {
        t.Fatal(err)
    }
    // Manifest dishonestly points at the torn generation 2.
    if err := writeManifestFsync(dir, Manifest{Epoch: 1, BaseGen: 2, BaseSeq: 1}); err != nil {
        t.Fatal(err)
    }

    s2, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer s2.Close()
    v, ok := s2.Get("k")
    if !ok || v != "durable" {
        t.Fatalf("recovery picked torn generation: got %q ok=%v", v, ok)
    }
    // seq must not jump to the phantom generation.
    if s2.Seq() != 1 {
        t.Fatalf("seq=%d want 1", s2.Seq())
    }
    // New flush must not reuse generation 2 (it exists on disk).
    gen, err := s2.Flush()
    if err != nil {
        t.Fatal(err)
    }
    if gen <= 2 {
        t.Fatalf("reused generation %d", gen)
    }
}

// TestReopenRecoversAndTruncatesTornTail exercises the full Open path.
func TestReopenRecoversAndTruncatesTornTail(t *testing.T) {
    dir := t.TempDir()
    s, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    for i := 0; i < 5; i++ {
        if err := s.Append([]Record{{Key: "n", Val: fmt.Sprintf("%d", i)}}); err != nil {
            t.Fatal(err)
        }
    }
    if err := s.Close(); err != nil {
        t.Fatal(err)
    }

    // Append junk to the live WAL, then corrupt the last valid frame's CRC.
    walFile := filepath.Join(dir, "wal", "000001.wal")
    f, err := os.OpenFile(walFile, os.O_APPEND|os.O_WRONLY, 0o644)
    if err != nil {
        t.Fatal(err)
    }
    if _, err := f.Write([]byte("torn-garbage")); err != nil {
        t.Fatal(err)
    }
    f.Close()

    s2, err := Open(dir)
    if err != nil {
        t.Fatal(err)
    }
    defer s2.Close()
    if got, _ := s2.Get("n"); got != "4" {
        t.Fatalf("recovered n=%q want 4", got)
    }
    if s2.Seq() != 5 {
        t.Fatalf("seq=%d want 5", s2.Seq())
    }
    raw, err := os.ReadFile(walFile)
    if err != nil {
        t.Fatal(err)
    }
    // The torn garbage must be gone, and the 5 valid frames retained so they
    // can be replayed on the next open.
    if string(raw) == "torn-garbage" || len(raw) == 0 {
        t.Fatalf("torn tail not cleaned: %q", raw)
    }
    batches, _, _, torn := scanFrames(raw)
    if torn || len(batches) != 5 {
        t.Fatalf("expected 5 clean frames, got %d torn=%v", len(batches), torn)
    }
}

func equalKV(a, b map[string]string) bool {
    if len(a) != len(b) {
        return false
    }
    for k, v := range a {
        if b[k] != v {
            return false
        }
    }
    return true
}

5. Verification

go vet ./...
go test -race -count=3 ./...

Observed on go1.26.0 linux/amd64:

$ go test -v ./...
=== RUN   TestTornTailPrefixProperty
    store_test.go:124: verified 2370 truncation offsets across 8 random batch streams
--- PASS: TestTornTailPrefixProperty (0.02s)
=== RUN   TestAtomicBatchDiscard
--- PASS: TestAtomicBatchDiscard (0.00s)
=== RUN   TestStaleEpochWriterRejected
--- PASS: TestStaleEpochWriterRejected (0.00s)
=== RUN   TestNewestDurableGenerationSelected
--- PASS: TestNewestDurableGenerationSelected (0.00s)
=== RUN   TestReopenRecoversAndTruncatesTornTail
--- PASS: TestReopenRecoversAndTruncatesTornTail (0.00s)
PASS
ok      walstore        0.031s

Evidence map

Requirement Test
≥200 random truncation offsets × random batch boundaries; recovered state is a committed prefix TestTornTailPrefixProperty (2,370 offsets over 8 random streams)
Multi-key batch fully applied or fully discarded TestAtomicBatchDiscard
Stale-epoch writer rejected with a typed error TestStaleEpochWriterRejected (errors.As → *StaleEpochError)
Newest durable generation selected when manifest points at an un-fsynced SSTable TestNewestDurableGenerationSelected
End-to-end reopen drops the torn tail and replays the valid prefix TestReopenRecoversAndTruncatesTornTail

How the property test works

  1. For each random seed, generate a stream of batches with random sizes and random keys/values; encode them into one raw WAL image.
  2. Record each frame's end offset and the exact model state after each prefix.
  3. Cross 256 uniformly random truncation offsets with every frame boundary (±1 byte), yielding 2,370 evaluated images.
  4. recoverState must (a) apply exactly the frames fully contained before the cut, (b) never apply a partial batch, (c) keep seq contiguous and strictly increasing, and (d) produce a map equal to the corresponding model prefix. Any deviation is a hard t.Fatalf.

Why the fencing test is meaningful

Open fsyncs Manifest{Epoch: old+1} before returning. Writer A then holds epoch e. Opening writer B publishes e+1. A.Append re-reads the manifest, sees e+1, and returns *StaleEpochError before touching the WAL; B.Append succeeds. Fencing is therefore durable and checked on every mutation, not just at startup.


The artifact is saved at ~/walstore/ (wal.go, manifest.go, sst.go, store.go, store_test.go, go.mod, and SOLUTION.md). Run cd ~/walstore && go test -race ./... to reproduce.

Evidence & signatures

# Evidence
- Problem class: go-wal-torn-tail-recovery-manifest-fencing
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-27T16:32:27.219Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement crash recovery for a single-writer WAL plus SSTable store in Go: length-prefixed CRC32C frames, replay after the WAL is truncated at an arbitrary byte offset, multi-key batches that are either fully applied or fully discarded, and strictly monotonic sequence numbers. The manifest may point at an SSTable whose fsync never completed, so recovery must select the newest durable generation and reject a writer still holding a stale manifest epoch with a typed error. Evidence: a property test over at least 200 random truncation offsets crossed with random batch boundaries proving the recovered state equals some prefix of committed batches, plus a stale-epoch writer rejection test.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-wal-torn-tail-recovery-manifest-fencing", "provider": "openrouter", "solved_at": "2026-09-27T16:32:27.221Z", "version": "1.26"}
Generated from the verified corpus · MIT licensedBack to the catalog