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 all commits
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,213 @@ | ||
| 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. | ||
| // Keep n small: the body counts against defaultBodyLimit, and the server only ever | ||
| // decodes itemLimit+1 elements, so a few hundred exercises the same path as a few | ||
| // thousand. | ||
| 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, ",") + "]" | ||
| } | ||
|
|
||
| // makeTruncationProbe builds an oversize batch whose only call, id 99, sits just past the | ||
| // first itemLimit+1 elements. respondWithBatchTooLarge reports the first decoded call's | ||
| // id, so a null id means the decode stopped at the limit and an id of 99 means it did | ||
| // not. | ||
| func makeTruncationProbe(itemLimit int) string { | ||
| elems := make([]string, itemLimit+1) | ||
| for i := range elems { | ||
| elems[i] = `{"jsonrpc":"2.0","method":"test_echo","params":["x",99]}` | ||
| } | ||
| elems = append(elems, `{"jsonrpc":"2.0","id":99,"method":"test_echo","params":["x",99]}`) | ||
| return "[" + strings.Join(elems, ",") + "]" | ||
| } | ||
|
|
||
| // assertBatchTooLarge checks that resp is the single batch-too-large error response, and | ||
| // that it carries wantID ("null" for the id-less form). | ||
| func assertBatchTooLarge(t *testing.T, resp []jsonrpcMessage, wantID string) { | ||
| t.Helper() | ||
|
|
||
| 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]) | ||
| } | ||
| if id := string(resp[0].ID); id != wantID { | ||
| t.Fatalf("error id = %s, want %s", id, wantID) | ||
| } | ||
| } | ||
|
|
||
| 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 websocket codec reads | ||
| // with the handler's item limit, 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() | ||
|
|
||
| const itemLimit = 4 | ||
|
|
||
| srv := newTestServer() | ||
| srv.SetBatchLimits(itemLimit, 100000) | ||
| _, wsURL := startWSTestServer(t, srv) | ||
|
|
||
| conn := dialWS(t, wsURL) | ||
| defer conn.Close() | ||
|
|
||
| // A call just past the decoded prefix must not be found, which is only true if the | ||
| // codec passed the item limit into parseMessage. | ||
| writeWSJSON(t, conn, makeTruncationProbe(itemLimit)) | ||
| var probeResp []jsonrpcMessage | ||
| readWSJSON(t, conn, &probeResp) | ||
| assertBatchTooLarge(t, probeResp, "null") | ||
|
|
||
| // A call inside the decoded prefix is still found, so the null id above is truncation | ||
| // rather than ids going missing altogether. | ||
| writeWSJSON(t, conn, makeCallBatch(200)) | ||
| var callResp []jsonrpcMessage | ||
| readWSJSON(t, conn, &callResp) | ||
| assertBatchTooLarge(t, callResp, "1") | ||
|
|
||
| // A large cheap array, the shape from the report, is rejected the same way. | ||
| writeWSJSON(t, conn, makeScalarBatch(50000)) | ||
| var scalarResp []jsonrpcMessage | ||
| readWSJSON(t, conn, &scalarResp) | ||
| assertBatchTooLarge(t, scalarResp, "null") | ||
|
|
||
| // 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. The null id in the | ||
| // probe case is what proves that wiring is in place: without serveSingleRequest's | ||
| // attachHandler call the codec reads unlimited, finds the trailing call and reports its | ||
| // id. | ||
| 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() | ||
|
|
||
| const itemLimit = 4 | ||
|
|
||
| srv := newTestServer() | ||
| defer srv.Stop() | ||
| srv.SetBatchLimits(itemLimit, 100000) | ||
| httpsrv := httptest.NewServer(srv) | ||
| defer httpsrv.Close() | ||
|
|
||
| tests := []struct { | ||
| name string | ||
| batch string | ||
| wantID string | ||
| }{ | ||
| // A call past the decoded prefix cannot be reported. | ||
| {name: "truncation probe", batch: makeTruncationProbe(itemLimit), wantID: "null"}, | ||
| // A call inside the prefix still is, so the null id above is truncation rather | ||
| // than ids going missing altogether. | ||
| {name: "leading call", batch: makeCallBatch(200), wantID: "1"}, | ||
| // The large cheap array from the report. | ||
| {name: "scalar batch", batch: makeScalarBatch(50000), wantID: "null"}, | ||
| } | ||
| for _, test := range tests { | ||
| // Not parallel: the subtests must run before the deferred httpsrv.Close. | ||
| t.Run(test.name, func(t *testing.T) { | ||
| resp, err := http.Post(httpsrv.URL, "application/json", strings.NewReader(test.batch)) | ||
| 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) | ||
| } | ||
| assertBatchTooLarge(t, msgs, test.wantID) | ||
| }) | ||
| } | ||
| } | ||
| 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 |
|---|---|---|
|
|
@@ -127,10 +127,13 @@ func WithHTTPAuth(a HTTPAuth) ClientOption { | |
| // auth information to the request. | ||
| type HTTPAuth func(h http.Header) error | ||
|
|
||
| // WithBatchItemLimit changes the maximum number of items allowed in batch requests. | ||
| // WithBatchItemLimit changes the maximum number of items allowed in an incoming batch. | ||
| // | ||
| // Note: this option applies when processing incoming batch requests. It does not affect | ||
| // batch requests sent by the client. | ||
| // Note: this option applies to batches the client receives: both batch requests sent by | ||
| // the server on a bidirectional connection and batched responses to the client's own | ||
| // requests. A batch with more items than the limit is rejected instead of dispatched, and | ||
|
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] "rejected instead of dispatched" undersells what happens for the response case: when the over-limit batch is a batch of responses to the client's own calls, |
||
| // only its first limit+1 items are decoded. It does not cap the size of the batches the | ||
| // client itself sends. | ||
| func WithBatchItemLimit(limit int) ClientOption { | ||
| return optionFunc(func(cfg *clientConfig) { | ||
| cfg.batchItemLimit = limit | ||
|
|
||
| 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,20 @@ 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. | ||
| // | ||
| // The overshoot is load-bearing: handleBatch rejects on | ||
| // len(msgs) > h.batchRequestLimit, so returning exactly itemLimit elements here | ||
| // would make an over-limit batch look in-limit and be executed. Keep this break | ||
| // and that comparison in agreement. TestParseMessageBatchItemLimit pins the | ||
| // itemLimit+1 result, and testdata/reqresp-batch.js pins that a batch of exactly | ||
| // the limit is still served. | ||
| 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.