This commit is contained in:
2026-06-09 16:42:17 +00:00
parent 7007b1b6fd
commit 4089ca40c9
11 changed files with 597 additions and 406 deletions

204
database/wal_replayer.go Normal file
View File

@@ -0,0 +1,204 @@
package database
import (
"errors"
"fmt"
qb "gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/txlog"
)
type WALReplayer struct {
walReader *txlog.Reader
metrics map[uint32]*ReplayMetric
freeIndexPages []uint32
freeDataPages []uint32
indexPages []txlog.PageToWrite
dataPages []txlog.PageToWrite
}
type WALReplayerOptions struct {
WALReader *txlog.Reader
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
}
func (s *WALReplayer) Replay() error {
for {
records, isLast, err := s.walReader.NextPacket()
if err != nil {
return err
}
if isLast {
if len(records) > 0 {
var appendRecords []txlog.WALRecordAppendMeasures
for _, untyped := range records {
rec, ok := untyped.(txlog.WALRecordAppendMeasures)
if ok {
appendRecords = append(appendRecords, rec)
}
}
// prepare pages and rewrite
indexPages, dataPages := pagesToWriteFromWALRecords(metrics, appendRecords)
err = s.txlog.WritePagesToAtree(indexPages, dataPages)
if err != nil {
return err
}
}
// для останнього rec - повторюємо запис сторінок
// rewrite write pages
//s.txlog.
//return nil
}
for _, rec := range records {
if err = s.replayWALRecord(rec); err != nil {
return err
}
}
}
}
// func (s *WALReplayer) replayPacket(records []any, isLastPacket bool) {
// for _, untyped := range records {
// switch rec := untyped.(type) {
// case txlog.WALRecordAppendMeasures:
// s.onWALRecordAppendMeasures(rec, isLastPacket)
// }
// }
// return
// }
func (s *WALReplayer) replayWALRecord(untyped any) (err error) {
switch rec := untyped.(type) {
case txlog.AddedMetric:
if err = s.onAddedMetric(rec); err != nil {
return
}
case txlog.DeletedMetric:
if err = s.onDeletedMetric(rec); err != nil {
return
}
// case txlog.AppendedMeasure:
// metric, ok := s.metrics[rec.MetricID]
// if ok {
// metric.Timestamps.Append(rec.Timestamp)
// metric.Values.Append(rec.Value)
// if metric.Since == 0 {
// metric.Since = rec.Timestamp
// metric.SinceValue = rec.Value
// }
// metric.Until = rec.Timestamp
// metric.UntilValue = rec.Value
// }
case txlog.WALRecordAppendMeasures:
// metric, ok := s.metrics[rec.MetricID]
// if !ok {
// qb.Abort(qb.MetricNotFoundDuringReplay, nil)
// }
//metric.ReplayAppendedMeasures(rec) FIX
// case txlog.DeletedMeasures:
// metric, ok := s.metrics[rec.MetricID]
// if ok {
// metric.DeleteMeasures()
// if len(rec.FreePageNumbers) > 0 {
// s.freeList.AddPages(rec.FreePageNumbers)
// }
// }
default:
qb.Abort(qb.UnknownTxLogRecordTypeBug,
fmt.Errorf("bug: unknown record type %T in TransactionLog", rec))
}
return nil
}
func (s *WALReplayer) onAddedMetric(rec txlog.AddedMetric) (err error) {
_, ok := s.metrics[rec.MetricID]
if ok {
return fmt.Errorf("metric %d add failed: already added", rec.MetricID)
}
var (
buf = make([]byte, atree.DataPageSize)
)
s.metrics[rec.MetricID] = &ReplayMetric{
metricType: rec.MetricType,
fracDigits: byte(rec.FracDigits),
buf: buf,
databuf: buf[:atree.DataPagePayloadSize],
}
return
}
func (s *WALReplayer) onDeletedMetric(rec txlog.DeletedMetric) error {
_, 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
}
func (s *WALReplayer) onAppendMeasures(rec txlog.WALRecordAppendMeasures, isLastPacket bool) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.MetricNotFoundDuringReplay, nil)
}
if isLastPacket {
s.collectPagesToWrite(metric, rec)
}
result := metric.AppendedMeasures(rec, isLastPacket)
//
if isLastPacket {
s.indexPages = append(s.indexPages, result.IndexPages...)
s.dataPages = append(s.dataPages, result.DataPages...)
}
//
s.freeIndexPages = s.freeIndexPages[:len(s.freeIndexPages)-result.ReusedIndexPagesCount]
s.freeDataPages = s.freeDataPages[:len(s.freeDataPages)-result.ReusedDataPagesCount]
}
func (s *WALReplayer) collectPagesToWrite(metric *ReplayMetric, rec txlog.WALRecordAppendMeasures) {
metric.ApplyHeadDataPage(rec.CompletedDataPage)
s.dataPages = append(s.dataPages, txlog.PageToWrite{
PageNo: p.PageNo,
Content: metric.buf,
})
for _, x := range rec.DataPages {
s.dataPages = append(s.dataPages, txlog.PageToWrite{
PageNo: x.PageNo,
Content: x.Content,
})
}
//metric.ApplyHeadIndexPages(rec.ChangedIndexLevels)
}