wp
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
}
|
||||
Reference in New Issue
Block a user