Problem class: go-wal-torn-tail-recovery-manifest-fencing
The complete, verified solution is written to ~/walstore/SOLUTION.md (32 KB). Here it is in full.
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.
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:
fsync of the SSTable body never completed. Recovery trusts the pointer, reads a truncated file, and either errors forever or loads corrupt state.Every symptom is a durability-boundary bug: the code conflates "bytes were handed to the kernel" with "bytes are durable and complete".
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.
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).
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.
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.
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.
<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
magic u32 + len u32 + crc32c(payload) u32 + payload.seq u64 + count u32 + N × (klen u32,key,vlen u32,val).crc32c(body) u32 + magic u32.magic u32 + epoch u64 + baseGen u64 + baseSeq u64, followed by crc32c(body) u32.go.modmodule walstore
go 1.26
wal.gopackage 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.gopackage 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.gopackage 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.gopackage 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.gopackage 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
}
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
| 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 |
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.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 - 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"}The complete, verified solution is written to ~/walstore/SOLUTION.md (32 KB). Here it is in full.
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.
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:
fsync of the SSTable body never completed. Recovery trusts the pointer, reads a truncated file, and either errors forever or loads corrupt state.Every symptom is a durability-boundary bug: the code conflates "bytes were handed to the kernel" with "bytes are durable and complete".
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.
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).
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.
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.
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.
<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
magic u32 + len u32 + crc32c(payload) u32 + payload.seq u64 + count u32 + N × (klen u32,key,vlen u32,val).crc32c(body) u32 + magic u32.magic u32 + epoch u64 + baseGen u64 + baseSeq u64, followed by crc32c(body) u32.go.modmodule walstore
go 1.26
wal.gopackage 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.gopackage 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.gopackage 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.gopackage 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.gopackage 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
}
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
| 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 |
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.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 - 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"}