diff --git a/client/client.go b/client/client.go index 98b477d..8b695da 100644 --- a/client/client.go +++ b/client/client.go @@ -54,6 +54,7 @@ func (s *Connection) mustSuccess(reader *bufreader.BufferedReader) (err error) { } switch code { case proto.RespSuccess: + fmt.Println("SUCCESS") return nil // ok case proto.RespError: return s.onError() @@ -70,6 +71,7 @@ func (s *Connection) AddMetric(req proto.AddMetricReq) (err error) { req.FracDigits, } bin.PutUint32(arr[1:], req.MetricID) + fmt.Println("arr", arr) // req if _, err = s.conn.Write(arr); err != nil { return diff --git a/database/api.go b/database/api.go index c36aa6b..fcb01d7 100644 --- a/database/api.go +++ b/database/api.go @@ -45,6 +45,8 @@ func reply(conn io.Writer, errcode uint16) { func (s *Database) handleTCPConn(conn net.Conn) { defer conn.Close() + fmt.Println("client connected") + r := bufreader.New(conn, 128) for { @@ -59,6 +61,7 @@ func (s *Database) handleTCPConn(conn net.Conn) { } func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (err error) { + fmt.Println("process request") messageType, err := r.ReadByte() if err != nil { if err != io.EOF { @@ -68,6 +71,8 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } } + fmt.Println("messageType:", messageType) + switch messageType { case proto.TypeGetMetric: req, err := proto.ReadGetMetricReq(r) @@ -83,6 +88,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e if err != nil { return fmt.Errorf("proto.ReadAddMetricReq: %s", err) } + fmt.Println("ReadAddMetricReq:", req) reply(conn, s.AddMetric(req)) case proto.TypeDeleteMetric: @@ -154,6 +160,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } case proto.TypeDeleteMeasures: + //fmt.Println("delete metric") req, err := proto.ReadDeleteMeasuresReq(r) if err != nil { return fmt.Errorf("proto.ReadDeleteMeasuresReq: %s", err) @@ -180,6 +187,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } default: + fmt.Printf("unknown messageType: %d\n", messageType) return fmt.Errorf("unknown messageType: %d", messageType) } return @@ -188,6 +196,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e // API func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { + fmt.Println("database.AddMetric") // Валидация if req.MetricID == 0 { return proto.ErrEmptyMetricID @@ -204,26 +213,33 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { resultCh := make(chan byte, 1) - s.appendJobToWorkerQueue(tryAddMetricReq{ + fmt.Println("add job") + + s.worker.AddJobToQueue(tryAddMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) + fmt.Println("job added") + resultCode := <-resultCh switch resultCode { case Succeed: - s.storage.Append(storage.MetricAddRecord{ - MetricID: req.MetricID, - MetricType: req.MetricType, - FracDigits: req.FracDigits, - }) + fmt.Println("OK") + // s.storage.Append(storage.MetricAddRecord{ + // MetricID: req.MetricID, + // MetricType: req.MetricType, + // FracDigits: req.FracDigits, + // }) //<-waitCh case MetricDuplicate: + fmt.Println("ErrDuplicate") return proto.ErrDuplicate default: + fmt.Println("ErrWrongResultCodeBug") qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 @@ -238,7 +254,7 @@ type Metric struct { func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { resultCh := make(chan Metric, 1) - s.appendJobToWorkerQueue(tryGetMetricReq{ + s.worker.AddJobToQueue(tryGetMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) @@ -277,7 +293,7 @@ type tryDeleteMetricResult struct { func (s *Database) DeleteMetric(req proto.DeleteMetricReq) uint16 { resultCh := make(chan tryDeleteMetricResult, 1) - s.appendJobToWorkerQueue(tryDeleteMetricReq{ + s.worker.AddJobToQueue(tryDeleteMetricReq{ MetricID: req.MetricID, ResultCh: resultCh, }) @@ -441,7 +457,7 @@ type tryAppendMeasuresResult struct { func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { resultCh := make(chan tryAppendMeasuresResult, 1) - s.appendJobToWorkerQueue(tryAppendMeasuresReq{ + s.worker.AddJobToQueue(tryAppendMeasuresReq{ MetricID: req.MetricID, Measures: req.Measures, ResultCh: resultCh, @@ -486,7 +502,7 @@ type tryDeleteMeasuresResult struct { func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 { resultCh := make(chan tryDeleteMeasuresResult, 1) - s.appendJobToWorkerQueue(tryDeleteMeasuresReq{ + s.worker.AddJobToQueue(tryDeleteMeasuresReq{ MetricID: req.MetricID, Since: req.Since, ResultCh: resultCh, @@ -667,7 +683,7 @@ type rangeScanReq struct { func (s *Database) rangeScan(req rangeScanReq) error { resultCh := make(chan rangeScanResult, 1) - s.appendJobToWorkerQueue(tryRangeScanReq{ + s.worker.AddJobToQueue(tryRangeScanReq{ MetricID: req.MetricID, Since: req.Since, Until: req.Until, @@ -689,7 +705,7 @@ func (s *Database) rangeScan(req rangeScanReq) error { LastPageNo: result.LastPageNo, Since: req.Since, }) - s.metricRUnlock(req.MetricID) + //s.metricRUnlock(req.MetricID) if err != nil { reply(req.Conn, proto.ErrUnexpected) @@ -705,7 +721,7 @@ func (s *Database) rangeScan(req rangeScanReq) error { Since: req.Since, Until: req.Until, }) - s.metricRUnlock(req.MetricID) + //s.metricRUnlock(req.MetricID) if err != nil { reply(req.Conn, proto.ErrUnexpected) @@ -735,7 +751,7 @@ type fullScanReq struct { func (s *Database) fullScan(req fullScanReq) error { resultCh := make(chan fullScanResult, 1) - s.appendJobToWorkerQueue(tryFullScanReq{ + s.worker.AddJobToQueue(tryFullScanReq{ MetricID: req.MetricID, MetricType: req.MetricType, ResponseWriter: req.ResponseWriter, @@ -754,7 +770,7 @@ func (s *Database) fullScan(req fullScanReq) error { ResponseWriter: req.ResponseWriter, LastPageNo: result.LastPageNo, }) - s.metricRUnlock(req.MetricID) + s.worker.AddJobToQueue(req.MetricID) if err != nil { reply(req.Conn, proto.ErrUnexpected) } else { @@ -779,7 +795,7 @@ func (s *Database) ListCurrentValues(conn net.Conn, req proto.ListCurrentValuesR resultCh := make(chan struct{}) - s.appendJobToWorkerQueue(tryListCurrentValuesReq{ + s.worker.AddJobToQueue(tryListCurrentValuesReq{ MetricIDs: req.MetricIDs, ResponseWriter: responseWriter, ResultCh: resultCh, diff --git a/database/database.go b/database/database.go index db8d10b..382e149 100644 --- a/database/database.go +++ b/database/database.go @@ -10,8 +10,8 @@ import ( "sync" "time" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/freelist" "gordenko.dev/dima/qb/storage" ) @@ -30,24 +30,20 @@ func JoinWALFileName(dir string, snapshotNumber int) string { // } type Database struct { - mutex sync.Mutex - workerSignalCh chan struct{} - workerQueue []any - rLocksToRelease []uint32 - metrics map[uint32]*_metric + mutex sync.Mutex + worker *Worker //metricLockEntries map[uint32]*metricLockEntry - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList - dir string - databaseName string - snapshotNumber int - storage *storage.Writer - atree *atree.Atree - tcpPort int - logfile *os.File - logger *log.Logger - exitCh chan struct{} - waitGroup *sync.WaitGroup + //dataFreeList *freelist.FreeList + //indexFreeList *freelist.FreeList + dir string + databaseName string + storage *storage.Writer + atree *atree.Atree + tcpPort int + logfile *os.File + logger *log.Logger + exitCh chan struct{} + waitGroup *sync.WaitGroup } type Options struct { @@ -92,31 +88,41 @@ func New(opt Options) (_ *Database, err error) { // return nil, err // } - dataFreeList, err := freelist.New(freelist.Options{ - PageSize: 2048, - BaseFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data_base"), - DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data_delta"), - }) + // 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_base"), - DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_index_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{ - workerSignalCh: make(chan struct{}, 1), - dir: opt.Dir, - databaseName: opt.DatabaseName, - metrics: make(map[uint32]*_metric), + worker: worker, + dir: opt.Dir, + databaseName: opt.DatabaseName, //metricLockEntries: make(map[uint32]*metricLockEntry), - dataFreeList: dataFreeList, - indexFreeList: indexFreeList, - tcpPort: opt.TCPPort, - logfile: opt.Logfile, - logger: log.New(opt.Logfile, "", log.LstdFlags), - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, + //dataFreeList: dataFreeList, + //indexFreeList: indexFreeList, + tcpPort: opt.TCPPort, + logfile: opt.Logfile, + logger: log.New(opt.Logfile, "", log.LstdFlags), + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, } return s, nil } @@ -127,18 +133,32 @@ func (s *Database) ListenAndServe() (err error) { return fmt.Errorf("net.Listen: %s; port=%d", err, s.tcpPort) } - s.atree, err = atree.New(atree.Options{ - Dir: s.dir, - DatabaseName: s.databaseName, + 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 { - return fmt.Errorf("atree.New: %s", err) + qb.Abort(qb.CreateChangesWriterFailed, err) + } - s.atree.Run() + go s.storage.Run() - go s.worker() + // s.atree, err = atree.New(atree.Options{ + // Dir: s.dir, + // DatabaseName: s.databaseName, + // }) + // if err != nil { + // return fmt.Errorf("atree.New: %s", err) + // } + // s.atree.Run() - s.recovery() + s.waitGroup.Add(1) + go s.worker.Run() s.logger.Println("database started") for { @@ -153,92 +173,6 @@ func (s *Database) ListenAndServe() (err error) { } } -func (s *Database) 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) - // } - // } - // } -} - // // зробити object? @@ -283,9 +217,9 @@ func (s *Database) replayChanges(snapshotNumber int) (err error) { } func (s *Database) relayMetricsToMetrics(replayMetrics map[uint32]*ReplayMetric) { - for metricID, x := range replayMetrics { - s.metrics[metricID] = x.ToMetric() - } + //for metricID, x := range replayMetrics { + //s.metrics[metricID] = x.ToMetric() FIX + //} } //func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { diff --git a/database/helpers.go b/database/helpers.go index 3e9d503..f0a546b 100644 --- a/database/helpers.go +++ b/database/helpers.go @@ -44,28 +44,6 @@ func isFileExist(fileName string) (bool, error) { } } -func (s *Database) appendJobToWorkerQueue(job any) { - s.mutex.Lock() - s.workerQueue = append(s.workerQueue, job) - s.mutex.Unlock() - - select { - case s.workerSignalCh <- struct{}{}: - default: - } -} - -func (s *Database) metricRUnlock(metricID uint32) { - s.mutex.Lock() - s.rLocksToRelease = append(s.rLocksToRelease, metricID) - s.mutex.Unlock() - - select { - case s.workerSignalCh <- struct{}{}: - default: - } -} - func correctToFHD(since, until uint32, firstHourOfDay int) (uint32, uint32) { duration := time.Duration(firstHourOfDay) * time.Hour since = uint32(time.Unix(int64(since), 0).Add(duration).Unix()) diff --git a/database/metric.go b/database/metric.go index d3b9c4d..42d0028 100644 --- a/database/metric.go +++ b/database/metric.go @@ -70,7 +70,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, sendToStorage func(any)) { +func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, saveToStorage func(any)) { if s.capturedState != nil { s.WaitQueue = append(s.WaitQueue, req) } @@ -174,7 +174,7 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send if len(pages) > 0 { // пишу в storage довгим шляхом через redo файл і запис в data файл - sendToStorage(storage.MeasuresAppendWithGrow{ + saveToStorage(storage.MeasuresAppendWithGrow{ MetricID: req.MetricID, LastPageNo: s.lastPageNo, TimestampsRewindOffset: timestampsRewindOffset, @@ -191,7 +191,7 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send }) } else { // короткий шлях - запис лише в storage - sendToStorage(storage.MeasuresAppend{ + saveToStorage(storage.MeasuresAppend{ MetricID: req.MetricID, TimestampsRewindOffset: timestampsRewindOffset, ValuesRewindOffset: valuesRewindOffset, diff --git a/database/proc.go b/database/worker.go similarity index 66% rename from database/proc.go rename to database/worker.go index bb68369..a29ec2e 100644 --- a/database/proc.go +++ b/database/worker.go @@ -2,6 +2,7 @@ package database import ( "fmt" + "sync" qb "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" @@ -35,50 +36,109 @@ const ( DeleteFromAtreeRequired = 14 ) -func (s *Database) worker() { +type Worker struct { + mutex sync.Mutex + queue []any + signalCh chan struct{} + dir string + snapshotNumber int + tmp []byte + 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 +} + +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: + } +} + +func (s *Worker) ReleaseRLock(metricID uint32) { + s.mutex.Lock() + //s.rLocksToRelease = append(s.rLocksToRelease, metricID) + s.mutex.Unlock() + + select { + case s.signalCh <- struct{}{}: + default: + } +} + +func (s *Worker) Run() { for { select { - case <-s.workerSignalCh: - s.DoWork() + case <-s.signalCh: + s.doWork() + + case <-s.exitCh: + fmt.Println("worker done") + s.waitGroup.Done() + return } } } -func (s *Database) DoWork() { +func (s *Worker) doWork() { + s.mutex.Lock() - rLocksToRelease := s.rLocksToRelease - workerQueue := s.workerQueue - s.rLocksToRelease = nil - s.workerQueue = nil + //rLocksToRelease := s.rLocksToRelease + queue := s.queue + //s.rLocksToRelease = nil + s.queue = nil s.mutex.Unlock() - for _, metricID := range rLocksToRelease { - metric, ok := s.metrics[metricID] - if !ok { - qb.Abort(qb.NoMetricBug, - fmt.Errorf("drainQueues: metric %d not found", metricID)) - } + // for _, metricID := range rLocksToRelease { + // metric, ok := s.metrics[metricID] + // if !ok { + // qb.Abort(qb.NoMetricBug, + // fmt.Errorf("drainQueues: metric %d not found", metricID)) + // } - if metric.XLock { - qb.Abort(qb.XLockBug, - fmt.Errorf("drainQueues: xlock is set for the metric %d", - metricID)) - } + // if metric.XLock { + // qb.Abort(qb.XLockBug, + // fmt.Errorf("drainQueues: xlock is set for the metric %d", + // metricID)) + // } - if metric.RLocks <= 0 { - qb.Abort(qb.NoRLockBug, - fmt.Errorf("drainQueues: rlock not set for the metric %d", - metricID)) - } + // if metric.RLocks <= 0 { + // qb.Abort(qb.NoRLockBug, + // fmt.Errorf("drainQueues: rlock not set for the metric %d", + // metricID)) + // } - metric.RLocks-- + // metric.RLocks-- - if len(metric.WaitQueue) > 0 { - s.processMetricQueue(metricID, metric) - } - } + // if len(metric.WaitQueue) > 0 { + // s.processMetricQueue(metricID, metric) + // } + // } - for _, untyped := range workerQueue { + for _, untyped := range queue { switch req := untyped.(type) { case tryAppendMeasuresReq: s.tryAppendMeasures(req) @@ -115,7 +175,7 @@ func (s *Database) DoWork() { } // суть у тому що треба запускати запити, пока не зустріну XLock -func (s *Database) processMetricQueue(metricID uint32, metric *_metric) { +func (s *Worker) processMetricQueue(metricID uint32, metric *_metric, tmp []byte) { if len(metric.WaitQueue) == 0 { return } @@ -132,7 +192,7 @@ func (s *Database) processMetricQueue(metricID uint32, metric *_metric) { s.tryGetMetric(req) case tryAppendMeasuresReq: - //metric.StartAppendMeasures(req, s.storage.Append) FIX + metric.StartAppendMeasures(req, tmp, s.saveToStorage) case tryDeleteMetricReq: s.startDeleteMetric(metric, req) @@ -152,12 +212,13 @@ type tryAddMetricReq struct { ResultCh chan byte } -func (s *Database) tryAddMetric(req tryAddMetricReq) { +func (s *Worker) tryAddMetric(req tryAddMetricReq) { _, ok := s.metrics[req.MetricID] if ok { req.ResultCh <- MetricDuplicate return } + fmt.Println("metric not found. Success") req.ResultCh <- Succeed // new // lockEntry, ok := s.metricLockEntries[req.MetricID] @@ -171,7 +232,7 @@ func (s *Database) tryAddMetric(req tryAddMetricReq) { // } } -func (s *Database) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetricReq) { +func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetricReq) { if len(reqs) == 0 { return } @@ -197,7 +258,7 @@ type tryGetMetricReq struct { ResultCh chan Metric } -func (s *Database) tryGetMetric(req tryGetMetricReq) { +func (s *Worker) tryGetMetric(req tryGetMetricReq) { metric, ok := s.metrics[req.MetricID] if ok { req.ResultCh <- Metric{ @@ -217,7 +278,7 @@ type tryDeleteMetricReq struct { ResultCh chan tryDeleteMetricResult } -func (s *Database) tryDeleteMetric(req tryDeleteMetricReq) { +func (s *Worker) tryDeleteMetric(req tryDeleteMetricReq) { // FIX // metric, ok := s.metrics[req.MetricID] // if !ok { @@ -235,7 +296,7 @@ func (s *Database) tryDeleteMetric(req tryDeleteMetricReq) { // } } -func (s *Database) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) { +func (s *Worker) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) { // FIX // s.metricLockEntries[req.MetricID] = &metricLockEntry{ // XLock: true, @@ -252,7 +313,7 @@ type tryDeleteMeasuresReq struct { ResultCh chan tryDeleteMeasuresResult } -func (s *Database) tryDeleteMeasures(req tryDeleteMeasuresReq) { +func (s *Worker) tryDeleteMeasures(req tryDeleteMeasuresReq) { // FIX // metric, ok := s.metrics[req.MetricID] // if !ok { @@ -276,7 +337,7 @@ func (s *Database) tryDeleteMeasures(req tryDeleteMeasuresReq) { // } } -func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq) { +func (s *Worker) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq) { // FIX // s.metricLockEntries[req.MetricID] = &metricLockEntry{ // XLock: true, @@ -294,7 +355,7 @@ func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq // } } -func (s *Database) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) { +func (s *Worker) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, @@ -304,7 +365,7 @@ func (s *Database) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) metric.OnMeasuresAppendCommited(rec) } -func (s *Database) onMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) { +func (s *Worker) onMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, @@ -320,7 +381,7 @@ type tryAppendMeasuresReq struct { ResultCh chan tryAppendMeasuresResult } -func (s *Database) tryAppendMeasures(req tryAppendMeasuresReq) { +func (s *Worker) tryAppendMeasures(req tryAppendMeasuresReq) { metric, ok := s.metrics[req.MetricID] if !ok { req.ResultCh <- tryAppendMeasuresResult{ @@ -332,7 +393,7 @@ func (s *Database) tryAppendMeasures(req tryAppendMeasuresReq) { metric.WaitQueue = append(metric.WaitQueue, req) return } - //metric.StartAppendMeasures(req, s.storage.Append) FIX + metric.StartAppendMeasures(req, s.tmp, s.saveToStorage) } type tryRangeScanReq struct { @@ -344,7 +405,7 @@ type tryRangeScanReq struct { ResultCh chan rangeScanResult } -func (s *Database) tryRangeScan(req tryRangeScanReq) { +func (s *Worker) tryRangeScan(req tryRangeScanReq) { metric, ok := s.metrics[req.MetricID] if !ok { req.ResultCh <- rangeScanResult{ @@ -374,7 +435,7 @@ type tryFullScanReq struct { ResultCh chan fullScanResult } -func (s *Database) tryFullScan(req tryFullScanReq) { +func (s *Worker) tryFullScan(req tryFullScanReq) { metric, ok := s.metrics[req.MetricID] if !ok { req.ResultCh <- fullScanResult{ @@ -402,7 +463,7 @@ type tryListCurrentValuesReq struct { ResultCh chan struct{} } -func (s *Database) tryListCurrentValues(req tryListCurrentValuesReq) { +func (s *Worker) tryListCurrentValues(req tryListCurrentValuesReq) { for _, metricID := range req.MetricIDs { metric, ok := s.metrics[metricID] if ok { @@ -418,7 +479,7 @@ func (s *Database) tryListCurrentValues(req tryListCurrentValuesReq) { /////////////////////////////////////////////////////// -func (s *Database) applyCommits(req storage.Changes) { +func (s *Worker) applyCommits(req storage.Changes) { for _, untyped := range req.Commits { switch rec := untyped.(type) { case storage.MetricAddCommited: @@ -458,7 +519,7 @@ func (s *Database) applyCommits(req storage.Changes) { } } -func (s *Database) onMetricAddCommited(rec storage.MetricAddCommited) { +func (s *Worker) onMetricAddCommited(rec storage.MetricAddCommited) { // fix lock _, ok := s.metrics[rec.MetricID] if ok { @@ -477,7 +538,7 @@ func (s *Database) onMetricAddCommited(rec storage.MetricAddCommited) { // delete(s.metricLockEntries, rec.MetricID) } -func (s *Database) addMetric(rec storage.MetricAddRecord) { +func (s *Worker) addMetric(rec storage.MetricAddRecord) { var ( buf = make([]byte, storage.DataPageSize) databuf = buf[:storage.DataPagePayloadSize] @@ -491,7 +552,7 @@ func (s *Database) addMetric(rec storage.MetricAddRecord) { } } -func (s *Database) onMetricDeleteCommited(rec storage.MetricDeleteCommited) { +func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, @@ -562,7 +623,7 @@ func (s *Database) onMetricDeleteCommited(rec storage.MetricDeleteCommited) { } } -func (s *Database) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { +func (s *Worker) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, @@ -581,11 +642,97 @@ func (s *Database) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) // if len(rec.FreePageNumbers) > 0 { // s.freeList.AddPages(rec.FreePageNumbers) // } - s.doAfterReleaseXLock(rec.MetricID, metric) + //s.doAfterReleaseXLock(rec.MetricID, metric) } -func (s *Database) doAfterReleaseXLock(metricID uint32, metric *_metric) { - if len(metric.WaitQueue) > 0 { - s.processMetricQueue(metricID, metric) - } +// func (s *Database) doAfterReleaseXLock(metricID uint32, metric *_metric) { +// if len(metric.WaitQueue) > 0 { +// 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) + // } + // } + // } } diff --git a/database_linux b/database_linux new file mode 100755 index 0000000..c09b0f8 Binary files /dev/null and b/database_linux differ diff --git a/examples/database/main.go b/examples/database/main.go index aaa884d..7fc2b9f 100644 --- a/examples/database/main.go +++ b/examples/database/main.go @@ -58,7 +58,7 @@ func main() { } }() - wg.Add(1) + //wg.Add(1) fmt.Fprintf(logfile, "database %q started on port %d.\n", config.DatabaseName, config.TcpPort) diff --git a/examples/play/main.go b/examples/play/main.go new file mode 100644 index 0000000..418842e --- /dev/null +++ b/examples/play/main.go @@ -0,0 +1,567 @@ +package main + +import ( + "flag" + "fmt" + "log" + "math/rand/v2" + "os" + "time" + + "gopkg.in/ini.v1" + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/client" + "gordenko.dev/dima/qb/proto" +) + +var ( + metricTypeToName = []string{ + "", + "cumulative", + "instant", + } +) + +func main() { + var ( + iniFileName string + ) + + flag.Usage = func() { + fmt.Fprint(flag.CommandLine.Output(), helpMessage) + fmt.Fprint(flag.CommandLine.Output(), configExample) + fmt.Fprintf(flag.CommandLine.Output(), mainUsage, os.Args[0]) + flag.PrintDefaults() + } + flag.StringVar(&iniFileName, "c", "requests.ini", "path to *.ini config file") + flag.Parse() + + config, err := loadConfig(iniFileName) + if err != nil { + log.Fatalln(err) + } + + conn, err := client.Connect(config.DatabaseAddr) + if err != nil { + log.Fatalf("client.Connect(%s): %s\n", config.DatabaseAddr, err) + } else { + fmt.Println("Connected to database") + } + + sendRequests(conn) +} + +// CONFIG FILE + +const mainUsage = `Usage: + %s -c path/to/config.ini + +` + +const helpMessage = `Gordenko project. Example requests. Version: 1.0 +created by Dmytro Gordenko, 1.e4.kc6@gmail.com +` + +const configExample = ` +requests.ini example: + +databaseAddr = :12345 + +` + +type Config struct { + DatabaseAddr string +} + +func (s Config) String() string { + return fmt.Sprintf(`starting options: +databaseAddr = %s +`, + s.DatabaseAddr) +} + +func loadConfig(iniFileName string) (_ Config, err error) { + file, err := ini.Load(iniFileName) + if err != nil { + return + } + + conf := Config{} + top := file.Section("") + + conf.DatabaseAddr = top.Key("databaseAddr").String() + return conf, nil +} + +func sendRequests(conn *client.Connection) { + var ( + //instantMetricID uint32 = 10000 + cumulativeMetricID uint32 = 10001 + fracDigits byte = 2 + err error + ) + + //conn.DeleteMetric(instantMetricID) + //conn.DeleteMetric(cumulativeMetricID) + + // ADD INSTANT METRIC + + err = conn.AddMetric(proto.AddMetricReq{ + MetricID: cumulativeMetricID, + MetricType: qb.Instant, + FracDigits: fracDigits, + }) + if err != nil { + log.Fatalf("conn.AddMetric: %s\n", err) + } else { + fmt.Printf("\nInstant metric %d added\n", cumulativeMetricID) + } + + // GET INSTANT METRIC + + // iMetric, err := conn.GetMetric(cumulativeMetricID) + // if err != nil { + // log.Fatalf("conn.GetMetric: %s\n", err) + // } else { + // fmt.Printf(` + // GetMetric: + // metricID: %d + // metricType: %s + // fracDigits: %d + // `, + // iMetric.MetricID, metricTypeToName[iMetric.MetricType], fracDigits) + //} + + // APPEND MEASURES + + // instantMeasures := GenerateInstantMeasures(62, 220) + + // err = conn.AppendMeasures(proto.AppendMeasuresReq{ + // MetricID: instantMetricID, + // Measures: instantMeasures, + // }) + // if err != nil { + // log.Fatalf("conn.AppendMeasures: %s\n", err) + // } else { + // fmt.Printf("\nAppended %d measures for the metric %d\n", + // len(instantMeasures), instantMetricID) + // } + + // // LIST INSTANT MEASURES + + // lastTimestamp := instantMeasures[len(instantMeasures)-1].Timestamp + // until := time.Unix(int64(lastTimestamp), 0) + // since := until.Add(-5 * time.Hour) + + // instantList, err := conn.ListInstantMeasures(proto.ListInstantMeasuresReq{ + // MetricID: instantMetricID, + // Since: uint32(since.Unix()), + // Until: uint32(until.Unix()), + // }) + // if err != nil { + // log.Fatalf("conn.ListInstantMeasures: %s\n", err) + // } else { + // fmt.Printf("\nListInstantMeasures %s - %s:\n", + // formatTime(uint32(since.Unix())), formatTime(uint32(until.Unix()))) + // for _, item := range instantList { + // fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value) + // } + // } + + // // LIST ALL INSTANT MEASURES + + // instantList, err = conn.ListAllInstantMeasures(instantMetricID) + // if err != nil { + // log.Fatalf("conn.ListAllInstantMeasures: %s\n", err) + // } else { + // fmt.Printf("\nListAllInstantMeasures (last 15 items):\n") + // for _, item := range instantList[:15] { + // fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value) + // } + // } + + // // LIST INSTANT PERIODS (group by hour) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.Add(-24 * time.Hour) + + // instantPeriods, err := conn.ListInstantPeriods(proto.ListInstantPeriodsReq{ + // MetricID: instantMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByHour, + // AggregateFuncs: qb.AggregateMin | qb.AggregateMax | qb.AggregateAvg, + // }) + // if err != nil { + // log.Fatalf("conn.ListInstantPeriods: %s\n", err) + // } else { + // fmt.Printf("\nListInstantPeriods (1 day, group by hour):\n") + // for _, item := range instantPeriods { + // fmt.Printf(" %s => min %.2f, max %.2f, avg %.2f\n", formatHourPeriod(item.Period), item.Min, item.Max, item.Avg) + // } + // } + + // // LIST INSTANT PERIODS (group by day) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.AddDate(0, 0, -7) + + // instantPeriods, err = conn.ListInstantPeriods(proto.ListInstantPeriodsReq{ + // MetricID: instantMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByDay, + // AggregateFuncs: qb.AggregateMin | qb.AggregateMax | qb.AggregateAvg, + // }) + // if err != nil { + // log.Fatalf("conn.ListInstantPeriods: %s\n", err) + // } else { + // fmt.Printf("\nListInstantPeriods (7 days, group by day):\n") + // for _, item := range instantPeriods { + // fmt.Printf(" %s => min %.2f, max %.2f, avg %.2f\n", formatDayPeriod(item.Period), item.Min, item.Max, item.Avg) + // } + // } + + // // LIST INSTANT PERIODS (group by month) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.AddDate(0, 0, -62) + + // instantPeriods, err = conn.ListInstantPeriods(proto.ListInstantPeriodsReq{ + // MetricID: instantMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByMonth, + // AggregateFuncs: qb.AggregateMin | qb.AggregateMax | qb.AggregateAvg, + // }) + // if err != nil { + // log.Fatalf("conn.ListInstantPeriods: %s\n", err) + // } else { + // fmt.Printf("\nListInstantPeriods (62 days, group by month):\n") + // for _, item := range instantPeriods { + // fmt.Printf(" %s => min %.2f, max %.2f, avg %.2f\n", formatMonthPeriod(item.Period), item.Min, item.Max, item.Avg) + // } + // } + + // // DELETE INSTANT METRIC MEASURES + + // err = conn.DeleteMeasures(proto.DeleteMeasuresReq{ + // MetricID: instantMetricID, + // }) + // if err != nil { + // log.Fatalf("conn.DeleteMeasures: %s\n", err) + // } else { + // fmt.Printf("\nInstant metric %d measures deleted\n", instantMetricID) + // } + + // // DELETE INSTANT METRIC + + // err = conn.DeleteMetric(instantMetricID) + // if err != nil { + // log.Fatalf("conn.DeleteMetric: %s\n", err) + // } else { + // fmt.Printf("\nInstant metric %d deleted\n", instantMetricID) + // } + + // // ADD CUMULATIVE METRIC + + // err = conn.AddMetric(proto.AddMetricReq{ + // MetricID: cumulativeMetricID, + // MetricType: qb.Cumulative, + // FracDigits: fracDigits, + // }) + // if err != nil { + // log.Fatalf("conn.AddMetric: %s\n", err) + // } else { + // fmt.Printf("\nCumulative metric %d added\n", cumulativeMetricID) + // } + + // // GET CUMULATIVE METRIC + + // cMetric, err := conn.GetMetric(cumulativeMetricID) + // if err != nil { + // log.Fatalf("conn.GetMetric: %s\n", err) + // } else { + // fmt.Printf(` + // GetMetric: + // metricID: %d + // metricType: %s + // fracDigits: %d + // `, + // cMetric.MetricID, metricTypeToName[cMetric.MetricType], fracDigits) + // } + + // // APPEND MEASURES + + // cumulativeMeasures := GenerateCumulativeMeasures(62) + + // err = conn.AppendMeasures(proto.AppendMeasuresReq{ + // MetricID: cumulativeMetricID, + // Measures: cumulativeMeasures, + // }) + // if err != nil { + // log.Fatalf("conn.AppendMeasures: %s\n", err) + // } else { + // fmt.Printf("\nAppended %d measures for the metric %d\n", + // len(cumulativeMeasures), cumulativeMetricID) + // } + + // // LIST CUMULATIVE MEASURES + + // lastTimestamp = cumulativeMeasures[len(cumulativeMeasures)-1].Timestamp + // until = time.Unix(int64(lastTimestamp), 0) + // since = until.Add(-5 * time.Hour) + + // cumulativeList, err := conn.ListCumulativeMeasures(proto.ListCumulativeMeasuresReq{ + // MetricID: cumulativeMetricID, + // Since: uint32(since.Unix()), + // Until: uint32(until.Unix()), + // }) + // if err != nil { + // log.Fatalf("conn.ListCumulativeMeasures: %s\n", err) + // } else { + // fmt.Printf("\nListCumulativeMeasures %s - %s:\n", + // formatTime(uint32(since.Unix())), formatTime(uint32(until.Unix()))) + + // for _, item := range cumulativeList { + // fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value) + // } + // } + + // // LIST ALL CUMULATIVE MEASURES + + // cumulativeList, err = conn.ListAllCumulativeMeasures(cumulativeMetricID) + // if err != nil { + // log.Fatalf("conn.ListAllCumulativeMeasures: %s\n", err) + // } else { + // fmt.Printf("\nListAllCumulativeMeasures (last 15 items):\n") + // for _, item := range cumulativeList[:15] { + // fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value) + // } + // } + + // // LIST CUMULATIVE PERIODS (group by hour) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.Add(-24 * time.Hour) + + // cumulativePeriods, err := conn.ListCumulativePeriods(proto.ListCumulativePeriodsReq{ + // MetricID: cumulativeMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByHour, + // }) + // if err != nil { + // log.Fatalf("conn.ListCumulativePeriods: %s\n", err) + // } else { + // fmt.Printf("\nListCumulativePeriods (1 day, group by hour):\n") + // for _, item := range cumulativePeriods { + // fmt.Printf(" %s => end value %.2f, total %.2f\n", formatHourPeriod(item.Period), item.EndValue, item.Total) + // } + // } + + // // LIST CUMULATIVE PERIODS (group by day) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.AddDate(0, 0, -7) + + // cumulativePeriods, err = conn.ListCumulativePeriods(proto.ListCumulativePeriodsReq{ + // MetricID: cumulativeMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByDay, + // }) + // if err != nil { + // log.Fatalf("conn.ListCumulativePeriods: %s\n", err) + // } else { + // fmt.Printf("\nListCumulativePeriods (7 days, group by day):\n") + // for _, item := range cumulativePeriods { + // fmt.Printf(" %s => end value %.2f, total %.2f\n", formatDayPeriod(item.Period), item.EndValue, item.Total) + // } + // } + + // // LIST CUMULATIVE PERIODS (group by day) + + // until = time.Unix(int64(lastTimestamp+1), 0) + // since = until.AddDate(0, 0, -62) + + // cumulativePeriods, err = conn.ListCumulativePeriods(proto.ListCumulativePeriodsReq{ + // MetricID: cumulativeMetricID, + // Since: proto.TimeBound{ + // Year: since.Year(), + // Month: since.Month(), + // Day: since.Day(), + // }, + // Until: proto.TimeBound{ + // Year: until.Year(), + // Month: until.Month(), + // Day: until.Day(), + // }, + // GroupBy: qb.GroupByMonth, + // }) + // if err != nil { + // log.Fatalf("conn.ListCumulativePeriods: %s\n", err) + // } else { + // fmt.Printf("\nListCumulativePeriods (62 days, group by month):\n") + // for _, item := range cumulativePeriods { + // fmt.Printf(" %s => end value %.2f, total %.2f\n", formatMonthPeriod(item.Period), item.EndValue, item.Total) + // } + // } + + // // DELETE CUMULATIVE METRIC MEASURES + + // err = conn.DeleteMeasures(proto.DeleteMeasuresReq{ + // MetricID: cumulativeMetricID, + // }) + // if err != nil { + // log.Fatalf("conn.DeleteMeasures: %s\n", err) + // } else { + // fmt.Printf("\nCumulative metric %d measures deleted\n", cumulativeMetricID) + // } + + // // DELETE CUMULATIVE METRIC + + // err = conn.DeleteMetric(cumulativeMetricID) + // + // if err != nil { + // log.Fatalf("conn.DeleteMetric: %s\n", err) + // } else { + // + // fmt.Printf("\nCumulative metric %d deleted\n", cumulativeMetricID) + // } +} + +const datetimeLayout = "2006-01-02 15:04:05" + +func formatTime(timestamp uint32) string { + tm := time.Unix(int64(timestamp), 0) + return tm.Format(datetimeLayout) +} + +func formatHourPeriod(period uint32) string { + tm := time.Unix(int64(period), 0) + return tm.Format("2006-01-02 15:00 - 15") + ":59" +} + +func formatDayPeriod(period uint32) string { + tm := time.Unix(int64(period), 0) + return tm.Format("2006-01-02") +} + +func formatMonthPeriod(period uint32) string { + tm := time.Unix(int64(period), 0) + return tm.Format("2006-01") +} + +func GenerateCumulativeMeasures(days int) []proto.Measure { + var ( + measures []proto.Measure + minutes = []int{14, 29, 44, 59} + hoursPerDay = 24 + totalHours = days * hoursPerDay + since = time.Now().AddDate(0, 0, -days) + totalValue float64 + ) + + for i := range totalHours { + hourTime := since.Add(time.Duration(i) * time.Hour) + for _, m := range minutes { + measureTime := time.Date( + hourTime.Year(), + hourTime.Month(), + hourTime.Day(), + hourTime.Hour(), + m, // minutes + 0, // seconds + 0, // nanoseconds + time.Local, + ) + + measure := proto.Measure{ + Timestamp: uint32(measureTime.Unix()), + Value: totalValue, + } + measures = append(measures, measure) + + totalValue += rand.Float64() + } + } + return measures +} + +func GenerateInstantMeasures(days int, baseValue float64) []proto.Measure { + var ( + measures []proto.Measure + minutes = []int{14, 29, 44, 59} + hoursPerDay = 24 + totalHours = days * hoursPerDay + since = time.Now().AddDate(0, 0, -days) + ) + + for i := range totalHours { + hourTime := since.Add(time.Duration(i) * time.Hour) + for _, m := range minutes { + measureTime := time.Date( + hourTime.Year(), + hourTime.Month(), + hourTime.Day(), + hourTime.Hour(), + m, // minutes + 0, // seconds + 0, // nanoseconds + time.Local, + ) + + // value = +-10% from base value + fluctuation := baseValue * 0.1 + value := baseValue + (rand.Float64()*2-1)*fluctuation + + measure := proto.Measure{ + Timestamp: uint32(measureTime.Unix()), + Value: value, + } + measures = append(measures, measure) + } + } + return measures +} diff --git a/examples/play/play b/examples/play/play new file mode 100755 index 0000000..b4b8dcf Binary files /dev/null and b/examples/play/play differ diff --git a/examples/play/requests.ini b/examples/play/requests.ini new file mode 100644 index 0000000..9c314ac --- /dev/null +++ b/examples/play/requests.ini @@ -0,0 +1 @@ +databaseAddr = :12345 diff --git a/linux_build.sh b/linux_build.sh index 0747f5d..6b64115 100755 --- a/linux_build.sh +++ b/linux_build.sh @@ -2,10 +2,10 @@ cd examples/database env CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o ../../database_linux cd - -cd examples/loadtest -env CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o ../../loadtest_linux -cd - +# cd examples/loadtest +# env CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o ../../loadtest_linux +# cd - -cd examples/requests -env CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o ../../requests_linux -cd - \ No newline at end of file +# cd examples/requests +# env CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o ../../requests_linux +# cd - \ No newline at end of file diff --git a/recovery/advisor.go b/recovery/advisor.go index af08d1e..8013f44 100644 --- a/recovery/advisor.go +++ b/recovery/advisor.go @@ -11,18 +11,18 @@ import ( ) var ( - reChanges = regexp.MustCompile(`(\d+)\.changes`) + reChanges = regexp.MustCompile(`(\d+)\.wal`) reSnapshot = regexp.MustCompile(`(\d+)\.snapshot`) ) -func joinChangesFileName(dir string, logNumber int) string { - return filepath.Join(dir, fmt.Sprintf("%d.changes", logNumber)) +func joinWALFileName(dir string, snapshotNumber int) string { + return filepath.Join(dir, fmt.Sprintf("%d.wal", snapshotNumber)) } type RecoveryRecipe struct { Snapshot string Changes []string - LogNumber int + SnapshotNumber int ToDelete []string CompleteSnapshot bool // флаг - что нужно завершить создание снапшота } @@ -52,8 +52,8 @@ func NewRecoveryAdvisor(opt RecoveryAdvisorOptions) (*RecoveryAdvisor, error) { type SnapshotChangesPair struct { SnapshotFileName string - ChangesFileName string - LogNumber int + WALFileName string + SnapshotNumber int } func (s *RecoveryAdvisor) getSnapshotChangesPairs() (*SnapshotChangesPair, *SnapshotChangesPair, error) { @@ -89,21 +89,21 @@ func (s *RecoveryAdvisor) getSnapshotChangesPairs() (*SnapshotChangesPair, *Snap } } - for logNumber := range numSet { + for snapshotNumber := range numSet { var ( snapshotFileName string changesFileName string ) - if changesSet[logNumber] { - changesFileName = joinChangesFileName(s.dir, logNumber) + if changesSet[snapshotNumber] { + changesFileName = joinWALFileName(s.dir, snapshotNumber) } - if snapshotsSet[logNumber] { - snapshotFileName = filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", logNumber)) + if snapshotsSet[snapshotNumber] { + snapshotFileName = filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", snapshotNumber)) } pairs = append(pairs, SnapshotChangesPair{ - ChangesFileName: changesFileName, + WALFileName: changesFileName, SnapshotFileName: snapshotFileName, - LogNumber: logNumber, + SnapshotNumber: snapshotNumber, }) } @@ -112,22 +112,22 @@ func (s *RecoveryAdvisor) getSnapshotChangesPairs() (*SnapshotChangesPair, *Snap } sort.Slice(pairs, func(i, j int) bool { - return pairs[i].LogNumber > pairs[j].LogNumber + return pairs[i].SnapshotNumber > pairs[j].SnapshotNumber }) pair := pairs[0] - if pair.ChangesFileName == "" { + if pair.WALFileName == "" { return nil, nil, fmt.Errorf("has %d.shapshot file, but %d.changes file not found", - pair.LogNumber, pair.LogNumber) + pair.SnapshotNumber, pair.SnapshotNumber) } if len(pairs) > 1 { prevPair := pairs[1] - if prevPair.SnapshotFileName == "" && prevPair.LogNumber != 1 { + if prevPair.SnapshotFileName == "" && prevPair.SnapshotNumber != 1 { return &pair, nil, nil } - if prevPair.ChangesFileName == "" && pair.SnapshotFileName == "" { + if prevPair.WALFileName == "" && pair.SnapshotFileName == "" { return &pair, nil, nil } return &pair, &prevPair, nil @@ -157,14 +157,14 @@ func (s *RecoveryAdvisor) GetRecipe() (*RecoveryRecipe, error) { recipe := &RecoveryRecipe{ Snapshot: pair.SnapshotFileName, Changes: []string{ - pair.ChangesFileName, + pair.WALFileName, }, - LogNumber: pair.LogNumber, + SnapshotNumber: pair.SnapshotNumber, } if prevPair != nil { - if prevPair.ChangesFileName != "" { - recipe.ToDelete = append(recipe.ToDelete, prevPair.ChangesFileName) + if prevPair.WALFileName != "" { + recipe.ToDelete = append(recipe.ToDelete, prevPair.WALFileName) } if prevPair.SnapshotFileName != "" { recipe.ToDelete = append(recipe.ToDelete, prevPair.SnapshotFileName) @@ -175,55 +175,55 @@ func (s *RecoveryAdvisor) GetRecipe() (*RecoveryRecipe, error) { if prevPair != nil { return s.tryPrevPair(pair, prevPair) } - return nil, fmt.Errorf("%d.shapshot is corrupted", pair.LogNumber) + return nil, fmt.Errorf("%d.shapshot is corrupted", pair.SnapshotNumber) } else { if prevPair != nil { return s.tryPrevPair(pair, prevPair) } else { - if pair.LogNumber == 1 { + if pair.SnapshotNumber == 1 { return &RecoveryRecipe{ Changes: []string{ - pair.ChangesFileName, + pair.WALFileName, }, - LogNumber: pair.LogNumber, + SnapshotNumber: pair.SnapshotNumber, }, nil } else { - return nil, fmt.Errorf("%d.snapshot not found", pair.LogNumber) + return nil, fmt.Errorf("%d.snapshot not found", pair.SnapshotNumber) } } } } func (s *RecoveryAdvisor) tryPrevPair(pair, prevPair *SnapshotChangesPair) (*RecoveryRecipe, error) { - if prevPair.ChangesFileName == "" { + if prevPair.WALFileName == "" { if pair.SnapshotFileName != "" { return nil, fmt.Errorf("%d.shapshot is corrupted and %d.changes not found", - pair.LogNumber, prevPair.LogNumber) + pair.SnapshotNumber, prevPair.SnapshotNumber) } else { - return nil, fmt.Errorf("%d.changes not found", prevPair.LogNumber) + return nil, fmt.Errorf("%d.changes not found", prevPair.SnapshotNumber) } } if prevPair.SnapshotFileName == "" { - if prevPair.LogNumber == 1 { + if prevPair.SnapshotNumber == 1 { recipe := &RecoveryRecipe{ Changes: []string{ - prevPair.ChangesFileName, - pair.ChangesFileName, + prevPair.WALFileName, + pair.WALFileName, }, - LogNumber: pair.LogNumber, + SnapshotNumber: pair.SnapshotNumber, CompleteSnapshot: true, ToDelete: []string{ - prevPair.ChangesFileName, + prevPair.WALFileName, }, } return recipe, nil } else { if pair.SnapshotFileName != "" { return nil, fmt.Errorf("%d.shapshot is corrupted and %d.snapshot not found", - pair.LogNumber, prevPair.LogNumber) + pair.SnapshotNumber, prevPair.SnapshotNumber) } else { - return nil, fmt.Errorf("%d.snapshot not found", pair.LogNumber) + return nil, fmt.Errorf("%d.snapshot not found", pair.SnapshotNumber) } } } @@ -235,19 +235,19 @@ func (s *RecoveryAdvisor) tryPrevPair(pair, prevPair *SnapshotChangesPair) (*Rec } if !isVerified { - return nil, fmt.Errorf("%d.shapshot is corrupted", prevPair.LogNumber) + return nil, fmt.Errorf("%d.shapshot is corrupted", prevPair.SnapshotNumber) } recipe := &RecoveryRecipe{ Snapshot: prevPair.SnapshotFileName, Changes: []string{ - prevPair.ChangesFileName, - pair.ChangesFileName, + prevPair.WALFileName, + pair.WALFileName, }, - LogNumber: pair.LogNumber, + SnapshotNumber: pair.SnapshotNumber, CompleteSnapshot: true, ToDelete: []string{ - prevPair.ChangesFileName, + prevPair.WALFileName, prevPair.SnapshotFileName, }, } diff --git a/storage/wal_writer.go b/storage/wal_writer.go index 5ff0712..62d3d05 100644 --- a/storage/wal_writer.go +++ b/storage/wal_writer.go @@ -11,7 +11,6 @@ import ( "sync" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/freelist" ) @@ -63,21 +62,19 @@ type Changes struct { } type Writer struct { - mutex sync.Mutex - dataPagesCount uint32 - indexPagesCount uint32 - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList - atree *atree.Atree + mutex sync.Mutex + 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 - indexPageSize int - dataPageSize int input []any - dataPreparer *WritePreparer + writePreparer *WritePreparer appendToWorkerQueue func(any) written int64 isExited bool @@ -87,25 +84,23 @@ type Writer struct { } type WriterOptions struct { - IndexPageSize int - DataPageSize int - Dir string - SnapshotNumber int // номер журнала + Dir string + //SnapshotNumber int // номер журнала AppendToWorkerQueue func(any) DataFreeList *freelist.FreeList IndexFreeList *freelist.FreeList - Atree *atree.Atree - ExitCh chan struct{} - WaitGroup *sync.WaitGroup + //Atree *atree.Atree + ExitCh chan struct{} + WaitGroup *sync.WaitGroup } func NewWriter(opt WriterOptions) (*Writer, error) { - if (opt.IndexPageSize % 2) != 0 { - return nil, errors.New("IndexPageSize must be multiple of 2") - } - if (opt.DataPageSize % 2) != 0 { - return nil, errors.New("DataPageSize must be multiple of 2") - } + // if (opt.IndexPageSize % 2) != 0 { + // return nil, errors.New("IndexPageSize must be multiple of 2") + // } + // if (opt.DataPageSize % 2) != 0 { + // return nil, errors.New("DataPageSize must be multiple of 2") + // } if opt.Dir == "" { return nil, errors.New("Dir option is required") } @@ -118,9 +113,9 @@ func NewWriter(opt WriterOptions) (*Writer, error) { if opt.IndexFreeList == nil { return nil, errors.New("IndexFreeList option is required") } - if opt.Atree == nil { - return nil, errors.New("Atree option is required") - } + // if opt.Atree == nil { + // return nil, errors.New("Atree option is required") + // } if opt.ExitCh == nil { return nil, errors.New("ExitCh option is required") } @@ -128,21 +123,31 @@ func NewWriter(opt WriterOptions) (*Writer, error) { return nil, errors.New("WaitGroup option is required") } + pageIngester, err := NewPageIngester(PageIngesterOptions{ + GetIndexPageNumber: s.getIndexPageNumber, + GetDataPageNumber: s.getDataPageNumber, + }) + if err != nil { + return nil, err + } + + writePreparer := NewWritePreparer(WritePreparerOptions{ + //MinWALBufferSize: 128, + PageIngester: pageIngester, + }) + s := &Writer{ - indexPageSize: opt.IndexPageSize, - dataPageSize: opt.DataPageSize, dir: opt.Dir, appendToWorkerQueue: opt.AppendToWorkerQueue, dataFreeList: opt.DataFreeList, indexFreeList: opt.IndexFreeList, - atree: opt.Atree, - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, - signalCh: make(chan struct{}, 1), + writePreparer: writePreparer, + //atree: opt.Atree, + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, + signalCh: make(chan struct{}, 1), } - var err error - if opt.SnapshotNumber <= 0 { log.Fatalln("FIX FUCK NUMBER") } @@ -214,7 +219,7 @@ func (s *Writer) packAndWrite() (err error) { // } s.mutex.Unlock() - prepared := s.dataPreparer.Prepare(input) + prepared := s.writePreparer.Prepare(input) // 3. Пишу на диск WAL (append) n, err := s.wal.Write(prepared.Packet) @@ -348,35 +353,35 @@ type PageToWrite struct { // writePagesToAtree - записує сторінки в .data та .index файли func (s *Writer) WritePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) { for _, p := range dataPages { - if len(p.Content) != s.dataPageSize { + if len(p.Content) != DataPageSize { return fmt.Errorf("wrong data page size: %d", len(p.Content)) } var ( - off = int(p.PageNo-1) * s.dataPageSize + off = int(p.PageNo-1) * DataPageSize n int ) n, err = s.dataFile.WriteAt(p.Content, int64(off)) if err != nil { return } - if n != s.dataPageSize { - return fmt.Errorf("write %d instead of %d", n, s.dataPageSize) + if n != DataPageSize { + return fmt.Errorf("write %d instead of %d", n, DataPageSize) } } for _, p := range indexPages { - if len(p.Content) != s.indexPageSize { + if len(p.Content) != IndexPageSize { return fmt.Errorf("wrong index page size: %d", len(p.Content)) } var ( - off = int(p.PageNo-1) * s.indexPageSize + off = int(p.PageNo-1) * IndexPageSize n int ) n, err = s.indexFile.WriteAt(p.Content, int64(off)) if err != nil { return } - if n != s.indexPageSize { - return fmt.Errorf("write %d instead of %d", n, s.indexPageSize) + if n != IndexPageSize { + return fmt.Errorf("write %d instead of %d", n, IndexPageSize) } } return nil diff --git a/testdir/test.free_data b/testdir/test.free_data new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.free_data_delta b/testdir/test.free_data_delta new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.free_index b/testdir/test.free_index new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.free_index_delta b/testdir/test.free_index_delta new file mode 100644 index 0000000..e69de29