This commit is contained in:
2026-02-23 19:37:05 +00:00
parent 6011535509
commit 271046231c
15 changed files with 948 additions and 816 deletions

View File

@@ -3,41 +3,64 @@ package atree
import ( import (
"errors" "errors"
"fmt" "fmt"
"hash/crc32"
"os" "os"
"path/filepath" "path/filepath"
"sync" "sync"
diploma "gordenko.dev/dima/qb" diploma "gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree/redo"
"gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/bin"
) )
const ( const (
PageTypeData = 1
PageTypeIndex = 2
filePerm = 0770 filePerm = 0770
// common
crc32Idx = PageSize - 4
pageType = PageSize - 5
// index page // index page
indexRecordsQtyIdx = IndexPageSize - 7 indexRecordsQtyIdx = PageSize - 8
isDataPageNumbersIdx = IndexPageSize - 5 isDataPageNumbersIdx = PageSize - 6
indexCRC32Idx = IndexPageSize - 4
// data page // data page
timestampsSizeIdx = DataPageSize - 12 timestampsSizeIdx = PageSize - 13
valuesSizeIdx = DataPageSize - 10 valuesSizeIdx = PageSize - 11
prevPageIdx = DataPageSize - 8 prevPageIdx = PageSize - 9
dataCRC32Idx = DataPageSize - 4
timestampSize = 4 timestampSize = 4
pairSize = timestampSize + PageNoSize pairSize = timestampSize + PageNoSize
indexFooterIdx = indexRecordsQtyIdx indexFooterIdx = indexRecordsQtyIdx
dataFooterIdx = timestampsSizeIdx dataFooterIdx = timestampsSizeIdx
DataPageSize = 8192 PageSize = 8192
IndexPageSize = 1024
PageNoSize = 4 PageNoSize = 4
// //
DataPagePayloadSize int = dataFooterIdx 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 { type FreeList interface {
// використовується в allocPage // використовується в allocPage
ReservePage() uint32 ReservePage() uint32
@@ -50,20 +73,13 @@ type _page struct {
} }
type Atree struct { type Atree struct {
redoDir string freelist FreeList
indexFreelist FreeList file *os.File
dataFreelist FreeList
dataFile *os.File
indexFile *os.File
mutex sync.Mutex mutex sync.Mutex
allocatedIndexPagesQty uint32 allocatedPagesQty uint32
allocatedDataPagesQty uint32 pages map[uint32]*_page
indexPages map[uint32]*_page pageWaits map[uint32][]chan readResult
dataPages map[uint32]*_page pagesToRead []uint32
indexWaits map[uint32][]chan readResult
dataWaits map[uint32][]chan readResult
indexPagesToRead []uint32
dataPagesToRead []uint32
readSignalCh chan struct{} readSignalCh chan struct{}
writeSignalCh chan struct{} writeSignalCh chan struct{}
writeTasksQueue []WriteTask writeTasksQueue []WriteTask
@@ -71,97 +87,55 @@ type Atree struct {
type Options struct { type Options struct {
Dir string Dir string
RedoDir string
DatabaseName string DatabaseName string
DataFreeList FreeList FreeList FreeList
IndexFreeList FreeList
} }
func New(opt Options) (*Atree, error) { func New(opt Options) (*Atree, error) {
if opt.Dir == "" { if opt.Dir == "" {
return nil, errors.New("Dir option is required") return nil, errors.New("Dir option is required")
} }
if opt.RedoDir == "" { // if opt.RedoDir == "" {
return nil, errors.New("RedoDir option is required") // return nil, errors.New("RedoDir option is required")
} // }
if opt.DatabaseName == "" { if opt.DatabaseName == "" {
return nil, errors.New("DatabaseName option is required") return nil, errors.New("DatabaseName option is required")
} }
if opt.DataFreeList == nil { if opt.FreeList == nil {
return nil, errors.New("DataFreeList option is required") return nil, errors.New("FreeList option is required")
}
if opt.IndexFreeList == nil {
return nil, errors.New("IndexFreeList option is required")
} }
// открываю или создаю dbName.data и dbName.index файлы // открываю или создаю dbName.data и dbName.index файлы
var ( var (
indexFileName = filepath.Join(opt.Dir, opt.DatabaseName+".index") fileName = filepath.Join(opt.Dir, opt.DatabaseName+".db")
dataFileName = filepath.Join(opt.Dir, opt.DatabaseName+".data") file *os.File
allocatedPagesQty uint32
indexFile *os.File
dataFile *os.File
allocatedIndexPagesQty uint32
allocatedDataPagesQty uint32
) )
// При создании data файла сразу создается индекс, поэтому корректное // При создании data файла сразу создается индекс, поэтому корректное
// состояние БД: либо оба файла есть, либо ни одного файла нет. // состояние БД: либо оба файла есть, либо ни одного файла нет.
isIndexExist, err := isFileExist(indexFileName) isDataExist, err := isFileExist(fileName)
if err != nil {
return nil, fmt.Errorf("check index file is exist: %s", err)
}
isDataExist, err := isFileExist(dataFileName)
if err != nil { if err != nil {
return nil, fmt.Errorf("check data file is exist: %s", err) return nil, fmt.Errorf("check data file is exist: %s", err)
} }
if isIndexExist {
if isDataExist { if isDataExist {
// открываю оба файла file, allocatedPagesQty, err = openFile(fileName, PageSize)
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 { if err != nil {
return nil, fmt.Errorf("open data file: %s", err) return nil, fmt.Errorf("open data file: %s", err)
} }
} else { } else {
// нет data файла // нет файла
return nil, errors.New("not found data file") file, err = os.OpenFile(fileName, os.O_CREATE|os.O_RDWR, filePerm)
}
} 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 { if err != nil {
return nil, err return nil, err
} }
dataFile, err = os.OpenFile(dataFileName, os.O_CREATE|os.O_RDWR, filePerm)
if err != nil {
return nil, err
}
}
} }
tree := &Atree{ tree := &Atree{
redoDir: opt.RedoDir, //freelist: opt.FreeList,
indexFreelist: opt.IndexFreeList, file: file,
dataFreelist: opt.DataFreeList, allocatedPagesQty: allocatedPagesQty,
indexFile: indexFile, pages: make(map[uint32]*_page),
dataFile: dataFile, pageWaits: make(map[uint32][]chan readResult),
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), readSignalCh: make(chan struct{}, 1),
writeSignalCh: make(chan struct{}, 1), writeSignalCh: make(chan struct{}, 1),
} }
@@ -185,7 +159,7 @@ func (s *Atree) findDataPage(rootPageNo uint32, timestamp uint32) (uint32, []byt
} }
foundPageNo := findPageNo(buf, timestamp) foundPageNo := findPageNo(buf, timestamp)
s.releaseIndexPage(indexPageNo) s.releasePage(indexPageNo)
if buf[isDataPageNumbersIdx] == 1 { if buf[isDataPageNumbersIdx] == 1 {
buf, err := s.fetchDataPage(foundPageNo) buf, err := s.fetchDataPage(foundPageNo)
@@ -256,16 +230,37 @@ type AppendDataPageReq struct {
ValuesChunks [][]byte ValuesChunks [][]byte
ValuesSize uint16 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 ( var (
flags byte pagesToRelease []uint32
dataPagesToRelease []uint32 report Report
indexPagesToRelease []uint32
) )
newDataPage := s.allocDataPage() newDataPage := s.allocPage()
dataPagesToRelease = append(dataPagesToRelease, newDataPage.PageNo) pagesToRelease = append(pagesToRelease, newDataPage.PageNo)
chunksToDataPage(newDataPage.Data, chunksToDataPageReq{ chunksToDataPage(newDataPage.Data, chunksToDataPageReq{
PrevPageNo: req.PrevPageNo, PrevPageNo: req.PrevPageNo,
@@ -275,18 +270,23 @@ func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error)
ValuesSize: req.ValuesSize, ValuesSize: req.ValuesSize,
}) })
redoWriter, err := redo.NewWriter(redo.WriterOptions{ report.Pages = append(report.Pages, PageToWrite{
Dir: s.redoDir, PageNo: newDataPage.PageNo,
MetricID: req.MetricID, Data: newDataPage.Data,
Timestamp: req.Timestamp, IsReused: newDataPage.IsReused,
Value: req.Value,
IsDataPageReused: newDataPage.IsReused,
DataPageNo: newDataPage.PageNo,
Page: newDataPage.Data,
}) })
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 { if req.RootPageNo > 0 {
var path pathToDataPage var path pathToDataPage
@@ -295,7 +295,7 @@ func (s *Atree) AppendDataPage(req AppendDataPageReq) (_ redo.Report, err error)
return return
} }
for _, leg := range path.Legs { for _, leg := range path.Legs {
indexPagesToRelease = append(indexPagesToRelease, leg.PageNo) pagesToRelease = append(pagesToRelease, leg.PageNo)
} }
if path.LastPageNo != req.PrevPageNo { 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) ok := appendPair(leg.Data, req.Since, newPageNo)
if ok { if ok {
err = redoWriter.AppendIndexPage(leg.PageNo, leg.Data, 0) // index
if err != nil { report.Pages = append(report.Pages, PageToWrite{
return PageNo: leg.PageNo,
} Data: leg.Data,
})
// err = redoWriter.AppendIndexPage(leg.PageNo, leg.Data, 0)
// if err != nil {
// return
// }
break break
} }
newIndexPage := s.allocIndexPage() newIndexPage := s.allocPage()
indexPagesToRelease = append(indexPagesToRelease, newIndexPage.PageNo) pagesToRelease = append(pagesToRelease, newIndexPage.PageNo)
appendPair(newIndexPage.Data, req.Since, newPageNo) appendPair(newIndexPage.Data, req.Since, newPageNo)
// ставлю мітку що всі pageNo на сторінці - це data pageNo // ставлю мітку що всі pageNo на сторінці - це data pageNo
if legIdx == lastIdx { if legIdx == lastIdx {
newIndexPage.Data[isDataPageNumbersIdx] = 1 newIndexPage.Data[isDataPageNumbersIdx] = 1
} }
flags = 0 // flags = 0
if newIndexPage.IsReused { // if newIndexPage.IsReused {
flags |= redo.FlagReused // flags |= FlagReused
} // }
err = redoWriter.AppendIndexPage(newIndexPage.PageNo, newIndexPage.Data, flags) report.Pages = append(report.Pages, PageToWrite{
if err != nil { PageNo: newIndexPage.PageNo,
return Data: newIndexPage.Data,
} IsReused: newIndexPage.IsReused,
})
// err = redoWriter.AppendIndexPage(newIndexPage.PageNo, newIndexPage.Data, flags)
// if err != nil {
// return
// }
// //
newPageNo = newIndexPage.PageNo newPageNo = newIndexPage.PageNo
if legIdx == 0 { if legIdx == 0 {
newRoot := s.allocIndexPage() newRoot := s.allocPage()
indexPagesToRelease = append(indexPagesToRelease, newRoot.PageNo) pagesToRelease = append(pagesToRelease, newRoot.PageNo)
appendPair(newRoot.Data, getSince(leg.Data), leg.PageNo) // old rootPageNo appendPair(newRoot.Data, getSince(leg.Data), leg.PageNo) // old rootPageNo
appendPair(newRoot.Data, req.Since, newIndexPage.PageNo) appendPair(newRoot.Data, req.Since, newIndexPage.PageNo)
// Фиксирую новый root в REDO логе // Фиксирую новый root в REDO логе
flags = redo.FlagNewRoot
if newRoot.IsReused { report.Pages = append(report.Pages, PageToWrite{
flags |= redo.FlagReused PageNo: newRoot.PageNo,
} Data: newRoot.Data,
err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) IsReused: newRoot.IsReused,
if err != nil { })
return report.NewRootPageNo = newRoot.PageNo
} // flags = FlagNewRoot
// if newRoot.IsReused {
// flags |= FlagReused
// }
// err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags)
// if err != nil {
// return
// }
break break
} }
} }
} else { } else {
newRoot := s.allocIndexPage() newRoot := s.allocPage()
indexPagesToRelease = append(indexPagesToRelease, newRoot.PageNo) pagesToRelease = append(pagesToRelease, newRoot.PageNo)
newRoot.Data[isDataPageNumbersIdx] = 1 newRoot.Data[isDataPageNumbersIdx] = 1
appendPair(newRoot.Data, req.Since, newDataPage.PageNo) appendPair(newRoot.Data, req.Since, newDataPage.PageNo)
flags = redo.FlagNewRoot report.Pages = append(report.Pages, PageToWrite{
if newRoot.IsReused { PageNo: newRoot.PageNo,
flags |= redo.FlagReused Data: newRoot.Data,
} IsReused: newRoot.IsReused,
err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags) })
if err != nil { report.NewRootPageNo = newRoot.PageNo
return // flags = FlagNewRoot
} // if newRoot.IsReused {
// flags |= FlagReused
// }
// err = redoWriter.AppendIndexPage(newRoot.PageNo, newRoot.Data, flags)
// if err != nil {
// return
// }
} }
err = redoWriter.Close() // err = redoWriter.Close()
if err != nil { // if err != nil {
return // return
} // }
// На данний момен схема - наступна. Всі сторінки - data та index - зафіксовані в кеші. // На данний момен схема - наступна. Всі сторінки - data та index - зафіксовані в кеші.
// Отже запис на диск пройде максимально швидко. Після цього ReferenceCount кожної // Отже запис на диск пройде максимально швидко. Після цього ReferenceCount кожної
// сторінки зменшиться на 1. Оскільки на метрику утримується XLock, сторінки мають // сторінки зменшиться на 1. Оскільки на метрику утримується XLock, сторінки мають
// ReferenceCount = 1 (немає інших читачів). // ReferenceCount = 1 (немає інших читачів).
waitCh := make(chan struct{}) // waitCh := make(chan struct{})
task := WriteTask{ // task := WriteTask{
WaitCh: waitCh, // WaitCh: waitCh,
DataPage: redo.PageToWrite{ // Pages: report.Pages,
PageNo: newDataPage.PageNo, // }
Data: newDataPage.Data,
},
IndexPages: redoWriter.IndexPagesToWrite(),
}
s.appendWriteTaskToQueue(task) // s.appendWriteTaskToQueue(task)
<-waitCh // <-waitCh
for _, pageNo := range dataPagesToRelease { // for _, pageNo := range dataPagesToRelease {
s.releaseDataPage(pageNo) // s.releasePage(pageNo)
} // }
for _, pageNo := range indexPagesToRelease { // for _, pageNo := range indexPagesToRelease {
s.releaseIndexPage(pageNo) // s.releasePage(pageNo)
} // }
return redoWriter.GetReport(), nil return report, nil
} }
// DELETE // DELETE
@@ -439,7 +458,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) {
pageNumbers := listPageNumbers(buf) pageNumbers := listPageNumbers(buf)
dataPages = append(dataPages, pageNumbers...) dataPages = append(dataPages, pageNumbers...)
s.releaseIndexPage(rootPageNo) s.releasePage(rootPageNo)
return PageLists{ return PageLists{
DataPages: dataPages, DataPages: dataPages,
@@ -481,7 +500,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) {
pageNumbers := listPageNumbers(buf) pageNumbers := listPageNumbers(buf)
dataPages = append(dataPages, pageNumbers...) dataPages = append(dataPages, pageNumbers...)
s.releaseIndexPage(pageNo) s.releasePage(pageNo)
} else { } else {
levels = append(levels, &Level{ levels = append(levels, &Level{
PageNo: pageNo, PageNo: pageNo,
@@ -491,7 +510,7 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ PageLists, err error) {
}) })
} }
} else { } else {
s.releaseIndexPage(level.PageNo) s.releasePage(level.PageNo)
levels = levels[:lastIdx] levels = levels[:lastIdx]
} }
} }

