This commit is contained in:
2026-06-09 08:16:16 +03:00
parent 964fe4b4a8
commit 7007b1b6fd
18 changed files with 1334 additions and 1026 deletions

View File

@@ -109,7 +109,7 @@ func (s *DataPreparer) Prepare(input []any) PreparedData {
})
indexLevelTails = append(indexLevelTails, IndexLevelTail{
Records: level.TailRecords,
Buffer: level.TailRecords,
RecordsCount: level.TailRecordsCount,
})
}
@@ -181,9 +181,9 @@ func (s *DataPreparer) Reset() {
type AppendedMeasures struct {
MetricID uint32
LastPageNo uint32
TimestampsOffset int
Timestamps []byte // offset on 1st page
ValuesOffset int // offset on 1st page
TimestampsOffset int // offset on 1st page
Timestamps []byte
ValuesOffset int // offset on 1st page
Values []byte
IndexLevelTails []IndexLevelTail
DataPages []DataPayload

View File

@@ -122,7 +122,7 @@ func (s *PagePreparer) sealIndexPage(in sealIndexPageIn) (checksum uint32) {
}
type IndexLevelTail struct {
Records []byte // розмір більший за кількість
Buffer []byte // розмір більший за кількість
RecordsCount int
}
@@ -185,7 +185,7 @@ func (s *PagePreparer) SealPages(in SealPagesIn) (result SealResult) {
for _, tail := range in.IndexLevelTails {
sealedIndexLevels = append(sealedIndexLevels, &SealedIndexLevel{
SkipRecords: tail.RecordsCount,
TailRecords: tail.Records,
TailRecords: tail.Buffer,
TailRecordsCount: tail.RecordsCount,
})
}

View File

@@ -5,16 +5,20 @@ import (
"bytes"
"errors"
"fmt"
"hash/crc32"
"io"
"os"
"gordenko.dev/dima/qb/bin"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb/util"
)
const readBufferSize = 8 * 1024 * 1024
type Reader struct {
file *os.File
reader *bufio.Reader
file *os.File
reader *bufio.Reader
current []any
next []any
}
type ReaderOptions struct {
@@ -37,7 +41,7 @@ func NewReader(opt ReaderOptions) (*Reader, error) {
return &Reader{
file: file,
reader: bufio.NewReaderSize(file, 1024*1024),
reader: bufio.NewReaderSize(file, readBufferSize),
}, nil
}
@@ -45,42 +49,62 @@ func (s *Reader) Close() {
s.file.Close()
}
func (s *Reader) ReadPacket() (uint32, []any, bool, error) {
prefix := make([]byte, packetPrefixSize)
n, err := s.reader.Read(prefix)
if err != nil {
if err == io.EOF && n == 0 {
return 0, nil, true, nil
} else {
return 0, nil, false, fmt.Errorf("read packet prefix: %s", err)
func (s *Reader) NextPacket() ([]any, bool, error) {
var err error
if s.current == nil {
s.current, err = s.readPacket()
if err != nil {
return nil, false, err
}
if len(s.current) > 0 {
s.next, err = s.readPacket()
if err != nil {
return nil, false, err
}
}
}
length := bin.GetUint32(prefix[lengthIdx:])
storedCRC := bin.GetUint32(prefix[checksumIdx:])
lsn := bin.GetUint32(prefix[lsnIdx:])
current := s.current
done := s.next == nil
body, err := bin.ReadN(s.reader, int(length))
s.current = s.next
s.next, err = s.readPacket()
if err != nil {
return 0, nil, false, fmt.Errorf("read packet body: %s", err)
return nil, false, err
}
return current, done, nil
}
func (s *Reader) readPacket() (_ []any, err error) {
payloadSize, err := bin.ReadVarSize(s.reader)
if err != nil {
if err == io.EOF {
return nil, nil
} else {
return
}
}
hasher := crc32.NewIEEE()
hasher.Write(prefix[lsnIdx:])
storedCRC, err := bin.ReadUint32(s.reader)
if err != nil {
return
}
body, err := bin.ReadN(s.reader, int(payloadSize))
if err != nil {
return
}
hasher := util.NewHasher()
hasher.Write(body)
calculatedCRC := hasher.Sum32()
if calculatedCRC != storedCRC {
return 0, nil, false, fmt.Errorf("stored CRC %d != calculated CRC %d",
storedCRC, calculatedCRC)
if calculatedCRC == storedCRC {
return s.parseRecords(body)
}
records, err := s.parseRecords(body)
if err != nil {
return 0, nil, false, err
}
return lsn, records, false, nil
//s.logger.Printf("stored CRC %d != calculated CRC %d", storedCRC, calculatedCRC)
return nil, nil
}
func (s *Reader) parseRecords(body []byte) ([]any, error) {

View File

@@ -1,122 +0,0 @@
package txlog
import (
"fmt"
"io"
"os"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/util"
)
func LoadSnapshot(fileName string) (snapshot Snapshot, err error) {
var (
hasher = util.NewHasher()
header = make([]byte, metricHeaderSize)
body = make([]byte, atree.PageSize)
)
file, err := os.Open(fileName)
if err != nil {
return
}
src := io.TeeReader(file, hasher)
metricsQty, err := bin.ReadVarSize(src)
if err != nil {
return
}
for range metricsQty {
var metric Metric
err = bin.ReadNInto(src, header)
if err != nil {
return
}
metric.MetricID, _ = bin.GetUint32(header[0:])
metric.MetricType = qb.MetricType(header[4])
metric.FracDigits = header[5]
metric.LastPageNo, _ = bin.GetUint32(header[6:])
metric.Since, _ = bin.GetUint32(header[10:])
metric.SinceValue, _ = bin.GetFloat64(header[14:])
metric.Until, _ = bin.GetUint32(header[22:])
metric.UntilValue, _ = bin.GetFloat64(header[26:])
tSize, _ := bin.GetUint16(header[34:])
vSize, _ := bin.GetUint16(header[36:])
buf := body[:tSize]
err = bin.ReadNInto(src, buf)
if err != nil {
return
}
// FIX
// metric.Timestamps = chunkenc.NewReverseTimeDeltaCompressor(
// conbuf.NewFromBuffer(buf),
// int(tSize),
// )
buf = body[:vSize]
err = bin.ReadNInto(src, buf)
if err != nil {
return
}
// FIX
// if metric.MetricType == qb.Cumulative {
// metric.Values = chunkenc.NewReverseCumulativeDeltaCompressor(
// conbuf.NewFromBuffer(buf),
// int(vSize),
// metric.FracDigits,
// )
// } else {
// metric.Values = chunkenc.NewReverseInstantDeltaCompressor(
// conbuf.NewFromBuffer(buf),
// int(vSize),
// metric.FracDigits,
// )
// }
// s.metrics[metricID] = &metric
}
// FIX
// err = restoreFreeList(s.dataFreeList, src)
// if err != nil {
// return fmt.Errorf("restore dataFreeList: %s", err)
// }
// FIX
// err = restoreFreeList(s.indexFreeList, src)
// if err != nil {
// return fmt.Errorf("restore indexFreeList: %s", err)
// }
calculatedChecksum := hasher.Sum32()
writtenChecksum, err := bin.ReadUint32(file)
if err != nil {
return
}
if calculatedChecksum != writtenChecksum {
err = fmt.Errorf("calculated checksum %d not equal written checksum %d",
calculatedChecksum, writtenChecksum)
return
}
return
}
// serialized + frozen pages count
// func restoreFreeList(freeList *freelist.FreeList, src io.Reader) error {
// size, err := bin.ReadVarSize(src)
// if err != nil {
// return err
// }
// serialized, err := bin.ReadN(src, size)
// if err != nil {
// return err
// }
// freeList.Restore(serialized)
// return nil
// }

View File

@@ -5,7 +5,6 @@ import (
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/conbuf"
)
var (
@@ -79,49 +78,6 @@ func (s *DeletedMetric) Parse(src io.Reader) (err error) {
return nil
}
func writeChunkedPayloadTo(chunks [][]byte, size int, w io.Writer) {
bin.WriteVarSize(w, size)
for _, chunk := range chunks {
if size >= len(chunk) {
w.Write(chunk)
size -= len(chunk)
} else {
w.Write(chunk[:size])
return
}
}
}
func readChunkedPayload(src io.Reader) (_ [][]byte, _ int, err error) {
var (
chunks [][]byte
size int
remainingSize = size
)
size, err = bin.ReadVarSize(src)
if err != nil {
return
}
for remainingSize > 0 {
var chunk []byte
if remainingSize >= conbuf.ChunkSize {
chunk, err = bin.ReadN(src, conbuf.ChunkSize)
if err != nil {
return
}
} else {
chunk = make([]byte, conbuf.ChunkSize)
err = bin.ReadNInto(src, chunk[:remainingSize])
if err != nil {
return
}
}
chunks = append(chunks, chunk)
remainingSize -= conbuf.ChunkSize
}
return chunks, size, nil
}
// fix - add changed pages
type DeletedMeasures struct {
MetricID uint32

View File

@@ -1,21 +1,18 @@
package txlog
import (
"bufio"
"bytes"
"errors"
"fmt"
"io"
"log"
"math"
"os"
"path/filepath"
"sync"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/freelist"
"gordenko.dev/dima/qb/util"
)
const (
@@ -52,23 +49,27 @@ func JoinChangesFileName(dir string, logNumber int) string {
}
type MetricsState struct {
Metrics []Metric
WaitCh chan struct{}
//Metrics []Metric
WaitCh chan struct{}
}
type Changes struct {
Records []any
MetricsCh chan MetricsState
Records []any
SnapshotNumberCh chan int // log number
FrozenIndexPagesCount int
IndexPageNumbers []uint32
FrozenDataPagesCount int
DataPageNumbers []uint32
}
type Writer struct {
mutex sync.Mutex
dataPagesCount uint32
indexPagesCount uint32
dataFreeList *freelist.FreeList
indexFreeList *freelist.FreeList
atree *atree.Atree
logNumber int
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
@@ -91,7 +92,7 @@ type WriterOptions struct {
IndexPageSize int
DataPageSize int
Dir string
LogNumber int // номер журнала
SnapshotNumber int // номер журнала
AppendToWorkerQueue func(any)
DataFreeList *freelist.FreeList
IndexFreeList *freelist.FreeList
@@ -137,7 +138,6 @@ func NewWriter(opt WriterOptions) (*Writer, error) {
dataFreeList: opt.DataFreeList,
indexFreeList: opt.IndexFreeList,
atree: opt.Atree,
logNumber: opt.LogNumber,
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
signalCh: make(chan struct{}, 1),
@@ -145,25 +145,16 @@ func NewWriter(opt WriterOptions) (*Writer, error) {
var err error
if opt.LogNumber > 0 {
s.wal, err = os.OpenFile(
JoinChangesFileName(opt.Dir, s.logNumber),
os.O_APPEND|os.O_WRONLY,
filePerm,
)
if err != nil {
return nil, err
}
} else {
s.logNumber = 1
s.wal, err = os.OpenFile(
JoinChangesFileName(opt.Dir, s.logNumber),
os.O_CREATE|os.O_WRONLY,
filePerm,
)
if err != nil {
return nil, err
}
if opt.SnapshotNumber <= 0 {
log.Fatalln("FIX FUCK NUMBER")
}
s.wal, err = os.OpenFile(
JoinChangesFileName(opt.Dir, opt.SnapshotNumber),
os.O_APPEND|os.O_WRONLY,
filePerm,
)
if err != nil {
return nil, err
}
return s, nil
}
@@ -240,7 +231,7 @@ func (s *Writer) packAndWrite() (err error) {
}
// 4. Пишу в atree сторінки
err = s.writePagesToAtree(prepared.WriteToIndex, prepared.WriteToData)
err = s.WritePagesToAtree(prepared.WriteToIndex, prepared.WriteToData)
if err != nil {
return
}
@@ -260,64 +251,49 @@ func (s *Writer) packAndWrite() (err error) {
// }
if forceSnapshot {
metricsStateCh := make(chan MetricsState)
if err := s.wal.Close(); err != nil {
return fmt.Errorf("close wal file: %s", err)
}
s.wal = nil
s.written = 0
snapshotNumberCh := make(chan int, 1)
s.appendToWorkerQueue(Changes{
Records: prepared.WriteResults,
MetricsCh: metricsStateCh,
Records: prepared.WriteResults,
SnapshotNumberCh: snapshotNumberCh,
FrozenIndexPagesCount: s.indexFreeList.Pages(),
IndexPageNumbers: s.indexFreeList.Cached(),
FrozenDataPagesCount: s.dataFreeList.Pages(),
DataPageNumbers: s.dataFreeList.Cached(),
})
if err := s.wal.Close(); err != nil {
return fmt.Errorf("close changes file: %s", err)
}
snapshotNumber := <-snapshotNumberCh
state := <-metricsStateCh
s.logNumber++
// write snapshot FIX
err = writeSnapshot(Snapshot{
LogNumber: s.logNumber,
Dir: s.dir,
WriteBufferSize: writeBufferSize,
Metrics: state.Metrics,
DataFreeList: s.dataFreeList,
IndexFreeList: s.indexFreeList,
})
// копіюю сторінки із delta файла в base файл і потім роблю Truncate
err = s.indexFreeList.Merge()
if err != nil {
return fmt.Errorf("write snapshot file: %s", err)
return
}
err = s.dataFreeList.Merge()
if err != nil {
return
}
var err error
s.wal, err = os.OpenFile(
JoinChangesFileName(s.dir, s.logNumber),
JoinChangesFileName(s.dir, snapshotNumber),
os.O_CREATE|os.O_WRONLY,
filePerm,
)
if err != nil {
return fmt.Errorf("create new changes file: %s", err)
}
s.written = 0
// release Worker
close(state.WaitCh)
} else {
s.appendToWorkerQueue(Changes{
Records: prepared.WriteResults,
})
}
// Якщо потрібен снапшот - відправляю канал із буфером розміру 1,
// в який воркер має покласти снапшот отриманий після застосування змін (records).
// Після чого Writer створює нові файли snapshot і changes і працює далі.
// flush
// err = s.flush(w)
// if err != nil {
// return err
// }
// FIX send to worker workerReqs
return nil
}
@@ -372,7 +348,7 @@ type PageToWrite struct {
}
// writePagesToAtree - записує сторінки в .data та .index файли
func (s *Writer) writePagesToAtree(indexPages []PageToWrite, dataPages []PageToWrite) (err error) {
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))
@@ -457,203 +433,33 @@ func (s *Writer) sendSignal() {
// helpers
/*
Формат:
metricsQty - varuint
[metric]*
де metric - це:
metricID - 4b
metricType - 1b
fracDigits - 1b
lastPageNo - 4b
since - 4b
sinceValue - 8b
until - 4b
untilValue - 8b
timestamps size - 2b
values size - 2b
timestams payload - Nb
values payload - Nb
data free list frozen pages - varsize
dataFreeList size - varsize
dataFreeList - Nb
index free list frozen pages - varsize
indexFreeList size - varsize
indexFreeList - Nb
CRC32 - 4b
*/
// type Metric struct {
// MetricID uint32
// MetricType qb.MetricType
// FracDigits byte
// LastPageNo uint32
// Since uint32
// SinceValue float64
// Until uint32
// UntilValue float64
// Timestamps []byte // payload FIX
// Values []byte // payload FIX add lastDelta + h ?
// }
const metricHeaderSize = 38
type Metric struct {
MetricID uint32
MetricType qb.MetricType
FracDigits byte
LastPageNo uint32
Since uint32
SinceValue float64
Until uint32
UntilValue float64
Timestamps [][]byte
TimestampsSize int
Values [][]byte
ValuesSize int
}
type Snapshot struct {
LogNumber int
Dir string
WriteBufferSize int
Metrics []Metric
IndexFreeList *freelist.FreeList
DataFreeList *freelist.FreeList
}
func writeSnapshot(snapshot Snapshot) (err error) {
var (
fileName = filepath.Join(snapshot.Dir, fmt.Sprintf("%d.snapshot", snapshot.LogNumber))
hasher = util.NewHasher()
prefix = make([]byte, metricHeaderSize)
)
file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0770)
if err != nil {
return
}
dst := io.MultiWriter(bufio.NewWriterSize(file, snapshot.WriteBufferSize), hasher)
_, err = bin.WriteVarSize(dst, len(snapshot.Metrics))
if err != nil {
return
}
for _, metric := range snapshot.Metrics {
tSize := metric.TimestampsSize
vSize := metric.ValuesSize
bin.PutUint32(prefix[0:], metric.MetricID)
prefix[4] = byte(metric.MetricType)
prefix[5] = metric.FracDigits
bin.PutUint32(prefix[6:], metric.LastPageNo)
bin.PutUint32(prefix[10:], metric.Since)
bin.PutFloat64(prefix[14:], metric.SinceValue)
bin.PutUint32(prefix[22:], metric.Until)
bin.PutFloat64(prefix[26:], metric.UntilValue)
bin.PutUint16(prefix[34:], uint16(tSize))
bin.PutUint16(prefix[36:], uint16(vSize))
_, err = dst.Write(prefix)
if err != nil {
return
}
// copy timestamps
writeChunks(dst, metric.Timestamps, tSize)
// copy values
writeChunks(dst, metric.Values, vSize)
}
// free data pages
_, err = bin.WriteVarSize(dst, snapshot.DataFreeList.Pages())
if err != nil {
return
}
err = freeListWriteTo(snapshot.DataFreeList, dst)
if err != nil {
return
}
// free index pages
_, err = bin.WriteVarSize(dst, snapshot.IndexFreeList.Pages())
if err != nil {
return
}
err = freeListWriteTo(snapshot.IndexFreeList, dst)
if err != nil {
return
}
bin.WriteUint32(file, hasher.Sum32())
err = file.Sync()
if err != nil {
return
}
err = file.Close()
if err != nil {
return
}
// копіюю сторінки із delta файла в base файл і потім роблю Truncate
err = snapshot.DataFreeList.Merge() // fix - get frozen pages count + unfilled
if err != nil {
return
}
err = snapshot.IndexFreeList.Merge()
if err != nil {
return
}
// prevLogNumber := logNumber - 1
// prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber))
// prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber))
// isExist, err := isFileExist(prevChanges)
// if err != nil {
// return
// }
// if isExist {
// err = os.Remove(prevChanges)
// if err != nil {
// qb.Abort(qb.DeletePrevChangesFileFailed, err)
// }
// }
// isExist, err = isFileExist(prevSnapshot)
// if err != nil {
// return
// }
// if isExist {
// err = os.Remove(prevSnapshot)
// if err != nil {
// qb.Abort(qb.DeletePrevSnapshotFileFailed, err)
// }
// }
return
}
// HELPERS
func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error {
serialized, err := freeList.Serialize()
if err != nil {
qb.Abort(qb.FailedFreeListSerialize, err)
}
_, err = bin.WriteVarSize(dst, len(serialized))
if err != nil {
return err
}
_, err = dst.Write(serialized)
if err != nil {
return err
}
return nil
}
func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) {
remaining := size
for _, buf := range chunks {
if remaining < len(buf) {
buf = buf[:remaining]
}
_, err = dst.Write(buf)
if err != nil {
return
}
remaining -= len(buf)
if remaining == 0 {
break
}
}
return
}
// func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) {
// remaining := size
// for _, buf := range chunks {
// if remaining < len(buf) {
// buf = buf[:remaining]
// }
// _, err = dst.Write(buf)
// if err != nil {
// return
// }
// remaining -= len(buf)
// if remaining == 0 {
// break
// }
// }
// return
// }