diff --git a/database/database.go b/database/database.go index 96d37ff..7777e46 100644 --- a/database/database.go +++ b/database/database.go @@ -3,7 +3,6 @@ package database import ( "errors" "fmt" - "hash/crc32" "io" "log" "net" @@ -12,8 +11,9 @@ import ( "sync" "time" + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/freelist" "gordenko.dev/dima/qb/txlog" ) @@ -35,17 +35,18 @@ type Database struct { rLocksToRelease []uint32 metrics map[uint32]*_metric //metricLockEntries map[uint32]*metricLockEntry - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList - dir string - databaseName string - txlog *txlog.Writer - atree *atree.Atree - tcpPort int - logfile *os.File - logger *log.Logger - exitCh chan struct{} - waitGroup *sync.WaitGroup + dataFreeList *freelist.FreeList + indexFreeList *freelist.FreeList + dir string + databaseName string + snapshotNumber int + txlog *txlog.Writer + atree *atree.Atree + tcpPort int + logfile *os.File + logger *log.Logger + exitCh chan struct{} + waitGroup *sync.WaitGroup } type Options struct { @@ -237,107 +238,9 @@ func (s *Database) recovery() { // } } -// func (s *Database) searchREDOFiles() ([]string, error) { -// var ( -// reREDO = regexp.MustCompile(`a\d+\.redo`) -// fileNames []string -// ) +// -// entries, err := os.ReadDir(s.redoDir) -// if err != nil { -// return nil, err -// } - -// for _, entry := range entries { -// if entry.Type().IsRegular() { -// baseName := entry.Name() -// if reREDO.MatchString(baseName) { -// fileNames = append(fileNames, filepath.Join(s.redoDir, baseName)) -// } -// } -// } -// return fileNames, nil -// } - -// func (s *Database) replayREDOFile(fileName string) error { -// redoFile, err := redo.ReadREDOFile(redo.ReadREDOFileReq{ -// FileName: fileName, -// DataPageSize: atree.DataPageSize, -// IndexPageSize: atree.IndexPageSize, -// }) -// if err != nil { -// return fmt.Errorf("can't read REDO file %s: %s", fileName, err) -// } - -// metric, ok := s.metrics[redoFile.MetricID] -// if !ok { -// return fmt.Errorf("has REDOFile, metric %d not found", redoFile.MetricID) -// } - -// if metric.Until < redoFile.Timestamp { -// waitCh := make(chan struct{}) -// s.atree.ApplyREDO(atree.WriteTask{ -// DataPage: redoFile.DataPage, -// IndexPages: redoFile.IndexPages, -// }) -// <-waitCh - -// waitCh = s.txlog.WriteAppendedMeasureWithOverflow( -// txlog.AppendedMeasureWithOverflow{ -// MetricID: redoFile.MetricID, -// Timestamp: redoFile.Timestamp, -// Value: redoFile.Value, -// IsDataPageReused: redoFile.IsDataPageReused, -// DataPageNo: redoFile.DataPage.PageNo, -// IsRootChanged: redoFile.IsRootChanged, -// RootPageNo: redoFile.RootPageNo, -// ReusedIndexPages: redoFile.ReusedIndexPages, -// }, -// fileName, -// false, -// ) -// <-waitCh -// } -// return nil -// } - -func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { - file, err := os.Open(fileName) - if err != nil { - return - } - defer file.Close() - - stat, err := file.Stat() - if err != nil { - return - } - - if stat.Size() <= 4 { - return false, nil - } - - var ( - payloadSize = stat.Size() - 4 - hash = crc32.NewIEEE() - ) - - _, err = io.CopyN(hash, file, payloadSize) - if err != nil { - return - } - calculatedCRC := hash.Sum32() - - storedCRC, err := bin.ReadUint32(file) - if err != nil { - return - } - if storedCRC != calculatedCRC { - return false, fmt.Errorf("strored CRC %d not equal calculated CRC %d", - storedCRC, calculatedCRC) - } - return true, nil -} +// зробити object? func (s *Database) replayChanges(fileName string) error { walReader, err := txlog.NewReader(txlog.ReaderOptions{ @@ -348,120 +251,391 @@ func (s *Database) replayChanges(fileName string) error { return err } + var metrics map[uint32]*ReplayMetric + for { - lsn, records, done, err := walReader.ReadPacket() + records, isLast, err := walReader.NextPacket() if err != nil { return err } - _ = lsn - if done { - return nil + if isLast { + if len(records) > 0 { + var appendRecords []txlog.WALRecordAppendMeasures + for _, untyped := range records { + rec, ok := untyped.(txlog.WALRecordAppendMeasures) + if ok { + appendRecords = append(appendRecords, rec) + } + } + // prepare pages and rewrite + indexPages, dataPages := pagesToWriteFromWALRecords(metrics, appendRecords) + + err = s.txlog.WritePagesToAtree(indexPages, dataPages) + if err != nil { + return err + } + } + // для останнього rec - повторюємо запис сторінок + // rewrite write pages + //s.txlog. + //return nil } for _, record := range records { - if err = s.replayChangesRecord(record); err != nil { + if err = s.replayWALRecord(record); err != nil { return err } } } } -func (s *Database) replayChangesRecord(untyped any) error { - // switch rec := untyped.(type) { - // case txlog.AddedMetric: - // var ( - // values qb.ValueCompressor - // timestampsBuf = conbuf.New(nil) - // valuesBuf = conbuf.New(nil) - // ) +func (s *Database) replayWALRecord(untyped any) error { + switch rec := untyped.(type) { + case txlog.AddedMetric: + s.addMetric(rec) - // if rec.MetricType == qb.Cumulative { - // values = chunkenc.NewReverseCumulativeDeltaCompressor( - // valuesBuf, 0, byte(rec.FracDigits)) - // } else { - // values = chunkenc.NewReverseInstantDeltaCompressor( - // valuesBuf, 0, byte(rec.FracDigits)) + case txlog.DeletedMetric: + // pages add and delete to tmp file + + // delete(s.metrics, rec.MetricID) + // if len(rec.FreePageNumbers) > 0 { + // s.freeList.AddPages(rec.FreePageNumbers) + // } + + // case txlog.AppendedMeasure: + // metric, ok := s.metrics[rec.MetricID] + // if ok { + // metric.Timestamps.Append(rec.Timestamp) + // metric.Values.Append(rec.Value) + + // if metric.Since == 0 { + // metric.Since = rec.Timestamp + // metric.SinceValue = rec.Value // } - // s.metrics[rec.MetricID] = &_metric{ - // MetricType: rec.MetricType, - // FracDigits: byte(rec.FracDigits), - // TimestampsBuf: timestampsBuf, - // ValuesBuf: valuesBuf, - // Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), - // Values: values, - // } + // metric.Until = rec.Timestamp + // metric.UntilValue = rec.Value + // } - // case txlog.DeletedMetric: - // delete(s.metrics, rec.MetricID) + case txlog.WALRecordAppendMeasures: + metric, ok := s.metrics[rec.MetricID] + if !ok { + qb.Abort(qb.MetricNotFoundDuringReplay, nil) + } + //metric.ReplayAppendedMeasures(rec) FIX + + // case txlog.DeletedMeasures: + // metric, ok := s.metrics[rec.MetricID] + // if ok { + // metric.DeleteMeasures() // if len(rec.FreePageNumbers) > 0 { // s.freeList.AddPages(rec.FreePageNumbers) // } + // } - // case txlog.AppendedMeasure: + default: + qb.Abort(qb.UnknownTxLogRecordTypeBug, + fmt.Errorf("bug: unknown record type %T in TransactionLog", rec)) + } + return nil +} + +type ReplayMetric struct { + metricType qb.MetricType + fracDigits byte + lastPageNo uint32 + buf []byte + databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер) + tSize int + vSize int + indexLevelTails []txlog.IndexLevelTail // root - last element +} + +func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { + metricType, err := bin.ReadByte(r) + if err != nil { + return + } + s.metricType = qb.MetricType(metricType) + s.fracDigits, err = bin.ReadByte(r) + if err != nil { + return + } + s.lastPageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.tSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + s.vSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + s.buf, s.databuf = allocateBuffers(s.tSize + s.vSize) + // + err = bin.ReadNInto(r, s.databuf[len(s.databuf)-s.tSize:]) + if err != nil { + return + } + err = bin.ReadNInto(r, s.databuf[:s.vSize]) + if err != nil { + return + } + // + levelsCount, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range levelsCount { + var ( + buf = make([]byte, atree.IndexPageSize) + count int + ) + count, err = bin.ReadVarSize(r) + if err != nil { + return + } + err = bin.ReadNInto(r, buf[:count*indexRecordSize]) + if err != nil { + return + } + s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ + Buffer: buf, + RecordsCount: count, + }) + } + return +} + +func (s *ReplayMetric) ReplayAppendedMeasures(rec txlog.WALRecordAppendMeasures) { + // Додати в index level tails недостаючі дані, або замінити + for levelIdx, change := range rec.ChangedIndexLevels { + if levelIdx < len(s.indexLevelTails) { + level := s.indexLevelTails[levelIdx] + if change.CompletedIndexPage != nil { + // попередній хвіст індексного рівня перетворився на сторінку, + // отже TailRecords - це новий хвіст + // або копіюю дані, або алокую більший буфер? + copy(level.Buffer, change.TailRecords) + level.RecordsCount = len(change.TailRecords) / indexRecordSize + } else { + // TailRecords - це нові дані, які треба додати + pos := level.RecordsCount * indexRecordSize + copy(level.Buffer[pos:], change.TailRecords) + } + s.indexLevelTails[levelIdx] = level + } else { + var ( + recordsCount = len(change.TailRecords) / indexRecordSize + buffer = make([]byte, atree.IndexPageSize) + ) + copy(buffer, change.TailRecords) + s.indexLevelTails = append(s.indexLevelTails, txlog.IndexLevelTail{ + Buffer: buffer, + RecordsCount: recordsCount, + }) + } + } + + // Перезаписати в buffer timestamps та values. Якщо розмір буферу недостатній - + // спочатку створити більший + + var ( + // CompletedDataPage є полюбому, оскільки це Record із мінімум одною заповненою + // data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня. + + requiredSpace = len(rec.TailTimestamps) + len(rec.TailValues) + databuf []byte + ) + if requiredSpace < len(s.buf) { + databuf = s.buf + // fix - update sizes? + } else { + var buf []byte + buf, databuf = allocateBuffers(requiredSpace) + copy(databuf, rec.TailValues) + pos := len(databuf) - len(rec.TailTimestamps) + copy(databuf[pos:], rec.TailTimestamps) + // fix - replace and remeber sizes? + + s.buf = buf + } + + //rec.TailTimestamps + //rec.TailValues +} + +func (s *Database) rewritePagesFromLastWALPacket(list []any) error { + // var ( + // indexPages []txlog.PageToWrite + // dataPages []txlog.PageToWrite + // ) + // for _, untyped := range list { + // switch rec := untyped.(type) { + // case txlog.WALRecordAppendMeasures: // metric, ok := s.metrics[rec.MetricID] - // if ok { - // metric.Timestamps.Append(rec.Timestamp) - // metric.Values.Append(rec.Value) - - // if metric.Since == 0 { - // metric.Since = rec.Timestamp - // metric.SinceValue = rec.Value - // } - - // metric.Until = rec.Timestamp - // metric.UntilValue = rec.Value - // } - - // case txlog.AppendedMeasures: - // metric, ok := s.metrics[rec.MetricID] - // if ok { - // for _, measure := range rec.Measures { - // metric.Timestamps.Append(measure.Timestamp) - // metric.Values.Append(measure.Value) - - // if metric.Since == 0 { - // metric.Since = measure.Timestamp - // metric.SinceValue = measure.Value - // } - - // metric.Until = measure.Timestamp - // metric.UntilValue = measure.Value - // } - // } - - // // case txlog.AppendedMeasureWithOverflow: - // // metric, ok := s.metrics[rec.MetricID] - // // if ok { - // // metric.ReinitBy(rec.Timestamp, rec.Value) - // // if rec.IsRootChanged { - // // metric.RootPageNo = rec.RootPageNo - // // } - // // metric.LastPageNo = rec.DataPageNo - // // // delete free pages - // // if rec.IsDataPageReused { - // // s.freeList.DeleteReservedPages([]uint32{ - // // rec.DataPageNo, - // // }) - // // } - // // if len(rec.ReusedIndexPages) > 0 { - // // s.freeList.DeleteReservedPages(rec.ReusedIndexPages) - // // } - // // } - - // case txlog.DeletedMeasures: - // metric, ok := s.metrics[rec.MetricID] - // if ok { - // metric.DeleteMeasures() - // if len(rec.FreePageNumbers) > 0 { - // s.freeList.AddPages(rec.FreePageNumbers) - // } + // if !ok { + // qb.Abort(qb.MetricNotFoundDuringReplay, nil) // } + // //metric.ReplayAppendedMeasures(rec) // default: // qb.Abort(qb.UnknownTxLogRecordTypeBug, // fmt.Errorf("bug: unknown record type %T in TransactionLog", rec)) // } + // } return nil } + +func pagesToWriteFromWALRecords(metrics map[uint32]*ReplayMetric, records []txlog.WALRecordAppendMeasures) (indexPages []txlog.PageToWrite, dataPages []txlog.PageToWrite) { + for _, rec := range records { + metric, ok := metrics[rec.MetricID] + if !ok { + qb.Abort(qb.MetricNotFoundDuringReplay, nil) + } + p := rec.CompletedDataPage + // буфер має бути розміру сторінки + if len(metric.buf) != atree.DataPageSize { + metric.buf, metric.databuf = growBuffers(growBuffersIn{ + databuf: metric.databuf, + tSize: metric.tSize, + vSize: metric.vSize, + requiredSpace: atree.DataPageSize, + }) + } + // дописую дані по зміщенню + pos := len(metric.databuf) - p.TimestampsOffset - len(p.Timestamps) + copy(metric.databuf[pos:], p.Timestamps) + copy(metric.databuf[p.ValuesOffset:], p.Values) + // + metric.tSize = p.TimestampsOffset + len(p.Timestamps) + metric.tSize = p.ValuesOffset + len(p.Values) + + // fix - заполнить страницу + // fix - check checksum + + dataPages = append(dataPages, txlog.PageToWrite{ + PageNo: p.PageNo, + Content: metric.buf, + }) + + for _, x := range rec.DataPages { + dataPages = append(dataPages, txlog.PageToWrite{ + PageNo: x.PageNo, + Content: x.Content, + }) + } + + for levelIdx, changedLevel := range rec.ChangedIndexLevels { + p := changedLevel.CompletedIndexPage + if p != nil { + level := metric.indexLevelTails[levelIdx] + copy(level.Buffer[level.RecordsCount*indexRecordSize:], p.Records) + + // fix - заполнить страницу + // fix - check checksum + + indexPages = append(dataPages, txlog.PageToWrite{ + PageNo: p.PageNo, + Content: level.Buffer, + }) + } + + for _, x := range changedLevel.IndexPages { + indexPages = append(dataPages, txlog.PageToWrite{ + PageNo: x.PageNo, + Content: x.Content, + }) + } + } + } + return +} + +// FIX +// УВАГА! +// Якщо буфер metric.buffer < розміру сторінки - для timestamps доступний весь буфер. +// Якщо буфер досяг розміру сторінки - в timestamps я передаю buffer[:DataPagePayloadSize] + +var dataBufferSizes = []int{ + 1024, + 2048, + 4096, + 8192, +} + +func calculateDataBufferSize(payloadSize int) int { + for _, bufferSize := range dataBufferSizes { + if payloadSize <= bufferSize { + return bufferSize + } + } + return atree.DataPageSize +} + +func allocateBuffers(requiredSpace int) (buf, databuf []byte) { + bufferSize := calculateDataBufferSize(requiredSpace) + buf = make([]byte, bufferSize) + if bufferSize == atree.DataPageSize { + databuf = buf[:atree.DataPagePayloadSize] + } else { + databuf = buf + } + return +} + +// src - databuf +type growBuffersIn struct { + databuf []byte + tSize int + vSize int + requiredSpace int +} + +func growBuffers(in growBuffersIn) (buf, databuf []byte) { + buf, databuf = allocateBuffers(in.requiredSpace) + copy(databuf[len(databuf)-in.tSize:], in.databuf[len(in.databuf)-in.tSize:]) + copy(databuf, in.databuf[:in.vSize]) + return +} + +//func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { +// file, err := os.Open(fileName) +// if err != nil { +// return +// } +// defer file.Close() + +// stat, err := file.Stat() +// if err != nil { +// return +// } + +// if stat.Size() <= 4 { +// return false, nil +// } + +// var ( +// payloadSize = stat.Size() - 4 +// hash = crc32.NewIEEE() +// ) + +// _, err = io.CopyN(hash, file, payloadSize) +// if err != nil { +// return +// } +// calculatedCRC := hash.Sum32() + +// storedCRC, err := bin.ReadUint32(file) +// if err != nil { +// return +// } +// if storedCRC != calculatedCRC { +// return false, fmt.Errorf("strored CRC %d not equal calculated CRC %d", +// storedCRC, calculatedCRC) +// } +// return true, nil +// } diff --git a/database/metric.go b/database/metric.go index ec87527..a5329d4 100644 --- a/database/metric.go +++ b/database/metric.go @@ -1,6 +1,9 @@ package database import ( + "io" + + bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/txlog" @@ -10,26 +13,25 @@ import ( const minBufferSize = 1024 -type IndexLevelTail struct { - Payload []byte - RecordsCount int -} +var ( + indexRecordSize = 8 +) type _metric struct { - MetricType qb.MetricType - FracDigits byte - LastPageNo uint32 - SinceValue float64 - Since uint32 - UntilValue float64 - Until uint32 - Buffer []byte - Timestamps qb.TimestampCompressor - Values qb.ValueCompressor + metricType qb.MetricType + fracDigits byte + lastPageNo uint32 + //SinceValue float64 + //Since uint32 + lastValue float64 + //Until uint32 + buffer []byte + timestamps qb.TimestampCompressor + values qb.ValueCompressor XLock bool RLocks int WaitQueue []any - IndexLevelTails []txlog.IndexLevelTail // root - last element + indexLevelTails []txlog.IndexLevelTail // root - last element } //IndexLevels [][]IndexRec // root - last element @@ -73,96 +75,105 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu // DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page // } - timestamps = s.Timestamps - values = s.Values + timestamps = s.timestamps + values = s.values - indexLevels []txlog.IndexLevelTail - dataPages []txlog.DataPayload + dataPages []txlog.DataPayload //written int //resultCode byte ) - s.Values.CaptureState() - s.Timestamps.CaptureState() + s.values.CaptureState() + s.timestamps.CaptureState() for idx, measure := range req.Measures { - if s.Since == 0 { - s.Since = measure.Timestamp - } else { - // FIX - у випадку помилки треба в транзакції зберегти що помилка, але також зафіксувати скільки елементів збережено. - // якщо idx == 0 - одразу знімаю блокування і нічого не відправляю в txlog - if measure.Timestamp <= s.Until { - if idx == 0 { - s.Values.ForgetCapturedState() - s.Timestamps.ForgetCapturedState() + // if s.Since == 0 { + // s.Since = measure.Timestamp + // } else { + // FIX - у випадку помилки треба в транзакції зберегти що помилка, але також зафіксувати скільки елементів збережено. + // якщо idx == 0 - одразу знімаю блокування і нічого не відправляю в txlog + if measure.Timestamp <= s.timestamps.LastTimestamp() { + if idx == 0 { + s.values.ForgetCapturedState() + s.timestamps.ForgetCapturedState() - req.ResultCh <- tryAppendMeasuresResult{ - ResultCode: ExpiredMeasure, - } - return + req.ResultCh <- tryAppendMeasuresResult{ + ResultCode: ExpiredMeasure, } - - //resultCode = ExpiredMeasure - //written = idx - break + return } - if s.MetricType == qb.Cumulative && measure.Value < s.UntilValue { - if idx == 0 { - s.Values.ForgetCapturedState() - s.Timestamps.ForgetCapturedState() - - req.ResultCh <- tryAppendMeasuresResult{ - ResultCode: NonMonotonicValue, - } - return - } - //resultCode = NonMonotonicValue - //written = idx - break - } + //resultCode = ExpiredMeasure + //written = idx + break } + if s.metricType == qb.Cumulative && measure.Value < s.lastValue { + if idx == 0 { + s.values.ForgetCapturedState() + s.timestamps.ForgetCapturedState() + + req.ResultCh <- tryAppendMeasuresResult{ + ResultCode: NonMonotonicValue, + } + return + } + //resultCode = NonMonotonicValue + //written = idx + break + } + //} + // fix - 1 + 8 bytes timestampCompressionWay, timestampRequiredSpace := timestamps.Evaluate(measure.Timestamp) valueCompressionWay, valueRequiredSpace := values.Evaluate(measure.Value) - totalSpace := timestampRequiredSpace + valueRequiredSpace + totalRequiredSpace := timestampRequiredSpace + valueRequiredSpace - if totalSpace <= atree.DataPagePayloadSize { + if totalRequiredSpace <= len(s.buffer) { // накопичую timestamps.Compress(timestampCompressionWay, measure.Timestamp) values.Compress(valueCompressionWay, measure.Value) + } else if len(s.buffer) < atree.DataPagePayloadSize { + // allocate bigger buffer + buffer := make([]byte, len(s.buffer)*2) + // copy timestamps + // copy values + // replace buffer in timestamps and values + + s.buffer = buffer + + timestamps.Compress(timestampCompressionWay, measure.Timestamp) + values.Compress(valueCompressionWay, measure.Value) } else { // сторінка заповнена - buffer := make([]byte, minBufferSize) + since := s.timestamps.ReplaceSinceWithUntil() - timestampsSize := timestamps.Rotate(buffer) - valuesSize := values.Rotate(buffer) - // prevPageNo - виставляю в txlog, коли забираю номер сторінки із freeList або генерую новий dataPages = append(dataPages, txlog.DataPayload{ - Since: s.Since, - Content: s.Buffer, - TimestampsSize: timestampsSize, - ValuesSize: valuesSize, + Since: since, + Content: s.buffer, + TimestampsSize: timestamps.Size(), + ValuesSize: values.Size(), }) + buffer := make([]byte, minBufferSize) + + timestamps.Rotate(buffer) + values.Rotate(buffer) + // renew - s.Buffer = buffer + s.buffer = buffer timestampCompressionWay, _ = timestamps.Evaluate(measure.Timestamp) valueCompressionWay, _ = values.Evaluate(measure.Value) timestamps.Compress(timestampCompressionWay, measure.Timestamp) values.Compress(valueCompressionWay, measure.Value) - - s.Since = measure.Timestamp } - - s.Until = measure.Timestamp - s.UntilValue = measure.Value + // + s.lastValue = measure.Value } // виділити змінені байти. @@ -170,20 +181,13 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu if len(dataPages) > 0 { // пишу в txlog довгим шляхом через redo файл і запис в data файл - for _, tail := range s.IndexLevelTails { - indexLevels = append(indexLevels, txlog.IndexLevelTail{ - Records: tail.Records, - RecordsCount: tail.RecordsCount, - }) - } - sendToStorage(txlog.AppendedMeasures{ MetricID: req.MetricID, - LastPageNo: s.LastPageNo, + LastPageNo: s.lastPageNo, TimestampsOffset: timestamps.Offset(), // state.Pos() з якої позиції дописувати дані на сторінку 0 (при відновленні) ValuesOffset: values.Offset(), // state.Pos() //Payload: s.payload, // fix - timestamps + values ? or timestamps and values (for WAL) - IndexLevelTails: indexLevels, + IndexLevelTails: s.indexLevelTails, DataPages: dataPages, Timestamps: nil, Values: nil, @@ -199,19 +203,15 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu func (s *_metric) FinAppendMeasures(rec txlog.AppendMeasuresSummary) { // Видаляю state. Оригінальні Timestamps і Values вже мають останню версію - s.Values.ForgetCapturedState() - s.Timestamps.ForgetCapturedState() - // fix write index levels - // update prev pageNo - // if len(rec.DataPages) > 0 { - // s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo - // } + s.values.ForgetCapturedState() + s.timestamps.ForgetCapturedState() + if rec.LastPageNo > 0 { - s.LastPageNo = rec.LastPageNo + s.lastPageNo = rec.LastPageNo } // В txlog я передав повний індекс. У нього додали елементи (можливо нові рівні). // Тому проста заміна - s.IndexLevelTails = rec.Index + s.indexLevelTails = rec.Index } // READ @@ -288,38 +288,131 @@ func (s *_metric) StartRangeScan(req tryRangeScanReq) { } func (s *_metric) StartFullScan(req tryFullScanReq) { - if s.Since == 0 { - req.ResultCh <- fullScanResult{ - ResultCode: QueryDone, - } + // if s.Since == 0 { + // req.ResultCh <- fullScanResult{ + // ResultCode: QueryDone, + // } + // return + // } + + // timestampDecompressor := s.Timestamps.CreateDecompressor() + // valueDecompressor := s.Values.CreateDecompressor() + + // for { + // timestamp, done := timestampDecompressor.NextValue() + // if done { + // break + // } + // value, done := valueDecompressor.NextValue() + // if done { + // qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + // } + // req.ResponseWriter.FeedNoSend(timestamp, value) + // } + + // if s.LastPageNo > 0 { + // req.ResultCh <- fullScanResult{ + // ResultCode: UntilFound, + // LastPageNo: s.LastPageNo, + // FracDigits: s.FracDigits, + // } + // s.RLocks++ + // } else { + // req.ResultCh <- fullScanResult{ + // ResultCode: QueryDone, + // } + // } +} + +// індекси +// Metric encode format: +// metricID - 4b +// metricType - 1b +// fracDigits - 1b +// lastPageNo - 4b +// until - 4b +// timestamps size - 2b +// timestams payload - Nb +// values size - 2b +// values payload - Nb +// index levels count - varsize +// [ +// records qty - varsize +// records - Nb +// ] + +func (s *_metric) WriteTo(w io.Writer) (err error) { + _, err = w.Write([]byte{ + byte(s.metricType), + s.fracDigits, + }) + if err != nil { return } - - timestampDecompressor := s.Timestamps.CreateDecompressor() - valueDecompressor := s.Values.CreateDecompressor(s.FracDigits) - - for { - timestamp, done := timestampDecompressor.NextValue() - if done { - break - } - value, done := valueDecompressor.NextValue() - if done { - qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - } - req.ResponseWriter.FeedNoSend(timestamp, value) + err = bin.WriteUint32(w, s.lastPageNo) + if err != nil { + return } - - if s.LastPageNo > 0 { - req.ResultCh <- fullScanResult{ - ResultCode: UntilFound, - LastPageNo: s.LastPageNo, - FracDigits: s.FracDigits, - } - s.RLocks++ - } else { - req.ResultCh <- fullScanResult{ - ResultCode: QueryDone, - } + // err = bin.WriteUint32(w, s.Since) // fix unixtime in timestamps + // if err != nil { + // return + // } + // err = bin.WriteFloat64(w, s.SinceValue) + // if err != nil { + // return + // } + err = bin.WriteUint32(w, s.timestamps.LastTimestamp()) + if err != nil { + return } + // err = bin.WriteFloat64(w, s.UntilValue) + // if err != nil { + // return + // } + // FIX - write sizes, then payloads + // copy timestamps payload + err = s.timestamps.WritePayloadTo(w) + if err != nil { + return + } + // copy values payload + err = s.values.WritePayloadTo(w) + if err != nil { + return + } + // indexes + _, err = bin.WriteVarSize(w, len(s.indexLevelTails)) + if err != nil { + return + } + for _, level := range s.indexLevelTails { + _, err = bin.WriteVarSize(w, level.RecordsCount) + if err != nil { + return + } + _, err = w.Write(level.Buffer) + if err != nil { + return + } + + } + return } + +// var ( +// values qb.ValueCompressor +// ) +// if s.metricType == qb.Cumulative { +// values = enc.NewCumulativeDeltaCompressor(s.buffer, s.fracDigits) +// } else { +// values = enc.NewInstantDeltaCompressor(s.buffer, s.fracDigits) +// } +// values.RestoreState(valuesSize) +// s.timestamps = enc.NewTimeDeltaCompressor(s.buffer) +// s.timestamps.RestoreState(timestampsSize, until) +// s.values = values +// s.lastValue = s.values.LastValue() + +// since - 4b +// sinceValue - 8b - +// untilValue - 8b diff --git a/database/proc.go b/database/proc.go index fcb45e2..122ee62 100644 --- a/database/proc.go +++ b/database/proc.go @@ -5,8 +5,7 @@ import ( qb "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/chunkenc" - "gordenko.dev/dima/qb/conbuf" + "gordenko.dev/dima/qb/enc" "gordenko.dev/dima/qb/proto" "gordenko.dev/dima/qb/transform" "gordenko.dev/dima/qb/txlog" @@ -203,8 +202,8 @@ func (s *Database) tryGetMetric(req tryGetMetricReq) { if ok { req.ResultCh <- Metric{ ResultCode: Succeed, - MetricType: metric.MetricType, - FracDigits: metric.FracDigits, + MetricType: metric.metricType, + FracDigits: metric.fracDigits, } } else { req.ResultCh <- Metric{ @@ -343,7 +342,7 @@ func (s *Database) tryRangeScan(req tryRangeScanReq) { } return } - if metric.MetricType != req.MetricType { + if metric.metricType != req.MetricType { req.ResultCh <- rangeScanResult{ ResultCode: WrongMetricType, } @@ -373,7 +372,7 @@ func (s *Database) tryFullScan(req tryFullScanReq) { } return } - if metric.MetricType != req.MetricType { + if metric.metricType != req.MetricType { req.ResultCh <- fullScanResult{ ResultCode: WrongMetricType, } @@ -399,8 +398,8 @@ func (s *Database) tryListCurrentValues(req tryListCurrentValuesReq) { if ok { req.ResponseWriter.BufferValue(transform.CurrentValue{ MetricID: metricID, - Timestamp: metric.Until, - Value: metric.UntilValue, + Timestamp: metric.timestamps.LastTimestamp(), + Value: metric.lastValue, }) } } @@ -413,7 +412,7 @@ func (s *Database) applyChanges(req txlog.Changes) { for _, untyped := range req.Records { switch rec := untyped.(type) { case txlog.AddedMetric: - s.addMetric(rec) + s.applyAddMetric(rec) case txlog.DeletedMetric: s.deleteMetric(rec) @@ -432,37 +431,27 @@ func (s *Database) applyChanges(req txlog.Changes) { } } - if req.MetricsCh != nil { - waitCh := make(chan struct{}) - req.MetricsCh <- txlog.MetricsState{ - Metrics: s.createMetricsState(), - WaitCh: waitCh, + if req.SnapshotNumberCh != nil { + s.snapshotNumber++ + err := writeSnapshot(writeSnapshotIn{ + SnapshotNumber: s.snapshotNumber, + Dir: s.dir, + WriteBufferSize: 1 * 1024 * 1024, // 1mb + Metrics: s.metrics, + FrozenIndexPagesCount: req.FrozenIndexPagesCount, + IndexPageNumbers: req.IndexPageNumbers, + FrozenDataPagesCount: req.FrozenDataPagesCount, + DataPageNumbers: req.DataPageNumbers, + }) + if err != nil { + qb.Abort(qb.WriteSnapshotFailed, err) } - // чекаю поки txlog запише снапшот, а отже можна змінювати буфери - <-waitCh + + req.SnapshotNumberCh <- s.snapshotNumber } } -func (s *Database) createMetricsState() (metrics []txlog.Metric) { - for metricID, metric := range s.metrics { - x := txlog.Metric{ - MetricID: metricID, - MetricType: metric.MetricType, - FracDigits: metric.FracDigits, - LastPageNo: metric.LastPageNo, - Since: metric.Since, - SinceValue: metric.SinceValue, - Until: metric.Until, // ? - UntilValue: metric.UntilValue, // ? - } - x.Timestamps, x.TimestampsSize = metric.Timestamps.Snapshot() - x.Values, x.ValuesSize = metric.Values.Snapshot() - metrics = append(metrics, x) - } - return -} - -func (s *Database) addMetric(rec txlog.AddedMetric) { +func (s *Database) applyAddMetric(rec txlog.AddedMetric) { // fix lock _, ok := s.metrics[rec.MetricID] if ok { @@ -477,29 +466,29 @@ func (s *Database) addMetric(rec txlog.AddedMetric) { // rec.MetricID)) // } + // lockEntry.XLock = false + // delete(s.metricLockEntries, rec.MetricID) +} + +func (s *Database) addMetric(rec txlog.AddedMetric) { var ( - values qb.ValueCompressor - timestampsBuf = conbuf.New(nil) - valuesBuf = conbuf.New(nil) + values qb.ValueCompressor + buffer = make([]byte, minBufferSize) ) if rec.MetricType == qb.Cumulative { - values = chunkenc.NewCumulativeDeltaCompressor( - valuesBuf, 0, byte(rec.FracDigits)) + values = enc.NewCumulativeDeltaCompressor(buffer, byte(rec.FracDigits)) } else { - values = chunkenc.NewInstantDeltaCompressor( - valuesBuf, 0, byte(rec.FracDigits)) + values = enc.NewInstantDeltaCompressor(buffer, byte(rec.FracDigits)) } s.metrics[rec.MetricID] = &_metric{ - MetricType: rec.MetricType, - FracDigits: byte(rec.FracDigits), - Timestamps: chunkenc.NewTimeDeltaCompressor(timestampsBuf, 0), - Values: values, + metricType: rec.MetricType, + fracDigits: byte(rec.FracDigits), + buffer: buffer, + timestamps: enc.NewTimeDeltaCompressor(buffer), + values: values, } - - lockEntry.XLock = false - delete(s.metricLockEntries, rec.MetricID) } func (s *Database) deleteMetric(rec txlog.DeletedMetric) { diff --git a/database/snapshot.go b/database/snapshot.go new file mode 100644 index 0000000..0d94280 --- /dev/null +++ b/database/snapshot.go @@ -0,0 +1,269 @@ +package database + +import ( + "bufio" + "fmt" + "io" + "os" + "path/filepath" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/util" +) + +/* +Формат: +metricsQty - varsize +[metric]* +де metric - це: + +index free list frozen pages - varsize +indexFreeList size - varsize +indexFreeList - Nb + +data free list frozen pages - varsize +dataFreeList size - varsize +dataFreeList - Nb + +CRC32 - 4b +*/ + +const readBufferSize = 8 * 1024 * 1024 // 8mb + +type writeSnapshotIn struct { + SnapshotNumber int + Dir string + WriteBufferSize int + Metrics map[uint32]*_metric + FrozenIndexPagesCount int + IndexPageNumbers []uint32 + FrozenDataPagesCount int + DataPageNumbers []uint32 +} + +func writeSnapshot(in writeSnapshotIn) (err error) { + var ( + fileName = filepath.Join(in.Dir, fmt.Sprintf("%d.snapshot", in.SnapshotNumber)) + hasher = util.NewHasher() + ) + + file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0666) //0770) + if err != nil { + return + } + + dst := io.MultiWriter(bufio.NewWriterSize(file, in.WriteBufferSize), hasher) + + _, err = bin.WriteVarSize(dst, len(in.Metrics)) + if err != nil { + return + } + + for metricID, metric := range in.Metrics { + err = bin.WriteUint32(dst, metricID) + if err != nil { + return + } + err = metric.WriteTo(dst) + if err != nil { + return + } + } + // free index pages + _, err = bin.WriteVarSize(dst, in.FrozenIndexPagesCount) + if err != nil { + return + } + _, err = bin.WriteVarSize(dst, len(in.IndexPageNumbers)) + if err != nil { + return + } + for _, pageNo := range in.IndexPageNumbers { + err = bin.WriteUint32(dst, pageNo) + if err != nil { + return + } + } + // free data pages + _, err = bin.WriteVarSize(dst, in.FrozenDataPagesCount) + if err != nil { + return + } + _, err = bin.WriteVarSize(dst, len(in.DataPageNumbers)) + if err != nil { + return + } + for _, pageNo := range in.DataPageNumbers { + err = bin.WriteUint32(dst, pageNo) + if err != nil { + return + } + } + // CRC32 + bin.WriteUint32(file, hasher.Sum32()) + + err = file.Close() + if err != nil { + return + } + + // prevLogNumber := logNumber - 1 + // prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber)) + // prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber)) + + // isExist, err := isFileExist(prevChanges) + // if err != nil { + // return + // } + + // if isExist { + // err = os.Remove(prevChanges) + // if err != nil { + // qb.Abort(qb.DeletePrevChangesFileFailed, err) + // } + // } + + // isExist, err = isFileExist(prevSnapshot) + // if err != nil { + // return + // } + + // if isExist { + // err = os.Remove(prevSnapshot) + // if err != nil { + // qb.Abort(qb.DeletePrevSnapshotFileFailed, err) + // } + // } + return +} + +type readSnapshotOut struct { + Metrics map[uint32]*ReplayMetric + FrozenIndexPagesCount int + IndexPageNumbers []uint32 + FrozenDataPagesCount int + DataPageNumbers []uint32 +} + +func readSnapshot(fileName string) (out readSnapshotOut, err error) { + var ( + metrics = make(map[uint32]*ReplayMetric) + frozenIndexPagesCount int + indexPageNumbers []uint32 + frozenDataPagesCount int + dataPageNumbers []uint32 + hasher = util.NewHasher() + ) + + file, err := os.Open(fileName) + if err != nil { + return + } + defer file.Close() + + stat, err := file.Stat() + if err != nil { + return + } + fileSize := stat.Size() + + if fileSize == 0 { + return readSnapshotOut{ + Metrics: metrics, + }, nil + } + + if fileSize < 4 { + err = fmt.Errorf("%s is corrupted", fileName) + return + } + payloadSize := fileSize - 4 + + // 2. Створюємо великий буфер для NVMe (4 MB) + bufferedReader := bufio.NewReaderSize(file, readBufferSize) + + // 3. Обмежуємо читання лише розміром Payload + limitReader := io.LimitReader(bufferedReader, payloadSize) + + // 4. Створюємо хеш та обгортку TeeReader, яка рахує CRC32 на льоту + src := io.TeeReader(limitReader, hasher) + + // читаю payload + metricsQty, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range metricsQty { + var ( + metricID uint32 + metric = &ReplayMetric{ + buffer: make([]byte, minBufferSize), + } + ) + metricID, err = bin.ReadUint32(src) + if err != nil { + return + } + err = metric.ReadFrom(src) + if err != nil { + return + } + metrics[metricID] = metric + } + // index pages + frozenIndexPagesCount, err = bin.ReadVarSize(src) + if err != nil { + return + } + indexPageNumbersCount, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range indexPageNumbersCount { + var pageNo uint32 + pageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + indexPageNumbers = append(indexPageNumbers, pageNo) + } + // data pages + frozenDataPagesCount, err = bin.ReadVarSize(src) + if err != nil { + return + } + dataPageNumbersCount, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range dataPageNumbersCount { + var pageNo uint32 + pageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + dataPageNumbers = append(dataPageNumbers, pageNo) + } + + // verify CRC32 + expectedCRC, err := bin.ReadUint32(bufferedReader) + if err != nil { + return + } + + calculatedCRC := hasher.Sum32() + + if expectedCRC != calculatedCRC { + err = fmt.Errorf("%s is corrupted. Calculated CRC %d not equal expected CRC %d", + fileName, calculatedCRC, expectedCRC) + return + } + + return readSnapshotOut{ + Metrics: metrics, + FrozenIndexPagesCount: frozenIndexPagesCount, + IndexPageNumbers: indexPageNumbers, + FrozenDataPagesCount: frozenDataPagesCount, + DataPageNumbers: dataPageNumbers, + }, nil +} diff --git a/enc/cumdelta.go b/enc/cumdelta.go index f140421..456f344 100644 --- a/enc/cumdelta.go +++ b/enc/cumdelta.go @@ -1,6 +1,7 @@ package enc import ( + "io" "log" "math" @@ -47,32 +48,32 @@ type CumulativeDeltaCompressor struct { } // Після відновлення із снапшота -func NewCumulativeDeltaCompressor(buf []byte, payloadSize int, fracDigits byte) *CumulativeDeltaCompressor { +func NewCumulativeDeltaCompressor(buf []byte, fracDigits byte) *CumulativeDeltaCompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) } - s := &CumulativeDeltaCompressor{ + return &CumulativeDeltaCompressor{ buf: buf, - pos: payloadSize, // перший вільний байт coef: coef, } +} + +func (s *CumulativeDeltaCompressor) RestoreState(payloadSize int) { if payloadSize > 0 { - u64, _, err := bin.GetVarUint64(buf) + s.pos = payloadSize - 1 // завжди показує на h + // base value на початку + u64, _, err := bin.GetVarUint64(s.buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } s.baseValue = float64(u64) / s.coef - s.pos-- s.h = s.buf[s.pos] - var n int - s.lastDelta, n, err = bin.ReverseGetVarUint64(s.buf[:s.pos]) + s.lastDelta, _, err = bin.ReverseGetVarUint64(s.buf[:s.pos-1]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } - s.pos -= n } - return s } // func (s *CumulativeDeltaCompressor) Size() int { @@ -80,83 +81,73 @@ func NewCumulativeDeltaCompressor(buf []byte, payloadSize int, fracDigits byte) // } func (s *CumulativeDeltaCompressor) Evaluate(value float64) (compressionWay int, requiredSpace int) { - var ( - delta = uint64((value-s.baseValue)*s.coef + eps) - ) + delta := uint64((value-s.baseValue)*s.coef + eps) // fix - if delta 0, no eps + requiredSpace = s.pos + // fix - requiredSpace - all space if s.pos > 0 { if s.h < 128 { // run if delta == s.lastDelta && s.h < 127 { compressionWay = incrementRun - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // current delta - hSize } else { compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - hSize + // h - end of run - bin.CountVarUint64(uint64(delta)) + // new literal - hSize // new h + requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize } } else { // literal if delta != s.lastDelta { if s.h < 255 { compressionWay = incrementLiteral - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - bin.CountVarUint64(uint64(delta)) + // new delta - hSize + requiredSpace += bin.CountVarUint64(uint64(delta)) } else { compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - hSize + // h - end of literal - bin.CountVarUint64(uint64(delta)) + // new literal - hSize // new h + requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize } } else { compressionWay = startRun if s.h > 128 { - requiredSpace += hSize // h - end of literal + requiredSpace += hSize // h - end of run } - requiredSpace += bin.CountVarUint64(uint64(delta)) + // new delta - hSize // new h } } } else { // encode base value compressionWay = addBaseValue requiredSpace = bin.CountVarUint64(uint64(value*s.coef)) + // base value - bin.CountVarUint64(uint64(delta)) + // new delta - +hSize // new h + bin.CountVarUint64(0) + // new delta + hSize // new h } return } +// pos завжди вказує на h func (s *CumulativeDeltaCompressor) Compress(compressionWay int, value float64) { delta := uint64((value-s.baseValue)*s.coef + eps) switch compressionWay { case incrementRun: s.h++ case incrementLiteral: - // write previous delta and increment counter - n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], uint64(s.lastDelta)) - s.pos += n s.lastDelta = delta s.h++ - case endSeries: - n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], uint64(s.lastDelta)) + // write previous delta and increment counter + n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], delta) s.pos += n - // write h - end of run - s.buf[s.pos] = s.h - s.pos++ - // start new literal (length=1) + case endSeries: s.lastDelta = delta - s.h = 128 + n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], delta) + s.pos += n + s.h = 128 // start new literal (length=1) case startRun: if s.h > 128 { // write h - end of literal (because length > 1) + s.pos -= 1 + bin.CountVarUint64(s.lastDelta) s.h-- s.buf[s.pos] = s.h s.pos++ + // run delta + n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], delta) + s.pos += n + } else { } // start new run (length=2) s.h = 0 @@ -165,10 +156,14 @@ func (s *CumulativeDeltaCompressor) Compress(compressionWay int, value float64) n, _ := bin.PutVarUint64(s.buf[s.pos:], uint64(value*s.coef)) s.pos += n s.baseValue = value - // start new literal (length=1) + s.lastDelta = 0 + n, _ = bin.ReversePutVarUint64(s.buf[s.pos:], s.lastDelta) + s.pos += n + // start new literal (length=1) s.h = 128 } + s.buf[s.pos] = s.h } func (s *CumulativeDeltaCompressor) DeleteLast() { @@ -186,10 +181,11 @@ func (s *CumulativeDeltaCompressor) CaptureState() { qb.Abort(qb.RepeatableLock, nil) } // позиція посувається вліво, отже може перескочити на попередній chunk + pos := s.pos - hSize - bin.CountVarUint64(s.lastDelta) s.state = &CumulativeDeltaCapturedState{ H: s.h, LastDelta: s.lastDelta, - Payload: s.buf[:s.pos], + Payload: s.buf[:pos], } } @@ -205,31 +201,65 @@ func (s *CumulativeDeltaCompressor) Offset() int { return 0 } -// Snapshot - для створення снапшота. -func (s *CumulativeDeltaCompressor) Snapshot() (left []byte, right []byte) { +func (s *CumulativeDeltaCompressor) Size() int { if s.state == nil { - left = s.buf[:s.pos] - right = s.encodeTail(s.lastDelta, s.h) + return s.pos } else { - left = s.state.Payload - right = s.encodeTail(s.state.LastDelta, s.state.H) + return len(s.state.Payload) } - return } -func (s *CumulativeDeltaCompressor) encodeTail(lastDelta uint64, h byte) []byte { - tail := make([]byte, 10) // max var uint64 + h - n, _ := bin.PutVarUint64(tail, lastDelta) - tail[n] = h - return tail[:n+1] +// Snapshot - для створення снапшота. +func (s *CumulativeDeltaCompressor) Payload() []byte { + if s.state == nil { + return s.buf[:s.pos] + } else { + return s.state.Payload + } } -func (s *CumulativeDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { - n, _ := bin.PutVarUint64(s.buf[s.pos:], s.lastDelta) - s.pos += n - s.buf[s.pos] = s.h - s.pos++ - payloadSize = s.pos +func (s *CumulativeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { + if s.state == nil { + err = bin.WriteUint16(w, uint16(s.pos)) + if err != nil { + return + } + _, err = w.Write(s.buf[:s.pos]) + return + } else { + size := bin.CountVarUint64(s.state.LastDelta) + hSize + len(s.state.Payload) + err = bin.WriteUint16(w, uint16(size)) + if err != nil { + return + } + _, err = w.Write(s.state.Payload) + if err != nil { + return + } + _, err = bin.WriteVarUint64(w, s.state.LastDelta) + if err != nil { + return + } + w.Write([]byte{ + s.state.H, + }) + return + } +} + +// func (s *CumulativeDeltaCompressor) encodeTail(lastDelta uint64, h byte) []byte { +// tail := make([]byte, 10) // max var uint64 + h +// n, _ := bin.PutVarUint64(tail, lastDelta) +// tail[n] = h +// return tail[:n+1] +// } + +func (s *CumulativeDeltaCompressor) Rotate(newbuf []byte) { + // n, _ := bin.PutVarUint64(s.buf[s.pos:], s.lastDelta) + // s.pos += n + // s.buf[s.pos] = s.h + // s.pos++ + // payloadSize = s.pos // УВАГА! // state не чіпаємо s.buf = newbuf @@ -237,7 +267,14 @@ func (s *CumulativeDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { s.baseValue = 0 s.lastDelta = 0 s.h = 0 - return +} + +func (s *CumulativeDeltaCompressor) LastValue() float64 { + if s.state == nil { + return s.baseValue + float64(s.lastDelta)*s.coef + } else { + return s.baseValue + float64(s.state.LastDelta)*s.coef + } } func (s *CumulativeDeltaCompressor) CreateDecompressor() qb.ValueDecompressor { @@ -259,10 +296,6 @@ func (s *CumulativeDeltaCompressor) CreateDecompressor() qb.ValueDecompressor { }) } -func (s *CumulativeDeltaCompressor) Chunks() []byte { - return s.buf -} - // DECOMPRESSOR type CumulativeDeltaDecompressor struct { diff --git a/enc/enc_test.go b/enc/enc_test.go index c1ba3d3..ccdebe6 100644 --- a/enc/enc_test.go +++ b/enc/enc_test.go @@ -165,11 +165,10 @@ func TestCumulativeDeltaCompressor(t *testing.T) { var ( fracDigits byte = 1 buf = make([]byte, 8) - payloadSize = 0 compressionWay int requiredSpace int ) - c := NewCumulativeDeltaCompressor(buf, payloadSize, fracDigits) + c := NewCumulativeDeltaCompressor(buf, fracDigits) for _, num := range testCase.Nums { compressionWay, requiredSpace = c.Evaluate(num) c.Compress(compressionWay, num) @@ -282,10 +281,9 @@ func TestCumulativeDeltaDecompressorFromState(t *testing.T) { var ( fracDigits byte = 1 buf = make([]byte, 16) - payloadSize = 0 decodedNums []float64 ) - c := NewCumulativeDeltaCompressor(buf, payloadSize, fracDigits) + c := NewCumulativeDeltaCompressor(buf, fracDigits) for _, num := range testCase.Nums { compressionWay, _ := c.Evaluate(num) c.Compress(compressionWay, num) @@ -526,11 +524,10 @@ func TestInstantDeltaCompressor(t *testing.T) { var ( fracDigits byte = 1 buf = make([]byte, 8) - payloadSize = 0 compressionWay int requiredSpace int ) - c := NewInstantDeltaCompressor(buf, payloadSize, fracDigits) + c := NewInstantDeltaCompressor(buf, fracDigits) for _, num := range testCase.Nums { compressionWay, requiredSpace = c.Evaluate(num) c.Compress(compressionWay, num) @@ -643,10 +640,9 @@ func TestInstantDeltaDecompressorFromState(t *testing.T) { var ( fracDigits byte = 1 buf = make([]byte, 16) - payloadSize = 0 decodedNums []float64 ) - c := NewInstantDeltaCompressor(buf, payloadSize, fracDigits) + c := NewInstantDeltaCompressor(buf, fracDigits) for _, num := range testCase.Nums { compressionWay, _ := c.Evaluate(num) c.Compress(compressionWay, num) @@ -847,7 +843,7 @@ func TestTimeDeltaCompressor(t *testing.T) { compressionWay int requiredSpace int ) - c := NewTimeDeltaCompressor(buf, 0) + c := NewTimeDeltaCompressor(buf) for _, num := range testCase.Nums { compressionWay, requiredSpace = c.Evaluate(num) c.Compress(compressionWay, num) @@ -963,7 +959,7 @@ func TestTimeDeltaDecompressorFromState(t *testing.T) { buf = make([]byte, 16) decodedNums []uint32 ) - c := NewTimeDeltaCompressor(buf, 0) + c := NewTimeDeltaCompressor(buf) for _, num := range testCase.Nums { compressionWay, _ := c.Evaluate(num) c.Compress(compressionWay, num) diff --git a/enc/insdelta.go b/enc/insdelta.go index 4782378..8814a48 100644 --- a/enc/insdelta.go +++ b/enc/insdelta.go @@ -1,6 +1,7 @@ package enc import ( + "io" "log" "math" @@ -18,32 +19,33 @@ type InstantDeltaCompressor struct { state *InstantDeltaCapturedState } -func NewInstantDeltaCompressor(buf []byte, payloadSize int, fracDigits byte) *InstantDeltaCompressor { +func NewInstantDeltaCompressor(buf []byte, fracDigits byte) *InstantDeltaCompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) } s := &InstantDeltaCompressor{ buf: buf, - pos: payloadSize, coef: coef, } + return s +} + +func (s *InstantDeltaCompressor) RestoreState(payloadSize int) { if payloadSize > 0 { - u64, _, err := bin.GetVarInt64(buf) + s.pos = payloadSize - 1 // завжди показує на h + // base value на початку + u64, _, err := bin.GetVarUint64(s.buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } s.baseValue = float64(u64) / s.coef - s.pos-- s.h = s.buf[s.pos] - var n int - s.lastDelta, n, err = bin.ReverseGetVarInt64(s.buf[:s.pos]) + s.lastDelta, _, err = bin.ReverseGetVarInt64(s.buf[:s.pos-1]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } - s.pos -= n } - return s } // func (s *InstantDeltaCompressor) Size() int { @@ -184,31 +186,60 @@ func (s *InstantDeltaCompressor) Offset() int { return 0 } -// Snapshot - для створення снапшота. -func (s *InstantDeltaCompressor) Snapshot() (left []byte, right []byte) { +func (s *InstantDeltaCompressor) Size() int { if s.state == nil { - left = s.buf[:s.pos] - right = s.encodeTail(s.lastDelta, s.h) + return s.pos } else { - left = s.state.Payload - right = s.encodeTail(s.state.LastDelta, s.state.H) + return len(s.state.Payload) } - return } -func (s *InstantDeltaCompressor) encodeTail(lastDelta int64, h byte) []byte { - tail := make([]byte, 10) // max var uint64 + h - n, _ := bin.PutVarInt64(tail, lastDelta) - tail[n] = h - return tail[:n+1] +// Snapshot - для створення снапшота. +func (s *InstantDeltaCompressor) Payload() []byte { + if s.state == nil { + return s.buf[:s.pos] + } else { + return s.state.Payload + } } -func (s *InstantDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { - n, _ := bin.PutVarInt64(s.buf[s.pos:], s.lastDelta) - s.pos += n - s.buf[s.pos] = s.h - s.pos++ - payloadSize = s.pos +func (s *InstantDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { + if s.state == nil { + err = bin.WriteUint16(w, uint16(s.pos)) + if err != nil { + return + } + _, err = w.Write(s.buf[:s.pos]) + return + } else { + size := bin.CountVarInt64(s.state.LastDelta) + hSize + len(s.state.Payload) + err = bin.WriteUint16(w, uint16(size)) + if err != nil { + return + } + _, err = w.Write(s.state.Payload) + if err != nil { + return + } + _, err = bin.WriteVarInt64(w, s.state.LastDelta) + if err != nil { + return + } + w.Write([]byte{ + s.state.H, + }) + return + } +} + +// func (s *InstantDeltaCompressor) encodeTail(lastDelta int64, h byte) []byte { +// tail := make([]byte, 10) // max var uint64 + h +// n, _ := bin.PutVarInt64(tail, lastDelta) +// tail[n] = h +// return tail[:n+1] +// } + +func (s *InstantDeltaCompressor) Rotate(newbuf []byte) { // УВАГА! // state не чіпаємо s.buf = newbuf @@ -216,7 +247,14 @@ func (s *InstantDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { s.baseValue = 0 s.lastDelta = 0 s.h = 0 - return +} + +func (s *InstantDeltaCompressor) LastValue() float64 { + if s.state == nil { + return s.baseValue + float64(s.lastDelta)*s.coef + } else { + return s.baseValue + float64(s.state.LastDelta)*s.coef + } } func (s *InstantDeltaCompressor) CreateDecompressor() qb.ValueDecompressor { diff --git a/enc/time_delta.go b/enc/time_delta.go index d221a1c..bc942d9 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -1,6 +1,7 @@ package enc import ( + "io" "log" bin "gordenko.dev/dima/bin/little" @@ -27,125 +28,125 @@ type TimeDeltaCompressor struct { state *TimeDeltaCapturedState } -// Кодує значення справа наліво футнкціями TailPut*. Читає значення зліва направо функціями Get* -func NewTimeDeltaCompressor(buf []byte, payloadSize int) *TimeDeltaCompressor { - s := &TimeDeltaCompressor{ +// Початкове створення, коли даних немає +// Створення після page filled, коли одразу є початкове значення +// Відновлення стану із снапшота. LastTimestamp передаю окермо від payload, +// а lastDelta і h треба прочитати із payload (якщо є) + +// Кодує значення справа наліво функціями TailPut*. Читає значення зліва направо функціями Get* +func NewTimeDeltaCompressor(buf []byte) *TimeDeltaCompressor { + return &TimeDeltaCompressor{ buf: buf, - pos: len(buf) - payloadSize, } +} + +func (s *TimeDeltaCompressor) RestoreState(payloadSize int, lastTimestamp uint32) { if payloadSize > 0 { - var err error - s.lastUnixtime, err = bin.GetUint32(s.buf[s.pos:]) - if err != nil { - log.Fatalf("bug: get last unixtime: %s", err) - } - s.pos += 4 - if s.pos < len(buf) { + s.lastUnixtime = lastTimestamp + s.pos = len(s.buf) - payloadSize + if payloadSize > 4 { s.h = s.buf[s.pos] - s.pos++ - u64, n, err := bin.GetVarUint64(s.buf[s.pos:]) + u64, _, err := bin.GetVarUint64(s.buf[s.pos+1:]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } s.lastDelta = uint32(u64) - s.pos += n } } - return s } -// func (s *TimeDeltaCompressor) Size() int { -// return len(s.buf) - s.pos -// } - // retrun func (s *TimeDeltaCompressor) Evaluate(unixtime uint32) (compressionWay int, requiredSpace int) { - requiredSpace = 4 // last unixtime - if s.lastUnixtime > 0 { + requiredSpace = len(s.buf) - s.pos + if requiredSpace > 0 { delta := unixtime - s.lastUnixtime if s.lastDelta > 0 { - // 3rd value if s.h < 128 { // run if delta == s.lastDelta && s.h < 127 { compressionWay = incrementRun - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // current delta - hSize } else { compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - hSize + // h - end of run - bin.CountVarUint64(uint64(delta)) + // new literal - hSize // new h + requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize } } else { // literal if delta != s.lastDelta { if s.h < 255 { compressionWay = incrementLiteral - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - bin.CountVarUint64(uint64(delta)) + // new delta - hSize + requiredSpace += bin.CountVarUint64(uint64(delta)) } else { compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(s.lastDelta)) + // previous delta - hSize + // h - end of literal - bin.CountVarUint64(uint64(delta)) + // new literal - hSize // new h + requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize } } else { compressionWay = startRun if s.h > 128 { - requiredSpace += hSize // h - end of literal + requiredSpace += hSize // h - end of run } - requiredSpace += bin.CountVarUint64(uint64(delta)) + // new delta - hSize // new h } } } else { compressionWay = add1stDelta - requiredSpace += bin.CountVarUint64(uint64(delta)) + // new literal - hSize // new h + requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize } } else { compressionWay = addUnixtime + requiredSpace += 4 } return } +// pos вказує на h byte func (s *TimeDeltaCompressor) Compress(compressionWay int, unixtime uint32) { delta := unixtime - s.lastUnixtime switch compressionWay { case incrementRun: s.h++ + s.buf[s.pos] = s.h case incrementLiteral: - // write previous delta and increment counter - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(s.lastDelta)) - s.pos -= n s.lastDelta = delta s.h++ - case endSeries: - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(s.lastDelta)) + // тому що перезаписую h + s.pos++ + n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) s.pos -= n - // write h - end of run s.pos-- s.buf[s.pos] = s.h + case endSeries: // start new literal (length=1) s.lastDelta = delta s.h = 128 + n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) + s.pos -= n + s.pos-- + s.buf[s.pos] = s.h case startRun: if s.h > 128 { - // write h - end of literal (because length > 1) s.h-- + s.pos += bin.CountVarUint64(uint64(s.lastDelta)) // забираю останню дельту у literal + s.buf[s.pos] = s.h // write h - end of literal (because length > 1) + s.pos-- + // run delta + n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) + s.pos -= n + // для h s.pos-- - s.buf[s.pos] = s.h } // start new run (length=2) s.h = 0 + s.buf[s.pos] = s.h case add1stDelta: // start new literal (length=1) s.lastDelta = delta s.h = 128 + n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) + s.pos -= n + s.pos-- + s.buf[s.pos] = s.h + case addUnixtime: + s.pos -= 4 + bin.PutUint32(s.buf[s.pos:], unixtime) } // 1st value or update after every update // do not write immediatelly because of update after every update @@ -184,6 +185,14 @@ func (s *TimeDeltaCompressor) ForgetCapturedState() { s.state = nil } +// коли сторінка заповнена і відправляється в txlog, since замінюю на until +func (s *TimeDeltaCompressor) ReplaceSinceWithUntil() uint32 { + pos := len(s.buf) - 4 + since, _ := bin.GetUint32(s.buf[pos:]) + bin.PutUint32(s.buf[pos:], s.lastUnixtime) + return since +} + // fix func (s *TimeDeltaCompressor) Offset() int { // if s.state != nil { @@ -192,35 +201,64 @@ func (s *TimeDeltaCompressor) Offset() int { return 0 } -// Snapshot - для створення снапшота. -func (s *TimeDeltaCompressor) Snapshot() (left []byte, right []byte) { +func (s *TimeDeltaCompressor) Size() int { if s.state == nil { - left = s.encodeTail(s.lastUnixtime, s.lastDelta, s.h) - right = s.buf[s.pos:] + return len(s.buf) - s.pos } else { - left = s.encodeTail(s.state.LastUnixtime, s.state.LastDelta, s.state.H) - right = s.state.Payload + return len(s.state.Payload) } - return } -func (s *TimeDeltaCompressor) encodeTail(lastUnixtime, lastDelta uint32, h byte) []byte { - tail := make([]byte, 9) - bin.PutUint32(tail, lastUnixtime) - tail[4] = h - n, _ := bin.PutVarUint64(tail[5:], uint64(lastDelta)) - return tail[:n+5] +// Snapshot - для створення снапшота. +// FIX - encode firstUnixtime in 1st 4 bytes +// +// while page not completed - has first unixtime, after - last unixtime +func (s *TimeDeltaCompressor) Payload() []byte { + if s.state == nil { + return s.buf[s.pos:] + } else { + return s.state.Payload + } } -func (s *TimeDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { - var n int - n, _ = bin.TailPutVarUint64(s.buf[:s.pos], uint64(s.lastDelta)) - s.pos -= n - s.pos-- - s.buf[s.pos] = s.h - n, _ = bin.TailPutVarUint64(s.buf[:s.pos], uint64(s.lastUnixtime)) - s.pos -= n - payloadSize = len(s.buf) - s.pos +func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { + if s.state == nil { + err = bin.WriteUint16(w, uint16(s.pos)) + if err != nil { + return + } + _, err = w.Write(s.buf[:s.pos]) + return + } else { + size := bin.CountVarUint64(uint64(s.state.LastDelta)) + hSize + len(s.state.Payload) + err = bin.WriteUint16(w, uint16(size)) + if err != nil { + return + } + _, err = w.Write(s.state.Payload) + if err != nil { + return + } + _, err = bin.WriteVarUint64(w, uint64(s.state.LastDelta)) + if err != nil { + return + } + w.Write([]byte{ + s.state.H, + }) + return + } +} + +// func (s *TimeDeltaCompressor) encodeTail(lastUnixtime, lastDelta uint32, h byte) []byte { +// tail := make([]byte, 9) +// bin.PutUint32(tail, lastUnixtime) +// tail[4] = h +// n, _ := bin.PutVarUint64(tail[5:], uint64(lastDelta)) +// return tail[:n+5] +// } + +func (s *TimeDeltaCompressor) Rotate(newbuf []byte) { // УВАГА! // state не чіпаємо s.buf = newbuf @@ -228,7 +266,6 @@ func (s *TimeDeltaCompressor) Rotate(newbuf []byte) (payloadSize int) { s.lastUnixtime = 0 s.lastDelta = 0 s.h = 0 - return } func (s *TimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompressor { @@ -252,8 +289,8 @@ func (s *TimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompressor { }) } -func (s *TimeDeltaCompressor) Chunks() []byte { - return s.buf +func (s *TimeDeltaCompressor) LastTimestamp() uint32 { + return s.lastUnixtime } // DECOMPRESSOR diff --git a/freelist/freelist.go b/freelist/freelist.go index 7342c6d..a337c13 100644 --- a/freelist/freelist.go +++ b/freelist/freelist.go @@ -227,6 +227,10 @@ func (s *FreeList) Serialize() ([]byte, error) { return w.Bytes(), nil } +func (s *FreeList) Cached() []uint32 { + return s.free +} + func (s *FreeList) Merge() (err error) { var ( off = int64(s.basePagesCount * s.pageSize) diff --git a/go.mod b/go.mod index abff683..9b6cb11 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-20260606153512-101b10d92a39 + gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696 gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 ) diff --git a/go.sum b/go.sum index 0f356ec..f2d3751 100644 --- a/go.sum +++ b/go.sum @@ -20,5 +20,7 @@ gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gordenko.dev/dima/bin v0.0.0-20260606153512-101b10d92a39 h1:qW3SnQ9HAcJpfj8ottBfPGTzHX+cX2vM8B6TVOgnypE= 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/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/qb.go b/qb.go index 6ddd515..eeb2723 100644 --- a/qb.go +++ b/qb.go @@ -2,6 +2,7 @@ package qb import ( "fmt" + "io" "os" ) @@ -23,37 +24,41 @@ const ( ) type TimestampCompressor interface { + RestoreState(int, uint32) // (payload size, last timestamp) // (timestamp) => compressionWay, requiredSpace Evaluate(uint32) (int, int) Compress(int, uint32) - //Size() int + Size() int //Chunks() [][]byte //DeleteLast() CaptureState() ForgetCapturedState() CreateDecompressor() TimestampDecompressor - Snapshot() ([]byte, []byte) // tail, head + //Payload() []byte // для снапшота Offset() int - Rotate([]byte) int // payloadSize - //LastTimestamp() uint32 + Rotate([]byte) + WritePayloadTo(io.Writer) error + LastTimestamp() uint32 + ReplaceSinceWithUntil() uint32 } type ValueCompressor interface { + RestoreState(int) // payload size // (value) => compressionWay, requiredSpace Evaluate(float64) (int, int) Compress(int, float64) - //Size() int + Size() int //Chunks() [][]byte //DeleteLast() - Renew() // створює новий conbuf, але не чіпає state CaptureState() ForgetCapturedState() // fracDigits - CreateDecompressor(byte) ValueDecompressor - Snapshot() ([]byte, []byte) // tail, head + CreateDecompressor() ValueDecompressor + //Payload() []byte // для снапшота Offset() int - Rotate([]byte) int // payloadSize - //LastValue() float64 + Rotate([]byte) + WritePayloadTo(io.Writer) error + LastValue() float64 } type TimestampDecompressor interface { @@ -90,6 +95,8 @@ const ( UnknownMetricWaitQueueItemBug AbortCode = 21 RepeatableLock AbortCode = 22 NoSpaceOnIndexPage AbortCode = 23 + + WriteSnapshotFailed AbortCode = 24 // GetRecoveryRecipeFailed AbortCode = 26 LoadSnapshotFailed AbortCode = 27 @@ -101,6 +108,8 @@ const ( ReplayREDOFileFailed AbortCode = 33 DeletePrevChangesFileFailed AbortCode = 34 DeletePrevSnapshotFileFailed AbortCode = 35 + + MetricNotFoundDuringReplay AbortCode = 36 ) func Abort(code AbortCode, err error) { diff --git a/txlog/data_preparer.go b/txlog/data_preparer.go index 8b43037..4da5ddf 100644 --- a/txlog/data_preparer.go +++ b/txlog/data_preparer.go @@ -109,7 +109,7 @@ func (s *DataPreparer) Prepare(input []any) PreparedData { }) indexLevelTails = append(indexLevelTails, IndexLevelTail{ - Records: level.TailRecords, + Buffer: level.TailRecords, RecordsCount: level.TailRecordsCount, }) } @@ -181,9 +181,9 @@ func (s *DataPreparer) Reset() { type AppendedMeasures struct { MetricID uint32 LastPageNo uint32 - TimestampsOffset int - Timestamps []byte // offset on 1st page - ValuesOffset int // offset on 1st page + TimestampsOffset int // offset on 1st page + Timestamps []byte + ValuesOffset int // offset on 1st page Values []byte IndexLevelTails []IndexLevelTail DataPages []DataPayload diff --git a/txlog/page_preparer.go b/txlog/page_preparer.go index 5555b29..26deb4e 100644 --- a/txlog/page_preparer.go +++ b/txlog/page_preparer.go @@ -122,7 +122,7 @@ func (s *PagePreparer) sealIndexPage(in sealIndexPageIn) (checksum uint32) { } type IndexLevelTail struct { - Records []byte // розмір більший за кількість + Buffer []byte // розмір більший за кількість RecordsCount int } @@ -185,7 +185,7 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { for _, tail := range in.IndexLevelTails { sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{ SkipRecords: tail.RecordsCount, - TailRecords: tail.Records, + TailRecords: tail.Buffer, TailRecordsCount: tail.RecordsCount, }) } diff --git a/txlog/reader.go b/txlog/reader.go index a6f56e0..a6ece2e 100644 --- a/txlog/reader.go +++ b/txlog/reader.go @@ -5,16 +5,20 @@ import ( "bytes" "errors" "fmt" - "hash/crc32" "io" "os" - "gordenko.dev/dima/qb/bin" + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/util" ) +const readBufferSize = 8 * 1024 * 1024 + type Reader struct { - file *os.File - reader *bufio.Reader + file *os.File + reader *bufio.Reader + current []any + next []any } type ReaderOptions struct { @@ -37,7 +41,7 @@ func NewReader(opt ReaderOptions) (*Reader, error) { return &Reader{ file: file, - reader: bufio.NewReaderSize(file, 1024*1024), + reader: bufio.NewReaderSize(file, readBufferSize), }, nil } @@ -45,42 +49,62 @@ func (s *Reader) Close() { s.file.Close() } -func (s *Reader) ReadPacket() (uint32, []any, bool, error) { - prefix := make([]byte, packetPrefixSize) - n, err := s.reader.Read(prefix) - if err != nil { - if err == io.EOF && n == 0 { - return 0, nil, true, nil - } else { - return 0, nil, false, fmt.Errorf("read packet prefix: %s", err) +func (s *Reader) NextPacket() ([]any, bool, error) { + var err error + if s.current == nil { + s.current, err = s.readPacket() + if err != nil { + return nil, false, err + } + if len(s.current) > 0 { + s.next, err = s.readPacket() + if err != nil { + return nil, false, err + } } } - length := bin.GetUint32(prefix[lengthIdx:]) - storedCRC := bin.GetUint32(prefix[checksumIdx:]) - lsn := bin.GetUint32(prefix[lsnIdx:]) + current := s.current + done := s.next == nil - body, err := bin.ReadN(s.reader, int(length)) + s.current = s.next + s.next, err = s.readPacket() if err != nil { - return 0, nil, false, fmt.Errorf("read packet body: %s", err) + return nil, false, err + } + return current, done, nil +} + +func (s *Reader) readPacket() (_ []any, err error) { + payloadSize, err := bin.ReadVarSize(s.reader) + if err != nil { + if err == io.EOF { + return nil, nil + } else { + return + } } - hasher := crc32.NewIEEE() - hasher.Write(prefix[lsnIdx:]) + storedCRC, err := bin.ReadUint32(s.reader) + if err != nil { + return + } + + body, err := bin.ReadN(s.reader, int(payloadSize)) + if err != nil { + return + } + + hasher := util.NewHasher() hasher.Write(body) calculatedCRC := hasher.Sum32() - if calculatedCRC != storedCRC { - return 0, nil, false, fmt.Errorf("stored CRC %d != calculated CRC %d", - storedCRC, calculatedCRC) + if calculatedCRC == storedCRC { + return s.parseRecords(body) } - - records, err := s.parseRecords(body) - if err != nil { - return 0, nil, false, err - } - return lsn, records, false, nil + //s.logger.Printf("stored CRC %d != calculated CRC %d", storedCRC, calculatedCRC) + return nil, nil } func (s *Reader) parseRecords(body []byte) ([]any, error) { diff --git a/txlog/snapshot.go b/txlog/snapshot.go deleted file mode 100644 index 942d0a8..0000000 --- a/txlog/snapshot.go +++ /dev/null @@ -1,122 +0,0 @@ -package txlog - -import ( - "fmt" - "io" - "os" - - bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/util" -) - -func LoadSnapshot(fileName string) (snapshot Snapshot, err error) { - var ( - hasher = util.NewHasher() - header = make([]byte, metricHeaderSize) - body = make([]byte, atree.PageSize) - ) - - file, err := os.Open(fileName) - if err != nil { - return - } - - src := io.TeeReader(file, hasher) - metricsQty, err := bin.ReadVarSize(src) - if err != nil { - return - } - - for range metricsQty { - var metric Metric - err = bin.ReadNInto(src, header) - if err != nil { - return - } - - metric.MetricID, _ = bin.GetUint32(header[0:]) - metric.MetricType = qb.MetricType(header[4]) - metric.FracDigits = header[5] - metric.LastPageNo, _ = bin.GetUint32(header[6:]) - metric.Since, _ = bin.GetUint32(header[10:]) - metric.SinceValue, _ = bin.GetFloat64(header[14:]) - metric.Until, _ = bin.GetUint32(header[22:]) - metric.UntilValue, _ = bin.GetFloat64(header[26:]) - tSize, _ := bin.GetUint16(header[34:]) - vSize, _ := bin.GetUint16(header[36:]) - - buf := body[:tSize] - err = bin.ReadNInto(src, buf) - if err != nil { - return - } - // FIX - // metric.Timestamps = chunkenc.NewReverseTimeDeltaCompressor( - // conbuf.NewFromBuffer(buf), - // int(tSize), - // ) - - buf = body[:vSize] - err = bin.ReadNInto(src, buf) - if err != nil { - return - } - - // FIX - // if metric.MetricType == qb.Cumulative { - // metric.Values = chunkenc.NewReverseCumulativeDeltaCompressor( - // conbuf.NewFromBuffer(buf), - // int(vSize), - // metric.FracDigits, - // ) - // } else { - // metric.Values = chunkenc.NewReverseInstantDeltaCompressor( - // conbuf.NewFromBuffer(buf), - // int(vSize), - // metric.FracDigits, - // ) - // } - // s.metrics[metricID] = &metric - } - // FIX - // err = restoreFreeList(s.dataFreeList, src) - // if err != nil { - // return fmt.Errorf("restore dataFreeList: %s", err) - // } - // FIX - // err = restoreFreeList(s.indexFreeList, src) - // if err != nil { - // return fmt.Errorf("restore indexFreeList: %s", err) - // } - - calculatedChecksum := hasher.Sum32() - - writtenChecksum, err := bin.ReadUint32(file) - if err != nil { - return - } - - if calculatedChecksum != writtenChecksum { - err = fmt.Errorf("calculated checksum %d not equal written checksum %d", - calculatedChecksum, writtenChecksum) - return - } - return -} - -// serialized + frozen pages count - -// func restoreFreeList(freeList *freelist.FreeList, src io.Reader) error { -// size, err := bin.ReadVarSize(src) -// if err != nil { -// return err -// } -// serialized, err := bin.ReadN(src, size) -// if err != nil { -// return err -// } -// freeList.Restore(serialized) -// return nil -// } diff --git a/txlog/txlog.go b/txlog/txlog.go index 95c3a90..014818d 100644 --- a/txlog/txlog.go +++ b/txlog/txlog.go @@ -5,7 +5,6 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/conbuf" ) var ( @@ -79,49 +78,6 @@ func (s *DeletedMetric) Parse(src io.Reader) (err error) { return nil } -func writeChunkedPayloadTo(chunks [][]byte, size int, w io.Writer) { - bin.WriteVarSize(w, size) - for _, chunk := range chunks { - if size >= len(chunk) { - w.Write(chunk) - size -= len(chunk) - } else { - w.Write(chunk[:size]) - return - } - } -} - -func readChunkedPayload(src io.Reader) (_ [][]byte, _ int, err error) { - var ( - chunks [][]byte - size int - remainingSize = size - ) - size, err = bin.ReadVarSize(src) - if err != nil { - return - } - for remainingSize > 0 { - var chunk []byte - if remainingSize >= conbuf.ChunkSize { - chunk, err = bin.ReadN(src, conbuf.ChunkSize) - if err != nil { - return - } - } else { - chunk = make([]byte, conbuf.ChunkSize) - err = bin.ReadNInto(src, chunk[:remainingSize]) - if err != nil { - return - } - } - chunks = append(chunks, chunk) - remainingSize -= conbuf.ChunkSize - } - return chunks, size, nil -} - // fix - add changed pages type DeletedMeasures struct { MetricID uint32 diff --git a/txlog/writer.go b/txlog/writer.go index 2392e72..567e80c 100644 --- a/txlog/writer.go +++ b/txlog/writer.go @@ -1,21 +1,18 @@ package txlog import ( - "bufio" "bytes" "errors" "fmt" - "io" + "log" "math" "os" "path/filepath" "sync" - bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/freelist" - "gordenko.dev/dima/qb/util" ) const ( @@ -52,23 +49,27 @@ func JoinChangesFileName(dir string, logNumber int) string { } type MetricsState struct { - Metrics []Metric - WaitCh chan struct{} + //Metrics []Metric + WaitCh chan struct{} } type Changes struct { - Records []any - MetricsCh chan MetricsState + Records []any + SnapshotNumberCh chan int // log number + FrozenIndexPagesCount int + IndexPageNumbers []uint32 + FrozenDataPagesCount int + DataPageNumbers []uint32 } type Writer struct { - mutex sync.Mutex - dataPagesCount uint32 - indexPagesCount uint32 - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList - atree *atree.Atree - logNumber int + mutex sync.Mutex + dataPagesCount uint32 + indexPagesCount uint32 + dataFreeList *freelist.FreeList + indexFreeList *freelist.FreeList + atree *atree.Atree + //logNumber int dir string w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск wal *os.File @@ -91,7 +92,7 @@ type WriterOptions struct { IndexPageSize int DataPageSize int Dir string - LogNumber int // номер журнала + SnapshotNumber int // номер журнала AppendToWorkerQueue func(any) DataFreeList *freelist.FreeList IndexFreeList *freelist.FreeList @@ -137,7 +138,6 @@ func NewWriter(opt WriterOptions) (*Writer, error) { dataFreeList: opt.DataFreeList, indexFreeList: opt.IndexFreeList, atree: opt.Atree, - logNumber: opt.LogNumber, exitCh: opt.ExitCh, waitGroup: opt.WaitGroup, signalCh: make(chan struct{}, 1), @@ -145,25 +145,16 @@ func NewWriter(opt WriterOptions) (*Writer, error) { var err error - if opt.LogNumber > 0 { - s.wal, err = os.OpenFile( - JoinChangesFileName(opt.Dir, s.logNumber), - os.O_APPEND|os.O_WRONLY, - filePerm, - ) - if err != nil { - return nil, err - } - } else { - s.logNumber = 1 - s.wal, err = os.OpenFile( - JoinChangesFileName(opt.Dir, s.logNumber), - os.O_CREATE|os.O_WRONLY, - filePerm, - ) - if err != nil { - return nil, err - } + if opt.SnapshotNumber <= 0 { + log.Fatalln("FIX FUCK NUMBER") + } + s.wal, err = os.OpenFile( + JoinChangesFileName(opt.Dir, opt.SnapshotNumber), + os.O_APPEND|os.O_WRONLY, + filePerm, + ) + if err != nil { + return nil, err } return s, nil } @@ -240,7 +231,7 @@ func (s *Writer) packAndWrite() (err error) { } // 4. Пишу в atree сторінки - err = s.writePagesToAtree(prepared.WriteToIndex, prepared.WriteToData) + err = s.WritePagesToAtree(prepared.WriteToIndex, prepared.WriteToData) if err != nil { return } @@ -260,64 +251,49 @@ func (s *Writer) packAndWrite() (err error) { // } if forceSnapshot { - metricsStateCh := make(chan MetricsState) + if err := s.wal.Close(); err != nil { + return fmt.Errorf("close wal file: %s", err) + } + s.wal = nil + s.written = 0 + + snapshotNumberCh := make(chan int, 1) s.appendToWorkerQueue(Changes{ - Records: prepared.WriteResults, - MetricsCh: metricsStateCh, + Records: prepared.WriteResults, + SnapshotNumberCh: snapshotNumberCh, + FrozenIndexPagesCount: s.indexFreeList.Pages(), + IndexPageNumbers: s.indexFreeList.Cached(), + FrozenDataPagesCount: s.dataFreeList.Pages(), + DataPageNumbers: s.dataFreeList.Cached(), }) - if err := s.wal.Close(); err != nil { - return fmt.Errorf("close changes file: %s", err) - } + snapshotNumber := <-snapshotNumberCh - state := <-metricsStateCh - - s.logNumber++ - - // write snapshot FIX - err = writeSnapshot(Snapshot{ - LogNumber: s.logNumber, - Dir: s.dir, - WriteBufferSize: writeBufferSize, - Metrics: state.Metrics, - DataFreeList: s.dataFreeList, - IndexFreeList: s.indexFreeList, - }) + // копіюю сторінки із delta файла в base файл і потім роблю Truncate + err = s.indexFreeList.Merge() if err != nil { - return fmt.Errorf("write snapshot file: %s", err) + return + } + err = s.dataFreeList.Merge() + if err != nil { + return } var err error s.wal, err = os.OpenFile( - JoinChangesFileName(s.dir, s.logNumber), + JoinChangesFileName(s.dir, snapshotNumber), os.O_CREATE|os.O_WRONLY, filePerm, ) if err != nil { return fmt.Errorf("create new changes file: %s", err) } - - s.written = 0 - // release Worker - close(state.WaitCh) } else { s.appendToWorkerQueue(Changes{ Records: prepared.WriteResults, }) } - - // Якщо потрібен снапшот - відправляю канал із буфером розміру 1, - // в який воркер має покласти снапшот отриманий після застосування змін (records). - // Після чого Writer створює нові файли snapshot і changes і працює далі. - - // flush - // err = s.flush(w) - // if err != nil { - // return err - // } - - // FIX send to worker workerReqs return nil } @@ -372,7 +348,7 @@ type PageToWrite struct { } // writePagesToAtree - записує сторінки в .data та .index файли -func (s *Writer) writePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) { +func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) { for _, p := range dataPages { if len(p.Content) != s.dataPageSize { return fmt.Errorf("wrong data page size: %d", len(p.Content)) @@ -457,203 +433,33 @@ func (s *Writer) sendSignal() { // helpers -/* -Формат: -metricsQty - varuint -[metric]* -де metric - це: -metricID - 4b -metricType - 1b -fracDigits - 1b -lastPageNo - 4b -since - 4b -sinceValue - 8b -until - 4b -untilValue - 8b -timestamps size - 2b -values size - 2b -timestams payload - Nb -values payload - Nb -data free list frozen pages - varsize -dataFreeList size - varsize -dataFreeList - Nb -index free list frozen pages - varsize -indexFreeList size - varsize -indexFreeList - Nb -CRC32 - 4b -*/ +// type Metric struct { +// MetricID uint32 +// MetricType qb.MetricType +// FracDigits byte +// LastPageNo uint32 +// Since uint32 +// SinceValue float64 +// Until uint32 +// UntilValue float64 +// Timestamps []byte // payload FIX +// Values []byte // payload FIX add lastDelta + h ? +// } -const metricHeaderSize = 38 - -type Metric struct { - MetricID uint32 - MetricType qb.MetricType - FracDigits byte - LastPageNo uint32 - Since uint32 - SinceValue float64 - Until uint32 - UntilValue float64 - Timestamps [][]byte - TimestampsSize int - Values [][]byte - ValuesSize int -} - -type Snapshot struct { - LogNumber int - Dir string - WriteBufferSize int - Metrics []Metric - IndexFreeList *freelist.FreeList - DataFreeList *freelist.FreeList -} - -func writeSnapshot(snapshot Snapshot) (err error) { - var ( - fileName = filepath.Join(snapshot.Dir, fmt.Sprintf("%d.snapshot", snapshot.LogNumber)) - hasher = util.NewHasher() - prefix = make([]byte, metricHeaderSize) - ) - - file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0770) - if err != nil { - return - } - - dst := io.MultiWriter(bufio.NewWriterSize(file, snapshot.WriteBufferSize), hasher) - - _, err = bin.WriteVarSize(dst, len(snapshot.Metrics)) - if err != nil { - return - } - - for _, metric := range snapshot.Metrics { - tSize := metric.TimestampsSize - vSize := metric.ValuesSize - - bin.PutUint32(prefix[0:], metric.MetricID) - prefix[4] = byte(metric.MetricType) - prefix[5] = metric.FracDigits - bin.PutUint32(prefix[6:], metric.LastPageNo) - bin.PutUint32(prefix[10:], metric.Since) - bin.PutFloat64(prefix[14:], metric.SinceValue) - bin.PutUint32(prefix[22:], metric.Until) - bin.PutFloat64(prefix[26:], metric.UntilValue) - bin.PutUint16(prefix[34:], uint16(tSize)) - bin.PutUint16(prefix[36:], uint16(vSize)) - - _, err = dst.Write(prefix) - if err != nil { - return - } - // copy timestamps - writeChunks(dst, metric.Timestamps, tSize) - // copy values - writeChunks(dst, metric.Values, vSize) - } - // free data pages - _, err = bin.WriteVarSize(dst, snapshot.DataFreeList.Pages()) - if err != nil { - return - } - err = freeListWriteTo(snapshot.DataFreeList, dst) - if err != nil { - return - } - // free index pages - _, err = bin.WriteVarSize(dst, snapshot.IndexFreeList.Pages()) - if err != nil { - return - } - err = freeListWriteTo(snapshot.IndexFreeList, dst) - if err != nil { - return - } - - bin.WriteUint32(file, hasher.Sum32()) - - err = file.Sync() - if err != nil { - return - } - - err = file.Close() - if err != nil { - return - } - // копіюю сторінки із delta файла в base файл і потім роблю Truncate - err = snapshot.DataFreeList.Merge() // fix - get frozen pages count + unfilled - if err != nil { - return - } - err = snapshot.IndexFreeList.Merge() - if err != nil { - return - } - - // prevLogNumber := logNumber - 1 - // prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber)) - // prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber)) - - // isExist, err := isFileExist(prevChanges) - // if err != nil { - // return - // } - - // if isExist { - // err = os.Remove(prevChanges) - // if err != nil { - // qb.Abort(qb.DeletePrevChangesFileFailed, err) - // } - // } - - // isExist, err = isFileExist(prevSnapshot) - // if err != nil { - // return - // } - - // if isExist { - // err = os.Remove(prevSnapshot) - // if err != nil { - // qb.Abort(qb.DeletePrevSnapshotFileFailed, err) - // } - // } - return -} - -// HELPERS - -func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error { - serialized, err := freeList.Serialize() - if err != nil { - qb.Abort(qb.FailedFreeListSerialize, err) - } - _, err = bin.WriteVarSize(dst, len(serialized)) - if err != nil { - return err - } - _, err = dst.Write(serialized) - if err != nil { - return err - } - return nil -} - -func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) { - remaining := size - for _, buf := range chunks { - if remaining < len(buf) { - buf = buf[:remaining] - } - _, err = dst.Write(buf) - if err != nil { - return - } - remaining -= len(buf) - if remaining == 0 { - break - } - } - return -} +// func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) { +// remaining := size +// for _, buf := range chunks { +// if remaining < len(buf) { +// buf = buf[:remaining] +// } +// _, err = dst.Write(buf) +// if err != nil { +// return +// } +// remaining -= len(buf) +// if remaining == 0 { +// break +// } +// } +// return +// }