Files
qb/worker/worker.go
2026-06-19 07:25:46 +03:00

443 lines
11 KiB
Go
Raw Permalink 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/enc"
"gordenko.dev/dima/qb/inbox"
"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
ExpiredMeasure = 8
NonMonotonicValue = 9
WrongMetricType = 11
NoMeasuresToDelete = 12
//DeleteFromAtreeNotNeeded = 13
//DeleteFromAtreeRequired = 14
)
type Worker struct {
inbox *inbox.Inbox
storageInbox *inbox.Inbox
dir string
databaseName 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
DatabaseName string
ExitCh chan struct{}
WaitGroup *sync.WaitGroup
}
func New(opt Options) *Worker {
s := &Worker{
inbox: opt.Inbox,
storageInbox: opt.StorageInbox,
dir: opt.Dir,
databaseName: opt.DatabaseName,
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,
}
//fmt.Printf("index: %#v\n", x.IndexLevelTails)
}
return s
}
type ReleaseRLock struct {
MetricID uint32
}
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()
for _, untyped := range queue {
switch req := untyped.(type) {
case AppendMeasuresReq:
s.appendMeasures(req)
case ReleaseRLock:
s.releaseRLock(req.MetricID)
case storage.Changes:
s.applyCommits(req) // all metrics only
case ListCurrentValuesReq:
s.listCurrentValues(req) // all metrics only
case RangeScanReq:
s.rangeScan(req)
case FullScanReq:
s.fullScan(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))
}
}
}
func (s *Worker) releaseRLock(metricID uint32) {
metric, ok := s.metrics[metricID]
if ok {
metric.ReleaseRLock()
}
}
type AddMetricReq struct {
MetricID uint32
MetricType qb.MetricType
FracDigits byte
ResultCh chan byte
}
func (s *Worker) addMetric(req AddMetricReq) {
_, ok := s.metrics[req.MetricID]
if ok {
req.ResultCh <- MetricDuplicate
return
}
var (
buf = make([]byte, storage.DataPageSize)
databuf = buf[:storage.DataPagePayloadSize]
)
s.metrics[req.MetricID] = &Metric{
metricType: req.MetricType,
fracDigits: req.FracDigits,
buffer: buf,
timestamps: enc.NewTimeDeltaCompressor(databuf, 0),
values: enc.NewValueDeltaCompressor(req.MetricType, req.FracDigits, databuf, 0),
xLock: true,
}
s.storageInbox.Push(storage.MetricAdd{
MetricID: req.MetricID,
MetricType: req.MetricType,
FracDigits: req.FracDigits,
ResultCh: req.ResultCh,
})
}
func (s *Worker) getMetric(req GetMetricReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- GetMetricResult{
ResultCode: NoMetric,
}
}
metric.GetMetric(req)
}
type DeleteMetricReq struct {
MetricID uint32
ResultCh chan byte
}
func (s *Worker) deleteMetric(req DeleteMetricReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- NoMetric
return
}
metric.DeleteMetric(req)
}
func (s *Worker) deleteMeasures(req DeleteMeasuresReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- NoMetric
return
}
metric.DeleteMeasures(req)
}
func (s *Worker) appendMeasures(req AppendMeasuresReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- storage.MeasuresAppendResult{
ResultCode: NoMetric,
}
return
}
if metric.xLock {
metric.waitQueue = append(metric.waitQueue, req)
} else {
metric.AppendMeasures(req, s.tmp, s.storageInbox)
}
}
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)
} else {
metric.RangeScan(req)
}
}
func (s *Worker) fullScan(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)
} else {
metric.FullScan(req)
}
}
type ListCurrentValuesReq struct {
MetricIDs []uint32
ResponseWriter *transform.CurrentValueWriter
ResultCh chan struct{}
}
func (s *Worker) listCurrentValues(req ListCurrentValuesReq) {
for _, metricID := range req.MetricIDs {
metric, ok := s.metrics[metricID]
if ok {
req.ResponseWriter.BufferValue(transform.CurrentValue{
MetricID: metricID,
Timestamp: metric.timestamps.CommitedUntil(),
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(
qb.GetSnapshotFilePath(s.dir, s.databaseName, s.snapshotNumber),
WriteSnapshotIn{
SnapshotNumber: s.snapshotNumber,
WriteBufferSize: 4 * 1024 * 1024, // 4mb
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) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.MetricAddedBug,
fmt.Errorf("onMetricAddCommited: metric %d not found", rec.MetricID))
}
metric.xLock = false
rec.ResultCh <- Succeed
}
func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.NoMetricBug,
fmt.Errorf("onMetricDeleteCommited: metric %d not found", rec.MetricID))
}
if len(metric.waitQueue) > 0 {
for _, untyped := range metric.waitQueue {
switch req := untyped.(type) {
case AppendMeasuresReq:
req.ResultCh <- storage.MeasuresAppendResult{
ResultCode: NoMetric,
WrittenCount: 0,
}
case RangeScanReq:
req.ResultCh <- RangeScanResult{
ResultCode: NoMetric,
}
case FullScanReq:
req.ResultCh <- FullScanResult{
ResultCode: NoMetric,
}
case DeleteMetricReq:
req.ResultCh <- NoMetric
case DeleteMeasuresReq:
req.ResultCh <- 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)
//
rec.ResultCh <- Succeed
}
func (s *Worker) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.NoMetricBug,
fmt.Errorf("onMeasuresAppendCommited: 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("onMeasuresAppendWithGrowCommited: metric %d not found",
rec.MetricID))
}
metric.OnMeasuresAppendWithGrowCommited(rec)
}
func (s *Worker) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.NoMetricBug,
fmt.Errorf("onMeasuresDeleteCommited: metric %d not found", rec.MetricID))
}
metric.OnMeasuresDeleteCommited(rec)
metric.xLock = false
//s.doAfterReleaseXLock(rec.MetricID, metric)
}
// rLock сумісний із append measures,
// xLock ні з чим
// після xLock - вся черга
// після rLock, якщо capturedState == nil - вся черга, оскільки rLock блокує лише xLock задачі
// після capturedState = nil, якщо перша задача - appendMeasures - беру, інакше перевірка rLock
// суть у тому що треба запускати запити, пока не зустріну XLock
func (s *Worker) ProcessQueue(metric *Metric, tmp []byte, storageInbox *inbox.Inbox) {
// for _, untyped := range metric.waitQueue {
// switch req := untyped.(type) {
// case RangeScanReq:
// metric.RangeScan(req)
// case FullScanReq:
// metric.FullScan(req)
// case GetMetricReq:
// metric.GetMetric(req)
// case AppendMeasuresReq:
// metric.AppendMeasures(req, tmp, s.storageInbox)
// case DeleteMetricReq:
// metric.DeleteMetric(req)
// case DeleteMeasuresReq:
// metric.DeleteMeasures(req)
// default:
// qb.Abort(qb.UnknownMetricWaitQueueItemBug,
// fmt.Errorf("bug: unknown metric wait queue item type %T", req))
// }
// }
}