View File

@@ -85,7 +85,7 @@ func (s *BackwardCursor) Prev() (uint32, float64, bool, error) {
if prevPageNo == 0 { if prevPageNo == 0 {
return 0, 0, true, nil return 0, 0, true, nil
} }
s.atree.releaseDataPage(s.pageNo) s.atree.releasePage(s.pageNo)
s.pageNo = prevPageNo s.pageNo = prevPageNo
s.pageData, err = s.atree.fetchDataPage(s.pageNo) s.pageData, err = s.atree.fetchDataPage(s.pageNo)
@@ -114,7 +114,7 @@ func (s *BackwardCursor) Prev() (uint32, float64, bool, error) {
} }
func (s *BackwardCursor) Close() { func (s *BackwardCursor) Close() {
s.atree.releaseDataPage(s.pageNo) s.atree.releasePage(s.pageNo)
} }
// HELPER // HELPER

View File

@@ -3,16 +3,15 @@ package atree
import ( import (
"errors" "errors"
"fmt" "fmt"
"hash/crc32"
"io/fs" "io/fs"
"math" "math"
"os" "os"
"gordenko.dev/dima/qb" "gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree/redo"
"gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/bin"
) )
// fix - додати ID, щоб потім перемістити із тимчасового буфера в pages, або звільнити
type AllocatedPage struct { type AllocatedPage struct {
PageNo uint32 PageNo uint32
Data []byte Data []byte
@@ -26,114 +25,38 @@ type readResult struct {
// INDEX PAGES // INDEX PAGES
func (s *Atree) DeleteIndexPages(pageNumbers []uint32) { func (s *Atree) DeletePages(pageNumbers []uint32) {
s.mutex.Lock() s.mutex.Lock()
for _, pageNo := range pageNumbers { for _, pageNo := range pageNumbers {
delete(s.indexPages, pageNo) delete(s.pages, pageNo)
} }
s.mutex.Unlock() s.mutex.Unlock()
} }
func (s *Atree) fetchIndexPage(pageNo uint32) ([]byte, error) { func (s *Atree) fetchIndexPage(pageNo uint32) ([]byte, error) {
s.mutex.Lock() return s.fetchPage(pageNo, PageTypeIndex)
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()
} }
func (s *Atree) fetchDataPage(pageNo uint32) ([]byte, error) { 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() s.mutex.Lock()
p, ok := s.dataPages[pageNo] p, ok := s.pages[pageNo]
if ok { if ok {
if p.Buf[pageTypeIdx] != pageTypeIdx {
return nil, fmt.Errorf("wrong pageType %d instead of %d", p.Buf[pageTypeIdx], pageTypeIdx)
}
p.ReferenceCount++ p.ReferenceCount++
s.mutex.Unlock() s.mutex.Unlock()
return p.Buf, nil return p.Buf, nil
} }
resultCh := make(chan readResult, 1) resultCh := make(chan readResult, 1)
s.dataWaits[pageNo] = append(s.dataWaits[pageNo], resultCh) s.pageWaits[pageNo] = append(s.pageWaits[pageNo], resultCh)
if len(s.dataWaits[pageNo]) == 1 { if len(s.pageWaits[pageNo]) == 1 {
s.dataPagesToRead = append(s.dataPagesToRead, pageNo) s.pagesToRead = append(s.pagesToRead, pageNo)
s.mutex.Unlock() s.mutex.Unlock()
select { select {
@@ -143,18 +66,19 @@ func (s *Atree) fetchDataPage(pageNo uint32) ([]byte, error) {
} else { } else {
s.mutex.Unlock() s.mutex.Unlock()
} }
result := <-resultCh result := <-resultCh
if result.Err == nil { if result.Err == nil {
result.Err = s.verifyCRC(result.Data, DataPageSize) result.Err = s.verifyCRC(result.Data, PageSize)
} }
return result.Data, result.Err return result.Data, result.Err
} }
func (s *Atree) releaseDataPage(pageNo uint32) { func (s *Atree) releasePage(pageNo uint32) {
s.mutex.Lock() s.mutex.Lock()
defer s.mutex.Unlock() defer s.mutex.Unlock()
p, ok := s.dataPages[pageNo] p, ok := s.pages[pageNo]
if ok { if ok {
if p.ReferenceCount > 0 { if p.ReferenceCount > 0 {
p.ReferenceCount-- p.ReferenceCount--
@@ -162,43 +86,102 @@ func (s *Atree) releaseDataPage(pageNo uint32) {
} else { } else {
qb.Abort( qb.Abort(
qb.ReferenceCountBug, 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), pageNo, p.ReferenceCount),
) )
} }
} }
} }
func (s *Atree) allocDataPage() AllocatedPage { func (s *Atree) allocPage() AllocatedPage {
var ( var (
allocated = AllocatedPage{ allocated = AllocatedPage{
Data: make([]byte, DataPageSize), Data: make([]byte, PageSize),
} }
) )
allocated.PageNo = s.freelist.ReservePage()
allocated.PageNo = s.dataFreelist.ReservePage() s.mutex.Lock()
if allocated.PageNo > 0 { if allocated.PageNo > 0 {
allocated.IsReused = true allocated.IsReused = true
s.mutex.Lock()
} else { } else {
s.mutex.Lock() if s.allocatedPagesQty == math.MaxUint32 {
if s.allocatedDataPagesQty == math.MaxUint32 { qb.Abort(qb.MaxAtreeSizeExceeded, errors.New("no space in Atree index"))
qb.Abort(qb.MaxAtreeSizeExceeded,
errors.New("no space in Atree index"))
} }
s.allocatedDataPagesQty++ s.allocatedPagesQty++
allocated.PageNo = s.allocatedDataPagesQty allocated.PageNo = s.allocatedPagesQty
} }
s.dataPages[allocated.PageNo] = &_page{ // s.pages[allocated.PageNo] = &_page{
PageNo: allocated.PageNo, // // fix pageType
Buf: allocated.Data, // PageNo: allocated.PageNo,
ReferenceCount: 1, // Buf: allocated.Data,
} // ReferenceCount: 1,
// }
s.mutex.Unlock() s.mutex.Unlock()
return allocated 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 // READ
func (s *Atree) pageReader() { func (s *Atree) pageReader() {
@@ -212,27 +195,25 @@ func (s *Atree) pageReader() {
func (s *Atree) readPages() { func (s *Atree) readPages() {
s.mutex.Lock() s.mutex.Lock()
if len(s.indexPagesToRead) == 0 && len(s.dataPagesToRead) == 0 { if len(s.pagesToRead) == 0 {
s.mutex.Unlock() s.mutex.Unlock()
return return
} }
indexPagesToRead := s.indexPagesToRead pagesToRead := s.pagesToRead
s.indexPagesToRead = nil s.pagesToRead = nil
dataPagesToRead := s.dataPagesToRead
s.dataPagesToRead = nil
s.mutex.Unlock() s.mutex.Unlock()
for _, pageNo := range dataPagesToRead { for _, pageNo := range pagesToRead {
buf := make([]byte, DataPageSize) buf := make([]byte, PageSize)
off := (pageNo - 1) * DataPageSize off := (pageNo - 1) * PageSize
n, err := s.dataFile.ReadAt(buf, int64(off)) n, err := s.file.ReadAt(buf, int64(off))
if n != DataPageSize { if n != PageSize {
err = fmt.Errorf("read %d instead of %d", n, DataPageSize) err = fmt.Errorf("read %d instead of %d", n, PageSize)
} }
s.mutex.Lock() s.mutex.Lock()
resultChannels := s.dataWaits[pageNo] resultChannels := s.pageWaits[pageNo]
delete(s.dataWaits, pageNo) delete(s.pageWaits, pageNo)
if err != nil { if err != nil {
s.mutex.Unlock() s.mutex.Unlock()
@@ -242,7 +223,8 @@ func (s *Atree) readPages() {
} }
} }
} else { } else {
s.dataPages[pageNo] = &_page{ s.pages[pageNo] = &_page{
// fix - page type
PageNo: pageNo, PageNo: pageNo,
Buf: buf, Buf: buf,
ReferenceCount: len(resultChannels), 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 // WRITE
@@ -308,8 +255,7 @@ func (s *Atree) pageWriter() {
type WriteTask struct { type WriteTask struct {
WaitCh chan struct{} WaitCh chan struct{}
DataPage redo.PageToWrite Pages []PageToWrite
IndexPages []redo.PageToWrite
} }
func (s *Atree) appendWriteTaskToQueue(task WriteTask) { func (s *Atree) appendWriteTaskToQueue(task WriteTask) {
@@ -330,31 +276,15 @@ func (s *Atree) writeTasks() error {
s.mutex.Unlock() s.mutex.Unlock()
for _, task := range tasks { for _, task := range tasks {
// data page for _, p := range task.Pages {
p := task.DataPage if len(p.Data) != PageSize {
if len(p.Data) != DataPageSize { return fmt.Errorf("wrong page %d size: %d",
return fmt.Errorf("wrong data page %d size: %d",
p.PageNo, len(p.Data)) p.PageNo, len(p.Data))
} }
off := (p.PageNo - 1) * DataPageSize bin.PutUint32(p.Data[crc32Idx:], calcChecksum(p.Data[:crc32Idx]))
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 off := (p.PageNo - 1) * PageSize
for _, p := range task.IndexPages { n, err := s.file.WriteAt(p.Data, int64(off))
if len(p.Data) != IndexPageSize {
return fmt.Errorf("wrong index page %d size: %d",
p.PageNo, len(p.Data))
}
bin.PutUint32(p.Data[indexCRC32Idx:], crc32.ChecksumIEEE(p.Data[:indexCRC32Idx]))
off := (p.PageNo - 1) * IndexPageSize
n, err := s.indexFile.WriteAt(p.Data, int64(off))
if err != nil { if err != nil {
return err return err
} }
@@ -416,7 +346,7 @@ func (s *Atree) ApplyREDO(task WriteTask) {
func (s *Atree) verifyCRC(data []byte, pageSize int) error { func (s *Atree) verifyCRC(data []byte, pageSize int) error {
var ( var (
pos = pageSize - 4 pos = pageSize - 4
calculatedCRC = crc32.ChecksumIEEE(data[:pos]) calculatedCRC = calcChecksum(data[:pos])
storedCRC = bin.GetUint32(data[pos:]) storedCRC = bin.GetUint32(data[pos:])
) )
if calculatedCRC != storedCRC { if calculatedCRC != storedCRC {

21
atree/io_test.go Normal file
View File

@@ -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)
}
}

View File

@@ -1,8 +1,6 @@
package atree package atree
import ( import (
"hash/crc32"
"gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/bin"
) )
@@ -98,7 +96,7 @@ func chunksToDataPage(buf []byte, req chunksToDataPageReq) {
break break
} }
} }
bin.PutUint32(buf[dataCRC32Idx:], crc32.ChecksumIEEE(buf[:dataCRC32Idx])) bin.PutUint32(buf[crc32Idx:], calcChecksum(buf[:crc32Idx]))
} }
func setPrevPageNo(buf []byte, pageNo uint32) { func setPrevPageNo(buf []byte, pageNo uint32) {

View File

@@ -1,4 +1,4 @@
package redo package redox
import ( import (
"fmt" "fmt"

View File

@@ -1,4 +1,4 @@
package redo package redox
import ( import (
"errors" "errors"

164
atree/writer.go Normal file
View File

@@ -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,
// }
// }

View File

@@ -372,18 +372,19 @@ func (s *Database) AppendMeasure(req proto.AppendMeasureReq) uint16 {
if err != nil { if err != nil {
diploma.Abort(diploma.WriteToAtreeFailed, err) diploma.Abort(diploma.WriteToAtreeFailed, err)
} }
_ = report
waitCh := s.txlog.WriteAppendedMeasureWithOverflow( waitCh := s.txlog.WriteAppendedMeasureWithOverflow(
txlog.AppendedMeasureWithOverflow{ txlog.AppendedMeasureWithOverflow{
MetricID: req.MetricID, // FIX
Timestamp: req.Timestamp, // MetricID: req.MetricID,
Value: req.Value, // Timestamp: req.Timestamp,
IsDataPageReused: report.IsDataPageReused, // Value: req.Value,
DataPageNo: report.DataPageNo, // IsDataPageReused: report.IsDataPageReused,
IsRootChanged: report.IsRootChanged, // DataPageNo: report.DataPageNo,
RootPageNo: report.NewRootPageNo, // IsRootChanged: report.IsRootChanged,
ReusedIndexPages: report.ReusedIndexPages, // RootPageNo: report.NewRootPageNo,
// ReusedIndexPages: report.ReusedIndexPages,
}, },
report.FileName,
false, false,
) )
<-waitCh <-waitCh
@@ -445,6 +446,7 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 {
) )
for idx, measure := range req.Measures { for idx, measure := range req.Measures {
_ = idx
if since == 0 { if since == 0 {
since = measure.Timestamp since = measure.Timestamp
} else { } else {
@@ -514,26 +516,26 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 {
if err != nil { if err != nil {
diploma.Abort(diploma.WriteToAtreeFailed, err) diploma.Abort(diploma.WriteToAtreeFailed, err)
} }
_ = report
prevPageNo = report.DataPageNo // FIX
if report.IsRootChanged { // prevPageNo = report.DataPageNo
rootPageNo = report.NewRootPageNo // if report.IsRootChanged {
} // rootPageNo = report.NewRootPageNo
waitCh := s.txlog.WriteAppendedMeasureWithOverflow( // }
txlog.AppendedMeasureWithOverflow{ // waitCh := s.txlog.WriteAppendedMeasureWithOverflow(
MetricID: req.MetricID, // txlog.AppendedMeasureWithOverflow{
Timestamp: measure.Timestamp, // MetricID: req.MetricID,
Value: measure.Value, // Timestamp: measure.Timestamp,
IsDataPageReused: report.IsDataPageReused, // Value: measure.Value,
DataPageNo: report.DataPageNo, // IsDataPageReused: report.IsDataPageReused,
IsRootChanged: report.IsRootChanged, // DataPageNo: report.DataPageNo,
RootPageNo: report.NewRootPageNo, // IsRootChanged: report.IsRootChanged,
ReusedIndexPages: report.ReusedIndexPages, // RootPageNo: report.NewRootPageNo,
}, // ReusedIndexPages: report.ReusedIndexPages,
report.FileName, // },
(idx+1) < len(req.Measures), // (idx+1) < len(req.Measures),
) // )
<-waitCh // <-waitCh
timestampsBuf = conbuf.New(nil) timestampsBuf = conbuf.New(nil)
valuesBuf = conbuf.New(nil) valuesBuf = conbuf.New(nil)

View File

@@ -9,13 +9,11 @@ import (
"net" "net"
"os" "os"
"path/filepath" "path/filepath"
"regexp"
"sync" "sync"
"time" "time"
diploma "gordenko.dev/dima/qb" diploma "gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/atree/redo"
"gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/bin"
"gordenko.dev/dima/qb/chunkenc" "gordenko.dev/dima/qb/chunkenc"
"gordenko.dev/dima/qb/conbuf" "gordenko.dev/dima/qb/conbuf"
@@ -41,11 +39,9 @@ type Database struct {
rLocksToRelease []uint32 rLocksToRelease []uint32
metrics map[uint32]*_metric metrics map[uint32]*_metric
metricLockEntries map[uint32]*metricLockEntry metricLockEntries map[uint32]*metricLockEntry
dataFreeList *freelist.FreeList freeList *freelist.FreeList
indexFreeList *freelist.FreeList
dir string dir string
databaseName string databaseName string
redoDir string
txlog *txlog.Writer txlog *txlog.Writer
atree *atree.Atree atree *atree.Atree
tcpPort int tcpPort int
@@ -92,11 +88,9 @@ func New(opt Options) (_ *Database, err error) {
workerSignalCh: make(chan struct{}, 1), workerSignalCh: make(chan struct{}, 1),
dir: opt.Dir, dir: opt.Dir,
databaseName: opt.DatabaseName, databaseName: opt.DatabaseName,
redoDir: opt.RedoDir,
metrics: make(map[uint32]*_metric), metrics: make(map[uint32]*_metric),
metricLockEntries: make(map[uint32]*metricLockEntry), metricLockEntries: make(map[uint32]*metricLockEntry),
dataFreeList: freelist.New(), freeList: freelist.New(),
indexFreeList: freelist.New(),
tcpPort: opt.TCPPort, tcpPort: opt.TCPPort,
logfile: opt.Logfile, logfile: opt.Logfile,
logger: log.New(opt.Logfile, "", log.LstdFlags), logger: log.New(opt.Logfile, "", log.LstdFlags),
@@ -115,9 +109,7 @@ func (s *Database) ListenAndServe() (err error) {
s.atree, err = atree.New(atree.Options{ s.atree, err = atree.New(atree.Options{
Dir: s.dir, Dir: s.dir,
DatabaseName: s.databaseName, DatabaseName: s.databaseName,
RedoDir: s.redoDir, FreeList: s.freeList,
DataFreeList: s.dataFreeList,
IndexFreeList: s.indexFreeList,
}) })
if err != nil { if err != nil {
return fmt.Errorf("atree.New: %s", err) return fmt.Errorf("atree.New: %s", err)
@@ -186,26 +178,26 @@ func (s *Database) recovery() {
} }
go s.txlog.Run() go s.txlog.Run()
fileNames, err := s.searchREDOFiles() // fileNames, err := s.searchREDOFiles()
if err != nil { // if err != nil {
diploma.Abort(diploma.SearchREDOFilesFailed, err) // diploma.Abort(diploma.SearchREDOFilesFailed, err)
} // }
if len(fileNames) > 0 { // if len(fileNames) > 0 {
for _, fileName := range fileNames { // for _, fileName := range fileNames {
err = s.replayREDOFile(fileName) // err = s.replayREDOFile(fileName)
if err != nil { // if err != nil {
diploma.Abort(diploma.ReplayREDOFileFailed, err) // diploma.Abort(diploma.ReplayREDOFileFailed, err)
} // }
} // }
for _, fileName := range fileNames { // for _, fileName := range fileNames {
err = os.Remove(fileName) // err = os.Remove(fileName)
if err != nil { // if err != nil {
diploma.Abort(diploma.RemoveREDOFileFailed, err) // diploma.Abort(diploma.RemoveREDOFileFailed, err)
} // }
} // }
} // }
if recipe != nil { if recipe != nil {
if recipe.CompleteSnapshot { if recipe.CompleteSnapshot {
@@ -224,69 +216,69 @@ func (s *Database) recovery() {
} }
} }
func (s *Database) searchREDOFiles() ([]string, error) { // func (s *Database) searchREDOFiles() ([]string, error) {
var ( // var (
reREDO = regexp.MustCompile(`a\d+\.redo`) // reREDO = regexp.MustCompile(`a\d+\.redo`)
fileNames []string // fileNames []string
) // )
entries, err := os.ReadDir(s.redoDir) // entries, err := os.ReadDir(s.redoDir)
if err != nil { // if err != nil {
return nil, err // return nil, err
} // }
for _, entry := range entries { // for _, entry := range entries {
if entry.Type().IsRegular() { // if entry.Type().IsRegular() {
baseName := entry.Name() // baseName := entry.Name()
if reREDO.MatchString(baseName) { // if reREDO.MatchString(baseName) {
fileNames = append(fileNames, filepath.Join(s.redoDir, baseName)) // fileNames = append(fileNames, filepath.Join(s.redoDir, baseName))
} // }
} // }
} // }
return fileNames, nil // return fileNames, nil
} // }
func (s *Database) replayREDOFile(fileName string) error { // func (s *Database) replayREDOFile(fileName string) error {
redoFile, err := redo.ReadREDOFile(redo.ReadREDOFileReq{ // redoFile, err := redo.ReadREDOFile(redo.ReadREDOFileReq{
FileName: fileName, // FileName: fileName,
DataPageSize: atree.DataPageSize, // DataPageSize: atree.DataPageSize,
IndexPageSize: atree.IndexPageSize, // IndexPageSize: atree.IndexPageSize,
}) // })
if err != nil { // if err != nil {
return fmt.Errorf("can't read REDO file %s: %s", fileName, err) // return fmt.Errorf("can't read REDO file %s: %s", fileName, err)
} // }
metric, ok := s.metrics[redoFile.MetricID] // metric, ok := s.metrics[redoFile.MetricID]
if !ok { // if !ok {
return fmt.Errorf("has REDOFile, metric %d not found", redoFile.MetricID) // return fmt.Errorf("has REDOFile, metric %d not found", redoFile.MetricID)
} // }
if metric.Until < redoFile.Timestamp { // if metric.Until < redoFile.Timestamp {
waitCh := make(chan struct{}) // waitCh := make(chan struct{})
s.atree.ApplyREDO(atree.WriteTask{ // s.atree.ApplyREDO(atree.WriteTask{
DataPage: redoFile.DataPage, // DataPage: redoFile.DataPage,
IndexPages: redoFile.IndexPages, // IndexPages: redoFile.IndexPages,
}) // })
<-waitCh // <-waitCh
waitCh = s.txlog.WriteAppendedMeasureWithOverflow( // waitCh = s.txlog.WriteAppendedMeasureWithOverflow(
txlog.AppendedMeasureWithOverflow{ // txlog.AppendedMeasureWithOverflow{
MetricID: redoFile.MetricID, // MetricID: redoFile.MetricID,
Timestamp: redoFile.Timestamp, // Timestamp: redoFile.Timestamp,
Value: redoFile.Value, // Value: redoFile.Value,
IsDataPageReused: redoFile.IsDataPageReused, // IsDataPageReused: redoFile.IsDataPageReused,
DataPageNo: redoFile.DataPage.PageNo, // DataPageNo: redoFile.DataPage.PageNo,
IsRootChanged: redoFile.IsRootChanged, // IsRootChanged: redoFile.IsRootChanged,
RootPageNo: redoFile.RootPageNo, // RootPageNo: redoFile.RootPageNo,
ReusedIndexPages: redoFile.ReusedIndexPages, // ReusedIndexPages: redoFile.ReusedIndexPages,
}, // },
fileName, // fileName,
false, // false,
) // )
<-waitCh // <-waitCh
} // }
return nil // return nil
} // }
func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { func (s *Database) verifySnapshot(fileName string) (_ bool, err error) {
file, err := os.Open(fileName) file, err := os.Open(fileName)
@@ -383,10 +375,10 @@ func (s *Database) replayChangesRecord(untyped any) error {
case txlog.DeletedMetric: case txlog.DeletedMetric:
delete(s.metrics, rec.MetricID) delete(s.metrics, rec.MetricID)
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.dataFreeList.AddPages(rec.FreeDataPages) s.freeList.AddPages(rec.FreeDataPages)
} }
if len(rec.FreeIndexPages) > 0 { if len(rec.FreeIndexPages) > 0 {
s.indexFreeList.AddPages(rec.FreeIndexPages) s.freeList.AddPages(rec.FreeIndexPages)
} }
case txlog.AppendedMeasure: case txlog.AppendedMeasure:
@@ -431,12 +423,12 @@ func (s *Database) replayChangesRecord(untyped any) error {
metric.LastPageNo = rec.DataPageNo metric.LastPageNo = rec.DataPageNo
// delete free pages // delete free pages
if rec.IsDataPageReused { if rec.IsDataPageReused {
s.dataFreeList.DeleteReservedPages([]uint32{ s.freeList.DeleteReservedPages([]uint32{
rec.DataPageNo, rec.DataPageNo,
}) })
} }
if len(rec.ReusedIndexPages) > 0 { 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 { if ok {
metric.DeleteMeasures() metric.DeleteMeasures()
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.dataFreeList.AddPages(rec.FreeDataPages) s.freeList.AddPages(rec.FreeDataPages)
} }
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.indexFreeList.AddPages(rec.FreeIndexPages) s.freeList.AddPages(rec.FreeIndexPages)
} }
} }

View File

@@ -535,13 +535,13 @@ func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWith
metric.LastPageNo = rec.DataPageNo metric.LastPageNo = rec.DataPageNo
if rec.IsDataPageReused { if rec.IsDataPageReused {
s.dataFreeList.DeleteReservedPages([]uint32{ s.freeList.DeleteReservedPages([]uint32{
rec.DataPageNo, rec.DataPageNo,
}) })
} }
if len(rec.ReusedIndexPages) > 0 { if len(rec.ReusedIndexPages) > 0 {
s.indexFreeList.DeleteReservedPages(rec.ReusedIndexPages) s.freeList.DeleteReservedPages(rec.ReusedIndexPages)
} }
if !extended.HoldLock { if !extended.HoldLock {
@@ -1020,10 +1020,10 @@ func (s *Database) deleteMetric(rec txlog.DeletedMetric) {
delete(s.metricLockEntries, rec.MetricID) delete(s.metricLockEntries, rec.MetricID)
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.dataFreeList.AddPages(rec.FreeDataPages) s.freeList.AddPages(rec.FreeDataPages)
} }
if len(rec.FreeIndexPages) > 0 { if len(rec.FreeIndexPages) > 0 {
s.indexFreeList.AddPages(rec.FreeIndexPages) s.freeList.AddPages(rec.FreeIndexPages)
} }
if len(addMetricReqs) > 0 { if len(addMetricReqs) > 0 {
@@ -1054,10 +1054,10 @@ func (s *Database) deleteMeasures(rec txlog.DeletedMeasures) {
metric.DeleteMeasures() metric.DeleteMeasures()
lockEntry.XLock = false lockEntry.XLock = false
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.dataFreeList.AddPages(rec.FreeDataPages) s.freeList.AddPages(rec.FreeDataPages)
} }
if len(rec.FreeDataPages) > 0 { if len(rec.FreeDataPages) > 0 {
s.indexFreeList.AddPages(rec.FreeIndexPages) s.freeList.AddPages(rec.FreeIndexPages)
} }
s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry)
} }

View File

@@ -114,12 +114,12 @@ func (s *Database) dumpSnapshot(logNumber int) (err error) {
} }
} }
// free data pages // free data pages
err = freeListWriteTo(s.dataFreeList, dst) err = freeListWriteTo(s.freeList, dst)
if err != nil { if err != nil {
return return
} }
// free index pages // free index pages
err = freeListWriteTo(s.indexFreeList, dst) err = freeListWriteTo(s.freeList, dst)
if err != nil { if err != nil {
return return
} }
@@ -171,7 +171,7 @@ func (s *Database) loadSnapshot(fileName string) (err error) {
hasher = crc32.NewIEEE() hasher = crc32.NewIEEE()
metricsQty int metricsQty int
header = make([]byte, metricHeaderSize) header = make([]byte, metricHeaderSize)
body = make([]byte, atree.DataPageSize) body = make([]byte, atree.PageSize)
) )
file, err := os.Open(fileName) file, err := os.Open(fileName)
@@ -232,12 +232,12 @@ func (s *Database) loadSnapshot(fileName string) (err error) {
s.metrics[metricID] = &metric s.metrics[metricID] = &metric
} }
err = restoreFreeList(s.dataFreeList, src) err = restoreFreeList(s.freeList, src)
if err != nil { if err != nil {
return fmt.Errorf("restore dataFreeList: %s", err) return fmt.Errorf("restore dataFreeList: %s", err)
} }
err = restoreFreeList(s.indexFreeList, src) err = restoreFreeList(s.freeList, src)
if err != nil { if err != nil {
return fmt.Errorf("restore indexFreeList: %s", err) return fmt.Errorf("restore indexFreeList: %s", err)
} }

View File

@@ -54,7 +54,6 @@ type Writer struct {
dir string dir string
file *os.File file *os.File
buf *bytes.Buffer buf *bytes.Buffer
redoFilesToDelete []string
workerReqs []any workerReqs []any
waitCh chan struct{} waitCh chan struct{}
appendToWorkerQueue func(any) appendToWorkerQueue func(any)
@@ -148,7 +147,6 @@ func (s *Writer) reset() {
0, 0, 0, 0, // lsn 0, 0, 0, 0, // lsn
}) })
s.redoFilesToDelete = nil
s.workerReqs = nil s.workerReqs = nil
s.waitCh = make(chan struct{}) s.waitCh = make(chan struct{})
} }
@@ -166,7 +164,6 @@ func (s *Writer) flush() error {
} }
if s.buf.Len() > packetPrefixSize { if s.buf.Len() > packetPrefixSize {
redoFilesToDelete := s.redoFilesToDelete
s.lsn++ s.lsn++
lsn := s.lsn lsn := s.lsn
packet := make([]byte, s.buf.Len()) packet := make([]byte, s.buf.Len())
@@ -192,13 +189,6 @@ func (s *Writer) flush() error {
if err := s.file.Sync(); err != nil { if err := s.file.Sync(); err != nil {
return fmt.Errorf("TxLog sync: %s", err) return fmt.Errorf("TxLog sync: %s", err)
} }
for _, fileName := range redoFilesToDelete {
err = os.Remove(fileName)
if err != nil {
octopus.Abort(octopus.RemoveREDOFileFailed, err)
}
}
} else { } else {
s.waitCh = make(chan struct{}) s.waitCh = make(chan struct{})
s.mutex.Unlock() s.mutex.Unlock()
@@ -392,8 +382,25 @@ type AppendedMeasureWithOverflowExtended struct {
[4b] newRootPageNo [4b] newRootPageNo
1b reusedIndexPages length 1b reusedIndexPages length
[N * 4b] reusedIndexPages [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 size := 24 + len(req.ReusedIndexPages)*4
if req.IsRootChanged { if req.IsRootChanged {
size += 4 size += 4
@@ -433,7 +440,6 @@ func (s *Writer) WriteAppendedMeasureWithOverflow(req AppendedMeasureWithOverflo
Record: req, Record: req,
HoldLock: holdLock, HoldLock: holdLock,
}) })
s.redoFilesToDelete = append(s.redoFilesToDelete, redoFileName)
s.mutex.Unlock() s.mutex.Unlock()
s.sendSignal() s.sendSignal()