wp
This commit is contained in:
293
database/api.go
293
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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
// }
|
||||
|
||||
@@ -1,397 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
bin "gordenko.dev/dima/bin/little"
|
||||
"gordenko.dev/dima/qb"
|
||||
"gordenko.dev/dima/qb/storage"
|
||||
)
|
||||
|
||||
// METRIC
|
||||
|
||||
//const minBufferSize = 1024
|
||||
|
||||
var (
|
||||
indexRecordSize = 8
|
||||
)
|
||||
|
||||
type CapturedState struct {
|
||||
LastTimestamp uint32
|
||||
LastValue float64
|
||||
}
|
||||
|
||||
type _metric struct {
|
||||
metricType qb.MetricType
|
||||
fracDigits byte
|
||||
lastPageNo uint32
|
||||
//SinceValue float64
|
||||
//Since uint32
|
||||
lastValue float64
|
||||
//Until uint32
|
||||
buffer []byte
|
||||
timestamps qb.TimestampCompressor
|
||||
values qb.ValueCompressor
|
||||
XLock bool
|
||||
RLocks int
|
||||
WaitQueue []any
|
||||
indexLevelTails []storage.IndexLevelTail // root - last element
|
||||
capturedState *CapturedState
|
||||
}
|
||||
|
||||
func (s *_metric) LastValue() float64 {
|
||||
if s.capturedState != nil {
|
||||
return s.capturedState.LastValue
|
||||
}
|
||||
return s.lastValue
|
||||
}
|
||||
|
||||
func (s *_metric) LastTimestamp() uint32 {
|
||||
if s.capturedState != nil {
|
||||
return s.capturedState.LastTimestamp
|
||||
}
|
||||
return s.timestamps.LastTimestamp()
|
||||
}
|
||||
|
||||
func (s *_metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
|
||||
//s.timestamps.Reset()
|
||||
//s.values.Reset()
|
||||
s.XLock = false
|
||||
// s.Timestamps.Renew()
|
||||
// s.Values.Renew()
|
||||
|
||||
// s.LastPageNo = 0
|
||||
// s.Since = 0
|
||||
// s.SinceValue = 0
|
||||
// s.Until = 0
|
||||
s.indexLevelTails = nil
|
||||
s.lastPageNo = 0
|
||||
s.lastValue = 0
|
||||
}
|
||||
|
||||
// func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) {
|
||||
func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, saveToStorage func(any)) {
|
||||
if s.capturedState != nil {
|
||||
s.WaitQueue = append(s.WaitQueue, req)
|
||||
}
|
||||
var (
|
||||
timestamps = s.timestamps
|
||||
values = s.values
|
||||
|
||||
// офсети на head сторінці
|
||||
timestampsOffset int
|
||||
timestampsRewindOffset int
|
||||
valuesOffset int
|
||||
valuesRewindOffset int
|
||||
|
||||
headTimestamps []byte
|
||||
headValues []byte
|
||||
|
||||
pages []storage.DataPayload
|
||||
|
||||
written int
|
||||
resultCode byte
|
||||
)
|
||||
|
||||
s.capturedState = &CapturedState{
|
||||
LastTimestamp: timestamps.LastTimestamp(),
|
||||
LastValue: s.lastValue,
|
||||
}
|
||||
|
||||
s.timestamps.CaptureState()
|
||||
s.values.CaptureState()
|
||||
|
||||
for idx, measure := range req.Measures {
|
||||
if measure.Timestamp <= s.timestamps.LastTimestamp() {
|
||||
resultCode = ExpiredMeasure
|
||||
break
|
||||
}
|
||||
if s.metricType == qb.Cumulative && measure.Value < s.lastValue {
|
||||
resultCode = NonMonotonicValue
|
||||
break
|
||||
}
|
||||
|
||||
tReport := timestamps.Evaluate(tmp, measure.Timestamp)
|
||||
vReport := values.Evaluate(tmp, measure.Value)
|
||||
|
||||
totalRequiredSpace := tReport.TotalSpace + vReport.TotalSpace
|
||||
|
||||
if totalRequiredSpace <= len(s.buffer) {
|
||||
// якщо на сторінці є місце
|
||||
timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp)
|
||||
values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta)
|
||||
|
||||
if idx == 0 {
|
||||
timestampsOffset = tReport.Offset
|
||||
timestampsRewindOffset = tReport.RewindOffset
|
||||
valuesOffset = vReport.Offset
|
||||
valuesRewindOffset = vReport.RewindOffset
|
||||
}
|
||||
} else {
|
||||
// сторінка заповнена
|
||||
since := s.timestamps.ReplaceSinceWithUntil()
|
||||
|
||||
if len(pages) == 0 && idx > 0 {
|
||||
// idx > 0 required because page may overflows without append any data
|
||||
headTimestamps = timestamps.Tail(timestampsOffset)
|
||||
headValues = values.Tail(valuesOffset)
|
||||
}
|
||||
|
||||
pages = append(pages, storage.DataPayload{
|
||||
Since: since,
|
||||
Content: s.buffer,
|
||||
TimestampsSize: timestamps.Size(),
|
||||
ValuesSize: values.Size(),
|
||||
})
|
||||
|
||||
buf := make([]byte, storage.DataPageSize)
|
||||
databuf := buf[:storage.DataPagePayloadSize]
|
||||
|
||||
timestamps.ReplaceBuffer(databuf)
|
||||
values.ReplaceBuffer(databuf)
|
||||
|
||||
// renew
|
||||
s.buffer = buf
|
||||
|
||||
timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp)
|
||||
values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta)
|
||||
}
|
||||
//
|
||||
s.lastValue = measure.Value
|
||||
written++
|
||||
}
|
||||
|
||||
if written == 0 {
|
||||
s.capturedState = nil
|
||||
req.ResultCh <- tryAppendMeasuresResult{
|
||||
ResultCode: resultCode,
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// виділити змінені байти.
|
||||
// скопіювати. Причому можна скопіювати зрізи chunks
|
||||
|
||||
if len(pages) > 0 {
|
||||
// пишу в storage довгим шляхом через redo файл і запис в data файл
|
||||
saveToStorage(storage.MeasuresAppendWithGrow{
|
||||
MetricID: req.MetricID,
|
||||
LastPageNo: s.lastPageNo,
|
||||
TimestampsRewindOffset: timestampsRewindOffset,
|
||||
Timestamps: headTimestamps,
|
||||
ValuesRewindOffset: valuesRewindOffset,
|
||||
Values: headValues,
|
||||
IndexLevelTails: s.indexLevelTails,
|
||||
DataPages: pages,
|
||||
TailTimestamps: timestamps.Tail(0), // payload
|
||||
TailValues: values.Tail(0),
|
||||
ResultCode: resultCode,
|
||||
WrittenCount: written,
|
||||
//ResultCh: req.ResultCh,
|
||||
})
|
||||
} else {
|
||||
// короткий шлях - запис лише в storage
|
||||
saveToStorage(storage.MeasuresAppend{
|
||||
MetricID: req.MetricID,
|
||||
TimestampsRewindOffset: timestampsRewindOffset,
|
||||
ValuesRewindOffset: valuesRewindOffset,
|
||||
Timestamps: timestamps.Tail(timestampsOffset),
|
||||
Values: values.Tail(valuesOffset),
|
||||
ResultCode: resultCode,
|
||||
WrittenCount: written,
|
||||
//ResultCh: req.ResultCh,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
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) {
|
||||
// Видаляю state. Оригінальні Timestamps і Values вже мають останню версію
|
||||
s.capturedState = nil
|
||||
s.timestamps.ForgetCapturedState()
|
||||
s.values.ForgetCapturedState()
|
||||
|
||||
if rec.LastPageNo > 0 {
|
||||
s.lastPageNo = rec.LastPageNo
|
||||
}
|
||||
// В storage я передав повний індекс. У нього додали елементи (можливо нові рівні).
|
||||
// Тому проста заміна
|
||||
s.indexLevelTails = rec.Index
|
||||
}
|
||||
|
||||
// READ
|
||||
|
||||
func (s *_metric) StartRangeScan(req tryRangeScanReq) {
|
||||
// if s.Since == 0 {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
|
||||
// if req.Since > s.Until {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
|
||||
// if req.Until < s.Since {
|
||||
// if s.RootPageNo > 0 {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: UntilNotFound,
|
||||
// RootPageNo: s.RootPageNo,
|
||||
// FracDigits: s.FracDigits,
|
||||
// }
|
||||
// s.RLocks++
|
||||
// return
|
||||
// } else {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
// }
|
||||
|
||||
// timestampDecompressor := s.timestamps.CreateDecompressor()
|
||||
// valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits)
|
||||
|
||||
// for {
|
||||
// timestamp, done := timestampDecompressor.NextValue()
|
||||
// if done {
|
||||
// break
|
||||
// }
|
||||
// value, done := valueDecompressor.NextValue()
|
||||
// if done {
|
||||
// qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug)
|
||||
// }
|
||||
// if timestamp <= req.Until {
|
||||
// req.ResponseWriter.FeedNoSend(timestamp, value)
|
||||
// if timestamp < req.Since {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
// if s.lastPageNo > 0 {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: UntilFound,
|
||||
// LastPageNo: s.lastPageNo,
|
||||
// FracDigits: s.fracDigits,
|
||||
// }
|
||||
// s.RLocks++
|
||||
// } else {
|
||||
// req.ResultCh <- rangeScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *_metric) StartFullScan(req tryFullScanReq) {
|
||||
// if s.Since == 0 {
|
||||
// req.ResultCh <- fullScanResult{
|
||||
// ResultCode: QueryDone,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
|
||||
timestampDecompressor := s.timestamps.CreateDecompressor()
|
||||
valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits)
|
||||
|
||||
for {
|
||||
timestamp, done := timestampDecompressor.NextValue()
|
||||
if done {
|
||||
break
|
||||
}
|
||||
value, done := valueDecompressor.NextValue()
|
||||
if done {
|
||||
qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug)
|
||||
}
|
||||
req.ResponseWriter.FeedNoSend(timestamp, value)
|
||||
}
|
||||
|
||||
if s.lastPageNo > 0 {
|
||||
req.ResultCh <- fullScanResult{
|
||||
ResultCode: UntilFound,
|
||||
LastPageNo: s.lastPageNo,
|
||||
FracDigits: s.fracDigits,
|
||||
}
|
||||
s.RLocks++
|
||||
} else {
|
||||
req.ResultCh <- fullScanResult{
|
||||
ResultCode: QueryDone,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// індекси
|
||||
// Metric encode format:
|
||||
// metricID - 4b
|
||||
// metricType - 1b
|
||||
// fracDigits - 1b
|
||||
// lastPageNo - 4b
|
||||
// timestamps size - 2b
|
||||
// values size - 2b
|
||||
// timestams payload - Nb
|
||||
// values payload - Nb
|
||||
// index levels count - varsize
|
||||
// [
|
||||
// records qty - varsize
|
||||
// records - Nb
|
||||
// ]
|
||||
|
||||
func (s *_metric) WriteTo(w io.Writer) (err error) {
|
||||
_, err = w.Write([]byte{
|
||||
byte(s.metricType),
|
||||
s.fracDigits,
|
||||
})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.WriteUint32(w, s.lastPageNo)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.WriteUint16(w, uint16(s.timestamps.StoredSize()))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.WriteUint16(w, uint16(s.values.StoredSize()))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// timestamps payload
|
||||
err = s.timestamps.WriteStoredTo(w)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// values payload
|
||||
err = s.values.WriteStoredTo(w)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// indexes
|
||||
_, err = bin.WriteVarSize(w, len(s.indexLevelTails))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for _, level := range s.indexLevelTails {
|
||||
_, err = bin.WriteVarSize(w, level.RecordsCount)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_, err = w.Write(level.Buffer)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
@@ -1,287 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
bin "gordenko.dev/dima/bin/little"
|
||||
"gordenko.dev/dima/qb"
|
||||
"gordenko.dev/dima/qb/enc"
|
||||
"gordenko.dev/dima/qb/storage"
|
||||
)
|
||||
|
||||
type ReplayMetric struct {
|
||||
metricType qb.MetricType
|
||||
fracDigits byte
|
||||
lastPageNo uint32
|
||||
buf []byte
|
||||
databuf []byte // якщо розмір buf == DataPageSize, то databuf буде менше (мінус футер)
|
||||
tSize int
|
||||
vSize int
|
||||
indexLevelTails []storage.IndexLevelTail // root - last element
|
||||
}
|
||||
|
||||
func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) {
|
||||
metricType, err := bin.ReadByte(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
s.metricType = qb.MetricType(metricType)
|
||||
s.fracDigits, err = bin.ReadByte(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
s.lastPageNo, err = bin.ReadUint32(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
s.tSize, err = bin.ReadUint16AsInt(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
s.vSize, err = bin.ReadUint16AsInt(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.ReadNInto(r, s.databuf[len(s.databuf)-s.tSize:])
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.ReadNInto(r, s.databuf[:s.vSize])
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
//
|
||||
levelsCount, err := bin.ReadVarSize(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for range levelsCount {
|
||||
var (
|
||||
buf = make([]byte, storage.IndexPageSize)
|
||||
count int
|
||||
)
|
||||
count, err = bin.ReadVarSize(r)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = bin.ReadNInto(r, buf[:count*indexRecordSize])
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
s.indexLevelTails = append(s.indexLevelTails, storage.IndexLevelTail{
|
||||
Buffer: buf,
|
||||
RecordsCount: count,
|
||||
})
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// func appendNewRecordsToIndexTail(level storage.IndexLevelTail, records []byte) {
|
||||
|
||||
// }
|
||||
|
||||
// func countReusedIndexPages(changedIndexLevels []storage.ChangedIndexLevel) (count int) {
|
||||
|
||||
// return
|
||||
// }
|
||||
|
||||
func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) {
|
||||
timestampsSize := rec.TimestampsRewindOffset - len(rec.Timestamps)
|
||||
// копіюю нові дані із WAL з урахуванням offset-ів
|
||||
copy(s.databuf[rec.ValuesRewindOffset:], rec.Values)
|
||||
copy(s.databuf[len(s.databuf)-timestampsSize:], rec.Timestamps)
|
||||
//
|
||||
s.tSize = timestampsSize
|
||||
s.vSize = rec.ValuesRewindOffset - len(rec.Values)
|
||||
}
|
||||
|
||||
type AppendMeasuresResult struct {
|
||||
IndexPages []storage.PageToWrite
|
||||
DataPages []storage.PageToWrite
|
||||
ReusedIndexPagesCount int
|
||||
ReusedDataPagesCount int
|
||||
}
|
||||
|
||||
func (s *ReplayMetric) MeasuresAppendWithGrow(rec storage.MeasuresAppendWithGrowRecord, isLastPacket bool) (_ AppendMeasuresResult) {
|
||||
var (
|
||||
reusedIndexPages int
|
||||
reusedDataPages int
|
||||
indexPages []storage.PageToWrite
|
||||
dataPages []storage.PageToWrite
|
||||
)
|
||||
// Додати в index level tails недостаючі дані, або замінити
|
||||
for levelIdx, change := range rec.ChangedIndexLevels {
|
||||
var (
|
||||
head = change.IndexPageTail
|
||||
newRecordsCount = len(change.TailRecords) / indexRecordSize
|
||||
)
|
||||
// 1. набиваю сторінки для перезапису в index файлі
|
||||
if isLastPacket {
|
||||
if head != nil {
|
||||
indexPages = append(indexPages,
|
||||
composeHeadIndexPage(levelIdx, s.indexLevelTails[levelIdx], head))
|
||||
}
|
||||
for _, x := range change.IndexPages {
|
||||
indexPages = append(indexPages, storage.PageToWrite{
|
||||
PageNo: x.PageNo,
|
||||
Content: x.Content,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Вношу зміни в поточні індексні хвости
|
||||
if levelIdx < len(s.indexLevelTails) {
|
||||
level := s.indexLevelTails[levelIdx]
|
||||
if head != nil {
|
||||
// попередній хвіст індексного рівня перетворився на head сторінку,
|
||||
// отже TailRecords - це новий хвіст
|
||||
copy(level.Buffer, change.TailRecords)
|
||||
level.RecordsCount = newRecordsCount
|
||||
} else {
|
||||
// TailRecords - це нові дані, які треба додати
|
||||
pos := level.RecordsCount * indexRecordSize
|
||||
copy(level.Buffer[pos:], change.TailRecords)
|
||||
level.RecordsCount += newRecordsCount
|
||||
}
|
||||
s.indexLevelTails[levelIdx] = level
|
||||
} else {
|
||||
// додаю новий індексний рівень
|
||||
buf := make([]byte, storage.IndexPageSize)
|
||||
copy(buf, change.TailRecords)
|
||||
//
|
||||
s.indexLevelTails = append(s.indexLevelTails, storage.IndexLevelTail{
|
||||
Buffer: buf,
|
||||
RecordsCount: newRecordsCount,
|
||||
})
|
||||
}
|
||||
|
||||
// 3. рахую кількість reused індексних сторінок
|
||||
if head != nil {
|
||||
if head.Reused {
|
||||
reusedIndexPages++
|
||||
}
|
||||
}
|
||||
for _, x := range change.IndexPages {
|
||||
if x.Reused {
|
||||
reusedIndexPages++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 4. набиваю сторінки для перезапису в data файлі
|
||||
if isLastPacket {
|
||||
dataPages = append(dataPages, composeHeadDataPage(s.databuf, rec.DataPageTail))
|
||||
|
||||
for _, x := range rec.DataPages {
|
||||
dataPages = append(dataPages, storage.PageToWrite{
|
||||
PageNo: x.PageNo,
|
||||
Content: x.Content,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// рахую кількість reused дата сторінок
|
||||
if rec.DataPageTail.Reused {
|
||||
reusedDataPages++
|
||||
}
|
||||
for _, x := range rec.DataPages {
|
||||
if x.Reused {
|
||||
reusedDataPages++
|
||||
}
|
||||
}
|
||||
// HeadDataPage є полюбому, оскільки це транзакція із мінімум однією заповненою
|
||||
// data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня.
|
||||
copy(s.databuf, rec.TailValues)
|
||||
pos := len(s.databuf) - len(rec.TailTimestamps)
|
||||
copy(s.databuf[pos:], rec.TailTimestamps)
|
||||
|
||||
s.tSize = len(rec.TailTimestamps)
|
||||
s.vSize = len(rec.TailValues)
|
||||
|
||||
if len(rec.DataPages) == 0 {
|
||||
s.lastPageNo = rec.DataPageTail.PageNo
|
||||
} else {
|
||||
s.lastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo
|
||||
}
|
||||
|
||||
return AppendMeasuresResult{
|
||||
IndexPages: indexPages,
|
||||
DataPages: dataPages,
|
||||
ReusedIndexPagesCount: reusedIndexPages,
|
||||
ReusedDataPagesCount: reusedDataPages,
|
||||
}
|
||||
}
|
||||
|
||||
// HELPERS
|
||||
|
||||
func composeHeadIndexPage(levelIdx int, level storage.IndexLevelTail, head *storage.IndexPageTail) storage.PageToWrite {
|
||||
// розраховую pos, з якого буду дописувати хвіст
|
||||
pos := level.RecordsCount * indexRecordSize
|
||||
// створюю копію сторінки
|
||||
page := make([]byte, storage.IndexPageSize)
|
||||
copy(page, level.Buffer[:pos]) // поточні дані
|
||||
copy(page[pos:], head.Records)
|
||||
// запечатати сторінку
|
||||
calculatedCRC := storage.SealIndexPage(storage.SealIndexPageIn{
|
||||
Content: page,
|
||||
RecordsCount: level.RecordsCount + len(head.Records)/indexRecordSize,
|
||||
ZeroLevel: levelIdx == 0,
|
||||
})
|
||||
// перевірка CRC
|
||||
if calculatedCRC != head.CRC32 {
|
||||
qb.Abort(qb.WALReplayFailed,
|
||||
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{
|
||||
PageNo: head.PageNo,
|
||||
Content: page,
|
||||
}
|
||||
}
|
||||
|
||||
func composeHeadDataPage(databuf []byte, head storage.DataPageTail) storage.PageToWrite {
|
||||
// створюю копію сторінки
|
||||
var (
|
||||
page = make([]byte, storage.DataPageSize)
|
||||
timestampsSize = head.TimestampsRewindOffset - len(head.Timestamps)
|
||||
)
|
||||
// копіюю поточні дані
|
||||
copy(page, databuf)
|
||||
// копіюю нові дані із WAL з урахуванням offset-ів
|
||||
copy(page[head.ValuesRewindOffset:], head.Values)
|
||||
copy(page[len(page)-timestampsSize:], head.Timestamps)
|
||||
// запечатати сторінку
|
||||
calculatedCRC := storage.SealDataPage(storage.SealDataPageIn{
|
||||
Content: page,
|
||||
PrevPageNo: head.PrevPageNo,
|
||||
TimestampsSize: timestampsSize,
|
||||
ValuesSize: head.ValuesRewindOffset + len(head.Values),
|
||||
})
|
||||
// перевірка CRC
|
||||
if calculatedCRC != head.CRC32 {
|
||||
qb.Abort(qb.WALReplayFailed,
|
||||
fmt.Errorf("calculated CRC %d not equal expected %d of head data page %d",
|
||||
calculatedCRC, head.CRC32, head.PageNo))
|
||||
}
|
||||
return storage.PageToWrite{
|
||||
PageNo: head.PageNo,
|
||||
Content: page,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *ReplayMetric) ToMetric() *_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,
|
||||
}
|
||||
}
|
||||
@@ -1,272 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
||||
bin "gordenko.dev/dima/bin/little"
|
||||
"gordenko.dev/dima/qb/storage"
|
||||
"gordenko.dev/dima/qb/util"
|
||||
)
|
||||
|
||||
/*
|
||||
Формат:
|
||||
metricsQty - varsize
|
||||
[metric]*
|
||||
де metric - це:
|
||||
|
||||
index free list frozen pages - varsize
|
||||
indexFreeList size - varsize
|
||||
indexFreeList - Nb
|
||||
|
||||
data free list frozen pages - varsize
|
||||
dataFreeList size - varsize
|
||||
dataFreeList - Nb
|
||||
|
||||
CRC32 - 4b
|
||||
*/
|
||||
|
||||
const readBufferSize = 8 * 1024 * 1024 // 8mb
|
||||
|
||||
type writeSnapshotIn struct {
|
||||
SnapshotNumber int
|
||||
Dir string
|
||||
WriteBufferSize int
|
||||
Metrics map[uint32]*_metric
|
||||
FrozenIndexPagesCount int
|
||||
IndexPageNumbers []uint32
|
||||
FrozenDataPagesCount int
|
||||
DataPageNumbers []uint32
|
||||
}
|
||||
|
||||
func writeSnapshot(in writeSnapshotIn) (err error) {
|
||||
var (
|
||||
fileName = filepath.Join(in.Dir, fmt.Sprintf("%d.snapshot", in.SnapshotNumber))
|
||||
hasher = util.NewHasher()
|
||||
)
|
||||
|
||||
file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0666) //0770)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
dst := io.MultiWriter(bufio.NewWriterSize(file, in.WriteBufferSize), hasher)
|
||||
|
||||
_, err = bin.WriteVarSize(dst, len(in.Metrics))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
for metricID, metric := range in.Metrics {
|
||||
err = bin.WriteUint32(dst, metricID)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = metric.WriteTo(dst)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
// free index pages
|
||||
err = writeFreePages(dst, in.FrozenIndexPagesCount, in.IndexPageNumbers)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// free data pages
|
||||
err = writeFreePages(dst, in.FrozenDataPagesCount, in.DataPageNumbers)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// CRC32
|
||||
bin.WriteUint32(file, hasher.Sum32())
|
||||
|
||||
err = file.Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// prevLogNumber := logNumber - 1
|
||||
// prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber))
|
||||
// prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber))
|
||||
|
||||
// isExist, err := isFileExist(prevChanges)
|
||||
// if err != nil {
|
||||
// return
|
||||
// }
|
||||
|
||||
// if isExist {
|
||||
// err = os.Remove(prevChanges)
|
||||
// if err != nil {
|
||||
// qb.Abort(qb.DeletePrevChangesFileFailed, err)
|
||||
// }
|
||||
// }
|
||||
|
||||
// isExist, err = isFileExist(prevSnapshot)
|
||||
// if err != nil {
|
||||
// return
|
||||
// }
|
||||
|
||||
// if isExist {
|
||||
// err = os.Remove(prevSnapshot)
|
||||
// if err != nil {
|
||||
// qb.Abort(qb.DeletePrevSnapshotFileFailed, err)
|
||||
// }
|
||||
// }
|
||||
return
|
||||
}
|
||||
|
||||
type readSnapshotOut struct {
|
||||
Metrics map[uint32]*ReplayMetric
|
||||
FrozenIndexPagesCount int
|
||||
IndexPageNumbers []uint32
|
||||
FrozenDataPagesCount int
|
||||
DataPageNumbers []uint32
|
||||
}
|
||||
|
||||
func readSnapshot(fileName string) (out readSnapshotOut, err error) {
|
||||
var (
|
||||
metrics = make(map[uint32]*ReplayMetric)
|
||||
frozenIndexPagesCount int
|
||||
indexPageNumbers []uint32
|
||||
frozenDataPagesCount int
|
||||
dataPageNumbers []uint32
|
||||
hasher = util.NewHasher()
|
||||
)
|
||||
|
||||
file, err := os.Open(fileName)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
stat, err := file.Stat()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
fileSize := stat.Size()
|
||||
|
||||
if fileSize == 0 {
|
||||
return readSnapshotOut{
|
||||
Metrics: metrics,
|
||||
}, nil
|
||||
}
|
||||
|
||||
if fileSize < 4 {
|
||||
err = fmt.Errorf("%s is corrupted", fileName)
|
||||
return
|
||||
}
|
||||
payloadSize := fileSize - 4
|
||||
|
||||
// 2. Створюємо великий буфер для NVMe (4 MB)
|
||||
bufferedReader := bufio.NewReaderSize(file, readBufferSize)
|
||||
|
||||
// 3. Обмежуємо читання лише розміром Payload
|
||||
limitReader := io.LimitReader(bufferedReader, payloadSize)
|
||||
|
||||
// 4. Створюємо хеш та обгортку TeeReader, яка рахує CRC32 на льоту
|
||||
src := io.TeeReader(limitReader, hasher)
|
||||
|
||||
// читаю payload
|
||||
metricsQty, err := bin.ReadVarSize(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for range metricsQty {
|
||||
var (
|
||||
metricID uint32
|
||||
buf = make([]byte, storage.DataPageSize)
|
||||
// для простоти створюю буфер максимального розміру
|
||||
metric = &ReplayMetric{
|
||||
buf: buf,
|
||||
databuf: buf[:storage.DataPagePayloadSize],
|
||||
}
|
||||
)
|
||||
metricID, err = bin.ReadUint32(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = metric.ReadFrom(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
metrics[metricID] = metric
|
||||
}
|
||||
// index pages
|
||||
frozenIndexPagesCount, indexPageNumbers, err = readFreePages(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// data pages
|
||||
frozenDataPagesCount, dataPageNumbers, err = readFreePages(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
// verify CRC32
|
||||
expectedCRC, err := bin.ReadUint32(bufferedReader)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
calculatedCRC := hasher.Sum32()
|
||||
|
||||
if expectedCRC != calculatedCRC {
|
||||
err = fmt.Errorf("%s is corrupted. Calculated CRC %d not equal expected CRC %d",
|
||||
fileName, calculatedCRC, expectedCRC)
|
||||
return
|
||||
}
|
||||
|
||||
return readSnapshotOut{
|
||||
Metrics: metrics,
|
||||
FrozenIndexPagesCount: frozenIndexPagesCount,
|
||||
IndexPageNumbers: indexPageNumbers,
|
||||
FrozenDataPagesCount: frozenDataPagesCount,
|
||||
DataPageNumbers: dataPageNumbers,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// HELPERS
|
||||
|
||||
func writeFreePages(dst io.Writer, frozenPagesCount int, pageNumbers []uint32) (err error) {
|
||||
_, err = bin.WriteVarSize(dst, frozenPagesCount)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_, err = bin.WriteVarSize(dst, len(pageNumbers))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for _, pageNo := range pageNumbers {
|
||||
err = bin.WriteUint32(dst, pageNo)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func readFreePages(src io.Reader) (_ int, _ []uint32, err error) {
|
||||
var (
|
||||
frozenPagesCount int
|
||||
pageNumbers []uint32
|
||||
)
|
||||
frozenPagesCount, err = bin.ReadVarSize(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
pageNumbersCount, err := bin.ReadVarSize(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for range pageNumbersCount {
|
||||
var pageNo uint32
|
||||
pageNo, err = bin.ReadUint32(src)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
pageNumbers = append(pageNumbers, pageNo)
|
||||
}
|
||||
return frozenPagesCount, pageNumbers, nil
|
||||
}
|
||||
@@ -1,177 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"gordenko.dev/dima/qb/storage"
|
||||
)
|
||||
|
||||
type WALReplayer struct {
|
||||
walReader *storage.WALReader
|
||||
metrics map[uint32]*ReplayMetric
|
||||
freeIndexPages []uint32
|
||||
freeDataPages []uint32
|
||||
indexPages []storage.PageToWrite
|
||||
dataPages []storage.PageToWrite
|
||||
}
|
||||
|
||||
type WALReplayerOptions struct {
|
||||
WALReader *storage.WALReader
|
||||
Metrics map[uint32]*ReplayMetric // from snapshot
|
||||
FreeIndexPages []uint32 // from snapshot
|
||||
FreeDataPages []uint32 // from snapshot
|
||||
}
|
||||
|
||||
func NewWALReplayer(opt WALReplayerOptions) (*WALReplayer, error) {
|
||||
if opt.WALReader == nil {
|
||||
return nil, errors.New("required option: WALReader")
|
||||
}
|
||||
if opt.Metrics == nil {
|
||||
return nil, errors.New("required option: Metrics")
|
||||
}
|
||||
return &WALReplayer{
|
||||
walReader: opt.WALReader,
|
||||
metrics: opt.Metrics,
|
||||
freeIndexPages: opt.FreeIndexPages,
|
||||
freeDataPages: opt.FreeDataPages,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *WALReplayer) IndexPages() []storage.PageToWrite {
|
||||
return s.indexPages
|
||||
}
|
||||
|
||||
func (s *WALReplayer) DataPages() []storage.PageToWrite {
|
||||
return s.dataPages
|
||||
}
|
||||
|
||||
func (s *WALReplayer) Replay() error {
|
||||
for {
|
||||
records, isLastPacket, err := s.walReader.NextPacket()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, rec := range records {
|
||||
if err = s.replayRecord(rec, isLastPacket); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *WALReplayer) replayRecord(untyped any, isLastPacket bool) (err error) {
|
||||
switch rec := untyped.(type) {
|
||||
case storage.MeasuresAppendRecord:
|
||||
if err = s.onMeasuresAppend(rec); err != nil {
|
||||
return
|
||||
}
|
||||
case storage.MeasuresAppendWithGrowRecord:
|
||||
if err = s.onMeasuresAppendWithGrow(rec, isLastPacket); err != nil {
|
||||
return
|
||||
}
|
||||
case storage.MetricAddRecord:
|
||||
if err = s.onMetricAdd(rec); err != nil {
|
||||
return
|
||||
}
|
||||
case storage.MetricDeleteRecord:
|
||||
if err = s.onMetricDelete(rec); err != nil {
|
||||
return
|
||||
}
|
||||
// case storage.DeletedMeasures:
|
||||
// metric, ok := s.metrics[rec.MetricID]
|
||||
// if ok {
|
||||
// metric.DeleteMeasures()
|
||||
// if len(rec.FreePageNumbers) > 0 {
|
||||
// s.freeList.AddPages(rec.FreePageNumbers)
|
||||
// }
|
||||
// }
|
||||
|
||||
// default:
|
||||
// qb.Abort(qb.UnknownstorageRecordTypeBug,
|
||||
// fmt.Errorf("bug: unknown record type %T in TransactionLog", rec))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *WALReplayer) onMetricAdd(rec storage.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)
|
||||
)
|
||||
s.metrics[rec.MetricID] = &ReplayMetric{
|
||||
metricType: rec.MetricType,
|
||||
fracDigits: byte(rec.FracDigits),
|
||||
buf: buf,
|
||||
databuf: buf[:storage.DataPagePayloadSize],
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *WALReplayer) onMetricDelete(rec storage.MetricDeleteRecord) error {
|
||||
_, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
return fmt.Errorf("metric %d deletion failed: not found", rec.MetricID)
|
||||
}
|
||||
delete(s.metrics, rec.MetricID)
|
||||
s.freeIndexPages = append(s.freeIndexPages, rec.FreeIndexPages...)
|
||||
s.freeDataPages = append(s.freeDataPages, rec.FreeDataPages...)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *WALReplayer) onMeasuresAppend(rec storage.MeasuresAppendRecord) (err error) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
return fmt.Errorf("append measures failed: metric %d not found", rec.MetricID)
|
||||
}
|
||||
metric.MeasuresAppend(rec)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *WALReplayer) onMeasuresAppendWithGrow(rec storage.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)
|
||||
}
|
||||
result := metric.MeasuresAppendWithGrow(rec, isLastPacket)
|
||||
//
|
||||
if isLastPacket {
|
||||
s.indexPages = append(s.indexPages, result.IndexPages...)
|
||||
s.dataPages = append(s.dataPages, result.DataPages...)
|
||||
}
|
||||
//
|
||||
s.freeIndexPages = s.freeIndexPages[:len(s.freeIndexPages)-result.ReusedIndexPagesCount]
|
||||
s.freeDataPages = s.freeDataPages[:len(s.freeDataPages)-result.ReusedDataPagesCount]
|
||||
return
|
||||
}
|
||||
|
||||
//func (s *WALReplayer) collectPagesToWrite(metric *ReplayMetric, rec storage.WALRecordAppendMeasures) {
|
||||
// metric.ApplyHeadDataPage(rec.CompletedDataPage)
|
||||
|
||||
// s.dataPages = append(s.dataPages, storage.PageToWrite{
|
||||
// PageNo: p.PageNo,
|
||||
// Content: metric.buf,
|
||||
// })
|
||||
|
||||
// for _, x := range rec.DataPages {
|
||||
// s.dataPages = append(s.dataPages, storage.PageToWrite{
|
||||
// PageNo: x.PageNo,
|
||||
// Content: x.Content,
|
||||
// })
|
||||
// }
|
||||
//metric.ApplyHeadIndexPages(rec.ChangedIndexLevels)
|
||||
|
||||
//}
|
||||
|
||||
// func (s *WALReplayer) replayPacket(records []any, isLastPacket bool) {
|
||||
// for _, untyped := range records {
|
||||
// switch rec := untyped.(type) {
|
||||
// case storage.WALRecordAppendMeasures:
|
||||
// s.onWALRecordAppendMeasures(rec, isLastPacket)
|
||||
// }
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
@@ -1,738 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
qb "gordenko.dev/dima/qb"
|
||||
"gordenko.dev/dima/qb/atree"
|
||||
"gordenko.dev/dima/qb/enc"
|
||||
"gordenko.dev/dima/qb/proto"
|
||||
"gordenko.dev/dima/qb/storage"
|
||||
"gordenko.dev/dima/qb/transform"
|
||||
)
|
||||
|
||||
// AddMetric - create metric in map and set Xlock, rwad - append to queue
|
||||
// DeleteMetric - set Xlock
|
||||
// DeleteSince - xLock (тупо редкая операция)
|
||||
// Якийсь прапор поставити для pendingDelete, щоб read операції вставали в чергу
|
||||
|
||||
const (
|
||||
QueryDone = 1
|
||||
UntilFound = 2
|
||||
UntilNotFound = 3
|
||||
RangeFound = 15
|
||||
NoMeasures = 16
|
||||
NoMetric = 4
|
||||
MetricDuplicate = 5
|
||||
Succeed = 6
|
||||
NewPage = 7
|
||||
ExpiredMeasure = 8
|
||||
NonMonotonicValue = 9
|
||||
CanAppend = 10
|
||||
WrongMetricType = 11
|
||||
NoMeasuresToDelete = 12
|
||||
DeleteFromAtreeNotNeeded = 13
|
||||
DeleteFromAtreeRequired = 14
|
||||
)
|
||||
|
||||
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.signalCh:
|
||||
s.doWork()
|
||||
|
||||
case <-s.exitCh:
|
||||
fmt.Println("worker done")
|
||||
s.waitGroup.Done()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Worker) doWork() {
|
||||
|
||||
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]
|
||||
// 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.RLocks <= 0 {
|
||||
// qb.Abort(qb.NoRLockBug,
|
||||
// fmt.Errorf("drainQueues: rlock not set for the metric %d",
|
||||
// metricID))
|
||||
// }
|
||||
|
||||
// metric.RLocks--
|
||||
|
||||
// if len(metric.WaitQueue) > 0 {
|
||||
// s.processMetricQueue(metricID, metric)
|
||||
// }
|
||||
// }
|
||||
|
||||
for _, untyped := range queue {
|
||||
switch req := untyped.(type) {
|
||||
case tryAppendMeasuresReq:
|
||||
s.tryAppendMeasures(req)
|
||||
|
||||
case storage.Changes:
|
||||
s.applyCommits(req) // all metrics only
|
||||
|
||||
case tryListCurrentValuesReq:
|
||||
s.tryListCurrentValues(req) // all metrics only
|
||||
|
||||
case tryRangeScanReq:
|
||||
s.tryRangeScan(req)
|
||||
|
||||
case tryFullScanReq:
|
||||
s.tryFullScan(req)
|
||||
|
||||
case tryAddMetricReq:
|
||||
s.tryAddMetric(req)
|
||||
|
||||
case tryDeleteMetricReq:
|
||||
s.tryDeleteMetric(req)
|
||||
|
||||
case tryDeleteMeasuresReq:
|
||||
s.tryDeleteMeasures(req)
|
||||
|
||||
case tryGetMetricReq:
|
||||
s.tryGetMetric(req)
|
||||
|
||||
default:
|
||||
qb.Abort(qb.UnknownWorkerQueueItemBug,
|
||||
fmt.Errorf("bug: unknown worker queue item type %T", req))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// суть у тому що треба запускати запити, пока не зустріну XLock
|
||||
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:
|
||||
metric.StartRangeScan(req)
|
||||
|
||||
case tryFullScanReq:
|
||||
metric.StartFullScan(req)
|
||||
|
||||
case tryGetMetricReq:
|
||||
s.tryGetMetric(req)
|
||||
|
||||
case tryAppendMeasuresReq:
|
||||
metric.StartAppendMeasures(req, tmp, s.saveToStorage)
|
||||
|
||||
case tryDeleteMetricReq:
|
||||
s.startDeleteMetric(metric, req)
|
||||
|
||||
case tryDeleteMeasuresReq:
|
||||
s.startDeleteMeasures(metric, req)
|
||||
|
||||
default:
|
||||
qb.Abort(qb.UnknownMetricWaitQueueItemBug,
|
||||
fmt.Errorf("bug: unknown metric wait queue item type %T", req))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type tryAddMetricReq struct {
|
||||
MetricID uint32
|
||||
ResultCh chan byte
|
||||
}
|
||||
|
||||
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]
|
||||
// if ok {
|
||||
// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req)
|
||||
// } else {
|
||||
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
|
||||
// XLock: true,
|
||||
// }
|
||||
// req.ResultCh <- Succeed
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []tryAddMetricReq) {
|
||||
if len(reqs) == 0 {
|
||||
return
|
||||
}
|
||||
var (
|
||||
req = reqs[0]
|
||||
waitQueue []any
|
||||
)
|
||||
if len(reqs) > 1 {
|
||||
for _, req := range reqs[1:] {
|
||||
waitQueue = append(waitQueue, req)
|
||||
}
|
||||
}
|
||||
// FIX
|
||||
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
|
||||
// XLock: true,
|
||||
// WaitQueue: waitQueue,
|
||||
// }
|
||||
req.ResultCh <- Succeed
|
||||
}
|
||||
|
||||
type tryGetMetricReq struct {
|
||||
MetricID uint32
|
||||
ResultCh chan Metric
|
||||
}
|
||||
|
||||
func (s *Worker) tryGetMetric(req tryGetMetricReq) {
|
||||
metric, ok := s.metrics[req.MetricID]
|
||||
if ok {
|
||||
req.ResultCh <- Metric{
|
||||
ResultCode: Succeed,
|
||||
MetricType: metric.metricType,
|
||||
FracDigits: metric.fracDigits,
|
||||
}
|
||||
} else {
|
||||
req.ResultCh <- Metric{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type tryDeleteMetricReq struct {
|
||||
MetricID uint32
|
||||
ResultCh chan tryDeleteMetricResult
|
||||
}
|
||||
|
||||
func (s *Worker) tryDeleteMetric(req tryDeleteMetricReq) {
|
||||
// FIX
|
||||
// metric, ok := s.metrics[req.MetricID]
|
||||
// if !ok {
|
||||
// req.ResultCh <- tryDeleteMetricResult{
|
||||
// ResultCode: NoMetric,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
|
||||
// lockEntry, ok := s.metricLockEntries[req.MetricID]
|
||||
// if ok {
|
||||
// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req)
|
||||
// } else {
|
||||
// s.startDeleteMetric(metric, req)
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *Worker) startDeleteMetric(metric *_metric, req tryDeleteMetricReq) {
|
||||
// FIX
|
||||
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
|
||||
// XLock: true,
|
||||
// }
|
||||
// req.ResultCh <- tryDeleteMetricResult{
|
||||
// ResultCode: Succeed,
|
||||
// RootPageNo: metric.RootPageNo,
|
||||
// }
|
||||
}
|
||||
|
||||
type tryDeleteMeasuresReq struct {
|
||||
MetricID uint32
|
||||
Since uint32
|
||||
ResultCh chan tryDeleteMeasuresResult
|
||||
}
|
||||
|
||||
func (s *Worker) tryDeleteMeasures(req tryDeleteMeasuresReq) {
|
||||
// FIX
|
||||
// metric, ok := s.metrics[req.MetricID]
|
||||
// if !ok {
|
||||
// req.ResultCh <- tryDeleteMeasuresResult{
|
||||
// ResultCode: NoMetric,
|
||||
// }
|
||||
// return
|
||||
// }
|
||||
|
||||
// if metric.Since == 0 || (req.Since > 0 && metric.Until < req.Since) {
|
||||
// req.ResultCh <- tryDeleteMeasuresResult{
|
||||
// ResultCode: NoMeasuresToDelete,
|
||||
// }
|
||||
// }
|
||||
|
||||
// lockEntry, ok := s.metricLockEntries[req.MetricID]
|
||||
// if ok {
|
||||
// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req)
|
||||
// } else {
|
||||
// s.startDeleteMeasures(metric, req)
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *Worker) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq) {
|
||||
// FIX
|
||||
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
|
||||
// XLock: true,
|
||||
// }
|
||||
|
||||
// if metric.RootPageNo > 0 {
|
||||
// req.ResultCh <- tryDeleteMeasuresResult{
|
||||
// ResultCode: DeleteFromAtreeRequired,
|
||||
// RootPageNo: metric.RootPageNo,
|
||||
// }
|
||||
// } else {
|
||||
// req.ResultCh <- tryDeleteMeasuresResult{
|
||||
// ResultCode: DeleteFromAtreeNotNeeded,
|
||||
// }
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *Worker) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
fmt.Errorf("finAppendMeasures: metric %d not found",
|
||||
rec.MetricID))
|
||||
}
|
||||
metric.OnMeasuresAppendCommited(rec)
|
||||
}
|
||||
|
||||
func (s *Worker) onMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
fmt.Errorf("finAppendMeasures: metric %d not found",
|
||||
rec.MetricID))
|
||||
}
|
||||
metric.OnMeasuresAppendWithGrowCommited(rec)
|
||||
}
|
||||
|
||||
type tryAppendMeasuresReq struct {
|
||||
MetricID uint32
|
||||
Measures []proto.Measure
|
||||
ResultCh chan tryAppendMeasuresResult
|
||||
}
|
||||
|
||||
func (s *Worker) tryAppendMeasures(req tryAppendMeasuresReq) {
|
||||
metric, ok := s.metrics[req.MetricID]
|
||||
if !ok {
|
||||
req.ResultCh <- tryAppendMeasuresResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
return
|
||||
}
|
||||
if metric.XLock {
|
||||
metric.WaitQueue = append(metric.WaitQueue, req)
|
||||
return
|
||||
}
|
||||
metric.StartAppendMeasures(req, s.tmp, s.saveToStorage)
|
||||
}
|
||||
|
||||
type tryRangeScanReq struct {
|
||||
MetricID uint32
|
||||
Since uint32
|
||||
Until uint32
|
||||
MetricType qb.MetricType
|
||||
ResponseWriter atree.WorkerMeasureConsumer
|
||||
ResultCh chan rangeScanResult
|
||||
}
|
||||
|
||||
func (s *Worker) tryRangeScan(req tryRangeScanReq) {
|
||||
metric, ok := s.metrics[req.MetricID]
|
||||
if !ok {
|
||||
req.ResultCh <- rangeScanResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
return
|
||||
}
|
||||
if metric.metricType != req.MetricType {
|
||||
req.ResultCh <- rangeScanResult{
|
||||
ResultCode: WrongMetricType,
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if metric.XLock {
|
||||
metric.WaitQueue = append(metric.WaitQueue, req)
|
||||
return
|
||||
}
|
||||
|
||||
metric.StartRangeScan(req)
|
||||
}
|
||||
|
||||
type tryFullScanReq struct {
|
||||
MetricID uint32
|
||||
MetricType qb.MetricType
|
||||
ResponseWriter atree.WorkerMeasureConsumer
|
||||
ResultCh chan fullScanResult
|
||||
}
|
||||
|
||||
func (s *Worker) tryFullScan(req tryFullScanReq) {
|
||||
metric, ok := s.metrics[req.MetricID]
|
||||
if !ok {
|
||||
req.ResultCh <- fullScanResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
return
|
||||
}
|
||||
if metric.metricType != req.MetricType {
|
||||
req.ResultCh <- fullScanResult{
|
||||
ResultCode: WrongMetricType,
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if metric.XLock {
|
||||
metric.WaitQueue = append(metric.WaitQueue, req)
|
||||
return
|
||||
}
|
||||
metric.StartFullScan(req)
|
||||
}
|
||||
|
||||
type tryListCurrentValuesReq struct {
|
||||
MetricIDs []uint32
|
||||
ResponseWriter *transform.CurrentValueWriter
|
||||
ResultCh chan struct{}
|
||||
}
|
||||
|
||||
func (s *Worker) tryListCurrentValues(req tryListCurrentValuesReq) {
|
||||
for _, metricID := range req.MetricIDs {
|
||||
metric, ok := s.metrics[metricID]
|
||||
if ok {
|
||||
req.ResponseWriter.BufferValue(transform.CurrentValue{
|
||||
MetricID: metricID,
|
||||
Timestamp: metric.LastTimestamp(),
|
||||
Value: metric.LastValue(),
|
||||
})
|
||||
}
|
||||
}
|
||||
req.ResultCh <- struct{}{}
|
||||
}
|
||||
|
||||
///////////////////////////////////////////////////////
|
||||
|
||||
func (s *Worker) applyCommits(req storage.Changes) {
|
||||
for _, untyped := range req.Commits {
|
||||
switch rec := untyped.(type) {
|
||||
case storage.MetricAddCommited:
|
||||
s.onMetricAddCommited(rec)
|
||||
|
||||
case storage.MetricDeleteCommited:
|
||||
s.onMetricDeleteCommited(rec)
|
||||
|
||||
case storage.MeasuresAppendCommited:
|
||||
s.onMeasuresAppendCommited(rec)
|
||||
|
||||
case storage.MeasuresAppendWithGrowCommited:
|
||||
s.onMeasuresAppendWithGrowCommited(rec)
|
||||
|
||||
case storage.MeasuresDeleteCommited:
|
||||
s.onMeasuresDeleteCommited(rec)
|
||||
}
|
||||
}
|
||||
|
||||
if req.SnapshotNumberCh != nil {
|
||||
s.snapshotNumber++
|
||||
err := writeSnapshot(writeSnapshotIn{
|
||||
SnapshotNumber: s.snapshotNumber,
|
||||
Dir: s.dir,
|
||||
WriteBufferSize: 4 * 1024 * 1024, // 1mb
|
||||
Metrics: s.metrics,
|
||||
FrozenIndexPagesCount: req.FrozenIndexPagesCount,
|
||||
IndexPageNumbers: req.IndexPageNumbers,
|
||||
FrozenDataPagesCount: req.FrozenDataPagesCount,
|
||||
DataPageNumbers: req.DataPageNumbers,
|
||||
})
|
||||
if err != nil {
|
||||
qb.Abort(qb.WriteSnapshotFailed, err)
|
||||
}
|
||||
|
||||
req.SnapshotNumberCh <- s.snapshotNumber
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Worker) onMetricAddCommited(rec storage.MetricAddCommited) {
|
||||
// fix lock
|
||||
_, ok := s.metrics[rec.MetricID]
|
||||
if ok {
|
||||
qb.Abort(qb.MetricAddedBug,
|
||||
fmt.Errorf("addMetric: metric %d already added",
|
||||
rec.MetricID))
|
||||
}
|
||||
|
||||
// if !lockEntry.XLock {
|
||||
// qb.Abort(qb.NoXLockBug,
|
||||
// fmt.Errorf("addMetric: xlock not set for the metric %d",
|
||||
// rec.MetricID))
|
||||
// }
|
||||
|
||||
// lockEntry.XLock = false
|
||||
// delete(s.metricLockEntries, rec.MetricID)
|
||||
}
|
||||
|
||||
func (s *Worker) addMetric(rec storage.MetricAddRecord) {
|
||||
var (
|
||||
buf = make([]byte, storage.DataPageSize)
|
||||
databuf = buf[:storage.DataPagePayloadSize]
|
||||
)
|
||||
s.metrics[rec.MetricID] = &_metric{
|
||||
metricType: rec.MetricType,
|
||||
fracDigits: rec.FracDigits,
|
||||
buffer: buf,
|
||||
timestamps: enc.NewTimeDeltaCompressor(databuf, 0),
|
||||
values: enc.NewValueDeltaCompressor(rec.MetricType, rec.FracDigits, databuf, 0),
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
fmt.Errorf("deleteMetric: metric %d not found",
|
||||
rec.MetricID))
|
||||
}
|
||||
|
||||
if !metric.XLock {
|
||||
qb.Abort(qb.NoXLockBug,
|
||||
fmt.Errorf("deleteMetric: xlock not set for the metric %d",
|
||||
rec.MetricID))
|
||||
}
|
||||
|
||||
var addMetricReqs []tryAddMetricReq
|
||||
|
||||
if len(metric.WaitQueue) > 0 {
|
||||
for _, untyped := range metric.WaitQueue {
|
||||
switch req := untyped.(type) {
|
||||
// case tryAppendMeasureReq:
|
||||
// req.ResultCh <- tryAppendMeasureResult{
|
||||
// MetricID: req.MetricID,
|
||||
// ResultCode: NoMetric,
|
||||
// }
|
||||
|
||||
case tryRangeScanReq:
|
||||
req.ResultCh <- rangeScanResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
|
||||
case tryFullScanReq:
|
||||
req.ResultCh <- fullScanResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
|
||||
case tryAddMetricReq:
|
||||
addMetricReqs = append(addMetricReqs, req)
|
||||
|
||||
case tryDeleteMetricReq:
|
||||
req.ResultCh <- tryDeleteMetricResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
|
||||
case tryDeleteMeasuresReq:
|
||||
req.ResultCh <- tryDeleteMeasuresResult{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
|
||||
case tryGetMetricReq:
|
||||
req.ResultCh <- Metric{
|
||||
ResultCode: NoMetric,
|
||||
}
|
||||
|
||||
default:
|
||||
qb.Abort(qb.UnknownMetricWaitQueueItemBug,
|
||||
fmt.Errorf("bug: unknown metric wait queue item type %T", req))
|
||||
}
|
||||
}
|
||||
}
|
||||
delete(s.metrics, rec.MetricID)
|
||||
|
||||
// ADD in storage
|
||||
// if len(rec.FreePageNumbers) > 0 {
|
||||
// s.freeList.AddPages(rec.FreePageNumbers)
|
||||
// }
|
||||
|
||||
if len(addMetricReqs) > 0 {
|
||||
s.processTryAddMetricReqsImmediatelyAfterDelete(addMetricReqs)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Worker) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
fmt.Errorf("deleteMeasures: metric %d not found",
|
||||
rec.MetricID))
|
||||
}
|
||||
|
||||
if !metric.XLock {
|
||||
qb.Abort(qb.NoXLockBug,
|
||||
fmt.Errorf("deleteMeasures: xlock not set for the metric %d",
|
||||
rec.MetricID))
|
||||
}
|
||||
metric.OnMeasuresDeleteCommited(rec)
|
||||
metric.XLock = false
|
||||
// FIX add in storage
|
||||
// if len(rec.FreePageNumbers) > 0 {
|
||||
// s.freeList.AddPages(rec.FreePageNumbers)
|
||||
// }
|
||||
//s.doAfterReleaseXLock(rec.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