wal

package
v0.16.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 9, 2026 License: MIT Imports: 24 Imported by: 0

Documentation

Overview

Package wal implements a versioned, length-prefixed, CRC32C-checksummed Write-Ahead Log for the gograph durability stack.

The on-disk format is documented in FORMAT.md alongside this package. A store's log is a directory of numbered segment files described by a durable control file; every frame carries its logical position, a link to its predecessor and the store's identity, so a reader refuses a frame from another store, a frame moved to another position, and a frame from an older generation. Readers stop cleanly at the first torn or corrupted frame.

Example

Example shows the core write-ahead-log loop: open a writer, append a few opaque payload frames, Sync them durably, then reopen the file with a Reader and replay every frame back in append order.

package main

import (
	"fmt"
	"os"
	"path/filepath"

	"github.com/FlavioCFOliveira/GoGraph/store/wal"
)

func main() {
	dir, err := os.MkdirTemp("", "wal-example")
	if err != nil {
		panic(err)
	}
	defer func() { _ = os.RemoveAll(dir) }()

	path := filepath.Join(dir, "wal")

	// Append three records. A WAL payload is opaque bytes as far as the
	// log is concerned; the durability stack above it (store/txn) gives
	// them meaning. Group-commit: several Appends, then a single Sync.
	w, err := wal.Open(path)
	if err != nil {
		panic(err)
	}
	for _, rec := range [][]byte{[]byte("alpha"), []byte("bravo"), []byte("charlie")} {
		if err := w.Append(rec); err != nil {
			panic(err)
		}
	}
	if err := w.Sync(); err != nil {
		panic(err)
	}
	if err := w.Close(); err != nil {
		panic(err)
	}

	// Replay: a fresh Reader iterates the frames in the order they were
	// appended, stopping cleanly at the first torn frame (none here).
	r, err := wal.OpenReader(path)
	if err != nil {
		panic(err)
	}
	defer func() { _ = r.Close() }()

	count := 0
	err = r.Replay(func(f wal.Frame) error {
		count++
		fmt.Printf("frame %d: %s\n", count, f.Payload)
		return nil
	})
	if err != nil {
		panic(err)
	}
	fmt.Printf("replayed %d frames\n", count)

}
Output:
frame 1: alpha
frame 2: bravo
frame 3: charlie
replayed 3 frames

Index

Examples

Constants

View Source
const (
	// ControlPrefixTruncated records that a checkpoint has discarded a prefix
	// of the log, so recovery requires the snapshot that folded it.
	ControlPrefixTruncated uint32 = 1 << 0
	// ControlLegacyV1Pending records that the legacy single-file log still
	// holds history and has not yet been replaced by its seal stub.
	ControlLegacyV1Pending uint32 = 1 << 1
)

Control flags.

View Source
const (
	// ControlRecordTag is the leading payload byte of a control record.
	ControlRecordTag byte = 0xFC
	// CtlReserveIDs is the control-record kind of an id reservation: every
	// intra index below the limit in the shard may have been issued (WAL v2
	// step 4). Body: u8 shard | u64 limit.
	CtlReserveIDs byte = 1
	// CtlNextIDsExact is the control-record kind written at a clean close: the
	// exact per-shard high-water marks (WAL v2 step 4). Body: 256 × uvarint next.
	CtlNextIDsExact byte = 2
	// CtlLegacySeal is the control-record kind that seals the legacy
	// single-file log: no frame of that file follows it.
	CtlLegacySeal byte = 3

	// ReserveIDsSize is the payload length of a [CtlReserveIDs] record.
	ReserveIDsSize = 2 + 1 + 8
)

Control records ride in a frame payload whose leading byte is ControlRecordTag (docs/design-wal-v2.md §1.2).

View Source
const (
	// DefaultSegmentSize is the segment size target when [Options.SegmentSize]
	// is 0: 16 MiB, PostgreSQL's default wal_segment_size.
	DefaultSegmentSize int64 = 16 << 20
	// MinSegmentSize is the smallest segment size target [Options] accepts.
	MinSegmentSize int64 = 1 << 20
)

Segment sizes (docs/design-wal-v2.md §2.3).

View Source
const CurrentVersion uint16 = 2

CurrentVersion is the WAL frame version this package writes: the 36-byte header of HeaderSizeV2, which carries the frame's logical position, the distance to its predecessor and the store id. Readers accept every version <= CurrentVersion, so a fresh build replays the single-file logs (LegacyVersion) of previous releases.

View Source
const HeaderSize = 4 + 2 + 4 + 4

HeaderSize is the fixed number of bytes occupying a LegacyVersion frame header (magic + version + length + crc32c).

View Source
const HeaderSizeV2 = 36

HeaderSizeV2 is the fixed number of bytes occupying a CurrentVersion frame header: magic (4), version (2), flags (2), length (4), position (8), prevLen (4), store id (8) and crc32c (4).

View Source
const LegacyVersion uint16 = 1

LegacyVersion is the frame version of the single-file log written by releases before WAL v2: the 14-byte header of HeaderSize, with no position and no store identity. Encode still writes it for a Frame whose Version is 0 or LegacyVersion, which is how test fixtures of the legacy format are built.

View Source
const NoFramePos = ^uint64(0)

NoFramePos marks "no frame": the prev-frame position of the store's first frame, as recorded in Control.PrevFramePosAtOR.

Variables

