Files
qb/worker/worker.go
2026-06-14 07:57:01 +03:00

691 lines
15 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package worker
import (
"fmt"
"sync"
qb "gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/enc"
"gordenko.dev/dima/qb/inbox"
"gordenko.dev/dima/qb/proto"
"gordenko.dev/dima/qb/storage"
"gordenko.dev/dima/qb/transform"
)
// 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 {
inbox *inbox.Inbox
storageInbox *inbox.Inbox
dir string
snapshotNumber int
tmp []byte
metrics map[uint32]*Metric
exitCh chan struct{}
waitGroup *sync.WaitGroup
}
type Options struct {
Inbox *inbox.Inbox
StorageInbox *inbox.Inbox
ReplayMetrics map[uint32]*storage.ReplayMetric
Dir string
ExitCh chan struct{}
WaitGroup *sync.WaitGroup
}
func New(opt Options) *Worker {
s := &Worker{
inbox: opt.Inbox,
storageInbox: opt.StorageInbox,
dir: opt.Dir,
tmp: make([]byte, 26),
metrics: make(map[uint32]*Metric),
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
}
for metricID, x := range opt.ReplayMetrics {
databuf := x.Buf[:storage.DataPagePayloadSize]
values := enc.NewValueDeltaCompressor(x.MetricType, x.FracDigits, databuf, x.ValuesSize)
s.metrics[metricID] = &Metric{
metricType: x.MetricType,
fracDigits: x.FracDigits,
lastPageNo: x.LastPageNo,
lastValue: values.LastValue(),
buffer: x.Buf,
timestamps: enc.NewTimeDeltaCompressor(databuf, x.TimestampsSize),
values: values,
indexLevelTails: x.IndexLevelTails,
}
}
return s
}
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.inbox.Ready():
s.doWork()
case <-s.exitCh:
fmt.Println("worker done")
s.waitGroup.Done()
return
}
}
}
func (s *Worker) doWork() {
queue := s.inbox.Drain()
//rLocksToRelease := s.rLocksToRelease
//s.rLocksToRelease = nil
// 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 AppendMeasuresReq:
s.AppendMeasures(req)
case storage.Changes:
s.applyCommits(req) // all metrics only
case ListCurrentValuesReq:
s.tryListCurrentValues(req) // all metrics only
case RangeScanReq:
s.RangeScan(req)
case FullScanReq:
s.tryFullScan(req)
case AddMetricReq:
s.AddMetric(req)
case DeleteMetricReq:
s.DeleteMetric(req)
case DeleteMeasuresReq:
s.DeleteMeasures(req)
case GetMetricReq:
s.GetMetric(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 RangeScanReq:
metric.StartRangeScan(req)
case FullScanReq:
metric.StartFullScan(req)
case GetMetricReq:
s.GetMetric(req)
case AppendMeasuresReq:
metric.StartAppendMeasures(req, tmp, s.storageInbox)
case DeleteMetricReq:
s.startDeleteMetric(metric, req)
case DeleteMeasuresReq:
s.startDeleteMeasures(metric, req)
default:
qb.Abort(qb.UnknownMetricWaitQueueItemBug,
fmt.Errorf("bug: unknown metric wait queue item type %T", req))
}
}
}
type AddMetricReq struct {
MetricID uint32
ResultCh chan byte
}
func (s *Worker) AddMetric(req AddMetricReq) {
_, 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 []AddMetricReq) {
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 GetMetricResult struct {
MetricType qb.MetricType
FracDigits byte
ResultCode byte
}
type GetMetricReq struct {
MetricID uint32
ResultCh chan GetMetricResult
}
func (s *Worker) GetMetric(req GetMetricReq) {
metric, ok := s.metrics[req.MetricID]
if ok {
req.ResultCh <- GetMetricResult{
ResultCode: Succeed,
MetricType: metric.MetricType(),
FracDigits: metric.FracDigits(),
}
} else {
req.ResultCh <- GetMetricResult{
ResultCode: NoMetric,
}
}
}
type DeleteMetricResult struct {
ResultCode byte
RootPageNo uint32
}
type DeleteMetricReq struct {
MetricID uint32
ResultCh chan DeleteMetricResult
}
func (s *Worker) DeleteMetric(req DeleteMetricReq) {
// 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 DeleteMetricReq) {
// FIX
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
// XLock: true,
// }
// req.ResultCh <- tryDeleteMetricResult{
// ResultCode: Succeed,
// RootPageNo: metric.RootPageNo,
// }
}
type DeleteMeasuresResult struct {
ResultCode byte
RootPageNo uint32
}
type DeleteMeasuresReq struct {
MetricID uint32
Since uint32
ResultCh chan DeleteMeasuresResult
}
func (s *Worker) DeleteMeasures(req DeleteMeasuresReq) {
// FIX
// metric, ok := s.metrics[req.MetricID]
// if !ok {
// 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 DeleteMeasuresReq) {
// 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 AppendMeasuresResult struct {
ResultCode byte
Written int
}
type AppendMeasuresReq struct {
MetricID uint32
Measures []proto.Measure
ResultCh chan AppendMeasuresResult
}
func (s *Worker) AppendMeasures(req AppendMeasuresReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- AppendMeasuresResult{
ResultCode: NoMetric,
}
return
}
if metric.XLock {
metric.WaitQueue = append(metric.WaitQueue, req)
return
}
metric.StartAppendMeasures(req, s.tmp, s.storageInbox)
}
type RangeScanResult struct {
ResultCode byte
FracDigits byte
RootPageNo uint32
LastPageNo uint32
}
type RangeScanReq struct {
MetricID uint32
Since uint32
Until uint32
MetricType qb.MetricType
ResponseWriter atree.WorkerMeasureConsumer
ResultCh chan RangeScanResult
}
func (s *Worker) RangeScan(req RangeScanReq) {
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 FullScanResult struct {
ResultCode byte
FracDigits byte
LastPageNo uint32
}
type FullScanReq struct {
MetricID uint32
MetricType qb.MetricType
ResponseWriter atree.WorkerMeasureConsumer
ResultCh chan FullScanResult
}
func (s *Worker) tryFullScan(req FullScanReq) {
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 ListCurrentValuesReq struct {
MetricIDs []uint32
ResponseWriter *transform.CurrentValueWriter
ResultCh chan struct{}
}
func (s *Worker) tryListCurrentValues(req ListCurrentValuesReq) {
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 []AddMetricReq
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 RangeScanReq:
req.ResultCh <- RangeScanResult{
ResultCode: NoMetric,
}
case FullScanReq:
req.ResultCh <- FullScanResult{
ResultCode: NoMetric,
}
case AddMetricReq:
addMetricReqs = append(addMetricReqs, req)
case DeleteMetricReq:
req.ResultCh <- DeleteMetricResult{
ResultCode: NoMetric,
}
case DeleteMeasuresReq:
req.ResultCh <- DeleteMeasuresResult{
ResultCode: NoMetric,
}
case GetMetricReq:
req.ResultCh <- GetMetricResult{
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)
// }
// }