diff --git a/sei-tendermint/internal/autobahn/consensus/persist/blocks.go b/sei-tendermint/internal/autobahn/consensus/persist/blocks.go index 23b894113c..ed238bbad3 100644 --- a/sei-tendermint/internal/autobahn/consensus/persist/blocks.go +++ b/sei-tendermint/internal/autobahn/consensus/persist/blocks.go @@ -1,4 +1,3 @@ -// TODO: add Prometheus metrics for blocks written and truncated. package persist import ( @@ -41,6 +40,8 @@ type LoadedBlock struct { type laneWALState struct { wal seiwal.WAL[*types.Signed[*types.LaneProposal]] nextBlockNum types.BlockNumber + // appended is the number of blocks accepted for append since the last successful flush. + appended uint64 } // persistBlock schedules a proposal for the WAL and advances nextBlockNum. The block is not durable @@ -52,10 +53,12 @@ func (s *laneWALState) persistBlock(proposal *types.Signed[*types.LaneProposal]) if s.nextBlockNum > 0 && h.BlockNumber() != s.nextBlockNum { return fmt.Errorf("block %s/%d out of sequence (next=%d)", h.Lane(), h.BlockNumber(), s.nextBlockNum) } + addMetricsRecords(walBlocks, stageAsked, 1) if err := s.wal.Append(uint64(h.BlockNumber()), proposal); err != nil { return fmt.Errorf("persist block %s/%d: %w", h.Lane(), h.BlockNumber(), err) } s.nextBlockNum = h.BlockNumber() + 1 + s.appended++ return nil } @@ -64,6 +67,8 @@ func (s *laneWALState) flush(lane types.LaneID) error { if err := s.wal.Flush(); err != nil { return fmt.Errorf("flush lane %s WAL: %w", lane, err) } + addMetricsRecords(walBlocks, stagePersisted, s.appended) + s.appended = 0 return nil } diff --git a/sei-tendermint/internal/autobahn/consensus/persist/commitqcs.go b/sei-tendermint/internal/autobahn/consensus/persist/commitqcs.go index 3692a3872f..313fe056fc 100644 --- a/sei-tendermint/internal/autobahn/consensus/persist/commitqcs.go +++ b/sei-tendermint/internal/autobahn/consensus/persist/commitqcs.go @@ -1,4 +1,3 @@ -// TODO: add Prometheus metrics for commitQCs written and truncated. package persist import ( @@ -26,6 +25,8 @@ type commitQCState struct { // Whether a QC has been appended since the last flush, so a prune that re-persists its anchor is // still made durable while a run of duplicates costs no fsync. unflushed bool + // appended is the number of QCs accepted for append since the last successful flush. + appended uint64 } // persistCommitQC schedules a CommitQC for the WAL under its own road index. The QC is not durable @@ -41,10 +42,12 @@ func (s *commitQCState) persist(qc *types.CommitQC) error { return fmt.Errorf("commitqc %d out of sequence (next=%d)", idx, s.persisted.Next) } if w, ok := s.wal.Get(); ok { + addMetricsRecords(walCommitQCs, stageAsked, 1) if err := w.Append(uint64(idx), qc); err != nil { return fmt.Errorf("persist commitqc %d: %w", idx, err) } s.unflushed = true + s.appended++ } s.persisted.Next += 1 return nil @@ -59,6 +62,8 @@ func (s *commitQCState) flush() error { if err := w.Flush(); err != nil { return fmt.Errorf("flush commitqc WAL: %w", err) } + addMetricsRecords(walCommitQCs, stagePersisted, s.appended) + s.appended = 0 s.unflushed = false return nil } diff --git a/sei-tendermint/internal/autobahn/consensus/persist/metrics.gen.go b/sei-tendermint/internal/autobahn/consensus/persist/metrics.gen.go new file mode 100644 index 0000000000..b1784520a0 --- /dev/null +++ b/sei-tendermint/internal/autobahn/consensus/persist/metrics.gen.go @@ -0,0 +1,31 @@ +// Code generated by metricsgen. DO NOT EDIT. + +package persist + +import ( + "github.com/prometheus/client_golang/prometheus" + tmprometheus "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/prometheus" +) + +var Global = newMetrics() + +func init() { + prometheus.MustRegister( + Global.records, + ) +} + +func newMetrics() *metrics { + return &metrics{ + records: tmprometheus.NewCounterIntVec(prometheus.CounterOpts{ + Namespace: MetricsNamespace, + Subsystem: MetricsSubsystem, + Name: "records", + Help: "Records asked for append, or made durable.", + }, []string{"wal", "stage"}), + } +} + +func (m *metrics) recordsAt(wal string, stage string) *tmprometheus.CounterInt { + return m.records.WithLabelValues(wal, stage) +} diff --git a/sei-tendermint/internal/autobahn/consensus/persist/metrics.go b/sei-tendermint/internal/autobahn/consensus/persist/metrics.go new file mode 100644 index 0000000000..03e4dc1a09 --- /dev/null +++ b/sei-tendermint/internal/autobahn/consensus/persist/metrics.go @@ -0,0 +1,29 @@ +package persist + +import ( + "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/prometheus" +) + +const MetricsNamespace = "tendermint" +const MetricsSubsystem = "internal_autobahn_consensus_persist" + +const ( + walBlocks = "blocks" + walCommitQCs = "commitqcs" + + stageAsked = "asked" + stagePersisted = "persisted" +) + +//go:generate go run github.com/sei-protocol/sei-chain/sei-tendermint/scripts/metricsgen -struct=metrics +type metrics struct { + // Records asked for append, or made durable. + records prometheus.CounterIntVec `metrics_labels:"wal,stage"` +} + +func addMetricsRecords(wal, stage string, n uint64) { + if n == 0 { + return + } + Global.recordsAt(wal, stage).Add(int64(n)) //nolint:gosec // a record count fits in int64 +} diff --git a/sei-tendermint/internal/autobahn/consensus/persist/metrics_test.go b/sei-tendermint/internal/autobahn/consensus/persist/metrics_test.go new file mode 100644 index 0000000000..64367d0aa8 --- /dev/null +++ b/sei-tendermint/internal/autobahn/consensus/persist/metrics_test.go @@ -0,0 +1,71 @@ +package persist + +import ( + "testing" + + dto "github.com/prometheus/client_model/go" + + "github.com/sei-protocol/sei-chain/sei-tendermint/autobahn/types" + "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils" + "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/require" +) + +func recordCount(t *testing.T, wal, stage string) int64 { + t.Helper() + var m dto.Metric + require.NoError(t, Global.recordsAt(wal, stage).Write(&m)) + return int64(m.GetCounter().GetValue()) +} + +func TestPersistRecordCounters(t *testing.T) { + t.Run("blocks", func(t *testing.T) { + rng := utils.TestRng() + dir := t.TempDir() + key := types.GenSecretKey(rng) + lane := types.LaneID{Validator: key.Public(), Joined: 0} + bp, _, err := NewBlockPersister(utils.Some(dir)) + require.NoError(t, err) + t.Cleanup(func() { _ = bp.Close() }) + + asked := recordCount(t, walBlocks, stageAsked) + persisted := recordCount(t, walBlocks, stagePersisted) + + require.NoError(t, bp.PruneAndPersist(lane, 0, []*types.Signed[*types.LaneProposal]{ + testSignedProposal(rng, key, 0), + testSignedProposal(rng, key, 1), + })) + require.Equal(t, asked+2, recordCount(t, walBlocks, stageAsked)) + require.Equal(t, persisted+2, recordCount(t, walBlocks, stagePersisted)) + + err = bp.PruneAndPersist(lane, 0, []*types.Signed[*types.LaneProposal]{testSignedProposal(rng, key, 0)}) + require.Error(t, err) + require.Equal(t, asked+2, recordCount(t, walBlocks, stageAsked)) + require.Equal(t, persisted+2, recordCount(t, walBlocks, stagePersisted)) + }) + + t.Run("commitqcs", func(t *testing.T) { + rng := utils.TestRng() + dir := t.TempDir() + committee, keys := genTestCommittee(rng, 4) + qcs := makeSequentialCommitQCs(committee, keys, 2) + cp, _, err := NewCommitQCPersister(utils.Some(dir)) + require.NoError(t, err) + t.Cleanup(func() { _ = cp.Close() }) + + asked := recordCount(t, walCommitQCs, stageAsked) + persisted := recordCount(t, walCommitQCs, stagePersisted) + + err = cp.PruneAndPersist(0, []*types.CommitQC{qcs[1]}) + require.Error(t, err) + require.Equal(t, asked, recordCount(t, walCommitQCs, stageAsked)) + require.Equal(t, persisted, recordCount(t, walCommitQCs, stagePersisted)) + + require.NoError(t, cp.PruneAndPersist(0, qcs)) + require.Equal(t, asked+2, recordCount(t, walCommitQCs, stageAsked)) + require.Equal(t, persisted+2, recordCount(t, walCommitQCs, stagePersisted)) + + require.NoError(t, cp.PruneAndPersist(0, []*types.CommitQC{qcs[0]})) + require.Equal(t, asked+2, recordCount(t, walCommitQCs, stageAsked)) + require.Equal(t, persisted+2, recordCount(t, walCommitQCs, stagePersisted)) + }) +}