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
14 changes: 14 additions & 0 deletions cmd/juno/juno.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ const (
rpcRequestTimeoutF = "rpc-request-timeout"
rpcMaxConcurrentRequestsF = "rpc-max-concurrent-requests"
rpcMaxRequestQueueF = "rpc-max-request-queue"
rpcMaxWSConnectionsF = "rpc-max-ws-connections"
rpcMaxSubscriptionsF = "rpc-max-subscriptions"
rpcMaxBatchSizeF = "rpc-max-batch-size"
rpcMaxBatchResponseSizeF = "rpc-max-batch-response-size"
rpcBatchConcurrencyF = "rpc-batch-concurrency"
Expand Down Expand Up @@ -182,6 +184,8 @@ const (
defaultRPCRequestTimeout = 1 * time.Minute
defaultRPCMaxConcurrentRequests = 256000
defaultRPCMaxQueuedRequests = 256000
defaultRPCMaxWSConnections = 1024
defaultRPCMaxSubscriptions = 128
defaultRPCMaxBatchSize = 1000
defaultRPCMaxBatchResponseSize = 64 // MB
defaultRPCBatchConcurrency = uint(0)
Expand Down Expand Up @@ -278,6 +282,12 @@ const (
"together; 0 disables the limit."
rpcMaxRequestQueueUsage = "Maximum number of HTTP RPC requests to queue after " +
"reaching rpc-max-concurrent-requests limit. Websocket requests are never queued."
rpcMaxWSConnectionsUsage = "Maximum concurrent websocket connections, across all " +
"RPC versions. 0 disables the limit."
rpcMaxSubscriptionsUsage = "Maximum subscriptions one websocket connection may hold. " +
"The limit is per connection so that one client cannot take the slots another needs; " +
"multiplied by rpc-max-ws-connections it is the most the process will carry. " +
"0 disables the limit."
rpcMaxBatchSizeUsage = "Maximum number of calls in a single batch request. " +
"0 disables the limit."
rpcMaxBatchResponseSizeUsage = "Size (in MBs) at which a batch stops being processed. " +
Expand Down Expand Up @@ -528,6 +538,8 @@ func NewCmd(config *node.Config, run func(*cobra.Command, []string) error) *cobr
defaultRPCMaxQueuedRequests,
rpcMaxRequestQueueUsage,
)
junoCmd.Flags().Uint(rpcMaxWSConnectionsF, defaultRPCMaxWSConnections, rpcMaxWSConnectionsUsage)
junoCmd.Flags().Uint(rpcMaxSubscriptionsF, defaultRPCMaxSubscriptions, rpcMaxSubscriptionsUsage)
junoCmd.Flags().Uint(rpcMaxBatchSizeF, defaultRPCMaxBatchSize, rpcMaxBatchSizeUsage)
junoCmd.Flags().Uint(
rpcMaxBatchResponseSizeF,
Expand All @@ -553,6 +565,8 @@ func NewCmd(config *node.Config, run func(*cobra.Command, []string) error) *cobr
rpcRequestTimeoutF,
rpcMaxConcurrentRequestsF,
rpcMaxRequestQueueF,
rpcMaxWSConnectionsF,
rpcMaxSubscriptionsF,
rpcMaxBatchSizeF,
rpcMaxBatchResponseSizeF,
rpcBatchConcurrencyF,
Expand Down
28 changes: 28 additions & 0 deletions cmd/juno/juno_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ func TestConfigPrecedence(t *testing.T) {
defaultMaxVMs := uint(3 * runtime.GOMAXPROCS(0))
defaultRPCMaxConcurrentRequests := uint(256000)
defaultRPCMaxRequestQueue := uint(256000)
defaultRPCMaxWSConnections := uint(1024)
defaultRPCMaxSubscriptions := uint(128)
defaultRPCMaxBatchSize := uint(1000)
defaultRPCMaxBatchResponseSize := uint(64)
defaultRPCMaxBlockScan := uint(math.MaxUint)
Expand Down Expand Up @@ -125,6 +127,8 @@ func TestConfigPrecedence(t *testing.T) {
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -177,6 +181,8 @@ func TestConfigPrecedence(t *testing.T) {
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -338,6 +344,8 @@ pprof: true
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -396,6 +404,8 @@ http-port: 4576
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -452,6 +462,8 @@ http-port: 4576
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -510,6 +522,8 @@ http-port: 4576
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -591,6 +605,8 @@ db-cache-size: 1024
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -651,6 +667,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -707,6 +725,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -761,6 +781,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -816,6 +838,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -871,6 +895,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down Expand Up @@ -925,6 +951,8 @@ network: sepolia
MaxVMQueue: 2 * defaultMaxVMs,
RPCMaxConcurrentRequests: defaultRPCMaxConcurrentRequests,
RPCMaxRequestQueue: defaultRPCMaxRequestQueue,
RPCMaxWSConnections: defaultRPCMaxWSConnections,
RPCMaxSubscriptions: defaultRPCMaxSubscriptions,
RPCMaxBatchSize: defaultRPCMaxBatchSize,
RPCMaxBatchResponseSize: defaultRPCMaxBatchResponseSize,
RPCMaxBlockScan: defaultRPCMaxBlockScan,
Expand Down
29 changes: 27 additions & 2 deletions jsonrpc/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -287,10 +287,12 @@ type Conn interface {
io.Writer
Equal(Conn) bool
Context() context.Context
SubscriptionSlots
}

type connection struct {
w io.Writer
slots SubscriptionSlots
activated <-chan struct{}
ctx context.Context

Expand Down Expand Up @@ -320,6 +322,28 @@ func (c *connection) Context() context.Context {
return c.ctx
}

// SubscriptionSlots caps how many subscriptions one connection carries at once.
type SubscriptionSlots interface {
// TryAcquireSubscription reserves a slot without waiting
TryAcquireSubscription() bool
// ReleaseSubscription returns a slot taken by TryAcquireSubscription
ReleaseSubscription()
}

// Transport is what HandleReadWriter needs of a connection
type Transport interface {
io.ReadWriter
SubscriptionSlots
}

func (c *connection) TryAcquireSubscription() bool {
return c.slots.TryAcquireSubscription()
}

func (c *connection) ReleaseSubscription() {
c.slots.ReleaseSubscription()
}

// ConnKey the key used to retrieve the connection from the context passed to a handler.
// It is exported to allow transports to set it manually if they decide not to use HandleReadWriter, which sets it automatically.
// Manually setting the connection can be especially useful when testing handlers.
Expand All @@ -343,12 +367,13 @@ func ConnFromContext(ctx context.Context) (Conn, bool) {
func (s *Server) HandleReadWriter(
connCtx context.Context,
requestTimeout time.Duration,
rw io.ReadWriter,
rw Transport,
) error {
activated := make(chan struct{})
defer close(activated)
conn := &connection{
w: rw.(io.Writer),
w: rw,
slots: rw,
activated: activated,
ctx: connCtx,
}
Expand Down
19 changes: 17 additions & 2 deletions jsonrpc/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -710,6 +710,17 @@ func TestCannotWriteToConnInHandler(t *testing.T) {
require.NotNil(t, header)
}

// uncappedTransport hands HandleReadWriter a stream with no subscription cap,
// which these tests have no use for. The cap is part of jsonrpc.Transport so
// that a real transport cannot forget it, and saying so here is the price.
type uncappedTransport struct {
net.Conn
}

func (uncappedTransport) TryAcquireSubscription() bool { return true }

func (uncappedTransport) ReleaseSubscription() {}

type fakeConn struct {
ctx context.Context
}
Expand All @@ -726,6 +737,10 @@ func (fc *fakeConn) Context() context.Context {
return fc.ctx
}

func (fc *fakeConn) TryAcquireSubscription() bool { return true }

func (fc *fakeConn) ReleaseSubscription() {}

func TestWriteToConnInHandler(t *testing.T) {
testBytes := "written from handler"
server := jsonrpc.NewServer(1, log.NewNopZapLogger())
Expand Down Expand Up @@ -754,7 +769,7 @@ func TestWriteToConnInHandler(t *testing.T) {
})

wg.Go(func() {
err := server.HandleReadWriter(t.Context(), 0, serverConn)
err := server.HandleReadWriter(t.Context(), 0, uncappedTransport{serverConn})
require.NoError(t, err)
})

Expand Down Expand Up @@ -791,7 +806,7 @@ func TestWriteToClosedConnInHandler(t *testing.T) {
})

wg.Go(func() {
err := server.HandleReadWriter(t.Context(), 0, serverConn)
err := server.HandleReadWriter(t.Context(), 0, uncappedTransport{serverConn})
require.ErrorIs(t, err, io.ErrClosedPipe)
})

Expand Down
Loading
Loading