View Source
var (
	// ErrMissingControl indicates segments exist but the control file does
	// not, so neither the store id nor the oldest retained position is known.
	ErrMissingControl = errors.New("wal: segments exist but the control file is missing")
	// ErrControlCorrupt indicates the control file has a bad magic, version,
	// length or CRC.
	ErrControlCorrupt = errors.New("wal: control file is corrupt")
	// ErrSegmentHeader indicates a segment whose header has a bad magic,
	// version, length or CRC, or whose segment number differs from its name.
	ErrSegmentHeader = errors.New("wal: segment header is corrupt")
	// ErrForeignStore indicates a segment or a frame that carries a store id
	// other than the control file's: it belongs to another store.
	ErrForeignStore = errors.New("wal: segment or frame belongs to another store")
	// ErrSegmentGap indicates missing segments: segment numbers that are not
	// consecutive, an empty segment followed by one holding frames, or no
	// retained segment holding the oldest retained position.
	ErrSegmentGap = errors.New("wal: segments are missing")
	// ErrTornSegment indicates a torn frame anywhere but at the end of the
	// last segment holding frames. A non-tail segment is fully durable before
	// its successor receives a byte, so a tear there is corruption.
	ErrTornSegment = errors.New("wal: torn frame in a segment that is not the tail")
	// ErrFramePosition indicates a frame whose logical position is not the end
	// of its predecessor: a frame moved to another position or a hole.
	ErrFramePosition = errors.New("wal: frame position does not follow its predecessor")
	// ErrPrevLink indicates a frame whose prev-link does not name its
	// predecessor's position: a frame from an older generation of the log.
	ErrPrevLink = errors.New("wal: frame prev-link does not name its predecessor")
	// ErrLegacyNotSealed indicates a single-file legacy log that does not end
	// with its seal although segment frames exist, so the boundary between the
	// two histories is unknown.
	ErrLegacyNotSealed = errors.New("wal: legacy log is not sealed but segment frames exist")
)

Errors of the WAL v2 container. Every one of them is corruption: recovery refuses to open a directory that reports one.

View Source
var (
	// ErrBadMagic indicates the next four bytes did not match Magic.
	ErrBadMagic = errors.New("wal: bad frame magic")
	// ErrUnsupportedVersion indicates the frame version is newer
	// than this build knows how to parse.
	ErrUnsupportedVersion = errors.New("wal: unsupported frame version")
	// ErrCRCMismatch indicates the frame's CRC32C did not match the
	// re-computed value.
	ErrCRCMismatch = errors.New("wal: crc32c mismatch")
	// ErrTornFrame indicates the underlying reader returned EOF
	// before the frame was fully read.
	ErrTornFrame = errors.New("wal: torn frame at end of input")
	// ErrTornFrameMasksData indicates a frame's declared payload length
	// over-declared past the end of input AND the bytes it would have
	// consumed contain at least one further valid (CRC-checking) frame.
	// This is genuine mid-stream corruption masquerading as a benign torn
	// tail: a corrupt length field swallowed durable frames that follow it.
	// Unlike [ErrTornFrame] (a benign final partial write), this is a hard
	// error — it MUST fail-stop so the durable frames the bad length hid are
	// never silently dropped. It is a DISTINCT sentinel (it deliberately does
	// not wrap [ErrTornFrame]) so recovery's corruption classifier treats it
	// as corruption rather than a benign tail.
	ErrTornFrameMasksData = errors.New("wal: torn frame hides later valid frames (corrupt length)")
	// ErrFrameTooLarge indicates a frame payload longer than maxFrameSize.
	//
	// On DECODE the length is a declared one and is treated as corruption: the
	// frame is rejected before any allocation, so a crafted or corrupted length
	// cannot force a large one-shot make.
	//
	// On ENCODE it is a real payload, and the rejection is a durability guard
	// (rmp #2742): a frame this large could be written but never read back, so
	// [Encode] refuses it before any byte reaches the writer and the commit
	// fails rather than being acknowledged.
	ErrFrameTooLarge = errors.New("wal: frame payload length exceeds maximum")
)

Errors returned by the reader.

View Source
var ErrControlRecordCorrupt = errors.New("wal: control record corrupt")

ErrControlRecordCorrupt is returned by DecodeReserveIDs and DecodeNextIDsExact for a payload of their kind that does not parse.

View Source
var ErrDurabilityFailed = errors.New("wal: durability failed; the un-synced suffix was discarded and this writer is poisoned")

ErrDurabilityFailed marks every error a POISONED writer returns: the write-ahead log could not be made durable, the un-synced suffix has been discarded, and this writer will refuse every further append and sync.

It wraps the underlying I/O error, so `errors.Is(err, ErrDurabilityFailed)` identifies the class and `errors.Unwrap` still reaches the cause.

Why it exists, and what it is NOT

It is NOT retriable, and that is the whole point of naming it. A group commit is FAIL-ALL: when the leader's fsync fails, every member's frames and OpCommit markers are discarded together, so a transaction that did nothing wrong fails because another transaction's I/O failed. A caller needs to tell "MY transaction lost a conflict, retry it" from "the storage substrate failed, everything in flight is gone, and retrying will not help" (rmp #2306).

Fail-all is kept rather than softened, because the alternative is to acknowledge a commit whose durability is unknown, which the module's ACID mandate forbids outright. It is also the LENIENT end of the prior art: PostgreSQL does not fail the transaction, it fails the PROCESS — `issue_xlog_fsync` carries the comment "PANIC if failed to fsync" (postgres/postgres, master, read 2026-08-04 at commit 69ed7fd7e9da1cff2f04af04f630287971fe99fe; src/backend/access/transam/xlog.c). GoGraph is a library embedded in the caller's process, so the handle dies and says so, which is PostgreSQL's conclusion scoped to what a library owns.

View Source
var ErrSegmentsUnsupported = errors.New("wal: operation requires a segmented writer (use Open or OpenFS, not OpenWith)")

