diff --git a/client/client.go b/client/client.go index 8b695da..2e78d7f 100644 --- a/client/client.go +++ b/client/client.go @@ -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) diff --git a/database/api.go b/database/api.go index 3ec29b2..8a309fc 100644 --- a/database/api.go +++ b/database/api.go @@ -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: diff --git a/database/database.go b/database/database.go index 7e404cb..be3f084 100644 --- a/database/database.go +++ b/database/database.go @@ -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{ diff --git a/database_linux b/database_linux index c09b0f8..d14f7d9 100755 Binary files a/database_linux and b/database_linux differ diff --git a/enc/enc_test.go b/enc/enc_test.go index ec03ae6..052142d 100644 --- a/enc/enc_test.go +++ b/enc/enc_test.go @@ -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{ diff --git a/enc/time_delta.go b/enc/time_delta.go index 57010fe..3a9a33f 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -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 } } diff --git a/enc/value_delta.go b/enc/value_delta.go index 6e116a9..6169b6b 100644 --- a/enc/value_delta.go +++ b/enc/value_delta.go @@ -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() { diff --git a/examples/play/main.go b/examples/play/main.go index 418842e..fb6f541 100644 --- a/examples/play/main.go +++ b/examples/play/main.go @@ -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) diff --git a/examples/play/play b/examples/play/play index b4b8dcf..dd5e877 100755 Binary files a/examples/play/play and b/examples/play/play differ diff --git a/examples/requests/requests.go b/examples/requests/requests.go index 12c2467..02cdfe0 100644 --- a/examples/requests/requests.go +++ b/examples/requests/requests.go @@ -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 diff --git a/go.sum b/go.sum index 52f17ce..0d2d171 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/inbox/inbox.go b/inbox/inbox.go index 6c83a8d..3e43068 100644 --- a/inbox/inbox.go +++ b/inbox/inbox.go @@ -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 } diff --git a/qb.go b/qb.go index 37955b7..0f76350 100644 --- a/qb.go +++ b/qb.go @@ -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 } diff --git a/recovery/recovery.go b/recovery/recovery.go index 502e74f..6fee097 100644 --- a/recovery/recovery.go +++ b/recovery/recovery.go @@ -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 diff --git a/storage/replay_metric.go b/storage/replay_metric.go index 768064e..afadbe0 100644 --- a/storage/replay_metric.go +++ b/storage/replay_metric.go @@ -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{ diff --git a/storage/storage.go b/storage/storage.go index 6dbeeea..49bd577 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -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 } diff --git a/storage/storage_test.go b/storage/storage_test.go index 6951072..b07bb0a 100644 --- a/storage/storage_test.go +++ b/storage/storage_test.go @@ -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) +} diff --git a/storage/wal_reader.go b/storage/wal_reader.go index 99866ee..571fd96 100644 --- a/storage/wal_reader.go +++ b/storage/wal_reader.go @@ -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) diff --git a/storage/wal_replayer.go b/storage/wal_replayer.go index 656013b..f9d8296 100644 --- a/storage/wal_replayer.go +++ b/storage/wal_replayer.go @@ -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 diff --git a/storage/write_preparer.go b/storage/write_preparer.go index 69a5640..ca31e75 100644 --- a/storage/write_preparer.go +++ b/storage/write_preparer.go @@ -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:] +} diff --git a/storage/wal_writer.go b/storage/writer.go similarity index 93% rename from storage/wal_writer.go rename to storage/writer.go index 07233d1..b1abafe 100644 --- a/storage/wal_writer.go +++ b/storage/writer.go @@ -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 diff --git a/testdir/test.data b/testdir/test.data new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.data_free b/testdir/test.data_free new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.data_free_delta b/testdir/test.data_free_delta new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.index b/testdir/test.index new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.index_free b/testdir/test.index_free new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.index_free_delta b/testdir/test.index_free_delta new file mode 100644 index 0000000..e69de29 diff --git a/testdir/test.wal_0 b/testdir/test.wal_0 new file mode 100755 index 0000000..974036d Binary files /dev/null and b/testdir/test.wal_0 differ diff --git a/transform/raw.go b/transform/raw.go index 4fa9264..9b149de 100644 --- a/transform/raw.go +++ b/transform/raw.go @@ -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() diff --git a/transform/responder.go b/transform/responder.go index 9e3636d..988e06c 100644 --- a/transform/responder.go +++ b/transform/responder.go @@ -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 { diff --git a/worker/metric.go b/worker/metric.go index 679255a..7203f26 100644 --- a/worker/metric.go +++ b/worker/metric.go @@ -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 } diff --git a/worker/snapshot.go b/worker/snapshot.go index 3b906a3..da9f7fb 100644 --- a/worker/snapshot.go +++ b/worker/snapshot.go @@ -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 diff --git a/worker/snapshot_test.go b/worker/snapshot_test.go index 82d97ae..b768841 100644 --- a/worker/snapshot_test.go +++ b/worker/snapshot_test.go @@ -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 } diff --git a/worker/worker.go b/worker/worker.go index 99efb43..5114354 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -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)