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/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/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..38a6e9d --- /dev/null +++ b/metrics/balance.go @@ -0,0 +1,106 @@ +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 errgroup with the metrics server and returning +// one would cancel the group, taking the /metrics endpoint down over a +// transient RPC failure. +func (p *BalancePoller) poll(parent context.Context) { + ctx, cancel := context.WithTimeout(parent, balancePollTimeout) + defer cancel() + + wei, err := p.client.BalanceAt(ctx, p.address, nil) + if err != nil { + // Cancellations should not count as failed RPC calls. They are + // detected by checking if the parent context is done. + if parent.Err() != nil { + return + } + 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 b1cb12c..1a1f09b 100644 --- a/metrics/metrics.go +++ b/metrics/metrics.go @@ -2,32 +2,33 @@ package metrics import "github.com/prometheus/client_golang/prometheus" -var TotalSuccessfulIdentityRegistration = prometheus.NewGauge( - prometheus.GaugeOpts{ +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.NewGauge( - prometheus.GaugeOpts{ +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.NewGauge( - prometheus.GaugeOpts{ +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) + initBalanceMetrics() } 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() } } }