Files
qb/storage/wal_records.go

656 lines
14 KiB
Go
Raw Permalink Normal View History

2026-06-10 06:18:45 +03:00
package storage
import (
"io"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
)
2026-06-13 08:01:42 +03:00
type MetricAddRecord struct {
MetricID uint32
MetricType qb.MetricType
FracDigits byte
2026-06-10 06:18:45 +03:00
}
2026-06-13 08:01:42 +03:00
func (s MetricAddRecord) Pack(w io.Writer) {
arr := []byte{
CodeMetricAdd,
0, 0, 0, 0, //
byte(s.MetricType),
s.FracDigits,
}
bin.PutUint32(arr[1:], s.MetricID)
w.Write(arr)
}
func (s *MetricAddRecord) Parse(src io.Reader) (err error) {
arr, err := bin.ReadN(src, 6)
if err != nil {
return
}
s.MetricID, _ = bin.GetUint32(arr)
s.MetricType = qb.MetricType(arr[4])
s.FracDigits = arr[5]
return nil
}
type MetricDeleteRecord struct {
MetricID uint32
FreeIndexPages []uint32
FreeDataPages []uint32
}
func (s MetricDeleteRecord) Pack(w io.Writer) {
arr := []byte{
CodeMetricDelete,
0, 0, 0, 0, // metricID
}
bin.PutUint32(arr[1:], s.MetricID)
w.Write(arr)
packPageNumbers(w, s.FreeIndexPages)
packPageNumbers(w, s.FreeDataPages)
}
func (s *MetricDeleteRecord) Parse(src io.Reader) (err error) {
s.MetricID, err = bin.ReadUint32(src)
if err != nil {
return
}
s.FreeIndexPages, err = parsePageNumbers(src)
if err != nil {
return
}
s.FreeDataPages, err = parsePageNumbers(src)
if err != nil {
return
}
return
}
type DataPageTail struct {
PageNo uint32
Reused bool
PrevPageNo uint32
CRC32 uint32
TimestampsRewindOffset int // offset on 1st page
Timestamps []byte
ValuesRewindOffset int // offset on 1st page
Values []byte
}
type IndexPageTail struct {
PageNo uint32
Reused bool
CRC32 uint32
Records []byte // offset не потрібен, оскільки Records додаються в кінець
2026-06-10 06:18:45 +03:00
}
type ChangedIndexLevel struct {
2026-06-13 08:01:42 +03:00
IndexPageTail *IndexPageTail
2026-06-10 06:18:45 +03:00
IndexPages []SealedIndexPage
TailRecords []byte // payload only
}
type MeasuresAppendWithGrowRecord struct {
MetricID uint32
2026-06-13 08:01:42 +03:00
DataPageTail DataPageTail
2026-06-10 06:18:45 +03:00
DataPages []SealedDataPage
TailTimestamps []byte // data level tail
TailValues []byte // data level tail
2026-06-13 08:01:42 +03:00
ChangedIndexLevels []ChangedIndexLevel
2026-06-10 06:18:45 +03:00
}
/*
Format appended measures:
1b - tx type
4b - metricID
4b - pageNo
1b - reused
4b - prevPageNo
4b - page crc32
2026-06-13 08:01:42 +03:00
1b - timestamps & values offsets
Nb - (varsize) timestamps tail size on 1st filled data page
2026-06-10 06:18:45 +03:00
Nb - timestamps tail on 1st filled data page
Nb - (varsize) values tail size on 1st filled data page
Nb - values tail on 1st filled data page
2026-06-13 08:01:42 +03:00
Nb - (varsize) full data pages qty
2026-06-10 06:18:45 +03:00
[
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
]
*/
2026-06-13 08:01:42 +03:00
func (s MeasuresAppendWithGrowRecord) Pack(w io.Writer) {
arr := []byte{
CodeMeasuresAppendWithGrow,
0, 0, 0, 0, //
0, 0, 0, 0, // tail pageNo
0, // tail reused
0, 0, 0, 0, // tail prevPageNo
0, 0, 0, 0, // tail CRC32
0, // tail offsets
}
bin.PutUint32(arr[1:], s.MetricID)
// data page tail
tail := s.DataPageTail
bin.PutUint32(arr[5:], tail.PageNo)
if tail.Reused {
arr[9] = 1
}
bin.PutUint32(arr[10:], tail.PrevPageNo)
bin.PutUint32(arr[14:], tail.CRC32)
arr[18] = byte(tail.TimestampsRewindOffset) | (byte(tail.ValuesRewindOffset) << 4)
w.Write(arr)
bin.WriteVarSized(w, tail.Timestamps)
bin.WriteVarSized(w, tail.Values)
2026-06-10 06:18:45 +03:00
// data pages
bin.WriteVarSize(w, len(s.DataPages))
for _, p := range s.DataPages {
bin.WriteUint32(w, p.PageNo)
bin.WriteBool(w, p.Reused)
w.Write(p.Content)
}
// data tail
2026-06-13 08:01:42 +03:00
bin.WriteVarSized(w, s.TailTimestamps)
bin.WriteVarSized(w, s.TailValues)
2026-06-10 06:18:45 +03:00
// changed levels
bin.WriteVarSize(w, len(s.ChangedIndexLevels))
for _, level := range s.ChangedIndexLevels {
2026-06-13 08:01:42 +03:00
// index page tail
tail := level.IndexPageTail
if tail != nil {
bin.WriteBool(w, true)
bin.WriteUint32(w, tail.PageNo)
bin.WriteBool(w, tail.Reused)
bin.WriteUint32(w, tail.CRC32)
2026-06-13 22:43:17 +00:00
bin.WriteVarSize(w, len(tail.Records)/IndexRecordSize)
2026-06-13 08:01:42 +03:00
w.Write(tail.Records)
} else {
bin.WriteBool(w, false)
}
2026-06-10 06:18:45 +03:00
// full index pages
2026-06-13 08:01:42 +03:00
bin.WriteVarSize(w, len(level.IndexPages))
2026-06-10 06:18:45 +03:00
for _, p := range level.IndexPages {
bin.WriteUint32(w, p.PageNo)
bin.WriteBool(w, p.Reused)
w.Write(p.Content)
}
// level tail
2026-06-13 22:43:17 +00:00
bin.WriteVarSize(w, len(level.TailRecords)/IndexRecordSize)
2026-06-10 06:18:45 +03:00
w.Write(level.TailRecords)
}
}
2026-06-13 08:01:42 +03:00
func (s *MeasuresAppendWithGrowRecord) Parse(r io.Reader) (err error) {
2026-06-10 06:18:45 +03:00
s.MetricID, err = bin.ReadUint32(r)
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.PageNo, err = bin.ReadUint32(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.Reused, err = bin.ReadBool(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.PrevPageNo, err = bin.ReadUint32(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.CRC32, err = bin.ReadUint32(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
offsets, err := bin.ReadByte(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.TimestampsRewindOffset = int(offsets & 15) // low 4 bits
s.DataPageTail.ValuesRewindOffset = int(offsets >> 4) // high 4 bits
//
s.DataPageTail.Timestamps, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.DataPageTail.Values, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
// data pages
dataPagesQty, err := bin.ReadVarSize(r)
if err != nil {
return
}
for range dataPagesQty {
2026-06-13 08:01:42 +03:00
var p SealedDataPage
2026-06-10 06:18:45 +03:00
p.PageNo, err = bin.ReadUint32(r)
if err != nil {
return
}
p.Reused, err = bin.ReadBool(r)
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
p.Content, err = bin.ReadN(r, DataPageSize)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
s.DataPages = append(s.DataPages, p)
}
2026-06-13 08:01:42 +03:00
// // data tail
s.TailTimestamps, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.TailValues, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
// changed levels
changedIndexLevelsCount, err := bin.ReadVarSize(r)
if err != nil {
return
}
for range changedIndexLevelsCount {
var (
2026-06-13 08:01:42 +03:00
hasIndexPageTail bool
recordsCount int
indexPagesCount int
level ChangedIndexLevel
2026-06-10 06:18:45 +03:00
)
2026-06-13 08:01:42 +03:00
hasIndexPageTail, err = bin.ReadBool(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
if hasIndexPageTail {
//index page tail
var tail IndexPageTail
tail.PageNo, err = bin.ReadUint32(r)
if err != nil {
return
}
tail.Reused, err = bin.ReadBool(r)
if err != nil {
return
}
tail.CRC32, err = bin.ReadUint32(r)
if err != nil {
return
}
recordsCount, err = bin.ReadVarSize(r)
if err != nil {
return
}
2026-06-13 22:43:17 +00:00
tail.Records, err = bin.ReadN(r, recordsCount*IndexRecordSize)
2026-06-13 08:01:42 +03:00
if err != nil {
return
}
level.IndexPageTail = &tail
2026-06-10 06:18:45 +03:00
}
//
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
}
2026-06-13 08:01:42 +03:00
p.Content, err = bin.ReadN(r, IndexPageSize)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
level.IndexPages = append(level.IndexPages, p)
}
recordsCount, err = bin.ReadVarSize(r)
if err != nil {
return
}
2026-06-13 22:43:17 +00:00
level.TailRecords, err = bin.ReadN(r, recordsCount*IndexRecordSize)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
s.ChangedIndexLevels = append(s.ChangedIndexLevels, level)
}
return
}
type MeasuresAppendRecord struct {
2026-06-12 14:21:14 +00:00
MetricID uint32
TimestampsRewindOffset int
Timestamps []byte
ValuesRewindOffset int
Values []byte
2026-06-10 06:18:45 +03:00
}
2026-06-13 08:01:42 +03:00
func (s MeasuresAppendRecord) Pack(w io.Writer) {
arr := []byte{
CodeMeasuresAppend,
0, 0, 0, 0, // metricID
0, // offsets
}
bin.PutUint32(arr[1:], s.MetricID)
arr[5] = byte(s.TimestampsRewindOffset) | (byte(s.ValuesRewindOffset) << 4)
w.Write(arr)
bin.WriteVarSized(w, s.Timestamps)
bin.WriteVarSized(w, s.Values)
2026-06-10 06:18:45 +03:00
}
2026-06-13 08:01:42 +03:00
func (s *MeasuresAppendRecord) Parse(r io.Reader) (err error) {
2026-06-10 06:18:45 +03:00
s.MetricID, err = bin.ReadUint32(r)
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
offsets, err := bin.ReadByte(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.TimestampsRewindOffset = int(offsets & 15) // low 4 bits
s.ValuesRewindOffset = int(offsets >> 4) // high 4 bits
//
s.Timestamps, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.Values, err = bin.ReadVarSized(r)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
return
}
2026-06-13 08:01:42 +03:00
// fix - add changed pages
type MeasuresDeleteRecord struct {
2026-06-10 06:18:45 +03:00
MetricID uint32
FreeIndexPages []uint32
FreeDataPages []uint32
}
2026-06-13 08:01:42 +03:00
func (s MeasuresDeleteRecord) Pack(w io.Writer) {
2026-06-10 06:18:45 +03:00
arr := []byte{
CodeMeasuresDelete,
0, 0, 0, 0,
}
bin.PutUint32(arr[1:], s.MetricID)
w.Write(arr)
2026-06-13 08:01:42 +03:00
packPageNumbers(w, s.FreeIndexPages)
packPageNumbers(w, s.FreeDataPages)
2026-06-10 06:18:45 +03:00
}
2026-06-13 08:01:42 +03:00
func (s *MeasuresDeleteRecord) Parse(src io.Reader) (err error) {
2026-06-10 06:18:45 +03:00
s.MetricID, err = bin.ReadUint32(src)
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.FreeIndexPages, err = parsePageNumbers(src)
2026-06-10 06:18:45 +03:00
if err != nil {
return
}
2026-06-13 08:01:42 +03:00
s.FreeDataPages, err = parsePageNumbers(src)
if err != nil {
return
2026-06-10 06:18:45 +03:00
}
2026-06-13 08:01:42 +03:00
return
2026-06-10 06:18:45 +03:00
}
// type AppendedMeasure struct {
// MetricID uint32
// Timestamp uint32
// Value float64
// }
// func (s AppendedMeasure) Pack(w io.Writer) {
// arr := []byte{
// CodeAppendedMeasure,
// 0, 0, 0, 0, // metricID
// 0, 0, 0, 0, // timestamp
// 0, 0, 0, 0, 0, 0, 0, 0, // value
// }
// bin.PutUint32(arr[1:], s.MetricID)
// bin.PutUint32(arr[5:], s.Timestamp)
// bin.PutFloat64(arr[9:], s.Value)
// w.Write(arr)
// }
// func (s *AppendedMeasure) Parse(src io.Reader) (err error) {
// arr, err := bin.ReadN(src, 16)
// if err != nil {
// return
// }
// s.MetricID, _ = bin.GetUint32(arr[0:])
// s.Timestamp, _ = bin.GetUint32(arr[4:])
// s.Value, _ = bin.GetFloat64(arr[8:])
// return nil
// }
// type AppendedMeasures struct {
// MetricID uint32
// Measures []proto.Measure
// }
// func (s AppendedMeasures) Pack(w io.Writer) {
// arr := []byte{
// CodeAppendedMeasures,
// 0, 0, 0, 0, // metricID
// 0, 0, // qty
// }
// bin.PutUint32(arr[1:], s.MetricID)
// bin.PutUint16(arr[5:], uint16(len(s.Measures)))
// w.Write(arr)
// for _, measure := range s.Measures {
// bin.WriteUint32(w, measure.Timestamp)
// bin.WriteFloat64(w, measure.Value)
// }
// }
// func (s *AppendedMeasures) Parse(src io.Reader) (err error) {
// s.MetricID, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// qty, err := bin.ReadUint16(src)
// if err != nil {
// return
// }
// for range qty {
// var measure proto.Measure
// measure.Timestamp, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// measure.Value, err = bin.ReadFloat64(src)
// if err != nil {
// return
// }
// s.Measures = append(s.Measures, measure)
// }
// return nil
// }
// type AppendedPagesReq struct {
// Legs []atree.PathLeg
// MetricID uint32
// Timestamp uint32 // last measure
// Value float64 // last measure
// //RootPageNo uint32
// LastPageNo uint32
// Pages []atree.NotLinkedDataPage
// TimestampsChunks [][]byte
// TimestampsSize uint16
// ValuesChunks [][]byte
// ValuesSize uint16
// }
/*
4b metricID
4b timestamp
8b value
4b LastPageNo (0 if no new root page)
4b newRootPageNo (0 if no new root page)
Nb - varsize (pages qty)
[
4b pageNo
1b isReused
Xb page (page size)
]
2b - timestamps size (not filled page)
Nb - timestamps payload
2b - values size (not filled page)
Nb - values payload
*/
// type AppendedPages struct {
// MetricID uint32
// Timestamp uint32 // last measure
// Value float64 // last measure
// NewRootPageNo uint32
// LastPageNo uint32
// Pages []atree.PageToWrite
// TimestampsChunks [][]byte
// TimestampsSize int
// ValuesChunks [][]byte
// ValuesSize int
// }
// func (s *AppendedPages) Pack(w io.Writer) {
// w.Write([]byte{
// CodeAppendedPages,
// })
// bin.WriteUint32(w, s.MetricID)
// bin.WriteUint32(w, s.Timestamp)
// bin.WriteFloat64(w, s.Value)
// bin.WriteUint32(w, s.LastPageNo)
// bin.WriteUint32(w, s.NewRootPageNo)
// //
// bin.WriteVarSize(w, len(s.Pages))
// for _, p := range s.Pages {
// bin.WriteUint32(w, p.PageNo)
// if p.IsReused {
// w.Write([]byte{1})
// } else {
// w.Write([]byte{0})
// }
// w.Write(p.Data)
// }
// // timestamps
// writeChunkedPayloadTo(s.TimestampsChunks, int(s.TimestampsSize), w)
// // values
// writeChunkedPayloadTo(s.ValuesChunks, int(s.ValuesSize), w)
// }
// func (s *AppendedPages) Parse(src io.Reader) (err error) {
// s.MetricID, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// s.Timestamp, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// s.Value, err = bin.ReadFloat64(src)
// if err != nil {
// return
// }
// s.LastPageNo, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// s.NewRootPageNo, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// pagesQty, err := bin.ReadVarSize(src)
// if err != nil {
// return
// }
// for range pagesQty {
// var page atree.PageToWrite
// page.PageNo, err = bin.ReadUint32(src)
// if err != nil {
// return
// }
// page.IsReused, err = bin.ReadBool(src)
// if err != nil {
// return
// }
// page.Data, err = bin.ReadN(src, atree.PageSize)
// if err != nil {
// return
// }
// s.Pages = append(s.Pages, page)
// }
// s.TimestampsChunks, s.TimestampsSize, err = readChunkedPayload(src)
// if err != nil {
// return
// }
// s.ValuesChunks, s.ValuesSize, err = readChunkedPayload(src)
// if err != nil {
// return
// }
// return nil
// }
2026-06-13 08:01:42 +03:00
// HELPERS
func packPageNumbers(w io.Writer, pageNumbers []uint32) {
bin.WriteVarSize(w, len(pageNumbers))
for _, pageNo := range pageNumbers {
bin.WriteUint32(w, pageNo)
}
}
func parsePageNumbers(src io.Reader) (pageNumbers []uint32, err error) {
count, err := bin.ReadVarSize(src)
if err != nil {
return
}
for range count {
var pageNo uint32
pageNo, err = bin.ReadUint32(src)
if err != nil {
return
}
pageNumbers = append(pageNumbers, pageNo)
}
return
}