◐ Off-By-One · answer catalog

go-bounded-rollback-streaming-parser

1 answer(s)gogo1.26

A naive streaming parser fails in five independent ways, and the design must make each failure distinguishable rather than collapsing them into one generic "bad frame":

📦 Source in repository (JSON)

Answer

All 14 tests pass under -race, and 4.5M fuzz executions found no panics. Here is the verified solution.


Bounded, Rollback-Capable Streaming Parser for a Nested Length-Prefixed Binary Protocol

Root-cause analysis

A naive streaming parser fails in five independent ways, and the design must make each failure distinguishable rather than collapsing them into one generic "bad frame":

Failure mode Why it is dangerous Required signal
Fragmented input Length header may be split across Write calls; a partial header must not be parsed as a full one ErrIncomplete, no consumption
Length-field overflow A uint64 length added to the header size can wrap to a small value, defeating bounds checks ErrIntegerOverflow
Checksum failure Corrupt payload must never be interpreted, and must not desync the stream ErrChecksumMismatch + resync
Decompression bomb A tiny compressed frame can expand without bound ErrDecompressionBomb
Corrupt compression Truncated/invalid DEFLATE is not a bomb ErrDecompressionCorrupt

The central correctness rule: checksum is verified over the on-wire payload before decompression, so we never inflate data that failed integrity. A second rule: every size is validated with overflow-safe arithmetic before any allocation or buffering. A third: on a malformed frame the parser skips exactly the declared frame extent (or one byte when the extent is unknown) and hunts for the next magic, so one bad frame cannot poison following valid frames.

Wire format

magic     2 bytes    0xB7 0x01
flags     1 byte     bit0 compressed, bit1 nested container
onwireLen 8 bytes    physical payload bytes after the header
logicalLen 8 bytes   payload bytes after decompression
crc32c    4 bytes    Castagnoli CRC of the on-wire payload
payload   onwireLen bytes

Header is fixed 23 bytes. 64-bit lengths make hostile inputs detectable as overflow.

Classification decision procedure

  1. onwireLen > MaxFramePayload → ErrFrameTooLarge.
  2. SafeAdd(headerSize, onwireLen) overflows uint64 or exceeds MaxInt → ErrIntegerOverflow.
  3. Fewer than total bytes buffered → ErrIncomplete (wait, unless retention exceeds MaxBufferBytes → ErrBufferBudget).
  4. crc32c(payload) != crc → ErrChecksumMismatch.
  5. If compressed: logicalLen > MaxDecompressed → ErrDecompressionBomb; else inflate exactly logicalLen:
  6. produced fewer bytes or DEFLATE error → ErrDecompressionCorrupt
  7. produced extra bytes → ErrDecompressionBomb
  8. If uncompressed: logicalLen != onwireLen → ErrMalformed.
  9. Nested flag with depth+1 > MaxDepth → ErrDepthExceeded; nested container fully present but truncated → ErrMalformed.

errors.Is(err, ErrChecksumMismatch) etc. distinguish them programmatically.


Exact fix

Create a module and drop in the files below.

mkdir broll && cd broll
go mod init broll

errors.go

package broll

import (
    "errors"
    "fmt"
)

type Kind int

const (
    KindMalformed Kind = iota
    KindIncomplete
    KindChecksum
    KindIntegerOverflow
    KindDecompressionBomb
    KindDecompressionCorrupt
    KindFrameTooLarge
    KindBufferBudget
    KindCPUBudget
    KindDepthExceeded
)

var (
    ErrMalformed            = errors.New("malformed frame")
    ErrIncomplete           = errors.New("incomplete frame")
    ErrChecksumMismatch     = errors.New("checksum mismatch")
    ErrIntegerOverflow      = errors.New("integer overflow")
    ErrDecompressionBomb    = errors.New("decompression bomb")
    ErrDecompressionCorrupt = errors.New("decompression corrupt")
    ErrFrameTooLarge        = errors.New("frame too large")
    ErrBufferBudget         = errors.New("buffer budget exceeded")
    ErrCPUBudget            = errors.New("cpu budget exceeded")
    ErrDepthExceeded        = errors.New("nesting depth exceeded")
)

