Files
qb/storage/writer.go

475 lines
11 KiB
Go
Raw Normal View History

2026-06-10 06:18:45 +03:00
package storage
2026-02-10 14:02:11 +00:00
import (
"bytes"
"errors"
"fmt"
"os"
"path/filepath"
"sync"
2026-05-21 23:38:52 +03:00
"gordenko.dev/dima/qb"
2026-05-31 20:01:28 +00:00
"gordenko.dev/dima/qb/freelist"
2026-06-13 22:43:17 +00:00
"gordenko.dev/dima/qb/inbox"
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-06-10 06:18:45 +03:00
CodeMetricAdd byte = 1
CodeMetricDelete byte = 2
CodeMeasuresAppendWithGrow byte = 5
2026-06-13 08:01:42 +03:00
CodeMeasuresAppend byte = 6
2026-06-10 06:18:45 +03:00
CodeMeasuresDelete 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-05-31 20:01:28 +00:00
type MetricsState struct {
2026-06-09 08:16:16 +03:00
//Metrics []Metric
WaitCh chan struct{}
2026-05-31 20:01:28 +00:00
}
2026-02-10 14:02:11 +00:00
type Changes struct {
2026-06-13 07:13:03 +00:00
Commits []any
// for snapshot only
2026-06-09 08:16:16 +03:00
SnapshotNumberCh chan int // log number
FrozenIndexPagesCount int
IndexPageNumbers []uint32
FrozenDataPagesCount int
DataPageNumbers []uint32
2026-02-10 14:02:11 +00:00
}
type Writer struct {
2026-06-13 19:53:27 +03:00
mutex sync.Mutex
2026-06-13 22:43:17 +00:00
inbox *inbox.Inbox
workerInbox *inbox.Inbox
2026-06-13 19:53:27 +03:00
dataPagesCount uint32
indexPagesCount uint32
dataFreeList *freelist.FreeList
indexFreeList *freelist.FreeList
//atree *atree.Atree
2026-06-13 22:43:17 +00:00
dir string
databaseName string
w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск
wal *os.File
dataFile *os.File
indexFile *os.File
input []any
writePreparer *WritePreparer
written int64
isExited bool
exitCh chan struct{}
waitGroup *sync.WaitGroup
//signalCh chan struct{}
2026-02-10 14:02:11 +00:00
}
type WriterOptions struct {
2026-06-13 22:43:17 +00:00
Inbox *inbox.Inbox
WorkerInbox *inbox.Inbox
Dir string
DatabaseName string
2026-06-13 19:53:27 +03:00
//SnapshotNumber int // номер журнала
2026-06-13 22:43:17 +00:00
WAL string
DataFreeList *freelist.FreeList
IndexFreeList *freelist.FreeList
2026-06-15 01:20:30 +03:00
DataFile *os.File
IndexFile *os.File
2026-06-13 19:53:27 +03:00
//Atree *atree.Atree
ExitCh chan struct{}
WaitGroup *sync.WaitGroup
2026-02-10 14:02:11 +00:00
}
func NewWriter(opt WriterOptions) (*Writer, error) {
2026-06-13 19:53:27 +03:00
// 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")
// }
2026-06-13 22:43:17 +00:00
if opt.Inbox == nil {
return nil, errors.New("Inbox option is required")
}
if opt.WorkerInbox == nil {
return nil, errors.New("WorkerInbox option is required")
}
2026-02-10 14:02:11 +00:00
if opt.Dir == "" {
return nil, errors.New("Dir option is required")
}
2026-06-13 22:43:17 +00:00
if opt.DatabaseName == "" {
return nil, errors.New("DatabaseName option is required")
2026-02-10 14:02:11 +00:00
}
2026-06-13 22:43:17 +00:00
// 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
}
2026-06-15 01:20:30 +03:00
if opt.DataFile == nil {
return nil, errors.New("DataFile option is required")
}
if opt.IndexFile == nil {
return nil, errors.New("IndexFile 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")
}
2026-06-13 22:43:17 +00:00
indexPageNoProvider, err := NewPageNoProvider(PageNoProviderOptions{
FilePath: qb.GetIndexFilePath(opt.Dir, opt.DatabaseName),
PageSize: IndexPageSize,
FreeList: opt.IndexFreeList,
})
dataPageNoProvider, err := NewPageNoProvider(PageNoProviderOptions{
FilePath: qb.GetDataFilePath(opt.Dir, opt.DatabaseName),
PageSize: DataPageSize,
FreeList: opt.DataFreeList,
})
2026-06-13 19:53:27 +03:00
pageIngester, err := NewPageIngester(PageIngesterOptions{
2026-06-13 22:43:17 +00:00
GetIndexPageNumber: indexPageNoProvider.GetPageNumber,
GetDataPageNumber: dataPageNoProvider.GetPageNumber,
2026-06-13 19:53:27 +03:00
})
if err != nil {
return nil, err
}
writePreparer := NewWritePreparer(WritePreparerOptions{
//MinWALBufferSize: 128,
PageIngester: pageIngester,
})
2026-02-10 14:02:11 +00:00
s := &Writer{
2026-06-13 22:43:17 +00:00
inbox: opt.Inbox,
workerInbox: opt.WorkerInbox,
dir: opt.Dir,
databaseName: opt.DatabaseName,
2026-06-15 01:20:30 +03:00
dataFile: opt.DataFile,
indexFile: opt.IndexFile,
2026-06-13 22:43:17 +00:00
dataFreeList: opt.DataFreeList,
indexFreeList: opt.IndexFreeList,
writePreparer: writePreparer,
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
2026-02-10 14:02:11 +00:00
}
2026-06-13 22:43:17 +00:00
if opt.WAL != "" {
s.wal, err = os.OpenFile(filepath.Join(opt.Dir, opt.WAL), os.O_APPEND|os.O_WRONLY, filePerm)
if err != nil {
return nil, err
}
} else {
s.wal, err = os.OpenFile(
qb.GetWALFilePath(opt.Dir, opt.DatabaseName, 0), // .wal_0
os.O_CREATE|os.O_APPEND|os.O_WRONLY,
filePerm,
)
if err != nil {
return nil, err
}
2026-02-10 14:02:11 +00:00
}
2026-06-15 01:20:30 +03:00
2026-02-10 14:02:11 +00:00
return s, nil
}
func (s *Writer) Run() {
for {
select {
2026-06-13 22:43:17 +00:00
case <-s.inbox.Ready():
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-21 23:38:52 +03:00
func (s *Writer) packAndWrite() (err error) {
2026-06-14 23:12:03 +03:00
//fmt.Println("packAndWrite")
2026-06-13 22:43:17 +00:00
//s.mutex.Lock()
2026-05-31 20:01:28 +00:00
//isExited := s.isExited
// var exitWaitGroup *sync.WaitGroup
// if s.isExited {
// exitWaitGroup = s.waitGroup
// }
2026-06-13 22:43:17 +00:00
//s.mutex.Unlock()
input := s.inbox.Drain()
2026-05-10 00:59:47 +00:00
2026-06-14 23:12:03 +03:00
//fmt.Printf("Drain: %#v\n", input)
2026-06-13 19:53:27 +03:00
prepared := s.writePreparer.Prepare(input)
2026-05-21 23:38:52 +03:00
2026-06-14 23:12:03 +03:00
//fmt.Printf("prepared: %#v\n", prepared)
2026-06-15 01:20:30 +03:00
fmt.Println("prepared")
2026-06-14 23:12:03 +03:00
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
}
2026-06-15 01:20:30 +03:00
//fmt.Printf("written to wal: % x\n", prepared.Packet)
2026-05-31 20:01:28 +00:00
if n != len(prepared.Packet) {
2026-06-14 23:12:03 +03:00
return fmt.Errorf("written %d != packet size %d", n, len(prepared.Packet))
2026-05-31 20:01:28 +00:00
}
if err = s.wal.Sync(); err != nil {
return
}
2026-06-15 01:20:30 +03:00
fmt.Println("synced to wal")
2026-05-21 23:38:52 +03:00
// 4. Пишу в atree сторінки
2026-06-13 22:43:17 +00:00
err = WriteDataPages(s.dataFile, prepared.WriteToData)
if err != nil {
return
}
2026-06-15 01:20:30 +03:00
fmt.Println("written to data")
2026-06-13 22:43:17 +00:00
err = WriteIndexPages(s.indexFile, prepared.WriteToIndex)
2026-05-31 20:01:28 +00:00
if err != nil {
return
}
2026-06-15 01:20:30 +03:00
fmt.Println("written to index")
2026-05-31 20:01:28 +00:00
// 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 {
2026-06-09 08:16:16 +03:00
return fmt.Errorf("close wal file: %s", err)
2026-05-31 20:01:28 +00:00
}
2026-06-09 08:16:16 +03:00
s.wal = nil
s.written = 0
2026-05-31 20:01:28 +00:00
2026-06-09 08:16:16 +03:00
snapshotNumberCh := make(chan int, 1)
2026-05-31 20:01:28 +00:00
2026-06-13 22:43:17 +00:00
s.workerInbox.Push(Changes{
2026-06-13 07:13:03 +00:00
Commits: prepared.Commits,
2026-06-09 08:16:16 +03:00
SnapshotNumberCh: snapshotNumberCh,
FrozenIndexPagesCount: s.indexFreeList.Pages(),
IndexPageNumbers: s.indexFreeList.Cached(),
FrozenDataPagesCount: s.dataFreeList.Pages(),
DataPageNumbers: s.dataFreeList.Cached(),
2026-05-31 20:01:28 +00:00
})
2026-06-09 08:16:16 +03:00
snapshotNumber := <-snapshotNumberCh
// копіюю сторінки із delta файла в base файл і потім роблю Truncate
err = s.indexFreeList.Merge()
2026-05-21 23:38:52 +03:00
if err != nil {
2026-06-09 08:16:16 +03:00
return
}
err = s.dataFreeList.Merge()
if err != nil {
return
2026-05-21 23:38:52 +03:00
}
2026-05-31 20:01:28 +00:00
var err error
s.wal, err = os.OpenFile(
2026-06-13 22:43:17 +00:00
qb.GetWALFilePath(s.dir, s.databaseName, snapshotNumber),
2026-05-31 20:01:28 +00:00
os.O_CREATE|os.O_WRONLY,
filePerm,
)
if err != nil {
return fmt.Errorf("create new changes file: %s", err)
}
} else {
2026-06-13 22:43:17 +00:00
s.workerInbox.Push(Changes{
2026-06-13 07:13:03 +00:00
Commits: prepared.Commits,
2026-05-31 20:01:28 +00:00
})
2026-06-15 01:20:30 +03:00
fmt.Println("pushed commits")
2026-05-31 20:01:28 +00:00
}
2026-05-21 23:38:52 +03:00
return nil
2026-02-10 14:02:11 +00:00
}
2026-05-10 00:59:47 +00:00
// 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-06-05 19:43:01 +00:00
// type IndexPayloadToWrite struct {
// PageNo uint32
// Payload []byte
// ZeroLevel bool
// }
2026-05-31 20:01:28 +00:00
2026-06-05 19:43:01 +00:00
type PageToWrite struct {
PageNo uint32
Content []byte
2026-05-31 20:01:28 +00:00
}
2026-06-13 22:43:17 +00:00
func WriteDataPages(file *os.File, pages []PageToWrite) (err error) {
for _, p := range pages {
2026-06-15 01:20:30 +03:00
fmt.Println("pageNo: %d\n", p.PageNo)
}
for _, p := range pages {
2026-06-13 19:53:27 +03:00
if len(p.Content) != DataPageSize {
2026-06-05 19:43:01 +00:00
return fmt.Errorf("wrong data page size: %d", len(p.Content))
}
2026-05-31 20:01:28 +00:00
var (
2026-06-13 19:53:27 +03:00
off = int(p.PageNo-1) * DataPageSize
2026-05-31 20:01:28 +00:00
n int
)
2026-06-13 22:43:17 +00:00
n, err = file.WriteAt(p.Content, int64(off))
2026-05-10 00:59:47 +00:00
if err != nil {
2026-06-15 01:20:30 +03:00
return fmt.Errorf("file.WriteAt: %s; off=%d", err, off)
2026-05-10 00:59:47 +00:00
}
2026-06-13 19:53:27 +03:00
if n != DataPageSize {
return fmt.Errorf("write %d instead of %d", n, DataPageSize)
2026-05-21 23:38:52 +03:00
}
}
2026-06-13 22:43:17 +00:00
return
}
func WriteIndexPages(file *os.File, pages []PageToWrite) (err error) {
for _, p := range pages {
2026-06-13 19:53:27 +03:00
if len(p.Content) != IndexPageSize {
2026-06-05 19:43:01 +00:00
return fmt.Errorf("wrong index page size: %d", len(p.Content))
}
2026-05-31 20:01:28 +00:00
var (
2026-06-13 19:53:27 +03:00
off = int(p.PageNo-1) * IndexPageSize
2026-05-31 20:01:28 +00:00
n int
)
2026-06-13 22:43:17 +00:00
n, err = file.WriteAt(p.Content, int64(off))
2026-05-31 20:01:28 +00:00
if err != nil {
return
}
2026-06-13 19:53:27 +03:00
if n != IndexPageSize {
return fmt.Errorf("write %d instead of %d", n, IndexPageSize)
2026-05-10 00:59:47 +00:00
}
2026-02-10 14:02:11 +00:00
}
2026-06-13 22:43:17 +00:00
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
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
2026-06-09 08:16:16 +03:00
// 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 ?
// }
2026-05-31 20:01:28 +00:00
2026-06-09 08:16:16 +03:00
// 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
// }