wal

package module
v1.0.4 Latest Latest
Warning

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

Go to latest
Published: Sep 24, 2026 License: MIT Imports: 22 Imported by: 0

README

wal

A write-ahead log implementation for an event collector application.

Install

go get codeberg.org/kchan/wal

Usage

To create a WAL (wal.Log)

Use wal.NewLog to create a new write-ahead log:

import (
    "log/slog"
    "os"
    
    "codeberg.org/kchan/wal"
)

...

logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))

config := wal.Config{
    Dir: "/path/to/wal-data",      // directory in which to store WAL files
    Segment: wal.SegmentConfig{
        MaxStoreBytes: 16 * 1024,  // maximum individual store file size 
        MaxIndexBytes: 8 * 1024,   // maximum individual index file size
    },
}
encoder := wal.NewEncoder()  // use default encoder
w, err := wal.NewLog(config, encoder, logger)
if err != nil {
    return err
}
defer w.Close()
...
To write data to the WAL
// some data to write to the log -- this should be encoded as []byte
data := []byte("some data")

// wal.Log.Write will return a LSN (log sequence number) and an error
lsn, err := w.Write(data)
if err != nil {
    return err
}

To read entries from the WAL

Create an Iterator and call Next to process each entry from the log segments:

import (
    "context"
    "time"
    
    "codeberg.org/kchan/wal"
)

...
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

// each entry has a "LSN" (log sequence number)
// set starting index by calculating a "start" value
start := w.LowestIndex()       // the lowest index in the WAL
checkpoint := w.Checkpoint.C() // the last stored checkpoint
start = max(start, checkpoint) // start the sync process using the latest index

// create iterator
iter, err := w.Iterator(start)
if err != nil {
    return err
}

// iterate through each entry starting with "start"
// until the iterator returns either wal.ContextCancelled,
// wal.StopIteration, or an error
for {
    // lsn is the log sequence number
    // data is the []byte data stored in the log
    lsn, data, err := iter.Next(ctx) // supply context as parameter
    if errors.Is(err, wal.ContextCancelled) {
        return err
    }
    if errors.Is(err, wal.StopIteration) {
        return err
    }
    if err != nil {
        return err
    }

    // perform sync
    // ...

    // then update checkpoint
    err = w.Checkpoint.Write(lsn)
    if err != nil {
        return err
    }
}
...
To purge WAL segments

Run a goroutine that calls the Purge method:

import (
    "context"
    "time"
    
    "codeberg.org/kchan/wal"
)

...
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

purgeInterval := 15 // purge the WAL segments every 15 seconds
ticker := time.NewTicker(time.Duration(purgeInterval) * time.Second)

for {
    select {
    case <-ctx.Done():
        return nil
    case <-ticker.C:
        // all segments prior to the latest checkpoint will be purged (deleted)
        err := w.Purge(ctx, func(segNum uint32, segment *wal.Segment) error {
            return nil
        })
        if err != nil {
            return err
        }
    }
}
...

License

MIT License. See LICENSE.

Documentation

Overview

The LSN (log sequence number) is an uint64 and structured as follows:

------------------------------------------------------ high 32-bits | low 32-bits segment number (segNum) | 0-based index of entry (idx) ------------------------------------------------------

Example:

Entry with segment number (segNum) = 1 and index (idx) of 1 will have an LSN of (decimal) 4294967297:

----------------------------------------------------------------------------- high 32-bits | low 32-bits 0b00000000 00000000 00000000 00000001 | 0b00000000 00000000 00000000 00000001 -----------------------------------------------------------------------------

Index

Constants

View Source
const (
	SizeBytes       = 8
	MaxEntrySize    = 100 * 1024 * 1024
	BufioWriterSize = 1024
)
View Source
const (
	LsnSize = 8
)

Variables