func (k Kind) sentinel() error {
    switch k {
    case KindMalformed:
        return ErrMalformed
    case KindIncomplete:
        return ErrIncomplete
    case KindChecksum:
        return ErrChecksumMismatch
    case KindIntegerOverflow:
        return ErrIntegerOverflow
    case KindDecompressionBomb:
        return ErrDecompressionBomb
    case KindDecompressionCorrupt:
        return ErrDecompressionCorrupt
    case KindFrameTooLarge:
        return ErrFrameTooLarge
    case KindBufferBudget:
        return ErrBufferBudget
    case KindCPUBudget:
        return ErrCPUBudget
    case KindDepthExceeded:
        return ErrDepthExceeded
    }
    return ErrMalformed
}

func (k Kind) String() string {
    switch k {
    case KindMalformed:
        return "malformed"
    case KindIncomplete:
        return "incomplete"
    case KindChecksum:
        return "checksum"
    case KindIntegerOverflow:
        return "integer-overflow"
    case KindDecompressionBomb:
        return "decompression-bomb"
    case KindDecompressionCorrupt:
        return "decompression-corrupt"
    case KindFrameTooLarge:
        return "frame-too-large"
    case KindBufferBudget:
        return "buffer-budget"
    case KindCPUBudget:
        return "cpu-budget"
    case KindDepthExceeded:
        return "depth-exceeded"
    }
    return "unknown"
}

type ParseError struct {
    Kind   Kind
    Offset int64
    Detail string
}

func (e *ParseError) Error() string {
    if e.Detail == "" {
        return fmt.Sprintf("broll: %s frame at offset %d", e.Kind, e.Offset)
    }
    return fmt.Sprintf("broll: %s frame at offset %d: %s", e.Kind, e.Offset, e.Detail)
}

func (e *ParseError) Is(target error) bool { return target == e.Kind.sentinel() }
func (e *ParseError) Unwrap() error         { return e.Kind.sentinel() }

func perr(k Kind, off int64, format string, args ...any) *ParseError {
    return &ParseError{Kind: k, Offset: off, Detail: fmt.Sprintf(format, args...)}
}

protocol.go

package broll

import (
    "bytes"
    "compress/flate"
    "encoding/binary"
    "hash/crc32"
    "io"
    "math"
)

const (
    headerSize = 2 + 1 + 8 + 8 + 4

    magic0 = 0xB7
    magic1 = 0x01

    flagCompressed = 1 << 0
    flagNested     = 1 << 1
)

var castagnoli = crc32.MakeTable(crc32.Castagnoli)

type Config struct {
    MaxFramePayload uint64
    MaxDecompressed uint64
    MaxBufferBytes  int
    MaxDepth        int
    MaxWork         uint64
}

func DefaultConfig() Config {
    return Config{
        MaxFramePayload: 8 << 20,
        MaxDecompressed: 64 << 20,
        MaxBufferBytes:  16 << 20,
        MaxDepth:        32,
        MaxWork:         1 << 32,
    }
}

func (c Config) withDefaults() Config {
    d := DefaultConfig()
    if c.MaxFramePayload == 0 {
        c.MaxFramePayload = d.MaxFramePayload
    }
    if c.MaxDecompressed == 0 {
        c.MaxDecompressed = d.MaxDecompressed
    }
    if c.MaxBufferBytes == 0 {
        c.MaxBufferBytes = d.MaxBufferBytes
    }
    if c.MaxDepth == 0 {
        c.MaxDepth = d.MaxDepth
    }
    if c.MaxWork == 0 {
        c.MaxWork = d.MaxWork
    }
    return c
}

