Skip to content

Commit 48839fb

Browse files
authored
Merge pull request #154 from Astn/async-scratch-cache
ProcessAsync: per-thread scratch cache, scaling gate, async rows, frame fix
2 parents 805c64f + 7c83fd3 commit 48839fb

36 files changed

Lines changed: 1402 additions & 359 deletions
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
# Every `lock (`, `Interlocked.`, `Volatile.Write`, `[ThreadStatic]` field and writable static field in the
2+
# request-path files (core and both companion serializers), with the reason it is allowed there.
3+
# Format: file | the code line, trimmed, without its trailing comment | reason, starting with its tag:
4+
# per-thread the field is thread-static, so no line is shared
5+
# miss-path touched only when a per-thread cache misses or overflows, never by a warm inline document
6+
# registration-only written when methods or sessions are bound, never by a request
7+
# read-only-after-init written once at startup, then only read
8+
# shared-write a request can write it and every core reads it: the reason must say why it is tolerated
9+
# Checked by .github/scripts/check_request_path_sync.py in the pull-request workflow. The measurement that backs this
10+
# list is `TestServer_Console --scale` (README, Benchmarks). A `[ThreadStatic]` attribute on its own line is not
11+
# listed; the field line under it is.
12+
Json-Rpc/JsonRpcProcessor.Async.cs | private static int _count; | miss-path: the shared pool's fill count, read and written only under lock (Pool)
13+
Json-Rpc/JsonRpcProcessor.Async.cs | [ThreadStatic] private static AsyncScratch _slot; | per-thread: the one-slot cache of an idle async scratch
14+
Json-Rpc/JsonRpcProcessor.Async.cs | lock (Pool) | miss-path: AsyncScratch.Rent, taken only when the thread's slot is empty (cold thread, nested or overlapping document)
15+
Json-Rpc/JsonRpcProcessor.Async.cs | lock (Pool) | miss-path: AsyncScratch.Return, taken only when the completing thread's slot is already occupied; the array is the bounded overflow
16+
Json-Rpc/JsonRpcProcessor.cs | [ThreadStatic] private static Scratch _current; | per-thread: the synchronous scratch
17+
Json-Rpc/Handler.Async.cs | [ThreadStatic] private static InvocationState __unflowedState; | per-thread: the invocation frame of a RpcContextFlow.None call
18+
Json-Rpc/Handler.cs | private static int _sessionHandlerMasterVersion = 1; | registration-only: written by Interlocked.Increment when a session is created or destroyed; requests read it to validate their per-thread snapshot
19+
Json-Rpc/Handler.cs | Interlocked.Increment(ref _sessionHandlerMasterVersion); | registration-only: GetSessionHandler creating a session and DestroySession; never on the request path
20+
Json-Rpc/Handler.cs | private static Dictionary<string, Handler> _sessionHandlersLocal; | per-thread ([ThreadStatic] on the line above): the thread's snapshot of the session registry
21+
Json-Rpc/Handler.cs | private static int _sessionHandlerLocalVersion = 0; | per-thread ([ThreadStatic] on the line above): the snapshot's version
22+
Json-Rpc/Handler.cs | private static string _lastSessionId; | per-thread ([ThreadStatic] on the line above): the last-session cache
23+
Json-Rpc/Handler.cs | private static Handler _lastSessionHandler; | per-thread ([ThreadStatic] on the line above): the last-session cache
24+
Json-Rpc/Handler.cs | private static InvocationState __state; | per-thread ([ThreadStatic] on the line above): the current invocation frame
25+
Json-Rpc/Jsmn/JsmnSerializer.cs | [ThreadStatic] private static JsmnTokenizer _scratch; | per-thread: the tokenizer scratch
26+
Json-Rpc/Jsmn/JsmnSerializer.cs | [ThreadStatic] private static bool _scratchInUse; | per-thread: re-entrancy flag for the tokenizer scratch
27+
Json-Rpc/Serialization/Utf8KeyTable.cs | lock (_writeLock) | registration-only: Set, Remove, Clear and ReplaceAll rebuild a snapshot under the lock; requests read the published snapshot without it
28+
Json-Rpc/Serialization/Utf8KeyTable.cs | System.Threading.Volatile.Write(ref _snapshot, snapshot); | registration-only: publishes a rebuilt snapshot; read-only afterwards
29+
Json-Rpc/RpcBinding.cs | lock (_sync) | registration-only: RpcBinding.Dispose unbinding an interface tree
30+
AustinHarris.JsonRpc.SystemTextJson/SystemTextJsonRpcSerializer.cs | [ThreadStatic] private static Utf8JsonWriter _cachedWriter; | per-thread: the cached writer
31+
AustinHarris.JsonRpc.SystemTextJson/SystemTextJsonRpcSerializer.cs | [ThreadStatic] private static JsonSerializerOptions _cachedWriterOptions; | per-thread: the options the cached writer was built with
32+
AustinHarris.JsonRpc.SystemTextJson/SystemTextJsonRpcSerializer.cs | [ThreadStatic] private static bool _cachedWriterInUse; | per-thread: re-entrancy flag for the cached writer
33+
AustinHarris.JsonRpc.SystemTextJson/SystemTextJsonRpcSerializer.cs | [ThreadStatic] private static byte[] _scratch; | per-thread: transcoding scratch
34+
AustinHarris.JsonRpc.SystemTextJson/SystemTextJsonRpcSerializer.cs | private static Entry _last; | shared-write: TypeInfo<T>'s last-options cache, rewritten whenever the options differ from the previous call, so a per-request write when two serializers with different options serve requests concurrently; tolerated because one options object is the common case; the per-thread copy is the first candidate of the 2026-09-25 P6 resolution, pending its A/B
35+
AustinHarris.JsonRpc.Newtonsoft/NewtonsoftJsonRpcSerializer.cs | [ThreadStatic] private static Scratch _current; | per-thread: the Json.NET scratch
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
#!/usr/bin/env python3
2+
"""Every synchronization point and every writable static on the request path must be listed, with a reason, in
3+
.github/request-path-sync.allowlist: `lock (`, `Interlocked.`, `Volatile.Write`, `[ThreadStatic]` fields and static
4+
fields that are not readonly, in the request-path files of the core and of both companion serializers.
5+
6+
A process-wide serialization point on the per-document path caps ProcessAsync at a few million requests per second
7+
on every core count (the AsyncScratch pool lock, 2026-09-25), and the single-threaded micro-benchmarks cannot see
8+
it; a static that one request writes and every core reads costs the same kind of cache-line traffic without a lock.
9+
This check is a review aid: it catches a new lock, atomic or writable static on those files and asks for a written
10+
reason (per-thread, miss-path, registration-only, read-only-after-init). It is not the measurement;
11+
`TestServer_Console --scale` is (a probe moved behind an allowed lock would pass here and fail there). Mutable
12+
objects reached through readonly references are outside its reach and belong to review.
13+
14+
Usage: python3 .github/scripts/check_request_path_sync.py (from the repository root; exit 1 on any unlisted use)
15+
"""
16+
import glob
17+
import os
18+
import re
19+
import sys
20+
21+
ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
22+
ALLOWLIST = os.path.join(ROOT, ".github", "request-path-sync.allowlist")
23+
PATTERNS = ["Json-Rpc/JsonRpcProcessor*.cs", "Json-Rpc/Handler*.cs", "Json-Rpc/Invocation/*.cs", "Json-Rpc/Jsmn/*.cs",
24+
"Json-Rpc/Serialization/*.cs", "Json-Rpc/JsonRpcContext.cs", "Json-Rpc/JsonRpcRequestId.cs", "Json-Rpc/RpcBinding.cs",
25+
"AustinHarris.JsonRpc.SystemTextJson/*.cs", "AustinHarris.JsonRpc.Newtonsoft/*.cs"]
26+
USE = re.compile(r"\block\s*\(|\bInterlocked\.|\bVolatile\.Write")
27+
# A static field declaration that is not readonly or const: `[ThreadStatic] private static T name;` or `static T name = ...;`.
28+
# Expression-bodied members (`=>`), methods, classes, events and operators are not fields.
29+
FIELD = re.compile(r"^(?:\[ThreadStatic\]\s*)?(?:(?:public|private|internal|protected|new|volatile)\s+)*static\s+"
30+
r"(?!readonly\b|const\b|class\b|void\b|partial\b|event\b|explicit\b|implicit\b|operator\b)"
31+
r"(?!.*=>)(?!.*\()[\w<>\[\],.?\s]+?\s+\w+\s*(?:=[^;]*)?;$")
32+
33+
34+
def load_allowlist():
35+
"""{(file, code line stripped): reason}; lines are `file | code | reason`, `#` comments and blanks ignored."""
36+
allowed = {}
37+
with open(ALLOWLIST, encoding="utf-8") as f:
38+
for n, line in enumerate(f, 1):
39+
line = line.strip()
40+
if not line or line.startswith("#"):
41+
continue
42+
parts = [p.strip() for p in line.split("|", 2)]
43+
if len(parts) != 3 or not all(parts):
44+
sys.exit(f"{ALLOWLIST}:{n}: expected `file | code | reason`")
45+
allowed[(parts[0].replace("\\", "/"), parts[1])] = parts[2]
46+
return allowed
47+
48+
49+
def main():
50+
allowed = load_allowlist()
51+
seen = set()
52+
problems = []
53+
for pattern in PATTERNS:
54+
for path in sorted(glob.glob(os.path.join(ROOT, pattern))):
55+
rel = os.path.relpath(path, ROOT).replace("\\", "/")
56+
with open(path, encoding="utf-8-sig") as f:
57+
for n, line in enumerate(f, 1):
58+
code = line.split("//", 1)[0].strip()
59+
if not (USE.search(code) or FIELD.match(code)):
60+
continue
61+
key = (rel, code)
62+
if key in allowed:
63+
seen.add(key)
64+
else:
65+
problems.append(f"{rel}:{n}: `{code}` is not in {os.path.relpath(ALLOWLIST, ROOT)}; add it with a reason or take it off the request path")
66+
for key in allowed:
67+
if key not in seen:
68+
problems.append(f"{os.path.relpath(ALLOWLIST, ROOT)}: `{key[1]}` in {key[0]} no longer exists; remove the entry")
69+
for p in problems:
70+
print(p)
71+
if problems:
72+
return 1
73+
print(f"request-path synchronization: {len(seen)} listed uses, none unlisted")
74+
return 0
75+
76+
77+
if __name__ == "__main__":
78+
sys.exit(main())

