mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
Merge pull request #4 from flashbots/feat/max-concurrent-build-sessions
Limit concurrent builder sessions
This commit is contained in:
commit
33e520a14c
4 changed files with 75 additions and 14 deletions
|
|
@ -8,7 +8,7 @@ import (
|
||||||
|
|
||||||
// SessionManager is the backend that manages the session state of the builder API.
|
// SessionManager is the backend that manages the session state of the builder API.
|
||||||
type SessionManager interface {
|
type SessionManager interface {
|
||||||
NewSession() (string, error)
|
NewSession(context.Context) (string, error)
|
||||||
AddTransaction(sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error)
|
AddTransaction(sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -24,7 +24,7 @@ type Server struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) NewSession(ctx context.Context) (string, error) {
|
func (s *Server) NewSession(ctx context.Context) (string, error) {
|
||||||
return s.sessionMngr.NewSession()
|
return s.sessionMngr.NewSession(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) AddTransaction(ctx context.Context, sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error) {
|
func (s *Server) AddTransaction(ctx context.Context, sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error) {
|
||||||
|
|
|
||||||
|
|
@ -30,10 +30,10 @@ func TestAPI(t *testing.T) {
|
||||||
|
|
||||||
type nullSessionManager struct{}
|
type nullSessionManager struct{}
|
||||||
|
|
||||||
func (n *nullSessionManager) NewSession() (string, error) {
|
func (nullSessionManager) NewSession(ctx context.Context) (string, error) {
|
||||||
return "1", nil
|
return "1", ctx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (n *nullSessionManager) AddTransaction(sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error) {
|
func (nullSessionManager) AddTransaction(sessionId string, tx *types.Transaction) (*types.SimulateTransactionResult, error) {
|
||||||
return &types.SimulateTransactionResult{Logs: []*types.SimulatedLog{}}, nil
|
return &types.SimulateTransactionResult{Logs: []*types.SimulatedLog{}}, nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
package builder
|
package builder
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
@ -31,11 +32,13 @@ type blockchain interface {
|
||||||
}
|
}
|
||||||
|
|
||||||
type Config struct {
|
type Config struct {
|
||||||
GasCeil uint64
|
GasCeil uint64
|
||||||
SessionIdleTimeout time.Duration
|
SessionIdleTimeout time.Duration
|
||||||
|
MaxConcurrentSessions int
|
||||||
}
|
}
|
||||||
|
|
||||||
type SessionManager struct {
|
type SessionManager struct {
|
||||||
|
sem chan struct{}
|
||||||
sessions map[string]*builder
|
sessions map[string]*builder
|
||||||
sessionTimers map[string]*time.Timer
|
sessionTimers map[string]*time.Timer
|
||||||
sessionsLock sync.RWMutex
|
sessionsLock sync.RWMutex
|
||||||
|
|
@ -50,8 +53,17 @@ func NewSessionManager(blockchain blockchain, config *Config) *SessionManager {
|
||||||
if config.SessionIdleTimeout == 0 {
|
if config.SessionIdleTimeout == 0 {
|
||||||
config.SessionIdleTimeout = 5 * time.Second
|
config.SessionIdleTimeout = 5 * time.Second
|
||||||
}
|
}
|
||||||
|
if config.MaxConcurrentSessions <= 0 {
|
||||||
|
config.MaxConcurrentSessions = 16 // chosen arbitrarily
|
||||||
|
}
|
||||||
|
|
||||||
|
sem := make(chan struct{}, config.MaxConcurrentSessions)
|
||||||
|
for len(sem) < cap(sem) {
|
||||||
|
sem <- struct{}{} // fill 'er up
|
||||||
|
}
|
||||||
|
|
||||||
s := &SessionManager{
|
s := &SessionManager{
|
||||||
|
sem: sem,
|
||||||
sessions: make(map[string]*builder),
|
sessions: make(map[string]*builder),
|
||||||
sessionTimers: make(map[string]*time.Timer),
|
sessionTimers: make(map[string]*time.Timer),
|
||||||
blockchain: blockchain,
|
blockchain: blockchain,
|
||||||
|
|
@ -61,12 +73,17 @@ func NewSessionManager(blockchain blockchain, config *Config) *SessionManager {
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewSession creates a new builder session and returns the session id
|
// NewSession creates a new builder session and returns the session id
|
||||||
func (s *SessionManager) NewSession() (string, error) {
|
func (s *SessionManager) NewSession(ctx context.Context) (string, error) {
|
||||||
s.sessionsLock.Lock()
|
// Wait for session to become available
|
||||||
defer s.sessionsLock.Unlock()
|
select {
|
||||||
|
case <-s.sem:
|
||||||
|
s.sessionsLock.Lock()
|
||||||
|
defer s.sessionsLock.Unlock()
|
||||||
|
case <-ctx.Done():
|
||||||
|
return "", ctx.Err()
|
||||||
|
}
|
||||||
|
|
||||||
parent := s.blockchain.CurrentHeader()
|
parent := s.blockchain.CurrentHeader()
|
||||||
|
|
||||||
chainConfig := s.blockchain.Config()
|
chainConfig := s.blockchain.Config()
|
||||||
|
|
||||||
header := &types.Header{
|
header := &types.Header{
|
||||||
|
|
@ -111,6 +128,14 @@ func (s *SessionManager) NewSession() (string, error) {
|
||||||
delete(s.sessionTimers, id)
|
delete(s.sessionTimers, id)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
// Technically, we are certain that there is an open slot in the semaphore
|
||||||
|
// channel, but let's be defensive and panic if the invariant is violated.
|
||||||
|
select {
|
||||||
|
case s.sem <- struct{}{}:
|
||||||
|
default:
|
||||||
|
panic("released more sessions than are open") // unreachable
|
||||||
|
}
|
||||||
|
|
||||||
return id, nil
|
return id, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
package builder
|
package builder
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"crypto/ecdsa"
|
"crypto/ecdsa"
|
||||||
"math/big"
|
"math/big"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
@ -21,7 +22,7 @@ func TestSessionManager_SessionTimeout(t *testing.T) {
|
||||||
SessionIdleTimeout: 500 * time.Millisecond,
|
SessionIdleTimeout: 500 * time.Millisecond,
|
||||||
})
|
})
|
||||||
|
|
||||||
id, err := mngr.NewSession()
|
id, err := mngr.NewSession(context.TODO())
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
time.Sleep(1 * time.Second)
|
time.Sleep(1 * time.Second)
|
||||||
|
|
@ -30,12 +31,47 @@ func TestSessionManager_SessionTimeout(t *testing.T) {
|
||||||
require.Error(t, err)
|
require.Error(t, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSessionManager_MaxConcurrentSessions(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
const d = time.Millisecond * 100
|
||||||
|
|
||||||
|
mngr, _ := newSessionManager(t, &Config{
|
||||||
|
MaxConcurrentSessions: 1,
|
||||||
|
SessionIdleTimeout: d,
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("SessionAvailable", func(t *testing.T) {
|
||||||
|
sess, err := mngr.NewSession(context.TODO())
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotZero(t, sess)
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("ContextExpired", func(t *testing.T) {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
sess, err := mngr.NewSession(ctx)
|
||||||
|
require.Zero(t, sess)
|
||||||
|
require.ErrorIs(t, err, context.Canceled)
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("SessionExpired", func(t *testing.T) {
|
||||||
|
time.Sleep(d) // Wait for the session to expire.
|
||||||
|
|
||||||
|
// We should be able to open a session again.
|
||||||
|
sess, err := mngr.NewSession(context.TODO())
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotZero(t, sess)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
func TestSessionManager_SessionRefresh(t *testing.T) {
|
func TestSessionManager_SessionRefresh(t *testing.T) {
|
||||||
mngr, _ := newSessionManager(t, &Config{
|
mngr, _ := newSessionManager(t, &Config{
|
||||||
SessionIdleTimeout: 500 * time.Millisecond,
|
SessionIdleTimeout: 500 * time.Millisecond,
|
||||||
})
|
})
|
||||||
|
|
||||||
id, err := mngr.NewSession()
|
id, err := mngr.NewSession(context.TODO())
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
// if we query the session under the idle timeout,
|
// if we query the session under the idle timeout,
|
||||||
|
|
@ -60,7 +96,7 @@ func TestSessionManager_StartSession(t *testing.T) {
|
||||||
// test that the session starts and it can simulate transactions
|
// test that the session starts and it can simulate transactions
|
||||||
mngr, bMock := newSessionManager(t, &Config{})
|
mngr, bMock := newSessionManager(t, &Config{})
|
||||||
|
|
||||||
id, err := mngr.NewSession()
|
id, err := mngr.NewSession(context.TODO())
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
txn := bMock.state.newTransfer(t, common.Address{}, big.NewInt(1))
|
txn := bMock.state.newTransfer(t, common.Address{}, big.NewInt(1))
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue