This commit is contained in:
2026-06-14 23:12:03 +03:00
parent 181031a753
commit b32279deaf
35 changed files with 1395 additions and 1062 deletions

View File

@@ -149,26 +149,32 @@ func (s *Connection) DeleteMetric(metricID uint32) (err error) {
return s.mustSuccess(s.src)
}
func (s *Connection) AppendMeasure(req proto.AppendMeasureReq) (err error) {
arr := []byte{
proto.TypeAppendMeasure,
0, 0, 0, 0, // metricID
0, 0, 0, 0, // timestamp
0, 0, 0, 0, 0, 0, 0, 0, // value
}
bin.PutUint32(arr[1:], req.MetricID)
bin.PutUint32(arr[5:], req.Timestamp)
bin.PutFloat64(arr[9:], req.Value)
// req
if _, err = s.conn.Write(arr); err != nil {
return
}
return s.mustSuccess(s.src)
// func (s *Connection) AppendMeasure(req proto.AppendMeasureReq) (err error) {
// arr := []byte{
// proto.TypeAppendMeasure,
// 0, 0, 0, 0, // metricID
// 0, 0, 0, 0, // timestamp
// 0, 0, 0, 0, 0, 0, 0, 0, // value
// }
// bin.PutUint32(arr[1:], req.MetricID)
// bin.PutUint32(arr[5:], req.Timestamp)
// bin.PutFloat64(arr[9:], req.Value)
// // req
// if _, err = s.conn.Write(arr); err != nil {
// return
// }
// return s.mustSuccess(s.src)
// }
type AppendMeasuresResult struct {
WrittenCount int
ErrorCode byte
}
func (s *Connection) AppendMeasures(req proto.AppendMeasuresReq) (err error) {
func (s *Connection) AppendMeasures(req proto.AppendMeasuresReq) (result AppendMeasuresResult, err error) {
if len(req.Measures) > 65535 {
return fmt.Errorf("wrong measures qty: %d", len(req.Measures))
err = fmt.Errorf("wrong measures qty: %d", len(req.Measures))
return
}
var (
prefixSize = 7
@@ -188,37 +194,6 @@ func (s *Connection) AppendMeasures(req proto.AppendMeasuresReq) (err error) {
if _, err = s.conn.Write(arr); err != nil {
return
}
return s.mustSuccess(s.src)
}
// type AppendMeasurePerMetricReq struct {
// MetricID uint32
// Measures []Measure
// }
func (s *Connection) AppendMeasurePerMetric(list []proto.MetricMeasure) (_ []proto.AppendError, err error) {
if len(list) > 65535 {
return nil, fmt.Errorf("wrong measures qty: %d", len(list))
}
var (
// 3 bytes: 1b message type + 2b records qty
fixedSize = 3
recordSize = 16
arr = make([]byte, fixedSize+len(list)*recordSize)
)
arr[0] = proto.TypeAppendMeasures
bin.PutUint16(arr[1:], uint16(len(list)))
pos := fixedSize
for _, item := range list {
bin.PutUint32(arr[pos:], item.MetricID)
bin.PutUint32(arr[pos+4:], item.Timestamp)
bin.PutFloat64(arr[pos+8:], item.Value)
pos += recordSize
}
// req
if _, err = s.conn.Write(arr); err != nil {
return
}
// answer
code, err := s.src.ReadByte()
if err != nil {
@@ -226,32 +201,82 @@ func (s *Connection) AppendMeasurePerMetric(list []proto.MetricMeasure) (_ []pro
}
switch code {
case proto.RespValue:
var (
count int
appendErrors []proto.AppendError
)
count, err = bin.ReadUint16AsInt(s.src)
result.WrittenCount, err = bin.ReadUint16AsInt(s.src)
if err != nil {
return
}
for range count {
var ae proto.AppendError
ae.MetricID, err = bin.ReadUint32(s.src)
if err != nil {
return
}
ae.ErrorCode, err = bin.ReadUint16(s.src)
if err != nil {
return
}
appendErrors = append(appendErrors, ae)
result.ErrorCode, err = bin.ReadByte(s.src)
if err != nil {
return
}
return appendErrors, nil
default:
return nil, fmt.Errorf("unknown reponse code %d", code)
err = fmt.Errorf("unknown reponse code %d", code)
return
}
return
}
// type AppendMeasurePerMetricReq struct {
// MetricID uint32
// Measures []Measure
// }
// func (s *Connection) AppendMeasurePerMetric(list []proto.MetricMeasure) (_ []proto.AppendError, err error) {
// if len(list) > 65535 {
// return nil, fmt.Errorf("wrong measures qty: %d", len(list))
// }
// var (
// // 3 bytes: 1b message type + 2b records qty
// fixedSize = 3
// recordSize = 16
// arr = make([]byte, fixedSize+len(list)*recordSize)
// )
// arr[0] = proto.TypeAppendMeasures
// bin.PutUint16(arr[1:], uint16(len(list)))
// pos := fixedSize
// for _, item := range list {
// bin.PutUint32(arr[pos:], item.MetricID)
// bin.PutUint32(arr[pos+4:], item.Timestamp)
// bin.PutFloat64(arr[pos+8:], item.Value)
// pos += recordSize
// }
// // req
// if _, err = s.conn.Write(arr); err != nil {
// return
// }
// // answer
// code, err := s.src.ReadByte()
// if err != nil {
// return
// }
// switch code {
// case proto.RespValue:
// var (
// count int
// appendErrors []proto.AppendError
// )
// count, err = bin.ReadUint16AsInt(s.src)
// if err != nil {
// return
// }
// for range count {
// var ae proto.AppendError
// ae.MetricID, err = bin.ReadUint32(s.src)
// if err != nil {
// return
// }
// ae.ErrorCode, err = bin.ReadUint16(s.src)
// if err != nil {
// return
// }
// appendErrors = append(appendErrors, ae)
// }
// return appendErrors, nil
// default:
// return nil, fmt.Errorf("unknown reponse code %d", code)
// }
// }
func (s *Connection) ListAllInstantMeasures(metricID uint32) (_ []proto.InstantMeasure, err error) {
arr := []byte{
proto.TypeListAllInstantMeasures,
@@ -361,8 +386,10 @@ func (s *Connection) readCumulativeMeasures() (_ []proto.CumulativeMeasure, err
if err != nil {
return nil, fmt.Errorf("read response code: %s", err)
}
fmt.Println("code", code)
switch code {
case proto.RespPartOfValue:
fmt.Println("RespPartOfValue")
var count int
count, err = bin.ReadUint32AsInt(s.src)
if err != nil {
@@ -385,8 +412,10 @@ func (s *Connection) readCumulativeMeasures() (_ []proto.CumulativeMeasure, err
result = append(result, measure)
}
case proto.RespEndOfValue:
fmt.Println("RespEndOfValue")
return result, nil
case proto.RespError:
fmt.Println("RespError")
return nil, s.onError()
default:
return nil, fmt.Errorf("unknown reponse code %d", code)

View File

@@ -11,6 +11,7 @@ import (
"gordenko.dev/dima/qb/atree"
"gordenko.dev/dima/qb/bufreader"
"gordenko.dev/dima/qb/proto"
"gordenko.dev/dima/qb/storage"
"gordenko.dev/dima/qb/transform"
"gordenko.dev/dima/qb/worker"
)
@@ -60,7 +61,7 @@ func (s *Database) handleTCPConn(conn net.Conn) {
}
func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (err error) {
fmt.Println("process request")
//fmt.Println("process request")
messageType, err := r.ReadByte()
if err != nil {
if err != io.EOF {
@@ -70,7 +71,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e
}
}
fmt.Println("messageType:", messageType)
//fmt.Println("messageType:", messageType)
switch messageType {
case proto.TypeGetMetric:
@@ -111,7 +112,9 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e
return fmt.Errorf("proto.ReadAppendMeasuresReq: %s", err)
}
//fmt.Println("append measure", req.MetricID, conn.RemoteAddr().String())
reply(conn, s.AppendMeasures(req))
if err = s.AppendMeasures(conn, req); err != nil {
return fmt.Errorf("AppendMeasures: %s", err)
}
case proto.TypeListInstantMeasures:
req, err := proto.ReadListInstantMeasuresReq(r)
@@ -195,7 +198,7 @@ func (s *Database) processRequest(conn net.Conn, r *bufreader.BufferedReader) (e
// API
func (s *Database) AddMetric(req proto.AddMetricReq) uint16 {
fmt.Println("database.AddMetric")
//fmt.Println("database.AddMetric")
// Валидация
if req.MetricID == 0 {
return proto.ErrEmptyMetricID
@@ -215,8 +218,10 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 {
fmt.Println("add job")
s.workerInbox.Push(worker.AddMetricReq{
MetricID: req.MetricID,
ResultCh: resultCh,
MetricID: req.MetricID,
MetricType: req.MetricType,
FracDigits: req.FracDigits,
ResultCh: resultCh,
})
fmt.Println("job added")
@@ -226,19 +231,13 @@ func (s *Database) AddMetric(req proto.AddMetricReq) uint16 {
switch resultCode {
case worker.Succeed:
fmt.Println("OK")
// s.storage.Append(storage.MetricAddRecord{
// MetricID: req.MetricID,
// MetricType: req.MetricType,
// FracDigits: req.FracDigits,
// })
//<-waitCh
case worker.MetricDuplicate:
fmt.Println("ErrDuplicate")
//fmt.Println("ErrDuplicate")
return proto.ErrDuplicate
default:
fmt.Println("ErrWrongResultCodeBug")
//fmt.Println("ErrWrongResultCodeBug")
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
}
return 0
@@ -325,8 +324,8 @@ type FilledPage struct {
ValuesSize uint16
}
func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 {
resultCh := make(chan worker.AppendMeasuresResult, 1)
func (s *Database) AppendMeasures(conn io.Writer, req proto.AppendMeasuresReq) error {
resultCh := make(chan storage.MeasuresAppendResult, 1)
s.workerInbox.Push(worker.AppendMeasuresReq{
MetricID: req.MetricID,
@@ -337,32 +336,30 @@ func (s *Database) AppendMeasures(req proto.AppendMeasuresReq) uint16 {
result := <-resultCh
switch result.ResultCode {
case worker.CanAppend:
// report, err := s.atree.AppendDataPage(atree.AppendDataPageReq{})
// if err != nil {
// qb.Abort(qb.WriteToAtreeFailed, err)
// }
// _ = report
// if len(toAppendMeasures) > 0 {
// waitCh := s.storage.WriteAppendMeasures(
// storage.AppendedMeasures{
// MetricID: req.MetricID,
// Measures: toAppendMeasures,
// },
// false,
// )
// <-waitCh
// }
case worker.Succeed, worker.NonMonotonicValue, worker.ExpiredMeasure:
var errorCode byte
switch result.ResultCode {
case worker.NonMonotonicValue:
errorCode = proto.ErrNonMonotonicValue
case worker.ExpiredMeasure:
errorCode = proto.ErrExpiredMeasure
}
answer := []byte{
proto.RespValue,
0, 0, // written count
errorCode,
}
bin.PutIntAsUint16(answer[1:], result.WrittenCount)
_, err := conn.Write(answer)
if err != nil {
return err
}
case worker.NoMetric:
return proto.ErrNoMetric
reply(conn, proto.ErrNoMetric)
default:
qb.Abort(qb.WrongResultCodeBug, ErrWrongResultCodeBug)
}
return 0
return nil
}
func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 {
@@ -380,25 +377,25 @@ func (s *Database) DeleteMeasures(req proto.DeleteMeasuresReq) uint16 {
case worker.NoMeasuresToDelete:
// ok
case worker.DeleteFromAtreeNotNeeded:
// регистрирую удаление в TransactionLog
// s.storage.Append(storage.MeasuresDeleteRecord{
// MetricID: req.MetricID,
// })
//<-waitCh
//case worker.DeleteFromAtreeNotNeeded:
// регистрирую удаление в TransactionLog
// s.storage.Append(storage.MeasuresDeleteRecord{
// MetricID: req.MetricID,
// })
//<-waitCh
case worker.DeleteFromAtreeRequired:
// собираю номера всех data и index страниц метрики (типа запись REDO лога).
// pageNumbers, err := s.atree.GetAllPages(req.MetricID)
// if err != nil {
// qb.Abort(qb.FailedAtreeRequest, err)
// }
// // регистрирую удаление в TransactionLog
// s.storage.Append(storage.DeletedMeasures{
// MetricID: req.MetricID,
// FreePageNumbers: pageNumbers,
// })
//<-waitCh
//case worker.DeleteFromAtreeRequired:
// собираю номера всех data и index страниц метрики (типа запись REDO лога).
// pageNumbers, err := s.atree.GetAllPages(req.MetricID)
// if err != nil {
// qb.Abort(qb.FailedAtreeRequest, err)
// }
// // регистрирую удаление в TransactionLog
// s.storage.Append(storage.DeletedMeasures{
// MetricID: req.MetricID,
// FreePageNumbers: pageNumbers,
// })
//<-waitCh
case worker.NoMetric:
return proto.ErrNoMetric
@@ -615,6 +612,7 @@ func (s *Database) fullScan(req fullScanReq) error {
switch result.ResultCode {
case worker.QueryDone:
fmt.Printf("query done")
req.ResponseWriter.Close()
case worker.UntilFound:

View File

@@ -9,7 +9,6 @@ import (
"sync"
"time"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/inbox"
"gordenko.dev/dima/qb/recovery"
"gordenko.dev/dima/qb/storage"
@@ -69,21 +68,20 @@ func New(opt Options) (_ *Database, err error) {
s := &Database{
dir: opt.Dir,
databaseName: opt.DatabaseName,
//metricLockEntries: make(map[uint32]*metricLockEntry),
//dataFreeList: dataFreeList,
//indexFreeList: indexFreeList,
tcpPort: opt.TCPPort,
logfile: opt.Logfile,
logger: log.New(opt.Logfile, "", log.LstdFlags),
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
tcpPort: opt.TCPPort,
logfile: opt.Logfile,
logger: log.New(opt.Logfile, "", log.LstdFlags),
exitCh: opt.ExitCh,
waitGroup: opt.WaitGroup,
}
recoveryReport, err := recovery.Recovery(s.dir, s.databaseName)
if err != nil {
return
return nil, fmt.Errorf("recovery.Recovery: %s", err)
}
fmt.Printf("%#v\n", recoveryReport)
storageInbox := inbox.New()
s.workerInbox = inbox.New()
@@ -99,18 +97,21 @@ func New(opt Options) (_ *Database, err error) {
WaitGroup: s.waitGroup,
})
if err != nil {
qb.Abort(qb.CreateChangesWriterFailed, err)
return nil, fmt.Errorf("storage.NewWriter: %s", err)
}
fmt.Println("storage created")
s.worker = worker.New(worker.Options{
Inbox: s.workerInbox,
StorageInbox: storageInbox,
//Metrics: make(map[uint32]*_metric), FIX
Dir: opt.Dir,
ExitCh: opt.ExitCh,
WaitGroup: opt.WaitGroup,
Inbox: s.workerInbox,
StorageInbox: storageInbox,
ReplayMetrics: recoveryReport.ReplayMetrics,
Dir: opt.Dir,
DatabaseName: opt.DatabaseName,
ExitCh: opt.ExitCh,
WaitGroup: opt.WaitGroup,
})
fmt.Println("worker created")
return s, nil
}
@@ -120,6 +121,7 @@ func (s *Database) ListenAndServe() (err error) {
return fmt.Errorf("net.Listen: %s; port=%d", err, s.tcpPort)
}
//s.waitGroup.Add(1)
go s.storage.Run()
// s.atree, err = atree.New(atree.Options{

Binary file not shown.

View File

@@ -31,9 +31,25 @@ var (
}{
{
Nums: []float64{
0,
},
Name: "add 1st value zero",
Offset: 0,
ChangeSize: 3,
Buf: []byte{
0x8f, // base value
0x80, // delta 0
0x80, // literal (len = 1)
0x00, 0x00, 0x00, 0x00, 0x00,
},
BaseValue: 0,
LastDelta: 0,
},
{
Nums: []float64{
1.5,
},
Name: "add 1st value",
Name: "add 1st value non zero",
Offset: 0,
ChangeSize: 3,
Buf: []byte{

View File

@@ -1,6 +1,7 @@
package enc
import (
"fmt"
"io"
"log"
@@ -159,7 +160,7 @@ func (s *TimeDeltaCompressor) Evaluate(tmp []byte, timestamp uint32) qb.TimeEval
}
return qb.TimeEvaluationReport{
RewindOffset: rewindOffset,
Offset: s.pos + rewindOffset,
Offset: len(s.buf) - s.pos - rewindOffset,
ChangeSize: i,
TotalSpace: s.pos - rewindOffset + i,
}
@@ -185,13 +186,21 @@ func (s *TimeDeltaCompressor) Append(rewindOffset int, change []byte, timestamp
// }
func (s *TimeDeltaCompressor) getState() qb.TimeDeltaCapturedState {
bound := s.pos + 1 + bin.CountVarUint64(uint64(s.lastDelta))
return qb.TimeDeltaCapturedState{
H: s.buf[s.pos],
state := qb.TimeDeltaCapturedState{
LastUnixtime: s.lastUnixtime,
LastDelta: s.lastDelta,
Payload: s.buf[bound:],
}
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:]
}
return state
}
func (s *TimeDeltaCompressor) CaptureState() {
@@ -204,6 +213,7 @@ func (s *TimeDeltaCompressor) ForgetCapturedState() {
}
func (s *TimeDeltaCompressor) Tail(offset int) []byte {
fmt.Println("time tail:", s.pos, len(s.buf), offset)
return s.buf[s.pos : len(s.buf)-offset]
}
@@ -220,11 +230,19 @@ func (s *TimeDeltaCompressor) Size() int {
return len(s.buf) - s.pos
}
func (s *TimeDeltaCompressor) StoredSize() int {
func (s *TimeDeltaCompressor) CommitedSize() int {
if s.state == nil {
return len(s.buf) - s.pos
} else {
return bin.CountVarUint64(uint64(s.state.LastDelta)) + 1 + len(s.state.Payload)
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
}
}
}
@@ -240,9 +258,9 @@ func (s *TimeDeltaCompressor) StoredSize() int {
// }
// }
func (s *TimeDeltaCompressor) WriteStoredTo(w io.Writer) (err error) {
func (s *TimeDeltaCompressor) WriteCommitedTo(w io.Writer) (err error) {
if s.state == nil {
_, err = w.Write(s.buf[:s.pos])
_, err = w.Write(s.buf[s.pos:])
return
} else {
_, err = w.Write(s.state.Payload)
@@ -308,14 +326,18 @@ func NewTimeDeltaDecompressor() *TimeDeltaDecompressor {
// викликається для повних сторінок
func (s *TimeDeltaDecompressor) RestoreFromEnd(buf []byte) {
s.buf = buf
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()
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
}
// fmt.Printf("h: % x\n", s.buf[s.pos])
// fmt.Printf("lastUnixtime: %d\n", s.lastUnixtime)
@@ -326,10 +348,14 @@ func (s *TimeDeltaDecompressor) RestoreFromEnd(buf []byte) {
// викликається для data level tail
func (s *TimeDeltaDecompressor) RestoreFromState(state qb.TimeDeltaCapturedState) {
s.buf = state.Payload
s.lastUnixtime = state.LastUnixtime
if state.LastDelta > 0 {
s.lastDelta = state.LastDelta
s.decodeHeaderByte(state.H)
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
}
}

View File

@@ -1,6 +1,7 @@
package enc
import (
"fmt"
"io"
"log"
"math"
@@ -48,9 +49,6 @@ func NewValueDeltaCompressor(metricType qb.MetricType, fracDigits byte, buf []by
s.toFloat64 = toInstantFloat64
}
if payloadSize > 0 {
// fmt.Printf("% x\n", buf)
// fmt.Printf("% x\n", buf[:s.pos])
// fmt.Printf("pos: %d\n", s.pos)
u64, _, err := bin.GetVarUint64(s.buf)
if err != nil {
log.Fatalf("bug: get base value: %s", err)
@@ -145,13 +143,21 @@ func (s *ValueDeltaCompressor) Evaluate(tmp []byte, value float64) qb.ValueEvalu
// pos завжди вказує на h
func (s *ValueDeltaCompressor) Append(rewindOffset int, change []byte, value float64, delta uint64) {
// fmt.Printf("append value:\n")
// fmt.Printf("rewind offset: %d\n", rewindOffset)
// fmt.Printf("change: % x\n", change)
// fmt.Printf("value: %v\n", value)
// fmt.Printf("delta: %d\n", delta)
if s.pos > 0 {
s.lastDelta = delta
} else {
s.baseValue = value
}
//fmt.Printf("buf before: % x\n", s.buf[:s.pos])
copy(s.buf[s.pos-rewindOffset:], change)
s.pos += len(change) - rewindOffset
//fmt.Printf("buf after: % x\n", s.buf[:s.pos])
//fmt.Printf("buf after: % x\n", s.buf[:s.pos])
}
func (s *ValueDeltaCompressor) DeleteLast() {
@@ -159,11 +165,20 @@ func (s *ValueDeltaCompressor) DeleteLast() {
}
func (s *ValueDeltaCompressor) getState() qb.ValueDeltaCapturedState {
bound := s.pos - 1 - bin.CountVarUint64(s.lastDelta)
return qb.ValueDeltaCapturedState{
H: s.buf[s.pos-1],
LastDelta: s.lastDelta,
Payload: s.buf[:bound],
if s.pos > 0 {
// fmt.Println("getState")
// fmt.Printf("pos: %d\n", s.pos)
// fmt.Printf("buf: % x\n", s.buf[:s.pos])
// fmt.Printf("h: %d\n", s.buf[s.pos-1])
// fmt.Printf("lastDelta: %d\n", s.lastDelta)
bound := s.pos - 1 - bin.CountVarUint64(s.lastDelta)
return qb.ValueDeltaCapturedState{
H: s.buf[s.pos-1],
LastDelta: s.lastDelta,
Payload: s.buf[:bound],
}
} else {
return qb.ValueDeltaCapturedState{}
}
}
@@ -185,11 +200,15 @@ func (s *ValueDeltaCompressor) Size() int {
return s.pos
}
func (s *ValueDeltaCompressor) StoredSize() int {
func (s *ValueDeltaCompressor) CommitedSize() int {
if s.state == nil {
return len(s.buf) - s.pos
return s.pos
} else {
return bin.CountVarUint64(uint64(s.state.LastDelta)) + 1 + len(s.state.Payload)
if len(s.state.Payload) > 0 {
return bin.CountVarUint64(uint64(s.state.LastDelta)) + 1 + len(s.state.Payload)
} else {
return 0
}
}
}
@@ -202,7 +221,7 @@ func (s *ValueDeltaCompressor) StoredSize() int {
// }
// }
func (s *ValueDeltaCompressor) WriteStoredTo(w io.Writer) (err error) {
func (s *ValueDeltaCompressor) WriteCommitedTo(w io.Writer) (err error) {
if s.state == nil {
_, err = w.Write(s.buf[:s.pos])
return
@@ -289,17 +308,26 @@ func NewValueDeltaDecompressor(metricType qb.MetricType, fracDigits byte) *Value
func (s *ValueDeltaDecompressor) RestoreFromEnd(buf []byte) {
s.buf = buf
s.pos = len(buf) // first free
s.readBaseValue()
s.readHeader()
s.readValue()
if len(buf) > 0 {
s.readBaseValue()
s.readHeader()
s.readValue()
} else {
s.done = true
}
}
func (s *ValueDeltaDecompressor) RestoreFromState(state qb.ValueDeltaCapturedState) {
fmt.Printf("%#v\n", state)
s.buf = state.Payload
s.pos = len(s.buf) // first free
s.readBaseValue()
s.lastValue = s.baseValue + s.toFloat64(state.LastDelta, s.coef)
s.decodeHeaderByte(state.H)
if len(s.buf) > 0 {
s.readBaseValue()
s.lastValue = s.baseValue + s.toFloat64(state.LastDelta, s.coef)
s.decodeHeaderByte(state.H)
} else {
s.done = true
}
}
func (s *ValueDeltaDecompressor) readBaseValue() {

View File

@@ -101,51 +101,54 @@ func sendRequests(conn *client.Connection) {
err error
)
_ = cumulativeMetricID
_ = fracDigits
//conn.DeleteMetric(instantMetricID)
//conn.DeleteMetric(cumulativeMetricID)
// ADD INSTANT METRIC
err = conn.AddMetric(proto.AddMetricReq{
MetricID: cumulativeMetricID,
MetricType: qb.Instant,
FracDigits: fracDigits,
})
if err != nil {
log.Fatalf("conn.AddMetric: %s\n", err)
} else {
fmt.Printf("\nInstant metric %d added\n", cumulativeMetricID)
}
// err = conn.AddMetric(proto.AddMetricReq{
// MetricID: cumulativeMetricID,
// MetricType: qb.Instant,
// FracDigits: fracDigits,
// })
// if err != nil {
// log.Fatalf("conn.AddMetric: %s\n", err)
// } else {
// fmt.Printf("\nInstant metric %d added\n", cumulativeMetricID)
// }
// GET INSTANT METRIC
// iMetric, err := conn.GetMetric(cumulativeMetricID)
// if err != nil {
// log.Fatalf("conn.GetMetric: %s\n", err)
// } else {
// fmt.Printf(`
// iMetric, err := conn.GetMetric(cumulativeMetricID)
// if err != nil {
// log.Fatalf("conn.GetMetric: %s\n", err)
// } else {
// fmt.Printf(`
// GetMetric:
// metricID: %d
// metricType: %s
// fracDigits: %d
// `,
// iMetric.MetricID, metricTypeToName[iMetric.MetricType], fracDigits)
//}
// iMetric.MetricID, metricTypeToName[iMetric.MetricType], fracDigits)
// }
// APPEND MEASURES
// instantMeasures := GenerateInstantMeasures(62, 220)
// instantMeasures := GenerateInstantMeasures(62, 220)
// err = conn.AppendMeasures(proto.AppendMeasuresReq{
// MetricID: instantMetricID,
// Measures: instantMeasures,
// })
// if err != nil {
// log.Fatalf("conn.AppendMeasures: %s\n", err)
// } else {
// fmt.Printf("\nAppended %d measures for the metric %d\n",
// len(instantMeasures), instantMetricID)
// }
// err = conn.AppendMeasures(proto.AppendMeasuresReq{
// MetricID: instantMetricID,
// Measures: instantMeasures,
// })
// if err != nil {
// log.Fatalf("conn.AppendMeasures: %s\n", err)
// } else {
// fmt.Printf("\nAppended %d measures for the metric %d\n",
// len(instantMeasures), instantMetricID)
// }
// // LIST INSTANT MEASURES
@@ -287,50 +290,62 @@ func sendRequests(conn *client.Connection) {
// fmt.Printf("\nInstant metric %d deleted\n", instantMetricID)
// }
// // ADD CUMULATIVE METRIC
// ADD CUMULATIVE METRIC
// err = conn.AddMetric(proto.AddMetricReq{
// MetricID: cumulativeMetricID,
// MetricType: qb.Cumulative,
// FracDigits: fracDigits,
// })
// if err != nil {
// log.Fatalf("conn.AddMetric: %s\n", err)
// } else {
// fmt.Printf("\nCumulative metric %d added\n", cumulativeMetricID)
err = conn.AddMetric(proto.AddMetricReq{
MetricID: cumulativeMetricID,
MetricType: qb.Cumulative,
FracDigits: fracDigits,
})
if err != nil {
log.Fatalf("conn.AddMetric: %s\n", err)
} else {
fmt.Printf("\nCumulative metric %d added\n", cumulativeMetricID)
}
// GET CUMULATIVE METRIC
cMetric, err := conn.GetMetric(cumulativeMetricID)
if err != nil {
log.Fatalf("conn.GetMetric: %s\n", err)
} else {
fmt.Printf(`
GetMetric:
metricID: %d
metricType: %s
fracDigits: %d
`,
cMetric.MetricID, metricTypeToName[cMetric.MetricType], cMetric.FracDigits)
}
// APPEND MEASURES
cumulativeMeasures := GenerateCumulativeMeasures(62)
result, err := conn.AppendMeasures(proto.AppendMeasuresReq{
MetricID: cumulativeMetricID,
Measures: cumulativeMeasures[:1],
})
if err != nil {
log.Fatalf("conn.AppendMeasures: %s\n", err)
} else {
fmt.Printf("\nAppended %d measures for the metric %d: count=%d, errorCode=%d\n",
len(cumulativeMeasures), cumulativeMetricID, result.WrittenCount, result.ErrorCode)
}
// currentValues, err := conn.ListCurrentValues([]uint32{
// cumulativeMetricID,
// })
// if err != nil {
// log.Fatalf("conn.ListCurrentValues: %s\n", err)
// } else {
// for _, v := range currentValues {
// fmt.Println(v.MetricID, v.Timestamp, v.Value)
// }
// }
// // GET CUMULATIVE METRIC
// cMetric, err := conn.GetMetric(cumulativeMetricID)
// if err != nil {
// log.Fatalf("conn.GetMetric: %s\n", err)
// } else {
// fmt.Printf(`
// GetMetric:
// metricID: %d
// metricType: %s
// fracDigits: %d
// `,
// cMetric.MetricID, metricTypeToName[cMetric.MetricType], fracDigits)
// }
// // APPEND MEASURES
// cumulativeMeasures := GenerateCumulativeMeasures(62)
// err = conn.AppendMeasures(proto.AppendMeasuresReq{
// MetricID: cumulativeMetricID,
// Measures: cumulativeMeasures,
// })
// if err != nil {
// log.Fatalf("conn.AppendMeasures: %s\n", err)
// } else {
// fmt.Printf("\nAppended %d measures for the metric %d\n",
// len(cumulativeMeasures), cumulativeMetricID)
// }
// // LIST CUMULATIVE MEASURES
// LIST CUMULATIVE MEASURES
var cumulativeList []proto.CumulativeMeasure
// lastTimestamp = cumulativeMeasures[len(cumulativeMeasures)-1].Timestamp
// until = time.Unix(int64(lastTimestamp), 0)
@@ -352,17 +367,18 @@ func sendRequests(conn *client.Connection) {
// }
// }
// // LIST ALL CUMULATIVE MEASURES
// LIST ALL CUMULATIVE MEASURES
// cumulativeList, err = conn.ListAllCumulativeMeasures(cumulativeMetricID)
// if err != nil {
// log.Fatalf("conn.ListAllCumulativeMeasures: %s\n", err)
// } else {
// fmt.Printf("\nListAllCumulativeMeasures (last 15 items):\n")
// for _, item := range cumulativeList[:15] {
// fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value)
// }
// }
cumulativeList, err = conn.ListAllCumulativeMeasures(cumulativeMetricID)
if err != nil {
log.Fatalf("conn.ListAllCumulativeMeasures: %s\n", err)
} else {
fmt.Printf("\nListAllCumulativeMeasures (last 15 items):\n")
fmt.Printf("%#v\n", cumulativeList)
for _, item := range cumulativeList {
fmt.Printf(" %s => %.2f\n", formatTime(item.Timestamp), item.Value)
}
}
// // LIST CUMULATIVE PERIODS (group by hour)

Binary file not shown.

View File

@@ -53,15 +53,15 @@ GetMetric:
instantMeasures := GenerateInstantMeasures(62, 220)
err = conn.AppendMeasures(proto.AppendMeasuresReq{
result, err := conn.AppendMeasures(proto.AppendMeasuresReq{
MetricID: instantMetricID,
Measures: instantMeasures,
})
if err != nil {
log.Fatalf("conn.AppendMeasures: %s\n", err)
} else {
fmt.Printf("\nAppended %d measures for the metric %d\n",
len(instantMeasures), instantMetricID)
fmt.Printf("\nAppended %d measures for the metric %d: writtenCount=%d, errorCode=%d\n",
len(instantMeasures), instantMetricID, result.WrittenCount, result.ErrorCode)
}
// LIST INSTANT MEASURES
@@ -236,15 +236,15 @@ GetMetric:
cumulativeMeasures := GenerateCumulativeMeasures(62)
err = conn.AppendMeasures(proto.AppendMeasuresReq{
result, err = conn.AppendMeasures(proto.AppendMeasuresReq{
MetricID: cumulativeMetricID,
Measures: cumulativeMeasures,
})
if err != nil {
log.Fatalf("conn.AppendMeasures: %s\n", err)
} else {
fmt.Printf("\nAppended %d measures for the metric %d\n",
len(cumulativeMeasures), cumulativeMetricID)
fmt.Printf("\nAppended %d measures for the metric %d: writtenCount=%d, errorCode=%d\n",
len(cumulativeMeasures), cumulativeMetricID, result.WrittenCount, result.ErrorCode)
}
// LIST CUMULATIVE MEASURES

4
go.sum
View File

@@ -18,10 +18,6 @@ gopkg.in/ini.v1 v1.67.1/go.mod h1:x/cyOwCgZqOkJoDIJ3c1KNHMo10+nLGAhh+kn3Zizss=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gordenko.dev/dima/bin v0.0.0-20260606153512-101b10d92a39 h1:qW3SnQ9HAcJpfj8ottBfPGTzHX+cX2vM8B6TVOgnypE=
gordenko.dev/dima/bin v0.0.0-20260606153512-101b10d92a39/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ=
gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696 h1:JY91JYvv1g+XbaKdGmLFDBueUTx51RdW9mnbNYhiUdw=
gordenko.dev/dima/bin v0.0.0-20260608125602-78363e903696/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ=
gordenko.dev/dima/bin v0.0.0-20260612161453-4ee9be3474fb h1:XBYZbK5Z+yfpKmSR2wuPLyW14ReCRSNdOlL5MWro8ns=
gordenko.dev/dima/bin v0.0.0-20260612161453-4ee9be3474fb/go.mod h1:up64wpJp9xI+HqACtMDBIfPNUH6BYY/b3inKpUfrqDQ=
gordenko.dev/dima/pretty v0.0.0-20221225212746-0c27d8c0ac69 h1:nyJ3mzTQ46yUeMZCdLyYcs7B5JCS54c67v84miyhq2E=

View File

@@ -1,6 +1,8 @@
package inbox
import "sync"
import (
"sync"
)
type Inbox struct {
mutex *sync.Mutex
@@ -20,9 +22,11 @@ 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:
@@ -30,9 +34,11 @@ 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
}

8
qb.go
View File

@@ -30,7 +30,7 @@ type TimestampCompressor interface {
// (rewindOffset, change, timestamp)
Append(int, []byte, uint32)
Size() int
StoredSize() int
CommitedSize() int
//DeleteLast()
CaptureState()
ForgetCapturedState()
@@ -40,7 +40,7 @@ type TimestampCompressor interface {
//Payload() []byte // для снапшота
// Offset() int
ReplaceBuffer([]byte)
WriteStoredTo(io.Writer) error
WriteCommitedTo(io.Writer) error
LastTimestamp() uint32
ReplaceSinceWithUntil() uint32
}
@@ -65,7 +65,7 @@ type ValueCompressor interface {
// (rewindOffset, change, value, delta)
Append(int, []byte, float64, uint64)
Size() int
StoredSize() int
CommitedSize() int
//Chunks() [][]byte
//DeleteLast()
CaptureState()
@@ -77,7 +77,7 @@ type ValueCompressor interface {
//Payload() []byte // для снапшота
//Offset() int
ReplaceBuffer([]byte)
WriteStoredTo(io.Writer) error
WriteCommitedTo(io.Writer) error
LastValue() float64
}

View File

@@ -2,6 +2,7 @@ package recovery
import (
"errors"
"fmt"
"os"
"path/filepath"
"strconv"
@@ -113,6 +114,7 @@ func resolveRecoveryFiles(dir string) (_ RecoveryFiles, err error) {
type RecoveryReport struct {
SnapshotNumber int
WAL string
ReplayMetrics map[uint32]*storage.ReplayMetric
IndexFreeList *freelist.FreeList
DataFreeList *freelist.FreeList
}
@@ -122,6 +124,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) {
if err != nil {
return
}
fmt.Printf("recovery files: %#v\n", recoveryFiles)
var (
frozenIndexPageCount int
frozenDataPageCount int
@@ -131,27 +134,33 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) {
)
if recoveryFiles.Snapshot != "" {
var snapshot worker.ReadSnapshotOut
snapshot, err = worker.ReadSnapshot(recoveryFiles.Snapshot)
snapshot, err = worker.ReadSnapshot(filepath.Join(dir, recoveryFiles.Snapshot))
if err != nil {
return
}
//
fmt.Printf("snapshot: %#v\n", snapshot)
frozenIndexPageCount = snapshot.FrozenIndexPagesCount
frozenDataPageCount = snapshot.FrozenDataPagesCount
freeIndexPages = snapshot.IndexPageNumbers
freeDataPages = snapshot.DataPageNumbers
metrics = snapshot.Metrics
} else {
metrics = make(map[uint32]*storage.ReplayMetric)
}
if recoveryFiles.WAL != "" {
var (
wal *os.File
walReader *storage.WALReader
walReplayer *storage.WALReplayer
)
walReader, err = storage.NewWALReader(storage.WALReaderOptions{
FileName: recoveryFiles.WAL,
BufferSize: 8 * 1024 * 1024,
})
wal, err = os.Open(filepath.Join(dir, recoveryFiles.WAL))
if err != nil {
return
}
defer wal.Close()
walReader, err = storage.NewWALReader(wal)
if err != nil {
return
}
@@ -166,13 +175,15 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) {
return
}
fmt.Printf("before replay\n")
err = walReplayer.Replay()
if err != nil {
return
}
fmt.Printf("after replay\n")
// перезаписую сторінки із останнього комміта
var dataFile *os.File
dataFile, err = os.OpenFile(qb.GetDataFilePath(dir, databaseName), os.O_WRONLY, 0666)
dataFile, err = os.OpenFile(qb.GetDataFilePath(dir, databaseName), os.O_CREATE|os.O_WRONLY, 0666)
if err != nil {
return
}
@@ -185,7 +196,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) {
return
}
var indexFile *os.File
indexFile, err = os.OpenFile(qb.GetIndexFilePath(dir, databaseName), os.O_WRONLY, 0666)
indexFile, err = os.OpenFile(qb.GetIndexFilePath(dir, databaseName), os.O_CREATE|os.O_WRONLY, 0666)
if err != nil {
return
}
@@ -221,6 +232,7 @@ func Recovery(dir string, databaseName string) (_ RecoveryReport, err error) {
return RecoveryReport{
SnapshotNumber: recoveryFiles.SnapshotNumber,
WAL: recoveryFiles.WAL,
ReplayMetrics: metrics,
IndexFreeList: indexFreeList,
DataFreeList: dataFreeList,
}, nil

View File

@@ -43,14 +43,14 @@ func (s *ReplayMetric) MeasuresAppend(rec MeasuresAppendRecord) {
rec.Values, rec.ValuesRewindOffset)
}
type AppendMeasuresResult struct {
type ReplayMeasuresAppendWithGrowResult struct {
IndexPagesToRewrite []PageToWrite
DataPagesToRewrite []PageToWrite
ReusedIndexPagesCount int
ReusedDataPagesCount int
}
func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (_ AppendMeasuresResult) {
func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord, isLastPacket bool) (_ ReplayMeasuresAppendWithGrowResult) {
var (
reusedIndexPages int
reusedDataPages int
@@ -82,6 +82,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord,
if head != nil {
// попередній хвіст індексного рівня перетворився на head сторінку,
// отже TailRecords - це новий хвіст
clear(level.Buffer)
copy(level.Buffer, change.TailRecords)
level.RecordsCount = newRecordsCount
} else {
@@ -138,6 +139,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord,
}
// HeadDataPage є полюбому, оскільки це транзакція із мінімум однією заповненою
// data сторінкою. Отже TailTimestamps і TailValues - це нові хвости data рівня.
clear(s.Buf)
copy(s.Buf, rec.TailValues)
s.ValuesSize = len(rec.TailValues)
@@ -151,7 +153,7 @@ func (s *ReplayMetric) MeasuresAppendWithGrow(rec MeasuresAppendWithGrowRecord,
s.LastPageNo = rec.DataPages[len(rec.DataPages)-1].PageNo
}
return AppendMeasuresResult{
return ReplayMeasuresAppendWithGrowResult{
IndexPagesToRewrite: indexPages,
DataPagesToRewrite: dataPages,
ReusedIndexPagesCount: reusedIndexPages,
@@ -192,9 +194,9 @@ func (s *ReplayMetric) composeHeadDataPage(head DataPageTail) PageToWrite {
page := make([]byte, DataPageSize)
copy(page, s.Buf)
// append data
timestampsSize := writeTimestampsWithRewind(s.Buf, s.TimestampsSize,
timestampsSize := writeTimestampsWithRewind(page, s.TimestampsSize,
head.Timestamps, head.TimestampsRewindOffset)
valuesSize := writeValuesWithRewind(s.Buf, s.ValuesSize,
valuesSize := writeValuesWithRewind(page, s.ValuesSize,
head.Values, head.ValuesRewindOffset)
// запечатати сторінку
calculatedCRC := SealDataPage(SealDataPageIn{

View File

@@ -49,14 +49,19 @@ type MetricAdd struct {
MetricID uint32
MetricType qb.MetricType
FracDigits byte
ResultCh chan struct{}
ResultCh chan byte
}
type MetricDelete struct {
MetricID uint32
FreeIndexPages []uint32
FreeDataPages []uint32
ResultCh chan struct{}
ResultCh chan byte
}
type MeasuresAppendResult struct {
ResultCode byte
WrittenCount int
}
type MeasuresAppendWithGrow struct {
@@ -72,7 +77,7 @@ type MeasuresAppendWithGrow struct {
IndexLevelTails []IndexLevelTail
ResultCode byte
WrittenCount int
ResultCh chan AppendMeasuresResult
ResultCh chan MeasuresAppendResult
}
type MeasuresAppend struct {
@@ -83,14 +88,14 @@ type MeasuresAppend struct {
Values []byte
ResultCode byte
WrittenCount int
ResultCh chan AppendMeasuresResult
ResultCh chan MeasuresAppendResult
}
type MeasuresDelete struct {
MetricID uint32
FreeIndexPages []uint32
FreeDataPages []uint32
ResultCh chan struct{}
ResultCh chan byte
}
// WRITE RESULTS
@@ -100,7 +105,7 @@ type MeasuresAppendCommited struct {
MetricID uint32
ResultCode byte
WrittenCount int
ResultCh chan AppendMeasuresResult
ResultCh chan MeasuresAppendResult
}
type MeasuresAppendWithGrowCommited struct {
@@ -109,20 +114,20 @@ type MeasuresAppendWithGrowCommited struct {
Index []IndexLevelTail
ResultCode byte
WrittenCount int
ResultCh chan AppendMeasuresResult
ResultCh chan MeasuresAppendResult
}
type MeasuresDeleteCommited struct {
MetricID uint32
ResultCh chan struct{}
ResultCh chan byte
}
type MetricAddCommited struct {
MetricID uint32
ResultCh chan struct{}
ResultCh chan byte
}
type MetricDeleteCommited struct {
MetricID uint32
ResultCh chan struct{}
ResultCh chan byte
}

View File

@@ -3,6 +3,7 @@ package storage
import (
"bytes"
"fmt"
"io"
"math"
"reflect"
"slices"
@@ -914,3 +915,169 @@ func TestWriteValuesWithRewind(t *testing.T) {
t.Fatalf("got buf % x are not equal expected % x", before, after)
}
}
func TestWALReader(t *testing.T) {
rec1 := MetricAddRecord{
MetricID: 100,
MetricType: qb.Cumulative,
FracDigits: 4,
}
rec2 := MetricAddRecord{
MetricID: 200,
MetricType: qb.Instant,
FracDigits: 1,
}
rec3 := MetricDeleteRecord{
MetricID: 100,
FreeIndexPages: []uint32{
1, 2, 3,
},
FreeDataPages: []uint32{
10, 20,
},
}
rec4 := MetricDeleteRecord{
MetricID: 200,
}
packet1 := createPacket([]Packable{rec1, rec2})
packet2 := createPacket([]Packable{rec3, rec4})
src := bytes.NewBuffer(nil)
src.Write(packet1)
src.Write(packet2)
walReader, err := NewWALReader(src)
if err != nil {
t.Fatal(err)
}
i := 0
for {
records, isLastPacket, err := walReader.NextPacket()
if err != nil {
t.Fatal(err)
}
if records == nil {
break
}
switch i {
case 0:
if isLastPacket {
t.Fatalf("1st packet can't be the last packet in WAL")
}
if !reflect.DeepEqual(records, []any{rec1, rec2}) {
t.Fatalf("1st packet decoded incorrectly")
}
case 1:
if !isLastPacket {
t.Fatalf("2nd packet must be the last packet in WAL")
}
if !reflect.DeepEqual(records, []any{rec3, rec4}) {
t.Fatalf("2nd packet decoded incorrectly")
}
}
i++
}
}
func createPacket(records []Packable) []byte {
w := bytes.NewBuffer(nil)
w.Write([]byte{
0, 0, 0, 0, 0, 0, 0, 0, 0, // size (max 9 byte)
})
hasher := util.NewHasher()
multi := io.MultiWriter(w, hasher)
for _, rec := range records {
rec.Pack(multi)
}
return finalizePacket(w, hasher.Sum32())
}
func TestReplayMetric(t *testing.T) {
// data
DataPageSize = 24
DataPagePayloadSize = 12
dataCRC32Idx = DataPageSize - 4
timestampsSizeIdx = DataPageSize - 6
valuesSizeIdx = DataPageSize - 8
prevPageIdx = DataPageSize - 12
// index
IndexPageSize = 24
indexCRC32Idx = IndexPageSize - 4
indexRecordsCountIdx = IndexPageSize - 6
isZeroLevelIdx = IndexPageSize - 7
maxRecordsOnIndexPage = 2
//
metric := &ReplayMetric{
Buf: []byte{
0xa1, 0xa2, 0xa3, 0xa4,
0x00, 0x00, 0x00, 0x00,
0xb4, 0xb3, 0xb2, 0xb1,
0x01, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
},
TimestampsSize: 4,
ValuesSize: 4,
IndexLevelTails: []IndexLevelTail{
{
Buffer: []byte{
0x01, 0x01, 0x01, 0x01,
0x02, 0x02, 0x02, 0x02,
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
},
RecordsCount: 1,
},
},
}
result := metric.MeasuresAppendWithGrow(MeasuresAppendWithGrowRecord{
DataPageTail: DataPageTail{
PrevPageNo: 1,
CRC32: 3235260875,
TimestampsRewindOffset: 1,
Timestamps: []byte{0x55, 0x44},
ValuesRewindOffset: 2,
Values: []byte{0x33, 0x44},
},
TailTimestamps: []byte{
0xf3, 0xf2, 0xf1,
},
TailValues: []byte{
0xd1, 0xd2,
},
ChangedIndexLevels: []ChangedIndexLevel{
{
IndexPageTail: &IndexPageTail{
CRC32: 2138044007,
Records: []byte{
0x03, 0x03, 0x03, 0x03,
0x04, 0x04, 0x04, 0x04,
},
},
TailRecords: []byte{
0x08, 0x08, 0x08, 0x08,
0x09, 0x09, 0x09, 0x09,
},
},
{
TailRecords: []byte{
0x01, 0x01, 0x01, 0x01,
0x02, 0x02, 0x02, 0x02,
},
},
},
}, true) // last packet
fmt.Printf("% x\n", metric.Buf)
fmt.Printf("% x\n", result.DataPagesToRewrite[0].Content)
fmt.Printf("% x\n", result.IndexPagesToRewrite[0].Content)
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)
}

View File

@@ -3,10 +3,8 @@ package storage
import (
"bufio"
"bytes"
"errors"
"fmt"
"io"
"os"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb/util"
@@ -15,155 +13,133 @@ import (
const readBufferSize = 8 * 1024 * 1024
type WALReader struct {
file *os.File
reader *bufio.Reader
current []any
next []any
//file *os.File
reader *bufio.Reader
next []any
done bool
}
type WALReaderOptions struct {
FileName string
BufferSize int
}
// type WALReaderOptions struct {
// File *os.File
// BufferSize int
// }
func NewWALReader(opt WALReaderOptions) (*WALReader, error) {
if opt.FileName == "" {
return nil, errors.New("FileName option is required")
}
if opt.BufferSize <= 0 {
return nil, errors.New("BufferSize option is required")
}
file, err := os.Open(opt.FileName)
if err != nil {
return nil, err
}
func NewWALReader(src io.Reader) (*WALReader, error) {
// if opt.FileName == "" {
// return nil, errors.New("FileName option is required")
// }
// if opt.BufferSize <= 0 {
// return nil, errors.New("BufferSize option is required")
// }
return &WALReader{
file: file,
reader: bufio.NewReaderSize(file, readBufferSize),
//file: file,
reader: bufio.NewReaderSize(src, readBufferSize),
}, nil
}
func (s *WALReader) Close() {
s.file.Close()
}
// func (s *WALReader) Close() {
// s.file.Close()
// }
func (s *WALReader) NextPacket() ([]any, bool, error) {
var err error
if s.current == nil {
s.current, err = s.readPacket()
// (packet records, isLastPacket, error)
func (s *WALReader) NextPacket() (_ []any, _ bool, err error) {
if s.done {
return
}
var current []any
if s.next != nil {
current = s.next
} else {
current, err = s.readPacket()
if err != nil {
return nil, false, err
return
}
if len(s.current) > 0 {
s.next, err = s.readPacket()
if err != nil {
return nil, false, err
}
}
}
current := s.current
done := s.next == nil
s.current = s.next
s.next, err = s.readPacket()
if err != nil {
return nil, false, err
}
return current, done, nil
}
func (s *WALReader) readPacket() (_ []any, err error) {
payloadSize, err := bin.ReadVarSize(s.reader)
if err != nil {
if err == io.EOF {
return nil, nil
} else {
if current == nil {
return
}
}
storedCRC, err := bin.ReadUint32(s.reader)
s.next, err = s.readPacket()
if err != nil {
return
}
body, err := bin.ReadN(s.reader, int(payloadSize))
if err != nil {
return
}
hasher := util.NewHasher()
hasher.Write(body)
calculatedCRC := hasher.Sum32()
if calculatedCRC == storedCRC {
return s.parseRecords(body)
}
//s.logger.Printf("stored CRC %d != calculated CRC %d", storedCRC, calculatedCRC)
return nil, nil
return current, s.next == nil, nil
}
func (s *WALReader) parseRecords(body []byte) ([]any, error) {
func (s *WALReader) readPacket() (_ []any, err error) {
bodySize, err := bin.ReadVarSize(s.reader) // fix add n
if err != nil {
if err == io.EOF {
s.done = true
return nil, nil
}
return
}
body, err := bin.ReadN(s.reader, int(bodySize))
if err != nil {
return
}
writtenCRC, err := bin.ReadUint32(s.reader)
if err != nil {
return
}
calculatedCRC := util.CalculateCRC32(body)
if calculatedCRC != writtenCRC {
return nil, fmt.Errorf("written CRC %d != calculated CRC %d",
writtenCRC, calculatedCRC)
}
return s.parseRecords(body)
}
func (s *WALReader) parseRecords(body []byte) (_ []any, err error) {
var (
src = bytes.NewBuffer(body)
records []any
)
for {
recordType, err := src.ReadByte()
var recordType byte
recordType, err = src.ReadByte()
if err != nil {
if err == io.EOF {
return records, nil
}
return nil, err
return
}
switch recordType {
case CodeMetricAdd:
rec := new(MetricAddRecord)
if err = rec.Parse(src); err != nil {
return nil, err
return
}
records = append(records, rec)
records = append(records, *rec)
case CodeMetricDelete:
rec := new(MetricDeleteRecord)
if err = rec.Parse(src); err != nil {
return nil, err
return
}
records = append(records, rec)
records = append(records, *rec)
// case CodeAppendedMeasure:
// rec := new(AppendedMeasure)
// if err = rec.Parse(src); err != nil {
// return nil, err
// }
// records = append(records, rec)
case CodeMeasuresAppend:
rec := new(MeasuresAppendRecord)
if err = rec.Parse(src); err != nil {
return
}
records = append(records, *rec)
// case CodeAppendedMeasures:
// rec := new(AppendedMeasures)
// if err = rec.Parse(src); err != nil {
// return nil, err
// }
// records = append(records, rec)
// case CodeAppendedPages:
// rec := new(AppendedPages)
// if err = rec.Parse(src); err != nil {
// return nil, err
// }
// records = append(records, rec)
case CodeMeasuresAppendWithGrow:
rec := new(MeasuresAppendWithGrowRecord)
if err = rec.Parse(src); err != nil {
return
}
records = append(records, *rec)
case CodeMeasuresDelete:
rec := new(MeasuresDeleteRecord)
if err = rec.Parse(src); err != nil {
return nil, err
return
}
records = append(records, rec)
records = append(records, *rec)
default:
return nil, fmt.Errorf("unknown record type code: %d", recordType)

View File

@@ -58,6 +58,9 @@ func (s *WALReplayer) Replay() error {
if err != nil {
return err
}
if len(records) == 0 {
return nil
}
for _, rec := range records {
if err = s.replayRecord(rec, isLastPacket); err != nil {
return err

View File

@@ -9,6 +9,10 @@ import (
"gordenko.dev/dima/qb/util"
)
type Packable interface {
Pack(io.Writer)
}
type WritePreparer struct {
pagePreparer *PageIngester
//dataPagePayloadBound int
@@ -39,15 +43,19 @@ type PreparedData struct {
Commits []any
}
// const (
// prefixSize = 9
// crc32Size = 4
// )
func (s *WritePreparer) Prepare(input []any) PreparedData {
s.w.Write([]byte{
0, 0, 0, 0, 0, 0, 0, 0, 0, // size (max 9 byte)
0, 0, 0, 0, // crc32
})
hasher := util.NewHasher()
w := io.MultiWriter(s.w, hasher)
multi := io.MultiWriter(s.w, hasher)
// 1. Пакую всі дані в WAL буфер запису
for _, untyped := range input {
@@ -58,7 +66,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
MetricType: x.MetricType,
FracDigits: x.FracDigits,
}
rec.Pack(w)
rec.Pack(multi)
s.commits = append(s.commits, MetricAddCommited{
MetricID: x.MetricID,
ResultCh: x.ResultCh,
@@ -76,7 +84,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
FreeIndexPages: x.FreeIndexPages,
FreeDataPages: x.FreeDataPages,
}
rec.Pack(w)
rec.Pack(multi)
s.commits = append(s.commits, MetricDeleteCommited{
MetricID: x.MetricID,
ResultCh: x.ResultCh,
@@ -90,7 +98,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
ValuesRewindOffset: x.ValuesRewindOffset,
Values: x.Values,
}
rec.Pack(w)
rec.Pack(multi)
// Дані для Worker
s.commits = append(s.commits, MeasuresAppendCommited{
MetricID: x.MetricID,
@@ -188,9 +196,9 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
ChangedIndexLevels: changedIndexPages,
}
fmt.Printf("%#v\ns", rec.ChangedIndexLevels[0].IndexPageTail)
//fmt.Printf("%#v\ns", rec.ChangedIndexLevels[0].IndexPageTail)
rec.Pack(w)
rec.Pack(multi)
// Дані для Worker
s.commits = append(s.commits, MeasuresAppendWithGrowCommited{
MetricID: x.MetricID,
@@ -213,7 +221,7 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
FreeIndexPages: x.FreeIndexPages,
FreeDataPages: x.FreeDataPages,
}
rec.Pack(w)
rec.Pack(multi)
s.commits = append(s.commits, MeasuresDeleteCommited{
MetricID: x.MetricID,
ResultCh: x.ResultCh,
@@ -222,19 +230,11 @@ func (s *WritePreparer) Prepare(input []any) PreparedData {
}
}
// 2. Додаю розмір пакету і чексуму
packet := finalizePacket(s.w, hasher.Sum32())
//s.written += int64(len(packet)) + 12
packet := s.w.Bytes()
// write size
payloadSize := len(packet) - 13
start := 9 - bin.CountVarSize(payloadSize)
bin.PutVarSize(packet[start:], payloadSize)
// CRC32
bin.PutUint32(packet[9:], hasher.Sum32())
//
return PreparedData{
Packet: packet[start:],
Packet: packet,
WriteToIndex: s.indexPages,
WriteToData: s.dataPages,
Commits: s.commits,
@@ -260,3 +260,13 @@ func getIndexRecords(buf []byte, count int) (list []IndexRecord) {
}
return
}
func finalizePacket(w *bytes.Buffer, checksum uint32) []byte {
bin.WriteUint32(w, checksum)
packet := w.Bytes()
bodySize := len(packet) - 13 // 9 bytes prefix + 4 bytes crc32
// cut empty bytes at start
start := 9 - bin.CountVarSize(bodySize)
bin.PutVarSize(packet[start:], bodySize)
return packet[start:]
}

View File

@@ -4,7 +4,6 @@ import (
"bytes"
"errors"
"fmt"
"math"
"os"
"path/filepath"
"sync"
@@ -203,37 +202,8 @@ func (s *Writer) Run() {
}
}
func (s *Writer) getDataPageNumber() (uint32, bool, error) {
pageNo, err := s.dataFreeList.GetPageNumber()
if err != nil {
return 0, false, err
}
if pageNo > 0 {
return pageNo, true, nil
}
if s.dataPagesCount < math.MaxUint32 {
s.dataPagesCount++
return s.dataPagesCount, false, nil
}
return 0, false, errors.New("no space")
}
func (s *Writer) getIndexPageNumber() (uint32, bool, error) {
pageNo, err := s.indexFreeList.GetPageNumber()
if err != nil {
return 0, false, err
}
if pageNo > 0 {
return pageNo, true, nil
}
if s.indexPagesCount < math.MaxUint32 {
s.indexPagesCount++
return s.indexPagesCount, false, nil
}
return 0, false, errors.New("no space")
}
func (s *Writer) packAndWrite() (err error) {
//fmt.Println("packAndWrite")
//s.mutex.Lock()
//isExited := s.isExited
@@ -245,15 +215,19 @@ func (s *Writer) packAndWrite() (err error) {
input := s.inbox.Drain()
//fmt.Printf("Drain: %#v\n", input)
prepared := s.writePreparer.Prepare(input)
//fmt.Printf("prepared: %#v\n", prepared)
// 3. Пишу на диск WAL (append)
n, err := s.wal.Write(prepared.Packet)
if err != nil {
return
}
if n != len(prepared.Packet) {
return fmt.Errorf("written %d != total size %d", n, len(prepared.Packet))
return fmt.Errorf("written %d != packet size %d", n, len(prepared.Packet))
}
if err = s.wal.Sync(); err != nil {
return

0
testdir/test.data Normal file
View File

0
testdir/test.data_free Normal file
View File

View File

0
testdir/test.index Normal file
View File

0
testdir/test.index_free Normal file
View File

View File

BIN
testdir/test.wal_0 Executable file

Binary file not shown.

View File

@@ -1,6 +1,7 @@
package transform
import (
"fmt"
"io"
bin "gordenko.dev/dima/bin/little"
@@ -101,6 +102,7 @@ type CumulativeMeasureWriter struct {
}
func NewCumulativeMeasureWriter(dst io.Writer, since uint32) *CumulativeMeasureWriter {
fmt.Println("NewCumulativeMeasureWriter")
// 20 - это timestamp, value, total
return &CumulativeMeasureWriter{
arr: make([]byte, 20),
@@ -126,11 +128,13 @@ func (s *CumulativeMeasureWriter) feed(timestamp uint32, value float64, isBuffer
s.responder.AppendRecord(s.arr)
}
}
fmt.Println("feed:", timestamp, value, isBuffer)
s.endTimestamp = timestamp
s.endValue = value
}
func (s *CumulativeMeasureWriter) pack(total float64) {
fmt.Println("PACK")
bin.PutUint32(s.arr[0:], s.endTimestamp)
bin.PutFloat64(s.arr[4:], s.endValue)
bin.PutFloat64(s.arr[12:], total)
@@ -141,6 +145,7 @@ func (s *CumulativeMeasureWriter) Close() error {
// endTimestamp внутри заданного периода. Других показаний нет,
// поэтому время добавляю, но накопленную сумму ставлю 0.
s.pack(0)
s.responder.BufferRecord(s.arr)
// Если < since - ничего делать не нужно, ибо накопленная сумма уже добавлена
}
return s.responder.Flush()

View File

@@ -38,11 +38,13 @@ func NewChunkedResponder(dst io.Writer) *ChunkedResponder {
func (s *ChunkedResponder) BufferRecord(rec []byte) {
s.buf.Write(rec)
s.recordsQty++
fmt.Println("BufferRecord", s.recordsQty)
}
func (s *ChunkedResponder) AppendRecord(rec []byte) error {
s.buf.Write(rec)
s.recordsQty++
fmt.Println("AppendRecord", s.recordsQty)
if s.buf.Len() < 1500 {
return nil
@@ -61,6 +63,7 @@ func (s *ChunkedResponder) AppendRecord(rec []byte) error {
}
func (s *ChunkedResponder) Flush() error {
fmt.Println("Flush", s.recordsQty)
if s.recordsQty > 0 {
if err := s.sendBuffered(); err != nil {
return err
@@ -75,13 +78,14 @@ func (s *ChunkedResponder) Flush() error {
}
func (s *ChunkedResponder) sendBuffered() (err error) {
fmt.Println("sendBuffered", s.recordsQty)
msg := s.buf.Bytes()
bin.PutUint32(msg[1:], uint32(s.recordsQty))
//fmt.Printf("put uint16: %d\n", msg[:3])
fmt.Printf("put uint16: %d\n", msg[:3])
//fmt.Printf("send %d records\n", s.recordsQty)
fmt.Printf("send %d records\n", s.recordsQty)
//fmt.Printf("send buffered: %d, qty: %d\n", msg, s.recordsQty)
fmt.Printf("send buffered: %d, qty: %d\n", msg, s.recordsQty)
n, err := s.dst.Write(msg)
if err != nil {

View File

@@ -2,6 +2,7 @@ package worker
import (
"errors"
"fmt"
"io"
bin "gordenko.dev/dima/bin/little"
@@ -30,8 +31,8 @@ type Metric struct {
buffer []byte
timestamps qb.TimestampCompressor
values qb.ValueCompressor
XLock bool
RLocks int
xLock bool
rLocks int
WaitQueue []any
indexLevelTails []storage.IndexLevelTail // root - last element
capturedState *CapturedState
@@ -62,7 +63,7 @@ func (s *Metric) LastTimestamp() uint32 {
func (s *Metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
//s.timestamps.Reset()
//s.values.Reset()
s.XLock = false
s.xLock = false
// s.Timestamps.Renew()
// s.Values.Renew()
@@ -75,8 +76,7 @@ func (s *Metric) OnMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
s.lastValue = 0
}
// func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) {
func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox *inbox.Inbox) {
func (s *Metric) AppendMeasures(req AppendMeasuresReq, tmp []byte, storageInbox *inbox.Inbox) {
if s.capturedState != nil {
s.WaitQueue = append(s.WaitQueue, req)
}
@@ -96,7 +96,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageI
pages []storage.DataPayload
written int
resultCode byte
resultCode byte = Succeed
)
s.capturedState = &CapturedState{
@@ -117,8 +117,10 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageI
break
}
tReport := timestamps.Evaluate(tmp, measure.Timestamp)
vReport := values.Evaluate(tmp, measure.Value)
tReport := timestamps.Evaluate(tmp[:7], measure.Timestamp)
fmt.Printf("tReport: %#v\n", tReport)
vReport := values.Evaluate(tmp[7:], measure.Value)
fmt.Printf("vReport: %#v\n", vReport)
totalRequiredSpace := tReport.TotalSpace + vReport.TotalSpace
@@ -169,7 +171,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageI
if written == 0 {
s.capturedState = nil
req.ResultCh <- AppendMeasuresResult{
req.ResultCh <- storage.MeasuresAppendResult{
ResultCode: resultCode,
}
return
@@ -193,7 +195,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageI
TailValues: values.Tail(0),
ResultCode: resultCode,
WrittenCount: written,
//ResultCh: req.ResultCh,
ResultCh: req.ResultCh,
})
} else {
// короткий шлях - запис лише в storage
@@ -205,7 +207,7 @@ func (s *Metric) StartAppendMeasures(req AppendMeasuresReq, tmp []byte, storageI
Values: values.Tail(valuesOffset),
ResultCode: resultCode,
WrittenCount: written,
//ResultCh: req.ResultCh,
ResultCh: req.ResultCh,
})
}
}
@@ -215,6 +217,11 @@ func (s *Metric) OnMeasuresAppendCommited(rec storage.MeasuresAppendCommited) {
s.capturedState = nil
s.timestamps.ForgetCapturedState()
s.values.ForgetCapturedState()
rec.ResultCh <- storage.MeasuresAppendResult{
ResultCode: rec.ResultCode,
WrittenCount: rec.WrittenCount,
}
}
func (s *Metric) OnMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWithGrowCommited) {
@@ -229,6 +236,11 @@ func (s *Metric) OnMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWith
// В storage я передав повний індекс. У нього додали елементи (можливо нові рівні).
// Тому проста заміна
s.indexLevelTails = rec.Index
rec.ResultCh <- storage.MeasuresAppendResult{
ResultCode: rec.ResultCode,
WrittenCount: rec.WrittenCount,
}
}
// READ
@@ -317,10 +329,12 @@ func (s *Metric) StartFullScan(req FullScanReq) {
if done {
break
}
fmt.Println("ts:", timestamp)
value, done := valueDecompressor.NextValue()
if done {
qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug)
}
fmt.Println("value:", value)
req.ResponseWriter.FeedNoSend(timestamp, value)
}
@@ -330,7 +344,7 @@ func (s *Metric) StartFullScan(req FullScanReq) {
LastPageNo: s.lastPageNo,
FracDigits: s.fracDigits,
}
s.RLocks++
s.rLocks++
} else {
req.ResultCh <- FullScanResult{
ResultCode: QueryDone,
@@ -366,21 +380,21 @@ func (s *Metric) WriteTo(w io.Writer) (err error) {
if err != nil {
return
}
err = bin.WriteUint16(w, uint16(s.timestamps.StoredSize()))
err = bin.WriteUint16(w, uint16(s.timestamps.CommitedSize()))
if err != nil {
return
}
err = bin.WriteUint16(w, uint16(s.values.StoredSize()))
err = bin.WriteUint16(w, uint16(s.values.CommitedSize()))
if err != nil {
return
}
// timestamps payload
err = s.timestamps.WriteStoredTo(w)
err = s.timestamps.WriteCommitedTo(w)
if err != nil {
return
}
// values payload
err = s.values.WriteStoredTo(w)
err = s.values.WriteCommitedTo(w)
if err != nil {
return
}
@@ -394,7 +408,7 @@ func (s *Metric) WriteTo(w io.Writer) (err error) {
if err != nil {
return
}
_, err = w.Write(level.Buffer)
_, err = w.Write(level.Buffer[:level.RecordsCount*storage.IndexRecordSize])
if err != nil {
return
}

View File

@@ -5,7 +5,6 @@ import (
"fmt"
"io"
"os"
"path/filepath"
bin "gordenko.dev/dima/bin/little"
"gordenko.dev/dima/qb"
@@ -30,11 +29,10 @@ dataFreeList - Nb
CRC32 - 4b
*/
const readBufferSize = 4 * 1024 * 1024 // 8mb
const readBufferSize = 4 * 1024 * 1024
type WriteSnapshotIn struct {
SnapshotNumber int
Dir string
WriteBufferSize int
Metrics map[uint32]*Metric
FrozenIndexPagesCount int
@@ -43,24 +41,21 @@ type WriteSnapshotIn struct {
DataPageNumbers []uint32
}
func WriteSnapshot(in WriteSnapshotIn) (err error) {
var (
fileName = filepath.Join(in.Dir, fmt.Sprintf("%d.snapshot", in.SnapshotNumber))
hasher = util.NewHasher()
)
file, err := os.OpenFile(fileName, os.O_CREATE|os.O_WRONLY, 0666) //0770)
func writeSnapshot(filePath string, in WriteSnapshotIn) (err error) {
file, err := os.OpenFile(filePath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0666)
if err != nil {
return
qb.Abort(qb.WriteSnapshotFailed, err)
}
dst := io.MultiWriter(bufio.NewWriterSize(file, in.WriteBufferSize), hasher)
var (
bufferedWriter = bufio.NewWriterSize(file, in.WriteBufferSize)
hasher = util.NewHasher()
)
dst := io.MultiWriter(bufferedWriter, hasher)
//
_, err = bin.WriteVarSize(dst, len(in.Metrics))
if err != nil {
return
}
for metricID, metric := range in.Metrics {
err = bin.WriteUint32(dst, metricID)
if err != nil {
@@ -81,13 +76,16 @@ func WriteSnapshot(in WriteSnapshotIn) (err error) {
if err != nil {
return
}
// CRC32
bin.WriteUint32(file, hasher.Sum32())
err = file.Close()
err = bufferedWriter.Flush()
if err != nil {
return
}
// CRC32
err = bin.WriteUint32(file, hasher.Sum32())
if err != nil {
return
}
return file.Close()
// prevLogNumber := logNumber - 1
// prevChanges := filepath.Join(s.dir, fmt.Sprintf("%d.changes", prevLogNumber))
@@ -116,7 +114,6 @@ func WriteSnapshot(in WriteSnapshotIn) (err error) {
// qb.Abort(qb.DeletePrevSnapshotFileFailed, err)
// }
// }
return
}
type ReadSnapshotOut struct {
@@ -127,7 +124,7 @@ type ReadSnapshotOut struct {
DataPageNumbers []uint32
}
func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
func ReadSnapshot(filePath string) (out ReadSnapshotOut, err error) {
var (
metrics = make(map[uint32]*storage.ReplayMetric)
frozenIndexPagesCount int
@@ -137,7 +134,7 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
hasher = util.NewHasher()
)
file, err := os.Open(fileName)
file, err := os.Open(filePath)
if err != nil {
return
}
@@ -148,25 +145,16 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
return
}
fileSize := stat.Size()
if fileSize == 0 {
return ReadSnapshotOut{
Metrics: metrics,
}, nil
}
if fileSize < 4 {
err = fmt.Errorf("%s is corrupted", fileName)
err = fmt.Errorf("snapshot %s is corrupted", filePath)
return
}
payloadSize := fileSize - 4
// 2. Створюємо великий буфер для NVMe (4 MB)
bufferedReader := bufio.NewReaderSize(file, readBufferSize)
// 3. Обмежуємо читання лише розміром Payload
limitReader := io.LimitReader(bufferedReader, payloadSize)
// 4. Створюємо хеш та обгортку TeeReader, яка рахує CRC32 на льоту
src := io.TeeReader(limitReader, hasher)
@@ -189,6 +177,7 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
fmt.Println("metricID", metricID)
metricType, err = bin.ReadByte(src)
if err != nil {
return
@@ -198,10 +187,13 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
fmt.Println("m.MetricType", m.MetricType)
fmt.Println("m.FracDigits", m.FracDigits)
m.LastPageNo, err = bin.ReadUint32(src)
if err != nil {
return
}
fmt.Println("m.LastPageNo", m.LastPageNo)
m.TimestampsSize, err = bin.ReadUint16AsInt(src)
if err != nil {
return
@@ -210,7 +202,8 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
err = bin.ReadNInto(src, buf[len(buf)-m.TimestampsSize:])
fmt.Println("m.TimestampsSize", m.TimestampsSize, "m.ValuesSize", m.ValuesSize)
err = bin.ReadNInto(src, buf[storage.DataPagePayloadSize-m.TimestampsSize:storage.DataPagePayloadSize])
if err != nil {
return
}
@@ -218,11 +211,13 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
fmt.Printf("% x\n", buf)
//
levelsCount, err = bin.ReadVarSize(src)
if err != nil {
return
}
fmt.Println("levels count:", levelsCount)
for range levelsCount {
var (
buf = make([]byte, storage.IndexPageSize)
@@ -232,10 +227,12 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
fmt.Println("records count:", count)
err = bin.ReadNInto(src, buf[:count*storage.IndexRecordSize])
if err != nil {
return
}
fmt.Printf("% x\n", buf)
m.IndexLevelTails = append(m.IndexLevelTails, storage.IndexLevelTail{
Buffer: buf,
RecordsCount: count,
@@ -258,12 +255,10 @@ func ReadSnapshot(fileName string) (out ReadSnapshotOut, err error) {
if err != nil {
return
}
calculatedCRC := hasher.Sum32()
if expectedCRC != calculatedCRC {
err = fmt.Errorf("%s is corrupted. Calculated CRC %d not equal expected CRC %d",
fileName, calculatedCRC, expectedCRC)
filePath, calculatedCRC, expectedCRC)
return
}
@@ -283,6 +278,7 @@ func writeFreePages(dst io.Writer, frozenPagesCount int, pageNumbers []uint32) (
if err != nil {
return
}
fmt.Println("write frozenPagesCount:", frozenPagesCount)
_, err = bin.WriteVarSize(dst, len(pageNumbers))
if err != nil {
return
@@ -305,6 +301,7 @@ func readFreePages(src io.Reader) (_ int, _ []uint32, err error) {
if err != nil {
return
}
fmt.Println("read frozenPagesCount:", frozenPagesCount)
pageNumbersCount, err := bin.ReadVarSize(src)
if err != nil {
return

View File

@@ -1,25 +1,74 @@
package worker
import (
"bytes"
"fmt"
"os"
"reflect"
"slices"
"testing"
"gordenko.dev/dima/qb"
"gordenko.dev/dima/qb/enc"
"gordenko.dev/dima/qb/storage"
)
func TestWriteReadSnapshot(t *testing.T) {
metrics := make(map[uint32]*Metric)
metrics[1] = &Metric{
metricType: qb.Cumulative,
fracDigits: 3,
lastPageNo: 1000,
buffer: []byte{},
storage.DataPageSize = 32
storage.DataPagePayloadSize = 20
storage.IndexPageSize = 16
buf := []byte{
// values
0x8f, // base value
0x80, // delta 0
0x01, // run (len = 3)
0x81, // delta 1
0x81, // literal (len = 2)
// empty space
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00,
// timestamps
0x80, // h-byte (literal, len=1)
0x94, // delta 20
0x01, // h-byte (run, len=3)
0xbc, // delta 60
0x28, 0x80, 0x24, 0x6a, // since
// footer
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00,
}
databuf := buf[:storage.DataPagePayloadSize]
timestampsSize := 8
valuesSize := 5
metrics := map[uint32]*Metric{
1: {
metricType: qb.Cumulative,
fracDigits: 3,
lastPageNo: 1000,
buffer: buf,
timestamps: enc.NewTimeDeltaCompressor(databuf, timestampsSize),
values: enc.NewValueDeltaCompressor(qb.Cumulative, 3, databuf, valuesSize),
indexLevelTails: []storage.IndexLevelTail{
{
Buffer: []byte{
0x01, 0x01, 0x01, 0x01,
0x02, 0x02, 0x02, 0x02,
0x00,
// footer
0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00,
},
RecordsCount: 1,
},
},
},
}
in := WriteSnapshotIn{
SnapshotNumber: 1,
Dir: ".",
WriteBufferSize: 8 * 1024 * 1024,
Metrics: metrics,
FrozenIndexPagesCount: 7,
@@ -31,14 +80,17 @@ func TestWriteReadSnapshot(t *testing.T) {
6, 7, 8,
},
}
err := WriteSnapshot(in)
fileName := "test.snapshot"
err := writeSnapshot(fileName, in)
if err != nil {
t.Fatalf("writeSnapshot: %s", err)
}
out, err := ReadSnapshot("1.snapshot")
out, err := ReadSnapshot(fileName)
if err != nil {
return
t.Fatal(err)
}
if out.FrozenIndexPagesCount != in.FrozenIndexPagesCount {
@@ -67,6 +119,7 @@ func TestWriteReadSnapshot(t *testing.T) {
t.Fatalf("metric %d: %s", metricID, err)
}
}
os.Remove(fileName)
}
func cmpMetric(replay *storage.ReplayMetric, origin *Metric) error {
@@ -82,9 +135,21 @@ func cmpMetric(replay *storage.ReplayMetric, origin *Metric) error {
return fmt.Errorf("lastPageNo: got %d are not equal expected %d",
replay.LastPageNo, origin.lastPageNo)
}
// if !reflect.DeepEqual() {
// return fmt.Errorf("lastPageNo: got %d are not equal expected %d",
// replay.LastPageNo, origin.lastPageNo)
// }
if !reflect.DeepEqual(replay.IndexLevelTails, origin.indexLevelTails) {
return fmt.Errorf("indexLevelTails: got %#v are not equal expected %#v",
replay.IndexLevelTails, origin.indexLevelTails)
}
if !bytes.Equal(replay.Buf, origin.buffer) {
return fmt.Errorf("buffer: got % x are not equal expected % x",
replay.Buf, origin.buffer)
}
if replay.TimestampsSize != origin.timestamps.Size() {
return fmt.Errorf("timestampsSize: got %d are not equal expected %d",
replay.TimestampsSize, origin.timestamps.Size())
}
if replay.ValuesSize != origin.values.Size() {
return fmt.Errorf("valuesSize: got %d are not equal expected %d",
replay.ValuesSize, origin.values.Size())
}
return nil
}

View File

@@ -19,28 +19,27 @@ import (
// Якийсь прапор поставити для pendingDelete, щоб read операції вставали в чергу
const (
QueryDone = 1
UntilFound = 2
UntilNotFound = 3
RangeFound = 15
NoMeasures = 16
NoMetric = 4
MetricDuplicate = 5
Succeed = 6
NewPage = 7
ExpiredMeasure = 8
NonMonotonicValue = 9
CanAppend = 10
WrongMetricType = 11
NoMeasuresToDelete = 12
DeleteFromAtreeNotNeeded = 13
DeleteFromAtreeRequired = 14
QueryDone = 1
UntilFound = 2
UntilNotFound = 3
RangeFound = 15
NoMeasures = 16
NoMetric = 4
MetricDuplicate = 5
Succeed = 6
ExpiredMeasure = 8
NonMonotonicValue = 9
WrongMetricType = 11
NoMeasuresToDelete = 12
//DeleteFromAtreeNotNeeded = 13
//DeleteFromAtreeRequired = 14
)
type Worker struct {
inbox *inbox.Inbox
storageInbox *inbox.Inbox
dir string
databaseName string
snapshotNumber int
tmp []byte
metrics map[uint32]*Metric
@@ -53,6 +52,7 @@ type Options struct {
StorageInbox *inbox.Inbox
ReplayMetrics map[uint32]*storage.ReplayMetric
Dir string
DatabaseName string
ExitCh chan struct{}
WaitGroup *sync.WaitGroup
}
@@ -62,6 +62,7 @@ func New(opt Options) *Worker {
inbox: opt.Inbox,
storageInbox: opt.StorageInbox,
dir: opt.Dir,
databaseName: opt.DatabaseName,
tmp: make([]byte, 26),
metrics: make(map[uint32]*Metric),
exitCh: opt.ExitCh,
@@ -151,13 +152,13 @@ func (s *Worker) doWork() {
s.applyCommits(req) // all metrics only
case ListCurrentValuesReq:
s.tryListCurrentValues(req) // all metrics only
s.ListCurrentValues(req) // all metrics only
case RangeScanReq:
s.RangeScan(req)
case FullScanReq:
s.tryFullScan(req)
s.FullScan(req)
case AddMetricReq:
s.AddMetric(req)
@@ -196,7 +197,7 @@ func (s *Worker) processMetricQueue(metricID uint32, metric *Metric, tmp []byte)
s.GetMetric(req)
case AppendMeasuresReq:
metric.StartAppendMeasures(req, tmp, s.storageInbox)
metric.AppendMeasures(req, tmp, s.storageInbox)
case DeleteMetricReq:
s.startDeleteMetric(metric, req)
@@ -212,28 +213,37 @@ func (s *Worker) processMetricQueue(metricID uint32, metric *Metric, tmp []byte)
}
type AddMetricReq struct {
MetricID uint32
ResultCh chan byte
MetricID uint32
MetricType qb.MetricType
FracDigits byte
ResultCh chan byte
}
func (s *Worker) AddMetric(req AddMetricReq) {
_, ok := s.metrics[req.MetricID]
if ok {
fmt.Println("add metric duplicate")
req.ResultCh <- MetricDuplicate
return
}
fmt.Println("metric not found. Success")
req.ResultCh <- Succeed // new
// lockEntry, ok := s.metricLockEntries[req.MetricID]
// if ok {
// lockEntry.WaitQueue = append(lockEntry.WaitQueue, req)
// } else {
// s.metricLockEntries[req.MetricID] = &metricLockEntry{
// XLock: true,
// }
// req.ResultCh <- Succeed
// }
var (
buf = make([]byte, storage.DataPageSize)
databuf = buf[:storage.DataPagePayloadSize]
)
s.metrics[req.MetricID] = &Metric{
metricType: req.MetricType,
fracDigits: req.FracDigits,
buffer: buf,
timestamps: enc.NewTimeDeltaCompressor(databuf, 0),
values: enc.NewValueDeltaCompressor(req.MetricType, req.FracDigits, databuf, 0),
xLock: true,
}
s.storageInbox.Push(storage.MetricAdd{
MetricID: req.MetricID,
MetricType: req.MetricType,
FracDigits: req.FracDigits,
ResultCh: req.ResultCh,
})
}
func (s *Worker) processTryAddMetricReqsImmediatelyAfterDelete(reqs []AddMetricReq) {
@@ -379,7 +389,7 @@ func (s *Worker) onMeasuresAppendCommited(rec storage.MeasuresAppendCommited) {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.NoMetricBug,
fmt.Errorf("finAppendMeasures: metric %d not found",
fmt.Errorf("onMeasuresAppendCommited: metric %d not found",
rec.MetricID))
}
metric.OnMeasuresAppendCommited(rec)
@@ -395,30 +405,25 @@ func (s *Worker) onMeasuresAppendWithGrowCommited(rec storage.MeasuresAppendWith
metric.OnMeasuresAppendWithGrowCommited(rec)
}
type AppendMeasuresResult struct {
ResultCode byte
Written int
}
type AppendMeasuresReq struct {
MetricID uint32
Measures []proto.Measure
ResultCh chan AppendMeasuresResult
ResultCh chan storage.MeasuresAppendResult
}
func (s *Worker) AppendMeasures(req AppendMeasuresReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- AppendMeasuresResult{
req.ResultCh <- storage.MeasuresAppendResult{
ResultCode: NoMetric,
}
return
}
if metric.XLock {
if metric.xLock {
metric.WaitQueue = append(metric.WaitQueue, req)
return
}
metric.StartAppendMeasures(req, s.tmp, s.storageInbox)
metric.AppendMeasures(req, s.tmp, s.storageInbox)
}
type RangeScanResult struct {
@@ -452,7 +457,7 @@ func (s *Worker) RangeScan(req RangeScanReq) {
return
}
if metric.XLock {
if metric.xLock {
metric.WaitQueue = append(metric.WaitQueue, req)
return
}
@@ -473,7 +478,7 @@ type FullScanReq struct {
ResultCh chan FullScanResult
}
func (s *Worker) tryFullScan(req FullScanReq) {
func (s *Worker) FullScan(req FullScanReq) {
metric, ok := s.metrics[req.MetricID]
if !ok {
req.ResultCh <- FullScanResult{
@@ -488,7 +493,7 @@ func (s *Worker) tryFullScan(req FullScanReq) {
return
}
if metric.XLock {
if metric.xLock {
metric.WaitQueue = append(metric.WaitQueue, req)
return
}
@@ -501,7 +506,7 @@ type ListCurrentValuesReq struct {
ResultCh chan struct{}
}
func (s *Worker) tryListCurrentValues(req ListCurrentValuesReq) {
func (s *Worker) ListCurrentValues(req ListCurrentValuesReq) {
for _, metricID := range req.MetricIDs {
metric, ok := s.metrics[metricID]
if ok {
@@ -539,55 +544,32 @@ func (s *Worker) applyCommits(req storage.Changes) {
if req.SnapshotNumberCh != nil {
s.snapshotNumber++
err := WriteSnapshot(WriteSnapshotIn{
SnapshotNumber: s.snapshotNumber,
Dir: s.dir,
WriteBufferSize: 4 * 1024 * 1024, // 1mb
Metrics: s.metrics,
FrozenIndexPagesCount: req.FrozenIndexPagesCount,
IndexPageNumbers: req.IndexPageNumbers,
FrozenDataPagesCount: req.FrozenDataPagesCount,
DataPageNumbers: req.DataPageNumbers,
})
err := writeSnapshot(
qb.GetSnapshotFilePath(s.dir, s.databaseName, s.snapshotNumber),
WriteSnapshotIn{
SnapshotNumber: s.snapshotNumber,
WriteBufferSize: 4 * 1024 * 1024, // 4mb
Metrics: s.metrics,
FrozenIndexPagesCount: req.FrozenIndexPagesCount,
IndexPageNumbers: req.IndexPageNumbers,
FrozenDataPagesCount: req.FrozenDataPagesCount,
DataPageNumbers: req.DataPageNumbers,
})
if err != nil {
qb.Abort(qb.WriteSnapshotFailed, err)
}
req.SnapshotNumberCh <- s.snapshotNumber
}
}
func (s *Worker) onMetricAddCommited(rec storage.MetricAddCommited) {
// fix lock
_, ok := s.metrics[rec.MetricID]
if ok {
metric, ok := s.metrics[rec.MetricID]
if !ok {
qb.Abort(qb.MetricAddedBug,
fmt.Errorf("addMetric: metric %d already added",
rec.MetricID))
}
// if !lockEntry.XLock {
// qb.Abort(qb.NoXLockBug,
// fmt.Errorf("addMetric: xlock not set for the metric %d",
// rec.MetricID))
// }
// lockEntry.XLock = false
// delete(s.metricLockEntries, rec.MetricID)
}
func (s *Worker) addMetric(rec storage.MetricAddRecord) {
var (
buf = make([]byte, storage.DataPageSize)
databuf = buf[:storage.DataPagePayloadSize]
)
s.metrics[rec.MetricID] = &Metric{
metricType: rec.MetricType,
fracDigits: rec.FracDigits,
buffer: buf,
timestamps: enc.NewTimeDeltaCompressor(databuf, 0),
values: enc.NewValueDeltaCompressor(rec.MetricType, rec.FracDigits, databuf, 0),
fmt.Errorf("metric %d not found after commit", rec.MetricID))
}
metric.xLock = false
rec.ResultCh <- Succeed // new
}
func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) {
@@ -598,7 +580,7 @@ func (s *Worker) onMetricDeleteCommited(rec storage.MetricDeleteCommited) {
rec.MetricID))
}
if !metric.XLock {
if !metric.xLock {
qb.Abort(qb.NoXLockBug,
fmt.Errorf("deleteMetric: xlock not set for the metric %d",
rec.MetricID))
@@ -669,13 +651,13 @@ func (s *Worker) onMeasuresDeleteCommited(rec storage.MeasuresDeleteCommited) {
rec.MetricID))
}
if !metric.XLock {
if !metric.xLock {
qb.Abort(qb.NoXLockBug,
fmt.Errorf("deleteMeasures: xlock not set for the metric %d",
rec.MetricID))
}
metric.OnMeasuresDeleteCommited(rec)
metric.XLock = false
metric.xLock = false
// FIX add in storage
// if len(rec.FreePageNumbers) > 0 {
// s.freeList.AddPages(rec.FreePageNumbers)