View Source
var (

	// WAL defaults
	InitialSegNum        uint32 = 1
	DefaultMinStoreBytes uint64 = 1024
	DefaultMinIndexBytes uint64 = 1024
	DefaultMaxStoreBytes uint64 = 16 * 1024
	DefaultMaxIndexBytes uint64 = 4 * 1024

	// WAL file parameters
	StoreExtension     = ".store"
	IndexExtension     = ".index"
	DamagedExtension   = ".damaged"
	CheckpointFileName = "checkpoint.sync"
	LogFilePerms       = os.FileMode(0640)

	// Errors
	EntryTooBig          = errors.New("entry exceeds maximum size of 100MB")
	LsnEntryNotFound     = errors.New("entry with LSN not found")
	IndexLsnError        = errors.New("incorrect LSN supplied to index")
	WriteHeaderError     = errors.New("written bytes to index header not equal to expected")
	WriteIndexEntryError = errors.New("written bytes to index not equal to expected")
	InvalidLogEntry      = errors.New("segment detected invalid log entry")
	InvalidIndexValues   = errors.New("segment detected invalid index values")
	WriteCheckpointError = errors.New("written bytes to checkpoint file not equal to expected")
	StopIteration        = errors.New("stop iteration")
	InvalidSyncIndex     = errors.New("invalid sync entry index")
	ContextCancelled     = errors.New("context cancelled")
	SegmentNotFound      = errors.New("no segment matches LSN")
)
View Source
var File_entry_proto protoreflect.FileDescriptor
View Source
var (
	InvalidEntryCRC = errors.New("log entry CRC is invalid")
)
View Source
var (
	LsnMask = uint64(0b11111111111111111111111111111111)
)

Functions

func ComputeCRC

func ComputeCRC(lsn uint64, data []byte) uint32

Compute CRC32 using LSN (log sequence number) + []byte data

func LsnDecode

func LsnDecode(lsn uint64) (uint32, uint32)

func LsnEncode

func LsnEncode(segNum, idx uint32) uint64

func Marshal

func Marshal(entry *LogEntry) ([]byte, error)

func Unmarshal

func Unmarshal(data []byte, entry *LogEntry) error

func ValidateCRC

func ValidateCRC(entry *LogEntry) error

func ValidateSegment

func ValidateSegment(s *Segment) error

Validate segment by reading all entries from index and store If any entries cannot be validated, the Read method will return an error. This function is called at Log initialization.

Types

type Checkpoint

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

The Checkpoint file header (uint64, 8 bytes) contains the latest checkpoint value. At initialization, either read the existing checkpoint value into "c" or write a starting value into the header and store it in "c" for subsequent reads.

func NewCheckpoint

func NewCheckpoint(dir string, val uint64) (*Checkpoint, error)

func (*Checkpoint) C

func (c *Checkpoint) C() uint64

func (*Checkpoint) Close

func (c *Checkpoint) Close() error

func (*Checkpoint) Open

func (c *Checkpoint) Open() error

func (*Checkpoint) Read

func (c *Checkpoint) Read() (uint64, error)

func (*Checkpoint) Remove

func (c *Checkpoint) Remove() error

func (*Checkpoint) Write

func (c *Checkpoint) Write(checkpoint uint64) error

type Config

type Config struct {
	Dir     string        `json:"dir"`
	Segment SegmentConfig `json:"segment"`
}

type Encoder

type Encoder interface {
	Encode(lsn uint64, data []byte) ([]byte, error)
	Decode(data []byte) (uint64, []byte, error)
	Valid(data []byte) error
}

func NewEncoder

func NewEncoder() Encoder

type Index

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

func NewIndex

func NewIndex(f *os.File, segNum uint32, c Config) (*Index, error)

func (*Index) Close

func (i *Index) Close() error

func (*Index) FileSize

func (i *Index) FileSize() uint64

func (*Index) IsEmpty

func (i *Index) IsEmpty() bool

func (*Index) Last

func (i *Index) Last() uint64

func (*Index) Name

func (i *Index) Name() string

func (*Index) Read

func (i *Index) Read(lsn uint64) (stored uint64, pos uint64, err error)

Read the entry in the current index corresponding to "lsn"

func (*Index) Start

func (i *Index) Start() uint64

func (*Index) Sync

func (i *Index) Sync() error

func (*Index) Write

func (i *Index) Write(lsn uint64, pos uint64) error

Write entry into index

type Log

type Log struct {
	Dir        string
	Config     Config
	Encoder    Encoder
	Checkpoint *Checkpoint // checkpoint to indicate latest synced LSN
	// contains filtered or unexported fields
}

func NewLog

func NewLog(c Config, e Encoder, logger *slog.Logger) (*Log, error)

func (*Log) Close

func (l *Log) Close() error

func (*Log) HighestIndex

func (l *Log) HighestIndex() uint64

func (*Log) Iterator

func (l *Log) Iterator(start uint64) (LogEntryIterator, error)

func (*Log) LowestIndex

func (l *Log) LowestIndex() uint64

func (*Log) Purge

func (l *Log) Purge(ctx context.Context, fn func(segNum uint32, segment *Segment) error) error

func (*Log) Read

func (l *Log) Read(lsn uint64) ([]byte, error)

