Files
qb/storage/page_ingester.go
2026-06-13 22:43:17 +00:00

252 lines
7.1 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package storage
import (
"errors"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/util"
)
// Задача метода:
// - для data сторінок отримати pageNo та reused, задати sizes та prevPageNo, порахувати CRC32
// - додати пари до index, згенерувати нові index сторінки
// levels []*IndexLevel, dataPages []*DataPage
// fix - reduce leves
type PageIngester struct {
getDataPageNumber func() (uint32, bool, error)
getIndexPageNumber func() (uint32, bool, error)
}
type PageIngesterOptions struct {
GetDataPageNumber func() (uint32, bool, error)
GetIndexPageNumber func() (uint32, bool, error)
}
func NewPageIngester(opt PageIngesterOptions) (*PageIngester, 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")
}
s := &PageIngester{
getIndexPageNumber: opt.GetIndexPageNumber,
getDataPageNumber: opt.GetDataPageNumber,
}
return s, nil
}
type IndexLevelTail struct {
Buffer []byte // розмір більший за кількість
RecordsCount int
}
type DataPayload struct {
Since uint32
Content []byte
TimestampsSize int
ValuesSize int
}
type IngestIn struct {
LastPageNo uint32
IndexLevelTails []IndexLevelTail
DataPages []DataPayload
}
type SealedDataPage struct {
PageNo uint32
Reused bool
Content []byte
}
type SealedIndexPage struct {
PageNo uint32
Reused bool
Content []byte
}
// 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
}
func (s *PageIngester) Ingest(in IngestIn) ([]SealedDataPage, []*SealedIndexLevel) {
var (
sealedDataPages []SealedDataPage
sealedIndexLevels []*SealedIndexLevel
)
for _, tail := range in.IndexLevelTails {
sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{
SkipRecords: tail.RecordsCount,
TailRecords: tail.Buffer,
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 = sealedDataPages[dataPageIdx-1].PageNo
}
pageNo, reused, err := s.getDataPageNumber()
if err != nil {
qb.Abort(qb.FailedGetPageNumber, err)
}
SealDataPage(SealDataPageIn{
Content: d.Content,
PrevPageNo: prevPageNo,
TimestampsSize: d.TimestampsSize,
ValuesSize: d.ValuesSize,
})
sealedDataPages = append(sealedDataPages, SealedDataPage{
Content: d.Content,
PageNo: pageNo,
Reused: reused,
})
upPageNo = pageNo
for {
if levelIdx == len(sealedIndexLevels) {
// новий root
sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{
TailRecords: appendIndexRecord(appendIndexRecordIn{
Timestamp: upTimestamp,
PageNo: upPageNo,
}),
TailRecordsCount: 1,
})
break
}
level := sealedIndexLevels[levelIdx]
if level.TailRecordsCount < maxRecordsOnIndexPage {
level.TailRecords = appendIndexRecord(appendIndexRecordIn{
Records: level.TailRecords,
RecordsCount: level.TailRecordsCount,
Timestamp: upTimestamp,
PageNo: upPageNo,
})
level.TailRecordsCount++
break
}
// в буфері немає місця для додавання нової пари, отже це
// заповнена сторінка
pageNo, reused, err := s.getIndexPageNumber()
if err != nil {
qb.Abort(qb.FailedGetPageNumber, err)
}
SealIndexPage(SealIndexPageIn{
Content: level.TailRecords,
RecordsCount: level.TailRecordsCount,
ZeroLevel: levelIdx == 0,
})
filled := SealedIndexPage{
Content: level.TailRecords,
PageNo: pageNo,
Reused: reused,
}
level.IndexPages = append(level.IndexPages, filled)
// новий tail на індексному рівні
level.TailRecords = appendIndexRecord(appendIndexRecordIn{
Timestamp: upTimestamp,
PageNo: upPageNo,
})
level.TailRecordsCount = 1
//
upPageNo = pageNo
upTimestamp, _ = bin.GetUint32(filled.Content)
levelIdx++
}
}
return sealedDataPages, sealedIndexLevels
}
// HELPERS
type SealDataPageIn struct {
Content []byte
PrevPageNo uint32
TimestampsSize int
ValuesSize int
}
func SealDataPage(in SealDataPageIn) (checksum uint32) {
bin.PutUint16(in.Content[timestampsSizeIdx:], uint16(in.TimestampsSize))
bin.PutUint16(in.Content[valuesSizeIdx:], uint16(in.ValuesSize))
bin.PutUint32(in.Content[prevPageIdx:], in.PrevPageNo)
checksum = util.CalculateCRC32(in.Content[:dataCRC32Idx])
bin.PutUint32(in.Content[dataCRC32Idx:], checksum)
return
}
type SealIndexPageIn struct {
Content []byte
RecordsCount int
ZeroLevel bool
}
func SealIndexPage(in SealIndexPageIn) (checksum uint32) {
bin.PutUint16(in.Content[indexRecordsCountIdx:], uint16(in.RecordsCount))
if in.ZeroLevel {
in.Content[isZeroLevelIdx] = 1
}
checksum = util.CalculateCRC32(in.Content[:indexCRC32Idx])
bin.PutUint32(in.Content[indexCRC32Idx:], checksum)
return
}
type appendIndexRecordIn struct {
Records []byte
RecordsCount int
Timestamp uint32
PageNo uint32
}
func appendIndexRecord(in appendIndexRecordIn) []byte {
//fmt.Printf("appendIndexRecord: % x\n", in.Records)
//fmt.Printf("timestamp: %d, pageNo: %d\n", in.Timestamp, in.PageNo)
if in.RecordsCount < maxRecordsOnIndexPage {
var (
pos int
buf = in.Records
)
if in.RecordsCount > 0 {
pos = in.RecordsCount * IndexRecordSize
//fmt.Printf("pos: %d\n", pos)
} else {
buf = make([]byte, IndexPageSize) // IndexPageIncSize
}
bin.PutUint32(buf[pos:], in.Timestamp)
bin.PutUint32(buf[pos+4:], in.PageNo)
//fmt.Printf("appendIndexRecord after: % x\n", buf)
return buf
}
qb.Abort(qb.NoSpaceOnIndexPage, nil)
return nil
}