‎.github/workflows/build_pull_request.yml‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,3 +34,41 @@ jobs:
3434
run: dotnet test AustinHarris.JsonRpcTestN --configuration Release --no-build
3535
# Pull requests do not publish. Packages go out from the master workflow through Trusted Publishing, whose
3636
# nuget.org policy is bound to build_publish_master.yml; nothing else can push.
37+
38+
# Every lock and Interlocked on the request-path files must be listed with a reason (a process-wide lock on the
39+
# per-document path once capped ProcessAsync at 4 M RPC/s on every core count). A review aid, not the measurement.
40+
request-path-sync:
41+
runs-on: ubuntu-latest
42+
steps:
43+
- uses: actions/checkout@v7
44+
- name: Synchronization on the request path is allowlisted
45+
run: python3 .github/scripts/check_request_path_sync.py
46+
47+
# The measurement: ProcessAsync must scale from 1 to 4 workers. A process-wide serialization point holds the ratio
48+
# near 1.3 on any core count. Diagnostic on the shared runner (its core count and isolation are not promised, so
49+
# this job is not required and continues on error); the release gate is `--scale 3 16 4.0` on the reference machine.
50+
scaling:
51+
runs-on: ubuntu-latest
52+
continue-on-error: true
53+
steps:
54+
- uses: actions/checkout@v7
55+
- name: Setup .NET
56+
uses: actions/setup-dotnet@v6
57+
with:
58+
global-json-file: global.json
59+
dotnet-version: 10.0.x
60+
- name: Build the harness
61+
run: dotnet build TestServer_Console --configuration Release
62+
- name: ProcessAsync scales from 1 to 4 workers (4/1 at least 2.0)
63+
run: |
64+
set +e
65+
dotnet run -c Release --no-build --project TestServer_Console -- --scale 3 4 2.0 | tee scale.txt
66+
status=${PIPESTATUS[0]}
67+
{
68+
echo "## ProcessAsync scaling (diagnostic, not required)"
69+
echo
70+
grep -E '^\|' scale.txt
71+
echo
72+
[ "$status" -eq 0 ] && echo "pass: 4/1 at least 2.0" || echo "**flag: 4/1 below 2.0 on this runner; reproduce with --scale 3 16 4.0 on the reference machine before reading it as a regression**"
73+
} >> "$GITHUB_STEP_SUMMARY"
74+
exit $status

