From e7a48a5c08ad2db6bfa171e46407670af0e8430f Mon Sep 17 00:00:00 2001 From: aodhgan Date: Wed, 1 Oct 2025 22:41:01 -0700 Subject: [PATCH] move to event driven --- internal/ethapi/api.go | 56 ++++++++++++++++++++++++++----------- internal/ethapi/api_test.go | 5 ++-- 2 files changed, 43 insertions(+), 18 deletions(-) diff --git a/internal/ethapi/api.go b/internal/ethapi/api.go index 8b28bb7092..e35ffcceac 100644 --- a/internal/ethapi/api.go +++ b/internal/ethapi/api.go @@ -41,6 +41,7 @@ import ( "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/eth/gasestimator" "github.com/ethereum/go-ethereum/eth/tracers/logger" + "github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/internal/ethapi/override" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/p2p" @@ -1653,7 +1654,7 @@ func (api *TransactionAPI) SendRawTransaction(ctx context.Context, input hexutil } // SendRawTransactionSync will add the signed transaction to the transaction pool -// and wait for the transaction to be mined until the timeout (in milliseconds) is reached. +// and wait until the transaction has been included in a block and return the receipt, or the timeout. func (api *TransactionAPI) SendRawTransactionSync(ctx context.Context, input hexutil.Bytes, timeoutMs *hexutil.Uint64) (map[string]interface{}, error) { tx := new(types.Transaction) if err := tx.UnmarshalBinary(input); err != nil { @@ -1664,7 +1665,7 @@ func (api *TransactionAPI) SendRawTransactionSync(ctx context.Context, input hex return nil, err } - // compute effective timeout + // effective timeout: min(requested, max), default if none max := api.b.RPCTxSyncMaxTimeout() def := api.b.RPCTxSyncDefaultTimeout() @@ -1674,40 +1675,63 @@ func (api *TransactionAPI) SendRawTransactionSync(ctx context.Context, input hex if req > max { eff = max } else { - eff = req // allow shorter than default + eff = req } } - // Wait for receipt until timeout receiptCtx, cancel := context.WithTimeout(ctx, eff) defer cancel() - ticker := time.NewTicker(time.Millisecond * 100) - defer ticker.Stop() + // Fast path: maybe already mined/indexed + if r, err := api.GetTransactionReceipt(receiptCtx, hash); err == nil && r != nil { + return r, nil + } + + // Subscribe to head changes and re-check on each new block + heads := make(chan core.ChainHeadEvent, 16) + var sub event.Subscription = api.b.SubscribeChainHeadEvent(heads) + if sub == nil { + return nil, errors.New("chain head subscription unavailable") + } + defer sub.Unsubscribe() for { select { - case <-ticker.C: - receipt, getErr := api.GetTransactionReceipt(receiptCtx, hash) - // If tx-index still building, keep polling - if getErr != nil && !errors.Is(getErr, NewTxIndexingError()) { - // transient or other error: just keep polling - } - if receipt != nil { - return receipt, nil - } case <-receiptCtx.Done(): + // Distinguish upstream cancellation from our timeout if err := ctx.Err(); err != nil { - return nil, err // context canceled / deadline exceeded upstream + return nil, err } return nil, &txSyncTimeoutError{ msg: fmt.Sprintf("The transaction was added to the transaction pool but wasn't processed in %v.", eff), hash: hash, } + + case err := <-sub.Err(): + return nil, err + + case <-heads: + receipt, getErr := api.GetTransactionReceipt(receiptCtx, hash) + // If tx-index still building, keep waiting; ignore transient errors + if getErr != nil && !errors.Is(getErr, NewTxIndexingError()) { + // ignore and wait for next head + continue + } + if receipt != nil { + return receipt, nil + } } } } +func (api *TransactionAPI) tryGetReceipt(ctx context.Context, hash common.Hash) (map[string]interface{}, error) { + receipt, err := api.GetTransactionReceipt(ctx, hash) + if err != nil { + return nil, err + } + return receipt, nil +} + // Sign calculates an ECDSA signature for: // keccak256("\x19Ethereum Signed Message:\n" + len(message) + message). // diff --git a/internal/ethapi/api_test.go b/internal/ethapi/api_test.go index 04ac249e4b..5ee3453a1a 100644 --- a/internal/ethapi/api_test.go +++ b/internal/ethapi/api_test.go @@ -448,6 +448,7 @@ type testBackend struct { syncMaxTO time.Duration sentTx *types.Transaction sentTxHash common.Hash + headFeed *event.Feed } func newTestBackend(t *testing.T, n int, gspec *core.Genesis, engine consensus.Engine, generator func(i int, b *core.BlockGen)) *testBackend { @@ -474,6 +475,7 @@ func newTestBackend(t *testing.T, n int, gspec *core.Genesis, engine consensus.E acc: acc, pending: blocks[n], pendingReceipts: receipts[n], + headFeed: new(event.Feed), } return backend } @@ -598,7 +600,7 @@ func (b testBackend) SubscribeChainEvent(ch chan<- core.ChainEvent) event.Subscr panic("implement me") } func (b testBackend) SubscribeChainHeadEvent(ch chan<- core.ChainHeadEvent) event.Subscription { - panic("implement me") + return b.headFeed.Subscribe(ch) } func (b *testBackend) SendTx(ctx context.Context, tx *types.Transaction) error { b.sentTx = tx @@ -3959,7 +3961,6 @@ func makeSignedRaw(t *testing.T, api *TransactionAPI, from, to common.Address, v func makeSelfSignedRaw(t *testing.T, api *TransactionAPI, addr common.Address) (hexutil.Bytes, *types.Transaction) { return makeSignedRaw(t, api, addr, addr, big.NewInt(0)) } - func TestSendRawTransactionSync_Success(t *testing.T) { t.Parallel() genesis := &core.Genesis{