read the entry corresponding to the Lsn

func (*Log) Remove

func (l *Log) Remove() error

func (*Log) Reset

func (l *Log) Reset() error

func (*Log) SegmentData

func (l *Log) SegmentData() []*SegmentData

func (*Log) Sync

func (l *Log) Sync() error

func (*Log) Write

func (l *Log) Write(data []byte) (uint64, error)

Write a new entry and return the index of the stored entry

func (*Log) WriteWithSync

func (l *Log) WriteWithSync(data []byte) (uint64, error)

Write a new entry with immediate sync

type LogEntry

type LogEntry struct {
	Lsn  uint64 `protobuf:"varint,1,opt,name=lsn,proto3" json:"lsn,omitempty"`
	Data []byte `protobuf:"bytes,2,opt,name=data,proto3" json:"data,omitempty"`
	CRC  uint32 `protobuf:"varint,3,opt,name=CRC,proto3" json:"CRC,omitempty"`
	// contains filtered or unexported fields
}

func (*LogEntry) Descriptor deprecated

func (*LogEntry) Descriptor() ([]byte, []int)

Deprecated: Use LogEntry.ProtoReflect.Descriptor instead.

func (*LogEntry) GetCRC

func (x *LogEntry) GetCRC() uint32

func (*LogEntry) GetData

func (x *LogEntry) GetData() []byte

func (*LogEntry) GetLsn

func (x *LogEntry) GetLsn() uint64

func (*LogEntry) ProtoMessage

func (*LogEntry) ProtoMessage()

func (*LogEntry) ProtoReflect

func (x *LogEntry) ProtoReflect() protoreflect.Message

func (*LogEntry) Reset

func (x *LogEntry) Reset()

func (*LogEntry) String

func (x *LogEntry) String() string

type LogEntryIterator

type LogEntryIterator interface {
	Next(ctx context.Context) (lsn uint64, data []byte, err error)
}

type Segment

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

func NewSegment

func NewSegment(dir string, segNum uint32, config Config, encoder Encoder) (*Segment, error)

func RepairSegment

func RepairSegment(s *Segment, logger *slog.Logger) (*Segment, error)

Repair a damaged segment by migrating entries to a new segment

func (*Segment) Close

func (s *Segment) Close() error

func (*Segment) IsMaxed

func (s *Segment) IsMaxed() bool

func (*Segment) Last

func (s *Segment) Last() uint64

func (*Segment) Next

func (s *Segment) Next() uint64

func (*Segment) Open

func (s *Segment) Open() error

func (*Segment) Read

func (s *Segment) Read(lsn uint64) ([]byte, error)

func (*Segment) Remove

func (s *Segment) Remove() error

func (*Segment) Start

func (s *Segment) Start() uint64

func (*Segment) Sync

func (s *Segment) Sync() error

func (*Segment) Write

func (s *Segment) Write(data []byte) (lsn uint64, err error)

Write entry to the log file

type SegmentConfig

type SegmentConfig struct {
	MaxStoreBytes uint64 `json:"max_store_bytes"`
	MaxIndexBytes uint64 `json:"max_index_bytes"`
}

type SegmentData

type SegmentData struct {
	SegNum    uint32
	Start     uint64
	Last      uint64
	StartIdx  uint32
	LastIdx   uint32
	Entries   uint32
	IndexSize uint64
	StoreSize uint64
}

type Store

type Store struct {
	*os.File
	// contains filtered or unexported fields
}

func NewStore

func NewStore(f *os.File) (*Store, error)

func (*Store) Close

func (s *Store) Close() error

func (*Store) FileSize

func (s *Store) FileSize() uint64

func (*Store) Name

func (s *Store) Name() string

func (*Store) Read

func (s *Store) Read(pos uint64) ([]byte, error)

When reading an entry: - read the size of the data first (8 bytes -> SizeBytes) - then read "size" bytes of data starting at offset ("pos" + SizeBytes)

func (*Store) ReadAt

func (s *Store) ReadAt(p []byte, off int64) (int, error)

func (*Store) Sync

func (s *Store) Sync() error

func (*Store) Write

func (s *Store) Write(p []byte) (n uint64, pos uint64, err error)

Each data entry is stored in file as: ------------------ size | data 8 bytes | N bytes ------------------

When appending, write: - the size of the data (as uint64) - then the actual data (as bytes)

Returns: n: number of bytes written pos: start position of entry err: error if error occurred, or nil

Jump to

Keyboard shortcuts

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