From db5ecd3dfd3f1b524322bbb1dcb8f7e8170459a0 Mon Sep 17 00:00:00 2001 From: dima <1.e4.kc6@gmail.com> Date: Sun, 31 May 2026 20:01:28 +0000 Subject: [PATCH] tmp --- atree/atree.go | 403 ++++++++++++------------ atree/cursor.go | 28 +- atree/io.go | 62 ++-- atree/misc.go | 112 +++---- atree/redox/reader.go | 96 ------ atree/redox/writer.go | 207 ------------- atree/writer.go | 298 +++++++++--------- atree/x.go | 34 ++- chunkenc/chunkenc_test.go | 28 +- chunkenc/cumdelta.go | 56 +++- chunkenc/insdelta.go | 72 ++++- chunkenc/time_delta.go | 113 +++++-- database/database.go | 212 +++++++------ database/metric.go | 272 +++++++++++++++++ database/proc.go | 762 +++++++--------------------------------------- database/snapshot.go | 279 ----------------- freelist/freelist.go | 193 +++++++++--- freelist/freelist_test.go | 88 +++++- freelist/test.basefree | Bin 0 -> 80 bytes freelist/test.deltafree | 0 freelist/test.free | Bin 20 -> 0 bytes go.mod | 2 +- go.sum | 4 +- qb.go | 9 + txlog/snapshot.go | 122 ++++++++ txlog/txlog.go | 205 ++++++++----- txlog/writer.go | 583 +++++++++++++++++++++++++---------- util/until.go | 5 + 28 files changed, 2090 insertions(+), 2155 deletions(-) delete mode 100644 atree/redox/reader.go delete mode 100644 atree/redox/writer.go delete mode 100644 database/snapshot.go create mode 100644 freelist/test.basefree create mode 100644 freelist/test.deltafree delete mode 100644 freelist/test.free create mode 100644 txlog/snapshot.go diff --git a/atree/atree.go b/atree/atree.go index b53d781..369cdf5 100644 --- a/atree/atree.go +++ b/atree/atree.go @@ -8,7 +8,6 @@ import ( "sync" "gordenko.dev/dima/qb/bin" - "gordenko.dev/dima/qb/util" ) const ( @@ -47,17 +46,6 @@ const ( DataPagePayloadSize int = dataFooterIdx ) -const ( - FlagReused byte = 1 // сторінка із FreeList - FlagNewRoot byte = 2 // новая страница -) - -type PageToWrite struct { - PageNo uint32 - Data []byte - IsReused bool -} - type _page struct { PageNo uint32 Buf []byte @@ -85,13 +73,9 @@ 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.DatabaseName == "" { return nil, errors.New("DatabaseName option is required") } - // открываю или создаю dbName.data и dbName.index файлы var ( fileName = filepath.Join(opt.Dir, opt.DatabaseName+".db") @@ -205,193 +189,6 @@ func (s *Atree) FindPathToLastPage(rootPageNo uint32) (_ PathToDataPage, err err } } -// APPEND DATA PAGE - -//MetricID uint32 -//Timestamp uint32 -//Value float64 -//RootPageNo uint32 -//PrevPageNo uint32 - -type AppendDataPageReq struct { - Since uint32 - TimestampsChunks [][]byte - TimestampsSize uint16 - ValuesChunks [][]byte - ValuesSize uint16 -} - -// 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 будуть відбуватись в одному потоці -// NotLinkedDataPage - це data сторінка із payload, але без встановленого prevPageNo та без розрахованого CRC32 -type NotLinkedDataPage struct { - Since uint32 - Data []byte -} - -func (s NotLinkedDataPage) SetPrevPageNo(prevPageNo uint32) { - bin.PutUint32(s.Data[prevPageIdx:], prevPageNo) - bin.PutUint32(s.Data[crc32Idx:], util.CalcChecksum(s.Data[:crc32Idx])) -} - -type AppendDataPagesReq struct { - LastPageNo uint32 - Legs []PathLeg - DataPages []NotLinkedDataPage -} - -type Report struct { - NewRootPageNo uint32 - LastPageNo uint32 - Pages []PageToWrite -} - -// Можливо atree не треба блокувати при читанні. Я можу створити копії index сторінок, -// а потім під 1 мутексом замінити в pages. -func (s *Atree) AppendDataPages(req AppendDataPagesReq) Report { - var ( - //pagesToRelease []uint32 - newRootPageNo uint32 - lastPageNo = req.LastPageNo - pages []PageToWrite - legs = req.Legs - changed = make(map[uint32]PathLeg) - ) - - for _, p := range req.DataPages { - newDataPageNo, isReused := s.allocPageNumber() // alloc only number - // pagesToRelease = append(pagesToRelease, newDataPage.PageNo) - - // set prevPageNo - p.SetPrevPageNo(lastPageNo) - - pages = append(pages, PageToWrite{ - PageNo: newDataPageNo, - Data: p.Data, - IsReused: isReused, - }) - - // FIX - після додавання index page потрібно модифікувати path, - // або будувати щось типу дерева знизу вверх - - if len(legs) > 0 { - newPageNo := newDataPageNo - lastIdx := len(legs) - 1 - - for legIdx := lastIdx; legIdx >= 0; legIdx-- { - leg := req.Legs[legIdx] - ok := appendPair(leg.Data, p.Since, newPageNo) - if ok { - // index FIX - // потрібно запам'ятати змінені сторінки, але їх можуть змінювати - // кілька ітерацій, тому додавати в Pages не можна. - // на індексній сторінці достатньо місця. Запис вставлено. - changed[leg.PageNo] = leg - break - } - // на індексній сторінці НЕ достатньо місця. Створюю нову. - newIndexPage := s.allocPage() - //pagesToRelease = append(pagesToRelease, newIndexPage.PageNo) - appendPair(newIndexPage.Data, p.Since, newPageNo) - // ставлю мітку що всі pageNo на сторінці - це data pageNo - // fix - єдина оптимізація від існування isDataPageNumbersIdx - getAllPages не завантажує data pages - // if legIdx == lastIdx { - // newIndexPage.Data[isDataPageNumbersIdx] = 1 - // } - pages = append(pages, PageToWrite{ - PageNo: newIndexPage.PageNo, - Data: newIndexPage.Data, - IsReused: newIndexPage.IsReused, - }) - // замінюю крок в path на новий - legs[legIdx] = PathLeg{ - PageNo: newIndexPage.PageNo, - Data: newIndexPage.Data, - } - // - newPageNo = newIndexPage.PageNo - - if legIdx == 0 { - newRoot := s.allocPage() - //pagesToRelease = append(pagesToRelease, newRoot.PageNo) - appendPair(newRoot.Data, getSince(leg.Data), leg.PageNo) // old rootPageNo - appendPair(newRoot.Data, p.Since, newIndexPage.PageNo) - - // Фиксирую новый root в REDO логе - pages = append(pages, PageToWrite{ - PageNo: newRoot.PageNo, - Data: newRoot.Data, - IsReused: newRoot.IsReused, - }) - newRootPageNo = newRoot.PageNo - // додаю в початок списку кроків нову root сторінку - legs = append([]PathLeg{ - { - PageNo: newRoot.PageNo, - Data: newRoot.Data, - }, - }, legs...) - break - } - } - } else { - // індексну root сторінку створюю одразу для першої data сторінки, - // root data сторінки, як в B+Tree не буває. - newRoot := s.allocPage() - //pagesToRelease = append(pagesToRelease, newRoot.PageNo) - //newRoot.Data[isDataPageNumbersIdx] = 1 - appendPair(newRoot.Data, p.Since, newDataPageNo) - - pages = append(pages, PageToWrite{ - PageNo: newRoot.PageNo, - Data: newRoot.Data, - IsReused: newRoot.IsReused, - }) - newRootPageNo = newRoot.PageNo - // додаю в початок списку кроків нову root сторінку - legs = append([]PathLeg{ - { - PageNo: newRoot.PageNo, - Data: newRoot.Data, - }, - }, legs...) - } - - // На данний момен схема - наступна. Всі сторінки - data та index - зафіксовані в кеші. - // Отже запис на диск пройде максимально швидко. Після цього ReferenceCount кожної - // сторінки зменшиться на 1. Оскільки на метрику утримується XLock, сторінки мають - // ReferenceCount = 1 (немає інших читачів). - // for _, pageNo := range indexPagesToRelease { - // s.releasePage(pageNo) - // } - lastPageNo = newDataPageNo - } - for _, leg := range changed { - // fix recalc checksum - pages = append(pages, PageToWrite{ - PageNo: leg.PageNo, - Data: leg.Data, - }) - } - return Report{ - NewRootPageNo: newRootPageNo, - LastPageNo: lastPageNo, - Pages: pages, - } -} - // DELETE type Level struct { @@ -469,3 +266,203 @@ func (s *Atree) GetAllPages(rootPageNo uint32) (_ []uint32, err error) { } } } + +// const ( +// FlagReused byte = 1 // сторінка із FreeList +// FlagNewRoot byte = 2 // новая страница +// ) + +// type PageToWrite struct { +// PageNo uint32 +// Data []byte +// IsReused bool +// } + +//type AppendDataPagesReq struct { +// LastPageNo uint32 +// Legs []PathLeg +// DataPages []NotLinkedDataPage +// } + +// type Report struct { +// NewRootPageNo uint32 +// LastPageNo uint32 +// Pages []PageToWrite +// } + +// // Можливо atree не треба блокувати при читанні. Я можу створити копії index сторінок, +// // а потім під 1 мутексом замінити в pages. +// func (s *Atree) AppendDataPages(req AppendDataPagesReq) Report { +// var ( +// //pagesToRelease []uint32 +// newRootPageNo uint32 +// lastPageNo = req.LastPageNo +// pages []PageToWrite +// legs = req.Legs +// changed = make(map[uint32]PathLeg) +// ) + +// for _, p := range req.DataPages { +// newDataPageNo, isReused := s.allocPageNumber() // alloc only number +// // pagesToRelease = append(pagesToRelease, newDataPage.PageNo) + +// // set prevPageNo +// p.SetPrevPageNo(lastPageNo) + +// pages = append(pages, PageToWrite{ +// PageNo: newDataPageNo, +// Data: p.Data, +// IsReused: isReused, +// }) + +// // FIX - після додавання index page потрібно модифікувати path, +// // або будувати щось типу дерева знизу вверх + +// if len(legs) > 0 { +// newPageNo := newDataPageNo +// lastIdx := len(legs) - 1 + +// for legIdx := lastIdx; legIdx >= 0; legIdx-- { +// leg := req.Legs[legIdx] +// ok := appendPair(leg.Data, p.Since, newPageNo) +// if ok { +// // index FIX +// // потрібно запам'ятати змінені сторінки, але їх можуть змінювати +// // кілька ітерацій, тому додавати в Pages не можна. +// // на індексній сторінці достатньо місця. Запис вставлено. +// changed[leg.PageNo] = leg +// break +// } +// // на індексній сторінці НЕ достатньо місця. Створюю нову. +// newIndexPage := s.allocPage() +// //pagesToRelease = append(pagesToRelease, newIndexPage.PageNo) +// appendPair(newIndexPage.Data, p.Since, newPageNo) +// // ставлю мітку що всі pageNo на сторінці - це data pageNo +// // fix - єдина оптимізація від існування isDataPageNumbersIdx - getAllPages не завантажує data pages +// // if legIdx == lastIdx { +// // newIndexPage.Data[isDataPageNumbersIdx] = 1 +// // } +// pages = append(pages, PageToWrite{ +// PageNo: newIndexPage.PageNo, +// Data: newIndexPage.Data, +// IsReused: newIndexPage.IsReused, +// }) +// // замінюю крок в path на новий +// legs[legIdx] = PathLeg{ +// PageNo: newIndexPage.PageNo, +// Data: newIndexPage.Data, +// } +// // +// newPageNo = newIndexPage.PageNo + +// if legIdx == 0 { +// newRoot := s.allocPage() +// //pagesToRelease = append(pagesToRelease, newRoot.PageNo) +// appendPair(newRoot.Data, getSince(leg.Data), leg.PageNo) // old rootPageNo +// appendPair(newRoot.Data, p.Since, newIndexPage.PageNo) + +// // Фиксирую новый root в REDO логе +// pages = append(pages, PageToWrite{ +// PageNo: newRoot.PageNo, +// Data: newRoot.Data, +// IsReused: newRoot.IsReused, +// }) +// newRootPageNo = newRoot.PageNo +// // додаю в початок списку кроків нову root сторінку +// legs = append([]PathLeg{ +// { +// PageNo: newRoot.PageNo, +// Data: newRoot.Data, +// }, +// }, legs...) +// break +// } +// } +// } else { +// // індексну root сторінку створюю одразу для першої data сторінки, +// // root data сторінки, як в B+Tree не буває. +// newRoot := s.allocPage() +// //pagesToRelease = append(pagesToRelease, newRoot.PageNo) +// //newRoot.Data[isDataPageNumbersIdx] = 1 +// appendPair(newRoot.Data, p.Since, newDataPageNo) + +// pages = append(pages, PageToWrite{ +// PageNo: newRoot.PageNo, +// Data: newRoot.Data, +// IsReused: newRoot.IsReused, +// }) +// newRootPageNo = newRoot.PageNo +// // додаю в початок списку кроків нову root сторінку +// legs = append([]PathLeg{ +// { +// PageNo: newRoot.PageNo, +// Data: newRoot.Data, +// }, +// }, legs...) +// } + +// // На данний момен схема - наступна. Всі сторінки - data та index - зафіксовані в кеші. +// // Отже запис на диск пройде максимально швидко. Після цього ReferenceCount кожної +// // сторінки зменшиться на 1. Оскільки на метрику утримується XLock, сторінки мають +// // ReferenceCount = 1 (немає інших читачів). +// // for _, pageNo := range indexPagesToRelease { +// // s.releasePage(pageNo) +// // } +// lastPageNo = newDataPageNo +// } +// for _, leg := range changed { +// // fix recalc checksum +// pages = append(pages, PageToWrite{ +// PageNo: leg.PageNo, +// Data: leg.Data, +// }) +// } +// return Report{ +// NewRootPageNo: newRootPageNo, +// LastPageNo: lastPageNo, +// Pages: pages, +// } +// } + +// APPEND DATA PAGE + +//MetricID uint32 +//Timestamp uint32 +//Value float64 +//RootPageNo uint32 +//PrevPageNo uint32 + +// type AppendDataPageReq struct { +// Since uint32 +// TimestampsChunks [][]byte +// TimestampsSize uint16 +// ValuesChunks [][]byte +// ValuesSize uint16 +// } + +// 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 будуть відбуватись в одному потоці +// NotLinkedDataPage - це data сторінка із payload, але без встановленого prevPageNo та без розрахованого CRC32 +// type NotLinkedDataPage struct { +// Since uint32 +// Data []byte +// } + +// func (s NotLinkedDataPage) SetPrevPageNo(prevPageNo uint32) { +// bin.PutUint32(s.Data[prevPageIdx:], prevPageNo) +// bin.PutUint32(s.Data[crc32Idx:], util.CalcChecksum(s.Data[:crc32Idx])) +// } + +// diff --git a/atree/cursor.go b/atree/cursor.go index a2a9157..84f4960 100644 --- a/atree/cursor.go +++ b/atree/cursor.go @@ -4,23 +4,23 @@ import ( "errors" "fmt" - octopus "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/enc" ) type BackwardCursor struct { - metricType octopus.MetricType + metricType qb.MetricType fracDigits byte atree *Atree pageNo uint32 pageData []byte - timestampDecompressor octopus.TimestampDecompressor - valueDecompressor octopus.ValueDecompressor + timestampDecompressor qb.TimestampDecompressor + valueDecompressor qb.ValueDecompressor } type BackwardCursorOptions struct { - MetricType octopus.MetricType + MetricType qb.MetricType FracDigits byte PageNo uint32 PageData []byte @@ -29,12 +29,12 @@ type BackwardCursorOptions struct { func NewBackwardCursor(opt BackwardCursorOptions) (*BackwardCursor, error) { switch opt.MetricType { - case octopus.Instant, octopus.Cumulative: + case qb.Instant, qb.Cumulative: // ok default: return nil, fmt.Errorf("MetricType option has wrong value: %d", opt.MetricType) } - if opt.FracDigits > octopus.MaxFracDigits { + if opt.FracDigits > qb.MaxFracDigits { return nil, errors.New("FracDigits option is required") } if opt.Atree == nil { @@ -137,11 +137,11 @@ func (s *BackwardCursor) makeDecompressors() error { vbuf := s.pageData[timestampsPayloadSize : timestampsPayloadSize+valuesPayloadSize] switch s.metricType { - case octopus.Instant: + case qb.Instant: s.valueDecompressor = enc.NewReverseInstantDeltaDecompressor( vbuf, s.fracDigits) - case octopus.Cumulative: + case qb.Cumulative: s.valueDecompressor = enc.NewReverseCumulativeDeltaDecompressor( vbuf, s.fracDigits) @@ -151,8 +151,8 @@ func (s *BackwardCursor) makeDecompressors() error { return nil } -func makeDecompressors(pageData []byte, metricType octopus.MetricType, fracDigits byte) ( - octopus.TimestampDecompressor, octopus.ValueDecompressor, error, +func makeDecompressors(pageData []byte, metricType qb.MetricType, fracDigits byte) ( + qb.TimestampDecompressor, qb.ValueDecompressor, error, ) { timestampsPayloadSize := bin.GetUint16(pageData[timestampsSizeIdx:]) valuesPayloadSize := bin.GetUint16(pageData[valuesSizeIdx:]) @@ -170,13 +170,13 @@ func makeDecompressors(pageData []byte, metricType octopus.MetricType, fracDigit vbuf := pageData[timestampsPayloadSize : timestampsPayloadSize+valuesPayloadSize] - var valueDecompressor octopus.ValueDecompressor + var valueDecompressor qb.ValueDecompressor switch metricType { - case octopus.Instant: + case qb.Instant: valueDecompressor = enc.NewReverseInstantDeltaDecompressor( vbuf, fracDigits) - case octopus.Cumulative: + case qb.Cumulative: valueDecompressor = enc.NewReverseCumulativeDeltaDecompressor( vbuf, fracDigits) diff --git a/atree/io.go b/atree/io.go index d4d6f8e..6f7e3f7 100644 --- a/atree/io.go +++ b/atree/io.go @@ -13,11 +13,11 @@ import ( ) // fix - додати ID, щоб потім перемістити із тимчасового буфера в pages, або звільнити -type AllocatedPage struct { - PageNo uint32 - Data []byte - IsReused bool -} +// type AllocatedPage struct { +// PageNo uint32 +// Data []byte +// IsReused bool +// } type readResult struct { Data []byte @@ -110,33 +110,33 @@ func (s *Atree) allocPageNumber() (uint32, bool) { return 0, true } -func (s *Atree) allocPage() 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 - // } +// func (s *Atree) allocPage() 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 -} +// // 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 { diff --git a/atree/misc.go b/atree/misc.go index fe73400..17868c8 100644 --- a/atree/misc.go +++ b/atree/misc.go @@ -55,51 +55,51 @@ func BinarySearch(qty int, keyComparator bin.KeyComparator) (elemIdx int, isFoun } } -type ChunksToNotLinkedDataPageReq struct { - // PrevPageNo uint32 - TimestampsChunks [][]byte - TimestampsSize uint16 - ValuesChunks [][]byte - ValuesSize uint16 -} +// type ChunksToNotLinkedDataPageReq struct { +// // PrevPageNo uint32 +// TimestampsChunks [][]byte +// TimestampsSize uint16 +// ValuesChunks [][]byte +// ValuesSize uint16 +// } -func ChunksToNotLinkedDataPage(buf []byte, req ChunksToNotLinkedDataPageReq) { - bin.PutUint16(buf[timestampsSizeIdx:], req.TimestampsSize) - bin.PutUint16(buf[valuesSizeIdx:], req.ValuesSize) +// func ChunksToNotLinkedDataPage(buf []byte, req ChunksToNotLinkedDataPageReq) { +// bin.PutUint16(buf[timestampsSizeIdx:], req.TimestampsSize) +// bin.PutUint16(buf[valuesSizeIdx:], req.ValuesSize) - var ( - remainingSize = int(req.TimestampsSize) - pos = 0 - ) - for _, chunk := range req.TimestampsChunks { - if remainingSize >= len(chunk) { - copy(buf[pos:], chunk) - remainingSize -= len(chunk) - pos += len(chunk) - } else { - copy(buf[pos:], chunk[:remainingSize]) - break - } - } +// var ( +// remainingSize = int(req.TimestampsSize) +// pos = 0 +// ) +// for _, chunk := range req.TimestampsChunks { +// if remainingSize >= len(chunk) { +// copy(buf[pos:], chunk) +// remainingSize -= len(chunk) +// pos += len(chunk) +// } else { +// copy(buf[pos:], chunk[:remainingSize]) +// break +// } +// } - remainingSize = int(req.ValuesSize) - pos = int(req.TimestampsSize) +// remainingSize = int(req.ValuesSize) +// pos = int(req.TimestampsSize) - for _, chunk := range req.ValuesChunks { - if remainingSize >= len(chunk) { - copy(buf[pos:], chunk) - remainingSize -= len(chunk) - pos += len(chunk) - } else { - copy(buf[pos:], chunk[:remainingSize]) - break - } - } -} +// for _, chunk := range req.ValuesChunks { +// if remainingSize >= len(chunk) { +// copy(buf[pos:], chunk) +// remainingSize -= len(chunk) +// pos += len(chunk) +// } else { +// copy(buf[pos:], chunk[:remainingSize]) +// break +// } +// } +// } -func setPrevPageNo(buf []byte, pageNo uint32) { - bin.PutUint32(buf[prevPageIdx:], pageNo) -} +// func setPrevPageNo(buf []byte, pageNo uint32) { +// bin.PutUint32(buf[prevPageIdx:], pageNo) +// } func getPrevPageNo(buf []byte) uint32 { return bin.GetUint32(buf[prevPageIdx:]) @@ -191,20 +191,20 @@ func listPageNumbersSince(buf []byte, timestamp uint32) (pageNumbers []uint32) { return } -func getSince(buf []byte) uint32 { - return bin.GetUint32(buf[0:]) -} +// func getSince(buf []byte) uint32 { +// return bin.GetUint32(buf[0:]) +// } -func appendPair(buf []byte, timestamp uint32, pageNo uint32) bool { - qty := bin.GetUint16AsInt(buf[indexRecordsQtyIdx:]) - free := indexFooterIdx - qty*pairSize - if free < pairSize { - return false - } - pos := qty * timestampSize - bin.PutUint32(buf[pos:], timestamp) - pos = indexFooterIdx - (qty+1)*PageNoSize - bin.PutUint32(buf[pos:], pageNo) - bin.PutIntAsUint16(buf[indexRecordsQtyIdx:], qty+1) - return true -} +// func appendPair(buf []byte, timestamp uint32, pageNo uint32) bool { +// qty := bin.GetUint16AsInt(buf[indexRecordsQtyIdx:]) +// free := indexFooterIdx - qty*pairSize +// if free < pairSize { +// return false +// } +// pos := qty * timestampSize +// bin.PutUint32(buf[pos:], timestamp) +// pos = indexFooterIdx - (qty+1)*PageNoSize +// bin.PutUint32(buf[pos:], pageNo) +// bin.PutIntAsUint16(buf[indexRecordsQtyIdx:], qty+1) +// return true +// } diff --git a/atree/redox/reader.go b/atree/redox/reader.go deleted file mode 100644 index e0e9a19..0000000 --- a/atree/redox/reader.go +++ /dev/null @@ -1,96 +0,0 @@ -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 deleted file mode 100644 index a36da25..0000000 --- a/atree/redox/writer.go +++ /dev/null @@ -1,207 +0,0 @@ -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 index 413d45c..ce98d01 100644 --- a/atree/writer.go +++ b/atree/writer.go @@ -1,164 +1,158 @@ 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 +// 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 // } -/* -Формат -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 +// 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.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) +// s := &Writer{ +// metricID: opt.MetricID, +// timestamp: opt.Timestamp, +// value: opt.Value, +// tmp: make([]byte, 21), +// isDataPageReused: opt.IsDataPageReused, +// dataPageNo: opt.DataPageNo, +// hasher: crc32.NewIEEE(), // } -// if (flags & FlagNewRoot) == FlagNewRoot { -// s.newRootPageNo = indexPageNo -// s.isRootChanged = true -// } - -// s.indexPagesToWrite = append(s.indexPagesToWrite, -// PageToWrite{ -// PageNo: indexPageNo, -// Data: indexPage, -// }) -// return nil +// // var err error +// // err = s.init(opt.Page) +// // if err != nil { +// // return nil, err +// // } +// return s, nil // } -// func (s *Writer) IndexPagesToWrite() []PageToWrite { -// return s.indexPagesToWrite -// } +// /* +// Формат: +// 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) -// 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() -// } +// // _, err := s.file.Write(s.tmp) +// // if err != nil { +// // return err +// // } -// 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, -// } -// } +// // _, 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/atree/x.go b/atree/x.go index 8c3a88f..092a4fa 100644 --- a/atree/x.go +++ b/atree/x.go @@ -1,6 +1,7 @@ package atree import ( + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/util" ) @@ -19,7 +20,7 @@ type DataPage struct { ValuesSize int PageNo uint32 IsReused bool - Checksum uint32 + //Checksum uint32 } type IndexPage struct { @@ -27,7 +28,6 @@ type IndexPage struct { Data []byte PageNo uint32 IsReused bool - Checksum uint32 } type IndexLevel struct { // вже записані дані (лише для першої в списку індексної сторінки, беремо із _metric) @@ -57,7 +57,7 @@ type ChunksToDataPageReq struct { } // return checksum -func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) (checksum uint32) { +func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) { var ( remainingSize = int(req.TimestampsSize) pos = 0 @@ -90,30 +90,29 @@ func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) (checksum uint32) { bin.PutUint16(buf[valuesSizeIdx:], uint16(req.ValuesSize)) bin.PutUint32(buf[prevPageIdx:], req.PrevPageNo) - checksum = util.CalcChecksum(buf[:dataCRC32Idx]) + checksum := util.CalcChecksum(buf[:dataCRC32Idx]) bin.PutUint32(buf[dataCRC32Idx:], checksum) - return } -func DataToIndexPage(buf []byte, data []byte, lastLevel bool) (checksum uint32) { +func DataToIndexPage(buf []byte, data []byte, lastLevel bool) { copy(buf, data) bin.PutUint16(buf[indexRecordsQtyIdx:], uint16(len(data)/indexRecordSize)) if lastLevel { buf[isLastLevelIdx] = 1 } - checksum = util.CalcChecksum(buf[:indexCRC32Idx]) + checksum := util.CalcChecksum(buf[:indexCRC32Idx]) bin.PutUint32(buf[indexCRC32Idx:], checksum) - return -} - -type Pager interface { - GetPageNumber() (uint32, bool) } // fix - reduce leves -func AppendDataPagesToTree(pager Pager, levels []*IndexLevel, dataPages []*DataPage) { +func AppendDataPagesToTree(getDataPageNumber func() (uint32, bool, error), getIndexPageNumber func() (uint32, bool, error), levels []*IndexLevel, dataPages []*DataPage) { for i, d := range dataPages[1:] { - d.PageNo, d.IsReused = pager.GetPageNumber() + pageNo, isReused, err := getDataPageNumber() + if err != nil { + qb.Abort(qb.FailedGetPageNumber, err) + } + d.PageNo = pageNo + d.IsReused = isReused d.PrevPageNo = dataPages[i-1].PageNo } @@ -143,7 +142,12 @@ func AppendDataPagesToTree(pager Pager, levels []*IndexLevel, dataPages []*DataP LowerTimestamp: level.LowerTimestamp, Data: level.Data, } - filled.PageNo, filled.IsReused = pager.GetPageNumber() + pageNo, isReused, err := getIndexPageNumber() + if err != nil { + qb.Abort(qb.FailedGetPageNumber, err) + } + filled.PageNo = pageNo + filled.IsReused = isReused level.Filled = append(level.Filled, filled) // diff --git a/chunkenc/chunkenc_test.go b/chunkenc/chunkenc_test.go index 7e62abf..d8fb5cb 100644 --- a/chunkenc/chunkenc_test.go +++ b/chunkenc/chunkenc_test.go @@ -98,17 +98,17 @@ func TestTimeDelta(t *testing.T) { //c.Sync() - bound := c.GetState() + // bound := c.GetState() - fmt.Println("AFTER SYNC") - fmt.Printf("pos: %d\n", c.pos) - chunks := buf.Chunks() - if len(chunks) > 0 { - fmt.Printf("%d\n", chunks[0]) - } + // fmt.Println("AFTER SYNC") + // fmt.Printf("pos: %d\n", c.pos) + // chunks := buf.Chunks() + // if len(chunks) > 0 { + // fmt.Printf("%d\n", chunks[0]) + // } d := NewReverseTimeDeltaDecompressor(buf, c.Size()) - d.RestoreFromBound(bound) + // d.RestoreFromBound(bound) //d.RestoreFromEnd() for range 12 { @@ -210,9 +210,9 @@ func TestInsdeltaBound(t *testing.T) { fmt.Printf("%d\n", buf.Chunks()[0]) ////////////////////////// - bound := c.GetState() + //bound := c.GetState() - fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) + //fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) c.Append(248) c.Append(305) @@ -229,7 +229,7 @@ func TestInsdeltaBound(t *testing.T) { } boundDecompressor := NewReverseInstantDeltaDecompressor(buf, c.Size(), fracDigits) - boundDecompressor.RestoreFromBound(bound) + //boundDecompressor.RestoreFromBound(bound) fmt.Println("from bound:") for range 8 { @@ -270,9 +270,9 @@ func TestInsdeltaBound2(t *testing.T) { fmt.Printf("%d\n", buf.Chunks()[0]) ////////////////////////// - bound := c.GetState() + //bound := c.GetState() - fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) + //fmt.Printf("bound: pos=%d, h=%d\n", bound.Pos, bound.H) c.Append(70) c.Append(10) @@ -289,7 +289,7 @@ func TestInsdeltaBound2(t *testing.T) { } boundDecompressor := NewReverseInstantDeltaDecompressor(buf, c.Size(), fracDigits) - boundDecompressor.RestoreFromBound(bound) + //sboundDecompressor.RestoreFromBound(bound) fmt.Println("from bound:") for range 8 { diff --git a/chunkenc/cumdelta.go b/chunkenc/cumdelta.go index 84bad8c..74717ce 100644 --- a/chunkenc/cumdelta.go +++ b/chunkenc/cumdelta.go @@ -5,6 +5,7 @@ import ( "log" "math" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" ) @@ -180,15 +181,58 @@ type CumulativeDeltaBound struct { // delta h func (s *ReverseCumulativeDeltaCompressor) Lock() { + if s.state != nil { + qb.Abort(qb.RepeatableLock, nil) + } + // позиція посувається вліво, отже може перескочити на попередній chunk + pos := s.pos - 1 - s.lastDeltaSize + chunksQty := pos / conbuf.ChunkSize + if (pos % conbuf.ChunkSize) > 0 { + chunksQty++ + } s.state = &CumulativeDeltaBound{ - Pos: s.pos - 1 - s.lastDeltaSize, + Pos: pos, H: s.h, LastDelta: s.lastDelta, - Chunks: s.buf.Chunks(), + Chunks: s.buf.Chunks()[:chunksQty], } } -func (s *ReverseCumulativeDeltaCompressor) CreateDecompressor(fracDigits byte) *ReverseCumulativeDeltaDecompressor { +// fix - повернути в Pool буфери +func (s *ReverseCumulativeDeltaCompressor) Unlock() { + s.state = nil +} + +func (s *ReverseCumulativeDeltaCompressor) Offset() int { + if s.state != nil { + return s.state.Pos + } + return 0 +} + +func (s *ReverseCumulativeDeltaCompressor) Snapshot() ([][]byte, int) { + if s.state == nil { + return s.buf.Chunks(), s.Size() + } + // ВАЖЛИВО! + // Треба відтворити стан останнього чанка + var ( + pos = s.state.Pos + chunk = make([]byte, conbuf.ChunkSize) + lastChunkIdx = len(s.state.Chunks) - 1 + qtyToCopy = pos % conbuf.ChunkSize + ) + copy(chunk, s.state.Chunks[lastChunkIdx][:qtyToCopy]) + chunks := append(s.state.Chunks[:lastChunkIdx], chunk) + + buf := conbuf.New(chunks) + pos += buf.ReversePutVarUint64(pos, s.state.LastDelta) + buf.SetByte(pos, s.state.H) + pos++ + return chunks, pos +} + +func (s *ReverseCumulativeDeltaCompressor) CreateDecompressor(fracDigits byte) qb.ValueDecompressor { if s.state == nil { d := NewReverseCumulativeDeltaDecompressor(s.buf, s.Size(), fracDigits) d.RestoreFromEnd() @@ -199,11 +243,9 @@ func (s *ReverseCumulativeDeltaCompressor) CreateDecompressor(fracDigits byte) * return d } -func (s *ReverseCumulativeDeltaCompressor) Unlock() { - s.state = nil -} - func (s *ReverseCumulativeDeltaCompressor) Renew() { + // УВАГА! + // state не чіпаємо s.buf = conbuf.New(nil) s.pos = 0 // diff --git a/chunkenc/insdelta.go b/chunkenc/insdelta.go index fa54cee..1c3dee2 100644 --- a/chunkenc/insdelta.go +++ b/chunkenc/insdelta.go @@ -5,6 +5,7 @@ import ( "log" "math" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" ) @@ -16,6 +17,7 @@ type ReverseInstantDeltaCompressor struct { lastDelta int64 lastDeltaSize int h byte + state *InstantDeltaBound } func NewReverseInstantDeltaCompressor(buf *conbuf.ContinuousBuffer, size int, fracDigits byte) *ReverseInstantDeltaCompressor { @@ -147,30 +149,76 @@ type InstantDeltaBound struct { Pos int H byte LastDelta int64 + Chunks [][]byte } -func (s *ReverseInstantDeltaCompressor) GetState() InstantDeltaBound { - return InstantDeltaBound{ - Pos: s.pos - 1 - s.lastDeltaSize, +// delta h +func (s *ReverseInstantDeltaCompressor) Lock() { + if s.state != nil { + qb.Abort(qb.RepeatableLock, nil) + } + // позиція посувається вліво, отже може перескочити на попередній chunk + pos := s.pos - 1 - s.lastDeltaSize + chunksQty := pos / conbuf.ChunkSize + if (pos % conbuf.ChunkSize) > 0 { + chunksQty++ + } + s.state = &InstantDeltaBound{ + Pos: pos, H: s.h, LastDelta: s.lastDelta, + Chunks: s.buf.Chunks()[:chunksQty], } } -func (s *ReverseInstantDeltaCompressor) Lock() { - // s.state = &CumulativeDeltaBound{ - // Pos: s.pos - 1 - s.lastDeltaSize, - // H: s.h, - // LastDelta: s.lastDelta, - // Chunks: s.buf.Chunks(), - // } +// fix - повернути в Pool буфери +func (s *ReverseInstantDeltaCompressor) Unlock() { + s.state = nil } -func (s *ReverseInstantDeltaCompressor) Unlock() { - //s.state = nil +func (s *ReverseInstantDeltaCompressor) Offset() int { + if s.state != nil { + return s.state.Pos + } + return 0 +} + +func (s *ReverseInstantDeltaCompressor) Snapshot() ([][]byte, int) { + if s.state == nil { + return s.buf.Chunks(), s.Size() + } + // ВАЖЛИВО! + // Треба відтворити стан останнього чанка + var ( + pos = s.state.Pos + chunk = make([]byte, conbuf.ChunkSize) + lastChunkIdx = len(s.state.Chunks) - 1 + qtyToCopy = pos % conbuf.ChunkSize + ) + copy(chunk, s.state.Chunks[lastChunkIdx][:qtyToCopy]) + chunks := append(s.state.Chunks[:lastChunkIdx], chunk) + + buf := conbuf.New(chunks) + pos += buf.ReversePutVarInt64(pos, s.state.LastDelta) + buf.SetByte(pos, s.state.H) + pos++ + return chunks, pos +} + +func (s *ReverseInstantDeltaCompressor) CreateDecompressor(fracDigits byte) qb.ValueDecompressor { + if s.state == nil { + d := NewReverseInstantDeltaDecompressor(s.buf, s.Size(), fracDigits) + d.RestoreFromEnd() + return d + } + d := NewReverseInstantDeltaDecompressor(s.buf, s.Size(), fracDigits) + d.RestoreFromBound(*s.state) + return d } func (s *ReverseInstantDeltaCompressor) Renew() { + // УВАГА! + // state не чіпаємо s.buf = conbuf.New(nil) s.pos = 0 // diff --git a/chunkenc/time_delta.go b/chunkenc/time_delta.go index 6bd8b49..05d2914 100644 --- a/chunkenc/time_delta.go +++ b/chunkenc/time_delta.go @@ -5,6 +5,7 @@ import ( "log" "gordenko.dev/dima/pretty" + "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" ) @@ -15,6 +16,7 @@ type ReverseTimeDeltaCompressor struct { lastDelta uint32 lastDeltaSize int h byte + state *TimeDeltaBound } func NewReverseTimeDeltaCompressor(buf *conbuf.ContinuousBuffer, size int) *ReverseTimeDeltaCompressor { @@ -45,10 +47,6 @@ func (s *ReverseTimeDeltaCompressor) Size() int { return s.pos } -func (s *ReverseTimeDeltaCompressor) Chunks() [][]byte { - return s.buf.Chunks() -} - func (s *ReverseTimeDeltaCompressor) CalcRequiredSpace(value uint32) int { return 0 } @@ -151,42 +149,109 @@ func (s *ReverseTimeDeltaCompressor) Sync() { s.pos += s.buf.ReversePutVarUint64(s.pos, uint64(s.lastUnixtime)) } +// FIX - check methods + type TimeDeltaBound struct { Pos int H byte LastUnixtime uint32 LastDelta uint32 + Chunks [][]byte } // delta h -func (s *ReverseTimeDeltaCompressor) GetState() TimeDeltaBound { // fix replace by Lock - bound := TimeDeltaBound{ - LastUnixtime: s.lastUnixtime, - } - if s.pos > 0 { - bound.Pos = s.pos - 1 - s.lastDeltaSize - bound.H = s.h - bound.LastDelta = s.lastDelta - } - return bound -} +// func (s *ReverseTimeDeltaCompressor) GetState() TimeDeltaBound { // fix replace by Lock +// bound := TimeDeltaBound{ +// LastUnixtime: s.lastUnixtime, +// } +// if s.pos > 0 { +// bound.Pos = s.pos - 1 - s.lastDeltaSize +// bound.H = s.h +// bound.LastDelta = s.lastDelta +// } +// return bound +// } +// delta h func (s *ReverseTimeDeltaCompressor) Lock() { - // fix + if s.state != nil { + qb.Abort(qb.RepeatableLock, nil) + } + // позиція посувається вліво, отже може перескочити на попередній chunk + pos := s.pos - 1 - s.lastDeltaSize + chunksQty := pos / conbuf.ChunkSize + if (pos % conbuf.ChunkSize) > 0 { + chunksQty++ + } + s.state = &TimeDeltaBound{ + Pos: pos, + H: s.h, + LastUnixtime: s.lastUnixtime, // fix + LastDelta: s.lastDelta, + Chunks: s.buf.Chunks()[:chunksQty], + } } +// fix - повернути в Pool буфери func (s *ReverseTimeDeltaCompressor) Unlock() { - //s.state = nil // fix + s.state = nil +} + +func (s *ReverseTimeDeltaCompressor) Offset() int { + if s.state != nil { + return s.state.Pos + } + return 0 +} + +// fix - +func (s *ReverseTimeDeltaCompressor) Snapshot() ([][]byte, int) { + if s.state == nil { + return s.buf.Chunks(), s.Size() + } + // ВАЖЛИВО! + // Треба відтворити стан останнього чанка + var ( + pos = s.state.Pos + chunk = make([]byte, conbuf.ChunkSize) + lastChunkIdx = len(s.state.Chunks) - 1 + qtyToCopy = pos % conbuf.ChunkSize + ) + copy(chunk, s.state.Chunks[lastChunkIdx][:qtyToCopy]) + chunks := append(s.state.Chunks[:lastChunkIdx], chunk) + + buf := conbuf.New(chunks) + pos += buf.ReversePutVarUint64(pos, uint64(s.state.LastDelta)) + buf.SetByte(pos, s.state.H) + pos++ + return chunks, pos +} + +func (s *ReverseTimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompressor { + if s.state == nil { + d := NewReverseTimeDeltaDecompressor(s.buf, s.Size()) + d.RestoreFromEnd() + return d + } + d := NewReverseTimeDeltaDecompressor(s.buf, s.Size()) + d.RestoreFromBound(*s.state) + return d } func (s *ReverseTimeDeltaCompressor) Renew() { - // s.buf = conbuf.New(nil) - // s.pos = 0 - // // - // s.baseValue = 0 - // s.lastDelta = 0 - // s.lastDeltaSize = 0 - // s.h = 0 + // УВАГА! + // state не чіпаємо + s.buf = conbuf.New(nil) + s.pos = 0 + // + //s.baseValue = 0 + s.lastDelta = 0 + s.lastDeltaSize = 0 + s.h = 0 +} + +func (s *ReverseTimeDeltaCompressor) Chunks() [][]byte { + return s.buf.Chunks() } // DECOMPRESSOR diff --git a/database/database.go b/database/database.go index 77c2702..96d37ff 100644 --- a/database/database.go +++ b/database/database.go @@ -12,11 +12,9 @@ import ( "sync" "time" - "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/freelist" - "gordenko.dev/dima/qb/recovery" "gordenko.dev/dima/qb/txlog" ) @@ -24,11 +22,11 @@ func JoinSnapshotFileName(dir string, logNumber int) string { return filepath.Join(dir, fmt.Sprintf("%d.snapshot", logNumber)) } -type metricLockEntry struct { - XLock bool - RLocks int - WaitQueue []any -} +// type metricLockEntry struct { +// XLock bool +// RLocks int +// WaitQueue []any +// } type Database struct { mutex sync.Mutex @@ -37,16 +35,17 @@ type Database struct { rLocksToRelease []uint32 metrics map[uint32]*_metric //metricLockEntries map[uint32]*metricLockEntry - freeList *freelist.FreeList - dir string - databaseName string - txlog *txlog.Writer - atree *atree.Atree - tcpPort int - logfile *os.File - logger *log.Logger - exitCh chan struct{} - waitGroup *sync.WaitGroup + dataFreeList *freelist.FreeList + indexFreeList *freelist.FreeList + dir string + databaseName string + txlog *txlog.Writer + atree *atree.Atree + tcpPort int + logfile *os.File + logger *log.Logger + exitCh chan struct{} + waitGroup *sync.WaitGroup } type Options struct { @@ -82,16 +81,25 @@ func New(opt Options) (_ *Database, err error) { return nil, errors.New("WaitGroup option is required") } - filePath := filepath.Join(opt.Dir, opt.DatabaseName+".free") + // deltaFreeListFile, err := os.OpenFile( + // filepath.Join(opt.Dir, opt.DatabaseName+".deltafree"), + // os.O_RDWR|os.O_CREATE, + // 0666, + // ) + // if err != nil { + // return nil, err + // } - freeListFile, err := os.OpenFile(filePath, os.O_RDWR|os.O_CREATE, 0666) - if err != nil { - return nil, err - } + dataFreeList, err := freelist.New(freelist.Options{ + PageSize: 2048, + BaseFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data_base"), + DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_data_delta"), + }) - freeList, err := freelist.New(freelist.Options{ - PageSize: 2048, - File: freeListFile, + indexFreeList, err := freelist.New(freelist.Options{ + PageSize: 2048, + BaseFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_index_base"), + DeltaFilePath: filepath.Join(opt.Dir, opt.DatabaseName+".free_index_delta"), }) s := &Database{ @@ -100,12 +108,13 @@ func New(opt Options) (_ *Database, err error) { databaseName: opt.DatabaseName, metrics: make(map[uint32]*_metric), //metricLockEntries: make(map[uint32]*metricLockEntry), - freeList: freeList, - tcpPort: opt.TCPPort, - logfile: opt.Logfile, - logger: log.New(opt.Logfile, "", log.LstdFlags), - exitCh: opt.ExitCh, - waitGroup: opt.WaitGroup, + dataFreeList: dataFreeList, + indexFreeList: indexFreeList, + tcpPort: opt.TCPPort, + logfile: opt.Logfile, + logger: log.New(opt.Logfile, "", log.LstdFlags), + exitCh: opt.ExitCh, + waitGroup: opt.WaitGroup, } return s, nil } @@ -143,88 +152,89 @@ func (s *Database) ListenAndServe() (err error) { } func (s *Database) recovery() { - advisor, err := recovery.NewRecoveryAdvisor(recovery.RecoveryAdvisorOptions{ - Dir: s.dir, - VerifySnapshot: s.verifySnapshot, - }) - if err != nil { - panic(err) - } - - recipe, err := advisor.GetRecipe() - if err != nil { - qb.Abort(qb.GetRecoveryRecipeFailed, err) - } - - var logNumber int - - if recipe != nil { - if recipe.Snapshot != "" { - err = s.loadSnapshot(recipe.Snapshot) - if err != nil { - qb.Abort(qb.LoadSnapshotFailed, err) - } - } - for _, changesFileName := range recipe.Changes { - err = s.replayChanges(changesFileName) - if err != nil { - qb.Abort(qb.ReplayChangesFailed, err) - } - } - logNumber = recipe.LogNumber - } - - s.txlog, err = txlog.NewWriter(txlog.WriterOptions{ - Dir: s.dir, - LogNumber: logNumber, - AppendToWorkerQueue: s.appendJobToWorkerQueue, - FreeList: s.freeList, - Atree: s.atree, - ExitCh: s.exitCh, - WaitGroup: s.waitGroup, - }) - if err != nil { - qb.Abort(qb.CreateChangesWriterFailed, err) - - } - go s.txlog.Run() - - // fileNames, err := s.searchREDOFiles() + // FIX + // advisor, err := recovery.NewRecoveryAdvisor(recovery.RecoveryAdvisorOptions{ + // Dir: s.dir, + // VerifySnapshot: s.verifySnapshot, + // }) // if err != nil { - // qb.Abort(qb.SearchREDOFilesFailed, err) + // panic(err) // } - // if len(fileNames) > 0 { - // for _, fileName := range fileNames { - // err = s.replayREDOFile(fileName) + // recipe, err := advisor.GetRecipe() + // if err != nil { + // qb.Abort(qb.GetRecoveryRecipeFailed, err) + // } + + // var logNumber int + + // if recipe != nil { + // if recipe.Snapshot != "" { + // err = s.loadSnapshot(recipe.Snapshot) // if err != nil { - // qb.Abort(qb.ReplayREDOFileFailed, err) + // qb.Abort(qb.LoadSnapshotFailed, err) + // } + // } + // for _, changesFileName := range recipe.Changes { + // err = s.replayChanges(changesFileName) + // if err != nil { + // qb.Abort(qb.ReplayChangesFailed, err) + // } + // } + // logNumber = recipe.LogNumber + // } + + // s.txlog, err = txlog.NewWriter(txlog.WriterOptions{ + // Dir: s.dir, + // LogNumber: logNumber, + // AppendToWorkerQueue: s.appendJobToWorkerQueue, + // FreeList: s.freeList, + // Atree: s.atree, + // ExitCh: s.exitCh, + // WaitGroup: s.waitGroup, + // }) + // if err != nil { + // qb.Abort(qb.CreateChangesWriterFailed, err) + + // } + // go s.txlog.Run() + + // // fileNames, err := s.searchREDOFiles() + // // if err != nil { + // // qb.Abort(qb.SearchREDOFilesFailed, err) + // // } + + // // if len(fileNames) > 0 { + // // for _, fileName := range fileNames { + // // err = s.replayREDOFile(fileName) + // // if err != nil { + // // qb.Abort(qb.ReplayREDOFileFailed, err) + // // } + // // } + + // // for _, fileName := range fileNames { + // // err = os.Remove(fileName) + // // if err != nil { + // // qb.Abort(qb.RemoveREDOFileFailed, err) + // // } + // // } + // // } + + // if recipe != nil { + // if recipe.CompleteSnapshot { + // err = s.dumpSnapshot(logNumber) + // if err != nil { + // qb.Abort(qb.DumpSnapshotFailed, err) // } // } - // for _, fileName := range fileNames { + // for _, fileName := range recipe.ToDelete { // err = os.Remove(fileName) // if err != nil { - // qb.Abort(qb.RemoveREDOFileFailed, err) + // qb.Abort(qb.RemoveRecipeFileFailed, err) // } // } // } - - if recipe != nil { - if recipe.CompleteSnapshot { - err = s.dumpSnapshot(logNumber) - if err != nil { - qb.Abort(qb.DumpSnapshotFailed, err) - } - } - - for _, fileName := range recipe.ToDelete { - err = os.Remove(fileName) - if err != nil { - qb.Abort(qb.RemoveRecipeFileFailed, err) - } - } - } } // func (s *Database) searchREDOFiles() ([]string, error) { diff --git a/database/metric.go b/database/metric.go index b0a6337..8230763 100644 --- a/database/metric.go +++ b/database/metric.go @@ -2,10 +2,18 @@ package database import ( "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/bin" + "gordenko.dev/dima/qb/txlog" ) // METRIC +type IndexRec struct { + Since uint32 + PageNo uint32 +} + type _metric struct { MetricType qb.MetricType FracDigits byte @@ -16,6 +24,12 @@ type _metric struct { Until uint32 Timestamps qb.TimestampCompressor Values qb.ValueCompressor + // fix - add index levels + XLock bool + RLocks int + WaitQueue []any + //IndexLevels [][]IndexRec // root - last element + IndexLevels [][]byte // root - last element } func (s *_metric) ReinitBy(timestamp uint32, value float64) { @@ -41,3 +55,261 @@ func (s *_metric) DeleteMeasures() { s.Until = 0 s.UntilValue = 0 } + +// func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(txlog.AppendedMeasures)) { +func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(any)) { + var ( + // AppendedMeasures struct { + // MetricID uint32 + // TimestampsOffset int // (заповнені одразу) + // ValuesOffset int // (заповнені одразу) + // Timestamps [][]byte + // TimestampsSize int + // Values [][]byte + // ValuesSize int + // IndexLevels []*atree.IndexLevel + // DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page + // } + + timestamps = s.Timestamps + values = s.Values + + indexLevels []*atree.IndexLevel + dataPages []*atree.DataPage + + //written int + //resultCode byte + ) + + s.Values.Lock() + s.Timestamps.Lock() + + for idx, measure := range req.Measures { + if s.Since == 0 { + s.Since = measure.Timestamp + } else { + // FIX - у випадку помилки треба в транзакції зберегти що помилка, але також зафіксувати скільки елементів збережено. + // якщо idx == 0 - одразу знімаю блокування і нічого не відправляю в txlog + if measure.Timestamp <= s.Until { + if idx == 0 { + s.Values.Unlock() + s.Timestamps.Unlock() + + req.ResultCh <- tryAppendMeasuresResult{ + ResultCode: ExpiredMeasure, + } + return + } + + //resultCode = ExpiredMeasure + //written = idx + break + } + + if s.MetricType == qb.Cumulative && measure.Value < s.UntilValue { + if idx == 0 { + s.Values.Unlock() + s.Timestamps.Unlock() + + req.ResultCh <- tryAppendMeasuresResult{ + ResultCode: NonMonotonicValue, + } + return + } + //resultCode = NonMonotonicValue + //written = idx + break + } + } + + // fix - 1 + 8 bytes + extraSpace := timestamps.CalcRequiredSpace(measure.Timestamp) + + values.CalcRequiredSpace(measure.Value) + + totalSpace := timestamps.Size() + values.Size() + extraSpace + + if totalSpace <= atree.DataPagePayloadSize { + // накопичую + timestamps.Append(measure.Timestamp) + values.Append(measure.Value) + } else { + // сторінка заповнена + // prevPageNo - виставляю в txlog, коли забираю номер сторінки із freeList або генерую новий + dataPages = append(dataPages, &atree.DataPage{ + PrevPageNo: 0, // ? + LowerTimestamp: s.Since, + Timestamps: timestamps.Chunks(), + TimestampsSize: timestamps.Size(), + Values: values.Chunks(), + ValuesSize: values.Size(), + }) + + // обнуляються буфери, але state незмінний + timestamps.Renew() + values.Renew() + + timestamps.Append(measure.Timestamp) + values.Append(measure.Value) + + s.Since = measure.Timestamp + } + + s.Until = measure.Timestamp + s.UntilValue = measure.Value + } + + // виділити змінені байти. + // скопіювати. Причому можна скопіювати зрізи chunks + + if len(dataPages) > 0 { + // пишу в txlog довгим шляхом через redo файл і запис в data файл + for _, unfilled := range s.IndexLevels { + indexLevels = append(indexLevels, &atree.IndexLevel{ + Offset: len(unfilled), + LowerTimestamp: bin.GetUint32(unfilled), + Data: unfilled, + }) + } + } else { + // короткий шлях - запис лише в txlog + } + + sendToStorage(txlog.AppendedMeasures{ + MetricID: req.MetricID, + TimestampsOffset: timestamps.Offset(), // state.Pos() з якої позиції дописувати дані на сторінку 0 (при відновленні) + ValuesOffset: values.Offset(), // state.Pos() + Timestamps: timestamps.Chunks(), + TimestampsSize: timestamps.Size(), + Values: values.Chunks(), + ValuesSize: values.Size(), + IndexLevels: indexLevels, + DataPages: dataPages, + }) +} + +func (s *_metric) FinAppendMeasures(rec txlog.AppendMeasuresSummary) { + // Видаляю state. Оригінальні Timestamps і Values вже мають останню версію + s.Values.Unlock() + s.Timestamps.Unlock() + // fix write index levels + // update prev pageNo + // if len(rec.DataPages) > 0 { + // s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo + // } + if rec.LastPageNo > 0 { + s.LastPageNo = rec.LastPageNo + } + // В txlog я передав повний індекс. У нього додали елементи (можливо нові рівні). + // Тому проста заміна + s.IndexLevels = rec.Index +} + +// READ + +func (s *_metric) StartRangeScan(req tryRangeScanReq) { + if s.Since == 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + + if req.Since > s.Until { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + + if req.Until < s.Since { + if s.RootPageNo > 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: UntilNotFound, + RootPageNo: s.RootPageNo, + FracDigits: s.FracDigits, + } + s.RLocks++ + return + } else { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + } + + timestampDecompressor := s.Timestamps.CreateDecompressor() + valueDecompressor := s.Values.CreateDecompressor(s.FracDigits) + + for { + timestamp, done := timestampDecompressor.NextValue() + if done { + break + } + + value, done := valueDecompressor.NextValue() + if done { + qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + } + + if timestamp <= req.Until { + req.ResponseWriter.FeedNoSend(timestamp, value) + if timestamp < req.Since { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + } + } + + if s.LastPageNo > 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: UntilFound, + LastPageNo: s.LastPageNo, + FracDigits: s.FracDigits, + } + s.RLocks++ + } else { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + } +} + +func (s *_metric) StartFullScan(req tryFullScanReq) { + if s.Since == 0 { + req.ResultCh <- fullScanResult{ + ResultCode: QueryDone, + } + return + } + + timestampDecompressor := s.Timestamps.CreateDecompressor() + valueDecompressor := s.Values.CreateDecompressor(s.FracDigits) + + for { + timestamp, done := timestampDecompressor.NextValue() + if done { + break + } + value, done := valueDecompressor.NextValue() + if done { + qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + } + req.ResponseWriter.FeedNoSend(timestamp, value) + } + + if s.LastPageNo > 0 { + req.ResultCh <- fullScanResult{ + ResultCode: UntilFound, + LastPageNo: s.LastPageNo, + FracDigits: s.FracDigits, + } + s.RLocks++ + } else { + req.ResultCh <- fullScanResult{ + ResultCode: QueryDone, + } + } +} diff --git a/database/proc.go b/database/proc.go index 3a58b52..df50a8d 100644 --- a/database/proc.go +++ b/database/proc.go @@ -12,6 +12,11 @@ import ( "gordenko.dev/dima/qb/txlog" ) +// AddMetric - create metric in map and set Xlock, rwad - append to queue +// DeleteMetric - set Xlock +// DeleteSince - xLock (тупо редкая операция) +// Якийсь прапор поставити для pendingDelete, щоб read операції вставали в чергу + const ( QueryDone = 1 UntilFound = 2 @@ -49,46 +54,33 @@ func (s *Database) DoWork() { s.mutex.Unlock() for _, metricID := range rLocksToRelease { - lockEntry, ok := s.metricLockEntries[metricID] + metric, ok := s.metrics[metricID] if !ok { - qb.Abort(qb.NoLockEntryBug, - fmt.Errorf("drainQueues: lockEntry not found for the metric %d", - metricID)) + qb.Abort(qb.NoMetricBug, + fmt.Errorf("drainQueues: metric %d not found", metricID)) } - if lockEntry.XLock { + if metric.XLock { qb.Abort(qb.XLockBug, fmt.Errorf("drainQueues: xlock is set for the metric %d", metricID)) } - if lockEntry.RLocks <= 0 { + if metric.RLocks <= 0 { qb.Abort(qb.NoRLockBug, fmt.Errorf("drainQueues: rlock not set for the metric %d", metricID)) } - lockEntry.RLocks-- + metric.RLocks-- - if len(lockEntry.WaitQueue) > 0 { - metric, ok := s.metrics[metricID] - if !ok { - qb.Abort(qb.NoMetricBug, - fmt.Errorf("drainQueues: metric %d not found", metricID)) - } - s.processMetricQueue(metricID, metric, lockEntry) - } else { - if lockEntry.RLocks == 0 { - delete(s.metricLockEntries, metricID) - } + if len(metric.WaitQueue) > 0 { + s.processMetricQueue(metricID, metric) } } for _, untyped := range workerQueue { switch req := untyped.(type) { - //case tryAppendMeasureReq: - // s.tryAppendMeasure(req) - case tryAppendMeasuresReq: s.tryAppendMeasures(req) @@ -123,64 +115,35 @@ func (s *Database) DoWork() { } } -func (s *Database) processMetricQueue(metricID uint32, metric *_metric, lockEntry *metricLockEntry) { - if len(lockEntry.WaitQueue) == 0 { +// суть у тому що треба запускати запити, пока не зустріну XLock +func (s *Database) processMetricQueue(metricID uint32, metric *_metric) { + if len(metric.WaitQueue) == 0 { return } - var modificationReqs []any - - for _, untyped := range lockEntry.WaitQueue { - var rLockRequired bool + for _, untyped := range metric.WaitQueue { switch req := untyped.(type) { case tryRangeScanReq: - rLockRequired = s.startRangeScan(metric, req) + metric.StartRangeScan(req) case tryFullScanReq: - rLockRequired = s.startFullScan(metric, req) + metric.StartFullScan(req) case tryGetMetricReq: s.tryGetMetric(req) + case tryAppendMeasuresReq: + metric.StartAppendMeasures(req, s.txlog.Append) + + case tryDeleteMetricReq: + s.startDeleteMetric(metric, req) + + case tryDeleteMeasuresReq: + s.startDeleteMeasures(metric, req) + default: - modificationReqs = append(modificationReqs, untyped) - } - - if rLockRequired { - lockEntry.RLocks++ - } - } - lockEntry.WaitQueue = nil - if lockEntry.RLocks > 0 { - lockEntry.WaitQueue = modificationReqs - } else { - for idx, untyped := range modificationReqs { - switch req := untyped.(type) { - //case tryAppendMeasureReq: - // s.startAppendMeasure(metric, req, nil) - - case tryAppendMeasuresReq: - s.startAppendMeasures(metric, req, nil) - - case tryDeleteMetricReq: - s.startDeleteMetric(metric, req) - - case tryDeleteMeasuresReq: - s.startDeleteMeasures(metric, req) - - default: - qb.Abort(qb.UnknownMetricWaitQueueItemBug, - fmt.Errorf("bug: unknown metric wait queue item type %T", req)) - } - - lockEntry, ok := s.metricLockEntries[metricID] - if ok { - start := idx + 1 - if start < len(modificationReqs) { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, modificationReqs[start:]...) - } - break - } + qb.Abort(qb.UnknownMetricWaitQueueItemBug, + fmt.Errorf("bug: unknown metric wait queue item type %T", req)) } } } @@ -332,222 +295,16 @@ func (s *Database) startDeleteMeasures(metric *_metric, req tryDeleteMeasuresReq // } } -// type tryAppendMeasureReq struct { -// MetricID uint32 -// Timestamp uint32 -// Value float64 -// ResultCh chan tryAppendMeasureResult -// } - -// func (s *Database) tryAppendMeasure(req tryAppendMeasureReq) { -// metric, ok := s.metrics[req.MetricID] -// if !ok { -// req.ResultCh <- tryAppendMeasureResult{ -// MetricID: req.MetricID, -// ResultCode: NoMetric, -// } -// return -// } - -// lockEntry, ok := s.metricLockEntries[req.MetricID] -// if ok { -// if lockEntry.XLock { -// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) -// return -// } -// } -// s.startAppendMeasure(metric, req, lockEntry) -// } - -type ReadBound struct { - ValuesPos int - ValuesLastValue float64 -} - -// func (s *Database) startAppendMeasure(metric *_metric, req tryAppendMeasureReq, lockEntry *metricLockEntry) { -// if req.Timestamp <= metric.Until { -// req.ResultCh <- tryAppendMeasureResult{ -// MetricID: req.MetricID, -// ResultCode: ExpiredMeasure, -// } -// return -// } - -// if metric.MetricType == qb.Cumulative && req.Value < metric.UntilValue { -// req.ResultCh <- tryAppendMeasureResult{ -// MetricID: req.MetricID, -// ResultCode: NonMonotonicValue, -// } -// return -// } - -// extraSpace := metric.Timestamps.CalcRequiredSpace(req.Timestamp) + -// metric.Values.CalcRequiredSpace(req.Value) - -// totalSpace := metric.Timestamps.Size() + metric.Values.Size() + extraSpace - -// if totalSpace <= atree.DataPagePayloadSize { -// if lockEntry != nil { -// lockEntry.RLocks++ -// } else { -// s.metricLockEntries[req.MetricID] = &metricLockEntry{ -// RLocks: 1, -// } -// } -// req.ResultCh <- tryAppendMeasureResult{ -// MetricID: req.MetricID, -// Timestamp: req.Timestamp, -// Value: req.Value, -// ResultCode: CanAppend, -// } -// } else { -// if lockEntry != nil { -// if lockEntry.RLocks > 0 { -// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) -// return -// } -// lockEntry.XLock = true -// } else { -// s.metricLockEntries[req.MetricID] = &metricLockEntry{ -// XLock: true, -// } -// } - -// req.ResultCh <- tryAppendMeasureResult{ -// MetricID: req.MetricID, -// Timestamp: req.Timestamp, -// Value: req.Value, -// ResultCode: NewPage, -// FilledPage: &FilledPage{ -// Since: metric.Since, -// RootPageNo: metric.RootPageNo, -// PrevPageNo: metric.LastPageNo, -// TimestampsChunks: metric.TimestampsBuf.Chunks(), -// TimestampsSize: uint16(metric.Timestamps.Size()), -// ValuesChunks: metric.ValuesBuf.Chunks(), -// ValuesSize: uint16(metric.Values.Size()), -// }, -// } -// } -// } - -// func (s *Database) appendMeasure(rec txlog.AppendedMeasure) { -// metric, ok := s.metrics[rec.MetricID] -// if !ok { -// qb.Abort(qb.NoMetricBug, -// fmt.Errorf("appendMeasure: metric %d not found", -// rec.MetricID)) -// } - -// lockEntry, ok := s.metricLockEntries[rec.MetricID] -// if !ok { -// qb.Abort(qb.NoLockEntryBug, -// fmt.Errorf("appendMeasure: lockEntry not found for the metric %d", -// rec.MetricID)) -// } - -// if lockEntry.XLock { -// qb.Abort(qb.XLockBug, -// fmt.Errorf("appendMeasure: xlock is set for the metric %d", -// rec.MetricID)) -// } - -// if lockEntry.RLocks <= 0 { -// qb.Abort(qb.NoRLockBug, -// fmt.Errorf("appendMeasure: rlock not set for the metric %d", -// rec.MetricID)) -// } - -// if metric.Since == 0 { -// metric.Since = rec.Timestamp -// metric.SinceValue = rec.Value -// } - -// metric.Timestamps.Append(rec.Timestamp) -// metric.Values.Append(rec.Value) -// metric.Until = rec.Timestamp -// metric.UntilValue = rec.Value - -// lockEntry.RLocks-- -// if len(lockEntry.WaitQueue) > 0 { -// s.processMetricQueue(rec.MetricID, metric, lockEntry) -// } else { -// if lockEntry.RLocks == 0 { -// delete(s.metricLockEntries, rec.MetricID) -// } -// } -// } - -func (s *Database) appendMeasures(rec txlog.AppendedMeasures) { +func (s *Database) finAppendMeasures(rec txlog.AppendMeasuresSummary) { metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, - fmt.Errorf("appendMeasureAfterOverflow: metric %d not found", + fmt.Errorf("finAppendMeasures: metric %d not found", rec.MetricID)) } - - metric.Values.Unlock() - metric.Timestamps.Unlock() - - // lockEntry, ok := s.metricLockEntries[rec.MetricID] - // if !ok { - // qb.Abort(qb.NoLockEntryBug, - // fmt.Errorf("appendMeasureAfterOverflow: lockEntry not found for the metric %d", - // rec.MetricID)) - // } - - // if !lockEntry.XLock { - // qb.Abort(qb.NoXLockBug, - // fmt.Errorf("appendMeasureAfterOverflow: xlock not set for the metric %d", - // rec.MetricID)) - // } - - // lockEntry.XLock = false - // s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) + metric.FinAppendMeasures(rec) } -// func (s *Database) appendMeasureAfterOverflow(extended txlog.AppendedMeasureWithOverflowExtended) { -// rec := extended.Record -// metric, ok := s.metrics[rec.MetricID] -// if !ok { -// qb.Abort(qb.NoMetricBug, -// fmt.Errorf("appendMeasureAfterOverflow: metric %d not found", -// rec.MetricID)) -// } - -// lockEntry, ok := s.metricLockEntries[rec.MetricID] -// if !ok { -// qb.Abort(qb.NoLockEntryBug, -// fmt.Errorf("appendMeasureAfterOverflow: lockEntry not found for the metric %d", -// rec.MetricID)) -// } - -// if !lockEntry.XLock { -// qb.Abort(qb.NoXLockBug, -// fmt.Errorf("appendMeasureAfterOverflow: xlock not set for the metric %d", -// rec.MetricID)) -// } - -// metric.ReinitBy(rec.Timestamp, rec.Value) -// if rec.IsRootChanged { -// metric.RootPageNo = rec.RootPageNo -// } -// metric.LastPageNo = rec.DataPageNo - -// if rec.IsDataPageReused { -// s.freeList.DeleteReservedPages([]uint32{ -// rec.DataPageNo, -// }) -// } - -// if len(rec.ReusedIndexPages) > 0 { -// s.freeList.DeleteReservedPages(rec.ReusedIndexPages) -// } - -// lockEntry.XLock = false -// s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) -// } - type tryAppendMeasuresReq struct { MetricID uint32 Measures []proto.Measure @@ -562,175 +319,11 @@ func (s *Database) tryAppendMeasures(req tryAppendMeasuresReq) { } return } - - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - if lockEntry.XLock { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - return - } + if metric.XLock { + metric.WaitQueue = append(metric.WaitQueue, req) + return } - s.startAppendMeasures(metric, req, lockEntry) -} - -func (s *Database) startAppendMeasures(metric *_metric, req tryAppendMeasuresReq, lockEntry *metricLockEntry) { - var ( - timestamps = metric.Timestamps - values = metric.Values - - dataPages []atree.NotLinkedDataPage - - resultCode byte - written int - ) - - metric.Values.Lock() - metric.Timestamps.Lock() - - for idx, measure := range req.Measures { - if metric.Since == 0 { - metric.Since = measure.Timestamp - } else { - // FIX - у випадку помилки треба в транзакції зберегти що помилка, але також зафіксувати скільки елементів збережено. - // якщо idx == 0 - одразу знімаю блокування і нічого не відправляю в txlog - if measure.Timestamp <= metric.Until { - if idx == 0 { - metric.Values.Unlock() - metric.Timestamps.Unlock() - - req.ResultCh <- tryAppendMeasuresResult{ - ResultCode: ExpiredMeasure, - } - return - } - - resultCode = ExpiredMeasure - written = idx - break - } - - if metric.MetricType == qb.Cumulative && measure.Value < metric.UntilValue { - if idx == 0 { - metric.Values.Unlock() - metric.Timestamps.Unlock() - - req.ResultCh <- tryAppendMeasuresResult{ - ResultCode: NonMonotonicValue, - } - return - } - resultCode = NonMonotonicValue - written = idx - break - } - } - - // fix - 1 + 8 bytes - extraSpace := timestamps.CalcRequiredSpace(measure.Timestamp) + - values.CalcRequiredSpace(measure.Value) - - totalSpace := timestamps.Size() + values.Size() + extraSpace - - if totalSpace <= atree.DataPagePayloadSize { - // накопичую - timestamps.Append(measure.Timestamp) - values.Append(measure.Value) - } else { - // сторінка заповнена - pageData := make([]byte, atree.PageSize) - - // prevPageNo - виставляю в txlog, коли забираю номер сторінки із freeList або генерую новий - atree.ChunksToNotLinkedDataPage(pageData, atree.ChunksToNotLinkedDataPageReq{ - TimestampsChunks: timestamps.Chunks(), - TimestampsSize: uint16(timestamps.Size()), - ValuesChunks: values.Chunks(), - ValuesSize: uint16(values.Size()), - }) - - dataPages = append(dataPages, atree.NotLinkedDataPage{ - Since: metric.Since, - Data: pageData, - }) - - // 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 - - timestamps.Renew() - values.Renew() - - timestamps.Append(measure.Timestamp) - values.Append(measure.Value) - - metric.Since = measure.Timestamp - } - - metric.Until = measure.Timestamp - metric.UntilValue = measure.Value - } - - // виділити змінені байти. - // скопіювати. Причому можна скопіювати зрізи chunks - - if len(dataPages) > 0 { - // пишу в txlog довгим шляхом через redo файл і запис в data файл - } else { - // короткий шлях - запис лише в txlog - } - - // waitCh := s.txlog.Append(txlog.AppendedPagesReq{ - // MetricID: req.MetricID, - // Timestamp: until, - // Value: untilValue, - // LastPageNo: prevPageNo, - // //Measures: toAppendMeasures, - // }, - // ) - // <-waitCh - - // if lockEntry != nil { - // if lockEntry.RLocks > 0 { - // lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - // return - // } - // lockEntry.XLock = true - // } else { - // s.metricLockEntries[req.MetricID] = &metricLockEntry{ - // XLock: true, - // } - // } - -} - -type PartialChange struct { - Data [][]byte - Size int - Offset int -} - -type AppendedMeasures struct { - RootPageNo uint32 - LastPageNo uint32 - DataPages []atree.ChunksToNotLinkedDataPageReq - Values PartialChange - Timestamps PartialChange - ResultCh chan tryAppendMeasuresResult + metric.StartAppendMeasures(req, s.txlog.Append) } type tryRangeScanReq struct { @@ -757,102 +350,12 @@ func (s *Database) tryRangeScan(req tryRangeScanReq) { return } - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - if lockEntry.XLock { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - return - } + if metric.XLock { + metric.WaitQueue = append(metric.WaitQueue, req) + return } - if s.startRangeScan(metric, req) { - if lockEntry != nil { - lockEntry.RLocks++ - } else { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - RLocks: 1, - } - } - } -} - -func (*Database) startRangeScan(metric *_metric, req tryRangeScanReq) bool { - if metric.Since == 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return false - } - - if req.Since > metric.Until { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return false - } - - if req.Until < metric.Since { - if metric.RootPageNo > 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: UntilNotFound, - RootPageNo: metric.RootPageNo, - FracDigits: metric.FracDigits, - } - return true - } else { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return false - } - } - - timestampDecompressor := chunkenc.NewReverseTimeDeltaDecompressor( - metric.TimestampsBuf, - metric.Timestamps.Size(), - ) - - valueDecompressor := chunkenc.NewReverseInstantDeltaDecompressor( - metric.ValuesBuf, - metric.Values.Size(), - metric.FracDigits, - ) - - for { - timestamp, done := timestampDecompressor.NextValue() - if done { - break - } - - value, done := valueDecompressor.NextValue() - if done { - qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - } - - if timestamp <= req.Until { - req.ResponseWriter.FeedNoSend(timestamp, value) - if timestamp < req.Since { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return false - } - } - } - - if metric.LastPageNo > 0 { - req.ResultCh <- rangeScanResult{ - ResultCode: UntilFound, - LastPageNo: metric.LastPageNo, - FracDigits: metric.FracDigits, - } - return true - } else { - req.ResultCh <- rangeScanResult{ - ResultCode: QueryDone, - } - return false - } + metric.StartRangeScan(req) } type tryFullScanReq struct { @@ -877,68 +380,11 @@ func (s *Database) tryFullScan(req tryFullScanReq) { return } - lockEntry, ok := s.metricLockEntries[req.MetricID] - if ok { - if lockEntry.XLock { - lockEntry.WaitQueue = append(lockEntry.WaitQueue, req) - return - } - } - - if s.startFullScan(metric, req) { - if lockEntry != nil { - lockEntry.RLocks++ - } else { - s.metricLockEntries[req.MetricID] = &metricLockEntry{ - RLocks: 1, - } - } - } -} - -func (*Database) startFullScan(metric *_metric, req tryFullScanReq) bool { - if metric.Since == 0 { - req.ResultCh <- fullScanResult{ - ResultCode: QueryDone, - } - return false - } - - timestampDecompressor := chunkenc.NewReverseTimeDeltaDecompressor( - metric.TimestampsBuf, - metric.Timestamps.Size(), - ) - valueDecompressor := chunkenc.NewReverseInstantDeltaDecompressor( - metric.ValuesBuf, - metric.Values.Size(), - metric.FracDigits, - ) - - for { - timestamp, done := timestampDecompressor.NextValue() - if done { - break - } - value, done := valueDecompressor.NextValue() - if done { - qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - } - req.ResponseWriter.FeedNoSend(timestamp, value) - } - - if metric.LastPageNo > 0 { - req.ResultCh <- fullScanResult{ - ResultCode: UntilFound, - LastPageNo: metric.LastPageNo, - FracDigits: metric.FracDigits, - } - return true - } else { - req.ResultCh <- fullScanResult{ - ResultCode: QueryDone, - } - return false + if metric.XLock { + metric.WaitQueue = append(metric.WaitQueue, req) + return } + metric.StartFullScan(req) } type tryListCurrentValuesReq struct { @@ -975,8 +421,8 @@ func (s *Database) applyChanges(req txlog.Changes) { // case txlog.AppendedMeasure: // s.appendMeasure(rec) - case txlog.AppendedMeasures: - s.appendMeasures(rec) + case txlog.AppendMeasuresSummary: + s.finAppendMeasures(rec) // case txlog.AppendedMeasureWithOverflowExtended: // s.appendMeasureAfterOverflow(rec) @@ -986,18 +432,38 @@ func (s *Database) applyChanges(req txlog.Changes) { } } - if req.ForceSnapshot || req.ExitWaitGroup != nil { - s.dumpSnapshot(req.LogNumber) - } - - close(req.WaitCh) - - if req.ExitWaitGroup != nil { - req.ExitWaitGroup.Done() + if req.MetricsCh != nil { + waitCh := make(chan struct{}) + req.MetricsCh <- txlog.MetricsState{ + Metrics: s.createMetricsState(), + WaitCh: waitCh, + } + // чекаю поки txlog запише снапшот, а отже можна змінювати буфери + <-waitCh } } +func (s *Database) createMetricsState() (metrics []txlog.Metric) { + for metricID, metric := range s.metrics { + x := txlog.Metric{ + MetricID: metricID, + MetricType: metric.MetricType, + FracDigits: metric.FracDigits, + LastPageNo: metric.LastPageNo, + Since: metric.Since, + SinceValue: metric.SinceValue, + Until: metric.Until, // ? + UntilValue: metric.UntilValue, // ? + } + x.Timestamps, x.TimestampsSize = metric.Timestamps.Snapshot() + x.Values, x.ValuesSize = metric.Values.Snapshot() + metrics = append(metrics, x) + } + return +} + func (s *Database) addMetric(rec txlog.AddedMetric) { + // fix lock _, ok := s.metrics[rec.MetricID] if ok { qb.Abort(qb.MetricAddedBug, @@ -1005,18 +471,11 @@ func (s *Database) addMetric(rec txlog.AddedMetric) { rec.MetricID)) } - lockEntry, ok := s.metricLockEntries[rec.MetricID] - if !ok { - qb.Abort(qb.NoLockEntryBug, - fmt.Errorf("addMetric: lockEntry not found for the metric %d", - rec.MetricID)) - } - - if !lockEntry.XLock { - qb.Abort(qb.NoXLockBug, - fmt.Errorf("addMetric: xlock not set for the metric %d", - rec.MetricID)) - } + // if !lockEntry.XLock { + // qb.Abort(qb.NoXLockBug, + // fmt.Errorf("addMetric: xlock not set for the metric %d", + // rec.MetricID)) + // } var ( values qb.ValueCompressor @@ -1033,12 +492,10 @@ func (s *Database) addMetric(rec txlog.AddedMetric) { } s.metrics[rec.MetricID] = &_metric{ - MetricType: rec.MetricType, - FracDigits: byte(rec.FracDigits), - TimestampsBuf: timestampsBuf, - ValuesBuf: valuesBuf, - Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), - Values: values, + MetricType: rec.MetricType, + FracDigits: byte(rec.FracDigits), + Timestamps: chunkenc.NewReverseTimeDeltaCompressor(timestampsBuf, 0), + Values: values, } lockEntry.XLock = false @@ -1046,21 +503,14 @@ func (s *Database) addMetric(rec txlog.AddedMetric) { } func (s *Database) deleteMetric(rec txlog.DeletedMetric) { - _, ok := s.metrics[rec.MetricID] + metric, ok := s.metrics[rec.MetricID] if !ok { qb.Abort(qb.NoMetricBug, fmt.Errorf("deleteMetric: metric %d not found", rec.MetricID)) } - lockEntry, ok := s.metricLockEntries[rec.MetricID] - if !ok { - qb.Abort(qb.NoLockEntryBug, - fmt.Errorf("deleteMetric: lockEntry not found for the metric %d", - rec.MetricID)) - } - - if !lockEntry.XLock { + if !metric.XLock { qb.Abort(qb.NoXLockBug, fmt.Errorf("deleteMetric: xlock not set for the metric %d", rec.MetricID)) @@ -1068,8 +518,8 @@ func (s *Database) deleteMetric(rec txlog.DeletedMetric) { var addMetricReqs []tryAddMetricReq - if len(lockEntry.WaitQueue) > 0 { - for _, untyped := range lockEntry.WaitQueue { + if len(metric.WaitQueue) > 0 { + for _, untyped := range metric.WaitQueue { switch req := untyped.(type) { // case tryAppendMeasureReq: // req.ResultCh <- tryAppendMeasureResult{ @@ -1112,11 +562,11 @@ func (s *Database) deleteMetric(rec txlog.DeletedMetric) { } } delete(s.metrics, rec.MetricID) - delete(s.metricLockEntries, rec.MetricID) - if len(rec.FreePageNumbers) > 0 { - s.freeList.AddPages(rec.FreePageNumbers) - } + // ADD in txlog + // if len(rec.FreePageNumbers) > 0 { + // s.freeList.AddPages(rec.FreePageNumbers) + // } if len(addMetricReqs) > 0 { s.processTryAddMetricReqsImmediatelyAfterDelete(addMetricReqs) @@ -1131,30 +581,22 @@ func (s *Database) deleteMeasures(rec txlog.DeletedMeasures) { rec.MetricID)) } - lockEntry, ok := s.metricLockEntries[rec.MetricID] - if !ok { - qb.Abort(qb.NoLockEntryBug, - fmt.Errorf("deleteMeasures: lockEntry not found for the metric %d", - rec.MetricID)) - } - - if !lockEntry.XLock { + if !metric.XLock { qb.Abort(qb.NoXLockBug, fmt.Errorf("deleteMeasures: xlock not set for the metric %d", rec.MetricID)) } metric.DeleteMeasures() - lockEntry.XLock = false - if len(rec.FreePageNumbers) > 0 { - s.freeList.AddPages(rec.FreePageNumbers) - } - s.doAfterReleaseXLock(rec.MetricID, metric, lockEntry) + metric.XLock = false + // FIX add in txlog + // if len(rec.FreePageNumbers) > 0 { + // s.freeList.AddPages(rec.FreePageNumbers) + // } + s.doAfterReleaseXLock(rec.MetricID, metric) } -func (s *Database) doAfterReleaseXLock(metricID uint32, metric *_metric, lockEntry *metricLockEntry) { - if len(lockEntry.WaitQueue) == 0 { - delete(s.metricLockEntries, metricID) - } else { - s.processMetricQueue(metricID, metric, lockEntry) +func (s *Database) doAfterReleaseXLock(metricID uint32, metric *_metric) { + if len(metric.WaitQueue) > 0 { + s.processMetricQueue(metricID, metric) } } diff --git a/database/snapshot.go b/database/snapshot.go deleted file mode 100644 index e4a7df2..0000000 --- a/database/snapshot.go +++ /dev/null @@ -1,279 +0,0 @@ -package database - -import ( - "fmt" - "hash/crc32" - "io" - "os" - "path/filepath" - - bin "gordenko.dev/dima/bin/little" - "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" - "gordenko.dev/dima/qb/chunkenc" - "gordenko.dev/dima/qb/conbuf" - "gordenko.dev/dima/qb/freelist" -) - -/* -Формат: -metricsQty - varuint -[metric]* -де metric - це: -metricID - 4b -metricType - 1b -fracDigits - 1b -lastPageNo - 4b -since - 4b -sinceValue - 8b -until - 4b -untilValue - 8b -timestamps size - 2b -values size - 2b -timestams payload - Nb -values payload - Nb -dataFreeList size - varsize -dataFreeList - Nb -indexFreeList size - varsize -indexFreeList - Nb -CRC32 - 4b -*/ - -const metricHeaderSize = 38 - -func (s *Database) dumpSnapshot(logNumber int) (err error) { - var ( - fileName = filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", logNumber)) - hasher = crc32.NewIEEE() - prefix = make([]byte, metricHeaderSize) - ) - - file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0770) - if err != nil { - return - } - - dst := io.MultiWriter(file, hasher) - - _, err = bin.WriteVarUint64(dst, uint64(len(s.metrics))) - if err != nil { - return - } - - for metricID, metric := range s.metrics { - tSize := metric.Timestamps.Size() - vSize := metric.Values.Size() - - bin.PutUint32(prefix[0:], metricID) - prefix[4] = byte(metric.MetricType) - prefix[5] = metric.FracDigits - bin.PutUint32(prefix[6:], metric.LastPageNo) - bin.PutUint32(prefix[10:], metric.Since) - bin.PutFloat64(prefix[14:], metric.SinceValue) - bin.PutUint32(prefix[22:], metric.Until) - bin.PutFloat64(prefix[26:], metric.UntilValue) - bin.PutUint16(prefix[34:], uint16(tSize)) - bin.PutUint16(prefix[36:], uint16(vSize)) - - _, err = dst.Write(prefix) - if err != nil { - return - } - // copy timestamps - writeChunks(dst, metric.Timestamps.Chunks(), tSize) - // copy values - writeChunks(dst, metric.Values.Chunks(), vSize) - - } - // free data pages - err = freeListWriteTo(s.freeList, dst) - if err != nil { - return - } - // free index pages - err = freeListWriteTo(s.freeList, dst) - if err != nil { - return - } - - bin.WriteUint32(file, hasher.Sum32()) - - err = file.Sync() - if err != nil { - return - } - - err = file.Close() - if err != nil { - return - } - - // prevLogNumber := logNumber - 1 - // prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber)) - // prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber)) - - // isExist, err := isFileExist(prevChanges) - // if err != nil { - // return - // } - - // if isExist { - // err = os.Remove(prevChanges) - // if err != nil { - // qb.Abort(qb.DeletePrevChangesFileFailed, err) - // } - // } - - // isExist, err = isFileExist(prevSnapshot) - // if err != nil { - // return - // } - - // if isExist { - // err = os.Remove(prevSnapshot) - // if err != nil { - // qb.Abort(qb.DeletePrevSnapshotFileFailed, err) - // } - // } - return -} - -func (s *Database) loadSnapshot(fileName string) (err error) { - var ( - hasher = crc32.NewIEEE() - header = make([]byte, metricHeaderSize) - body = make([]byte, atree.PageSize) - ) - - file, err := os.Open(fileName) - if err != nil { - return - } - - src := io.TeeReader(file, hasher) - metricsQty, err := bin.ReadVarSize(src) - if err != nil { - return - } - - for range metricsQty { - var metric _metric - err = bin.ReadNInto(src, header) - if err != nil { - return - } - - metricID, _ := bin.GetUint32(header[0:]) - metric.MetricType = qb.MetricType(header[4]) - metric.FracDigits = header[5] - metric.LastPageNo, _ = bin.GetUint32(header[6:]) - metric.Since, _ = bin.GetUint32(header[10:]) - metric.SinceValue, _ = bin.GetFloat64(header[14:]) - metric.Until, _ = bin.GetUint32(header[22:]) - metric.UntilValue, _ = bin.GetFloat64(header[26:]) - tSize, _ := bin.GetUint16(header[34:]) - vSize, _ := bin.GetUint16(header[36:]) - - buf := body[:tSize] - err = bin.ReadNInto(src, buf) - if err != nil { - return - } - metric.Timestamps = chunkenc.NewReverseTimeDeltaCompressor( - conbuf.NewFromBuffer(buf), - int(tSize), - ) - - buf = body[:vSize] - err = bin.ReadNInto(src, buf) - if err != nil { - return - } - - if metric.MetricType == qb.Cumulative { - metric.Values = chunkenc.NewReverseCumulativeDeltaCompressor( - conbuf.NewFromBuffer(buf), - int(vSize), - metric.FracDigits, - ) - } else { - metric.Values = chunkenc.NewReverseInstantDeltaCompressor( - conbuf.NewFromBuffer(buf), - int(vSize), - metric.FracDigits, - ) - } - s.metrics[metricID] = &metric - } - - err = restoreFreeList(s.freeList, src) - if err != nil { - return fmt.Errorf("restore dataFreeList: %s", err) - } - - err = restoreFreeList(s.freeList, src) - if err != nil { - return fmt.Errorf("restore indexFreeList: %s", err) - } - - calculatedChecksum := hasher.Sum32() - - writtenChecksum, err := bin.ReadUint32(file) - if err != nil { - return - } - - if calculatedChecksum != writtenChecksum { - return fmt.Errorf("calculated checksum %d not equal written checksum %d", calculatedChecksum, writtenChecksum) - } - return -} - -// HELPERS - -func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error { - serialized, err := freeList.Serialize() - if err != nil { - qb.Abort(qb.FailedFreeListSerialize, err) - } - _, err = bin.WriteVarSize(dst, len(serialized)) - if err != nil { - return err - } - _, err = dst.Write(serialized) - if err != nil { - return err - } - return nil -} - -func restoreFreeList(freeList *freelist.FreeList, src io.Reader) error { - size, err := bin.ReadVarSize(src) - if err != nil { - return err - } - serialized, err := bin.ReadN(src, size) - if err != nil { - return err - } - freeList.Restore(serialized) - return nil -} - -func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) { - remaining := size - for _, buf := range chunks { - if remaining < len(buf) { - buf = buf[:remaining] - } - _, err = dst.Write(buf) - if err != nil { - return - } - remaining -= len(buf) - if remaining == 0 { - break - } - } - return -} diff --git a/freelist/freelist.go b/freelist/freelist.go index 88dc273..7342c6d 100644 --- a/freelist/freelist.go +++ b/freelist/freelist.go @@ -3,52 +3,82 @@ package freelist import ( "bytes" "fmt" + "io" "os" "gordenko.dev/dima/qb/bin" "gordenko.dev/dima/qb/util" ) +/* +Ідея - я фіксую кількість сторінок записаних в снапшоті. +Додаю і видаляю сторінки із основного файла, поки їх кількість > кількості зафіксованих. +Якщо pageNumbers видалено і кількість сторінок треба зменшити - нові сторінки пишу в delta файл. +При створенні снапшота - сторінки із delta файла записую в основний файл і фіксую в снапшоті +їх кількість. PageNumbers із буфера пишу в снапшот. +*/ + const ( crcSize = 4 ptrSize = 4 ) type FreeList struct { - file *os.File - pageSize int - pointersOnPage int - flushTreshold int // кількість вільних сторінок коли вже треба писати на диск - pages int - free []uint32 + baseFile *os.File + deltaFile *os.File + pageSize int + pointersOnPage int + flushTreshold int // кількість вільних сторінок коли вже треба писати на диск + frozenPageCount int + basePagesCount int + deltaPagesCount int + free []uint32 } type Options struct { - PageSize int - File *os.File + PageSize int + BaseFilePath string + DeltaFilePath string + FrozenPageCount int } func New(opt Options) (*FreeList, error) { s := &FreeList{ - file: opt.File, - pageSize: opt.PageSize, - pointersOnPage: (opt.PageSize - crcSize) / ptrSize, + pageSize: opt.PageSize, + pointersOnPage: (opt.PageSize - crcSize) / ptrSize, + frozenPageCount: opt.FrozenPageCount, + basePagesCount: opt.FrozenPageCount, } s.flushTreshold = int(float64(s.pointersOnPage) * 1.3) - info, err := s.file.Stat() + + var err error + s.baseFile, err = os.OpenFile(opt.BaseFilePath, os.O_RDWR|os.O_CREATE, 0666) if err != nil { return nil, err } - fileSize := info.Size() - if fileSize > 0 { - if (fileSize % int64(s.pageSize)) > 0 { - return nil, fmt.Errorf("file size %d is not a multiple of page size %d", fileSize, s.pageSize) - } - s.pages = int(fileSize) / s.pageSize + s.deltaFile, err = os.OpenFile(opt.DeltaFilePath, os.O_RDWR|os.O_CREATE, 0666) + if err != nil { + return nil, err } + // info, err := s.baseFile.Stat() + // if err != nil { + // return nil, err + // } + // fileSize := info.Size() + // if fileSize > 0 { + // if (fileSize % int64(s.pageSize)) > 0 { + // return nil, fmt.Errorf("file size %d is not a multiple of page size %d", + // fileSize, s.pageSize) + // } + // s.basePagesCount = int(fileSize) / s.pageSize + // } return s, nil } +func (s *FreeList) Pages() int { + return s.basePagesCount + s.deltaPagesCount +} + func (s *FreeList) AddPageNumbers(pageNumbers []uint32) (err error) { if len(pageNumbers) == 0 { return @@ -63,11 +93,13 @@ func (s *FreeList) AddPageNumbers(pageNumbers []uint32) (err error) { // викидаю номери сторінок, які записав на диск copy(s.free, s.free[s.pointersOnPage:]) s.free = s.free[:len(s.free)-s.pointersOnPage] + fmt.Println("free:", s.free) } return } func (s *FreeList) save() error { + fmt.Println("save:", s.free[:s.pointersOnPage]) buf := make([]byte, s.pageSize) i := crcSize for _, pageNo := range s.free[:s.pointersOnPage] { @@ -76,28 +108,49 @@ func (s *FreeList) save() error { } bin.PutUint32(buf, util.CalcChecksum(buf[crcSize:])) // - off := int64(s.pages * s.pageSize) - n, err := s.file.WriteAt(buf, off) + var ( + file *os.File + off int64 + ) + if s.deltaPagesCount > 0 || s.basePagesCount < s.frozenPageCount { + file = s.deltaFile + off = int64(s.deltaPagesCount * s.pageSize) + s.deltaPagesCount++ + } else { + file = s.baseFile + off = int64(s.basePagesCount * s.pageSize) + s.basePagesCount++ + } + fmt.Println("write at:", off) + n, err := file.WriteAt(buf, off) if err != nil { return err } if n != s.pageSize { return fmt.Errorf("written size %d bytes not equal page size %d", n, s.pageSize) } - err = s.file.Sync() + err = file.Sync() if err != nil { return err } - s.pages++ return nil } func (s *FreeList) loadLastPage() error { var ( - buf = make([]byte, s.pageSize) - off = int64((s.pages - 1) * s.pageSize) + file *os.File + off int64 + buf = make([]byte, s.pageSize) ) - n, err := s.file.ReadAt(buf, off) + if s.deltaPagesCount > 0 { + file = s.deltaFile + off = int64((s.deltaPagesCount - 1) * s.pageSize) + } else { + file = s.baseFile + off = int64((s.basePagesCount - 1) * s.pageSize) + } + + n, err := file.ReadAt(buf, off) if err != nil { return fmt.Errorf("file.Seek: %s", err) } @@ -107,32 +160,49 @@ func (s *FreeList) loadLastPage() error { // check crc checksum := bin.GetUint32(buf[0:]) if util.CalcChecksum(buf[crcSize:]) != checksum { - return fmt.Errorf("page %d is corrupted", s.pages) + return fmt.Errorf("page %d is corrupted", s.basePagesCount) } for i := crcSize; i < len(buf); i += ptrSize { s.free = append(s.free, bin.GetUint32(buf[i:])) } - s.pages-- - err = s.file.Truncate(int64(s.pages * s.pageSize)) - if err != nil { - return err + if s.deltaPagesCount > 0 { + s.deltaPagesCount-- + err = file.Truncate(off) + if err != nil { + return err + } + return file.Sync() + } else { + s.basePagesCount-- + if s.basePagesCount >= s.frozenPageCount { + err = s.baseFile.Truncate(off) + if err != nil { + return err + } + return file.Sync() + } } return nil } func (s *FreeList) GetPageNumber() (pageNo uint32, err error) { if len(s.free) == 0 { - if s.pages == 0 { - return - } - err = s.loadLastPage() - if err != nil { + if s.basePagesCount > 0 || s.deltaPagesCount > 0 { + // якщо немає вільних сторінок у буфері, але є у freelist файлах - + // завантажую сторінку із диску + err = s.loadLastPage() + if err != nil { + return + } + } else { return } } - lastIdx := len(s.free) - 1 - pageNo = s.free[lastIdx] - s.free = s.free[:lastIdx] + if len(s.free) > 0 { + lastIdx := len(s.free) - 1 + pageNo = s.free[lastIdx] + s.free = s.free[:lastIdx] + } return } @@ -144,6 +214,7 @@ func (s *FreeList) Restore(serialized []byte) error { for i := 0; i < len(serialized); i += ptrSize { s.free = append(s.free, bin.GetUint32(serialized[i:])) } + fmt.Println(s.free) return nil } @@ -155,3 +226,49 @@ func (s *FreeList) Serialize() ([]byte, error) { } return w.Bytes(), nil } + +func (s *FreeList) Merge() (err error) { + var ( + off = int64(s.basePagesCount * s.pageSize) + deltaSize = int64(s.deltaPagesCount * s.pageSize) + ) + if s.deltaPagesCount > 0 { + var ( + ret int64 + copied int64 + ) + ret, err = s.baseFile.Seek(off, 0) + if err != nil { + return + } + if ret != off { + return fmt.Errorf("Seek failed") + } + copied, err = io.Copy(s.baseFile, s.deltaFile) + if err != nil { + return + } + if copied != deltaSize { + return fmt.Errorf("copy failed") + } + err = s.baseFile.Sync() + if err != nil { + return + } + err = s.deltaFile.Truncate(0) + if err != nil { + return + } + s.basePagesCount += s.deltaPagesCount + s.deltaPagesCount = 0 + } + if s.basePagesCount < s.frozenPageCount { + off = int64(s.basePagesCount * s.pageSize) + err = s.baseFile.Truncate(off) + if err != nil { + return + } + } + s.frozenPageCount = s.basePagesCount + return nil +} diff --git a/freelist/freelist_test.go b/freelist/freelist_test.go index a3f5883..4c6e7e7 100644 --- a/freelist/freelist_test.go +++ b/freelist/freelist_test.go @@ -2,22 +2,39 @@ package freelist import ( "fmt" - "os" "testing" ) func createFreeList() (_ *FreeList, err error) { - file, err := os.OpenFile("test.free", os.O_CREATE|os.O_RDWR, 0666) - if err != nil { - return - } return New(Options{ - PageSize: 20, - File: file, + PageSize: 20, + BaseFilePath: "test.basefree", + DeltaFilePath: "test.deltafree", + FrozenPageCount: 2, }) } -func TestAddToFreeList(t *testing.T) { +func TestBuffer(t *testing.T) { + freeList, err := createFreeList() + if err != nil { + t.Fatal(err) + } + + pageNumbers := []uint32{1, 2, 3, 4} + //pageNumbers := []uint32{1, 2, 3, 4, 5, 6, 7, 8, 9, 10} + + err = freeList.AddPageNumbers(pageNumbers) + if err != nil { + t.Fatal(err) + } + + for range len(pageNumbers) + 1 { + pageNo, err := freeList.GetPageNumber() + fmt.Println(pageNo, err) + } +} + +func TestSave(t *testing.T) { freeList, err := createFreeList() if err != nil { t.Fatal(err) @@ -29,16 +46,67 @@ func TestAddToFreeList(t *testing.T) { if err != nil { t.Fatal(err) } + + err = freeList.Merge() + if err != nil { + t.Fatal(err) + } + + // for range len(pageNumbers) + 1 { + // pageNo, err := freeList.GetPageNumber() + // fmt.Println(pageNo, err) + // } } -func TestGetFromFreeList(t *testing.T) { +func TestDelta(t *testing.T) { freeList, err := createFreeList() if err != nil { t.Fatal(err) } - for range 4 { + serialized := []byte{9, 0, 0, 0, 10, 0, 0, 0} + + err = freeList.Restore(serialized) + if err != nil { + t.Fatal(err) + } + + for range 3 { pageNo, err := freeList.GetPageNumber() fmt.Println(pageNo, err) } + + err = freeList.AddPageNumbers([]uint32{20, 21, 22}) + if err != nil { + t.Fatal(err) + } + + err = freeList.AddPageNumbers([]uint32{23, 24, 25, 26}) + if err != nil { + t.Fatal(err) + } + + serialized, err = freeList.Serialize() + if err != nil { + t.Fatal(err) + } + + fmt.Println("free serialized:", serialized) + + err = freeList.Merge() + if err != nil { + t.Fatal(err) + } } + +// func TestGetFromFreeList(t *testing.T) { +// freeList, err := createFreeList() +// if err != nil { +// t.Fatal(err) +// } + +// for range 4 { +// pageNo, err := freeList.GetPageNumber() +// fmt.Println(pageNo, err) +// } +// } diff --git a/freelist/test.basefree b/freelist/test.basefree new file mode 100644 index 0000000000000000000000000000000000000000..923d25640b4c97052666ab8d34220d7800f74cbb GIT binary patch literal 80 ocmZQzAP(4C?b;y%R4NL@Vn8eo#Cw-I-em>SY(UHo#2i2j0D=?*+5i9m literal 0 HcmV?d00001 diff --git a/freelist/test.deltafree b/freelist/test.deltafree new file mode 100644 index 0000000..e69de29 diff --git a/freelist/test.free b/freelist/test.free deleted file mode 100644 index d2ebc93c2271c196c4ae2f786e7af25680238d13..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 20 WcmaFY`Ll?Tfq{Vuh?#+y1&9GQI|A(h diff --git a/go.mod b/go.mod index cbdd8af..45f442c 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,6 @@ go 1.24.2 require ( gopkg.in/ini.v1 v1.67.1 - gordenko.dev/dima/bin v0.0.0-20251101161620-7fa2d0f2923a + gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248 gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 ) diff --git a/go.sum b/go.sum index 157cb07..611f732 100644 --- a/go.sum +++ b/go.sum @@ -18,7 +18,7 @@ gopkg.in/ini.v1 v1.67.1/go.mod h1:x/cyOwCgZqOkJoDIJ3c1KNHMo10+nLGAhh+kn3Zizss= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gordenko.dev/dima/bin v0.0.0-20251101161620-7fa2d0f2923a h1:IAJvA+4IxPAUVXqXAIGrlfljoAbbzek4jPgrLwWJm9E= -gordenko.dev/dima/bin v0.0.0-20251101161620-7fa2d0f2923a/go.mod h1:/I+9fvRUzXHgXSwGEwOK1mBBM4BJ3AtCmoCZ2upDiwk= +gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248 h1:APYOkErjDz3hykfcddvmABsexoyKUgdRPU3XiM16F34= +gordenko.dev/dima/bin v0.0.0-20260528204801-12b890585248/go.mod h1:/I+9fvRUzXHgXSwGEwOK1mBBM4BJ3AtCmoCZ2upDiwk= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 h1:nyJ3mzTQ46yUeMZCdLyYcs7B5JCS54c67v84miyhq2E= gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69/go.mod h1:AxgKDktpqBVyIOhIcP+nlCpK+EsJyjN5kPdqyd8euVU= diff --git a/qb.go b/qb.go index 448a3df..5f48e8e 100644 --- a/qb.go +++ b/qb.go @@ -31,6 +31,9 @@ type TimestampCompressor interface { Renew() Lock() Unlock() + CreateDecompressor() TimestampDecompressor + Snapshot() ([][]byte, int) // chunks, size + Offset() int //LastTimestamp() uint32 } @@ -43,6 +46,10 @@ type ValueCompressor interface { Renew() // створює новий conbuf, але не чіпає state Lock() Unlock() + // fracDigits + CreateDecompressor(byte) ValueDecompressor + Snapshot() ([][]byte, int) // chunks, size + Offset() int //LastValue() float64 } @@ -66,6 +73,7 @@ const ( WrongResultCodeBug AbortCode = 6 RemoveREDOFileFailed AbortCode = 7 FailedAtreeRequest AbortCode = 8 + FailedGetPageNumber AbortCode = 9 UnknownTxLogRecordTypeBug AbortCode = 11 HasTimestampNoValueBug AbortCode = 12 NoMetricBug AbortCode = 13 @@ -77,6 +85,7 @@ const ( FailedFreeListSerialize AbortCode = 19 UnknownWorkerQueueItemBug AbortCode = 20 UnknownMetricWaitQueueItemBug AbortCode = 21 + RepeatableLock AbortCode = 22 // GetRecoveryRecipeFailed AbortCode = 26 LoadSnapshotFailed AbortCode = 27 diff --git a/txlog/snapshot.go b/txlog/snapshot.go new file mode 100644 index 0000000..942d0a8 --- /dev/null +++ b/txlog/snapshot.go @@ -0,0 +1,122 @@ +package txlog + +import ( + "fmt" + "io" + "os" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/util" +) + +func LoadSnapshot(fileName string) (snapshot Snapshot, err error) { + var ( + hasher = util.NewHasher() + header = make([]byte, metricHeaderSize) + body = make([]byte, atree.PageSize) + ) + + file, err := os.Open(fileName) + if err != nil { + return + } + + src := io.TeeReader(file, hasher) + metricsQty, err := bin.ReadVarSize(src) + if err != nil { + return + } + + for range metricsQty { + var metric Metric + err = bin.ReadNInto(src, header) + if err != nil { + return + } + + metric.MetricID, _ = bin.GetUint32(header[0:]) + metric.MetricType = qb.MetricType(header[4]) + metric.FracDigits = header[5] + metric.LastPageNo, _ = bin.GetUint32(header[6:]) + metric.Since, _ = bin.GetUint32(header[10:]) + metric.SinceValue, _ = bin.GetFloat64(header[14:]) + metric.Until, _ = bin.GetUint32(header[22:]) + metric.UntilValue, _ = bin.GetFloat64(header[26:]) + tSize, _ := bin.GetUint16(header[34:]) + vSize, _ := bin.GetUint16(header[36:]) + + buf := body[:tSize] + err = bin.ReadNInto(src, buf) + if err != nil { + return + } + // FIX + // metric.Timestamps = chunkenc.NewReverseTimeDeltaCompressor( + // conbuf.NewFromBuffer(buf), + // int(tSize), + // ) + + buf = body[:vSize] + err = bin.ReadNInto(src, buf) + if err != nil { + return + } + + // FIX + // if metric.MetricType == qb.Cumulative { + // metric.Values = chunkenc.NewReverseCumulativeDeltaCompressor( + // conbuf.NewFromBuffer(buf), + // int(vSize), + // metric.FracDigits, + // ) + // } else { + // metric.Values = chunkenc.NewReverseInstantDeltaCompressor( + // conbuf.NewFromBuffer(buf), + // int(vSize), + // metric.FracDigits, + // ) + // } + // s.metrics[metricID] = &metric + } + // FIX + // err = restoreFreeList(s.dataFreeList, src) + // if err != nil { + // return fmt.Errorf("restore dataFreeList: %s", err) + // } + // FIX + // err = restoreFreeList(s.indexFreeList, src) + // if err != nil { + // return fmt.Errorf("restore indexFreeList: %s", err) + // } + + calculatedChecksum := hasher.Sum32() + + writtenChecksum, err := bin.ReadUint32(file) + if err != nil { + return + } + + if calculatedChecksum != writtenChecksum { + err = fmt.Errorf("calculated checksum %d not equal written checksum %d", + calculatedChecksum, writtenChecksum) + return + } + return +} + +// serialized + frozen pages count + +// func restoreFreeList(freeList *freelist.FreeList, src io.Reader) error { +// size, err := bin.ReadVarSize(src) +// if err != nil { +// return err +// } +// serialized, err := bin.ReadN(src, size) +// if err != nil { +// return err +// } +// freeList.Restore(serialized) +// return nil +// } diff --git a/txlog/txlog.go b/txlog/txlog.go index 21e1873..415b14f 100644 --- a/txlog/txlog.go +++ b/txlog/txlog.go @@ -1,11 +1,13 @@ package txlog import ( + "bytes" "io" bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/conbuf" + "gordenko.dev/dima/qb/util" ) type AddedMetric struct { @@ -75,6 +77,89 @@ func (s *DeletedMetric) Parse(src io.Reader) (err error) { return nil } +func writeChunkedPayloadTo(chunks [][]byte, size int, w io.Writer) { + bin.WriteVarSize(w, size) + for _, chunk := range chunks { + if size >= len(chunk) { + w.Write(chunk) + size -= len(chunk) + } else { + w.Write(chunk[:size]) + return + } + } +} + +func readChunkedPayload(src io.Reader) (_ [][]byte, _ int, err error) { + var ( + chunks [][]byte + size int + remainingSize = size + ) + size, err = bin.ReadVarSize(src) + if err != nil { + return + } + for remainingSize > 0 { + var chunk []byte + if remainingSize >= conbuf.ChunkSize { + chunk, err = bin.ReadN(src, conbuf.ChunkSize) + if err != nil { + return + } + } else { + chunk = make([]byte, conbuf.ChunkSize) + err = bin.ReadNInto(src, chunk[:remainingSize]) + if err != nil { + return + } + } + chunks = append(chunks, chunk) + remainingSize -= conbuf.ChunkSize + } + return chunks, size, nil +} + +// fix - add changed pages +type DeletedMeasures struct { + MetricID uint32 + FreePageNumbers []uint32 +} + +func (s DeletedMeasures) Pack(w io.Writer) { + arr := []byte{ + CodeDeletedMeasures, + 0, 0, 0, 0, + } + bin.PutUint32(arr[1:], s.MetricID) + w.Write(arr) + bin.WriteVarSize(w, len(s.FreePageNumbers)) + for _, pageNo := range s.FreePageNumbers { + bin.WriteUint32(w, pageNo) + } +} + +func (s *DeletedMeasures) Parse(src io.Reader) (err error) { + s.MetricID, err = bin.ReadUint32(src) + if err != nil { + return + } + // free pages + qty, err := bin.ReadVarSize(src) + if err != nil { + return + } + for range qty { + var pageNo uint32 + pageNo, err = bin.ReadUint32(src) + if err != nil { + return + } + s.FreePageNumbers = append(s.FreePageNumbers, pageNo) + } + return nil +} + // type AppendedMeasure struct { // MetricID uint32 // Timestamp uint32 @@ -273,85 +358,63 @@ Nb - values payload // return nil // } -func writeChunkedPayloadTo(chunks [][]byte, size int, w io.Writer) { - bin.WriteVarSize(w, size) - for _, chunk := range chunks { - if size >= len(chunk) { - w.Write(chunk) - size -= len(chunk) - } else { - w.Write(chunk[:size]) - return - } - } +// PACK + +type PreparedData struct { + Packet []byte // пакет для запису в WAL + WriteToIndex []IndexPayloadToWrite + WriteToData []DataPayloadToWrite + ToWorker []any } -func readChunkedPayload(src io.Reader) (_ [][]byte, _ int, err error) { +func prepareData(input []any) PreparedData { + buffer := bytes.NewBuffer(nil) + buffer.Write([]byte{ + 0, 0, 0, 0, 0, 0, 0, 0, // size (max 8 byte) + 0, 0, 0, 0, // crc32 + }) + + hasher := util.NewHasher() + + w := io.MultiWriter(buffer, hasher) + var ( - chunks [][]byte - size int - remainingSize = size + writeToIndex []IndexPayloadToWrite + writeToData []DataPayloadToWrite ) - size, err = bin.ReadVarSize(src) - if err != nil { - return - } - for remainingSize > 0 { - var chunk []byte - if remainingSize >= conbuf.ChunkSize { - chunk, err = bin.ReadN(src, conbuf.ChunkSize) - if err != nil { - return - } - } else { - chunk = make([]byte, conbuf.ChunkSize) - err = bin.ReadNInto(src, chunk[:remainingSize]) - if err != nil { - return - } + + // 1. Пакую всі дані в WAL буфер запису + for _, untyped := range input { + switch x := untyped.(type) { + case AddedMetric: + x.Pack(w) + // case DeletedMetric: + // if len(x.FreePageNumbers) > 0 { + // s.freeList.AddPageNumbers(x.FreePageNumbers) + // } + case AppendedMeasures: + //s.packMeasuresIntoWALBuffer(x) + //case DeletedMeasures: + //case DeletedMeasuresSince: } - chunks = append(chunks, chunk) - remainingSize -= conbuf.ChunkSize } - return chunks, size, nil -} + // 2. Додаю розмір пакету і чексуму -// fix - add changed pages -type DeletedMeasures struct { - MetricID uint32 - FreePageNumbers []uint32 -} + //s.written += int64(len(packet)) + 12 -func (s DeletedMeasures) Pack(w io.Writer) { - arr := []byte{ - CodeDeletedMeasures, - 0, 0, 0, 0, - } - bin.PutUint32(arr[1:], s.MetricID) - w.Write(arr) - bin.WriteVarSize(w, len(s.FreePageNumbers)) - for _, pageNo := range s.FreePageNumbers { - bin.WriteUint32(w, pageNo) - } -} + packet := buffer.Bytes() + // write size + payloadSize := len(packet) - 12 + start := 8 - bin.CountVarSize(payloadSize) + bin.PutVarSize(packet[start:], payloadSize) + // checksum + bin.PutUint32(packet[8:], hasher.Sum32()) + // + //totalSize := (12 - start) + payloadSize -func (s *DeletedMeasures) Parse(src io.Reader) (err error) { - s.MetricID, err = bin.ReadUint32(src) - if err != nil { - return + return PreparedData{ + Packet: packet[start:], + WriteToIndex: writeToIndex, + WriteToData: writeToData, } - // free pages - qty, err := bin.ReadVarSize(src) - if err != nil { - return - } - for range qty { - var pageNo uint32 - pageNo, err = bin.ReadUint32(src) - if err != nil { - return - } - s.FreePageNumbers = append(s.FreePageNumbers, pageNo) - } - return nil } diff --git a/txlog/writer.go b/txlog/writer.go index 732311a..19debde 100644 --- a/txlog/writer.go +++ b/txlog/writer.go @@ -1,11 +1,13 @@ package txlog import ( + "bufio" "bytes" "errors" "fmt" - "hash/crc32" + "io" "log" + "math" "os" "path/filepath" "sync" @@ -13,6 +15,8 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/atree" + "gordenko.dev/dima/qb/freelist" + "gordenko.dev/dima/qb/util" ) const ( @@ -26,6 +30,8 @@ const ( filePerm = 0770 dumpSnapshotAfterNBytes = 1024 * 1024 * 1024 // 1 GB + + writeBufferSize = 4 * 1024 * 1024 ) const ( @@ -46,31 +52,32 @@ func JoinChangesFileName(dir string, logNumber int) string { return filepath.Join(dir, fmt.Sprintf("%d.changes", logNumber)) } +type MetricsState struct { + Metrics []Metric + WaitCh chan struct{} +} + type Changes struct { - Records []any - LogNumber int - ForceSnapshot bool - ExitWaitGroup *sync.WaitGroup - WaitCh chan struct{} + Records []any + MetricsCh chan MetricsState } type Writer struct { - mutex sync.Mutex - freelist FreeList - atree *atree.Atree - logNumber int - dir string - w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск - wal *os.File - dataFile *os.File - indexFile *os.File - dataPage []byte - indexPage []byte - input []any - //buf *bytes.Buffer - //pagesToWrite []atree.PageToWrite - workerReqs []any - //waitCh chan struct{} + mutex sync.Mutex + dataPagesCount uint32 + indexPagesCount uint32 + dataFreeList *freelist.FreeList + indexFreeList *freelist.FreeList + atree *atree.Atree + logNumber int + dir string + w *bytes.Buffer // wal буфер для упаковки даних перед записом на диск + wal *os.File + dataFile *os.File + indexFile *os.File + // буфер для складання payload на сторінку, розрахунку CRC32 та інше + pageBuffer []byte + input []any appendToWorkerQueue func(any) //lsn uint32 written int64 @@ -84,7 +91,8 @@ type WriterOptions struct { Dir string LogNumber int // номер журнала AppendToWorkerQueue func(any) - FreeList FreeList + DataFreeList *freelist.FreeList + IndexFreeList *freelist.FreeList Atree *atree.Atree ExitCh chan struct{} WaitGroup *sync.WaitGroup @@ -97,8 +105,11 @@ func NewWriter(opt WriterOptions) (*Writer, error) { if opt.AppendToWorkerQueue == nil { return nil, errors.New("AppendToWorkerQueue option is required") } - if opt.FreeList == nil { - return nil, errors.New("FreeList option is required") + if opt.DataFreeList == nil { + return nil, errors.New("DataFreeList option is required") + } + if opt.IndexFreeList == nil { + return nil, errors.New("IndexFreeList option is required") } if opt.Atree == nil { return nil, errors.New("Atree option is required") @@ -113,7 +124,8 @@ func NewWriter(opt WriterOptions) (*Writer, error) { s := &Writer{ dir: opt.Dir, appendToWorkerQueue: opt.AppendToWorkerQueue, - freelist: opt.FreeList, + dataFreeList: opt.DataFreeList, + indexFreeList: opt.IndexFreeList, atree: opt.Atree, logNumber: opt.LogNumber, exitCh: opt.ExitCh, @@ -161,53 +173,140 @@ func (s *Writer) Run() { } } +func (s *Writer) getDataPageNumber() (uint32, bool, error) { + pageNo, err := s.dataFreeList.GetPageNumber() + if err != nil { + return 0, false, err + } + if pageNo > 0 { + return pageNo, true, nil + } + if s.dataPagesCount < math.MaxUint32 { + s.dataPagesCount++ + return s.dataPagesCount, false, nil + } + return 0, false, errors.New("no space") +} + +func (s *Writer) getIndexPageNumber() (uint32, bool, error) { + pageNo, err := s.indexFreeList.GetPageNumber() + if err != nil { + return 0, false, err + } + if pageNo > 0 { + return pageNo, true, nil + } + if s.indexPagesCount < math.MaxUint32 { + s.indexPagesCount++ + return s.indexPagesCount, false, nil + } + return 0, false, errors.New("no space") +} + func (s *Writer) packAndWrite() (err error) { s.mutex.Lock() + //isExited := s.isExited input := s.input s.input = nil + + // var exitWaitGroup *sync.WaitGroup + // if s.isExited { + // exitWaitGroup = s.waitGroup + // } s.mutex.Unlock() - w := bytes.NewBuffer(nil) - - // 1. Пакую всі дані в WAL буфер запису - for _, untyped := range input { - switch x := untyped.(type) { - case AddedMetric: - x.Pack(w) - // case DeletedMetric: - // if len(x.FreePageNumbers) > 0 { - // s.freeList.AddPageNumbers(x.FreePageNumbers) - // } - case AppendedMeasures: - // - - s.packMeasuresIntoWALBuffer(x) - //case DeletedMeasures: - //case DeletedMeasuresSince: - } - } - // 2. Додаю розмір пакету і чексуму + prepared := prepareData(input) // 3. Пишу на диск WAL (append) - - // 4. Пишу в atree сторінки - for range 1 { - err = s.writePagesToAtree(nil, nil) - if err != nil { - return - } + n, err := s.wal.Write(prepared.Packet) + if err != nil { + return + } + if n != len(prepared.Packet) { + return fmt.Errorf("written %d != total size %d", n, len(prepared.Packet)) + } + if err = s.wal.Sync(); err != nil { + return } - // 5. Пишу зміни в free list + // 4. Пишу в atree сторінки + err = s.writePagesToAtree(prepared.WriteToIndex, prepared.WriteToData) + if err != nil { + return + } // 6. відправляю input - worker-у - // flush - err = s.flush(w) - if err != nil { - return err + var ( + forceSnapshot bool + ) + + if s.written > dumpSnapshotAfterNBytes { + forceSnapshot = true } + // if isExited && s.written > 0 { + // forceSnapshot = true + // } + + if forceSnapshot { + metricsStateCh := make(chan MetricsState) + + s.appendToWorkerQueue(Changes{ + Records: prepared.ToWorker, + MetricsCh: metricsStateCh, + }) + + if err := s.wal.Close(); err != nil { + return fmt.Errorf("close changes file: %s", err) + } + + state := <-metricsStateCh + + s.logNumber++ + + // write snapshot FIX + err = writeSnapshot(Snapshot{ + LogNumber: s.logNumber, + Dir: s.dir, + WriteBufferSize: writeBufferSize, + Metrics: state.Metrics, + DataFreeList: s.dataFreeList, + IndexFreeList: s.indexFreeList, + }) + if err != nil { + return fmt.Errorf("write snapshot file: %s", err) + } + + var err error + s.wal, err = os.OpenFile( + JoinChangesFileName(s.dir, s.logNumber), + os.O_CREATE|os.O_WRONLY, + filePerm, + ) + if err != nil { + return fmt.Errorf("create new changes file: %s", err) + } + + s.written = 0 + // release Worker + close(state.WaitCh) + } else { + s.appendToWorkerQueue(Changes{ + Records: prepared.ToWorker, + }) + } + + // Якщо потрібен снапшот - відправляю канал із буфером розміру 1, + // в який воркер має покласти снапшот отриманий після застосування змін (records). + // Після чого Writer створює нові файли snapshot і changes і працює далі. + + // flush + // err = s.flush(w) + // if err != nil { + // return err + // } + // FIX send to worker workerReqs return nil } @@ -215,7 +314,6 @@ func (s *Writer) packAndWrite() (err error) { func (s *Writer) Append(req any) { s.mutex.Lock() s.input = append(s.input, req) - s.workerReqs = append(s.workerReqs, req) s.mutex.Unlock() select { @@ -237,87 +335,10 @@ func (s *Writer) Append(req any) { // s.waitCh = make(chan struct{}) // } -func (s *Writer) flush(w *bytes.Buffer) error { - s.mutex.Lock() +// func (s *Writer) flush(w *bytes.Buffer) error { - //pagesToWrite := s.pagesToWrite - workerReqs := s.workerReqs - isExited := s.isExited - - var exitWaitGroup *sync.WaitGroup - if s.isExited { - exitWaitGroup = s.waitGroup - } - - if w.Len() > packetPrefixSize { - //s.lsn++ - //lsn := s.lsn - packet := make([]byte, w.Len()) - copy(packet, w.Bytes()) - //s.reset() fix - - s.written += int64(len(packet)) + 12 - s.mutex.Unlock() - - bin.PutUint32(packet[lengthIdx:], uint32(len(packet)-packetPrefixSize)) - //bin.PutUint32(packet[lsnIdx:], lsn) - bin.PutUint32(packet[checksumIdx:], crc32.ChecksumIEEE(packet[8:])) - - n, err := s.wal.Write(packet) - if err != nil { - return fmt.Errorf("TxLog write: %s", err) - } - if n != len(packet) { - return fmt.Errorf("TxLog written %d != packet size %d", n, len(packet)) - } - if err := s.wal.Sync(); err != nil { - return fmt.Errorf("TxLog sync: %s", err) - } - - // err = s.writePagesToAtree(pagesToWrite) - // if err != nil { - // return fmt.Errorf("TxLog writePagesToAtree: %s", err) - // } - } else { - s.mutex.Unlock() - } - - var forceSnapshot bool - - if s.written > dumpSnapshotAfterNBytes { - forceSnapshot = true - } - - if isExited && s.written > 0 { - forceSnapshot = true - } - - if forceSnapshot { - if err := s.wal.Close(); err != nil { - return fmt.Errorf("close changes file: %s", err) - } - s.logNumber++ - var err error - s.wal, err = os.OpenFile( - JoinChangesFileName(s.dir, s.logNumber), - os.O_CREATE|os.O_WRONLY, - filePerm, - ) - if err != nil { - return fmt.Errorf("create new changes file: %s", err) - } - - s.written = 0 - } - - s.appendToWorkerQueue(Changes{ - Records: workerReqs, - ForceSnapshot: forceSnapshot, - LogNumber: s.logNumber, - ExitWaitGroup: exitWaitGroup, - }) - return nil -} +// return nil +// } func (s *Writer) exit() { s.mutex.Lock() @@ -329,45 +350,79 @@ func (s *Writer) exit() { } } +type IndexPayloadToWrite struct { + PageNo uint32 + Payload []byte + ZeroLevel bool +} + +type DataPayloadToWrite struct { + PrevPageNo uint32 + Timestamps [][]byte // chunks + TimestampsSize int + Values [][]byte // chunks + ValuesSize int + PageNo uint32 +} + // writePagesToAtree - записує сторінки в .data та .index файли -func (s *Writer) writePagesToAtree(levels []*atree.IndexLevel, dataPages []*atree.DataPage) (err error) { - for _, p := range dataPages { - p.Checksum = atree.ChunksToDataPage(s.dataPage, atree.ChunksToDataPageReq{ +func (s *Writer) writePagesToAtree(indexItems []IndexPayloadToWrite, dataItems []DataPayloadToWrite) (err error) { + for _, p := range dataItems { + atree.ChunksToDataPage(s.pageBuffer, atree.ChunksToDataPageReq{ PrevPageNo: p.PrevPageNo, Timestamps: p.Timestamps, TimestampsSize: p.TimestampsSize, Values: p.Values, ValuesSize: p.ValuesSize, }) - off := (p.PageNo - 1) * atree.DataPageSize - n, err := s.dataFile.WriteAt(s.dataPage, int64(off)) + var ( + off = (p.PageNo - 1) * atree.DataPageSize + n int + ) + n, err = s.dataFile.WriteAt(s.pageBuffer, int64(off)) if err != nil { - return err + return } if n != atree.DataPageSize { return fmt.Errorf("write %d instead of %d", n, atree.DataPageSize) } } - for levelIdx, level := range levels { - for _, p := range level.Filled { - // if len(p.Data) != atree.PageSize { - // return fmt.Errorf("wrong page %d size: %d", - // p.PageNo, len(p.Data)) - // } - p.Checksum = atree.DataToIndexPage(s.indexPage, p.Data, levelIdx == 0) - off := (p.PageNo - 1) * atree.IndexPageSize - n, err := s.indexFile.WriteAt(s.indexPage, int64(off)) - if err != nil { - return err - } - if n != atree.IndexPageSize { - return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize) - } + for _, p := range indexItems { + // if len(p.Payload) > atree.MaxIndexPayload { + // return fmt.Errorf("wrong index payload size: %d", len(p.Payload)) + // } + indexPageBuffer := s.pageBuffer[:atree.IndexPageSize] + atree.DataToIndexPage(indexPageBuffer, p.Payload, p.ZeroLevel) + var ( + off = (p.PageNo - 1) * atree.IndexPageSize + n int + ) + n, err = s.indexFile.WriteAt(indexPageBuffer, int64(off)) + if err != nil { + return + } + if n != atree.IndexPageSize { + return fmt.Errorf("write %d instead of %d", n, atree.IndexPageSize) } } return nil } +type AppendMeasuresSummary struct { + MetricID uint32 + //TimestampsOffset int // (заповнені одразу) + //ValuesOffset int // (заповнені одразу) + // Timestamps [][]byte + // TimestampsSize int + // Values [][]byte + // ValuesSize int + LastPageNo uint32 + Index [][]byte + ResultCode int + WrittenCount int + ResultCh chan struct{} +} + type AppendedMeasures struct { MetricID uint32 TimestampsOffset int // (заповнені одразу) @@ -378,6 +433,9 @@ type AppendedMeasures struct { ValuesSize int IndexLevels []*atree.IndexLevel DataPages []*atree.DataPage // fix - prevPageNo for the 1st data page + ResultCode int + WrittenCount int + ResultCh chan struct{} } func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) { @@ -392,7 +450,7 @@ func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) { } s.w.WriteByte(reused) bin.WriteUint32(s.w, p.PrevPageNo) - bin.WriteUint32(s.w, p.Checksum) + //bin.WriteUint32(s.w, p.Checksum) var ( timestampsOffset int valuesOffset int @@ -430,7 +488,7 @@ func (s *Writer) packMeasuresIntoWALBuffer(req AppendedMeasures) (err error) { reused = 1 } s.w.WriteByte(reused) - bin.WriteUint32(s.w, p.Checksum) + //bin.WriteUint32(s.w, p.Checksum) records := p.Data if idx == 0 { records = p.Data[level.Offset:] @@ -502,3 +560,204 @@ func (s *Writer) sendSignal() { // } // helpers + +/* +Формат: +metricsQty - varuint +[metric]* +де metric - це: +metricID - 4b +metricType - 1b +fracDigits - 1b +lastPageNo - 4b +since - 4b +sinceValue - 8b +until - 4b +untilValue - 8b +timestamps size - 2b +values size - 2b +timestams payload - Nb +values payload - Nb +data free list frozen pages - varsize +dataFreeList size - varsize +dataFreeList - Nb +index free list frozen pages - varsize +indexFreeList size - varsize +indexFreeList - Nb +CRC32 - 4b +*/ + +const metricHeaderSize = 38 + +type Metric struct { + MetricID uint32 + MetricType qb.MetricType + FracDigits byte + LastPageNo uint32 + Since uint32 + SinceValue float64 + Until uint32 + UntilValue float64 + Timestamps [][]byte + TimestampsSize int + Values [][]byte + ValuesSize int +} + +type Snapshot struct { + LogNumber int + Dir string + WriteBufferSize int + Metrics []Metric + IndexFreeList *freelist.FreeList + DataFreeList *freelist.FreeList +} + +func writeSnapshot(snapshot Snapshot) (err error) { + var ( + fileName = filepath.Join(snapshot.Dir, fmt.Sprintf("%d.snapshot", snapshot.LogNumber)) + hasher = util.NewHasher() + prefix = make([]byte, metricHeaderSize) + ) + + file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0770) + if err != nil { + return + } + + dst := io.MultiWriter(bufio.NewWriterSize(file, snapshot.WriteBufferSize), hasher) + + _, err = bin.WriteVarSize(dst, len(snapshot.Metrics)) + if err != nil { + return + } + + for _, metric := range snapshot.Metrics { + tSize := metric.TimestampsSize + vSize := metric.ValuesSize + + bin.PutUint32(prefix[0:], metric.MetricID) + prefix[4] = byte(metric.MetricType) + prefix[5] = metric.FracDigits + bin.PutUint32(prefix[6:], metric.LastPageNo) + bin.PutUint32(prefix[10:], metric.Since) + bin.PutFloat64(prefix[14:], metric.SinceValue) + bin.PutUint32(prefix[22:], metric.Until) + bin.PutFloat64(prefix[26:], metric.UntilValue) + bin.PutUint16(prefix[34:], uint16(tSize)) + bin.PutUint16(prefix[36:], uint16(vSize)) + + _, err = dst.Write(prefix) + if err != nil { + return + } + // copy timestamps + writeChunks(dst, metric.Timestamps, tSize) + // copy values + writeChunks(dst, metric.Values, vSize) + } + // free data pages + _, err = bin.WriteVarSize(dst, snapshot.DataFreeList.Pages()) + if err != nil { + return + } + err = freeListWriteTo(snapshot.DataFreeList, dst) + if err != nil { + return + } + // free index pages + _, err = bin.WriteVarSize(dst, snapshot.IndexFreeList.Pages()) + if err != nil { + return + } + err = freeListWriteTo(snapshot.IndexFreeList, dst) + if err != nil { + return + } + + bin.WriteUint32(file, hasher.Sum32()) + + err = file.Sync() + if err != nil { + return + } + + err = file.Close() + if err != nil { + return + } + // копіюю сторінки із delta файла в base файл і потім роблю Truncate + err = snapshot.DataFreeList.Merge() // fix - get frozen pages count + unfilled + if err != nil { + return + } + err = snapshot.IndexFreeList.Merge() + if err != nil { + return + } + + // prevLogNumber := logNumber - 1 + // prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber)) + // prevSnapshot := filepath.Join(s.dir, fmt.Sprintf("%d.snapshot", prevLogNumber)) + + // isExist, err := isFileExist(prevChanges) + // if err != nil { + // return + // } + + // if isExist { + // err = os.Remove(prevChanges) + // if err != nil { + // qb.Abort(qb.DeletePrevChangesFileFailed, err) + // } + // } + + // isExist, err = isFileExist(prevSnapshot) + // if err != nil { + // return + // } + + // if isExist { + // err = os.Remove(prevSnapshot) + // if err != nil { + // qb.Abort(qb.DeletePrevSnapshotFileFailed, err) + // } + // } + return +} + +// HELPERS + +func freeListWriteTo(freeList *freelist.FreeList, dst io.Writer) error { + serialized, err := freeList.Serialize() + if err != nil { + qb.Abort(qb.FailedFreeListSerialize, err) + } + _, err = bin.WriteVarSize(dst, len(serialized)) + if err != nil { + return err + } + _, err = dst.Write(serialized) + if err != nil { + return err + } + return nil +} + +func writeChunks(dst io.Writer, chunks [][]byte, size int) (err error) { + remaining := size + for _, buf := range chunks { + if remaining < len(buf) { + buf = buf[:remaining] + } + _, err = dst.Write(buf) + if err != nil { + return + } + remaining -= len(buf) + if remaining == 0 { + break + } + } + return +} diff --git a/util/until.go b/util/until.go index 1ca881d..4c824eb 100644 --- a/util/until.go +++ b/util/until.go @@ -1,6 +1,7 @@ package util import ( + "hash" "hash/crc32" ) @@ -15,3 +16,7 @@ func CalcChecksum(page []byte) uint32 { func IsValidChecksum(page []byte, checksum uint32) bool { return crc32.Checksum(page, castagnoliTable) == checksum } + +func NewHasher() hash.Hash32 { + return crc32.New(castagnoliTable) +}