ErrSegmentsUnsupported is returned by Writer.MarkCheckpoint and Writer.ReclaimSegments on a Writer created by OpenWith: a single-file test writer has no control file and no segments to reclaim.

View Source
var ErrWALLocked = errors.New("wal: WAL directory is locked by another process")

ErrWALLocked is returned by Open when another process already holds the exclusive OS-level lock on the WAL directory. It signals that the WAL is in active use and the caller must not open a second writer against it — doing so would silently interleave frames and corrupt the log.

View Source
var ErrWriterClosed = errors.New("wal: writer is closed")

ErrWriterClosed is returned by methods on a Writer that has already been closed.

View Source
var Magic = [4]byte{'G', 'G', 'W', 'A'}

Magic is the 4-byte identifier prefix of every WAL frame: ASCII "GGWA".

Functions

func AppendNextIDsExact added in v0.16.0

func AppendNextIDsExact(dst []byte, next *[idShards]uint64) []byte

AppendNextIDsExact appends the payload of a CtlNextIDsExact record carrying next to dst and returns it.

func AppendReserveIDs added in v0.16.0

func AppendReserveIDs(dst []byte, shard uint8, limit uint64) []byte

AppendReserveIDs appends the payload of a CtlReserveIDs record for shard and limit to dst and returns it.

func ControlPath added in v0.16.0

func ControlPath(walPath string) string

ControlPath returns the path of the control file of the WAL at walPath. It is a pure function and safe for concurrent use.

func DecodeNextIDsExact added in v0.16.0

func DecodeNextIDsExact(payload []byte) (next [idShards]uint64, ok bool, err error)

DecodeNextIDsExact parses a CtlNextIDsExact payload. ok is false for a payload of another kind; err is ErrControlRecordCorrupt for one of this kind that does not hold exactly 256 marks.

func DecodeReserveIDs added in v0.16.0

func DecodeReserveIDs(payload []byte) (shard uint8, limit uint64, ok bool, err error)

DecodeReserveIDs parses a CtlReserveIDs payload. ok is false for a payload of another kind; err is ErrControlRecordCorrupt for one of this kind with the wrong length.

func Encode

func Encode(w io.Writer, f Frame) (int, error)

Encode writes f to w as a single binary frame. It returns the number of bytes written and any underlying writer error.

A Frame whose Version is 0 or LegacyVersion is written in the legacy 14-byte layout; a Frame whose Version is CurrentVersion is written in the 36-byte layout with its Pos, PrevLen and StoreID. Any other version is refused with ErrUnsupportedVersion. The Writer always writes CurrentVersion frames; the legacy layout exists so fixtures of logs written by earlier releases can still be built.

func EncodeLegacySeal added in v0.16.0

func EncodeLegacySeal(s LegacySeal) []byte

EncodeLegacySeal returns the payload of a CtlLegacySeal record.

func FrameSize added in v0.16.0

func FrameSize(f Frame) int

FrameSize returns the number of bytes f occupies on disk: its header (by version) plus its payload. A Version of 0 counts as LegacyVersion.

func PrefixTruncatedMarkerPath added in v0.16.0

func PrefixTruncatedMarkerPath(walPath string) string

PrefixTruncatedMarkerPath returns the path of the legacy prefix-truncation marker next to the WAL at walPath: the record a single-file log used to state that its history requires a snapshot (rmp #2990). Writers of the segmented format record that fact in the control file (ControlPrefixTruncated) and never write the marker; recovery still honours a marker it finds, and writes one for a legacy store loaded from a self-sufficient snapshot before that store is migrated (rmp #3002).

It is a pure function of walPath and is safe for concurrent use.

func SegmentDir added in v0.16.0

func SegmentDir(walPath string) string

SegmentDir returns the directory holding the segments of the WAL at walPath. It is a pure function and safe for concurrent use.

func SegmentPath added in v0.16.0

func SegmentPath(walPath string, segNo uint64) string

SegmentPath returns the path of segment segNo of the WAL at walPath. It is a pure function and safe for concurrent use.

func WritePrefixMarker added in v0.16.0

func WritePrefixMarker(walPath string) error

WritePrefixMarker makes the legacy prefix-truncation marker for the WAL at walPath durable on the operating-system filesystem (temp, fsync, rename, parent-directory fsync). Recovery uses it for a legacy store (rmp #3002).

Safe for concurrent use with respect to other directories; two concurrent calls for the same walPath share a temp file and must not overlap.

func WritePrefixMarkerFS added in v0.16.0

func WritePrefixMarkerFS(fsys walFS, walPath string) error

WritePrefixMarkerFS is WritePrefixMarker over a caller-supplied filesystem backend, whose ParentDirSync makes the rename durable. It is the seam the deterministic-simulation harness uses.

Types

type Control added in v0.16.0

type Control struct {
	// StoreID is the store's random 64-bit identity, created once.
	StoreID uint64
	// OldestRetainedPos (OR) is the first position of the oldest retained
	// segment, or the end of the log when the retained segments hold no frame.
	OldestRetainedPos uint64
	// PrevFramePosAtOR is the position of the frame preceding OR, or
	// [NoFramePos] when there is none.
	PrevFramePosAtOR uint64
	// CheckpointRedoPos is the redo position of the last truncating
	// checkpoint; diagnostic.
	CheckpointRedoPos uint64
	// CreatedUnixNano is when the store id was created; diagnostic.
	CreatedUnixNano uint64
	// Flags holds [ControlPrefixTruncated] and [ControlLegacyV1Pending].
	Flags uint32
}

Control is the decoded control file of a WAL (docs/design-wal-v2.md §2.2), the analogue of PostgreSQL's pg_control. It is a plain value, safe to copy and to read concurrently.

type Frame

type Frame struct {
	Payload []byte
	Pos     uint64
	StoreID uint64
	PrevLen uint32
	Version uint16
}

Frame is the in-memory representation of one WAL frame.

A Frame carries no synchronisation of its own, so its concurrency contract follows the ownership of Payload. A Frame returned by Decode owns its Payload outright — the decoder allocates a fresh slice per frame and never aliases the reader's buffer — so it is safe to hand to another goroutine and to read concurrently. A Frame passed to Encode merely borrows the caller's Payload for the duration of that call, which is what lets the transaction layer re-use one pooled scratch buffer for every op; such a Frame must not be retained or shared past the call that consumed it.

Pos, PrevLen and StoreID are meaningful only for a CurrentVersion frame: Pos is the frame's logical position (the count of frame bytes the store's log held before it, independent of segments and never reset), PrevLen is Pos minus the predecessor's Pos (0 only for the store's first frame), and StoreID is the store's identity from the control file. A LegacyVersion frame carries none of them and decodes with all three zero.

