Files
qb/database/proc.go
2026-06-10 06:18:45 +03:00

590 lines
13 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 database
import (
"fmt"
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
)
func (s *Database) worker() {
for {
select {
case <-s.workerSignalCh:
s.DoWork()
}
}
}
func (s *Database) DoWork() {
s.mutex.Lock()
rLocksToRelease := s.rLocksToRelease
workerQueue := s.workerQueue
s.rLocksToRelease = nil
s.workerQueue = 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 workerQueue {
switch req := untyped.(type) {
case tryAppendMeasuresReq:
s.tryAppendMeasures(req)
case storage.Changes:
s.applyChanges(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 *Database) processMetricQueue(metricID uint32, metric *_metric) {
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, s.storage.Append)
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 *Database) tryAddMetric(req tryAddMetricReq) {
_, ok := s.metrics[req.MetricID]
if ok {
req.ResultCh <- MetricDuplicate
return
}
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 *Database) 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 *Database) 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 *Database) 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 *Database) 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 *Database) 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 *Database) 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 *Database) finAppendMeasures(rec storage.AppendMeasuresSummary) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.NoMetricBug,
fmt.Errorf("finAppendMeasures: metric %d not found",
rec.MetricID))
}
metric.FinAppendMeasures(rec)
}
type tryAppendMeasuresReq struct {
MetricID uint32
Measures []proto.Measure
ResultCh chan tryAppendMeasuresResult
}
func (s *Database) 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.storage.Append)
}
type tryRangeScanReq struct {
MetricID uint32
Since uint32
Until uint32
MetricType qb.MetricType
ResponseWriter atree.WorkerMeasureConsumer
ResultCh chan rangeScanResult
}
func (s *Database) 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 *Database) 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 *Database) tryListCurrentValues(req tryListCurrentValuesReq) {
for _, metricID := range req.MetricIDs {
metric, ok := s.metrics[metricID]
if ok {
req.ResponseWriter.BufferValue(transform.CurrentValue{
MetricID: metricID,
Timestamp: metric.timestamps.LastTimestamp(),
Value: metric.lastValue,
})
}
}
req.ResultCh <- struct{}{}
}
///////////////////////////////////////////////////////
func (s *Database) applyChanges(req storage.Changes) {
for _, untyped := range req.Records {
switch rec := untyped.(type) {
case storage.MetricAddRecord:
s.applyAddMetric(rec)
case storage.MetricDeleteRecord:
s.deleteMetric(rec)
// case storage.AppendedMeasure:
// s.appendMeasure(rec)
case storage.AppendMeasuresSummary:
s.finAppendMeasures(rec)
// case storage.AppendedMeasureWithOverflowExtended:
// s.appendMeasureAfterOverflow(rec)
case storage.DeletedMeasures:
s.deleteMeasures(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 *Database) applyAddMetric(rec storage.MetricAddRecord) {
// 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 *Database) addMetric(rec storage.MetricAddRecord) {
var (
values qb.ValueCompressor
buffer = make([]byte, minBufferSize)
)
if rec.MetricType == qb.Cumulative {
values = enc.NewCumulativeDeltaCompressor(byte(rec.FracDigits), buffer, 0)
} else {
values = enc.NewInstantDeltaCompressor(byte(rec.FracDigits), buffer, 0)
}
s.metrics[rec.MetricID] = &_metric{
metricType: rec.MetricType,
fracDigits: byte(rec.FracDigits),
buffer: buffer,
timestamps: enc.NewTimeDeltaCompressor(buffer, 0),
values: values,
}
}
func (s *Database) deleteMetric(rec storage.MetricDeleteRecord) {
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 *Database) deleteMeasures(rec storage.DeletedMeasures) {
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.DeleteMeasures()
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)
}
}