Skip to content

Commit 0f77172

Browse files
committed
fixup! feat(jsonrpc): add support for JSON-RPC 2.0 batch requests
1 parent 2633472 commit 0f77172

2 files changed

Lines changed: 70 additions & 26 deletions

File tree

internal/jsonrpc/batchcalls_test.go

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -283,6 +283,45 @@ func TestJSONRPCBatchStopsBetweenEntriesWhenContextIsCanceled(t *testing.T) {
283283
}
284284
}
285285

286+
func TestJSONRPCBatchReturnsErrorsForIDDRequestsAfterDeadline(t *testing.T) {
287+
s := newBatchTestService()
288+
var calls atomic.Int32
289+
const method = "test_deadline_batch"
290+
withTestRPCHandler(t, method, func(_ *Service, _ *http.Request, _ RPCRequest) (any, error) {
291+
calls.Add(1)
292+
return true, nil
293+
})
294+
295+
body := []byte(fmt.Sprintf(`[
296+
{"jsonrpc":"2.0","method":%q,"id":1},
297+
{"jsonrpc":"2.0","method":%q},
298+
{"jsonrpc":"2.0","method":%q,"id":"three"},
299+
false,
300+
{"jsonrpc":"2.0","method":%q,"id":null}
301+
]`, method, method, method, method))
302+
ctx, cancel := context.WithDeadline(context.Background(), time.Now().Add(-time.Second))
303+
defer cancel()
304+
req := httptest.NewRequest(http.MethodPost, "/rpc", bytes.NewReader(body)).WithContext(ctx)
305+
rr := httptest.NewRecorder()
306+
s.handleRPC(rr, req)
307+
308+
require.Zero(t, calls.Load(), "expired batch entries must not be dispatched")
309+
responses := decodeRPCBatch(t, rr.Body.Bytes())
310+
require.Len(t, responses, 5, "not all entries receive deadline errors")
311+
requireRPCError(t, responses[0], float64(1), JSONRPC_RATE_LIMIT_EXCEEDED)
312+
requireRPCError(t, responses[1], nil, JSONRPC_RATE_LIMIT_EXCEEDED)
313+
requireRPCError(t, responses[2], "three", JSONRPC_RATE_LIMIT_EXCEEDED)
314+
requireRPCError(t, responses[3], nil, JSONRPC_INVALID_REQUEST)
315+
requireRPCError(t, responses[4], nil, JSONRPC_RATE_LIMIT_EXCEEDED)
316+
for i, response := range responses {
317+
if i == 3 {
318+
require.Equal(t, "invalid request", response.Error.Message)
319+
} else {
320+
require.Equal(t, "batch deadline exceeded", response.Error.Message)
321+
}
322+
}
323+
}
324+
286325
func TestJSONRPCBatchUsesOneAdmissionPermit(t *testing.T) {
287326
s := newBatchTestService()
288327
s.admission = service.NewSemaphoreAdmission(1)

internal/jsonrpc/jsonrpc.go

Lines changed: 31 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -42,10 +42,11 @@ const (
4242
// but the requested entity does not; unknown applications use
4343
// JSONRPC_APPLICATION_NOT_FOUND. For forward-looking keys, this can be the
4444
// "not created yet" signal and may be safe to poll depending on the method.
45-
JSONRPC_RESOURCE_NOT_FOUND int = -32001 //nolint: revive
45+
JSONRPC_RESOURCE_NOT_FOUND int = -31001 //nolint: revive
4646
// Application not found: the application identifier itself is unknown to
4747
// this node. A configuration error that will not resolve by retrying.
48-
JSONRPC_APPLICATION_NOT_FOUND int = -32002 //nolint: revive
48+
JSONRPC_APPLICATION_NOT_FOUND int = -31002 //nolint: revive
49+
JSONRPC_RATE_LIMIT_EXCEEDED int = -32002 //nolint: revive
4950
JSONRPC_PARSE_ERROR int = -32700 //nolint: revive
5051
JSONRPC_INVALID_REQUEST int = -32600 //nolint: revive
5152
JSONRPC_METHOD_NOT_FOUND int = -32601 //nolint: revive
@@ -195,40 +196,44 @@ func (s *Service) handleRPC(w http.ResponseWriter, r *http.Request) {
195196
}
196197

197198
budgetResp := newBudgetWriter(w, MAX_RESPONSE_SIZE)
198-
reqLoop:
199199
for i, rawReq := range reqSeq {
200200

201-
switch r.Context().Err() {
202-
case context.DeadlineExceeded:
203-
s.Logger.Warn("RPC method dispatch timeout")
204-
fallthrough
205-
case context.Canceled:
206-
break reqLoop
207-
}
208-
209201
if i > 0 && !s.writeByte(w, ',') {
210202
return
211203
}
212204

213205
var responded bool
214206
var req RPCRequest
215-
if err := json.Unmarshal(rawReq, &req); err != nil {
216-
responded = s.writeRPCError(w, nil, JSONRPC_INVALID_REQUEST, "invalid request")
217-
} else {
218-
s.Logger.Debug(fmt.Sprintf("Dispatching RPC request: %s", req.Method))
219-
buffer := budgetResp.NewLimitedWriter()
220-
if buffer == nil {
221-
responded = s.writeRPCError(w, req.ID, JSONRPC_INVALID_REQUEST, "response size budget exceeded")
207+
208+
switch r.Context().Err() {
209+
case context.Canceled:
210+
return
211+
case context.DeadlineExceeded:
212+
s.Logger.Warn("RPC method dispatch timeout")
213+
if err := json.Unmarshal(rawReq, &req); err != nil {
214+
responded = s.writeRPCError(w, nil, JSONRPC_INVALID_REQUEST, "invalid request")
215+
} else {
216+
responded = s.writeRPCError(w, req.ID, JSONRPC_RATE_LIMIT_EXCEEDED, "batch deadline exceeded")
217+
}
218+
default:
219+
if err := json.Unmarshal(rawReq, &req); err != nil {
220+
responded = s.writeRPCError(w, nil, JSONRPC_INVALID_REQUEST, "invalid request")
222221
} else {
223-
err := s.dispatchOneRequest(buffer, r, req)
224-
switch {
225-
case err == nil:
226-
responded = s.handleResponseResult(w, buffer.Flush())
227-
case errors.Is(err, io.ErrShortBuffer):
222+
s.Logger.Debug(fmt.Sprintf("Dispatching RPC request: %s", req.Method))
223+
buffer := budgetResp.NewLimitedWriter()
224+
if buffer == nil {
228225
responded = s.writeRPCError(w, req.ID, JSONRPC_INVALID_REQUEST, "response size budget exceeded")
229-
default:
230-
s.Logger.Error("RPC method response encode failed", "method", req.Method, "error", err)
231-
responded = s.writeRPCError(w, req.ID, JSONRPC_INTERNAL_ERROR, "Internal server error")
226+
} else {
227+
err := s.dispatchOneRequest(buffer, r, req)
228+
switch {
229+
case err == nil:
230+
responded = s.handleResponseResult(w, buffer.Flush())
231+
case errors.Is(err, io.ErrShortBuffer):
232+
responded = s.writeRPCError(w, req.ID, JSONRPC_INVALID_REQUEST, "response size budget exceeded")
233+
default:
234+
s.Logger.Error("RPC method response encode failed", "method", req.Method, "error", err)
235+
responded = s.writeRPCError(w, req.ID, JSONRPC_INTERNAL_ERROR, "Internal server error")
236+
}
232237
}
233238
}
234239
}

0 commit comments

Comments
 (0)