Files
qb/storage/replay_metric.go
2026-06-19 07:25:46 +03:00

237 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 (
"fmt"
"gordenko.dev/dima/qb"
)
type ReplayMetric struct {
MetricType qb.MetricType
FracDigits byte
LastPageNo uint32
Buf []byte
TimestampsSize int
ValuesSize int
IndexLevelTails []IndexLevelTail // root - last element
}
// return occupied after write
func writeTimestampsWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int {
// timestamps writes right-to-left
pos := DataPagePayloadSize - occupied
// timestamps rewind offset moves pos left-to-right
pos += rewindOffset
pos -= len(payload)
copy(buf[pos:], payload)
return DataPagePayloadSize - pos
}
func writeValuesWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int {
// values writes left-to-right
pos := occupied
// values rewind offset moves pos right-to-left
pos -= rewindOffset
copy(buf[pos:], payload)
return pos + len(payload)
}
func (s *ReplayMetric) MeasuresAppend(rec MeasuresAppendRecord) {
s.TimestampsSize = writeTimestampsWithRewind(s.Buf, s.TimestampsSize,
rec.Timestamps, rec.TimestampsRewindOffset)
s.ValuesSize = writeValuesWithRewind(s.Buf, s.ValuesSize,
rec.Values, rec.ValuesRewindOffset)
}
type ReplayMeasuresAppendWithGrowResult struct {
IndexPagesToRewrite []PageToWrite
DataPagesToRewrite []PageToWrite
ReusedIndexPagesCount int
ReusedDataPagesCount int
}
func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (_ ReplayMeasuresAppendWithGrowResult) {
//fmt.Printf("%#v\n", rec)
var (
reusedIndexPages int
reusedDataPages int
indexPages []PageToWrite
dataPages []PageToWrite
)
// Додати в index level tails недостаючі дані, або замінити
for levelIdx, change := range rec.ChangedIndexLevels {
var (
head = change.IndexPageTail
newRecordsCount = len(change.TailRecords) / IndexRecordSize
)
// 1. набиваю сторінки для перезапису в index файлі
if isLastPacket {
if head != nil {
indexPages = append(indexPages, s.composeHeadIndexPage(levelIdx, head))
}
for _, x := range change.IndexPages {
indexPages = append(indexPages, PageToWrite{
PageNo: x.PageNo,
Content: x.Content,
})
}
}
// 2. Вношу зміни в поточні індексні хвости
if levelIdx < len(s.IndexLevelTails) {
level := s.IndexLevelTails[levelIdx]
if head != nil {
// попередній хвіст індексного рівня перетворився на head сторінку,
// отже TailRecords - це новий хвіст
clear(level.Buffer)
copy(level.Buffer, change.TailRecords)
level.RecordsCount = newRecordsCount
} else {
// TailRecords - це нові дані, які треба додати
pos := level.RecordsCount * IndexRecordSize
copy(level.Buffer[pos:], change.TailRecords)
level.RecordsCount += newRecordsCount
}
s.IndexLevelTails[levelIdx] = level
} else {
// додаю новий індексний рівень
buf := make([]byte, IndexPageSize)
copy(buf, change.TailRecords)
//
s.IndexLevelTails = append(s.IndexLevelTails, IndexLevelTail{
Buffer: buf,
RecordsCount: newRecordsCount,
})
}
// 3. рахую кількість reused індексних сторінок
if head != nil {
if head.Reused {
reusedIndexPages++
}
}
for _, x := range change.IndexPages {
if x.Reused {
reusedIndexPages++
}
}
}
// 4. набиваю сторінки для перезапису в data файлі
if isLastPacket {
dataPages = append(dataPages, s.composeHeadDataPage(rec.DataPageTail))
for _, x := range rec.DataPages {
dataPages = append(dataPages, PageToWrite{
PageNo: x.PageNo,
Content: x.Content,
})
}
}
// рахую кількість reused дата сторінок
if rec.DataPageTail.Reused {
reusedDataPages++
}
for _, x := range rec.DataPages {
if x.Reused {
reusedDataPages++
}
}
// HeadDataPage є полюбому, оскільки це транзакція із мінімум однією заповненою
// data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня.
clear(s.Buf)
copy(s.Buf, rec.TailValues)
s.ValuesSize = len(rec.TailValues)
pos := DataPagePayloadSize - len(rec.TailTimestamps)
copy(s.Buf[pos:], rec.TailTimestamps)
s.TimestampsSize = len(rec.TailTimestamps)
if len(rec.DataPages) == 0 {
s.LastPageNo = rec.DataPageTail.PageNo
} else {
s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo
}
return ReplayMeasuresAppendWithGrowResult{
IndexPagesToRewrite: indexPages,
DataPagesToRewrite: dataPages,
ReusedIndexPagesCount: reusedIndexPages,
ReusedDataPagesCount: reusedDataPages,
}
}
// HELPERS
func (s *ReplayMetric) composeHeadIndexPage(levelIdx int, head *IndexPageTail) PageToWrite {
var (
page = make([]byte, IndexPageSize)
pos int
prevCount int
)
//fmt.Printf("composeHeadIndexPage: levels=%d\n", len(s.IndexLevelTails))
//fmt.Printf("head: % x\n", head.Records)
if levelIdx < len(s.IndexLevelTails) {
level := s.IndexLevelTails[levelIdx]
//fmt.Printf("level (%d): % x\n", level.RecordsCount, level.Buffer)
// розраховую pos, з якого буду дописувати хвіст
pos = level.RecordsCount * IndexRecordSize
// створюю копію сторінки
copy(page, level.Buffer[:pos]) // поточні дані
prevCount = level.RecordsCount
} else {
//
s.IndexLevelTails = append(s.IndexLevelTails, IndexLevelTail{
Buffer: page,
RecordsCount: len(head.Records) / IndexRecordSize,
})
}
copy(page[pos:], head.Records)
// запечатати сторінку
calculatedCRC := SealIndexPage(SealIndexPageIn{
Content: page,
RecordsCount: prevCount + len(head.Records)/IndexRecordSize,
ZeroLevel: levelIdx == 0,
})
// перевірка CRC
if calculatedCRC != head.CRC32 {
qb.Abort(qb.WALReplayFailed,
fmt.Errorf("calculated CRC %d not equal expected %d of head page %d on index level %d",
calculatedCRC, head.CRC32, head.PageNo, levelIdx))
}
return PageToWrite{
PageNo: head.PageNo,
Content: page,
}
}
func (s *ReplayMetric) composeHeadDataPage(head DataPageTail) PageToWrite {
// copy page
page := make([]byte, DataPageSize)
copy(page, s.Buf)
// append data
timestampsSize := writeTimestampsWithRewind(page, s.TimestampsSize,
head.Timestamps, head.TimestampsRewindOffset)
valuesSize := writeValuesWithRewind(page, s.ValuesSize,
head.Values, head.ValuesRewindOffset)
// запечатати сторінку
calculatedCRC := SealDataPage(SealDataPageIn{
Content: page,
PrevPageNo: head.PrevPageNo,
TimestampsSize: timestampsSize,
ValuesSize: valuesSize,
})
// перевірка CRC
if calculatedCRC != head.CRC32 {
qb.Abort(qb.WALReplayFailed,
fmt.Errorf("calculated CRC %d not equal expected %d of head data page %d",
calculatedCRC, head.CRC32, head.PageNo))
}
return PageToWrite{
PageNo: head.PageNo,
Content: page,
}
}