From 4089ca40c9db4eca4d2a991fbf5f752d342ecdfb Mon Sep 17 00:00:00 2001 From: dima <1.e4.kc6@gmail.com> Date: Tue, 9 Jun 2026 16:42:17 +0000 Subject: [PATCH] wp --- database/database.go | 309 ++++------------------------------------------ database/database_test.go | 89 +++++++++++++ database/metric.go | 27 ++-- database/replay_metric.go | 190 ++++++++++++++++++++++++++++ database/snapshot.go | 105 ++++++++-------- database/wal_replayer.go | 204 ++++++++++++++++++++++++++++++ enc/cumdelta.go | 13 +- enc/insdelta.go | 13 +- enc/time_delta.go | 13 +- txlog/txlog.go | 39 +++--- txlog/wal_records.go | 1 - 11 files changed, 597 insertions(+), 406 deletions(-) create mode 100644 database/database_test.go create mode 100644 database/replay_metric.go create mode 100644 database/wal_replayer.go diff --git a/database/database.go b/database/database.go index 7777e46..5d97dbb 100644 --- a/database/database.go +++ b/database/database.go @@ -3,7 +3,6 @@ package database import ( "errors" "fmt" - "io" "log" "net" "os" @@ -11,15 +10,17 @@ import ( "sync" "time" - bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/freelist" "gordenko.dev/dima/qb/txlog" ) -func JoinSnapshotFileName(dir string, logNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.snapshot", logNumber)) +func JoinSnapshotFileName(dir string, snapshotNumber int) string { + return filepath.Join(dir, fmt.Sprintf("%d.snapshot", snapshotNumber)) +} + +func JoinWALFileName(dir string, snapshotNumber int) string { + return filepath.Join(dir, fmt.Sprintf("%d.wal", snapshotNumber)) } // type metricLockEntry struct { @@ -242,232 +243,37 @@ func (s *Database) recovery() { // зробити object? -func (s *Database) replayChanges(fileName string) error { +func (s *Database) replayChanges(snapshotNumber int) (err error) { + snapshot, err := readSnapshot(JoinSnapshotFileName(s.dir, snapshotNumber)) + if err != nil { + return + } + walReader, err := txlog.NewReader(txlog.ReaderOptions{ - FileName: fileName, - BufferSize: 1024 * 1024, + FileName: JoinWALFileName(s.dir, snapshotNumber), + BufferSize: 8 * 1024 * 1024, }) if err != nil { return err } - var metrics map[uint32]*ReplayMetric - - for { - records, isLast, err := 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 _, record := range records { - if err = s.replayWALRecord(record); err != nil { - return err - } - } + walReplayer, err := NewWALReplayer(WALReplayerOptions{ + WALReader: walReader, + Metrics: snapshot.Metrics, + FreeIndexPages: snapshot.IndexPageNumbers, + FreeDataPages: snapshot.DataPageNumbers, + }) + if err != nil { + return } -} -func (s *Database) replayWALRecord(untyped any) error { - switch rec := untyped.(type) { - case txlog.AddedMetric: - s.addMetric(rec) - - case txlog.DeletedMetric: - // pages add and delete to tmp file - - // delete(s.metrics, rec.MetricID) - // if len(rec.FreePageNumbers) > 0 { - // s.freeList.AddPages(rec.FreePageNumbers) - // } - - // 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)) + err = walReplayer.Replay() + if err != nil { + return } return nil } -type ReplayMetric struct { - metricType qb.MetricType - fracDigits byte - lastPageNo uint32 - buf []byte - databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер) - tSize int - vSize int - indexLevelTails []txlog.IndexLevelTail // root - last element -} - -func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { - metricType, err := bin.ReadByte(r) - if err != nil { - return - } - s.metricType = qb.MetricType(metricType) - s.fracDigits, err = bin.ReadByte(r) - if err != nil { - return - } - s.lastPageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - s.tSize, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - s.vSize, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - s.buf, s.databuf = allocateBuffers(s.tSize + s.vSize) - // - err = bin.ReadNInto(r, s.databuf[len(s.databuf)-s.tSize:]) - if err != nil { - return - } - err = bin.ReadNInto(r, s.databuf[:s.vSize]) - if err != nil { - return - } - // - levelsCount, err := bin.ReadVarSize(r) - if err != nil { - return - } - for range levelsCount { - var ( - buf = make([]byte, atree.IndexPageSize) - count int - ) - count, err = bin.ReadVarSize(r) - if err != nil { - return - } - err = bin.ReadNInto(r, buf[:count*indexRecordSize]) - if err != nil { - return - } - s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ - Buffer: buf, - RecordsCount: count, - }) - } - return -} - -func (s *ReplayMetric) ReplayAppendedMeasures(rec txlog.WALRecordAppendMeasures) { - // Додати в index level tails недостаючі дані, або замінити - for levelIdx, change := range rec.ChangedIndexLevels { - if levelIdx < len(s.indexLevelTails) { - level := s.indexLevelTails[levelIdx] - if change.CompletedIndexPage != nil { - // попередній хвіст індексного рівня перетворився на сторінку, - // отже TailRecords - це новий хвіст - // або копіюю дані, або алокую більший буфер? - copy(level.Buffer, change.TailRecords) - level.RecordsCount = len(change.TailRecords) / indexRecordSize - } else { - // TailRecords - це нові дані, які треба додати - pos := level.RecordsCount * indexRecordSize - copy(level.Buffer[pos:], change.TailRecords) - } - s.indexLevelTails[levelIdx] = level - } else { - var ( - recordsCount = len(change.TailRecords) / indexRecordSize - buffer = make([]byte, atree.IndexPageSize) - ) - copy(buffer, change.TailRecords) - s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ - Buffer: buffer, - RecordsCount: recordsCount, - }) - } - } - - // Перезаписати в buffer timestamps та values. Якщо розмір буферу недостатній - - // спочатку створити більший - - var ( - // CompletedDataPage є полюбому, оскільки це Record із мінімум одною заповненою - // data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня. - - requiredSpace = len(rec.TailTimestamps) + len(rec.TailValues) - databuf []byte - ) - if requiredSpace < len(s.buf) { - databuf = s.buf - // fix - update sizes? - } else { - var buf []byte - buf, databuf = allocateBuffers(requiredSpace) - copy(databuf, rec.TailValues) - pos := len(databuf) - len(rec.TailTimestamps) - copy(databuf[pos:], rec.TailTimestamps) - // fix - replace and remeber sizes? - - s.buf = buf - } - - //rec.TailTimestamps - //rec.TailValues -} - func (s *Database) rewritePagesFromLastWALPacket(list []any) error { // var ( // indexPages []txlog.PageToWrite @@ -490,71 +296,6 @@ func (s *Database) rewritePagesFromLastWALPacket(list []any) error { return nil } -func pagesToWriteFromWALRecords(metrics map[uint32]*ReplayMetric, records []txlog.WALRecordAppendMeasures) (indexPages []txlog.PageToWrite, dataPages []txlog.PageToWrite) { - for _, rec := range records { - metric, ok := metrics[rec.MetricID] - if !ok { - qb.Abort(qb.MetricNotFoundDuringReplay, nil) - } - p := rec.CompletedDataPage - // буфер має бути розміру сторінки - if len(metric.buf) != atree.DataPageSize { - metric.buf, metric.databuf = growBuffers(growBuffersIn{ - databuf: metric.databuf, - tSize: metric.tSize, - vSize: metric.vSize, - requiredSpace: atree.DataPageSize, - }) - } - // дописую дані по зміщенню - pos := len(metric.databuf) - p.TimestampsOffset - len(p.Timestamps) - copy(metric.databuf[pos:], p.Timestamps) - copy(metric.databuf[p.ValuesOffset:], p.Values) - // - metric.tSize = p.TimestampsOffset + len(p.Timestamps) - metric.tSize = p.ValuesOffset + len(p.Values) - - // fix - заполнить страницу - // fix - check checksum - - dataPages = append(dataPages, txlog.PageToWrite{ - PageNo: p.PageNo, - Content: metric.buf, - }) - - for _, x := range rec.DataPages { - dataPages = append(dataPages, txlog.PageToWrite{ - PageNo: x.PageNo, - Content: x.Content, - }) - } - - for levelIdx, changedLevel := range rec.ChangedIndexLevels { - p := changedLevel.CompletedIndexPage - if p != nil { - level := metric.indexLevelTails[levelIdx] - copy(level.Buffer[level.RecordsCount*indexRecordSize:], p.Records) - - // fix - заполнить страницу - // fix - check checksum - - indexPages = append(dataPages, txlog.PageToWrite{ - PageNo: p.PageNo, - Content: level.Buffer, - }) - } - - for _, x := range changedLevel.IndexPages { - indexPages = append(dataPages, txlog.PageToWrite{ - PageNo: x.PageNo, - Content: x.Content, - }) - } - } - } - return -} - // FIX // УВАГА! // Якщо буфер metric.buffer < розміру сторінки - для timestamps доступний весь буфер. diff --git a/database/database_test.go b/database/database_test.go new file mode 100644 index 0000000..4ea59db --- /dev/null +++ b/database/database_test.go @@ -0,0 +1,89 @@ +package database + +import ( + "fmt" + "slices" + "testing" + + "gordenko.dev/dima/qb" +) + +func TestWriteReadSnapshot(t *testing.T) { + metrics := make(map[uint32]*_metric) + metrics[1] = &_metric{ + metricType: qb.Cumulative, + fracDigits: 3, + lastPageNo: 1000, + buffer: []byte{}, + } + in := writeSnapshotIn{ + SnapshotNumber: 1, + Dir: ".", + WriteBufferSize: 8 * 1024 * 1024, + Metrics: metrics, + FrozenIndexPagesCount: 7, + IndexPageNumbers: []uint32{ + 1, 2, 3, 4, 5, + }, + FrozenDataPagesCount: 8, + DataPageNumbers: []uint32{ + 6, 7, 8, + }, + } + err := writeSnapshot(in) + if err != nil { + t.Fatalf("writeSnapshot: %s", err) + } + + out, err := readSnapshot("1.snapshot") + if err != nil { + return + } + + if out.FrozenIndexPagesCount != in.FrozenIndexPagesCount { + t.Fatalf("FrozenIndexPagesCount: got %d are not equal expected %d", + out.FrozenIndexPagesCount, in.FrozenIndexPagesCount) + } + if out.FrozenDataPagesCount != in.FrozenDataPagesCount { + t.Fatalf("FrozenDataPagesCount: got %d are not equal expected %d", + out.FrozenDataPagesCount, in.FrozenDataPagesCount) + } + if !slices.Equal(out.IndexPageNumbers, in.IndexPageNumbers) { + t.Fatalf("IndexPageNumbers: got %v are not equal expected %v", + out.IndexPageNumbers, in.IndexPageNumbers) + } + if !slices.Equal(out.DataPageNumbers, in.DataPageNumbers) { + t.Fatalf("DataPageNumbers: got %v are not equal expected %v", + out.DataPageNumbers, in.DataPageNumbers) + } + for metricID, replay := range out.Metrics { + origin, ok := metrics[metricID] + if !ok { + t.Fatalf("decoded metricID %d not found in metrics", metricID) + } + err = cmpMetric(replay, origin) + if err != nil { + t.Fatalf("metric %d: %s", metricID, err) + } + } +} + +func cmpMetric(replay *ReplayMetric, origin *_metric) error { + if replay.metricType != origin.metricType { + return fmt.Errorf("metricType: got %d are not equal expected %d", + replay.metricType, origin.metricType) + } + if replay.fracDigits != origin.fracDigits { + return fmt.Errorf("fracDigits: got %d are not equal expected %d", + replay.fracDigits, origin.fracDigits) + } + if replay.lastPageNo != origin.lastPageNo { + return fmt.Errorf("lastPageNo: got %d are not equal expected %d", + replay.lastPageNo, origin.lastPageNo) + } + // if !refect.DeepEq{ + // return fmt.Errorf("lastPageNo: got %d are not equal expected %d", + // replay.lastPageNo, origin.lastPageNo) + // } + return nil +} diff --git a/database/metric.go b/database/metric.go index a5329d4..6479d4a 100644 --- a/database/metric.go +++ b/database/metric.go @@ -330,10 +330,9 @@ func (s *_metric) StartFullScan(req tryFullScanReq) { // metricType - 1b // fracDigits - 1b // lastPageNo - 4b -// until - 4b // timestamps size - 2b -// timestams payload - Nb // values size - 2b +// timestams payload - Nb // values payload - Nb // index levels count - varsize // [ @@ -353,29 +352,20 @@ func (s *_metric) WriteTo(w io.Writer) (err error) { if err != nil { return } - // err = bin.WriteUint32(w, s.Since) // fix unixtime in timestamps - // if err != nil { - // return - // } - // err = bin.WriteFloat64(w, s.SinceValue) - // if err != nil { - // return - // } - err = bin.WriteUint32(w, s.timestamps.LastTimestamp()) + err = bin.WriteUint16(w, uint16(s.timestamps.Size())) if err != nil { return } - // err = bin.WriteFloat64(w, s.UntilValue) - // if err != nil { - // return - // } - // FIX - write sizes, then payloads - // copy timestamps payload + err = bin.WriteUint16(w, uint16(s.values.Size())) + if err != nil { + return + } + // timestamps payload err = s.timestamps.WritePayloadTo(w) if err != nil { return } - // copy values payload + // values payload err = s.values.WritePayloadTo(w) if err != nil { return @@ -394,7 +384,6 @@ func (s *_metric) WriteTo(w io.Writer) (err error) { if err != nil { return } - } return } diff --git a/database/replay_metric.go b/database/replay_metric.go new file mode 100644 index 0000000..8a9fbbc --- /dev/null +++ b/database/replay_metric.go @@ -0,0 +1,190 @@ +package database + +import ( + "io" + + bin "gordenko.dev/dima/bin/little" + qb "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/txlog" +) + +type ReplayMetric struct { + metricType qb.MetricType + fracDigits byte + lastPageNo uint32 + buf []byte + databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер) + tSize int + vSize int + indexLevelTails []txlog.IndexLevelTail // root - last element +} + +func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { + metricType, err := bin.ReadByte(r) + if err != nil { + return + } + s.metricType = qb.MetricType(metricType) + s.fracDigits, err = bin.ReadByte(r) + if err != nil { + return + } + s.lastPageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.tSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + s.vSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + err = bin.ReadNInto(r, s.databuf[len(s.databuf)-s.tSize:]) + if err != nil { + return + } + err = bin.ReadNInto(r, s.databuf[:s.vSize]) + if err != nil { + return + } + // + levelsCount, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range levelsCount { + var ( + buf = make([]byte, atree.IndexPageSize) + count int + ) + count, err = bin.ReadVarSize(r) + if err != nil { + return + } + err = bin.ReadNInto(r, buf[:count*indexRecordSize]) + if err != nil { + return + } + s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ + Buffer: buf, + RecordsCount: count, + }) + } + return +} + +type AppendMeasuresResult struct { + IndexPages []txlog.PageToWrite + DataPages []txlog.PageToWrite + ReusedIndexPagesCount int + ReusedDataPagesCount int +} + +func (s *ReplayMetric) AppendedMeasures(rec txlog.WALRecordAppendMeasures, pagesNeeded bool) (result AppendMeasuresResult) { + // Додати в index level tails недостаючі дані, або замінити + for levelIdx, change := range rec.ChangedIndexLevels { + if levelIdx < len(s.indexLevelTails) { + level := s.indexLevelTails[levelIdx] + if change.CompletedIndexPage != nil { + // попередній хвіст індексного рівня перетворився на сторінку, + // отже TailRecords - це новий хвіст + // або копіюю дані, або алокую більший буфер? + copy(level.Buffer, change.TailRecords) + level.RecordsCount = len(change.TailRecords) / indexRecordSize + } else { + // TailRecords - це нові дані, які треба додати + pos := level.RecordsCount * indexRecordSize + copy(level.Buffer[pos:], change.TailRecords) + } + s.indexLevelTails[levelIdx] = level + } else { + var ( + recordsCount = len(change.TailRecords) / indexRecordSize + buffer = make([]byte, atree.IndexPageSize) + ) + copy(buffer, change.TailRecords) + s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ + Buffer: buffer, + RecordsCount: recordsCount, + }) + } + } + + // Перезаписати в buffer timestamps та values. Якщо розмір буферу недостатній - + // спочатку створити більший + + var ( + // CompletedDataPage є полюбому, оскільки це Record із мінімум одною заповненою + // data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня. + + requiredSpace = len(rec.TailTimestamps) + len(rec.TailValues) + databuf []byte + ) + if requiredSpace < len(s.buf) { + databuf = s.buf + // fix - update sizes? + } else { + var buf []byte + buf, databuf = allocateBuffers(requiredSpace) + copy(databuf, rec.TailValues) + pos := len(databuf) - len(rec.TailTimestamps) + copy(databuf[pos:], rec.TailTimestamps) + // fix - replace and remeber sizes? + + s.buf = buf + } + + //rec.TailTimestamps + //rec.TailValues + return +} + +func (s *ReplayMetric) ApplyHeadDataPage(p txlog.CompletedDataPage) { + // буфер має бути розміру сторінки + if len(s.buf) != atree.DataPageSize { + s.buf, s.databuf = growBuffers(growBuffersIn{ + databuf: s.databuf, + tSize: s.tSize, + vSize: s.vSize, + requiredSpace: atree.DataPageSize, + }) + } + // дописую дані по зміщенню + pos := len(s.databuf) - p.TimestampsOffset - len(p.Timestamps) + copy(s.databuf[pos:], p.Timestamps) + copy(s.databuf[p.ValuesOffset:], p.Values) + // + s.tSize = p.TimestampsOffset + len(p.Timestamps) + s.tSize = p.ValuesOffset + len(p.Values) + + // fix - заполнить страницу + // fix - check checksum +} + +// func (s *ReplayMetric) ApplyHeadIndexPages(changedIndexLevels txlog.ChangedIndexLevel) { +// for levelIdx, changedLevel := range changedIndexLevels { +// p := changedLevel.CompletedIndexPage +// if p != nil { +// level := s.indexLevelTails[levelIdx] +// copy(level.Buffer[level.RecordsCount*indexRecordSize:], p.Records) + +// // fix - заполнить страницу +// // fix - check checksum + +// s.indexPages = append(s.indexPages, txlog.PageToWrite{ +// PageNo: p.PageNo, +// Content: level.Buffer, +// }) +// } + +// for _, x := range changedLevel.IndexPages { +// s.indexPages = append(s.indexPages, txlog.PageToWrite{ +// PageNo: x.PageNo, +// Content: x.Content, +// }) +// } +// } +// } diff --git a/database/snapshot.go b/database/snapshot.go index 0d94280..69d90a7 100644 --- a/database/snapshot.go +++ b/database/snapshot.go @@ -8,6 +8,7 @@ import ( "path/filepath" bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/util" ) @@ -70,35 +71,15 @@ func writeSnapshot(in writeSnapshotIn) (err error) { } } // free index pages - _, err = bin.WriteVarSize(dst, in.FrozenIndexPagesCount) + err = writeFreePages(dst, in.FrozenIndexPagesCount, in.IndexPageNumbers) if err != nil { return } - _, err = bin.WriteVarSize(dst, len(in.IndexPageNumbers)) - if err != nil { - return - } - for _, pageNo := range in.IndexPageNumbers { - err = bin.WriteUint32(dst, pageNo) - if err != nil { - return - } - } // free data pages - _, err = bin.WriteVarSize(dst, in.FrozenDataPagesCount) + err = writeFreePages(dst, in.FrozenDataPagesCount, in.DataPageNumbers) if err != nil { return } - _, err = bin.WriteVarSize(dst, len(in.DataPageNumbers)) - if err != nil { - return - } - for _, pageNo := range in.DataPageNumbers { - err = bin.WriteUint32(dst, pageNo) - if err != nil { - return - } - } // CRC32 bin.WriteUint32(file, hasher.Sum32()) @@ -196,8 +177,11 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { for range metricsQty { var ( metricID uint32 - metric = &ReplayMetric{ - buffer: make([]byte, minBufferSize), + buf = make([]byte, atree.DataPageSize) + // для простоти створюю буфер максимального розміру + metric = &ReplayMetric{ + buf: buf, + databuf: buf[:atree.DataPagePayloadSize], } ) metricID, err = bin.ReadUint32(src) @@ -211,40 +195,15 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { metrics[metricID] = metric } // index pages - frozenIndexPagesCount, err = bin.ReadVarSize(src) + frozenIndexPagesCount, indexPageNumbers, err = readFreePages(src) if err != nil { return } - indexPageNumbersCount, err := bin.ReadVarSize(src) - if err != nil { - return - } - for range indexPageNumbersCount { - var pageNo uint32 - pageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - indexPageNumbers = append(indexPageNumbers, pageNo) - } // data pages - frozenDataPagesCount, err = bin.ReadVarSize(src) + frozenDataPagesCount, dataPageNumbers, err = readFreePages(src) if err != nil { return } - dataPageNumbersCount, err := bin.ReadVarSize(src) - if err != nil { - return - } - for range dataPageNumbersCount { - var pageNo uint32 - pageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - dataPageNumbers = append(dataPageNumbers, pageNo) - } - // verify CRC32 expectedCRC, err := bin.ReadUint32(bufferedReader) if err != nil { @@ -267,3 +226,47 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { DataPageNumbers: dataPageNumbers, }, nil } + +// HELPERS + +func writeFreePages(dst io.Writer, frozenPagesCount int, pageNumbers []uint32) (err error) { + _, err = bin.WriteVarSize(dst, frozenPagesCount) + if err != nil { + return + } + _, err = bin.WriteVarSize(dst, len(pageNumbers)) + if err != nil { + return + } + for _, pageNo := range pageNumbers { + err = bin.WriteUint32(dst, pageNo) + if err != nil { + return + } + } + return +} + +func readFreePages(src io.Reader) (_ int, _ []uint32, err error) { + var ( + frozenPagesCount int + pageNumbers []uint32 + ) + frozenPagesCount, err = bin.ReadVarSize(src) + if err != nil { + return + } + pageNumbersCount, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range pageNumbersCount { + var pageNo uint32 + pageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + pageNumbers = append(pageNumbers, pageNo) + } + return frozenPagesCount, pageNumbers, nil +} diff --git a/database/wal_replayer.go b/database/wal_replayer.go new file mode 100644 index 0000000..bd4d80c --- /dev/null +++ b/database/wal_replayer.go @@ -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) + +} diff --git a/enc/cumdelta.go b/enc/cumdelta.go index 456f344..d3fde53 100644 --- a/enc/cumdelta.go +++ b/enc/cumdelta.go @@ -205,7 +205,7 @@ func (s *CumulativeDeltaCompressor) Size() int { if s.state == nil { return s.pos } else { - return len(s.state.Payload) + return bin.CountVarUint64(s.state.LastDelta) + hSize + len(s.state.Payload) } } @@ -220,18 +220,9 @@ func (s *CumulativeDeltaCompressor) Payload() []byte { func (s *CumulativeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if s.state == nil { - err = bin.WriteUint16(w, uint16(s.pos)) - if err != nil { - return - } _, err = w.Write(s.buf[:s.pos]) return } else { - size := bin.CountVarUint64(s.state.LastDelta) + hSize + len(s.state.Payload) - err = bin.WriteUint16(w, uint16(size)) - if err != nil { - return - } _, err = w.Write(s.state.Payload) if err != nil { return @@ -240,7 +231,7 @@ func (s *CumulativeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if err != nil { return } - w.Write([]byte{ + _, err = w.Write([]byte{ s.state.H, }) return diff --git a/enc/insdelta.go b/enc/insdelta.go index 8814a48..991ad33 100644 --- a/enc/insdelta.go +++ b/enc/insdelta.go @@ -190,7 +190,7 @@ func (s *InstantDeltaCompressor) Size() int { if s.state == nil { return s.pos } else { - return len(s.state.Payload) + return bin.CountVarInt64(s.state.LastDelta) + hSize + len(s.state.Payload) } } @@ -205,18 +205,9 @@ func (s *InstantDeltaCompressor) Payload() []byte { func (s *InstantDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if s.state == nil { - err = bin.WriteUint16(w, uint16(s.pos)) - if err != nil { - return - } _, err = w.Write(s.buf[:s.pos]) return } else { - size := bin.CountVarInt64(s.state.LastDelta) + hSize + len(s.state.Payload) - err = bin.WriteUint16(w, uint16(size)) - if err != nil { - return - } _, err = w.Write(s.state.Payload) if err != nil { return @@ -225,7 +216,7 @@ func (s *InstantDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if err != nil { return } - w.Write([]byte{ + _, err = w.Write([]byte{ s.state.H, }) return diff --git a/enc/time_delta.go b/enc/time_delta.go index bc942d9..adc5f2d 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -205,7 +205,7 @@ func (s *TimeDeltaCompressor) Size() int { if s.state == nil { return len(s.buf) - s.pos } else { - return len(s.state.Payload) + return bin.CountVarUint64(uint64(s.state.LastDelta)) + hSize + len(s.state.Payload) } } @@ -223,18 +223,9 @@ func (s *TimeDeltaCompressor) Payload() []byte { func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if s.state == nil { - err = bin.WriteUint16(w, uint16(s.pos)) - if err != nil { - return - } _, err = w.Write(s.buf[:s.pos]) return } else { - size := bin.CountVarUint64(uint64(s.state.LastDelta)) + hSize + len(s.state.Payload) - err = bin.WriteUint16(w, uint16(size)) - if err != nil { - return - } _, err = w.Write(s.state.Payload) if err != nil { return @@ -243,7 +234,7 @@ func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if err != nil { return } - w.Write([]byte{ + _, err = w.Write([]byte{ s.state.H, }) return diff --git a/txlog/txlog.go b/txlog/txlog.go index 014818d..03be1fa 100644 --- a/txlog/txlog.go +++ b/txlog/txlog.go @@ -40,8 +40,9 @@ func (s *AddedMetric) Parse(src io.Reader) (err error) { } type DeletedMetric struct { - MetricID uint32 - FreePageNumbers []uint32 + MetricID uint32 + FreeIndexPages []uint32 + FreeDataPages []uint32 } func (s DeletedMetric) Pack(w io.Writer) { @@ -51,10 +52,11 @@ func (s DeletedMetric) Pack(w io.Writer) { } bin.PutUint32(arr[1:], s.MetricID) w.Write(arr) - bin.WriteVarSize(w, len(s.FreePageNumbers)) - for _, pageNo := range s.FreePageNumbers { - bin.WriteUint32(w, pageNo) - } + // fix + // bin.WriteVarSize(w, len(s.FreePageNumbers)) + // for _, pageNo := range s.FreePageNumbers { + // bin.WriteUint32(w, pageNo) + // } } func (s *DeletedMetric) Parse(src io.Reader) (err error) { @@ -63,18 +65,19 @@ func (s *DeletedMetric) Parse(src io.Reader) (err error) { return } // free pages - qty, err := bin.ReadVarSize(src) - if err != nil { - return - } - for range qty { - var pageNo uint32 - pageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - s.FreePageNumbers = append(s.FreePageNumbers, pageNo) - } + // fix + // qty, err := bin.ReadVarSize(src) + // if err != nil { + // return + // } + // for range qty { + // var pageNo uint32 + // pageNo, err = bin.ReadUint32(src) + // if err != nil { + // return + // } + // //s.FreePageNumbers = append(s.FreePageNumbers, pageNo) + // } return nil } diff --git a/txlog/wal_records.go b/txlog/wal_records.go index 29e289c..abf491d 100644 --- a/txlog/wal_records.go +++ b/txlog/wal_records.go @@ -139,7 +139,6 @@ func (s WALRecordAppendMeasures) WriteTo(w *bytes.Buffer) { bin.WriteVarSize(w, len(level.TailRecords)/indexRecordSize) w.Write(level.TailRecords) } - return } func (s *WALRecordAppendMeasures) Parse(r *bytes.Buffer) (err error) {