diff --git a/chunkenc/chunkenc_test.go b/chunkenc/chunkenc_test.go index 7d56f2f..7e62abf 100644 --- a/chunkenc/chunkenc_test.go +++ b/chunkenc/chunkenc_test.go @@ -149,9 +149,9 @@ func TestCumdeltaBound(t *testing.T) { fmt.Printf("%d\n", buf.Chunks()[0]) ////////////////////////// - bound := c.GetState() + c.Lock() - fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) + //fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) c.Append(248) c.Append(305) @@ -167,8 +167,9 @@ func TestCumdeltaBound(t *testing.T) { fmt.Println(value, done) } - boundDecompressor := NewReverseCumulativeDeltaDecompressor(buf, c.Size(), fracDigits) - boundDecompressor.RestoreFromBound(bound) + //boundDecompressor := NewReverseCumulativeDeltaDecompressor(buf, c.Size(), fracDigits) + //boundDecompressor.RestoreFromBound(bound) + boundDecompressor := c.CreateDecompressor(fracDigits) fmt.Println("from bound:") for range 8 { diff --git a/database/api.go b/database/api.go index 2792f88..ca0991e 100644 --- a/database/api.go +++ b/database/api.go @@ -94,13 +94,13 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } reply(conn, s.DeleteMetric(req)) - case proto.TypeAppendMeasure: - req, err := proto.ReadAppendMeasureReq(r) - if err != nil { - return fmt.Errorf("proto.ReadAppendMeasureReq: %s", err) - } - //fmt.Println("append measure", req.MetricID, conn.RemoteAddr().String()) - reply(conn, s.AppendMeasure(req)) + // case proto.TypeAppendMeasure: + // req, err := proto.ReadAppendMeasureReq(r) + // if err != nil { + // return fmt.Errorf("proto.ReadAppendMeasureReq: %s", err) + // } + // //fmt.Println("append measure", req.MetricID, conn.RemoteAddr().String()) + // reply(conn, s.AppendMeasure(req)) case proto.TypeAppendMeasures: req, err := proto.ReadAppendMeasuresReq(r) @@ -323,120 +323,120 @@ type FilledPage struct { ValuesSize uint16 } -type tryAppendMeasureResult struct { - MetricID uint32 - Timestamp uint32 - Value float64 - FilledPage *FilledPage - ResultCode byte -} +// type tryAppendMeasureResult struct { +// MetricID uint32 +// Timestamp uint32 +// Value float64 +// FilledPage *FilledPage +// ResultCode byte +// } -func (s *Database) AppendMeasure(req proto.AppendMeasureReq) uint16 { - resultCh := make(chan tryAppendMeasureResult, 1) +// func (s *Database) AppendMeasure(req proto.AppendMeasureReq) uint16 { +// resultCh := make(chan tryAppendMeasureResult, 1) - s.appendJobToWorkerQueue(tryAppendMeasureReq{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - ResultCh: resultCh, - }) +// s.appendJobToWorkerQueue(tryAppendMeasureReq{ +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// ResultCh: resultCh, +// }) - result := <-resultCh +// result := <-resultCh - switch result.ResultCode { - case CanAppend: - waitCh := s.txlog.Append(txlog.AppendedMeasure{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - }) - <-waitCh +// switch result.ResultCode { +// case CanAppend: +// waitCh := s.txlog.Append(txlog.AppendedMeasure{ +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// }) +// <-waitCh - case NewPage: - filled := result.FilledPage +// case NewPage: +// filled := result.FilledPage - var path atree.PathToDataPage - if filled.RootPageNo > 0 { - path, err := s.atree.FindPathToLastPage(filled.RootPageNo) - if err != nil { - // FIX - diploma.Abort(diploma.WriteToAtreeFailed, err) - } - // for _, leg := range path.Legs { - // pagesToRelease = append(pagesToRelease, leg.PageNo) - // } +// var path atree.PathToDataPage +// if filled.RootPageNo > 0 { +// path, err := s.atree.FindPathToLastPage(filled.RootPageNo) +// if err != nil { +// // FIX +// diploma.Abort(diploma.WriteToAtreeFailed, err) +// } +// // for _, leg := range path.Legs { +// // pagesToRelease = append(pagesToRelease, leg.PageNo) +// // } - if path.LastPageNo != filled.PrevPageNo { - diploma.Abort( - diploma.WrongPrevPageNo, - fmt.Errorf("bug: last pageNo %d in tree != prev pageNo %d in _metric", - path.LastPageNo, filled.PrevPageNo), - ) - } - } +// if path.LastPageNo != filled.PrevPageNo { +// diploma.Abort( +// diploma.WrongPrevPageNo, +// fmt.Errorf("bug: last pageNo %d in tree != prev pageNo %d in _metric", +// path.LastPageNo, filled.PrevPageNo), +// ) +// } +// } - // report, err := s.atree.AppendDataPages(atree.AppendDataPageReq{ - // MetricID: req.MetricID, - // Timestamp: req.Timestamp, - // Value: req.Value, - // Since: filled.Since, - // RootPageNo: filled.RootPageNo, - // PrevPageNo: filled.PrevPageNo, - // TimestampsChunks: filled.TimestampsChunks, - // TimestampsSize: filled.TimestampsSize, - // ValuesChunks: filled.ValuesChunks, - // ValuesSize: filled.ValuesSize, - // }) - // if err != nil { - // diploma.Abort(diploma.WriteToAtreeFailed, err) - // } - // _ = report +// // report, err := s.atree.AppendDataPages(atree.AppendDataPageReq{ +// // MetricID: req.MetricID, +// // Timestamp: req.Timestamp, +// // Value: req.Value, +// // Since: filled.Since, +// // RootPageNo: filled.RootPageNo, +// // PrevPageNo: filled.PrevPageNo, +// // TimestampsChunks: filled.TimestampsChunks, +// // TimestampsSize: filled.TimestampsSize, +// // ValuesChunks: filled.ValuesChunks, +// // ValuesSize: filled.ValuesSize, +// // }) +// // if err != nil { +// // diploma.Abort(diploma.WriteToAtreeFailed, err) +// // } +// // _ = report - //waitCh := s.txlog.WriteAppended(path, report) +// //waitCh := s.txlog.WriteAppended(path, report) - waitCh := s.txlog.Append(txlog.AppendedPagesReq{ - Legs: path.Legs, - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - LastPageNo: filled.PrevPageNo, - //Pages: , - TimestampsChunks: filled.TimestampsChunks, - TimestampsSize: filled.TimestampsSize, - ValuesChunks: filled.ValuesChunks, - ValuesSize: filled.ValuesSize, - }, - ) - <-waitCh +// waitCh := s.txlog.Append(txlog.AppendedPagesReq{ +// Legs: path.Legs, +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// LastPageNo: filled.PrevPageNo, +// //Pages: , +// TimestampsChunks: filled.TimestampsChunks, +// TimestampsSize: filled.TimestampsSize, +// ValuesChunks: filled.ValuesChunks, +// ValuesSize: filled.ValuesSize, +// }, +// ) +// <-waitCh - case NoMetric: - return proto.ErrNoMetric +// case NoMetric: +// return proto.ErrNoMetric - case ExpiredMeasure: - return proto.ErrExpiredMeasure +// case ExpiredMeasure: +// return proto.ErrExpiredMeasure - case NonMonotonicValue: - return proto.ErrNonMonotonicValue +// case NonMonotonicValue: +// return proto.ErrNonMonotonicValue - default: - diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) - } - return 0 -} +// default: +// diploma.Abort(diploma.WrongResultCodeBug, ErrWrongResultCodeBug) +// } +// return 0 +// } type tryAppendMeasuresResult struct { - ResultCode byte - MetricType diploma.MetricType - FracDigits byte - Since uint32 - Until uint32 - UntilValue float64 - RootPageNo uint32 - PrevPageNo uint32 - TimestampsBuf *conbuf.ContinuousBuffer - ValuesBuf *conbuf.ContinuousBuffer - Timestamps diploma.TimestampCompressor - Values diploma.ValueCompressor + ResultCode byte + // MetricType diploma.MetricType + // FracDigits byte + // Since uint32 + // Until uint32 + // UntilValue float64 + // RootPageNo uint32 + // PrevPageNo uint32 + // TimestampsBuf *conbuf.ContinuousBuffer + // ValuesBuf *conbuf.ContinuousBuffer + // Timestamps diploma.TimestampCompressor + // Values diploma.ValueCompressor } func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { @@ -444,6 +444,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { s.appendJobToWorkerQueue(tryAppendMeasuresReq{ MetricID: req.MetricID, + Measures: req.Measures, ResultCh: resultCh, }) diff --git a/database/database.go b/database/database.go index 7d12bcf..4c49794 100644 --- a/database/database.go +++ b/database/database.go @@ -33,22 +33,22 @@ type metricLockEntry struct { } type Database struct { - mutex sync.Mutex - workerSignalCh chan struct{} - workerQueue []any - rLocksToRelease []uint32 - metrics map[uint32]*_metric - metricLockEntries map[uint32]*metricLockEntry - freeList *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 + mutex sync.Mutex + workerSignalCh chan struct{} + workerQueue []any + rLocksToRelease []uint32 + metrics map[uint32]*_metric + //metricLockEntries map[uint32]*metricLockEntry + freeList *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 } type Options struct { @@ -85,17 +85,17 @@ func New(opt Options) (_ *Database, err error) { } 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(), - tcpPort: opt.TCPPort, - logfile: opt.Logfile, - logger: log.New(opt.Logfile, "", log.LstdFlags), - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, + workerSignalCh: make(chan struct{}, 1), + dir: opt.Dir, + databaseName: opt.DatabaseName, + metrics: make(map[uint32]*_metric), + //metricLockEntries: make(map[uint32]*metricLockEntry), + freeList: freelist.New(), + tcpPort: opt.TCPPort, + logfile: opt.Logfile, + logger: log.New(opt.Logfile, "", log.LstdFlags), + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, } return s, nil } @@ -369,7 +369,7 @@ func (s *Database) replayChangesRecord(untyped any) error { FracDigits: byte(rec.FracDigits), TimestampsBuf: timestampsBuf, ValuesBuf: valuesBuf, - Timestamps: chunkenc.NewReverseTimeDeltaOfDeltaCompressor(timestampsBuf, 0), + Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), Values: values, } diff --git a/database/proc.go b/database/proc.go index f6c478e..ec758ce 100644 --- a/database/proc.go +++ b/database/proc.go @@ -7,6 +7,7 @@ import ( "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/chunkenc" "gordenko.dev/dima/qb/conbuf" + "gordenko.dev/dima/qb/proto" "gordenko.dev/dima/qb/transform" "gordenko.dev/dima/qb/txlog" ) @@ -85,8 +86,8 @@ func (s *Database) DoWork() { for _, untyped := range workerQueue { switch req := untyped.(type) { - case tryAppendMeasureReq: - s.tryAppendMeasure(req) + //case tryAppendMeasureReq: + // s.tryAppendMeasure(req) case tryAppendMeasuresReq: s.tryAppendMeasures(req) @@ -155,8 +156,8 @@ func (s *Database) processMetricQueue(metricID uint32, metric *_metric, lockEntr } else { for idx, untyped := range modificationReqs { switch req := untyped.(type) { - case tryAppendMeasureReq: - s.startAppendMeasure(metric, req, nil) + //case tryAppendMeasureReq: + // s.startAppendMeasure(metric, req, nil) case tryAppendMeasuresReq: s.startAppendMeasures(metric, req, nil) @@ -195,16 +196,17 @@ func (s *Database) tryAddMetric(req tryAddMetricReq) { req.ResultCh <- MetricDuplicate return } + req.ResultCh <- Succeed // new - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - } else { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - XLock: true, - } - req.ResultCh <- Succeed - } + // lockEntry, ok := s.metricLockEntries[req.MetricID] + // if ok { + // lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) + // } else { + // s.metricLockEntries[req.MetricID] = &metricLockEntry{ + // XLock: true, + // } + // req.ResultCh <- Succeed + // } } func (s *Database) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetricReq) { @@ -220,10 +222,11 @@ func (s *Database) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMe waitQueue = append(waitQueue, req) } } - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - XLock: true, - WaitQueue: waitQueue, - } + // FIX + // s.metricLockEntries[req.MetricID] = &metricLockEntry{ + // XLock: true, + // WaitQueue: waitQueue, + // } req.ResultCh <- Succeed } @@ -253,30 +256,32 @@ type tryDeleteMetricReq struct { } func (s *Database) tryDeleteMetric(req tryDeleteMetricReq) { - metric, ok := s.metrics[req.MetricID] - if !ok { - req.ResultCh <- tryDeleteMetricResult{ - ResultCode: NoMetric, - } - return - } + // FIX + // metric, ok := s.metrics[req.MetricID] + // if !ok { + // req.ResultCh <- tryDeleteMetricResult{ + // ResultCode: NoMetric, + // } + // return + // } - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - } else { - s.startDeleteMetric(metric, req) - } + // lockEntry, ok := s.metricLockEntries[req.MetricID] + // if ok { + // lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) + // } else { + // s.startDeleteMetric(metric, req) + // } } func (s *Database) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - XLock: true, - } - req.ResultCh <- tryDeleteMetricResult{ - ResultCode: Succeed, - RootPageNo: metric.RootPageNo, - } + // FIX + // s.metricLockEntries[req.MetricID] = &metricLockEntry{ + // XLock: true, + // } + // req.ResultCh <- tryDeleteMetricResult{ + // ResultCode: Succeed, + // RootPageNo: metric.RootPageNo, + // } } type tryDeleteMeasuresReq struct { @@ -286,190 +291,192 @@ type tryDeleteMeasuresReq struct { } func (s *Database) tryDeleteMeasures(req tryDeleteMeasuresReq) { - metric, ok := s.metrics[req.MetricID] - if !ok { - req.ResultCh <- tryDeleteMeasuresResult{ - ResultCode: NoMetric, - } - return - } + // FIX + // metric, ok := s.metrics[req.MetricID] + // if !ok { + // req.ResultCh <- tryDeleteMeasuresResult{ + // ResultCode: NoMetric, + // } + // return + // } - if metric.Since == 0 || (req.Since > 0 && metric.Until < req.Since) { - req.ResultCh <- tryDeleteMeasuresResult{ - ResultCode: NoMeasuresToDelete, - } - } + // if metric.Since == 0 || (req.Since > 0 && metric.Until < req.Since) { + // req.ResultCh <- tryDeleteMeasuresResult{ + // ResultCode: NoMeasuresToDelete, + // } + // } - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - } else { - s.startDeleteMeasures(metric, req) - } + // lockEntry, ok := s.metricLockEntries[req.MetricID] + // if ok { + // lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) + // } else { + // s.startDeleteMeasures(metric, req) + // } } func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq) { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - XLock: true, - } + // FIX + // s.metricLockEntries[req.MetricID] = &metricLockEntry{ + // XLock: true, + // } - if metric.RootPageNo > 0 { - req.ResultCh <- tryDeleteMeasuresResult{ - ResultCode: DeleteFromAtreeRequired, - RootPageNo: metric.RootPageNo, - } - } else { - req.ResultCh <- tryDeleteMeasuresResult{ - ResultCode: DeleteFromAtreeNotNeeded, - } - } + // if metric.RootPageNo > 0 { + // req.ResultCh <- tryDeleteMeasuresResult{ + // ResultCode: DeleteFromAtreeRequired, + // RootPageNo: metric.RootPageNo, + // } + // } else { + // req.ResultCh <- tryDeleteMeasuresResult{ + // ResultCode: DeleteFromAtreeNotNeeded, + // } + // } } -type tryAppendMeasureReq struct { - MetricID uint32 - Timestamp uint32 - Value float64 - ResultCh chan tryAppendMeasureResult -} +// type tryAppendMeasureReq struct { +// MetricID uint32 +// Timestamp uint32 +// Value float64 +// ResultCh chan tryAppendMeasureResult +// } -func (s *Database) tryAppendMeasure(req tryAppendMeasureReq) { - metric, ok := s.metrics[req.MetricID] - if !ok { - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - ResultCode: NoMetric, - } - return - } +// func (s *Database) tryAppendMeasure(req tryAppendMeasureReq) { +// metric, ok := s.metrics[req.MetricID] +// if !ok { +// req.ResultCh <- tryAppendMeasureResult{ +// MetricID: req.MetricID, +// ResultCode: NoMetric, +// } +// return +// } - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - if lockEntry.XLock { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - return - } - } - s.startAppendMeasure(metric, req, lockEntry) -} +// lockEntry, ok := s.metricLockEntries[req.MetricID] +// if ok { +// if lockEntry.XLock { +// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) +// return +// } +// } +// s.startAppendMeasure(metric, req, lockEntry) +// } type ReadBound struct { ValuesPos int ValuesLastValue float64 } -func (s *Database) startAppendMeasure(metric *_metric, req tryAppendMeasureReq, lockEntry *metricLockEntry) { - if req.Timestamp <= metric.Until { - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - ResultCode: ExpiredMeasure, - } - return - } +// func (s *Database) startAppendMeasure(metric *_metric, req tryAppendMeasureReq, lockEntry *metricLockEntry) { +// if req.Timestamp <= metric.Until { +// req.ResultCh <- tryAppendMeasureResult{ +// MetricID: req.MetricID, +// ResultCode: ExpiredMeasure, +// } +// return +// } - if metric.MetricType == diploma.Cumulative && req.Value < metric.UntilValue { - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - ResultCode: NonMonotonicValue, - } - return - } +// if metric.MetricType == diploma.Cumulative && req.Value < metric.UntilValue { +// req.ResultCh <- tryAppendMeasureResult{ +// MetricID: req.MetricID, +// ResultCode: NonMonotonicValue, +// } +// return +// } - extraSpace := metric.Timestamps.CalcRequiredSpace(req.Timestamp) + - metric.Values.CalcRequiredSpace(req.Value) +// extraSpace := metric.Timestamps.CalcRequiredSpace(req.Timestamp) + +// metric.Values.CalcRequiredSpace(req.Value) - totalSpace := metric.Timestamps.Size() + metric.Values.Size() + extraSpace +// totalSpace := metric.Timestamps.Size() + metric.Values.Size() + extraSpace - if totalSpace <= atree.DataPagePayloadSize { - if lockEntry != nil { - lockEntry.RLocks++ - } else { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - RLocks: 1, - } - } - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - ResultCode: CanAppend, - } - } else { - if lockEntry != nil { - if lockEntry.RLocks > 0 { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - return - } - lockEntry.XLock = true - } else { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - XLock: true, - } - } +// if totalSpace <= atree.DataPagePayloadSize { +// if lockEntry != nil { +// lockEntry.RLocks++ +// } else { +// s.metricLockEntries[req.MetricID] = &metricLockEntry{ +// RLocks: 1, +// } +// } +// req.ResultCh <- tryAppendMeasureResult{ +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// ResultCode: CanAppend, +// } +// } else { +// if lockEntry != nil { +// if lockEntry.RLocks > 0 { +// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) +// return +// } +// lockEntry.XLock = true +// } else { +// s.metricLockEntries[req.MetricID] = &metricLockEntry{ +// XLock: true, +// } +// } - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - ResultCode: NewPage, - FilledPage: &FilledPage{ - Since: metric.Since, - RootPageNo: metric.RootPageNo, - PrevPageNo: metric.LastPageNo, - TimestampsChunks: metric.TimestampsBuf.Chunks(), - TimestampsSize: uint16(metric.Timestamps.Size()), - ValuesChunks: metric.ValuesBuf.Chunks(), - ValuesSize: uint16(metric.Values.Size()), - }, - } - } -} +// req.ResultCh <- tryAppendMeasureResult{ +// MetricID: req.MetricID, +// Timestamp: req.Timestamp, +// Value: req.Value, +// ResultCode: NewPage, +// FilledPage: &FilledPage{ +// Since: metric.Since, +// RootPageNo: metric.RootPageNo, +// PrevPageNo: metric.LastPageNo, +// TimestampsChunks: metric.TimestampsBuf.Chunks(), +// TimestampsSize: uint16(metric.Timestamps.Size()), +// ValuesChunks: metric.ValuesBuf.Chunks(), +// ValuesSize: uint16(metric.Values.Size()), +// }, +// } +// } +// } -func (s *Database) appendMeasure(rec txlog.AppendedMeasure) { - metric, ok := s.metrics[rec.MetricID] - if !ok { - diploma.Abort(diploma.NoMetricBug, - fmt.Errorf("appendMeasure: metric %d not found", - rec.MetricID)) - } +// func (s *Database) appendMeasure(rec txlog.AppendedMeasure) { +// metric, ok := s.metrics[rec.MetricID] +// if !ok { +// diploma.Abort(diploma.NoMetricBug, +// fmt.Errorf("appendMeasure: metric %d not found", +// rec.MetricID)) +// } - lockEntry, ok := s.metricLockEntries[rec.MetricID] - if !ok { - diploma.Abort(diploma.NoLockEntryBug, - fmt.Errorf("appendMeasure: lockEntry not found for the metric %d", - rec.MetricID)) - } +// lockEntry, ok := s.metricLockEntries[rec.MetricID] +// if !ok { +// diploma.Abort(diploma.NoLockEntryBug, +// fmt.Errorf("appendMeasure: lockEntry not found for the metric %d", +// rec.MetricID)) +// } - if lockEntry.XLock { - diploma.Abort(diploma.XLockBug, - fmt.Errorf("appendMeasure: xlock is set for the metric %d", - rec.MetricID)) - } +// if lockEntry.XLock { +// diploma.Abort(diploma.XLockBug, +// fmt.Errorf("appendMeasure: xlock is set for the metric %d", +// rec.MetricID)) +// } - if lockEntry.RLocks <= 0 { - diploma.Abort(diploma.NoRLockBug, - fmt.Errorf("appendMeasure: rlock not set for the metric %d", - rec.MetricID)) - } +// if lockEntry.RLocks <= 0 { +// diploma.Abort(diploma.NoRLockBug, +// fmt.Errorf("appendMeasure: rlock not set for the metric %d", +// rec.MetricID)) +// } - if metric.Since == 0 { - metric.Since = rec.Timestamp - metric.SinceValue = rec.Value - } +// if metric.Since == 0 { +// metric.Since = rec.Timestamp +// metric.SinceValue = rec.Value +// } - metric.Timestamps.Append(rec.Timestamp) - metric.Values.Append(rec.Value) - metric.Until = rec.Timestamp - metric.UntilValue = rec.Value +// metric.Timestamps.Append(rec.Timestamp) +// metric.Values.Append(rec.Value) +// metric.Until = rec.Timestamp +// metric.UntilValue = rec.Value - lockEntry.RLocks-- - if len(lockEntry.WaitQueue) > 0 { - s.processMetricQueue(rec.MetricID, metric, lockEntry) - } else { - if lockEntry.RLocks == 0 { - delete(s.metricLockEntries, rec.MetricID) - } - } -} +// lockEntry.RLocks-- +// if len(lockEntry.WaitQueue) > 0 { +// s.processMetricQueue(rec.MetricID, metric, lockEntry) +// } else { +// if lockEntry.RLocks == 0 { +// delete(s.metricLockEntries, rec.MetricID) +// } +// } +// } func (s *Database) appendMeasures(rec txlog.AppendedMeasures) { metric, ok := s.metrics[rec.MetricID] @@ -508,50 +515,51 @@ func (s *Database) appendMeasures(rec txlog.AppendedMeasures) { s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) } -func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWithOverflowExtended) { - rec := extended.Record - metric, ok := s.metrics[rec.MetricID] - if !ok { - diploma.Abort(diploma.NoMetricBug, - fmt.Errorf("appendMeasureAfterOverflow: metric %d not found", - rec.MetricID)) - } +// func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWithOverflowExtended) { +// rec := extended.Record +// metric, ok := s.metrics[rec.MetricID] +// if !ok { +// diploma.Abort(diploma.NoMetricBug, +// fmt.Errorf("appendMeasureAfterOverflow: metric %d not found", +// rec.MetricID)) +// } - lockEntry, ok := s.metricLockEntries[rec.MetricID] - if !ok { - diploma.Abort(diploma.NoLockEntryBug, - fmt.Errorf("appendMeasureAfterOverflow: lockEntry not found for the metric %d", - rec.MetricID)) - } +// lockEntry, ok := s.metricLockEntries[rec.MetricID] +// if !ok { +// diploma.Abort(diploma.NoLockEntryBug, +// fmt.Errorf("appendMeasureAfterOverflow: lockEntry not found for the metric %d", +// rec.MetricID)) +// } - if !lockEntry.XLock { - diploma.Abort(diploma.NoXLockBug, - fmt.Errorf("appendMeasureAfterOverflow: xlock not set for the metric %d", - rec.MetricID)) - } +// if !lockEntry.XLock { +// diploma.Abort(diploma.NoXLockBug, +// fmt.Errorf("appendMeasureAfterOverflow: xlock not set for the metric %d", +// rec.MetricID)) +// } - metric.ReinitBy(rec.Timestamp, rec.Value) - if rec.IsRootChanged { - metric.RootPageNo = rec.RootPageNo - } - metric.LastPageNo = rec.DataPageNo +// metric.ReinitBy(rec.Timestamp, rec.Value) +// if rec.IsRootChanged { +// metric.RootPageNo = rec.RootPageNo +// } +// metric.LastPageNo = rec.DataPageNo - if rec.IsDataPageReused { - s.freeList.DeleteReservedPages([]uint32{ - rec.DataPageNo, - }) - } +// if rec.IsDataPageReused { +// s.freeList.DeleteReservedPages([]uint32{ +// rec.DataPageNo, +// }) +// } - if len(rec.ReusedIndexPages) > 0 { - s.freeList.DeleteReservedPages(rec.ReusedIndexPages) - } +// if len(rec.ReusedIndexPages) > 0 { +// s.freeList.DeleteReservedPages(rec.ReusedIndexPages) +// } - lockEntry.XLock = false - s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) -} +// lockEntry.XLock = false +// s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) +// } type tryAppendMeasuresReq struct { MetricID uint32 + Measures []proto.Measure ResultCh chan tryAppendMeasuresResult } @@ -623,18 +631,18 @@ func (s *Database) startAppendMeasures(metric *_metric, req tryAppendMeasuresReq } req.ResultCh <- tryAppendMeasuresResult{ - ResultCode: CanAppend, - MetricType: metric.MetricType, - FracDigits: metric.FracDigits, - Since: metric.Since, - Until: metric.Until, - UntilValue: metric.UntilValue, - RootPageNo: metric.RootPageNo, - PrevPageNo: metric.LastPageNo, - TimestampsBuf: timestampsBuf, - ValuesBuf: valuesBuf, - Timestamps: timestamps, - Values: values, + ResultCode: CanAppend, + // MetricType: metric.MetricType, + // FracDigits: metric.FracDigits, + // Since: metric.Since, + // Until: metric.Until, + // UntilValue: metric.UntilValue, + // RootPageNo: metric.RootPageNo, + // PrevPageNo: metric.LastPageNo, + // TimestampsBuf: timestampsBuf, + // ValuesBuf: valuesBuf, + // Timestamps: timestamps, + // Values: values, } } @@ -712,7 +720,7 @@ func (*Database) startRangeScan(metric *_metric, req tryRangeScanReq) bool { } } - timestampDecompressor := chunkenc.NewReverseTimeDeltaOfDeltaDecompressor( + timestampDecompressor := chunkenc.NewReverseTimeDeltaDecompressor( metric.TimestampsBuf, metric.Timestamps.Size(), ) @@ -809,7 +817,7 @@ func (*Database) startFullScan(metric *_metric, req tryFullScanReq) bool { return false } - timestampDecompressor := chunkenc.NewReverseTimeDeltaOfDeltaDecompressor( + timestampDecompressor := chunkenc.NewReverseTimeDeltaDecompressor( metric.TimestampsBuf, metric.Timestamps.Size(), ) @@ -877,14 +885,14 @@ func (s *Database) applyChanges(req txlog.Changes) { case txlog.DeletedMetric: s.deleteMetric(rec) - case txlog.AppendedMeasure: - s.appendMeasure(rec) + // case txlog.AppendedMeasure: + // s.appendMeasure(rec) case txlog.AppendedMeasures: s.appendMeasures(rec) - case txlog.AppendedMeasureWithOverflowExtended: - s.appendMeasureAfterOverflow(rec) + // case txlog.AppendedMeasureWithOverflowExtended: + // s.appendMeasureAfterOverflow(rec) case txlog.DeletedMeasures: s.deleteMeasures(rec) @@ -942,7 +950,7 @@ func (s *Database) addMetric(rec txlog.AddedMetric) { FracDigits: byte(rec.FracDigits), TimestampsBuf: timestampsBuf, ValuesBuf: valuesBuf, - Timestamps: chunkenc.NewReverseTimeDeltaOfDeltaCompressor(timestampsBuf, 0), + Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), Values: values, } @@ -976,11 +984,11 @@ func (s *Database) deleteMetric(rec txlog.DeletedMetric) { if len(lockEntry.WaitQueue) > 0 { for _, untyped := range lockEntry.WaitQueue { switch req := untyped.(type) { - case tryAppendMeasureReq: - req.ResultCh <- tryAppendMeasureResult{ - MetricID: req.MetricID, - ResultCode: NoMetric, - } + // case tryAppendMeasureReq: + // req.ResultCh <- tryAppendMeasureResult{ + // MetricID: req.MetricID, + // ResultCode: NoMetric, + // } case tryRangeScanReq: req.ResultCh <- rangeScanResult{ diff --git a/database/snapshot.go b/database/snapshot.go index 74627cc..91987f9 100644 --- a/database/snapshot.go +++ b/database/snapshot.go @@ -8,10 +8,7 @@ import ( "path/filepath" octopus "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" ) @@ -167,91 +164,92 @@ func (s *Database) dumpSnapshot(logNumber int) (err error) { } func (s *Database) loadSnapshot(fileName string) (err error) { - var ( - hasher = crc32.NewIEEE() - metricsQty int - header = make([]byte, metricHeaderSize) - body = make([]byte, atree.PageSize) - ) + // FIX + // var ( + // hasher = crc32.NewIEEE() + // metricsQty int + // 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) + // u64, _, err := bin.ReadVarUint64(src) + // if err != nil { + // return + // } + // metricsQty = int(u64) - 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 = 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:]) - 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.TimestampsBuf = conbuf.NewFromBuffer(buf) - 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.ValuesBuf = conbuf.NewFromBuffer(buf) - metric.Timestamps = chunkenc.NewReverseTimeDeltaOfDeltaCompressor( - metric.TimestampsBuf, int(tSize)) + // metric.Timestamps = chunkenc.NewReverseTimeDeltaOfDeltaCompressor( + // metric.TimestampsBuf, int(tSize)) - 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 - } + // 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 }