diff --git a/database/metric.go b/database/metric.go index 8ddda8b..5f9a4c2 100644 --- a/database/metric.go +++ b/database/metric.go @@ -61,9 +61,8 @@ func (s *_metric) DeleteMeasures() { } type CapturedState struct { - Pages []storage.DataPayload - TimeState qb.TimeDeltaCapturedState - ValueState qb.ValueDeltaCapturedState + LastTimestamp uint32 + LastValue float64 } // func (s *_metric) startAppendMeasures(req tryAppendMeasuresReq, sendToStorage func(storage.AppendedMeasures)) { @@ -76,19 +75,24 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send values = s.values // офсети на head сторінці - timestampsOffset int - valuesOffset int + timestampsOffset int + timestampsRewindOffset int + valuesOffset int + valuesRewindOffset int headTimestamps []byte headValues []byte + pages []storage.DataPayload + written int resultCode byte ) s.capturedState = &CapturedState{ - // - + LastTimestamp: timestamps.LastTimestamp(), + LastValue: s.lastValue, + //LastValue: values.LastValue(), } s.timestamps.CaptureState() @@ -111,20 +115,25 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send if totalRequiredSpace <= len(s.buffer) { // якщо на сторінці є місце + timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp) + values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta) + if idx == 0 { timestampsOffset = tReport.Offset + timestampsRewindOffset = tReport.RewindOffset valuesOffset = vReport.Offset + valuesRewindOffset = vReport.RewindOffset } } else { // сторінка заповнена since := s.timestamps.ReplaceSinceWithUntil() - if len(s.capturedState.Pages) == 0 { + if len(pages) == 0 { headTimestamps = timestamps.Tail(timestampsOffset) headValues = values.Tail(valuesOffset) } - s.capturedState.Pages = append(s.capturedState.Pages, storage.DataPayload{ + pages = append(pages, storage.DataPayload{ Since: since, Content: s.buffer, TimestampsSize: timestamps.Size(), @@ -139,10 +148,11 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send // renew s.buffer = buf + + timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp) + values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta) } - timestamps.Append(tReport.RewindOffset, tmp[:tReport.ChangeSize], measure.Timestamp) - values.Append(vReport.RewindOffset, tmp[7:7+vReport.ChangeSize], measure.Value, vReport.Delta) // s.lastValue = measure.Value written++ @@ -159,26 +169,26 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send // виділити змінені байти. // скопіювати. Причому можна скопіювати зрізи chunks - if len(s.capturedState.Pages) > 0 { + if len(pages) > 0 { // пишу в storage довгим шляхом через redo файл і запис в data файл - sendToStorage(storage.AppendedMeasures{ + sendToStorage(storage.AppendedMeasuresWithGrow{ MetricID: req.MetricID, LastPageNo: s.lastPageNo, - TimestampsOffset: timestampsOffset, + TimestampsOffset: timestampsRewindOffset, Timestamps: headTimestamps, - ValuesOffset: valuesOffset, + ValuesOffset: valuesRewindOffset, Values: headValues, IndexLevelTails: s.indexLevelTails, - DataPages: s.capturedState.Pages, + DataPages: pages, TailTimestamps: timestamps.Tail(0), // payload TailValues: values.Tail(0), ResultCode: resultCode, WrittenCount: written, - //ResultCh: req.ResultCh, + ResultCh: req.ResultCh, }) } else { // короткий шлях - запис лише в storage - sendToStorage(storage.AppendedMeasures{ + sendToStorage(storage.AppendedMeasuresWithGrow{ MetricID: req.MetricID, TimestampsOffset: timestampsOffset, ValuesOffset: valuesOffset, @@ -186,7 +196,7 @@ func (s *_metric) StartAppendMeasures(req tryAppendMeasuresReq, tmp []byte, send Values: values.Tail(valuesOffset), ResultCode: resultCode, WrittenCount: written, - //ResultCh: req.ResultCh, + ResultCh: req.ResultCh, }) } @@ -207,74 +217,71 @@ func (s *_metric) FinAppendMeasures(rec storage.AppendMeasuresSummary) { // READ func (s *_metric) StartRangeScan(req tryRangeScanReq) { - // if s.Since == 0 { - // req.ResultCh <- rangeScanResult{ - // ResultCode: QueryDone, - // } - // return - // } + if s.Since == 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } - // if req.Since > s.Until { - // req.ResultCh <- rangeScanResult{ - // ResultCode: QueryDone, - // } - // return - // } + if req.Since > s.Until { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } - // if req.Until < s.Since { - // if s.RootPageNo > 0 { - // req.ResultCh <- rangeScanResult{ - // ResultCode: UntilNotFound, - // RootPageNo: s.RootPageNo, - // FracDigits: s.FracDigits, - // } - // s.RLocks++ - // return - // } else { - // req.ResultCh <- rangeScanResult{ - // ResultCode: QueryDone, - // } - // return - // } - // } + if req.Until < s.Since { + if s.RootPageNo > 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: UntilNotFound, + RootPageNo: s.RootPageNo, + FracDigits: s.FracDigits, + } + s.RLocks++ + return + } else { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + } - // timestampDecompressor := s.Timestamps.CreateDecompressor() - // valueDecompressor := s.Values.CreateDecompressor(s.FracDigits) + timestampDecompressor := s.timestamps.CreateDecompressor() + valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits) - // for { - // timestamp, done := timestampDecompressor.NextValue() - // if done { - // break - // } - - // value, done := valueDecompressor.NextValue() - // if done { - // qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - // } - - // if timestamp <= req.Until { - // req.ResponseWriter.FeedNoSend(timestamp, value) - // if timestamp < req.Since { - // req.ResultCh <- rangeScanResult{ - // ResultCode: QueryDone, - // } - // return - // } - // } - // } - - // if s.LastPageNo > 0 { - // req.ResultCh <- rangeScanResult{ - // ResultCode: UntilFound, - // LastPageNo: s.LastPageNo, - // FracDigits: s.FracDigits, - // } - // s.RLocks++ - // } else { - // req.ResultCh <- rangeScanResult{ - // ResultCode: QueryDone, - // } - // } + for { + timestamp, done := timestampDecompressor.NextValue() + if done { + break + } + value, done := valueDecompressor.NextValue() + if done { + qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + } + if timestamp <= req.Until { + req.ResponseWriter.FeedNoSend(timestamp, value) + if timestamp < req.Since { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + return + } + } + } + if s.lastPageNo > 0 { + req.ResultCh <- rangeScanResult{ + ResultCode: UntilFound, + LastPageNo: s.lastPageNo, + FracDigits: s.fracDigits, + } + s.RLocks++ + } else { + req.ResultCh <- rangeScanResult{ + ResultCode: QueryDone, + } + } } func (s *_metric) StartFullScan(req tryFullScanReq) { @@ -285,33 +292,33 @@ func (s *_metric) StartFullScan(req tryFullScanReq) { // return // } - // timestampDecompressor := s.Timestamps.CreateDecompressor() - // valueDecompressor := s.Values.CreateDecompressor() + timestampDecompressor := s.timestamps.CreateDecompressor() + valueDecompressor := s.values.CreateDecompressor(s.metricType, s.fracDigits) - // for { - // timestamp, done := timestampDecompressor.NextValue() - // if done { - // break - // } - // value, done := valueDecompressor.NextValue() - // if done { - // qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) - // } - // req.ResponseWriter.FeedNoSend(timestamp, value) - // } + for { + timestamp, done := timestampDecompressor.NextValue() + if done { + break + } + value, done := valueDecompressor.NextValue() + if done { + qb.Abort(qb.HasTimestampNoValueBug, ErrNoValueBug) + } + req.ResponseWriter.FeedNoSend(timestamp, value) + } - // if s.LastPageNo > 0 { - // req.ResultCh <- fullScanResult{ - // ResultCode: UntilFound, - // LastPageNo: s.LastPageNo, - // FracDigits: s.FracDigits, - // } - // s.RLocks++ - // } else { - // req.ResultCh <- fullScanResult{ - // ResultCode: QueryDone, - // } - // } + if s.lastPageNo > 0 { + req.ResultCh <- fullScanResult{ + ResultCode: UntilFound, + LastPageNo: s.lastPageNo, + FracDigits: s.fracDigits, + } + s.RLocks++ + } else { + req.ResultCh <- fullScanResult{ + ResultCode: QueryDone, + } + } } // індекси @@ -342,29 +349,21 @@ func (s *_metric) WriteTo(w io.Writer) (err error) { if err != nil { return } - err = bin.WriteUint16(w, uint16(s.timestamps.Size())) + err = bin.WriteUint16(w, uint16(s.timestamps.StoredSize())) if err != nil { return } - err = bin.WriteUint16(w, uint16(s.values.Size())) + err = bin.WriteUint16(w, uint16(s.values.StoredSize())) if err != nil { return } - var ( - timestampsCapturedState *qb.TimeDeltaCapturedState - valuesCapturedState *qb.ValueDeltaCapturedState - ) - if s.capturedState != nil { - timestampsCapturedState = &s.capturedState.TimeState - valuesCapturedState = &s.capturedState.ValueState - } // timestamps payload - err = s.timestamps.WritePayloadTo(w, timestampsCapturedState) + err = s.timestamps.WriteStoredTo(w) if err != nil { return } // values payload - err = s.values.WritePayloadTo(w, valuesCapturedState) + err = s.values.WriteStoredTo(w) if err != nil { return } diff --git a/database/replay_metric.go b/database/replay_metric.go index 1c34d18..d56649a 100644 --- a/database/replay_metric.go +++ b/database/replay_metric.go @@ -87,13 +87,13 @@ func (s *ReplayMetric) ReadFrom(r io.Reader) (err error) { // } func (s *ReplayMetric) MeasuresAppend(rec storage.MeasuresAppendRecord) { - timestampsSize := rec.TimestampsOffset - len(rec.Timestamps) + timestampsSize := rec.TimestampsRewindOffset - len(rec.Timestamps) // копіюю нові дані із WAL з урахуванням offset-ів - copy(s.databuf[rec.ValuesOffset:], rec.Values) + copy(s.databuf[rec.ValuesRewindOffset:], rec.Values) copy(s.databuf[len(s.databuf)-timestampsSize:], rec.Timestamps) // s.tSize = timestampsSize - s.vSize = rec.ValuesOffset - len(rec.Values) + s.vSize = rec.ValuesRewindOffset - len(rec.Values) } type AppendMeasuresResult struct { diff --git a/enc/time_delta.go b/enc/time_delta.go index 9bea0d0..77eb208 100644 --- a/enc/time_delta.go +++ b/enc/time_delta.go @@ -236,21 +236,21 @@ func (s *TimeDeltaCompressor) StoredSize() int { // } // } -func (s *TimeDeltaCompressor) WritePayloadTo(w io.Writer, state *qb.TimeDeltaCapturedState) (err error) { - if state == nil { +func (s *TimeDeltaCompressor) WriteStoredTo(w io.Writer) (err error) { + if s.state == nil { _, err = w.Write(s.buf[:s.pos]) return } else { - _, err = w.Write(state.Payload) + _, err = w.Write(s.state.Payload) if err != nil { return } - _, err = bin.WriteVarUint64(w, uint64(state.LastDelta)) + _, err = bin.WriteVarUint64(w, uint64(s.state.LastDelta)) if err != nil { return } _, err = w.Write([]byte{ - state.H, + s.state.H, }) return } diff --git a/enc/value_delta.go b/enc/value_delta.go index f42b266..53c1471 100644 --- a/enc/value_delta.go +++ b/enc/value_delta.go @@ -198,21 +198,21 @@ func (s *ValueDeltaCompressor) StoredSize() int { // } // } -func (s *ValueDeltaCompressor) WritePayloadTo(w io.Writer, state *qb.ValueDeltaCapturedState) (err error) { - if state == nil { +func (s *ValueDeltaCompressor) WriteStoredTo(w io.Writer) (err error) { + if s.state == nil { _, err = w.Write(s.buf[:s.pos]) return } else { - _, err = w.Write(state.Payload) + _, err = w.Write(s.state.Payload) if err != nil { return } - _, err = bin.WriteVarUint64(w, state.LastDelta) + _, err = bin.WriteVarUint64(w, s.state.LastDelta) if err != nil { return } _, err = w.Write([]byte{ - state.H, + s.state.H, }) return } diff --git a/qb.go b/qb.go index 2ddb421..8394a00 100644 --- a/qb.go +++ b/qb.go @@ -29,7 +29,7 @@ type TimestampCompressor interface { // (rewindOffset, change, timestamp) Append(int, []byte, uint32) Size() int - //Chunks() [][]byte + StoredSize() int //DeleteLast() CaptureState() // (offset) => payload @@ -38,7 +38,7 @@ type TimestampCompressor interface { //Payload() []byte // для снапшота // Offset() int ReplaceBuffer([]byte) - WritePayloadTo(io.Writer, *TimeDeltaCapturedState) error + WriteStoredTo(io.Writer) error LastTimestamp() uint32 ReplaceSinceWithUntil() uint32 } @@ -63,6 +63,7 @@ type ValueCompressor interface { // (rewindOffset, change, value, delta) Append(int, []byte, float64, uint64) Size() int + StoredSize() int //Chunks() [][]byte //DeleteLast() CaptureState() @@ -73,7 +74,7 @@ type ValueCompressor interface { //Payload() []byte // для снапшота //Offset() int ReplaceBuffer([]byte) - WritePayloadTo(io.Writer, *ValueDeltaCapturedState) error + WriteStoredTo(io.Writer) error LastValue() float64 } diff --git a/storage/data_preparer.go b/storage/data_preparer.go index 37e724d..79367a2 100644 --- a/storage/data_preparer.go +++ b/storage/data_preparer.go @@ -57,7 +57,28 @@ func (s *DataPreparer) Prepare(input []any) PreparedData { // if len(x.FreePageNumbers) > 0 { // s.freeList.AddPageNumbers(x.FreePageNumbers) // } - case AppendedMeasures: + case MeasuresAppendRecord: + rec := MeasuresAppendRecord{ + MetricID: x.MetricID, + TimestampsRewindOffset: x.TimestampsRewindOffset, + Timestamps: x.Timestamps, + ValuesRewindOffset: x.ValuesRewindOffset, + Values: x.Values, + } + + rec.WriteTo(s.w) + + // Дані для Worker + s.writeResults = append(s.writeResults, AppendMeasuresSummary{ + // MetricID: x.MetricID, + // LastPageNo: sealResult.DataPages[len(sealResult.DataPages)-1].PageNo, + // Index: indexLevelTails, + // ResultCode: x.ResultCode, + // WrittenCount: x.WrittenCount, + // ResultCh: x.ResultCh, + }) + + case AppendedMeasuresWithGrow: sealResult := s.pagePreparer.SealPages(SealPagesIn{ LastPageNo: x.LastPageNo, IndexLevelTails: x.IndexLevelTails, @@ -178,7 +199,7 @@ func (s *DataPreparer) Reset() { } // MeasuresToWrite -type AppendedMeasures struct { +type AppendedMeasuresWithGrow struct { MetricID uint32 LastPageNo uint32 TimestampsOffset int // offset on 1st page diff --git a/storage/wal_records.go b/storage/wal_records.go index 5e479cd..c8671d4 100644 --- a/storage/wal_records.go +++ b/storage/wal_records.go @@ -306,20 +306,20 @@ func (s *MeasuresAppendWithGrowRecord) Parse(r *bytes.Buffer) (err error) { } type MeasuresAppendRecord struct { - MetricID uint32 - TimestampsOffset int - Timestamps []byte - ValuesOffset int - Values []byte + MetricID uint32 + TimestampsRewindOffset int + Timestamps []byte + ValuesRewindOffset int + Values []byte } func (s MeasuresAppendRecord) WriteTo(w *bytes.Buffer) { w.WriteByte(CodeMeasuresAppendWithGrow) bin.WriteUint32(w, s.MetricID) - bin.WriteVarSize(w, s.TimestampsOffset) + bin.WriteVarSize(w, s.TimestampsRewindOffset) bin.WriteVarSize(w, len(s.Timestamps)) w.Write(s.Timestamps) - bin.WriteVarSize(w, s.ValuesOffset) + bin.WriteVarSize(w, s.ValuesRewindOffset) bin.WriteVarSize(w, len(s.Values)) w.Write(s.Values) } @@ -330,7 +330,7 @@ func (s *MeasuresAppendRecord) Parse(r *bytes.Buffer) (err error) { if err != nil { return } - s.TimestampsOffset, err = bin.ReadVarSize(r) + s.TimestampsRewindOffset, err = bin.ReadVarSize(r) if err != nil { return } @@ -342,7 +342,7 @@ func (s *MeasuresAppendRecord) Parse(r *bytes.Buffer) (err error) { if err != nil { return } - s.ValuesOffset, err = bin.ReadVarSize(r) + s.ValuesRewindOffset, err = bin.ReadVarSize(r) if err != nil { return }