From 05eea3b27986cf4bb2ff4259574b776a322cfe0e Mon Sep 17 00:00:00 2001 From: dima <1.e4.kc6@gmail.com> Date: Fri, 5 Jun 2026 19:43:01 +0000 Subject: [PATCH] tmp --- atree/x.go | 163 --------------------- chunkenc/chunkenc.go | 5 +- chunkenc/chunkenc_test.go | 38 +++-- chunkenc/cumdelta.go | 60 ++++---- chunkenc/insdelta.go | 59 ++++---- chunkenc/time_delta.go | 70 +++++---- database/metric.go | 84 +++++------ go.mod | 2 +- go.sum | 2 + qb.go | 1 + txlog/data_preparer.go | 207 ++++++++++++++++++++++++++ txlog/helpers.go | 360 +++++++++++++++++++++------------------------- txlog/page_preparer.go | 279 +++++++++++++++++++++++++++++++++++ txlog/txlog.go | 65 +-------- txlog/wal_records.go | 306 +++++++++++++++++++++++++++++++++++++++ txlog/writer.go | 210 +++++++-------------------- 16 files changed, 1187 insertions(+), 724 deletions(-) delete mode 100644 atree/x.go create mode 100644 txlog/data_preparer.go create mode 100644 txlog/page_preparer.go create mode 100644 txlog/wal_records.go diff --git a/atree/x.go b/atree/x.go deleted file mode 100644 index 092a4fa..0000000 --- a/atree/x.go +++ /dev/null @@ -1,163 +0,0 @@ -package atree - -import ( - "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/bin" - "gordenko.dev/dima/qb/util" -) - -const ( - payloadSize = 24 - recordSize = 4 + 4 // timestamp + pageNo -) - -type DataPage struct { - PrevPageNo uint32 - LowerTimestamp uint32 - Timestamps [][]byte // chunks - TimestampsSize int - Values [][]byte // chunks - ValuesSize int - PageNo uint32 - IsReused bool - //Checksum uint32 -} - -type IndexPage struct { - LowerTimestamp uint32 - Data []byte - PageNo uint32 - IsReused bool -} -type IndexLevel struct { - // вже записані дані (лише для першої в списку індексної сторінки, беремо із _metric) - // дані потрібні, якщо сторінка буде заповнена і піде на запис в .index файл - Offset int // (заповнені одразу) - LowerTimestamp uint32 // (заповнені одразу) - Data []byte // поточні дані index сторінки (заповнені одразу) - Filled []*IndexPage // пусто -} - -func appendIndexRecord(data []byte, timestamp uint32, pageNo uint32) []byte { - var tmp = []byte{ - 0, 0, 0, 0, // timestamp - 0, 0, 0, 0, // pageNo - } - bin.PutUint32(tmp, timestamp) - bin.PutUint32(tmp[4:], pageNo) - return append(data, tmp...) -} - -type ChunksToDataPageReq struct { - PrevPageNo uint32 - Timestamps [][]byte // chunks - TimestampsSize int - Values [][]byte // chunks - ValuesSize int -} - -// return checksum -func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) { - var ( - remainingSize = int(req.TimestampsSize) - pos = 0 - ) - for _, chunk := range req.Timestamps { - if remainingSize >= len(chunk) { - copy(buf[pos:], chunk) - remainingSize -= len(chunk) - pos += len(chunk) - } else { - copy(buf[pos:], chunk[:remainingSize]) - break - } - } - - remainingSize = int(req.ValuesSize) - pos = int(req.TimestampsSize) - - for _, chunk := range req.Values { - if remainingSize >= len(chunk) { - copy(buf[pos:], chunk) - remainingSize -= len(chunk) - pos += len(chunk) - } else { - copy(buf[pos:], chunk[:remainingSize]) - break - } - } - bin.PutUint16(buf[timestampsSizeIdx:], uint16(req.TimestampsSize)) - bin.PutUint16(buf[valuesSizeIdx:], uint16(req.ValuesSize)) - bin.PutUint32(buf[prevPageIdx:], req.PrevPageNo) - - checksum := util.CalcChecksum(buf[:dataCRC32Idx]) - bin.PutUint32(buf[dataCRC32Idx:], checksum) -} - -func DataToIndexPage(buf []byte, data []byte, lastLevel bool) { - copy(buf, data) - bin.PutUint16(buf[indexRecordsQtyIdx:], uint16(len(data)/indexRecordSize)) - if lastLevel { - buf[isLastLevelIdx] = 1 - } - checksum := util.CalcChecksum(buf[:indexCRC32Idx]) - bin.PutUint32(buf[indexCRC32Idx:], checksum) -} - -// fix - reduce leves -func AppendDataPagesToTree(getDataPageNumber func() (uint32, bool, error), getIndexPageNumber func() (uint32, bool, error), levels []*IndexLevel, dataPages []*DataPage) { - for i, d := range dataPages[1:] { - pageNo, isReused, err := getDataPageNumber() - if err != nil { - qb.Abort(qb.FailedGetPageNumber, err) - } - d.PageNo = pageNo - d.IsReused = isReused - d.PrevPageNo = dataPages[i-1].PageNo - } - - for _, d := range dataPages { - var ( - upTimestamp = d.LowerTimestamp - upPageNo = d.PageNo - levelIdx = 0 - ) - for { - if levelIdx == len(levels) { - // encode prev page record + up* - levels = append(levels, &IndexLevel{ - LowerTimestamp: upTimestamp, - Data: appendIndexRecord(nil, upTimestamp, upPageNo), - }) - break - } - level := levels[levelIdx] - if len(level.Data)+recordSize <= payloadSize { - level.Data = appendIndexRecord(level.Data, upTimestamp, upPageNo) - break - } - // в буфері немає місця для додавання нової пари, отже це - // заповнена сторінка - filled := &IndexPage{ - LowerTimestamp: level.LowerTimestamp, - Data: level.Data, - } - pageNo, isReused, err := getIndexPageNumber() - if err != nil { - qb.Abort(qb.FailedGetPageNumber, err) - } - filled.PageNo = pageNo - filled.IsReused = isReused - - level.Filled = append(level.Filled, filled) - // - level.LowerTimestamp = upTimestamp - level.Data = appendIndexRecord(nil, upTimestamp, upPageNo) - - // - upPageNo = filled.PageNo - upTimestamp = filled.LowerTimestamp - levelIdx++ - } - } -} diff --git a/chunkenc/chunkenc.go b/chunkenc/chunkenc.go index 0432cca..d6192d8 100644 --- a/chunkenc/chunkenc.go +++ b/chunkenc/chunkenc.go @@ -2,4 +2,7 @@ package chunkenc const eps = 0.000001 -const flagLiteral = 128 +const ( + flagLiteral = 128 + minBufferSize = 1024 +) diff --git a/chunkenc/chunkenc_test.go b/chunkenc/chunkenc_test.go index d8fb5cb..c2bdd44 100644 --- a/chunkenc/chunkenc_test.go +++ b/chunkenc/chunkenc_test.go @@ -3,14 +3,12 @@ package chunkenc import ( "fmt" "testing" - - "gordenko.dev/dima/qb/conbuf" ) func TestCumdelta(t *testing.T) { var ( fracDigits byte = 0 - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value float64 done bool ) @@ -41,7 +39,7 @@ func TestCumdelta(t *testing.T) { func TestInsdelta(t *testing.T) { var ( fracDigits byte = 2 - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value float64 done bool ) @@ -70,7 +68,7 @@ func TestInsdelta(t *testing.T) { func TestTimeDelta(t *testing.T) { var ( - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value uint32 done bool ) @@ -120,7 +118,7 @@ func TestTimeDelta(t *testing.T) { func TestCumdeltaBound(t *testing.T) { var ( fracDigits byte = 0 - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value float64 done bool ) @@ -181,7 +179,7 @@ func TestCumdeltaBound(t *testing.T) { func TestInsdeltaBound(t *testing.T) { var ( fracDigits byte = 0 - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value float64 done bool ) @@ -241,7 +239,7 @@ func TestInsdeltaBound(t *testing.T) { func TestInsdeltaBound2(t *testing.T) { var ( fracDigits byte = 0 - buf = conbuf.New(nil) + buf = make([]byte, minBufferSize) value float64 done bool ) @@ -297,3 +295,27 @@ func TestInsdeltaBound2(t *testing.T) { fmt.Println(value, done) } } + +func TestCap(t *testing.T) { + a := make([]byte, 0, 10) + fmt.Println("len:", len(a)) + fmt.Println("cap:", cap(a)) + + a = append(a, 1) + fmt.Println("after append:") + fmt.Println("len:", len(a)) + fmt.Println("cap:", cap(a)) + fmt.Println("a:", a) + + a[1] = 55 + fmt.Println("after []:") + fmt.Println("len:", len(a)) + fmt.Println("cap:", cap(a)) + fmt.Println("a:", a) + + a = append(a, 2) + fmt.Println("after 2nd append:") + fmt.Println("len:", len(a)) + fmt.Println("cap:", cap(a)) + fmt.Println("a:", a) +} diff --git a/chunkenc/cumdelta.go b/chunkenc/cumdelta.go index 74717ce..0f8e8b9 100644 --- a/chunkenc/cumdelta.go +++ b/chunkenc/cumdelta.go @@ -5,6 +5,7 @@ import ( "log" "math" + bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" ) @@ -38,7 +39,7 @@ v1 v2 h-byte(literal, 2) v3 h-byte(run, 2) */ type ReverseCumulativeDeltaCompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte coef float64 pos int baseValue float64 @@ -48,7 +49,7 @@ type ReverseCumulativeDeltaCompressor struct { state *CumulativeDeltaBound } -func NewReverseCumulativeDeltaCompressor(buf *conbuf.ContinuousBuffer, size int, fracDigits byte) *ReverseCumulativeDeltaCompressor { +func NewReverseCumulativeDeltaCompressor(buf []byte, size int, fracDigits byte) *ReverseCumulativeDeltaCompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) @@ -59,13 +60,13 @@ func NewReverseCumulativeDeltaCompressor(buf *conbuf.ContinuousBuffer, size int, coef: coef, } if size > 0 { - u64, _, err := s.buf.GetVarUint64(0) + u64, _, err := bin.GetVarUint64(s.buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } s.baseValue = float64(u64) / s.coef - s.h = s.buf.GetByte(s.pos - 1) - s.lastDelta, s.lastDeltaSize, err = s.buf.ReverseGetVarUint64(s.pos - 2) + s.h = s.buf[s.pos-1] + s.lastDelta, s.lastDeltaSize, err = bin.ReverseGetVarUint64(s.buf[:s.pos-2]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } @@ -84,7 +85,8 @@ func (s *ReverseCumulativeDeltaCompressor) CalcRequiredSpace(value float64) int func (s *ReverseCumulativeDeltaCompressor) Append(value float64) { if s.pos == 0 { // base value - s.pos += s.buf.PutVarUint64(s.pos, uint64(value*s.coef)) + n, _ := bin.PutVarUint64(s.buf[s.pos:], uint64(value*s.coef)) + s.pos += n s.baseValue = value s.appendNewLiteral(0) } else { @@ -95,7 +97,7 @@ func (s *ReverseCumulativeDeltaCompressor) Append(value float64) { if s.h < 127 { // increase counter s.h++ - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } else { // не можу збільшити - буде переповнення. Додаю новий literal блок // counter overflow @@ -134,37 +136,37 @@ func (s *ReverseCumulativeDeltaCompressor) convertLastFromLiteralToRun() { // Зменшую кількість елементів в literal блоці s.h-- s.pos -= 1 + s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) // закриваю literal блок + s.buf[s.pos] = s.h // закриваю literal блок s.pos++ - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, s.lastDelta) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], s.lastDelta) s.pos += s.lastDeltaSize s.h = 0 // run блок, довжини 2 - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } func (s *ReverseCumulativeDeltaCompressor) convertLiteralToRun() { // Знімаю flagLiteral, а лічильник 0 дорівнює 2 елементам в серії. s.h = 0 - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } func (s *ReverseCumulativeDeltaCompressor) appendDeltaToLiteral(delta uint64) { s.h++ // збільшую к-сть дельт s.lastDelta = delta s.pos-- - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, delta) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], delta) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } func (s *ReverseCumulativeDeltaCompressor) appendNewLiteral(delta uint64) { s.h = flagLiteral s.lastDelta = delta - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, delta) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], delta) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, flagLiteral) // literal, length = 1 + s.buf[s.pos] = flagLiteral // literal, length = 1 s.pos++ } @@ -176,7 +178,7 @@ type CumulativeDeltaBound struct { Pos int H byte LastDelta uint64 - Chunks [][]byte + Chunks []byte } // delta h @@ -186,15 +188,11 @@ func (s *ReverseCumulativeDeltaCompressor) Lock() { } // позиція посувається вліво, отже може перескочити на попередній chunk pos := s.pos - 1 - s.lastDeltaSize - chunksQty := pos / conbuf.ChunkSize - if (pos % conbuf.ChunkSize) > 0 { - chunksQty++ - } s.state = &CumulativeDeltaBound{ Pos: pos, H: s.h, LastDelta: s.lastDelta, - Chunks: s.buf.Chunks()[:chunksQty], + Chunks: s.buf[:s.pos], // fix check ? } } @@ -210,9 +208,9 @@ func (s *ReverseCumulativeDeltaCompressor) Offset() int { return 0 } -func (s *ReverseCumulativeDeltaCompressor) Snapshot() ([][]byte, int) { +func (s *ReverseCumulativeDeltaCompressor) Snapshot() ([]byte, int) { if s.state == nil { - return s.buf.Chunks(), s.Size() + return s.buf, s.Size() } // ВАЖЛИВО! // Треба відтворити стан останнього чанка @@ -246,7 +244,7 @@ func (s *ReverseCumulativeDeltaCompressor) CreateDecompressor(fracDigits byte) q func (s *ReverseCumulativeDeltaCompressor) Renew() { // УВАГА! // state не чіпаємо - s.buf = conbuf.New(nil) + s.buf = make([]byte, minBufferSize) s.pos = 0 // s.baseValue = 0 @@ -255,14 +253,14 @@ func (s *ReverseCumulativeDeltaCompressor) Renew() { s.h = 0 } -func (s *ReverseCumulativeDeltaCompressor) Chunks() [][]byte { - return s.buf.Chunks() +func (s *ReverseCumulativeDeltaCompressor) Chunks() []byte { + return s.buf } // DECOMPRESSOR type ReverseCumulativeDeltaDecompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte coef float64 pos int bound int @@ -273,12 +271,12 @@ type ReverseCumulativeDeltaDecompressor struct { done bool } -func NewReverseCumulativeDeltaDecompressor(buf *conbuf.ContinuousBuffer, size int, fracDigits byte) *ReverseCumulativeDeltaDecompressor { +func NewReverseCumulativeDeltaDecompressor(buf []byte, size int, fracDigits byte) *ReverseCumulativeDeltaDecompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) } - u64, n, err := buf.GetVarUint64(0) + u64, n, err := bin.GetVarUint64(buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } @@ -335,7 +333,7 @@ func (s *ReverseCumulativeDeltaDecompressor) NextValue() (value float64, done bo } func (s *ReverseCumulativeDeltaDecompressor) readHeader() { - h := s.buf.GetByte(s.pos) + h := s.buf[s.pos] s.pos-- s.decodeHeaderByte(h) // fmt.Println("h:", h) @@ -343,7 +341,7 @@ func (s *ReverseCumulativeDeltaDecompressor) readHeader() { // fmt.Println("pending:", s.pending) } func (s *ReverseCumulativeDeltaDecompressor) readValue() { - u64, n, err := s.buf.ReverseGetVarUint64(s.pos) + u64, n, err := bin.ReverseGetVarUint64(s.buf[:s.pos]) if err != nil { log.Fatalln(err) } diff --git a/chunkenc/insdelta.go b/chunkenc/insdelta.go index 1c3dee2..459963b 100644 --- a/chunkenc/insdelta.go +++ b/chunkenc/insdelta.go @@ -6,11 +6,12 @@ import ( "math" "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/conbuf" ) type ReverseInstantDeltaCompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte coef float64 pos int baseValue float64 @@ -20,7 +21,7 @@ type ReverseInstantDeltaCompressor struct { state *InstantDeltaBound } -func NewReverseInstantDeltaCompressor(buf *conbuf.ContinuousBuffer, size int, fracDigits byte) *ReverseInstantDeltaCompressor { +func NewReverseInstantDeltaCompressor(buf []byte, size int, fracDigits byte) *ReverseInstantDeltaCompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) @@ -31,13 +32,13 @@ func NewReverseInstantDeltaCompressor(buf *conbuf.ContinuousBuffer, size int, fr coef: coef, } if size > 0 { - i64, _, err := s.buf.GetVarInt64(0) + i64, _, err := bin.GetVarInt64(s.buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } s.baseValue = float64(i64) / s.coef - s.h = s.buf.GetByte(s.pos - 1) - s.lastDelta, s.lastDeltaSize, err = s.buf.ReverseGetVarInt64(s.pos - 2) + s.h = s.buf[s.pos-1] + s.lastDelta, s.lastDeltaSize, err = bin.ReverseGetVarInt64(s.buf[:s.pos-2]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } @@ -52,7 +53,7 @@ func (s *ReverseInstantDeltaCompressor) Size() int { func (s *ReverseInstantDeltaCompressor) Append(value float64) { if s.pos == 0 { // base value - n := s.buf.PutVarInt64(s.pos, int64(value*s.coef)) + n, _ := bin.PutVarInt64(s.buf[s.pos:], int64(value*s.coef)) s.pos += n s.baseValue = value s.appendNewLiteral(0) @@ -70,7 +71,7 @@ func (s *ReverseInstantDeltaCompressor) Append(value float64) { if s.h < 127 { // increase counter s.h++ - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } else { // не можу збільшити - буде переповнення. Додаю новий literal блок // counter overflow @@ -109,37 +110,37 @@ func (s *ReverseInstantDeltaCompressor) convertLastFromLiteralToRun() { // Зменшую кількість елементів в literal блоці s.h-- s.pos -= 1 + s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) // закриваю literal блок + s.buf[s.pos] = s.h // закриваю literal блок s.pos++ - s.lastDeltaSize = s.buf.ReversePutVarInt64(s.pos, s.lastDelta) + s.lastDeltaSize, _ = bin.ReversePutVarInt64(s.buf[:s.pos], s.lastDelta) s.pos += s.lastDeltaSize s.h = 0 // run блок, довжини 2 - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } func (s *ReverseInstantDeltaCompressor) convertLiteralToRun() { // Знімаю flagLiteral, а лічильник 0 дорівнює 2 елементам в серії. s.h = 0 - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } func (s *ReverseInstantDeltaCompressor) appendDeltaToLiteral(delta int64) { s.h++ // збільшую к-сть дельт s.lastDelta = delta s.pos-- - s.lastDeltaSize = s.buf.ReversePutVarInt64(s.pos, delta) + s.lastDeltaSize, _ = bin.ReversePutVarInt64(s.buf[s.pos:], delta) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } func (s *ReverseInstantDeltaCompressor) appendNewLiteral(delta int64) { s.h = flagLiteral s.lastDelta = delta - s.lastDeltaSize = s.buf.ReversePutVarInt64(s.pos, delta) + s.lastDeltaSize, _ = bin.ReversePutVarInt64(s.buf[s.pos:], delta) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, flagLiteral) // literal, length = 1 + s.buf[s.pos] = flagLiteral // literal, length = 1 s.pos++ } @@ -149,7 +150,7 @@ type InstantDeltaBound struct { Pos int H byte LastDelta int64 - Chunks [][]byte + Chunks []byte } // delta h @@ -159,15 +160,11 @@ func (s *ReverseInstantDeltaCompressor) Lock() { } // позиція посувається вліво, отже може перескочити на попередній chunk pos := s.pos - 1 - s.lastDeltaSize - chunksQty := pos / conbuf.ChunkSize - if (pos % conbuf.ChunkSize) > 0 { - chunksQty++ - } s.state = &InstantDeltaBound{ Pos: pos, H: s.h, LastDelta: s.lastDelta, - Chunks: s.buf.Chunks()[:chunksQty], + Chunks: s.buf[:s.pos], // fix check ? } } @@ -183,9 +180,9 @@ func (s *ReverseInstantDeltaCompressor) Offset() int { return 0 } -func (s *ReverseInstantDeltaCompressor) Snapshot() ([][]byte, int) { +func (s *ReverseInstantDeltaCompressor) Snapshot() ([]byte, int) { if s.state == nil { - return s.buf.Chunks(), s.Size() + return s.buf, s.Size() } // ВАЖЛИВО! // Треба відтворити стан останнього чанка @@ -219,7 +216,7 @@ func (s *ReverseInstantDeltaCompressor) CreateDecompressor(fracDigits byte) qb.V func (s *ReverseInstantDeltaCompressor) Renew() { // УВАГА! // state не чіпаємо - s.buf = conbuf.New(nil) + s.buf = make([]byte, minBufferSize) s.pos = 0 // s.baseValue = 0 @@ -232,14 +229,14 @@ func (s *ReverseInstantDeltaCompressor) CalcRequiredSpace(value float64) int { return 0 } -func (s *ReverseInstantDeltaCompressor) Chunks() [][]byte { - return s.buf.Chunks() +func (s *ReverseInstantDeltaCompressor) Chunks() []byte { + return s.buf } // DECOMPRESSOR type ReverseInstantDeltaDecompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte coef float64 pos int bound int @@ -250,12 +247,12 @@ type ReverseInstantDeltaDecompressor struct { done bool } -func NewReverseInstantDeltaDecompressor(buf *conbuf.ContinuousBuffer, size int, fracDigits byte) *ReverseInstantDeltaDecompressor { +func NewReverseInstantDeltaDecompressor(buf []byte, size int, fracDigits byte) *ReverseInstantDeltaDecompressor { var coef float64 = 1 if fracDigits > 0 { coef = math.Pow(10, float64(fracDigits)) } - i64, n, err := buf.GetVarInt64(0) + i64, n, err := bin.GetVarInt64(buf) if err != nil { log.Fatalf("bug: get base value: %s", err) } @@ -312,7 +309,7 @@ func (s *ReverseInstantDeltaDecompressor) NextValue() (value float64, done bool) } func (s *ReverseInstantDeltaDecompressor) readHeader() { - h := s.buf.GetByte(s.pos) + h := s.buf[s.pos] s.pos-- s.decodeHeaderByte(h) // fmt.Println("h:", h) @@ -320,7 +317,7 @@ func (s *ReverseInstantDeltaDecompressor) readHeader() { // fmt.Println("pending:", s.pending) } func (s *ReverseInstantDeltaDecompressor) readValue() { - i64, n, err := s.buf.ReverseGetVarInt64(s.pos) + i64, n, err := bin.ReverseGetVarInt64(s.buf[:s.pos]) if err != nil { log.Fatalln(err) } diff --git a/chunkenc/time_delta.go b/chunkenc/time_delta.go index 05d2914..6350321 100644 --- a/chunkenc/time_delta.go +++ b/chunkenc/time_delta.go @@ -4,13 +4,24 @@ import ( "fmt" "log" + bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/pretty" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" ) +/* +Payload data файла має таку структуру: +vvvvvvvv-> <-ttttttt +де v - це values, які додаються зліва направо, а читаються справа наліво; +t - це timestamps, які додаються справа наліво, а читаються зліва направо. + +Timestamps читаються звичайними методами Get*, а писатись мають типу ByRightBound + +*/ + type ReverseTimeDeltaCompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte pos int lastUnixtime uint32 lastDelta uint32 @@ -19,21 +30,21 @@ type ReverseTimeDeltaCompressor struct { state *TimeDeltaBound } -func NewReverseTimeDeltaCompressor(buf *conbuf.ContinuousBuffer, size int) *ReverseTimeDeltaCompressor { +func NewReverseTimeDeltaCompressor(buf []byte, size int) *ReverseTimeDeltaCompressor { s := &ReverseTimeDeltaCompressor{ buf: buf, pos: size, // перший вільний байт } if size > 0 { - u64, n, err := s.buf.ReverseGetVarUint64(s.pos - 1) + u64, n, err := bin.ReverseGetVarUint64(s.buf[:s.pos-1]) if err != nil { log.Fatalf("bug: get last unixtime: %s", err) } s.lastUnixtime = uint32(u64) s.pos -= n if s.pos > 0 { - s.h = s.buf.GetByte(s.pos - 1) - u64, s.lastDeltaSize, err = s.buf.ReverseGetVarUint64(s.pos - 2) + s.h = s.buf[s.pos-1] + u64, s.lastDeltaSize, err = bin.ReverseGetVarUint64(s.buf[:s.pos-2]) if err != nil { log.Fatalf("bug: get last delta: %s", err) } @@ -66,7 +77,7 @@ func (s *ReverseTimeDeltaCompressor) Append(unixtime uint32) { if s.h < 127 { // increase counter s.h++ - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } else { // не можу збільшити - буде переповнення. Додаю новий literal блок // counter overflow @@ -106,28 +117,28 @@ func (s *ReverseTimeDeltaCompressor) convertLastFromLiteralToRun() { // Зменшую кількість елементів в literal блоці s.h-- s.pos -= 1 + s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) // закриваю literal блок + s.buf[s.pos] = s.h // закриваю literal блок s.pos++ - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, uint64(s.lastDelta)) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], uint64(s.lastDelta)) s.pos += s.lastDeltaSize s.h = 0 // run блок, довжини 2 - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } func (s *ReverseTimeDeltaCompressor) convertLiteralToRun() { // Знімаю flagLiteral, а лічильник 0 дорівнює 2 елементам в серії. s.h = 0 - s.buf.SetByte(s.pos-1, s.h) + s.buf[s.pos-1] = s.h } func (s *ReverseTimeDeltaCompressor) appendDeltaToLiteral(delta uint32) { s.h++ // збільшую к-сть дельт s.lastDelta = delta s.pos-- - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, uint64(delta)) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], uint64(delta)) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, s.h) + s.buf[s.pos] = s.h s.pos++ } @@ -135,9 +146,9 @@ func (s *ReverseTimeDeltaCompressor) appendNewLiteral(delta uint32) { //fmt.Println("appendNewLiteral", delta) s.h = flagLiteral s.lastDelta = delta - s.lastDeltaSize = s.buf.ReversePutVarUint64(s.pos, uint64(delta)) + s.lastDeltaSize, _ = bin.ReversePutVarUint64(s.buf[s.pos:], uint64(delta)) s.pos += s.lastDeltaSize - s.buf.SetByte(s.pos, flagLiteral) // literal, length = 1 + s.buf[s.pos] = flagLiteral // literal, length = 1 s.pos++ } @@ -146,7 +157,8 @@ func (s *ReverseTimeDeltaCompressor) DeleteLast() { } func (s *ReverseTimeDeltaCompressor) Sync() { - s.pos += s.buf.ReversePutVarUint64(s.pos, uint64(s.lastUnixtime)) + n, _ := bin.ReversePutVarUint64(s.buf[s.pos:], uint64(s.lastUnixtime)) + s.pos += n } // FIX - check methods @@ -156,7 +168,7 @@ type TimeDeltaBound struct { H byte LastUnixtime uint32 LastDelta uint32 - Chunks [][]byte + Chunks []byte } // delta h @@ -179,16 +191,12 @@ func (s *ReverseTimeDeltaCompressor) Lock() { } // позиція посувається вліво, отже може перескочити на попередній chunk pos := s.pos - 1 - s.lastDeltaSize - chunksQty := pos / conbuf.ChunkSize - if (pos % conbuf.ChunkSize) > 0 { - chunksQty++ - } s.state = &TimeDeltaBound{ Pos: pos, H: s.h, LastUnixtime: s.lastUnixtime, // fix LastDelta: s.lastDelta, - Chunks: s.buf.Chunks()[:chunksQty], + Chunks: s.buf[:s.pos], // fix check pos? } } @@ -205,9 +213,9 @@ func (s *ReverseTimeDeltaCompressor) Offset() int { } // fix - -func (s *ReverseTimeDeltaCompressor) Snapshot() ([][]byte, int) { +func (s *ReverseTimeDeltaCompressor) Snapshot() ([]byte, int) { if s.state == nil { - return s.buf.Chunks(), s.Size() + return s.buf, s.Size() } // ВАЖЛИВО! // Треба відтворити стан останнього чанка @@ -241,7 +249,7 @@ func (s *ReverseTimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompress func (s *ReverseTimeDeltaCompressor) Renew() { // УВАГА! // state не чіпаємо - s.buf = conbuf.New(nil) + s.buf = make([]byte, minBufferSize) s.pos = 0 // //s.baseValue = 0 @@ -250,14 +258,14 @@ func (s *ReverseTimeDeltaCompressor) Renew() { s.h = 0 } -func (s *ReverseTimeDeltaCompressor) Chunks() [][]byte { - return s.buf.Chunks() +func (s *ReverseTimeDeltaCompressor) Chunks() []byte { + return s.buf } // DECOMPRESSOR type ReverseTimeDeltaDecompressor struct { - buf *conbuf.ContinuousBuffer + buf []byte pos int lastDelta uint32 lastUnixtime uint32 @@ -267,7 +275,7 @@ type ReverseTimeDeltaDecompressor struct { done bool } -func NewReverseTimeDeltaDecompressor(buf *conbuf.ContinuousBuffer, size int) *ReverseTimeDeltaDecompressor { +func NewReverseTimeDeltaDecompressor(buf []byte, size int) *ReverseTimeDeltaDecompressor { return &ReverseTimeDeltaDecompressor{ buf: buf, pos: size, @@ -278,7 +286,7 @@ func (s *ReverseTimeDeltaDecompressor) RestoreFromEnd() { fmt.Println("RestoreFromEnd", s.pos) if s.pos > 0 { s.pos-- // перший байт даних - u64, n, err := s.buf.ReverseGetVarUint64(s.pos) + u64, n, err := bin.ReverseGetVarUint64(s.buf[:s.pos]) if err != nil { log.Fatalf("bug: get last unixtime: %s", err) } @@ -344,7 +352,7 @@ func (s *ReverseTimeDeltaDecompressor) NextValue() (value uint32, done bool) { } func (s *ReverseTimeDeltaDecompressor) readHeader() { - h := s.buf.GetByte(s.pos) + h := s.buf[s.pos] s.pos-- s.decodeHeaderByte(h) // fmt.Println("h:", h) @@ -362,7 +370,7 @@ func (s *ReverseTimeDeltaDecompressor) decodeHeaderByte(h byte) { } func (s *ReverseTimeDeltaDecompressor) readDelta() { - u64, n, err := s.buf.ReverseGetVarUint64(s.pos) + u64, n, err := bin.ReverseGetVarUint64(s.buf[:s.pos]) if err != nil { log.Fatalln(err) } diff --git a/database/metric.go b/database/metric.go index 8230763..6ba426f 100644 --- a/database/metric.go +++ b/database/metric.go @@ -3,35 +3,34 @@ package database import ( "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/txlog" ) // METRIC -type IndexRec struct { - Since uint32 - PageNo uint32 +type IndexLevelTail struct { + Payload []byte + RecordsCount int } type _metric struct { - MetricType qb.MetricType - FracDigits byte - LastPageNo uint32 - SinceValue float64 - Since uint32 - UntilValue float64 - Until uint32 - Timestamps qb.TimestampCompressor - Values qb.ValueCompressor - // fix - add index levels - XLock bool - RLocks int - WaitQueue []any - //IndexLevels [][]IndexRec // root - last element - IndexLevels [][]byte // root - last element + MetricType qb.MetricType + FracDigits byte + LastPageNo uint32 + SinceValue float64 + Since uint32 + UntilValue float64 + Until uint32 + Timestamps qb.TimestampCompressor + Values qb.ValueCompressor + XLock bool + RLocks int + WaitQueue []any + IndexLevelTails []IndexLevelTail // root - last element } +//IndexLevels [][]IndexRec // root - last element + func (s *_metric) ReinitBy(timestamp uint32, value float64) { s.Timestamps.Renew() s.Values.Renew() @@ -74,8 +73,8 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu timestamps = s.Timestamps values = s.Values - indexLevels []*atree.IndexLevel - dataPages []*atree.DataPage + indexLevels []atree.IndexLevelTail + dataPages []atree.DataPayload //written int //resultCode byte @@ -135,12 +134,10 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu } else { // сторінка заповнена // prevPageNo - виставляю в txlog, коли забираю номер сторінки із freeList або генерую новий - dataPages = append(dataPages, &atree.DataPage{ - PrevPageNo: 0, // ? - LowerTimestamp: s.Since, - Timestamps: timestamps.Chunks(), + dataPages = append(dataPages, atree.DataPayload{ + Since: s.Since, + Content: nil, TimestampsSize: timestamps.Size(), - Values: values.Chunks(), ValuesSize: values.Size(), }) @@ -163,28 +160,31 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage fu if len(dataPages) > 0 { // пишу в txlog довгим шляхом через redo файл і запис в data файл - for _, unfilled := range s.IndexLevels { - indexLevels = append(indexLevels, &atree.IndexLevel{ - Offset: len(unfilled), - LowerTimestamp: bin.GetUint32(unfilled), - Data: unfilled, + for _, tail := range s.IndexLevelTails { + indexLevels = append(indexLevels, atree.IndexLevelTail{ + Payload: tail.Payload, + RecordsCount: tail.RecordsCount, }) } + + sendToStorage(txlog.AppendedMeasures{ + MetricID: req.MetricID, + LastPageNo: s.LastPageNo, + TimestampsOffset: timestamps.Offset(), // state.Pos() з якої позиції дописувати дані на сторінку 0 (при відновленні) + ValuesOffset: values.Offset(), // state.Pos() + //Payload: s.payload, // fix - timestamps + values ? or timestamps and values (for WAL) + IndexLevelTails: indexLevels, + DataPages: dataPages, + Timestamps: nil, + Values: nil, + //ResultCode: resultCode, + //WrittenCount: wri, + ResultCh: nil, + }) } else { // короткий шлях - запис лише в txlog } - sendToStorage(txlog.AppendedMeasures{ - MetricID: req.MetricID, - TimestampsOffset: timestamps.Offset(), // state.Pos() з якої позиції дописувати дані на сторінку 0 (при відновленні) - ValuesOffset: values.Offset(), // state.Pos() - Timestamps: timestamps.Chunks(), - TimestampsSize: timestamps.Size(), - Values: values.Chunks(), - ValuesSize: values.Size(), - IndexLevels: indexLevels, - DataPages: dataPages, - }) } func (s *_metric) FinAppendMeasures(rec txlog.AppendMeasuresSummary) { diff --git a/go.mod b/go.mod index 45f442c..e57b9dc 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,6 @@ go 1.24.2 require ( gopkg.in/ini.v1 v1.67.1 - gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248 + gordenko.dev/dima/bin v0.0.0-20260604235618-bff16d774d98 gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 ) diff --git a/go.sum b/go.sum index 611f732..7a7aa11 100644 --- a/go.sum +++ b/go.sum @@ -20,5 +20,7 @@ gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248 h1:APYOkErjDz3hykfcddvmABsexoyKUgdRPU3XiM16F34= gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248/go.mod h1:/I+9fvRUzXHgXSwGEwOK1mBBM4BJ3AtCmoCZ2upDiwk= +gordenko.dev/dima/bin v0.0.0-20260604235618-bff16d774d98 h1:pQzJ4wnSrFXn9Hu+v6KvWZ4GQVBmuTljJGZ20Pb+BJA= +gordenko.dev/dima/bin v0.0.0-20260604235618-bff16d774d98/go.mod h1:/I+9fvRUzXHgXSwGEwOK1mBBM4BJ3AtCmoCZ2upDiwk= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 h1:nyJ3mzTQ46yUeMZCdLyYcs7B5JCS54c67v84miyhq2E= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69/go.mod h1:AxgKDktpqBVyIOhIcP+nlCpK+EsJyjN5kPdqyd8euVU= diff --git a/qb.go b/qb.go index 5f48e8e..337e2b4 100644 --- a/qb.go +++ b/qb.go @@ -86,6 +86,7 @@ const ( UnknownWorkerQueueItemBug AbortCode = 20 UnknownMetricWaitQueueItemBug AbortCode = 21 RepeatableLock AbortCode = 22 + NoSpaceOnIndexPage AbortCode = 23 // GetRecoveryRecipeFailed AbortCode = 26 LoadSnapshotFailed AbortCode = 27 diff --git a/txlog/data_preparer.go b/txlog/data_preparer.go new file mode 100644 index 0000000..8b43037 --- /dev/null +++ b/txlog/data_preparer.go @@ -0,0 +1,207 @@ +package txlog + +import ( + "bytes" + "io" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/util" +) + +type DataPreparer struct { + pagePreparer *PagePreparer + //dataPagePayloadBound int + w *bytes.Buffer + indexPages []PageToWrite + dataPages []PageToWrite + writeResults []any +} + +type DataPreparerOptions struct { + MinWALBufferSize int + PagePreparer *PagePreparer +} + +func NewDataPreparer(opt DataPreparerOptions) *DataPreparer { + buf := make([]byte, opt.MinWALBufferSize) + s := &DataPreparer{ + pagePreparer: opt.PagePreparer, + w: bytes.NewBuffer(buf), + } + return s +} + +type PreparedData struct { + Packet []byte // пакет для запису в WAL + WriteToIndex []PageToWrite + WriteToData []PageToWrite + WriteResults []any +} + +func (s *DataPreparer) Prepare(input []any) PreparedData { + s.w.Write([]byte{ + 0, 0, 0, 0, 0, 0, 0, 0, 0, // size (max 9 byte) + 0, 0, 0, 0, // crc32 + }) + + hasher := util.NewHasher() + + w := io.MultiWriter(s.w, hasher) + + // 1. Пакую всі дані в WAL буфер запису + for _, untyped := range input { + switch x := untyped.(type) { + case AddedMetric: + x.Pack(w) + // case DeletedMetric: + // if len(x.FreePageNumbers) > 0 { + // s.freeList.AddPageNumbers(x.FreePageNumbers) + // } + case AppendedMeasures: + sealResult := s.pagePreparer.SealPages(SealPagesIn{ + LastPageNo: x.LastPageNo, + IndexLevelTails: x.IndexLevelTails, + DataPages: x.DataPages, + }) + + // сторінки для запису в index та data файли + for _, p := range sealResult.DataPages { + s.dataPages = append(s.dataPages, PageToWrite{ + PageNo: p.PageNo, + Content: p.Content, + }) + } + for _, level := range sealResult.IndexLevels { + for _, p := range level.IndexPages { + s.indexPages = append(s.indexPages, PageToWrite{ + PageNo: p.PageNo, + Content: p.Content, + }) + } + } + + var ( + // дані для запису в WAL + changedIndexPages []ChangedIndexLevel + // дані для Воркера + indexLevelTails []IndexLevelTail + ) + + for _, level := range sealResult.IndexLevels { + if len(level.IndexPages) == 0 && level.SkipRecords == level.TailRecordsCount { + break + } + var completedIndexPage *CompletedIndexPage + if len(level.IndexPages) > 0 { + first := level.IndexPages[0] + skipSize := level.SkipRecords * indexRecordSize + completedIndexPage = &CompletedIndexPage{ + PageNo: first.PageNo, + Reused: first.Reused, + Checksum: first.Checksum, + Records: first.Content[skipSize:], + } + } + changedIndexPages = append(changedIndexPages, ChangedIndexLevel{ + CompletedIndexPage: completedIndexPage, + IndexPages: level.IndexPages[1:], + TailRecords: level.TailRecords[:level.TailRecordsCount*indexRecordSize], + }) + + indexLevelTails = append(indexLevelTails, IndexLevelTail{ + Records: level.TailRecords, + RecordsCount: level.TailRecordsCount, + }) + } + + first := sealResult.DataPages[0] + + rec := WALRecordAppendMeasures{ + MetricID: x.MetricID, + CompletedDataPage: CompletedDataPage{ + PageNo: first.PageNo, + Reused: first.Reused, + PrevPageNo: first.PrevPageNo, + Checksum: first.Checksum, + TimestampsOffset: x.TimestampsOffset, + Timestamps: x.Timestamps, + ValuesOffset: x.ValuesOffset, + Values: x.Values, + }, + DataPages: sealResult.DataPages[1:], + ChangedIndexLevels: changedIndexPages, + TailTimestamps: x.Timestamps, + TailValues: x.Values, + } + + rec.WriteTo(s.w) + + // Дані для Worker + s.writeResults = append(s.writeResults, AppendMeasuresSummary{ + MetricID: x.MetricID, + LastPageNo: sealResult.DataPages[len(sealResult.DataPages)-1].PageNo, + Index: indexLevelTails, + ResultCode: x.ResultCode, + WrittenCount: x.WrittenCount, + ResultCh: x.ResultCh, + }) + + //case DeletedMeasures: + //case DeletedMeasuresSince: + } + } + // 2. Додаю розмір пакету і чексуму + + //s.written += int64(len(packet)) + 12 + + packet := s.w.Bytes() + // write size + payloadSize := len(packet) - 13 + start := 9 - bin.CountVarSize(payloadSize) + bin.PutVarSize(packet[start:], payloadSize) + // checksum + bin.PutUint32(packet[9:], hasher.Sum32()) + // + return PreparedData{ + Packet: packet[start:], + WriteToIndex: s.indexPages, + WriteToData: s.dataPages, + WriteResults: s.writeResults, + } +} + +func (s *DataPreparer) Reset() { + s.w.Reset() + s.indexPages = nil + s.dataPages = nil + s.writeResults = nil +} + +// MeasuresToWrite +type AppendedMeasures struct { + MetricID uint32 + LastPageNo uint32 + TimestampsOffset int + Timestamps []byte // offset on 1st page + ValuesOffset int // offset on 1st page + Values []byte + IndexLevelTails []IndexLevelTail + DataPages []DataPayload + TailTimestamps []byte // data level tail + TailValues []byte // data level tail + ResultCode int + WrittenCount int + ResultCh chan struct{} +} + +// WRITE RESULTS + +// результат +type AppendMeasuresSummary struct { + MetricID uint32 + LastPageNo uint32 + Index []IndexLevelTail + ResultCode int + WrittenCount int + ResultCh chan struct{} +} diff --git a/txlog/helpers.go b/txlog/helpers.go index 08ce092..975b8c0 100644 --- a/txlog/helpers.go +++ b/txlog/helpers.go @@ -9,42 +9,6 @@ import ( const recordSize = 8 const chunkSize = 24 -/* -Format appended measures: -1b - tx type -4b - metricID -Nb - qty of filled data pages (varsize) -[ - 4b - pageNo - 1b - reused - 4b - prevPageNo - 4b - page crc32 - 2b - timestamps size - Nb - timestamps payload - 2b - values size - Nb - values payload -] -2b - timestamps offset (for the 1st data page only) -2b - values offset (for the 1st data page only) -Nb - timestamps size unfilled (varsize) -Nb - timestamps payload unfilled -Nb - values size unfilled (varsize) -Nb - values payload unfilled -Nb - qty of index levels (varsize) -[ - Nb - qty of level filled pages (varsize) - [ - 4b - pageNo - 1b - reused - 4b - page crc32 - Nb - records qty (varsize) - Nb - records payload - ] - Nb - records qty unfilled (varsize) - Nb - records payload unfilled -] -*/ - func copyPayloadFromChunks(w *bytes.Buffer, chunks [][]byte, size int, offset int) { bin.WriteUint16(w, uint16(size)) var ( @@ -69,170 +33,170 @@ func copyPayloadFromChunks(w *bytes.Buffer, chunks [][]byte, size int, offset in ///////////////////////////////////////// -type DataPage struct { - PageNo uint32 - Reused bool - PrevPageNo uint32 - Checksum uint32 - Timestamps []byte - Values []byte -} +// type DataPage struct { +// PageNo uint32 +// Reused bool +// PrevPageNo uint32 +// Checksum uint32 +// Timestamps []byte +// Values []byte +// } -type IndexPage struct { - PageNo uint32 - Reused bool - Checksum uint32 - Records []byte -} +// type IndexPage struct { +// PageNo uint32 +// Reused bool +// Checksum uint32 +// Records []byte +// } -type IndexLevel struct { - Pages []IndexPage - Records []byte // unfilled -} +// type IndexLevel struct { +// Pages []IndexPage +// Records []byte // unfilled +// } // Для декодінга -type TxAppendedMeasures struct { - MetricID uint32 - DataPages []DataPage - // offsets потрібні тому що дані не додаються в кінець, а перезаписують кілька - // останніх байтів unfilled даних. Потрібно для коректного recovery. - // Якщо є DataPages, то застосовуються для 1-ї Data сторінки, інакше - до timestamps і values (unfilled) - TimestampsOffset int - ValuesOffset int - Timestamps []byte // unfilled - Values []byte // unfilled - IndexLevels []IndexLevel -} +// type TxAppendedMeasures struct { +// MetricID uint32 +// DataPages []DataPage +// // offsets потрібні тому що дані не додаються в кінець, а перезаписують кілька +// // останніх байтів unfilled даних. Потрібно для коректного recovery. +// // Якщо є DataPages, то застосовуються для 1-ї Data сторінки, інакше - до timestamps і values (unfilled) +// TimestampsOffset int +// ValuesOffset int +// Timestamps []byte // unfilled +// Values []byte // unfilled +// IndexLevels []IndexLevel +// } -func (s *TxAppendedMeasures) Read(r *bytes.Buffer) (err error) { - s.MetricID, err = bin.ReadUint32(r) - if err != nil { - return - } - dataPagesQty, err := bin.ReadVarSize(r) - if err != nil { - return - } - for range dataPagesQty { - var ( - p DataPage - timestampsSize, valuesSize int - //b byte - ) - p.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - p.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - p.PrevPageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - p.Checksum, err = bin.ReadUint32(r) - if err != nil { - return - } - timestampsSize, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - p.Timestamps, err = bin.ReadN(r, timestampsSize) - if err != nil { - return - } - valuesSize, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - p.Values, err = bin.ReadN(r, valuesSize) - if err != nil { - return - } - s.DataPages = append(s.DataPages, p) - } - s.TimestampsOffset, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - s.ValuesOffset, err = bin.ReadUint16AsInt(r) - if err != nil { - return - } - var ( - timestampsSize, valuesSize int - ) - timestampsSize, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.Timestamps, err = bin.ReadN(r, timestampsSize) - if err != nil { - return - } - valuesSize, err = bin.ReadVarSize(r) - if err != nil { - return - } - s.Values, err = bin.ReadN(r, valuesSize) - if err != nil { - return - } - if len(s.DataPages) == 0 { - return - } - indexLevelsQty, err := bin.ReadVarSize(r) - if err != nil { - return - } - for range indexLevelsQty { - var ( - indexPagesQty int - level IndexLevel - recordsQty int - ) - indexPagesQty, err = bin.ReadVarSize(r) - if err != nil { - return - } - for range indexPagesQty { - var ( - p IndexPage - recordsQty int - ) - p.PageNo, err = bin.ReadUint32(r) - if err != nil { - return - } - p.Reused, err = bin.ReadBool(r) - if err != nil { - return - } - p.Checksum, err = bin.ReadUint32(r) - if err != nil { - return - } - recordsQty, err = bin.ReadVarSize(r) - if err != nil { - return - } - p.Records, err = bin.ReadN(r, recordsQty*recordSize) - if err != nil { - return - } - level.Pages = append(level.Pages, p) - } - recordsQty, err = bin.ReadVarSize(r) - if err != nil { - return - } - level.Records, err = bin.ReadN(r, recordsQty*recordSize) - if err != nil { - return - } - s.IndexLevels = append(s.IndexLevels, level) - } - return -} +// func (s *TxAppendedMeasures) Read(r *bytes.Buffer) (err error) { +// s.MetricID, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// dataPagesQty, err := bin.ReadVarSize(r) +// if err != nil { +// return +// } +// for range dataPagesQty { +// var ( +// p DataPage +// timestampsSize, valuesSize int +// //b byte +// ) +// p.PageNo, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// p.Reused, err = bin.ReadBool(r) +// if err != nil { +// return +// } +// p.PrevPageNo, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// p.Checksum, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// timestampsSize, err = bin.ReadUint16AsInt(r) +// if err != nil { +// return +// } +// p.Timestamps, err = bin.ReadN(r, timestampsSize) +// if err != nil { +// return +// } +// valuesSize, err = bin.ReadUint16AsInt(r) +// if err != nil { +// return +// } +// p.Values, err = bin.ReadN(r, valuesSize) +// if err != nil { +// return +// } +// s.DataPages = append(s.DataPages, p) +// } +// s.TimestampsOffset, err = bin.ReadUint16AsInt(r) +// if err != nil { +// return +// } +// s.ValuesOffset, err = bin.ReadUint16AsInt(r) +// if err != nil { +// return +// } +// var ( +// timestampsSize, valuesSize int +// ) +// timestampsSize, err = bin.ReadVarSize(r) +// if err != nil { +// return +// } +// s.Timestamps, err = bin.ReadN(r, timestampsSize) +// if err != nil { +// return +// } +// valuesSize, err = bin.ReadVarSize(r) +// if err != nil { +// return +// } +// s.Values, err = bin.ReadN(r, valuesSize) +// if err != nil { +// return +// } +// if len(s.DataPages) == 0 { +// return +// } +// indexLevelsQty, err := bin.ReadVarSize(r) +// if err != nil { +// return +// } +// for range indexLevelsQty { +// var ( +// indexPagesQty int +// level IndexLevel +// recordsQty int +// ) +// indexPagesQty, err = bin.ReadVarSize(r) +// if err != nil { +// return +// } +// for range indexPagesQty { +// var ( +// p IndexPage +// recordsQty int +// ) +// p.PageNo, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// p.Reused, err = bin.ReadBool(r) +// if err != nil { +// return +// } +// p.Checksum, err = bin.ReadUint32(r) +// if err != nil { +// return +// } +// recordsQty, err = bin.ReadVarSize(r) +// if err != nil { +// return +// } +// p.Records, err = bin.ReadN(r, recordsQty*recordSize) +// if err != nil { +// return +// } +// level.Pages = append(level.Pages, p) +// } +// recordsQty, err = bin.ReadVarSize(r) +// if err != nil { +// return +// } +// level.Records, err = bin.ReadN(r, recordsQty*recordSize) +// if err != nil { +// return +// } +// s.IndexLevels = append(s.IndexLevels, level) +// } +// return +// } diff --git a/txlog/page_preparer.go b/txlog/page_preparer.go new file mode 100644 index 0000000..5555b29 --- /dev/null +++ b/txlog/page_preparer.go @@ -0,0 +1,279 @@ +package txlog + +import ( + "errors" + + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/util" +) + +// Задача метода: +// - для data сторінок отримати pageNo та reused, задати sizes та prevPageNo, порахувати CRC32 +// - додати пари до index, згенерувати нові index сторінки +// levels []*IndexLevel, dataPages []*DataPage +// fix - reduce leves + +type PagePreparer struct { + dataChecksumIdx int + maxRecordsOnIndexPage int + indexPageIncSize int // кратно indexPageSize + timestampsSizeIdx int + valuesSizeIdx int + prevPageIdx int + indexRecordsCountIdx int + isZeroLevelIdx int + indexChecksumIdx int + getDataPageNumber func() (uint32, bool, error) + getIndexPageNumber func() (uint32, bool, error) +} + +type PagePreparerOptions struct { + IndexPageSize int + IndexPageIncSize int + DataPageSize int + GetDataPageNumber func() (uint32, bool, error) + GetIndexPageNumber func() (uint32, bool, error) +} + +func NewPagePreparer(opt PagePreparerOptions) (*PagePreparer, error) { + if opt.GetIndexPageNumber == nil { + return nil, errors.New("missing required option: GetIndexPageNumber") + } + if opt.GetDataPageNumber == nil { + return nil, errors.New("missing required option: GetDataPageNumber") + } + if (opt.IndexPageSize % opt.IndexPageIncSize) != 0 { + return nil, errors.New("IndexPageIncSize must multiple of IndexPageSize") + } + s := &PagePreparer{ + dataChecksumIdx: opt.DataPageSize - 4, + indexPageIncSize: opt.IndexPageIncSize, + timestampsSizeIdx: opt.DataPageSize - 6, + valuesSizeIdx: opt.DataPageSize - 8, + prevPageIdx: opt.DataPageSize - 12, + indexRecordsCountIdx: opt.IndexPageSize - 6, + isZeroLevelIdx: opt.IndexPageSize - 7, + indexChecksumIdx: opt.IndexPageSize - 4, + getIndexPageNumber: opt.GetIndexPageNumber, + getDataPageNumber: opt.GetDataPageNumber, + } + s.maxRecordsOnIndexPage = (opt.IndexPageSize - 7) / indexRecordSize + return s, nil +} + +type appendIndexRecordIn struct { + Records []byte + RecordsCount int + Timestamp uint32 + PageNo uint32 +} + +func (s *PagePreparer) appendIndexRecord(in appendIndexRecordIn) []byte { + if in.RecordsCount < s.maxRecordsOnIndexPage { + var ( + pos int + buf = in.Records + ) + if in.RecordsCount > 0 { + pos = in.RecordsCount * indexRecordSize + } else { + buf = make([]byte, s.indexPageIncSize) + } + bin.PutUint32(buf[pos:], in.Timestamp) + bin.PutUint32(buf[pos+4:], in.PageNo) + return buf + } + qb.Abort(qb.NoSpaceOnIndexPage, nil) + return nil +} + +type sealDataPageIn struct { + Content []byte + PrevPageNo uint32 + TimestampsSize int + ValuesSize int +} + +func (s *PagePreparer) sealDataPage(in sealDataPageIn) (checksum uint32) { + bin.PutUint16(in.Content[s.timestampsSizeIdx:], uint16(in.TimestampsSize)) + bin.PutUint16(in.Content[s.valuesSizeIdx:], uint16(in.ValuesSize)) + bin.PutUint32(in.Content[s.prevPageIdx:], in.PrevPageNo) + + checksum = util.CalcChecksum(in.Content[:s.dataChecksumIdx]) + bin.PutUint32(in.Content[s.dataChecksumIdx:], checksum) + return +} + +type sealIndexPageIn struct { + Content []byte + RecordsCount int + LastLevel bool +} + +func (s *PagePreparer) sealIndexPage(in sealIndexPageIn) (checksum uint32) { + bin.PutUint16(in.Content[s.indexRecordsCountIdx:], uint16(in.RecordsCount)) + if in.LastLevel { + in.Content[s.isZeroLevelIdx] = 1 + } + checksum = util.CalcChecksum(in.Content[:s.indexChecksumIdx]) + bin.PutUint32(in.Content[s.indexChecksumIdx:], checksum) + return +} + +type IndexLevelTail struct { + Records []byte // розмір більший за кількість + RecordsCount int +} + +type DataPayload struct { + Since uint32 + Content []byte + TimestampsSize int + ValuesSize int +} + +type SealPagesIn struct { + LastPageNo uint32 + IndexLevelTails []IndexLevelTail + DataPages []DataPayload +} + +type SealedDataPage struct { + PageNo uint32 + Reused bool + PrevPageNo uint32 + Checksum uint32 + Content []byte +} + +type SealedIndexPage struct { + Content []byte + PageNo uint32 + Reused bool + Checksum uint32 +} + +// SealedIndexLevel - як зрозуміти чи потрібно щось писати в WAL чи index файл, +// чи додавання нових data сторінок не зачепило індексний рівень? +// Якщо є IndexPages - пишемо весь рівень в WAL і готові сторінки в index файл. +// Якщо IndexPages пустий - порівнюємо Offset і len(Payload) - якщо довжина +// більша за змішення - дані додані на незаповнену сторінку рівня, +// отже різницю пишемо в WAL. +type SealedIndexLevel struct { + // Якщо є IndexPages, offset - це зміщення із якого починаються дані першої сторінки, + // які ще не зписані в WAL. + // Якщо IndexPages пустий, offset - це зміщення із якого починаються дані payload, + // які ще не зписані в WAL. + // В payload в будь-якому разі дані останньої індексної сторінки рівня (незаповненої). + SkipRecords int + IndexPages []SealedIndexPage + TailRecords []byte + TailRecordsCount int +} + +type SealResult struct { + DataPages []SealedDataPage + IndexLevels []*SealedIndexLevel +} + +func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) { + var ( + sealedDataPages []SealedDataPage + sealedIndexLevels []*SealedIndexLevel + ) + for _, tail := range in.IndexLevelTails { + sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{ + SkipRecords: tail.RecordsCount, + TailRecords: tail.Records, + TailRecordsCount: tail.RecordsCount, + }) + } + for dataPageIdx, d := range in.DataPages { + var ( + upTimestamp = d.Since + upPageNo uint32 + levelIdx = 0 + prevPageNo uint32 + ) + if dataPageIdx == 0 { + prevPageNo = in.LastPageNo + } else { + prevPageNo = result.DataPages[dataPageIdx-1].PageNo + } + pageNo, reused, err := s.getDataPageNumber() + if err != nil { + qb.Abort(qb.FailedGetPageNumber, err) + } + sealedDataPages = append(sealedDataPages, SealedDataPage{ + Content: d.Content, + PrevPageNo: prevPageNo, + PageNo: pageNo, + Reused: reused, + Checksum: s.sealDataPage(sealDataPageIn{ + Content: d.Content, + PrevPageNo: prevPageNo, + TimestampsSize: d.TimestampsSize, + ValuesSize: d.ValuesSize, + }), + }) + upPageNo = pageNo + for { + if levelIdx == len(sealedIndexLevels) { + // новий root + sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{ + TailRecords: s.appendIndexRecord(appendIndexRecordIn{ + Records: nil, + Timestamp: upTimestamp, + PageNo: upPageNo, + }), + TailRecordsCount: 1, + }) + break + } + level := sealedIndexLevels[levelIdx] + if level.TailRecordsCount < s.maxRecordsOnIndexPage { + level.TailRecords = s.appendIndexRecord(appendIndexRecordIn{ + Records: level.TailRecords, + Timestamp: upTimestamp, + PageNo: upPageNo, + }) + level.TailRecordsCount++ + break + } + // в буфері немає місця для додавання нової пари, отже це + // заповнена сторінка + pageNo, reused, err := s.getIndexPageNumber() + if err != nil { + qb.Abort(qb.FailedGetPageNumber, err) + } + + filled := SealedIndexPage{ + Content: level.TailRecords, + PageNo: pageNo, + Reused: reused, + Checksum: s.sealIndexPage(sealIndexPageIn{ + Content: level.TailRecords, + RecordsCount: level.TailRecordsCount, + LastLevel: levelIdx == 0, + }), + } + level.IndexPages = append(level.IndexPages, filled) + // новий tail на індексному рівні + level.TailRecords = s.appendIndexRecord(appendIndexRecordIn{ + Records: nil, + Timestamp: upTimestamp, + PageNo: upPageNo, + }) + level.TailRecordsCount = 1 + // + upPageNo = pageNo + upTimestamp = bin.GetUint32(filled.Content) + levelIdx++ + } + } + return SealResult{ + DataPages: sealedDataPages, + IndexLevels: sealedIndexLevels, + } +} diff --git a/txlog/txlog.go b/txlog/txlog.go index 415b14f..95c3a90 100644 --- a/txlog/txlog.go +++ b/txlog/txlog.go @@ -1,13 +1,15 @@ package txlog import ( - "bytes" "io" bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" - "gordenko.dev/dima/qb/util" +) + +var ( + indexRecordSize = 8 ) type AddedMetric struct { @@ -359,62 +361,3 @@ Nb - values payload // } // PACK - -type PreparedData struct { - Packet []byte // пакет для запису в WAL - WriteToIndex []IndexPayloadToWrite - WriteToData []DataPayloadToWrite - ToWorker []any -} - -func prepareData(input []any) PreparedData { - buffer := bytes.NewBuffer(nil) - buffer.Write([]byte{ - 0, 0, 0, 0, 0, 0, 0, 0, // size (max 8 byte) - 0, 0, 0, 0, // crc32 - }) - - hasher := util.NewHasher() - - w := io.MultiWriter(buffer, hasher) - - var ( - writeToIndex []IndexPayloadToWrite - writeToData []DataPayloadToWrite - ) - - // 1. Пакую всі дані в WAL буфер запису - for _, untyped := range input { - switch x := untyped.(type) { - case AddedMetric: - x.Pack(w) - // case DeletedMetric: - // if len(x.FreePageNumbers) > 0 { - // s.freeList.AddPageNumbers(x.FreePageNumbers) - // } - case AppendedMeasures: - //s.packMeasuresIntoWALBuffer(x) - //case DeletedMeasures: - //case DeletedMeasuresSince: - } - } - // 2. Додаю розмір пакету і чексуму - - //s.written += int64(len(packet)) + 12 - - packet := buffer.Bytes() - // write size - payloadSize := len(packet) - 12 - start := 8 - bin.CountVarSize(payloadSize) - bin.PutVarSize(packet[start:], payloadSize) - // checksum - bin.PutUint32(packet[8:], hasher.Sum32()) - // - //totalSize := (12 - start) + payloadSize - - return PreparedData{ - Packet: packet[start:], - WriteToIndex: writeToIndex, - WriteToData: writeToData, - } -} diff --git a/txlog/wal_records.go b/txlog/wal_records.go new file mode 100644 index 0000000..29e289c --- /dev/null +++ b/txlog/wal_records.go @@ -0,0 +1,306 @@ +package txlog + +import ( + "bytes" + + bin "gordenko.dev/dima/bin/little" +) + +type CompletedDataPage struct { + PageNo uint32 + Reused bool + PrevPageNo uint32 + Checksum uint32 + TimestampsOffset int // offset on 1st page + Timestamps []byte + ValuesOffset int // offset on 1st page + Values []byte +} + +type CompletedIndexPage struct { + PageNo uint32 + Reused bool + Checksum uint32 + Records []byte // offset не потрібен, оскільки Records додаються в кінець +} + +type ChangedIndexLevel struct { + CompletedIndexPage *CompletedIndexPage + IndexPages []SealedIndexPage + TailRecords []byte // payload only +} + +type WALRecordAppendMeasures struct { + MetricID uint32 + CompletedDataPage CompletedDataPage + DataPages []SealedDataPage + ChangedIndexLevels []ChangedIndexLevel + TailTimestamps []byte // data level tail + TailValues []byte // data level tail +} + +/* +Format appended measures: +1b - tx type +4b - metricID +4b - pageNo +1b - reused +4b - prevPageNo +4b - page crc32 +Nb - (varsize) timestamps offset +Nb - (varsize) timestamps tail size on 1st filled data page +Nb - timestamps tail on 1st filled data page +Nb - (varsize) values offset +Nb - (varsize) values tail size on 1st filled data page +Nb - values tail on 1st filled data page + +Nb - (varsize) filled data pages qty +[ + + 4b - pageNo + 1b - reused + Nb - (data page size) page content + +] +Nb - (varsize) timestamps payload size on tail +Nb - timestamps payload +Nb - (varsize) values payload size on tail +Nb - values payload +Nb - (varsize) qty of index levels +[ + + // NOTE! + // - idx = 0 - zeroLevel + // - skipped = maxRecordsOnIndexPage - records count + 4b - pageNo of 1st index page + 1b - reused + 4b - page crc32 + Nb - (varsize) records tail count on 1st filled index page + Nb - records tail on 1st filled index page + + + Nb - (varsize) qty of level filled pages + [ + 4b - pageNo + 1b - reused + Nb - (index page size) page content + ] + Nb - (varsize) size of records level tail + Nb - level tail records + +] +*/ +func (s WALRecordAppendMeasures) WriteTo(w *bytes.Buffer) { + w.WriteByte(CodeAppendedMeasures) + bin.WriteUint32(w, s.MetricID) + // completed data page + completed := s.CompletedDataPage + bin.WriteUint32(w, completed.PageNo) + bin.WriteBool(w, completed.Reused) + bin.WriteUint32(w, completed.PrevPageNo) + bin.WriteUint32(w, completed.Checksum) + bin.WriteVarSize(w, completed.TimestampsOffset) + bin.WriteVarSize(w, len(completed.Timestamps)) + w.Write(completed.Timestamps) + bin.WriteVarSize(w, completed.ValuesOffset) + bin.WriteVarSize(w, len(completed.Values)) + w.Write(completed.Values) + // data pages + bin.WriteVarSize(w, len(s.DataPages)) + for _, p := range s.DataPages { + bin.WriteUint32(w, p.PageNo) + bin.WriteBool(w, p.Reused) + bin.WriteVarSize(w, len(p.Content)) + w.Write(p.Content) + } + // data tail + bin.WriteVarSize(w, len(s.TailTimestamps)) + w.Write(s.TailTimestamps) + bin.WriteVarSize(w, len(s.TailValues)) + w.Write(s.TailValues) + // changed levels + bin.WriteVarSize(w, len(s.ChangedIndexLevels)) + for _, level := range s.ChangedIndexLevels { + // completed index page + completed := level.CompletedIndexPage + bin.WriteUint32(w, completed.PageNo) + bin.WriteBool(w, completed.Reused) + bin.WriteUint32(w, completed.Checksum) + bin.WriteVarSize(w, len(completed.Records)/indexRecordSize) + w.Write(completed.Records) + // full index pages + for _, p := range level.IndexPages { + bin.WriteUint32(w, p.PageNo) + bin.WriteBool(w, p.Reused) + bin.WriteVarSize(w, len(p.Content)) + w.Write(p.Content) + } + // level tail + bin.WriteVarSize(w, len(level.TailRecords)/indexRecordSize) + w.Write(level.TailRecords) + } + return +} + +func (s *WALRecordAppendMeasures) Parse(r *bytes.Buffer) (err error) { + var ( + size int + recordsCount int + ) + s.MetricID, err = bin.ReadUint32(r) + if err != nil { + return + } + s.CompletedDataPage.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.CompletedDataPage.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + s.CompletedDataPage.PrevPageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + s.CompletedDataPage.Checksum, err = bin.ReadUint32(r) + if err != nil { + return + } + s.CompletedDataPage.TimestampsOffset, err = bin.ReadVarSize(r) + if err != nil { + return + } + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.CompletedDataPage.Timestamps, err = bin.ReadN(r, size) + if err != nil { + return + } + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.CompletedDataPage.Values, err = bin.ReadN(r, size) + if err != nil { + return + } + // data pages + dataPagesQty, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range dataPagesQty { + var ( + p SealedDataPage + ) + p.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + p.Content, err = bin.ReadN(r, size) + if err != nil { + return + } + s.DataPages = append(s.DataPages, p) + } + // data tail + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.TailTimestamps, err = bin.ReadN(r, size) + if err != nil { + return + } + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + s.TailValues, err = bin.ReadN(r, size) + if err != nil { + return + } + // changed levels + changedIndexLevelsCount, err := bin.ReadVarSize(r) + if err != nil { + return + } + for range changedIndexLevelsCount { + var ( + indexPagesCount int + level ChangedIndexLevel + ) + // completed index page + level.CompletedIndexPage.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + level.CompletedIndexPage.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + level.CompletedIndexPage.Checksum, err = bin.ReadUint32(r) + if err != nil { + return + } + recordsCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + level.CompletedIndexPage.Records, err = bin.ReadN(r, recordsCount*indexRecordSize) + if err != nil { + return + } + // + indexPagesCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + for range indexPagesCount { + var p SealedIndexPage + + p.PageNo, err = bin.ReadUint32(r) + if err != nil { + return + } + p.Reused, err = bin.ReadBool(r) + if err != nil { + return + } + p.Checksum, err = bin.ReadUint32(r) + if err != nil { + return + } + size, err = bin.ReadVarSize(r) + if err != nil { + return + } + p.Content, err = bin.ReadN(r, size) + if err != nil { + return + } + level.IndexPages = append(level.IndexPages, p) + } + recordsCount, err = bin.ReadVarSize(r) + if err != nil { + return + } + level.TailRecords, err = bin.ReadN(r, recordsCount*indexRecordSize) + if err != nil { + return + } + s.ChangedIndexLevels = append(s.ChangedIndexLevels, level) + } + return +} diff --git a/txlog/writer.go b/txlog/writer.go index 19debde..2392e72 100644 --- a/txlog/writer.go +++ b/txlog/writer.go @@ -6,7 +6,6 @@ import ( "errors" "fmt" "io" - "log" "math" "os" "path/filepath" @@ -63,21 +62,22 @@ type Changes struct { } type Writer struct { - mutex sync.Mutex - dataPagesCount uint32 - indexPagesCount uint32 - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList - atree *atree.Atree - logNumber int - dir string - w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск - wal *os.File - dataFile *os.File - indexFile *os.File - // буфер для складання payload на сторінку, розрахунку CRC32 та інше - pageBuffer []byte + mutex sync.Mutex + dataPagesCount uint32 + indexPagesCount uint32 + dataFreeList *freelist.FreeList + indexFreeList *freelist.FreeList + atree *atree.Atree + logNumber int + dir string + w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск + wal *os.File + dataFile *os.File + indexFile *os.File + indexPageSize int + dataPageSize int input []any + dataPreparer *DataPreparer appendToWorkerQueue func(any) //lsn uint32 written int64 @@ -88,6 +88,8 @@ type Writer struct { } type WriterOptions struct { + IndexPageSize int + DataPageSize int Dir string LogNumber int // номер журнала AppendToWorkerQueue func(any) @@ -99,6 +101,12 @@ type WriterOptions struct { } func NewWriter(opt WriterOptions) (*Writer, error) { + if (opt.IndexPageSize % 2) != 0 { + return nil, errors.New("IndexPageSize must be multiple of 2") + } + if (opt.DataPageSize % 2) != 0 { + return nil, errors.New("DataPageSize must be multiple of 2") + } if opt.Dir == "" { return nil, errors.New("Dir option is required") } @@ -122,6 +130,8 @@ func NewWriter(opt WriterOptions) (*Writer, error) { } s := &Writer{ + indexPageSize: opt.IndexPageSize, + dataPageSize: opt.DataPageSize, dir: opt.Dir, appendToWorkerQueue: opt.AppendToWorkerQueue, dataFreeList: opt.DataFreeList, @@ -215,7 +225,7 @@ func (s *Writer) packAndWrite() (err error) { // } s.mutex.Unlock() - prepared := prepareData(input) + prepared := s.dataPreparer.Prepare(input) // 3. Пишу на диск WAL (append) n, err := s.wal.Write(prepared.Packet) @@ -253,7 +263,7 @@ func (s *Writer) packAndWrite() (err error) { metricsStateCh := make(chan MetricsState) s.appendToWorkerQueue(Changes{ - Records: prepared.ToWorker, + Records: prepared.WriteResults, MetricsCh: metricsStateCh, }) @@ -293,7 +303,7 @@ func (s *Writer) packAndWrite() (err error) { close(state.WaitCh) } else { s.appendToWorkerQueue(Changes{ - Records: prepared.ToWorker, + Records: prepared.WriteResults, }) } @@ -350,168 +360,54 @@ func (s *Writer) exit() { } } -type IndexPayloadToWrite struct { - PageNo uint32 - Payload []byte - ZeroLevel bool -} +// type IndexPayloadToWrite struct { +// PageNo uint32 +// Payload []byte +// ZeroLevel bool +// } -type DataPayloadToWrite struct { - PrevPageNo uint32 - Timestamps [][]byte // chunks - TimestampsSize int - Values [][]byte // chunks - ValuesSize int - PageNo uint32 +type PageToWrite struct { + PageNo uint32 + Content []byte } // writePagesToAtree - записує сторінки в .data та .index файли -func (s *Writer) writePagesToAtree(indexItems []IndexPayloadToWrite, dataItems []DataPayloadToWrite) (err error) { - for _, p := range dataItems { - atree.ChunksToDataPage(s.pageBuffer, atree.ChunksToDataPageReq{ - PrevPageNo: p.PrevPageNo, - Timestamps: p.Timestamps, - TimestampsSize: p.TimestampsSize, - Values: p.Values, - ValuesSize: p.ValuesSize, - }) +func (s *Writer) writePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) { + for _, p := range dataPages { + if len(p.Content) != s.dataPageSize { + return fmt.Errorf("wrong data page size: %d", len(p.Content)) + } var ( - off = (p.PageNo - 1) * atree.DataPageSize + off = int(p.PageNo-1) * s.dataPageSize n int ) - n, err = s.dataFile.WriteAt(s.pageBuffer, int64(off)) + n, err = s.dataFile.WriteAt(p.Content, int64(off)) if err != nil { return } - if n != atree.DataPageSize { - return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize) + if n != s.dataPageSize { + return fmt.Errorf("write %d instead of %d", n, s.dataPageSize) } } - for _, p := range indexItems { - // if len(p.Payload) > atree.MaxIndexPayload { - // return fmt.Errorf("wrong index payload size: %d", len(p.Payload)) - // } - indexPageBuffer := s.pageBuffer[:atree.IndexPageSize] - atree.DataToIndexPage(indexPageBuffer, p.Payload, p.ZeroLevel) + for _, p := range indexPages { + if len(p.Content) != s.indexPageSize { + return fmt.Errorf("wrong index page size: %d", len(p.Content)) + } var ( - off = (p.PageNo - 1) * atree.IndexPageSize + off = int(p.PageNo-1) * s.indexPageSize n int ) - n, err = s.indexFile.WriteAt(indexPageBuffer, int64(off)) + n, err = s.indexFile.WriteAt(p.Content, int64(off)) if err != nil { return } - if n != atree.IndexPageSize { - return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize) + if n != s.indexPageSize { + return fmt.Errorf("write %d instead of %d", n, s.indexPageSize) } } return nil } -type AppendMeasuresSummary struct { - MetricID uint32 - //TimestampsOffset int // (заповнені одразу) - //ValuesOffset int // (заповнені одразу) - // Timestamps [][]byte - // TimestampsSize int - // Values [][]byte - // ValuesSize int - LastPageNo uint32 - Index [][]byte - ResultCode int - WrittenCount int - ResultCh chan struct{} -} - -type AppendedMeasures struct { - MetricID uint32 - TimestampsOffset int // (заповнені одразу) - ValuesOffset int // (заповнені одразу) - Timestamps [][]byte - TimestampsSize int - Values [][]byte - ValuesSize int - IndexLevels []*atree.IndexLevel - DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page - ResultCode int - WrittenCount int - ResultCh chan struct{} -} - -func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) { - // fix write code - bin.WriteUint32(s.w, req.MetricID) - bin.WriteVarSize(s.w, len(req.DataPages)) - for idx, p := range req.DataPages { - bin.WriteUint32(s.w, p.PageNo) - var reused byte - if p.IsReused { - reused = 1 - } - s.w.WriteByte(reused) - bin.WriteUint32(s.w, p.PrevPageNo) - //bin.WriteUint32(s.w, p.Checksum) - var ( - timestampsOffset int - valuesOffset int - ) - if idx == 0 { - timestampsOffset = req.TimestampsOffset - valuesOffset = req.ValuesOffset - } - copyPayloadFromChunks(s.w, p.Timestamps, p.TimestampsSize, timestampsOffset) - copyPayloadFromChunks(s.w, p.Values, p.ValuesSize, valuesOffset) - } - bin.WriteUint16(s.w, uint16(req.TimestampsOffset)) - bin.WriteUint16(s.w, uint16(req.ValuesOffset)) - var ( - timestampsOffset int - valuesOffset int - ) - if len(req.DataPages) == 0 { - timestampsOffset = req.TimestampsOffset - valuesOffset = req.ValuesOffset - } - copyPayloadFromChunks(s.w, req.Timestamps, req.TimestampsSize, timestampsOffset) - copyPayloadFromChunks(s.w, req.Values, req.ValuesSize, valuesOffset) - if len(req.DataPages) == 0 { - return - } - // fix - levels reduced - bin.WriteVarSize(s.w, len(req.IndexLevels)) - for _, level := range req.IndexLevels { - bin.WriteVarSize(s.w, len(level.Filled)) - for idx, p := range level.Filled { - bin.WriteUint32(s.w, p.PageNo) - var reused byte - if p.IsReused { - reused = 1 - } - s.w.WriteByte(reused) - //bin.WriteUint32(s.w, p.Checksum) - records := p.Data - if idx == 0 { - records = p.Data[level.Offset:] - } - if (len(records) % recordSize) != 0 { - log.Fatalf("wrong records length: %d", len(records)) - } - bin.WriteVarSize(s.w, len(records)/recordSize) - s.w.Write(records) - } - records := level.Data - if len(level.Filled) == 0 { - records = level.Data[level.Offset:] - } - if (len(records) % recordSize) != 0 { - log.Fatalf("wrong unfilled records length: %d", len(records)) - } - bin.WriteVarSize(s.w, len(records)/recordSize) - s.w.Write(records) - } - return -} - // API // FIX - add