Files
qb/txlog/writer.go
2026-06-05 19:43:01 +00:00

660 lines
15 KiB
Go
Raw 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 txlog
import (
"bufio"
"bytes"
"errors"
"fmt"
"io"
"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 (
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 (
CodeAddedMetric byte = 1
CodeDeletedMetric byte = 2
CodeAppendedMeasure byte = 4
CodeAppendedMeasures byte = 5
CodeAppendedPages byte = 6
CodeDeletedMeasures 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
MetricsCh chan MetricsState
}
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
LogNumber 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,
logNumber: opt.LogNumber,
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
signalCh: make(chan struct{}, 1),
}
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
}
}
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 {
metricsStateCh := make(chan MetricsState)
s.appendToWorkerQueue(Changes{
Records: prepared.WriteResults,
MetricsCh: metricsStateCh,
})
if err := s.wal.Close(); err != nil {
return fmt.Errorf("close changes file: %s", err)
}
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,
})
if err != nil {
return fmt.Errorf("write snapshot file: %s", err)
}
var err error
s.wal, err = os.OpenFile(
JoinChangesFileName(s.dir, s.logNumber),
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
}
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
/*
Формат:
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
*/
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
}