type Frame struct {
    Flags  uint8
    Data   []byte
    Nested []Frame
    Offset int64
}

func (f Frame) Compressed() bool { return f.Flags&flagCompressed != 0 }
func (f Frame) NestedFlag() bool { return f.Flags&flagNested != 0 }

func SafeAdd(a, b uint64) (uint64, bool) {
    if a > math.MaxUint64-b {
        return 0, false
    }
    return a + b, true
}

func inflate(src []byte, declared, limit uint64) ([]byte, error) {
    if declared > limit {
        return nil, perr(KindDecompressionBomb, 0, "declared %d > limit %d", declared, limit)
    }
    fr := flate.NewReader(bytes.NewReader(src))
    defer fr.Close()

    out := make([]byte, declared)
    if _, err := io.ReadFull(fr, out); err != nil {
        if err == io.EOF || err == io.ErrUnexpectedEOF {
            return nil, perr(KindDecompressionCorrupt, 0, "stream ended before %d bytes", declared)
        }
        return nil, perr(KindDecompressionCorrupt, 0, "%v", err)
    }
    var extra [1]byte
    if n, err := fr.Read(extra[:]); n > 0 {
        return nil, perr(KindDecompressionBomb, 0, "stream produced more than %d bytes", declared)
    } else if err != nil && err != io.EOF {
        return nil, perr(KindDecompressionCorrupt, 0, "%v", err)
    }
    return out, nil
}

func parseFrameAt(cfg Config, b []byte, depth int, absOff int64) (Frame, int, error) {
    if len(b) < headerSize {
        return Frame{}, 0, perr(KindIncomplete, absOff, "need %d header bytes, have %d", headerSize, len(b))
    }
    if b[0] != magic0 || b[1] != magic1 {
        return Frame{}, 0, perr(KindMalformed, absOff, "bad magic %02x%02x", b[0], b[1])
    }
    flags := b[2]
    onwire := binary.BigEndian.Uint64(b[3:11])
    logical := binary.BigEndian.Uint64(b[11:19])
    crc := binary.BigEndian.Uint32(b[19:23])

    if onwire > cfg.MaxFramePayload {
        return Frame{}, 0, perr(KindFrameTooLarge, absOff, "payload %d > max %d", onwire, cfg.MaxFramePayload)
    }
    total, ok := SafeAdd(headerSize, onwire)
    if !ok || total > math.MaxInt {
        return Frame{}, 0, perr(KindIntegerOverflow, absOff, "header+payload overflow (%d+%d)", headerSize, onwire)
    }
    if uint64(len(b)) < total {
        return Frame{}, 0, perr(KindIncomplete, absOff, "need %d bytes, have %d", total, len(b))
    }
    payload := b[headerSize:total]
    if got := crc32.Checksum(payload, castagnoli); got != crc {
        return Frame{}, int(total), perr(KindChecksum, absOff, "have %08x want %08x", got, crc)
    }

    var data []byte
    if flags&flagCompressed != 0 {
        if logical > cfg.MaxDecompressed {
            return Frame{}, int(total), perr(KindDecompressionBomb, absOff, "declared %d > max %d", logical, cfg.MaxDecompressed)
        }
        d, err := inflate(payload, logical, cfg.MaxDecompressed)
        if err != nil {
            if pe, ok := err.(*ParseError); ok {
                pe.Offset = absOff
                return Frame{}, int(total), pe
            }
            return Frame{}, int(total), err
        }
        data = d
    } else {
        if logical != onwire {
            return Frame{}, int(total), perr(KindMalformed, absOff, "uncompressed logical %d != onwire %d", logical, onwire)
        }
        data = append([]byte(nil), payload...)
    }

    f := Frame{Flags: flags, Data: data, Offset: absOff}
    if flags&flagNested != 0 {
        if depth+1 > cfg.MaxDepth {
            return Frame{}, int(total), perr(KindDepthExceeded, absOff, "depth %d > max %d", depth+1, cfg.MaxDepth)
        }
        children, err := parseNested(cfg, data, depth+1, absOff+headerSize)
        if err != nil {
            return Frame{}, int(total), err
        }
        f.Nested = children
    }
    return f, int(total), nil
}

