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 }