This commit is contained in:
2026-06-05 19:43:01 +00:00
parent db5ecd3dfd
commit 05eea3b279
16 changed files with 1187 additions and 724 deletions

View File

@@ -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++
}
}
}

View File

@@ -2,4 +2,7 @@ package chunkenc
const eps = 0.000001
const flagLiteral = 128
const (
flagLiteral = 128
minBufferSize = 1024
)

View File

@@ -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)
}

View File

@@ -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)
}

View File

@@ -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)
}

View File

@@ -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)
}

View File

@@ -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) {

2
go.mod
View File

@@ -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
)

2
go.sum
View File

@@ -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=

1
qb.go
View File

@@ -86,6 +86,7 @@ const (
UnknownWorkerQueueItemBug AbortCode = 20
UnknownMetricWaitQueueItemBug AbortCode = 21
RepeatableLock AbortCode = 22
NoSpaceOnIndexPage AbortCode = 23
//
GetRecoveryRecipeFailed AbortCode = 26
LoadSnapshotFailed AbortCode = 27

207
txlog/data_preparer.go Normal file
View File

@@ -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{}
}

View File

@@ -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
// }

279
txlog/page_preparer.go Normal file
View File

@@ -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,
}
}

View File

@@ -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,
}
}

306
txlog/wal_records.go Normal file
View File

@@ -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
}

View File

@@ -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