package storage import ( "fmt" "gordenko.dev/dima/qb" ) type ReplayMetric struct { MetricType qb.MetricType FracDigits byte LastPageNo uint32 Buf []byte TimestampsSize int ValuesSize int IndexLevelTails []IndexLevelTail // root - last element } // return occupied after write func writeTimestampsWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int { // timestamps writes right-to-left pos := DataPagePayloadSize - occupied // timestamps rewind offset moves pos left-to-right pos += rewindOffset pos -= len(payload) copy(buf[pos:], payload) return DataPagePayloadSize - pos } func writeValuesWithRewind(buf []byte, occupied int, payload []byte, rewindOffset int) int { // values writes left-to-right pos := occupied // values rewind offset moves pos right-to-left pos -= rewindOffset copy(buf[pos:], payload) return pos + len(payload) } func (s *ReplayMetric) MeasuresAppend(rec MeasuresAppendRecord) { s.TimestampsSize = writeTimestampsWithRewind(s.Buf, s.TimestampsSize, rec.Timestamps, rec.TimestampsRewindOffset) s.ValuesSize = writeValuesWithRewind(s.Buf, s.ValuesSize, rec.Values, rec.ValuesRewindOffset) } type ReplayMeasuresAppendWithGrowResult struct { IndexPagesToRewrite []PageToWrite DataPagesToRewrite []PageToWrite ReusedIndexPagesCount int ReusedDataPagesCount int } func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (_ ReplayMeasuresAppendWithGrowResult) { //fmt.Printf("%#v\n", rec) var ( reusedIndexPages int reusedDataPages int indexPages []PageToWrite dataPages []PageToWrite ) // Додати в index level tails недостаючі дані, або замінити for levelIdx, change := range rec.ChangedIndexLevels { var ( head = change.IndexPageTail newRecordsCount = len(change.TailRecords) / IndexRecordSize ) // 1. набиваю сторінки для перезапису в index файлі if isLastPacket { if head != nil { indexPages = append(indexPages, s.composeHeadIndexPage(levelIdx, head)) } for _, x := range change.IndexPages { indexPages = append(indexPages, PageToWrite{ PageNo: x.PageNo, Content: x.Content, }) } } // 2. Вношу зміни в поточні індексні хвости if levelIdx < len(s.IndexLevelTails) { level := s.IndexLevelTails[levelIdx] if head != nil { // попередній хвіст індексного рівня перетворився на head сторінку, // отже TailRecords - це новий хвіст clear(level.Buffer) copy(level.Buffer, change.TailRecords) level.RecordsCount = newRecordsCount } else { // TailRecords - це нові дані, які треба додати pos := level.RecordsCount * IndexRecordSize copy(level.Buffer[pos:], change.TailRecords) level.RecordsCount += newRecordsCount } s.IndexLevelTails[levelIdx] = level } else { // додаю новий індексний рівень buf := make([]byte, IndexPageSize) copy(buf, change.TailRecords) // s.IndexLevelTails = append(s.IndexLevelTails, IndexLevelTail{ Buffer: buf, RecordsCount: newRecordsCount, }) } // 3. рахую кількість reused індексних сторінок if head != nil { if head.Reused { reusedIndexPages++ } } for _, x := range change.IndexPages { if x.Reused { reusedIndexPages++ } } } // 4. набиваю сторінки для перезапису в data файлі if isLastPacket { dataPages = append(dataPages, s.composeHeadDataPage(rec.DataPageTail)) for _, x := range rec.DataPages { dataPages = append(dataPages, PageToWrite{ PageNo: x.PageNo, Content: x.Content, }) } } // рахую кількість reused дата сторінок if rec.DataPageTail.Reused { reusedDataPages++ } for _, x := range rec.DataPages { if x.Reused { reusedDataPages++ } } // HeadDataPage є полюбому, оскільки це транзакція із мінімум однією заповненою // data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня. clear(s.Buf) copy(s.Buf, rec.TailValues) s.ValuesSize = len(rec.TailValues) pos := DataPagePayloadSize - len(rec.TailTimestamps) copy(s.Buf[pos:], rec.TailTimestamps) s.TimestampsSize = len(rec.TailTimestamps) if len(rec.DataPages) == 0 { s.LastPageNo = rec.DataPageTail.PageNo } else { s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo } return ReplayMeasuresAppendWithGrowResult{ IndexPagesToRewrite: indexPages, DataPagesToRewrite: dataPages, ReusedIndexPagesCount: reusedIndexPages, ReusedDataPagesCount: reusedDataPages, } } // HELPERS func (s *ReplayMetric) composeHeadIndexPage(levelIdx int, head *IndexPageTail) PageToWrite { var ( page = make([]byte, IndexPageSize) pos int prevCount int ) //fmt.Printf("composeHeadIndexPage: levels=%d\n", len(s.IndexLevelTails)) //fmt.Printf("head: % x\n", head.Records) if levelIdx < len(s.IndexLevelTails) { level := s.IndexLevelTails[levelIdx] //fmt.Printf("level (%d): % x\n", level.RecordsCount, level.Buffer) // розраховую pos, з якого буду дописувати хвіст pos = level.RecordsCount * IndexRecordSize // створюю копію сторінки copy(page, level.Buffer[:pos]) // поточні дані prevCount = level.RecordsCount } else { // s.IndexLevelTails = append(s.IndexLevelTails, IndexLevelTail{ Buffer: page, RecordsCount: len(head.Records) / IndexRecordSize, }) } copy(page[pos:], head.Records) // запечатати сторінку calculatedCRC := SealIndexPage(SealIndexPageIn{ Content: page, RecordsCount: prevCount + len(head.Records)/IndexRecordSize, ZeroLevel: levelIdx == 0, }) // перевірка CRC if calculatedCRC != head.CRC32 { qb.Abort(qb.WALReplayFailed, fmt.Errorf("calculated CRC %d not equal expected %d of head page %d on index level %d", calculatedCRC, head.CRC32, head.PageNo, levelIdx)) } return PageToWrite{ PageNo: head.PageNo, Content: page, } } func (s *ReplayMetric) composeHeadDataPage(head DataPageTail) PageToWrite { // copy page page := make([]byte, DataPageSize) copy(page, s.Buf) // append data timestampsSize := writeTimestampsWithRewind(page, s.TimestampsSize, head.Timestamps, head.TimestampsRewindOffset) valuesSize := writeValuesWithRewind(page, s.ValuesSize, head.Values, head.ValuesRewindOffset) // запечатати сторінку calculatedCRC := SealDataPage(SealDataPageIn{ Content: page, PrevPageNo: head.PrevPageNo, TimestampsSize: timestampsSize, ValuesSize: valuesSize, }) // перевірка CRC if calculatedCRC != head.CRC32 { qb.Abort(qb.WALReplayFailed, fmt.Errorf("calculated CRC %d not equal expected %d of head data page %d", calculatedCRC, head.CRC32, head.PageNo)) } return PageToWrite{ PageNo: head.PageNo, Content: page, } }