363 lines
7.4 KiB
Go
363 lines
7.4 KiB
Go
package storage
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
|
|
bin "gordenko.dev/dima/bin/little"
|
|
"gordenko.dev/dima/qb/util"
|
|
)
|
|
|
|
const readBufferSize = 8 * 1024 * 1024
|
|
|
|
type WALReader struct {
|
|
file *os.File
|
|
reader *bufio.Reader
|
|
current []any
|
|
next []any
|
|
}
|
|
|
|
type WALReaderOptions struct {
|
|
FileName string
|
|
BufferSize int
|
|
}
|
|
|
|
func NewWALReader(opt WALReaderOptions) (*WALReader, error) {
|
|
if opt.FileName == "" {
|
|
return nil, errors.New("FileName option is required")
|
|
}
|
|
if opt.BufferSize <= 0 {
|
|
return nil, errors.New("BufferSize option is required")
|
|
}
|
|
|
|
file, err := os.Open(opt.FileName)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &WALReader{
|
|
file: file,
|
|
reader: bufio.NewReaderSize(file, readBufferSize),
|
|
}, nil
|
|
}
|
|
|
|
func (s *WALReader) Close() {
|
|
s.file.Close()
|
|
}
|
|
|
|
func (s *WALReader) 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
|
|
}
|
|
}
|
|
}
|
|
|
|
current := s.current
|
|
done := s.next == nil
|
|
|
|
s.current = s.next
|
|
s.next, err = s.readPacket()
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
return current, done, nil
|
|
}
|
|
|
|
func (s *WALReader) readPacket() (_ []any, err error) {
|
|
payloadSize, err := bin.ReadVarSize(s.reader)
|
|
if err != nil {
|
|
if err == io.EOF {
|
|
return nil, nil
|
|
} else {
|
|
return
|
|
}
|
|
}
|
|
|
|
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 s.parseRecords(body)
|
|
}
|
|
//s.logger.Printf("stored CRC %d != calculated CRC %d", storedCRC, calculatedCRC)
|
|
return nil, nil
|
|
}
|
|
|
|
func (s *WALReader) parseRecords(body []byte) ([]any, error) {
|
|
var (
|
|
src = bytes.NewBuffer(body)
|
|
records []any
|
|
)
|
|
|
|
for {
|
|
recordType, err := src.ReadByte()
|
|
if err != nil {
|
|
if err == io.EOF {
|
|
return records, nil
|
|
}
|
|
return nil, err
|
|
}
|
|
|
|
switch recordType {
|
|
case CodeMetricAdd:
|
|
rec := new(MetricAddRecord)
|
|
if err = rec.Parse(src); err != nil {
|
|
return nil, err
|
|
}
|
|
records = append(records, rec)
|
|
|
|
case CodeMetricDelete:
|
|
rec := new(MetricDeleteRecord)
|
|
if err = rec.Parse(src); err != nil {
|
|
return nil, err
|
|
}
|
|
records = append(records, rec)
|
|
|
|
// case CodeAppendedMeasure:
|
|
// rec := new(AppendedMeasure)
|
|
// if err = rec.Parse(src); err != nil {
|
|
// return nil, err
|
|
// }
|
|
// records = append(records, rec)
|
|
|
|
// case CodeAppendedMeasures:
|
|
// rec := new(AppendedMeasures)
|
|
// if err = rec.Parse(src); err != nil {
|
|
// return nil, err
|
|
// }
|
|
// records = append(records, rec)
|
|
|
|
// case CodeAppendedPages:
|
|
// rec := new(AppendedPages)
|
|
// if err = rec.Parse(src); err != nil {
|
|
// return nil, err
|
|
// }
|
|
// records = append(records, rec)
|
|
|
|
case CodeMeasuresDelete:
|
|
rec := new(DeletedMeasures)
|
|
if err = rec.Parse(src); err != nil {
|
|
return nil, err
|
|
}
|
|
records = append(records, rec)
|
|
|
|
default:
|
|
return nil, fmt.Errorf("unknown record type code: %d", recordType)
|
|
}
|
|
}
|
|
}
|
|
|
|
// func (s *Reader) readAddedMetric(src *bytes.Buffer) (_ AddedMetric, err error) {
|
|
// arr, err := bin.ReadN(src, 6)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// return AddedMetric{
|
|
// MetricID: bin.GetUint32(arr),
|
|
// MetricType: diploma.MetricType(arr[4]),
|
|
// FracDigits: int(arr[5]),
|
|
// }, nil
|
|
// }
|
|
|
|
// func (s *Reader) readDeletedMetric(src *bytes.Buffer) (_ DeletedMetric, err error) {
|
|
// var rec DeletedMetric
|
|
// rec.MetricID, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// // free data pages
|
|
// dataQty, _, err := bin.ReadVarUint64(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// for range dataQty {
|
|
// var pageNo uint32
|
|
// pageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.FreeDataPages = append(rec.FreeDataPages, pageNo)
|
|
// }
|
|
// // free index pages
|
|
// indexQty, _, err := bin.ReadVarUint64(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// for range indexQty {
|
|
// var pageNo uint32
|
|
// pageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.FreeIndexPages = append(rec.FreeIndexPages, pageNo)
|
|
// }
|
|
// return rec, nil
|
|
// }
|
|
|
|
// func (s *Reader) readAppendedMeasure(src *bytes.Buffer) (_ AppendedMeasure, err error) {
|
|
// arr, err := bin.ReadN(src, 16)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// return AppendedMeasure{
|
|
// MetricID: bin.GetUint32(arr[0:]),
|
|
// Timestamp: bin.GetUint32(arr[4:]),
|
|
// Value: bin.GetFloat64(arr[8:]),
|
|
// }, nil
|
|
// }
|
|
|
|
// func (s *Reader) readAppendedMeasures(src *bytes.Buffer) (_ AppendedMeasures, err error) {
|
|
// var rec AppendedMeasures
|
|
// rec.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
|
|
// }
|
|
// rec.Measures = append(rec.Measures, measure)
|
|
// }
|
|
// return rec, nil
|
|
// }
|
|
|
|
// func (s *Reader) readAppendedMeasureWithOverflow(src *bytes.Buffer) (_ AppendedMeasureWithOverflow, err error) {
|
|
// var (
|
|
// b byte
|
|
// rec AppendedMeasureWithOverflow
|
|
// )
|
|
// rec.MetricID, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.Timestamp, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.Value, err = bin.ReadFloat64(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// b, err = src.ReadByte()
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.IsDataPageReused = b == 1
|
|
|
|
// rec.DataPageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// b, err = src.ReadByte()
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// if b == 1 {
|
|
// rec.IsRootChanged = true
|
|
// rec.RootPageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// }
|
|
// // index pages
|
|
// indexQty, err := src.ReadByte()
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// for range indexQty {
|
|
// var pageNo uint32
|
|
// pageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// rec.ReusedIndexPages = append(rec.ReusedIndexPages, pageNo)
|
|
// }
|
|
// return rec, nil
|
|
// }
|
|
|
|
// func (s *Reader) readDeletedMeasures(src *bytes.Buffer) (_ DeletedMeasures, err error) {
|
|
// var (
|
|
// rec DeletedMeasures
|
|
// )
|
|
// rec.MetricID, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// // free data pages
|
|
// rec.FreeDataPages, err = s.readFreePageNumbers(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// // free index pages
|
|
// rec.FreeIndexPages, err = s.readFreePageNumbers(src)
|
|
// if err != nil {
|
|
// return
|
|
// }
|
|
// return rec, nil
|
|
// }
|
|
|
|
// HELPERS
|
|
|
|
// func (s *Reader) readFreePageNumbers(src *bytes.Buffer) ([]uint32, error) {
|
|
// var freePages []uint32
|
|
// qty, _, err := bin.ReadVarUint64(src)
|
|
// if err != nil {
|
|
// return nil, err
|
|
// }
|
|
// for range qty {
|
|
// var pageNo uint32
|
|
// pageNo, err = bin.ReadUint32(src)
|
|
// if err != nil {
|
|
// return nil, err
|
|
// }
|
|
// freePages = append(freePages, pageNo)
|
|
// }
|
|
// return freePages, nil
|
|
// }
|
|
|
|
func (s *WALReader) Seek(offset int64) error {
|
|
ret, err := s.file.Seek(offset, 0)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if ret != offset {
|
|
return fmt.Errorf("ret %d != offset %d", ret, offset)
|
|
}
|
|
return nil
|
|
}
|