Files
qb/enc/time_delta.go

441 lines
10 KiB
Go
Raw Permalink Normal View History

2026-06-07 21:27:18 +03:00
package enc
2026-05-10 00:59:47 +00:00
import (
2026-06-09 08:16:16 +03:00
"io"
2026-05-10 00:59:47 +00:00
"log"
2026-06-05 19:43:01 +00:00
bin "gordenko.dev/dima/bin/little"
2026-05-31 20:01:28 +00:00
"gordenko.dev/dima/qb"
2026-05-10 00:59:47 +00:00
)
2026-06-05 19:43:01 +00:00
/*
Payload data файла має таку структуру:
vvvvvvvv-> <-ttttttt
де v - це values, які додаються зліва направо, а читаються справа наліво;
t - це timestamps, які додаються справа наліво, а читаються зліва направо.
Timestamps читаються звичайними методами Get*, а писатись мають типу ByRightBound
*/
2026-06-11 01:46:25 +00:00
const (
tmpTimeSize = 7 // max when start run
)
2026-06-07 01:06:16 +00:00
type TimeDeltaCompressor struct {
buf []byte
// останній записаний байт, якщо рахувати справа наліво
pos int
lastUnixtime uint32
lastDelta uint32
2026-06-12 06:11:57 +00:00
state *qb.TimeDeltaCapturedState
2026-05-10 00:59:47 +00:00
}
2026-06-09 08:16:16 +03:00
// Початкове створення, коли даних немає
// Створення після page filled, коли одразу є початкове значення
// Відновлення стану із снапшота. LastTimestamp передаю окермо від payload,
// а lastDelta і h треба прочитати із payload (якщо є)
// Кодує значення справа наліво функціями TailPut*. Читає значення зліва направо функціями Get*
2026-06-10 06:18:45 +03:00
func NewTimeDeltaCompressor(buf []byte, payloadSize int) *TimeDeltaCompressor {
s := &TimeDeltaCompressor{
2026-05-10 00:59:47 +00:00
buf: buf,
2026-06-11 01:46:25 +00:00
pos: len(buf) - payloadSize,
2026-06-07 01:06:16 +00:00
}
2026-06-07 21:27:18 +03:00
if payloadSize > 0 {
2026-06-12 06:11:57 +00:00
s.restore(payloadSize)
2026-05-10 00:59:47 +00:00
}
2026-06-10 06:18:45 +03:00
return s
2026-05-10 00:59:47 +00:00
}
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) restore(payloadSize int) {
bound := len(s.buf) - 4
since, err := bin.GetUint32(s.buf[bound:])
if err != nil {
log.Fatalf("bug: get since: %s", err)
}
if payloadSize == 4 {
s.lastUnixtime = since
return
}
var (
i = s.pos
totalDelta uint64
)
for i < bound {
h := s.buf[i]
i++
if h < 128 {
// run
count := h + 2
delta, n, _ := bin.GetVarUint64(s.buf[i:])
i += n
totalDelta += delta * uint64(count)
} else {
// literal
count := (h & 127) + 1
for range count {
delta, n, _ := bin.GetVarUint64(s.buf[i:])
i += n
totalDelta += delta
}
}
}
u64, _, err := bin.GetVarUint64(s.buf[s.pos+1:])
if err != nil {
log.Fatalf("bug: get last delta: %s", err)
}
s.lastDelta = uint32(u64)
s.lastUnixtime = since + uint32(totalDelta)
}
2026-06-11 01:46:25 +00:00
func (s *TimeDeltaCompressor) Evaluate(tmp []byte, timestamp uint32) qb.TimeEvaluationReport {
var (
2026-06-12 06:11:57 +00:00
rewindOffset int
i int
2026-06-11 01:46:25 +00:00
)
if s.pos < len(s.buf) {
2026-06-19 07:25:46 +03:00
delta := timestamp - s.lastUnixtime
2026-06-06 15:36:27 +03:00
if s.lastDelta > 0 {
2026-06-11 01:46:25 +00:00
h := s.buf[s.pos]
if h < 128 {
2026-06-06 15:36:27 +03:00
// run
2026-06-11 01:46:25 +00:00
if delta == s.lastDelta && h < 127 {
// incrementRun
tmp[i] = h + 1
i++
2026-06-12 06:11:57 +00:00
rewindOffset = 1 // перезапис h
2026-05-10 00:59:47 +00:00
} else {
2026-06-11 01:46:25 +00:00
// endSeries
tmp[i] = 128 // start new literal (length=1)
i++
2026-06-19 07:25:46 +03:00
n, _ := bin.PutVarUint64(tmp[i:], uint64(delta))
i += n
2026-05-10 00:59:47 +00:00
}
} else {
2026-06-06 15:36:27 +03:00
// literal
if delta != s.lastDelta {
2026-06-11 01:46:25 +00:00
if h < 255 {
// incrementLiteral
tmp[i] = h + 1
i++
2026-06-19 07:25:46 +03:00
n, _ := bin.PutVarUint64(tmp[i:], uint64(delta))
i += n
2026-06-12 06:11:57 +00:00
rewindOffset = 1 // перезапис h
2026-06-06 15:36:27 +03:00
} else {
2026-06-11 01:46:25 +00:00
// endSeries
tmp[i] = 128 // start new literal (length=1)
i++
2026-06-19 07:25:46 +03:00
n, _ := bin.PutVarUint64(tmp[i:], uint64(delta))
i += n
2026-06-06 15:36:27 +03:00
}
2026-05-10 00:59:47 +00:00
} else {
2026-06-11 01:46:25 +00:00
// startRun
if h > 128 {
2026-06-19 07:25:46 +03:00
tmp[i] = 0 // start new run (length=2)
2026-06-11 01:46:25 +00:00
i++
n, _ := bin.PutVarUint64(tmp[i:], uint64(delta))
i += n
2026-06-19 07:25:46 +03:00
tmp[i] = h - 1 // зменшую довжину попередньої серії на 1
2026-06-11 01:46:25 +00:00
i++
2026-06-12 06:11:57 +00:00
rewindOffset = 1 + n // перезапис пари delta/h
2026-06-11 01:46:25 +00:00
} else {
tmp[i] = 0 // change literal (length=1) to run (length=2)
i++
2026-06-12 06:11:57 +00:00
rewindOffset = 1
2026-06-07 01:06:16 +00:00
}
2026-05-10 00:59:47 +00:00
}
}
} else {
2026-06-11 01:46:25 +00:00
tmp[i] = 128 // start new literal (length=1)
i++
2026-06-19 07:25:46 +03:00
n, _ := bin.PutVarUint64(tmp[i:], uint64(delta))
i += n
2026-06-06 15:36:27 +03:00
}
} else {
2026-06-11 01:46:25 +00:00
bin.PutUint32(tmp, timestamp) // 1st timestamp (since)
i += 4
}
return qb.TimeEvaluationReport{
2026-06-12 06:11:57 +00:00
RewindOffset: rewindOffset,
2026-06-14 23:12:03 +03:00
Offset: len(s.buf) - s.pos - rewindOffset,
2026-06-12 06:11:57 +00:00
ChangeSize: i,
2026-06-15 01:20:30 +03:00
TotalSpace: len(s.buf) - s.pos - rewindOffset + i,
2026-06-06 15:36:27 +03:00
}
2026-05-10 00:59:47 +00:00
}
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) Append(rewindOffset int, change []byte, timestamp uint32) {
2026-06-19 07:25:46 +03:00
// fmt.Printf("----\nchange: % x\n", change)
// fmt.Printf("rewindOffset: %d\n", rewindOffset)
// fmt.Printf("pos: %d\n", s.pos)
// fmt.Printf("buf: %d\n", len(s.buf))
if s.lastUnixtime > 0 {
2026-06-11 01:46:25 +00:00
s.lastDelta = timestamp - s.lastUnixtime
2026-06-06 15:36:27 +03:00
}
2026-06-19 07:25:46 +03:00
idx := s.pos - len(change) + rewindOffset
copy(s.buf[idx:], change)
2026-06-12 06:11:57 +00:00
s.pos -= len(change) - rewindOffset
2026-06-11 01:46:25 +00:00
s.lastUnixtime = timestamp
2026-05-10 00:59:47 +00:00
}
2026-06-11 01:46:25 +00:00
// func (s *TimeDeltaCompressor) DeleteLast() {
2026-05-10 00:59:47 +00:00
2026-06-11 01:46:25 +00:00
// }
2026-05-31 20:01:28 +00:00
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) getState() qb.TimeDeltaCapturedState {
2026-06-14 23:12:03 +03:00
state := qb.TimeDeltaCapturedState{
2026-06-07 01:06:16 +00:00
LastUnixtime: s.lastUnixtime,
2026-05-31 20:01:28 +00:00
LastDelta: s.lastDelta,
2026-05-10 00:59:47 +00:00
}
2026-06-14 23:12:03 +03:00
if s.pos < len(s.buf) {
var bound int
if s.lastDelta > 0 {
state.H = s.buf[s.pos]
bound = s.pos + 1 + bin.CountVarUint64(uint64(s.lastDelta))
} else {
bound = s.pos // only since encoded
}
state.Payload = s.buf[bound:]
2026-06-19 07:25:46 +03:00
// fmt.Printf("% x\n", s.buf)
// fmt.Printf("% x\n", s.buf[bound:])
// fmt.Printf("h: % x\n", state.H)
// fmt.Printf("lastDelta: %d\n", s.lastDelta)
2026-06-14 23:12:03 +03:00
}
return state
2026-05-10 00:59:47 +00:00
}
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) CaptureState() {
state := s.getState()
s.state = &state
}
2026-06-13 07:13:03 +00:00
func (s *TimeDeltaCompressor) ForgetCapturedState() {
s.state = nil
}
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) Tail(offset int) []byte {
2026-06-19 07:25:46 +03:00
//fmt.Println("time tail:", s.pos, len(s.buf), offset)
2026-06-12 06:11:57 +00:00
return s.buf[s.pos : len(s.buf)-offset]
}
2026-06-09 08:16:16 +03:00
// коли сторінка заповнена і відправляється в txlog, since замінюю на until
func (s *TimeDeltaCompressor) ReplaceSinceWithUntil() uint32 {
pos := len(s.buf) - 4
since, _ := bin.GetUint32(s.buf[pos:])
bin.PutUint32(s.buf[pos:], s.lastUnixtime)
return since
}
2026-06-11 14:27:38 +00:00
// для зростання буфера під час вставки даних. State не цікавить
func (s *TimeDeltaCompressor) Size() int {
return len(s.buf) - s.pos
}
2026-06-14 23:12:03 +03:00
func (s *TimeDeltaCompressor) CommitedSize() int {
2026-06-12 06:11:57 +00:00
if s.state == nil {
return len(s.buf) - s.pos
} else {
2026-06-14 23:12:03 +03:00
if len(s.state.Payload) > 0 {
if s.state.LastDelta > 0 {
return bin.CountVarUint64(uint64(s.state.LastDelta)) + 1 + len(s.state.Payload)
} else {
return 4 // since
}
} else {
return 0
}
2026-06-12 06:11:57 +00:00
}
}
2026-06-09 08:16:16 +03:00
2026-06-07 01:06:16 +00:00
// Snapshot - для створення снапшота.
2026-06-09 08:16:16 +03:00
// FIX - encode firstUnixtime in 1st 4 bytes
//
// while page not completed - has first unixtime, after - last unixtime
2026-06-11 01:46:25 +00:00
// func (s *TimeDeltaCompressor) Payload() []byte {
// if s.state == nil {
// return s.buf[s.pos:]
// } else {
// return s.state.Payload
// }
// }
2026-06-07 01:06:16 +00:00
2026-06-14 23:12:03 +03:00
func (s *TimeDeltaCompressor) WriteCommitedTo(w io.Writer) (err error) {
2026-06-12 14:21:14 +00:00
if s.state == nil {
2026-06-14 23:12:03 +03:00
_, err = w.Write(s.buf[s.pos:])
2026-06-09 08:16:16 +03:00
return
} else {
2026-06-12 14:21:14 +00:00
_, err = w.Write(s.state.Payload)
2026-06-09 08:16:16 +03:00
if err != nil {
return
}
2026-06-12 14:21:14 +00:00
_, err = bin.WriteVarUint64(w, uint64(s.state.LastDelta))
2026-06-09 08:16:16 +03:00
if err != nil {
return
}
2026-06-09 16:42:17 +00:00
_, err = w.Write([]byte{
2026-06-12 14:21:14 +00:00
s.state.H,
2026-06-09 08:16:16 +03:00
})
return
}
2026-05-31 20:01:28 +00:00
}
2026-06-12 06:11:57 +00:00
func (s *TimeDeltaCompressor) ReplaceBuffer(newbuf []byte) {
2026-06-15 01:20:30 +03:00
//fmt.Printf("time replace buffer: new size %d\n", len(newbuf))
2026-06-07 21:27:18 +03:00
s.buf = newbuf
2026-06-07 01:06:16 +00:00
s.pos = len(s.buf)
s.lastUnixtime = 0
2026-05-31 20:01:28 +00:00
s.lastDelta = 0
2026-06-07 21:27:18 +03:00
}
func (s *TimeDeltaCompressor) CreateDecompressor() qb.TimestampDecompressor {
2026-06-11 01:46:25 +00:00
if s.pos == len(s.buf) {
return nil
}
2026-06-12 06:11:57 +00:00
d := NewTimeDeltaDecompressor()
if s.state != nil {
d.RestoreFromState(*s.state)
} else {
d.RestoreFromState(s.getState())
}
// fmt.Printf("h: % x\n", s.buf[s.pos])
// fmt.Printf("lastUnixtime: % x\n", s.lastUnixtime)
// fmt.Printf("lastDelta: % x\n", s.lastDelta)
// fmt.Printf("buf: % x\n", s.buf)
// fmt.Printf("pay: % x\n", s.buf[bound:])
return d
2026-05-31 20:01:28 +00:00
}
2026-06-15 05:32:15 +03:00
func (s *TimeDeltaCompressor) FirstTimestamp() uint32 {
pos := len(s.buf) - 4
timestamp, _ := bin.GetUint32(s.buf[pos:])
return timestamp
}
2026-06-15 10:47:24 +00:00
func (s *TimeDeltaCompressor) Until() uint32 {
2026-06-09 08:16:16 +03:00
return s.lastUnixtime
2026-05-21 23:38:52 +03:00
}
2026-06-15 10:47:24 +00:00
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
}
2026-05-10 00:59:47 +00:00
// DECOMPRESSOR
2026-06-07 01:06:16 +00:00
type TimeDeltaDecompressor struct {
2026-06-05 19:43:01 +00:00
buf []byte
2026-05-10 00:59:47 +00:00
pos int
lastDelta uint32
lastUnixtime uint32
isRun bool
pending int
done bool
}
2026-06-12 06:11:57 +00:00
func NewTimeDeltaDecompressor() *TimeDeltaDecompressor {
return new(TimeDeltaDecompressor)
2026-05-10 00:59:47 +00:00
}
2026-06-12 06:11:57 +00:00
// викликається для повних сторінок
func (s *TimeDeltaDecompressor) RestoreFromEnd(buf []byte) {
s.buf = buf
2026-06-14 23:12:03 +03:00
if len(s.buf) > 0 {
var err error
s.lastUnixtime, err = bin.GetUint32(buf[len(buf)-4:])
if err != nil {
log.Fatalf("bug: get last unixtime: %s", err)
}
if len(buf) > 4 {
s.readHeader()
s.readDelta()
}
} else {
s.done = true
2026-06-12 06:11:57 +00:00
}
// fmt.Printf("h: % x\n", s.buf[s.pos])
// fmt.Printf("lastUnixtime: %d\n", s.lastUnixtime)
// fmt.Printf("lastDelta: %d\n", s.lastDelta)
// fmt.Printf("buf: % x\n", s.buf)
2026-06-07 21:27:18 +03:00
}
2026-06-12 06:11:57 +00:00
// викликається для data level tail
func (s *TimeDeltaDecompressor) RestoreFromState(state qb.TimeDeltaCapturedState) {
s.buf = state.Payload
2026-06-14 23:12:03 +03:00
if len(s.buf) > 0 {
s.lastUnixtime = state.LastUnixtime
if state.LastDelta > 0 {
s.lastDelta = state.LastDelta
s.decodeHeaderByte(state.H)
}
} else {
s.done = true
2026-05-10 00:59:47 +00:00
}
}
2026-06-07 01:06:16 +00:00
func (s *TimeDeltaDecompressor) NextValue() (value uint32, done bool) {
2026-05-10 00:59:47 +00:00
if s.done {
return 0, true
}
value = s.lastUnixtime
2026-06-07 01:06:16 +00:00
if s.lastDelta > 0 {
s.lastUnixtime -= s.lastDelta
2026-05-10 00:59:47 +00:00
s.pending--
if s.pending > 0 {
if !s.isRun {
s.readDelta()
}
2026-06-12 06:11:57 +00:00
} else if s.pos < len(s.buf)-4 {
2026-05-10 00:59:47 +00:00
s.readHeader()
s.readDelta()
} else {
2026-06-07 01:06:16 +00:00
s.lastDelta = 0
2026-05-10 00:59:47 +00:00
}
} else {
2026-06-07 01:06:16 +00:00
s.done = true
2026-05-10 00:59:47 +00:00
}
return value, false
}
2026-06-07 01:06:16 +00:00
func (s *TimeDeltaDecompressor) readHeader() {
2026-06-05 19:43:01 +00:00
h := s.buf[s.pos]
2026-06-07 01:06:16 +00:00
s.pos++
2026-05-10 00:59:47 +00:00
s.decodeHeaderByte(h)
}
2026-06-07 01:06:16 +00:00
func (s *TimeDeltaDecompressor) decodeHeaderByte(h byte) {
2026-05-10 00:59:47 +00:00
s.isRun = h < 128
if s.isRun {
2026-06-07 01:06:16 +00:00
s.pending = int(h) + 2
2026-05-10 00:59:47 +00:00
} else {
s.pending = int(h&127) + 1
}
}
2026-06-07 01:06:16 +00:00
func (s *TimeDeltaDecompressor) readDelta() {
u64, n, err := bin.GetVarUint64(s.buf[s.pos:])
2026-05-10 00:59:47 +00:00
if err != nil {
log.Fatalln(err)
}
2026-06-07 01:06:16 +00:00
s.pos += n
2026-05-10 00:59:47 +00:00
s.lastDelta = uint32(u64)
}