diff --git a/atree/atree.go b/atree/atree.go index 3280648..4d30d65 100644 --- a/atree/atree.go +++ b/atree/atree.go @@ -22,21 +22,27 @@ const ( pageTypeIdx = PageSize - 5 // index page - indexRecordsQtyIdx = PageSize - 8 - isDataPageNumbersIdx = PageSize - 6 + indexCRC32Idx = IndexPageSize - 4 + isLastLevelIdx = IndexPageSize - 5 + indexRecordsQtyIdx = IndexPageSize - 7 + + indexRecordSize = 8 // data page - timestampsSizeIdx = PageSize - 13 - valuesSizeIdx = PageSize - 11 - prevPageIdx = PageSize - 9 + dataCRC32Idx = DataPageSize - 4 + prevPageIdx = DataPageSize - 8 + valuesSizeIdx = DataPageSize - 10 + timestampsSizeIdx = DataPageSize - 12 timestampSize = 4 pairSize = timestampSize + PageNoSize indexFooterIdx = indexRecordsQtyIdx dataFooterIdx = timestampsSizeIdx - PageSize = 8192 - PageNoSize = 4 + PageSize = 8192 + DataPageSize = 8192 + IndexPageSize = 1024 + PageNoSize = 4 // DataPagePayloadSize int = dataFooterIdx ) @@ -161,7 +167,7 @@ func (s *Atree) findDataPage(rootPageNo uint32, timestamp uint32) (uint32, []byt s.releasePage(indexPageNo) // fix - if buf[isDataPageNumbersIdx] == 1 { + if buf[isLastLevelIdx] == 1 { buf, err := s.fetchDataPage(foundPageNo) if err != nil { return 0, nil, fmt.Errorf("fetchDataPage(%d): %s", foundPageNo, err) @@ -206,7 +212,7 @@ func (s *Atree) FindPathToLastPage(rootPageNo uint32) (_ PathToDataPage, err err foundPageNo := getLastPageNo(buf) // fix - if buf[isDataPageNumbersIdx] == 1 { + if buf[isLastLevelIdx] == 1 { return PathToDataPage{ Legs: legs, LastPageNo: foundPageNo, diff --git a/atree/x.go b/atree/x.go new file mode 100644 index 0000000..d9c2072 --- /dev/null +++ b/atree/x.go @@ -0,0 +1,147 @@ +package atree + +import "gordenko.dev/dima/qb/bin" + +const ( + payloadSize = 24 + recordSize = 4 + 4 // timestamp + pageNo +) + +type DataPage struct { + PrevPageNo uint32 + LowerTimestamp uint32 + Timestamps [][]byte // chunks + TimestampsSize int + Values [][]byte // chunks + ValuesSize int + PageNo uint32 + IsReused bool +} + +type IndexPage struct { + LowerTimestamp uint32 + Data []byte + PageNo uint32 + IsReused bool +} +type IndexLevel struct { + // вже записані дані (лише для першої в списку індексної сторінки, беремо із _metric) + // дані потрібні, якщо сторінка буде заповнена і піде на запис в .index файл + Offset int // (заповнені одразу) + LowerTimestamp uint32 // (заповнені одразу) + Data []byte // поточні дані index сторінки (заповнені одразу) + Filled []IndexPage // пусто +} + +func appendIndexRecord(data []byte, timestamp uint32, pageNo uint32) []byte { + var tmp = []byte{ + 0, 0, 0, 0, // timestamp + 0, 0, 0, 0, // pageNo + } + bin.PutUint32(tmp, timestamp) + bin.PutUint32(tmp[4:], pageNo) + return append(data, tmp...) +} + +type ChunksToDataPageReq struct { + PrevPageNo uint32 + Timestamps [][]byte // chunks + TimestampsSize uint16 + Values [][]byte // chunks + ValuesSize uint16 +} + +func ChunksToDataPage(buf []byte, req ChunksToDataPageReq) { + var ( + remainingSize = int(req.TimestampsSize) + pos = 0 + ) + for _, chunk := range req.Timestamps { + 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) + + for _, chunk := range req.Values { + if remainingSize >= len(chunk) { + copy(buf[pos:], chunk) + remainingSize -= len(chunk) + pos += len(chunk) + } else { + copy(buf[pos:], chunk[:remainingSize]) + break + } + } + bin.PutUint16(buf[timestampsSizeIdx:], req.TimestampsSize) + bin.PutUint16(buf[valuesSizeIdx:], req.ValuesSize) + bin.PutUint32(buf[prevPageIdx:], req.PrevPageNo) + bin.PutUint32(buf[dataCRC32Idx:], calcChecksum(buf[:dataCRC32Idx])) +} + +func DataToIndexPage(buf []byte, data []byte, lastLevel bool) { + copy(buf, data) + bin.PutUint16(buf[indexRecordsQtyIdx:], uint16(len(data)/indexRecordSize)) + if lastLevel { + buf[isLastLevelIdx] = 1 + } + bin.PutUint32(buf[indexCRC32Idx:], calcChecksum(buf[:indexCRC32Idx])) +} + +type Pager interface { + GetPageNumber() (uint32, bool) +} + +func AppendDataPagesToTree(pager Pager, levels []*IndexLevel, dataPages []*DataPage) { + for i, d := range dataPages[1:] { + d.PageNo, d.IsReused = pager.GetPageNumber() + d.PrevPageNo = dataPages[i-1].PageNo + } + + for _, d := range dataPages { + var ( + upTimestamp = d.LowerTimestamp + upPageNo = d.PageNo + levelIdx = 0 + ) + for { + if levelIdx == len(levels) { + // encode prev page record + up* + levels = append(levels, &IndexLevel{ + LowerTimestamp: upTimestamp, + Data: appendIndexRecord(nil, upTimestamp, upPageNo), + }) + break + } + level := levels[levelIdx] + if len(level.Data)+recordSize <= payloadSize { + level.Data = appendIndexRecord(level.Data, upTimestamp, upPageNo) + break + } + // в буфері немає місця для додавання нової пари, отже це + // заповнена сторінка + filled := IndexPage{ + LowerTimestamp: level.LowerTimestamp, + Data: level.Data, + } + filled.PageNo, filled.IsReused = pager.GetPageNumber() + + level.Filled = append(level.Filled, filled) + // + level.LowerTimestamp = upTimestamp + level.Data = appendIndexRecord(nil, upTimestamp, upPageNo) + + // + upPageNo = filled.PageNo + upTimestamp = filled.LowerTimestamp + levelIdx++ + } + } +} diff --git a/txlog/x.go b/txlog/x.go index 97ab852..66d9426 100644 --- a/txlog/x.go +++ b/txlog/x.go @@ -1,83 +1,233 @@ package txlog -const ( - payloadSize = 24 - recordSize = 4 + 4 // timestamp + pageNo +import ( + "bytes" + "fmt" + "log" + "os" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb/atree" ) -type Page struct { - //Head []byte - //Tail []byte - Data []byte - PageNo uint32 - LowerTimestamp uint32 +const recordSize = 8 + +type LogWriter struct { + w *bytes.Buffer + wal *os.File + dataFile *os.File + indexFile *os.File + dataPage []byte + indexPage []byte } -type Level struct { - // вже записані дані (лише для першої в списку індексної сторінки, беремо із _metric) - // дані потрібні, якщо сторінка буде заповнена і піде на запис в .index файл - AlreadyWritten int // (заповнені одразу) - Unfilled []byte // поточні дані index сторінки (заповнені одразу) - LowerTimestamp uint32 // (заповнені одразу) - PageNo uint32 // пусто - Filled []Page // пусто -} - -type Pager interface { - GetPageNumber() uint32 -} - -func xxx(pager Pager, levels []*Level, dataPages []*Page) { - var ( - //upTimestamp uint32 - lastDataIdx = len(dataPages) - 1 - ) - for idx, d := range dataPages { - d.PageNo = pager.GetPageNumber() - if idx < lastDataIdx { - // setPrevPageNo(dataPages[idx+1], d.PageNo) - // calc CRC32 - } else { - +func (s *LogWriter) writePagesToAtree(levels []*atree.IndexLevel, dataPages []*atree.DataPage) (err error) { + for _, p := range dataPages { + atree.ChunksToDataPage(s.dataPage, 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)) + if err != nil { + return err } + 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)) + // } + 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) + } + } + } + return nil +} + +type ToWrite 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 +} + +type DataPage struct { + PageNo uint32 + Reused bool + PrevPageNo uint32 + Checksum uint32 + Timestamps []byte + Values []byte +} + +type IndexPage struct { + PageNo uint32 + Reused bool + Checksum uint32 + Records []byte +} + +type IndexLevel struct { + Pages []IndexPage + // offset потрібен тому що дані не додаються в кінець, а перезаписують кілька + // останніх байтів unfilled даних. Потрібно для коректного recovery. + // Якщо є Pages, то застосовуються для 1-ї Index сторінки, інакше - до records (unfilled) + RecordsOffset int + Records []byte // unfilled +} + +type TxAppendedMeasures struct { + MetricID uint32 + DataPages []DataPage + // offsets потрібні тому що дані не додаються в кінець, а перезаписують кілька + // останніх байтів unfilled даних. Потрібно для коректного recovery. + // Якщо є DataPages, то застосовуються для 1-ї Data сторінки, інакше - до timestamps і values (unfilled) + TimestampsOffset int + ValuesOffset int + Timestamps []byte // unfilled + Values []byte // unfilled + IndexLevels []DataPage +} + +/* +Format appended measures: +1b - tx type +4b - metricID +Nb - qty of filled data pages (varsize) +[ + 4b - pageNo + 1b - reused + 4b - prevPageNo + 4b - page crc32 + Nb - timestamps size (varsize) + Nb - timestamps payload + Nb - values size (varsize) + Nb - values payload +] +2b - timestamps offset (for the 1st data page only) +2b - values offset (for the 1st data page only) +Nb - timestamps size unfilled (varsize) +Nb - timestamps payload unfilled +Nb - values size unfilled (varsize) +Nb - values payload unfilled +Nb - qty of index levels (varsize) +[ + Nb - qty of level filled pages (varsize) + [ + 4b - pageNo + 1b - reused + 4b - page crc32 + Nb - records qty (varsize) + Nb - records payload + ] + Nb - records offset (for the 1st index page only) + Nb - records qty unfilled (varsize) + Nb - records payload unfilled +] +*/ + +func (s *LogWriter) writeToWAL(req ToWrite) (err error) { + // fix write code + bin.WriteUint32(s.w, req.MetricID) + bin.WriteVarSize(s.w, len(req.DataPages)) + for idx, p := range req.DataPages { + bin.WriteUint32(s.w, p.PageNo) + var reused byte + if p.IsReused { + reused = 1 + } + s.w.WriteByte(reused) + bin.WriteUint32(s.w, p.PrevPageNo) + //bin.WriteUint32(s.w, p.CRC32) fix var ( - upTimestamp = d.LowerTimestamp - upPageNo uint32 = d.PageNo - levelIdx = 0 + timestampsOffset int + valuesOffset int ) - for { - if levelIdx == len(levels) { - // - newRoot := &Level{ - Unfilled: nil, - //LowerTimestamp: prevPage . LowerTimestamp, - PageNo: pager.GetPageNumber(), - } - // encode prev page record + up* - levels = append(levels, newRoot) - break + if idx == 0 { + timestampsOffset = req.TimestampsOffset + valuesOffset = req.ValuesOffset + } + copyPayloadFromChunks(s.w, p.Timestamps, p.TimestampsSize, timestampsOffset) + copyPayloadFromChunks(s.w, p.Values, p.ValuesSize, valuesOffset) + } + bin.WriteUint16(s.w, uint16(req.TimestampsOffset)) + bin.WriteUint16(s.w, uint16(req.ValuesOffset)) + // fix offset + copyPayloadFromChunks(s.w, req.Timestamps, req.TimestampsSize, req.TimestampsOffset) + copyPayloadFromChunks(s.w, req.Values, req.ValuesSize, req.ValuesOffset) + if len(req.DataPages) == 0 { + return + } + bin.WriteVarSize(s.w, len(req.IndexLevels)) + for _, level := range req.IndexLevels { + bin.WriteVarSize(s.w, len(level.Filled)) + for idx, p := range level.Filled { + bin.WriteUint32(s.w, p.PageNo) + var reused byte + if p.IsReused { + reused = 1 } - level := levels[levelIdx] - if len(level.Unfilled)+recordSize <= payloadSize { - // append - break + s.w.WriteByte(reused) + //bin.WriteUint32(s.w, p.CRC32) fix + records := p.Data + if idx == 0 { + records = p.Data[level.Offset:] } - level.Filled = append(level.Filled, Page{ - Data: level.Unfilled, - LowerTimestamp: level.LowerTimestamp, - PageNo: level.PageNo, - }) - // - level.Unfilled = nil - // fix - pack record - level.PageNo = pager.GetPageNumber() - level.LowerTimestamp = upTimestamp + if (len(records) % recordSize) != 0 { + log.Fatalf("wrong records length: %d", len(records)) + } + bin.WriteVarSize(s.w, len(records)/recordSize) + s.w.Write(records) + } + if len(level.Filled) > 0 { + s.w.Write(level.Data) + } else { + s.w.Write(level.Data[level.Offset:]) + } + } + return +} - // - upPageNo = level.PageNo - upTimestamp = level.LowerTimestamp - levelIdx++ +const chunkSize = 24 +func copyPayloadFromChunks(w *bytes.Buffer, chunks [][]byte, size int, offset int) { + bin.WriteVarSize(w, size) + var ( + chunkIdx = offset / chunkSize + byteIdx = offset % chunkSize + ) + for _, chunk := range chunks[chunkIdx:] { + available := len(chunk) - byteIdx + if available <= size { + w.Write(chunk[byteIdx:]) + size -= available + if size == 0 { + return + } + byteIdx = 0 + } else { + w.Write(chunk[byteIdx : byteIdx+size]) + return } } }