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
59 changes: 12 additions & 47 deletions zetaclient/maintenance/tss_listener.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,9 @@ import (
observertypes "github.com/zeta-chain/node/x/observer/types"
)

const tssListenerTicker = 5 * time.Second
// tssListenerTicker is how often the watchers re-query zetacore. A var rather than a const so
// tests can shrink it; nothing in production writes to it.
var tssListenerTicker = 5 * time.Second

// TSSListener is a struct that listens for TSS updates, new keygen, and new TSS key generation.
type TSSListener struct {
Expand All @@ -32,6 +34,15 @@ func NewTSSListener(client ZetacoreClient, logger zerolog.Logger) *TSSListener {
}

// Listen listens for any maintenance regarding TSS and calls action specified. Works in the background.
//
// Both watchers key off the TSS itself: the address changing, or a new key landing in history.
// The keygen record is deliberately not watched. It is reset to "pending at block MaxInt64" on
// any observer set change, and restarting on that reset is what turns a routine validator
// unbonding into an outage — every signer shuts down at once and none of them can start again.
//
// Note this also removes the only trigger that restarted zetaclient for a scheduled keygen, so
// while a finalized key exists a rotation ceremony will not start. That is deliberate; see the
// step 5 comment in zetaclient/tss/setup.go.
func (tl *TSSListener) Listen(ctx context.Context, action func()) {
var (
withLogger = bg.WithLogger(tl.logger)
Expand All @@ -40,7 +51,6 @@ func (tl *TSSListener) Listen(ctx context.Context, action func()) {

bg.Work(ctx, tl.waitForUpdate, bg.WithName("tss.wait_for_update"), withLogger, onComplete)
bg.Work(ctx, tl.waitForNewKeyGeneration, bg.WithName("tss.wait_for_generation"), withLogger, onComplete)
bg.Work(ctx, tl.waitForNewKeygen, bg.WithName("tss.wait_for_keygen"), withLogger, onComplete)
}

// waitForUpdate listens for TSS updates. Returns `nil` when the TSS address is updated
Expand Down Expand Up @@ -126,48 +136,3 @@ func (tl *TSSListener) waitForNewKeyGeneration(ctx context.Context) error {
}
}
}

// waitForNewKeygen is a background thread that listens for new keygen; it returns when a new keygen is set
func (tl *TSSListener) waitForNewKeygen(ctx context.Context) error {
// Initial Keygen retrieval
keygen, err := tl.client.GetKeyGen(ctx)
if err != nil {
return errors.Wrap(err, "failed to get initial TSS history")
}

ticker := time.NewTicker(tssListenerTicker)
defer ticker.Stop()

for {
select {
case <-ticker.C:
keygenUpdated, err := tl.client.GetKeyGen(ctx)
switch {
case err != nil:
tl.logger.Warn().Err(err).Msg("unable to get keygen")
continue
// Keygen is not pending it has already been successfully generated, continue loop
case keygenUpdated.Status == observertypes.KeygenStatus_KeyGenSuccess:
continue
// Keygen failed we to need to wait until a new keygen is set, continue loop
case keygenUpdated.Status == observertypes.KeygenStatus_KeyGenFailed:
continue
// Keygen is pending but block number is not updated, continue loop.
// Most likely the zetaclient is waiting for the keygen block to arrive.
case keygenUpdated.Status == observertypes.KeygenStatus_PendingKeygen &&
keygenUpdated.BlockNumber <= keygen.BlockNumber:
continue
}

// Trigger restart only when the following conditions are met:
// 1. Keygen is pending
// 2. Block number is updated

tl.logger.Info().Int64("block_number", keygenUpdated.BlockNumber).Msg("got new keygen")
return nil
case <-ctx.Done():
tl.logger.Info().Msg("stopped waiting for new keygen in the TSS listener")
return nil
}
}
}
124 changes: 124 additions & 0 deletions zetaclient/maintenance/tss_listener_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
package maintenance

import (
"context"
"io"
"testing"
"time"

"github.com/rs/zerolog"

observertypes "github.com/zeta-chain/node/x/observer/types"
"github.com/zeta-chain/node/zetaclient/testutils/mocks"
)

const (
tssPubkeyOld = "zetapub1addwnpepqtadxdyt037h86z60nl98t6zk56mw5zpnm79tsmvspln3hgt5phdc79kvfc"
tssPubkeyNew = "zetapub1addwnpepqglunjrgl3qg08duxq9pf28jmvrer3crwnnfzp6m0u0yh9jk9mnn5p76utc"
)

// useFastTicker shrinks the watcher poll interval for the duration of one test. The assertions
// are about which events cause a shutdown, not about the real 5s cadence, so waiting on it only
// bought ~31s of sleep across this file.
func useFastTicker(t *testing.T) {
original := tssListenerTicker
tssListenerTicker = 10 * time.Millisecond
t.Cleanup(func() { tssListenerTicker = original })
}

// waitForListenerTick gives the listener room for at least one tick so a shutdown it was
// going to signal has actually had the chance to fire.
func waitForListenerTick() {
time.Sleep(20 * tssListenerTicker)
}

func TestTSSListener(t *testing.T) {
// Deliberately not zerolog.NewTestWriter(t): Listen's workers keep running after the test
// returns, and logging into a finished t panics with "Log in goroutine after test completed".
logger := zerolog.New(io.Discard)

oldTSS := observertypes.TSS{TssPubkey: tssPubkeyOld}

// A blanked keygen record must not restart the client.
//
// zetacore writes exactly this on any observer set change (x/observer/abci.go
// BeginBlocker): status back to pending, grantees erased, block set to MaxInt64. On
// mainnet that reset restarted every signer at once, and none could start again because
// the erased grantee list left them with an empty p2p whitelist.
t.Run("blanked keygen record does not trigger a shutdown", func(t *testing.T) {
useFastTicker(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

client := mocks.NewZetacoreClient(t)
client.Mock.On("GetTSS", ctx).Return(oldTSS, nil)
client.Mock.On("GetTSSHistory", ctx).Return([]observertypes.TSS{oldTSS}, nil)

complete := make(chan interface{})
NewTSSListener(client, logger).Listen(ctx, func() { close(complete) })

waitForListenerTick()
assertChannelNotClosed(t, complete)
})

t.Run("TSS address change still triggers a shutdown", func(t *testing.T) {
useFastTicker(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

client := mocks.NewZetacoreClient(t)
client.Mock.On("GetTSSHistory", ctx).Return([]observertypes.TSS{oldTSS}, nil)
client.Mock.On("GetTSS", ctx).Return(oldTSS, nil).Once()
client.Mock.On("GetTSS", ctx).Return(observertypes.TSS{TssPubkey: tssPubkeyNew}, nil)

complete := make(chan interface{})
NewTSSListener(client, logger).Listen(ctx, func() { close(complete) })

<-complete
})

t.Run("new key in TSS history still triggers a shutdown", func(t *testing.T) {
useFastTicker(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

newTSS := observertypes.TSS{TssPubkey: tssPubkeyNew}

client := mocks.NewZetacoreClient(t)
client.Mock.On("GetTSS", ctx).Return(oldTSS, nil)
client.Mock.On("GetTSSHistory", ctx).Return([]observertypes.TSS{oldTSS}, nil).Once()
client.Mock.On("GetTSSHistory", ctx).Return([]observertypes.TSS{oldTSS, newTSS}, nil)

complete := make(chan interface{})
NewTSSListener(client, logger).Listen(ctx, func() { close(complete) })

<-complete
})
}

// TestTSSListenerIgnoresKeygen pins that the listener never reads the keygen record. The mock
// fails the test on any unexpected call, so a reintroduced keygen watcher shows up here rather
// than as an outage.
func TestTSSListenerIgnoresKeygen(t *testing.T) {
useFastTicker(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

// Deliberately not zerolog.NewTestWriter(t): Listen's workers keep running after the test
// returns, and logging into a finished t panics with "Log in goroutine after test completed".
logger := zerolog.New(io.Discard)
oldTSS := observertypes.TSS{TssPubkey: tssPubkeyOld}

client := mocks.NewZetacoreClient(t)
client.Mock.On("GetTSS", ctx).Return(oldTSS, nil)
client.Mock.On("GetTSSHistory", ctx).Return([]observertypes.TSS{oldTSS}, nil)

NewTSSListener(client, logger).Listen(ctx, func() {})
waitForListenerTick()

client.Mock.AssertNotCalled(t, "GetKeyGen", ctx)
}
8 changes: 7 additions & 1 deletion zetaclient/tss/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"github.com/zeta-chain/go-tss/keysign"

"github.com/zeta-chain/node/pkg/chains"
"github.com/zeta-chain/node/pkg/retry"
observertypes "github.com/zeta-chain/node/x/observer/types"
keyinterfaces "github.com/zeta-chain/node/zetaclient/keys/interfaces"
"github.com/zeta-chain/node/zetaclient/logs"
Expand Down Expand Up @@ -106,7 +107,12 @@ func WithPostBlame(postBlame bool) Opt {
// Otherwise, no metrics will be collected.
func WithMetrics(ctx context.Context, zetacore Zetacore, m *Metrics) Opt {
return func(cfg *serviceConfig, _ zerolog.Logger) error {
keygen, err := zetacore.GetKeyGen(ctx)
// Retried like the other startup queries: this one is easy to miss because it sits in
// an option rather than in Setup, but it runs on the same path and is just as fatal.
keygen, err := retry.DoTypedWithBackoffAndRetry(
func() (observertypes.Keygen, error) { return zetacore.GetKeyGen(ctx) },
retry.DefaultConstantBackoff(),
)
if err != nil {
return errors.Wrap(err, "failed to get keygen (WithMetrics)")
}
Expand Down
Loading
Loading