‎AustinHarris.JsonRpc.AspNetCore/README.md‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -156,9 +156,15 @@ invoked.
156156

157157
- **HTTP:** the call is cancelled when the client disconnects (`HttpContext.RequestAborted`). Notifications are
158158
awaited and still answer `204`. The body reader stays leased until the invocation finishes.
159-
- **Raw connections:** documents are processed one at a time, in order. Replies already finished are flushed
160-
before the connection waits on a slow method. When the connection closes, the running method is waited for and
161-
its response discarded.
159+
- **Raw connections:** documents are processed one at a time, in order. On one connection, 256 pipelined requests
160+
are 256 sequential invocations, not 256 concurrent suspensions. Concurrency comes from connections.
161+
Replies already finished are flushed before the connection waits on a slow method. When the connection closes,
162+
the handler waits for the running method to finish and discards its response.
163+
- **Cost:** every document then goes through `ProcessAsync`. With methods that complete inline, the TCP row measured
164+
15.0 M to 15.4 M against 14.3 M to 16.5 M for the synchronous mode on the same day, inside its spread. A method that
165+
suspends pays for its own async state, the library's completion state and a continuation per request; the main README's
166+
Async table measured 559 B per request for the yielding None row in process at one worker, including the service's own
167+
allocations. The main README's Kestrel table has both rows, measured with `TestServer_Console --kestrel 3 async`.
162168

