From 42e83452f363d6bedda5f6993069ee85f0a2c813 Mon Sep 17 00:00:00 2001 From: Jannik Luhn Date: Wed, 12 Aug 2026 19:04:25 +0200 Subject: [PATCH 1/3] metrics: use counters instead of gauges for monotonic totals TotalSuccessfulIdentityRegistration, TotalDecryptionKeysReceived and TotalFailedRPCCalls are only ever incremented, so Counter is the correct type. Series names are unchanged, so existing queries keep working, and rate() over them is now legitimate rather than an accident. --- metrics/metrics.go | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/metrics/metrics.go b/metrics/metrics.go index b1cb12c..5266130 100644 --- a/metrics/metrics.go +++ b/metrics/metrics.go @@ -2,24 +2,24 @@ package metrics import "github.com/prometheus/client_golang/prometheus" -var TotalSuccessfulIdentityRegistration = prometheus.NewGauge( - prometheus.GaugeOpts{ +var TotalSuccessfulIdentityRegistration = prometheus.NewCounter( + prometheus.CounterOpts{ Namespace: "shutter_api", Name: "total_successful_identities_registration", Help: "counter of successful identity registration", }, ) -var TotalDecryptionKeysReceived = prometheus.NewGauge( - prometheus.GaugeOpts{ +var TotalDecryptionKeysReceived = prometheus.NewCounter( + prometheus.CounterOpts{ Namespace: "shutter_api", Name: "total_decryption_keys_received", Help: "counter of total dec keys received", }, ) -var TotalFailedRPCCalls = prometheus.NewGauge( - prometheus.GaugeOpts{ +var TotalFailedRPCCalls = prometheus.NewCounter( + prometheus.CounterOpts{ Namespace: "shutter_api", Name: "total_failed_rpc_calls", Help: "Counter of failed rpc calls", From adc5105130b3606c2dcc4700878a792d2ef58414 Mon Sep 17 00:00:00 2001 From: Jannik Luhn Date: Fri, 14 Aug 2026 21:31:39 +0200 Subject: [PATCH 2/3] metrics: rename counters to the _total suffix convention Prometheus counters carry _total as a suffix, not a prefix, and the plural belongs on the thing being counted: shutter_api_total_successful_identities_registration -> shutter_api_successful_identity_registrations_total shutter_api_total_decryption_keys_received -> shutter_api_decryption_keys_received_total shutter_api_total_failed_rpc_calls -> shutter_api_failed_rpc_calls_total Go identifiers drop their now-redundant Total prefix to match. Existing series keep their history under the old names but stop being written to, so any dashboard or alert rule querying them needs updating. --- internal/usecase/crypto.go | 22 +++++++++++----------- internal/usecase/eventtrigger.go | 16 ++++++++-------- metrics/metrics.go | 24 ++++++++++++------------ watcher/watcher.go | 2 +- 4 files changed, 32 insertions(+), 32 deletions(-) diff --git a/internal/usecase/crypto.go b/internal/usecase/crypto.go index 215697f..883aa64 100644 --- a/internal/usecase/crypto.go +++ b/internal/usecase/crypto.go @@ -120,7 +120,7 @@ func (uc *CryptoUsecase) getSigner(ctx context.Context) (*bind.TransactOpts, *ht chainID, err := uc.ethClient.ChainID(ctx) if err != nil { log.Err(err).Msg("err encountered while querying chain id") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying chain id", "", @@ -168,7 +168,7 @@ func (uc *CryptoUsecase) GetDecryptionKey(ctx context.Context, identity string) registrationData, err := uc.shutterRegistryContract.Registrations(nil, [32]byte(identityBytes)) if err != nil { log.Err(err).Msg("err encountered while querying contract") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error while querying for identity from the contract", "", @@ -301,7 +301,7 @@ func (uc *CryptoUsecase) GetDataForEncryption(ctx context.Context, address strin blockNumber, err := uc.ethClient.BlockNumber(ctx) if err != nil { log.Err(err).Msg("err encountered while querying for recent block") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for recent block", "", @@ -313,7 +313,7 @@ func (uc *CryptoUsecase) GetDataForEncryption(ctx context.Context, address strin eon, err := uc.keyperSetManagerContract.GetKeyperSetIndexByBlock(nil, blockNumber) if err != nil { log.Err(err).Msg("err encountered while querying keyper set index") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for keyper set index", "", @@ -325,7 +325,7 @@ func (uc *CryptoUsecase) GetDataForEncryption(ctx context.Context, address strin eonKeyBytes, err := uc.keyBroadcastContract.GetEonKey(nil, eon) if err != nil { log.Err(err).Msg("err encountered while querying for eon key") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for eon key", "", @@ -442,7 +442,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta blockNumber, err := uc.ethClient.BlockNumber(ctx) if err != nil { log.Err(err).Msg("err encountered while querying for recent block") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for recent block", "", @@ -454,7 +454,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta eon, err := uc.keyperSetManagerContract.GetKeyperSetIndexByBlock(nil, blockNumber) if err != nil { log.Err(err).Msg("err encountered while querying keyper set index") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for keyper set index", "", @@ -466,7 +466,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta eonKeyBytes, err := uc.keyBroadcastContract.GetEonKey(nil, eon) if err != nil { log.Err(err).Msg("err encountered while querying for eon key") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for eon key", "", @@ -496,7 +496,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta registrationData, err := uc.shutterRegistryContract.Registrations(nil, [32]byte(identity)) if err != nil { log.Err(err).Msg("err encountered while querying contract") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error while querying for registrations from the contract", "", @@ -525,7 +525,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta tx, err := uc.shutterRegistryContract.Register(&opts, eon, identityPrefix, decryptionTimestamp) if err != nil { log.Err(err).Msg("failed to send transaction") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "failed to register identity", "", @@ -537,7 +537,7 @@ func (uc *CryptoUsecase) RegisterIdentity(ctx context.Context, decryptionTimesta // we return the transaction hash in response to allow // users the ability to monitor it themselves - metrics.TotalSuccessfulIdentityRegistration.Inc() + metrics.SuccessfulIdentityRegistrations.Inc() return &RegisterIdentityResponse{ Eon: eon, Identity: common.PrefixWith0x(hex.EncodeToString(identity)), diff --git a/internal/usecase/eventtrigger.go b/internal/usecase/eventtrigger.go index 4d29e7f..f3a08e5 100644 --- a/internal/usecase/eventtrigger.go +++ b/internal/usecase/eventtrigger.go @@ -288,7 +288,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger blockNumber, err := uc.ethClient.BlockNumber(ctx) if err != nil { log.Err(err).Msg("err encountered while querying for recent block") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for recent block", "", @@ -300,7 +300,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger eon, err := uc.keyperSetManagerContract.GetKeyperSetIndexByBlock(nil, blockNumber) if err != nil { log.Err(err).Msg("err encountered while querying keyper set index") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for keyper set index", "", @@ -312,7 +312,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger eonKeyBytes, err := uc.keyBroadcastContract.GetEonKey(nil, eon) if err != nil { log.Err(err).Msg("err encountered while querying for eon key") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for eon key", "", @@ -335,7 +335,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger chainId, err := uc.ethClient.ChainID(ctx) if err != nil { log.Err(err).Msg("err encountered while quering chain id") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying chain id", "", @@ -411,7 +411,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger tx, err := uc.shutterEventRegistryContract.Register(&opts, eon, identityPrefix, eventTriggerDefinition, ttl) if err != nil { log.Err(err).Msg("failed to send transaction") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "failed to register identity", "", @@ -441,7 +441,7 @@ func (uc *CryptoUsecase) RegisterEventIdentity(ctx context.Context, eventTrigger go uc.updateEventIdentityExpirationBlockNumber(tx.Hash(), eon, identity, ttl) - metrics.TotalSuccessfulIdentityRegistration.Inc() + metrics.SuccessfulIdentityRegistrations.Inc() return &RegisterIdentityResponse{ Eon: eon, Identity: common.PrefixWith0x(hex.EncodeToString(identity)), @@ -560,7 +560,7 @@ func (uc *CryptoUsecase) GetEventDecryptionKey(ctx context.Context, identity str blockNumber, err := uc.ethClient.BlockNumber(ctx) if err != nil { log.Err(err).Msg("err encountered while querying for recent block") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying for recent block", "", @@ -572,7 +572,7 @@ func (uc *CryptoUsecase) GetEventDecryptionKey(ctx context.Context, identity str eonUint, err := uc.keyperSetManagerContract.GetKeyperSetIndexByBlock(nil, blockNumber) if err != nil { log.Err(err).Msg("err encountered while querying current eon") - metrics.TotalFailedRPCCalls.Inc() + metrics.FailedRPCCalls.Inc() err := httpError.NewHttpError( "error encountered while querying current eon", "", diff --git a/metrics/metrics.go b/metrics/metrics.go index 5266130..5d4488c 100644 --- a/metrics/metrics.go +++ b/metrics/metrics.go @@ -2,32 +2,32 @@ package metrics import "github.com/prometheus/client_golang/prometheus" -var TotalSuccessfulIdentityRegistration = prometheus.NewCounter( +var SuccessfulIdentityRegistrations = prometheus.NewCounter( prometheus.CounterOpts{ Namespace: "shutter_api", - Name: "total_successful_identities_registration", - Help: "counter of successful identity registration", + Name: "successful_identity_registrations_total", + Help: "Count of successful identity registrations.", }, ) -var TotalDecryptionKeysReceived = prometheus.NewCounter( +var DecryptionKeysReceived = prometheus.NewCounter( prometheus.CounterOpts{ Namespace: "shutter_api", - Name: "total_decryption_keys_received", - Help: "counter of total dec keys received", + Name: "decryption_keys_received_total", + Help: "Count of decryption keys received from the keypers.", }, ) -var TotalFailedRPCCalls = prometheus.NewCounter( +var FailedRPCCalls = prometheus.NewCounter( prometheus.CounterOpts{ Namespace: "shutter_api", - Name: "total_failed_rpc_calls", - Help: "Counter of failed rpc calls", + Name: "failed_rpc_calls_total", + Help: "Count of failed RPC calls.", }, ) func InitMetrics() { - prometheus.MustRegister(TotalSuccessfulIdentityRegistration) - prometheus.MustRegister(TotalDecryptionKeysReceived) - prometheus.MustRegister(TotalFailedRPCCalls) + prometheus.MustRegister(SuccessfulIdentityRegistrations) + prometheus.MustRegister(DecryptionKeysReceived) + prometheus.MustRegister(FailedRPCCalls) } diff --git a/watcher/watcher.go b/watcher/watcher.go index d5603dc..a282fd9 100644 --- a/watcher/watcher.go +++ b/watcher/watcher.go @@ -46,7 +46,7 @@ func (w *Watcher) Start(ctx context.Context, runner service.Runner) error { }); err != nil { log.Err(err).Msg("failed to insert decryption key") } - metrics.TotalDecryptionKeysReceived.Inc() + metrics.DecryptionKeysReceived.Inc() } } } From c9096d5f3a01448788d8c07478e7ead6d7d651e8 Mon Sep 17 00:00:00 2001 From: Jannik Luhn Date: Sun, 16 Aug 2026 16:55:01 +0200 Subject: [PATCH 3/3] metrics: add a gauge for the signer's balance The signer pays gas for every identity registration, but the service never queried its balance, so an account draining to empty was only visible once registrations started failing. At that point it surfaced as failed_rpc_calls_total from the transaction send sites, indistinguishable from an RPC outage. shutter_api_signer_balance_ether is published by a new BalancePoller running as a service alongside the metrics server, so it is gated on METRICS_ENABLED. It reads once at startup and then every 60s, each read bounded by a 10s timeout so a hung endpoint cannot stall the loop. A failed read logs, increments failed_rpc_calls_total and leaves the gauge alone; it never returns an error, because a transient RPC failure must not bring the service down through the error group. The metric is a GaugeVec with no labels rather than a plain Gauge. A plain Gauge is registered holding 0, so a restart while the RPC endpoint was down would publish 0 ether and fire the low-balance alert this metric exists to raise. With no value set the series is simply absent, and the exposed series is otherwise identical. Balance is reported in ether rather than wei so that alert thresholds are readable; the conversion goes through big.Float, as an integer quotient would truncate everything below 1 ether to zero. Known gap: a reading that stops being refreshed keeps its last value indefinitely, so a dead poller looks healthy. Detection relies on failed_rpc_calls_total, which now also moves for background polls and can therefore rise with no traffic. Co-Authored-By: Claude --- go.mod | 1 + main.go | 5 +- metrics/balance.go | 100 ++++++++++++++++++++++++++++++++++++++++ metrics/balance_test.go | 83 +++++++++++++++++++++++++++++++++ metrics/metrics.go | 1 + 5 files changed, 189 insertions(+), 1 deletion(-) create mode 100644 metrics/balance.go create mode 100644 metrics/balance_test.go diff --git a/go.mod b/go.mod index 4351d9b..f77acd4 100644 --- a/go.mod +++ b/go.mod @@ -102,6 +102,7 @@ require ( github.com/koron/go-ssdp v0.0.5 // indirect github.com/kr/pretty v0.3.1 // indirect github.com/kr/text v0.2.0 // indirect + github.com/kylelemons/godebug v1.1.0 // indirect github.com/libp2p/go-buffer-pool v0.1.0 // indirect github.com/libp2p/go-cidranger v1.1.0 // indirect github.com/libp2p/go-flow-metrics v0.2.0 // indirect diff --git a/main.go b/main.go index ed850af..0d26815 100644 --- a/main.go +++ b/main.go @@ -200,7 +200,10 @@ func main() { defer deferFn() if metricsConfig.Enabled { - group, deferFn := service.RunBackground(ctx, metricsServer) + signerAddress := crypto.PubkeyToAddress(*config.PublicKey) + balancePoller := metrics.NewBalancePoller(client, signerAddress) + + group, deferFn := service.RunBackground(ctx, metricsServer, balancePoller) defer deferFn() go func() { if err := group.Wait(); err != nil { diff --git a/metrics/balance.go b/metrics/balance.go new file mode 100644 index 0000000..cf79c74 --- /dev/null +++ b/metrics/balance.go @@ -0,0 +1,100 @@ +package metrics + +import ( + "context" + "math/big" + "time" + + ecommon "github.com/ethereum/go-ethereum/common" + "github.com/prometheus/client_golang/prometheus" + "github.com/rs/zerolog/log" + "github.com/shutter-network/rolling-shutter/rolling-shutter/medley/service" +) + +const ( + balancePollInterval = 60 * time.Second + balancePollTimeout = 10 * time.Second +) + +// weiPerEther is the divisor turning a wei balance into ether. +var weiPerEther = new(big.Float).SetFloat64(1e18) + +var SignerBalanceEther = newSignerBalanceGauge() + +// A GaugeVec with no labels, so the series stays absent until a balance has +// been read. A plain Gauge would be registered holding 0, and a restart while +// the RPC endpoint is down would publish 0 ether and fire the low-balance +// alert. The exposed series is the same either way. +func newSignerBalanceGauge() *prometheus.GaugeVec { + return prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: "shutter_api", + Name: "signer_balance_ether", + Help: "Balance of the signer account in ether.", + }, + []string{}, + ) +} + +func initBalanceMetrics() { + prometheus.MustRegister(SignerBalanceEther) +} + +// BalanceReader reads an account balance from the chain. It is the single +// method of ethclient.Client that the poller needs. +type BalanceReader interface { + BalanceAt(ctx context.Context, account ecommon.Address, blockNumber *big.Int) (*big.Int, error) +} + +// BalancePoller publishes the signer's balance as a gauge. The signer pays gas +// for every identity registration, so an empty account takes registration down; +// without this the service spends from an account it cannot observe. +type BalancePoller struct { + client BalanceReader + address ecommon.Address +} + +func NewBalancePoller(client BalanceReader, address ecommon.Address) *BalancePoller { + return &BalancePoller{client: client, address: address} +} + +func (p *BalancePoller) Start(ctx context.Context, runner service.Runner) error { + runner.Go(func() error { + ticker := time.NewTicker(balancePollInterval) + defer ticker.Stop() + + // Publish once up front so the series exists from the first scrape + // rather than only after a full interval. + p.poll(ctx) + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + p.poll(ctx) + } + } + }) + return nil +} + +// poll reads the balance and updates the gauge. A failed read leaves the gauge +// at its previous value: publishing a zero would look like an empty account and +// fire the very alert this metric exists to raise. Errors are never returned, +// because the poller shares an error group with the API and a transient RPC +// failure must not shut the service down. +func (p *BalancePoller) poll(ctx context.Context) { + ctx, cancel := context.WithTimeout(ctx, balancePollTimeout) + defer cancel() + + wei, err := p.client.BalanceAt(ctx, p.address, nil) + if err != nil { + log.Err(err).Str("address", p.address.Hex()).Msg("failed to query signer balance") + FailedRPCCalls.Inc() + return + } + + ether, _ := new(big.Float).Quo(new(big.Float).SetInt(wei), weiPerEther).Float64() + SignerBalanceEther.WithLabelValues().Set(ether) +} diff --git a/metrics/balance_test.go b/metrics/balance_test.go new file mode 100644 index 0000000..601a4f4 --- /dev/null +++ b/metrics/balance_test.go @@ -0,0 +1,83 @@ +package metrics + +import ( + "context" + "errors" + "math/big" + "testing" + "time" + + ecommon "github.com/ethereum/go-ethereum/common" + "github.com/prometheus/client_golang/prometheus/testutil" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" +) + +// mockBalanceReader is local to this test: the poller needs a single method, +// while tests/mock.MockEthClient serves the usecase's wider interface. +type mockBalanceReader struct { + mock.Mock +} + +func (m *mockBalanceReader) BalanceAt(ctx context.Context, account ecommon.Address, blockNumber *big.Int) (*big.Int, error) { + args := m.Called(ctx, account, blockNumber) + balance, _ := args.Get(0).(*big.Int) + return balance, args.Error(1) +} + +// Deliberately not a real signer address. +var testSignerAddress = ecommon.HexToAddress("0x1111111111111111111111111111111111111111") + +// A registered-but-never-read balance must expose nothing at all. Publishing a +// zero would be read as an empty signer account. +func TestBalanceIsAbsentBeforeTheFirstRead(t *testing.T) { + require.Zero(t, testutil.CollectAndCount(newSignerBalanceGauge())) +} + +func TestPollPublishesBalanceInEther(t *testing.T) { + // 1.5 ether, to catch a big.Int quotient truncating the fraction away. + wei := new(big.Int).Add( + new(big.Int).SetUint64(1e18), + new(big.Int).SetUint64(5e17), + ) + + client := &mockBalanceReader{} + client.On("BalanceAt", mock.Anything, testSignerAddress, (*big.Int)(nil)).Return(wei, nil) + + NewBalancePoller(client, testSignerAddress).poll(context.Background()) + + require.InDelta(t, 1.5, testutil.ToFloat64(SignerBalanceEther.WithLabelValues()), 1e-9) + client.AssertExpectations(t) +} + +func TestPollLeavesGaugeUntouchedOnError(t *testing.T) { + SignerBalanceEther.WithLabelValues().Set(2) + failuresBefore := testutil.ToFloat64(FailedRPCCalls) + + client := &mockBalanceReader{} + client.On("BalanceAt", mock.Anything, testSignerAddress, (*big.Int)(nil)). + Return(nil, errors.New("rpc unavailable")) + + NewBalancePoller(client, testSignerAddress).poll(context.Background()) + + // A zero here would read as an empty account and fire a false alert. + require.InDelta(t, 2.0, testutil.ToFloat64(SignerBalanceEther.WithLabelValues()), 1e-9) + require.Equal(t, failuresBefore+1, testutil.ToFloat64(FailedRPCCalls)) +} + +// The read must carry its own deadline, so a hung RPC cannot stall the loop +// and stop every later poll. +func TestPollBoundsTheReadWithATimeout(t *testing.T) { + var got context.Context + + client := &mockBalanceReader{} + client.On("BalanceAt", mock.Anything, testSignerAddress, (*big.Int)(nil)). + Run(func(args mock.Arguments) { got = args.Get(0).(context.Context) }). + Return(big.NewInt(0), nil) + + NewBalancePoller(client, testSignerAddress).poll(context.Background()) + + deadline, ok := got.Deadline() + require.True(t, ok, "read was made without a deadline") + require.LessOrEqual(t, time.Until(deadline), balancePollTimeout) +} diff --git a/metrics/metrics.go b/metrics/metrics.go index 5d4488c..1a1f09b 100644 --- a/metrics/metrics.go +++ b/metrics/metrics.go @@ -30,4 +30,5 @@ func InitMetrics() { prometheus.MustRegister(SuccessfulIdentityRegistrations) prometheus.MustRegister(DecryptionKeysReceived) prometheus.MustRegister(FailedRPCCalls) + initBalanceMetrics() }