Files
qb/txlog/writer.go
2026-05-10 00:59:47 +00:00

379 lines
8.1 KiB
Go

package txlog
import (
"bytes"
"errors"
"fmt"
"hash/crc32"
"os"
"path/filepath"
"sync"
bin "gordenko.dev/dima/bin/little"
octopus "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
ReservePage() uint32
}
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
dataFile *os.File
file *os.File
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),
waitCh: make(chan struct{}),
}
var err error
if opt.LogNumber > 0 {
s.file, 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.file, 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 {
octopus.Abort(octopus.FailedWriteToTxLog, err)
}
case <-s.exitCh:
s.exit()
return
}
}
}
func (s *Writer) packAndWrite() error {
s.mutex.Lock()
input := s.input
waitCh := s.waitCh
s.input = nil
s.waitCh = make(chan struct{})
s.mutex.Unlock()
w := bytes.NewBuffer(nil)
for _, untyped := range input {
switch x := untyped.(type) {
case AppendedPagesReq:
s.packWriteAppended(w, x)
case AddedMetric:
case DeletedMetric:
if len(x.FreePageNumbers) > 0 {
//s.freeList.AddPages(x.FreePageNumbers)
}
case AppendedMeasure:
case AppendedMeasures:
case DeletedMeasures:
//case DeletedMeasuresSince:
}
}
// flush
close(waitCh)
return s.flush(w)
}
func (s *Writer) Append(req any) chan struct{} {
s.mutex.Lock()
s.input = append(s.input, req)
s.workerReqs = append(s.workerReqs, req)
waitCh := s.waitCh
s.mutex.Unlock()
select {
case s.signalCh <- struct{}{}:
default:
}
return waitCh
}
// 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
waitCh := s.waitCh
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.file.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.file.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.waitCh = make(chan struct{})
s.mutex.Unlock()
}
var forceSnapshot bool
if s.written > dumpSnapshotAfterNBytes {
forceSnapshot = true
}
if isExited && s.written > 0 {
forceSnapshot = true
}
if forceSnapshot {
if err := s.file.Close(); err != nil {
return fmt.Errorf("close changes file: %s", err)
}
s.logNumber++
var err error
s.file, 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,
WaitCh: waitCh,
ExitWaitGroup: exitWaitGroup,
})
return nil
}
func (s *Writer) writePagesToAtree(pages []atree.PageToWrite) error {
for _, p := range pages {
if len(p.Data) != atree.PageSize {
return fmt.Errorf("wrong page %d size: %d",
p.PageNo, len(p.Data))
}
off := (p.PageNo - 1) * atree.PageSize
n, err := s.dataFile.WriteAt(p.Data, int64(off))
if err != nil {
return err
}
if n != len(p.Data) {
return fmt.Errorf("write %d instead of %d", n, len(p.Data))
}
}
return nil
}
func (s *Writer) exit() {
s.mutex.Lock()
s.isExited = true
s.mutex.Unlock()
if err := s.packAndWrite(); err != nil {
octopus.Abort(octopus.FailedWriteToTxLog, err)
}
}
// 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