func Decode

func Decode(r io.Reader) (Frame, error)

Decode reads the next frame from r, of either version. It returns ErrTornFrame when the reader ends mid-frame (clean tail truncation), ErrBadMagic on a missing magic, ErrUnsupportedVersion on a version this build does not know, and ErrCRCMismatch on integrity failure. Any other error is propagated from the underlying reader.

type FrameSource added in v0.16.0

type FrameSource interface {
	Frames() iter.Seq[Frame]
	TailError() error
	TailOffset() int64
}

FrameSource is what a replay consumes: an iterator over frames, and after it ends the reason it stopped and the end of the last consumed frame. *Reader (one file) and *Log (a store's segmented log) both implement it.

Concurrency: implementations are NOT safe for concurrent use; one goroutine iterates a source and then reads its tail.

type LegacySeal added in v0.16.0

type LegacySeal struct {
	StoreID    uint64
	V2StartPos uint64
}

LegacySeal is the decoded CtlLegacySeal record: the store id the legacy log was sealed for and the segment position at which its history continues. It is a plain value, safe to copy and to read concurrently.

func DecodeLegacySeal added in v0.16.0

func DecodeLegacySeal(payload []byte) (LegacySeal, bool)

DecodeLegacySeal parses a CtlLegacySeal payload; ok is false for any other payload.

type Log added in v0.16.0

type Log struct {
	// contains filtered or unexported fields
}

Log reads a store's write-ahead log: the control file, the segments from the one holding the oldest retained position, and the legacy single-file log. It validates every frame's CRC, store id, position and prev-link.

A Log is not safe for concurrent use; open one per goroutine.

func OpenLog added in v0.16.0

func OpenLog(walPath string) (*Log, error)

OpenLog opens the write-ahead log at walPath on the operating-system filesystem; see OpenLogFS.

func OpenLogFS added in v0.16.0

func OpenLogFS(fsys LogFS, walPath string) (*Log, error)

OpenLogFS opens the write-ahead log at walPath through fsys. It reads the control file and enumerates and validates the segments; frames are read lazily by Log.Frames.

It returns ErrMissingControl when segments exist without a control file, ErrControlCorrupt, ErrSegmentHeader, ErrForeignStore and ErrSegmentGap for a damaged segment set. A directory with neither control file nor segments returns a Log whose Log.Control reports false: a store written only in the legacy single-file format.

func (*Log) Close added in v0.16.0

func (l *Log) Close() error

Close releases nothing today; it exists so a Log can be handled like a Reader. Segment files are opened and closed by Log.Frames itself.

func (*Log) Control added in v0.16.0

func (l *Log) Control() (Control, bool)

Control returns the control file, and false for a legacy-only store.

func (*Log) Frames added in v0.16.0

func (l *Log) Frames() iter.Seq[Frame]

Frames iterates the segment frames from the first frame of the segment that holds OR. Frames below OR in that segment are yielded too, so a caller can skip what a snapshot already covers. Iteration stops at the first error; Log.TailError then reports it: ErrTornFrame for a benign torn tail at the end of the last segment holding frames, and a corruption sentinel otherwise. A legacy-only Log yields nothing.

func (*Log) LastFramePos added in v0.16.0

func (l *Log) LastFramePos() int64

LastFramePos returns the position of the last valid frame, or -1 when none is known.

func (*Log) LegacyReader added in v0.16.0

func (l *Log) LegacyReader() (*Reader, error)

LegacyReader opens the legacy single-file log at walPath for reading, or returns nil when it does not exist. The caller closes it.

func (*Log) OldestRetainedPos added in v0.16.0

func (l *Log) OldestRetainedPos() int64

OldestRetainedPos returns OR, the first retained position (0 for a legacy-only store).

func (*Log) TailError added in v0.16.0

func (l *Log) TailError() error

TailError returns why Log.Frames stopped: nil at a clean end, ErrTornFrame for a benign torn tail, or a corruption sentinel.

func (*Log) TailOffset added in v0.16.0

func (l *Log) TailOffset() int64

TailOffset returns E, the logical position just past the last valid frame (OR when no retained frame exists). It is final once Log.Frames has run.

func (*Log) TailSegment added in v0.16.0

func (l *Log) TailSegment() (segNo uint64, offset int64)

TailSegment returns the segment number and the file offset just past the last valid frame; zero when no frame was read.

type LogFS added in v0.16.0

type LogFS interface {
	Open(path string) (io.ReadCloser, error)
	ReadDir(dir string) ([]string, error)
}

LogFS is the read-only filesystem surface OpenLogFS reads a WAL through: open a file by path and list a directory's entry names. Both return an error satisfying errors.Is(err, os.ErrNotExist) when the path does not exist.

Production code uses OpenLog, which reads the operating-system filesystem; the deterministic-simulation harness supplies its in-memory disk.

Concurrency: a Log calls its LogFS from one goroutine at a time; an implementation shared by several Logs must be safe for concurrent use.

type Options added in v0.16.0

type Options struct {
	// SegmentSize is the target size of a segment in bytes: a segment at or
	// above it is rolled over at the next run boundary, so one run larger
	// than the target still lands in one segment. 0 selects
	// [DefaultSegmentSize]; a value below [MinSegmentSize] is refused.
	SegmentSize int64
	// SyncLatency, when non-nil, delays every data fsync and directory fsync
	// of the writer; a testing instrument, see [SyncLatency].
	SyncLatency *SyncLatency
}

Options configures OpenWithOptions and OpenFSWithOptions. It is a plain configuration value: populate it before the call and do not mutate it while the call is in flight; read-only, it is safe for concurrent use.

type Reader

type Reader struct {
	// contains filtered or unexported fields
}

Reader iterates the frames of a WAL file. It is read-only and stops cleanly at the first torn or corrupted frame, reporting the byte offset where the cut occurred via Reader.TailOffset.

A Reader opened by OpenReader on the base path of a segmented log iterates that store's whole log instead: see OpenReader.

Reader is not safe for concurrent use; create one Reader per goroutine that wishes to iterate.

func NewReader

func NewReader(r io.Reader, closer io.Closer) *Reader

NewReader builds a Reader over an io.Reader. closer may be nil if the caller owns the resource.

func OpenReader

func OpenReader(path string) (*Reader, error)

OpenReader opens path for read-only frame iteration.

When path is the base path of a segmented log (a control file exists at ControlPath(path)), the Reader iterates the store's log: the frames of the legacy single-file log except its seal, then the segment frames from the oldest retained segment, validated as Log.Frames validates them; and Reader.TailOffset reports the end position of the log. Otherwise it reads path as one file of frames.

func (*Reader) Close

func (r *Reader) Close() error

Close releases any underlying resource passed to NewReader or OpenReader.

func (*Reader) Frames

func (r *Reader) Frames() iter.Seq[Frame]

Frames returns an iterator over every frame in the WAL. The iterator stops at the first error; call Reader.TailError / Reader.TailOffset after iteration to inspect why.

func (*Reader) Replay

func (r *Reader) Replay(apply func(Frame) error) error

Replay applies apply to every frame in the WAL in order. If apply returns an error, replay stops with that error returned to the caller. After Replay returns, TailOffset/TailError describe where and why iteration stopped (frame-level errors).

func (*Reader) TailError

func (r *Reader) TailError() error

TailError returns the error that ended iteration (typically ErrTornFrame, ErrCRCMismatch, or ErrBadMagic), or nil when iteration ended at clean EOF.

func (*Reader) TailOffset

func (r *Reader) TailOffset() int64

TailOffset returns the byte offset (from the start of the input) where iteration stopped. After a successful iteration to EOF this equals the file size; after a torn frame this equals the start of the torn frame.

type Stats

type Stats struct {
	Frames uint64 // total frames appended
	Bytes  uint64 // total frame bytes appended (header + payload)
	Syncs  uint64 // total successful data syncs (commit path and rollover)
	// SyncFailed counts sync rounds that failed at the flush/fsync
	// I/O layer. Calls rejected because the writer was already
	// poisoned by an earlier failure are not counted (mirroring how
	// context-cancelled calls are not counted).
	SyncFailed uint64
	// ControlFrames and ControlBytes count the subset of Frames and Bytes whose
	// payload is a control record ([ControlRecordTag]): node id reservations and
	// the clean-close id marks (WAL v2 step 4). Frames − ControlFrames is the
	// number of transaction frames.
	ControlFrames uint64
	ControlBytes  uint64
}

Stats is a snapshot of a Writer's lifetime counters. Counters are monotonic; subtract two snapshots to compute deltas. Values are read with sync/atomic.LoadUint64, so they may race slightly behind in-flight operations but never observe a torn value. The four counters are loaded one at a time, so a Stats is a per-field snapshot rather than a single atomic view across all four.

The value Writer.Stats returns is a detached copy of plain integers, so a Stats is safe for concurrent reads and Writer.Stats is safe to call concurrently with Writer.Append and Writer.Sync.

type SyncLatency added in v0.16.0

type SyncLatency struct {
	// contains filtered or unexported fields
}

SyncLatency adds a seeded, randomised delay before every WAL data fsync and every directory fsync a Writer performs. It is a TESTING instrument (rmp #3022), not a production setting.

Why it exists

A RAM drive shrinks an fsync to microseconds, and with it the window between a commit's timestamp allocation and its visibility publish, which is where commit-ordering races live. A test that needs the RAM drive's speed but must still open that window as a real device does installs a SyncLatency: each fsync then waits a duration drawn uniformly from [min, max] before it runs. Measured on rmp #2971's reproduction (the sessionless body of cypher.TestIndexBuild_CreateDropUnderCommittingWriters, 30 runs each): 0 failures on the RAM drive without it, 3 with 1-5 ms, 2 on an APFS SSD.

How it is activated

Explicitly, per Writer: OpenWithSyncLatency, or store.Options.SyncLatency for store.Open. There is no global switch and no environment variable in this package. A Writer opened any other way carries a nil SyncLatency and pays one nil-pointer test per fsync, nothing more.

Reproducibility

The delays are drawn from a PCG generator seeded with SyncLatency.Seed, so a seed reproduces the sequence of delays. It does not reproduce goroutine scheduling, so it makes a failing interleaving likely again rather than certain.

Concurrency: safe for concurrent use; one SyncLatency may be shared by several Writers, which then draw from one sequence.

func NewSyncLatency added in v0.16.0

func NewSyncLatency(seed uint64, lo, hi time.Duration) *SyncLatency

NewSyncLatency returns a SyncLatency that delays each fsync by a duration drawn uniformly from [lo, hi], seeded with seed. A negative bound is treated as 0, and hi below lo is raised to lo.

func (*SyncLatency) Bounds added in v0.16.0

func (l *SyncLatency) Bounds() (lo, hi time.Duration)

Bounds returns the delay interval [lo, hi].

func (*SyncLatency) Seed added in v0.16.0

func (l *SyncLatency) Seed() uint64

Seed returns the seed the delays are drawn from, for a test to report on failure.

type WALFile added in v0.6.0

type WALFile interface {
	io.Writer
	// Reader lets the writer scan the files it reopens (the legacy log it
	// seals, the tail segment whose torn tail it discards).
	io.Reader
	io.Seeker
	// Sync flushes OS write buffers to durable storage.
	Sync() error
	// Truncate resizes the file to size bytes.
	Truncate(size int64) error
	// Close releases underlying OS resources.
	Close() error
}

WALFile is the minimal open-handle interface that Writer requires of its underlying file. *os.File and *testfs.FaultFile both satisfy it, which lets fault-injection tests substitute a synthetic file without touching production paths.

It is exported so an external filesystem backend (the deterministic- simulation harness, internal/sim) can name it as the return type of its [walFS].OpenFile method and thereby satisfy the unexported walFS interface, exactly as github.com/FlavioCFOliveira/GoGraph/store/snapshot.File is exported for the snapshot seam. Production callers open WAL files via Open and never reference this type directly; tests and the simulator reach for OpenWith / OpenFS.

Concurrency: a WALFile is used by a single Writer whose own mutex serialises every access; any implementation's further guarantees are its own.

type Writer

type Writer struct {
	// contains filtered or unexported fields
}

Writer appends frames to a store's write-ahead log. Callers append frames with Writer.Append / Writer.AppendRun and durably commit them with Writer.Sync / Writer.SyncGroup; group-commit is achieved by appending several frames before a single sync.

A Writer created by Open or OpenFS writes a segmented log (docs/design-wal-v2.md): a control file, numbered segments under walPath+".d", and every frame stamped with its logical position, a link to its predecessor and the store id. Positions never reset; Writer.DurableOffset and the watermark Writer.AppendRun returns are positions. A segment at or above the target size is rolled over at the next run boundary, after it has been made durable, so no transaction spans two segments and every non-tail segment is fully durable before its successor receives a byte. A checkpoint reclaims space by unlinking whole segments (Writer.MarkCheckpoint, Writer.ReclaimSegments) without taking the append lock.

A Writer created by OpenWith writes the same frames into one file with no control file and no rollover; it is a test writer.

Concurrency

Writer is safe for concurrent use by any number of goroutines. Appends, syncs and rollover serialise on one internal mutex; a group-commit leader releases it across its fsync. Writer.MarkCheckpoint and Writer.ReclaimSegments never take that mutex: they touch only the control file and segments the writer no longer appends to, under their own locks. The writer of an [Open]ed log may run one short-lived background goroutine that prepares the next segment; Writer.Close waits for it.

Fail-stop

A Writer fail-stops on commit failure: the first flush or fsync error permanently poisons it. The un-synced suffix of the active segment — which may hold the flushed frames (including the commit marker) of the very transaction whose sync just failed — is physically discarded, and every subsequent Append/Sync returns the original error. Without the poison, a later transaction's successful fsync would make the failed transaction's frames durable even though its commit was never acknowledged: a phantom commit violating Atomicity and Durability. A poisoned Writer accepts only Writer.Close; the owner must discard it and re-open the WAL.

Which method answers "is this Writer healthy" — rmp #2525

Writer.Poisoned does, and it is the only member that does. The set splits this way:

After Writer.Close every method in the first two groups returns ErrWriterClosed instead of the sticky error, while Writer.Poisoned still reports the sticky error.

func Open

func Open(walPath string) (*Writer, error)

Open opens or creates the write-ahead log at walPath (dir/wal in a store directory) for appending; see OpenWithOptions.

func OpenFS added in v0.6.0

func OpenFS(fsys walFS, walPath string) (*Writer, error)

OpenFS opens or creates the segmented write-ahead log at walPath over a caller-supplied filesystem backend; see OpenFSWithOptions.

func OpenFSWithOptions added in v0.16.0

func OpenFSWithOptions(fsys walFS, walPath string, opts Options) (*Writer, error)

OpenFSWithOptions is OpenWithOptions over a caller-supplied filesystem backend. It is the seam the deterministic-simulation harness (internal/sim) uses to run the full snapshot + WAL + checkpoint stack against its in-memory disk; production code uses Open.

It differs from OpenWithOptions in two ways. It takes no OS lock (flock has no analogue on an injected backend, and callers are single-writer by contract), and it prepares the next segment synchronously at rollover rather than in a background goroutine, so the backend sees a deterministic sequence of operations. Directory fsyncs go through fsys.ParentDirSync, delayed by opts.SyncLatency when set.

func OpenWith

func OpenWith(f WALFile) (*Writer, error)

OpenWith builds a single-file test Writer over an already-open file handle. The caller transfers ownership: Writer.Close will call f.Close().

The Writer appends CurrentVersion frames with store id 0 after the file's existing bytes, never rolls over and has no control file; positions equal file offsets. It exists for tests that inject a *testfs.FaultFile; production code uses Open.

func OpenWithOptions added in v0.16.0

func OpenWithOptions(walPath string, opts Options) (*Writer, error)

OpenWithOptions opens or creates the write-ahead log at walPath for appending.

It takes an exclusive OS lock on walPath+".lock" for the Writer's lifetime (ErrWALLocked when another process holds it). Then:

  • A fresh directory gets a control file holding a new random store id, a first segment, and a one-frame seal stub at walPath, so a build that predates the segmented format refuses the directory instead of ignoring its segments.
  • A directory holding a legacy single-file log at walPath and no control file is migrated: a control file is written with a new store id and the legacy-pending flag, the first segment is created, and a seal frame is appended to the legacy file and fsynced. The legacy file's history is kept and replays as before; its benign torn tail is discarded first, and a legacy file whose scan stops at corruption is refused.
  • An existing segmented log is reopened: interrupted spare creations and segments wholly below the oldest retained position are deleted, the tail segment's benign torn tail is truncated and fsynced so new frames are never appended behind junk, and appending resumes at the end of the last valid frame. A tail that stops at genuine corruption is left byte-for-byte intact for recovery to report.

Every created file has mode 0o600. Open is OpenWithOptions with zero Options.

func OpenWithSyncLatency added in v0.16.0

func OpenWithSyncLatency(path string, lat *SyncLatency) (*Writer, error)

OpenWithSyncLatency is Open with lat installed on the returned Writer: every data fsync of the WAL and every directory fsync it performs, including those of the open itself, first waits a delay drawn from lat. A nil lat makes it exactly Open. It is a testing entry point; see SyncLatency.

func (*Writer) Append

func (w *Writer) Append(payload []byte) error

Append writes one frame with the given opaque payload. The frame is buffered in process memory; call Writer.Sync to durably commit.

On a writer poisoned by an earlier Sync failure, Append rejects the frame and returns the original sync error; see the Writer type documentation.

func (*Writer) AppendCtx

func (w *Writer) AppendCtx(ctx context.Context, payload []byte) error

AppendCtx is the context-aware variant of Writer.Append. ctx.Err() is checked before acquiring the internal mutex and again before writing; on cancellation returns the wrapped ctx.Err.

func (*Writer) AppendRun added in v0.11.0

func (w *Writer) AppendRun(fn func(emit func([]byte) error) error) (int64, error)

AppendRun appends every frame fn emits as ONE CONTIGUOUS RUN: no other appender's frame can land between them, and a run never spans two segments.

Why this exists — rmp #2302, audit finding E5

Crash recovery commits the ops carrying a marker's own TxnSeq and discards the buffered prefix as orphaned, and that reading is correct ONLY IF a transaction's frames are contiguous. Contiguity therefore lives in the component that owns the log.

The lock this holds

w.mu is taken once, before fn, and released after it — so fn runs with the writer exclusively held. Group commit is unaffected: Writer.SyncGroup coalesces on sync, not on append. Rollover happens at the start of the run, before fn, so the run lands whole in one segment.

Contract

The emit closure handed to fn is valid ONLY for the duration of the call; retaining it and calling it later panics. fn MUST NOT call any other method on this Writer — w.mu is not re-entrant and doing so deadlocks.

An error from fn is returned unchanged, and frames already emitted stay in the buffer: they are an un-marked, incomplete transaction, which recovery discards for atomicity. An error from append itself is the same fail-stop as Writer.AppendCtx's.

Prior art: PostgreSQL's XLogInsertRecord does the expensive work (assembling and CRCing the record) outside its insertion lock and holds it only for the copy (postgres/postgres, master, src/backend/access/transam/xlog.c).

Return value — the run's own durability watermark

AppendRun returns the log position immediately after the run's last frame. That position is the run's OWN watermark, and it is what the caller must hand to Writer.SyncGroup to make this run durable. The caller cannot derive it afterwards: the accepted position is shared mutable state that another appender advances and a poison REWINDS (rmp #2322).

func (*Writer) Close

func (w *Writer) Close() error

Close flushes any buffered frames, fsyncs the active segment, waits for the segment preparer, and releases every file and the lock.

On a writer poisoned by an earlier sync failure, Close skips the flush and performs a second-chance truncation of the un-synced suffix followed by a best-effort fsync, then returns the sticky error.

func (*Writer) DurableOffset added in v0.3.1

func (w *Writer) DurableOffset() int64

DurableOffset returns the log position covered by the last successful fsync: the end of every frame durably committed so far. It always lands on a frame boundary, because a transaction commits only after its marker has been appended and fsynced.

It is the redo position a non-blocking checkpoint captures under the store's quiesce boundary ([txn.Store.RunUnderCommitLock], which drains in-flight group commits), where it equals the accepted position: exactly the prefix the checkpoint's snapshot folds.

Concurrency: safe for concurrent use; it reads under the internal mutex.

func (*Writer) MarkCheckpoint added in v0.16.0

func (w *Writer) MarkCheckpoint(redoPos int64) error

MarkCheckpoint records in the control file that the log may from now on begin at the segment holding redoPos: it sets ControlPrefixTruncated, the oldest retained position OR to the first position of the segment containing redoPos, the prev-frame position at OR, and the checkpoint redo position. It is phase 2 of a truncating checkpoint, called after the snapshot covering redoPos is published and read back, and before Writer.ReclaimSegments unlinks anything: PostgreSQL's UpdateControlFile before RemoveOldXlogFiles.

redoPos must be a value Writer.DurableOffset returned. The control file is written temp, fsync, rename, directory fsync, so a crash leaves the old or the new one, and both are consistent with the segments on disk.

It never takes the append lock and never stalls a commit. Safe for concurrent use; control writes serialise among themselves.

func (*Writer) Poisoned added in v0.8.0

func (w *Writer) Poisoned() error

Poisoned reports the writer's fail-stop state: it returns the sticky commit-failure error when a prior flush or fsync has permanently poisoned the writer, and nil while the writer is healthy. It is the Writer's ONLY health probe (rmp #2525) and remains legible after Writer.Close.

It is the WAL-health probe a non-blocking checkpoint consults under the store's quiesce boundary BEFORE it captures and publishes a snapshot (rmp #1919): a writer observed poisoned there means a schema DDL whose commit failed may still be reflected in the engine's registry.

Concurrency: safe for concurrent use; it reads syncErr under the internal mutex.

func (*Writer) ReclaimSegments added in v0.16.0

func (w *Writer) ReclaimSegments() (int64, error)

ReclaimSegments unlinks every segment whose frames all lie below the oldest retained position the control file records (never the active segment), oldest first, and fsyncs the segment directory. When the control file still marks the legacy single-file log as pending, it then replaces that file with the one-frame seal stub and clears the flag. It returns the number of log bytes (positions) reclaimed.

It is phase 3 of a truncating checkpoint and needs NO commit lock: the writer never touches a non-active segment, and every frame it removes lies below the redo position the published snapshot covers. It never takes the append lock. Safe for concurrent use; it serialises with Writer.MarkCheckpoint.

func (*Writer) Stats

func (w *Writer) Stats() Stats

Stats returns a snapshot of the writer's lifetime counters.

func (*Writer) StoreID added in v0.16.0

func (w *Writer) StoreID() uint64

StoreID returns the store id stamped on every frame, or 0 for an OpenWith writer. Safe for concurrent use.

func (*Writer) Sync

func (w *Writer) Sync() error

Sync flushes the buffered frames to the OS and then issues the per-commit data sync (fdatasync(2) on Linux, os.File.Sync / fsync elsewhere; see dataSync) so the appended frames and the grown file size reach durable storage before returning.

The first flush or fsync failure permanently poisons the writer; see the Writer type documentation.

func (*Writer) SyncBuffered added in v0.11.0

func (w *Writer) SyncBuffered() error

SyncBuffered makes durable everything the writer has accepted at the moment of the call, coalescing with any concurrent group round exactly as Writer.SyncGroup does.

It is a FLUSH, not a commit acknowledgement: the accepted position is shared mutable state, so a committer must use Writer.AppendRun's watermark with Writer.SyncGroup to learn whether its own frames are durable (rmp #2322).

A nil return does NOT mean the Writer is healthy — rmp #2525

On a POISONED Writer SyncBuffered returns nil: the poison rewinds the accepted position to the durable one, so the already-durable fast path fires before the sticky error is tested. Writer.Poisoned is the health probe. After Writer.Close SyncBuffered returns ErrWriterClosed.

func (*Writer) SyncCtx

func (w *Writer) SyncCtx(ctx context.Context) error

SyncCtx is the context-aware variant of Writer.Sync. ctx.Err() is checked before acquiring the internal mutex; on cancellation returns the wrapped ctx.Err without flushing.

func (*Writer) SyncGroup added in v0.3.1

func (w *Writer) SyncGroup(target int64) error

SyncGroup durably commits the caller's already-appended frames, coalescing the fsync with those of every other committer whose frames are buffered at the same time — PostgreSQL-XLogFlush-style group commit. It returns nil only after a data sync has made durable every byte up to and including the caller's last appended frame (its OpCommit marker).

Contract

SyncGroup must be called AFTER the caller has appended all of its frames, with target set to the watermark Writer.AppendRun returned for that run. Then:

  • If a previous sync has already advanced the durable position to target, it returns nil without any I/O — the follower fast path.
  • Otherwise, if the writer is poisoned, it returns the sticky error.
  • Otherwise, if no leader is flushing, the caller becomes the LEADER: it flushes the buffer and fsyncs once, covering every buffered committer's frames, publishes the new durable position, and wakes the followers. If a leader is already flushing, the caller waits until the durable position covers its watermark or the writer poisons.

A run never spans segments and a rollover makes the old segment durable before switching, so the one fsync of the active segment covers every buffered frame.

Failure semantics

FAIL-ALL: if the leader's flush or fsync fails, [Writer.poison] discards the entire un-synced suffix (every group member's frames and markers) and broadcasts; every waiter then observes the sticky error and fails its own commit.

Cancellation

SyncGroup is intentionally NOT context-aware. Once a committer's frames are in the shared buffer they cannot be un-appended, so abandoning the wait while the group still fsyncs them would make the transaction durable while returning an error. This matches PostgreSQL: a backend cannot un-write WAL it has already inserted.

Concurrency: safe for concurrent calls; it serialises on the same internal mutex as Append/Sync and guarantees a single leader per fsync round.

func (*Writer) Truncate

func (w *Writer) Truncate() (int64, error)

Truncate discards every frame written so far. It is a test and maintenance helper, not a checkpoint path; the caller must already hold a snapshot that covers every frame.

On a segmented writer it rolls the active segment over (making it durable), writes the control file with ControlPrefixTruncated and OR at the current position, and unlinks every older segment. On an OpenWith writer it empties the file. It returns the number of log bytes discarded.

A poisoned or closed writer returns its sticky error or ErrWriterClosed and touches nothing.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL