wp
This commit is contained in:
@@ -53,7 +53,10 @@ func (s *_metric) LastTimestamp() uint32 {
|
||||
return s.timestamps.LastTimestamp()
|
||||
}
|
||||
|
||||
func (s *_metric) DeleteMeasures() {
|
||||
func (s *_metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
|
||||
//s.timestamps.Reset()
|
||||
//s.values.Reset()
|
||||
s.XLock = false
|
||||
// s.Timestamps.Renew()
|
||||
// s.Values.Renew()
|
||||
|
||||
@@ -61,7 +64,9 @@ func (s *_metric) DeleteMeasures() {
|
||||
// s.Since = 0
|
||||
// s.SinceValue = 0
|
||||
// s.Until = 0
|
||||
// s.UntilValue = 0
|
||||
s.indexLevelTails = nil
|
||||
s.lastPageNo = 0
|
||||
s.lastValue = 0
|
||||
}
|
||||
|
||||
// func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) {
|
||||
@@ -197,12 +202,20 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send
|
||||
//ResultCh: req.ResultCh,
|
||||
})
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func (s *_metric) FinAppendMeasures(rec storage.MeasuresAppendWithGrowCommited) {
|
||||
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
|
||||
|
||||
@@ -84,7 +84,7 @@ func (s *Database) DoWork() {
|
||||
s.tryAppendMeasures(req)
|
||||
|
||||
case storage.Changes:
|
||||
s.applyChanges(req) // all metrics only
|
||||
s.applyCommits(req) // all metrics only
|
||||
|
||||
case tryListCurrentValuesReq:
|
||||
s.tryListCurrentValues(req) // all metrics only
|
||||
@@ -294,14 +294,24 @@ func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq
|
||||
// }
|
||||
}
|
||||
|
||||
func (s *Database) finAppendMeasures(rec storage.MeasuresAppendWithGrowCommited) {
|
||||
func (s *Database) 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.FinAppendMeasures(rec)
|
||||
metric.OnMeasuresAppendCommited(rec)
|
||||
}
|
||||
|
||||
func (s *Database) 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 {
|
||||
@@ -408,26 +418,23 @@ func (s *Database) tryListCurrentValues(req tryListCurrentValuesReq) {
|
||||
|
||||
///////////////////////////////////////////////////////
|
||||
|
||||
func (s *Database) applyChanges(req storage.Changes) {
|
||||
for _, untyped := range req.Records {
|
||||
func (s *Database) applyCommits(req storage.Changes) {
|
||||
for _, untyped := range req.Commits {
|
||||
switch rec := untyped.(type) {
|
||||
case storage.MetricAddRecord:
|
||||
s.applyAddMetric(rec)
|
||||
case storage.MetricAddCommited:
|
||||
s.onMetricAddCommited(rec)
|
||||
|
||||
case storage.MetricDeleteRecord:
|
||||
s.deleteMetric(rec)
|
||||
case storage.MetricDeleteCommited:
|
||||
s.onMetricDeleteCommited(rec)
|
||||
|
||||
// case storage.AppendedMeasure:
|
||||
// s.appendMeasure(rec)
|
||||
case storage.MeasuresAppendCommited:
|
||||
s.onMeasuresAppendCommited(rec)
|
||||
|
||||
case storage.MeasuresAppendWithGrowCommited:
|
||||
s.finAppendMeasures(rec)
|
||||
s.onMeasuresAppendWithGrowCommited(rec)
|
||||
|
||||
// case storage.AppendedMeasureWithOverflowExtended:
|
||||
// s.appendMeasureAfterOverflow(rec)
|
||||
|
||||
case storage.MeasuresDeleteRecord:
|
||||
s.deleteMeasures(rec)
|
||||
case storage.MeasuresDeleteCommited:
|
||||
s.onMeasuresDeleteCommited(rec)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -451,7 +458,7 @@ func (s *Database) applyChanges(req storage.Changes) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Database) applyAddMetric(rec storage.MetricAddRecord) {
|
||||
func (s *Database) onMetricAddCommited(rec storage.MetricAddCommited) {
|
||||
// fix lock
|
||||
_, ok := s.metrics[rec.MetricID]
|
||||
if ok {
|
||||
@@ -484,7 +491,7 @@ func (s *Database) addMetric(rec storage.MetricAddRecord) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Database) deleteMetric(rec storage.MetricDeleteRecord) {
|
||||
func (s *Database) onMetricDeleteCommited(rec storage.MetricDeleteCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
@@ -555,7 +562,7 @@ func (s *Database) deleteMetric(rec storage.MetricDeleteRecord) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Database) deleteMeasures(rec storage.MeasuresDeleteRecord) {
|
||||
func (s *Database) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
|
||||
metric, ok := s.metrics[rec.MetricID]
|
||||
if !ok {
|
||||
qb.Abort(qb.NoMetricBug,
|
||||
@@ -568,7 +575,7 @@ func (s *Database) deleteMeasures(rec storage.MeasuresDeleteRecord) {
|
||||
fmt.Errorf("deleteMeasures: xlock not set for the metric %d",
|
||||
rec.MetricID))
|
||||
}
|
||||
metric.DeleteMeasures()
|
||||
metric.OnMeasuresDeleteCommited(rec)
|
||||
metric.XLock = false
|
||||
// FIX add in storage
|
||||
// if len(rec.FreePageNumbers) > 0 {
|
||||
|
||||
@@ -199,6 +199,10 @@ func (s *TimeDeltaCompressor) CaptureState() {
|
||||
s.state = &state
|
||||
}
|
||||
|
||||
func (s *TimeDeltaCompressor) ForgetCapturedState() {
|
||||
s.state = nil
|
||||
}
|
||||
|
||||
func (s *TimeDeltaCompressor) Tail(offset int) []byte {
|
||||
return s.buf[s.pos : len(s.buf)-offset]
|
||||
}
|
||||
|
||||
@@ -172,6 +172,10 @@ func (s *ValueDeltaCompressor) CaptureState() {
|
||||
s.state = &state
|
||||
}
|
||||
|
||||
func (s *ValueDeltaCompressor) ForgetCapturedState() {
|
||||
s.state = nil
|
||||
}
|
||||
|
||||
func (s *ValueDeltaCompressor) Tail(offset int) []byte {
|
||||
return s.buf[offset:s.pos]
|
||||
}
|
||||
|
||||
2
qb.go
2
qb.go
@@ -32,6 +32,7 @@ type TimestampCompressor interface {
|
||||
StoredSize() int
|
||||
//DeleteLast()
|
||||
CaptureState()
|
||||
ForgetCapturedState()
|
||||
// (offset) => payload
|
||||
Tail(int) []byte
|
||||
CreateDecompressor() TimestampDecompressor
|
||||
@@ -67,6 +68,7 @@ type ValueCompressor interface {
|
||||
//Chunks() [][]byte
|
||||
//DeleteLast()
|
||||
CaptureState()
|
||||
ForgetCapturedState()
|
||||
// (offset) => payload
|
||||
Tail(int) []byte
|
||||
// fracDigits
|
||||
|
||||
@@ -228,6 +228,8 @@ type appendIndexRecordIn struct {
|
||||
}
|
||||
|
||||
func appendIndexRecord(in appendIndexRecordIn) []byte {
|
||||
//fmt.Printf("appendIndexRecord: % x\n", in.Records)
|
||||
//fmt.Printf("timestamp: %d, pageNo: %d\n", in.Timestamp, in.PageNo)
|
||||
if in.RecordsCount < maxRecordsOnIndexPage {
|
||||
var (
|
||||
pos int
|
||||
@@ -235,11 +237,13 @@ func appendIndexRecord(in appendIndexRecordIn) []byte {
|
||||
)
|
||||
if in.RecordsCount > 0 {
|
||||
pos = in.RecordsCount * indexRecordSize
|
||||
//fmt.Printf("pos: %d\n", pos)
|
||||
} else {
|
||||
buf = make([]byte, IndexPageSize) // IndexPageIncSize
|
||||
}
|
||||
bin.PutUint32(buf[pos:], in.Timestamp)
|
||||
bin.PutUint32(buf[pos+4:], in.PageNo)
|
||||
//fmt.Printf("appendIndexRecord after: % x\n", buf)
|
||||
return buf
|
||||
}
|
||||
qb.Abort(qb.NoSpaceOnIndexPage, nil)
|
||||
|
||||
@@ -7,8 +7,8 @@ var (
|
||||
//PageNoSize = 4
|
||||
|
||||
// data page
|
||||
DataPageSize = 8192
|
||||
DataPageFooterSize int = 12
|
||||
DataPageSize = 8192
|
||||
|
||||
DataPagePayloadSize int = DataPageSize - DataPageFooterSize
|
||||
|
||||
dataCRC32Idx = DataPageSize - 4
|
||||
@@ -17,13 +17,13 @@ var (
|
||||
prevPageIdx = DataPageSize - 12
|
||||
|
||||
// index page
|
||||
IndexPageSize = 1024
|
||||
IndexPageFooterSize = 7
|
||||
IndexPageSize = 1024
|
||||
|
||||
//indexPageIncSize = IndexPageIncSize
|
||||
indexCRC32Idx = IndexPageSize - 4
|
||||
indexRecordsCountIdx = IndexPageSize - 6
|
||||
isZeroLevelIdx = IndexPageSize - 7
|
||||
indexRecordSize = 8
|
||||
indexCRC32Idx = IndexPageSize - 4
|
||||
indexRecordsCountIdx = IndexPageSize - 6
|
||||
isZeroLevelIdx = IndexPageSize - 7
|
||||
|
||||
maxRecordsOnIndexPage = (IndexPageSize - IndexPageFooterSize) / indexRecordSize
|
||||
|
||||
// timestampSize = 4
|
||||
@@ -32,6 +32,12 @@ var (
|
||||
// dataFooterIdx = timestampsSizeIdx
|
||||
)
|
||||
|
||||
const (
|
||||
DataPageFooterSize = 12
|
||||
IndexPageFooterSize = 7
|
||||
indexRecordSize = 8
|
||||
)
|
||||
|
||||
type IndexRecord struct {
|
||||
Timestamp uint32
|
||||
PageNo uint32
|
||||
|
||||
@@ -623,7 +623,6 @@ func TestMeasuresAppendRecord(t *testing.T) {
|
||||
func TestMeasuresAppendWithGrowRecord(t *testing.T) {
|
||||
DataPageSize = 8
|
||||
IndexPageSize = 4
|
||||
indexRecordSize = 2
|
||||
//
|
||||
rec := MeasuresAppendWithGrowRecord{
|
||||
MetricID: 12345,
|
||||
@@ -735,8 +734,7 @@ func TestWritePreparer(t *testing.T) {
|
||||
prevPageIdx = DataPageSize - 12
|
||||
|
||||
// index page
|
||||
IndexPageSize = 16
|
||||
indexRecordSize = 4 // 2 records per index page
|
||||
IndexPageSize = 24
|
||||
maxRecordsOnIndexPage = 2
|
||||
indexCRC32Idx = IndexPageSize - 4
|
||||
indexRecordsCountIdx = IndexPageSize - 6
|
||||
@@ -835,9 +833,9 @@ func TestWritePreparer(t *testing.T) {
|
||||
IndexLevelTails: []IndexLevelTail{
|
||||
{
|
||||
Buffer: []byte{
|
||||
0x01, 0x01, // since
|
||||
0x0a, 0x00, // pageNo 10
|
||||
0x00, 0x00, 0x00, 0x00, 0x00, // empty space
|
||||
0x01, 0x01, 0x00, 0x00, // since
|
||||
0x0a, 0x00, 0x00, 0x00, // pageNo 10
|
||||
0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // empty space
|
||||
0x00, // zero level
|
||||
0x00, 0x00, // records count
|
||||
0x00, 0x00, 0x00, 0x00, // crc32
|
||||
|
||||
@@ -53,7 +53,8 @@ type MetricsState struct {
|
||||
}
|
||||
|
||||
type Changes struct {
|
||||
Records []any
|
||||
Commits []any
|
||||
// for snapshot only
|
||||
SnapshotNumberCh chan int // log number
|
||||
FrozenIndexPagesCount int
|
||||
IndexPageNumbers []uint32
|
||||
@@ -62,13 +63,12 @@ type Changes struct {
|
||||
}
|
||||
|
||||
type Writer struct {
|
||||
mutex sync.Mutex
|
||||
dataPagesCount uint32
|
||||
indexPagesCount uint32
|
||||
dataFreeList *freelist.FreeList
|
||||
indexFreeList *freelist.FreeList
|
||||
atree *atree.Atree
|
||||
//logNumber int
|
||||
mutex sync.Mutex
|
||||
dataPagesCount uint32
|
||||
indexPagesCount uint32
|
||||
dataFreeList *freelist.FreeList
|
||||
indexFreeList *freelist.FreeList
|
||||
atree *atree.Atree
|
||||
dir string
|
||||
w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск
|
||||
wal *os.File
|
||||
@@ -79,12 +79,11 @@ type Writer struct {
|
||||
input []any
|
||||
dataPreparer *WritePreparer
|
||||
appendToWorkerQueue func(any)
|
||||
//lsn uint32
|
||||
written int64
|
||||
isExited bool
|
||||
exitCh chan struct{}
|
||||
waitGroup *sync.WaitGroup
|
||||
signalCh chan struct{}
|
||||
written int64
|
||||
isExited bool
|
||||
exitCh chan struct{}
|
||||
waitGroup *sync.WaitGroup
|
||||
signalCh chan struct{}
|
||||
}
|
||||
|
||||
type WriterOptions struct {
|
||||
@@ -259,7 +258,7 @@ func (s *Writer) packAndWrite() (err error) {
|
||||
snapshotNumberCh := make(chan int, 1)
|
||||
|
||||
s.appendToWorkerQueue(Changes{
|
||||
Records: prepared.Commits,
|
||||
Commits: prepared.Commits,
|
||||
SnapshotNumberCh: snapshotNumberCh,
|
||||
FrozenIndexPagesCount: s.indexFreeList.Pages(),
|
||||
IndexPageNumbers: s.indexFreeList.Cached(),
|
||||
@@ -290,7 +289,7 @@ func (s *Writer) packAndWrite() (err error) {
|
||||
}
|
||||
} else {
|
||||
s.appendToWorkerQueue(Changes{
|
||||
Records: prepared.Commits,
|
||||
Commits: prepared.Commits,
|
||||
})
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -115,7 +115,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
|
||||
}
|
||||
for _, level := range sealedIndexLevels {
|
||||
for _, p := range level.IndexPages {
|
||||
fmt.Printf("records: %v\n", getIndexRecords(p.Content, 2))
|
||||
//fmt.Printf("records: %v\n", getIndexRecords(p.Content, 2))
|
||||
s.indexPages = append(s.indexPages, PageToWrite{
|
||||
PageNo: p.PageNo,
|
||||
Content: p.Content,
|
||||
@@ -139,10 +139,12 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
|
||||
)
|
||||
if len(level.IndexPages) > 0 {
|
||||
first := level.IndexPages[0]
|
||||
firstCRC32, _ := bin.GetUint32(first.Content[indexCRC32Idx:])
|
||||
skipSize := level.SkipRecords * indexRecordSize
|
||||
indexPageTail = &IndexPageTail{
|
||||
PageNo: first.PageNo,
|
||||
Reused: first.Reused,
|
||||
CRC32: firstCRC32,
|
||||
// в WAL файл попадає лише payload
|
||||
Records: first.Content[skipSize : maxRecordsOnIndexPage*indexRecordSize],
|
||||
}
|
||||
@@ -181,11 +183,13 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
|
||||
Values: x.Values,
|
||||
},
|
||||
DataPages: sealedDataPages[1:],
|
||||
TailTimestamps: x.TailTimestamps,
|
||||
TailValues: x.TailValues,
|
||||
ChangedIndexLevels: changedIndexPages,
|
||||
TailTimestamps: x.Timestamps,
|
||||
TailValues: x.Values,
|
||||
}
|
||||
|
||||
fmt.Printf("%#v\ns", rec.ChangedIndexLevels[0].IndexPageTail)
|
||||
|
||||
rec.Pack(w)
|
||||
// Дані для Worker
|
||||
s.commits = append(s.commits, MeasuresAppendWithGrowCommited{
|
||||
|
||||
Reference in New Issue
Block a user