diff --git a/enc/cumdelta.go b/enc/cumdelta.go index 7558063..07d9e9d 100644 --- a/enc/cumdelta.go +++ b/enc/cumdelta.go @@ -86,9 +86,14 @@ type EvaluationReport struct { CompressionWay int TotalSpace int AdditionalSpace int + Offset int + BytesCount int Delta uint64 } +// можна виділити буфер максимального розміру, що може бути змінено під час кодування +// закодувати на етапі evaluate значення і повернути buf + offset. Причому буфер може бути +// спільний на всі metrics func (s *CumulativeDeltaCompressor) Evaluate(value float64) (compressionWay int, requiredSpace int) { delta := uint64((value-s.baseValue)*s.coef + eps) // fix - if delta 0, no eps requiredSpace = s.pos @@ -97,9 +102,13 @@ func (s *CumulativeDeltaCompressor) Evaluate(value float64) (compressionWay int, // run if delta == s.lastDelta && s.h < 127 { compressionWay = incrementRun + // offset = s.pos + // bytesCount = 1 } else { compressionWay = endSeries requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize + // offset = s.pos + 1 + // bytesCount = bin.CountVarUint64(uint64(delta)) + hSize } } else { // literal @@ -107,6 +116,8 @@ func (s *CumulativeDeltaCompressor) Evaluate(value float64) (compressionWay int, if s.h < 255 { compressionWay = incrementLiteral requiredSpace += bin.CountVarUint64(uint64(delta)) + // offset = s.pos + // bytesCount = bin.CountVarUint64(uint64(delta)) // hSize просто переміщається } else { compressionWay = endSeries requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize @@ -114,6 +125,7 @@ func (s *CumulativeDeltaCompressor) Evaluate(value float64) (compressionWay int, } else { compressionWay = startRun if s.h > 128 { + // offset = s.pos - 1 - bin.CountVarUint64(uint64(delta)) requiredSpace += hSize // h - end of run } } diff --git a/enc/enc_test.go b/enc/enc_test.go index 0b49485..8d52cab 100644 --- a/enc/enc_test.go +++ b/enc/enc_test.go @@ -8,6 +8,7 @@ import ( "testing" bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" ) func equalFloatSlices(a, b []float64, epsilon float64) bool { @@ -18,44 +19,49 @@ func equalFloatSlices(a, b []float64, epsilon float64) bool { // CUMULATIVE -func TestCumulativeDeltaCompressor(t *testing.T) { +func TestValueDeltaCompressor(t *testing.T) { var ( testCases = []struct { - Nums []float64 - RepeatLastDeltaNTimes int - Code int - RequiredSpace int - Buf []byte - BaseValue float64 - LastDelta uint64 - H byte + Nums []float64 + Name string + Offset int + ChangeSize int + Buf []byte + BaseValue float64 + LastDelta uint64 }{ { Nums: []float64{ 1.5, }, - Code: addBaseValue, - RequiredSpace: 3, // base value, delta, h + Name: "add 1st value", + Offset: 0, + ChangeSize: 3, Buf: []byte{ - 143, 0, 0, 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x80, // literal (len = 1) + 0x00, 0x00, 0x00, 0x00, 0x00, }, BaseValue: 1.5, LastDelta: 0, - H: 128, // literal, length = 1 }, { Nums: []float64{ 1.5, 1.5, }, - Code: startRun, // literal changed to run - RequiredSpace: 2, // delta, h + Name: "literal switch to run", + Offset: 1, + ChangeSize: 1, Buf: []byte{ - 143, 0, 0, 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x00, // run (len = 2) + 0x00, 0x00, 0x00, 0x00, 0x00, }, BaseValue: 1.5, LastDelta: 0, - H: 0, // run, length = 2 }, { Nums: []float64{ @@ -63,17 +69,19 @@ func TestCumulativeDeltaCompressor(t *testing.T) { 1.6, 1.6, }, - Code: startRun, // literal decreased by 1 - RequiredSpace: 3, // end of literal, new delta, new h + Name: "literal decrease by 1 and switch to run", + Offset: 2, + ChangeSize: 3, Buf: []byte{ - 143, // base value - 128, // literal 1st delta (0) - 128, // h - end of literal, length = 1 - 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x80, // literal (len = 1) + 0x81, // delta 10 + 0x00, // run (len = 2) + 0x00, 0x00, 0x00, }, BaseValue: 1.5, LastDelta: 1, - H: 0, // run, length = 2 }, { Nums: []float64{ @@ -81,15 +89,17 @@ func TestCumulativeDeltaCompressor(t *testing.T) { 1.5, 1.5, }, - Code: incrementRun, - RequiredSpace: 2, + Name: "increment run", + Offset: 1, + ChangeSize: 1, Buf: []byte{ - 143, // base value - 0, 0, 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x01, // run (len = 3) + 0x00, 0x00, 0x00, 0x00, 0x00, }, BaseValue: 1.5, LastDelta: 0, - H: 1, // run, length = 3 }, { Nums: []float64{ @@ -98,112 +108,107 @@ func TestCumulativeDeltaCompressor(t *testing.T) { 1.5, 1.6, }, - Code: endSeries, // end of run - RequiredSpace: 4, + Name: "run switch to literal", + Offset: 0, + ChangeSize: 2, Buf: []byte{ - 143, // base value - 128, // run delta (0) - 1, // h of run, length = 3 - 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x01, // run (len = 3) + 0x81, // delta 10 + 0x80, // literal (len = 1) + 0x00, 0x00, 0x00, }, BaseValue: 1.5, LastDelta: 1, - H: 128, // literal, length = 1 - }, - { - Nums: []float64{ - 1.5, - }, - RepeatLastDeltaNTimes: 128, - Code: incrementRun, // h byte full filled - RequiredSpace: 2, - Buf: []byte{ - 143, // base value - 0, 0, 0, 0, 0, 0, 0, - }, - BaseValue: 1.5, - LastDelta: 0, - H: 127, // run, length = 129 - }, - { - Nums: []float64{ - 1.5, - }, - RepeatLastDeltaNTimes: 129, - Code: endSeries, // run, h byte overflowed - RequiredSpace: 4, // run delta, h of run, 1st literal delta, h of literal - Buf: []byte{ - 143, // base value - 128, // run delta (0) - 127, // h of run (length = 129) - 0, 0, 0, 0, 0, - }, - BaseValue: 1.5, - LastDelta: 0, - H: 128, // literal, length = 1 }, + // { + // Nums: []float64{ + // 1.5, + // }, + // RepeatLastDeltaNTimes: 128, + // Code: incrementRun, // h byte full filled + // RequiredSpace: 2, + // Buf: []byte{ + // 143, // base value + // 0, 0, 0, 0, 0, 0, 0, + // }, + // BaseValue: 1.5, + // LastDelta: 0, + // H: 127, // run, length = 129 + // }, + // { + // Nums: []float64{ + // 1.5, + // }, + // RepeatLastDeltaNTimes: 129, + // Code: endSeries, // run, h byte overflowed + // RequiredSpace: 4, // run delta, h of run, 1st literal delta, h of literal + // Buf: []byte{ + // 143, // base value + // 128, // run delta (0) + // 127, // h of run (length = 129) + // 0, 0, 0, 0, 0, + // }, + // BaseValue: 1.5, + // LastDelta: 0, + // H: 128, // literal, length = 1 + // }, { Nums: []float64{ 1.5, 1.6, + 1.7, }, - Code: incrementLiteral, - RequiredSpace: 3, // previous delta, new delta, h + Name: "increment literal", + Offset: 1, + ChangeSize: 2, Buf: []byte{ - 143, // base value - 128, // literal 1st delta (0) - 0, 0, 0, 0, 0, 0, + 0x8f, // base value + 0x80, // delta 0 + 0x81, // delta 10 + 0x82, // delta 20 + 0x82, // literal (len = 3) + 0x00, 0x00, 0x00, }, BaseValue: 1.5, - LastDelta: 1, - H: 129, // literal, length = 2 + LastDelta: 2, }, } ) - for caseIdx, testCase := range testCases { + for _, testCase := range testCases { var ( - fracDigits byte = 1 - buf = make([]byte, 8) - payloadSize = 0 - compressionWay int - requiredSpace int + tmp = make([]byte, tmpValueSize) + fracDigits byte = 1 + buf = make([]byte, 8) + payloadSize = 0 + report qb.ValueEvaluationReport ) - c := NewCumulativeDeltaCompressor(fracDigits, buf, payloadSize) + c := NewValueDeltaCompressor(fracDigits, buf, payloadSize) for _, num := range testCase.Nums { - compressionWay, requiredSpace = c.Evaluate(num) - c.Compress(compressionWay, num) + report = c.Evaluate(tmp, num) + c.Append(report.Offset, tmp[:report.ChangeSize], num, report.Delta) } - if testCase.RepeatLastDeltaNTimes > 0 { - num := testCase.Nums[len(testCase.Nums)-1] - for range testCase.RepeatLastDeltaNTimes { - compressionWay, requiredSpace = c.Evaluate(num) - c.Compress(compressionWay, num) - } + if report.Offset != testCase.Offset { + t.Fatalf("%s: got offset %d are not equal expected %d", + testCase.Name, report.Offset, testCase.Offset) } - if compressionWay != testCase.Code { - t.Fatalf("%d: got code %d are not equal expected %d", - caseIdx, compressionWay, testCase.Code) - } - if requiredSpace != testCase.RequiredSpace { - t.Fatalf("%d: got requiredSpace %d are not equal expected %d", - caseIdx, requiredSpace, testCase.RequiredSpace) + if report.ChangeSize != testCase.ChangeSize { + t.Fatalf("%s: got changeSize %d are not equal expected %d", + testCase.Name, report.ChangeSize, testCase.ChangeSize) } if !bytes.Equal(buf, testCase.Buf) { - t.Fatalf("%d: got buf %v are not equal expected %v", - caseIdx, buf, testCase.Buf) + t.Fatalf("%s: got buf % x are not equal expected % x", + testCase.Name, buf, testCase.Buf) } if c.baseValue != testCase.BaseValue { - t.Fatalf("%d: got lastUnixtime %v are not equal expected %v", - caseIdx, c.baseValue, testCase.BaseValue) + t.Fatalf("%s: got baseValue %v are not equal expected %v", + testCase.Name, c.baseValue, testCase.BaseValue) } if c.lastDelta != testCase.LastDelta { - t.Fatalf("%d: got lastDelta %d are not equal expected %d", - caseIdx, c.lastDelta, testCase.LastDelta) - } - if c.h != testCase.H { - t.Fatalf("%d: got h %d are not equal expected %d", - caseIdx, c.h, testCase.H) + t.Fatalf("%s: got lastDelta %d are not equal expected %d", + testCase.Name, c.lastDelta, testCase.LastDelta) } } } @@ -687,41 +692,44 @@ func TestInstantDeltaDecompressorFromState(t *testing.T) { func TestTimeDeltaCompressor(t *testing.T) { var ( testCases = []struct { - Nums []uint32 - RepeatLastDeltaNTimes int - CompressionWay int - RequiredSpace int - Buf []byte - LastUnixtime uint32 - LastDelta uint32 - H byte + Name string + Nums []uint32 + Buf []byte + LastUnixtime uint32 + LastDelta uint32 + Offset int + ChangeSize int }{ { Nums: []uint32{ 1780777000, }, - CompressionWay: addUnixtime, - RequiredSpace: 4, + Name: "add 1st value", + Offset: 0, + ChangeSize: 4, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, 0, + 0x00, 0x00, 0x00, 0x00, + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777000, LastDelta: 0, - H: 0, }, { Nums: []uint32{ 1780777000, 1780777060, // +60 }, - CompressionWay: add1stDelta, - RequiredSpace: 6, + Name: "add 1st delta", + Offset: 0, + ChangeSize: 2, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, 0, + 0x00, 0x00, + 0x80, // h-byte (literal, len=1) + 0xbc, // delta (60) + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777060, LastDelta: 60, - H: 128, }, { Nums: []uint32{ @@ -729,14 +737,17 @@ func TestTimeDeltaCompressor(t *testing.T) { 1780777060, // +60 1780777120, // +60 }, - CompressionWay: startRun, // literal changed by run - RequiredSpace: 6, + Name: "literal changed to run", + Offset: 1, + ChangeSize: 1, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, 0, + 0x00, 0x00, + 0x00, // h-byte (run, len=2) + 0xbc, // delta + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777120, LastDelta: 60, - H: 0, // run, length = 2 }, { Nums: []uint32{ @@ -745,14 +756,18 @@ func TestTimeDeltaCompressor(t *testing.T) { 1780777130, // +70 1780777200, // +70 }, - CompressionWay: startRun, // literal decreased by 1 - RequiredSpace: 7, // unixtime, end of literal, new delta, new h + Name: "literal decrease by 1 and switch to run", + Offset: 2, + ChangeSize: 3, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 128, 188, + 0x00, // h-byte (run, len=2) + 0xc6, // delta 70 + 0x80, // h-byte (literal, len=1) + 0xbc, // delta 60 + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777200, LastDelta: 70, - H: 0, // run, length = 2 }, { Nums: []uint32{ @@ -761,14 +776,17 @@ func TestTimeDeltaCompressor(t *testing.T) { 1780777120, // +60 1780777180, // +60 }, - CompressionWay: incrementRun, - RequiredSpace: 6, + Name: "increment run", + Offset: 1, + ChangeSize: 1, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, 0, + 0x00, 0x00, + 0x01, // h-byte (run, len=2) + 0xbc, // delta 60 + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777180, LastDelta: 60, - H: 1, // run, length = 3 }, { Nums: []uint32{ @@ -778,111 +796,114 @@ func TestTimeDeltaCompressor(t *testing.T) { 1780777180, // +60 1780777200, // +20 }, - CompressionWay: endSeries, // run - RequiredSpace: 8, + Name: "switch run to literal", + Offset: 0, + ChangeSize: 2, Buf: []byte{ - 0, 0, 0, 0, 0, 0, - 1, // h of run, length = 3 - 188, // varUint64(60) + 0x80, // h-byte (literal, len=1) + 0x94, // delta 20 + 0x01, // h-byte (run, len=3) + 0xbc, // delta 60 + 0x28, 0x80, 0x24, 0x6a, // since }, LastUnixtime: 1780777200, LastDelta: 20, - H: 128, // literal, length = 1 }, + // { + // Nums: []uint32{ + // 1780777000, + // 1780777060, // +60 + // }, + // RepeatLastDeltaNTimes: 128, + // CompressionWay: incrementRun, // h byte full filled + // RequiredSpace: 6, + // Buf: []byte{ + // 0x80, // h-byte (literal, len=1) + // 0x94, // delta 20 + // 0x01, // h-byte (run, len=3) + // 0xbc, // delta 60 + // 0x28, 0x80, 0x24, 0x6a, // since + // }, + // LastUnixtime: 1780777060 + 128*60, + // LastDelta: 60, + // H: 127, // run, length = 129 + // }, + // { + // Nums: []uint32{ + // 1780777000, + // 1780777060, // +60 + // }, + // RepeatLastDeltaNTimes: 129, + // CompressionWay: endSeries, // run, h byte overflowed + // RequiredSpace: 8, + // Buf: []byte{ + // 0, 0, 0, 0, 0, 0, + // 127, // h - end of run (length = 129) + // 188, // varUint64(60) + // }, + // LastUnixtime: 1780777060 + 129*60, + // LastDelta: 60, + // H: 128, // literal, length = 1 + // }, { Nums: []uint32{ 1780777000, 1780777060, // +60 + 1780777130, // +70 }, - RepeatLastDeltaNTimes: 128, - CompressionWay: incrementRun, // h byte full filled - RequiredSpace: 6, + Name: "increment literal", + Offset: 1, + ChangeSize: 2, Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, 0, + 0x00, + 0x81, // h-byte (literal, len=2) + 0xc6, // delta 70 + 0xbc, // delta 60 + 0x28, 0x80, 0x24, 0x6a, // since }, - LastUnixtime: 1780777060 + 128*60, - LastDelta: 60, - H: 127, // run, length = 129 - }, - { - Nums: []uint32{ - 1780777000, - 1780777060, // +60 - }, - RepeatLastDeltaNTimes: 129, - CompressionWay: endSeries, // run, h byte overflowed - RequiredSpace: 8, - Buf: []byte{ - 0, 0, 0, 0, 0, 0, - 127, // h - end of run (length = 129) - 188, // varUint64(60) - }, - LastUnixtime: 1780777060 + 129*60, - LastDelta: 60, - H: 128, // literal, length = 1 - }, - { - Nums: []uint32{ - 1780777000, - 1780777060, // +60 - 1780777122, // +62 - }, - CompressionWay: incrementLiteral, - RequiredSpace: 7, - Buf: []byte{ - 0, 0, 0, 0, 0, 0, 0, - 188, // varUint64(60) - }, - LastUnixtime: 1780777122, - LastDelta: 62, - H: 129, // literal, length = 2 + LastUnixtime: 1780777130, + LastDelta: 70, }, } ) - for caseIdx, testCase := range testCases { + for _, testCase := range testCases { //fmt.Println("----------") var ( - buf = make([]byte, 8) - payloadSize = 0 - compressionWay int - requiredSpace int + tmp = make([]byte, tmpTimeSize) + buf = make([]byte, 8) + payloadSize = 0 + report qb.TimeEvaluationReport ) c := NewTimeDeltaCompressor(buf, payloadSize) for _, num := range testCase.Nums { - compressionWay, requiredSpace = c.Evaluate(num) - c.Compress(compressionWay, num) + report = c.Evaluate(tmp, num) + + // fmt.Println("----------\npos", c.pos) + // fmt.Printf("buf: % x\n", c.buf[c.pos:]) + // pretty.Println(report) + // fmt.Printf("change: % x\n", tmp[:report.ChangeSize]) + c.Append(report.Offset, tmp[:report.ChangeSize], num) } - if testCase.RepeatLastDeltaNTimes > 0 { - for range testCase.RepeatLastDeltaNTimes { - num := c.lastUnixtime + c.lastDelta - compressionWay, requiredSpace = c.Evaluate(num) - c.Compress(compressionWay, num) - } + if report.Offset != testCase.Offset { + t.Fatalf("%s: got offset %d are not equal expected %d", + testCase.Name, report.Offset, testCase.Offset) } - if compressionWay != testCase.CompressionWay { - t.Fatalf("%d: got code %d are not equal expected %d", - caseIdx, compressionWay, testCase.CompressionWay) - } - if requiredSpace != testCase.RequiredSpace { - t.Fatalf("%d: got requiredSpace %d are not equal expected %d", - caseIdx, requiredSpace, testCase.RequiredSpace) + if report.ChangeSize != testCase.ChangeSize { + t.Fatalf("%s: got changeSize %d are not equal expected %d", + testCase.Name, report.ChangeSize, testCase.ChangeSize) } if !bytes.Equal(buf, testCase.Buf) { - t.Fatalf("%d: got buf %v are not equal expected %v", - caseIdx, buf, testCase.Buf) + t.Fatalf("%s: got buf % x are not equal expected % x", + testCase.Name, buf, testCase.Buf) } if c.lastUnixtime != testCase.LastUnixtime { - t.Fatalf("%d: got lastUnixtime %d are not equal expected %d", - caseIdx, c.lastUnixtime, testCase.LastUnixtime) + t.Fatalf("%s: got lastUnixtime %d are not equal expected %d", + testCase.Name, c.lastUnixtime, testCase.LastUnixtime) } if c.lastDelta != testCase.LastDelta { - t.Fatalf("%d: got lastDelta %d are not equal expected %d", - caseIdx, c.lastDelta, testCase.LastDelta) - } - if c.h != testCase.H { - t.Fatalf("%d: got h %d are not equal expected %d", - caseIdx, c.h, testCase.H) + t.Fatalf("%s: got lastDelta %d are not equal expected %d", + testCase.Name, c.lastDelta, testCase.LastDelta) } } } @@ -961,14 +982,15 @@ func TestTimeDeltaDecompressorFromState(t *testing.T) { ) for caseIdx, testCase := range testCases { var ( + tmp = make([]byte, tmpTimeSize) buf = make([]byte, 16) payloadSize = 0 decodedNums []uint32 ) c := NewTimeDeltaCompressor(buf, payloadSize) for _, num := range testCase.Nums { - compressionWay, _ := c.Evaluate(num) - c.Compress(compressionWay, num) + report := c.Evaluate(tmp, num) + c.Append(report.Offset, tmp[:report.ChangeSize], num) } // diff --git a/enc/time_delta.go b/enc/time_delta.go index 86d0909..c510084 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -18,13 +18,16 @@ Timestamps читаються звичайними методами Get*, а п */ +const ( + tmpTimeSize = 7 // max when start run +) + type TimeDeltaCompressor struct { buf []byte // останній записаний байт, якщо рахувати справа наліво pos int lastUnixtime uint32 lastDelta uint32 - h byte state *TimeDeltaCapturedState } @@ -37,17 +40,22 @@ type TimeDeltaCompressor struct { func NewTimeDeltaCompressor(buf []byte, payloadSize int) *TimeDeltaCompressor { s := &TimeDeltaCompressor{ buf: buf, + pos: len(buf) - payloadSize, } if payloadSize > 0 { - //s.lastUnixtime = lastTimestamp - s.pos = len(s.buf) - payloadSize + var err error + s.lastUnixtime, err = bin.GetUint32(s.buf[len(s.buf)-4:]) + if err != nil { + log.Fatalf("bug: get since: %s", err) + } + //s.lastUnixtime = lastTimestamp // fix - порахувати + if payloadSize > 4 { - s.h = s.buf[s.pos] - u64, _, err := bin.GetVarUint64(s.buf[s.pos+1:]) - if err != nil { - log.Fatalf("bug: get last delta: %s", err) - } - s.lastDelta = uint32(u64) + // u64, _, err := bin.GetVarUint64(s.buf[s.pos+1:]) + // if err != nil { + // log.Fatalf("bug: get last delta: %s", err) + // } + // s.lastDelta = uint32(u64) } } return s @@ -58,109 +66,99 @@ func NewTimeDeltaCompressor(buf []byte, payloadSize int) *TimeDeltaCompressor { // } -// retrun -func (s *TimeDeltaCompressor) Evaluate(unixtime uint32) (compressionWay int, requiredSpace int) { - requiredSpace = len(s.buf) - s.pos - if requiredSpace > 0 { - delta := unixtime - s.lastUnixtime +func (s *TimeDeltaCompressor) Evaluate(tmp []byte, timestamp uint32) qb.TimeEvaluationReport { + var ( + delta = timestamp - s.lastUnixtime + offset int + i int + ) + if s.pos < len(s.buf) { if s.lastDelta > 0 { - if s.h < 128 { + h := s.buf[s.pos] + if h < 128 { // run - if delta == s.lastDelta && s.h < 127 { - compressionWay = incrementRun + if delta == s.lastDelta && h < 127 { + // incrementRun + tmp[i] = h + 1 + i++ + offset = 1 // перезапис h } else { - compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize + // endSeries + n, _ := bin.PutVarUint64(tmp, uint64(delta)) + i += n + tmp[i] = 128 // start new literal (length=1) + i++ } } else { // literal if delta != s.lastDelta { - if s.h < 255 { - compressionWay = incrementLiteral - requiredSpace += bin.CountVarUint64(uint64(delta)) + if h < 255 { + // incrementLiteral + n, _ := bin.PutVarUint64(tmp, uint64(delta)) + i += n + tmp[i] = h + 1 + i++ + offset = 1 // перезапис h } else { - compressionWay = endSeries - requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize + // endSeries + n, _ := bin.PutVarUint64(tmp, uint64(delta)) + i += n + tmp[i] = 128 // start new literal (length=1) + i++ } } else { - compressionWay = startRun - if s.h > 128 { - requiredSpace += hSize // h - end of run + // startRun + if h > 128 { + tmp[i] = h - 1 // зменшую довжину попередньої серії на 1 + i++ + n, _ := bin.PutVarUint64(tmp[i:], uint64(delta)) + i += n + tmp[i] = 0 // start new run (length=2) + i++ + offset = 1 + n // перезапис пари delta/h + } else { + tmp[i] = 0 // change literal (length=1) to run (length=2) + i++ + offset = 1 } } } } else { - compressionWay = add1stDelta - requiredSpace += bin.CountVarUint64(uint64(delta)) + hSize + // add 1st delta + n, _ := bin.PutVarUint64(tmp, uint64(delta)) + i += n + tmp[i] = 128 // start new literal (length=1) + i++ } } else { - compressionWay = addUnixtime - requiredSpace += 4 + bin.PutUint32(tmp, timestamp) // 1st timestamp (since) + i += 4 + } + return qb.TimeEvaluationReport{ + Offset: offset, + ChangeSize: i, + TotalSpace: s.pos - offset + i, } - return } -// pos вказує на h byte -func (s *TimeDeltaCompressor) Compress(compressionWay int, unixtime uint32) { - delta := unixtime - s.lastUnixtime - switch compressionWay { - case incrementRun: - s.h++ - s.buf[s.pos] = s.h - case incrementLiteral: - s.lastDelta = delta - s.h++ - // тому що перезаписую h - s.pos++ - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) - s.pos -= n - s.pos-- - s.buf[s.pos] = s.h - case endSeries: - // start new literal (length=1) - s.lastDelta = delta - s.h = 128 - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) - s.pos -= n - s.pos-- - s.buf[s.pos] = s.h - case startRun: - if s.h > 128 { - s.h-- - s.pos += bin.CountVarUint64(uint64(s.lastDelta)) // забираю останню дельту у literal - s.buf[s.pos] = s.h // write h - end of literal (because length > 1) - s.pos-- - // run delta - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) - s.pos -= n - // для h - s.pos-- +func (s *TimeDeltaCompressor) Append(offset int, change []byte, timestamp uint32) { + if s.pos < len(s.buf) { + i := s.pos + offset - 1 // -1, because s.pos always points to h byte + for _, b := range change { + s.buf[i] = b + i-- } - // start new run (length=2) - s.h = 0 - s.buf[s.pos] = s.h - case add1stDelta: - // start new literal (length=1) - s.lastDelta = delta - s.h = 128 - n, _ := bin.TailPutVarUint64(s.buf[:s.pos], uint64(delta)) - s.pos -= n - s.pos-- - s.buf[s.pos] = s.h - case addUnixtime: - s.pos -= 4 - bin.PutUint32(s.buf[s.pos:], unixtime) + s.lastDelta = timestamp - s.lastUnixtime + } else { + copy(s.buf[len(s.buf)-4:], change) // 4b since } - // 1st value or update after every update - // do not write immediatelly because of update after every update - s.lastUnixtime = unixtime + s.pos -= len(change) - offset + s.lastUnixtime = timestamp } -func (s *TimeDeltaCompressor) DeleteLast() { +// func (s *TimeDeltaCompressor) DeleteLast() { -} - -// FIX - check methods +// } type TimeDeltaCapturedState struct { H byte @@ -176,7 +174,7 @@ func (s *TimeDeltaCompressor) CaptureState() { } // позиція посувається вліво, отже може перескочити на попередній chunk s.state = &TimeDeltaCapturedState{ - H: s.h, + H: s.buf[s.pos], LastUnixtime: s.lastUnixtime, LastDelta: s.lastDelta, Payload: s.buf[s.pos:], @@ -197,32 +195,25 @@ func (s *TimeDeltaCompressor) ReplaceSinceWithUntil() uint32 { } // fix -func (s *TimeDeltaCompressor) Offset() int { - // if s.state != nil { - // return s.state.Pos - // } - return 0 -} - -func (s *TimeDeltaCompressor) Size() int { - if s.state == nil { - return len(s.buf) - s.pos - } else { - return bin.CountVarUint64(uint64(s.state.LastDelta)) + hSize + len(s.state.Payload) - } -} +// func (s *TimeDeltaCompressor) Size() int { +// if s.state == nil { +// return len(s.buf) - s.pos +// } else { +// return bin.CountVarUint64(uint64(s.state.LastDelta)) + hSize + len(s.state.Payload) +// } +// } // Snapshot - для створення снапшота. // FIX - encode firstUnixtime in 1st 4 bytes // // while page not completed - has first unixtime, after - last unixtime -func (s *TimeDeltaCompressor) Payload() []byte { - if s.state == nil { - return s.buf[s.pos:] - } else { - return s.state.Payload - } -} +// func (s *TimeDeltaCompressor) Payload() []byte { +// if s.state == nil { +// return s.buf[s.pos:] +// } else { +// return s.state.Payload +// } +// } func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { if s.state == nil { @@ -259,12 +250,14 @@ func (s *TimeDeltaCompressor) Rotate(newbuf []byte) { s.pos = len(s.buf) s.lastUnixtime = 0 s.lastDelta = 0 - s.h = 0 } func (s *TimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompressor { + if s.pos == len(s.buf) { + return nil + } var ( - h = s.h + h = s.buf[s.pos] lastUnixtime = s.lastUnixtime lastDelta = s.lastDelta payload = s.buf[s.pos:] diff --git a/enc/value_delta.go b/enc/value_delta.go new file mode 100644 index 0000000..c08b134 --- /dev/null +++ b/enc/value_delta.go @@ -0,0 +1,393 @@ +package enc + +import ( + "io" + "log" + "math" + + bin "gordenko.dev/dima/bin/little" + "gordenko.dev/dima/qb" +) + +const ( + tmpValueSize = 19 +) + +type ValueDeltaCompressor struct { + buf []byte + coef float64 + pos int // payload size + baseValue float64 + lastDelta uint64 + state *ValueDeltaCapturedState +} + +// Після відновлення із снапшота +func NewValueDeltaCompressor(fracDigits byte, buf []byte, payloadSize int) *ValueDeltaCompressor { + var coef float64 = 1 + if fracDigits > 0 { + coef = math.Pow(10, float64(fracDigits)) + } + s := &ValueDeltaCompressor{ + buf: buf, + coef: coef, + pos: payloadSize, + } + if payloadSize > 0 { + // base value на початку + u64, _, err := bin.GetVarUint64(s.buf) + if err != nil { + log.Fatalf("bug: get base value: %s", err) + } + s.baseValue = float64(u64) / s.coef + s.lastDelta, _, err = bin.ReverseGetVarUint64(s.buf[:s.pos-2]) // skip h byte + if err != nil { + log.Fatalf("bug: get last delta: %s", err) + } + } + return s +} + +// можна виділити буфер максимального розміру, що може бути змінено під час кодування +// закодувати на етапі evaluate значення і повернути buf + offset. Причому буфер може бути +// спільний на всі metrics +// arr - +func (s *ValueDeltaCompressor) Evaluate(tmp []byte, value float64) qb.ValueEvaluationReport { + var ( + delta uint64 + offset int + i int + ) + if s.pos > 0 { + delta = uint64((value-s.baseValue)*s.coef + eps) + h := s.buf[s.pos-1] + if h < 128 { + // run + if delta == s.lastDelta && h < 127 { + // incrementRun + tmp[i] = h + 1 + i++ + offset = 1 // перезапис h + } else { + // endSeries + n, _ := bin.ReversePutVarUint64(tmp, delta) + i += n + tmp[i] = 128 // start new literal (length=1) + i++ + } + } else { + // literal + if delta != s.lastDelta { + if h < 255 { + // incrementLiteral + n, _ := bin.ReversePutVarUint64(tmp, delta) + i += n + tmp[i] = h + 1 + i++ + offset = 1 // перезапис h + } else { + // endSeries + n, _ := bin.ReversePutVarUint64(tmp, delta) + i += n + tmp[i] = 128 // start new literal (length=1) + i++ + } + } else { + // startRun + if h > 128 { + tmp[i] = h - 1 // зменшую довжину попередньої серії на 1 + i++ + n, _ := bin.ReversePutVarUint64(tmp[i:], delta) + i += n + tmp[i] = 0 // start new run (length=2) + i++ + offset = 1 + n // перезапис пари delta/h + } else { + tmp[i] = 0 // change literal (length=1) to run (length=2) + i++ + offset = 1 + } + } + } + } else { + n, _ := bin.PutVarUint64(tmp, uint64(value*s.coef)) // base value + i += n + n, _ = bin.ReversePutVarUint64(tmp[i:], 0) // delta + i += n + tmp[i] = 128 // start new literal (length=1) + i++ + } + return qb.ValueEvaluationReport{ + Offset: offset, + ChangeSize: i, + TotalSpace: s.pos + i - offset, + Delta: delta, + } +} + +// pos завжди вказує на h +func (s *ValueDeltaCompressor) Append(offset int, change []byte, value float64, delta uint64) { + if s.pos > 0 { + s.lastDelta = delta + } else { + s.baseValue = value + } + copy(s.buf[s.pos-offset:], change) + s.pos += len(change) - offset +} + +func (s *ValueDeltaCompressor) DeleteLast() { + +} + +type ValueDeltaCapturedState struct { + H byte + LastDelta uint64 + Payload []byte +} + +func (s *ValueDeltaCompressor) CaptureState() { + if s.state != nil { + qb.Abort(qb.RepeatableLock, nil) + } + // позиція посувається вліво, отже може перескочити на попередній chunk + pos := s.pos - 1 - bin.CountVarUint64(s.lastDelta) + s.state = &ValueDeltaCapturedState{ + H: s.buf[s.pos-1], + LastDelta: s.lastDelta, + Payload: s.buf[:pos], + } +} + +// fix - повернути в Pool буфери +func (s *ValueDeltaCompressor) ForgetCapturedState() { + s.state = nil +} + +func (s *ValueDeltaCompressor) Size() int { + if s.state == nil { + return s.pos + } else { + return bin.CountVarUint64(s.state.LastDelta) + hSize + len(s.state.Payload) + } +} + +// // Snapshot - для створення снапшота. +// func (s *ValueDeltaCompressor) Payload() []byte { +// if s.state == nil { +// return s.buf[:s.pos] +// } else { +// return s.state.Payload +// } +// } + +func (s *ValueDeltaCompressor) WritePayloadTo(w io.Writer) (err error) { + if s.state == nil { + _, err = w.Write(s.buf[:s.pos]) + return + } else { + _, err = w.Write(s.state.Payload) + if err != nil { + return + } + _, err = bin.WriteVarUint64(w, s.state.LastDelta) + if err != nil { + return + } + _, err = w.Write([]byte{ + s.state.H, + }) + return + } +} + +func (s *ValueDeltaCompressor) Rotate(newbuf []byte) { + // УВАГА! + // state не чіпаємо + s.buf = newbuf + s.pos = 0 + s.baseValue = 0 + s.lastDelta = 0 +} + +func (s *ValueDeltaCompressor) LastValue() float64 { + var delta uint64 + if s.state == nil { + delta = s.lastDelta + } else { + delta = s.state.LastDelta + } + return s.baseValue + float64(delta)*s.coef +} + +func (s *ValueDeltaCompressor) CreateDecompressor() qb.ValueDecompressor { + if s.pos == 0 { + return nil + } + var ( + h = s.buf[s.pos-1] + lastDelta = s.lastDelta + payload = s.buf[:s.pos] + ) + if s.state != nil { + h = s.state.H + lastDelta = s.state.LastDelta + payload = s.state.Payload + } + return NewValueDeltaDecompressorFromState(ValueDeltaDecompressorFromStateOptions{ + Coef: s.coef, + H: h, + LastDelta: lastDelta, + Payload: payload, + }) +} + +// DECOMPRESSOR + +type ValueDeltaDecompressor struct { + buf []byte + coef float64 + pos int + bound int + baseValue float64 + lastValue float64 + isRun bool + pending int + done bool +} + +func NewValueDeltaDecompressor(buf []byte, fracDigits byte) *ValueDeltaDecompressor { + var coef float64 = 1 + if fracDigits > 0 { + coef = math.Pow(10, float64(fracDigits)) + } + u64, n, err := bin.GetVarUint64(buf) + if err != nil { + log.Fatalf("bug: get base value: %s", err) + } + //fmt.Printf("baseValue: %.2f, bound: %d, pos: %d\n", float64(u64)/coef, n, size) + s := &ValueDeltaDecompressor{ + buf: buf, + coef: coef, + pos: len(buf), // first free + baseValue: float64(u64) / coef, + bound: n, + } + // читаю заголовок наступної серії + s.readHeader() + s.readValue() + return s +} + +type ValueDeltaDecompressorFromStateOptions struct { + Coef float64 + H byte + LastDelta uint64 + Payload []byte +} + +func NewValueDeltaDecompressorFromState(opt ValueDeltaDecompressorFromStateOptions) *ValueDeltaDecompressor { + u64, n, err := bin.GetVarUint64(opt.Payload) + if err != nil { + log.Fatalf("bug: get base value: %s", err) + } + //fmt.Printf("baseValue: %.2f, bound: %d, pos: %d\n", float64(u64)/opt.Coef, n, size) + s := &ValueDeltaDecompressor{ + buf: opt.Payload, + coef: opt.Coef, + pos: len(opt.Payload), // first free + baseValue: float64(u64) / opt.Coef, + bound: n, + } + s.lastValue = s.baseValue + float64(opt.LastDelta)/s.coef + s.decodeHeaderByte(opt.H) + //fmt.Printf("payload: %d\n", opt.Payload) + //fmt.Printf("restore from state: bound=%d, isRun=%t, pending=%d, baseValue=%v, lastValue=%v\n", + // n, s.isRun, s.pending, s.baseValue, s.lastValue) + return s +} + +func (s *ValueDeltaDecompressor) NextValue() (value float64, done bool) { + //fmt.Printf("NextValue(): bound: %d, pos: %d, pending: %d\n", s.bound, s.pos, s.pending) + if s.done { + return 0, true + } + // метод працює як do while - спочатку значення, а потім перевірка умови + value = s.lastValue + s.pending-- + if s.pending > 0 { + // якщо в серії залишаються елементи + if !s.isRun { + s.readValue() + } + } else if s.pos > s.bound { + // в серії більше немає елементів, отже перевіряє чи є ще дані в буфері. + // дані є - читаю заголовок наступної серії + s.readHeader() + s.readValue() + } else { + s.done = true + } + // серія завершена - перевіряю чи є ще серії + return value, false +} + +func (s *ValueDeltaDecompressor) readHeader() { + //fmt.Println("read from pos:", s.pos) + s.pos-- + h := s.buf[s.pos] + s.decodeHeaderByte(h) + + // fmt.Println("h:", h) + // fmt.Println("isRun:", s.isRun) + // fmt.Println("pending:", s.pending) +} +func (s *ValueDeltaDecompressor) readValue() { + u64, n, err := bin.ReverseGetVarUint64(s.buf[:s.pos]) + if err != nil { + log.Fatalln(err) + } + // fmt.Println() + // fmt.Println("read from pos:", s.pos) + // fmt.Println("read delta:", u64) + // fmt.Println("read delta n:", n) + // fmt.Println() + s.pos -= n + s.lastValue = s.baseValue + float64(u64)/s.coef + //fmt.Println(s.baseValue, float64(u64)/s.coef) +} + +func (s *ValueDeltaDecompressor) decodeHeaderByte(h byte) { + s.isRun = h < 128 + if s.isRun { + s.pending = int(h) + 2 + } else { + s.pending = int(h&127) + 1 + } +} + +/* +Формат: +base value (var u64) +delta (var u64) +q (qty of previous deltas; msb=0 - series, msb=1 - non series, low 7 bits - qty 1..128) +delta +delta +delta +q + +Дельта рахується від base value. +Декодування у зворотньому порядку. + +Після base value слідують run або literal блоки. +Run блок - це delta + header byte в кінці. +Literal блок - це від одної до N дельт + header byte в кінці. +Спочатку створюється literal блок. +Якщо для останної дельти додається дублікат, literal блок модифікується - +лічильник зменшується до 1. А остання дельта переміщюється в новий run блок. +Причому лічильник 0 - означає 2 елементи. Приклад: +До: +v1 v2 v3 h-byte(literal, 3) <- v3 +Після: +v1 v2 h-byte(literal, 2) v3 h-byte(run, 2) +*/ diff --git a/qb.go b/qb.go index 8f6f35c..5e29b6a 100644 --- a/qb.go +++ b/qb.go @@ -24,28 +24,36 @@ const ( ) type TimestampCompressor interface { - // (timestamp) => compressionWay, requiredSpace - Evaluate(uint32) (int, int) - Compress(int, uint32) - Size() int + // (tmp, timestamp) + Evaluate([]byte, uint32) TimeEvaluationReport + // (offset, tmp) + Append(int, []byte) + //Size() int //Chunks() [][]byte //DeleteLast() CaptureState() ForgetCapturedState() CreateDecompressor() TimestampDecompressor //Payload() []byte // для снапшота - Offset() int + // Offset() int Rotate([]byte) WritePayloadTo(io.Writer) error LastTimestamp() uint32 ReplaceSinceWithUntil() uint32 } +type TimeEvaluationReport struct { + TotalSpace int + Offset int + ChangeSize int +} + type ValueCompressor interface { - // (value) => compressionWay, requiredSpace - Evaluate(float64) (int, int) - Compress(int, float64) - Size() int + // (tmp, value) + Evaluate([]byte, float64) ValueEvaluationReport + // (offset, tmp, delta) + Append(int, []byte, uint64) + //Size() int //Chunks() [][]byte //DeleteLast() CaptureState() @@ -53,12 +61,19 @@ type ValueCompressor interface { // fracDigits CreateDecompressor() ValueDecompressor //Payload() []byte // для снапшота - Offset() int + //Offset() int Rotate([]byte) WritePayloadTo(io.Writer) error LastValue() float64 } +type ValueEvaluationReport struct { + TotalSpace int + Offset int + ChangeSize int + Delta uint64 +} + type TimestampDecompressor interface { NextValue() (uint32, bool) }