Skip to content
Draft
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
23 changes: 16 additions & 7 deletions ethfinalizer/ethfinalizer.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,10 @@ import (
// Type parameters:
// - T: transaction metadata type
type FinalizerOptions[T any] struct {
// Wallet is the wallet to be managed by this finalizer, required.
// Wallet is the local wallet to manage. Set exactly one of Wallet or Signer.
Wallet *ethwallet.Wallet
// Signer signs transactions, including replacements, using the operation context.
Signer Signer
// Chain is the provider for the chain where transactions will be sent, required.
// See NewEthkitChain for an implementation using ethkit components.
Chain Chain
Expand Down Expand Up @@ -63,8 +65,8 @@ type FinalizerOptions[T any] struct {
}

func (o FinalizerOptions[T]) IsValid() error {
if o.Wallet == nil {
return fmt.Errorf("no wallet")
if (o.Wallet == nil) == (o.Signer == nil) {
return fmt.Errorf("exactly one of wallet or signer is required")
}

if o.Chain == nil {
Expand Down Expand Up @@ -116,6 +118,7 @@ func (o FinalizerOptions[T]) IsValid() error {
// - T: transaction metadata type
type Finalizer[T any] struct {
FinalizerOptions[T]
signer Signer

isRunning, isStuck atomic.Bool

Expand Down Expand Up @@ -173,7 +176,13 @@ func NewFinalizer[T any](options FinalizerOptions[T]) (*Finalizer[T], error) {
options.Logger = slog.New(slog.DiscardHandler)
}

signer := options.Signer
if signer == nil {
signer = NewWalletSigner(options.Wallet)
}

return &Finalizer[T]{
signer: signer,
FinalizerOptions: options,

subscriptions: map[chan Event[T]]struct{}{},
Expand Down Expand Up @@ -280,7 +289,7 @@ func (f *Finalizer[T]) Run(ctx context.Context) error {

f.Logger.DebugContext(ctx, "polling", slog.Duration("interval", f.PollInterval), slog.Duration("timeout", f.PollTimeout))

chainNonce, err := f.Chain.LatestNonce(ctx, f.Wallet.Address())
chainNonce, err := f.Chain.LatestNonce(ctx, f.signer.Address())
if err != nil {
return fmt.Errorf("unable to read chain nonce: %w", err)
}
Expand Down Expand Up @@ -494,7 +503,7 @@ func (f *Finalizer[T]) Run(ctx context.Context) error {
f.Logger.ErrorContext(ctx, "unable to resend transaction to chain", slog.Any("error", err), slog.String("transaction", transaction.Hash().String()))
}
} else {
if replacement, err = f.Wallet.SignTransaction(replacement, f.Chain.ChainID()); err == nil {
if replacement, err = f.signer.SignTransaction(ctx, replacement, f.Chain.ChainID()); err == nil {
if err := f.Mempool.Commit(ctx, replacement, transaction.Metadata); err != nil {
f.Logger.ErrorContext(ctx, "unable to commit replacement transaction to mempool", slog.Any("error", err))
continue
Expand Down Expand Up @@ -666,7 +675,7 @@ func (f *Finalizer[T]) Send(ctx context.Context, transaction *types.Transaction,
return nil, fmt.Errorf("unable to read mempool nonce: %w", err)
}

chainNonce, err := f.Chain.PendingNonce(ctx, f.Wallet.Address())
chainNonce, err := f.Chain.PendingNonce(ctx, f.signer.Address())
if err != nil {
return nil, fmt.Errorf("unable to read chain nonce: %w", err)
}
Expand All @@ -677,7 +686,7 @@ func (f *Finalizer[T]) Send(ctx context.Context, transaction *types.Transaction,

transaction = withNonce(transaction, max(mempoolNonce, chainNonce))

transaction, err = f.Wallet.SignTransaction(transaction, f.Chain.ChainID())
transaction, err = f.signer.SignTransaction(ctx, transaction, f.Chain.ChainID())
if err != nil {
return nil, fmt.Errorf("unable to sign transaction: %w", err)
}
Expand Down
36 changes: 36 additions & 0 deletions ethfinalizer/signer.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package ethfinalizer

import (
"context"
"math/big"

"github.com/0xsequence/ethkit/ethwallet"
"github.com/0xsequence/ethkit/go-ethereum/common"
"github.com/0xsequence/ethkit/go-ethereum/core/types"
)

// Signer signs transactions without requiring access to private key material.
// Implementations must respect cancellation and return a transaction signed by Address.
type Signer interface {
Address() common.Address
SignTransaction(context.Context, *types.Transaction, *big.Int) (*types.Transaction, error)
}

type walletSigner struct{ wallet *ethwallet.Wallet }

// NewWalletSigner adapts a local wallet to Signer. A nil wallet returns a nil Signer.
func NewWalletSigner(wallet *ethwallet.Wallet) Signer {
if wallet == nil {
return nil
}
return &walletSigner{wallet: wallet}
}

func (s *walletSigner) Address() common.Address { return s.wallet.Address() }

func (s *walletSigner) SignTransaction(ctx context.Context, tx *types.Transaction, chainID *big.Int) (*types.Transaction, error) {
if err := ctx.Err(); err != nil {
return nil, err
}
return s.wallet.SignTransaction(tx, chainID)
}
147 changes: 147 additions & 0 deletions ethfinalizer/signer_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,147 @@
package ethfinalizer

import (
"context"
"errors"
"math/big"
"sync/atomic"
"testing"
"time"

"github.com/0xsequence/ethkit/ethwallet"
"github.com/0xsequence/ethkit/go-ethereum/common"
"github.com/0xsequence/ethkit/go-ethereum/core/types"
"github.com/stretchr/testify/require"
)

type externalTestSigner struct {
Signer
fail atomic.Bool
calls atomic.Int32
}

func (s *externalTestSigner) SignTransaction(ctx context.Context, tx *types.Transaction, chain *big.Int) (*types.Transaction, error) {
s.calls.Add(1)
if s.fail.Load() {
return nil, errors.New("signer unavailable")
}
return s.Signer.SignTransaction(ctx, tx, chain)
}

type externalTestChain struct {
gas atomic.Int64
sent chan *types.Transaction
}

func (c *externalTestChain) ChainID() *big.Int { return big.NewInt(1) }
func (c *externalTestChain) IsEIP1559() bool { return false }
func (c *externalTestChain) LatestNonce(context.Context, common.Address) (uint64, error) {
return 0, nil
}
func (c *externalTestChain) PendingNonce(context.Context, common.Address) (uint64, error) {
return 0, nil
}
func (c *externalTestChain) GasPrice(context.Context) (*big.Int, error) {
return big.NewInt(c.gas.Load()), nil
}
func (c *externalTestChain) BaseFee(context.Context) (*big.Int, error) { return nil, nil }
func (c *externalTestChain) PriorityFee(context.Context) (*big.Int, error) { return nil, nil }
func (c *externalTestChain) Send(ctx context.Context, tx *types.Transaction) error {
select {
case c.sent <- tx:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (c *externalTestChain) Subscribe(ctx context.Context) (<-chan Diff, error) {
ch := make(chan Diff)
go func() { <-ctx.Done(); close(ch) }()
return ch, nil
}

func TestNewWalletSignerNil(t *testing.T) {
signer := NewWalletSigner(nil)
require.True(t, signer == nil, "nil wallet must produce a nil interface")

finalizer, err := NewFinalizer(FinalizerOptions[string]{
Signer: signer,
Chain: &externalTestChain{},
Mempool: NewMemoryMempool[string](),
PollInterval: time.Second,
PollTimeout: time.Second,
RetryDelay: time.Second,
})
require.ErrorContains(t, err, "exactly one of wallet or signer is required")
require.Nil(t, finalizer)
}

func TestExternalSignerLifecycle(t *testing.T) {
wallet, err := ethwallet.NewWalletFromRandomEntropy()
require.NoError(t, err)
signer := &externalTestSigner{Signer: NewWalletSigner(wallet)}
chain := &externalTestChain{sent: make(chan *types.Transaction, 100)}
chain.gas.Store(10)
mempool := NewMemoryMempool[string]()
options := FinalizerOptions[string]{Signer: signer, Chain: chain, Mempool: mempool, PollInterval: 5 * time.Millisecond, PollTimeout: time.Second, RetryDelay: time.Millisecond, PriceBump: 15}
f, err := NewFinalizer(options)
require.NoError(t, err)
options.Wallet = wallet
_, err = NewFinalizer(options)
require.Error(t, err)
options.Signer = nil
_, err = NewFinalizer(options)
require.NoError(t, err)
options.Wallet = nil
_, err = NewFinalizer(options)
require.Error(t, err)
tx := types.NewTx(&types.LegacyTx{Gas: 21000, GasPrice: big.NewInt(10)})
signer.fail.Store(true)
_, err = f.Send(context.Background(), tx, "metadata")
require.Error(t, err)
nonce, err := mempool.Nonce(context.Background())
require.NoError(t, err)
require.Zero(t, nonce)
require.Empty(t, chain.sent)
signer.fail.Store(false)
ctx, cancel := context.WithCancel(context.Background())
cancel()
_, err = f.Send(ctx, tx, "metadata")
require.ErrorIs(t, err, context.Canceled)
signed, err := f.Send(context.Background(), tx, "metadata")
require.NoError(t, err)
require.Equal(t, signed.Hash(), (<-chain.sent).Hash())
calls := signer.calls.Load()
ctx, cancel = context.WithCancel(context.Background())
defer cancel()
done := make(chan error, 1)
go func() { done <- f.Run(ctx) }()
select {
case resent := <-chain.sent:
require.Equal(t, signed.Hash(), resent.Hash())
case <-time.After(time.Second):
t.Fatal("no rebroadcast")
}
require.Equal(t, calls, signer.calls.Load(), "rebroadcast must reuse persisted signature")
signer.fail.Store(true)
chain.gas.Store(20)
require.Eventually(t, func() bool { return signer.calls.Load() > calls }, time.Second, time.Millisecond)
_, latest, err := mempool.Status(context.Background(), 0)
require.NoError(t, err)
require.Equal(t, signed.Hash(), latest.Transaction.Hash())
signer.fail.Store(false)
require.Eventually(t, func() bool {
_, latest, err := mempool.Status(context.Background(), 0)
return err == nil && latest.Transaction.GasPrice().Int64() == 20
}, time.Second, time.Millisecond)
_, latest, err = mempool.Status(context.Background(), 0)
require.NoError(t, err)
require.Equal(t, "metadata", latest.Transaction.Metadata)
require.Zero(t, latest.Transaction.Nonce())
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("finalizer did not stop")
}
}
Loading