diff --git a/atree/atree.go b/atree/atree.go index 4d30d65..b53d781 100644 --- a/atree/atree.go +++ b/atree/atree.go @@ -3,12 +3,12 @@ package atree import ( "errors" "fmt" - "hash/crc32" "os" "path/filepath" "sync" "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/util" ) const ( @@ -52,25 +52,12 @@ const ( FlagNewRoot byte = 2 // новая страница ) -var ( - castagnoliTable = crc32.MakeTable(crc32.Castagnoli) -) - -func calcChecksum(page []byte) uint32 { - return crc32.Checksum(page, castagnoliTable) -} - type PageToWrite struct { PageNo uint32 Data []byte IsReused bool } -type FreeList interface { - // використовується в allocPage - ReservePage() uint32 -} - type _page struct { PageNo uint32 Buf []byte @@ -80,7 +67,6 @@ type _page struct { type Atree struct { file *os.File mutex sync.Mutex - freelist FreeList allocatedPagesQty uint32 pages map[uint32]*_page pageWaits map[uint32][]chan readResult @@ -93,7 +79,6 @@ type Atree struct { type Options struct { Dir string DatabaseName string - FreeList FreeList } func New(opt Options) (*Atree, error) { @@ -106,9 +91,7 @@ func New(opt Options) (*Atree, error) { if opt.DatabaseName == "" { return nil, errors.New("DatabaseName option is required") } - if opt.FreeList == nil { - return nil, errors.New("FreeList option is required") - } + // открываю или создаю dbName.data и dbName.index файлы var ( fileName = filepath.Join(opt.Dir, opt.DatabaseName+".db") @@ -136,7 +119,6 @@ func New(opt Options) (*Atree, error) { } tree := &Atree{ - freelist: opt.FreeList, file: file, allocatedPagesQty: allocatedPagesQty, pages: make(map[uint32]*_page), @@ -261,7 +243,7 @@ type NotLinkedDataPage struct { func (s NotLinkedDataPage) SetPrevPageNo(prevPageNo uint32) { bin.PutUint32(s.Data[prevPageIdx:], prevPageNo) - bin.PutUint32(s.Data[crc32Idx:], calcChecksum(s.Data[:crc32Idx])) + bin.PutUint32(s.Data[crc32Idx:], util.CalcChecksum(s.Data[:crc32Idx])) } type AppendDataPagesReq struct { diff --git a/atree/io.go b/atree/io.go index 714036f..d4d6f8e 100644 --- a/atree/io.go +++ b/atree/io.go @@ -9,6 +9,7 @@ import ( "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/util" ) // fix - додати ID, щоб потім перемістити із тимчасового буфера в pages, або звільнити @@ -362,7 +363,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 = calcChecksum(data[:pos]) + calculatedCRC = util.CalcChecksum(data[:pos]) storedCRC = bin.GetUint32(data[pos:]) ) if calculatedCRC != storedCRC { diff --git a/atree/x.go b/atree/x.go index d9c2072..8c3a88f 100644 --- a/atree/x.go +++ b/atree/x.go @@ -1,6 +1,9 @@ package atree -import "gordenko.dev/dima/qb/bin" +import ( + "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/util" +) const ( payloadSize = 24 @@ -16,6 +19,7 @@ type DataPage struct { ValuesSize int PageNo uint32 IsReused bool + Checksum uint32 } type IndexPage struct { @@ -23,14 +27,15 @@ type IndexPage struct { Data []byte PageNo uint32 IsReused bool + Checksum uint32 } type IndexLevel struct { // вже записані дані (лише для першої в списку індексної сторінки, беремо із _metric) // дані потрібні, якщо сторінка буде заповнена і піде на запис в .index файл - Offset int // (заповнені одразу) - LowerTimestamp uint32 // (заповнені одразу) - Data []byte // поточні дані index сторінки (заповнені одразу) - Filled []IndexPage // пусто + Offset int // (заповнені одразу) + LowerTimestamp uint32 // (заповнені одразу) + Data []byte // поточні дані index сторінки (заповнені одразу) + Filled []*IndexPage // пусто } func appendIndexRecord(data []byte, timestamp uint32, pageNo uint32) []byte { @@ -46,12 +51,13 @@ func appendIndexRecord(data []byte, timestamp uint32, pageNo uint32) []byte { type ChunksToDataPageReq struct { PrevPageNo uint32 Timestamps [][]byte // chunks - TimestampsSize uint16 + TimestampsSize int Values [][]byte // chunks - ValuesSize uint16 + ValuesSize int } -func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) { +// return checksum +func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) (checksum uint32) { var ( remainingSize = int(req.TimestampsSize) pos = 0 @@ -80,25 +86,31 @@ func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) { break } } - bin.PutUint16(buf[timestampsSizeIdx:], req.TimestampsSize) - bin.PutUint16(buf[valuesSizeIdx:], req.ValuesSize) + bin.PutUint16(buf[timestampsSizeIdx:], uint16(req.TimestampsSize)) + bin.PutUint16(buf[valuesSizeIdx:], uint16(req.ValuesSize)) bin.PutUint32(buf[prevPageIdx:], req.PrevPageNo) - bin.PutUint32(buf[dataCRC32Idx:], calcChecksum(buf[:dataCRC32Idx])) + + checksum = util.CalcChecksum(buf[:dataCRC32Idx]) + bin.PutUint32(buf[dataCRC32Idx:], checksum) + return } -func DataToIndexPage(buf []byte, data []byte, lastLevel bool) { +func DataToIndexPage(buf []byte, data []byte, lastLevel bool) (checksum uint32) { copy(buf, data) bin.PutUint16(buf[indexRecordsQtyIdx:], uint16(len(data)/indexRecordSize)) if lastLevel { buf[isLastLevelIdx] = 1 } - bin.PutUint32(buf[indexCRC32Idx:], calcChecksum(buf[:indexCRC32Idx])) + checksum = util.CalcChecksum(buf[:indexCRC32Idx]) + bin.PutUint32(buf[indexCRC32Idx:], checksum) + return } type Pager interface { GetPageNumber() (uint32, bool) } +// fix - reduce leves func AppendDataPagesToTree(pager Pager, levels []*IndexLevel, dataPages []*DataPage) { for i, d := range dataPages[1:] { d.PageNo, d.IsReused = pager.GetPageNumber() @@ -127,7 +139,7 @@ func AppendDataPagesToTree(pager Pager, levels []*IndexLevel, dataPages []*DataP } // в буфері немає місця для додавання нової пари, отже це // заповнена сторінка - filled := IndexPage{ + filled := &IndexPage{ LowerTimestamp: level.LowerTimestamp, Data: level.Data, } diff --git a/chunkenc/cumdelta.go b/chunkenc/cumdelta.go index 4befc22..84bad8c 100644 --- a/chunkenc/cumdelta.go +++ b/chunkenc/cumdelta.go @@ -76,8 +76,9 @@ func (s *ReverseCumulativeDeltaCompressor) Size() int { return s.pos } -//func (s *ReverseCumulativeDeltaCompressor) CalcRequiredSpace(value float64) int { -//} +func (s *ReverseCumulativeDeltaCompressor) CalcRequiredSpace(value float64) int { + return 0 +} func (s *ReverseCumulativeDeltaCompressor) Append(value float64) { if s.pos == 0 { diff --git a/chunkenc/insdelta.go b/chunkenc/insdelta.go index bfa4033..fa54cee 100644 --- a/chunkenc/insdelta.go +++ b/chunkenc/insdelta.go @@ -157,6 +157,37 @@ func (s *ReverseInstantDeltaCompressor) GetState() InstantDeltaBound { } } +func (s *ReverseInstantDeltaCompressor) Lock() { + // s.state = &CumulativeDeltaBound{ + // Pos: s.pos - 1 - s.lastDeltaSize, + // H: s.h, + // LastDelta: s.lastDelta, + // Chunks: s.buf.Chunks(), + // } +} + +func (s *ReverseInstantDeltaCompressor) Unlock() { + //s.state = nil +} + +func (s *ReverseInstantDeltaCompressor) Renew() { + s.buf = conbuf.New(nil) + s.pos = 0 + // + s.baseValue = 0 + s.lastDelta = 0 + s.lastDeltaSize = 0 + s.h = 0 +} + +func (s *ReverseInstantDeltaCompressor) CalcRequiredSpace(value float64) int { + return 0 +} + +func (s *ReverseInstantDeltaCompressor) Chunks() [][]byte { + return s.buf.Chunks() +} + // DECOMPRESSOR type ReverseInstantDeltaDecompressor struct { diff --git a/chunkenc/time_delta.go b/chunkenc/time_delta.go index 255e59f..6bd8b49 100644 --- a/chunkenc/time_delta.go +++ b/chunkenc/time_delta.go @@ -45,8 +45,13 @@ func (s *ReverseTimeDeltaCompressor) Size() int { return s.pos } -//func (s *ReverseTimeDeltaCompressor) CalcRequiredSpace(value float64) int { -//} +func (s *ReverseTimeDeltaCompressor) Chunks() [][]byte { + return s.buf.Chunks() +} + +func (s *ReverseTimeDeltaCompressor) CalcRequiredSpace(value uint32) int { + return 0 +} func (s *ReverseTimeDeltaCompressor) Append(unixtime uint32) { if s.lastUnixtime > 0 { @@ -154,7 +159,7 @@ type TimeDeltaBound struct { } // delta h -func (s *ReverseTimeDeltaCompressor) GetState() TimeDeltaBound { +func (s *ReverseTimeDeltaCompressor) GetState() TimeDeltaBound { // fix replace by Lock bound := TimeDeltaBound{ LastUnixtime: s.lastUnixtime, } @@ -166,6 +171,24 @@ func (s *ReverseTimeDeltaCompressor) GetState() TimeDeltaBound { return bound } +func (s *ReverseTimeDeltaCompressor) Lock() { + // fix +} + +func (s *ReverseTimeDeltaCompressor) Unlock() { + //s.state = nil // fix +} + +func (s *ReverseTimeDeltaCompressor) Renew() { + // s.buf = conbuf.New(nil) + // s.pos = 0 + // // + // s.baseValue = 0 + // s.lastDelta = 0 + // s.lastDeltaSize = 0 + // s.h = 0 +} + // DECOMPRESSOR type ReverseTimeDeltaDecompressor struct { diff --git a/database/api.go b/database/api.go index 590d109..378d97a 100644 --- a/database/api.go +++ b/database/api.go @@ -6,7 +6,7 @@ import ( "io" "net" - diploma "gordenko.dev/dima/qb" + qb "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/bufreader" @@ -192,11 +192,11 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { if req.MetricID == 0 { return proto.ErrEmptyMetricID } - if byte(req.FracDigits) > diploma.MaxFracDigits { + if byte(req.FracDigits) > qb.MaxFracDigits { return proto.ErrWrongFracDigits } switch req.MetricType { - case diploma.Cumulative, diploma.Instant: + case qb.Cumulative, qb.Instant: // ok default: return proto.ErrWrongMetricType @@ -213,24 +213,24 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { switch resultCode { case Succeed: - waitCh := s.txlog.Append(txlog.AddedMetric{ + s.txlog.Append(txlog.AddedMetric{ MetricID: req.MetricID, MetricType: req.MetricType, FracDigits: req.FracDigits, }) - <-waitCh + //<-waitCh case MetricDuplicate: return proto.ErrDuplicate default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 } type Metric struct { - MetricType diploma.MetricType + MetricType qb.MetricType FracDigits byte ResultCode byte } @@ -264,7 +264,7 @@ func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { reply(conn, proto.ErrNoMetric) default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return nil } @@ -293,20 +293,20 @@ func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 { var err error freePageNumbers, err = s.atree.GetAllPages(result.RootPageNo) if err != nil { - diploma.Abort(diploma.FailedAtreeRequest, err) + qb.Abort(qb.FailedAtreeRequest, err) } } - waitCh := s.txlog.Append(txlog.DeletedMetric{ + s.txlog.Append(txlog.DeletedMetric{ MetricID: req.MetricID, FreePageNumbers: freePageNumbers, }) - <-waitCh + //<-waitCh case NoMetric: return proto.ErrNoMetric default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 } @@ -358,15 +358,15 @@ type FilledPage struct { // path, err := s.atree.FindPathToLastPage(filled.RootPageNo) // if err != nil { // // FIX -// diploma.Abort(diploma.WriteToAtreeFailed, err) +// qb.Abort(qb.WriteToAtreeFailed, err) // } // // for _, leg := range path.Legs { // // pagesToRelease = append(pagesToRelease, leg.PageNo) // // } // if path.LastPageNo != filled.PrevPageNo { -// diploma.Abort( -// diploma.WrongPrevPageNo, +// qb.Abort( +// qb.WrongPrevPageNo, // fmt.Errorf("bug: last pageNo %d in tree != prev pageNo %d in _metric", // path.LastPageNo, filled.PrevPageNo), // ) @@ -386,7 +386,7 @@ type FilledPage struct { // // ValuesSize: filled.ValuesSize, // // }) // // if err != nil { -// // diploma.Abort(diploma.WriteToAtreeFailed, err) +// // qb.Abort(qb.WriteToAtreeFailed, err) // // } // // _ = report @@ -417,7 +417,7 @@ type FilledPage struct { // return proto.ErrNonMonotonicValue // default: -// diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) +// qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) // } // return 0 // } @@ -425,7 +425,7 @@ type FilledPage struct { type tryAppendMeasuresResult struct { ResultCode byte Written int - // MetricType diploma.MetricType + // MetricType qb.MetricType // FracDigits byte // Since uint32 // Until uint32 @@ -434,8 +434,8 @@ type tryAppendMeasuresResult struct { // PrevPageNo uint32 // TimestampsBuf *conbuf.ContinuousBuffer // ValuesBuf *conbuf.ContinuousBuffer - // Timestamps diploma.TimestampCompressor - // Values diploma.ValueCompressor + // Timestamps qb.TimestampCompressor + // Values qb.ValueCompressor } func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { @@ -454,7 +454,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { // report, err := s.atree.AppendDataPage(atree.AppendDataPageReq{}) // if err != nil { - // diploma.Abort(diploma.WriteToAtreeFailed, err) + // qb.Abort(qb.WriteToAtreeFailed, err) // } // _ = report @@ -473,7 +473,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { return proto.ErrNoMetric default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 } @@ -500,29 +500,29 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { case DeleteFromAtreeNotNeeded: // регистрирую удаление в TransactionLog - waitCh := s.txlog.Append(txlog.DeletedMeasures{ + s.txlog.Append(txlog.DeletedMeasures{ MetricID: req.MetricID, }) - <-waitCh + //<-waitCh case DeleteFromAtreeRequired: // собираю номера всех data и index страниц метрики (типа запись REDO лога). pageNumbers, err := s.atree.GetAllPages(req.MetricID) if err != nil { - diploma.Abort(diploma.FailedAtreeRequest, err) + qb.Abort(qb.FailedAtreeRequest, err) } // регистрирую удаление в TransactionLog - waitCh := s.txlog.Append(txlog.DeletedMeasures{ + s.txlog.Append(txlog.DeletedMeasures{ MetricID: req.MetricID, FreePageNumbers: pageNumbers, }) - <-waitCh + //<-waitCh case NoMetric: return proto.ErrNoMetric default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 } @@ -540,7 +540,7 @@ func (s *Database) ListAllInstantMeasures(conn net.Conn, req proto.ListAllInstan return s.fullScan(fullScanReq{ MetricID: req.MetricID, - MetricType: diploma.Instant, + MetricType: qb.Instant, Conn: conn, ResponseWriter: responseWriter, }) @@ -551,7 +551,7 @@ func (s *Database) ListAllCumulativeMeasures(conn io.Writer, req proto.ListAllCu return s.fullScan(fullScanReq{ MetricID: req.MetricID, - MetricType: diploma.Cumulative, + MetricType: qb.Cumulative, Conn: conn, ResponseWriter: responseWriter, }) @@ -567,7 +567,7 @@ func (s *Database) ListInstantMeasures(conn net.Conn, req proto.ListInstantMeasu return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, - MetricType: diploma.Instant, + MetricType: qb.Instant, Since: req.Since, Until: req.Until - 1, Conn: conn, @@ -585,7 +585,7 @@ func (s *Database) ListCumulativeMeasures(conn net.Conn, req proto.ListCumulativ return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, - MetricType: diploma.Cumulative, + MetricType: qb.Cumulative, Since: req.Since, Until: req.Until - 1, Conn: conn, @@ -620,7 +620,7 @@ func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriod return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, - MetricType: diploma.Instant, + MetricType: qb.Instant, Since: uint32(since.Unix()), Until: uint32(until.Unix()), Conn: conn, @@ -647,7 +647,7 @@ func (s *Database) ListCumulativePeriods(conn net.Conn, req proto.ListCumulative return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, - MetricType: diploma.Cumulative, + MetricType: qb.Cumulative, Since: uint32(since.Unix()), Until: uint32(until.Unix()), Conn: conn, @@ -657,7 +657,7 @@ func (s *Database) ListCumulativePeriods(conn net.Conn, req proto.ListCumulative type rangeScanReq struct { MetricID uint32 - MetricType diploma.MetricType + MetricType qb.MetricType Since uint32 Until uint32 Conn io.Writer @@ -720,14 +720,14 @@ func (s *Database) rangeScan(req rangeScanReq) error { reply(req.Conn, proto.ErrWrongMetricType) default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return nil } type fullScanReq struct { MetricID uint32 - MetricType diploma.MetricType + MetricType qb.MetricType Conn io.Writer ResponseWriter atree.PeriodsWriter } @@ -768,7 +768,7 @@ func (s *Database) fullScan(req fullScanReq) error { reply(req.Conn, proto.ErrWrongMetricType) default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) + qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return nil } diff --git a/database/database.go b/database/database.go index 4c49794..77c2702 100644 --- a/database/database.go +++ b/database/database.go @@ -12,11 +12,9 @@ import ( "sync" "time" - diploma "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/bin" - "gordenko.dev/dima/qb/chunkenc" - "gordenko.dev/dima/qb/conbuf" "gordenko.dev/dima/qb/freelist" "gordenko.dev/dima/qb/recovery" "gordenko.dev/dima/qb/txlog" @@ -84,13 +82,25 @@ func New(opt Options) (_ *Database, err error) { return nil, errors.New("WaitGroup option is required") } + filePath := filepath.Join(opt.Dir, opt.DatabaseName+".free") + + freeListFile, err := os.OpenFile(filePath, os.O_RDWR|os.O_CREATE, 0666) + if err != nil { + return nil, err + } + + freeList, err := freelist.New(freelist.Options{ + PageSize: 2048, + File: freeListFile, + }) + s := &Database{ workerSignalCh: make(chan struct{}, 1), dir: opt.Dir, databaseName: opt.DatabaseName, metrics: make(map[uint32]*_metric), //metricLockEntries: make(map[uint32]*metricLockEntry), - freeList: freelist.New(), + freeList: freeList, tcpPort: opt.TCPPort, logfile: opt.Logfile, logger: log.New(opt.Logfile, "", log.LstdFlags), @@ -143,7 +153,7 @@ func (s *Database) recovery() { recipe, err := advisor.GetRecipe() if err != nil { - diploma.Abort(diploma.GetRecoveryRecipeFailed, err) + qb.Abort(qb.GetRecoveryRecipeFailed, err) } var logNumber int @@ -152,13 +162,13 @@ func (s *Database) recovery() { if recipe.Snapshot != "" { err = s.loadSnapshot(recipe.Snapshot) if err != nil { - diploma.Abort(diploma.LoadSnapshotFailed, err) + qb.Abort(qb.LoadSnapshotFailed, err) } } for _, changesFileName := range recipe.Changes { err = s.replayChanges(changesFileName) if err != nil { - diploma.Abort(diploma.ReplayChangesFailed, err) + qb.Abort(qb.ReplayChangesFailed, err) } } logNumber = recipe.LogNumber @@ -174,28 +184,28 @@ func (s *Database) recovery() { WaitGroup: s.waitGroup, }) if err != nil { - diploma.Abort(diploma.CreateChangesWriterFailed, err) + qb.Abort(qb.CreateChangesWriterFailed, err) } go s.txlog.Run() // fileNames, err := s.searchREDOFiles() // if err != nil { - // diploma.Abort(diploma.SearchREDOFilesFailed, err) + // qb.Abort(qb.SearchREDOFilesFailed, err) // } // if len(fileNames) > 0 { // for _, fileName := range fileNames { // err = s.replayREDOFile(fileName) // if err != nil { - // diploma.Abort(diploma.ReplayREDOFileFailed, err) + // qb.Abort(qb.ReplayREDOFileFailed, err) // } // } // for _, fileName := range fileNames { // err = os.Remove(fileName) // if err != nil { - // diploma.Abort(diploma.RemoveREDOFileFailed, err) + // qb.Abort(qb.RemoveREDOFileFailed, err) // } // } // } @@ -204,14 +214,14 @@ func (s *Database) recovery() { if recipe.CompleteSnapshot { err = s.dumpSnapshot(logNumber) if err != nil { - diploma.Abort(diploma.DumpSnapshotFailed, err) + qb.Abort(qb.DumpSnapshotFailed, err) } } for _, fileName := range recipe.ToDelete { err = os.Remove(fileName) if err != nil { - diploma.Abort(diploma.RemoveRecipeFileFailed, err) + qb.Abort(qb.RemoveRecipeFileFailed, err) } } } @@ -348,100 +358,100 @@ func (s *Database) replayChanges(fileName string) error { } func (s *Database) replayChangesRecord(untyped any) error { - switch rec := untyped.(type) { - case txlog.AddedMetric: - var ( - values diploma.ValueCompressor - timestampsBuf = conbuf.New(nil) - valuesBuf = conbuf.New(nil) - ) + // switch rec := untyped.(type) { + // case txlog.AddedMetric: + // var ( + // values qb.ValueCompressor + // timestampsBuf = conbuf.New(nil) + // valuesBuf = conbuf.New(nil) + // ) - if rec.MetricType == diploma.Cumulative { - values = chunkenc.NewReverseCumulativeDeltaCompressor( - valuesBuf, 0, byte(rec.FracDigits)) - } else { - values = chunkenc.NewReverseInstantDeltaCompressor( - valuesBuf, 0, byte(rec.FracDigits)) - } - - s.metrics[rec.MetricID] = &_metric{ - MetricType: rec.MetricType, - FracDigits: byte(rec.FracDigits), - TimestampsBuf: timestampsBuf, - ValuesBuf: valuesBuf, - Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), - Values: values, - } - - case txlog.DeletedMetric: - delete(s.metrics, rec.MetricID) - if len(rec.FreePageNumbers) > 0 { - s.freeList.AddPages(rec.FreePageNumbers) - } - - case txlog.AppendedMeasure: - metric, ok := s.metrics[rec.MetricID] - if ok { - metric.Timestamps.Append(rec.Timestamp) - metric.Values.Append(rec.Value) - - if metric.Since == 0 { - metric.Since = rec.Timestamp - metric.SinceValue = rec.Value - } - - metric.Until = rec.Timestamp - metric.UntilValue = rec.Value - } - - case txlog.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 + // if rec.MetricType == qb.Cumulative { + // values = chunkenc.NewReverseCumulativeDeltaCompressor( + // valuesBuf, 0, byte(rec.FracDigits)) + // } else { + // values = chunkenc.NewReverseInstantDeltaCompressor( + // valuesBuf, 0, byte(rec.FracDigits)) // } - // metric.LastPageNo = rec.DataPageNo - // // delete free pages - // if rec.IsDataPageReused { - // s.freeList.DeleteReservedPages([]uint32{ - // rec.DataPageNo, - // }) + + // s.metrics[rec.MetricID] = &_metric{ + // MetricType: rec.MetricType, + // FracDigits: byte(rec.FracDigits), + // TimestampsBuf: timestampsBuf, + // ValuesBuf: valuesBuf, + // Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), + // Values: values, // } - // if len(rec.ReusedIndexPages) > 0 { - // s.freeList.DeleteReservedPages(rec.ReusedIndexPages) + + // case txlog.DeletedMetric: + // delete(s.metrics, rec.MetricID) + // if len(rec.FreePageNumbers) > 0 { + // s.freeList.AddPages(rec.FreePageNumbers) // } + + // case txlog.AppendedMeasure: + // metric, ok := s.metrics[rec.MetricID] + // if ok { + // metric.Timestamps.Append(rec.Timestamp) + // metric.Values.Append(rec.Value) + + // if metric.Since == 0 { + // metric.Since = rec.Timestamp + // metric.SinceValue = rec.Value + // } + + // metric.Until = rec.Timestamp + // metric.UntilValue = rec.Value + // } + + // case txlog.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) + // } + // } + + // default: + // qb.Abort(qb.UnknownTxLogRecordTypeBug, + // fmt.Errorf("bug: unknown record type %T in TransactionLog", rec)) // } - - case txlog.DeletedMeasures: - metric, ok := s.metrics[rec.MetricID] - if ok { - metric.DeleteMeasures() - if len(rec.FreePageNumbers) > 0 { - s.freeList.AddPages(rec.FreePageNumbers) - } - } - - default: - diploma.Abort(diploma.UnknownTxLogRecordTypeBug, - fmt.Errorf("bug: unknown record type %T in TransactionLog", rec)) - } return nil } diff --git a/database/metric.go b/database/metric.go index 67c0019..b0a6337 100644 --- a/database/metric.go +++ b/database/metric.go @@ -2,7 +2,6 @@ package database import ( "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/conbuf" ) // METRIC @@ -10,59 +9,32 @@ import ( type _metric struct { MetricType qb.MetricType FracDigits byte - RootPageNo uint32 LastPageNo uint32 SinceValue float64 Since uint32 UntilValue float64 Until uint32 - //TimestampsBuf *conbuf.ContinuousBuffer - //ValuesBuf *conbuf.ContinuousBuffer Timestamps qb.TimestampCompressor Values qb.ValueCompressor } func (s *_metric) ReinitBy(timestamp uint32, value float64) { - // s.TimestampsBuf = conbuf.New(nil) - // s.ValuesBuf = conbuf.New(nil) - // // - // s.Timestamps = chunkenc.NewReverseTimeDeltaOfDeltaCompressor( - // s.TimestampsBuf, 0) + s.Timestamps.Renew() + s.Values.Renew() - // if s.MetricType == octopus.Cumulative { - // s.Values = chunkenc.NewReverseCumulativeDeltaCompressor( - // s.ValuesBuf, 0, s.FracDigits) - // } else { - // s.Values = chunkenc.NewReverseInstantDeltaCompressor( - // s.ValuesBuf, 0, s.FracDigits) - // } + s.Timestamps.Append(timestamp) + s.Values.Append(value) - // s.Timestamps.Append(timestamp) - // s.Values.Append(value) - - // s.Since = timestamp - // s.SinceValue = value - // s.Until = timestamp - // s.UntilValue = value + s.Since = timestamp + s.SinceValue = value + s.Until = timestamp + s.UntilValue = value } func (s *_metric) DeleteMeasures() { - s.TimestampsBuf = conbuf.New(nil) - s.ValuesBuf = conbuf.New(nil) - // fix - // s.Timestamps = chunkenc.NewReverseTimeDeltaOfDeltaCompressor( - // s.TimestampsBuf, 0) + s.Timestamps.Renew() + s.Values.Renew() - // if s.MetricType == octopus.Cumulative { - // s.Values = chunkenc.NewReverseCumulativeDeltaCompressor( - // s.ValuesBuf, 0, s.FracDigits) - // } else { - // s.Values = chunkenc.NewReverseInstantDeltaCompressor( - // s.ValuesBuf, 0, s.FracDigits) - // } - // END OF FIX - - s.RootPageNo = 0 s.LastPageNo = 0 s.Since = 0 s.SinceValue = 0 diff --git a/database/snapshot.go b/database/snapshot.go index 91987f9..e4a7df2 100644 --- a/database/snapshot.go +++ b/database/snapshot.go @@ -7,21 +7,22 @@ import ( "os" "path/filepath" - octopus "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/bin" + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/chunkenc" + "gordenko.dev/dima/qb/conbuf" "gordenko.dev/dima/qb/freelist" ) /* Формат: -//lsn - varuint (останній LSN, що змінив дані у RAM) metricsQty - varuint [metric]* -где metric - це: +де metric - це: metricID - 4b metricType - 1b fracDigits - 1b -rootPageNo - 4b lastPageNo - 4b since - 4b sinceValue - 8b @@ -31,14 +32,14 @@ timestamps size - 2b values size - 2b timestams payload - Nb values payload - Nb -dataFreeList size - varuint +dataFreeList size - varsize dataFreeList - Nb -indexFreeList size - varuint +indexFreeList size - varsize indexFreeList - Nb CRC32 - 4b */ -const metricHeaderSize = 42 +const metricHeaderSize = 38 func (s *Database) dumpSnapshot(logNumber int) (err error) { var ( @@ -66,49 +67,23 @@ func (s *Database) dumpSnapshot(logNumber int) (err error) { bin.PutUint32(prefix[0:], metricID) prefix[4] = byte(metric.MetricType) prefix[5] = metric.FracDigits - bin.PutUint32(prefix[6:], metric.RootPageNo) - bin.PutUint32(prefix[10:], metric.LastPageNo) - bin.PutUint32(prefix[14:], metric.Since) - bin.PutFloat64(prefix[18:], metric.SinceValue) - bin.PutUint32(prefix[26:], metric.Until) - bin.PutFloat64(prefix[30:], metric.UntilValue) - bin.PutUint16(prefix[38:], uint16(tSize)) - bin.PutUint16(prefix[40:], uint16(vSize)) + 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 - remaining := tSize - for _, buf := range metric.TimestampsBuf.Chunks() { - if remaining < len(buf) { - buf = buf[:remaining] - } - _, err = dst.Write(buf) - if err != nil { - return - } - remaining -= len(buf) - if remaining == 0 { - break - } - } + writeChunks(dst, metric.Timestamps.Chunks(), tSize) // copy values - remaining = vSize - for _, buf := range metric.ValuesBuf.Chunks() { - if remaining < len(buf) { - buf = buf[:remaining] - } - _, err = dst.Write(buf) - if err != nil { - return - } - remaining -= len(buf) - if remaining == 0 { - break - } - } + writeChunks(dst, metric.Values.Chunks(), vSize) + } // free data pages err = freeListWriteTo(s.freeList, dst) @@ -133,123 +108,124 @@ func (s *Database) dumpSnapshot(logNumber int) (err error) { return } - prevLogNumber := logNumber - 1 - prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber)) - prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber)) + // 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 - } + // isExist, err := isFileExist(prevChanges) + // if err != nil { + // return + // } - if isExist { - err = os.Remove(prevChanges) - if err != nil { - octopus.Abort(octopus.DeletePrevChangesFileFailed, err) - } - } + // if isExist { + // err = os.Remove(prevChanges) + // if err != nil { + // qb.Abort(qb.DeletePrevChangesFileFailed, err) + // } + // } - isExist, err = isFileExist(prevSnapshot) - if err != nil { - return - } + // isExist, err = isFileExist(prevSnapshot) + // if err != nil { + // return + // } - if isExist { - err = os.Remove(prevSnapshot) - if err != nil { - octopus.Abort(octopus.DeletePrevSnapshotFileFailed, err) - } - } + // if isExist { + // err = os.Remove(prevSnapshot) + // if err != nil { + // qb.Abort(qb.DeletePrevSnapshotFileFailed, err) + // } + // } return } func (s *Database) loadSnapshot(fileName string) (err error) { - // FIX - // var ( - // hasher = crc32.NewIEEE() - // metricsQty int - // header = make([]byte, metricHeaderSize) - // body = make([]byte, atree.PageSize) - // ) + var ( + hasher = crc32.NewIEEE() + header = make([]byte, metricHeaderSize) + body = make([]byte, atree.PageSize) + ) - // file, err := os.Open(fileName) - // if err != nil { - // return - // } + file, err := os.Open(fileName) + if err != nil { + return + } - // src := io.TeeReader(file, hasher) - // u64, _, err := bin.ReadVarUint64(src) - // if err != nil { - // return - // } - // metricsQty = int(u64) + 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 - // } + for range metricsQty { + var metric _metric + err = bin.ReadNInto(src, header) + if err != nil { + return + } - // metricID := bin.GetUint32(header[0:]) - // metric.MetricType = octopus.MetricType(header[4]) - // metric.FracDigits = header[5] - // metric.RootPageNo = bin.GetUint32(header[6:]) - // metric.LastPageNo = bin.GetUint32(header[10:]) - // metric.Since = bin.GetUint32(header[14:]) - // metric.SinceValue = bin.GetFloat64(header[18:]) - // metric.Until = bin.GetUint32(header[26:]) - // metric.UntilValue = bin.GetFloat64(header[30:]) - // tSize := bin.GetUint16(header[38:]) - // vSize := bin.GetUint16(header[40:]) + 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 - // } - // metric.TimestampsBuf = conbuf.NewFromBuffer(buf) + buf := body[:tSize] + err = bin.ReadNInto(src, buf) + if err != nil { + return + } + metric.Timestamps = chunkenc.NewReverseTimeDeltaCompressor( + conbuf.NewFromBuffer(buf), + int(tSize), + ) - // buf = body[:vSize] - // err = bin.ReadNInto(src, buf) - // if err != nil { - // return - // } - // metric.ValuesBuf = conbuf.NewFromBuffer(buf) + buf = body[:vSize] + err = bin.ReadNInto(src, buf) + if err != nil { + return + } - // metric.Timestamps = chunkenc.NewReverseTimeDeltaOfDeltaCompressor( - // metric.TimestampsBuf, int(tSize)) + 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 + } - // if metric.MetricType == octopus.Cumulative { - // metric.Values = chunkenc.NewReverseCumulativeDeltaCompressor( - // metric.ValuesBuf, int(vSize), metric.FracDigits) - // } else { - // metric.Values = chunkenc.NewReverseInstantDeltaCompressor( - // metric.ValuesBuf, int(vSize), metric.FracDigits) - // } - // s.metrics[metricID] = &metric - // } + err = restoreFreeList(s.freeList, src) + if err != nil { + return fmt.Errorf("restore dataFreeList: %s", err) + } - // err = restoreFreeList(s.freeList, src) - // if err != nil { - // return fmt.Errorf("restore dataFreeList: %s", err) - // } + err = restoreFreeList(s.freeList, src) + if err != nil { + return fmt.Errorf("restore indexFreeList: %s", err) + } - // err = restoreFreeList(s.freeList, src) - // if err != nil { - // return fmt.Errorf("restore indexFreeList: %s", err) - // } + calculatedChecksum := hasher.Sum32() - // calculatedChecksum := hasher.Sum32() + writtenChecksum, err := bin.ReadUint32(file) + if err != nil { + return + } - // writtenChecksum, err := bin.ReadUint32(file) - // if err != nil { - // return - // } - - // if calculatedChecksum != writtenChecksum { - // return fmt.Errorf("calculated checksum %d not equal written checksum %d", calculatedChecksum, writtenChecksum) - // } + if calculatedChecksum != writtenChecksum { + return fmt.Errorf("calculated checksum %d not equal written checksum %d", calculatedChecksum, writtenChecksum) + } return } @@ -258,9 +234,9 @@ func (s *Database) loadSnapshot(fileName string) (err error) { func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error { serialized, err := freeList.Serialize() if err != nil { - octopus.Abort(octopus.FailedFreeListSerialize, err) + qb.Abort(qb.FailedFreeListSerialize, err) } - _, err = bin.WriteVarUint64(dst, uint64(len(serialized))) + _, err = bin.WriteVarSize(dst, len(serialized)) if err != nil { return err } @@ -272,14 +248,32 @@ func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error { } func restoreFreeList(freeList *freelist.FreeList, src io.Reader) error { - size, _, err := bin.ReadVarUint64(src) + size, err := bin.ReadVarSize(src) if err != nil { return err } - serialized, err := bin.ReadN(src, int(size)) + serialized, err := bin.ReadN(src, size) if err != nil { return err } freeList.Restore(serialized) 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 +} diff --git a/freelist/freelist.go b/freelist/freelist.go index 4b68fcf..88dc273 100644 --- a/freelist/freelist.go +++ b/freelist/freelist.go @@ -4,37 +4,139 @@ import ( "bytes" "fmt" "os" - "slices" - "sync" "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/util" ) const ( - pageSize = 1024 - crcSize = 4 - ptrSize = 4 - pointersOnPage = (pageSize - crcSize) / ptrSize - flushTreshold = 1000 //= int(float(pageSize/ptrSize) * 1.3) // кількість вільних сторінок коли вже треба писати на диск + crcSize = 4 + ptrSize = 4 ) type FreeList struct { - mutex sync.Mutex - file *os.File - pages int - //free *roaring.Bitmap - //reserved *roaring.Bitmap - free []uint32 - reserved []uint32 + file *os.File + pageSize int + pointersOnPage int + flushTreshold int // кількість вільних сторінок коли вже треба писати на диск + pages int + free []uint32 } -func New() *FreeList { - return &FreeList{ - //free: roaring.New(), - //reserved: roaring.New(), +type Options struct { + PageSize int + File *os.File +} + +func New(opt Options) (*FreeList, error) { + s := &FreeList{ + file: opt.File, + pageSize: opt.PageSize, + pointersOnPage: (opt.PageSize - crcSize) / ptrSize, } + s.flushTreshold = int(float64(s.pointersOnPage) * 1.3) + info, err := s.file.Stat() + if err != nil { + return nil, err + } + fileSize := info.Size() + if fileSize > 0 { + if (fileSize % int64(s.pageSize)) > 0 { + return nil, fmt.Errorf("file size %d is not a multiple of page size %d", fileSize, s.pageSize) + } + s.pages = int(fileSize) / s.pageSize + } + return s, nil } +func (s *FreeList) AddPageNumbers(pageNumbers []uint32) (err error) { + if len(pageNumbers) == 0 { + return + } + s.free = append(s.free, pageNumbers...) + + for len(s.free) > s.flushTreshold { + err = s.save() + if err != nil { + return + } + // викидаю номери сторінок, які записав на диск + copy(s.free, s.free[s.pointersOnPage:]) + s.free = s.free[:len(s.free)-s.pointersOnPage] + } + return +} + +func (s *FreeList) save() error { + buf := make([]byte, s.pageSize) + i := crcSize + for _, pageNo := range s.free[:s.pointersOnPage] { + bin.PutUint32(buf[i:], pageNo) + i += ptrSize + } + bin.PutUint32(buf, util.CalcChecksum(buf[crcSize:])) + // + off := int64(s.pages * s.pageSize) + n, err := s.file.WriteAt(buf, off) + if err != nil { + return err + } + if n != s.pageSize { + return fmt.Errorf("written size %d bytes not equal page size %d", n, s.pageSize) + } + err = s.file.Sync() + if err != nil { + return err + } + s.pages++ + return nil +} + +func (s *FreeList) loadLastPage() error { + var ( + buf = make([]byte, s.pageSize) + off = int64((s.pages - 1) * s.pageSize) + ) + n, err := s.file.ReadAt(buf, off) + if err != nil { + return fmt.Errorf("file.Seek: %s", err) + } + if n != s.pageSize { + return fmt.Errorf("read size %d bytes not equal page size %d", n, s.pageSize) + } + // check crc + checksum := bin.GetUint32(buf[0:]) + if util.CalcChecksum(buf[crcSize:]) != checksum { + return fmt.Errorf("page %d is corrupted", s.pages) + } + for i := crcSize; i < len(buf); i += ptrSize { + s.free = append(s.free, bin.GetUint32(buf[i:])) + } + s.pages-- + err = s.file.Truncate(int64(s.pages * s.pageSize)) + if err != nil { + return err + } + return nil +} + +func (s *FreeList) GetPageNumber() (pageNo uint32, err error) { + if len(s.free) == 0 { + if s.pages == 0 { + return + } + err = s.loadLastPage() + if err != nil { + return + } + } + lastIdx := len(s.free) - 1 + pageNo = s.free[lastIdx] + s.free = s.free[:lastIdx] + return +} + +// для відновлення недозаповненнох сторінки із снапшота func (s *FreeList) Restore(serialized []byte) error { if (len(serialized) % ptrSize) != 0 { return fmt.Errorf("wrong size") @@ -45,112 +147,11 @@ func (s *FreeList) Restore(serialized []byte) error { return nil } -func (s *FreeList) AddPages(pageNumbers []uint32) { - if len(pageNumbers) == 0 { - return - } - s.mutex.Lock() - s.free = append(s.free, pageNumbers...) - if len(s.free) > flushTreshold { - err := s.save() - if err != nil { - // ABORT - } - } - s.mutex.Unlock() -} - -func (s *FreeList) save() error { - buf := make([]byte, pageSize) - //w := bytes.NewBuffer(nil) - //s.mutex.Lock() - i := crcSize - for _, pageNo := range s.free[:pointersOnPage] { - bin.PutUint32(buf[i:], pageNo) - i += ptrSize - } - // fix checksum - fileOffset := int64(s.pages-1) * pageSize - n, err := s.file.WriteAt(buf, fileOffset) - if err != nil { - return err - } - if n != pageSize { - return fmt.Errorf("written size %d bytes not equal page size %d", n, pageSize) - } - s.free = s.free[pointersOnPage:] - s.pages++ - return nil -} - -func (s *FreeList) loadLastPage() error { - var ( - buf = make([]byte, pageSize) - offset = int64(s.pages-1) * pageSize - ) - n, err := s.file.ReadAt(buf, offset) - if err != nil { - return fmt.Errorf("file.Seek: %s", err) - } - if n != s.pages { - return fmt.Errorf("read size %d bytes not equal page size %d", n, pageSize) - } - // check crc - //var pageNumbers []uint32 - for i := crcSize; i < len(buf); i += ptrSize { - s.free = append(s.free, bin.GetUint32(buf[i:])) - } - s.pages-- - return nil -} - -// ReserveDataPage - аллокатор резервирует страницу, но не удаляет до визова -// DeleteFromFree, ибо транзакция может не завершится, а между віделением страници -// и падением транзакции - будет создан init файл. -func (s *FreeList) ReservePage() (pageNo uint32) { - s.mutex.Lock() - defer s.mutex.Unlock() - - if len(s.free) == 0 { - if s.pages == 0 { - return - } - if err := s.loadLastPage(); err != nil { - // ABORT - return - } - } - lastIdx := len(s.free) - pageNo = s.free[lastIdx] - s.free = s.free[:lastIdx] - s.reserved = append(s.reserved, pageNo) - return -} - -// Удаляет ранее зарезервированные страницы -func (s *FreeList) DeleteReservedPages(pageNumbers []uint32) { - var cleared []uint32 - s.mutex.Lock() - for _, reservedPageNo := range s.reserved { - if slices.Contains(pageNumbers, reservedPageNo) { - // delete - } else { - cleared = append(cleared, reservedPageNo) - } - } - s.reserved = cleared - s.mutex.Unlock() -} - +// для запису в snapshot func (s *FreeList) Serialize() ([]byte, error) { w := bytes.NewBuffer(nil) - s.mutex.Lock() for _, pageNo := range s.free { bin.WriteUint32(w, pageNo) } - for _, pageNo := range s.reserved { - bin.WriteUint32(w, pageNo) - } - s.mutex.Unlock() return w.Bytes(), nil } diff --git a/freelist/freelist_test.go b/freelist/freelist_test.go new file mode 100644 index 0000000..a3f5883 --- /dev/null +++ b/freelist/freelist_test.go @@ -0,0 +1,44 @@ +package freelist + +import ( + "fmt" + "os" + "testing" +) + +func createFreeList() (_ *FreeList, err error) { + file, err := os.OpenFile("test.free", os.O_CREATE|os.O_RDWR, 0666) + if err != nil { + return + } + return New(Options{ + PageSize: 20, + File: file, + }) +} + +func TestAddToFreeList(t *testing.T) { + freeList, err := createFreeList() + if err != nil { + t.Fatal(err) + } + + pageNumbers := []uint32{1, 2, 3, 4, 5, 6, 7, 8, 9, 10} + + err = freeList.AddPageNumbers(pageNumbers) + if err != nil { + t.Fatal(err) + } +} + +func TestGetFromFreeList(t *testing.T) { + freeList, err := createFreeList() + if err != nil { + t.Fatal(err) + } + + for range 4 { + pageNo, err := freeList.GetPageNumber() + fmt.Println(pageNo, err) + } +} diff --git a/freelist/test.free b/freelist/test.free new file mode 100644 index 0000000..d2ebc93 Binary files /dev/null and b/freelist/test.free differ diff --git a/txlog/helpers.go b/txlog/helpers.go new file mode 100644 index 0000000..08ce092 --- /dev/null +++ b/txlog/helpers.go @@ -0,0 +1,238 @@ +package txlog + +import ( + "bytes" + + bin "gordenko.dev/dima/bin/little" +) + +const recordSize = 8 +const chunkSize = 24 + +/* +Format appended measures: +1b - tx type +4b - metricID +Nb - qty of filled data pages (varsize) +[ + 4b - pageNo + 1b - reused + 4b - prevPageNo + 4b - page crc32 + 2b - timestamps size + Nb - timestamps payload + 2b - values size + Nb - values payload +] +2b - timestamps offset (for the 1st data page only) +2b - values offset (for the 1st data page only) +Nb - timestamps size unfilled (varsize) +Nb - timestamps payload unfilled +Nb - values size unfilled (varsize) +Nb - values payload unfilled +Nb - qty of index levels (varsize) +[ + Nb - qty of level filled pages (varsize) + [ + 4b - pageNo + 1b - reused + 4b - page crc32 + Nb - records qty (varsize) + Nb - records payload + ] + Nb - records qty unfilled (varsize) + Nb - records payload unfilled +] +*/ + +func copyPayloadFromChunks(w *bytes.Buffer, chunks [][]byte, size int, offset int) { + bin.WriteUint16(w, uint16(size)) + var ( + chunkIdx = offset / chunkSize + byteIdx = offset % chunkSize + ) + for _, chunk := range chunks[chunkIdx:] { + available := len(chunk) - byteIdx + if available <= size { + w.Write(chunk[byteIdx:]) + size -= available + if size == 0 { + return + } + byteIdx = 0 + } else { + w.Write(chunk[byteIdx : byteIdx+size]) + return + } + } +} + +///////////////////////////////////////// + +type DataPage struct { + PageNo uint32 + Reused bool + PrevPageNo uint32 + Checksum uint32 + Timestamps []byte + Values []byte +} + +type IndexPage struct { + PageNo uint32 + Reused bool + Checksum uint32 + Records []byte +} + +type IndexLevel struct { + Pages []IndexPage + Records []byte // unfilled +} + +// Для декодінга +type TxAppendedMeasures struct { + MetricID uint32 + DataPages []DataPage + // offsets потрібні тому що дані не додаються в кінець, а перезаписують кілька + // останніх байтів unfilled даних. Потрібно для коректного recovery. + // Якщо є DataPages, то застосовуються для 1-ї Data сторінки, інакше - до timestamps і values (unfilled) + TimestampsOffset int + ValuesOffset int + Timestamps []byte // unfilled + Values []byte // unfilled + IndexLevels []IndexLevel +} + +func (s *TxAppendedMeasures) Read(r *bytes.Buffer) (err error) { + s.MetricID, err = bin.ReadUint32(r) + if err != nil { + return + } + dataPagesQty, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range dataPagesQty { + var ( + p DataPage + timestampsSize, valuesSize int + //b byte + ) + p.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + p.PrevPageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Checksum, err = bin.ReadUint32(r) + if err != nil { + return + } + timestampsSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + p.Timestamps, err = bin.ReadN(r, timestampsSize) + if err != nil { + return + } + valuesSize, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + p.Values, err = bin.ReadN(r, valuesSize) + if err != nil { + return + } + s.DataPages = append(s.DataPages, p) + } + s.TimestampsOffset, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + s.ValuesOffset, err = bin.ReadUint16AsInt(r) + if err != nil { + return + } + var ( + timestampsSize, valuesSize int + ) + timestampsSize, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.Timestamps, err = bin.ReadN(r, timestampsSize) + if err != nil { + return + } + valuesSize, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.Values, err = bin.ReadN(r, valuesSize) + if err != nil { + return + } + if len(s.DataPages) == 0 { + return + } + indexLevelsQty, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range indexLevelsQty { + var ( + indexPagesQty int + level IndexLevel + recordsQty int + ) + indexPagesQty, err = bin.ReadVarSize(r) + if err != nil { + return + } + for range indexPagesQty { + var ( + p IndexPage + recordsQty int + ) + 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 + } + recordsQty, err = bin.ReadVarSize(r) + if err != nil { + return + } + p.Records, err = bin.ReadN(r, recordsQty*recordSize) + if err != nil { + return + } + level.Pages = append(level.Pages, p) + } + recordsQty, err = bin.ReadVarSize(r) + if err != nil { + return + } + level.Records, err = bin.ReadN(r, recordsQty*recordSize) + if err != nil { + return + } + s.IndexLevels = append(s.IndexLevels, level) + } + return +} diff --git a/txlog/reader.go b/txlog/reader.go index 849235c..a6f56e0 100644 --- a/txlog/reader.go +++ b/txlog/reader.go @@ -113,26 +113,26 @@ func (s *Reader) parseRecords(body []byte) ([]any, error) { } records = append(records, rec) - case CodeAppendedMeasure: - rec := new(AppendedMeasure) - if err = rec.Parse(src); err != nil { - return nil, err - } - records = append(records, rec) + // case CodeAppendedMeasure: + // rec := new(AppendedMeasure) + // if err = rec.Parse(src); err != nil { + // return nil, err + // } + // records = append(records, rec) - case CodeAppendedMeasures: - rec := new(AppendedMeasures) - if err = rec.Parse(src); err != nil { - return nil, err - } - records = append(records, rec) + // case CodeAppendedMeasures: + // rec := new(AppendedMeasures) + // if err = rec.Parse(src); err != nil { + // return nil, err + // } + // records = append(records, rec) - case CodeAppendedPages: - rec := new(AppendedPages) - if err = rec.Parse(src); err != nil { - return nil, err - } - records = append(records, rec) + // case CodeAppendedPages: + // rec := new(AppendedPages) + // if err = rec.Parse(src); err != nil { + // return nil, err + // } + // records = append(records, rec) case CodeDeletedMeasures: rec := new(DeletedMeasures) diff --git a/txlog/txlog.go b/txlog/txlog.go index d08ce1b..21e1873 100644 --- a/txlog/txlog.go +++ b/txlog/txlog.go @@ -5,9 +5,7 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/conbuf" - "gordenko.dev/dima/qb/proto" ) type AddedMetric struct { @@ -77,93 +75,93 @@ func (s *DeletedMetric) Parse(src io.Reader) (err error) { return nil } -type AppendedMeasure struct { - MetricID uint32 - Timestamp uint32 - Value float64 -} +// type AppendedMeasure struct { +// MetricID uint32 +// Timestamp uint32 +// Value float64 +// } -func (s AppendedMeasure) Pack(w io.Writer) { - arr := []byte{ - CodeAppendedMeasure, - 0, 0, 0, 0, // metricID - 0, 0, 0, 0, // timestamp - 0, 0, 0, 0, 0, 0, 0, 0, // value - } - bin.PutUint32(arr[1:], s.MetricID) - bin.PutUint32(arr[5:], s.Timestamp) - bin.PutFloat64(arr[9:], s.Value) - w.Write(arr) -} +// func (s AppendedMeasure) Pack(w io.Writer) { +// arr := []byte{ +// CodeAppendedMeasure, +// 0, 0, 0, 0, // metricID +// 0, 0, 0, 0, // timestamp +// 0, 0, 0, 0, 0, 0, 0, 0, // value +// } +// bin.PutUint32(arr[1:], s.MetricID) +// bin.PutUint32(arr[5:], s.Timestamp) +// bin.PutFloat64(arr[9:], s.Value) +// w.Write(arr) +// } -func (s *AppendedMeasure) Parse(src io.Reader) (err error) { - arr, err := bin.ReadN(src, 16) - if err != nil { - return - } - s.MetricID, _ = bin.GetUint32(arr[0:]) - s.Timestamp, _ = bin.GetUint32(arr[4:]) - s.Value, _ = bin.GetFloat64(arr[8:]) - return nil -} +// func (s *AppendedMeasure) Parse(src io.Reader) (err error) { +// arr, err := bin.ReadN(src, 16) +// if err != nil { +// return +// } +// s.MetricID, _ = bin.GetUint32(arr[0:]) +// s.Timestamp, _ = bin.GetUint32(arr[4:]) +// s.Value, _ = bin.GetFloat64(arr[8:]) +// return nil +// } -type AppendedMeasures struct { - MetricID uint32 - Measures []proto.Measure -} +// type AppendedMeasures struct { +// MetricID uint32 +// Measures []proto.Measure +// } -func (s AppendedMeasures) Pack(w io.Writer) { - arr := []byte{ - CodeAppendedMeasures, - 0, 0, 0, 0, // metricID - 0, 0, // qty - } - bin.PutUint32(arr[1:], s.MetricID) - bin.PutUint16(arr[5:], uint16(len(s.Measures))) - w.Write(arr) - for _, measure := range s.Measures { - bin.WriteUint32(w, measure.Timestamp) - bin.WriteFloat64(w, measure.Value) - } -} +// func (s AppendedMeasures) Pack(w io.Writer) { +// arr := []byte{ +// CodeAppendedMeasures, +// 0, 0, 0, 0, // metricID +// 0, 0, // qty +// } +// bin.PutUint32(arr[1:], s.MetricID) +// bin.PutUint16(arr[5:], uint16(len(s.Measures))) +// w.Write(arr) +// for _, measure := range s.Measures { +// bin.WriteUint32(w, measure.Timestamp) +// bin.WriteFloat64(w, measure.Value) +// } +// } -func (s *AppendedMeasures) Parse(src io.Reader) (err error) { - s.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 - } - s.Measures = append(s.Measures, measure) - } - return nil -} +// func (s *AppendedMeasures) Parse(src io.Reader) (err error) { +// s.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 +// } +// s.Measures = append(s.Measures, measure) +// } +// return nil +// } -type AppendedPagesReq struct { - Legs []atree.PathLeg - MetricID uint32 - Timestamp uint32 // last measure - Value float64 // last measure - //RootPageNo uint32 - LastPageNo uint32 - Pages []atree.NotLinkedDataPage - TimestampsChunks [][]byte - TimestampsSize uint16 - ValuesChunks [][]byte - ValuesSize uint16 -} +// type AppendedPagesReq struct { +// Legs []atree.PathLeg +// MetricID uint32 +// Timestamp uint32 // last measure +// Value float64 // last measure +// //RootPageNo uint32 +// LastPageNo uint32 +// Pages []atree.NotLinkedDataPage +// TimestampsChunks [][]byte +// TimestampsSize uint16 +// ValuesChunks [][]byte +// ValuesSize uint16 +// } /* 4b metricID @@ -184,96 +182,96 @@ Nb - timestamps payload 2b - values size (not filled page) Nb - values payload */ -type AppendedPages struct { - MetricID uint32 - Timestamp uint32 // last measure - Value float64 // last measure - NewRootPageNo uint32 - LastPageNo uint32 - Pages []atree.PageToWrite - TimestampsChunks [][]byte - TimestampsSize int - ValuesChunks [][]byte - ValuesSize int -} +// type AppendedPages struct { +// MetricID uint32 +// Timestamp uint32 // last measure +// Value float64 // last measure +// NewRootPageNo uint32 +// LastPageNo uint32 +// Pages []atree.PageToWrite +// TimestampsChunks [][]byte +// TimestampsSize int +// ValuesChunks [][]byte +// ValuesSize int +// } -func (s *AppendedPages) Pack(w io.Writer) { - w.Write([]byte{ - CodeAppendedPages, - }) - bin.WriteUint32(w, s.MetricID) - bin.WriteUint32(w, s.Timestamp) - bin.WriteFloat64(w, s.Value) - bin.WriteUint32(w, s.LastPageNo) - bin.WriteUint32(w, s.NewRootPageNo) - // - bin.WriteVarSize(w, len(s.Pages)) - for _, p := range s.Pages { - bin.WriteUint32(w, p.PageNo) - if p.IsReused { - w.Write([]byte{1}) - } else { - w.Write([]byte{0}) - } - w.Write(p.Data) - } - // timestamps - writeChunkedPayloadTo(s.TimestampsChunks, int(s.TimestampsSize), w) - // values - writeChunkedPayloadTo(s.ValuesChunks, int(s.ValuesSize), w) -} +// func (s *AppendedPages) Pack(w io.Writer) { +// w.Write([]byte{ +// CodeAppendedPages, +// }) +// bin.WriteUint32(w, s.MetricID) +// bin.WriteUint32(w, s.Timestamp) +// bin.WriteFloat64(w, s.Value) +// bin.WriteUint32(w, s.LastPageNo) +// bin.WriteUint32(w, s.NewRootPageNo) +// // +// bin.WriteVarSize(w, len(s.Pages)) +// for _, p := range s.Pages { +// bin.WriteUint32(w, p.PageNo) +// if p.IsReused { +// w.Write([]byte{1}) +// } else { +// w.Write([]byte{0}) +// } +// w.Write(p.Data) +// } +// // timestamps +// writeChunkedPayloadTo(s.TimestampsChunks, int(s.TimestampsSize), w) +// // values +// writeChunkedPayloadTo(s.ValuesChunks, int(s.ValuesSize), w) +// } -func (s *AppendedPages) Parse(src io.Reader) (err error) { - s.MetricID, err = bin.ReadUint32(src) - if err != nil { - return - } - s.Timestamp, err = bin.ReadUint32(src) - if err != nil { - return - } - s.Value, err = bin.ReadFloat64(src) - if err != nil { - return - } - s.LastPageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - s.NewRootPageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - pagesQty, err := bin.ReadVarSize(src) - if err != nil { - return - } - for range pagesQty { - var page atree.PageToWrite - page.PageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - page.IsReused, err = bin.ReadBool(src) - if err != nil { - return - } - page.Data, err = bin.ReadN(src, atree.PageSize) - if err != nil { - return - } - s.Pages = append(s.Pages, page) - } - s.TimestampsChunks, s.TimestampsSize, err = readChunkedPayload(src) - if err != nil { - return - } - s.ValuesChunks, s.ValuesSize, err = readChunkedPayload(src) - if err != nil { - return - } - return nil -} +// func (s *AppendedPages) Parse(src io.Reader) (err error) { +// s.MetricID, err = bin.ReadUint32(src) +// if err != nil { +// return +// } +// s.Timestamp, err = bin.ReadUint32(src) +// if err != nil { +// return +// } +// s.Value, err = bin.ReadFloat64(src) +// if err != nil { +// return +// } +// s.LastPageNo, err = bin.ReadUint32(src) +// if err != nil { +// return +// } +// s.NewRootPageNo, err = bin.ReadUint32(src) +// if err != nil { +// return +// } +// pagesQty, err := bin.ReadVarSize(src) +// if err != nil { +// return +// } +// for range pagesQty { +// var page atree.PageToWrite +// page.PageNo, err = bin.ReadUint32(src) +// if err != nil { +// return +// } +// page.IsReused, err = bin.ReadBool(src) +// if err != nil { +// return +// } +// page.Data, err = bin.ReadN(src, atree.PageSize) +// if err != nil { +// return +// } +// s.Pages = append(s.Pages, page) +// } +// s.TimestampsChunks, s.TimestampsSize, err = readChunkedPayload(src) +// if err != nil { +// return +// } +// s.ValuesChunks, s.ValuesSize, err = readChunkedPayload(src) +// if err != nil { +// return +// } +// return nil +// } func writeChunkedPayloadTo(chunks [][]byte, size int, w io.Writer) { bin.WriteVarSize(w, size) diff --git a/txlog/writer.go b/txlog/writer.go index 2f30d3d..732311a 100644 --- a/txlog/writer.go +++ b/txlog/writer.go @@ -5,12 +5,13 @@ import ( "errors" "fmt" "hash/crc32" + "log" "os" "path/filepath" "sync" bin "gordenko.dev/dima/bin/little" - octopus "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" ) @@ -38,7 +39,7 @@ const ( type FreeList interface { // використовується в allocPage - ReservePage() uint32 + GetPageNumber() (uint32, error) } func JoinChangesFileName(dir string, logNumber int) string { @@ -59,20 +60,24 @@ type Writer struct { atree *atree.Atree logNumber int dir string + w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск + wal *os.File dataFile *os.File - file *os.File + indexFile *os.File + dataPage []byte + indexPage []byte input []any //buf *bytes.Buffer - pagesToWrite []atree.PageToWrite - workerReqs []any - waitCh chan struct{} + //pagesToWrite []atree.PageToWrite + workerReqs []any + //waitCh chan struct{} appendToWorkerQueue func(any) - lsn uint32 - written int64 - isExited bool - exitCh chan struct{} - waitGroup *sync.WaitGroup - signalCh chan struct{} + //lsn uint32 + written int64 + isExited bool + exitCh chan struct{} + waitGroup *sync.WaitGroup + signalCh chan struct{} } type WriterOptions struct { @@ -114,13 +119,12 @@ func NewWriter(opt WriterOptions) (*Writer, error) { exitCh: opt.ExitCh, waitGroup: opt.WaitGroup, signalCh: make(chan struct{}, 1), - waitCh: make(chan struct{}), } var err error if opt.LogNumber > 0 { - s.file, err = os.OpenFile( + s.wal, err = os.OpenFile( JoinChangesFileName(opt.Dir, s.logNumber), os.O_APPEND|os.O_WRONLY, filePerm, @@ -130,7 +134,7 @@ func NewWriter(opt WriterOptions) (*Writer, error) { } } else { s.logNumber = 1 - s.file, err = os.OpenFile( + s.wal, err = os.OpenFile( JoinChangesFileName(opt.Dir, s.logNumber), os.O_CREATE|os.O_WRONLY, filePerm, @@ -147,7 +151,7 @@ func (s *Writer) Run() { select { case <-s.signalCh: if err := s.packAndWrite(); err != nil { - octopus.Abort(octopus.FailedWriteToTxLog, err) + qb.Abort(qb.FailedWriteToTxLog, err) } case <-s.exitCh: @@ -157,48 +161,67 @@ func (s *Writer) Run() { } } -func (s *Writer) packAndWrite() error { +func (s *Writer) packAndWrite() (err error) { s.mutex.Lock() input := s.input - waitCh := s.waitCh s.input = nil - s.waitCh = make(chan struct{}) s.mutex.Unlock() w := bytes.NewBuffer(nil) + // 1. Пакую всі дані в WAL буфер запису for _, untyped := range input { switch x := untyped.(type) { - case AppendedPagesReq: - s.packWriteAppended(w, x) case AddedMetric: - case DeletedMetric: - if len(x.FreePageNumbers) > 0 { - //s.freeList.AddPages(x.FreePageNumbers) - } - case AppendedMeasure: + x.Pack(w) + // case DeletedMetric: + // if len(x.FreePageNumbers) > 0 { + // s.freeList.AddPageNumbers(x.FreePageNumbers) + // } case AppendedMeasures: - case DeletedMeasures: + // + + s.packMeasuresIntoWALBuffer(x) + //case DeletedMeasures: //case DeletedMeasuresSince: } } + // 2. Додаю розмір пакету і чексуму + + // 3. Пишу на диск WAL (append) + + // 4. Пишу в atree сторінки + for range 1 { + err = s.writePagesToAtree(nil, nil) + if err != nil { + return + } + } + + // 5. Пишу зміни в free list + + // 6. відправляю input - worker-у + // flush - close(waitCh) - return s.flush(w) + err = s.flush(w) + if err != nil { + return err + } + + // FIX send to worker workerReqs + return nil } -func (s *Writer) Append(req any) chan struct{} { +func (s *Writer) Append(req any) { s.mutex.Lock() s.input = append(s.input, req) s.workerReqs = append(s.workerReqs, req) - waitCh := s.waitCh s.mutex.Unlock() select { case s.signalCh <- struct{}{}: default: } - return waitCh } // func (s *Writer) reset() { @@ -217,9 +240,8 @@ func (s *Writer) Append(req any) chan struct{} { func (s *Writer) flush(w *bytes.Buffer) error { s.mutex.Lock() - pagesToWrite := s.pagesToWrite + //pagesToWrite := s.pagesToWrite workerReqs := s.workerReqs - waitCh := s.waitCh isExited := s.isExited var exitWaitGroup *sync.WaitGroup @@ -228,8 +250,8 @@ func (s *Writer) flush(w *bytes.Buffer) error { } if w.Len() > packetPrefixSize { - s.lsn++ - lsn := s.lsn + //s.lsn++ + //lsn := s.lsn packet := make([]byte, w.Len()) copy(packet, w.Bytes()) //s.reset() fix @@ -238,26 +260,25 @@ func (s *Writer) flush(w *bytes.Buffer) error { s.mutex.Unlock() bin.PutUint32(packet[lengthIdx:], uint32(len(packet)-packetPrefixSize)) - bin.PutUint32(packet[lsnIdx:], lsn) + //bin.PutUint32(packet[lsnIdx:], lsn) bin.PutUint32(packet[checksumIdx:], crc32.ChecksumIEEE(packet[8:])) - n, err := s.file.Write(packet) + n, err := s.wal.Write(packet) if err != nil { return fmt.Errorf("TxLog write: %s", err) } if n != len(packet) { return fmt.Errorf("TxLog written %d != packet size %d", n, len(packet)) } - if err := s.file.Sync(); err != nil { + if err := s.wal.Sync(); err != nil { return fmt.Errorf("TxLog sync: %s", err) } - err = s.writePagesToAtree(pagesToWrite) - if err != nil { - return fmt.Errorf("TxLog writePagesToAtree: %s", err) - } + // err = s.writePagesToAtree(pagesToWrite) + // if err != nil { + // return fmt.Errorf("TxLog writePagesToAtree: %s", err) + // } } else { - s.waitCh = make(chan struct{}) s.mutex.Unlock() } @@ -272,12 +293,12 @@ func (s *Writer) flush(w *bytes.Buffer) error { } if forceSnapshot { - if err := s.file.Close(); err != nil { + if err := s.wal.Close(); err != nil { return fmt.Errorf("close changes file: %s", err) } s.logNumber++ var err error - s.file, err = os.OpenFile( + s.wal, err = os.OpenFile( JoinChangesFileName(s.dir, s.logNumber), os.O_CREATE|os.O_WRONLY, filePerm, @@ -293,41 +314,146 @@ func (s *Writer) flush(w *bytes.Buffer) error { Records: workerReqs, ForceSnapshot: forceSnapshot, LogNumber: s.logNumber, - WaitCh: waitCh, ExitWaitGroup: exitWaitGroup, }) return nil } -func (s *Writer) writePagesToAtree(pages []atree.PageToWrite) error { - for _, p := range pages { - if len(p.Data) != atree.PageSize { - return fmt.Errorf("wrong page %d size: %d", - p.PageNo, len(p.Data)) - } - - off := (p.PageNo - 1) * atree.PageSize - n, err := s.dataFile.WriteAt(p.Data, int64(off)) - if err != nil { - return err - } - if n != len(p.Data) { - return fmt.Errorf("write %d instead of %d", n, len(p.Data)) - } - } - return nil -} - func (s *Writer) exit() { s.mutex.Lock() s.isExited = true s.mutex.Unlock() if err := s.packAndWrite(); err != nil { - octopus.Abort(octopus.FailedWriteToTxLog, err) + qb.Abort(qb.FailedWriteToTxLog, err) } } +// writePagesToAtree - записує сторінки в .data та .index файли +func (s *Writer) writePagesToAtree(levels []*atree.IndexLevel, dataPages []*atree.DataPage) (err error) { + for _, p := range dataPages { + p.Checksum = atree.ChunksToDataPage(s.dataPage, atree.ChunksToDataPageReq{ + PrevPageNo: p.PrevPageNo, + Timestamps: p.Timestamps, + TimestampsSize: p.TimestampsSize, + Values: p.Values, + ValuesSize: p.ValuesSize, + }) + off := (p.PageNo - 1) * atree.DataPageSize + n, err := s.dataFile.WriteAt(s.dataPage, int64(off)) + if err != nil { + return err + } + if n != atree.DataPageSize { + return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize) + } + } + for levelIdx, level := range levels { + for _, p := range level.Filled { + // if len(p.Data) != atree.PageSize { + // return fmt.Errorf("wrong page %d size: %d", + // p.PageNo, len(p.Data)) + // } + p.Checksum = atree.DataToIndexPage(s.indexPage, p.Data, levelIdx == 0) + off := (p.PageNo - 1) * atree.IndexPageSize + n, err := s.indexFile.WriteAt(s.indexPage, int64(off)) + if err != nil { + return err + } + if n != atree.IndexPageSize { + return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize) + } + } + } + return nil +} + +type AppendedMeasures struct { + MetricID uint32 + TimestampsOffset int // (заповнені одразу) + ValuesOffset int // (заповнені одразу) + Timestamps [][]byte + TimestampsSize int + Values [][]byte + ValuesSize int + IndexLevels []*atree.IndexLevel + DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page +} + +func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) { + // fix write code + bin.WriteUint32(s.w, req.MetricID) + bin.WriteVarSize(s.w, len(req.DataPages)) + for idx, p := range req.DataPages { + bin.WriteUint32(s.w, p.PageNo) + var reused byte + if p.IsReused { + reused = 1 + } + s.w.WriteByte(reused) + bin.WriteUint32(s.w, p.PrevPageNo) + bin.WriteUint32(s.w, p.Checksum) + var ( + timestampsOffset int + valuesOffset int + ) + if idx == 0 { + timestampsOffset = req.TimestampsOffset + valuesOffset = req.ValuesOffset + } + copyPayloadFromChunks(s.w, p.Timestamps, p.TimestampsSize, timestampsOffset) + copyPayloadFromChunks(s.w, p.Values, p.ValuesSize, valuesOffset) + } + bin.WriteUint16(s.w, uint16(req.TimestampsOffset)) + bin.WriteUint16(s.w, uint16(req.ValuesOffset)) + var ( + timestampsOffset int + valuesOffset int + ) + if len(req.DataPages) == 0 { + timestampsOffset = req.TimestampsOffset + valuesOffset = req.ValuesOffset + } + copyPayloadFromChunks(s.w, req.Timestamps, req.TimestampsSize, timestampsOffset) + copyPayloadFromChunks(s.w, req.Values, req.ValuesSize, valuesOffset) + if len(req.DataPages) == 0 { + return + } + // fix - levels reduced + bin.WriteVarSize(s.w, len(req.IndexLevels)) + for _, level := range req.IndexLevels { + bin.WriteVarSize(s.w, len(level.Filled)) + for idx, p := range level.Filled { + bin.WriteUint32(s.w, p.PageNo) + var reused byte + if p.IsReused { + reused = 1 + } + s.w.WriteByte(reused) + bin.WriteUint32(s.w, p.Checksum) + records := p.Data + if idx == 0 { + records = p.Data[level.Offset:] + } + if (len(records) % recordSize) != 0 { + log.Fatalf("wrong records length: %d", len(records)) + } + bin.WriteVarSize(s.w, len(records)/recordSize) + s.w.Write(records) + } + records := level.Data + if len(level.Filled) == 0 { + records = level.Data[level.Offset:] + } + if (len(records) % recordSize) != 0 { + log.Fatalf("wrong unfilled records length: %d", len(records)) + } + bin.WriteVarSize(s.w, len(records)/recordSize) + s.w.Write(records) + } + return +} + // API // FIX - add @@ -351,28 +477,28 @@ func (s *Writer) sendSignal() { // Якщо (Pages) > 0 - незаповнену сторінку кодую в стисненому вигляді. // Якщо (Pages) = 0 - кодую просто пари timestamp / value -func (s *Writer) packWriteAppended(w *bytes.Buffer, req AppendedPagesReq) { - // завантажений path. Отже просто додаємо в дерево data сторінки, створюємо нові індексні - // без звернення до диску. - report := s.atree.AppendDataPages(atree.AppendDataPagesReq{ - LastPageNo: req.LastPageNo, - Legs: req.Legs, - DataPages: req.Pages, - }) +// func (s *Writer) packWriteAppended(w *bytes.Buffer, req AppendedPagesReq) { +// // завантажений path. Отже просто додаємо в дерево data сторінки, створюємо нові індексні +// // без звернення до диску. +// report := s.atree.AppendDataPages(atree.AppendDataPagesReq{ +// LastPageNo: req.LastPageNo, +// Legs: req.Legs, +// DataPages: req.Pages, +// }) - m := AppendedPages{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - NewRootPageNo: report.NewRootPageNo, - LastPageNo: report.LastPageNo, - Pages: report.Pages, - TimestampsChunks: req.TimestampsChunks, - TimestampsSize: int(req.TimestampsSize), - ValuesChunks: req.ValuesChunks, - ValuesSize: int(req.ValuesSize), - } - m.Pack(w) -} +// m := AppendedPages{ +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// NewRootPageNo: report.NewRootPageNo, +// LastPageNo: report.LastPageNo, +// Pages: report.Pages, +// TimestampsChunks: req.TimestampsChunks, +// TimestampsSize: int(req.TimestampsSize), +// ValuesChunks: req.ValuesChunks, +// ValuesSize: int(req.ValuesSize), +// } +// m.Pack(w) +// } // helpers diff --git a/txlog/x.go b/txlog/x.go deleted file mode 100644 index 66d9426..0000000 --- a/txlog/x.go +++ /dev/null @@ -1,233 +0,0 @@ -package txlog - -import ( - "bytes" - "fmt" - "log" - "os" - - bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb/atree" -) - -const recordSize = 8 - -type LogWriter struct { - w *bytes.Buffer - wal *os.File - dataFile *os.File - indexFile *os.File - dataPage []byte - indexPage []byte -} - -func (s *LogWriter) writePagesToAtree(levels []*atree.IndexLevel, dataPages []*atree.DataPage) (err error) { - for _, p := range dataPages { - atree.ChunksToDataPage(s.dataPage, atree.ChunksToDataPageReq{ - PrevPageNo: p.PrevPageNo, - Timestamps: p.Timestamps, - TimestampsSize: p.TimestampsSize, - Values: p.Values, - ValuesSize: p.ValuesSize, - }) - off := (p.PageNo - 1) * atree.DataPageSize - n, err := s.dataFile.WriteAt(s.dataPage, int64(off)) - if err != nil { - return err - } - if n != atree.DataPageSize { - return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize) - } - } - for levelIdx, level := range levels { - for _, p := range level.Filled { - // if len(p.Data) != atree.PageSize { - // return fmt.Errorf("wrong page %d size: %d", - // p.PageNo, len(p.Data)) - // } - atree.DataToIndexPage(s.indexPage, p.Data, levelIdx == 0) - off := (p.PageNo - 1) * atree.IndexPageSize - n, err := s.indexFile.WriteAt(s.indexPage, int64(off)) - if err != nil { - return err - } - if n != atree.IndexPageSize { - return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize) - } - } - } - return nil -} - -type ToWrite struct { - MetricID uint32 - TimestampsOffset int // (заповнені одразу) - ValuesOffset int // (заповнені одразу) - Timestamps [][]byte - TimestampsSize int - Values [][]byte - ValuesSize int - IndexLevels []*atree.IndexLevel - DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page -} - -type DataPage struct { - PageNo uint32 - Reused bool - PrevPageNo uint32 - Checksum uint32 - Timestamps []byte - Values []byte -} - -type IndexPage struct { - PageNo uint32 - Reused bool - Checksum uint32 - Records []byte -} - -type IndexLevel struct { - Pages []IndexPage - // offset потрібен тому що дані не додаються в кінець, а перезаписують кілька - // останніх байтів unfilled даних. Потрібно для коректного recovery. - // Якщо є Pages, то застосовуються для 1-ї Index сторінки, інакше - до records (unfilled) - RecordsOffset int - Records []byte // unfilled -} - -type TxAppendedMeasures struct { - MetricID uint32 - DataPages []DataPage - // offsets потрібні тому що дані не додаються в кінець, а перезаписують кілька - // останніх байтів unfilled даних. Потрібно для коректного recovery. - // Якщо є DataPages, то застосовуються для 1-ї Data сторінки, інакше - до timestamps і values (unfilled) - TimestampsOffset int - ValuesOffset int - Timestamps []byte // unfilled - Values []byte // unfilled - IndexLevels []DataPage -} - -/* -Format appended measures: -1b - tx type -4b - metricID -Nb - qty of filled data pages (varsize) -[ - 4b - pageNo - 1b - reused - 4b - prevPageNo - 4b - page crc32 - Nb - timestamps size (varsize) - Nb - timestamps payload - Nb - values size (varsize) - Nb - values payload -] -2b - timestamps offset (for the 1st data page only) -2b - values offset (for the 1st data page only) -Nb - timestamps size unfilled (varsize) -Nb - timestamps payload unfilled -Nb - values size unfilled (varsize) -Nb - values payload unfilled -Nb - qty of index levels (varsize) -[ - Nb - qty of level filled pages (varsize) - [ - 4b - pageNo - 1b - reused - 4b - page crc32 - Nb - records qty (varsize) - Nb - records payload - ] - Nb - records offset (for the 1st index page only) - Nb - records qty unfilled (varsize) - Nb - records payload unfilled -] -*/ - -func (s *LogWriter) writeToWAL(req ToWrite) (err error) { - // fix write code - bin.WriteUint32(s.w, req.MetricID) - bin.WriteVarSize(s.w, len(req.DataPages)) - for idx, p := range req.DataPages { - bin.WriteUint32(s.w, p.PageNo) - var reused byte - if p.IsReused { - reused = 1 - } - s.w.WriteByte(reused) - bin.WriteUint32(s.w, p.PrevPageNo) - //bin.WriteUint32(s.w, p.CRC32) fix - var ( - timestampsOffset int - valuesOffset int - ) - if idx == 0 { - timestampsOffset = req.TimestampsOffset - valuesOffset = req.ValuesOffset - } - copyPayloadFromChunks(s.w, p.Timestamps, p.TimestampsSize, timestampsOffset) - copyPayloadFromChunks(s.w, p.Values, p.ValuesSize, valuesOffset) - } - bin.WriteUint16(s.w, uint16(req.TimestampsOffset)) - bin.WriteUint16(s.w, uint16(req.ValuesOffset)) - // fix offset - copyPayloadFromChunks(s.w, req.Timestamps, req.TimestampsSize, req.TimestampsOffset) - copyPayloadFromChunks(s.w, req.Values, req.ValuesSize, req.ValuesOffset) - if len(req.DataPages) == 0 { - return - } - bin.WriteVarSize(s.w, len(req.IndexLevels)) - for _, level := range req.IndexLevels { - bin.WriteVarSize(s.w, len(level.Filled)) - for idx, p := range level.Filled { - bin.WriteUint32(s.w, p.PageNo) - var reused byte - if p.IsReused { - reused = 1 - } - s.w.WriteByte(reused) - //bin.WriteUint32(s.w, p.CRC32) fix - records := p.Data - if idx == 0 { - records = p.Data[level.Offset:] - } - if (len(records) % recordSize) != 0 { - log.Fatalf("wrong records length: %d", len(records)) - } - bin.WriteVarSize(s.w, len(records)/recordSize) - s.w.Write(records) - } - if len(level.Filled) > 0 { - s.w.Write(level.Data) - } else { - s.w.Write(level.Data[level.Offset:]) - } - } - return -} - -const chunkSize = 24 - -func copyPayloadFromChunks(w *bytes.Buffer, chunks [][]byte, size int, offset int) { - bin.WriteVarSize(w, size) - var ( - chunkIdx = offset / chunkSize - byteIdx = offset % chunkSize - ) - for _, chunk := range chunks[chunkIdx:] { - available := len(chunk) - byteIdx - if available <= size { - w.Write(chunk[byteIdx:]) - size -= available - if size == 0 { - return - } - byteIdx = 0 - } else { - w.Write(chunk[byteIdx : byteIdx+size]) - return - } - } -} diff --git a/atree/io_test.go b/util/io_test.go similarity index 95% rename from atree/io_test.go rename to util/io_test.go index 7ead11f..a58c073 100644 --- a/atree/io_test.go +++ b/util/io_test.go @@ -1,4 +1,4 @@ -package atree +package util import ( "hash/crc32" diff --git a/util/until.go b/util/until.go new file mode 100644 index 0000000..1ca881d --- /dev/null +++ b/util/until.go @@ -0,0 +1,17 @@ +package util + +import ( + "hash/crc32" +) + +var ( + castagnoliTable = crc32.MakeTable(crc32.Castagnoli) +) + +func CalcChecksum(page []byte) uint32 { + return crc32.Checksum(page, castagnoliTable) +} + +func IsValidChecksum(page []byte, checksum uint32) bool { + return crc32.Checksum(page, castagnoliTable) == checksum +}