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
- Variables
- func ComputeCRC(lsn uint64, data []byte) uint32
- func LsnDecode(lsn uint64) (uint32, uint32)
- func LsnEncode(segNum, idx uint32) uint64
- func Marshal(entry *LogEntry) ([]byte, error)
- func Unmarshal(data []byte, entry *LogEntry) error
- func ValidateCRC(entry *LogEntry) error
- func ValidateSegment(s *Segment) error
- type Checkpoint
- type Config
- type Encoder
- type Index
- func (i *Index) Close() error
- func (i *Index) FileSize() uint64
- func (i *Index) IsEmpty() bool
- func (i *Index) Last() uint64
- func (i *Index) Name() string
- func (i *Index) Read(lsn uint64) (stored uint64, pos uint64, err error)
- func (i *Index) Start() uint64
- func (i *Index) Sync() error
- func (i *Index) Write(lsn uint64, pos uint64) error
- type Log
- func (l *Log) Close() error
- func (l *Log) HighestIndex() uint64
- func (l *Log) Iterator(start uint64) (LogEntryIterator, error)
- func (l *Log) LowestIndex() uint64
- func (l *Log) Purge(ctx context.Context, fn func(segNum uint32, segment *Segment) error) error
- func (l *Log) Read(lsn uint64) ([]byte, error)
- func (l *Log) Remove() error
- func (l *Log) Reset() error
- func (l *Log) SegmentData() []*SegmentData
- func (l *Log) Sync() error
- func (l *Log) Write(data []byte) (uint64, error)
- func (l *Log) WriteWithSync(data []byte) (uint64, error)
- type LogEntry
- func (*LogEntry) Descriptor() ([]byte, []int)deprecated
- func (x *LogEntry) GetCRC() uint32
- func (x *LogEntry) GetData() []byte
- func (x *LogEntry) GetLsn() uint64
- func (*LogEntry) ProtoMessage()
- func (x *LogEntry) ProtoReflect() protoreflect.Message
- func (x *LogEntry) Reset()
- func (x *LogEntry) String() string
- type LogEntryIterator
- type Segment
- func (s *Segment) Close() error
- func (s *Segment) IsMaxed() bool
- func (s *Segment) Last() uint64
- func (s *Segment) Next() uint64
- func (s *Segment) Open() error
- func (s *Segment) Read(lsn uint64) ([]byte, error)
- func (s *Segment) Remove() error
- func (s *Segment) Start() uint64
- func (s *Segment) Sync() error
- func (s *Segment) Write(data []byte) (lsn uint64, err error)
- type SegmentConfig
- type SegmentData
- type Store
Constants ¶
const ( SizeBytes = 8 MaxEntrySize = 100 * 1024 * 1024 BufioWriterSize = 1024 )
const (
LsnSize = 8
)
Variables ¶
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") )
var File_entry_proto protoreflect.FileDescriptor
var (
InvalidEntryCRC = errors.New("log entry CRC is invalid")
)
var (
LsnMask = uint64(0b11111111111111111111111111111111)
)
Functions ¶
func ComputeCRC ¶
Compute CRC32 using LSN (log sequence number) + []byte data
func ValidateCRC ¶
func ValidateSegment ¶
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
}
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 (*Log) HighestIndex ¶
func (*Log) LowestIndex ¶
func (*Log) SegmentData ¶
func (l *Log) SegmentData() []*SegmentData
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) ProtoMessage ¶
func (*LogEntry) ProtoMessage()
func (*LogEntry) ProtoReflect ¶
func (x *LogEntry) ProtoReflect() protoreflect.Message
type LogEntryIterator ¶
type Segment ¶
type Segment struct {
// contains filtered or unexported fields
}
func NewSegment ¶
func RepairSegment ¶
Repair a damaged segment by migrating entries to a new segment
type SegmentConfig ¶
type SegmentData ¶
type Store ¶
func (*Store) Read ¶
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) Write ¶
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