163169
A method receives the token by declaring a `[JsonRpcCancellation] CancellationToken` parameter; see
164170
[Asynchronous methods and cancellation](https://github.com/Astn/JSON-RPC.NET#asynchronous-methods-and-cancellation)

‎AustinHarris.JsonRpcTestN/AsyncInvocationTests.cs‎

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,12 @@ public sealed class AsyncInvocationTests
2424
[TearDown] public void TearDown() => Handler.DestroySession(_session);
2525
private void Bind(string name, Delegate method, RpcContextFlow flow = RpcContextFlow.Flow) => ServiceBinder.BindMethod(_session, name, method, contextFlow: flow);
2626
private Task<string> Run(string json, JsonRpcSerializer serializer = null, object context = null, CancellationToken token = default) => JsonRpcProcessor.ProcessAsync(_session, json, context, serializer, token);
27+
private string Sync(string json, object context = null)
28+
{
29+
var output = new ArrayBufferWriter<byte>();
30+
JsonRpcProcessor.Process(_session, Encoding.UTF8.GetBytes(json).AsSpan(), output, context);
31+
return Encoding.UTF8.GetString(output.WrittenSpan);
32+
}
2733
private static TaskCompletionSource<int> Gate() => new TaskCompletionSource<int>(TaskCreationOptions.RunContinuationsAsynchronously);
2834
private static string Request(string method, string parameters = null, string id = "1") => "{\"method\":\"" + method + "\"" + (parameters == null ? "" : ",\"params\":" + parameters) + (id == null ? "" : ",\"id\":" + id) + "}";
2935
private static void Error(string json, int code) => Assert.AreEqual(code, (int)JObject.Parse(json)["error"]["code"], json);
@@ -309,6 +315,51 @@ public async Task FlowFrame_IsClearedForEscapedExecutionContext()
309315
await escaped;
310316
}
311317

318+
[Test]
319+
public async Task FlowScope_CompletedOnAnotherThread_LeavesThatThreadItsOwnFrame()
320+
{
321+
// A hooked Flow invocation suspends here and completes on a pool thread, where its scope is
322+
// disposed. That thread must keep its own frame afterwards: when it was handed this thread's
323+
// frame instead, a synchronous dispatch on each thread saved and restored the same frame, and
324+
// the interleaving below left this thread's method reading no id at all.
325+
var handler = Handler.GetSessionHandler(_session);
326+
handler.SetPreProcessHandler((request, context) => null);
327+
int completedOn = 0;
328+
handler.SetPostProcessHandler((request, response, context) => { if (request.Method == "suspend") completedOn = Environment.CurrentManagedThreadId; return null; });
329+
var suspended = new TaskCompletionSource<int>();
330+
Bind("suspend", new Func<Task<int>>(() => suspended.Task));
331+
var entered = new ManualResetEventSlim();
332+
var release = new ManualResetEventSlim();
333+
var left = new ManualResetEventSlim();
334+
Bind("hold", new Func<int>(() => { entered.Set(); release.Wait(); return 1; }));
335+
var mine = new object();
336+
Bind("peek", new Func<int>(() =>
337+
{
338+
release.Set();
339+
left.Wait();
340+
return (Handler.RpcRequestId().IsAbsent ? 0 : 1) + (ReferenceEquals(Handler.RpcContext(), mine) ? 2 : 0);
341+
}));
342+
Bind("warm", new Func<int>(() => 1));
343+
Sync(Request("warm")); // this thread owns a frame before the suspension, as any thread that has dispatched does
344+
var pending = Run(Request("suspend"));
345+
int completingThread = 0;
346+
var other = Task.Run(() =>
347+
{
348+
completingThread = Environment.CurrentManagedThreadId;
349+
suspended.SetResult(7);
350+
Sync(Request("hold", id: "2"));
351+
left.Set();
352+
});
353+
// Block, do not await, until the other thread is inside "hold": an awaited continuation could be
354+
// run on that thread, ahead of "hold", and wait for itself.
355+
entered.Wait();
356+
var response = Sync(Request("peek", id: "\"mine\""), mine);
357+
await other;
358+
Assert.AreEqual("{\"jsonrpc\":\"2.0\",\"result\":7,\"id\":1}", await pending);
359+
Assume.That(completedOn, Is.EqualTo(completingThread), "the completion did not run the scope's cleanup on the completing thread");
360+
Assert.AreEqual("{\"jsonrpc\":\"2.0\",\"result\":3,\"id\":\"mine\"}", response, "this thread's method sees its own id and context while the other thread dispatches");
361+
}
362+
312363
[TestCaseSource(nameof(Serializers))]
313364
public async Task Batch_IsSequential_AwaitsNotifications_AndIsolatesFaults(string name)
314365
{

0 commit comments

Comments
 (0)