505 lines
12 KiB
Go
505 lines
12 KiB
Go
package txlog
|
||
|
||
import (
|
||
"bytes"
|
||
"errors"
|
||
"fmt"
|
||
"hash/crc32"
|
||
"log"
|
||
"os"
|
||
"path/filepath"
|
||
"sync"
|
||
|
||
bin "gordenko.dev/dima/bin/little"
|
||
"gordenko.dev/dima/qb"
|
||
"gordenko.dev/dima/qb/atree"
|
||
)
|
||
|
||
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
|
||
)
|
||
|
||
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 Changes struct {
|
||
Records []any
|
||
LogNumber int
|
||
ForceSnapshot bool
|
||
ExitWaitGroup *sync.WaitGroup
|
||
WaitCh chan struct{}
|
||
}
|
||
|
||
type Writer struct {
|
||
mutex sync.Mutex
|
||
freelist FreeList
|
||
atree *atree.Atree
|
||
logNumber int
|
||
dir string
|
||
w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск
|
||
wal *os.File
|
||
dataFile *os.File
|
||
indexFile *os.File
|
||
dataPage []byte
|
||
indexPage []byte
|
||
input []any
|
||
//buf *bytes.Buffer
|
||
//pagesToWrite []atree.PageToWrite
|
||
workerReqs []any
|
||
//waitCh chan struct{}
|
||
appendToWorkerQueue func(any)
|
||
//lsn uint32
|
||
written int64
|
||
isExited bool
|
||
exitCh chan struct{}
|
||
waitGroup *sync.WaitGroup
|
||
signalCh chan struct{}
|
||
}
|
||
|
||
type WriterOptions struct {
|
||
Dir string
|
||
LogNumber int // номер журнала
|
||
AppendToWorkerQueue func(any)
|
||
FreeList FreeList
|
||
Atree *atree.Atree
|
||
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")
|
||
}
|
||
if opt.FreeList == nil {
|
||
return nil, errors.New("FreeList 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{
|
||
dir: opt.Dir,
|
||
appendToWorkerQueue: opt.AppendToWorkerQueue,
|
||
freelist: opt.FreeList,
|
||
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) packAndWrite() (err error) {
|
||
s.mutex.Lock()
|
||
input := s.input
|
||
s.input = nil
|
||
s.mutex.Unlock()
|
||
|
||
w := bytes.NewBuffer(nil)
|
||
|
||
// 1. Пакую всі дані в WAL буфер запису
|
||
for _, untyped := range input {
|
||
switch x := untyped.(type) {
|
||
case AddedMetric:
|
||
x.Pack(w)
|
||
// case DeletedMetric:
|
||
// if len(x.FreePageNumbers) > 0 {
|
||
// s.freeList.AddPageNumbers(x.FreePageNumbers)
|
||
// }
|
||
case AppendedMeasures:
|
||
//
|
||
|
||
s.packMeasuresIntoWALBuffer(x)
|
||
//case DeletedMeasures:
|
||
//case DeletedMeasuresSince:
|
||
}
|
||
}
|
||
// 2. Додаю розмір пакету і чексуму
|
||
|
||
// 3. Пишу на диск WAL (append)
|
||
|
||
// 4. Пишу в atree сторінки
|
||
for range 1 {
|
||
err = s.writePagesToAtree(nil, nil)
|
||
if err != nil {
|
||
return
|
||
}
|
||
}
|
||
|
||
// 5. Пишу зміни в free list
|
||
|
||
// 6. відправляю input - worker-у
|
||
|
||
// 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.workerReqs = append(s.workerReqs, 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 {
|
||
s.mutex.Lock()
|
||
|
||
//pagesToWrite := s.pagesToWrite
|
||
workerReqs := s.workerReqs
|
||
isExited := s.isExited
|
||
|
||
var exitWaitGroup *sync.WaitGroup
|
||
if s.isExited {
|
||
exitWaitGroup = s.waitGroup
|
||
}
|
||
|
||
if w.Len() > packetPrefixSize {
|
||
//s.lsn++
|
||
//lsn := s.lsn
|
||
packet := make([]byte, w.Len())
|
||
copy(packet, w.Bytes())
|
||
//s.reset() fix
|
||
|
||
s.written += int64(len(packet)) + 12
|
||
s.mutex.Unlock()
|
||
|
||
bin.PutUint32(packet[lengthIdx:], uint32(len(packet)-packetPrefixSize))
|
||
//bin.PutUint32(packet[lsnIdx:], lsn)
|
||
bin.PutUint32(packet[checksumIdx:], crc32.ChecksumIEEE(packet[8:]))
|
||
|
||
n, err := s.wal.Write(packet)
|
||
if err != nil {
|
||
return fmt.Errorf("TxLog write: %s", err)
|
||
}
|
||
if n != len(packet) {
|
||
return fmt.Errorf("TxLog written %d != packet size %d", n, len(packet))
|
||
}
|
||
if err := s.wal.Sync(); err != nil {
|
||
return fmt.Errorf("TxLog sync: %s", err)
|
||
}
|
||
|
||
// err = s.writePagesToAtree(pagesToWrite)
|
||
// if err != nil {
|
||
// return fmt.Errorf("TxLog writePagesToAtree: %s", err)
|
||
// }
|
||
} else {
|
||
s.mutex.Unlock()
|
||
}
|
||
|
||
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 changes file: %s", err)
|
||
}
|
||
s.logNumber++
|
||
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
|
||
}
|
||
|
||
s.appendToWorkerQueue(Changes{
|
||
Records: workerReqs,
|
||
ForceSnapshot: forceSnapshot,
|
||
LogNumber: s.logNumber,
|
||
ExitWaitGroup: exitWaitGroup,
|
||
})
|
||
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)
|
||
}
|
||
}
|
||
|
||
// writePagesToAtree - записує сторінки в .data та .index файли
|
||
func (s *Writer) writePagesToAtree(levels []*atree.IndexLevel, dataPages []*atree.DataPage) (err error) {
|
||
for _, p := range dataPages {
|
||
p.Checksum = atree.ChunksToDataPage(s.dataPage, atree.ChunksToDataPageReq{
|
||
PrevPageNo: p.PrevPageNo,
|
||
Timestamps: p.Timestamps,
|
||
TimestampsSize: p.TimestampsSize,
|
||
Values: p.Values,
|
||
ValuesSize: p.ValuesSize,
|
||
})
|
||
off := (p.PageNo - 1) * atree.DataPageSize
|
||
n, err := s.dataFile.WriteAt(s.dataPage, int64(off))
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if n != atree.DataPageSize {
|
||
return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize)
|
||
}
|
||
}
|
||
for levelIdx, level := range levels {
|
||
for _, p := range level.Filled {
|
||
// if len(p.Data) != atree.PageSize {
|
||
// return fmt.Errorf("wrong page %d size: %d",
|
||
// p.PageNo, len(p.Data))
|
||
// }
|
||
p.Checksum = atree.DataToIndexPage(s.indexPage, p.Data, levelIdx == 0)
|
||
off := (p.PageNo - 1) * atree.IndexPageSize
|
||
n, err := s.indexFile.WriteAt(s.indexPage, int64(off))
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if n != atree.IndexPageSize {
|
||
return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize)
|
||
}
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
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
|
||
}
|
||
|
||
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)
|
||
bin.WriteUint32(s.w, p.Checksum)
|
||
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
|
||
}
|
||
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)
|
||
bin.WriteUint32(s.w, p.Checksum)
|
||
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
|
||
}
|
||
|
||
// 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
|