diff --git a/atree/io.go b/atree/io.go index c9b5f1f..08ca562 100644 --- a/atree/io.go +++ b/atree/io.go @@ -366,7 +366,7 @@ func openFile(fileName string, pageSize int) (_ *os.File, _ uint32, err error) { func (s *Atree) verifyCRC(data []byte, pageSize int) error { var ( pos = pageSize - 4 - calculatedCRC = util.CalcChecksum(data[:pos]) + calculatedCRC = util.CalculateCRC32(data[:pos]) storedCRC, _ = bin.GetUint32(data[pos:]) ) if calculatedCRC != storedCRC { diff --git a/database/api.go b/database/api.go index 7f9c247..c36aa6b 100644 --- a/database/api.go +++ b/database/api.go @@ -500,22 +500,22 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { case DeleteFromAtreeNotNeeded: // регистрирую удаление в TransactionLog - s.storage.Append(storage.DeletedMeasures{ + s.storage.Append(storage.MeasuresDeleteRecord{ MetricID: req.MetricID, }) //<-waitCh case DeleteFromAtreeRequired: // собираю номера всех data и index страниц метрики (типа запись REDO лога). - pageNumbers, err := s.atree.GetAllPages(req.MetricID) - if err != nil { - qb.Abort(qb.FailedAtreeRequest, err) - } - // регистрирую удаление в TransactionLog - s.storage.Append(storage.DeletedMeasures{ - MetricID: req.MetricID, - FreePageNumbers: pageNumbers, - }) + // pageNumbers, err := s.atree.GetAllPages(req.MetricID) + // if err != nil { + // qb.Abort(qb.FailedAtreeRequest, err) + // } + // // регистрирую удаление в TransactionLog + // s.storage.Append(storage.DeletedMeasures{ + // MetricID: req.MetricID, + // FreePageNumbers: pageNumbers, + // }) //<-waitCh case NoMetric: diff --git a/database/database_test.go b/database/database_test.go index bc25a75..5ef4f5c 100644 --- a/database/database_test.go +++ b/database/database_test.go @@ -98,10 +98,10 @@ func TestComposeHeadIndexPage(t *testing.T) { }, RecordsCount: 2, } - head = &storage.HeadIndexPage{ - PageNo: 100, - Checksum: 12345, - Records: []byte{}, + head = &storage.IndexPageTail{ + PageNo: 100, + CRC32: 12345, + Records: []byte{}, } ) page := composeHeadIndexPage(levelIdx, level, head) diff --git a/database/metric.go b/database/metric.go index 5f9a4c2..f73cd12 100644 --- a/database/metric.go +++ b/database/metric.go @@ -16,6 +16,11 @@ var ( indexRecordSize = 8 ) +type CapturedState struct { + LastTimestamp uint32 + LastValue float64 +} + type _metric struct { metricType qb.MetricType fracDigits byte @@ -34,20 +39,19 @@ type _metric struct { capturedState *CapturedState } -//IndexLevels [][]IndexRec // root - last element +func (s *_metric) LastValue() float64 { + if s.capturedState != nil { + return s.capturedState.LastValue + } + return s.lastValue +} -// func (s *_metric) ReinitBy(timestamp uint32, value float64) { -// s.Timestamps.Renew() -// s.Values.Renew() - -// s.Timestamps.Append(timestamp) -// s.Values.Append(value) - -// s.Since = timestamp -// s.SinceValue = value -// s.Until = timestamp -// s.UntilValue = value -// } +func (s *_metric) LastTimestamp() uint32 { + if s.capturedState != nil { + return s.capturedState.LastTimestamp + } + return s.timestamps.LastTimestamp() +} func (s *_metric) DeleteMeasures() { // s.Timestamps.Renew() @@ -60,11 +64,6 @@ func (s *_metric) DeleteMeasures() { // s.UntilValue = 0 } -type CapturedState struct { - LastTimestamp uint32 - LastValue float64 -} - // func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) { func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, sendToStorage func(any)) { if s.capturedState != nil { @@ -92,7 +91,6 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send s.capturedState = &CapturedState{ LastTimestamp: timestamps.LastTimestamp(), LastValue: s.lastValue, - //LastValue: values.LastValue(), } s.timestamps.CaptureState() @@ -128,7 +126,8 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send // сторінка заповнена since := s.timestamps.ReplaceSinceWithUntil() - if len(pages) == 0 { + if len(pages) == 0 && idx > 0 { + // idx > 0 required because page may overflows without append any data headTimestamps = timestamps.Tail(timestampsOffset) headValues = values.Tail(valuesOffset) } @@ -152,7 +151,6 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp) values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta) } - // s.lastValue = measure.Value written++ @@ -171,38 +169,38 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send if len(pages) > 0 { // пишу в storage довгим шляхом через redo файл і запис в data файл - sendToStorage(storage.AppendedMeasuresWithGrow{ - MetricID: req.MetricID, - LastPageNo: s.lastPageNo, - TimestampsOffset: timestampsRewindOffset, - Timestamps: headTimestamps, - ValuesOffset: valuesRewindOffset, - Values: headValues, - IndexLevelTails: s.indexLevelTails, - DataPages: pages, - TailTimestamps: timestamps.Tail(0), // payload - TailValues: values.Tail(0), - ResultCode: resultCode, - WrittenCount: written, - ResultCh: req.ResultCh, + sendToStorage(storage.MeasuresAppendWithGrow{ + MetricID: req.MetricID, + LastPageNo: s.lastPageNo, + TimestampsRewindOffset: timestampsRewindOffset, + Timestamps: headTimestamps, + ValuesRewindOffset: valuesRewindOffset, + Values: headValues, + IndexLevelTails: s.indexLevelTails, + DataPages: pages, + TailTimestamps: timestamps.Tail(0), // payload + TailValues: values.Tail(0), + ResultCode: resultCode, + WrittenCount: written, + //ResultCh: req.ResultCh, }) } else { // короткий шлях - запис лише в storage - sendToStorage(storage.AppendedMeasuresWithGrow{ - MetricID: req.MetricID, - TimestampsOffset: timestampsOffset, - ValuesOffset: valuesOffset, - Timestamps: timestamps.Tail(timestampsOffset), - Values: values.Tail(valuesOffset), - ResultCode: resultCode, - WrittenCount: written, - ResultCh: req.ResultCh, + sendToStorage(storage.MeasuresAppend{ + MetricID: req.MetricID, + TimestampsRewindOffset: timestampsRewindOffset, + ValuesRewindOffset: valuesRewindOffset, + Timestamps: timestamps.Tail(timestampsOffset), + Values: values.Tail(valuesOffset), + ResultCode: resultCode, + WrittenCount: written, + //ResultCh: req.ResultCh, }) } } -func (s *_metric) FinAppendMeasures(rec storage.AppendMeasuresSummary) { +func (s *_metric) FinAppendMeasures(rec storage.MeasuresAppendWithGrowCommited) { // Видаляю state. Оригінальні Timestamps і Values вже мають останню версію s.capturedState = nil @@ -217,71 +215,71 @@ func (s *_metric) FinAppendMeasures(rec storage.AppendMeasuresSummary) { // READ func (s *_metric) StartRangeScan(req tryRangeScanReq) { - if s.Since == 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return - } + // if s.Since == 0 { + // req.ResultCh <- rangeScanResult{ + // ResultCode: QueryDone, + // } + // return + // } - if req.Since > s.Until { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return - } + // if req.Since > s.Until { + // req.ResultCh <- rangeScanResult{ + // ResultCode: QueryDone, + // } + // return + // } - if req.Until < s.Since { - if s.RootPageNo > 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: UntilNotFound, - RootPageNo: s.RootPageNo, - FracDigits: s.FracDigits, - } - s.RLocks++ - return - } else { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return - } - } + // if req.Until < s.Since { + // if s.RootPageNo > 0 { + // req.ResultCh <- rangeScanResult{ + // ResultCode: UntilNotFound, + // RootPageNo: s.RootPageNo, + // FracDigits: s.FracDigits, + // } + // s.RLocks++ + // return + // } else { + // req.ResultCh <- rangeScanResult{ + // ResultCode: QueryDone, + // } + // return + // } + // } - timestampDecompressor := s.timestamps.CreateDecompressor() - valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits) + // timestampDecompressor := s.timestamps.CreateDecompressor() + // valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits) - for { - timestamp, done := timestampDecompressor.NextValue() - if done { - break - } - value, done := valueDecompressor.NextValue() - if done { - qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - } - if timestamp <= req.Until { - req.ResponseWriter.FeedNoSend(timestamp, value) - if timestamp < req.Since { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return - } - } - } - if s.lastPageNo > 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: UntilFound, - LastPageNo: s.lastPageNo, - FracDigits: s.fracDigits, - } - s.RLocks++ - } else { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - } + // for { + // timestamp, done := timestampDecompressor.NextValue() + // if done { + // break + // } + // value, done := valueDecompressor.NextValue() + // if done { + // qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + // } + // if timestamp <= req.Until { + // req.ResponseWriter.FeedNoSend(timestamp, value) + // if timestamp < req.Since { + // req.ResultCh <- rangeScanResult{ + // ResultCode: QueryDone, + // } + // return + // } + // } + // } + // if s.lastPageNo > 0 { + // req.ResultCh <- rangeScanResult{ + // ResultCode: UntilFound, + // LastPageNo: s.lastPageNo, + // FracDigits: s.fracDigits, + // } + // s.RLocks++ + // } else { + // req.ResultCh <- rangeScanResult{ + // ResultCode: QueryDone, + // } + // } } func (s *_metric) StartFullScan(req tryFullScanReq) { @@ -384,36 +382,3 @@ func (s *_metric) WriteTo(w io.Writer) (err error) { } return } - -// since - 4b -// sinceValue - 8b - -// untilValue - 8b - -// time - -// if s.state != nil { -// h = s.state.H -// lastUnixtime = s.state.LastUnixtime -// lastDelta = s.state.LastDelta -// payload = s.state.Payload -// } -// return NewTimeDeltaDecompressorFromState(TimeDeltaDecompressorFromStateOptions{ -// H: h, -// LastUnixtime: lastUnixtime, -// LastDelta: lastDelta, -// Payload: payload, -// }) - -// value -// if s.state != nil { -// h = s.state.H -// lastDelta = s.state.LastDelta -// payload = s.state.Payload -// } -// return NewValueDeltaDecompressorFromState(ValueDeltaDecompressorFromStateOptions{ -// Coef: s.coef, -// ToFloat64: s.toFloat64, -// H: s.buf[s.pos-1], -// LastDelta: s.lastDelta, -// Payload: s.buf[:s.pos], -// }) diff --git a/database/proc.go b/database/proc.go index 3ee0768..42d114c 100644 --- a/database/proc.go +++ b/database/proc.go @@ -294,7 +294,7 @@ func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq // } } -func (s *Database) finAppendMeasures(rec storage.AppendMeasuresSummary) { +func (s *Database) finAppendMeasures(rec storage.MeasuresAppendWithGrowCommited) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, @@ -398,8 +398,8 @@ func (s *Database) tryListCurrentValues(req tryListCurrentValuesReq) { if ok { req.ResponseWriter.BufferValue(transform.CurrentValue{ MetricID: metricID, - Timestamp: metric.timestamps.LastTimestamp(), - Value: metric.lastValue, + Timestamp: metric.LastTimestamp(), + Value: metric.LastValue(), }) } } @@ -420,13 +420,13 @@ func (s *Database) applyChanges(req storage.Changes) { // case storage.AppendedMeasure: // s.appendMeasure(rec) - case storage.AppendMeasuresSummary: + case storage.MeasuresAppendWithGrowCommited: s.finAppendMeasures(rec) // case storage.AppendedMeasureWithOverflowExtended: // s.appendMeasureAfterOverflow(rec) - case storage.DeletedMeasures: + case storage.MeasuresDeleteRecord: s.deleteMeasures(rec) } } @@ -555,7 +555,7 @@ func (s *Database) deleteMetric(rec storage.MetricDeleteRecord) { } } -func (s *Database) deleteMeasures(rec storage.DeletedMeasures) { +func (s *Database) deleteMeasures(rec storage.MeasuresDeleteRecord) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, diff --git a/database/replay_metric.go b/database/replay_metric.go index d56649a..c38ae58 100644 --- a/database/replay_metric.go +++ b/database/replay_metric.go @@ -113,7 +113,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow // Додати в index level tails недостаючі дані, або замінити for levelIdx, change := range rec.ChangedIndexLevels { var ( - head = change.HeadIndexPage + head = change.IndexPageTail newRecordsCount = len(change.TailRecords) / indexRecordSize ) // 1. набиваю сторінки для перезапису в index файлі @@ -171,7 +171,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow // 4. набиваю сторінки для перезапису в data файлі if isLastPacket { - dataPages = append(dataPages, composeHeadDataPage(s.databuf, rec.HeadDataPage)) + dataPages = append(dataPages, composeHeadDataPage(s.databuf, rec.DataPageTail)) for _, x := range rec.DataPages { dataPages = append(dataPages, storage.PageToWrite{ @@ -182,7 +182,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow } // рахую кількість reused дата сторінок - if rec.HeadDataPage.Reused { + if rec.DataPageTail.Reused { reusedDataPages++ } for _, x := range rec.DataPages { @@ -200,7 +200,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow s.vSize = len(rec.TailValues) if len(rec.DataPages) == 0 { - s.lastPageNo = rec.HeadDataPage.PageNo + s.lastPageNo = rec.DataPageTail.PageNo } else { s.lastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo } @@ -215,7 +215,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow // HELPERS -func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *storage.HeadIndexPage) storage.PageToWrite { +func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *storage.IndexPageTail) storage.PageToWrite { // розраховую pos, з якого буду дописувати хвіст pos := level.RecordsCount * indexRecordSize // створюю копію сторінки @@ -229,10 +229,10 @@ func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *stor ZeroLevel: levelIdx == 0, }) // перевірка CRC - if calculatedCRC != head.Checksum { + if calculatedCRC != head.CRC32 { qb.Abort(qb.WALReplayFailed, fmt.Errorf("calculated CRC %d not equal expected %d of head page %d on index level %d", - calculatedCRC, head.Checksum, head.PageNo, levelIdx)) + calculatedCRC, head.CRC32, head.PageNo, levelIdx)) } return storage.PageToWrite{ PageNo: head.PageNo, @@ -240,29 +240,29 @@ func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *stor } } -func composeHeadDataPage(databuf []byte, head storage.HeadDataPage) storage.PageToWrite { +func composeHeadDataPage(databuf []byte, head storage.DataPageTail) storage.PageToWrite { // створюю копію сторінки var ( page = make([]byte, storage.DataPageSize) - timestampsSize = head.TimestampsOffset - len(head.Timestamps) + timestampsSize = head.TimestampsRewindOffset - len(head.Timestamps) ) // копіюю поточні дані copy(page, databuf) // копіюю нові дані із WAL з урахуванням offset-ів - copy(page[head.ValuesOffset:], head.Values) + copy(page[head.ValuesRewindOffset:], head.Values) copy(page[len(page)-timestampsSize:], head.Timestamps) // запечатати сторінку calculatedCRC := storage.SealDataPage(storage.SealDataPageIn{ Content: page, PrevPageNo: head.PrevPageNo, TimestampsSize: timestampsSize, - ValuesSize: head.ValuesOffset + len(head.Values), + ValuesSize: head.ValuesRewindOffset + len(head.Values), }) // перевірка CRC - if calculatedCRC != head.Checksum { + if calculatedCRC != head.CRC32 { qb.Abort(qb.WALReplayFailed, fmt.Errorf("calculated CRC %d not equal expected %d of head data page %d", - calculatedCRC, head.Checksum, head.PageNo)) + calculatedCRC, head.CRC32, head.PageNo)) } return storage.PageToWrite{ PageNo: head.PageNo, diff --git a/freelist/freelist.go b/freelist/freelist.go index 10dd8e8..173e1bd 100644 --- a/freelist/freelist.go +++ b/freelist/freelist.go @@ -105,7 +105,7 @@ func (s *FreeList) save() error { bin.PutUint32(buf[i:], pageNo) i += ptrSize } - bin.PutUint32(buf, util.CalcChecksum(buf[crcSize:])) + bin.PutUint32(buf, util.CalculateCRC32(buf[crcSize:])) // var ( file *os.File @@ -158,7 +158,7 @@ func (s *FreeList) loadLastPage() error { } // check crc checksum, _ := bin.GetUint32(buf[0:]) - if util.CalcChecksum(buf[crcSize:]) != checksum { + if util.CalculateCRC32(buf[crcSize:]) != checksum { return fmt.Errorf("page %d is corrupted", s.basePagesCount) } for i := crcSize; i < len(buf); i += ptrSize { diff --git a/go.mod b/go.mod index 9b6cb11..6b4b8df 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,6 @@ go 1.24.2 require ( gopkg.in/ini.v1 v1.67.1 - gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696 + gordenko.dev/dima/bin v0.0.0-20260612161453-4ee9be3474fb gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 ) diff --git a/go.sum b/go.sum index f2d3751..52f17ce 100644 --- a/go.sum +++ b/go.sum @@ -22,5 +22,7 @@ gordenko.dev/dima/bin v0.0.0-20260606153512-101b10d92a39 h1:qW3SnQ9HAcJpfj8ottBf gordenko.dev/dima/bin v0.0.0-20260606153512-101b10d92a39/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ= gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696 h1:JY91JYvv1g+XbaKdGmLFDBueUTx51RdW9mnbNYhiUdw= gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ= +gordenko.dev/dima/bin v0.0.0-20260612161453-4ee9be3474fb h1:XBYZbK5Z+yfpKmSR2wuPLyW14ReCRSNdOlL5MWro8ns= +gordenko.dev/dima/bin v0.0.0-20260612161453-4ee9be3474fb/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 h1:nyJ3mzTQ46yUeMZCdLyYcs7B5JCS54c67v84miyhq2E= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69/go.mod h1:AxgKDktpqBVyIOhIcP+nlCpK+EsJyjN5kPdqyd8euVU= diff --git a/storage/data_preparer.go b/storage/data_preparer.go deleted file mode 100644 index 79367a2..0000000 --- a/storage/data_preparer.go +++ /dev/null @@ -1,228 +0,0 @@ -package storage - -import ( - "bytes" - "io" - - bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb/util" -) - -type DataPreparer struct { - pagePreparer *PagePreparer - //dataPagePayloadBound int - w *bytes.Buffer - indexPages []PageToWrite - dataPages []PageToWrite - writeResults []any -} - -type DataPreparerOptions struct { - MinWALBufferSize int - PagePreparer *PagePreparer -} - -func NewDataPreparer(opt DataPreparerOptions) *DataPreparer { - buf := make([]byte, opt.MinWALBufferSize) - s := &DataPreparer{ - pagePreparer: opt.PagePreparer, - w: bytes.NewBuffer(buf), - } - return s -} - -type PreparedData struct { - Packet []byte // пакет для запису в WAL - WriteToIndex []PageToWrite - WriteToData []PageToWrite - WriteResults []any -} - -func (s *DataPreparer) Prepare(input []any) PreparedData { - s.w.Write([]byte{ - 0, 0, 0, 0, 0, 0, 0, 0, 0, // size (max 9 byte) - 0, 0, 0, 0, // crc32 - }) - - hasher := util.NewHasher() - - w := io.MultiWriter(s.w, hasher) - - // 1. Пакую всі дані в WAL буфер запису - for _, untyped := range input { - switch x := untyped.(type) { - case MetricAddRecord: - x.Pack(w) - // case DeletedMetric: - // if len(x.FreePageNumbers) > 0 { - // s.freeList.AddPageNumbers(x.FreePageNumbers) - // } - case MeasuresAppendRecord: - rec := MeasuresAppendRecord{ - MetricID: x.MetricID, - TimestampsRewindOffset: x.TimestampsRewindOffset, - Timestamps: x.Timestamps, - ValuesRewindOffset: x.ValuesRewindOffset, - Values: x.Values, - } - - rec.WriteTo(s.w) - - // Дані для Worker - s.writeResults = append(s.writeResults, AppendMeasuresSummary{ - // MetricID: x.MetricID, - // LastPageNo: sealResult.DataPages[len(sealResult.DataPages)-1].PageNo, - // Index: indexLevelTails, - // ResultCode: x.ResultCode, - // WrittenCount: x.WrittenCount, - // ResultCh: x.ResultCh, - }) - - case AppendedMeasuresWithGrow: - sealResult := s.pagePreparer.SealPages(SealPagesIn{ - LastPageNo: x.LastPageNo, - IndexLevelTails: x.IndexLevelTails, - DataPages: x.DataPages, - }) - - // сторінки для запису в index та data файли - for _, p := range sealResult.DataPages { - s.dataPages = append(s.dataPages, PageToWrite{ - PageNo: p.PageNo, - Content: p.Content, - }) - } - for _, level := range sealResult.IndexLevels { - for _, p := range level.IndexPages { - s.indexPages = append(s.indexPages, PageToWrite{ - PageNo: p.PageNo, - Content: p.Content, - }) - } - } - - var ( - // дані для запису в WAL - changedIndexPages []ChangedIndexLevel - // дані для Воркера - indexLevelTails []IndexLevelTail - ) - - for _, level := range sealResult.IndexLevels { - if len(level.IndexPages) == 0 && level.SkipRecords == level.TailRecordsCount { - break - } - var completedIndexPage *HeadIndexPage - if len(level.IndexPages) > 0 { - first := level.IndexPages[0] - skipSize := level.SkipRecords * indexRecordSize - completedIndexPage = &HeadIndexPage{ - PageNo: first.PageNo, - Reused: first.Reused, - Checksum: first.Checksum, - Records: first.Content[skipSize:], - } - } - changedIndexPages = append(changedIndexPages, ChangedIndexLevel{ - HeadIndexPage: completedIndexPage, - IndexPages: level.IndexPages[1:], - TailRecords: level.TailRecords[:level.TailRecordsCount*indexRecordSize], - }) - - indexLevelTails = append(indexLevelTails, IndexLevelTail{ - Buffer: level.TailRecords, - RecordsCount: level.TailRecordsCount, - }) - } - - first := sealResult.DataPages[0] - - rec := MeasuresAppendWithGrowRecord{ - MetricID: x.MetricID, - HeadDataPage: HeadDataPage{ - PageNo: first.PageNo, - Reused: first.Reused, - PrevPageNo: first.PrevPageNo, - Checksum: first.Checksum, - TimestampsOffset: x.TimestampsOffset, - Timestamps: x.Timestamps, - ValuesOffset: x.ValuesOffset, - Values: x.Values, - }, - DataPages: sealResult.DataPages[1:], - ChangedIndexLevels: changedIndexPages, - TailTimestamps: x.Timestamps, - TailValues: x.Values, - } - - rec.WriteTo(s.w) - - // Дані для Worker - s.writeResults = append(s.writeResults, AppendMeasuresSummary{ - MetricID: x.MetricID, - LastPageNo: sealResult.DataPages[len(sealResult.DataPages)-1].PageNo, - Index: indexLevelTails, - ResultCode: x.ResultCode, - WrittenCount: x.WrittenCount, - ResultCh: x.ResultCh, - }) - - //case DeletedMeasures: - //case DeletedMeasuresSince: - } - } - // 2. Додаю розмір пакету і чексуму - - //s.written += int64(len(packet)) + 12 - - packet := s.w.Bytes() - // write size - payloadSize := len(packet) - 13 - start := 9 - bin.CountVarSize(payloadSize) - bin.PutVarSize(packet[start:], payloadSize) - // checksum - bin.PutUint32(packet[9:], hasher.Sum32()) - // - return PreparedData{ - Packet: packet[start:], - WriteToIndex: s.indexPages, - WriteToData: s.dataPages, - WriteResults: s.writeResults, - } -} - -func (s *DataPreparer) Reset() { - s.w.Reset() - s.indexPages = nil - s.dataPages = nil - s.writeResults = nil -} - -// MeasuresToWrite -type AppendedMeasuresWithGrow struct { - MetricID uint32 - LastPageNo uint32 - TimestampsOffset int // offset on 1st page - Timestamps []byte - ValuesOffset int // offset on 1st page - Values []byte - IndexLevelTails []IndexLevelTail - DataPages []DataPayload - TailTimestamps []byte // data level tail - TailValues []byte // data level tail - ResultCode byte - WrittenCount int - ResultCh chan struct{} -} - -// WRITE RESULTS - -// результат -type AppendMeasuresSummary struct { - MetricID uint32 - LastPageNo uint32 - Index []IndexLevelTail - ResultCode byte - WrittenCount int - ResultCh chan struct{} -} diff --git a/storage/page_preparer.go b/storage/page_ingester.go similarity index 78% rename from storage/page_preparer.go rename to storage/page_ingester.go index d3f45da..0b61a63 100644 --- a/storage/page_preparer.go +++ b/storage/page_ingester.go @@ -14,56 +14,30 @@ import ( // levels []*IndexLevel, dataPages []*DataPage // fix - reduce leves -type PagePreparer struct { +type PageIngester struct { getDataPageNumber func() (uint32, bool, error) getIndexPageNumber func() (uint32, bool, error) } -type PagePreparerOptions struct { +type PageIngesterOptions struct { GetDataPageNumber func() (uint32, bool, error) GetIndexPageNumber func() (uint32, bool, error) } -func NewPagePreparer(opt PagePreparerOptions) (*PagePreparer, error) { +func NewPageIngester(opt PageIngesterOptions) (*PageIngester, error) { if opt.GetIndexPageNumber == nil { return nil, errors.New("missing required option: GetIndexPageNumber") } if opt.GetDataPageNumber == nil { return nil, errors.New("missing required option: GetDataPageNumber") } - s := &PagePreparer{ + s := &PageIngester{ getIndexPageNumber: opt.GetIndexPageNumber, getDataPageNumber: opt.GetDataPageNumber, } return s, nil } -type appendIndexRecordIn struct { - Records []byte - RecordsCount int - Timestamp uint32 - PageNo uint32 -} - -func (s *PagePreparer) appendIndexRecord(in appendIndexRecordIn) []byte { - if in.RecordsCount < maxRecordsOnIndexPage { - var ( - pos int - buf = in.Records - ) - if in.RecordsCount > 0 { - pos = in.RecordsCount * indexRecordSize - } else { - buf = make([]byte, IndexPageSize) // IndexPageIncSize - } - bin.PutUint32(buf[pos:], in.Timestamp) - bin.PutUint32(buf[pos+4:], in.PageNo) - return buf - } - qb.Abort(qb.NoSpaceOnIndexPage, nil) - return nil -} - type IndexLevelTail struct { Buffer []byte // розмір більший за кількість RecordsCount int @@ -76,25 +50,22 @@ type DataPayload struct { ValuesSize int } -type SealPagesIn struct { +type IngestIn struct { LastPageNo uint32 IndexLevelTails []IndexLevelTail DataPages []DataPayload } type SealedDataPage struct { - PageNo uint32 - Reused bool - PrevPageNo uint32 - Checksum uint32 - Content []byte + PageNo uint32 + Reused bool + Content []byte } type SealedIndexPage struct { - Content []byte - PageNo uint32 - Reused bool - Checksum uint32 + PageNo uint32 + Reused bool + Content []byte } // SealedIndexLevel - як зрозуміти чи потрібно щось писати в WAL чи index файл, @@ -115,12 +86,7 @@ type SealedIndexLevel struct { TailRecordsCount int } -type SealResult struct { - DataPages []SealedDataPage - IndexLevels []*SealedIndexLevel -} - -func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { +func (s *PageIngester) Ingest(in IngestIn) ([]SealedDataPage, []*SealedIndexLevel) { var ( sealedDataPages []SealedDataPage sealedIndexLevels []*SealedIndexLevel @@ -142,31 +108,31 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { if dataPageIdx == 0 { prevPageNo = in.LastPageNo } else { - prevPageNo = result.DataPages[dataPageIdx-1].PageNo + prevPageNo = sealedDataPages[dataPageIdx-1].PageNo } pageNo, reused, err := s.getDataPageNumber() if err != nil { qb.Abort(qb.FailedGetPageNumber, err) } + + SealDataPage(SealDataPageIn{ + Content: d.Content, + PrevPageNo: prevPageNo, + TimestampsSize: d.TimestampsSize, + ValuesSize: d.ValuesSize, + }) + sealedDataPages = append(sealedDataPages, SealedDataPage{ - Content: d.Content, - PrevPageNo: prevPageNo, - PageNo: pageNo, - Reused: reused, - Checksum: SealDataPage(SealDataPageIn{ - Content: d.Content, - PrevPageNo: prevPageNo, - TimestampsSize: d.TimestampsSize, - ValuesSize: d.ValuesSize, - }), + Content: d.Content, + PageNo: pageNo, + Reused: reused, }) upPageNo = pageNo for { if levelIdx == len(sealedIndexLevels) { // новий root sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{ - TailRecords: s.appendIndexRecord(appendIndexRecordIn{ - Records: nil, + TailRecords: appendIndexRecord(appendIndexRecordIn{ Timestamp: upTimestamp, PageNo: upPageNo, }), @@ -176,10 +142,11 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { } level := sealedIndexLevels[levelIdx] if level.TailRecordsCount < maxRecordsOnIndexPage { - level.TailRecords = s.appendIndexRecord(appendIndexRecordIn{ - Records: level.TailRecords, - Timestamp: upTimestamp, - PageNo: upPageNo, + level.TailRecords = appendIndexRecord(appendIndexRecordIn{ + Records: level.TailRecords, + RecordsCount: level.TailRecordsCount, + Timestamp: upTimestamp, + PageNo: upPageNo, }) level.TailRecordsCount++ break @@ -191,20 +158,20 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { qb.Abort(qb.FailedGetPageNumber, err) } + SealIndexPage(SealIndexPageIn{ + Content: level.TailRecords, + RecordsCount: level.TailRecordsCount, + ZeroLevel: levelIdx == 0, + }) + filled := SealedIndexPage{ Content: level.TailRecords, PageNo: pageNo, Reused: reused, - Checksum: SealIndexPage(SealIndexPageIn{ - Content: level.TailRecords, - RecordsCount: level.TailRecordsCount, - ZeroLevel: levelIdx == 0, - }), } level.IndexPages = append(level.IndexPages, filled) // новий tail на індексному рівні - level.TailRecords = s.appendIndexRecord(appendIndexRecordIn{ - Records: nil, + level.TailRecords = appendIndexRecord(appendIndexRecordIn{ Timestamp: upTimestamp, PageNo: upPageNo, }) @@ -215,10 +182,7 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { levelIdx++ } } - return SealResult{ - DataPages: sealedDataPages, - IndexLevels: sealedIndexLevels, - } + return sealedDataPages, sealedIndexLevels } // HELPERS @@ -235,7 +199,7 @@ func SealDataPage(in SealDataPageIn) (checksum uint32) { bin.PutUint16(in.Content[valuesSizeIdx:], uint16(in.ValuesSize)) bin.PutUint32(in.Content[prevPageIdx:], in.PrevPageNo) - checksum = util.CalcChecksum(in.Content[:dataCRC32Idx]) + checksum = util.CalculateCRC32(in.Content[:dataCRC32Idx]) bin.PutUint32(in.Content[dataCRC32Idx:], checksum) return } @@ -251,7 +215,33 @@ func SealIndexPage(in SealIndexPageIn) (checksum uint32) { if in.ZeroLevel { in.Content[isZeroLevelIdx] = 1 } - checksum = util.CalcChecksum(in.Content[:indexCRC32Idx]) + checksum = util.CalculateCRC32(in.Content[:indexCRC32Idx]) bin.PutUint32(in.Content[indexCRC32Idx:], checksum) return } + +type appendIndexRecordIn struct { + Records []byte + RecordsCount int + Timestamp uint32 + PageNo uint32 +} + +func appendIndexRecord(in appendIndexRecordIn) []byte { + if in.RecordsCount < maxRecordsOnIndexPage { + var ( + pos int + buf = in.Records + ) + if in.RecordsCount > 0 { + pos = in.RecordsCount * indexRecordSize + } else { + buf = make([]byte, IndexPageSize) // IndexPageIncSize + } + bin.PutUint32(buf[pos:], in.Timestamp) + bin.PutUint32(buf[pos+4:], in.PageNo) + return buf + } + qb.Abort(qb.NoSpaceOnIndexPage, nil) + return nil +} diff --git a/storage/storage.go b/storage/storage.go index 98a5aa3..2d5b7ea 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -1,6 +1,8 @@ package storage -const ( +import "gordenko.dev/dima/qb" + +var ( //PageNoSize = 4 @@ -15,18 +17,106 @@ const ( prevPageIdx = DataPageSize - 12 // index page - IndexPageSize = 1024 - IndexPageFooterSize = 7 + IndexPageSize = 1024 + IndexPageFooterSize = 7 + //indexPageIncSize = IndexPageIncSize + indexCRC32Idx = IndexPageSize - 4 + indexRecordsCountIdx = IndexPageSize - 6 + isZeroLevelIdx = IndexPageSize - 7 indexRecordSize = 8 maxRecordsOnIndexPage = (IndexPageSize - IndexPageFooterSize) / indexRecordSize - //indexPageIncSize = IndexPageIncSize - indexCRC32Idx = IndexPageSize - 4 - indexRecordsCountIdx = IndexPageSize - 6 - isZeroLevelIdx = IndexPageSize - 7 - // timestampSize = 4 // pairSize = timestampSize + PageNoSize // indexFooterIdx = indexRecordsQtyIdx // dataFooterIdx = timestampsSizeIdx ) + +type IndexRecord struct { + Timestamp uint32 + PageNo uint32 +} + +// TASKS FROM WORKER + +type MetricAdd struct { + MetricID uint32 + MetricType qb.MetricType + FracDigits byte + ResultCh chan struct{} +} + +type MetricDelete struct { + MetricID uint32 + FreeIndexPages []uint32 + FreeDataPages []uint32 + ResultCh chan struct{} +} + +type MeasuresAppendWithGrow struct { + MetricID uint32 + LastPageNo uint32 + TimestampsRewindOffset int // offset on 1st page + Timestamps []byte + ValuesRewindOffset int // offset on 1st page + Values []byte + DataPages []DataPayload + TailTimestamps []byte // data level tail + TailValues []byte // data level tail + IndexLevelTails []IndexLevelTail + ResultCode byte + WrittenCount int + ResultCh chan struct{} +} + +type MeasuresAppend struct { + MetricID uint32 + TimestampsRewindOffset int // offset on 1st page + Timestamps []byte + ValuesRewindOffset int // offset on 1st page + Values []byte + ResultCode byte + WrittenCount int + ResultCh chan struct{} +} + +type MeasuresDelete struct { + MetricID uint32 + FreeIndexPages []uint32 + FreeDataPages []uint32 + ResultCh chan struct{} +} + +// WRITE RESULTS + +// результат +type MeasuresAppendCommited struct { + MetricID uint32 + ResultCode byte + WrittenCount int + ResultCh chan struct{} +} + +type MeasuresAppendWithGrowCommited struct { + MetricID uint32 + LastPageNo uint32 + Index []IndexLevelTail + ResultCode byte + WrittenCount int + ResultCh chan struct{} +} + +type MeasuresDeleteCommited struct { + MetricID uint32 + ResultCh chan struct{} +} + +type MetricAddCommited struct { + MetricID uint32 + ResultCh chan struct{} +} + +type MetricDeleteCommited struct { + MetricID uint32 + ResultCh chan struct{} +} diff --git a/storage/storage_test.go b/storage/storage_test.go new file mode 100644 index 0000000..71c10ba --- /dev/null +++ b/storage/storage_test.go @@ -0,0 +1,871 @@ +package storage + +import ( + "bytes" + "fmt" + "math" + "reflect" + "slices" + "testing" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/util" +) + +func TestAppendIndexRecord(t *testing.T) { + IndexPageSize = 24 + var ( + after = []byte{ + 0x01, 0x00, 0x00, 0x00, 0x08, 0x00, 0x00, 0x00, + 0x02, 0x00, 0x00, 0x00, 0x09, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, + } + ) + // func creates index page buffer + buf := appendIndexRecord(appendIndexRecordIn{ + Timestamp: 1, + PageNo: 8, + }) + buf = appendIndexRecord(appendIndexRecordIn{ + Records: buf, + RecordsCount: 1, + Timestamp: 2, + PageNo: 9, + }) + + if !bytes.Equal(buf, after) { + t.Fatalf("before not equal after:\nbefore: % x\n after: % x", + buf, after) + } +} + +func TestSealIndexPage(t *testing.T) { + IndexPageSize = 16 + indexCRC32Idx = IndexPageSize - 4 + indexRecordsCountIdx = IndexPageSize - 6 + isZeroLevelIdx = IndexPageSize - 7 + var ( + before = []byte{ + 0x07, 0x07, 0x07, 0x07, 0x08, 0x08, 0x08, 0x08, 0x00, // payload + 0x00, // zero level + 0x00, 0x00, // records count + 0x00, 0x00, 0x00, 0x00, // crc32 + } + after = []byte{ + 0x07, 0x07, 0x07, 0x07, 0x08, 0x08, 0x08, 0x08, 0x00, // payload + 0x01, // zero level + 0x02, 0x00, // records count + 0x8d, 0xcc, 0x3c, 0xc0, // crc32 + } + ) + + calculatedCRC := util.CalculateCRC32(after[:indexCRC32Idx]) + + checksum := SealIndexPage(SealIndexPageIn{ + Content: before, + RecordsCount: 2, + ZeroLevel: true, + }) + + if calculatedCRC != checksum { + t.Fatalf("calculated CRC %d not equal returned %d", + calculatedCRC, checksum) + } + + if !bytes.Equal(before, after) { + t.Fatalf("before not equal after:\nbefore: % x\n after: % x", + before, after) + } +} + +func TestSealDataPage(t *testing.T) { + DataPageSize = 16 + dataCRC32Idx = DataPageSize - 4 + timestampsSizeIdx = DataPageSize - 6 + valuesSizeIdx = DataPageSize - 8 + prevPageIdx = DataPageSize - 12 + var ( + before = []byte{ + 0x07, 0x07, 0x07, 0x07, // payload + 0x00, 0x00, 0x00, 0x00, // prevPageNo + 0x00, 0x00, // values size + 0x00, 0x00, // timestamps size + 0x00, 0x00, 0x00, 0x00, // crc32 + } + after = []byte{ + 0x07, 0x07, 0x07, 0x07, // payload + 0x09, 0x00, 0x00, 0x00, // prevPageNo + 0x03, 0x00, // values size + 0x01, 0x00, // timestamps size + 0xa9, 0xec, 0xc8, 0x91, // crc32 + } + ) + + calculatedCRC := util.CalculateCRC32(after[:dataCRC32Idx]) + + checksum := SealDataPage(SealDataPageIn{ + Content: before, + PrevPageNo: 9, + TimestampsSize: 1, + ValuesSize: 3, + }) + + if calculatedCRC != checksum { + t.Fatalf("calculated CRC %d not equal returned %d", + calculatedCRC, checksum) + } + + if !bytes.Equal(before, after) { + t.Fatalf("before not equal after:\nbefore: % x\n after: % x", + before, after) + } +} + +func TestPageIngester(t *testing.T) { + // index + IndexPageSize = 24 + indexCRC32Idx = IndexPageSize - 4 + indexRecordsCountIdx = IndexPageSize - 6 + isZeroLevelIdx = IndexPageSize - 7 + maxRecordsOnIndexPage = 2 + // data + DataPageSize = 16 + DataPagePayloadSize = DataPageSize - DataPageFooterSize + dataCRC32Idx = DataPageSize - 4 + timestampsSizeIdx = DataPageSize - 6 + valuesSizeIdx = DataPageSize - 8 + prevPageIdx = DataPageSize - 12 + + var testCases = []struct { + Name string + DataPages []DataPayload + ExpectedIndexLevels []ExpectedIndexLevel + ExpectedDataPages []ExpectedDataPage + }{ + { + Name: "create root index tail", + DataPages: []DataPayload{ + { + Since: 100, + Content: makeDataPage([]byte{0x01, 0x00, 0x00, 0x00}), + }, + }, + ExpectedDataPages: []ExpectedDataPage{ + { + PageNo: 1, + PrevPageNo: 0, + Payload: []byte{0x01, 0x00, 0x00, 0x00}, + }, + }, + ExpectedIndexLevels: []ExpectedIndexLevel{ + { + TailRecords: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + }, + TailRecordsCount: 1, + }, + }, + }, + { + Name: "just append record to index tail", + DataPages: []DataPayload{ + { + Since: 100, + Content: makeDataPage([]byte{0x01, 0x00, 0x00, 0x00}), + }, + { + Since: 200, + Content: makeDataPage([]byte{0x02, 0x00, 0x00, 0x00}), + }, + }, + ExpectedDataPages: []ExpectedDataPage{ + { + PageNo: 1, + PrevPageNo: 0, + Payload: []byte{0x01, 0x00, 0x00, 0x00}, + }, + { + PageNo: 2, + PrevPageNo: 1, + Payload: []byte{0x02, 0x00, 0x00, 0x00}, + }, + }, + ExpectedIndexLevels: []ExpectedIndexLevel{ + { + TailRecords: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + {Timestamp: 200, PageNo: 2}, + }, + TailRecordsCount: 2, + }, + }, + }, + { + Name: "create index page and upper level (root)", + DataPages: []DataPayload{ + { + Since: 100, + Content: makeDataPage([]byte{0x01, 0x00, 0x00, 0x00}), + }, + { + Since: 200, + Content: makeDataPage([]byte{0x02, 0x00, 0x00, 0x00}), + }, + { + Since: 300, + Content: makeDataPage([]byte{0x03, 0x00, 0x00, 0x00}), + }, + }, + ExpectedDataPages: []ExpectedDataPage{ + { + PageNo: 1, + PrevPageNo: 0, + Payload: []byte{0x01, 0x00, 0x00, 0x00}, + }, + { + PageNo: 2, + PrevPageNo: 1, + Payload: []byte{0x02, 0x00, 0x00, 0x00}, + }, + { + PageNo: 3, + PrevPageNo: 2, + Payload: []byte{0x03, 0x00, 0x00, 0x00}, + }, + }, + ExpectedIndexLevels: []ExpectedIndexLevel{ + { + IndexPages: []ExpectedIndexPage{ + { + PageNo: 1, + ZeroLevel: true, + Records: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + {Timestamp: 200, PageNo: 2}, + }, + }, + }, + TailRecords: []IndexRecord{ + {Timestamp: 300, PageNo: 3}, + }, + TailRecordsCount: 1, + }, + { + TailRecords: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + }, + TailRecordsCount: 1, + }, + }, + }, + { + Name: "just append record to tail on zero level", + DataPages: []DataPayload{ + { + Since: 100, + Content: makeDataPage([]byte{0x01, 0x00, 0x00, 0x00}), + }, + { + Since: 200, + Content: makeDataPage([]byte{0x02, 0x00, 0x00, 0x00}), + }, + { + Since: 300, + Content: makeDataPage([]byte{0x03, 0x00, 0x00, 0x00}), + }, + { + Since: 400, + Content: makeDataPage([]byte{0x04, 0x00, 0x00, 0x00}), + }, + }, + ExpectedDataPages: []ExpectedDataPage{ + { + PageNo: 1, + PrevPageNo: 0, + Payload: []byte{0x01, 0x00, 0x00, 0x00}, + }, + { + PageNo: 2, + PrevPageNo: 1, + Payload: []byte{0x02, 0x00, 0x00, 0x00}, + }, + { + PageNo: 3, + PrevPageNo: 2, + Payload: []byte{0x03, 0x00, 0x00, 0x00}, + }, + { + PageNo: 4, + PrevPageNo: 3, + Payload: []byte{0x04, 0x00, 0x00, 0x00}, + }, + }, + ExpectedIndexLevels: []ExpectedIndexLevel{ + { + IndexPages: []ExpectedIndexPage{ + { + PageNo: 1, + ZeroLevel: true, + Records: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + {Timestamp: 200, PageNo: 2}, + }, + }, + }, + TailRecords: []IndexRecord{ + {Timestamp: 300, PageNo: 3}, + {Timestamp: 400, PageNo: 4}, + }, + TailRecordsCount: 2, + }, + { + TailRecords: []IndexRecord{ + {Timestamp: 100, PageNo: 1}, + }, + TailRecordsCount: 1, + }, + }, + }, + } + + for _, testCase := range testCases { + pageManager := new(PageNumbersMock) + + ingester, err := NewPageIngester(PageIngesterOptions{ + GetIndexPageNumber: pageManager.GetIndexPageNumber, + GetDataPageNumber: pageManager.GetDataPageNumber, + }) + if err != nil { + t.Fatal(err) + } + + sealedDataPages, sealedIndexLevels := ingester.Ingest(IngestIn{ + LastPageNo: 0, + IndexLevelTails: nil, + DataPages: testCase.DataPages, + }) + + if len(testCase.ExpectedDataPages) != len(sealedDataPages) { + t.Fatalf("%s: got %d data pages, but not %d", + testCase.Name, len(sealedDataPages), len(testCase.ExpectedDataPages)) + } + + for idx, page := range testCase.ExpectedDataPages { + err := cmpDataPage(page, sealedDataPages[idx]) + if err != nil { + t.Fatalf("%s: data page #%d: %s", testCase.Name, idx, err) + } + } + + if len(testCase.ExpectedIndexLevels) != len(sealedIndexLevels) { + t.Fatalf("%s: got %d index levels, but not %d", + testCase.Name, len(sealedIndexLevels), len(testCase.ExpectedIndexLevels)) + } + + for idx, level := range testCase.ExpectedIndexLevels { + err := cmpIndexLevel(level, *sealedIndexLevels[idx]) + if err != nil { + t.Fatalf("%s: index level #%d: %s", testCase.Name, idx, err) + } + } + } +} + +type PageNumbersMock struct { + indexPageNo uint32 + dataPageNo uint32 +} + +func (s *PageNumbersMock) GetIndexPageNumber() (uint32, bool, error) { + s.indexPageNo++ + return s.indexPageNo, false, nil +} + +func (s *PageNumbersMock) GetDataPageNumber() (uint32, bool, error) { + s.dataPageNo++ + return s.dataPageNo, false, nil +} + +// func getIndexRecords(buf []byte, count int) (list []IndexRecord) { +// i := 0 +// for range count { +// var rec IndexRecord +// rec.Timestamp, _ = bin.GetUint32(buf[i:]) +// rec.PageNo, _ = bin.GetUint32(buf[i+4:]) +// list = append(list, rec) +// i += indexRecordSize +// } +// return +// } + +type ExpectedIndexPage struct { + PageNo uint32 + Records []IndexRecord + ZeroLevel bool +} + +func cmpIndexPage(expected ExpectedIndexPage, page SealedIndexPage) error { + if expected.PageNo != page.PageNo { + return fmt.Errorf("got pageNo %d not equal expected %d", + page.PageNo, expected.PageNo) + } + buf := page.Content + calculatedCRC := util.CalculateCRC32(buf[:indexCRC32Idx]) + checksum, _ := bin.GetUint32(buf[indexCRC32Idx:]) + + if calculatedCRC != checksum { + return fmt.Errorf("calculated CRC %d not equal written %d", + calculatedCRC, checksum) + } + + zeroLevel, _ := bin.GetBool(buf[isZeroLevelIdx:]) + if expected.ZeroLevel != zeroLevel { + return fmt.Errorf("expected zero level %t not equal written %t", + expected.ZeroLevel, zeroLevel) + } + + count, _ := bin.GetUint16(buf[indexRecordsCountIdx:]) + + records := getIndexRecords(buf, int(count)) + if !slices.Equal(expected.Records, records) { + return fmt.Errorf("expected records %v not equal written %v", + expected.Records, records) + } + return nil +} + +type ExpectedIndexLevel struct { + SkipRecords int + IndexPages []ExpectedIndexPage + TailRecords []IndexRecord + TailRecordsCount int +} + +func cmpIndexLevel(expected ExpectedIndexLevel, level SealedIndexLevel) error { + if expected.SkipRecords != level.SkipRecords { + return fmt.Errorf("got skipRecords %d not equal expected %d", + level.SkipRecords, expected.SkipRecords) + } + + if expected.TailRecordsCount != level.TailRecordsCount { + return fmt.Errorf("got tailRecordsCount %d not equal expected %d", + level.TailRecordsCount, expected.TailRecordsCount) + } + + records := getIndexRecords(level.TailRecords, level.TailRecordsCount) + if !slices.Equal(expected.TailRecords, records) { + return fmt.Errorf("expected tailRecords %v not equal written %v", + expected.TailRecords, records) + } + + if len(expected.IndexPages) != len(level.IndexPages) { + return fmt.Errorf("got %d index pages, but not %d", + len(level.IndexPages), len(expected.IndexPages)) + } + + for idx, page := range expected.IndexPages { + err := cmpIndexPage(page, level.IndexPages[idx]) + if err != nil { + return fmt.Errorf("index page #%d: %s", idx, err) + } + } + return nil +} + +type ExpectedDataPage struct { + PageNo uint32 + Payload []byte + PrevPageNo uint32 +} + +func cmpDataPage(expected ExpectedDataPage, page SealedDataPage) error { + buf := page.Content + + if expected.PageNo != page.PageNo { + return fmt.Errorf("got pageNo %d not equal expected %d", + page.PageNo, expected.PageNo) + } + prevPageNo, _ := bin.GetUint32(buf[prevPageIdx:]) + if expected.PrevPageNo != prevPageNo { + return fmt.Errorf("written prevPageNo %d not equal expected %d", + prevPageNo, expected.PrevPageNo) + } + + calculatedCRC := util.CalculateCRC32(buf[:dataCRC32Idx]) + checksum, _ := bin.GetUint32(buf[dataCRC32Idx:]) + + if calculatedCRC != checksum { + return fmt.Errorf("calculated CRC %d not equal written %d", + calculatedCRC, checksum) + } + + writtenPayload := page.Content[:DataPagePayloadSize] + if !slices.Equal(expected.Payload, writtenPayload) { + return fmt.Errorf("expected payload % x not equal written % x", + expected.Payload, writtenPayload) + } + return nil +} + +func makeDataPage(payload []byte) []byte { + buf := make([]byte, DataPageSize) + copy(buf, payload) + return buf +} + +func TestMetricAddRecord(t *testing.T) { + rec := MetricAddRecord{ + MetricID: 12345, + MetricType: qb.Cumulative, + FracDigits: 5, + } + buf := bytes.NewBuffer(nil) + rec.Pack(buf) + + recordType, _ := buf.ReadByte() + + if recordType != CodeMetricAdd { + t.Fatalf("wrong record type: %d", recordType) + } + + decoded := new(MetricAddRecord) + err := decoded.Parse(buf) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rec, *decoded) { + t.Fatalf("decoded are not equal origin: %v", decoded) + } +} + +func TestMetricDeleteRecord(t *testing.T) { + rec := MetricDeleteRecord{ + MetricID: 12345, + FreeIndexPages: []uint32{1, math.MaxUint32}, + FreeDataPages: []uint32{1, math.MaxUint32}, + } + buf := bytes.NewBuffer(nil) + rec.Pack(buf) + + recordType, _ := buf.ReadByte() + + if recordType != CodeMetricDelete { + t.Fatalf("wrong record type: %d", recordType) + } + + decoded := new(MetricDeleteRecord) + err := decoded.Parse(buf) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rec, *decoded) { + t.Fatalf("decoded are not equal origin: %v", decoded) + } +} + +func TestMeasuresDeleteRecord(t *testing.T) { + rec := MeasuresDeleteRecord{ + MetricID: 12345, + FreeIndexPages: []uint32{1, math.MaxUint32}, + FreeDataPages: []uint32{1, math.MaxUint32}, + } + buf := bytes.NewBuffer(nil) + rec.Pack(buf) + + recordType, _ := buf.ReadByte() + + if recordType != CodeMeasuresDelete { + t.Fatalf("wrong record type: %d", recordType) + } + + decoded := new(MeasuresDeleteRecord) + err := decoded.Parse(buf) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rec, *decoded) { + t.Fatalf("decoded are not equal origin: %v", decoded) + } +} + +func TestMeasuresAppendRecord(t *testing.T) { + rec := MeasuresAppendRecord{ + MetricID: 12345, + TimestampsRewindOffset: 7, + Timestamps: []byte{ + 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, + }, + ValuesRewindOffset: 11, + Values: []byte{ + 0x08, 0x09, 0x0a, + }, + } + buf := bytes.NewBuffer(nil) + rec.Pack(buf) + + recordType, _ := buf.ReadByte() + + if recordType != CodeMeasuresAppend { + t.Fatalf("wrong record type: %d", recordType) + } + + decoded := new(MeasuresAppendRecord) + err := decoded.Parse(buf) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rec, *decoded) { + t.Fatalf("decoded are not equal origin: %v", decoded) + } +} + +func TestMeasuresAppendWithGrowRecord(t *testing.T) { + DataPageSize = 8 + IndexPageSize = 4 + indexRecordSize = 2 + // + rec := MeasuresAppendWithGrowRecord{ + MetricID: 12345, + DataPageTail: DataPageTail{ + PageNo: math.MaxUint32, + Reused: true, + PrevPageNo: 4127831, + CRC32: 316278321, + TimestampsRewindOffset: 7, + Timestamps: []byte{ + 0x07, 0x07, 0x07, 0x07, 0x07, + }, + Values: []byte{ + 0x08, 0x08, 0x08, 0x08, + }, + ValuesRewindOffset: 11, + }, + DataPages: []SealedDataPage{ + { + PageNo: 1000, + Reused: true, + //PrevPageNo: math.MaxUint32, + //CRC32: 312, + Content: []byte{ + 0x05, 0x05, 0x05, 0x05, 0x05, 0x05, 0x05, 0x05, + }, + }, + { + PageNo: 1001, + Reused: false, + //PrevPageNo: 1000, + //CRC32: 432432, + Content: []byte{ + 0x06, 0x06, 0x06, 0x06, 0x06, 0x06, 0x06, 0x06, + }, + }, + }, + TailTimestamps: []byte{ + 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, + }, + TailValues: []byte{ + 0x08, 0x09, 0x0a, + }, + ChangedIndexLevels: []ChangedIndexLevel{ + { + IndexPageTail: &IndexPageTail{ + PageNo: 20000, + Reused: true, + CRC32: 32897, + Records: []byte{ + 0x01, 0x00, 0x00, 0x00, + }, + }, + IndexPages: []SealedIndexPage{ + { + PageNo: 20001, + Reused: true, + //CRC32: 32897, + Content: []byte{ + 0x01, 0x00, 0x00, 0x01, + }, + }, + { + PageNo: 20002, + Reused: false, + //CRC32: 3284397, + Content: []byte{ + 0x01, 0x02, 0x03, 0x04, + }, + }, + }, + TailRecords: []byte{ + 0xaa, 0xbb, + }, + }, + { + IndexPageTail: nil, + TailRecords: []byte{ + 0xdd, 0xee, 0xff, 0x77, + }, + }, + }, + } + buf := bytes.NewBuffer(nil) + rec.Pack(buf) + + recordType, _ := buf.ReadByte() + + if recordType != CodeMeasuresAppendWithGrow { + t.Fatalf("wrong record type: %d", recordType) + } + + decoded := new(MeasuresAppendWithGrowRecord) + err := decoded.Parse(buf) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rec, *decoded) { + t.Fatalf("decoded are not equal origin: %v", decoded) + } +} + +func TestWritePreparer(t *testing.T) { + // data page + DataPageSize = 24 + dataCRC32Idx = DataPageSize - 4 + timestampsSizeIdx = DataPageSize - 6 + valuesSizeIdx = DataPageSize - 8 + prevPageIdx = DataPageSize - 12 + + // index page + IndexPageSize = 16 + indexRecordSize = 4 // 2 records per index page + maxRecordsOnIndexPage = 2 + indexCRC32Idx = IndexPageSize - 4 + indexRecordsCountIdx = IndexPageSize - 6 + isZeroLevelIdx = IndexPageSize - 7 + + pageManager := new(PageNumbersMock) + + pageIngester, err := NewPageIngester(PageIngesterOptions{ + GetIndexPageNumber: pageManager.GetIndexPageNumber, + GetDataPageNumber: pageManager.GetDataPageNumber, + }) + if err != nil { + t.Fatal(err) + } + + writePreparer := NewWritePreparer(WritePreparerOptions{ + MinWALBufferSize: 128, + PageIngester: pageIngester, + }) + + rec := MeasuresAppendWithGrow{ + MetricID: 12345, + LastPageNo: 1, + TimestampsRewindOffset: 2, + Timestamps: []byte{ + 0x04, 0x05, 0x06, + }, + ValuesRewindOffset: 1, + Values: []byte{ + 0x04, 0x03, + }, + DataPages: []DataPayload{ + // 1. fills existent index tail + { + Since: 513, // 0x01 0x02 + // full page + Content: []byte{ + 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x00, 0x00, 0x04, 0x03, 0x02, 0x01, + 0x00, 0x00, 0x00, 0x00, // prevPageNo + 0x00, 0x00, // values size + 0x00, 0x00, // timestamps size + 0x00, 0x00, 0x00, 0x00, // crc32 + }, + TimestampsSize: 6, + ValuesSize: 4, + }, + // 2. creates index page tail (zero level) and level 1 + { + Since: 769, // 0x01 0x03 + // full page + Content: []byte{ + 0xa1, 0xa2, 0xa3, 0xa4, 0xa5, 0x00, 0x00, 0xa5, 0xa4, 0xa3, 0xa2, 0xa1, + 0x00, 0x00, 0x00, 0x00, // prevPageNo + 0x00, 0x00, // values size + 0x00, 0x00, // timestamps size + 0x00, 0x00, 0x00, 0x00, // crc32 + }, + TimestampsSize: 5, + ValuesSize: 5, + }, + // 2. fills index tail records (zero level) + { + Since: 1025, // 0x01 0x04 + // full page + Content: []byte{ + 0xb1, 0xb2, 0xb3, 0xb4, 0xb5, 0xb6, 0xb7, 0xb5, 0xb4, 0xb3, 0xb2, 0xb1, + 0x00, 0x00, 0x00, 0x00, // prevPageNo + 0x00, 0x00, // values size + 0x00, 0x00, // timestamps size + 0x00, 0x00, 0x00, 0x00, // crc32 + }, + TimestampsSize: 7, + ValuesSize: 5, + }, + // 3. creates index page + { + Since: 1025, // 0x01 0x04 + // full page + Content: []byte{ + 0xc1, 0xc2, 0xc3, 0xc4, 0xc5, 0xc6, 0xc6, 0xc5, 0xc4, 0xc3, 0xc2, 0xc1, + 0x00, 0x00, 0x00, 0x00, // prevPageNo + 0x00, 0x00, // values size + 0x00, 0x00, // timestamps size + 0x00, 0x00, 0x00, 0x00, // crc32 + }, + TimestampsSize: 6, + ValuesSize: 6, + }, + }, + TailTimestamps: []byte{ + 0xc1, 0xc2, + }, + TailValues: []byte{ + 0xd1, 0xd2, + }, + IndexLevelTails: []IndexLevelTail{ + { + Buffer: []byte{ + 0x01, 0x01, // since + 0x0a, 0x00, // pageNo 10 + 0x00, 0x00, 0x00, 0x00, 0x00, // empty space + 0x00, // zero level + 0x00, 0x00, // records count + 0x00, 0x00, 0x00, 0x00, // crc32 + + }, + RecordsCount: 1, + }, + }, + WrittenCount: 10, + } + + prepared := writePreparer.Prepare([]any{rec}) + // fmt.Printf("%v\n", prepared) + // pretty.PPrintln("prepared", prepared) + + fmt.Printf("wal: % x\n", prepared.Packet) + + fmt.Println("index to write") + for _, p := range prepared.WriteToIndex { + fmt.Printf("%d: % x\n", p.PageNo, p.Content) + } + fmt.Println("data to write") + for _, p := range prepared.WriteToData { + fmt.Printf("%d: % x\n", p.PageNo, p.Content) + } + + fmt.Println("commits") + for _, x := range prepared.Commits { + fmt.Printf("%#v\n", x) + } +} diff --git a/storage/wal_reader.go b/storage/wal_reader.go index 4aac585..3c67172 100644 --- a/storage/wal_reader.go +++ b/storage/wal_reader.go @@ -159,7 +159,7 @@ func (s *WALReader) parseRecords(body []byte) ([]any, error) { // records = append(records, rec) case CodeMeasuresDelete: - rec := new(DeletedMeasures) + rec := new(MeasuresDeleteRecord) if err = rec.Parse(src); err != nil { return nil, err } @@ -171,184 +171,8 @@ func (s *WALReader) parseRecords(body []byte) ([]any, error) { } } -// func (s *Reader) readAddedMetric(src *bytes.Buffer) (_ AddedMetric, err error) { -// arr, err := bin.ReadN(src, 6) -// if err != nil { -// return -// } -// return AddedMetric{ -// MetricID: bin.GetUint32(arr), -// MetricType: diploma.MetricType(arr[4]), -// FracDigits: int(arr[5]), -// }, nil -// } - -// func (s *Reader) readDeletedMetric(src *bytes.Buffer) (_ DeletedMetric, err error) { -// var rec DeletedMetric -// rec.MetricID, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// // free data pages -// dataQty, _, err := bin.ReadVarUint64(src) -// if err != nil { -// return -// } -// for range dataQty { -// var pageNo uint32 -// pageNo, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// rec.FreeDataPages = append(rec.FreeDataPages, pageNo) -// } -// // free index pages -// indexQty, _, err := bin.ReadVarUint64(src) -// if err != nil { -// return -// } -// for range indexQty { -// var pageNo uint32 -// pageNo, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// rec.FreeIndexPages = append(rec.FreeIndexPages, pageNo) -// } -// return rec, nil -// } - -// func (s *Reader) readAppendedMeasure(src *bytes.Buffer) (_ AppendedMeasure, err error) { -// arr, err := bin.ReadN(src, 16) -// if err != nil { -// return -// } -// return AppendedMeasure{ -// MetricID: bin.GetUint32(arr[0:]), -// Timestamp: bin.GetUint32(arr[4:]), -// Value: bin.GetFloat64(arr[8:]), -// }, nil -// } - -// func (s *Reader) readAppendedMeasures(src *bytes.Buffer) (_ AppendedMeasures, err error) { -// var rec AppendedMeasures -// rec.MetricID, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// qty, err := bin.ReadUint16(src) -// if err != nil { -// return -// } -// for range qty { -// var measure proto.Measure -// measure.Timestamp, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// measure.Value, err = bin.ReadFloat64(src) -// if err != nil { -// return -// } -// rec.Measures = append(rec.Measures, measure) -// } -// return rec, nil -// } - -// func (s *Reader) readAppendedMeasureWithOverflow(src *bytes.Buffer) (_ AppendedMeasureWithOverflow, err error) { -// var ( -// b byte -// rec AppendedMeasureWithOverflow -// ) -// rec.MetricID, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// rec.Timestamp, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// rec.Value, err = bin.ReadFloat64(src) -// if err != nil { -// return -// } -// b, err = src.ReadByte() -// if err != nil { -// return -// } -// rec.IsDataPageReused = b == 1 - -// rec.DataPageNo, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// b, err = src.ReadByte() -// if err != nil { -// return -// } -// if b == 1 { -// rec.IsRootChanged = true -// rec.RootPageNo, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// } -// // index pages -// indexQty, err := src.ReadByte() -// if err != nil { -// return -// } -// for range indexQty { -// var pageNo uint32 -// pageNo, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// rec.ReusedIndexPages = append(rec.ReusedIndexPages, pageNo) -// } -// return rec, nil -// } - -// func (s *Reader) readDeletedMeasures(src *bytes.Buffer) (_ DeletedMeasures, err error) { -// var ( -// rec DeletedMeasures -// ) -// rec.MetricID, err = bin.ReadUint32(src) -// if err != nil { -// return -// } -// // free data pages -// rec.FreeDataPages, err = s.readFreePageNumbers(src) -// if err != nil { -// return -// } -// // free index pages -// rec.FreeIndexPages, err = s.readFreePageNumbers(src) -// if err != nil { -// return -// } -// return rec, nil -// } - // HELPERS -// func (s *Reader) readFreePageNumbers(src *bytes.Buffer) ([]uint32, error) { -// var freePages []uint32 -// qty, _, err := bin.ReadVarUint64(src) -// if err != nil { -// return nil, err -// } -// for range qty { -// var pageNo uint32 -// pageNo, err = bin.ReadUint32(src) -// if err != nil { -// return nil, err -// } -// freePages = append(freePages, pageNo) -// } -// return freePages, nil -// } - func (s *WALReader) Seek(offset int64) error { ret, err := s.file.Seek(offset, 0) if err != nil { diff --git a/storage/wal_records.go b/storage/wal_records.go index c8671d4..899b482 100644 --- a/storage/wal_records.go +++ b/storage/wal_records.go @@ -1,362 +1,12 @@ package storage import ( - "bytes" "io" bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" ) -type HeadDataPage struct { - PageNo uint32 - Reused bool - PrevPageNo uint32 - Checksum uint32 - TimestampsOffset int // offset on 1st page - Timestamps []byte - ValuesOffset int // offset on 1st page - Values []byte -} - -type HeadIndexPage struct { - PageNo uint32 - Reused bool - Checksum uint32 - Records []byte // offset не потрібен, оскільки Records додаються в кінець -} - -type ChangedIndexLevel struct { - HeadIndexPage *HeadIndexPage - IndexPages []SealedIndexPage - TailRecords []byte // payload only -} - -type MeasuresAppendWithGrowRecord struct { - MetricID uint32 - HeadDataPage HeadDataPage - DataPages []SealedDataPage - ChangedIndexLevels []ChangedIndexLevel - TailTimestamps []byte // data level tail - TailValues []byte // data level tail -} - -/* -Format appended measures: -1b - tx type -4b - metricID -4b - pageNo -1b - reused -4b - prevPageNo -4b - page crc32 -Nb - (varsize) timestamps offset -Nb - (varsize) timestamps tail size on 1st filled data page -Nb - timestamps tail on 1st filled data page -Nb - (varsize) values offset -Nb - (varsize) values tail size on 1st filled data page -Nb - values tail on 1st filled data page - -Nb - (varsize) filled data pages qty -[ - - 4b - pageNo - 1b - reused - Nb - (data page size) page content - -] -Nb - (varsize) timestamps payload size on tail -Nb - timestamps payload -Nb - (varsize) values payload size on tail -Nb - values payload -Nb - (varsize) qty of index levels -[ - - // NOTE! - // - idx = 0 - zeroLevel - // - skipped = maxRecordsOnIndexPage - records count - 4b - pageNo of 1st index page - 1b - reused - 4b - page crc32 - Nb - (varsize) records tail count on 1st filled index page - Nb - records tail on 1st filled index page - - - Nb - (varsize) qty of level filled pages - [ - 4b - pageNo - 1b - reused - Nb - (index page size) page content - ] - Nb - (varsize) size of records level tail - Nb - level tail records - -] -*/ -func (s MeasuresAppendWithGrowRecord) WriteTo(w *bytes.Buffer) { - w.WriteByte(CodeMeasuresAppendWithGrow) - bin.WriteUint32(w, s.MetricID) - // completed data page - completed := s.HeadDataPage - bin.WriteUint32(w, completed.PageNo) - bin.WriteBool(w, completed.Reused) - bin.WriteUint32(w, completed.PrevPageNo) - bin.WriteUint32(w, completed.Checksum) - bin.WriteVarSize(w, completed.TimestampsOffset) - bin.WriteVarSize(w, len(completed.Timestamps)) - w.Write(completed.Timestamps) - bin.WriteVarSize(w, completed.ValuesOffset) - bin.WriteVarSize(w, len(completed.Values)) - w.Write(completed.Values) - // data pages - bin.WriteVarSize(w, len(s.DataPages)) - for _, p := range s.DataPages { - bin.WriteUint32(w, p.PageNo) - bin.WriteBool(w, p.Reused) - bin.WriteVarSize(w, len(p.Content)) - w.Write(p.Content) - } - // data tail - bin.WriteVarSize(w, len(s.TailTimestamps)) - w.Write(s.TailTimestamps) - bin.WriteVarSize(w, len(s.TailValues)) - w.Write(s.TailValues) - // changed levels - bin.WriteVarSize(w, len(s.ChangedIndexLevels)) - for _, level := range s.ChangedIndexLevels { - // completed index page - completed := level.HeadIndexPage - bin.WriteUint32(w, completed.PageNo) - bin.WriteBool(w, completed.Reused) - bin.WriteUint32(w, completed.Checksum) - bin.WriteVarSize(w, len(completed.Records)/indexRecordSize) - w.Write(completed.Records) - // full index pages - for _, p := range level.IndexPages { - bin.WriteUint32(w, p.PageNo) - bin.WriteBool(w, p.Reused) - bin.WriteVarSize(w, len(p.Content)) - w.Write(p.Content) - } - // level tail - bin.WriteVarSize(w, len(level.TailRecords)/indexRecordSize) - w.Write(level.TailRecords) - } -} - -func (s *MeasuresAppendWithGrowRecord) Parse(r *bytes.Buffer) (err error) { - var ( - size int - recordsCount int - ) - s.MetricID, err = bin.ReadUint32(r) - if err != nil { - return - } - s.HeadDataPage.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - s.HeadDataPage.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - s.HeadDataPage.PrevPageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - s.HeadDataPage.Checksum, err = bin.ReadUint32(r) - if err != nil { - return - } - s.HeadDataPage.TimestampsOffset, err = bin.ReadVarSize(r) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.HeadDataPage.Timestamps, err = bin.ReadN(r, size) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.HeadDataPage.Values, err = bin.ReadN(r, size) - if err != nil { - return - } - // data pages - dataPagesQty, err := bin.ReadVarSize(r) - if err != nil { - return - } - for range dataPagesQty { - var ( - p SealedDataPage - ) - p.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - p.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - p.Content, err = bin.ReadN(r, size) - if err != nil { - return - } - s.DataPages = append(s.DataPages, p) - } - // data tail - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.TailTimestamps, err = bin.ReadN(r, size) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.TailValues, err = bin.ReadN(r, size) - if err != nil { - return - } - // changed levels - changedIndexLevelsCount, err := bin.ReadVarSize(r) - if err != nil { - return - } - for range changedIndexLevelsCount { - var ( - indexPagesCount int - level ChangedIndexLevel - ) - // completed index page - level.HeadIndexPage.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - level.HeadIndexPage.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - level.HeadIndexPage.Checksum, err = bin.ReadUint32(r) - if err != nil { - return - } - recordsCount, err = bin.ReadVarSize(r) - if err != nil { - return - } - level.HeadIndexPage.Records, err = bin.ReadN(r, recordsCount*indexRecordSize) - if err != nil { - return - } - // - indexPagesCount, err = bin.ReadVarSize(r) - if err != nil { - return - } - for range indexPagesCount { - var p SealedIndexPage - p.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - p.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - p.Checksum, err = bin.ReadUint32(r) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - p.Content, err = bin.ReadN(r, size) - if err != nil { - return - } - level.IndexPages = append(level.IndexPages, p) - } - recordsCount, err = bin.ReadVarSize(r) - if err != nil { - return - } - level.TailRecords, err = bin.ReadN(r, recordsCount*indexRecordSize) - if err != nil { - return - } - s.ChangedIndexLevels = append(s.ChangedIndexLevels, level) - } - return -} - -type MeasuresAppendRecord struct { - MetricID uint32 - TimestampsRewindOffset int - Timestamps []byte - ValuesRewindOffset int - Values []byte -} - -func (s MeasuresAppendRecord) WriteTo(w *bytes.Buffer) { - w.WriteByte(CodeMeasuresAppendWithGrow) - bin.WriteUint32(w, s.MetricID) - bin.WriteVarSize(w, s.TimestampsRewindOffset) - bin.WriteVarSize(w, len(s.Timestamps)) - w.Write(s.Timestamps) - bin.WriteVarSize(w, s.ValuesRewindOffset) - bin.WriteVarSize(w, len(s.Values)) - w.Write(s.Values) -} - -func (s *MeasuresAppendRecord) Parse(r *bytes.Buffer) (err error) { - var size int - s.MetricID, err = bin.ReadUint32(r) - if err != nil { - return - } - s.TimestampsRewindOffset, err = bin.ReadVarSize(r) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.Timestamps, err = bin.ReadN(r, size) - if err != nil { - return - } - s.ValuesRewindOffset, err = bin.ReadVarSize(r) - if err != nil { - return - } - size, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.Values, err = bin.ReadN(r, size) - if err != nil { - return - } - return -} - type MetricAddRecord struct { MetricID uint32 MetricType qb.MetricType @@ -398,11 +48,8 @@ func (s MetricDeleteRecord) Pack(w io.Writer) { } bin.PutUint32(arr[1:], s.MetricID) w.Write(arr) - // fix - // bin.WriteVarSize(w, len(s.FreePageNumbers)) - // for _, pageNo := range s.FreePageNumbers { - // bin.WriteUint32(w, pageNo) - // } + packPageNumbers(w, s.FreeIndexPages) + packPageNumbers(w, s.FreeDataPages) } func (s *MetricDeleteRecord) Parse(src io.Reader) (err error) { @@ -410,61 +57,378 @@ func (s *MetricDeleteRecord) Parse(src io.Reader) (err error) { if err != nil { return } - // free pages - // 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 + s.FreeIndexPages, err = parsePageNumbers(src) + if err != nil { + return + } + s.FreeDataPages, err = parsePageNumbers(src) + if err != nil { + return + } + return +} + +type DataPageTail struct { + PageNo uint32 + Reused bool + PrevPageNo uint32 + CRC32 uint32 + TimestampsRewindOffset int // offset on 1st page + Timestamps []byte + ValuesRewindOffset int // offset on 1st page + Values []byte +} + +type IndexPageTail struct { + PageNo uint32 + Reused bool + CRC32 uint32 + Records []byte // offset не потрібен, оскільки Records додаються в кінець +} + +type ChangedIndexLevel struct { + IndexPageTail *IndexPageTail + IndexPages []SealedIndexPage + TailRecords []byte // payload only +} + +type MeasuresAppendWithGrowRecord struct { + MetricID uint32 + DataPageTail DataPageTail + DataPages []SealedDataPage + TailTimestamps []byte // data level tail + TailValues []byte // data level tail + ChangedIndexLevels []ChangedIndexLevel +} + +/* +Format appended measures: +1b - tx type +4b - metricID +4b - pageNo +1b - reused +4b - prevPageNo +4b - page crc32 +1b - timestamps & values offsets +Nb - (varsize) timestamps tail size on 1st filled data page +Nb - timestamps tail on 1st filled data page +Nb - (varsize) values tail size on 1st filled data page +Nb - values tail on 1st filled data page + +Nb - (varsize) full data pages qty +[ + + 4b - pageNo + 1b - reused + Nb - (data page size) page content + +] +Nb - (varsize) timestamps payload size on tail +Nb - timestamps payload +Nb - (varsize) values payload size on tail +Nb - values payload +Nb - (varsize) qty of index levels +[ + + // NOTE! + // - idx = 0 - zeroLevel + // - skipped = maxRecordsOnIndexPage - records count + 4b - pageNo of 1st index page + 1b - reused + 4b - page crc32 + Nb - (varsize) records tail count on 1st filled index page + Nb - records tail on 1st filled index page + + + Nb - (varsize) qty of level filled pages + [ + 4b - pageNo + 1b - reused + Nb - (index page size) page content + ] + Nb - (varsize) size of records level tail + Nb - level tail records + +] +*/ +func (s MeasuresAppendWithGrowRecord) Pack(w io.Writer) { + arr := []byte{ + CodeMeasuresAppendWithGrow, + 0, 0, 0, 0, // + 0, 0, 0, 0, // tail pageNo + 0, // tail reused + 0, 0, 0, 0, // tail prevPageNo + 0, 0, 0, 0, // tail CRC32 + 0, // tail offsets + } + bin.PutUint32(arr[1:], s.MetricID) + // data page tail + tail := s.DataPageTail + bin.PutUint32(arr[5:], tail.PageNo) + if tail.Reused { + arr[9] = 1 + } + bin.PutUint32(arr[10:], tail.PrevPageNo) + bin.PutUint32(arr[14:], tail.CRC32) + arr[18] = byte(tail.TimestampsRewindOffset) | (byte(tail.ValuesRewindOffset) << 4) + w.Write(arr) + bin.WriteVarSized(w, tail.Timestamps) + bin.WriteVarSized(w, tail.Values) + // data pages + bin.WriteVarSize(w, len(s.DataPages)) + for _, p := range s.DataPages { + bin.WriteUint32(w, p.PageNo) + bin.WriteBool(w, p.Reused) + w.Write(p.Content) + } + // data tail + bin.WriteVarSized(w, s.TailTimestamps) + bin.WriteVarSized(w, s.TailValues) + // changed levels + bin.WriteVarSize(w, len(s.ChangedIndexLevels)) + for _, level := range s.ChangedIndexLevels { + // index page tail + tail := level.IndexPageTail + if tail != nil { + bin.WriteBool(w, true) + bin.WriteUint32(w, tail.PageNo) + bin.WriteBool(w, tail.Reused) + bin.WriteUint32(w, tail.CRC32) + bin.WriteVarSize(w, len(tail.Records)/indexRecordSize) + w.Write(tail.Records) + } else { + bin.WriteBool(w, false) + } + // full index pages + bin.WriteVarSize(w, len(level.IndexPages)) + for _, p := range level.IndexPages { + bin.WriteUint32(w, p.PageNo) + bin.WriteBool(w, p.Reused) + w.Write(p.Content) + } + // level tail + bin.WriteVarSize(w, len(level.TailRecords)/indexRecordSize) + w.Write(level.TailRecords) + } +} + +func (s *MeasuresAppendWithGrowRecord) Parse(r io.Reader) (err error) { + s.MetricID, err = bin.ReadUint32(r) + if err != nil { + return + } + s.DataPageTail.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.DataPageTail.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + s.DataPageTail.PrevPageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.DataPageTail.CRC32, err = bin.ReadUint32(r) + if err != nil { + return + } + offsets, err := bin.ReadByte(r) + if err != nil { + return + } + s.DataPageTail.TimestampsRewindOffset = int(offsets & 15) // low 4 bits + s.DataPageTail.ValuesRewindOffset = int(offsets >> 4) // high 4 bits + // + s.DataPageTail.Timestamps, err = bin.ReadVarSized(r) + if err != nil { + return + } + s.DataPageTail.Values, err = bin.ReadVarSized(r) + if err != nil { + return + } + // data pages + dataPagesQty, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range dataPagesQty { + var p SealedDataPage + p.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + p.Content, err = bin.ReadN(r, DataPageSize) + if err != nil { + return + } + s.DataPages = append(s.DataPages, p) + } + // // data tail + s.TailTimestamps, err = bin.ReadVarSized(r) + if err != nil { + return + } + s.TailValues, err = bin.ReadVarSized(r) + if err != nil { + return + } + // changed levels + changedIndexLevelsCount, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range changedIndexLevelsCount { + var ( + hasIndexPageTail bool + recordsCount int + indexPagesCount int + level ChangedIndexLevel + ) + hasIndexPageTail, err = bin.ReadBool(r) + if err != nil { + return + } + if hasIndexPageTail { + //index page tail + var tail IndexPageTail + tail.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + tail.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + tail.CRC32, err = bin.ReadUint32(r) + if err != nil { + return + } + recordsCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + tail.Records, err = bin.ReadN(r, recordsCount*indexRecordSize) + if err != nil { + return + } + level.IndexPageTail = &tail + } + // + indexPagesCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + for range indexPagesCount { + var p SealedIndexPage + p.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + p.Content, err = bin.ReadN(r, IndexPageSize) + if err != nil { + return + } + level.IndexPages = append(level.IndexPages, p) + } + recordsCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + level.TailRecords, err = bin.ReadN(r, recordsCount*indexRecordSize) + if err != nil { + return + } + s.ChangedIndexLevels = append(s.ChangedIndexLevels, level) + } + return +} + +type MeasuresAppendRecord struct { + MetricID uint32 + TimestampsRewindOffset int + Timestamps []byte + ValuesRewindOffset int + Values []byte +} + +func (s MeasuresAppendRecord) Pack(w io.Writer) { + arr := []byte{ + CodeMeasuresAppend, + 0, 0, 0, 0, // metricID + 0, // offsets + } + bin.PutUint32(arr[1:], s.MetricID) + arr[5] = byte(s.TimestampsRewindOffset) | (byte(s.ValuesRewindOffset) << 4) + w.Write(arr) + bin.WriteVarSized(w, s.Timestamps) + bin.WriteVarSized(w, s.Values) +} + +func (s *MeasuresAppendRecord) Parse(r io.Reader) (err error) { + s.MetricID, err = bin.ReadUint32(r) + if err != nil { + return + } + offsets, err := bin.ReadByte(r) + if err != nil { + return + } + s.TimestampsRewindOffset = int(offsets & 15) // low 4 bits + s.ValuesRewindOffset = int(offsets >> 4) // high 4 bits + // + s.Timestamps, err = bin.ReadVarSized(r) + if err != nil { + return + } + s.Values, err = bin.ReadVarSized(r) + if err != nil { + return + } + return } // fix - add changed pages -type DeletedMeasures struct { - MetricID uint32 - FreePageNumbers []uint32 +type MeasuresDeleteRecord struct { + MetricID uint32 + FreeIndexPages []uint32 + FreeDataPages []uint32 } -func (s DeletedMeasures) Pack(w io.Writer) { +func (s MeasuresDeleteRecord) Pack(w io.Writer) { arr := []byte{ CodeMeasuresDelete, 0, 0, 0, 0, } bin.PutUint32(arr[1:], s.MetricID) w.Write(arr) - bin.WriteVarSize(w, len(s.FreePageNumbers)) - for _, pageNo := range s.FreePageNumbers { - bin.WriteUint32(w, pageNo) - } + packPageNumbers(w, s.FreeIndexPages) + packPageNumbers(w, s.FreeDataPages) } -func (s *DeletedMeasures) Parse(src io.Reader) (err error) { +func (s *MeasuresDeleteRecord) Parse(src io.Reader) (err error) { s.MetricID, err = bin.ReadUint32(src) if err != nil { return } - // free pages - qty, err := bin.ReadVarSize(src) + s.FreeIndexPages, err = parsePageNumbers(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) + s.FreeDataPages, err = parsePageNumbers(src) + if err != nil { + return } - return nil + return } // type AppendedMeasure struct { @@ -664,3 +628,28 @@ Nb - values payload // } // return nil // } + +// HELPERS + +func packPageNumbers(w io.Writer, pageNumbers []uint32) { + bin.WriteVarSize(w, len(pageNumbers)) + for _, pageNo := range pageNumbers { + bin.WriteUint32(w, pageNo) + } +} + +func parsePageNumbers(src io.Reader) (pageNumbers []uint32, err error) { + count, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range count { + var pageNo uint32 + pageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + pageNumbers = append(pageNumbers, pageNo) + } + return +} diff --git a/storage/wal_writer.go b/storage/wal_writer.go index 8b674cf..bb51301 100644 --- a/storage/wal_writer.go +++ b/storage/wal_writer.go @@ -34,7 +34,7 @@ const ( CodeMetricAdd byte = 1 CodeMetricDelete byte = 2 CodeMeasuresAppendWithGrow byte = 5 - CodeMeasuresAppend = 6 + CodeMeasuresAppend byte = 6 CodeMeasuresDelete byte = 7 ) @@ -77,7 +77,7 @@ type Writer struct { indexPageSize int dataPageSize int input []any - dataPreparer *DataPreparer + dataPreparer *WritePreparer appendToWorkerQueue func(any) //lsn uint32 written int64 @@ -259,7 +259,7 @@ func (s *Writer) packAndWrite() (err error) { snapshotNumberCh := make(chan int, 1) s.appendToWorkerQueue(Changes{ - Records: prepared.WriteResults, + Records: prepared.Commits, SnapshotNumberCh: snapshotNumberCh, FrozenIndexPagesCount: s.indexFreeList.Pages(), IndexPageNumbers: s.indexFreeList.Cached(), @@ -290,7 +290,7 @@ func (s *Writer) packAndWrite() (err error) { } } else { s.appendToWorkerQueue(Changes{ - Records: prepared.WriteResults, + Records: prepared.Commits, }) } return nil diff --git a/storage/write_preparer.go b/storage/write_preparer.go new file mode 100644 index 0000000..3d21c2e --- /dev/null +++ b/storage/write_preparer.go @@ -0,0 +1,258 @@ +package storage + +import ( + "bytes" + "fmt" + "io" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/util" +) + +type WritePreparer struct { + pagePreparer *PageIngester + //dataPagePayloadBound int + w *bytes.Buffer + indexPages []PageToWrite + dataPages []PageToWrite + commits []any +} + +type WritePreparerOptions struct { + MinWALBufferSize int + PageIngester *PageIngester +} + +func NewWritePreparer(opt WritePreparerOptions) *WritePreparer { + //buf := make([]byte, opt.MinWALBufferSize) + s := &WritePreparer{ + pagePreparer: opt.PageIngester, + w: bytes.NewBuffer(nil), + } + return s +} + +type PreparedData struct { + Packet []byte // пакет для запису в WAL + WriteToIndex []PageToWrite + WriteToData []PageToWrite + Commits []any +} + +func (s *WritePreparer) Prepare(input []any) PreparedData { + s.w.Write([]byte{ + 0, 0, 0, 0, 0, 0, 0, 0, 0, // size (max 9 byte) + 0, 0, 0, 0, // crc32 + }) + + hasher := util.NewHasher() + + w := io.MultiWriter(s.w, hasher) + + // 1. Пакую всі дані в WAL буфер запису + for _, untyped := range input { + switch x := untyped.(type) { + case MetricAdd: + rec := MetricAddRecord{ + MetricID: x.MetricID, + MetricType: x.MetricType, + FracDigits: x.FracDigits, + } + rec.Pack(w) + s.commits = append(s.commits, MetricAddCommited{ + MetricID: x.MetricID, + ResultCh: x.ResultCh, + }) + + case MetricDelete: + // if len(x.FreeIndexPages) > 0 { + // s.indexFreeList.AddPageNumbers(x.FreeIndexPages) + // } + // if len(x.FreeIndexPages) > 0 { + // s.dataFreeList.AddPageNumbers(x.FreeDataPages) + // } + rec := MetricDeleteRecord{ + MetricID: x.MetricID, + FreeIndexPages: x.FreeIndexPages, + FreeDataPages: x.FreeDataPages, + } + rec.Pack(w) + s.commits = append(s.commits, MetricDeleteCommited{ + MetricID: x.MetricID, + ResultCh: x.ResultCh, + }) + + case MeasuresAppend: + rec := MeasuresAppendRecord{ + MetricID: x.MetricID, + TimestampsRewindOffset: x.TimestampsRewindOffset, + Timestamps: x.Timestamps, + ValuesRewindOffset: x.ValuesRewindOffset, + Values: x.Values, + } + rec.Pack(w) + // Дані для Worker + s.commits = append(s.commits, MeasuresAppendCommited{ + MetricID: x.MetricID, + ResultCode: x.ResultCode, + WrittenCount: x.WrittenCount, + ResultCh: x.ResultCh, + }) + + case MeasuresAppendWithGrow: + sealedDataPages, sealedIndexLevels := s.pagePreparer.Ingest(IngestIn{ + LastPageNo: x.LastPageNo, + IndexLevelTails: x.IndexLevelTails, + DataPages: x.DataPages, + }) + + // сторінки для запису в index та data файли + for _, p := range sealedDataPages { + s.dataPages = append(s.dataPages, PageToWrite{ + PageNo: p.PageNo, + Content: p.Content, + }) + } + for _, level := range sealedIndexLevels { + for _, p := range level.IndexPages { + fmt.Printf("records: %v\n", getIndexRecords(p.Content, 2)) + s.indexPages = append(s.indexPages, PageToWrite{ + PageNo: p.PageNo, + Content: p.Content, + }) + } + } + + var ( + // дані для запису в WAL + changedIndexPages []ChangedIndexLevel + // дані для Воркера + indexLevelTails []IndexLevelTail + ) + + for _, level := range sealedIndexLevels { + if len(level.IndexPages) > 0 || level.SkipRecords != level.TailRecordsCount { + // на індексному внесені зміни + var ( + indexPageTail *IndexPageTail + indexPages []SealedIndexPage + ) + if len(level.IndexPages) > 0 { + first := level.IndexPages[0] + skipSize := level.SkipRecords * indexRecordSize + indexPageTail = &IndexPageTail{ + PageNo: first.PageNo, + Reused: first.Reused, + // в WAL файл попадає лише payload + Records: first.Content[skipSize : maxRecordsOnIndexPage*indexRecordSize], + } + indexPages = level.IndexPages[1:] + } else { + indexPages = level.IndexPages + } + changedIndexPages = append(changedIndexPages, ChangedIndexLevel{ + IndexPageTail: indexPageTail, + IndexPages: indexPages, + // в WAL файл попадає лише payload + TailRecords: level.TailRecords[:level.TailRecordsCount*indexRecordSize], + }) + } + // Worker-у відправляю повний index + indexLevelTails = append(indexLevelTails, IndexLevelTail{ + // Worker отримує повний буфер + Buffer: level.TailRecords, + RecordsCount: level.TailRecordsCount, + }) + } + + first := sealedDataPages[0] + firstCRC32, _ := bin.GetUint32(first.Content[dataCRC32Idx:]) + + rec := MeasuresAppendWithGrowRecord{ + MetricID: x.MetricID, + DataPageTail: DataPageTail{ + PageNo: first.PageNo, + Reused: first.Reused, + PrevPageNo: x.LastPageNo, + CRC32: firstCRC32, + TimestampsRewindOffset: x.TimestampsRewindOffset, + Timestamps: x.Timestamps, + ValuesRewindOffset: x.ValuesRewindOffset, + Values: x.Values, + }, + DataPages: sealedDataPages[1:], + ChangedIndexLevels: changedIndexPages, + TailTimestamps: x.Timestamps, + TailValues: x.Values, + } + + rec.Pack(w) + // Дані для Worker + s.commits = append(s.commits, MeasuresAppendWithGrowCommited{ + MetricID: x.MetricID, + LastPageNo: sealedDataPages[len(sealedDataPages)-1].PageNo, + Index: indexLevelTails, + ResultCode: x.ResultCode, + WrittenCount: x.WrittenCount, + ResultCh: x.ResultCh, + }) + + case MeasuresDelete: + // if len(x.FreeIndexPages) > 0 { + // s.indexFreeList.AddPageNumbers(x.FreeIndexPages) + // } + // if len(x.FreeIndexPages) > 0 { + // s.dataFreeList.AddPageNumbers(x.FreeDataPages) + // } + rec := MeasuresDeleteRecord{ + MetricID: x.MetricID, + FreeIndexPages: x.FreeIndexPages, + FreeDataPages: x.FreeDataPages, + } + rec.Pack(w) + s.commits = append(s.commits, MeasuresDeleteCommited{ + MetricID: x.MetricID, + ResultCh: x.ResultCh, + }) + //case DeletedMeasuresSince: + } + } + // 2. Додаю розмір пакету і чексуму + + //s.written += int64(len(packet)) + 12 + + packet := s.w.Bytes() + // write size + payloadSize := len(packet) - 13 + start := 9 - bin.CountVarSize(payloadSize) + bin.PutVarSize(packet[start:], payloadSize) + // CRC32 + bin.PutUint32(packet[9:], hasher.Sum32()) + // + return PreparedData{ + Packet: packet[start:], + WriteToIndex: s.indexPages, + WriteToData: s.dataPages, + Commits: s.commits, + } +} + +func (s *WritePreparer) Reset() { + s.w.Reset() + s.indexPages = nil + s.dataPages = nil + s.commits = nil +} + +func getIndexRecords(buf []byte, count int) (list []IndexRecord) { + fmt.Printf("ipage: % x\n", buf) + i := 0 + for range count { + var rec IndexRecord + rec.Timestamp, _ = bin.GetUint32(buf[i:]) + rec.PageNo, _ = bin.GetUint32(buf[i+4:]) + list = append(list, rec) + i += indexRecordSize + } + return +} diff --git a/util/until.go b/util/until.go index 4c824eb..a143f52 100644 --- a/util/until.go +++ b/util/until.go @@ -9,7 +9,7 @@ var ( castagnoliTable = crc32.MakeTable(crc32.Castagnoli) ) -func CalcChecksum(page []byte) uint32 { +func CalculateCRC32(page []byte) uint32 { return crc32.Checksum(page, castagnoliTable) }