Skip to content
Open
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
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
22 changes: 11 additions & 11 deletions internal/usecase/crypto.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand All @@ -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",
"",
Expand All @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand All @@ -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",
"",
Expand All @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand All @@ -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)),
Expand Down
16 changes: 8 additions & 8 deletions internal/usecase/eventtrigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
"",
Expand All @@ -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",
"",
Expand All @@ -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",
"",
Expand All @@ -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",
"",
Expand Down Expand Up @@ -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",
"",
Expand Down Expand Up @@ -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)),
Expand Down Expand Up @@ -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",
"",
Expand All @@ -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",
"",
Expand Down
5 changes: 4 additions & 1 deletion main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
100 changes: 100 additions & 0 deletions metrics/balance.go
Original file line number Diff line number Diff line change
@@ -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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This gives the wrong reason for swallowing errors. Swallowing them is correct, but the stated reason isn't: the comment says the poller "shares an error group with the API and a transient RPC failure must not shut the service down". It shares the group at main.go:206 with the metrics server, and that group's error is only logged at :209, never cancelled.

The actual reason is: errgroup.WithContext cancels the group on the first non-nil error, so returning one would take metricsServer down with the poller and we'd lose the metrics entirely over a transient RPC blip. Suggested:

// Errors are never returned, because the poller shares an errgroup with the
// metrics server (main.go:206) and returning one would cancel the group,
// taking the whole /metrics endpoint down over a transient RPC failure.

// 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 {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This counts shutdown as an RPC failure. A cancelled parent context lands in the same error branch as a real failure, so a redeploy that catches a poll in flight logs an error and increments FailedRPCCalls. DeadlineExceeded should stay counted though, a node that doesn't answer in 10s is a genuine failure. So guard on cancellation only:

if err != nil {
      if errors.Is(err, context.Canceled) {
              return
      }
      log.Err(err).Str("address", p.address.Hex()).Msg("failed to query signer balance")
      FailedRPCCalls.Inc()
      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)
}
83 changes: 83 additions & 0 deletions metrics/balance_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading