Files
qb/txlog/writer.go

764 lines
18 KiB
Go
Raw Normal View History

2026-02-10 14:02:11 +00:00
package txlog
import (
2026-05-31 20:01:28 +00:00
"bufio"
2026-02-10 14:02:11 +00:00
"bytes"
"errors"
"fmt"
2026-05-31 20:01:28 +00:00
"io"
2026-05-21 23:38:52 +03:00
"log"
2026-05-31 20:01:28 +00:00
"math"
2026-02-10 14:02:11 +00:00
"os"
"path/filepath"
"sync"
2026-05-10 00:59:47 +00:00
bin "gordenko.dev/dima/bin/little"
2026-05-21 23:38:52 +03:00
"gordenko.dev/dima/qb"
2026-05-10 00:59:47 +00:00
"gordenko.dev/dima/qb/atree"
2026-05-31 20:01:28 +00:00
"gordenko.dev/dima/qb/freelist"
"gordenko.dev/dima/qb/util"
2026-02-10 14:02:11 +00:00
)
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
2026-05-31 20:01:28 +00:00
writeBufferSize = 4 * 1024 * 1024
2026-02-10 14:02:11 +00:00
)
const (
2026-05-10 00:59:47 +00:00
CodeAddedMetric byte = 1
CodeDeletedMetric byte = 2
CodeAppendedMeasure byte = 4
CodeAppendedMeasures byte = 5
CodeAppendedPages byte = 6
CodeDeletedMeasures byte = 7
2026-02-10 14:02:11 +00:00
)
2026-05-10 00:59:47 +00:00
type FreeList interface {
// використовується в allocPage
2026-05-21 23:38:52 +03:00
GetPageNumber() (uint32, error)
2026-05-10 00:59:47 +00:00
}
2026-02-10 14:02:11 +00:00
func JoinChangesFileName(dir string, logNumber int) string {
return filepath.Join(dir, fmt.Sprintf("%d.changes", logNumber))
}
2026-05-31 20:01:28 +00:00
type MetricsState struct {
Metrics []Metric
WaitCh chan struct{}
}
2026-02-10 14:02:11 +00:00
type Changes struct {
2026-05-31 20:01:28 +00:00
Records []any
MetricsCh chan MetricsState
2026-02-10 14:02:11 +00:00
}
type Writer struct {
2026-05-31 20:01:28 +00:00
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
// буфер для складання payload на сторінку, розрахунку CRC32 та інше
pageBuffer []byte
input []any
2026-02-10 14:02:11 +00:00
appendToWorkerQueue func(any)
2026-05-21 23:38:52 +03:00
//lsn uint32
written int64
isExited bool
exitCh chan struct{}
waitGroup *sync.WaitGroup
signalCh chan struct{}
2026-02-10 14:02:11 +00:00
}
type WriterOptions struct {
Dir string
LogNumber int // номер журнала
AppendToWorkerQueue func(any)
2026-05-31 20:01:28 +00:00
DataFreeList *freelist.FreeList
IndexFreeList *freelist.FreeList
2026-05-10 00:59:47 +00:00
Atree *atree.Atree
2026-02-10 14:02:11 +00:00
ExitCh chan struct{}
WaitGroup *sync.WaitGroup
}
func NewWriter(opt WriterOptions) (*Writer, error) {
if opt.Dir == "" {
return nil, errors.New("Dir option is required")
}
if opt.AppendToWorkerQueue == nil {
return nil, errors.New("AppendToWorkerQueue option is required")
}
2026-05-31 20:01:28 +00:00
if opt.DataFreeList == nil {
return nil, errors.New("DataFreeList option is required")
}
if opt.IndexFreeList == nil {
return nil, errors.New("IndexFreeList option is required")
2026-05-10 00:59:47 +00:00
}
if opt.Atree == nil {
return nil, errors.New("Atree option is required")
}
2026-02-10 14:02:11 +00:00
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{
dir: opt.Dir,
appendToWorkerQueue: opt.AppendToWorkerQueue,
2026-05-31 20:01:28 +00:00
dataFreeList: opt.DataFreeList,
indexFreeList: opt.IndexFreeList,
2026-05-10 00:59:47 +00:00
atree: opt.Atree,
2026-02-10 14:02:11 +00:00
logNumber: opt.LogNumber,
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
signalCh: make(chan struct{}, 1),
}
var err error
if opt.LogNumber > 0 {
2026-05-21 23:38:52 +03:00
s.wal, err = os.OpenFile(
2026-02-10 14:02:11 +00:00
JoinChangesFileName(opt.Dir, s.logNumber),
os.O_APPEND|os.O_WRONLY,
filePerm,
)
if err != nil {
return nil, err
}
} else {
s.logNumber = 1
2026-05-21 23:38:52 +03:00
s.wal, err = os.OpenFile(
2026-02-10 14:02:11 +00:00
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:
2026-05-10 00:59:47 +00:00
if err := s.packAndWrite(); err != nil {
2026-05-21 23:38:52 +03:00
qb.Abort(qb.FailedWriteToTxLog, err)
2026-02-10 14:02:11 +00:00
}
case <-s.exitCh:
s.exit()
return
}
}
}
2026-05-31 20:01:28 +00:00
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")
}
2026-05-21 23:38:52 +03:00
func (s *Writer) packAndWrite() (err error) {
2026-05-10 00:59:47 +00:00
s.mutex.Lock()
2026-05-31 20:01:28 +00:00
//isExited := s.isExited
2026-05-10 00:59:47 +00:00
input := s.input
s.input = nil
2026-05-31 20:01:28 +00:00
// var exitWaitGroup *sync.WaitGroup
// if s.isExited {
// exitWaitGroup = s.waitGroup
// }
2026-05-10 00:59:47 +00:00
s.mutex.Unlock()
2026-05-31 20:01:28 +00:00
prepared := prepareData(input)
2026-05-21 23:38:52 +03:00
// 3. Пишу на диск WAL (append)
2026-05-31 20:01:28 +00:00
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
}
2026-05-21 23:38:52 +03:00
// 4. Пишу в atree сторінки
2026-05-31 20:01:28 +00:00
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.ToWorker,
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,
})
2026-05-21 23:38:52 +03:00
if err != nil {
2026-05-31 20:01:28 +00:00
return fmt.Errorf("write snapshot file: %s", err)
2026-05-21 23:38:52 +03:00
}
2026-05-31 20:01:28 +00:00
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)
}
2026-05-21 23:38:52 +03:00
2026-05-31 20:01:28 +00:00
s.written = 0
// release Worker
close(state.WaitCh)
} else {
s.appendToWorkerQueue(Changes{
Records: prepared.ToWorker,
})
}
// Якщо потрібен снапшот - відправляю канал із буфером розміру 1,
// в який воркер має покласти снапшот отриманий після застосування змін (records).
// Після чого Writer створює нові файли snapshot і changes і працює далі.
2026-05-21 23:38:52 +03:00
2026-05-10 00:59:47 +00:00
// flush
2026-05-31 20:01:28 +00:00
// err = s.flush(w)
// if err != nil {
// return err
// }
2026-05-21 23:38:52 +03:00
// FIX send to worker workerReqs
return nil
2026-02-10 14:02:11 +00:00
}
2026-05-21 23:38:52 +03:00
func (s *Writer) Append(req any) {
2026-02-10 14:02:11 +00:00
s.mutex.Lock()
2026-05-10 00:59:47 +00:00
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{})
// }
2026-02-10 14:02:11 +00:00
2026-05-31 20:01:28 +00:00
// func (s *Writer) flush(w *bytes.Buffer) error {
2026-02-10 14:02:11 +00:00
2026-05-31 20:01:28 +00:00
// return nil
// }
2026-02-10 14:02:11 +00:00
2026-05-21 23:38:52 +03:00
func (s *Writer) exit() {
s.mutex.Lock()
s.isExited = true
s.mutex.Unlock()
if err := s.packAndWrite(); err != nil {
qb.Abort(qb.FailedWriteToTxLog, err)
}
}
2026-02-10 14:02:11 +00:00
2026-05-31 20:01:28 +00:00
type IndexPayloadToWrite struct {
PageNo uint32
Payload []byte
ZeroLevel bool
}
type DataPayloadToWrite struct {
PrevPageNo uint32
Timestamps [][]byte // chunks
TimestampsSize int
Values [][]byte // chunks
ValuesSize int
PageNo uint32
}
2026-05-21 23:38:52 +03:00
// writePagesToAtree - записує сторінки в .data та .index файли
2026-05-31 20:01:28 +00:00
func (s *Writer) writePagesToAtree(indexItems []IndexPayloadToWrite, dataItems []DataPayloadToWrite) (err error) {
for _, p := range dataItems {
atree.ChunksToDataPage(s.pageBuffer, atree.ChunksToDataPageReq{
2026-05-21 23:38:52 +03:00
PrevPageNo: p.PrevPageNo,
Timestamps: p.Timestamps,
TimestampsSize: p.TimestampsSize,
Values: p.Values,
ValuesSize: p.ValuesSize,
})
2026-05-31 20:01:28 +00:00
var (
off = (p.PageNo - 1) * atree.DataPageSize
n int
)
n, err = s.dataFile.WriteAt(s.pageBuffer, int64(off))
2026-05-10 00:59:47 +00:00
if err != nil {
2026-05-31 20:01:28 +00:00
return
2026-05-10 00:59:47 +00:00
}
2026-05-21 23:38:52 +03:00
if n != atree.DataPageSize {
return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize)
}
}
2026-05-31 20:01:28 +00:00
for _, p := range indexItems {
// if len(p.Payload) > atree.MaxIndexPayload {
// return fmt.Errorf("wrong index payload size: %d", len(p.Payload))
// }
indexPageBuffer := s.pageBuffer[:atree.IndexPageSize]
atree.DataToIndexPage(indexPageBuffer, p.Payload, p.ZeroLevel)
var (
off = (p.PageNo - 1) * atree.IndexPageSize
n int
)
n, err = s.indexFile.WriteAt(indexPageBuffer, int64(off))
if err != nil {
return
}
if n != atree.IndexPageSize {
return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize)
2026-05-10 00:59:47 +00:00
}
2026-02-10 14:02:11 +00:00
}
2026-05-10 00:59:47 +00:00
return nil
2026-02-10 14:02:11 +00:00
}
2026-05-31 20:01:28 +00:00
type AppendMeasuresSummary struct {
MetricID uint32
//TimestampsOffset int // (заповнені одразу)
//ValuesOffset int // (заповнені одразу)
// Timestamps [][]byte
// TimestampsSize int
// Values [][]byte
// ValuesSize int
LastPageNo uint32
Index [][]byte
ResultCode int
WrittenCount int
ResultCh chan struct{}
}
2026-05-21 23:38:52 +03:00
type AppendedMeasures struct {
MetricID uint32
TimestampsOffset int // (заповнені одразу)
ValuesOffset int // (заповнені одразу)
Timestamps [][]byte
TimestampsSize int
Values [][]byte
ValuesSize int
IndexLevels []*atree.IndexLevel
DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page
2026-05-31 20:01:28 +00:00
ResultCode int
WrittenCount int
ResultCh chan struct{}
2026-05-21 23:38:52 +03:00
}
2026-02-10 14:02:11 +00:00
2026-05-21 23:38:52 +03:00
func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) {
// fix write code
bin.WriteUint32(s.w, req.MetricID)
bin.WriteVarSize(s.w, len(req.DataPages))
for idx, p := range req.DataPages {
bin.WriteUint32(s.w, p.PageNo)
var reused byte
if p.IsReused {
reused = 1
}
s.w.WriteByte(reused)
bin.WriteUint32(s.w, p.PrevPageNo)
2026-05-31 20:01:28 +00:00
//bin.WriteUint32(s.w, p.Checksum)
2026-05-21 23:38:52 +03:00
var (
timestampsOffset int
valuesOffset int
)
if idx == 0 {
timestampsOffset = req.TimestampsOffset
valuesOffset = req.ValuesOffset
}
copyPayloadFromChunks(s.w, p.Timestamps, p.TimestampsSize, timestampsOffset)
copyPayloadFromChunks(s.w, p.Values, p.ValuesSize, valuesOffset)
}
bin.WriteUint16(s.w, uint16(req.TimestampsOffset))
bin.WriteUint16(s.w, uint16(req.ValuesOffset))
var (
timestampsOffset int
valuesOffset int
)
if len(req.DataPages) == 0 {
timestampsOffset = req.TimestampsOffset
valuesOffset = req.ValuesOffset
2026-02-10 14:02:11 +00:00
}
2026-05-21 23:38:52 +03:00
copyPayloadFromChunks(s.w, req.Timestamps, req.TimestampsSize, timestampsOffset)
copyPayloadFromChunks(s.w, req.Values, req.ValuesSize, valuesOffset)
if len(req.DataPages) == 0 {
return
}
// fix - levels reduced
bin.WriteVarSize(s.w, len(req.IndexLevels))
for _, level := range req.IndexLevels {
bin.WriteVarSize(s.w, len(level.Filled))
for idx, p := range level.Filled {
bin.WriteUint32(s.w, p.PageNo)
var reused byte
if p.IsReused {
reused = 1
}
s.w.WriteByte(reused)
2026-05-31 20:01:28 +00:00
//bin.WriteUint32(s.w, p.Checksum)
2026-05-21 23:38:52 +03:00
records := p.Data
if idx == 0 {
records = p.Data[level.Offset:]
}
if (len(records) % recordSize) != 0 {
log.Fatalf("wrong records length: %d", len(records))
}
bin.WriteVarSize(s.w, len(records)/recordSize)
s.w.Write(records)
}
records := level.Data
if len(level.Filled) == 0 {
records = level.Data[level.Offset:]
}
if (len(records) % recordSize) != 0 {
log.Fatalf("wrong unfilled records length: %d", len(records))
}
bin.WriteVarSize(s.w, len(records)/recordSize)
s.w.Write(records)
}
return
2026-02-10 14:02:11 +00:00
}
2026-05-10 00:59:47 +00:00
// API
2026-02-10 14:02:11 +00:00
2026-05-10 00:59:47 +00:00
// FIX - add
// type DeletedMeasuresSince struct {
// MetricID uint32
// LastPageNo uint32
// IsRootChanged bool
// RootPageNo uint32
// FreeDataPages []uint32
// FreeIndexPages []uint32
// TimestampsBuf []byte
// ValuesBuf []byte
// }
2026-02-10 14:02:11 +00:00
func (s *Writer) sendSignal() {
select {
case s.signalCh <- struct{}{}:
default:
}
}
2026-05-10 00:59:47 +00:00
// Якщо (Pages) > 0 - незаповнену сторінку кодую в стисненому вигляді.
// Якщо (Pages) = 0 - кодую просто пари timestamp / value
2026-05-21 23:38:52 +03:00
// 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,
// })
2026-02-10 14:02:11 +00:00
2026-05-21 23:38:52 +03:00
// 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)
// }
2026-05-10 00:59:47 +00:00
// helpers
2026-05-31 20:01:28 +00:00
/*
Формат:
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
}