Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion internal/hashlogcompare/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ const (
ReasonNotFound = "not_found"
ReasonUnavailable = "unavailable"
ReasonCoverageGap = "coverage_gap"
ReasonTornRow = "torn_row"
ReasonInvalidRow = "invalid_row"
)

Expand Down Expand Up @@ -103,7 +104,7 @@ func (m *Metrics) initPair(chain, migrating, reserve string, now time.Time) {
m.ComparedHeights.WithLabelValues(chain, migrating, reserve)
for _, role := range []string{RoleMigrating, RoleReserve} {
m.HeightGaps.WithLabelValues(chain, migrating, reserve, role)
for _, reason := range []string{ReasonNotFound, ReasonUnavailable, ReasonCoverageGap, ReasonInvalidRow} {
for _, reason := range []string{ReasonNotFound, ReasonUnavailable, ReasonCoverageGap, ReasonTornRow, ReasonInvalidRow} {
m.SourceErrors.WithLabelValues(chain, migrating, reserve, role, reason)
}
}
Expand Down
19 changes: 16 additions & 3 deletions internal/hashlogcompare/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ const (
// while it was being read), so the reader restarts at the tip.
var errCoverageGap = errors.New("hash log coverage gap")

// errTornRow means a sealed file ends in an incomplete line: a row that a
// crash tore mid-write. The reader drops it and moves to the next file.
var errTornRow = errors.New("hash log torn row")

// Source is the subset of the sidecar client the reader uses.
type Source interface {
ListHashLog(ctx context.Context) ([]sidecar.HashLogFile, error)
Expand All @@ -40,7 +44,11 @@ type Row struct {

// Reader tails one node's hash log in file-index order. It keeps a byte
// cursor into the current file and only consumes complete lines, so a line
// that is half written when read is picked up whole on the next poll.
// that is half written when read is picked up whole on the next poll. A
// sealed file gets no more bytes, so the reader drops an incomplete last
// line there, as the HashLogger's own reader does. The HashLogger seals only
// a file with a complete header and at least one complete row, so that last
// line is the only incomplete line a sealed file can hold.
type Reader struct {
src Source

Expand All @@ -60,7 +68,8 @@ func NewReader(src Source) *Reader {

// Poll returns the complete rows written since the last call, plus the number
// of lines that could not be parsed. Transport errors leave the cursor where
// it was; errCoverageGap resets it to the tip.
// it was; errCoverageGap resets it to the tip; errTornRow moves it past the
// dropped line.
func (r *Reader) Poll(ctx context.Context) (rows []Row, invalid int, err error) {
if !r.started {
ok, err := r.start(ctx)
Expand Down Expand Up @@ -111,9 +120,13 @@ func (r *Reader) Poll(ctx context.Context) (rows []Row, invalid int, err error)
rows = append(rows, parsed...)
invalid += bad

if !r.sealed || len(complete) < len(data) {
if !r.sealed {
return rows, invalid, nil
}
if torn := len(data) - len(complete); torn > 0 {
r.offset += int64(torn)
return rows, invalid, fmt.Errorf("%w: dropped %d trailing bytes of sealed file %s", errTornRow, torn, r.name)
}
more, err := r.nextFile(ctx)
if err != nil || !more {
return rows, invalid, err
Expand Down
21 changes: 21 additions & 0 deletions internal/hashlogcompare/reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,27 @@ func TestReader_FollowsSealAndNextFile(t *testing.T) {
g.Expect(got).To(BeEmpty())
}

func TestReader_DropsTornTailOfSealedFile(t *testing.T) {
g := NewWithT(t)
src := &fakeSource{}
f := src.add(1, testHeader+rows(10, 11))
r := NewReader(src)

got, _, err := r.Poll(t.Context())
g.Expect(err).NotTo(HaveOccurred())
g.Expect(heights(got)).To(Equal(span(10, 11)))

f.seal(10, 12, row(12, "")+row(13, "")[:5])
src.add(2, testHeader+rows(13, 14))
got, _, err = r.Poll(t.Context())
g.Expect(err).To(MatchError(errTornRow))
g.Expect(heights(got)).To(Equal([]int64{12}))

got, _, err = r.Poll(t.Context())
g.Expect(err).NotTo(HaveOccurred())
g.Expect(heights(got)).To(Equal(span(13, 14)))
}

func TestReader_SkipsRemovedEmptyFile(t *testing.T) {
g := NewWithT(t)
src := &fakeSource{}
Expand Down
2 changes: 2 additions & 0 deletions internal/hashlogcompare/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,8 @@ func (pr *pairRunner) read(ctx context.Context, role string, r *Reader, buf map[
reason = ReasonNotFound
case errors.Is(err, errCoverageGap):
reason = ReasonCoverageGap
case errors.Is(err, errTornRow):
reason = ReasonTornRow
}
m.SourceErrors.WithLabelValues(append(roleLabels, reason)...).Inc()
pr.c.Log.Warn("hash log read failed", pr.logAttrs(role, "reason", reason, "error", err)...)
Expand Down
24 changes: 23 additions & 1 deletion internal/hashlogcompare/run_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,28 @@ func TestPairRunner_MissingReserveLogIsSourceError(t *testing.T) {
g.Expect(value(t, m.ComparedHeights.WithLabelValues(pr.labels...))).To(BeZero())
}

func TestPairRunner_TornReserveRowIsSourceError(t *testing.T) {
g := NewWithT(t)
migrating, reserve := &fakeSource{}, &fakeSource{}
migrating.add(1, testHeader+rows(1, 6))
rf := reserve.add(1, testHeader+rows(1, 3)+row(4, "")[:5])
pr, m := newTestRunner(migrating, reserve)
reserveTorn := m.SourceErrors.WithLabelValues(append(pr.labels, RoleReserve, ReasonTornRow)...)

pr.poll(t.Context())
g.Expect(value(t, reserveTorn)).To(BeZero())

rf.seal(1, 3, "")
reserve.add(2, testHeader+rows(4, 6))
pr.poll(t.Context())
g.Expect(value(t, reserveTorn)).To(Equal(1.0))

pr.poll(t.Context())
g.Expect(value(t, reserveTorn)).To(Equal(1.0))
g.Expect(value(t, m.LastComparedHeight.WithLabelValues(pr.labels...))).To(Equal(6.0))
g.Expect(value(t, m.SourceHeight.WithLabelValues(append(pr.labels, RoleReserve)...))).To(Equal(6.0))
}

func TestInitPair_ExportsZeroSeries(t *testing.T) {
g := NewWithT(t)
reg := prometheus.NewRegistry()
Expand All @@ -123,5 +145,5 @@ func TestInitPair_ExportsZeroSeries(t *testing.T) {
g.Expect(series).To(HaveKeyWithValue("sei_hashlog_compare_mismatches_total", 1))
g.Expect(series).To(HaveKeyWithValue("sei_hashlog_compare_last_compared_timestamp_seconds", 1))
g.Expect(value(t, m.LastComparedTimestamp.WithLabelValues("c", "a", "b"))).To(Equal(1_800_000_000.0))
g.Expect(series).To(HaveKeyWithValue("sei_hashlog_compare_source_errors_total", 2*4))
g.Expect(series).To(HaveKeyWithValue("sei_hashlog_compare_source_errors_total", 2*5))
}
Loading