func parseNested(cfg Config, data []byte, depth int, absOff int64) ([]Frame, error) {
    var out []Frame
    pos := 0
    for pos < len(data) {
        f, n, err := parseFrameAt(cfg, data[pos:], depth, absOff+int64(pos))
        if err != nil {
            if pe, ok := err.(*ParseError); ok && pe.Kind == KindIncomplete {
                return nil, perr(KindMalformed, pe.Offset, "truncated nested frame")
            }
            return nil, err
        }
        out = append(out, f)
        pos += n
    }
    return out, nil
}

encode.go

package broll

import (
    "bytes"
    "compress/flate"
    "encoding/binary"
    "hash/crc32"
)

func Encode(data []byte, compressed, nested bool) []byte {
    onwire := data
    flags := uint8(0)
    if compressed {
        flags |= flagCompressed
        var buf bytes.Buffer
        zw, _ := flate.NewWriter(&buf, flate.BestSpeed)
        _, _ = zw.Write(data)
        _ = zw.Close()
        onwire = buf.Bytes()
    }
    if nested {
        flags |= flagNested
    }
    out := make([]byte, headerSize+len(onwire))
    out[0], out[1] = magic0, magic1
    out[2] = flags
    binary.BigEndian.PutUint64(out[3:11], uint64(len(onwire)))
    binary.BigEndian.PutUint64(out[11:19], uint64(len(data)))
    binary.BigEndian.PutUint32(out[19:23], crc32.Checksum(onwire, castagnoli))
    copy(out[headerSize:], onwire)
    return out
}

func EncodeContainer(children ...[]byte) []byte {
    var payload []byte
    for _, c := range children {
        payload = append(payload, c...)
    }
    return Encode(payload, false, true)
}

func EncodeCompressedContainer(children ...[]byte) []byte {
    var payload []byte
    for _, c := range children {
        payload = append(payload, c...)
    }
    return Encode(payload, true, true)
}

parser.go

package broll

import (
    "errors"
    "time"
)

type Parser struct {
    cfg Config

    buf          []byte
    base         int64
    cursor       int
    commitCursor int

    committed []Frame
    pending   []Frame
    errs      []*ParseError

    work        uint64
    deadline    time.Time
    budgetFired bool
}

func New(cfg Config) *Parser {
    return &Parser{cfg: cfg.withDefaults()}
}

func (p *Parser) SetDeadline(t time.Time) { p.deadline = t }

func (p *Parser) Write(b []byte) (int, error) {
    p.buf = append(p.buf, b...)
    p.Process()
    return len(b), nil
}

func (p *Parser) Process() {
    for {
        if !p.deadline.IsZero() && time.Now().After(p.deadline) {
            p.record(KindCPUBudget, p.base+int64(p.cursor), "wall-clock deadline exceeded")
            return
        }
        if p.work > p.cfg.MaxWork {
            p.record(KindCPUBudget, p.base+int64(p.cursor), "work %d > max %d", p.work, p.cfg.MaxWork)
            p.cursor = len(p.buf)
            return
        }

        i := findMagic(p.buf, p.cursor)
        if i < 0 {
            p.work += uint64(len(p.buf) - p.cursor)
            keep := 0
            if len(p.buf) > p.cursor && p.buf[len(p.buf)-1] == magic0 {
                keep = 1
            }
            p.cursor = len(p.buf) - keep
            return
        }
        if i > p.cursor {
            p.record(KindMalformed, p.base+int64(p.cursor), "dropped %d bytes before magic", i-p.cursor)
            p.work += uint64(i - p.cursor)
            p.cursor = i
        }

        f, n, err := parseFrameAt(p.cfg, p.buf[p.cursor:], 0, p.base+int64(p.cursor))
        if n > 0 {
            p.work += uint64(n)
        }
        if err != nil {
            if errors.Is(err, ErrIncomplete) {
                if len(p.buf)-p.commitCursor > p.cfg.MaxBufferBytes {
                    p.record(KindBufferBudget, p.base+int64(p.cursor),
                        "retained %d > max %d", len(p.buf)-p.commitCursor, p.cfg.MaxBufferBytes)
                    p.cursor++
                    if p.cursor > len(p.buf) {
                        p.cursor = len(p.buf)
                    }
                    continue
                }
                return
            }
            p.recordErr(err)
            skip := n
            if skip <= 0 {
                skip = 1
            }
            p.cursor += skip
            if p.cursor > len(p.buf) {
                p.cursor = len(p.buf)
            }
            continue
        }

        p.pending = append(p.pending, f)
        p.cursor += n
    }
}

func (p *Parser) Commit() {
    if len(p.pending) == 0 {
        p.commitCursor = p.cursor
        p.trim()
        return
    }
    p.committed = append(p.committed, p.pending...)
    p.pending = nil
    p.commitCursor = p.cursor
    p.trim()
}

func (p *Parser) trim() {
    if p.commitCursor <= 0 {
        return
    }
    n := copy(p.buf, p.buf[p.commitCursor:])
    p.buf = p.buf[:n]
    p.base += int64(p.commitCursor)
    p.cursor -= p.commitCursor
    p.commitCursor = 0
}

func (p *Parser) Rollback() {
    p.pending = nil
    p.cursor = p.commitCursor
}

func (p *Parser) DiscardUncommitted() {
    p.pending = nil
    p.cursor = len(p.buf)
    p.trim()
    p.Process()
}

func (p *Parser) Frames() []Frame            { return p.committed }
func (p *Parser) Pending() []Frame           { return p.pending }
func (p *Parser) Errors() []*ParseError      { return p.errs }
func (p *Parser) LastCommittedOffset() int64 { return p.base + int64(p.commitCursor) }

func (p *Parser) recordErr(err error) {
    var pe *ParseError
    if errors.As(err, &pe) {
        p.errs = append(p.errs, pe)
        return
    }
    p.errs = append(p.errs, perr(KindMalformed, p.base+int64(p.cursor), "%v", err))
}

func (p *Parser) record(k Kind, off int64, format string, args ...any) {
    p.errs = append(p.errs, perr(k, off, format, args...))
}

func findMagic(b []byte, from int) int {
    for i := from; i+1 < len(b); i++ {
        if b[i] == magic0 && b[i+1] == magic1 {
            return i
        }
    }
    return -1
}

Note: budgetFired is retained as a hook; the CPU budget is enforced via work and deadline.

API summary


Verification

Test suite

parser_test.go:

package broll

import (
    "bytes"
    "encoding/binary"
    "errors"
    "hash/crc32"
    "math"
    "testing"
)

func rawFrame(flags uint8, onwire, logical uint64, crc uint32, payload []byte) []byte {
    b := make([]byte, headerSize+len(payload))
    b[0], b[1], b[2] = magic0, magic1, flags
    binary.BigEndian.PutUint64(b[3:11], onwire)
    binary.BigEndian.PutUint64(b[11:19], logical)
    binary.BigEndian.PutUint32(b[19:23], crc)
    copy(b[headerSize:], payload)
    return b
}

func mustCommit(t *testing.T, p *Parser) []Frame {
    t.Helper()
    p.Process()
    p.Commit()
    return p.Frames()
}

