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