Repository navigation
Bound batch decode by item limit (PLT-820) #98
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 1 commit
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,150 @@ | ||
| package rpc | ||
|
|
||
| import ( | ||
| "encoding/json" | ||
| "fmt" | ||
| "net/http" | ||
| "net/http/httptest" | ||
| "strings" | ||
| "testing" | ||
| ) | ||
|
|
||
| // makeScalarBatch builds a compact batch array of n integer elements. Each element is a | ||
| // scalar rather than a JSON-RPC object, which is the shape that makes a small frame | ||
| // expand into one jsonrpcMessage per element. | ||
| func makeScalarBatch(n int) string { | ||
| var b strings.Builder | ||
| b.WriteByte('[') | ||
| for i := 0; i < n; i++ { | ||
| if i > 0 { | ||
| b.WriteByte(',') | ||
| } | ||
| b.WriteByte('1') | ||
| } | ||
| b.WriteByte(']') | ||
| return b.String() | ||
| } | ||
|
|
||
| // makeCallBatch builds a batch of n test_echo calls, the first of which carries id 1. | ||
| func makeCallBatch(n int) string { | ||
| elems := make([]string, n) | ||
| for i := range elems { | ||
| elems[i] = fmt.Sprintf(`{"jsonrpc":"2.0","id":%d,"method":"test_echo","params":["x",99]}`, i+1) | ||
| } | ||
| return "[" + strings.Join(elems, ",") + "]" | ||
| } | ||
|
|
||
| func TestParseMessageBatchItemLimit(t *testing.T) { | ||
| t.Parallel() | ||
|
|
||
| tests := []struct { | ||
| name string | ||
| raw string | ||
| itemLimit int | ||
| wantLen int | ||
| wantBatch bool | ||
| }{ | ||
| {name: "under limit", raw: makeScalarBatch(3), itemLimit: 5, wantLen: 3, wantBatch: true}, | ||
| {name: "at limit", raw: makeScalarBatch(5), itemLimit: 5, wantLen: 5, wantBatch: true}, | ||
| {name: "one over limit", raw: makeScalarBatch(6), itemLimit: 5, wantLen: 6, wantBatch: true}, | ||
| {name: "far over limit", raw: makeScalarBatch(100000), itemLimit: 5, wantLen: 6, wantBatch: true}, | ||
| {name: "no limit", raw: makeScalarBatch(1000), itemLimit: 0, wantLen: 1000, wantBatch: true}, | ||
| {name: "empty batch", raw: "[]", itemLimit: 5, wantLen: 0, wantBatch: true}, | ||
| {name: "single message", raw: `{"jsonrpc":"2.0","id":1,"method":"test_echo"}`, itemLimit: 5, wantLen: 1, wantBatch: false}, | ||
| } | ||
| for _, test := range tests { | ||
| t.Run(test.name, func(t *testing.T) { | ||
| t.Parallel() | ||
|
|
||
| msgs, batch := parseMessage(json.RawMessage(test.raw), test.itemLimit) | ||
| if len(msgs) != test.wantLen { | ||
| t.Fatalf("decoded %d messages, want %d", len(msgs), test.wantLen) | ||
| } | ||
| if batch != test.wantBatch { | ||
| t.Fatalf("batch = %v, want %v", batch, test.wantBatch) | ||
| } | ||
| }) | ||
| } | ||
| } | ||
|
|
||
| // TestParseMessageBatchAllocationsBounded is the regression guard for the report: a | ||
| // compact array of scalars must not allocate a jsonrpcMessage per element before the | ||
| // item limit is consulted. | ||
| // It cannot run in parallel: testing.AllocsPerRun panics in a parallel test. | ||
| func TestParseMessageBatchAllocationsBounded(t *testing.T) { | ||
| const ( | ||
| elements = 200000 | ||
| itemLimit = 100 | ||
| ) | ||
| raw := json.RawMessage(makeScalarBatch(elements)) | ||
|
|
||
| allocs := testing.AllocsPerRun(2, func() { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [suggestion] The 1000-alloc ceiling against an expected ~100-300 gives decent headroom, so this probably won't flake often — but when it does the failure will be confusing and unrelated to this code. Consider making it explicitly differential instead, e.g. measure |
||
| parseMessage(raw, itemLimit) | ||
| }) | ||
| // Bound generously: the point is the order of magnitude, not the exact count. Before | ||
| // the limit was pushed into parseMessage this exceeded `elements` allocations. | ||
| if maxAllocs := float64(10 * itemLimit); allocs > maxAllocs { | ||
| t.Fatalf("parseMessage made %.0f allocations for a %d-element batch limited to %d, want at most %.0f", | ||
| allocs, elements, itemLimit, maxAllocs) | ||
| } | ||
| } | ||
|
|
||
| // TestWSOversizeBatchRejectedAndConnectionSurvives checks that the truncated decode still | ||
| // produces the protocol-level rejection, and that the connection remains usable after it. | ||
| func TestWSOversizeBatchRejectedAndConnectionSurvives(t *testing.T) { | ||
| t.Parallel() | ||
|
|
||
| srv := newTestServer() | ||
| srv.SetBatchLimits(4, 100000) | ||
| _, wsURL := startWSTestServer(t, srv) | ||
|
|
||
| conn := dialWS(t, wsURL) | ||
| defer conn.Close() | ||
|
|
||
| writeWSJSON(t, conn, makeScalarBatch(50000)) | ||
| var resp []jsonrpcMessage | ||
| readWSJSON(t, conn, &resp) | ||
| if len(resp) != 1 { | ||
| t.Fatalf("got %d responses, want 1", len(resp)) | ||
| } | ||
| if resp[0].Error == nil || resp[0].Error.Message != errMsgBatchTooLarge { | ||
| t.Fatalf("wrong response to oversize batch: %+v", resp[0]) | ||
| } | ||
|
|
||
| // The connection must still serve the next request. | ||
| writeWSJSON(t, conn, `{"jsonrpc":"2.0","id":7,"method":"test_echo","params":["x",99]}`) | ||
| var next jsonrpcMessage | ||
| readWSJSON(t, conn, &next) | ||
| if next.Error != nil { | ||
| t.Fatalf("request after oversize batch failed: %v", next.Error) | ||
| } | ||
| } | ||
|
|
||
| // TestHTTPOversizeBatchRejected covers the single-request path, which builds its handler | ||
| // separately from ServeCodec and so wires the codec up on its own. | ||
| func TestHTTPOversizeBatchRejected(t *testing.T) { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Assert the codec actually observed the limit? Same for |
||
| t.Parallel() | ||
|
|
||
| srv := newTestServer() | ||
| defer srv.Stop() | ||
| srv.SetBatchLimits(4, 100000) | ||
| httpsrv := httptest.NewServer(srv) | ||
| defer httpsrv.Close() | ||
|
|
||
| resp, err := http.Post(httpsrv.URL, "application/json", strings.NewReader(makeCallBatch(50000))) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [nit] Since the server only ever decodes |
||
| if err != nil { | ||
| t.Fatalf("post batch: %v", err) | ||
| } | ||
| defer resp.Body.Close() | ||
|
|
||
| var msgs []jsonrpcMessage | ||
| if err := json.NewDecoder(resp.Body).Decode(&msgs); err != nil { | ||
| t.Fatalf("decode response: %v", err) | ||
| } | ||
| if len(msgs) != 1 { | ||
| t.Fatalf("got %d responses, want 1", len(msgs)) | ||
| } | ||
| if msgs[0].Error == nil || msgs[0].Error.Message != errMsgBatchTooLarge { | ||
| t.Fatalf("wrong response to oversize batch: %+v", msgs[0]) | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -126,19 +126,22 @@ func (c *Client) newClientConn(conn ServerCodec) *clientConn { | |
| ctx = context.WithValue(ctx, clientContextKey{}, c) | ||
| ctx = context.WithValue(ctx, peerInfoContextKey{}, conn.peerInfo()) | ||
| handler := newHandler(ctx, conn, c.idgen, c.services, c.batchItemLimit, c.batchResponseMaxSize, c.wsConcurrentBudget, c.readLimit, c.admissionEventHook, c.wsAdmissionTimeout) | ||
| attachBudgetHandler(conn, handler) | ||
| attachHandler(conn, handler) | ||
| return &clientConn{conn, handler} | ||
| } | ||
|
|
||
| // budgetHandlerSetter is implemented by codecs that need the handler wired in for | ||
| // byte-budget admission (jsonCodec and, via embedding, websocketCodec). | ||
| type budgetHandlerSetter interface { | ||
| setBudgetHandler(h *handler) | ||
| // handlerSetter is implemented by codecs that need the handler wired in to enforce its | ||
| // read limits (jsonCodec and, via embedding, websocketCodec). | ||
| type handlerSetter interface { | ||
| setHandler(h *handler) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [nit] Raised by the Codex pass, and it's a fair point though I'd rank it below blocking: Two things soften it: the shape is pre-existing (it was |
||
| } | ||
|
|
||
| func attachBudgetHandler(codec ServerCodec, h *handler) { | ||
| if c, ok := codec.(budgetHandlerSetter); ok { | ||
| c.setBudgetHandler(h) | ||
| // attachHandler wires h into codec so reads on it observe the handler's limits. A codec | ||
| // that does not implement handlerSetter reads unlimited, and the handler's own checks | ||
| // remain the backstop. | ||
| func attachHandler(codec ServerCodec, h *handler) { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [suggestion] Raised by Codex, kept with a lower severity. Two things temper this relative to Codex's "Medium":
So the exposure is external embedders only, and the handler-level check still rejects the batch — the DoS is degraded to "allocation happens, then rejection," not "batch executes." Still, since bounded decoding is now a security property rather than just a budgeting optimization, it's worth making it not depend on a private-interface assertion — e.g. plumb the item limit through codec construction, or fall back to a package-level default in |
||
| if c, ok := codec.(handlerSetter); ok { | ||
| c.setHandler(h) | ||
| } | ||
| } | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -393,7 +393,9 @@ func (h *handler) respondWithBatchTooLarge(cp *callProc, batch []*jsonrpcMessage | |
| resp := errorMessage(&invalidRequestError{errMsgBatchTooLarge}) | ||
| // Find the first call and add its "id" field to the error. | ||
| // This is the best we can do, given that the protocol doesn't have a way | ||
| // of reporting an error for the entire batch. | ||
| // of reporting an error for the entire batch. The batch is only decoded up to | ||
| // the item limit, so a batch whose every decoded element is a notification is | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [suggestion] Worth calling out in the PR description, not just this comment: this is a user-visible protocol regression, and the consequence is a client hang rather than a cosmetic id change. For a batch whose first Impact is limited in practice — geth's own |
||
| // answered with a null id even if a later element was a call. | ||
| for _, msg := range batch { | ||
| if msg.isCall() { | ||
| resp.ID = msg.ID | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -185,7 +185,7 @@ type jsonCodec struct { | |
| encMu sync.Mutex // guards the encoder | ||
| encode encodeFunc // encoder to allow multiple transports | ||
| conn deadlineCloser | ||
| handler *handler // set by the read loop for byte-budget admission | ||
| handler *handler // set by the read loop; source of this codec's read limits | ||
| } | ||
|
|
||
| type encodeFunc = func(v interface{}, isErrorResponse bool) error | ||
|
|
@@ -221,11 +221,20 @@ func NewCodec(conn Conn) ServerCodec { | |
| return NewFuncCodec(conn, encode, dec.Decode) | ||
| } | ||
|
|
||
| // setBudgetHandler wires in the handler used for byte-budget admission. | ||
| func (c *jsonCodec) setBudgetHandler(h *handler) { | ||
| // setHandler wires in the handler whose limits govern reads on this codec. | ||
| func (c *jsonCodec) setHandler(h *handler) { | ||
| c.handler = h | ||
| } | ||
|
|
||
| // batchItemLimit returns the number of batch elements a read may decode, or 0 when no | ||
| // limit applies. | ||
| func (c *jsonCodec) batchItemLimit() int { | ||
| if c.handler == nil { | ||
| return 0 | ||
| } | ||
| return c.handler.batchRequestLimit | ||
| } | ||
|
|
||
| func (c *jsonCodec) peerInfo() PeerInfo { | ||
| // This returns "ipc" because all other built-in transports have a separate codec type. | ||
| return PeerInfo{Transport: "ipc", RemoteAddr: c.remote} | ||
|
|
@@ -252,7 +261,7 @@ func (c *jsonCodec) readBatch() (messages []*jsonrpcMessage, batch bool, rawLen | |
| } | ||
| return nil, false, 0, err | ||
| } | ||
| messages, batch = parseMessage(rawmsg) | ||
| messages, batch = parseMessage(rawmsg, c.batchItemLimit()) | ||
| for i, msg := range messages { | ||
| if msg == nil { | ||
| // Message is JSON 'null'. Replace with zero value so it | ||
|
|
@@ -290,8 +299,9 @@ func (c *jsonCodec) closed() <-chan interface{} { | |
| // parseMessage parses raw bytes as a (batch of) JSON-RPC message(s). There are no error | ||
| // checks in this function because the raw message has already been syntax-checked when it | ||
| // is called. Any non-JSON-RPC messages in the input return the zero value of | ||
| // jsonrpcMessage. | ||
| func parseMessage(raw json.RawMessage) ([]*jsonrpcMessage, bool) { | ||
| // jsonrpcMessage. A batch is decoded to at most itemLimit+1 messages; an itemLimit of 0 | ||
| // means unlimited. | ||
| func parseMessage(raw json.RawMessage, itemLimit int) ([]*jsonrpcMessage, bool) { | ||
| if !isBatch(raw) { | ||
| msgs := []*jsonrpcMessage{{}} | ||
| json.Unmarshal(raw, &msgs[0]) | ||
|
|
@@ -301,6 +311,13 @@ func parseMessage(raw json.RawMessage) ([]*jsonrpcMessage, bool) { | |
| dec.Token() // skip '[' | ||
| var msgs []*jsonrpcMessage | ||
| for dec.More() { | ||
| // Stop one element past the limit rather than at it. That surplus element is | ||
| // what handleBatch's own count check reads to reject the batch, and decoding | ||
| // past it would allocate a jsonrpcMessage per element of an array the server | ||
| // has already decided not to serve. | ||
| if itemLimit > 0 && len(msgs) > itemLimit { | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [suggestion] The The comment explains the why but doesn't guard it. Two cheap options: reference (The limit source itself is safe by construction — |
||
| break | ||
| } | ||
| msgs = append(msgs, new(jsonrpcMessage)) | ||
| dec.Decode(&msgs[len(msgs)-1]) | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -223,6 +223,7 @@ func (s *Server) serveSingleRequest(ctx context.Context, codec ServerCodec) { | |
|
|
||
| h := newHandler(ctx, codec, s.idgen, &s.services, s.batchItemLimit, s.batchResponseLimit, nil, s.readLimit, nil, s.wsAdmissionTimeout) | ||
| h.allowSubscribe = false | ||
| attachHandler(codec, h) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [suggestion] Attaching the handler here also switches on the other behavior It is a latent trap though: nothing on the HTTP single-request path ever calls |
||
| defer h.close(io.EOF, nil) | ||
|
|
||
| reqs, batch, _, err := codec.readBatch() | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Maybe time to upgrade the go module?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Actually upgrading go.mod does not remove the need to keep this test non-parallel.
If this test calls t.Parallel():
AllocsPerRuncan panic if other parallel tests in the package are actively running — and rpc package has manyt.Parallel()tests.