func TestFragmentedByteByByte(t *testing.T) {
    want := []byte("hello streaming world")
    wire := Encode(want, false, false)
    p := New(Config{})

    for i := 0; i < len(wire); i++ {
        if _, err := p.Write(wire[i : i+1]); err != nil {
            t.Fatalf("write: %v", err)
        }
        if i < len(wire)-1 && len(p.Pending()) != 0 {
            t.Fatalf("frame committed early at byte %d", i)
        }
    }
    frames := mustCommit(t, p)
    if len(frames) != 1 {
        t.Fatalf("got %d frames, want 1; errs=%v", len(frames), p.Errors())
    }
    if !bytes.Equal(frames[0].Data, want) {
        t.Fatalf("payload mismatch: %q", frames[0].Data)
    }
    if len(p.Errors()) != 0 {
        t.Fatalf("unexpected errors: %v", p.Errors())
    }
}

func TestCompressedRoundTrip(t *testing.T) {
    payload := bytes.Repeat([]byte("abc123"), 5000)
    wire := Encode(payload, true, false)
    p := New(Config{})
    p.Write(wire)
    frames := mustCommit(t, p)
    if len(frames) != 1 {
        t.Fatalf("got %d frames; errs=%v", len(frames), p.Errors())
    }
    if !bytes.Equal(frames[0].Data, payload) {
        t.Fatal("decompressed payload mismatch")
    }
    if !frames[0].Compressed() {
        t.Fatal("compressed flag lost")
    }
}

func TestNestedContainer(t *testing.T) {
    inner1 := Encode([]byte("child-1"), false, false)
    inner2 := Encode([]byte("child-2"), true, false)
    wire := EncodeContainer(inner1, inner2)
    p := New(Config{})
    p.Write(wire)
    frames := mustCommit(t, p)
    if len(frames) != 1 || !frames[0].NestedFlag() {
        t.Fatalf("expected one nested frame; errs=%v", p.Errors())
    }
    if len(frames[0].Nested) != 2 {
        t.Fatalf("got %d children, want 2", len(frames[0].Nested))
    }
    if string(frames[0].Nested[0].Data) != "child-1" || string(frames[0].Nested[1].Data) != "child-2" {
        t.Fatalf("children wrong: %q %q", frames[0].Nested[0].Data, frames[0].Nested[1].Data)
    }
}

func TestChecksumFailureDoesNotCorruptFollowing(t *testing.T) {
    bad := Encode([]byte("corrupt me"), false, false)
    bad[headerSize+3] ^= 0xFF
    good := Encode([]byte("valid after corruption"), false, false)

    p := New(Config{})
    p.Write(append(append([]byte{}, bad...), good...))
    frames := mustCommit(t, p)

    if len(frames) != 1 {
        t.Fatalf("got %d committed frames, want 1; errs=%v", len(frames), p.Errors())
    }
    if string(frames[0].Data) != "valid after corruption" {
        t.Fatalf("wrong survivor: %q", frames[0].Data)
    }
    if len(p.Errors()) != 1 || !errors.Is(p.Errors()[0], ErrChecksumMismatch) {
        t.Fatalf("want one checksum error, got %v", p.Errors())
    }
}

func TestIntegerOverflow(t *testing.T) {
    wire := make([]byte, headerSize)
    wire[0], wire[1], wire[2] = magic0, magic1, 0
    binary.BigEndian.PutUint64(wire[3:11], math.MaxUint64)
    p := New(Config{MaxFramePayload: math.MaxUint64})
    p.Write(wire)
    if len(p.Errors()) != 1 || !errors.Is(p.Errors()[0], ErrIntegerOverflow) {
        t.Fatalf("want integer overflow, got %v", p.Errors())
    }
}

func TestDecompressionBombDeclared(t *testing.T) {
    wire := Encode(bytes.Repeat([]byte("x"), 1000), true, false)
    p := New(Config{MaxDecompressed: 16})
    p.Write(wire)
    if len(p.Errors()) != 1 || !errors.Is(p.Errors()[0], ErrDecompressionBomb) {
        t.Fatalf("want bomb, got %v", p.Errors())
    }
    if len(p.Frames()) != 0 {
        t.Fatal("bomb frame was committed")
    }
}

