wp
This commit is contained in:
464
storage/wal_writer.go
Normal file
464
storage/wal_writer.go
Normal file
@@ -0,0 +1,464 @@
|
||||
package storage
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"math"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
|
||||
"gordenko.dev/dima/qb"
|
||||
"gordenko.dev/dima/qb/atree"
|
||||
"gordenko.dev/dima/qb/freelist"
|
||||
)
|
||||
|
||||
const (
|
||||
lsnSize = 4
|
||||
packetPrefixSize = 12 // 4 lsn + 4 packet length + 4 crc32
|
||||
|
||||
lengthIdx = 0
|
||||
checksumIdx = 4
|
||||
lsnIdx = 8
|
||||
|
||||
filePerm = 0770
|
||||
|
||||
dumpSnapshotAfterNBytes = 1024 * 1024 * 1024 // 1 GB
|
||||
|
||||
writeBufferSize = 4 * 1024 * 1024
|
||||
)
|
||||
|
||||
const (
|
||||
CodeMetricAdd byte = 1
|
||||
CodeMetricDelete byte = 2
|
||||
CodeMeasuresAppendWithGrow byte = 5
|
||||
CodeMeasuresAppend = 6
|
||||
CodeMeasuresDelete byte = 7
|
||||
)
|
||||
|
||||
type FreeList interface {
|
||||
// використовується в allocPage
|
||||
GetPageNumber() (uint32, error)
|
||||
}
|
||||
|
||||
func JoinChangesFileName(dir string, logNumber int) string {
|
||||
return filepath.Join(dir, fmt.Sprintf("%d.changes", logNumber))
|
||||
}
|
||||
|
||||
type MetricsState struct {
|
||||
//Metrics []Metric
|
||||
WaitCh chan struct{}
|
||||
}
|
||||
|
||||
type Changes struct {
|
||||
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
|
||||
dir string
|
||||
w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск
|
||||
wal *os.File
|
||||
dataFile *os.File
|
||||
indexFile *os.File
|
||||
indexPageSize int
|
||||
dataPageSize int
|
||||
input []any
|
||||
dataPreparer *DataPreparer
|
||||
appendToWorkerQueue func(any)
|
||||
//lsn uint32
|
||||
written int64
|
||||
isExited bool
|
||||
exitCh chan struct{}
|
||||
waitGroup *sync.WaitGroup
|
||||
signalCh chan struct{}
|
||||
}
|
||||
|
||||
type WriterOptions struct {
|
||||
IndexPageSize int
|
||||
DataPageSize int
|
||||
Dir string
|
||||
SnapshotNumber int // номер журнала
|
||||
AppendToWorkerQueue func(any)
|
||||
DataFreeList *freelist.FreeList
|
||||
IndexFreeList *freelist.FreeList
|
||||
Atree *atree.Atree
|
||||
ExitCh chan struct{}
|
||||
WaitGroup *sync.WaitGroup
|
||||
}
|
||||
|
||||
func NewWriter(opt WriterOptions) (*Writer, error) {
|
||||
if (opt.IndexPageSize % 2) != 0 {
|
||||
return nil, errors.New("IndexPageSize must be multiple of 2")
|
||||
}
|
||||
if (opt.DataPageSize % 2) != 0 {
|
||||
return nil, errors.New("DataPageSize must be multiple of 2")
|
||||
}
|
||||
if opt.Dir == "" {
|
||||
return nil, errors.New("Dir option is required")
|
||||
}
|
||||
if opt.AppendToWorkerQueue == nil {
|
||||
return nil, errors.New("AppendToWorkerQueue option is required")
|
||||
}
|
||||
if opt.DataFreeList == nil {
|
||||
return nil, errors.New("DataFreeList option is required")
|
||||
}
|
||||
if opt.IndexFreeList == nil {
|
||||
return nil, errors.New("IndexFreeList option is required")
|
||||
}
|
||||
if opt.Atree == nil {
|
||||
return nil, errors.New("Atree option is required")
|
||||
}
|
||||
if opt.ExitCh == nil {
|
||||
return nil, errors.New("ExitCh option is required")
|
||||
}
|
||||
if opt.WaitGroup == nil {
|
||||
return nil, errors.New("WaitGroup option is required")
|
||||
}
|
||||
|
||||
s := &Writer{
|
||||
indexPageSize: opt.IndexPageSize,
|
||||
dataPageSize: opt.DataPageSize,
|
||||
dir: opt.Dir,
|
||||
appendToWorkerQueue: opt.AppendToWorkerQueue,
|
||||
dataFreeList: opt.DataFreeList,
|
||||
indexFreeList: opt.IndexFreeList,
|
||||
atree: opt.Atree,
|
||||
exitCh: opt.ExitCh,
|
||||
waitGroup: opt.WaitGroup,
|
||||
signalCh: make(chan struct{}, 1),
|
||||
}
|
||||
|
||||
var err error
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
func (s *Writer) Run() {
|
||||
for {
|
||||
select {
|
||||
case <-s.signalCh:
|
||||
if err := s.packAndWrite(); err != nil {
|
||||
qb.Abort(qb.FailedWriteToTxLog, err)
|
||||
}
|
||||
|
||||
case <-s.exitCh:
|
||||
s.exit()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Writer) getDataPageNumber() (uint32, bool, error) {
|
||||
pageNo, err := s.dataFreeList.GetPageNumber()
|
||||
if err != nil {
|
||||
return 0, false, err
|
||||
}
|
||||
if pageNo > 0 {
|
||||
return pageNo, true, nil
|
||||
}
|
||||
if s.dataPagesCount < math.MaxUint32 {
|
||||
s.dataPagesCount++
|
||||
return s.dataPagesCount, false, nil
|
||||
}
|
||||
return 0, false, errors.New("no space")
|
||||
}
|
||||
|
||||
func (s *Writer) getIndexPageNumber() (uint32, bool, error) {
|
||||
pageNo, err := s.indexFreeList.GetPageNumber()
|
||||
if err != nil {
|
||||
return 0, false, err
|
||||
}
|
||||
if pageNo > 0 {
|
||||
return pageNo, true, nil
|
||||
}
|
||||
if s.indexPagesCount < math.MaxUint32 {
|
||||
s.indexPagesCount++
|
||||
return s.indexPagesCount, false, nil
|
||||
}
|
||||
return 0, false, errors.New("no space")
|
||||
}
|
||||
|
||||
func (s *Writer) packAndWrite() (err error) {
|
||||
s.mutex.Lock()
|
||||
//isExited := s.isExited
|
||||
input := s.input
|
||||
s.input = nil
|
||||
|
||||
// var exitWaitGroup *sync.WaitGroup
|
||||
// if s.isExited {
|
||||
// exitWaitGroup = s.waitGroup
|
||||
// }
|
||||
s.mutex.Unlock()
|
||||
|
||||
prepared := s.dataPreparer.Prepare(input)
|
||||
|
||||
// 3. Пишу на диск WAL (append)
|
||||
n, err := s.wal.Write(prepared.Packet)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if n != len(prepared.Packet) {
|
||||
return fmt.Errorf("written %d != total size %d", n, len(prepared.Packet))
|
||||
}
|
||||
if err = s.wal.Sync(); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// 4. Пишу в atree сторінки
|
||||
err = s.WritePagesToAtree(prepared.WriteToIndex, prepared.WriteToData)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
// 6. відправляю input - worker-у
|
||||
|
||||
var (
|
||||
forceSnapshot bool
|
||||
)
|
||||
|
||||
if s.written > dumpSnapshotAfterNBytes {
|
||||
forceSnapshot = true
|
||||
}
|
||||
|
||||
// if isExited && s.written > 0 {
|
||||
// forceSnapshot = true
|
||||
// }
|
||||
|
||||
if forceSnapshot {
|
||||
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,
|
||||
SnapshotNumberCh: snapshotNumberCh,
|
||||
FrozenIndexPagesCount: s.indexFreeList.Pages(),
|
||||
IndexPageNumbers: s.indexFreeList.Cached(),
|
||||
FrozenDataPagesCount: s.dataFreeList.Pages(),
|
||||
DataPageNumbers: s.dataFreeList.Cached(),
|
||||
})
|
||||
|
||||
snapshotNumber := <-snapshotNumberCh
|
||||
|
||||
// копіюю сторінки із delta файла в base файл і потім роблю Truncate
|
||||
err = s.indexFreeList.Merge()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
err = s.dataFreeList.Merge()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
var err error
|
||||
s.wal, err = os.OpenFile(
|
||||
JoinChangesFileName(s.dir, snapshotNumber),
|
||||
os.O_CREATE|os.O_WRONLY,
|
||||
filePerm,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create new changes file: %s", err)
|
||||
}
|
||||
} else {
|
||||
s.appendToWorkerQueue(Changes{
|
||||
Records: prepared.WriteResults,
|
||||
})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Writer) Append(req any) {
|
||||
s.mutex.Lock()
|
||||
s.input = append(s.input, req)
|
||||
s.mutex.Unlock()
|
||||
|
||||
select {
|
||||
case s.signalCh <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// func (s *Writer) reset() {
|
||||
// s.buf.Reset()
|
||||
// s.buf.Write([]byte{
|
||||
// 0, 0, 0, 0, // packet length
|
||||
// 0, 0, 0, 0, // crc32
|
||||
// 0, 0, 0, 0, // lsn
|
||||
|
||||
// })
|
||||
// s.pagesToWrite = nil
|
||||
// s.workerReqs = nil
|
||||
// s.waitCh = make(chan struct{})
|
||||
// }
|
||||
|
||||
// func (s *Writer) flush(w *bytes.Buffer) error {
|
||||
|
||||
// return nil
|
||||
// }
|
||||
|
||||
func (s *Writer) exit() {
|
||||
s.mutex.Lock()
|
||||
s.isExited = true
|
||||
s.mutex.Unlock()
|
||||
|
||||
if err := s.packAndWrite(); err != nil {
|
||||
qb.Abort(qb.FailedWriteToTxLog, err)
|
||||
}
|
||||
}
|
||||
|
||||
// type IndexPayloadToWrite struct {
|
||||
// PageNo uint32
|
||||
// Payload []byte
|
||||
// ZeroLevel bool
|
||||
// }
|
||||
|
||||
type PageToWrite struct {
|
||||
PageNo uint32
|
||||
Content []byte
|
||||
}
|
||||
|
||||
// writePagesToAtree - записує сторінки в .data та .index файли
|
||||
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))
|
||||
}
|
||||
var (
|
||||
off = int(p.PageNo-1) * s.dataPageSize
|
||||
n int
|
||||
)
|
||||
n, err = s.dataFile.WriteAt(p.Content, int64(off))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if n != s.dataPageSize {
|
||||
return fmt.Errorf("write %d instead of %d", n, s.dataPageSize)
|
||||
}
|
||||
}
|
||||
for _, p := range indexPages {
|
||||
if len(p.Content) != s.indexPageSize {
|
||||
return fmt.Errorf("wrong index page size: %d", len(p.Content))
|
||||
}
|
||||
var (
|
||||
off = int(p.PageNo-1) * s.indexPageSize
|
||||
n int
|
||||
)
|
||||
n, err = s.indexFile.WriteAt(p.Content, int64(off))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if n != s.indexPageSize {
|
||||
return fmt.Errorf("write %d instead of %d", n, s.indexPageSize)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// API
|
||||
|
||||
// FIX - add
|
||||
// type DeletedMeasuresSince struct {
|
||||
// MetricID uint32
|
||||
// LastPageNo uint32
|
||||
// IsRootChanged bool
|
||||
// RootPageNo uint32
|
||||
// FreeDataPages []uint32
|
||||
// FreeIndexPages []uint32
|
||||
// TimestampsBuf []byte
|
||||
// ValuesBuf []byte
|
||||
// }
|
||||
|
||||
func (s *Writer) sendSignal() {
|
||||
select {
|
||||
case s.signalCh <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// Якщо (Pages) > 0 - незаповнену сторінку кодую в стисненому вигляді.
|
||||
// Якщо (Pages) = 0 - кодую просто пари timestamp / value
|
||||
// func (s *Writer) packWriteAppended(w *bytes.Buffer, req AppendedPagesReq) {
|
||||
// // завантажений path. Отже просто додаємо в дерево data сторінки, створюємо нові індексні
|
||||
// // без звернення до диску.
|
||||
// report := s.atree.AppendDataPages(atree.AppendDataPagesReq{
|
||||
// LastPageNo: req.LastPageNo,
|
||||
// Legs: req.Legs,
|
||||
// DataPages: req.Pages,
|
||||
// })
|
||||
|
||||
// m := AppendedPages{
|
||||
// MetricID: req.MetricID,
|
||||
// Timestamp: req.Timestamp,
|
||||
// Value: req.Value,
|
||||
// NewRootPageNo: report.NewRootPageNo,
|
||||
// LastPageNo: report.LastPageNo,
|
||||
// Pages: report.Pages,
|
||||
// TimestampsChunks: req.TimestampsChunks,
|
||||
// TimestampsSize: int(req.TimestampsSize),
|
||||
// ValuesChunks: req.ValuesChunks,
|
||||
// ValuesSize: int(req.ValuesSize),
|
||||
// }
|
||||
// m.Pack(w)
|
||||
// }
|
||||
|
||||
// helpers
|
||||
|
||||
// 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 ?
|
||||
// }
|
||||
|
||||
// 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
|
||||
// }
|
||||
Reference in New Issue
Block a user