Files
qb/storage/wal_replayer.go

158 lines
3.9 KiB
Go
Raw Permalink Normal View History

2026-06-13 22:43:17 +00:00
package storage
2026-06-09 16:42:17 +00:00
import (
"errors"
"fmt"
)
type WALReplayer struct {
2026-06-13 22:43:17 +00:00
walReader *WALReader
2026-06-09 16:42:17 +00:00
metrics map[uint32]*ReplayMetric
freeIndexPages []uint32
freeDataPages []uint32
2026-06-13 22:43:17 +00:00
indexPages []PageToWrite
dataPages []PageToWrite
2026-06-09 16:42:17 +00:00
}
type WALReplayerOptions struct {
2026-06-13 22:43:17 +00:00
WALReader *WALReader
2026-06-09 16:42:17 +00:00
Metrics map[uint32]*ReplayMetric // from snapshot
FreeIndexPages []uint32 // from snapshot
FreeDataPages []uint32 // from snapshot
}
func NewWALReplayer(opt WALReplayerOptions) (*WALReplayer, error) {
if opt.WALReader == nil {
return nil, errors.New("required option: WALReader")
}
if opt.Metrics == nil {
return nil, errors.New("required option: Metrics")
}
return &WALReplayer{
walReader: opt.WALReader,
metrics: opt.Metrics,
freeIndexPages: opt.FreeIndexPages,
freeDataPages: opt.FreeDataPages,
}, nil
}
2026-06-14 07:57:01 +03:00
func (s *WALReplayer) IndexPagesToRewrite() []PageToWrite {
2026-06-10 06:18:45 +03:00
return s.indexPages
}
2026-06-14 07:57:01 +03:00
func (s *WALReplayer) DataPagesToRewrite() []PageToWrite {
2026-06-10 06:18:45 +03:00
return s.dataPages
}
2026-06-14 07:57:01 +03:00
func (s *WALReplayer) FreeIndexPages() []uint32 {
return s.freeIndexPages
}
func (s *WALReplayer) FreeDataPages() []uint32 {
return s.freeDataPages
}
2026-06-09 16:42:17 +00:00
func (s *WALReplayer) Replay() error {
for {
2026-06-10 06:18:45 +03:00
records, isLastPacket, err := s.walReader.NextPacket()
2026-06-09 16:42:17 +00:00
if err != nil {
return err
}
2026-06-14 23:12:03 +03:00
if len(records) == 0 {
return nil
}
2026-06-09 16:42:17 +00:00
for _, rec := range records {
2026-06-10 06:18:45 +03:00
if err = s.replayRecord(rec, isLastPacket); err != nil {
2026-06-09 16:42:17 +00:00
return err
}
}
}
}
2026-06-10 06:18:45 +03:00
func (s *WALReplayer) replayRecord(untyped any, isLastPacket bool) (err error) {
2026-06-09 16:42:17 +00:00
switch rec := untyped.(type) {
2026-06-13 22:43:17 +00:00
case MeasuresAppendRecord:
2026-06-10 06:18:45 +03:00
if err = s.onMeasuresAppend(rec); err != nil {
2026-06-09 16:42:17 +00:00
return
}
2026-06-13 22:43:17 +00:00
case MeasuresAppendWithGrowRecord:
2026-06-10 06:18:45 +03:00
if err = s.onMeasuresAppendWithGrow(rec, isLastPacket); err != nil {
2026-06-09 16:42:17 +00:00
return
}
2026-06-13 22:43:17 +00:00
case MetricAddRecord:
2026-06-10 06:18:45 +03:00
if err = s.onMetricAdd(rec); err != nil {
return
}
2026-06-13 22:43:17 +00:00
case MetricDeleteRecord:
2026-06-10 06:18:45 +03:00
if err = s.onMetricDelete(rec); err != nil {
return
}
2026-06-13 22:43:17 +00:00
// case DeletedMeasures:
2026-06-10 06:18:45 +03:00
// metric, ok := s.metrics[rec.MetricID]
// if ok {
// metric.DeleteMeasures()
// if len(rec.FreePageNumbers) > 0 {
// s.freeList.AddPages(rec.FreePageNumbers)
// }
// }
// default:
// qb.Abort(qb.UnknownstorageRecordTypeBug,
// fmt.Errorf("bug: unknown record type %T in TransactionLog", rec))
2026-06-09 16:42:17 +00:00
}
return nil
}
2026-06-13 22:43:17 +00:00
func (s *WALReplayer) onMetricAdd(rec MetricAddRecord) (err error) {
2026-06-09 16:42:17 +00:00
_, ok := s.metrics[rec.MetricID]
if ok {
return fmt.Errorf("metric %d add failed: already added", rec.MetricID)
}
var (
2026-06-13 22:43:17 +00:00
buf = make([]byte, DataPageSize)
2026-06-09 16:42:17 +00:00
)
s.metrics[rec.MetricID] = &ReplayMetric{
2026-06-14 07:57:01 +03:00
MetricType: rec.MetricType,
FracDigits: rec.FracDigits,
Buf: buf,
2026-06-09 16:42:17 +00:00
}
return
}
2026-06-13 22:43:17 +00:00
func (s *WALReplayer) onMetricDelete(rec MetricDeleteRecord) error {
2026-06-09 16:42:17 +00:00
_, ok := s.metrics[rec.MetricID]
if !ok {
return fmt.Errorf("metric %d deletion failed: not found", rec.MetricID)
}
delete(s.metrics, rec.MetricID)
s.freeIndexPages = append(s.freeIndexPages, rec.FreeIndexPages...)
s.freeDataPages = append(s.freeDataPages, rec.FreeDataPages...)
return nil
}
2026-06-13 22:43:17 +00:00
func (s *WALReplayer) onMeasuresAppend(rec MeasuresAppendRecord) (err error) {
2026-06-09 16:42:17 +00:00
metric, ok := s.metrics[rec.MetricID]
if !ok {
2026-06-10 06:18:45 +03:00
return fmt.Errorf("append measures failed: metric %d not found", rec.MetricID)
2026-06-09 16:42:17 +00:00
}
2026-06-10 06:18:45 +03:00
metric.MeasuresAppend(rec)
return
}
2026-06-13 22:43:17 +00:00
func (s *WALReplayer) onMeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (err error) {
2026-06-10 06:18:45 +03:00
metric, ok := s.metrics[rec.MetricID]
if !ok {
return fmt.Errorf("append measures failed: metric %d not found", rec.MetricID)
2026-06-09 16:42:17 +00:00
}
2026-06-10 06:18:45 +03:00
result := metric.MeasuresAppendWithGrow(rec, isLastPacket)
2026-06-09 16:42:17 +00:00
//
if isLastPacket {
2026-06-14 07:57:01 +03:00
s.indexPages = append(s.indexPages, result.IndexPagesToRewrite...)
s.dataPages = append(s.dataPages, result.DataPagesToRewrite...)
2026-06-09 16:42:17 +00:00
}
//
s.freeIndexPages = s.freeIndexPages[:len(s.freeIndexPages)-result.ReusedIndexPagesCount]
s.freeDataPages = s.freeDataPages[:len(s.freeDataPages)-result.ReusedDataPagesCount]
2026-06-10 06:18:45 +03:00
return
2026-06-09 16:42:17 +00:00
}