func TestDecompressionBombActual(t *testing.T) {
    wire := Encode(bytes.Repeat([]byte("y"), 5000), true, false)
    binary.BigEndian.PutUint64(wire[11:19], 10)
    p := New(Config{})
    p.Write(wire)
    if len(p.Errors()) != 1 || !errors.Is(p.Errors()[0], ErrDecompressionBomb) {
        t.Fatalf("want actual bomb, got %v", p.Errors())
    }
}

func TestDecompressionCorruptNotBomb(t *testing.T) {
    wire := Encode(bytes.Repeat([]byte("z"), 1000), true, false)
    binary.BigEndian.PutUint64(wire[11:19], 5000)
    p := New(Config{})
    p.Write(wire)
    if len(p.Errors()) != 1 || !errors.Is(p.Errors()[0], ErrDecompressionCorrupt) {
        t.Fatalf("want corrupt, got %v", p.Errors())
    }
    if errors.Is(p.Errors()[0], ErrDecompressionBomb) {
        t.Fatal("corruption must not be reported as a bomb")
    }
}

func TestRollbackToLastCommitted(t *testing.T) {
    p := New(Config{})
    p.Write(Encode([]byte("one"), false, false))
    p.Write(Encode([]byte("two"), false, false))
    p.Process()
    p.Commit()
    if len(p.Frames()) != 2 {
        t.Fatalf("want 2 committed, got %d", len(p.Frames()))
    }
    committedOffset := p.LastCommittedOffset()

    p.Write(Encode([]byte("three"), false, false))
    if len(p.Pending()) != 1 {
        t.Fatalf("want 1 pending, got %d", len(p.Pending()))
    }

    p.Rollback()
    if len(p.Pending()) != 0 {
        t.Fatalf("rollback left %d pending", len(p.Pending()))
    }
    if len(p.Frames()) != 2 {
        t.Fatalf("rollback changed committed count: %d", len(p.Frames()))
    }
    if p.LastCommittedOffset() != committedOffset {
        t.Fatalf("offset moved: %d != %d", p.LastCommittedOffset(), committedOffset)
    }

    p.Process()
    if len(p.Pending()) != 1 || string(p.Pending()[0].Data) != "three" {
        t.Fatalf("retry did not recover frame: %v", p.Pending())
    }
    p.Commit()
    if len(p.Frames()) != 3 {
        t.Fatalf("want 3 committed after retry, got %d", len(p.Frames()))
    }
}

func TestBufferBudget(t *testing.T) {
    wire := Encode(bytes.Repeat([]byte("b"), 4096), false, false)
    p := New(Config{MaxFramePayload: 1 << 20, MaxBufferBytes: 128})
    p.Write(wire[:200])
    if len(p.Errors()) == 0 || !errors.Is(p.Errors()[0], ErrBufferBudget) {
        t.Fatalf("want buffer budget error, got %v", p.Errors())
    }
}

func TestDepthExceeded(t *testing.T) {
    inner := Encode([]byte("leaf"), false, false)
    wire := inner
    for i := 0; i < 5; i++ {
        wire = EncodeContainer(wire)
    }
    p := New(Config{MaxDepth: 3})
    p.Write(wire)
    found := false
    for _, e := range p.Errors() {
        if errors.Is(e, ErrDepthExceeded) {
            found = true
        }
    }
    if !found {
        t.Fatalf("want depth error, got %v", p.Errors())
    }
}

func TestResyncAfterGarbage(t *testing.T) {
    good := Encode([]byte("survivor"), false, false)
    junk := []byte{0x00, 0x11, 0x22, 0xB7, 0x99, 0x33}
    p := New(Config{})
    p.Write(append(append(append([]byte{}, junk...), good...), good...))
    frames := mustCommit(t, p)
    if len(frames) != 2 {
        t.Fatalf("want 2 survivors, got %d; errs=%v", len(frames), p.Errors())
    }
}

