Unverified Commit 073e93c7 by Brian Brazil Committed by GitHub

Gracefully handle unknown WAL record types. (#8004)

As we're looking to expand what's in the WAL,
having old Prometheus servers ignore the new record types
rather than treating them as corruption allows for better
upgrade/downgrade paths.

Adjust some tests accordingly, so they're still testing what they're
meant to test.
Signed-off-by: 's avatarBrian Brazil <brian.brazil@robustperception.io>
parent d2532512
......@@ -510,12 +510,7 @@ func (h *Head) loadWAL(r *wal.Reader, multiRef map[uint64]uint64, mmappedChunks
}
decoded <- tstones
default:
decodeErr = &wal.CorruptionErr{
Err: errors.Errorf("invalid record type %v", dec.Type(rec)),
Segment: r.Segment(),
Offset: r.Offset(),
}
return
// Noop.
}
}
}()
......
......@@ -319,6 +319,13 @@ func TestHead_WALMultiRef(t *testing.T) {
}}, series)
}
func TestHead_UnknownWALRecord(t *testing.T) {
head, w := newTestHead(t, 1000, false)
w.Log([]byte{255, 42})
testutil.Ok(t, head.Init(0))
testutil.Ok(t, head.Close())
}
func TestHead_Truncate(t *testing.T) {
h, _ := newTestHead(t, 1000, false)
defer func() {
......@@ -1208,18 +1215,6 @@ func TestWalRepair_DecodingError(t *testing.T) {
totalRecs int
expRecs int
}{
"invalid_record": {
func(rec []byte) []byte {
// Do not modify the base record because it is Logged multiple times.
res := make([]byte, len(rec))
copy(res, rec)
res[0] = byte(record.Invalid)
return res
},
enc.Series([]record.RefSeries{{Ref: 1, Labels: labels.FromStrings("a", "b")}}, []byte{}),
9,
5,
},
"decode_series": {
func(rec []byte) []byte {
return rec[:3]
......
......@@ -28,8 +28,8 @@ import (
type Type uint8
const (
// Invalid is returned for unrecognised WAL record types.
Invalid Type = 255
// Unknown is returned for unrecognised WAL record types.
Unknown Type = 255
// Series is used to match WAL records of type Series.
Series Type = 1
// Samples is used to match WAL records of type Samples.
......@@ -62,16 +62,16 @@ type Decoder struct {
}
// Type returns the type of the record.
// Returns RecordInvalid if no valid record type is found.
// Returns RecordUnknown if no valid record type is found.
func (d *Decoder) Type(rec []byte) Type {
if len(rec) < 1 {
return Invalid
return Unknown
}
switch t := Type(rec[0]); t {
case Series, Samples, Tombstones:
return t
}
return Invalid
return Unknown
}
// Series appends series in rec to the given slice.
......
......@@ -135,8 +135,8 @@ func TestRecord_Type(t *testing.T) {
testutil.Equals(t, Tombstones, recordType)
recordType = dec.Type(nil)
testutil.Equals(t, Invalid, recordType)
testutil.Equals(t, Unknown, recordType)
recordType = dec.Type([]byte{0})
testutil.Equals(t, Invalid, recordType)
testutil.Equals(t, Unknown, recordType)
}
......@@ -221,7 +221,8 @@ func Checkpoint(logger log.Logger, w *WAL, from, to int, keep func(id uint64) bo
stats.DroppedTombstones += len(tstones) - len(repl)
default:
return nil, errors.New("invalid record type")
// Unknown record type, probably from a future Prometheus version.
continue
}
if len(buf[start:]) == 0 {
continue // All contents discarded.
......
......@@ -140,6 +140,8 @@ func TestCheckpoint(t *testing.T) {
{Ref: 1, Labels: labels.FromStrings("a", "b", "c", "1")},
}, nil))
testutil.Ok(t, err)
// Log an unknown record, that might have come from a future Prometheus version.
testutil.Ok(t, w.Log([]byte{255}))
testutil.Ok(t, w.Close())
// Start a WAL and write records to it as usual.
......@@ -225,7 +227,7 @@ func TestCheckpoint(t *testing.T) {
}
func TestCheckpointNoTmpFolderAfterError(t *testing.T) {
// Create a new wal with an invalid records.
// Create a new wal with invalid data.
dir, err := ioutil.TempDir("", "test_checkpoint")
testutil.Ok(t, err)
defer func() {
......@@ -233,10 +235,19 @@ func TestCheckpointNoTmpFolderAfterError(t *testing.T) {
}()
w, err := NewSize(nil, nil, dir, 64*1024, false)
testutil.Ok(t, err)
testutil.Ok(t, w.Log([]byte{99}))
w.Close()
var enc record.Encoder
testutil.Ok(t, w.Log(enc.Series([]record.RefSeries{
{Ref: 0, Labels: labels.FromStrings("a", "b", "c", "2")}}, nil)))
testutil.Ok(t, w.Close())
// Run the checkpoint and since the wal contains an invalid records this should return an error.
// Corrupt data.
f, err := os.OpenFile(filepath.Join(w.Dir(), "00000000"), os.O_WRONLY, 0666)
testutil.Ok(t, err)
_, err = f.WriteAt([]byte{42}, 1)
testutil.Ok(t, err)
testutil.Ok(t, f.Close())
// Run the checkpoint and since the wal contains corrupt data this should return an error.
_, err = Checkpoint(log.NewNopLogger(), w, 0, 1, nil, 0)
testutil.NotOk(t, err)
......
......@@ -507,13 +507,10 @@ func (w *Watcher) readSegment(r *LiveReader, segmentNum int, tail bool) error {
}
case record.Tombstones:
// noop
case record.Invalid:
return errors.New("invalid record")
default:
// Could be corruption, or reading from a WAL from a newer Prometheus.
w.recordDecodeFailsMetric.Inc()
return errors.New("unknown TSDB record type")
}
}
return errors.Wrapf(r.Err(), "segment %d: %v", segmentNum, r.Err())
......@@ -526,8 +523,6 @@ func (w *Watcher) SetStartTime(t time.Time) {
func recordType(rt record.Type) string {
switch rt {
case record.Invalid:
return "invalid"
case record.Series:
return "series"
case record.Samples:
......
......@@ -278,6 +278,8 @@ func TestReadToEndWithCheckpoint(t *testing.T) {
},
}, nil)
testutil.Ok(t, w.Log(series))
// Add in an unknown record type, which should be ignored.
testutil.Ok(t, w.Log([]byte{255}))
for j := 0; j < samplesCount; j++ {
inner := rand.Intn(ref + 1)
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment