diff --git a/a b/a deleted file mode 100644 index 4fbaf93..0000000 --- a/a +++ /dev/null @@ -1 +0,0 @@ -01 11 27 00 00 01 02 11 4e 46 bd 00 00 00 00 00 00 00 00 00 06 11 27 00 00 00 9c 5e 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 98 5c dd 69 76 94 80 80 ca 81 20 81 59 82 07 82 56 82 75 83 2a 83 6c 84 11 84 54 85 02 85 0f 85 4e 86 24 87 04 87 55 87 56 87 5a 87 7b 88 01 88 36 89 0a 89 39 8a 02 8a 30 8a 36 8a 52 8a 7d 8b 5b 8b 5c 8c 11 8c 62 8d 22 8d 23 8d 2e 8d 30 8d 52 8d 7e 8e 32 8e 7f 8f 52 90 29 90 55 91 07 91 55 92 16 92 36 93 04 93 07 93 52 94 19 94 3b 94 73 95 2b 95 31 95 74 96 36 96 7c 97 18 97 54 98 28 98 77 99 35 9a 12 9a 4b 9b 04 9b 62 9c 40 9c 4a 9c 61 9d 24 9d 39 9d 78 9e 5c 9f 16 9f 6e a0 13 a0 20 a0 61 a0 7e a1 48 a2 09 a2 1e a2 6f a2 71 a3 55 a3 57 a4 0e a4 2a a4 7c a5 3c a5 74 a6 03 a6 50 a6 61 a6 76 a7 32 a7 54 a7 65 a8 25 a8 60 a9 18 a9 4b a9 77 aa 44 aa 56 ab 27 ab 61 ab 6f ac 35 ac 68 ad 06 ad 26 ad 2a ad 68 ae 35 af 05 af 4b b0 21 b0 53 b0 73 b1 0a b1 44 b1 55 b1 78 b2 4d b3 1f ff b3 76 b4 30 b4 51 b5 31 b5 70 b6 07 b6 52 b7 2e b8 0e b8 34 b9 03 b9 11 b9 26 b9 4d ba 1a ba 7d bb 16 bb 4c bb 79 bc 5a bc 72 bc 79 bd 01 bd 29 be 01 be 13 be 5e bf 34 c0 10 c0 4a c0 50 c1 1a c1 3f c1 64 c1 65 c1 7a c1 7d c2 27 c2 30 c2 73 c3 57 c3 7e c4 41 c4 42 c4 52 c5 00 c5 57 c6 1b c6 49 c6 5f c7 2d c7 73 c8 0b c8 42 c8 6a c9 33 ca 0e ca 20 ca 52 ca 63 cb 32 cb 72 cc 23 cc 3f cc 5a cc 60 cd 41 cd 5e ce 32 ce 4c cf 1d cf 67 d0 11 d0 6f d0 79 d1 2c d1 6a d2 1d d2 2f d2 5e d3 0a d3 12 d3 60 d4 09 d4 52 d4 77 d5 4b d5 62 d5 64 d5 7b d6 1e d6 3d d7 18 d7 67 d8 0f d8 4d d8 73 d9 03 d9 40 da 16 da 45 db 09 db 55 dc 15 dc 2f dc 41 dc 43 dc 6a dc 6d dd 33 dd 42 de 07 de 08 de 35 de 6a df 1f df 20 e0 02 e0 63 e1 29 e1 69 e2 01 e2 32 e2 7e e3 47 e4 19 e4 21 e4 28 ff e4 57 e5 35 e5 6a e6 2d e6 31 e6 67 e7 09 e7 65 e7 6c e8 41 e8 5f e8 71 e9 42 e9 51 ea 32 ea 4f eb 29 eb 33 eb 66 eb 7e ec 4c ec 7f ed 32 ed 3c ed 7f ee 01 ee 1f ee 20 ee 72 ef 39 f0 06 f0 44 f0 75 f1 40 f1 75 f2 2f f2 42 f2 5f f3 37 f3 7b f4 2d f4 3e f4 53 f5 05 f5 11 f5 4f f6 09 f6 43 f6 73 f7 02 f7 5c f8 30 f8 68 f8 79 f9 0e f9 5c f9 5e f9 78 fa 0f fa 4c fb 25 fb 62 fc 3e fd 07 fd 29 fd 7a fe 25 fe 74 ff 1a ff 38 ff 41 ff 7d 81 00 5d 81 00 60 81 01 32 81 01 3a 81 01 50 81 02 08 81 02 38 81 02 5f 81 03 01 81 03 19 81 03 48 81 04 07 81 04 4b 81 04 5e 81 05 05 81 05 21 81 05 28 81 05 47 81 05 5d 81 05 7d 81 06 0c 81 06 5c 81 07 0a 81 07 67 81 07 7b 81 08 2e 81 08 6c 81 08 76 81 09 0a 81 09 4a 81 09 50 81 09 64 81 0a 16 81 0a 56 81 0b 11 81 0b 28 81 0b 52 81 0b 7f 81 0c 0a 81 0c 5d 81 0d 3f 81 0e 0c 81 0e 33 81 0e 69 81 0f 06 81 0f 39 81 0f 3f 81 0f 40 81 0f 41 81 0f 74 81 10 51 81 10 67 81 11 1d 81 12 01 81 12 1a 81 12 3c ff 81 12 4f 81 12 7e 81 13 4c 81 13 56 81 13 6c 81 14 4e 81 14 65 81 15 41 81 15 49 81 16 07 81 16 18 81 16 76 81 17 56 81 18 00 81 18 12 81 18 6e 81 19 4b 81 1a 29 81 1a 33 81 1b 0e 81 1b 46 81 1b 6f 81 1b 78 81 1c 04 81 1c 37 81 1c 40 81 1c 7c 81 1d 12 81 1d 2e 81 1d 56 81 1d 59 81 1e 23 81 1e 46 81 1e 4e 81 1e 7b 81 1f 3d 81 1f 75 81 20 19 81 20 4a 81 20 58 81 21 07 81 21 4a 81 21 64 81 21 77 81 22 1b 81 22 2a 81 22 76 81 23 1f 81 23 40 81 23 66 81 24 0b 81 24 1d 81 24 75 81 25 51 81 25 76 81 26 33 81 26 62 81 27 29 81 27 61 81 28 3b 81 29 15 81 29 45 81 29 6c 81 2a 35 81 2a 7d 81 2b 50 81 2c 07 81 2c 68 81 2c 74 81 2d 4e 81 2e 2e 81 2e 7f 81 2f 07 81 2f 34 81 2f 7b 81 30 0c 81 30 5b 81 31 3e 81 31 60 81 32 2f 81 32 7b 81 33 30 81 34 0a 81 34 11 81 34 6a 81 35 18 81 35 5c 81 35 79 81 36 5a 81 37 0d 81 37 20 81 37 27 81 37 5b 81 38 3b 81 38 73 81 39 16 81 39 65 81 3a 48 81 3a 67 81 3b 24 81 3b 43 81 3c 26 81 3c 6a 81 3d 25 81 3d 4d 81 3d 55 81 3d 5d 81 3e 12 81 3e 24 81 3e 49 81 3f 0a 81 3f 17 81 3f 28 81 40 09 81 40 30 81 41 02 81 41 2a 81 41 52 81 41 64 81 41 75 81 42 33 81 42 59 81 43 1b 81 43 59 81 43 60 81 43 7f 81 44 26 81 44 29 ff 81 44 3b 81 44 68 81 45 16 81 45 5c 81 46 03 81 46 5b 81 46 7f 81 47 06 81 47 25 81 47 79 81 48 41 81 49 1a 81 49 1f 81 49 32 81 49 58 81 4a 35 81 4a 39 81 4a 5b 81 4a 72 81 4b 4e 81 4c 1b 81 4c 48 81 4c 53 81 4c 67 81 4d 00 81 4d 07 81 4d 57 81 4e 39 81 4e 64 81 4e 69 81 4f 04 81 4f 64 81 50 46 81 50 59 81 51 10 81 51 73 81 52 15 81 52 28 81 52 64 81 53 32 81 54 01 81 54 35 81 54 3c 81 54 4d 81 55 07 81 55 64 81 55 7e 81 56 39 81 56 54 81 56 73 81 57 42 81 57 58 81 58 2b 81 58 34 81 59 0f 81 59 5f 81 59 78 81 5a 02 81 5a 48 81 5a 7a 81 5b 44 81 5b 5e 81 5c 0f 81 5c 26 81 5c 6c 81 5d 1c 81 5d 3c 81 5d 5d 81 5e 30 81 5e 3f 81 5e 72 81 5f 06 81 5f 3a 81 5f 7a 81 5f 7f 81 60 0c 81 60 18 81 60 72 81 61 11 81 61 25 81 61 79 81 62 1c 81 62 69 81 63 4d 81 64 20 81 64 2e 81 64 30 81 64 4d 81 64 61 81 65 35 81 66 16 81 66 3e 81 67 11 81 67 20 81 67 7e 81 68 21 81 68 64 81 68 7c 81 69 1c 81 69 1e 81 69 76 81 6a 32 81 6a 76 81 6b 41 81 6b 5a 81 6c 30 81 6c 6d 81 6d 2e 81 6d 79 81 6e 0b 81 6e 29 81 6e 49 81 6e 7a 81 6f 06 81 6f 2a 81 6f 3b 81 6f 77 81 70 40 81 71 15 81 71 17 81 71 5b 81 71 69 81 71 7c 81 72 29 81 73 02 81 73 5e 81 74 36 81 75 08 ff 81 75 2b 81 75 5d 81 76 3c 81 77 02 81 77 64 81 78 20 81 78 51 81 78 79 81 79 2d 81 79 59 81 7a 20 81 7a 5b 81 7b 15 81 7b 22 81 7b 4a 81 7c 0a 81 7c 3b 81 7c 42 81 7d 06 81 7d 0b 81 7d 35 81 7d 7f 81 7e 61 81 7e 72 81 7f 36 81 7f 39 81 7f 71 82 00 1e 82 00 45 82 01 28 82 01 49 82 01 53 82 02 35 82 02 77 82 03 05 82 03 43 82 04 1b 82 04 6b 82 05 4d 82 06 29 82 06 34 82 06 3e 82 07 1d 82 07 29 82 07 3c 82 07 42 82 07 54 82 07 6f 82 08 28 82 08 5c 82 09 22 82 09 52 82 09 7b 82 0a 2a 82 0a 65 82 0a 6c 82 0b 19 82 0b 61 82 0c 0a 82 0c 3d 82 0d 21 82 0d 30 82 0d 71 82 0e 1f 82 0e 48 82 0f 07 82 0f 0f 82 0f 30 82 0f 4b 82 0f 6d 82 10 2e 82 10 67 82 11 10 82 11 40 82 11 7e 82 12 5c 82 13 3e 82 14 11 82 14 5b 82 15 35 82 16 10 82 16 30 82 17 11 82 17 53 82 17 69 82 18 1d 82 18 5b 82 18 5d 82 19 39 82 19 53 82 19 7e 82 1a 1c 82 1a 2c 82 1a 35 82 1a 74 82 1b 34 82 1b 6d 82 1b 75 82 1c 1e 82 1c 51 82 1d 1c 82 1d 65 82 1e 2f 82 1e 41 82 1f 0b 82 1f 5d 82 20 22 82 21 05 82 21 15 82 21 5f 82 21 73 82 22 1e 82 22 62 82 22 6b 82 23 2f 82 24 09 82 24 6c 82 25 32 82 25 7a 82 26 51 82 26 67 82 26 6d 82 27 23 82 27 4b 82 28 18 82 28 64 82 28 6c 82 29 3a ff 82 2a 14 82 2a 3a 82 2b 12 82 2b 72 82 2c 22 82 2c 65 82 2d 06 82 2d 1f 82 2d 4f 82 2e 13 82 2e 1b 82 2e 79 82 2f 4d 82 2f 67 82 30 44 82 30 70 82 31 47 82 32 13 82 32 5e 82 32 7c 82 33 57 82 34 02 82 34 1a 82 34 2b 82 34 37 82 34 6c 82 35 1f 82 35 4a 82 36 0b 82 36 17 82 36 49 82 37 1f 82 38 03 82 38 24 82 38 3f 82 39 01 82 39 12 82 39 20 82 39 24 82 39 59 82 3a 03 82 3a 19 82 3a 74 82 3b 2a 82 3b 37 82 3b 7e 82 3c 34 82 3c 4b 82 3c 69 82 3d 45 82 3e 20 82 3e 49 82 3f 0e 82 3f 1e 82 3f 20 82 3f 7c 82 40 34 82 40 75 82 41 39 82 42 0f 82 42 4a 82 43 0d 82 43 11 82 43 5f 82 43 69 82 44 2a 82 44 61 82 44 7a 82 45 34 82 46 10 82 46 11 82 46 2f 82 46 3d 82 46 61 82 47 07 82 47 0f 82 47 6b 82 48 20 82 48 42 82 48 65 82 48 68 82 49 1c 82 49 45 82 49 56 82 49 74 82 4a 1c 82 4a 65 82 4b 1f 82 4b 28 82 4b 7a 82 4b 7c 82 4c 10 82 4c 17 82 4c 35 82 4d 08 82 4d 69 82 4e 01 82 4e 18 82 4e 22 82 4e 4a 82 4e 72 82 4f 25 82 50 08 82 50 2b 82 50 2e 82 50 3c 82 50 6f 82 50 7e 82 51 2d 82 51 4e 82 52 20 82 52 62 82 53 11 82 53 45 82 54 23 82 54 78 82 55 1a 82 55 1d 82 55 33 82 55 67 82 56 27 82 56 7a 82 57 2f 82 58 12 82 58 55 82 58 68 82 59 10 82 59 6d ff 82 5a 09 82 5a 27 82 5a 2a 82 5b 04 82 5b 51 82 5c 1d 82 5c 42 82 5c 56 82 5d 15 82 5d 52 82 5d 53 82 5d 7b 82 5e 26 82 5e 2d 82 5e 7f 82 5f 46 82 5f 50 82 5f 6e 82 60 3c 82 61 0f 82 61 10 82 61 3b 82 61 4b 82 62 16 82 62 18 82 62 37 82 62 69 82 62 71 82 62 7f 82 63 0e 82 63 50 82 64 0b 82 64 49 82 65 15 82 65 3a 82 65 5a 82 65 71 82 66 3b 82 66 5a 82 66 72 82 67 0b 82 67 6c 82 67 7f 82 68 15 82 68 4d 82 69 1b 82 69 39 82 69 7b 82 6a 52 82 6a 72 82 6b 13 82 6b 3a 82 6b 5b 82 6c 29 82 6c 5d b6 82 6d 0e 00 82 6d 4e 82 6e 19 82 6e 6b 82 6f 42 82 6f 6d 82 70 21 82 70 6d 82 71 36 82 71 79 82 72 23 82 72 73 82 73 20 82 73 7a 82 74 06 82 74 43 82 75 26 82 75 46 82 76 1e 82 76 4f 82 77 26 82 77 4a 82 78 0c 82 78 61 82 79 0f 82 79 4a 82 7a 0c 82 7a 13 82 7a 44 82 7a 48 82 7a 7d 82 7b 16 82 7b 62 82 7b 6e 82 7c 49 82 7c 78 82 7d 50 82 7e 21 82 7e 34 82 7e 54 82 7e 62 82 7e 64 82 7f 24 82 7f 7c 83 00 18 83 00 5e 83 01 27 83 01 51 ae diff --git a/atree/atree.go b/atree/atree.go deleted file mode 100644 index efce9d8..0000000 --- a/atree/atree.go +++ /dev/null @@ -1,224 +0,0 @@ -package atree - -import ( - "errors" - "os" - "sync" -) - -const ( - filePerm = 0666 -) - -type PeriodsWriter interface { - Feed(uint32, float64) - FeedNoSend(uint32, float64) - Close() error -} - -type WorkerMeasureConsumer interface { - FeedNoSend(uint32, float64) -} - -type AtreeMeasureConsumer interface { - Feed(uint32, float64) -} - -// - -type _page struct { - PageNo uint32 - Buf []byte - ReferenceCount int -} - -type Atree struct { - indexFile *os.File - dataFile *os.File - mutex sync.Mutex - pages map[uint32]*_page - pageWaits map[uint32][]chan readResult - pagesToRead []uint32 - readSignalCh chan struct{} -} - -type Options struct { - IndexFile *os.File - DataFile *os.File -} - -func New(opt Options) (*Atree, error) { - if opt.IndexFile == nil { - return nil, errors.New("IndexFile option is required") - } - if opt.DataFile == nil { - return nil, errors.New("DataFile option is required") - } - s := &Atree{ - indexFile: opt.IndexFile, - dataFile: opt.DataFile, - pages: make(map[uint32]*_page), - pageWaits: make(map[uint32][]chan readResult), - readSignalCh: make(chan struct{}, 1), - } - return s, nil -} - -func (s *Atree) Run() { - //go s.pageWriter() - go s.pageReader() -} - -// FIND - -func (s *Atree) findDataPage(rootPageNo uint32, timestamp uint32) (uint32, []byte, error) { - // indexPageNo := rootPageNo - // for { - // buf, err := s.fetchIndexPage(indexPageNo) - // if err != nil { - // return 0, nil, fmt.Errorf("fetchIndexPage(%d): %s", indexPageNo, err) - // } - - // foundPageNo := findPageNo(buf, timestamp) - // s.releasePage(indexPageNo) - - // // fix - // if buf[isLastLevelIdx] == 1 { - // buf, err := s.fetchDataPage(foundPageNo) - // if err != nil { - // return 0, nil, fmt.Errorf("fetchDataPage(%d): %s", foundPageNo, err) - // } - // return foundPageNo, buf, nil - // } - // // вглубь - // indexPageNo = foundPageNo - // } - return 0, nil, nil -} - -type PathLeg struct { - PageNo uint32 - Data []byte -} - -type PathToDataPage struct { - Legs []PathLeg - LastPageNo uint32 -} - -func (s *Atree) FindPathToLastPage(rootPageNo uint32) (_ PathToDataPage, err error) { - // var ( - // pageNo = rootPageNo - // legs []PathLeg - // ) - - // for { - // var buf []byte - // buf, err = s.fetchIndexPage(pageNo) - // if err != nil { - // err = fmt.Errorf("FetchIndexPage(%d): %s", pageNo, err) - // return - // } - - // legs = append(legs, PathLeg{ - // PageNo: pageNo, - // Data: buf, - // // childIdx не нужен - // }) - - // foundPageNo := getLastPageNo(buf) - - // // fix - // if buf[isLastLevelIdx] == 1 { - // return PathToDataPage{ - // Legs: legs, - // LastPageNo: foundPageNo, - // }, nil - // } - // // вглубь - // pageNo = foundPageNo - // } - return -} - -// DELETE - -type Level struct { - PageNo uint32 - PageData []byte - Idx int - ChildQty int -} - -func (s *Atree) GetAllPages(rootPageNo uint32) (_ []uint32, err error) { - // var ( - // pageNumbers []uint32 - // levels []*Level - // ) - - // buf, err := s.fetchIndexPage(rootPageNo) - // if err != nil { - // return nil, fmt.Errorf("fetchIndexPage(%d): %s", rootPageNo, err) - // } - // pageNumbers = append(pageNumbers, rootPageNo) - - // // if buf[isDataPageNumbersIdx] == 1 { - // // pageNumbers := listPageNumbers(buf) - // // dataPages = append(dataPages, pageNumbers...) - - // // s.releasePage(rootPageNo) - - // // return PageLists{ - // // DataPages: dataPages, - // // IndexPages: indexPages, - // // }, nil - // // } - - // childQty, _ := bin.GetUint16(buf[indexRecordsQtyIdx:]) - - // levels = append(levels, &Level{ - // PageNo: rootPageNo, - // PageData: buf, - // Idx: 0, - // ChildQty: int(childQty), - // }) - - // for { - // if len(levels) == 0 { - // return pageNumbers, nil - // } - - // lastIdx := len(levels) - 1 - // level := levels[lastIdx] - - // if level.Idx < level.ChildQty { - // pageNo := getPageNo(level.PageData, level.Idx) - // level.Idx++ - - // var buf []byte - // buf, err = s.fetchPage(pageNo) - // if err != nil { - // return nil, fmt.Errorf("fetchPage(%d): %s", pageNo, err) - // } - // pageNumbers = append(pageNumbers, pageNo) - - // if buf[pageTypeIdx] == PageTypeData { - // //pageNumbers := listPageNumbers(buf) - // //dataPages = append(dataPages, pageNumbers...) - // s.releasePage(pageNo) - // } else { - // childQty, _ = bin.GetUint16(buf[indexRecordsQtyIdx:]) - // levels = append(levels, &Level{ - // PageNo: pageNo, - // PageData: buf, - // Idx: 0, - // ChildQty: int(childQty), - // }) - // } - // } else { - // s.releasePage(level.PageNo) - // levels = levels[:lastIdx] - // } - // } - return -} diff --git a/atree/cursor.go b/atree/cursor.go deleted file mode 100644 index ff7e24b..0000000 --- a/atree/cursor.go +++ /dev/null @@ -1,187 +0,0 @@ -package atree - -import ( - "errors" - "fmt" - - "gordenko.dev/dima/bin" - "gordenko.dev/dima/qb" -) - -type BackwardCursor struct { - metricType qb.MetricType - fracDigits byte - atree *Atree - pageNo uint32 - pageData []byte - timestampDecompressor qb.TimestampDecompressor - valueDecompressor qb.ValueDecompressor -} - -type BackwardCursorOptions struct { - MetricType qb.MetricType - FracDigits byte - PageNo uint32 - PageData []byte - Atree *Atree -} - -func NewBackwardCursor(opt BackwardCursorOptions) (*BackwardCursor, error) { - switch opt.MetricType { - case qb.Instant, qb.Cumulative: - // ok - default: - return nil, fmt.Errorf("MetricType option has wrong value: %d", opt.MetricType) - } - if opt.FracDigits > qb.MaxFracDigits { - return nil, errors.New("FracDigits option is required") - } - if opt.Atree == nil { - return nil, errors.New("Atree option is required") - } - if opt.PageNo == 0 { - return nil, errors.New("PageNo option is required") - } - if len(opt.PageData) == 0 { - return nil, errors.New("PageData option is required") - } - - s := &BackwardCursor{ - metricType: opt.MetricType, - fracDigits: opt.FracDigits, - atree: opt.Atree, - pageNo: opt.PageNo, - pageData: opt.PageData, - } - err := s.makeDecompressors() - if err != nil { - return nil, err - } - return s, nil -} - -// timestamp, value, done, error -func (s *BackwardCursor) Prev() (uint32, float64, bool, error) { - var ( - timestamp uint32 - value float64 - //done bool - //err error - ) - - // timestamp, done = s.timestampDecompressor.NextValue() - // if !done { - // value, done = s.valueDecompressor.NextValue() - // if done { - // return 0, 0, false, - // fmt.Errorf("corrupted data page %d: has timestamp, no value", - // s.pageNo) - // } - // return timestamp, value, false, nil - // } - - // prevPageNo, _ := bin.GetUint32(s.pageData[prevPageIdx:]) - // if prevPageNo == 0 { - // return 0, 0, true, nil - // } - // s.atree.releasePage(s.pageNo) - - // s.pageNo = prevPageNo - // s.pageData, err = s.atree.fetchDataPage(s.pageNo) - // if err != nil { - // return 0, 0, false, fmt.Errorf("atree.fetchDataPage(%d): %s", s.pageNo, err) - // } - - // err = s.makeDecompressors() - // if err != nil { - // return 0, 0, false, err - // } - - // timestamp, done = s.timestampDecompressor.NextValue() - // if done { - // return 0, 0, false, - // fmt.Errorf("corrupted data page %d: no timestamps", - // s.pageNo) - // } - // value, done = s.valueDecompressor.NextValue() - // if done { - // return 0, 0, false, - // fmt.Errorf("corrupted data page %d: no values", - // s.pageNo) - // } - return timestamp, value, false, nil -} - -func (s *BackwardCursor) Close() { - s.atree.releasePage(s.pageNo) -} - -// HELPER - -func (s *BackwardCursor) makeDecompressors() error { - timestampsSize, _ := bin.GetUint16(s.pageData[timestampsSizeIdx:]) - valuesSize, _ := bin.GetUint16(s.pageData[valuesSizeIdx:]) - - payloadSize := timestampsSize + valuesSize - - if payloadSize > dataFooterIdx { - return fmt.Errorf("corrupted data page %d: timestamps + values size %d gt payload size", - s.pageNo, payloadSize) - } - - s.timestampDecompressor = enc.NewTimeDeltaDecompressor( - s.pageData[:timestampsSize], - ) - - vbuf := s.pageData[timestampsSize : timestampsSize+valuesSize] - - switch s.metricType { - case qb.Instant: - s.valueDecompressor = enc.NewInstantDeltaDecompressor( - vbuf, s.fracDigits) - - case qb.Cumulative: - s.valueDecompressor = enc.NewCumulativeDeltaDecompressor( - vbuf, s.fracDigits) - - default: - return fmt.Errorf("bug: wrong metricType %d", s.metricType) - } - return nil -} - -// func makeDecompressors(pageData []byte, metricType qb.MetricType, fracDigits byte) ( -// qb.TimestampDecompressor, qb.ValueDecompressor, error, -// ) { -// timestampsSize, _ := bin.GetUint16(pageData[timestampsSizeIdx:]) -// valuesSize, _ := bin.GetUint16(pageData[valuesSizeIdx:]) - -// payloadSize := timestampsSize + valuesSize - -// if payloadSize > dataFooterIdx { -// return nil, nil, fmt.Errorf("corrupted: timestamps + values size %d > payload size", -// payloadSize) -// } - -// timestampDecompressor := enc.NewTimeDeltaDecompressor( -// pageData[:timestampsSize], -// ) - -// vbuf := pageData[timestampsSize : timestampsSize+valuesSize] - -// var valueDecompressor qb.ValueDecompressor -// switch metricType { -// case qb.Instant: -// valueDecompressor = enc.NewInstantDeltaDecompressor( -// vbuf, fracDigits) - -// case qb.Cumulative: -// valueDecompressor = enc.NewCumulativeDeltaDecompressor( -// vbuf, fracDigits) - -// default: -// return nil, nil, fmt.Errorf("bug: wrong metricType %d", metricType) -// } -//return timestampDecompressor, valueDecompressor, nil -// return nil, nil, nil -// } diff --git a/atree/select.go b/atree/select.go deleted file mode 100644 index ddb62d4..0000000 --- a/atree/select.go +++ /dev/null @@ -1,129 +0,0 @@ -package atree - -import ( - "fmt" - - "gordenko.dev/dima/qb" -) - -type ContinueFullScanReq struct { - FracDigits byte - ResponseWriter AtreeMeasureConsumer - LastPageNo uint32 -} - -func (s *Atree) ContinueFullScan(req ContinueFullScanReq) error { - buf, err := s.fetchDataPage(req.LastPageNo) - if err != nil { - return fmt.Errorf("fetchDataPage(%d): %s", req.LastPageNo, err) - } - - treeCursor, err := NewBackwardCursor(BackwardCursorOptions{ - PageNo: req.LastPageNo, - PageData: buf, - Atree: s, - FracDigits: req.FracDigits, - MetricType: qb.Instant, - }) - if err != nil { - return err - } - defer treeCursor.Close() - - for { - timestamp, value, done, err := treeCursor.Prev() - if err != nil { - return err - } - if done { - return nil - } - req.ResponseWriter.Feed(timestamp, value) - } -} - -type ContinueRangeScanReq struct { - FracDigits byte - ResponseWriter AtreeMeasureConsumer - LastPageNo uint32 - Since uint32 -} - -func (s *Atree) ContinueRangeScan(req ContinueRangeScanReq) error { - buf, err := s.fetchDataPage(req.LastPageNo) - if err != nil { - return fmt.Errorf("fetchDataPage(%d): %s", req.LastPageNo, err) - } - - treeCursor, err := NewBackwardCursor(BackwardCursorOptions{ - PageNo: req.LastPageNo, - PageData: buf, - Atree: s, - FracDigits: req.FracDigits, - MetricType: qb.Instant, - }) - if err != nil { - return err - } - defer treeCursor.Close() - - for { - timestamp, value, done, err := treeCursor.Prev() - if err != nil { - return err - } - if done { - return nil - } - req.ResponseWriter.Feed(timestamp, value) - if timestamp < req.Since { - return nil - } - } -} - -type RangeScanReq struct { - FracDigits byte - ResponseWriter AtreeMeasureConsumer - RootPageNo uint32 - Since uint32 - Until uint32 -} - -func (s *Atree) RangeScan(req RangeScanReq) error { - pageNo, buf, err := s.findDataPage(req.RootPageNo, req.Until) - if err != nil { - return err - } - - cursor, err := NewBackwardCursor(BackwardCursorOptions{ - PageNo: pageNo, - PageData: buf, - Atree: s, - FracDigits: req.FracDigits, - MetricType: qb.Instant, - }) - if err != nil { - return err - } - defer cursor.Close() - - for { - timestamp, value, done, err := cursor.Prev() - if err != nil { - return err - } - if done { - return nil - } - - if timestamp <= req.Until { - req.ResponseWriter.Feed(timestamp, value) - - if timestamp < req.Since { - // - записи, удовлетворяющие временным рамкам, закончились. - return nil - } - } - } -} diff --git a/b b/b deleted file mode 100644 index 4fbaf93..0000000 --- a/b +++ /dev/null @@ -1 +0,0 @@ -01 11 27 00 00 01 02 11 4e 46 bd 00 00 00 00 00 00 00 00 00 06 11 27 00 00 00 9c 5e 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 7f 87 04 98 5c dd 69 76 94 80 80 ca 81 20 81 59 82 07 82 56 82 75 83 2a 83 6c 84 11 84 54 85 02 85 0f 85 4e 86 24 87 04 87 55 87 56 87 5a 87 7b 88 01 88 36 89 0a 89 39 8a 02 8a 30 8a 36 8a 52 8a 7d 8b 5b 8b 5c 8c 11 8c 62 8d 22 8d 23 8d 2e 8d 30 8d 52 8d 7e 8e 32 8e 7f 8f 52 90 29 90 55 91 07 91 55 92 16 92 36 93 04 93 07 93 52 94 19 94 3b 94 73 95 2b 95 31 95 74 96 36 96 7c 97 18 97 54 98 28 98 77 99 35 9a 12 9a 4b 9b 04 9b 62 9c 40 9c 4a 9c 61 9d 24 9d 39 9d 78 9e 5c 9f 16 9f 6e a0 13 a0 20 a0 61 a0 7e a1 48 a2 09 a2 1e a2 6f a2 71 a3 55 a3 57 a4 0e a4 2a a4 7c a5 3c a5 74 a6 03 a6 50 a6 61 a6 76 a7 32 a7 54 a7 65 a8 25 a8 60 a9 18 a9 4b a9 77 aa 44 aa 56 ab 27 ab 61 ab 6f ac 35 ac 68 ad 06 ad 26 ad 2a ad 68 ae 35 af 05 af 4b b0 21 b0 53 b0 73 b1 0a b1 44 b1 55 b1 78 b2 4d b3 1f ff b3 76 b4 30 b4 51 b5 31 b5 70 b6 07 b6 52 b7 2e b8 0e b8 34 b9 03 b9 11 b9 26 b9 4d ba 1a ba 7d bb 16 bb 4c bb 79 bc 5a bc 72 bc 79 bd 01 bd 29 be 01 be 13 be 5e bf 34 c0 10 c0 4a c0 50 c1 1a c1 3f c1 64 c1 65 c1 7a c1 7d c2 27 c2 30 c2 73 c3 57 c3 7e c4 41 c4 42 c4 52 c5 00 c5 57 c6 1b c6 49 c6 5f c7 2d c7 73 c8 0b c8 42 c8 6a c9 33 ca 0e ca 20 ca 52 ca 63 cb 32 cb 72 cc 23 cc 3f cc 5a cc 60 cd 41 cd 5e ce 32 ce 4c cf 1d cf 67 d0 11 d0 6f d0 79 d1 2c d1 6a d2 1d d2 2f d2 5e d3 0a d3 12 d3 60 d4 09 d4 52 d4 77 d5 4b d5 62 d5 64 d5 7b d6 1e d6 3d d7 18 d7 67 d8 0f d8 4d d8 73 d9 03 d9 40 da 16 da 45 db 09 db 55 dc 15 dc 2f dc 41 dc 43 dc 6a dc 6d dd 33 dd 42 de 07 de 08 de 35 de 6a df 1f df 20 e0 02 e0 63 e1 29 e1 69 e2 01 e2 32 e2 7e e3 47 e4 19 e4 21 e4 28 ff e4 57 e5 35 e5 6a e6 2d e6 31 e6 67 e7 09 e7 65 e7 6c e8 41 e8 5f e8 71 e9 42 e9 51 ea 32 ea 4f eb 29 eb 33 eb 66 eb 7e ec 4c ec 7f ed 32 ed 3c ed 7f ee 01 ee 1f ee 20 ee 72 ef 39 f0 06 f0 44 f0 75 f1 40 f1 75 f2 2f f2 42 f2 5f f3 37 f3 7b f4 2d f4 3e f4 53 f5 05 f5 11 f5 4f f6 09 f6 43 f6 73 f7 02 f7 5c f8 30 f8 68 f8 79 f9 0e f9 5c f9 5e f9 78 fa 0f fa 4c fb 25 fb 62 fc 3e fd 07 fd 29 fd 7a fe 25 fe 74 ff 1a ff 38 ff 41 ff 7d 81 00 5d 81 00 60 81 01 32 81 01 3a 81 01 50 81 02 08 81 02 38 81 02 5f 81 03 01 81 03 19 81 03 48 81 04 07 81 04 4b 81 04 5e 81 05 05 81 05 21 81 05 28 81 05 47 81 05 5d 81 05 7d 81 06 0c 81 06 5c 81 07 0a 81 07 67 81 07 7b 81 08 2e 81 08 6c 81 08 76 81 09 0a 81 09 4a 81 09 50 81 09 64 81 0a 16 81 0a 56 81 0b 11 81 0b 28 81 0b 52 81 0b 7f 81 0c 0a 81 0c 5d 81 0d 3f 81 0e 0c 81 0e 33 81 0e 69 81 0f 06 81 0f 39 81 0f 3f 81 0f 40 81 0f 41 81 0f 74 81 10 51 81 10 67 81 11 1d 81 12 01 81 12 1a 81 12 3c ff 81 12 4f 81 12 7e 81 13 4c 81 13 56 81 13 6c 81 14 4e 81 14 65 81 15 41 81 15 49 81 16 07 81 16 18 81 16 76 81 17 56 81 18 00 81 18 12 81 18 6e 81 19 4b 81 1a 29 81 1a 33 81 1b 0e 81 1b 46 81 1b 6f 81 1b 78 81 1c 04 81 1c 37 81 1c 40 81 1c 7c 81 1d 12 81 1d 2e 81 1d 56 81 1d 59 81 1e 23 81 1e 46 81 1e 4e 81 1e 7b 81 1f 3d 81 1f 75 81 20 19 81 20 4a 81 20 58 81 21 07 81 21 4a 81 21 64 81 21 77 81 22 1b 81 22 2a 81 22 76 81 23 1f 81 23 40 81 23 66 81 24 0b 81 24 1d 81 24 75 81 25 51 81 25 76 81 26 33 81 26 62 81 27 29 81 27 61 81 28 3b 81 29 15 81 29 45 81 29 6c 81 2a 35 81 2a 7d 81 2b 50 81 2c 07 81 2c 68 81 2c 74 81 2d 4e 81 2e 2e 81 2e 7f 81 2f 07 81 2f 34 81 2f 7b 81 30 0c 81 30 5b 81 31 3e 81 31 60 81 32 2f 81 32 7b 81 33 30 81 34 0a 81 34 11 81 34 6a 81 35 18 81 35 5c 81 35 79 81 36 5a 81 37 0d 81 37 20 81 37 27 81 37 5b 81 38 3b 81 38 73 81 39 16 81 39 65 81 3a 48 81 3a 67 81 3b 24 81 3b 43 81 3c 26 81 3c 6a 81 3d 25 81 3d 4d 81 3d 55 81 3d 5d 81 3e 12 81 3e 24 81 3e 49 81 3f 0a 81 3f 17 81 3f 28 81 40 09 81 40 30 81 41 02 81 41 2a 81 41 52 81 41 64 81 41 75 81 42 33 81 42 59 81 43 1b 81 43 59 81 43 60 81 43 7f 81 44 26 81 44 29 ff 81 44 3b 81 44 68 81 45 16 81 45 5c 81 46 03 81 46 5b 81 46 7f 81 47 06 81 47 25 81 47 79 81 48 41 81 49 1a 81 49 1f 81 49 32 81 49 58 81 4a 35 81 4a 39 81 4a 5b 81 4a 72 81 4b 4e 81 4c 1b 81 4c 48 81 4c 53 81 4c 67 81 4d 00 81 4d 07 81 4d 57 81 4e 39 81 4e 64 81 4e 69 81 4f 04 81 4f 64 81 50 46 81 50 59 81 51 10 81 51 73 81 52 15 81 52 28 81 52 64 81 53 32 81 54 01 81 54 35 81 54 3c 81 54 4d 81 55 07 81 55 64 81 55 7e 81 56 39 81 56 54 81 56 73 81 57 42 81 57 58 81 58 2b 81 58 34 81 59 0f 81 59 5f 81 59 78 81 5a 02 81 5a 48 81 5a 7a 81 5b 44 81 5b 5e 81 5c 0f 81 5c 26 81 5c 6c 81 5d 1c 81 5d 3c 81 5d 5d 81 5e 30 81 5e 3f 81 5e 72 81 5f 06 81 5f 3a 81 5f 7a 81 5f 7f 81 60 0c 81 60 18 81 60 72 81 61 11 81 61 25 81 61 79 81 62 1c 81 62 69 81 63 4d 81 64 20 81 64 2e 81 64 30 81 64 4d 81 64 61 81 65 35 81 66 16 81 66 3e 81 67 11 81 67 20 81 67 7e 81 68 21 81 68 64 81 68 7c 81 69 1c 81 69 1e 81 69 76 81 6a 32 81 6a 76 81 6b 41 81 6b 5a 81 6c 30 81 6c 6d 81 6d 2e 81 6d 79 81 6e 0b 81 6e 29 81 6e 49 81 6e 7a 81 6f 06 81 6f 2a 81 6f 3b 81 6f 77 81 70 40 81 71 15 81 71 17 81 71 5b 81 71 69 81 71 7c 81 72 29 81 73 02 81 73 5e 81 74 36 81 75 08 ff 81 75 2b 81 75 5d 81 76 3c 81 77 02 81 77 64 81 78 20 81 78 51 81 78 79 81 79 2d 81 79 59 81 7a 20 81 7a 5b 81 7b 15 81 7b 22 81 7b 4a 81 7c 0a 81 7c 3b 81 7c 42 81 7d 06 81 7d 0b 81 7d 35 81 7d 7f 81 7e 61 81 7e 72 81 7f 36 81 7f 39 81 7f 71 82 00 1e 82 00 45 82 01 28 82 01 49 82 01 53 82 02 35 82 02 77 82 03 05 82 03 43 82 04 1b 82 04 6b 82 05 4d 82 06 29 82 06 34 82 06 3e 82 07 1d 82 07 29 82 07 3c 82 07 42 82 07 54 82 07 6f 82 08 28 82 08 5c 82 09 22 82 09 52 82 09 7b 82 0a 2a 82 0a 65 82 0a 6c 82 0b 19 82 0b 61 82 0c 0a 82 0c 3d 82 0d 21 82 0d 30 82 0d 71 82 0e 1f 82 0e 48 82 0f 07 82 0f 0f 82 0f 30 82 0f 4b 82 0f 6d 82 10 2e 82 10 67 82 11 10 82 11 40 82 11 7e 82 12 5c 82 13 3e 82 14 11 82 14 5b 82 15 35 82 16 10 82 16 30 82 17 11 82 17 53 82 17 69 82 18 1d 82 18 5b 82 18 5d 82 19 39 82 19 53 82 19 7e 82 1a 1c 82 1a 2c 82 1a 35 82 1a 74 82 1b 34 82 1b 6d 82 1b 75 82 1c 1e 82 1c 51 82 1d 1c 82 1d 65 82 1e 2f 82 1e 41 82 1f 0b 82 1f 5d 82 20 22 82 21 05 82 21 15 82 21 5f 82 21 73 82 22 1e 82 22 62 82 22 6b 82 23 2f 82 24 09 82 24 6c 82 25 32 82 25 7a 82 26 51 82 26 67 82 26 6d 82 27 23 82 27 4b 82 28 18 82 28 64 82 28 6c 82 29 3a ff 82 2a 14 82 2a 3a 82 2b 12 82 2b 72 82 2c 22 82 2c 65 82 2d 06 82 2d 1f 82 2d 4f 82 2e 13 82 2e 1b 82 2e 79 82 2f 4d 82 2f 67 82 30 44 82 30 70 82 31 47 82 32 13 82 32 5e 82 32 7c 82 33 57 82 34 02 82 34 1a 82 34 2b 82 34 37 82 34 6c 82 35 1f 82 35 4a 82 36 0b 82 36 17 82 36 49 82 37 1f 82 38 03 82 38 24 82 38 3f 82 39 01 82 39 12 82 39 20 82 39 24 82 39 59 82 3a 03 82 3a 19 82 3a 74 82 3b 2a 82 3b 37 82 3b 7e 82 3c 34 82 3c 4b 82 3c 69 82 3d 45 82 3e 20 82 3e 49 82 3f 0e 82 3f 1e 82 3f 20 82 3f 7c 82 40 34 82 40 75 82 41 39 82 42 0f 82 42 4a 82 43 0d 82 43 11 82 43 5f 82 43 69 82 44 2a 82 44 61 82 44 7a 82 45 34 82 46 10 82 46 11 82 46 2f 82 46 3d 82 46 61 82 47 07 82 47 0f 82 47 6b 82 48 20 82 48 42 82 48 65 82 48 68 82 49 1c 82 49 45 82 49 56 82 49 74 82 4a 1c 82 4a 65 82 4b 1f 82 4b 28 82 4b 7a 82 4b 7c 82 4c 10 82 4c 17 82 4c 35 82 4d 08 82 4d 69 82 4e 01 82 4e 18 82 4e 22 82 4e 4a 82 4e 72 82 4f 25 82 50 08 82 50 2b 82 50 2e 82 50 3c 82 50 6f 82 50 7e 82 51 2d 82 51 4e 82 52 20 82 52 62 82 53 11 82 53 45 82 54 23 82 54 78 82 55 1a 82 55 1d 82 55 33 82 55 67 82 56 27 82 56 7a 82 57 2f 82 58 12 82 58 55 82 58 68 82 59 10 82 59 6d ff 82 5a 09 82 5a 27 82 5a 2a 82 5b 04 82 5b 51 82 5c 1d 82 5c 42 82 5c 56 82 5d 15 82 5d 52 82 5d 53 82 5d 7b 82 5e 26 82 5e 2d 82 5e 7f 82 5f 46 82 5f 50 82 5f 6e 82 60 3c 82 61 0f 82 61 10 82 61 3b 82 61 4b 82 62 16 82 62 18 82 62 37 82 62 69 82 62 71 82 62 7f 82 63 0e 82 63 50 82 64 0b 82 64 49 82 65 15 82 65 3a 82 65 5a 82 65 71 82 66 3b 82 66 5a 82 66 72 82 67 0b 82 67 6c 82 67 7f 82 68 15 82 68 4d 82 69 1b 82 69 39 82 69 7b 82 6a 52 82 6a 72 82 6b 13 82 6b 3a 82 6b 5b 82 6c 29 82 6c 5d b6 82 6d 0e 00 82 6d 4e 82 6e 19 82 6e 6b 82 6f 42 82 6f 6d 82 70 21 82 70 6d 82 71 36 82 71 79 82 72 23 82 72 73 82 73 20 82 73 7a 82 74 06 82 74 43 82 75 26 82 75 46 82 76 1e 82 76 4f 82 77 26 82 77 4a 82 78 0c 82 78 61 82 79 0f 82 79 4a 82 7a 0c 82 7a 13 82 7a 44 82 7a 48 82 7a 7d 82 7b 16 82 7b 62 82 7b 6e 82 7c 49 82 7c 78 82 7d 50 82 7e 21 82 7e 34 82 7e 54 82 7e 62 82 7e 64 82 7f 24 82 7f 7c 83 00 18 83 00 5e 83 01 27 83 01 51 ae diff --git a/database/api.go b/database/api.go index 3a58038..e2f4610 100644 --- a/database/api.go +++ b/database/api.go @@ -8,10 +8,10 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/bufreader" "gordenko.dev/dima/qb/proto" "gordenko.dev/dima/qb/storage" + "gordenko.dev/dima/qb/timeutil" "gordenko.dev/dima/qb/transform" "gordenko.dev/dima/qb/worker" ) @@ -35,7 +35,6 @@ func reply(conn io.Writer, errcode uint16) { } bin.PutUint16(answer[1:], errcode) } - _, err := conn.Write(answer) if err != nil { return @@ -61,7 +60,6 @@ func (s *Database) handleTCPConn(conn net.Conn) { } func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (err error) { - //fmt.Println("process request") messageType, err := r.ReadByte() if err != nil { if err != io.EOF { @@ -71,8 +69,6 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } } - //fmt.Println("messageType:", messageType) - switch messageType { case proto.TypeGetMetric: req, err := proto.ReadGetMetricReq(r) @@ -111,7 +107,6 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e if err != nil { return fmt.Errorf("proto.ReadAppendMeasuresReq: %s", err) } - //fmt.Println("append measure", req.MetricID, conn.RemoteAddr().String()) if err = s.AppendMeasures(conn, req); err != nil { return fmt.Errorf("AppendMeasures: %s", err) } @@ -162,7 +157,6 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e } case proto.TypeDeleteMeasures: - //fmt.Println("delete metric") req, err := proto.ReadDeleteMeasuresReq(r) if err != nil { return fmt.Errorf("proto.ReadDeleteMeasuresReq: %s", err) @@ -183,7 +177,6 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e if err != nil { return fmt.Errorf("proto.ReadListAllCumulativeMeasuresReq: %s", err) } - if err = s.ListAllCumulativeMeasures(conn, req); err != nil { return fmt.Errorf("ListAllCumulativeMeasures: %s", err) } @@ -198,8 +191,6 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e // API func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { - //fmt.Println("database.AddMetric") - // Валидация if req.MetricID == 0 { return proto.ErrEmptyMetricID } @@ -215,8 +206,6 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { resultCh := make(chan byte, 1) - fmt.Println("add job") - s.workerInbox.Push(worker.AddMetricReq{ MetricID: req.MetricID, MetricType: req.MetricType, @@ -224,20 +213,14 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 { ResultCh: resultCh, }) - fmt.Println("job added") - resultCode := <-resultCh switch resultCode { case worker.Succeed: - fmt.Println("OK") - + // case worker.MetricDuplicate: - //fmt.Println("ErrDuplicate") return proto.ErrDuplicate - default: - //fmt.Println("ErrWrongResultCodeBug") qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } return 0 @@ -267,10 +250,8 @@ func (s *Database) GetMetric(conn io.Writer, req proto.GetMetricReq) error { if err != nil { return err } - case worker.NoMetric: reply(conn, proto.ErrNoMetric) - default: qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } @@ -469,12 +450,11 @@ func (s *Database) ListCumulativeMeasures(conn net.Conn, req proto.ListCumulativ } func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriodsReq) error { - since, until := timeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay) + since, until := timeutil.TimeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay) if since.After(until) { reply(conn, proto.ErrInvalidRange) return nil } - responseWriter, err := transform.NewInstantPeriodsWriter(transform.InstantPeriodsWriterOptions{ Dst: conn, GroupBy: req.GroupBy, @@ -485,7 +465,6 @@ func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriod reply(conn, proto.ErrUnexpected) return nil } - return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, MetricType: qb.Instant, @@ -497,12 +476,11 @@ func (s *Database) ListInstantPeriods(conn net.Conn, req proto.ListInstantPeriod } func (s *Database) ListCumulativePeriods(conn net.Conn, req proto.ListCumulativePeriodsReq) error { - since, until := timeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay) + since, until := timeutil.TimeBoundsOfAggregation(req.Since, req.Until, req.GroupBy, req.FirstHourOfDay) if since.After(until) { reply(conn, proto.ErrInvalidRange) return nil } - responseWriter, err := transform.NewCumulativePeriodsWriter(transform.CumulativePeriodsWriterOptions{ Dst: conn, GroupBy: req.GroupBy, @@ -512,7 +490,6 @@ func (s *Database) ListCumulativePeriods(conn net.Conn, req proto.ListCumulative reply(conn, proto.ErrUnexpected) return nil } - return s.rangeScan(rangeScanReq{ MetricID: req.MetricID, MetricType: qb.Cumulative, @@ -529,7 +506,7 @@ type rangeScanReq struct { Since uint32 Until uint32 Conn io.Writer - ResponseWriter atree.PeriodsWriter + ResponseWriter qb.PeriodsWriter } func (s *Database) rangeScan(req rangeScanReq) error { @@ -549,44 +526,40 @@ func (s *Database) rangeScan(req rangeScanReq) error { switch result.ResultCode { case worker.QueryDone: req.ResponseWriter.Close() - case worker.UntilFound: - // err := s.atree.ContinueRangeScan(atree.ContinueRangeScanReq{ - // FracDigits: result.FracDigits, - // ResponseWriter: req.ResponseWriter, - // LastPageNo: result.LastPageNo, - // Since: req.Since, - // }) - //s.metricRUnlock(req.MetricID) - - // if err != nil { - // reply(req.Conn, proto.ErrUnexpected) - // } else { - // req.ResponseWriter.Close() - // } - + err := s.ContinueRangeScan(ContinueRangeScanReq{ + MetricType: req.MetricType, + FracDigits: result.FracDigits, + ResponseWriter: req.ResponseWriter, + LastPageNo: result.LastPageNo, + Since: req.Since, + }) + //s.metricRUnlock(req.MetricID) fix release unlock + if err != nil { + reply(req.Conn, proto.ErrUnexpected) + } else { + req.ResponseWriter.Close() + } case worker.UntilNotFound: - // err := s.atree.RangeScan(atree.RangeScanReq{ - // FracDigits: result.FracDigits, - // ResponseWriter: req.ResponseWriter, - // RootPageNo: result.RootPageNo, - // Since: req.Since, - // Until: req.Until, - // }) - //s.metricRUnlock(req.MetricID) - - // if err != nil { - // reply(req.Conn, proto.ErrUnexpected) - // } else { - // req.ResponseWriter.Close() - // } - + err := s.RangeScan(RangeScanReq{ + MetricType: req.MetricType, + FracDigits: result.FracDigits, + ResponseWriter: req.ResponseWriter, + Since: req.Since, + Until: req.Until, + LastPageNo: result.LastPageNo, + IsDataPage: result.IsDataPage, + }) + //s.metricRUnlock(req.MetricID) // fix release + if err != nil { + reply(req.Conn, proto.ErrUnexpected) + } else { + req.ResponseWriter.Close() + } case worker.NoMetric: reply(req.Conn, proto.ErrNoMetric) - case worker.WrongMetricType: reply(req.Conn, proto.ErrWrongMetricType) - default: qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } @@ -597,7 +570,7 @@ type fullScanReq struct { MetricID uint32 MetricType qb.MetricType Conn io.Writer - ResponseWriter atree.PeriodsWriter + ResponseWriter qb.PeriodsWriter } func (s *Database) fullScan(req fullScanReq) error { @@ -614,11 +587,10 @@ func (s *Database) fullScan(req fullScanReq) error { switch result.ResultCode { case worker.QueryDone: - fmt.Printf("query done") req.ResponseWriter.Close() - case worker.UntilFound: - err := s.atree.ContinueFullScan(atree.ContinueFullScanReq{ + err := s.ContinueFullScan(ContinueFullScanReq{ + MetricType: req.MetricType, FracDigits: result.FracDigits, ResponseWriter: req.ResponseWriter, LastPageNo: result.LastPageNo, @@ -629,13 +601,10 @@ func (s *Database) fullScan(req fullScanReq) error { } else { req.ResponseWriter.Close() } - case worker.NoMetric: reply(req.Conn, proto.ErrNoMetric) - case worker.WrongMetricType: reply(req.Conn, proto.ErrWrongMetricType) - default: qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug) } diff --git a/database/database.go b/database/database.go index 11d466c..7fa66c1 100644 --- a/database/database.go +++ b/database/database.go @@ -10,8 +10,8 @@ import ( "time" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/inbox" + "gordenko.dev/dima/qb/pagecache" "gordenko.dev/dima/qb/recovery" "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/worker" @@ -24,7 +24,8 @@ type Database struct { workerInbox *inbox.Inbox worker *worker.Worker storage *storage.Writer - atree *atree.Atree + indexCache *pagecache.PageCache + dataCache *pagecache.PageCache tcpPort int logfile *os.File logger *log.Logger @@ -100,6 +101,24 @@ func New(opt Options) (_ *Database, err error) { return nil, err } + s.indexCache, err = pagecache.New(pagecache.Options{ + File: indexFile, + PageSize: storage.IndexPageSize, + VerifyPageCRC: storage.VerifyIndexPageCRC32, + }) + if err != nil { + return nil, fmt.Errorf("index pagecache.New: %s", err) + } + + s.dataCache, err = pagecache.New(pagecache.Options{ + File: dataFile, + PageSize: storage.DataPageSize, + VerifyPageCRC: storage.VerifyDataPageCRC32, + }) + if err != nil { + return nil, fmt.Errorf("data pagecache.New: %s", err) + } + s.storage, err = storage.NewWriter(storage.WriterOptions{ Inbox: storageInbox, WorkerInbox: s.workerInbox, @@ -126,14 +145,6 @@ func New(opt Options) (_ *Database, err error) { ExitCh: opt.ExitCh, WaitGroup: opt.WaitGroup, }) - - s.atree, err = atree.New(atree.Options{ - DataFile: dataFile, - IndexFile: indexFile, - }) - if err != nil { - return nil, fmt.Errorf("atree.New: %s", err) - } return s, nil } @@ -146,14 +157,11 @@ func (s *Database) ListenAndServe() (err error) { //s.waitGroup.Add(1) go s.storage.Run() - // s.atree.Run() - s.waitGroup.Add(1) go s.worker.Run() s.logger.Println("database started") for { - // Listen for an incoming connection. conn, err := listener.Accept() if err != nil { s.logger.Printf("listener.Accept: %s\n", err) @@ -164,53 +172,291 @@ func (s *Database) ListenAndServe() (err error) { } } -// +type ContinueFullScanReq struct { + MetricType qb.MetricType + FracDigits byte + ResponseWriter qb.AtreeMeasureConsumer + LastPageNo uint32 +} -// зробити object? +func (s *Database) ContinueFullScan(req ContinueFullScanReq) error { + buf, err := s.dataCache.FetchPage(req.LastPageNo) + if err != nil { + return fmt.Errorf("dataCache.FetchPage(%d): %s", req.LastPageNo, err) + } + treeCursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ + MetricType: req.MetricType, + FracDigits: req.FracDigits, + PageNo: req.LastPageNo, + PageData: buf, + FetchDataPage: s.dataCache.FetchPage, + ReleasePage: s.dataCache.ReleasePage, + }) + if err != nil { + return err + } + defer treeCursor.Close() + for { + timestamp, value, done, err := treeCursor.Prev() + if err != nil { + return err + } + if done { + return nil + } + req.ResponseWriter.Feed(timestamp, value) + } +} -// func (s *Database) replayChanges(snapshotNumber int) (err error) { -// snapshot, err := readSnapshot(JoinSnapshotFileName(s.dir, snapshotNumber)) -// if err != nil { -// return -// } +type ContinueRangeScanReq struct { + MetricType qb.MetricType + FracDigits byte + ResponseWriter qb.AtreeMeasureConsumer + LastPageNo uint32 + Since uint32 +} -// return nil +func (s *Database) ContinueRangeScan(req ContinueRangeScanReq) error { + buf, err := s.dataCache.FetchPage(req.LastPageNo) + if err != nil { + return fmt.Errorf("dataCache.FetchPage(%d): %s", req.LastPageNo, err) + } + treeCursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ + MetricType: req.MetricType, + FracDigits: req.FracDigits, + PageNo: req.LastPageNo, + PageData: buf, + FetchDataPage: s.dataCache.FetchPage, + ReleasePage: s.dataCache.ReleasePage, + }) + if err != nil { + return err + } + defer treeCursor.Close() + for { + timestamp, value, done, err := treeCursor.Prev() + if err != nil { + return err + } + if done { + return nil + } + req.ResponseWriter.Feed(timestamp, value) + if timestamp < req.Since { + return nil + } + } +} + +type RangeScanReq struct { + MetricType qb.MetricType + FracDigits byte + ResponseWriter qb.AtreeMeasureConsumer + Since uint32 + Until uint32 + LastPageNo uint32 + IsDataPage bool +} + +func (s *Database) RangeScan(req RangeScanReq) error { + var ( + pageNo uint32 + buf []byte + err error + ) + if req.IsDataPage { + pageNo = req.LastPageNo + buf, err = s.dataCache.FetchPage(pageNo) + if err != nil { + return fmt.Errorf("dataCache.FetchPage(%d): %s", pageNo, err) + } + } else { + pageNo, buf, err = s.findDataPage(req.LastPageNo, req.Until) + if err != nil { + return err + } + } + cursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ + MetricType: req.MetricType, + FracDigits: req.FracDigits, + PageNo: pageNo, + PageData: buf, + FetchDataPage: s.dataCache.FetchPage, + ReleasePage: s.dataCache.ReleasePage, + }) + if err != nil { + return err + } + defer cursor.Close() + for { + timestamp, value, done, err := cursor.Prev() + if err != nil { + return err + } + if done { + return nil + } + if timestamp <= req.Until { + req.ResponseWriter.Feed(timestamp, value) + if timestamp < req.Since { + // - записи, удовлетворяющие временным рамкам, закончились. + return nil + } + } + } +} + +// (records payload only, isZeroLevel, timestamp) +func (s *Database) findDataPage(foundPageNo uint32, timestamp uint32) (uint32, []byte, error) { + for { + buf, err := s.indexCache.FetchPage(foundPageNo) + if err != nil { + return 0, nil, fmt.Errorf("fetchIndexPage(%d): %s", foundPageNo, err) + } + toReleaseIndexPageNo := foundPageNo + foundPageNo = storage.FindPageOnIndexPage(buf, timestamp) + s.indexCache.ReleasePage(toReleaseIndexPageNo) + if storage.IsZeroLevelPage(buf) { + buf, err := s.dataCache.FetchPage(foundPageNo) + if err != nil { + return 0, nil, fmt.Errorf("fetchDataPage(%d): %s", foundPageNo, err) + } + return foundPageNo, buf, nil + } + } +} + +// type PathLeg struct { +// PageNo uint32 +// Data []byte // } -//func (s *Database) verifySnapshot(fileName string) (_ bool, err error) { -// file, err := os.Open(fileName) -// if err != nil { -// return -// } -// defer file.Close() - -// stat, err := file.Stat() -// if err != nil { -// return -// } - -// if stat.Size() <= 4 { -// return false, nil -// } - -// var ( -// payloadSize = stat.Size() - 4 -// hash = crc32.NewIEEE() -// ) - -// _, err = io.CopyN(hash, file, payloadSize) -// if err != nil { -// return -// } -// calculatedCRC := hash.Sum32() - -// storedCRC, err := bin.ReadUint32(file) -// if err != nil { -// return -// } -// if storedCRC != calculatedCRC { -// return false, fmt.Errorf("strored CRC %d not equal calculated CRC %d", -// storedCRC, calculatedCRC) -// } -// return true, nil +// type PathToDataPage struct { +// Legs []PathLeg +// LastPageNo uint32 +// } + +// func (s *Database) FindPathToLastPage(rootPageNo uint32) (_ PathToDataPage, err error) { +// // var ( +// // pageNo = rootPageNo +// // legs []PathLeg +// // ) + +// // for { +// // var buf []byte +// // buf, err = s.fetchIndexPage(pageNo) +// // if err != nil { +// // err = fmt.Errorf("FetchIndexPage(%d): %s", pageNo, err) +// // return +// // } + +// // legs = append(legs, PathLeg{ +// // PageNo: pageNo, +// // Data: buf, +// // // childIdx не нужен +// // }) + +// // foundPageNo := getLastPageNo(buf) + +// // // fix +// // if buf[isLastLevelIdx] == 1 { +// // return PathToDataPage{ +// // Legs: legs, +// // LastPageNo: foundPageNo, +// // }, nil +// // } +// // // вглубь +// // pageNo = foundPageNo +// // } +// return +// } + +// DELETE + +// func (s *Atree) DeletePages(pageNumbers []uint32) { +// s.mutex.Lock() +// for _, pageNo := range pageNumbers { +// delete(s.pages, pageNo) +// } +// s.mutex.Unlock() +// } + +// type Level struct { +// PageNo uint32 +// PageData []byte +// Idx int +// ChildQty int +// } + +// func (s *Atree) GetAllPages(rootPageNo uint32) (_ []uint32, err error) { +// // var ( +// // pageNumbers []uint32 +// // levels []*Level +// // ) + +// // buf, err := s.fetchIndexPage(rootPageNo) +// // if err != nil { +// // return nil, fmt.Errorf("fetchIndexPage(%d): %s", rootPageNo, err) +// // } +// // pageNumbers = append(pageNumbers, rootPageNo) + +// // // if buf[isDataPageNumbersIdx] == 1 { +// // // pageNumbers := listPageNumbers(buf) +// // // dataPages = append(dataPages, pageNumbers...) + +// // // s.releasePage(rootPageNo) + +// // // return PageLists{ +// // // DataPages: dataPages, +// // // IndexPages: indexPages, +// // // }, nil +// // // } + +// // childQty, _ := bin.GetUint16(buf[indexRecordsQtyIdx:]) + +// // levels = append(levels, &Level{ +// // PageNo: rootPageNo, +// // PageData: buf, +// // Idx: 0, +// // ChildQty: int(childQty), +// // }) + +// // for { +// // if len(levels) == 0 { +// // return pageNumbers, nil +// // } + +// // lastIdx := len(levels) - 1 +// // level := levels[lastIdx] + +// // if level.Idx < level.ChildQty { +// // pageNo := getPageNo(level.PageData, level.Idx) +// // level.Idx++ + +// // var buf []byte +// // buf, err = s.fetchPage(pageNo) +// // if err != nil { +// // return nil, fmt.Errorf("fetchPage(%d): %s", pageNo, err) +// // } +// // pageNumbers = append(pageNumbers, pageNo) + +// // if buf[pageTypeIdx] == PageTypeData { +// // //pageNumbers := listPageNumbers(buf) +// // //dataPages = append(dataPages, pageNumbers...) +// // s.releasePage(pageNo) +// // } else { +// // childQty, _ = bin.GetUint16(buf[indexRecordsQtyIdx:]) +// // levels = append(levels, &Level{ +// // PageNo: pageNo, +// // PageData: buf, +// // Idx: 0, +// // ChildQty: int(childQty), +// // }) +// // } +// // } else { +// // s.releasePage(level.PageNo) +// // levels = levels[:lastIdx] +// // } +// // } +// return // } diff --git a/database/database_test.go b/database/database_test.go deleted file mode 100644 index a50df16..0000000 --- a/database/database_test.go +++ /dev/null @@ -1,24 +0,0 @@ -package database - -// func TestComposeHeadIndexPage(t *testing.T) { -// var ( -// levelIdx = 0 -// level = storage.IndexLevelTail{ -// Buffer: []byte{ -// 1, 2, 3, -// }, -// RecordsCount: 2, -// } -// head = &storage.IndexPageTail{ -// PageNo: 100, -// CRC32: 12345, -// Records: []byte{}, -// } -// ) -// page := composeHeadIndexPage(levelIdx, level, head) -// if page.PageNo != head.PageNo { -// t.Fatalf("PageNo: got %d are not equal expected %v", -// page.PageNo, head.PageNo) -// } -// // fix compare pages -// } diff --git a/database/helpers.go b/database/helpers.go deleted file mode 100644 index f0a546b..0000000 --- a/database/helpers.go +++ /dev/null @@ -1,52 +0,0 @@ -package database - -import ( - "errors" - "io/fs" - "os" - "time" - - "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/proto" -) - -func timeBoundsOfAggregation(since, until proto.TimeBound, groupBy qb.GroupBy, firstHourOfDay int) (s time.Time, u time.Time) { - switch groupBy { - case qb.GroupByHour, qb.GroupByDay: - s = time.Date(since.Year, since.Month, since.Day, 0, 0, 0, 0, time.Local) - u = time.Date(until.Year, until.Month, until.Day, 0, 0, 0, 0, time.Local) - - case qb.GroupByMonth: - s = time.Date(since.Year, since.Month, 1, 0, 0, 0, 0, time.Local) - u = time.Date(until.Year, until.Month, 1, 0, 0, 0, 0, time.Local) - } - - if firstHourOfDay > 0 { - duration := time.Duration(firstHourOfDay) * time.Hour - s = s.Add(duration) - u = u.Add(duration) - } - - u = u.Add(-1 * time.Second) - return -} - -func isFileExist(fileName string) (bool, error) { - _, err := os.Stat(fileName) - if err != nil { - if errors.Is(err, fs.ErrNotExist) { - return false, nil - } else { - return false, err - } - } else { - return true, nil - } -} - -func correctToFHD(since, until uint32, firstHourOfDay int) (uint32, uint32) { - duration := time.Duration(firstHourOfDay) * time.Hour - since = uint32(time.Unix(int64(since), 0).Add(duration).Unix()) - until = uint32(time.Unix(int64(until), 0).Add(duration).Unix()) - return since, until -} diff --git a/enc/time_delta.go b/enc/time_delta.go index 9d83a70..7a00268 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -310,10 +310,33 @@ func (s *TimeDeltaCompressor) FirstTimestamp() uint32 { return timestamp } -func (s *TimeDeltaCompressor) LastTimestamp() uint32 { +func (s *TimeDeltaCompressor) Until() uint32 { return s.lastUnixtime } +func (s *TimeDeltaCompressor) CommitedSince() uint32 { + if s.state == nil { + if s.pos < len(s.buf) { + timestamp, _ := bin.GetUint32(s.buf[len(s.buf)-4:]) + return timestamp + } + } else { + payload := s.state.Payload + if len(payload) > 0 { + timestamp, _ := bin.GetUint32(payload[len(payload)-4:]) + return timestamp + } + } + return 0 +} + +func (s *TimeDeltaCompressor) CommitedUntil() uint32 { + if s.state == nil { + return s.lastUnixtime + } + return s.state.LastUnixtime +} + // DECOMPRESSOR type TimeDeltaDecompressor struct { diff --git a/inbox/inbox.go b/inbox/inbox.go index 3e43068..8b0054f 100644 --- a/inbox/inbox.go +++ b/inbox/inbox.go @@ -22,11 +22,9 @@ func (s *Inbox) Ready() chan struct{} { } func (s *Inbox) Push(x any) { - //fmt.Printf("inbox.Push: %#v\n", x) s.mutex.Lock() s.items = append(s.items, x) s.mutex.Unlock() - //fmt.Printf("inbox.Pushed\n") select { case s.signalCh <- struct{}{}: default: @@ -34,11 +32,9 @@ func (s *Inbox) Push(x any) { } func (s *Inbox) Drain() []any { - //fmt.Printf("inbox.Drain\n") s.mutex.Lock() items := s.items s.items = nil s.mutex.Unlock() - //fmt.Printf("inbox.Drained: %#v\n", items) return items } diff --git a/atree/io.go b/pagecache/pagecache.go similarity index 53% rename from atree/io.go rename to pagecache/pagecache.go index 69fa4cc..da2fc32 100644 --- a/atree/io.go +++ b/pagecache/pagecache.go @@ -1,12 +1,12 @@ -package atree +package pagecache import ( + "errors" "fmt" + "os" + "sync" - bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/storage" - "gordenko.dev/dima/qb/util" ) type readResult struct { @@ -16,35 +16,52 @@ type readResult struct { // INDEX PAGES -func (s *Atree) DeletePages(pageNumbers []uint32) { - s.mutex.Lock() - for _, pageNo := range pageNumbers { - delete(s.pages, pageNo) +type _page struct { + PageNo uint32 + Buf []byte + ReferenceCount int +} + +type PageCache struct { + mutex sync.Mutex + pageSize int + verifyPageCRC func([]byte) error + file *os.File + pages map[uint32]*_page + pageWaits map[uint32][]chan readResult + pagesToRead []uint32 + readSignalCh chan struct{} +} + +type Options struct { + File *os.File + PageSize int + VerifyPageCRC func([]byte) error +} + +func New(opt Options) (*PageCache, error) { + if opt.File == nil { + return nil, errors.New("File option is required") } - s.mutex.Unlock() -} - -func (s *Atree) fetchIndexPage(pageNo uint32) ([]byte, error) { - // buf, err := s.fetchPage(pageNo) - // if err != nil { - // return nil, err - // } - // if buf[pageTypeIdx] != PageTypeIndex { - // return nil, fmt.Errorf("wrong pageType %d instead of %d", buf[pageTypeIdx], PageTypeIndex) - // } - // return buf, nil - return nil, nil -} - -func (s *Atree) fetchDataPage(pageNo uint32) ([]byte, error) { - buf, err := s.fetchPage(pageNo) - if err != nil { - return nil, err + if opt.PageSize <= 0 { + return nil, errors.New("PageSize option is required") } - return buf, nil + if opt.VerifyPageCRC == nil { + return nil, errors.New("VerifyPageCRC option is required") + } + s := &PageCache{ + file: opt.File, + pageSize: opt.PageSize, + verifyPageCRC: opt.VerifyPageCRC, + pages: make(map[uint32]*_page), + pageWaits: make(map[uint32][]chan readResult), + readSignalCh: make(chan struct{}, 1), + } + go s.pageReader() + return s, nil } -func (s *Atree) fetchPage(pageNo uint32) ([]byte, error) { +func (s *PageCache) FetchPage(pageNo uint32) ([]byte, error) { s.mutex.Lock() p, ok := s.pages[pageNo] if ok { @@ -69,12 +86,12 @@ func (s *Atree) fetchPage(pageNo uint32) ([]byte, error) { result := <-resultCh if result.Err == nil { - result.Err = s.verifyCRC(result.Data, storage.DataPageSize) + result.Err = s.verifyPageCRC(result.Data) } return result.Data, result.Err } -func (s *Atree) releasePage(pageNo uint32) { +func (s *PageCache) ReleasePage(pageNo uint32) { s.mutex.Lock() defer s.mutex.Unlock() @@ -93,11 +110,7 @@ func (s *Atree) releasePage(pageNo uint32) { } } -// DATA PAGES - -// READ - -func (s *Atree) pageReader() { +func (s *PageCache) pageReader() { for { select { case <-s.readSignalCh: @@ -106,7 +119,7 @@ func (s *Atree) pageReader() { } } -func (s *Atree) readPages() { +func (s *PageCache) readPages() { s.mutex.Lock() if len(s.pagesToRead) == 0 { s.mutex.Unlock() @@ -117,17 +130,15 @@ func (s *Atree) readPages() { s.mutex.Unlock() for _, pageNo := range pagesToRead { - buf := make([]byte, storage.DataPageSize) - off := int(pageNo-1) * storage.DataPageSize + buf := make([]byte, s.pageSize) + off := int(pageNo-1) * s.pageSize n, err := s.file.ReadAt(buf, int64(off)) - if n != storage.DataPageSize { - err = fmt.Errorf("read %d instead of %d", n, storage.DataPageSize) + if n != s.pageSize { + err = fmt.Errorf("read %d instead of %d", n, s.pageSize) } - s.mutex.Lock() resultChannels := s.pageWaits[pageNo] delete(s.pageWaits, pageNo) - if err != nil { s.mutex.Unlock() for _, resultCh := range resultChannels { @@ -150,18 +161,3 @@ func (s *Atree) readPages() { } } } - -// WRITE - -func (s *Atree) verifyCRC(data []byte, pageSize int) error { - var ( - pos = pageSize - 4 - calculatedCRC = util.CalculateCRC32(data[:pos]) - storedCRC, _ = bin.GetUint32(data[pos:]) - ) - if calculatedCRC != storedCRC { - return fmt.Errorf("calculatedCRC %d not equal storedCRC %d", - calculatedCRC, storedCRC) - } - return nil -} diff --git a/qb.go b/qb.go index efc6c36..a816c37 100644 --- a/qb.go +++ b/qb.go @@ -24,6 +24,20 @@ const ( AggregateAvg byte = 4 ) +type AtreeMeasureConsumer interface { + Feed(uint32, float64) +} + +type PeriodsWriter interface { + Feed(uint32, float64) + FeedNoSend(uint32, float64) + Close() error +} + +type WorkerMeasureConsumer interface { + FeedNoSend(uint32, float64) +} + type TimestampCompressor interface { // (tmp, timestamp) Evaluate([]byte, uint32) TimeEvaluationReport @@ -37,12 +51,11 @@ type TimestampCompressor interface { // (offset) => payload Tail(int) []byte CreateDecompressor() TimestampDecompressor - //Payload() []byte // для снапшота - // Offset() int ReplaceBuffer([]byte) WriteCommitedTo(io.Writer) error - FirstTimestamp() uint32 - LastTimestamp() uint32 + Until() uint32 + CommitedSince() uint32 + CommitedUntil() uint32 ReplaceSinceWithUntil() uint32 } @@ -67,7 +80,6 @@ type ValueCompressor interface { Append(int, []byte, float64, uint64) Size() int CommitedSize() int - //Chunks() [][]byte //DeleteLast() CaptureState() ForgetCapturedState() @@ -75,8 +87,6 @@ type ValueCompressor interface { Tail(int) []byte // fracDigits CreateDecompressor(MetricType, byte) ValueDecompressor - //Payload() []byte // для снапшота - //Offset() int ReplaceBuffer([]byte) WriteCommitedTo(io.Writer) error LastValue() float64 diff --git a/recovery/test.wal_0 b/recovery/test.wal_0 deleted file mode 100644 index e69de29..0000000 diff --git a/recovery/test.wal_123 b/recovery/test.wal_123 deleted file mode 100644 index e69de29..0000000 diff --git a/storage/cursor.go b/storage/cursor.go new file mode 100644 index 0000000..8caded6 --- /dev/null +++ b/storage/cursor.go @@ -0,0 +1,136 @@ +package storage + +import ( + "errors" + "fmt" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" +) + +type BackwardCursor struct { + metricType qb.MetricType + fracDigits byte + fetchDataPage func(uint32) ([]byte, error) + releasePage func(uint32) + pageNo uint32 + pageData []byte + timestampDecompressor qb.TimestampDecompressor + valueDecompressor qb.ValueDecompressor +} + +type BackwardCursorOptions struct { + MetricType qb.MetricType + FracDigits byte + PageNo uint32 + PageData []byte + FetchDataPage func(uint32) ([]byte, error) + ReleasePage func(uint32) +} + +func NewBackwardCursor(opt BackwardCursorOptions) (*BackwardCursor, error) { + switch opt.MetricType { + case qb.Instant, qb.Cumulative: + // ok + default: + return nil, fmt.Errorf("MetricType option has wrong value: %d", opt.MetricType) + } + if opt.FracDigits > qb.MaxFracDigits { + return nil, errors.New("FracDigits option is required") + } + if opt.FetchDataPage == nil { + return nil, errors.New("FetchDataPage option is required") + } + if opt.ReleasePage == nil { + return nil, errors.New("ReleasePage option is required") + } + if opt.PageNo == 0 { + return nil, errors.New("PageNo option is required") + } + if len(opt.PageData) == 0 { + return nil, errors.New("PageData option is required") + } + s := &BackwardCursor{ + metricType: opt.MetricType, + fracDigits: opt.FracDigits, + fetchDataPage: opt.FetchDataPage, + releasePage: opt.ReleasePage, + pageNo: opt.PageNo, + pageData: opt.PageData, + } + var err error + s.timestampDecompressor, err = CreateTimestampDecompressor(s.pageData) + if err != nil { + return nil, fmt.Errorf("CreateTimestampDecompressor: %s", err) + } + s.valueDecompressor, err = CreateValueDecompressor(s.pageData, s.metricType, s.fracDigits) + if err != nil { + return nil, fmt.Errorf("CreateValueDecompressor: %s", err) + } + return s, nil +} + +// timestamp, value, done, error +func (s *BackwardCursor) Prev() (uint32, float64, bool, error) { + var ( + timestamp uint32 + value float64 + done bool + err error + ) + timestamp, done = s.timestampDecompressor.NextValue() + if !done { + value, done = s.valueDecompressor.NextValue() + if done { + return 0, 0, false, + fmt.Errorf("corrupted data page %d: has timestamp, no value", + s.pageNo) + } + return timestamp, value, false, nil + } + + prevPageNo, _ := bin.GetUint32(s.pageData[prevPageIdx:]) + if prevPageNo == 0 { + return 0, 0, true, nil + } + s.releasePage(s.pageNo) + + s.pageNo = prevPageNo + s.pageData, err = s.fetchDataPage(s.pageNo) + if err != nil { + return 0, 0, false, + fmt.Errorf("fetchDataPage(%d): %s", s.pageNo, err) + } + s.timestampDecompressor, err = CreateTimestampDecompressor(s.pageData) + if err != nil { + return 0, 0, false, + fmt.Errorf("CreateTimestampDecompressor: %s", err) + } + s.valueDecompressor, err = CreateValueDecompressor(s.pageData, s.metricType, s.fracDigits) + if err != nil { + return 0, 0, false, + fmt.Errorf("CreateValueDecompressor: %s", err) + } + // + timestamp, done = s.timestampDecompressor.NextValue() + if done { + return 0, 0, false, + fmt.Errorf("corrupted data page %d: no timestamps", s.pageNo) + } + value, done = s.valueDecompressor.NextValue() + if done { + return 0, 0, false, + fmt.Errorf("corrupted data page %d: no values", s.pageNo) + } + return timestamp, value, false, nil +} + +func (s *BackwardCursor) Close() { + s.releasePage(s.pageNo) +} + +// HELPER + +//func (s *BackwardCursor) makeDecompressors() error { + +//} diff --git a/storage/misc.go b/storage/misc.go index c07b28f..40d4302 100644 --- a/storage/misc.go +++ b/storage/misc.go @@ -6,87 +6,50 @@ import ( bin "gordenko.dev/dima/bin/little" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/enc" + "gordenko.dev/dima/qb/util" ) -// func (s *BackwardCursor) makeDecompressors() error { -// timestampsSize, _ := bin.GetUint16(s.pageData[timestampsSizeIdx:]) -// valuesSize, _ := bin.GetUint16(s.pageData[valuesSizeIdx:]) - -// payloadSize := timestampsSize + valuesSize - -// if payloadSize > dataFooterIdx { -// return fmt.Errorf("corrupted data page %d: timestamps + values size %d gt payload size", -// s.pageNo, payloadSize) -// } - -// s.timestampDecompressor = enc.NewTimeDeltaDecompressor( -// s.pageData[:timestampsSize], -// ) - -// vbuf := s.pageData[timestampsSize : timestampsSize+valuesSize] - -// switch s.metricType { -// case qb.Instant: -// s.valueDecompressor = enc.NewInstantDeltaDecompressor( -// vbuf, s.fracDigits) - -// case qb.Cumulative: -// s.valueDecompressor = enc.NewCumulativeDeltaDecompressor( -// vbuf, s.fracDigits) - -// default: -// return fmt.Errorf("bug: wrong metricType %d", s.metricType) -// } -// return nil -// } - -func makeDecompressors(pageData []byte, metricType qb.MetricType, fracDigits byte) ( - qb.TimestampDecompressor, qb.ValueDecompressor, error, -) { - - valuesSize, _ := bin.GetUint16(pageData[valuesSizeIdx:]) - - payloadSize := timestampsSize + valuesSize - - if payloadSize > dataFooterIdx { - return nil, nil, fmt.Errorf("corrupted: timestamps + values size %d > payload size", - payloadSize) - } - - timestampDecompressor := enc.NewTimeDeltaDecompressor( - pageData[:timestampsSize], - ) - - vbuf := pageData[timestampsSize : timestampsSize+valuesSize] - - var valueDecompressor qb.ValueDecompressor - switch metricType { - case qb.Instant: - valueDecompressor = enc.NewInstantDeltaDecompressor( - vbuf, fracDigits) - - case qb.Cumulative: - valueDecompressor = enc.NewCumulativeDeltaDecompressor( - vbuf, fracDigits) - - default: - return nil, nil, fmt.Errorf("bug: wrong metricType %d", metricType) - } - return timestampDecompressor, valueDecompressor, nil - return nil, nil, nil -} - -func CreateTimeDeltaDecompressor(page []byte) qb.TimestampDecompressor { +func CreateTimestampDecompressor(page []byte) (qb.TimestampDecompressor, error) { size, _ := bin.GetUint16(page[timestampsSizeIdx:]) + if size > DataPageFooterSize { + return nil, fmt.Errorf("bug: invalid timestamps size %d", size) + } pos := DataPagePayloadSize - int(size) d := enc.NewTimeDeltaDecompressor() d.RestoreFromEnd(page[pos:DataPagePayloadSize]) - return d + return d, nil } -func CreateValueDeltaDecompressor(page []byte, metricType qb.MetricType, fracDigits byte) qb.ValueDecompressor { +func CreateValueDecompressor(page []byte, metricType qb.MetricType, fracDigits byte) (qb.ValueDecompressor, error) { size, _ := bin.GetUint16(page[valuesSizeIdx:]) + if size > DataPageFooterSize { + return nil, fmt.Errorf("bug: invalid timestamps size %d", size) + } d := enc.NewValueDeltaDecompressor(metricType, fracDigits) d.RestoreFromEnd(page[:size]) - return d + return d, nil +} + +func VerifyDataPageCRC32(data []byte) error { + var ( + calculatedCRC = util.CalculateCRC32(data[:dataCRC32Idx]) + writtenCRC, _ = bin.GetUint32(data[dataCRC32Idx:]) + ) + if calculatedCRC != writtenCRC { + return fmt.Errorf("calculated CRC32 %d are not equal written CRC32 %d", + calculatedCRC, writtenCRC) + } + return nil +} + +func VerifyIndexPageCRC32(data []byte) error { + var ( + calculatedCRC = util.CalculateCRC32(data[:indexCRC32Idx]) + writtenCRC, _ = bin.GetUint32(data[indexCRC32Idx:]) + ) + if calculatedCRC != writtenCRC { + return fmt.Errorf("calculated CRC32 %d are not equal written CRC32 %d", + calculatedCRC, writtenCRC) + } + return nil } diff --git a/atree/misc.go b/storage/navigation.go similarity index 61% rename from atree/misc.go rename to storage/navigation.go index d8c1fb7..b24a0a8 100644 --- a/atree/misc.go +++ b/storage/navigation.go @@ -1,4 +1,10 @@ -package atree +package storage + +import bin "gordenko.dev/dima/bin/little" + +const ( + PageNoSize = 4 +) type KeyComparator interface { CompareTo(int) int @@ -10,19 +16,18 @@ type ValueAtComparator struct { } func (s ValueAtComparator) CompareTo(elemIdx int) int { - // var ( - // pos = elemIdx * timestampSize - // elem, _ = bin.GetUint32(s.buf[pos:]) - // ) + var ( + pos = elemIdx * IndexRecordSize + elem, _ = bin.GetUint32(s.buf[pos:]) + ) - // if s.timestamp < elem { - // return -1 - // } else if s.timestamp > elem { - // return 1 - // } else { - // return 0 - // } - return 213131132 + if s.timestamp < elem { + return -1 + } else if s.timestamp > elem { + return 1 + } else { + return 0 + } } func BinarySearch(qty int, keyComparator KeyComparator) (elemIdx int, isFound bool) { @@ -61,18 +66,63 @@ func BinarySearch(qty int, keyComparator KeyComparator) (elemIdx int, isFound bo // return pageNo // } -// func findPageNo(buf []byte, timestamp uint32) (pageNo uint32) { -// comparator := ValueAtComparator{ -// buf: buf, -// timestamp: timestamp, -// } -// qty, _ := bin.GetUint16(buf[indexRecordsQtyIdx:]) -// elemIdx, _ := BinarySearch(int(qty), comparator) -// pos := indexFooterIdx - (elemIdx+1)*PageNoSize -// pageNo, _ = bin.GetUint32(buf[pos:]) +// func GetIndexRecordsSince(buf []byte) (pageNo uint32) { +// pageNo, _ = bin.GetUint32(buf) // return // } +func FindPageOnIndexLevelTail(level IndexLevelTail, timestamp uint32) (pageNo uint32) { + comparator := ValueAtComparator{ + buf: level.Buffer, + timestamp: timestamp, + } + elemIdx, _ := BinarySearch(level.RecordsCount, comparator) + pos := elemIdx*IndexRecordSize + 4 // timestamp size + pageNo, _ = bin.GetUint32(level.Buffer[pos:]) + return +} + +// func FindPageOnRecords(records []byte, timestamp uint32) (pageNo uint32) { +// comparator := ValueAtComparator{ +// buf: records, +// timestamp: timestamp, +// } +// count := len(records) / IndexRecordSize +// elemIdx, _ := BinarySearch(count, comparator) +// pageNo, _ = bin.GetUint32(records[elemIdx*IndexRecordSize:]) +// return +// } + +func FindPageOnIndexPage(page []byte, timestamp uint32) (pageNo uint32) { + comparator := ValueAtComparator{ + buf: page, + timestamp: timestamp, + } + count, _ := bin.GetUint16(page[indexRecordsCountIdx:]) + elemIdx, _ := BinarySearch(int(count), comparator) + pos := elemIdx*IndexRecordSize + 4 // timestamp size + pageNo, _ = bin.GetUint32(page[pos:]) + return +} + +func IsZeroLevelPage(buf []byte) bool { + return buf[isZeroLevelIdx] == 1 +} + +// можна перевірити хвости від 0 до ... Перевіряю since кожного хвоста. +// Якщо timestamp >=, шукаю pageNo бінарним пошуком. І сторінка мені однозначно підходить. +// Якщо timestamp <, піднімаюсь вище. Якщо рівнів більше немає - until вказано за межами Range показань. +func FindPageOnIndexLevelTails(levels []IndexLevelTail, timestamp uint32) (pageNo uint32, isDataPage bool) { + for i, level := range levels { + tailSince, _ := bin.GetUint32(level.Buffer) + if timestamp >= tailSince { + pageNo = FindPageOnIndexLevelTail(level, timestamp) + return pageNo, i == 0 + } + } + return +} + // func findPageNoIdx(buf []byte, timestamp uint32) (idx int) { // comparator := ValueAtComparator{ // buf: buf, diff --git a/storage/storage.go b/storage/storage.go index 49bd577..07193ac 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -9,7 +9,7 @@ var ( // data page DataPageSize = 8192 - DataPagePayloadSize int = DataPageSize - DataPageFooterSize + DataPagePayloadSize = DataPageSize - DataPageFooterSize dataCRC32Idx = DataPageSize - 4 timestampsSizeIdx = DataPageSize - 6 @@ -19,6 +19,7 @@ var ( // index page IndexPageSize = 1024 + IndexPagePayloadSize = IndexPageSize - IndexPageFooterSize //indexPageIncSize = IndexPageIncSize indexCRC32Idx = IndexPageSize - 4 indexRecordsCountIdx = IndexPageSize - 6 @@ -28,7 +29,6 @@ var ( // timestampSize = 4 // pairSize = timestampSize + PageNoSize - // indexFooterIdx = indexRecordsQtyIdx // dataFooterIdx = timestampsSizeIdx ) diff --git a/storage/storage_test.go b/storage/storage_test.go index b07bb0a..3d3436a 100644 --- a/storage/storage_test.go +++ b/storage/storage_test.go @@ -1081,3 +1081,74 @@ func TestReplayMetric(t *testing.T) { fmt.Printf("%d: % x\n", metric.IndexLevelTails[0].RecordsCount, metric.IndexLevelTails[0].Buffer) fmt.Printf("%d: % x\n", metric.IndexLevelTails[1].RecordsCount, metric.IndexLevelTails[1].Buffer) } + +// func TestComposeHeadIndexPage(t *testing.T) { +// var ( +// levelIdx = 0 +// level = storage.IndexLevelTail{ +// Buffer: []byte{ +// 1, 2, 3, +// }, +// RecordsCount: 2, +// } +// head = &storage.IndexPageTail{ +// PageNo: 100, +// CRC32: 12345, +// Records: []byte{}, +// } +// ) +// page := composeHeadIndexPage(levelIdx, level, head) +// if page.PageNo != head.PageNo { +// t.Fatalf("PageNo: got %d are not equal expected %v", +// page.PageNo, head.PageNo) +// } +// // fix compare pages +// } + +func TestFindPageOnIndexTails(t *testing.T) { + levels := []IndexLevelTail{ + { + Buffer: []byte{ + 0x64, 0x00, 0x00, 0x00, 0x0a, 0x00, 0x00, 0x00, // 100 => 10 + 0x6e, 0x00, 0x00, 0x00, 0x0b, 0x00, 0x00, 0x00, // 110 => 11 + }, + RecordsCount: 2, + }, + { + Buffer: []byte{ + 0x0a, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, // 10 => 1 + 0x28, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x00, // 40 => 4 + 0x46, 0x00, 0x00, 0x00, 0x07, 0x00, 0x00, 0x00, // 70 => 7 + }, + RecordsCount: 3, + }, + } + testCases := []struct { + Timestamp uint32 + PageNo uint32 + IsDataPage bool + }{ + {Timestamp: 9, PageNo: 0, IsDataPage: false}, + {Timestamp: 10, PageNo: 1, IsDataPage: false}, + {Timestamp: 20, PageNo: 1, IsDataPage: false}, + {Timestamp: 40, PageNo: 4, IsDataPage: false}, + {Timestamp: 60, PageNo: 4, IsDataPage: false}, + {Timestamp: 70, PageNo: 7, IsDataPage: false}, + {Timestamp: 90, PageNo: 7, IsDataPage: false}, + {Timestamp: 100, PageNo: 10, IsDataPage: true}, + {Timestamp: 105, PageNo: 10, IsDataPage: true}, + {Timestamp: 110, PageNo: 11, IsDataPage: true}, + {Timestamp: 120, PageNo: 11, IsDataPage: true}, + } + for _, testCase := range testCases { + pageNo, isDataPage := FindPageOnIndexLevelTails(levels, testCase.Timestamp) + if pageNo != testCase.PageNo { + t.Fatalf("timestamp %d: got pageNo %d are not equal expected %d", + testCase.Timestamp, pageNo, testCase.PageNo) + } + if isDataPage != testCase.IsDataPage { + t.Fatalf("timestamp %d: got isDataPage %t are not equal expected %t", + testCase.Timestamp, isDataPage, testCase.IsDataPage) + } + } +} diff --git a/storage/writer.go b/storage/writer.go index 5fb8dfa..d0412d9 100644 --- a/storage/writer.go +++ b/storage/writer.go @@ -365,9 +365,6 @@ type PageToWrite struct { func WriteDataPages(file *os.File, pages []PageToWrite) (err error) { for _, p := range pages { - fmt.Println("pageNo: %d\n", p.PageNo) - } - for _, p := range pages { if len(p.Content) != DataPageSize { return fmt.Errorf("wrong data page size: %d", len(p.Content)) } diff --git a/timeutil/timeutil.go b/timeutil/timeutil.go index 5cd559d..12ac4d9 100644 --- a/timeutil/timeutil.go +++ b/timeutil/timeutil.go @@ -1,6 +1,11 @@ package timeutil -import "time" +import ( + "time" + + "gordenko.dev/dima/qb" + "gordenko.dev/dima/qb/proto" +) func FirstSecondInPeriod(since time.Time, period string) (_ time.Time) { y, m, d := since.Date() @@ -37,3 +42,24 @@ func LastSecondInPeriod(until time.Time, period string) (_ time.Time) { return until } } + +func TimeBoundsOfAggregation(since, until proto.TimeBound, groupBy qb.GroupBy, firstHourOfDay int) (s time.Time, u time.Time) { + switch groupBy { + case qb.GroupByHour, qb.GroupByDay: + s = time.Date(since.Year, since.Month, since.Day, 0, 0, 0, 0, time.Local) + u = time.Date(until.Year, until.Month, until.Day, 0, 0, 0, 0, time.Local) + + case qb.GroupByMonth: + s = time.Date(since.Year, since.Month, 1, 0, 0, 0, 0, time.Local) + u = time.Date(until.Year, until.Month, 1, 0, 0, 0, 0, time.Local) + } + + if firstHourOfDay > 0 { + duration := time.Duration(firstHourOfDay) * time.Hour + s = s.Add(duration) + u = u.Add(duration) + } + + u = u.Add(-1 * time.Second) + return +} diff --git a/database/tmp b/tmp similarity index 81% rename from database/tmp rename to tmp index 8013628..fabf43b 100644 --- a/database/tmp +++ b/tmp @@ -90,4 +90,24 @@ // ValuesSize: s.vSize, // NewBufferSize: bufferSize, // }) -// } \ No newline at end of file +// } + +// func isFileExist(fileName string) (bool, error) { +// _, err := os.Stat(fileName) +// if err != nil { +// if errors.Is(err, fs.ErrNotExist) { +// return false, nil +// } else { +// return false, err +// } +// } else { +// return true, nil +// } +// } + +// func correctToFHD(since, until uint32, firstHourOfDay int) (uint32, uint32) { +// duration := time.Duration(firstHourOfDay) * time.Hour +// since = uint32(time.Unix(int64(since), 0).Add(duration).Unix()) +// until = uint32(time.Unix(int64(until), 0).Add(duration).Unix()) +// return since, until +// } diff --git a/worker/metric.go b/worker/metric.go index 5888c77..10d568b 100644 --- a/worker/metric.go +++ b/worker/metric.go @@ -16,18 +16,14 @@ import ( var ErrNoValueBug = errors.New("has timestamp but no value") type CapturedState struct { - LastTimestamp uint32 - LastValue float64 + LastValue float64 } type Metric struct { - metricType qb.MetricType - fracDigits byte - lastPageNo uint32 - //SinceValue float64 - //Since uint32 - lastValue float64 - //Until uint32 + metricType qb.MetricType + fracDigits byte + lastPageNo uint32 + lastValue float64 buffer []byte timestamps qb.TimestampCompressor values qb.ValueCompressor @@ -117,14 +113,42 @@ func (s *Metric) LastValue() float64 { return s.lastValue } -func (s *Metric) LastTimestamp() uint32 { - if s.capturedState != nil { - return s.capturedState.LastTimestamp +// REQUESTS + +func (s *Metric) ReleaseRLock() { + s.rLocks-- + if s.rLocks == 0 { + if len(s.waitQueue) > 0 { + s.ProcessQueue() + } } - return s.timestamps.LastTimestamp() } -// REQUESTS +// суть у тому що треба запускати запити, пока не зустріну XLock +func (s *Metric) ProcessQueue(metricID uint32, tmp []byte) { + if len(s.waitQueue) == 0 { + return + } + for _, untyped := range s.waitQueue { + switch req := untyped.(type) { + case RangeScanReq: + s.StartRangeScan(req) + case FullScanReq: + s.StartFullScan(req) + case GetMetricReq: + s.GetMetric(req) + case AppendMeasuresReq: + metric.AppendMeasures(req, tmp, s.storageInbox) + case DeleteMetricReq: + s.DeleteMetric(req) + case DeleteMeasuresReq: + metric.DeleteMeasures(req) + default: + qb.Abort(qb.UnknownMetricWaitQueueItemBug, + fmt.Errorf("bug: unknown metric wait queue item type %T", req)) + } + } +} func (s *Metric) AppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox *inbox.Inbox) { if s.xLock || s.capturedState != nil { @@ -151,15 +175,14 @@ func (s *Metric) AppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox ) s.capturedState = &CapturedState{ - LastTimestamp: timestamps.LastTimestamp(), - LastValue: s.lastValue, + LastValue: s.lastValue, } s.timestamps.CaptureState() s.values.CaptureState() for idx, measure := range req.Measures { - if measure.Timestamp <= s.timestamps.LastTimestamp() { + if measure.Timestamp <= s.timestamps.Until() { resultCode = ExpiredMeasure break } @@ -270,15 +293,15 @@ func (s *Metric) AppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox } func (s *Metric) DeleteMeasures(req DeleteMeasuresReq) { - if s.xLock || s.capturedState != nil { - s.waitQueue = append(s.waitQueue, req) - return - } - since := s.timestamps.FirstTimestamp() - until := s.timestamps.LastTimestamp() - if since == 0 || (req.Since > 0 && until < req.Since) { - req.ResultCh <- NoMeasuresToDelete - } + // if s.xLock || s.capturedState != nil { + // s.waitQueue = append(s.waitQueue, req) + // return + // } + // since := s.timestamps.FirstTimestamp() + // until := s.timestamps.LastTimestamp() + // if since == 0 || (req.Since > 0 && until < req.Since) { + // req.ResultCh <- NoMeasuresToDelete + // } // if s.RootPageNo > 0 { // req.ResultCh <- tryDeleteMeasuresResult{ // ResultCode: DeleteFromAtreeRequired, @@ -296,36 +319,47 @@ func (s *Metric) StartRangeScan(req RangeScanReq) { s.waitQueue = append(s.waitQueue, req) return } - // if s.timestamps.CommitedSize() == 0 { - // req.ResultCh <- RangeScanResult{ - // ResultCode: QueryDone, - // } - // return - // } - - // if req.Since > s.timestamps.LastTimestamp() { - // 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 - // } - // } + since := s.timestamps.CommitedSince() + if since == 0 { + req.ResultCh <- RangeScanResult{ + ResultCode: QueryDone, + } + return + } + // range after + if req.Since > s.timestamps.CommitedUntil() { + req.ResultCh <- RangeScanResult{ + ResultCode: QueryDone, + } + return + } + // range before + if req.Until < since { + if len(s.indexLevelTails) > 0 { + pageNo, isDataPage := storage.FindPageOnIndexLevelTails(s.indexLevelTails, req.Until) + if pageNo > 0 { + req.ResultCh <- RangeScanResult{ + ResultCode: UntilNotFound, + FracDigits: s.fracDigits, + LastPageNo: pageNo, + IsDataPage: isDataPage, + } + s.rLocks++ + return + } else { + // range until before 1st measure timestamp + req.ResultCh <- RangeScanResult{ + ResultCode: QueryDone, + } + return + } + } else { + req.ResultCh <- RangeScanResult{ + ResultCode: QueryDone, + } + return + } + } timestampDecompressor := s.timestamps.CreateDecompressor() valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits) @@ -443,11 +477,6 @@ func (s *Metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) { s.xLock = false // s.Timestamps.Renew() // s.Values.Renew() - - // s.LastPageNo = 0 - // s.Since = 0 - // s.SinceValue = 0 - // s.Until = 0 s.indexLevelTails = nil s.lastPageNo = 0 s.lastValue = 0 diff --git a/worker/worker.go b/worker/worker.go index abb863d..c018257 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -5,7 +5,6 @@ import ( "sync" qb "gordenko.dev/dima/qb" - "gordenko.dev/dima/qb/atree" "gordenko.dev/dima/qb/enc" "gordenko.dev/dima/qb/inbox" "gordenko.dev/dima/qb/proto" @@ -86,15 +85,8 @@ func New(opt Options) *Worker { return s } -func (s *Worker) ReleaseRLock(metricID uint32) { - // s.mutex.Lock() - // //s.rLocksToRelease = append(s.rLocksToRelease, metricID) - // s.mutex.Unlock() - - // select { - // case s.signalCh <- struct{}{}: - // default: - // } +type ReleaseRLock struct { + MetricID uint32 } func (s *Worker) Run() { @@ -114,39 +106,12 @@ func (s *Worker) Run() { func (s *Worker) doWork() { queue := s.inbox.Drain() - //rLocksToRelease := s.rLocksToRelease - //s.rLocksToRelease = nil - - // for _, metricID := range rLocksToRelease { - // metric, ok := s.metrics[metricID] - // if !ok { - // qb.Abort(qb.NoMetricBug, - // fmt.Errorf("drainQueues: metric %d not found", metricID)) - // } - - // if metric.XLock { - // qb.Abort(qb.XLockBug, - // fmt.Errorf("drainQueues: xlock is set for the metric %d", - // metricID)) - // } - - // if metric.RLocks <= 0 { - // qb.Abort(qb.NoRLockBug, - // fmt.Errorf("drainQueues: rlock not set for the metric %d", - // metricID)) - // } - - // metric.RLocks-- - - // if len(metric.WaitQueue) > 0 { - // s.processMetricQueue(metricID, metric) - // } - // } - for _, untyped := range queue { switch req := untyped.(type) { case AppendMeasuresReq: s.AppendMeasures(req) + case ReleaseRLock: + s.releaseRLock(req.MetricID) case storage.Changes: s.applyCommits(req) // all metrics only case ListCurrentValuesReq: @@ -170,29 +135,10 @@ func (s *Worker) doWork() { } } -// суть у тому що треба запускати запити, пока не зустріну XLock -func (s *Worker) processMetricQueue(metricID uint32, metric *Metric, tmp []byte) { - if len(metric.waitQueue) == 0 { - return - } - for _, untyped := range metric.waitQueue { - switch req := untyped.(type) { - case RangeScanReq: - metric.StartRangeScan(req) - case FullScanReq: - metric.StartFullScan(req) - case GetMetricReq: - s.GetMetric(req) - case AppendMeasuresReq: - metric.AppendMeasures(req, tmp, s.storageInbox) - case DeleteMetricReq: - s.DeleteMetric(req) - case DeleteMeasuresReq: - metric.DeleteMeasures(req) - default: - qb.Abort(qb.UnknownMetricWaitQueueItemBug, - fmt.Errorf("bug: unknown metric wait queue item type %T", req)) - } +func (s *Worker) releaseRLock(metricID uint32) { + metric, ok := s.metrics[metricID] + if ok { + metric.ReleaseRLock() } } @@ -318,8 +264,8 @@ func (s *Worker) AppendMeasures(req AppendMeasuresReq) { type RangeScanResult struct { ResultCode byte FracDigits byte - RootPageNo uint32 LastPageNo uint32 + IsDataPage bool // for UntilNotFound only } type RangeScanReq struct { @@ -327,7 +273,7 @@ type RangeScanReq struct { Since uint32 Until uint32 MetricType qb.MetricType - ResponseWriter atree.WorkerMeasureConsumer + ResponseWriter qb.WorkerMeasureConsumer ResultCh chan RangeScanResult } @@ -361,7 +307,7 @@ type FullScanResult struct { type FullScanReq struct { MetricID uint32 MetricType qb.MetricType - ResponseWriter atree.WorkerMeasureConsumer + ResponseWriter qb.WorkerMeasureConsumer ResultCh chan FullScanResult } @@ -398,7 +344,7 @@ func (s *Worker) ListCurrentValues(req ListCurrentValuesReq) { if ok { req.ResponseWriter.BufferValue(transform.CurrentValue{ MetricID: metricID, - Timestamp: metric.LastTimestamp(), + Timestamp: metric.timestamps.CommitedUntil(), Value: metric.LastValue(), }) }