func TestNestedCorruptionIsolated(t *testing.T) {
    inner := Encode([]byte("nested"), false, false)
    container := EncodeContainer(inner)
    container[headerSize+headerSize+3] ^= 0xFF
    after := Encode([]byte("after"), false, false)

    p := New(Config{})
    p.Write(append(container, after...))
    frames := mustCommit(t, p)
    if len(frames) != 1 || string(frames[0].Data) != "after" {
        t.Fatalf("nested corruption leaked; frames=%v errs=%v", frames, p.Errors())
    }
    if !errors.Is(p.Errors()[0], ErrChecksumMismatch) {
        t.Fatalf("want nested checksum error, got %v", p.Errors())
    }
}

func TestCRCMatchesManual(t *testing.T) {
    payload := []byte("crc-check")
    wire := Encode(payload, false, false)
    want := crc32.Checksum(payload, crc32.MakeTable(crc32.Castagnoli))
    if got := binary.BigEndian.Uint32(wire[19:23]); got != want {
        t.Fatalf("crc mismatch: %08x != %08x", got, want)
    }
}

Fuzz test (fuzz_test.go)

package broll

import "testing"

func FuzzParser(f *testing.F) {
    f.Add(Encode([]byte("seed"), false, false))
    f.Add(Encode([]byte("compressed seed"), true, false))
    f.Add(EncodeContainer(Encode([]byte("child"), false, false)))
    f.Add([]byte{0xB7, 0x01})
    f.Add([]byte{0xB7})
    f.Add([]byte{0x00, 0x11, 0x22})

    overflow := make([]byte, headerSize)
    overflow[0], overflow[1] = magic0, magic1
    for i := 3; i < 11; i++ {
        overflow[i] = 0xFF
    }
    f.Add(overflow)

    f.Fuzz(func(t *testing.T, data []byte) {
        cfg := Config{
            MaxFramePayload: 1 << 16,
            MaxDecompressed: 1 << 20,
            MaxBufferBytes:  1 << 16,
            MaxDepth:        8,
            MaxWork:         1 << 24,
        }
        p := New(cfg)
        mid := len(data) / 2
        _, _ = p.Write(data[:mid])
        _, _ = p.Write(data[mid:])
        p.Process()
        p.Commit()

        for _, e := range p.Errors() {
            _ = e.Error()
        }

        var walk func(fs []Frame)
        walk = func(fs []Frame) {
            for _, fr := range fs {
                if uint64(len(fr.Data)) > cfg.MaxDecompressed {
                    t.Fatalf("payload %d exceeds budget", len(fr.Data))
                }
                walk(fr.Nested)
            }
        }
        walk(p.Frames())
    })
}

Commands and observed results

$ gofmt -l .
$ go vet ./...
$ go test ./... -count=1 -race
ok      broll   1.025s

$ go test ./... -count=1 -cover
ok      broll   0.005s  coverage: 78.6% of statements

$ go test -run=^$ -fuzz=FuzzParser -fuzztime=15s ./...
fuzz: elapsed: 15s, execs: 4494275 (269750/sec), new interesting: 65 (total: 72)
PASS
ok      broll   15.118s

All 14 unit tests pass under the race detector; 4.49 million fuzz executions produced no panic, no over-budget committed payload, and no non-termination.

What the tests prove

Evidence & signatures

# Evidence
- Problem class: go-bounded-rollback-streaming-parser
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-10-05T16:20:36.163Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Build a streaming parser for a nested binary protocol with length-prefixed frames, checksums, and optional compression, where input can arrive in arbitrary fragments and malformed frames must not corrupt subsequent valid frames. Enforce strict memory and CPU budgets, support rollback to the last committed frame, and specify how checksum failures, integer overflow, and decompression bombs are distinguished.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-bounded-rollback-streaming-parser", "provider": "openrouter", "solved_at": "2026-10-05T16:20:36.164Z", "version": "1.26"}
Generated from the verified corpus · MIT licensedBack to the catalog