From 5778782737d599669c3296d8132b10985ddf25d5 Mon Sep 17 00:00:00 2001 From: dima <1.e4.kc6@gmail.com> Date: Sat, 13 Jun 2026 07:13:03 +0000 Subject: [PATCH] wp --- database/metric.go | 21 ++++++++++++++++---- database/proc.go | 49 +++++++++++++++++++++++++++-------------------- enc/time_delta.go | 4 ++++ enc/value_delta.go | 4 ++++ qb.go | 2 ++ storage/page_ingester.go | 4 ++++ storage/storage.go | 22 +++++++++++++-------- storage/storage_test.go | 10 ++++------ storage/wal_writer.go | 31 +++++++++++++++--------------- storage/write_preparer.go | 10 +++++++--- 10 files changed, 99 insertions(+), 58 deletions(-) diff --git a/database/metric.go b/database/metric.go index f73cd12..d3b9c4d 100644 --- a/database/metric.go +++ b/database/metric.go @@ -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 diff --git a/database/proc.go b/database/proc.go index 42d114c..bb68369 100644 --- a/database/proc.go +++ b/database/proc.go @@ -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 { diff --git a/enc/time_delta.go b/enc/time_delta.go index 77eb208..57010fe 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -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] } diff --git a/enc/value_delta.go b/enc/value_delta.go index 53c1471..6e116a9 100644 --- a/enc/value_delta.go +++ b/enc/value_delta.go @@ -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] } diff --git a/qb.go b/qb.go index 8394a00..32a0acb 100644 --- a/qb.go +++ b/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 diff --git a/storage/page_ingester.go b/storage/page_ingester.go index 0b61a63..7eced42 100644 --- a/storage/page_ingester.go +++ b/storage/page_ingester.go @@ -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) diff --git a/storage/storage.go b/storage/storage.go index 2d5b7ea..3e9ffb2 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -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 diff --git a/storage/storage_test.go b/storage/storage_test.go index 71c10ba..c8adb61 100644 --- a/storage/storage_test.go +++ b/storage/storage_test.go @@ -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 diff --git a/storage/wal_writer.go b/storage/wal_writer.go index bb51301..5ff0712 100644 --- a/storage/wal_writer.go +++ b/storage/wal_writer.go @@ -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 diff --git a/storage/write_preparer.go b/storage/write_preparer.go index 3d21c2e..cca7020 100644 --- a/storage/write_preparer.go +++ b/storage/write_preparer.go @@ -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{