From 271046231c4a635958547b4ff1fba0f9f14c2712 Mon Sep 17 00:00:00 2001 From: dima <1.e4.kc6@gmail.com> Date: Mon, 23 Feb 2026 19:37:05 +0000 Subject: [PATCH] wp --- atree/atree.go | 369 ++++++++++++++++++++++++++------------------------ atree/cursor.go | 4 +- atree/io.go | 304 ++++++++++++++++------------------------- atree/io_test.go | 21 +++ atree/misc.go | 4 +- atree/redo/reader.go | 96 ------------- atree/redo/writer.go | 207 ---------------------------- atree/redox/reader.go | 96 +++++++++++++ atree/redox/writer.go | 207 ++++++++++++++++++++++++++++ atree/writer.go | 164 ++++++++++++++++++++++ database/api.go | 60 ++++---- database/database.go | 180 ++++++++++++------------ database/proc.go | 12 +- database/snapshot.go | 10 +- txlog/writer.go | 30 ++-- 15 files changed, 948 insertions(+), 816 deletions(-) create mode 100644 atree/io_test.go delete mode 100644 atree/redo/reader.go delete mode 100644 atree/redo/writer.go create mode 100644 atree/redox/reader.go create mode 100644 atree/redox/writer.go create mode 100644 atree/writer.go diff --git a/atree/atree.go b/atree/atree.go index c20cacf..3d6bf02 100644 --- a/atree/atree.go +++ b/atree/atree.go @@ -3,41 +3,64 @@ package atree import ( "errors" "fmt" + "hash/crc32" "os" "path/filepath" "sync" diploma "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree/redo" "gordenko.dev/dima/qb/bin" ) const ( + PageTypeData = 1 + PageTypeIndex = 2 + filePerm = 0770 + // common + crc32Idx = PageSize - 4 + pageType = PageSize - 5 + // index page - indexRecordsQtyIdx = IndexPageSize - 7 - isDataPageNumbersIdx = IndexPageSize - 5 - indexCRC32Idx = IndexPageSize - 4 + indexRecordsQtyIdx = PageSize - 8 + isDataPageNumbersIdx = PageSize - 6 // data page - timestampsSizeIdx = DataPageSize - 12 - valuesSizeIdx = DataPageSize - 10 - prevPageIdx = DataPageSize - 8 - dataCRC32Idx = DataPageSize - 4 + timestampsSizeIdx = PageSize - 13 + valuesSizeIdx = PageSize - 11 + prevPageIdx = PageSize - 9 timestampSize = 4 pairSize = timestampSize + PageNoSize indexFooterIdx = indexRecordsQtyIdx dataFooterIdx = timestampsSizeIdx - DataPageSize = 8192 - IndexPageSize = 1024 - PageNoSize = 4 + PageSize = 8192 + PageNoSize = 4 // DataPagePayloadSize int = dataFooterIdx ) +const ( + FlagReused byte = 1 // сторінка із FreeList + FlagNewRoot byte = 2 // новая страница +) + +var ( + castagnoliTable = crc32.MakeTable(crc32.Castagnoli) +) + +func calcChecksum(page []byte) uint32 { + return crc32.Checksum(page, castagnoliTable) +} + +type PageToWrite struct { + PageNo uint32 + Data []byte + IsReused bool +} + type FreeList interface { // використовується в allocPage ReservePage() uint32 @@ -50,120 +73,71 @@ type _page struct { } type Atree struct { - redoDir string - indexFreelist FreeList - dataFreelist FreeList - dataFile *os.File - indexFile *os.File - mutex sync.Mutex - allocatedIndexPagesQty uint32 - allocatedDataPagesQty uint32 - indexPages map[uint32]*_page - dataPages map[uint32]*_page - indexWaits map[uint32][]chan readResult - dataWaits map[uint32][]chan readResult - indexPagesToRead []uint32 - dataPagesToRead []uint32 - readSignalCh chan struct{} - writeSignalCh chan struct{} - writeTasksQueue []WriteTask + freelist FreeList + file *os.File + mutex sync.Mutex + allocatedPagesQty uint32 + pages map[uint32]*_page + pageWaits map[uint32][]chan readResult + pagesToRead []uint32 + readSignalCh chan struct{} + writeSignalCh chan struct{} + writeTasksQueue []WriteTask } type Options struct { - Dir string - RedoDir string - DatabaseName string - DataFreeList FreeList - IndexFreeList FreeList + Dir string + DatabaseName string + FreeList FreeList } func New(opt Options) (*Atree, error) { if opt.Dir == "" { return nil, errors.New("Dir option is required") } - if opt.RedoDir == "" { - return nil, errors.New("RedoDir option is required") - } + // if opt.RedoDir == "" { + // return nil, errors.New("RedoDir option is required") + // } if opt.DatabaseName == "" { return nil, errors.New("DatabaseName 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.FreeList == nil { + return nil, errors.New("FreeList option is required") } // открываю или создаю dbName.data и dbName.index файлы var ( - indexFileName = filepath.Join(opt.Dir, opt.DatabaseName+".index") - dataFileName = filepath.Join(opt.Dir, opt.DatabaseName+".data") - - indexFile *os.File - dataFile *os.File - allocatedIndexPagesQty uint32 - allocatedDataPagesQty uint32 + fileName = filepath.Join(opt.Dir, opt.DatabaseName+".db") + file *os.File + allocatedPagesQty uint32 ) // При создании data файла сразу создается индекс, поэтому корректное // состояние БД: либо оба файла есть, либо ни одного файла нет. - isIndexExist, err := isFileExist(indexFileName) - if err != nil { - return nil, fmt.Errorf("check index file is exist: %s", err) - } - - isDataExist, err := isFileExist(dataFileName) + isDataExist, err := isFileExist(fileName) if err != nil { return nil, fmt.Errorf("check data file is exist: %s", err) } - - if isIndexExist { - if isDataExist { - // открываю оба файла - indexFile, allocatedIndexPagesQty, err = openFile(indexFileName, IndexPageSize) - if err != nil { - return nil, fmt.Errorf("open index file: %s", err) - } - - dataFile, allocatedDataPagesQty, err = openFile(dataFileName, DataPageSize) - if err != nil { - return nil, fmt.Errorf("open data file: %s", err) - } - } else { - // нет data файла - return nil, errors.New("not found data file") + if isDataExist { + file, allocatedPagesQty, err = openFile(fileName, PageSize) + if err != nil { + return nil, fmt.Errorf("open data file: %s", err) } } else { - if isDataExist { - // index файла нет - return nil, errors.New("not found index file") - } else { - // нет обоих файлов - indexFile, err = os.OpenFile(indexFileName, os.O_CREATE|os.O_RDWR, filePerm) - if err != nil { - return nil, err - } - - dataFile, err = os.OpenFile(dataFileName, os.O_CREATE|os.O_RDWR, filePerm) - if err != nil { - return nil, err - } + // нет файла + file, err = os.OpenFile(fileName, os.O_CREATE|os.O_RDWR, filePerm) + if err != nil { + return nil, err } } tree := &Atree{ - redoDir: opt.RedoDir, - indexFreelist: opt.IndexFreeList, - dataFreelist: opt.DataFreeList, - indexFile: indexFile, - dataFile: dataFile, - allocatedIndexPagesQty: allocatedIndexPagesQty, - allocatedDataPagesQty: allocatedDataPagesQty, - indexPages: make(map[uint32]*_page), - dataPages: make(map[uint32]*_page), - indexWaits: make(map[uint32][]chan readResult), - dataWaits: make(map[uint32][]chan readResult), - readSignalCh: make(chan struct{}, 1), - writeSignalCh: make(chan struct{}, 1), + //freelist: opt.FreeList, + file: file, + allocatedPagesQty: allocatedPagesQty, + pages: make(map[uint32]*_page), + pageWaits: make(map[uint32][]chan readResult), + readSignalCh: make(chan struct{}, 1), + writeSignalCh: make(chan struct{}, 1), } return tree, nil @@ -185,7 +159,7 @@ func (s *Atree) findDataPage(rootPageNo uint32, timestamp uint32) (uint32, []byt } foundPageNo := findPageNo(buf, timestamp) - s.releaseIndexPage(indexPageNo) + s.releasePage(indexPageNo) if buf[isDataPageNumbersIdx] == 1 { buf, err := s.fetchDataPage(foundPageNo) @@ -256,16 +230,37 @@ type AppendDataPageReq struct { ValuesChunks [][]byte ValuesSize uint16 } +type Report struct { + //IsDataPageReused bool + //DataPageNo uint32 + IsRootChanged bool + NewRootPageNo uint32 + //ReusedIndexPages []uint32 + Pages []PageToWrite +} -func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error) { +// type ChangedPage struct { +// PageNo uint32 +// Data []byte +// IsReused bool +// } + +// AppendDataPage - метод не записує дані в data-файл, а лише змінює дані +// в page cache та freeList і повертає звіт що змінено. +// Цей звіт txlog має записати в transaction log і лише потім можна змінювати data файл. +// Є ідея - записати у index файли заглушки 255,255,255,255 замість номерів сторінок і зберегти зміщення. +// А потім одним викликом отримати із FreeList список вільних сторінок. +// Тому що є проблема із відновленням FreeList після збою, якщо з нього будуть паралельно +// забирати та добавляти номери сторінок інші потоки. +// Це буде працювати, якщо додавання в txlog і маніпуляції із freeList будуть відбуватись в одному потоці +func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ Report, err error) { var ( - flags byte - dataPagesToRelease []uint32 - indexPagesToRelease []uint32 + pagesToRelease []uint32 + report Report ) - newDataPage := s.allocDataPage() - dataPagesToRelease = append(dataPagesToRelease, newDataPage.PageNo) + newDataPage := s.allocPage() + pagesToRelease = append(pagesToRelease, newDataPage.PageNo) chunksToDataPage(newDataPage.Data, chunksToDataPageReq{ PrevPageNo: req.PrevPageNo, @@ -275,18 +270,23 @@ func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error) ValuesSize: req.ValuesSize, }) - redoWriter, err := redo.NewWriter(redo.WriterOptions{ - Dir: s.redoDir, - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - IsDataPageReused: newDataPage.IsReused, - DataPageNo: newDataPage.PageNo, - Page: newDataPage.Data, + report.Pages = append(report.Pages, PageToWrite{ + PageNo: newDataPage.PageNo, + Data: newDataPage.Data, + IsReused: newDataPage.IsReused, }) - if err != nil { - return - } + + // redoWriter, err := NewWriter(WriterOptions{ + // MetricID: req.MetricID, + // Timestamp: req.Timestamp, + // Value: req.Value, + // IsDataPageReused: newDataPage.IsReused, + // DataPageNo: newDataPage.PageNo, + // Page: newDataPage.Data, + // }) + // if err != nil { + // return + // } if req.RootPageNo > 0 { var path pathToDataPage @@ -295,7 +295,7 @@ func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error) return } for _, leg := range path.Legs { - indexPagesToRelease = append(indexPagesToRelease, leg.PageNo) + pagesToRelease = append(pagesToRelease, leg.PageNo) } if path.LastPageNo != req.PrevPageNo { @@ -314,97 +314,116 @@ func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error) ok := appendPair(leg.Data, req.Since, newPageNo) if ok { - err = redoWriter.AppendIndexPage(leg.PageNo, leg.Data, 0) - if err != nil { - return - } + // index + report.Pages = append(report.Pages, PageToWrite{ + PageNo: leg.PageNo, + Data: leg.Data, + }) + // err = redoWriter.AppendIndexPage(leg.PageNo, leg.Data, 0) + // if err != nil { + // return + // } break } - newIndexPage := s.allocIndexPage() - indexPagesToRelease = append(indexPagesToRelease, newIndexPage.PageNo) + newIndexPage := s.allocPage() + pagesToRelease = append(pagesToRelease, newIndexPage.PageNo) appendPair(newIndexPage.Data, req.Since, newPageNo) // ставлю мітку що всі pageNo на сторінці - це data pageNo if legIdx == lastIdx { newIndexPage.Data[isDataPageNumbersIdx] = 1 } - flags = 0 - if newIndexPage.IsReused { - flags |= redo.FlagReused - } - err = redoWriter.AppendIndexPage(newIndexPage.PageNo, newIndexPage.Data, flags) - if err != nil { - return - } + // flags = 0 + // if newIndexPage.IsReused { + // flags |= FlagReused + // } + report.Pages = append(report.Pages, PageToWrite{ + PageNo: newIndexPage.PageNo, + Data: newIndexPage.Data, + IsReused: newIndexPage.IsReused, + }) + // err = redoWriter.AppendIndexPage(newIndexPage.PageNo, newIndexPage.Data, flags) + // if err != nil { + // return + // } // newPageNo = newIndexPage.PageNo if legIdx == 0 { - newRoot := s.allocIndexPage() - indexPagesToRelease = append(indexPagesToRelease, newRoot.PageNo) + newRoot := s.allocPage() + pagesToRelease = append(pagesToRelease, newRoot.PageNo) appendPair(newRoot.Data, getSince(leg.Data), leg.PageNo) // old rootPageNo appendPair(newRoot.Data, req.Since, newIndexPage.PageNo) // Фиксирую новый root в REDO логе - flags = redo.FlagNewRoot - if newRoot.IsReused { - flags |= redo.FlagReused - } - err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) - if err != nil { - return - } + + report.Pages = append(report.Pages, PageToWrite{ + PageNo: newRoot.PageNo, + Data: newRoot.Data, + IsReused: newRoot.IsReused, + }) + report.NewRootPageNo = newRoot.PageNo + // flags = FlagNewRoot + // if newRoot.IsReused { + // flags |= FlagReused + // } + // err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) + // if err != nil { + // return + // } break } } } else { - newRoot := s.allocIndexPage() - indexPagesToRelease = append(indexPagesToRelease, newRoot.PageNo) + newRoot := s.allocPage() + pagesToRelease = append(pagesToRelease, newRoot.PageNo) newRoot.Data[isDataPageNumbersIdx] = 1 appendPair(newRoot.Data, req.Since, newDataPage.PageNo) - flags = redo.FlagNewRoot - if newRoot.IsReused { - flags |= redo.FlagReused - } - err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) - if err != nil { - return - } + report.Pages = append(report.Pages, PageToWrite{ + PageNo: newRoot.PageNo, + Data: newRoot.Data, + IsReused: newRoot.IsReused, + }) + report.NewRootPageNo = newRoot.PageNo + // flags = FlagNewRoot + // if newRoot.IsReused { + // flags |= FlagReused + // } + // err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) + // if err != nil { + // return + // } } - err = redoWriter.Close() - if err != nil { - return - } + // err = redoWriter.Close() + // if err != nil { + // return + // } // На данний момен схема - наступна. Всі сторінки - data та index - зафіксовані в кеші. // Отже запис на диск пройде максимально швидко. Після цього ReferenceCount кожної // сторінки зменшиться на 1. Оскільки на метрику утримується XLock, сторінки мають // ReferenceCount = 1 (немає інших читачів). - waitCh := make(chan struct{}) + // waitCh := make(chan struct{}) - task := WriteTask{ - WaitCh: waitCh, - DataPage: redo.PageToWrite{ - PageNo: newDataPage.PageNo, - Data: newDataPage.Data, - }, - IndexPages: redoWriter.IndexPagesToWrite(), - } + // task := WriteTask{ + // WaitCh: waitCh, + // Pages: report.Pages, + // } - s.appendWriteTaskToQueue(task) + // s.appendWriteTaskToQueue(task) - <-waitCh + // <-waitCh - for _, pageNo := range dataPagesToRelease { - s.releaseDataPage(pageNo) - } - for _, pageNo := range indexPagesToRelease { - s.releaseIndexPage(pageNo) - } - return redoWriter.GetReport(), nil + // for _, pageNo := range dataPagesToRelease { + // s.releasePage(pageNo) + // } + // for _, pageNo := range indexPagesToRelease { + // s.releasePage(pageNo) + // } + return report, nil } // DELETE @@ -439,7 +458,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) { pageNumbers := listPageNumbers(buf) dataPages = append(dataPages, pageNumbers...) - s.releaseIndexPage(rootPageNo) + s.releasePage(rootPageNo) return PageLists{ DataPages: dataPages, @@ -481,7 +500,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) { pageNumbers := listPageNumbers(buf) dataPages = append(dataPages, pageNumbers...) - s.releaseIndexPage(pageNo) + s.releasePage(pageNo) } else { levels = append(levels, &Level{ PageNo: pageNo, @@ -491,7 +510,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) { }) } } else { - s.releaseIndexPage(level.PageNo) + s.releasePage(level.PageNo) levels = levels[:lastIdx] } } diff --git a/atree/cursor.go b/atree/cursor.go index dbe5da5..a2a9157 100644 --- a/atree/cursor.go +++ b/atree/cursor.go @@ -85,7 +85,7 @@ func (s *BackwardCursor) Prev() (uint32, float64, bool, error) { if prevPageNo == 0 { return 0, 0, true, nil } - s.atree.releaseDataPage(s.pageNo) + s.atree.releasePage(s.pageNo) s.pageNo = prevPageNo s.pageData, err = s.atree.fetchDataPage(s.pageNo) @@ -114,7 +114,7 @@ func (s *BackwardCursor) Prev() (uint32, float64, bool, error) { } func (s *BackwardCursor) Close() { - s.atree.releaseDataPage(s.pageNo) + s.atree.releasePage(s.pageNo) } // HELPER diff --git a/atree/io.go b/atree/io.go index 6f486af..43251d3 100644 --- a/atree/io.go +++ b/atree/io.go @@ -3,16 +3,15 @@ package atree import ( "errors" "fmt" - "hash/crc32" "io/fs" "math" "os" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree/redo" "gordenko.dev/dima/qb/bin" ) +// fix - додати ID, щоб потім перемістити із тимчасового буфера в pages, або звільнити type AllocatedPage struct { PageNo uint32 Data []byte @@ -26,114 +25,38 @@ type readResult struct { // INDEX PAGES -func (s *Atree) DeleteIndexPages(pageNumbers []uint32) { +func (s *Atree) DeletePages(pageNumbers []uint32) { s.mutex.Lock() for _, pageNo := range pageNumbers { - delete(s.indexPages, pageNo) + delete(s.pages, pageNo) } s.mutex.Unlock() } func (s *Atree) fetchIndexPage(pageNo uint32) ([]byte, error) { - s.mutex.Lock() - p, ok := s.indexPages[pageNo] - if ok { - p.ReferenceCount++ - s.mutex.Unlock() - return p.Buf, nil - } - - resultCh := make(chan readResult, 1) - s.indexWaits[pageNo] = append(s.indexWaits[pageNo], resultCh) - if len(s.indexWaits[pageNo]) == 1 { - s.indexPagesToRead = append(s.indexPagesToRead, pageNo) - s.mutex.Unlock() - - select { - case s.readSignalCh <- struct{}{}: - default: - } - } else { - s.mutex.Unlock() - } - - result := <-resultCh - if result.Err == nil { - result.Err = s.verifyCRC(result.Data, IndexPageSize) - } - return result.Data, result.Err -} - -func (s *Atree) releaseIndexPage(pageNo uint32) { - s.mutex.Lock() - defer s.mutex.Unlock() - - p, ok := s.indexPages[pageNo] - if ok { - if p.ReferenceCount > 0 { - p.ReferenceCount-- - return - } else { - qb.Abort( - qb.ReferenceCountBug, - fmt.Errorf("call releaseIndexPage on page %d with reference count = %d", - pageNo, p.ReferenceCount), - ) - } - } -} - -func (s *Atree) allocIndexPage() AllocatedPage { - var ( - allocated = AllocatedPage{ - Data: make([]byte, IndexPageSize), - } - ) - - allocated.PageNo = s.indexFreelist.ReservePage() - s.mutex.Lock() - if allocated.PageNo > 0 { - allocated.IsReused = true - } else { - if s.allocatedIndexPagesQty == math.MaxUint32 { - qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index")) - } - s.allocatedIndexPagesQty++ - allocated.PageNo = s.allocatedIndexPagesQty - } - - s.indexPages[allocated.PageNo] = &_page{ - PageNo: allocated.PageNo, - Buf: allocated.Data, - ReferenceCount: 1, - } - s.mutex.Unlock() - return allocated -} - -// DATA PAGES - -func (s *Atree) DeleteDataPages(pageNumbers []uint32) { - s.mutex.Lock() - for _, pageNo := range pageNumbers { - delete(s.dataPages, pageNo) - } - s.mutex.Unlock() + return s.fetchPage(pageNo, PageTypeIndex) } func (s *Atree) fetchDataPage(pageNo uint32) ([]byte, error) { + return s.fetchPage(pageNo, PageTypeData) +} + +func (s *Atree) fetchPage(pageNo uint32, pageTypeIdx byte) ([]byte, error) { s.mutex.Lock() - p, ok := s.dataPages[pageNo] + p, ok := s.pages[pageNo] if ok { + if p.Buf[pageTypeIdx] != pageTypeIdx { + return nil, fmt.Errorf("wrong pageType %d instead of %d", p.Buf[pageTypeIdx], pageTypeIdx) + } p.ReferenceCount++ s.mutex.Unlock() return p.Buf, nil } resultCh := make(chan readResult, 1) - s.dataWaits[pageNo] = append(s.dataWaits[pageNo], resultCh) - if len(s.dataWaits[pageNo]) == 1 { - s.dataPagesToRead = append(s.dataPagesToRead, pageNo) + s.pageWaits[pageNo] = append(s.pageWaits[pageNo], resultCh) + if len(s.pageWaits[pageNo]) == 1 { + s.pagesToRead = append(s.pagesToRead, pageNo) s.mutex.Unlock() select { @@ -143,18 +66,19 @@ func (s *Atree) fetchDataPage(pageNo uint32) ([]byte, error) { } else { s.mutex.Unlock() } + result := <-resultCh if result.Err == nil { - result.Err = s.verifyCRC(result.Data, DataPageSize) + result.Err = s.verifyCRC(result.Data, PageSize) } return result.Data, result.Err } -func (s *Atree) releaseDataPage(pageNo uint32) { +func (s *Atree) releasePage(pageNo uint32) { s.mutex.Lock() defer s.mutex.Unlock() - p, ok := s.dataPages[pageNo] + p, ok := s.pages[pageNo] if ok { if p.ReferenceCount > 0 { p.ReferenceCount-- @@ -162,43 +86,102 @@ func (s *Atree) releaseDataPage(pageNo uint32) { } else { qb.Abort( qb.ReferenceCountBug, - fmt.Errorf("call releaseDataPage on page %d with reference count = %d", + fmt.Errorf("call releasePage on page %d with reference count = %d", pageNo, p.ReferenceCount), ) } } } -func (s *Atree) allocDataPage() AllocatedPage { +func (s *Atree) allocPage() AllocatedPage { var ( allocated = AllocatedPage{ - Data: make([]byte, DataPageSize), + Data: make([]byte, PageSize), } ) - - allocated.PageNo = s.dataFreelist.ReservePage() + allocated.PageNo = s.freelist.ReservePage() + s.mutex.Lock() if allocated.PageNo > 0 { allocated.IsReused = true - s.mutex.Lock() } else { - s.mutex.Lock() - if s.allocatedDataPagesQty == math.MaxUint32 { - qb.Abort(qb.MaxAtreeSizeExceeded, - errors.New("no space in Atree index")) + if s.allocatedPagesQty == math.MaxUint32 { + qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index")) } - s.allocatedDataPagesQty++ - allocated.PageNo = s.allocatedDataPagesQty + s.allocatedPagesQty++ + allocated.PageNo = s.allocatedPagesQty } - s.dataPages[allocated.PageNo] = &_page{ - PageNo: allocated.PageNo, - Buf: allocated.Data, - ReferenceCount: 1, - } + // s.pages[allocated.PageNo] = &_page{ + // // fix pageType + // PageNo: allocated.PageNo, + // Buf: allocated.Data, + // ReferenceCount: 1, + // } s.mutex.Unlock() return allocated } +// // fix - без freelist +// func (s *Atree) allocIndexPage() AllocatedPage { +// var ( +// allocated = AllocatedPage{ +// Data: make([]byte, PageSize), +// } +// ) + +// allocated.PageNo = s.freelist.ReservePage() +// s.mutex.Lock() +// if allocated.PageNo > 0 { +// allocated.IsReused = true +// } else { +// if s.allocatedPagesQty == math.MaxUint32 { +// qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index")) +// } +// s.allocatedPagesQty++ +// allocated.PageNo = s.allocatedPagesQty +// } + +// s.pages[allocated.PageNo] = &_page{ +// // fix pageType +// PageNo: allocated.PageNo, +// Buf: allocated.Data, +// ReferenceCount: 1, +// } +// s.mutex.Unlock() +// return allocated +// } + +// func (s *Atree) allocDataPage() AllocatedPage { +// var ( +// allocated = AllocatedPage{ +// Data: make([]byte, PageSize), +// } +// ) + +// allocated.PageNo = s.freelist.ReservePage() +// s.mutex.Lock() +// if allocated.PageNo > 0 { +// allocated.IsReused = true +// } else { +// if s.allocatedPagesQty == math.MaxUint32 { +// qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index")) +// } +// s.allocatedPagesQty++ +// allocated.PageNo = s.allocatedPagesQty +// } + +// s.pages[allocated.PageNo] = &_page{ +// // fix pageType +// PageNo: allocated.PageNo, +// Buf: allocated.Data, +// ReferenceCount: 1, +// } +// s.mutex.Unlock() +// return allocated +// } + +// DATA PAGES + // READ func (s *Atree) pageReader() { @@ -212,27 +195,25 @@ func (s *Atree) pageReader() { func (s *Atree) readPages() { s.mutex.Lock() - if len(s.indexPagesToRead) == 0 && len(s.dataPagesToRead) == 0 { + if len(s.pagesToRead) == 0 { s.mutex.Unlock() return } - indexPagesToRead := s.indexPagesToRead - s.indexPagesToRead = nil - dataPagesToRead := s.dataPagesToRead - s.dataPagesToRead = nil + pagesToRead := s.pagesToRead + s.pagesToRead = nil s.mutex.Unlock() - for _, pageNo := range dataPagesToRead { - buf := make([]byte, DataPageSize) - off := (pageNo - 1) * DataPageSize - n, err := s.dataFile.ReadAt(buf, int64(off)) - if n != DataPageSize { - err = fmt.Errorf("read %d instead of %d", n, DataPageSize) + for _, pageNo := range pagesToRead { + buf := make([]byte, PageSize) + off := (pageNo - 1) * PageSize + n, err := s.file.ReadAt(buf, int64(off)) + if n != PageSize { + err = fmt.Errorf("read %d instead of %d", n, PageSize) } s.mutex.Lock() - resultChannels := s.dataWaits[pageNo] - delete(s.dataWaits, pageNo) + resultChannels := s.pageWaits[pageNo] + delete(s.pageWaits, pageNo) if err != nil { s.mutex.Unlock() @@ -242,7 +223,8 @@ func (s *Atree) readPages() { } } } else { - s.dataPages[pageNo] = &_page{ + s.pages[pageNo] = &_page{ + // fix - page type PageNo: pageNo, Buf: buf, ReferenceCount: len(resultChannels), @@ -255,41 +237,6 @@ func (s *Atree) readPages() { } } } - - for _, pageNo := range indexPagesToRead { - buf := make([]byte, IndexPageSize) - off := (pageNo - 1) * IndexPageSize - n, err := s.indexFile.ReadAt(buf, int64(off)) - if n != IndexPageSize { - err = fmt.Errorf("read %d instead of %d", n, IndexPageSize) - } - - s.mutex.Lock() - resultChannels := s.indexWaits[pageNo] - delete(s.indexWaits, pageNo) - - if err != nil { - s.mutex.Unlock() - for _, resultCh := range resultChannels { - resultCh <- readResult{ - Err: err, - } - } - } else { - s.indexPages[pageNo] = &_page{ - PageNo: pageNo, - Buf: buf, - ReferenceCount: len(resultChannels), - } - s.mutex.Unlock() - for _, resultCh := range resultChannels { - resultCh <- readResult{ - Data: buf, - } - } - } - - } } // WRITE @@ -307,9 +254,8 @@ func (s *Atree) pageWriter() { } type WriteTask struct { - WaitCh chan struct{} - DataPage redo.PageToWrite - IndexPages []redo.PageToWrite + WaitCh chan struct{} + Pages []PageToWrite } func (s *Atree) appendWriteTaskToQueue(task WriteTask) { @@ -330,31 +276,15 @@ func (s *Atree) writeTasks() error { s.mutex.Unlock() for _, task := range tasks { - // data page - p := task.DataPage - if len(p.Data) != DataPageSize { - return fmt.Errorf("wrong data page %d size: %d", - p.PageNo, len(p.Data)) - } - off := (p.PageNo - 1) * DataPageSize - 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)) - } - - // index pages - for _, p := range task.IndexPages { - if len(p.Data) != IndexPageSize { - return fmt.Errorf("wrong index page %d size: %d", + for _, p := range task.Pages { + if len(p.Data) != PageSize { + return fmt.Errorf("wrong page %d size: %d", p.PageNo, len(p.Data)) } - bin.PutUint32(p.Data[indexCRC32Idx:], crc32.ChecksumIEEE(p.Data[:indexCRC32Idx])) + bin.PutUint32(p.Data[crc32Idx:], calcChecksum(p.Data[:crc32Idx])) - off := (p.PageNo - 1) * IndexPageSize - n, err := s.indexFile.WriteAt(p.Data, int64(off)) + off := (p.PageNo - 1) * PageSize + n, err := s.file.WriteAt(p.Data, int64(off)) if err != nil { return err } @@ -416,7 +346,7 @@ func (s *Atree) ApplyREDO(task WriteTask) { func (s *Atree) verifyCRC(data []byte, pageSize int) error { var ( pos = pageSize - 4 - calculatedCRC = crc32.ChecksumIEEE(data[:pos]) + calculatedCRC = calcChecksum(data[:pos]) storedCRC = bin.GetUint32(data[pos:]) ) if calculatedCRC != storedCRC { diff --git a/atree/io_test.go b/atree/io_test.go new file mode 100644 index 0000000..7ead11f --- /dev/null +++ b/atree/io_test.go @@ -0,0 +1,21 @@ +package atree + +import ( + "hash/crc32" + "testing" +) + +var page = make([]byte, 4096) +var ieeeTable = crc32.MakeTable(crc32.IEEE) + +func BenchmarkCRC32C(b *testing.B) { + for i := 0; i < b.N; i++ { + crc32.Checksum(page, castagnoliTable) + } +} + +func BenchmarkCRC32IEEE(b *testing.B) { + for i := 0; i < b.N; i++ { + crc32.Checksum(page, ieeeTable) + } +} diff --git a/atree/misc.go b/atree/misc.go index c609cba..b15826c 100644 --- a/atree/misc.go +++ b/atree/misc.go @@ -1,8 +1,6 @@ package atree import ( - "hash/crc32" - "gordenko.dev/dima/qb/bin" ) @@ -98,7 +96,7 @@ func chunksToDataPage(buf []byte, req chunksToDataPageReq) { break } } - bin.PutUint32(buf[dataCRC32Idx:], crc32.ChecksumIEEE(buf[:dataCRC32Idx])) + bin.PutUint32(buf[crc32Idx:], calcChecksum(buf[:crc32Idx])) } func setPrevPageNo(buf []byte, pageNo uint32) { diff --git a/atree/redo/reader.go b/atree/redo/reader.go deleted file mode 100644 index 9ae587b..0000000 --- a/atree/redo/reader.go +++ /dev/null @@ -1,96 +0,0 @@ -package redo - -import ( - "fmt" - "hash/crc32" - "io" - "os" - - "gordenko.dev/dima/qb/bin" -) - -type REDOFile struct { - MetricID uint32 - Timestamp uint32 - Value float64 - IsDataPageReused bool - DataPage PageToWrite - IsRootChanged bool - RootPageNo uint32 - ReusedIndexPages []uint32 - IndexPages []PageToWrite -} - -type ReadREDOFileReq struct { - FileName string - DataPageSize int - IndexPageSize int -} - -func ReadREDOFile(req ReadREDOFileReq) (*REDOFile, error) { - buf, err := os.ReadFile(req.FileName) - if err != nil { - return nil, err - } - - if len(buf) < 25 { - return nil, io.EOF - } - - var ( - end = len(buf) - 4 - payload = buf[:end] - checksum = bin.GetUint32(buf[end:]) - calculatedChecksum = crc32.ChecksumIEEE(payload) - ) - - // Помилка чексуми означає що файл або недописаний, або пошкодженний - if checksum != calculatedChecksum { - return nil, fmt.Errorf("written checksum %d not equal calculated checksum %d", - checksum, calculatedChecksum) - } - - var ( - redoLog = REDOFile{ - MetricID: bin.GetUint32(buf[0:]), - Timestamp: bin.GetUint32(buf[4:]), - Value: bin.GetFloat64(buf[8:]), - IsDataPageReused: buf[16] == 1, - DataPage: PageToWrite{ - PageNo: bin.GetUint32(buf[17:]), - Data: buf[21 : 21+req.DataPageSize], - }, - } - pos = 21 + req.DataPageSize - ) - - for { - if pos == len(payload) { - return &redoLog, nil - } - - if pos > len(payload) { - return nil, io.EOF - } - - flags := buf[pos] - - item := PageToWrite{ - PageNo: bin.GetUint32(buf[pos+1:]), - } - pos += 5 // flags + pageNo - item.Data = buf[pos : pos+req.IndexPageSize] - pos += req.IndexPageSize - - redoLog.IndexPages = append(redoLog.IndexPages, item) - - if (flags & FlagReused) == FlagReused { - redoLog.ReusedIndexPages = append(redoLog.ReusedIndexPages, item.PageNo) - } - - if (flags & FlagNewRoot) == FlagNewRoot { - redoLog.IsRootChanged = true - redoLog.RootPageNo = item.PageNo - } - } -} diff --git a/atree/redo/writer.go b/atree/redo/writer.go deleted file mode 100644 index 122e315..0000000 --- a/atree/redo/writer.go +++ /dev/null @@ -1,207 +0,0 @@ -package redo - -import ( - "errors" - "fmt" - "hash" - "hash/crc32" - "os" - "path/filepath" - - "gordenko.dev/dima/qb/bin" -) - -const ( - FlagReused byte = 1 // сторінка із FreeList - FlagNewRoot byte = 2 // новая страница -) - -type PageToWrite struct { - PageNo uint32 - Data []byte -} - -type Writer struct { - metricID uint32 - timestamp uint32 - value float64 - tmp []byte - fileName string - file *os.File - hasher hash.Hash32 - isDataPageReused bool - dataPageNo uint32 - isRootChanged bool - newRootPageNo uint32 - indexPages []uint32 - reusedIndexPages []uint32 - indexPagesToWrite []PageToWrite -} - -type WriterOptions struct { - Dir string - MetricID uint32 - Value float64 - Timestamp uint32 - IsDataPageReused bool - DataPageNo uint32 - Page []byte -} - -// dataPage можно записати 1 раз. Щоб не заплутувати інтерфейс - передаю data сторінку -// через Options. Index сторінок може бути від 1 до N, тому виділяю окремий метод -func NewWriter(opt WriterOptions) (*Writer, error) { - if opt.Dir == "" { - return nil, errors.New("Dir option is required") - } - if opt.MetricID == 0 { - return nil, errors.New("MetricID option is required") - } - if opt.DataPageNo == 0 { - return nil, errors.New("DataPageNo option is required") - } - // if len(opt.Page) != octopus.DataPageSize { - // return nil, fmt.Errorf("bug: wrong data page size %d", len(opt.Page)) - // } - - s := &Writer{ - fileName: JoinREDOFileName(opt.Dir, opt.MetricID), - metricID: opt.MetricID, - timestamp: opt.Timestamp, - value: opt.Value, - tmp: make([]byte, 21), - isDataPageReused: opt.IsDataPageReused, - dataPageNo: opt.DataPageNo, - hasher: crc32.NewIEEE(), - } - - var err error - s.file, err = os.OpenFile(s.fileName, os.O_CREATE|os.O_WRONLY, 0770) - if err != nil { - return nil, err - } - - err = s.init(opt.Page) - if err != nil { - return nil, err - } - return s, nil -} - -/* -Формат: -4b metricID -8b value -4b timestamp -1b flags (reused) -4b dataPageNo -8KB dataPage -*/ -func (s *Writer) init(dataPage []byte) error { - bin.PutUint32(s.tmp[0:], s.metricID) - bin.PutUint32(s.tmp[4:], s.timestamp) - bin.PutFloat64(s.tmp[8:], s.value) - if s.isDataPageReused { - s.tmp[16] = 1 - } - bin.PutUint32(s.tmp[17:], s.dataPageNo) - - _, err := s.file.Write(s.tmp) - if err != nil { - return err - } - - _, err = s.file.Write(dataPage) - if err != nil { - return err - } - - s.hasher.Write(s.tmp) - s.hasher.Write(dataPage) - return nil -} - -/* -Формат -1b index page flags -4b indexPageNo -Nb indexPage -*/ -func (s *Writer) AppendIndexPage(indexPageNo uint32, indexPage []byte, flags byte) error { - s.tmp[0] = flags - bin.PutUint32(s.tmp[1:], indexPageNo) - _, err := s.file.Write(s.tmp[:5]) - if err != nil { - return err - } - _, err = s.file.Write(indexPage) - if err != nil { - return err - } - - s.hasher.Write(s.tmp[:5]) - s.hasher.Write(indexPage) - - s.indexPages = append(s.indexPages, indexPageNo) - - if (flags & FlagReused) == FlagReused { - s.reusedIndexPages = append(s.reusedIndexPages, indexPageNo) - } - - if (flags & FlagNewRoot) == FlagNewRoot { - s.newRootPageNo = indexPageNo - s.isRootChanged = true - } - - s.indexPagesToWrite = append(s.indexPagesToWrite, - PageToWrite{ - PageNo: indexPageNo, - Data: indexPage, - }) - return nil -} - -func (s *Writer) IndexPagesToWrite() []PageToWrite { - return s.indexPagesToWrite -} - -func (s *Writer) Close() (err error) { - // финализирую запись - bin.PutUint32(s.tmp, s.hasher.Sum32()) - _, err = s.file.Write(s.tmp[:4]) - if err != nil { - return err - } - err = s.file.Sync() - if err != nil { - return - } - return s.file.Close() -} - -type Report struct { - FileName string - IsDataPageReused bool - DataPageNo uint32 - IsRootChanged bool - NewRootPageNo uint32 - ReusedIndexPages []uint32 -} - -func (s *Writer) GetReport() Report { - return Report{ - FileName: s.fileName, - IsDataPageReused: s.isDataPageReused, - DataPageNo: s.dataPageNo, - //IndexPages: s.indexPages, - IsRootChanged: s.isRootChanged, - NewRootPageNo: s.newRootPageNo, - ReusedIndexPages: s.reusedIndexPages, - } -} - -// HELPERS - -func JoinREDOFileName(dir string, metricID uint32) string { - return filepath.Join(dir, fmt.Sprintf("m%d.redo", metricID)) -} diff --git a/atree/redox/reader.go b/atree/redox/reader.go new file mode 100644 index 0000000..e0e9a19 --- /dev/null +++ b/atree/redox/reader.go @@ -0,0 +1,96 @@ +package redox + +import ( + "fmt" + "hash/crc32" + "io" + "os" + + "gordenko.dev/dima/qb/bin" +) + +type REDOFile struct { + MetricID uint32 + Timestamp uint32 + Value float64 + IsDataPageReused bool + DataPage PageToWrite + IsRootChanged bool + RootPageNo uint32 + ReusedIndexPages []uint32 + IndexPages []PageToWrite +} + +type ReadREDOFileReq struct { + FileName string + DataPageSize int + IndexPageSize int +} + +func ReadREDOFile(req ReadREDOFileReq) (*REDOFile, error) { + buf, err := os.ReadFile(req.FileName) + if err != nil { + return nil, err + } + + if len(buf) < 25 { + return nil, io.EOF + } + + var ( + end = len(buf) - 4 + payload = buf[:end] + checksum = bin.GetUint32(buf[end:]) + calculatedChecksum = crc32.ChecksumIEEE(payload) + ) + + // Помилка чексуми означає що файл або недописаний, або пошкодженний + if checksum != calculatedChecksum { + return nil, fmt.Errorf("written checksum %d not equal calculated checksum %d", + checksum, calculatedChecksum) + } + + var ( + redoLog = REDOFile{ + MetricID: bin.GetUint32(buf[0:]), + Timestamp: bin.GetUint32(buf[4:]), + Value: bin.GetFloat64(buf[8:]), + IsDataPageReused: buf[16] == 1, + DataPage: PageToWrite{ + PageNo: bin.GetUint32(buf[17:]), + Data: buf[21 : 21+req.DataPageSize], + }, + } + pos = 21 + req.DataPageSize + ) + + for { + if pos == len(payload) { + return &redoLog, nil + } + + if pos > len(payload) { + return nil, io.EOF + } + + flags := buf[pos] + + item := PageToWrite{ + PageNo: bin.GetUint32(buf[pos+1:]), + } + pos += 5 // flags + pageNo + item.Data = buf[pos : pos+req.IndexPageSize] + pos += req.IndexPageSize + + redoLog.IndexPages = append(redoLog.IndexPages, item) + + if (flags & FlagReused) == FlagReused { + redoLog.ReusedIndexPages = append(redoLog.ReusedIndexPages, item.PageNo) + } + + if (flags & FlagNewRoot) == FlagNewRoot { + redoLog.IsRootChanged = true + redoLog.RootPageNo = item.PageNo + } + } +} diff --git a/atree/redox/writer.go b/atree/redox/writer.go new file mode 100644 index 0000000..a36da25 --- /dev/null +++ b/atree/redox/writer.go @@ -0,0 +1,207 @@ +package redox + +import ( + "errors" + "fmt" + "hash" + "hash/crc32" + "os" + "path/filepath" + + "gordenko.dev/dima/qb/bin" +) + +const ( + FlagReused byte = 1 // сторінка із FreeList + FlagNewRoot byte = 2 // новая страница +) + +type PageToWrite struct { + PageNo uint32 + Data []byte +} + +type Writer struct { + metricID uint32 + timestamp uint32 + value float64 + tmp []byte + fileName string + file *os.File + hasher hash.Hash32 + isDataPageReused bool + dataPageNo uint32 + isRootChanged bool + newRootPageNo uint32 + indexPages []uint32 + reusedIndexPages []uint32 + indexPagesToWrite []PageToWrite +} + +type WriterOptions struct { + Dir string + MetricID uint32 + Value float64 + Timestamp uint32 + IsDataPageReused bool + DataPageNo uint32 + Page []byte +} + +// dataPage можно записати 1 раз. Щоб не заплутувати інтерфейс - передаю data сторінку +// через Options. Index сторінок може бути від 1 до N, тому виділяю окремий метод +func NewWriter(opt WriterOptions) (*Writer, error) { + if opt.Dir == "" { + return nil, errors.New("Dir option is required") + } + if opt.MetricID == 0 { + return nil, errors.New("MetricID option is required") + } + if opt.DataPageNo == 0 { + return nil, errors.New("DataPageNo option is required") + } + // if len(opt.Page) != octopus.DataPageSize { + // return nil, fmt.Errorf("bug: wrong data page size %d", len(opt.Page)) + // } + + s := &Writer{ + fileName: JoinREDOFileName(opt.Dir, opt.MetricID), + metricID: opt.MetricID, + timestamp: opt.Timestamp, + value: opt.Value, + tmp: make([]byte, 21), + isDataPageReused: opt.IsDataPageReused, + dataPageNo: opt.DataPageNo, + hasher: crc32.NewIEEE(), + } + + var err error + s.file, err = os.OpenFile(s.fileName, os.O_CREATE|os.O_WRONLY, 0770) + if err != nil { + return nil, err + } + + err = s.init(opt.Page) + if err != nil { + return nil, err + } + return s, nil +} + +/* +Формат: +4b metricID +8b value +4b timestamp +1b flags (reused) +4b dataPageNo +8KB dataPage +*/ +func (s *Writer) init(dataPage []byte) error { + bin.PutUint32(s.tmp[0:], s.metricID) + bin.PutUint32(s.tmp[4:], s.timestamp) + bin.PutFloat64(s.tmp[8:], s.value) + if s.isDataPageReused { + s.tmp[16] = 1 + } + bin.PutUint32(s.tmp[17:], s.dataPageNo) + + _, err := s.file.Write(s.tmp) + if err != nil { + return err + } + + _, err = s.file.Write(dataPage) + if err != nil { + return err + } + + s.hasher.Write(s.tmp) + s.hasher.Write(dataPage) + return nil +} + +/* +Формат +1b index page flags +4b indexPageNo +Nb indexPage +*/ +func (s *Writer) AppendIndexPage(indexPageNo uint32, indexPage []byte, flags byte) error { + s.tmp[0] = flags + bin.PutUint32(s.tmp[1:], indexPageNo) + _, err := s.file.Write(s.tmp[:5]) + if err != nil { + return err + } + _, err = s.file.Write(indexPage) + if err != nil { + return err + } + + s.hasher.Write(s.tmp[:5]) + s.hasher.Write(indexPage) + + s.indexPages = append(s.indexPages, indexPageNo) + + if (flags & FlagReused) == FlagReused { + s.reusedIndexPages = append(s.reusedIndexPages, indexPageNo) + } + + if (flags & FlagNewRoot) == FlagNewRoot { + s.newRootPageNo = indexPageNo + s.isRootChanged = true + } + + s.indexPagesToWrite = append(s.indexPagesToWrite, + PageToWrite{ + PageNo: indexPageNo, + Data: indexPage, + }) + return nil +} + +func (s *Writer) IndexPagesToWrite() []PageToWrite { + return s.indexPagesToWrite +} + +func (s *Writer) Close() (err error) { + // финализирую запись + bin.PutUint32(s.tmp, s.hasher.Sum32()) + _, err = s.file.Write(s.tmp[:4]) + if err != nil { + return err + } + err = s.file.Sync() + if err != nil { + return + } + return s.file.Close() +} + +type Report struct { + FileName string + IsDataPageReused bool + DataPageNo uint32 + IsRootChanged bool + NewRootPageNo uint32 + ReusedIndexPages []uint32 +} + +func (s *Writer) GetReport() Report { + return Report{ + FileName: s.fileName, + IsDataPageReused: s.isDataPageReused, + DataPageNo: s.dataPageNo, + //IndexPages: s.indexPages, + IsRootChanged: s.isRootChanged, + NewRootPageNo: s.newRootPageNo, + ReusedIndexPages: s.reusedIndexPages, + } +} + +// HELPERS + +func JoinREDOFileName(dir string, metricID uint32) string { + return filepath.Join(dir, fmt.Sprintf("m%d.redo", metricID)) +} diff --git a/atree/writer.go b/atree/writer.go new file mode 100644 index 0000000..413d45c --- /dev/null +++ b/atree/writer.go @@ -0,0 +1,164 @@ +package atree + +import ( + "errors" + "hash" + "hash/crc32" +) + +type Writer struct { + metricID uint32 + timestamp uint32 + value float64 + tmp []byte + hasher hash.Hash32 + isDataPageReused bool + dataPageNo uint32 + isRootChanged bool + newRootPageNo uint32 + indexPages []uint32 + reusedIndexPages []uint32 + indexPagesToWrite []PageToWrite +} + +type WriterOptions struct { + MetricID uint32 + Value float64 + Timestamp uint32 + IsDataPageReused bool + DataPageNo uint32 + Page []byte +} + +// dataPage можно записати 1 раз. Щоб не заплутувати інтерфейс - передаю data сторінку +// через Options. Index сторінок може бути від 1 до N, тому виділяю окремий метод +func NewWriter(opt WriterOptions) (*Writer, error) { + if opt.MetricID == 0 { + return nil, errors.New("MetricID option is required") + } + if opt.DataPageNo == 0 { + return nil, errors.New("DataPageNo option is required") + } + // if len(opt.Page) != octopus.DataPageSize { + // return nil, fmt.Errorf("bug: wrong data page size %d", len(opt.Page)) + // } + + s := &Writer{ + metricID: opt.MetricID, + timestamp: opt.Timestamp, + value: opt.Value, + tmp: make([]byte, 21), + isDataPageReused: opt.IsDataPageReused, + dataPageNo: opt.DataPageNo, + hasher: crc32.NewIEEE(), + } + + // var err error + // err = s.init(opt.Page) + // if err != nil { + // return nil, err + // } + return s, nil +} + +/* +Формат: +4b metricID +8b value +4b timestamp +1b flags (reused) +4b dataPageNo +8KB dataPage +*/ +// func (s *Writer) init(dataPage []byte) error { +// bin.PutUint32(s.tmp[0:], s.metricID) +// bin.PutUint32(s.tmp[4:], s.timestamp) +// bin.PutFloat64(s.tmp[8:], s.value) +// if s.isDataPageReused { +// s.tmp[16] = 1 +// } +// bin.PutUint32(s.tmp[17:], s.dataPageNo) + +// _, err := s.file.Write(s.tmp) +// if err != nil { +// return err +// } + +// _, err = s.file.Write(dataPage) +// if err != nil { +// return err +// } + +// s.hasher.Write(s.tmp) +// s.hasher.Write(dataPage) +// return nil +// } + +/* +Формат +1b index page flags +4b indexPageNo +Nb indexPage +*/ +// func (s *Writer) AppendIndexPage(indexPageNo uint32, indexPage []byte, flags byte) error { +// s.tmp[0] = flags +// bin.PutUint32(s.tmp[1:], indexPageNo) +// // _, err := s.file.Write(s.tmp[:5]) +// // if err != nil { +// // return err +// // } +// // _, err = s.file.Write(indexPage) +// // if err != nil { +// // return err +// // } + +// s.hasher.Write(s.tmp[:5]) +// s.hasher.Write(indexPage) + +// s.indexPages = append(s.indexPages, indexPageNo) + +// if (flags & FlagReused) == FlagReused { +// s.reusedIndexPages = append(s.reusedIndexPages, indexPageNo) +// } + +// if (flags & FlagNewRoot) == FlagNewRoot { +// s.newRootPageNo = indexPageNo +// s.isRootChanged = true +// } + +// s.indexPagesToWrite = append(s.indexPagesToWrite, +// PageToWrite{ +// PageNo: indexPageNo, +// Data: indexPage, +// }) +// return nil +// } + +// func (s *Writer) IndexPagesToWrite() []PageToWrite { +// return s.indexPagesToWrite +// } + +// func (s *Writer) Close() (err error) { +// // финализирую запись +// bin.PutUint32(s.tmp, s.hasher.Sum32()) +// _, err = s.file.Write(s.tmp[:4]) +// if err != nil { +// return err +// } +// err = s.file.Sync() +// if err != nil { +// return +// } +// return s.file.Close() +// } + +// func (s *Writer) GetReport() Report { +// return Report{ +// IsDataPageReused: s.isDataPageReused, +// DataPageNo: s.dataPageNo, +// //IndexPages: s.indexPages, +// IsRootChanged: s.isRootChanged, +// NewRootPageNo: s.newRootPageNo, +// ReusedIndexPages: s.reusedIndexPages, +// } +// } diff --git a/database/api.go b/database/api.go index a0b37f8..41e874f 100644 --- a/database/api.go +++ b/database/api.go @@ -372,18 +372,19 @@ func (s *Database) AppendMeasure(req proto.AppendMeasureReq) uint16 { if err != nil { diploma.Abort(diploma.WriteToAtreeFailed, err) } + _ = report waitCh := s.txlog.WriteAppendedMeasureWithOverflow( txlog.AppendedMeasureWithOverflow{ - MetricID: req.MetricID, - Timestamp: req.Timestamp, - Value: req.Value, - IsDataPageReused: report.IsDataPageReused, - DataPageNo: report.DataPageNo, - IsRootChanged: report.IsRootChanged, - RootPageNo: report.NewRootPageNo, - ReusedIndexPages: report.ReusedIndexPages, + // FIX + // MetricID: req.MetricID, + // Timestamp: req.Timestamp, + // Value: req.Value, + // IsDataPageReused: report.IsDataPageReused, + // DataPageNo: report.DataPageNo, + // IsRootChanged: report.IsRootChanged, + // RootPageNo: report.NewRootPageNo, + // ReusedIndexPages: report.ReusedIndexPages, }, - report.FileName, false, ) <-waitCh @@ -445,6 +446,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { ) for idx, measure := range req.Measures { + _ = idx if since == 0 { since = measure.Timestamp } else { @@ -514,26 +516,26 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 { if err != nil { diploma.Abort(diploma.WriteToAtreeFailed, err) } - - prevPageNo = report.DataPageNo - if report.IsRootChanged { - rootPageNo = report.NewRootPageNo - } - waitCh := s.txlog.WriteAppendedMeasureWithOverflow( - txlog.AppendedMeasureWithOverflow{ - MetricID: req.MetricID, - Timestamp: measure.Timestamp, - Value: measure.Value, - IsDataPageReused: report.IsDataPageReused, - DataPageNo: report.DataPageNo, - IsRootChanged: report.IsRootChanged, - RootPageNo: report.NewRootPageNo, - ReusedIndexPages: report.ReusedIndexPages, - }, - report.FileName, - (idx+1) < len(req.Measures), - ) - <-waitCh + _ = report + // FIX + // prevPageNo = report.DataPageNo + // if report.IsRootChanged { + // rootPageNo = report.NewRootPageNo + // } + // waitCh := s.txlog.WriteAppendedMeasureWithOverflow( + // txlog.AppendedMeasureWithOverflow{ + // MetricID: req.MetricID, + // Timestamp: measure.Timestamp, + // Value: measure.Value, + // IsDataPageReused: report.IsDataPageReused, + // DataPageNo: report.DataPageNo, + // IsRootChanged: report.IsRootChanged, + // RootPageNo: report.NewRootPageNo, + // ReusedIndexPages: report.ReusedIndexPages, + // }, + // (idx+1) < len(req.Measures), + // ) + // <-waitCh timestampsBuf = conbuf.New(nil) valuesBuf = conbuf.New(nil) diff --git a/database/database.go b/database/database.go index 5f7489c..890f64d 100644 --- a/database/database.go +++ b/database/database.go @@ -9,13 +9,11 @@ import ( "net" "os" "path/filepath" - "regexp" "sync" "time" diploma "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/atree/redo" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/chunkenc" "gordenko.dev/dima/qb/conbuf" @@ -41,11 +39,9 @@ type Database struct { rLocksToRelease []uint32 metrics map[uint32]*_metric metricLockEntries map[uint32]*metricLockEntry - dataFreeList *freelist.FreeList - indexFreeList *freelist.FreeList + freeList *freelist.FreeList dir string databaseName string - redoDir string txlog *txlog.Writer atree *atree.Atree tcpPort int @@ -92,11 +88,9 @@ func New(opt Options) (_ *Database, err error) { workerSignalCh: make(chan struct{}, 1), dir: opt.Dir, databaseName: opt.DatabaseName, - redoDir: opt.RedoDir, metrics: make(map[uint32]*_metric), metricLockEntries: make(map[uint32]*metricLockEntry), - dataFreeList: freelist.New(), - indexFreeList: freelist.New(), + freeList: freelist.New(), tcpPort: opt.TCPPort, logfile: opt.Logfile, logger: log.New(opt.Logfile, "", log.LstdFlags), @@ -113,11 +107,9 @@ func (s *Database) ListenAndServe() (err error) { } s.atree, err = atree.New(atree.Options{ - Dir: s.dir, - DatabaseName: s.databaseName, - RedoDir: s.redoDir, - DataFreeList: s.dataFreeList, - IndexFreeList: s.indexFreeList, + Dir: s.dir, + DatabaseName: s.databaseName, + FreeList: s.freeList, }) if err != nil { return fmt.Errorf("atree.New: %s", err) @@ -186,26 +178,26 @@ func (s *Database) recovery() { } go s.txlog.Run() - fileNames, err := s.searchREDOFiles() - if err != nil { - diploma.Abort(diploma.SearchREDOFilesFailed, err) - } + // fileNames, err := s.searchREDOFiles() + // if err != nil { + // diploma.Abort(diploma.SearchREDOFilesFailed, err) + // } - if len(fileNames) > 0 { - for _, fileName := range fileNames { - err = s.replayREDOFile(fileName) - if err != nil { - diploma.Abort(diploma.ReplayREDOFileFailed, err) - } - } + // if len(fileNames) > 0 { + // for _, fileName := range fileNames { + // err = s.replayREDOFile(fileName) + // if err != nil { + // diploma.Abort(diploma.ReplayREDOFileFailed, err) + // } + // } - for _, fileName := range fileNames { - err = os.Remove(fileName) - if err != nil { - diploma.Abort(diploma.RemoveREDOFileFailed, err) - } - } - } + // for _, fileName := range fileNames { + // err = os.Remove(fileName) + // if err != nil { + // diploma.Abort(diploma.RemoveREDOFileFailed, err) + // } + // } + // } if recipe != nil { if recipe.CompleteSnapshot { @@ -224,69 +216,69 @@ func (s *Database) recovery() { } } -func (s *Database) searchREDOFiles() ([]string, error) { - var ( - reREDO = regexp.MustCompile(`a\d+\.redo`) - fileNames []string - ) +// func (s *Database) searchREDOFiles() ([]string, error) { +// var ( +// reREDO = regexp.MustCompile(`a\d+\.redo`) +// fileNames []string +// ) - entries, err := os.ReadDir(s.redoDir) - if err != nil { - return nil, err - } +// entries, err := os.ReadDir(s.redoDir) +// if err != nil { +// return nil, err +// } - for _, entry := range entries { - if entry.Type().IsRegular() { - baseName := entry.Name() - if reREDO.MatchString(baseName) { - fileNames = append(fileNames, filepath.Join(s.redoDir, baseName)) - } - } - } - return fileNames, nil -} +// for _, entry := range entries { +// if entry.Type().IsRegular() { +// baseName := entry.Name() +// if reREDO.MatchString(baseName) { +// fileNames = append(fileNames, filepath.Join(s.redoDir, baseName)) +// } +// } +// } +// return fileNames, nil +// } -func (s *Database) replayREDOFile(fileName string) error { - redoFile, err := redo.ReadREDOFile(redo.ReadREDOFileReq{ - FileName: fileName, - DataPageSize: atree.DataPageSize, - IndexPageSize: atree.IndexPageSize, - }) - if err != nil { - return fmt.Errorf("can't read REDO file %s: %s", fileName, err) - } +// func (s *Database) replayREDOFile(fileName string) error { +// redoFile, err := redo.ReadREDOFile(redo.ReadREDOFileReq{ +// FileName: fileName, +// DataPageSize: atree.DataPageSize, +// IndexPageSize: atree.IndexPageSize, +// }) +// if err != nil { +// return fmt.Errorf("can't read REDO file %s: %s", fileName, err) +// } - metric, ok := s.metrics[redoFile.MetricID] - if !ok { - return fmt.Errorf("has REDOFile, metric %d not found", redoFile.MetricID) - } +// metric, ok := s.metrics[redoFile.MetricID] +// if !ok { +// return fmt.Errorf("has REDOFile, metric %d not found", redoFile.MetricID) +// } - if metric.Until < redoFile.Timestamp { - waitCh := make(chan struct{}) - s.atree.ApplyREDO(atree.WriteTask{ - DataPage: redoFile.DataPage, - IndexPages: redoFile.IndexPages, - }) - <-waitCh +// if metric.Until < redoFile.Timestamp { +// waitCh := make(chan struct{}) +// s.atree.ApplyREDO(atree.WriteTask{ +// DataPage: redoFile.DataPage, +// IndexPages: redoFile.IndexPages, +// }) +// <-waitCh - waitCh = s.txlog.WriteAppendedMeasureWithOverflow( - txlog.AppendedMeasureWithOverflow{ - MetricID: redoFile.MetricID, - Timestamp: redoFile.Timestamp, - Value: redoFile.Value, - IsDataPageReused: redoFile.IsDataPageReused, - DataPageNo: redoFile.DataPage.PageNo, - IsRootChanged: redoFile.IsRootChanged, - RootPageNo: redoFile.RootPageNo, - ReusedIndexPages: redoFile.ReusedIndexPages, - }, - fileName, - false, - ) - <-waitCh - } - return nil -} +// waitCh = s.txlog.WriteAppendedMeasureWithOverflow( +// txlog.AppendedMeasureWithOverflow{ +// MetricID: redoFile.MetricID, +// Timestamp: redoFile.Timestamp, +// Value: redoFile.Value, +// IsDataPageReused: redoFile.IsDataPageReused, +// DataPageNo: redoFile.DataPage.PageNo, +// IsRootChanged: redoFile.IsRootChanged, +// RootPageNo: redoFile.RootPageNo, +// ReusedIndexPages: redoFile.ReusedIndexPages, +// }, +// fileName, +// false, +// ) +// <-waitCh +// } +// return nil +// } func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { file, err := os.Open(fileName) @@ -383,10 +375,10 @@ func (s *Database) replayChangesRecord(untyped any) error { case txlog.DeletedMetric: delete(s.metrics, rec.MetricID) if len(rec.FreeDataPages) > 0 { - s.dataFreeList.AddPages(rec.FreeDataPages) + s.freeList.AddPages(rec.FreeDataPages) } if len(rec.FreeIndexPages) > 0 { - s.indexFreeList.AddPages(rec.FreeIndexPages) + s.freeList.AddPages(rec.FreeIndexPages) } case txlog.AppendedMeasure: @@ -431,12 +423,12 @@ func (s *Database) replayChangesRecord(untyped any) error { metric.LastPageNo = rec.DataPageNo // delete free pages if rec.IsDataPageReused { - s.dataFreeList.DeleteReservedPages([]uint32{ + s.freeList.DeleteReservedPages([]uint32{ rec.DataPageNo, }) } if len(rec.ReusedIndexPages) > 0 { - s.indexFreeList.DeleteReservedPages(rec.ReusedIndexPages) + s.freeList.DeleteReservedPages(rec.ReusedIndexPages) } } @@ -445,10 +437,10 @@ func (s *Database) replayChangesRecord(untyped any) error { if ok { metric.DeleteMeasures() if len(rec.FreeDataPages) > 0 { - s.dataFreeList.AddPages(rec.FreeDataPages) + s.freeList.AddPages(rec.FreeDataPages) } if len(rec.FreeDataPages) > 0 { - s.indexFreeList.AddPages(rec.FreeIndexPages) + s.freeList.AddPages(rec.FreeIndexPages) } } diff --git a/database/proc.go b/database/proc.go index 62cf9a6..2717554 100644 --- a/database/proc.go +++ b/database/proc.go @@ -535,13 +535,13 @@ func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWith metric.LastPageNo = rec.DataPageNo if rec.IsDataPageReused { - s.dataFreeList.DeleteReservedPages([]uint32{ + s.freeList.DeleteReservedPages([]uint32{ rec.DataPageNo, }) } if len(rec.ReusedIndexPages) > 0 { - s.indexFreeList.DeleteReservedPages(rec.ReusedIndexPages) + s.freeList.DeleteReservedPages(rec.ReusedIndexPages) } if !extended.HoldLock { @@ -1020,10 +1020,10 @@ func (s *Database) deleteMetric(rec txlog.DeletedMetric) { delete(s.metricLockEntries, rec.MetricID) if len(rec.FreeDataPages) > 0 { - s.dataFreeList.AddPages(rec.FreeDataPages) + s.freeList.AddPages(rec.FreeDataPages) } if len(rec.FreeIndexPages) > 0 { - s.indexFreeList.AddPages(rec.FreeIndexPages) + s.freeList.AddPages(rec.FreeIndexPages) } if len(addMetricReqs) > 0 { @@ -1054,10 +1054,10 @@ func (s *Database) deleteMeasures(rec txlog.DeletedMeasures) { metric.DeleteMeasures() lockEntry.XLock = false if len(rec.FreeDataPages) > 0 { - s.dataFreeList.AddPages(rec.FreeDataPages) + s.freeList.AddPages(rec.FreeDataPages) } if len(rec.FreeDataPages) > 0 { - s.indexFreeList.AddPages(rec.FreeIndexPages) + s.freeList.AddPages(rec.FreeIndexPages) } s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) } diff --git a/database/snapshot.go b/database/snapshot.go index 1d2ddb1..74627cc 100644 --- a/database/snapshot.go +++ b/database/snapshot.go @@ -114,12 +114,12 @@ func (s *Database) dumpSnapshot(logNumber int) (err error) { } } // free data pages - err = freeListWriteTo(s.dataFreeList, dst) + err = freeListWriteTo(s.freeList, dst) if err != nil { return } // free index pages - err = freeListWriteTo(s.indexFreeList, dst) + err = freeListWriteTo(s.freeList, dst) if err != nil { return } @@ -171,7 +171,7 @@ func (s *Database) loadSnapshot(fileName string) (err error) { hasher = crc32.NewIEEE() metricsQty int header = make([]byte, metricHeaderSize) - body = make([]byte, atree.DataPageSize) + body = make([]byte, atree.PageSize) ) file, err := os.Open(fileName) @@ -232,12 +232,12 @@ func (s *Database) loadSnapshot(fileName string) (err error) { s.metrics[metricID] = &metric } - err = restoreFreeList(s.dataFreeList, src) + err = restoreFreeList(s.freeList, src) if err != nil { return fmt.Errorf("restore dataFreeList: %s", err) } - err = restoreFreeList(s.indexFreeList, src) + err = restoreFreeList(s.freeList, src) if err != nil { return fmt.Errorf("restore indexFreeList: %s", err) } diff --git a/txlog/writer.go b/txlog/writer.go index 4c34d7d..d5cc646 100644 --- a/txlog/writer.go +++ b/txlog/writer.go @@ -54,7 +54,6 @@ type Writer struct { dir string file *os.File buf *bytes.Buffer - redoFilesToDelete []string workerReqs []any waitCh chan struct{} appendToWorkerQueue func(any) @@ -148,7 +147,6 @@ func (s *Writer) reset() { 0, 0, 0, 0, // lsn }) - s.redoFilesToDelete = nil s.workerReqs = nil s.waitCh = make(chan struct{}) } @@ -166,7 +164,6 @@ func (s *Writer) flush() error { } if s.buf.Len() > packetPrefixSize { - redoFilesToDelete := s.redoFilesToDelete s.lsn++ lsn := s.lsn packet := make([]byte, s.buf.Len()) @@ -192,13 +189,6 @@ func (s *Writer) flush() error { if err := s.file.Sync(); err != nil { return fmt.Errorf("TxLog sync: %s", err) } - - for _, fileName := range redoFilesToDelete { - err = os.Remove(fileName) - if err != nil { - octopus.Abort(octopus.RemoveREDOFileFailed, err) - } - } } else { s.waitCh = make(chan struct{}) s.mutex.Unlock() @@ -392,8 +382,25 @@ type AppendedMeasureWithOverflowExtended struct { [4b] newRootPageNo 1b reusedIndexPages length [N * 4b] reusedIndexPages +------------------ + +4b metricID +4b timestamp +8b value +4b newRootPageNo (0 if no new root page) + +Nb - varsize (N pages) +[ + + 4b pageNo + 1b isReused + Xb page (page size) + +] +Nb - varsize (last data page size) +Nb - payload */ -func (s *Writer) WriteAppendedMeasureWithOverflow(req AppendedMeasureWithOverflow, redoFileName string, holdLock bool) chan struct{} { +func (s *Writer) WriteAppendedMeasureWithOverflow(req AppendedMeasureWithOverflow, holdLock bool) chan struct{} { size := 24 + len(req.ReusedIndexPages)*4 if req.IsRootChanged { size += 4 @@ -433,7 +440,6 @@ func (s *Writer) WriteAppendedMeasureWithOverflow(req AppendedMeasureWithOverflo Record: req, HoldLock: holdLock, }) - s.redoFilesToDelete = append(s.redoFilesToDelete, redoFileName) s.mutex.Unlock() s.sendSignal()