diff --git a/database/database.go b/database/database.go index 6e3abb8..7e404cb 100644 --- a/database/database.go +++ b/database/database.go @@ -103,14 +103,11 @@ func New(opt Options) (_ *Database, err error) { } - s.worker = worker.NewWorker(worker.WorkerOptions{ + s.worker = worker.New(worker.Options{ Inbox: s.workerInbox, StorageInbox: storageInbox, //Metrics: make(map[uint32]*_metric), FIX - Dir: opt.Dir, - // SaveToStorage: func(in any) { - // fmt.Println("save to storage") - // }, + Dir: opt.Dir, ExitCh: opt.ExitCh, WaitGroup: opt.WaitGroup, }) diff --git a/freelist/freelist.go b/freelist/freelist.go index 173e1bd..5b77993 100644 --- a/freelist/freelist.go +++ b/freelist/freelist.go @@ -55,22 +55,10 @@ func New(opt Options) (*FreeList, error) { if err != nil { return nil, err } - s.deltaFile, err = os.OpenFile(opt.DeltaFilePath, os.O_RDWR|os.O_CREATE, 0666) + s.deltaFile, err = os.OpenFile(opt.DeltaFilePath, os.O_RDWR|os.O_CREATE|os.O_TRUNC, 0666) if err != nil { return nil, err } - // info, err := s.baseFile.Stat() - // if err != nil { - // return nil, err - // } - // fileSize := info.Size() - // if fileSize > 0 { - // if (fileSize % int64(s.pageSize)) > 0 { - // return nil, fmt.Errorf("file size %d is not a multiple of page size %d", - // fileSize, s.pageSize) - // } - // s.basePagesCount = int(fileSize) / s.pageSize - // } return s, nil } diff --git a/qb.go b/qb.go index 934f77a..37955b7 100644 --- a/qb.go +++ b/qb.go @@ -173,6 +173,14 @@ func GetIndexFreeListFilePath(dir string, databaseName string) string { return filepath.Join(dir, fmt.Sprintf("%s.index_free", databaseName)) } +func GetIndexFreeListDeltaFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.index_free_delta", databaseName)) +} + func GetDataFreeListFilePath(dir string, databaseName string) string { return filepath.Join(dir, fmt.Sprintf("%s.data_free", databaseName)) } + +func GetDataFreeListDeltaFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.data_free_delta", databaseName)) +} diff --git a/recovery/recovery.go b/recovery/recovery.go index 5051335..502e74f 100644 --- a/recovery/recovery.go +++ b/recovery/recovery.go @@ -10,6 +10,7 @@ import ( "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/freelist" "gordenko.dev/dima/qb/storage" + "gordenko.dev/dima/qb/worker" ) // Result зберігає результати для обох типів файлів @@ -110,7 +111,6 @@ func resolveRecoveryFiles(dir string) (_ RecoveryFiles, err error) { } type RecoveryReport struct { - //Metrics [] SnapshotNumber int WAL string IndexFreeList *freelist.FreeList @@ -118,32 +118,38 @@ type RecoveryReport struct { } func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { - files, err := resolveRecoveryFiles(dir) + recoveryFiles, err := resolveRecoveryFiles(dir) if err != nil { return } var ( - freeIndexPages []uint32 - freeDataPages []uint32 + frozenIndexPageCount int + frozenDataPageCount int + freeIndexPages []uint32 + freeDataPages []uint32 + metrics map[uint32]*storage.ReplayMetric ) - if files.Snapshot != "" { - var out storage.ReadSnapshotOut - out, err = storage.ReadSnapshot(files.Snapshot) + if recoveryFiles.Snapshot != "" { + var snapshot worker.ReadSnapshotOut + snapshot, err = worker.ReadSnapshot(recoveryFiles.Snapshot) if err != nil { return } // - freeIndexPages = out.IndexPageNumbers - freeDataPages = out.DataPageNumbers + frozenIndexPageCount = snapshot.FrozenIndexPagesCount + frozenDataPageCount = snapshot.FrozenDataPagesCount + freeIndexPages = snapshot.IndexPageNumbers + freeDataPages = snapshot.DataPageNumbers + metrics = snapshot.Metrics } - if files.WAL != "" { + if recoveryFiles.WAL != "" { var ( walReader *storage.WALReader walReplayer *storage.WALReplayer ) walReader, err = storage.NewWALReader(storage.WALReaderOptions{ - //FileName: JoinWALFileName(s.dir, files.WAL), + FileName: recoveryFiles.WAL, BufferSize: 8 * 1024 * 1024, }) if err != nil { @@ -152,7 +158,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { walReplayer, err = storage.NewWALReplayer(storage.WALReplayerOptions{ WALReader: walReader, - Metrics: snapshot.Metrics, + Metrics: metrics, FreeIndexPages: freeIndexPages, FreeDataPages: freeDataPages, }) @@ -170,7 +176,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { if err != nil { return } - err = storage.WriteDataPages(dataFile, walReplayer.DataPages()) + err = storage.WriteDataPages(dataFile, walReplayer.DataPagesToRewrite()) if err != nil { return } @@ -183,7 +189,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { if err != nil { return } - err = storage.WriteIndexPages(indexFile, walReplayer.IndexPages()) + err = storage.WriteIndexPages(indexFile, walReplayer.IndexPagesToRewrite()) if err != nil { return } @@ -191,31 +197,31 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { if err != nil { return } + + freeIndexPages = walReplayer.FreeIndexPages() + freeDataPages = walReplayer.FreeDataPages() } indexFreeList, err := freelist.New(freelist.Options{ - PageSize: 2048, - BaseFilePath: filepath.Join(dir, databaseName+".free_index"), - DeltaFilePath: filepath.Join(dir, databaseName+".free_index_delta"), + PageSize: 2048, + BaseFilePath: qb.GetIndexFreeListFilePath(dir, databaseName), + DeltaFilePath: qb.GetIndexFreeListDeltaFilePath(dir, databaseName), + FrozenPageCount: frozenIndexPageCount, }) + indexFreeList.AddPageNumbers(freeIndexPages) - // FIX free pages sync dataFreeList, err := freelist.New(freelist.Options{ - PageSize: 2048, - BaseFilePath: filepath.Join(dir, databaseName+".free_data"), - DeltaFilePath: filepath.Join(dir, databaseName+".free_data_delta"), + PageSize: 2048, + BaseFilePath: qb.GetDataFreeListFilePath(dir, databaseName), + DeltaFilePath: qb.GetDataFreeListDeltaFilePath(dir, databaseName), + FrozenPageCount: frozenDataPageCount, }) + dataFreeList.AddPageNumbers(freeDataPages) return RecoveryReport{ - SnapshotNumber: files.SnapshotNumber, - WAL: files.WAL, + SnapshotNumber: recoveryFiles.SnapshotNumber, + WAL: recoveryFiles.WAL, IndexFreeList: indexFreeList, DataFreeList: dataFreeList, }, nil } - -//func (s *Database) relayMetricsToMetrics(replayMetrics map[uint32]*ReplayMetric) { -//for metricID, x := range replayMetrics { -//s.metrics[metricID] = x.ToMetric() FIX -//} -//} diff --git a/storage/replay_metric.go b/storage/replay_metric.go index 8ec3c5b..768064e 100644 --- a/storage/replay_metric.go +++ b/storage/replay_metric.go @@ -2,101 +2,50 @@ package storage import ( "fmt" - "io" - bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" ) type ReplayMetric struct { - metricType qb.MetricType - fracDigits byte - lastPageNo uint32 - buf []byte - databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер) - tSize int - vSize int - indexLevelTails []IndexLevelTail // root - last element + MetricType qb.MetricType + FracDigits byte + LastPageNo uint32 + Buf []byte + TimestampsSize int + ValuesSize int + IndexLevelTails []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, 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, IndexLevelTail{ - Buffer: buf, - RecordsCount: count, - }) - } - return +// return occupied after write +func writeTimestampsWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int { + // timestamps writes right-to-left + pos := DataPagePayloadSize - occupied + // timestamps rewind offset moves pos left-to-right + pos += rewindOffset + pos -= len(payload) + copy(buf[pos:], payload) + return DataPagePayloadSize - pos } -// func appendNewRecordsToIndexTail(level IndexLevelTail, records []byte) { - -// } - -// func countReusedIndexPages(changedIndexLevels []ChangedIndexLevel) (count int) { - -// return -// } +func writeValuesWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int { + // values writes left-to-right + pos := occupied + // values rewind offset moves pos right-to-left + pos -= rewindOffset + copy(buf[pos:], payload) + return pos + len(payload) +} func (s *ReplayMetric) MeasuresAppend(rec MeasuresAppendRecord) { - timestampsSize := rec.TimestampsRewindOffset - len(rec.Timestamps) - // копіюю нові дані із WAL з урахуванням offset-ів - copy(s.databuf[rec.ValuesRewindOffset:], rec.Values) - copy(s.databuf[len(s.databuf)-timestampsSize:], rec.Timestamps) - // - s.tSize = timestampsSize - s.vSize = rec.ValuesRewindOffset - len(rec.Values) + s.TimestampsSize = writeTimestampsWithRewind(s.Buf, s.TimestampsSize, + rec.Timestamps, rec.TimestampsRewindOffset) + s.ValuesSize = writeValuesWithRewind(s.Buf, s.ValuesSize, + rec.Values, rec.ValuesRewindOffset) } type AppendMeasuresResult struct { - IndexPages []PageToWrite - DataPages []PageToWrite + IndexPagesToRewrite []PageToWrite + DataPagesToRewrite []PageToWrite ReusedIndexPagesCount int ReusedDataPagesCount int } @@ -117,8 +66,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, // 1. набиваю сторінки для перезапису в index файлі if isLastPacket { if head != nil { - indexPages = append(indexPages, - composeHeadIndexPage(levelIdx, s.indexLevelTails[levelIdx], head)) + indexPages = append(indexPages, s.composeHeadIndexPage(levelIdx, head)) } for _, x := range change.IndexPages { indexPages = append(indexPages, PageToWrite{ @@ -129,8 +77,8 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, } // 2. Вношу зміни в поточні індексні хвости - if levelIdx < len(s.indexLevelTails) { - level := s.indexLevelTails[levelIdx] + if levelIdx < len(s.IndexLevelTails) { + level := s.IndexLevelTails[levelIdx] if head != nil { // попередній хвіст індексного рівня перетворився на head сторінку, // отже TailRecords - це новий хвіст @@ -142,13 +90,13 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, copy(level.Buffer[pos:], change.TailRecords) level.RecordsCount += newRecordsCount } - s.indexLevelTails[levelIdx] = level + s.IndexLevelTails[levelIdx] = level } else { // додаю новий індексний рівень buf := make([]byte, IndexPageSize) copy(buf, change.TailRecords) // - s.indexLevelTails = append(s.indexLevelTails, IndexLevelTail{ + s.IndexLevelTails = append(s.IndexLevelTails, IndexLevelTail{ Buffer: buf, RecordsCount: newRecordsCount, }) @@ -169,7 +117,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, // 4. набиваю сторінки для перезапису в data файлі if isLastPacket { - dataPages = append(dataPages, composeHeadDataPage(s.databuf, rec.DataPageTail)) + dataPages = append(dataPages, s.composeHeadDataPage(rec.DataPageTail)) for _, x := range rec.DataPages { dataPages = append(dataPages, PageToWrite{ @@ -190,22 +138,22 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, } // HeadDataPage є полюбому, оскільки це транзакція із мінімум однією заповненою // data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня. - copy(s.databuf, rec.TailValues) - pos := len(s.databuf) - len(rec.TailTimestamps) - copy(s.databuf[pos:], rec.TailTimestamps) + copy(s.Buf, rec.TailValues) + s.ValuesSize = len(rec.TailValues) - s.tSize = len(rec.TailTimestamps) - s.vSize = len(rec.TailValues) + pos := DataPagePayloadSize - len(rec.TailTimestamps) + copy(s.Buf[pos:], rec.TailTimestamps) + s.TimestampsSize = len(rec.TailTimestamps) if len(rec.DataPages) == 0 { - s.lastPageNo = rec.DataPageTail.PageNo + s.LastPageNo = rec.DataPageTail.PageNo } else { - s.lastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo + s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo } return AppendMeasuresResult{ - IndexPages: indexPages, - DataPages: dataPages, + IndexPagesToRewrite: indexPages, + DataPagesToRewrite: dataPages, ReusedIndexPagesCount: reusedIndexPages, ReusedDataPagesCount: reusedDataPages, } @@ -213,7 +161,8 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, // HELPERS -func composeHeadIndexPage(levelIdx int, level IndexLevelTail, head *IndexPageTail) PageToWrite { +func (s *ReplayMetric) composeHeadIndexPage(levelIdx int, head *IndexPageTail) PageToWrite { + level := s.IndexLevelTails[levelIdx] // розраховую pos, з якого буду дописувати хвіст pos := level.RecordsCount * IndexRecordSize // створюю копію сторінки @@ -238,23 +187,21 @@ func composeHeadIndexPage(levelIdx int, level IndexLevelTail, head *IndexPageTai } } -func composeHeadDataPage(databuf []byte, head DataPageTail) PageToWrite { - // створюю копію сторінки - var ( - page = make([]byte, DataPageSize) - timestampsSize = head.TimestampsRewindOffset - len(head.Timestamps) - ) - // копіюю поточні дані - copy(page, databuf) - // копіюю нові дані із WAL з урахуванням offset-ів - copy(page[head.ValuesRewindOffset:], head.Values) - copy(page[len(page)-timestampsSize:], head.Timestamps) +func (s *ReplayMetric) composeHeadDataPage(head DataPageTail) PageToWrite { + // copy page + page := make([]byte, DataPageSize) + copy(page, s.Buf) + // append data + timestampsSize := writeTimestampsWithRewind(s.Buf, s.TimestampsSize, + head.Timestamps, head.TimestampsRewindOffset) + valuesSize := writeValuesWithRewind(s.Buf, s.ValuesSize, + head.Values, head.ValuesRewindOffset) // запечатати сторінку calculatedCRC := SealDataPage(SealDataPageIn{ Content: page, PrevPageNo: head.PrevPageNo, TimestampsSize: timestampsSize, - ValuesSize: head.ValuesRewindOffset + len(head.Values), + ValuesSize: valuesSize, }) // перевірка CRC if calculatedCRC != head.CRC32 { @@ -267,19 +214,3 @@ func composeHeadDataPage(databuf []byte, head DataPageTail) PageToWrite { Content: page, } } - -// func (s *ReplayMetric) ToMetric() *worker.metric { - -// databuf := s.buf[:storage.DataPagePayloadSize] -// values := enc.NewValueDeltaCompressor(s.metricType, s.fracDigits, databuf, s.vSize) -// return &_metric{ -// metricType: s.metricType, -// fracDigits: s.fracDigits, -// lastPageNo: s.lastPageNo, -// lastValue: values.LastValue(), -// buffer: s.buf, -// timestamps: enc.NewTimeDeltaCompressor(databuf, s.tSize), -// values: values, -// indexLevelTails: s.indexLevelTails, -// } -// } diff --git a/storage/storage.go b/storage/storage.go index d5356a3..6dbeeea 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -72,7 +72,7 @@ type MeasuresAppendWithGrow struct { IndexLevelTails []IndexLevelTail ResultCode byte WrittenCount int - ResultCh chan struct{} + ResultCh chan AppendMeasuresResult } type MeasuresAppend struct { @@ -83,7 +83,7 @@ type MeasuresAppend struct { Values []byte ResultCode byte WrittenCount int - ResultCh chan struct{} + ResultCh chan AppendMeasuresResult } type MeasuresDelete struct { @@ -100,7 +100,7 @@ type MeasuresAppendCommited struct { MetricID uint32 ResultCode byte WrittenCount int - ResultCh chan struct{} + ResultCh chan AppendMeasuresResult } type MeasuresAppendWithGrowCommited struct { @@ -109,7 +109,7 @@ type MeasuresAppendWithGrowCommited struct { Index []IndexLevelTail ResultCode byte WrittenCount int - ResultCh chan struct{} + ResultCh chan AppendMeasuresResult } type MeasuresDeleteCommited struct { diff --git a/storage/storage_test.go b/storage/storage_test.go index 20c92b5..6951072 100644 --- a/storage/storage_test.go +++ b/storage/storage_test.go @@ -868,82 +868,49 @@ func TestWritePreparer(t *testing.T) { } } -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) +func TestWriteTimestampsWithRewind(t *testing.T) { + DataPagePayloadSize = 8 + var ( + payload = []byte{ + 0xaa, 0xbb, 0xcc, } - err = cmpMetric(replay, origin) - if err != nil { - t.Fatalf("metric %d: %s", metricID, err) + before = []byte{ + 0x00, 0x00, 0x00, 0x00, 0x04, 0x03, 0x02, 0x01, } + occupied = 4 + after = []byte{ + 0x00, 0x00, 0x00, 0xaa, 0xbb, 0xcc, 0x02, 0x01, + } + occupiedAfter = 5 + ) + occupied = writeTimestampsWithRewind(before, occupied, payload, 2) + if occupied != occupiedAfter { + t.Fatalf("got occupied %d are not equal expected %d", occupied, occupiedAfter) + } + if !bytes.Equal(before, after) { + t.Fatalf("got buf % x are not equal expected % x", before, after) } } -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) +func TestWriteValuesWithRewind(t *testing.T) { + var ( + payload = []byte{ + 0xaa, 0xbb, 0xcc, + } + before = []byte{ + 0x01, 0x02, 0x03, 0x04, 0x00, 0x00, 0x00, 0x00, + } + occupied = 4 + after = []byte{ + 0x01, 0x02, 0xaa, 0xbb, 0xcc, 0x00, 0x00, 0x00, + } + occupiedAfter = 5 + ) + occupied = writeValuesWithRewind(before, occupied, payload, 2) + if occupied != occupiedAfter { + t.Fatalf("got occupied %d are not equal expected %d", occupied, occupiedAfter) } - if replay.fracDigits != origin.fracDigits { - return fmt.Errorf("fracDigits: got %d are not equal expected %d", - replay.fracDigits, origin.fracDigits) + if !bytes.Equal(before, after) { + t.Fatalf("got buf % x are not equal expected % x", before, after) } - 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/storage/wal_reader.go b/storage/wal_reader.go index 3c67172..99866ee 100644 --- a/storage/wal_reader.go +++ b/storage/wal_reader.go @@ -173,14 +173,14 @@ func (s *WALReader) parseRecords(body []byte) ([]any, error) { // HELPERS -func (s *WALReader) Seek(offset int64) error { - ret, err := s.file.Seek(offset, 0) - if err != nil { - return err - } +// func (s *WALReader) Seek(offset int64) error { +// ret, err := s.file.Seek(offset, 0) +// if err != nil { +// return err +// } - if ret != offset { - return fmt.Errorf("ret %d != offset %d", ret, offset) - } - return nil -} +// if ret != offset { +// return fmt.Errorf("ret %d != offset %d", ret, offset) +// } +// return nil +// } diff --git a/storage/wal_replayer.go b/storage/wal_replayer.go index 954b91a..656013b 100644 --- a/storage/wal_replayer.go +++ b/storage/wal_replayer.go @@ -36,14 +36,22 @@ func NewWALReplayer(opt WALReplayerOptions) (*WALReplayer, error) { }, nil } -func (s *WALReplayer) IndexPages() []PageToWrite { +func (s *WALReplayer) IndexPagesToRewrite() []PageToWrite { return s.indexPages } -func (s *WALReplayer) DataPages() []PageToWrite { +func (s *WALReplayer) DataPagesToRewrite() []PageToWrite { return s.dataPages } +func (s *WALReplayer) FreeIndexPages() []uint32 { + return s.freeIndexPages +} + +func (s *WALReplayer) FreeDataPages() []uint32 { + return s.freeDataPages +} + func (s *WALReplayer) Replay() error { for { records, isLastPacket, err := s.walReader.NextPacket() @@ -101,10 +109,9 @@ func (s *WALReplayer) onMetricAdd(rec MetricAddRecord) (err error) { buf = make([]byte, DataPageSize) ) s.metrics[rec.MetricID] = &ReplayMetric{ - metricType: rec.MetricType, - fracDigits: byte(rec.FracDigits), - buf: buf, - databuf: buf[:DataPagePayloadSize], + MetricType: rec.MetricType, + FracDigits: rec.FracDigits, + Buf: buf, } return } @@ -137,39 +144,11 @@ func (s *WALReplayer) onMeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, result := metric.MeasuresAppendWithGrow(rec, isLastPacket) // if isLastPacket { - s.indexPages = append(s.indexPages, result.IndexPages...) - s.dataPages = append(s.dataPages, result.DataPages...) + s.indexPages = append(s.indexPages, result.IndexPagesToRewrite...) + s.dataPages = append(s.dataPages, result.DataPagesToRewrite...) } // s.freeIndexPages = s.freeIndexPages[:len(s.freeIndexPages)-result.ReusedIndexPagesCount] s.freeDataPages = s.freeDataPages[:len(s.freeDataPages)-result.ReusedDataPagesCount] return } - -//func (s *WALReplayer) collectPagesToWrite(metric *ReplayMetric, rec WALRecordAppendMeasures) { -// metric.ApplyHeadDataPage(rec.CompletedDataPage) - -// s.dataPages = append(s.dataPages, PageToWrite{ -// PageNo: p.PageNo, -// Content: metric.buf, -// }) - -// for _, x := range rec.DataPages { -// s.dataPages = append(s.dataPages, PageToWrite{ -// PageNo: x.PageNo, -// Content: x.Content, -// }) -// } -//metric.ApplyHeadIndexPages(rec.ChangedIndexLevels) - -//} - -// func (s *WALReplayer) replayPacket(records []any, isLastPacket bool) { -// for _, untyped := range records { -// switch rec := untyped.(type) { -// case WALRecordAppendMeasures: -// s.onWALRecordAppendMeasures(rec, isLastPacket) -// } -// } -// return -// } diff --git a/worker/metric.go b/worker/metric.go index a33776b..679255a 100644 --- a/worker/metric.go +++ b/worker/metric.go @@ -6,6 +6,7 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/inbox" "gordenko.dev/dima/qb/storage" ) @@ -36,6 +37,14 @@ type Metric struct { capturedState *CapturedState } +func (s *Metric) MetricType() qb.MetricType { + return s.metricType +} + +func (s *Metric) FracDigits() byte { + return s.fracDigits +} + func (s *Metric) LastValue() float64 { if s.capturedState != nil { return s.capturedState.LastValue @@ -67,7 +76,7 @@ func (s *Metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { } // func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) { -func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, saveToStorage func(any)) { +func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox *inbox.Inbox) { if s.capturedState != nil { s.WaitQueue = append(s.WaitQueue, req) } @@ -171,7 +180,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, saveToSt if len(pages) > 0 { // пишу в storage довгим шляхом через redo файл і запис в data файл - saveToStorage(storage.MeasuresAppendWithGrow{ + storageInbox.Push(storage.MeasuresAppendWithGrow{ MetricID: req.MetricID, LastPageNo: s.lastPageNo, TimestampsRewindOffset: timestampsRewindOffset, @@ -188,7 +197,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, saveToSt }) } else { // короткий шлях - запис лише в storage - saveToStorage(storage.MeasuresAppend{ + storageInbox.Push(storage.MeasuresAppend{ MetricID: req.MetricID, TimestampsRewindOffset: timestampsRewindOffset, ValuesRewindOffset: valuesRewindOffset, diff --git a/storage/snapshot.go b/worker/snapshot.go similarity index 78% rename from storage/snapshot.go rename to worker/snapshot.go index ddaa3ee..3b906a3 100644 --- a/storage/snapshot.go +++ b/worker/snapshot.go @@ -1,4 +1,4 @@ -package storage +package worker import ( "bufio" @@ -8,6 +8,8 @@ import ( "path/filepath" bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/util" ) @@ -28,17 +30,13 @@ dataFreeList - Nb CRC32 - 4b */ -const ReadBufferSize = 8 * 1024 * 1024 // 8mb - -type WritableMetric interface { - WriteTo(io.Writer) error -} +const readBufferSize = 4 * 1024 * 1024 // 8mb type WriteSnapshotIn struct { SnapshotNumber int Dir string WriteBufferSize int - Metrics map[uint32]WritableMetric + Metrics map[uint32]*Metric FrozenIndexPagesCount int IndexPageNumbers []uint32 FrozenDataPagesCount int @@ -122,7 +120,7 @@ func WriteSnapshot(in WriteSnapshotIn) (err error) { } type ReadSnapshotOut struct { - Metrics map[uint32]*ReplayMetric + Metrics map[uint32]*storage.ReplayMetric FrozenIndexPagesCount int IndexPageNumbers []uint32 FrozenDataPagesCount int @@ -131,7 +129,7 @@ type ReadSnapshotOut struct { func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) { var ( - metrics = make(map[uint32]*ReplayMetric) + metrics = make(map[uint32]*storage.ReplayMetric) frozenIndexPagesCount int indexPageNumbers []uint32 frozenDataPagesCount int @@ -180,22 +178,70 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) { for range metricsQty { var ( metricID uint32 - buf = make([]byte, DataPageSize) - // для простоти створюю буфер максимального розміру - metric = &ReplayMetric{ - buf: buf, - databuf: buf[:DataPagePayloadSize], + buf = make([]byte, storage.DataPageSize) + m = &storage.ReplayMetric{ + Buf: buf, } + metricType byte + levelsCount int ) metricID, err = bin.ReadUint32(src) if err != nil { return } - err = metric.ReadFrom(src) + metricType, err = bin.ReadByte(src) if err != nil { return } - metrics[metricID] = metric + m.MetricType = qb.MetricType(metricType) + m.FracDigits, err = bin.ReadByte(src) + if err != nil { + return + } + m.LastPageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + m.TimestampsSize, err = bin.ReadUint16AsInt(src) + if err != nil { + return + } + m.ValuesSize, err = bin.ReadUint16AsInt(src) + if err != nil { + return + } + err = bin.ReadNInto(src, buf[len(buf)-m.TimestampsSize:]) + if err != nil { + return + } + err = bin.ReadNInto(src, buf[:m.ValuesSize]) + if err != nil { + return + } + // + levelsCount, err = bin.ReadVarSize(src) + if err != nil { + return + } + for range levelsCount { + var ( + buf = make([]byte, storage.IndexPageSize) + count int + ) + count, err = bin.ReadVarSize(src) + if err != nil { + return + } + err = bin.ReadNInto(src, buf[:count*storage.IndexRecordSize]) + if err != nil { + return + } + m.IndexLevelTails = append(m.IndexLevelTails, storage.IndexLevelTail{ + Buffer: buf, + RecordsCount: count, + }) + } + metrics[metricID] = m } // index pages frozenIndexPagesCount, indexPageNumbers, err = readFreePages(src) diff --git a/worker/snapshot_test.go b/worker/snapshot_test.go new file mode 100644 index 0000000..82d97ae --- /dev/null +++ b/worker/snapshot_test.go @@ -0,0 +1,90 @@ +package worker + +import ( + "fmt" + "slices" + "testing" + + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/storage" +) + +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 *storage.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 !reflect.DeepEqual() { + // return fmt.Errorf("lastPageNo: got %d are not equal expected %d", + // replay.LastPageNo, origin.lastPageNo) + // } + return nil +} diff --git a/worker/worker.go b/worker/worker.go index 3ae5d6a..99efb43 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -38,37 +38,51 @@ const ( ) type Worker struct { - //mutex sync.Mutex inbox *inbox.Inbox storageInbox *inbox.Inbox dir string snapshotNumber int tmp []byte metrics map[uint32]*Metric - saveToStorage func(any) exitCh chan struct{} waitGroup *sync.WaitGroup } -type WorkerOptions struct { - Inbox *inbox.Inbox - StorageInbox *inbox.Inbox - Metrics map[uint32]*Metric - Dir string - ExitCh chan struct{} - WaitGroup *sync.WaitGroup +type Options struct { + Inbox *inbox.Inbox + StorageInbox *inbox.Inbox + ReplayMetrics map[uint32]*storage.ReplayMetric + Dir string + ExitCh chan struct{} + WaitGroup *sync.WaitGroup } -func NewWorker(opt WorkerOptions) *Worker { - return &Worker{ +func New(opt Options) *Worker { + s := &Worker{ inbox: opt.Inbox, storageInbox: opt.StorageInbox, dir: opt.Dir, tmp: make([]byte, 26), - metrics: opt.Metrics, + metrics: make(map[uint32]*Metric), exitCh: opt.ExitCh, waitGroup: opt.WaitGroup, } + + for metricID, x := range opt.ReplayMetrics { + databuf := x.Buf[:storage.DataPagePayloadSize] + values := enc.NewValueDeltaCompressor(x.MetricType, x.FracDigits, databuf, x.ValuesSize) + s.metrics[metricID] = &Metric{ + metricType: x.MetricType, + fracDigits: x.FracDigits, + lastPageNo: x.LastPageNo, + lastValue: values.LastValue(), + buffer: x.Buf, + timestamps: enc.NewTimeDeltaCompressor(databuf, x.TimestampsSize), + values: values, + indexLevelTails: x.IndexLevelTails, + } + } + return s } func (s *Worker) ReleaseRLock(metricID uint32) { @@ -182,7 +196,7 @@ func (s *Worker) processMetricQueue(metricID uint32, metric *Metric, tmp []byte) s.GetMetric(req) case AppendMeasuresReq: - metric.StartAppendMeasures(req, tmp, s.saveToStorage) + metric.StartAppendMeasures(req, tmp, s.storageInbox) case DeleteMetricReq: s.startDeleteMetric(metric, req) @@ -259,8 +273,8 @@ func (s *Worker) GetMetric(req GetMetricReq) { if ok { req.ResultCh <- GetMetricResult{ ResultCode: Succeed, - MetricType: metric.metricType, - FracDigits: metric.fracDigits, + MetricType: metric.MetricType(), + FracDigits: metric.FracDigits(), } } else { req.ResultCh <- GetMetricResult{ @@ -404,7 +418,7 @@ func (s *Worker) AppendMeasures(req AppendMeasuresReq) { metric.WaitQueue = append(metric.WaitQueue, req) return } - metric.StartAppendMeasures(req, s.tmp, s.saveToStorage) + metric.StartAppendMeasures(req, s.tmp, s.storageInbox) } type RangeScanResult struct { @@ -431,7 +445,7 @@ func (s *Worker) RangeScan(req RangeScanReq) { } return } - if metric.metricType != req.MetricType { + if metric.MetricType() != req.MetricType { req.ResultCh <- RangeScanResult{ ResultCode: WrongMetricType, } @@ -467,7 +481,7 @@ func (s *Worker) tryFullScan(req FullScanReq) { } return } - if metric.metricType != req.MetricType { + if metric.MetricType() != req.MetricType { req.ResultCh <- FullScanResult{ ResultCode: WrongMetricType, } @@ -525,7 +539,7 @@ func (s *Worker) applyCommits(req storage.Changes) { if req.SnapshotNumberCh != nil { s.snapshotNumber++ - err := storage.WriteSnapshot(storage.WriteSnapshotIn{ + err := WriteSnapshot(WriteSnapshotIn{ SnapshotNumber: s.snapshotNumber, Dir: s.dir, WriteBufferSize: 4 * 1024 * 1024, // 1mb