diff --git a/database/api.go b/database/api.go index fcb01d7..3ec29b2 100644 --- a/database/api.go +++ b/database/api.go @@ -11,12 +11,11 @@ import ( "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/bufreader" "gordenko.dev/dima/qb/proto" - "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/transform" + "gordenko.dev/dima/qb/worker" ) var ( - ErrNoValueBug = errors.New("has timestamp but no value") ErrWrongResultCodeBug = errors.New("bug: wrong result code") successMsg = []byte{ @@ -215,7 +214,7 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { fmt.Println("add job") - s.worker.AddJobToQueue(tryAddMetricReq{ + s.workerInbox.Push(worker.AddMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) @@ -225,7 +224,7 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { resultCode := <-resultCh switch resultCode { - case Succeed: + case worker.Succeed: fmt.Println("OK") // s.storage.Append(storage.MetricAddRecord{ // MetricID: req.MetricID, @@ -234,7 +233,7 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { // }) //<-waitCh - case MetricDuplicate: + case worker.MetricDuplicate: fmt.Println("ErrDuplicate") return proto.ErrDuplicate @@ -245,16 +244,10 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { return 0 } -type Metric struct { - MetricType qb.MetricType - FracDigits byte - ResultCode byte -} - func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { - resultCh := make(chan Metric, 1) + resultCh := make(chan worker.GetMetricResult, 1) - s.worker.AddJobToQueue(tryGetMetricReq{ + s.workerInbox.Push(worker.GetMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) @@ -262,7 +255,7 @@ func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { result := <-resultCh switch result.ResultCode { - case Succeed: + case worker.Succeed: answer := []byte{ proto.RespValue, 0, 0, 0, 0, // metricID @@ -276,7 +269,7 @@ func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { return err } - case NoMetric: + case worker.NoMetric: reply(conn, proto.ErrNoMetric) default: @@ -285,15 +278,10 @@ func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { return nil } -type tryDeleteMetricResult struct { - ResultCode byte - RootPageNo uint32 -} - func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 { - resultCh := make(chan tryDeleteMetricResult, 1) + resultCh := make(chan worker.DeleteMetricResult, 1) - s.worker.AddJobToQueue(tryDeleteMetricReq{ + s.workerInbox.Push(worker.DeleteMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) @@ -301,7 +289,7 @@ func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 { result := <-resultCh switch result.ResultCode { - case Succeed: + case worker.Succeed: // var ( // //freePageNumbers []uint32 // ) @@ -318,7 +306,7 @@ func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 { // }) //<-waitCh - case NoMetric: + case worker.NoMetric: return proto.ErrNoMetric default: @@ -337,127 +325,10 @@ type FilledPage struct { ValuesSize uint16 } -// 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) - -// s.appendJobToWorkerQueue(tryAppendMeasureReq{ -// MetricID: req.MetricID, -// Timestamp: req.Timestamp, -// Value: req.Value, -// ResultCh: resultCh, -// }) - -// result := <-resultCh - -// switch result.ResultCode { -// case CanAppend: -// waitCh := s.storage.Append(storage.AppendedMeasure{ -// MetricID: req.MetricID, -// Timestamp: req.Timestamp, -// Value: req.Value, -// }) -// <-waitCh - -// case NewPage: -// filled := result.FilledPage - -// var path atree.PathToDataPage -// if filled.RootPageNo > 0 { -// path, err := s.atree.FindPathToLastPage(filled.RootPageNo) -// if err != nil { -// // FIX -// qb.Abort(qb.WriteToAtreeFailed, err) -// } -// // for _, leg := range path.Legs { -// // pagesToRelease = append(pagesToRelease, leg.PageNo) -// // } - -// if path.LastPageNo != filled.PrevPageNo { -// qb.Abort( -// qb.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 { -// // qb.Abort(qb.WriteToAtreeFailed, err) -// // } -// // _ = report - -// //waitCh := s.storage.WriteAppended(path, report) - -// waitCh := s.storage.Append(storage.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 ExpiredMeasure: -// return proto.ErrExpiredMeasure - -// case NonMonotonicValue: -// return proto.ErrNonMonotonicValue - -// default: -// qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) -// } -// return 0 -// } - -type tryAppendMeasuresResult struct { - ResultCode byte - Written int - // MetricType qb.MetricType - // FracDigits byte - // Since uint32 - // Until uint32 - // UntilValue float64 - // RootPageNo uint32 - // PrevPageNo uint32 - // TimestampsBuf *conbuf.ContinuousBuffer - // ValuesBuf *conbuf.ContinuousBuffer - // Timestamps qb.TimestampCompressor - // Values qb.ValueCompressor -} - func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { - resultCh := make(chan tryAppendMeasuresResult, 1) + resultCh := make(chan worker.AppendMeasuresResult, 1) - s.worker.AddJobToQueue(tryAppendMeasuresReq{ + s.workerInbox.Push(worker.AppendMeasuresReq{ MetricID: req.MetricID, Measures: req.Measures, ResultCh: resultCh, @@ -466,7 +337,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { result := <-resultCh switch result.ResultCode { - case CanAppend: + case worker.CanAppend: // report, err := s.atree.AppendDataPage(atree.AppendDataPageReq{}) // if err != nil { @@ -485,7 +356,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { // <-waitCh // } - case NoMetric: + case worker.NoMetric: return proto.ErrNoMetric default: @@ -494,15 +365,10 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { return 0 } -type tryDeleteMeasuresResult struct { - ResultCode byte - RootPageNo uint32 -} - func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { - resultCh := make(chan tryDeleteMeasuresResult, 1) + resultCh := make(chan worker.DeleteMeasuresResult, 1) - s.worker.AddJobToQueue(tryDeleteMeasuresReq{ + s.workerInbox.Push(worker.DeleteMeasuresReq{ MetricID: req.MetricID, Since: req.Since, ResultCh: resultCh, @@ -511,17 +377,17 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { result := <-resultCh switch result.ResultCode { - case NoMeasuresToDelete: + case worker.NoMeasuresToDelete: // ok - case DeleteFromAtreeNotNeeded: + case worker.DeleteFromAtreeNotNeeded: // регистрирую удаление в TransactionLog - s.storage.Append(storage.MeasuresDeleteRecord{ - MetricID: req.MetricID, - }) + // s.storage.Append(storage.MeasuresDeleteRecord{ + // MetricID: req.MetricID, + // }) //<-waitCh - case DeleteFromAtreeRequired: + case worker.DeleteFromAtreeRequired: // собираю номера всех data и index страниц метрики (типа запись REDO лога). // pageNumbers, err := s.atree.GetAllPages(req.MetricID) // if err != nil { @@ -534,7 +400,7 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { // }) //<-waitCh - case NoMetric: + case worker.NoMetric: return proto.ErrNoMetric default: @@ -545,12 +411,6 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { // SELECT -type fullScanResult struct { - ResultCode byte - FracDigits byte - LastPageNo uint32 -} - func (s *Database) ListAllInstantMeasures(conn net.Conn, req proto.ListAllInstantMetricMeasuresReq) error { responseWriter := transform.NewInstantMeasureWriter(conn, 0) @@ -609,13 +469,6 @@ func (s *Database) ListCumulativeMeasures(conn net.Conn, req proto.ListCumulativ }) } -type rangeScanResult struct { - ResultCode byte - FracDigits byte - RootPageNo uint32 - LastPageNo uint32 -} - func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriodsReq) error { since, until := timeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay) if since.After(until) { @@ -681,9 +534,9 @@ type rangeScanReq struct { } func (s *Database) rangeScan(req rangeScanReq) error { - resultCh := make(chan rangeScanResult, 1) + resultCh := make(chan worker.RangeScanResult, 1) - s.worker.AddJobToQueue(tryRangeScanReq{ + s.workerInbox.Push(worker.RangeScanReq{ MetricID: req.MetricID, Since: req.Since, Until: req.Until, @@ -695,44 +548,44 @@ func (s *Database) rangeScan(req rangeScanReq) error { result := <-resultCh switch result.ResultCode { - case QueryDone: + case worker.QueryDone: req.ResponseWriter.Close() - case UntilFound: - err := s.atree.ContinueRangeScan(atree.ContinueRangeScanReq{ - FracDigits: result.FracDigits, - ResponseWriter: req.ResponseWriter, - LastPageNo: result.LastPageNo, - Since: req.Since, - }) + case worker.UntilFound: + // err := s.atree.ContinueRangeScan(atree.ContinueRangeScanReq{ + // FracDigits: result.FracDigits, + // ResponseWriter: req.ResponseWriter, + // LastPageNo: result.LastPageNo, + // Since: req.Since, + // }) //s.metricRUnlock(req.MetricID) - if err != nil { - reply(req.Conn, proto.ErrUnexpected) - } else { - req.ResponseWriter.Close() - } + // if err != nil { + // reply(req.Conn, proto.ErrUnexpected) + // } else { + // req.ResponseWriter.Close() + // } - case UntilNotFound: - err := s.atree.RangeScan(atree.RangeScanReq{ - FracDigits: result.FracDigits, - ResponseWriter: req.ResponseWriter, - RootPageNo: result.RootPageNo, - Since: req.Since, - Until: req.Until, - }) + case worker.UntilNotFound: + // err := s.atree.RangeScan(atree.RangeScanReq{ + // FracDigits: result.FracDigits, + // ResponseWriter: req.ResponseWriter, + // RootPageNo: result.RootPageNo, + // Since: req.Since, + // Until: req.Until, + // }) //s.metricRUnlock(req.MetricID) - if err != nil { - reply(req.Conn, proto.ErrUnexpected) - } else { - req.ResponseWriter.Close() - } + // if err != nil { + // reply(req.Conn, proto.ErrUnexpected) + // } else { + // req.ResponseWriter.Close() + // } - case NoMetric: + case worker.NoMetric: reply(req.Conn, proto.ErrNoMetric) - case WrongMetricType: + case worker.WrongMetricType: reply(req.Conn, proto.ErrWrongMetricType) default: @@ -749,9 +602,9 @@ type fullScanReq struct { } func (s *Database) fullScan(req fullScanReq) error { - resultCh := make(chan fullScanResult, 1) + resultCh := make(chan worker.FullScanResult, 1) - s.worker.AddJobToQueue(tryFullScanReq{ + s.workerInbox.Push(worker.FullScanReq{ MetricID: req.MetricID, MetricType: req.MetricType, ResponseWriter: req.ResponseWriter, @@ -761,26 +614,26 @@ func (s *Database) fullScan(req fullScanReq) error { result := <-resultCh switch result.ResultCode { - case QueryDone: + case worker.QueryDone: req.ResponseWriter.Close() - case UntilFound: - err := s.atree.ContinueFullScan(atree.ContinueFullScanReq{ - FracDigits: result.FracDigits, - ResponseWriter: req.ResponseWriter, - LastPageNo: result.LastPageNo, - }) - s.worker.AddJobToQueue(req.MetricID) - if err != nil { - reply(req.Conn, proto.ErrUnexpected) - } else { - req.ResponseWriter.Close() - } + case worker.UntilFound: + // err := s.atree.ContinueFullScan(atree.ContinueFullScanReq{ + // FracDigits: result.FracDigits, + // ResponseWriter: req.ResponseWriter, + // LastPageNo: result.LastPageNo, + // }) + // s.worker.AddJobToQueue(req.MetricID) + // if err != nil { + // reply(req.Conn, proto.ErrUnexpected) + // } else { + // req.ResponseWriter.Close() + // } - case NoMetric: + case worker.NoMetric: reply(req.Conn, proto.ErrNoMetric) - case WrongMetricType: + case worker.WrongMetricType: reply(req.Conn, proto.ErrWrongMetricType) default: @@ -795,7 +648,7 @@ func (s *Database) ListCurrentValues(conn net.Conn, req proto.ListCurrentValuesR resultCh := make(chan struct{}) - s.worker.AddJobToQueue(tryListCurrentValuesReq{ + s.workerInbox.Push(worker.ListCurrentValuesReq{ MetricIDs: req.MetricIDs, ResponseWriter: responseWriter, ResultCh: resultCh, diff --git a/database/database.go b/database/database.go index 382e149..6e3abb8 100644 --- a/database/database.go +++ b/database/database.go @@ -6,39 +6,30 @@ import ( "log" "net" "os" - "path/filepath" "sync" "time" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/inbox" + "gordenko.dev/dima/qb/recovery" "gordenko.dev/dima/qb/storage" + "gordenko.dev/dima/qb/worker" ) -func JoinSnapshotFileName(dir string, snapshotNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.snapshot", snapshotNumber)) -} - -func JoinWALFileName(dir string, snapshotNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.wal", snapshotNumber)) -} - // type metricLockEntry struct { // XLock bool // RLocks int // WaitQueue []any // } +//metricLockEntries map[uint32]*metricLockEntry type Database struct { - mutex sync.Mutex - worker *Worker - //metricLockEntries map[uint32]*metricLockEntry - //dataFreeList *freelist.FreeList - //indexFreeList *freelist.FreeList + mutex sync.Mutex dir string databaseName string + workerInbox *inbox.Inbox + worker *worker.Worker storage *storage.Writer - atree *atree.Atree tcpPort int logfile *os.File logger *log.Logger @@ -50,7 +41,6 @@ type Options struct { TCPPort int Dir string DatabaseName string - RedoDir string Logfile *os.File ExitCh chan struct{} WaitGroup *sync.WaitGroup @@ -66,9 +56,6 @@ func New(opt Options) (_ *Database, err error) { if opt.DatabaseName == "" { return nil, errors.New("DatabaseName option is required") } - if opt.RedoDir == "" { - return nil, errors.New("RedoDir option is required") - } if opt.Logfile == nil { return nil, errors.New("Logfile option is required") } @@ -79,40 +66,7 @@ func New(opt Options) (_ *Database, err error) { return nil, errors.New("WaitGroup option is required") } - // deltaFreeListFile, err := os.OpenFile( - // filepath.Join(opt.Dir, opt.DatabaseName+".deltafree"), - // os.O_RDWR|os.O_CREATE, - // 0666, - // ) - // if err != nil { - // return nil, err - // } - - // dataFreeList, err := freelist.New(freelist.Options{ - // PageSize: 2048, - // BaseFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data"), - // DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data_delta"), - // }) - - // indexFreeList, err := freelist.New(freelist.Options{ - // PageSize: 2048, - // BaseFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_index"), - // DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_index_delta"), - // }) - - worker := NewWorker(WorkerOptions{ - // - Metrics: make(map[uint32]*_metric), - Dir: opt.Dir, - SaveToStorage: func(in any) { - fmt.Println("save to storage") - }, - ExitCh: opt.ExitCh, - WaitGroup: opt.WaitGroup, - }) - s := &Database{ - worker: worker, dir: opt.Dir, databaseName: opt.DatabaseName, //metricLockEntries: make(map[uint32]*metricLockEntry), @@ -124,6 +78,42 @@ func New(opt Options) (_ *Database, err error) { exitCh: opt.ExitCh, waitGroup: opt.WaitGroup, } + + recoveryReport, err := recovery.Recovery(s.dir, s.databaseName) + if err != nil { + return + } + + storageInbox := inbox.New() + s.workerInbox = inbox.New() + + s.storage, err = storage.NewWriter(storage.WriterOptions{ + Inbox: storageInbox, + WorkerInbox: s.workerInbox, + Dir: s.dir, + DatabaseName: s.databaseName, + WAL: recoveryReport.WAL, + IndexFreeList: recoveryReport.IndexFreeList, + DataFreeList: recoveryReport.DataFreeList, + ExitCh: s.exitCh, + WaitGroup: s.waitGroup, + }) + if err != nil { + qb.Abort(qb.CreateChangesWriterFailed, err) + + } + + s.worker = worker.NewWorker(worker.WorkerOptions{ + Inbox: s.workerInbox, + StorageInbox: storageInbox, + //Metrics: make(map[uint32]*_metric), FIX + Dir: opt.Dir, + // SaveToStorage: func(in any) { + // fmt.Println("save to storage") + // }, + ExitCh: opt.ExitCh, + WaitGroup: opt.WaitGroup, + }) return s, nil } @@ -133,19 +123,6 @@ func (s *Database) ListenAndServe() (err error) { return fmt.Errorf("net.Listen: %s; port=%d", err, s.tcpPort) } - s.storage, err = storage.NewWriter(storage.WriterOptions{ - Dir: s.dir, - //LogNumber: logNumber, - AppendToWorkerQueue: s.worker.AddJobToQueue, - //FreeList: s.freeList, - //Atree: s.atree, - ExitCh: s.exitCh, - WaitGroup: s.waitGroup, - }) - if err != nil { - qb.Abort(qb.CreateChangesWriterFailed, err) - - } go s.storage.Run() // s.atree, err = atree.New(atree.Options{ @@ -177,50 +154,14 @@ func (s *Database) ListenAndServe() (err error) { // зробити object? -func (s *Database) replayChanges(snapshotNumber int) (err error) { - snapshot, err := readSnapshot(JoinSnapshotFileName(s.dir, snapshotNumber)) - if err != nil { - return - } +// func (s *Database) replayChanges(snapshotNumber int) (err error) { +// snapshot, err := readSnapshot(JoinSnapshotFileName(s.dir, snapshotNumber)) +// if err != nil { +// return +// } - walReader, err := storage.NewWALReader(storage.WALReaderOptions{ - FileName: JoinWALFileName(s.dir, snapshotNumber), - BufferSize: 8 * 1024 * 1024, - }) - if err != nil { - return err - } - - walReplayer, err := NewWALReplayer(WALReplayerOptions{ - WALReader: walReader, - Metrics: snapshot.Metrics, - FreeIndexPages: snapshot.IndexPageNumbers, - FreeDataPages: snapshot.DataPageNumbers, - }) - if err != nil { - return - } - - err = walReplayer.Replay() - if err != nil { - return - } - - err = s.storage.WritePagesToAtree(walReplayer.IndexPages(), walReplayer.DataPages()) - if err != nil { - return err - } - - // FIX free pages sync - - return nil -} - -func (s *Database) relayMetricsToMetrics(replayMetrics map[uint32]*ReplayMetric) { - //for metricID, x := range replayMetrics { - //s.metrics[metricID] = x.ToMetric() FIX - //} -} +// return nil +// } //func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { // file, err := os.Open(fileName) diff --git a/database/database_test.go b/database/database_test.go index 5ef4f5c..a50df16 100644 --- a/database/database_test.go +++ b/database/database_test.go @@ -1,113 +1,24 @@ package database -import ( - "fmt" - "slices" - "testing" - - "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/storage" -) - -func TestWriteReadSnapshot(t *testing.T) { - metrics := make(map[uint32]*_metric) - metrics[1] = &_metric{ - metricType: qb.Cumulative, - fracDigits: 3, - lastPageNo: 1000, - buffer: []byte{}, - } - in := writeSnapshotIn{ - SnapshotNumber: 1, - Dir: ".", - WriteBufferSize: 8 * 1024 * 1024, - Metrics: metrics, - FrozenIndexPagesCount: 7, - IndexPageNumbers: []uint32{ - 1, 2, 3, 4, 5, - }, - FrozenDataPagesCount: 8, - DataPageNumbers: []uint32{ - 6, 7, 8, - }, - } - err := writeSnapshot(in) - if err != nil { - t.Fatalf("writeSnapshot: %s", err) - } - - out, err := readSnapshot("1.snapshot") - if err != nil { - return - } - - if out.FrozenIndexPagesCount != in.FrozenIndexPagesCount { - t.Fatalf("FrozenIndexPagesCount: got %d are not equal expected %d", - out.FrozenIndexPagesCount, in.FrozenIndexPagesCount) - } - if out.FrozenDataPagesCount != in.FrozenDataPagesCount { - t.Fatalf("FrozenDataPagesCount: got %d are not equal expected %d", - out.FrozenDataPagesCount, in.FrozenDataPagesCount) - } - if !slices.Equal(out.IndexPageNumbers, in.IndexPageNumbers) { - t.Fatalf("IndexPageNumbers: got %v are not equal expected %v", - out.IndexPageNumbers, in.IndexPageNumbers) - } - if !slices.Equal(out.DataPageNumbers, in.DataPageNumbers) { - t.Fatalf("DataPageNumbers: got %v are not equal expected %v", - out.DataPageNumbers, in.DataPageNumbers) - } - for metricID, replay := range out.Metrics { - origin, ok := metrics[metricID] - if !ok { - t.Fatalf("decoded metricID %d not found in metrics", metricID) - } - err = cmpMetric(replay, origin) - if err != nil { - t.Fatalf("metric %d: %s", metricID, err) - } - } -} - -func cmpMetric(replay *ReplayMetric, origin *_metric) error { - if replay.metricType != origin.metricType { - return fmt.Errorf("metricType: got %d are not equal expected %d", - replay.metricType, origin.metricType) - } - if replay.fracDigits != origin.fracDigits { - return fmt.Errorf("fracDigits: got %d are not equal expected %d", - replay.fracDigits, origin.fracDigits) - } - if replay.lastPageNo != origin.lastPageNo { - return fmt.Errorf("lastPageNo: got %d are not equal expected %d", - replay.lastPageNo, origin.lastPageNo) - } - // if !refect.DeepEq{ - // return fmt.Errorf("lastPageNo: got %d are not equal expected %d", - // replay.lastPageNo, origin.lastPageNo) - // } - return nil -} - -func TestComposeHeadIndexPage(t *testing.T) { - var ( - levelIdx = 0 - level = storage.IndexLevelTail{ - Buffer: []byte{ - 1, 2, 3, - }, - RecordsCount: 2, - } - head = &storage.IndexPageTail{ - PageNo: 100, - CRC32: 12345, - Records: []byte{}, - } - ) - page := composeHeadIndexPage(levelIdx, level, head) - if page.PageNo != head.PageNo { - t.Fatalf("PageNo: got %d are not equal expected %v", - page.PageNo, head.PageNo) - } - // fix compare pages -} +// func TestComposeHeadIndexPage(t *testing.T) { +// var ( +// levelIdx = 0 +// level = storage.IndexLevelTail{ +// Buffer: []byte{ +// 1, 2, 3, +// }, +// RecordsCount: 2, +// } +// head = &storage.IndexPageTail{ +// PageNo: 100, +// CRC32: 12345, +// Records: []byte{}, +// } +// ) +// page := composeHeadIndexPage(levelIdx, level, head) +// if page.PageNo != head.PageNo { +// t.Fatalf("PageNo: got %d are not equal expected %v", +// page.PageNo, head.PageNo) +// } +// // fix compare pages +// } diff --git a/examples/database/main.go b/examples/database/main.go index 7fc2b9f..e7cfc2e 100644 --- a/examples/database/main.go +++ b/examples/database/main.go @@ -42,7 +42,6 @@ func main() { TCPPort: config.TcpPort, Dir: config.Dir, DatabaseName: config.DatabaseName, - RedoDir: config.REDODir, Logfile: logfile, ExitCh: exitCh, WaitGroup: wg, diff --git a/inbox/inbox.go b/inbox/inbox.go new file mode 100644 index 0000000..6c83a8d --- /dev/null +++ b/inbox/inbox.go @@ -0,0 +1,38 @@ +package inbox + +import "sync" + +type Inbox struct { + mutex *sync.Mutex + signalCh chan struct{} + items []any +} + +func New() *Inbox { + return &Inbox{ + mutex: new(sync.Mutex), + signalCh: make(chan struct{}, 1), + } +} + +func (s *Inbox) Ready() chan struct{} { + return s.signalCh +} + +func (s *Inbox) Push(x any) { + s.mutex.Lock() + s.items = append(s.items, x) + s.mutex.Unlock() + select { + case s.signalCh <- struct{}{}: + default: + } +} + +func (s *Inbox) Drain() []any { + s.mutex.Lock() + items := s.items + s.items = nil + s.mutex.Unlock() + return items +} diff --git a/qb.go b/qb.go index 32a0acb..934f77a 100644 --- a/qb.go +++ b/qb.go @@ -4,6 +4,7 @@ import ( "fmt" "io" "os" + "path/filepath" ) type MetricType byte @@ -151,3 +152,27 @@ func Abort(code AbortCode, err error) { fmt.Println(err) os.Exit(int(code)) } + +func GetIndexFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.index", databaseName)) +} + +func GetDataFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.data", databaseName)) +} + +func GetWALFilePath(dir string, databaseName string, snapshotNumber int) string { + return filepath.Join(dir, fmt.Sprintf("%s.wal_%d", databaseName, snapshotNumber)) +} + +func GetSnapshotFilePath(dir string, databaseName string, snapshotNumber int) string { + return filepath.Join(dir, fmt.Sprintf("%s.snapshot_%d", databaseName, snapshotNumber)) +} + +func GetIndexFreeListFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.index_free", databaseName)) +} + +func GetDataFreeListFilePath(dir string, databaseName string) string { + return filepath.Join(dir, fmt.Sprintf("%s.data_free", databaseName)) +} diff --git a/recovery/advisor.go b/recovery/advisor.go deleted file mode 100644 index 8013f44..0000000 --- a/recovery/advisor.go +++ /dev/null @@ -1,255 +0,0 @@ -package recovery - -import ( - "errors" - "fmt" - "os" - "path/filepath" - "regexp" - "sort" - "strconv" -) - -var ( - reChanges = regexp.MustCompile(`(\d+)\.wal`) - reSnapshot = regexp.MustCompile(`(\d+)\.snapshot`) -) - -func joinWALFileName(dir string, snapshotNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.wal", snapshotNumber)) -} - -type RecoveryRecipe struct { - Snapshot string - Changes []string - SnapshotNumber int - ToDelete []string - CompleteSnapshot bool // флаг - что нужно завершить создание снапшота -} - -type RecoveryAdvisor struct { - dir string - verifySnapshot func(string) (bool, error) // (fileName) isVerified, error -} - -type RecoveryAdvisorOptions struct { - Dir string - VerifySnapshot func(string) (bool, error) -} - -func NewRecoveryAdvisor(opt RecoveryAdvisorOptions) (*RecoveryAdvisor, error) { - if opt.Dir == "" { - return nil, errors.New("Dir option is required") - } - if opt.VerifySnapshot == nil { - return nil, errors.New("VerifySnapshot option is required") - } - return &RecoveryAdvisor{ - dir: opt.Dir, - verifySnapshot: opt.VerifySnapshot, - }, nil -} - -type SnapshotChangesPair struct { - SnapshotFileName string - WALFileName string - SnapshotNumber int -} - -func (s *RecoveryAdvisor) getSnapshotChangesPairs() (*SnapshotChangesPair, *SnapshotChangesPair, error) { - var ( - numSet = make(map[int]bool) - changesSet = make(map[int]bool) - snapshotsSet = make(map[int]bool) - pairs []SnapshotChangesPair - ) - - entries, err := os.ReadDir(s.dir) - if err != nil { - return nil, nil, err - } - - for _, entry := range entries { - if entry.Type().IsRegular() { - baseName := entry.Name() - groups := reChanges.FindStringSubmatch(baseName) - if len(groups) == 2 { - num, _ := strconv.Atoi(groups[1]) - - numSet[num] = true - changesSet[num] = true - } - groups = reSnapshot.FindStringSubmatch(baseName) - if len(groups) == 2 { - num, _ := strconv.Atoi(groups[1]) - - numSet[num] = true - snapshotsSet[num] = true - } - } - } - - for snapshotNumber := range numSet { - var ( - snapshotFileName string - changesFileName string - ) - if changesSet[snapshotNumber] { - changesFileName = joinWALFileName(s.dir, snapshotNumber) - } - if snapshotsSet[snapshotNumber] { - snapshotFileName = filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", snapshotNumber)) - } - pairs = append(pairs, SnapshotChangesPair{ - WALFileName: changesFileName, - SnapshotFileName: snapshotFileName, - SnapshotNumber: snapshotNumber, - }) - } - - if len(pairs) == 0 { - return nil, nil, nil - } - - sort.Slice(pairs, func(i, j int) bool { - return pairs[i].SnapshotNumber > pairs[j].SnapshotNumber - }) - - pair := pairs[0] - if pair.WALFileName == "" { - return nil, nil, fmt.Errorf("has %d.shapshot file, but %d.changes file not found", - pair.SnapshotNumber, pair.SnapshotNumber) - } - - if len(pairs) > 1 { - prevPair := pairs[1] - if prevPair.SnapshotFileName == "" && prevPair.SnapshotNumber != 1 { - return &pair, nil, nil - } - - if prevPair.WALFileName == "" && pair.SnapshotFileName == "" { - return &pair, nil, nil - } - return &pair, &prevPair, nil - } else { - return &pair, nil, nil - } -} - -func (s *RecoveryAdvisor) GetRecipe() (*RecoveryRecipe, error) { - pair, prevPair, err := s.getSnapshotChangesPairs() - if err != nil { - return nil, err - } - - if pair == nil { - return nil, nil - } - - if pair.SnapshotFileName != "" { - isVerified, err := s.verifySnapshot(pair.SnapshotFileName) - if err != nil { - return nil, fmt.Errorf("verifySnapshot %s: %s", - pair.SnapshotFileName, err) - } - - if isVerified { - recipe := &RecoveryRecipe{ - Snapshot: pair.SnapshotFileName, - Changes: []string{ - pair.WALFileName, - }, - SnapshotNumber: pair.SnapshotNumber, - } - - if prevPair != nil { - if prevPair.WALFileName != "" { - recipe.ToDelete = append(recipe.ToDelete, prevPair.WALFileName) - } - if prevPair.SnapshotFileName != "" { - recipe.ToDelete = append(recipe.ToDelete, prevPair.SnapshotFileName) - } - } - return recipe, nil - } - if prevPair != nil { - return s.tryPrevPair(pair, prevPair) - } - return nil, fmt.Errorf("%d.shapshot is corrupted", pair.SnapshotNumber) - } else { - if prevPair != nil { - return s.tryPrevPair(pair, prevPair) - } else { - if pair.SnapshotNumber == 1 { - return &RecoveryRecipe{ - Changes: []string{ - pair.WALFileName, - }, - SnapshotNumber: pair.SnapshotNumber, - }, nil - } else { - return nil, fmt.Errorf("%d.snapshot not found", pair.SnapshotNumber) - } - } - } -} - -func (s *RecoveryAdvisor) tryPrevPair(pair, prevPair *SnapshotChangesPair) (*RecoveryRecipe, error) { - if prevPair.WALFileName == "" { - if pair.SnapshotFileName != "" { - return nil, fmt.Errorf("%d.shapshot is corrupted and %d.changes not found", - pair.SnapshotNumber, prevPair.SnapshotNumber) - } else { - return nil, fmt.Errorf("%d.changes not found", prevPair.SnapshotNumber) - } - } - - if prevPair.SnapshotFileName == "" { - if prevPair.SnapshotNumber == 1 { - recipe := &RecoveryRecipe{ - Changes: []string{ - prevPair.WALFileName, - pair.WALFileName, - }, - SnapshotNumber: pair.SnapshotNumber, - CompleteSnapshot: true, - ToDelete: []string{ - prevPair.WALFileName, - }, - } - return recipe, nil - } else { - if pair.SnapshotFileName != "" { - return nil, fmt.Errorf("%d.shapshot is corrupted and %d.snapshot not found", - pair.SnapshotNumber, prevPair.SnapshotNumber) - } else { - return nil, fmt.Errorf("%d.snapshot not found", pair.SnapshotNumber) - } - } - } - - isVerified, err := s.verifySnapshot(prevPair.SnapshotFileName) - if err != nil { - return nil, fmt.Errorf("verifySnapshot %s: %s", - prevPair.SnapshotFileName, err) - } - - if !isVerified { - return nil, fmt.Errorf("%d.shapshot is corrupted", prevPair.SnapshotNumber) - } - - recipe := &RecoveryRecipe{ - Snapshot: prevPair.SnapshotFileName, - Changes: []string{ - prevPair.WALFileName, - pair.WALFileName, - }, - SnapshotNumber: pair.SnapshotNumber, - CompleteSnapshot: true, - ToDelete: []string{ - prevPair.WALFileName, - prevPair.SnapshotFileName, - }, - } - return recipe, nil -} diff --git a/recovery/recovery.go b/recovery/recovery.go new file mode 100644 index 0000000..5051335 --- /dev/null +++ b/recovery/recovery.go @@ -0,0 +1,221 @@ +package recovery + +import ( + "errors" + "os" + "path/filepath" + "strconv" + "strings" + + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/freelist" + "gordenko.dev/dima/qb/storage" +) + +// Result зберігає результати для обох типів файлів +type FileVersion struct { + Name string + Version int +} + +// FindMaxSnapshotFiles шукає файли з розширеннями .snapshot_XXX та .wal_XXX +// і повертає структуру з максимальними номерами. +func findFileLatestVersion(dir string, extPrefix string) (FileVersion, error) { + var x FileVersion + if !strings.HasPrefix(extPrefix, ".") { + extPrefix = "." + extPrefix + } + entries, err := os.ReadDir(dir) + if err != nil { + return x, err + } + for _, entry := range entries { + if entry.IsDir() { + continue + } + name := entry.Name() + ext := filepath.Ext(name) + + if strings.HasPrefix(ext, extPrefix) { + version, err := strconv.Atoi(strings.TrimPrefix(ext, extPrefix)) + if err == nil { + if version > x.Version || x.Name == "" { + x.Name = name + x.Version = version + } + } + } + } + return x, nil +} + +var ( + ErrLostSnapshot = errors.New("lost snapshot") +) + +type RecoveryFiles struct { + Snapshot string + WAL string + SnapshotNumber int +} + +func resolveRecoveryFiles(dir string) (_ RecoveryFiles, err error) { + wal, err := findFileLatestVersion(dir, "wal_") + if err != nil { + return + } + snapshot, err := findFileLatestVersion(dir, "snapshot_") + if err != nil { + return + } + if snapshot.Name != "" { + if wal.Name != "" { + if snapshot.Version == wal.Version { + return RecoveryFiles{ + Snapshot: snapshot.Name, + WAL: wal.Name, + SnapshotNumber: snapshot.Version, + }, nil + } else if snapshot.Version > wal.Version { + return RecoveryFiles{ + Snapshot: snapshot.Name, + SnapshotNumber: snapshot.Version, + }, nil + } else { + // WAL version > snapshot version + err = ErrLostSnapshot + return + } + } else { + return RecoveryFiles{ + Snapshot: snapshot.Name, + SnapshotNumber: snapshot.Version, + }, nil + } + } else { + if wal.Name != "" { + if wal.Version == 0 { + return RecoveryFiles{ + WAL: wal.Name, + }, nil + } else { + err = ErrLostSnapshot + return + } + } else { + // no files - ok, 1st start of database + return + } + } +} + +type RecoveryReport struct { + //Metrics [] + SnapshotNumber int + WAL string + IndexFreeList *freelist.FreeList + DataFreeList *freelist.FreeList +} + +func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) { + files, err := resolveRecoveryFiles(dir) + if err != nil { + return + } + var ( + freeIndexPages []uint32 + freeDataPages []uint32 + ) + if files.Snapshot != "" { + var out storage.ReadSnapshotOut + out, err = storage.ReadSnapshot(files.Snapshot) + if err != nil { + return + } + // + freeIndexPages = out.IndexPageNumbers + freeDataPages = out.DataPageNumbers + } + + if files.WAL != "" { + var ( + walReader *storage.WALReader + walReplayer *storage.WALReplayer + ) + walReader, err = storage.NewWALReader(storage.WALReaderOptions{ + //FileName: JoinWALFileName(s.dir, files.WAL), + BufferSize: 8 * 1024 * 1024, + }) + if err != nil { + return + } + + walReplayer, err = storage.NewWALReplayer(storage.WALReplayerOptions{ + WALReader: walReader, + Metrics: snapshot.Metrics, + FreeIndexPages: freeIndexPages, + FreeDataPages: freeDataPages, + }) + if err != nil { + return + } + + err = walReplayer.Replay() + if err != nil { + return + } + // перезаписую сторінки із останнього комміта + var dataFile *os.File + dataFile, err = os.OpenFile(qb.GetDataFilePath(dir, databaseName), os.O_WRONLY, 0666) + if err != nil { + return + } + err = storage.WriteDataPages(dataFile, walReplayer.DataPages()) + if err != nil { + return + } + err = dataFile.Close() + if err != nil { + return + } + var indexFile *os.File + indexFile, err = os.OpenFile(qb.GetIndexFilePath(dir, databaseName), os.O_WRONLY, 0666) + if err != nil { + return + } + err = storage.WriteIndexPages(indexFile, walReplayer.IndexPages()) + if err != nil { + return + } + err = indexFile.Close() + if err != nil { + return + } + } + + indexFreeList, err := freelist.New(freelist.Options{ + PageSize: 2048, + BaseFilePath: filepath.Join(dir, databaseName+".free_index"), + DeltaFilePath: filepath.Join(dir, databaseName+".free_index_delta"), + }) + + // FIX free pages sync + dataFreeList, err := freelist.New(freelist.Options{ + PageSize: 2048, + BaseFilePath: filepath.Join(dir, databaseName+".free_data"), + DeltaFilePath: filepath.Join(dir, databaseName+".free_data_delta"), + }) + + return RecoveryReport{ + SnapshotNumber: files.SnapshotNumber, + WAL: files.WAL, + IndexFreeList: indexFreeList, + DataFreeList: dataFreeList, + }, nil +} + +//func (s *Database) relayMetricsToMetrics(replayMetrics map[uint32]*ReplayMetric) { +//for metricID, x := range replayMetrics { +//s.metrics[metricID] = x.ToMetric() FIX +//} +//} diff --git a/recovery/recovery_test.go b/recovery/recovery_test.go new file mode 100644 index 0000000..d7b459b --- /dev/null +++ b/recovery/recovery_test.go @@ -0,0 +1,18 @@ +package recovery + +import ( + "testing" +) + +func TestFindMaxVersion(t *testing.T) { + result, err := findFileLatestVersion(".", "wal_") + if err != nil { + t.Fatal(err) + } + if result.Name != "test.wal_123" { + t.Fatalf("wrong file name %s instead of test.wal_123", result.Name) + } + if result.Version != 123 { + t.Fatalf("wrong file version %d instead of 123", result.Version) + } +} diff --git a/recovery/test.wal_0 b/recovery/test.wal_0 new file mode 100644 index 0000000..e69de29 diff --git a/recovery/test.wal_123 b/recovery/test.wal_123 new file mode 100644 index 0000000..e69de29 diff --git a/storage/page_ingester.go b/storage/page_ingester.go index 7eced42..bf93560 100644 --- a/storage/page_ingester.go +++ b/storage/page_ingester.go @@ -236,7 +236,7 @@ func appendIndexRecord(in appendIndexRecordIn) []byte { buf = in.Records ) if in.RecordsCount > 0 { - pos = in.RecordsCount * indexRecordSize + pos = in.RecordsCount * IndexRecordSize //fmt.Printf("pos: %d\n", pos) } else { buf = make([]byte, IndexPageSize) // IndexPageIncSize diff --git a/storage/pageno_provider.go b/storage/pageno_provider.go new file mode 100644 index 0000000..8321676 --- /dev/null +++ b/storage/pageno_provider.go @@ -0,0 +1,55 @@ +package storage + +import ( + "errors" + "math" + "os" + + "gordenko.dev/dima/qb/freelist" +) + +type PageNoProvider struct { + freeList *freelist.FreeList + pagesCount uint32 +} + +type PageNoProviderOptions struct { + FilePath string + PageSize int + FreeList *freelist.FreeList +} + +func NewPageNoProvider(opt PageNoProviderOptions) (_ *PageNoProvider, err error) { + if opt.FreeList == nil { + return nil, errors.New("FreeList option is required") + } + file, err := os.Open(opt.FilePath) + if err != nil { + return + } + defer file.Close() + info, err := file.Stat() + if err != nil { + return + } + pagesCount := info.Size() / int64(opt.PageSize) + return &PageNoProvider{ + freeList: opt.FreeList, + pagesCount: uint32(pagesCount), + }, nil +} + +func (s *PageNoProvider) GetPageNumber() (uint32, bool, error) { + pageNo, err := s.freeList.GetPageNumber() + if err != nil { + return 0, false, err + } + if pageNo > 0 { + return pageNo, true, nil + } + if s.pagesCount < math.MaxUint32 { + s.pagesCount++ + return s.pagesCount, false, nil + } + return 0, false, errors.New("no space") +} diff --git a/database/replay_metric.go b/storage/replay_metric.go similarity index 72% rename from database/replay_metric.go rename to storage/replay_metric.go index c38ae58..8ec3c5b 100644 --- a/database/replay_metric.go +++ b/storage/replay_metric.go @@ -1,4 +1,4 @@ -package database +package storage import ( "fmt" @@ -6,8 +6,6 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/enc" - "gordenko.dev/dima/qb/storage" ) type ReplayMetric struct { @@ -18,7 +16,7 @@ type ReplayMetric struct { databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер) tSize int vSize int - indexLevelTails []storage.IndexLevelTail // root - last element + indexLevelTails []IndexLevelTail // root - last element } func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { @@ -58,18 +56,18 @@ func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { } for range levelsCount { var ( - buf = make([]byte, storage.IndexPageSize) + buf = make([]byte, IndexPageSize) count int ) count, err = bin.ReadVarSize(r) if err != nil { return } - err = bin.ReadNInto(r, buf[:count*indexRecordSize]) + err = bin.ReadNInto(r, buf[:count*IndexRecordSize]) if err != nil { return } - s.indexLevelTails = append(s.indexLevelTails, storage.IndexLevelTail{ + s.indexLevelTails = append(s.indexLevelTails, IndexLevelTail{ Buffer: buf, RecordsCount: count, }) @@ -77,16 +75,16 @@ func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { return } -// func appendNewRecordsToIndexTail(level storage.IndexLevelTail, records []byte) { +// func appendNewRecordsToIndexTail(level IndexLevelTail, records []byte) { // } -// func countReusedIndexPages(changedIndexLevels []storage.ChangedIndexLevel) (count int) { +// func countReusedIndexPages(changedIndexLevels []ChangedIndexLevel) (count int) { // return // } -func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) { +func (s *ReplayMetric) MeasuresAppend(rec MeasuresAppendRecord) { timestampsSize := rec.TimestampsRewindOffset - len(rec.Timestamps) // копіюю нові дані із WAL з урахуванням offset-ів copy(s.databuf[rec.ValuesRewindOffset:], rec.Values) @@ -97,24 +95,24 @@ func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) { } type AppendMeasuresResult struct { - IndexPages []storage.PageToWrite - DataPages []storage.PageToWrite + IndexPages []PageToWrite + DataPages []PageToWrite ReusedIndexPagesCount int ReusedDataPagesCount int } -func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrowRecord, isLastPacket bool) (_ AppendMeasuresResult) { +func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (_ AppendMeasuresResult) { var ( reusedIndexPages int reusedDataPages int - indexPages []storage.PageToWrite - dataPages []storage.PageToWrite + indexPages []PageToWrite + dataPages []PageToWrite ) // Додати в index level tails недостаючі дані, або замінити for levelIdx, change := range rec.ChangedIndexLevels { var ( head = change.IndexPageTail - newRecordsCount = len(change.TailRecords) / indexRecordSize + newRecordsCount = len(change.TailRecords) / IndexRecordSize ) // 1. набиваю сторінки для перезапису в index файлі if isLastPacket { @@ -123,7 +121,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow composeHeadIndexPage(levelIdx, s.indexLevelTails[levelIdx], head)) } for _, x := range change.IndexPages { - indexPages = append(indexPages, storage.PageToWrite{ + indexPages = append(indexPages, PageToWrite{ PageNo: x.PageNo, Content: x.Content, }) @@ -140,17 +138,17 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow level.RecordsCount = newRecordsCount } else { // TailRecords - це нові дані, які треба додати - pos := level.RecordsCount * indexRecordSize + pos := level.RecordsCount * IndexRecordSize copy(level.Buffer[pos:], change.TailRecords) level.RecordsCount += newRecordsCount } s.indexLevelTails[levelIdx] = level } else { // додаю новий індексний рівень - buf := make([]byte, storage.IndexPageSize) + buf := make([]byte, IndexPageSize) copy(buf, change.TailRecords) // - s.indexLevelTails = append(s.indexLevelTails, storage.IndexLevelTail{ + s.indexLevelTails = append(s.indexLevelTails, IndexLevelTail{ Buffer: buf, RecordsCount: newRecordsCount, }) @@ -174,7 +172,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow dataPages = append(dataPages, composeHeadDataPage(s.databuf, rec.DataPageTail)) for _, x := range rec.DataPages { - dataPages = append(dataPages, storage.PageToWrite{ + dataPages = append(dataPages, PageToWrite{ PageNo: x.PageNo, Content: x.Content, }) @@ -215,17 +213,17 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrow // HELPERS -func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *storage.IndexPageTail) storage.PageToWrite { +func composeHeadIndexPage(levelIdx int, level IndexLevelTail, head *IndexPageTail) PageToWrite { // розраховую pos, з якого буду дописувати хвіст - pos := level.RecordsCount * indexRecordSize + pos := level.RecordsCount * IndexRecordSize // створюю копію сторінки - page := make([]byte, storage.IndexPageSize) + page := make([]byte, IndexPageSize) copy(page, level.Buffer[:pos]) // поточні дані copy(page[pos:], head.Records) // запечатати сторінку - calculatedCRC := storage.SealIndexPage(storage.SealIndexPageIn{ + calculatedCRC := SealIndexPage(SealIndexPageIn{ Content: page, - RecordsCount: level.RecordsCount + len(head.Records)/indexRecordSize, + RecordsCount: level.RecordsCount + len(head.Records)/IndexRecordSize, ZeroLevel: levelIdx == 0, }) // перевірка CRC @@ -234,16 +232,16 @@ func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *stor fmt.Errorf("calculated CRC %d not equal expected %d of head page %d on index level %d", calculatedCRC, head.CRC32, head.PageNo, levelIdx)) } - return storage.PageToWrite{ + return PageToWrite{ PageNo: head.PageNo, Content: page, } } -func composeHeadDataPage(databuf []byte, head storage.DataPageTail) storage.PageToWrite { +func composeHeadDataPage(databuf []byte, head DataPageTail) PageToWrite { // створюю копію сторінки var ( - page = make([]byte, storage.DataPageSize) + page = make([]byte, DataPageSize) timestampsSize = head.TimestampsRewindOffset - len(head.Timestamps) ) // копіюю поточні дані @@ -252,7 +250,7 @@ func composeHeadDataPage(databuf []byte, head storage.DataPageTail) storage.Page copy(page[head.ValuesRewindOffset:], head.Values) copy(page[len(page)-timestampsSize:], head.Timestamps) // запечатати сторінку - calculatedCRC := storage.SealDataPage(storage.SealDataPageIn{ + calculatedCRC := SealDataPage(SealDataPageIn{ Content: page, PrevPageNo: head.PrevPageNo, TimestampsSize: timestampsSize, @@ -264,24 +262,24 @@ func composeHeadDataPage(databuf []byte, head storage.DataPageTail) storage.Page fmt.Errorf("calculated CRC %d not equal expected %d of head data page %d", calculatedCRC, head.CRC32, head.PageNo)) } - return storage.PageToWrite{ + return PageToWrite{ PageNo: head.PageNo, Content: page, } } -func (s *ReplayMetric) ToMetric() *_metric { +// func (s *ReplayMetric) ToMetric() *worker.metric { - databuf := s.buf[:storage.DataPagePayloadSize] - values := enc.NewValueDeltaCompressor(s.metricType, s.fracDigits, databuf, s.vSize) - return &_metric{ - metricType: s.metricType, - fracDigits: s.fracDigits, - lastPageNo: s.lastPageNo, - lastValue: values.LastValue(), - buffer: s.buf, - timestamps: enc.NewTimeDeltaCompressor(databuf, s.tSize), - values: values, - indexLevelTails: s.indexLevelTails, - } -} +// databuf := s.buf[:storage.DataPagePayloadSize] +// values := enc.NewValueDeltaCompressor(s.metricType, s.fracDigits, databuf, s.vSize) +// return &_metric{ +// metricType: s.metricType, +// fracDigits: s.fracDigits, +// lastPageNo: s.lastPageNo, +// lastValue: values.LastValue(), +// buffer: s.buf, +// timestamps: enc.NewTimeDeltaCompressor(databuf, s.tSize), +// values: values, +// indexLevelTails: s.indexLevelTails, +// } +// } diff --git a/database/snapshot.go b/storage/snapshot.go similarity index 91% rename from database/snapshot.go rename to storage/snapshot.go index c2a1cff..ddaa3ee 100644 --- a/database/snapshot.go +++ b/storage/snapshot.go @@ -1,4 +1,4 @@ -package database +package storage import ( "bufio" @@ -8,7 +8,6 @@ import ( "path/filepath" bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/util" ) @@ -29,20 +28,24 @@ dataFreeList - Nb CRC32 - 4b */ -const readBufferSize = 8 * 1024 * 1024 // 8mb +const ReadBufferSize = 8 * 1024 * 1024 // 8mb -type writeSnapshotIn struct { +type WritableMetric interface { + WriteTo(io.Writer) error +} + +type WriteSnapshotIn struct { SnapshotNumber int Dir string WriteBufferSize int - Metrics map[uint32]*_metric + Metrics map[uint32]WritableMetric FrozenIndexPagesCount int IndexPageNumbers []uint32 FrozenDataPagesCount int DataPageNumbers []uint32 } -func writeSnapshot(in writeSnapshotIn) (err error) { +func WriteSnapshot(in WriteSnapshotIn) (err error) { var ( fileName = filepath.Join(in.Dir, fmt.Sprintf("%d.snapshot", in.SnapshotNumber)) hasher = util.NewHasher() @@ -118,7 +121,7 @@ func writeSnapshot(in writeSnapshotIn) (err error) { return } -type readSnapshotOut struct { +type ReadSnapshotOut struct { Metrics map[uint32]*ReplayMetric FrozenIndexPagesCount int IndexPageNumbers []uint32 @@ -126,7 +129,7 @@ type readSnapshotOut struct { DataPageNumbers []uint32 } -func readSnapshot(fileName string) (out readSnapshotOut, err error) { +func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) { var ( metrics = make(map[uint32]*ReplayMetric) frozenIndexPagesCount int @@ -149,7 +152,7 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { fileSize := stat.Size() if fileSize == 0 { - return readSnapshotOut{ + return ReadSnapshotOut{ Metrics: metrics, }, nil } @@ -177,11 +180,11 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { for range metricsQty { var ( metricID uint32 - buf = make([]byte, storage.DataPageSize) + buf = make([]byte, DataPageSize) // для простоти створюю буфер максимального розміру metric = &ReplayMetric{ buf: buf, - databuf: buf[:storage.DataPagePayloadSize], + databuf: buf[:DataPagePayloadSize], } ) metricID, err = bin.ReadUint32(src) @@ -218,7 +221,7 @@ func readSnapshot(fileName string) (out readSnapshotOut, err error) { return } - return readSnapshotOut{ + return ReadSnapshotOut{ Metrics: metrics, FrozenIndexPagesCount: frozenIndexPagesCount, IndexPageNumbers: indexPageNumbers, diff --git a/storage/storage.go b/storage/storage.go index 3e9ffb2..d5356a3 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -24,7 +24,7 @@ var ( indexRecordsCountIdx = IndexPageSize - 6 isZeroLevelIdx = IndexPageSize - 7 - maxRecordsOnIndexPage = (IndexPageSize - IndexPageFooterSize) / indexRecordSize + maxRecordsOnIndexPage = (IndexPageSize - IndexPageFooterSize) / IndexRecordSize // timestampSize = 4 // pairSize = timestampSize + PageNoSize @@ -35,7 +35,7 @@ var ( const ( DataPageFooterSize = 12 IndexPageFooterSize = 7 - indexRecordSize = 8 + IndexRecordSize = 8 ) type IndexRecord struct { diff --git a/storage/storage_test.go b/storage/storage_test.go index c8adb61..20c92b5 100644 --- a/storage/storage_test.go +++ b/storage/storage_test.go @@ -867,3 +867,83 @@ func TestWritePreparer(t *testing.T) { fmt.Printf("%#v\n", x) } } + +func TestWriteReadSnapshot(t *testing.T) { + metrics := make(map[uint32]*_metric) + metrics[1] = &_metric{ + metricType: qb.Cumulative, + fracDigits: 3, + lastPageNo: 1000, + buffer: []byte{}, + } + in := WriteSnapshotIn{ + SnapshotNumber: 1, + Dir: ".", + WriteBufferSize: 8 * 1024 * 1024, + Metrics: metrics, + FrozenIndexPagesCount: 7, + IndexPageNumbers: []uint32{ + 1, 2, 3, 4, 5, + }, + FrozenDataPagesCount: 8, + DataPageNumbers: []uint32{ + 6, 7, 8, + }, + } + err := WriteSnapshot(in) + if err != nil { + t.Fatalf("writeSnapshot: %s", err) + } + + out, err := ReadSnapshot("1.snapshot") + if err != nil { + return + } + + if out.FrozenIndexPagesCount != in.FrozenIndexPagesCount { + t.Fatalf("FrozenIndexPagesCount: got %d are not equal expected %d", + out.FrozenIndexPagesCount, in.FrozenIndexPagesCount) + } + if out.FrozenDataPagesCount != in.FrozenDataPagesCount { + t.Fatalf("FrozenDataPagesCount: got %d are not equal expected %d", + out.FrozenDataPagesCount, in.FrozenDataPagesCount) + } + if !slices.Equal(out.IndexPageNumbers, in.IndexPageNumbers) { + t.Fatalf("IndexPageNumbers: got %v are not equal expected %v", + out.IndexPageNumbers, in.IndexPageNumbers) + } + if !slices.Equal(out.DataPageNumbers, in.DataPageNumbers) { + t.Fatalf("DataPageNumbers: got %v are not equal expected %v", + out.DataPageNumbers, in.DataPageNumbers) + } + for metricID, replay := range out.Metrics { + origin, ok := metrics[metricID] + if !ok { + t.Fatalf("decoded metricID %d not found in metrics", metricID) + } + err = cmpMetric(replay, origin) + if err != nil { + t.Fatalf("metric %d: %s", metricID, err) + } + } +} + +func cmpMetric(replay *ReplayMetric, origin *_metric) error { + if replay.metricType != origin.metricType { + return fmt.Errorf("metricType: got %d are not equal expected %d", + replay.metricType, origin.metricType) + } + if replay.fracDigits != origin.fracDigits { + return fmt.Errorf("fracDigits: got %d are not equal expected %d", + replay.fracDigits, origin.fracDigits) + } + if replay.lastPageNo != origin.lastPageNo { + return fmt.Errorf("lastPageNo: got %d are not equal expected %d", + replay.lastPageNo, origin.lastPageNo) + } + // if !refect.DeepEq{ + // return fmt.Errorf("lastPageNo: got %d are not equal expected %d", + // replay.lastPageNo, origin.lastPageNo) + // } + return nil +} diff --git a/storage/wal_records.go b/storage/wal_records.go index 899b482..265ebe9 100644 --- a/storage/wal_records.go +++ b/storage/wal_records.go @@ -194,7 +194,7 @@ func (s MeasuresAppendWithGrowRecord) Pack(w io.Writer) { bin.WriteUint32(w, tail.PageNo) bin.WriteBool(w, tail.Reused) bin.WriteUint32(w, tail.CRC32) - bin.WriteVarSize(w, len(tail.Records)/indexRecordSize) + bin.WriteVarSize(w, len(tail.Records)/IndexRecordSize) w.Write(tail.Records) } else { bin.WriteBool(w, false) @@ -207,7 +207,7 @@ func (s MeasuresAppendWithGrowRecord) Pack(w io.Writer) { w.Write(p.Content) } // level tail - bin.WriteVarSize(w, len(level.TailRecords)/indexRecordSize) + bin.WriteVarSize(w, len(level.TailRecords)/IndexRecordSize) w.Write(level.TailRecords) } } @@ -313,7 +313,7 @@ func (s *MeasuresAppendWithGrowRecord) Parse(r io.Reader) (err error) { if err != nil { return } - tail.Records, err = bin.ReadN(r, recordsCount*indexRecordSize) + tail.Records, err = bin.ReadN(r, recordsCount*IndexRecordSize) if err != nil { return } @@ -344,7 +344,7 @@ func (s *MeasuresAppendWithGrowRecord) Parse(r io.Reader) (err error) { if err != nil { return } - level.TailRecords, err = bin.ReadN(r, recordsCount*indexRecordSize) + level.TailRecords, err = bin.ReadN(r, recordsCount*IndexRecordSize) if err != nil { return } diff --git a/database/wal_replayer.go b/storage/wal_replayer.go similarity index 75% rename from database/wal_replayer.go rename to storage/wal_replayer.go index f332a4e..954b91a 100644 --- a/database/wal_replayer.go +++ b/storage/wal_replayer.go @@ -1,23 +1,21 @@ -package database +package storage import ( "errors" "fmt" - - "gordenko.dev/dima/qb/storage" ) type WALReplayer struct { - walReader *storage.WALReader + walReader *WALReader metrics map[uint32]*ReplayMetric freeIndexPages []uint32 freeDataPages []uint32 - indexPages []storage.PageToWrite - dataPages []storage.PageToWrite + indexPages []PageToWrite + dataPages []PageToWrite } type WALReplayerOptions struct { - WALReader *storage.WALReader + WALReader *WALReader Metrics map[uint32]*ReplayMetric // from snapshot FreeIndexPages []uint32 // from snapshot FreeDataPages []uint32 // from snapshot @@ -38,11 +36,11 @@ func NewWALReplayer(opt WALReplayerOptions) (*WALReplayer, error) { }, nil } -func (s *WALReplayer) IndexPages() []storage.PageToWrite { +func (s *WALReplayer) IndexPages() []PageToWrite { return s.indexPages } -func (s *WALReplayer) DataPages() []storage.PageToWrite { +func (s *WALReplayer) DataPages() []PageToWrite { return s.dataPages } @@ -62,23 +60,23 @@ func (s *WALReplayer) Replay() error { func (s *WALReplayer) replayRecord(untyped any, isLastPacket bool) (err error) { switch rec := untyped.(type) { - case storage.MeasuresAppendRecord: + case MeasuresAppendRecord: if err = s.onMeasuresAppend(rec); err != nil { return } - case storage.MeasuresAppendWithGrowRecord: + case MeasuresAppendWithGrowRecord: if err = s.onMeasuresAppendWithGrow(rec, isLastPacket); err != nil { return } - case storage.MetricAddRecord: + case MetricAddRecord: if err = s.onMetricAdd(rec); err != nil { return } - case storage.MetricDeleteRecord: + case MetricDeleteRecord: if err = s.onMetricDelete(rec); err != nil { return } - // case storage.DeletedMeasures: + // case DeletedMeasures: // metric, ok := s.metrics[rec.MetricID] // if ok { // metric.DeleteMeasures() @@ -94,24 +92,24 @@ func (s *WALReplayer) replayRecord(untyped any, isLastPacket bool) (err error) { return nil } -func (s *WALReplayer) onMetricAdd(rec storage.MetricAddRecord) (err error) { +func (s *WALReplayer) onMetricAdd(rec MetricAddRecord) (err error) { _, ok := s.metrics[rec.MetricID] if ok { return fmt.Errorf("metric %d add failed: already added", rec.MetricID) } var ( - buf = make([]byte, storage.DataPageSize) + buf = make([]byte, DataPageSize) ) s.metrics[rec.MetricID] = &ReplayMetric{ metricType: rec.MetricType, fracDigits: byte(rec.FracDigits), buf: buf, - databuf: buf[:storage.DataPagePayloadSize], + databuf: buf[:DataPagePayloadSize], } return } -func (s *WALReplayer) onMetricDelete(rec storage.MetricDeleteRecord) error { +func (s *WALReplayer) onMetricDelete(rec MetricDeleteRecord) error { _, ok := s.metrics[rec.MetricID] if !ok { return fmt.Errorf("metric %d deletion failed: not found", rec.MetricID) @@ -122,7 +120,7 @@ func (s *WALReplayer) onMetricDelete(rec storage.MetricDeleteRecord) error { return nil } -func (s *WALReplayer) onMeasuresAppend(rec storage.MeasuresAppendRecord) (err error) { +func (s *WALReplayer) onMeasuresAppend(rec MeasuresAppendRecord) (err error) { metric, ok := s.metrics[rec.MetricID] if !ok { return fmt.Errorf("append measures failed: metric %d not found", rec.MetricID) @@ -131,7 +129,7 @@ func (s *WALReplayer) onMeasuresAppend(rec storage.MeasuresAppendRecord) (err er return } -func (s *WALReplayer) onMeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrowRecord, isLastPacket bool) (err error) { +func (s *WALReplayer) onMeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (err error) { metric, ok := s.metrics[rec.MetricID] if !ok { return fmt.Errorf("append measures failed: metric %d not found", rec.MetricID) @@ -148,16 +146,16 @@ func (s *WALReplayer) onMeasuresAppendWithGrow(rec storage.MeasuresAppendWithGro return } -//func (s *WALReplayer) collectPagesToWrite(metric *ReplayMetric, rec storage.WALRecordAppendMeasures) { +//func (s *WALReplayer) collectPagesToWrite(metric *ReplayMetric, rec WALRecordAppendMeasures) { // metric.ApplyHeadDataPage(rec.CompletedDataPage) -// s.dataPages = append(s.dataPages, storage.PageToWrite{ +// s.dataPages = append(s.dataPages, PageToWrite{ // PageNo: p.PageNo, // Content: metric.buf, // }) // for _, x := range rec.DataPages { -// s.dataPages = append(s.dataPages, storage.PageToWrite{ +// s.dataPages = append(s.dataPages, PageToWrite{ // PageNo: x.PageNo, // Content: x.Content, // }) @@ -169,7 +167,7 @@ func (s *WALReplayer) onMeasuresAppendWithGrow(rec storage.MeasuresAppendWithGro // func (s *WALReplayer) replayPacket(records []any, isLastPacket bool) { // for _, untyped := range records { // switch rec := untyped.(type) { -// case storage.WALRecordAppendMeasures: +// case WALRecordAppendMeasures: // s.onWALRecordAppendMeasures(rec, isLastPacket) // } // } diff --git a/storage/wal_writer.go b/storage/wal_writer.go index 62d3d05..07233d1 100644 --- a/storage/wal_writer.go +++ b/storage/wal_writer.go @@ -4,7 +4,6 @@ import ( "bytes" "errors" "fmt" - "log" "math" "os" "path/filepath" @@ -12,6 +11,7 @@ import ( "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/freelist" + "gordenko.dev/dima/qb/inbox" ) const ( @@ -42,10 +42,6 @@ type FreeList interface { GetPageNumber() (uint32, error) } -func JoinChangesFileName(dir string, logNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.changes", logNumber)) -} - type MetricsState struct { //Metrics []Metric WaitCh chan struct{} @@ -63,32 +59,37 @@ type Changes struct { type Writer struct { mutex sync.Mutex + inbox *inbox.Inbox + workerInbox *inbox.Inbox dataPagesCount uint32 indexPagesCount uint32 dataFreeList *freelist.FreeList indexFreeList *freelist.FreeList //atree *atree.Atree - dir string - w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск - wal *os.File - dataFile *os.File - indexFile *os.File - input []any - writePreparer *WritePreparer - appendToWorkerQueue func(any) - written int64 - isExited bool - exitCh chan struct{} - waitGroup *sync.WaitGroup - signalCh chan struct{} + dir string + databaseName string + w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск + wal *os.File + dataFile *os.File + indexFile *os.File + input []any + writePreparer *WritePreparer + written int64 + isExited bool + exitCh chan struct{} + waitGroup *sync.WaitGroup + //signalCh chan struct{} } type WriterOptions struct { - Dir string + Inbox *inbox.Inbox + WorkerInbox *inbox.Inbox + Dir string + DatabaseName string //SnapshotNumber int // номер журнала - AppendToWorkerQueue func(any) - DataFreeList *freelist.FreeList - IndexFreeList *freelist.FreeList + WAL string + DataFreeList *freelist.FreeList + IndexFreeList *freelist.FreeList //Atree *atree.Atree ExitCh chan struct{} WaitGroup *sync.WaitGroup @@ -101,12 +102,21 @@ func NewWriter(opt WriterOptions) (*Writer, error) { // if (opt.DataPageSize % 2) != 0 { // return nil, errors.New("DataPageSize must be multiple of 2") // } + if opt.Inbox == nil { + return nil, errors.New("Inbox option is required") + } + if opt.WorkerInbox == nil { + return nil, errors.New("WorkerInbox option is required") + } if opt.Dir == "" { return nil, errors.New("Dir option is required") } - if opt.AppendToWorkerQueue == nil { - return nil, errors.New("AppendToWorkerQueue option is required") + if opt.DatabaseName == "" { + return nil, errors.New("DatabaseName option is required") } + // if opt.AppendToWorkerQueue == nil { + // return nil, errors.New("AppendToWorkerQueue option is required") + // } if opt.DataFreeList == nil { return nil, errors.New("DataFreeList option is required") } @@ -123,9 +133,21 @@ func NewWriter(opt WriterOptions) (*Writer, error) { return nil, errors.New("WaitGroup option is required") } + indexPageNoProvider, err := NewPageNoProvider(PageNoProviderOptions{ + FilePath: qb.GetIndexFilePath(opt.Dir, opt.DatabaseName), + PageSize: IndexPageSize, + FreeList: opt.IndexFreeList, + }) + + dataPageNoProvider, err := NewPageNoProvider(PageNoProviderOptions{ + FilePath: qb.GetDataFilePath(opt.Dir, opt.DatabaseName), + PageSize: DataPageSize, + FreeList: opt.DataFreeList, + }) + pageIngester, err := NewPageIngester(PageIngesterOptions{ - GetIndexPageNumber: s.getIndexPageNumber, - GetDataPageNumber: s.getDataPageNumber, + GetIndexPageNumber: indexPageNoProvider.GetPageNumber, + GetDataPageNumber: dataPageNoProvider.GetPageNumber, }) if err != nil { return nil, err @@ -137,27 +159,31 @@ func NewWriter(opt WriterOptions) (*Writer, error) { }) s := &Writer{ - dir: opt.Dir, - appendToWorkerQueue: opt.AppendToWorkerQueue, - dataFreeList: opt.DataFreeList, - indexFreeList: opt.IndexFreeList, - writePreparer: writePreparer, - //atree: opt.Atree, - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, - signalCh: make(chan struct{}, 1), + inbox: opt.Inbox, + workerInbox: opt.WorkerInbox, + dir: opt.Dir, + databaseName: opt.DatabaseName, + dataFreeList: opt.DataFreeList, + indexFreeList: opt.IndexFreeList, + writePreparer: writePreparer, + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, } - if opt.SnapshotNumber <= 0 { - log.Fatalln("FIX FUCK NUMBER") - } - s.wal, err = os.OpenFile( - JoinChangesFileName(opt.Dir, opt.SnapshotNumber), - os.O_APPEND|os.O_WRONLY, - filePerm, - ) - if err != nil { - return nil, err + if opt.WAL != "" { + s.wal, err = os.OpenFile(filepath.Join(opt.Dir, opt.WAL), os.O_APPEND|os.O_WRONLY, filePerm) + if err != nil { + return nil, err + } + } else { + s.wal, err = os.OpenFile( + qb.GetWALFilePath(opt.Dir, opt.DatabaseName, 0), // .wal_0 + os.O_CREATE|os.O_APPEND|os.O_WRONLY, + filePerm, + ) + if err != nil { + return nil, err + } } return s, nil } @@ -165,7 +191,7 @@ func NewWriter(opt WriterOptions) (*Writer, error) { func (s *Writer) Run() { for { select { - case <-s.signalCh: + case <-s.inbox.Ready(): if err := s.packAndWrite(); err != nil { qb.Abort(qb.FailedWriteToTxLog, err) } @@ -208,16 +234,16 @@ func (s *Writer) getIndexPageNumber() (uint32, bool, error) { } func (s *Writer) packAndWrite() (err error) { - s.mutex.Lock() + //s.mutex.Lock() //isExited := s.isExited - input := s.input - s.input = nil // var exitWaitGroup *sync.WaitGroup // if s.isExited { // exitWaitGroup = s.waitGroup // } - s.mutex.Unlock() + //s.mutex.Unlock() + + input := s.inbox.Drain() prepared := s.writePreparer.Prepare(input) @@ -234,7 +260,11 @@ func (s *Writer) packAndWrite() (err error) { } // 4. Пишу в atree сторінки - err = s.WritePagesToAtree(prepared.WriteToIndex, prepared.WriteToData) + err = WriteDataPages(s.dataFile, prepared.WriteToData) + if err != nil { + return + } + err = WriteIndexPages(s.indexFile, prepared.WriteToIndex) if err != nil { return } @@ -262,7 +292,7 @@ func (s *Writer) packAndWrite() (err error) { snapshotNumberCh := make(chan int, 1) - s.appendToWorkerQueue(Changes{ + s.workerInbox.Push(Changes{ Commits: prepared.Commits, SnapshotNumberCh: snapshotNumberCh, FrozenIndexPagesCount: s.indexFreeList.Pages(), @@ -285,7 +315,7 @@ func (s *Writer) packAndWrite() (err error) { var err error s.wal, err = os.OpenFile( - JoinChangesFileName(s.dir, snapshotNumber), + qb.GetWALFilePath(s.dir, s.databaseName, snapshotNumber), os.O_CREATE|os.O_WRONLY, filePerm, ) @@ -293,24 +323,13 @@ func (s *Writer) packAndWrite() (err error) { return fmt.Errorf("create new changes file: %s", err) } } else { - s.appendToWorkerQueue(Changes{ + s.workerInbox.Push(Changes{ Commits: prepared.Commits, }) } return nil } -func (s *Writer) Append(req any) { - s.mutex.Lock() - s.input = append(s.input, req) - s.mutex.Unlock() - - select { - case s.signalCh <- struct{}{}: - default: - } -} - // func (s *Writer) reset() { // s.buf.Reset() // s.buf.Write([]byte{ @@ -350,9 +369,8 @@ type PageToWrite struct { Content []byte } -// writePagesToAtree - записує сторінки в .data та .index файли -func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) { - for _, p := range dataPages { +func WriteDataPages(file *os.File, pages []PageToWrite) (err error) { + for _, p := range pages { if len(p.Content) != DataPageSize { return fmt.Errorf("wrong data page size: %d", len(p.Content)) } @@ -360,7 +378,7 @@ func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToW off = int(p.PageNo-1) * DataPageSize n int ) - n, err = s.dataFile.WriteAt(p.Content, int64(off)) + n, err = file.WriteAt(p.Content, int64(off)) if err != nil { return } @@ -368,7 +386,11 @@ func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToW return fmt.Errorf("write %d instead of %d", n, DataPageSize) } } - for _, p := range indexPages { + return +} + +func WriteIndexPages(file *os.File, pages []PageToWrite) (err error) { + for _, p := range pages { if len(p.Content) != IndexPageSize { return fmt.Errorf("wrong index page size: %d", len(p.Content)) } @@ -376,7 +398,7 @@ func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToW off = int(p.PageNo-1) * IndexPageSize n int ) - n, err = s.indexFile.WriteAt(p.Content, int64(off)) + n, err = file.WriteAt(p.Content, int64(off)) if err != nil { return } @@ -384,7 +406,7 @@ func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToW return fmt.Errorf("write %d instead of %d", n, IndexPageSize) } } - return nil + return } // API @@ -401,13 +423,6 @@ func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToW // ValuesBuf []byte // } -func (s *Writer) sendSignal() { - select { - case s.signalCh <- struct{}{}: - default: - } -} - // Якщо (Pages) > 0 - незаповнену сторінку кодую в стисненому вигляді. // Якщо (Pages) = 0 - кодую просто пари timestamp / value // func (s *Writer) packWriteAppended(w *bytes.Buffer, req AppendedPagesReq) { diff --git a/storage/write_preparer.go b/storage/write_preparer.go index cca7020..69a5640 100644 --- a/storage/write_preparer.go +++ b/storage/write_preparer.go @@ -140,13 +140,13 @@ func (s *WritePreparer) Prepare(input []any) PreparedData { if len(level.IndexPages) > 0 { first := level.IndexPages[0] firstCRC32, _ := bin.GetUint32(first.Content[indexCRC32Idx:]) - skipSize := level.SkipRecords * indexRecordSize + skipSize := level.SkipRecords * IndexRecordSize indexPageTail = &IndexPageTail{ PageNo: first.PageNo, Reused: first.Reused, CRC32: firstCRC32, // в WAL файл попадає лише payload - Records: first.Content[skipSize : maxRecordsOnIndexPage*indexRecordSize], + Records: first.Content[skipSize : maxRecordsOnIndexPage*IndexRecordSize], } indexPages = level.IndexPages[1:] } else { @@ -156,7 +156,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData { IndexPageTail: indexPageTail, IndexPages: indexPages, // в WAL файл попадає лише payload - TailRecords: level.TailRecords[:level.TailRecordsCount*indexRecordSize], + TailRecords: level.TailRecords[:level.TailRecordsCount*IndexRecordSize], }) } // Worker-у відправляю повний index @@ -256,7 +256,7 @@ func getIndexRecords(buf []byte, count int) (list []IndexRecord) { rec.Timestamp, _ = bin.GetUint32(buf[i:]) rec.PageNo, _ = bin.GetUint32(buf[i+4:]) list = append(list, rec) - i += indexRecordSize + i += IndexRecordSize } return } diff --git a/database/metric.go b/worker/metric.go similarity index 91% rename from database/metric.go rename to worker/metric.go index 42d0028..a33776b 100644 --- a/database/metric.go +++ b/worker/metric.go @@ -1,6 +1,7 @@ -package database +package worker import ( + "errors" "io" bin "gordenko.dev/dima/bin/little" @@ -10,18 +11,14 @@ import ( // METRIC -//const minBufferSize = 1024 - -var ( - indexRecordSize = 8 -) +var ErrNoValueBug = errors.New("has timestamp but no value") type CapturedState struct { LastTimestamp uint32 LastValue float64 } -type _metric struct { +type Metric struct { metricType qb.MetricType fracDigits byte lastPageNo uint32 @@ -39,21 +36,21 @@ type _metric struct { capturedState *CapturedState } -func (s *_metric) LastValue() float64 { +func (s *Metric) LastValue() float64 { if s.capturedState != nil { return s.capturedState.LastValue } return s.lastValue } -func (s *_metric) LastTimestamp() uint32 { +func (s *Metric) LastTimestamp() uint32 { if s.capturedState != nil { return s.capturedState.LastTimestamp } return s.timestamps.LastTimestamp() } -func (s *_metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { +func (s *Metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { //s.timestamps.Reset() //s.values.Reset() s.XLock = false @@ -70,7 +67,7 @@ func (s *_metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { } // func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) { -func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, saveToStorage func(any)) { +func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, saveToStorage func(any)) { if s.capturedState != nil { s.WaitQueue = append(s.WaitQueue, req) } @@ -163,7 +160,7 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, save if written == 0 { s.capturedState = nil - req.ResultCh <- tryAppendMeasuresResult{ + req.ResultCh <- AppendMeasuresResult{ ResultCode: resultCode, } return @@ -204,14 +201,14 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, save } } -func (s *_metric) OnMeasuresAppendCommited(rec storage.MeasuresAppendCommited) { +func (s *Metric) OnMeasuresAppendCommited(rec storage.MeasuresAppendCommited) { // Видаляю state. Оригінальні Timestamps і Values вже мають останню версію s.capturedState = nil s.timestamps.ForgetCapturedState() s.values.ForgetCapturedState() } -func (s *_metric) OnMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) { +func (s *Metric) OnMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) { // Видаляю state. Оригінальні Timestamps і Values вже мають останню версію s.capturedState = nil s.timestamps.ForgetCapturedState() @@ -227,7 +224,7 @@ func (s *_metric) OnMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWit // READ -func (s *_metric) StartRangeScan(req tryRangeScanReq) { +func (s *Metric) StartRangeScan(req RangeScanReq) { // if s.Since == 0 { // req.ResultCh <- rangeScanResult{ // ResultCode: QueryDone, @@ -295,7 +292,7 @@ func (s *_metric) StartRangeScan(req tryRangeScanReq) { // } } -func (s *_metric) StartFullScan(req tryFullScanReq) { +func (s *Metric) StartFullScan(req FullScanReq) { // if s.Since == 0 { // req.ResultCh <- fullScanResult{ // ResultCode: QueryDone, @@ -319,14 +316,14 @@ func (s *_metric) StartFullScan(req tryFullScanReq) { } if s.lastPageNo > 0 { - req.ResultCh <- fullScanResult{ + req.ResultCh <- FullScanResult{ ResultCode: UntilFound, LastPageNo: s.lastPageNo, FracDigits: s.fracDigits, } s.RLocks++ } else { - req.ResultCh <- fullScanResult{ + req.ResultCh <- FullScanResult{ ResultCode: QueryDone, } } @@ -348,7 +345,7 @@ func (s *_metric) StartFullScan(req tryFullScanReq) { // records - Nb // ] -func (s *_metric) WriteTo(w io.Writer) (err error) { +func (s *Metric) WriteTo(w io.Writer) (err error) { _, err = w.Write([]byte{ byte(s.metricType), s.fracDigits, diff --git a/database/worker.go b/worker/worker.go similarity index 68% rename from database/worker.go rename to worker/worker.go index a29ec2e..3ae5d6a 100644 --- a/database/worker.go +++ b/worker/worker.go @@ -1,4 +1,4 @@ -package database +package worker import ( "fmt" @@ -7,6 +7,7 @@ import ( qb "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/enc" + "gordenko.dev/dima/qb/inbox" "gordenko.dev/dima/qb/proto" "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/transform" @@ -37,62 +38,54 @@ const ( ) type Worker struct { - mutex sync.Mutex - queue []any - signalCh chan struct{} + //mutex sync.Mutex + inbox *inbox.Inbox + storageInbox *inbox.Inbox dir string snapshotNumber int tmp []byte - metrics map[uint32]*_metric + metrics map[uint32]*Metric saveToStorage func(any) exitCh chan struct{} waitGroup *sync.WaitGroup } type WorkerOptions struct { - SaveToStorage func(any) - Metrics map[uint32]*_metric - Dir string - ExitCh chan struct{} - WaitGroup *sync.WaitGroup + Inbox *inbox.Inbox + StorageInbox *inbox.Inbox + Metrics map[uint32]*Metric + Dir string + ExitCh chan struct{} + WaitGroup *sync.WaitGroup } func NewWorker(opt WorkerOptions) *Worker { return &Worker{ - signalCh: make(chan struct{}, 1), - dir: opt.Dir, - tmp: make([]byte, 26), - metrics: opt.Metrics, - saveToStorage: opt.SaveToStorage, - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, - } -} - -func (s *Worker) AddJobToQueue(job any) { - s.queue = append(s.queue, job) - - select { - case s.signalCh <- struct{}{}: - default: + inbox: opt.Inbox, + storageInbox: opt.StorageInbox, + dir: opt.Dir, + tmp: make([]byte, 26), + metrics: opt.Metrics, + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, } } func (s *Worker) ReleaseRLock(metricID uint32) { - s.mutex.Lock() - //s.rLocksToRelease = append(s.rLocksToRelease, metricID) - s.mutex.Unlock() + // s.mutex.Lock() + // //s.rLocksToRelease = append(s.rLocksToRelease, metricID) + // s.mutex.Unlock() - select { - case s.signalCh <- struct{}{}: - default: - } + // select { + // case s.signalCh <- struct{}{}: + // default: + // } } func (s *Worker) Run() { for { select { - case <-s.signalCh: + case <-s.inbox.Ready(): s.doWork() case <-s.exitCh: @@ -104,13 +97,10 @@ func (s *Worker) Run() { } func (s *Worker) doWork() { + queue := s.inbox.Drain() - s.mutex.Lock() //rLocksToRelease := s.rLocksToRelease - queue := s.queue //s.rLocksToRelease = nil - s.queue = nil - s.mutex.Unlock() // for _, metricID := range rLocksToRelease { // metric, ok := s.metrics[metricID] @@ -140,32 +130,32 @@ func (s *Worker) doWork() { for _, untyped := range queue { switch req := untyped.(type) { - case tryAppendMeasuresReq: - s.tryAppendMeasures(req) + case AppendMeasuresReq: + s.AppendMeasures(req) case storage.Changes: s.applyCommits(req) // all metrics only - case tryListCurrentValuesReq: + case ListCurrentValuesReq: s.tryListCurrentValues(req) // all metrics only - case tryRangeScanReq: - s.tryRangeScan(req) + case RangeScanReq: + s.RangeScan(req) - case tryFullScanReq: + case FullScanReq: s.tryFullScan(req) - case tryAddMetricReq: - s.tryAddMetric(req) + case AddMetricReq: + s.AddMetric(req) - case tryDeleteMetricReq: - s.tryDeleteMetric(req) + case DeleteMetricReq: + s.DeleteMetric(req) - case tryDeleteMeasuresReq: - s.tryDeleteMeasures(req) + case DeleteMeasuresReq: + s.DeleteMeasures(req) - case tryGetMetricReq: - s.tryGetMetric(req) + case GetMetricReq: + s.GetMetric(req) default: qb.Abort(qb.UnknownWorkerQueueItemBug, @@ -175,29 +165,29 @@ func (s *Worker) doWork() { } // суть у тому що треба запускати запити, пока не зустріну XLock -func (s *Worker) processMetricQueue(metricID uint32, metric *_metric, tmp []byte) { +func (s *Worker) processMetricQueue(metricID uint32, metric *Metric, tmp []byte) { if len(metric.WaitQueue) == 0 { return } for _, untyped := range metric.WaitQueue { switch req := untyped.(type) { - case tryRangeScanReq: + case RangeScanReq: metric.StartRangeScan(req) - case tryFullScanReq: + case FullScanReq: metric.StartFullScan(req) - case tryGetMetricReq: - s.tryGetMetric(req) + case GetMetricReq: + s.GetMetric(req) - case tryAppendMeasuresReq: + case AppendMeasuresReq: metric.StartAppendMeasures(req, tmp, s.saveToStorage) - case tryDeleteMetricReq: + case DeleteMetricReq: s.startDeleteMetric(metric, req) - case tryDeleteMeasuresReq: + case DeleteMeasuresReq: s.startDeleteMeasures(metric, req) default: @@ -207,12 +197,12 @@ func (s *Worker) processMetricQueue(metricID uint32, metric *_metric, tmp []byte } } -type tryAddMetricReq struct { +type AddMetricReq struct { MetricID uint32 ResultCh chan byte } -func (s *Worker) tryAddMetric(req tryAddMetricReq) { +func (s *Worker) AddMetric(req AddMetricReq) { _, ok := s.metrics[req.MetricID] if ok { req.ResultCh <- MetricDuplicate @@ -232,7 +222,7 @@ func (s *Worker) tryAddMetric(req tryAddMetricReq) { // } } -func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetricReq) { +func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []AddMetricReq) { if len(reqs) == 0 { return } @@ -253,32 +243,43 @@ func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetr req.ResultCh <- Succeed } -type tryGetMetricReq struct { - MetricID uint32 - ResultCh chan Metric +type GetMetricResult struct { + MetricType qb.MetricType + FracDigits byte + ResultCode byte } -func (s *Worker) tryGetMetric(req tryGetMetricReq) { +type GetMetricReq struct { + MetricID uint32 + ResultCh chan GetMetricResult +} + +func (s *Worker) GetMetric(req GetMetricReq) { metric, ok := s.metrics[req.MetricID] if ok { - req.ResultCh <- Metric{ + req.ResultCh <- GetMetricResult{ ResultCode: Succeed, MetricType: metric.metricType, FracDigits: metric.fracDigits, } } else { - req.ResultCh <- Metric{ + req.ResultCh <- GetMetricResult{ ResultCode: NoMetric, } } } -type tryDeleteMetricReq struct { - MetricID uint32 - ResultCh chan tryDeleteMetricResult +type DeleteMetricResult struct { + ResultCode byte + RootPageNo uint32 } -func (s *Worker) tryDeleteMetric(req tryDeleteMetricReq) { +type DeleteMetricReq struct { + MetricID uint32 + ResultCh chan DeleteMetricResult +} + +func (s *Worker) DeleteMetric(req DeleteMetricReq) { // FIX // metric, ok := s.metrics[req.MetricID] // if !ok { @@ -296,7 +297,7 @@ func (s *Worker) tryDeleteMetric(req tryDeleteMetricReq) { // } } -func (s *Worker) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) { +func (s *Worker) startDeleteMetric(metric *Metric, req DeleteMetricReq) { // FIX // s.metricLockEntries[req.MetricID] = &metricLockEntry{ // XLock: true, @@ -307,13 +308,18 @@ func (s *Worker) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) { // } } -type tryDeleteMeasuresReq struct { - MetricID uint32 - Since uint32 - ResultCh chan tryDeleteMeasuresResult +type DeleteMeasuresResult struct { + ResultCode byte + RootPageNo uint32 } -func (s *Worker) tryDeleteMeasures(req tryDeleteMeasuresReq) { +type DeleteMeasuresReq struct { + MetricID uint32 + Since uint32 + ResultCh chan DeleteMeasuresResult +} + +func (s *Worker) DeleteMeasures(req DeleteMeasuresReq) { // FIX // metric, ok := s.metrics[req.MetricID] // if !ok { @@ -337,7 +343,7 @@ func (s *Worker) tryDeleteMeasures(req tryDeleteMeasuresReq) { // } } -func (s *Worker) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq) { +func (s *Worker) startDeleteMeasures(metric *Metric, req DeleteMeasuresReq) { // FIX // s.metricLockEntries[req.MetricID] = &metricLockEntry{ // XLock: true, @@ -375,16 +381,21 @@ func (s *Worker) onMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWith metric.OnMeasuresAppendWithGrowCommited(rec) } -type tryAppendMeasuresReq struct { - MetricID uint32 - Measures []proto.Measure - ResultCh chan tryAppendMeasuresResult +type AppendMeasuresResult struct { + ResultCode byte + Written int } -func (s *Worker) tryAppendMeasures(req tryAppendMeasuresReq) { +type AppendMeasuresReq struct { + MetricID uint32 + Measures []proto.Measure + ResultCh chan AppendMeasuresResult +} + +func (s *Worker) AppendMeasures(req AppendMeasuresReq) { metric, ok := s.metrics[req.MetricID] if !ok { - req.ResultCh <- tryAppendMeasuresResult{ + req.ResultCh <- AppendMeasuresResult{ ResultCode: NoMetric, } return @@ -396,25 +407,32 @@ func (s *Worker) tryAppendMeasures(req tryAppendMeasuresReq) { metric.StartAppendMeasures(req, s.tmp, s.saveToStorage) } -type tryRangeScanReq struct { +type RangeScanResult struct { + ResultCode byte + FracDigits byte + RootPageNo uint32 + LastPageNo uint32 +} + +type RangeScanReq struct { MetricID uint32 Since uint32 Until uint32 MetricType qb.MetricType ResponseWriter atree.WorkerMeasureConsumer - ResultCh chan rangeScanResult + ResultCh chan RangeScanResult } -func (s *Worker) tryRangeScan(req tryRangeScanReq) { +func (s *Worker) RangeScan(req RangeScanReq) { metric, ok := s.metrics[req.MetricID] if !ok { - req.ResultCh <- rangeScanResult{ + req.ResultCh <- RangeScanResult{ ResultCode: NoMetric, } return } if metric.metricType != req.MetricType { - req.ResultCh <- rangeScanResult{ + req.ResultCh <- RangeScanResult{ ResultCode: WrongMetricType, } return @@ -428,23 +446,29 @@ func (s *Worker) tryRangeScan(req tryRangeScanReq) { metric.StartRangeScan(req) } -type tryFullScanReq struct { +type FullScanResult struct { + ResultCode byte + FracDigits byte + LastPageNo uint32 +} + +type FullScanReq struct { MetricID uint32 MetricType qb.MetricType ResponseWriter atree.WorkerMeasureConsumer - ResultCh chan fullScanResult + ResultCh chan FullScanResult } -func (s *Worker) tryFullScan(req tryFullScanReq) { +func (s *Worker) tryFullScan(req FullScanReq) { metric, ok := s.metrics[req.MetricID] if !ok { - req.ResultCh <- fullScanResult{ + req.ResultCh <- FullScanResult{ ResultCode: NoMetric, } return } if metric.metricType != req.MetricType { - req.ResultCh <- fullScanResult{ + req.ResultCh <- FullScanResult{ ResultCode: WrongMetricType, } return @@ -457,13 +481,13 @@ func (s *Worker) tryFullScan(req tryFullScanReq) { metric.StartFullScan(req) } -type tryListCurrentValuesReq struct { +type ListCurrentValuesReq struct { MetricIDs []uint32 ResponseWriter *transform.CurrentValueWriter ResultCh chan struct{} } -func (s *Worker) tryListCurrentValues(req tryListCurrentValuesReq) { +func (s *Worker) tryListCurrentValues(req ListCurrentValuesReq) { for _, metricID := range req.MetricIDs { metric, ok := s.metrics[metricID] if ok { @@ -501,7 +525,7 @@ func (s *Worker) applyCommits(req storage.Changes) { if req.SnapshotNumberCh != nil { s.snapshotNumber++ - err := writeSnapshot(writeSnapshotIn{ + err := storage.WriteSnapshot(storage.WriteSnapshotIn{ SnapshotNumber: s.snapshotNumber, Dir: s.dir, WriteBufferSize: 4 * 1024 * 1024, // 1mb @@ -543,7 +567,7 @@ func (s *Worker) addMetric(rec storage.MetricAddRecord) { buf = make([]byte, storage.DataPageSize) databuf = buf[:storage.DataPagePayloadSize] ) - s.metrics[rec.MetricID] = &_metric{ + s.metrics[rec.MetricID] = &Metric{ metricType: rec.MetricType, fracDigits: rec.FracDigits, buffer: buf, @@ -566,7 +590,7 @@ func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) { rec.MetricID)) } - var addMetricReqs []tryAddMetricReq + var addMetricReqs []AddMetricReq if len(metric.WaitQueue) > 0 { for _, untyped := range metric.WaitQueue { @@ -577,31 +601,31 @@ func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) { // ResultCode: NoMetric, // } - case tryRangeScanReq: - req.ResultCh <- rangeScanResult{ + case RangeScanReq: + req.ResultCh <- RangeScanResult{ ResultCode: NoMetric, } - case tryFullScanReq: - req.ResultCh <- fullScanResult{ + case FullScanReq: + req.ResultCh <- FullScanResult{ ResultCode: NoMetric, } - case tryAddMetricReq: + case AddMetricReq: addMetricReqs = append(addMetricReqs, req) - case tryDeleteMetricReq: - req.ResultCh <- tryDeleteMetricResult{ + case DeleteMetricReq: + req.ResultCh <- DeleteMetricResult{ ResultCode: NoMetric, } - case tryDeleteMeasuresReq: - req.ResultCh <- tryDeleteMeasuresResult{ + case DeleteMeasuresReq: + req.ResultCh <- DeleteMeasuresResult{ ResultCode: NoMetric, } - case tryGetMetricReq: - req.ResultCh <- Metric{ + case GetMetricReq: + req.ResultCh <- GetMetricResult{ ResultCode: NoMetric, } @@ -650,89 +674,3 @@ func (s *Worker) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { // s.processMetricQueue(metricID, metric) // } // } - -func (s *Worker) recovery() { - // FIX - // advisor, err := recovery.NewRecoveryAdvisor(recovery.RecoveryAdvisorOptions{ - // Dir: s.dir, - // VerifySnapshot: s.verifySnapshot, - // }) - // if err != nil { - // panic(err) - // } - - // recipe, err := advisor.GetRecipe() - // if err != nil { - // qb.Abort(qb.GetRecoveryRecipeFailed, err) - // } - - // var logNumber int - - // if recipe != nil { - // if recipe.Snapshot != "" { - // err = s.loadSnapshot(recipe.Snapshot) - // if err != nil { - // qb.Abort(qb.LoadSnapshotFailed, err) - // } - // } - // for _, changesFileName := range recipe.Changes { - // err = s.replayChanges(changesFileName) - // if err != nil { - // qb.Abort(qb.ReplayChangesFailed, err) - // } - // } - // logNumber = recipe.LogNumber - // } - - // s.storage, err = storage.NewWriter(storage.WriterOptions{ - // Dir: s.dir, - // LogNumber: logNumber, - // AppendToWorkerQueue: s.appendJobToWorkerQueue, - // FreeList: s.freeList, - // Atree: s.atree, - // ExitCh: s.exitCh, - // WaitGroup: s.waitGroup, - // }) - // if err != nil { - // qb.Abort(qb.CreateChangesWriterFailed, err) - - // } - // go s.storage.Run() - - // // fileNames, err := s.searchREDOFiles() - // // if err != nil { - // // qb.Abort(qb.SearchREDOFilesFailed, err) - // // } - - // // if len(fileNames) > 0 { - // // for _, fileName := range fileNames { - // // err = s.replayREDOFile(fileName) - // // if err != nil { - // // qb.Abort(qb.ReplayREDOFileFailed, err) - // // } - // // } - - // // for _, fileName := range fileNames { - // // err = os.Remove(fileName) - // // if err != nil { - // // qb.Abort(qb.RemoveREDOFileFailed, err) - // // } - // // } - // // } - - // if recipe != nil { - // if recipe.CompleteSnapshot { - // err = s.dumpSnapshot(logNumber) - // if err != nil { - // qb.Abort(qb.DumpSnapshotFailed, err) - // } - // } - - // for _, fileName := range recipe.ToDelete { - // err = os.Remove(fileName) - // if err != nil { - // qb.Abort(qb.RemoveRecipeFileFailed, err) - // } - // } - // } -}