|
| 1 | +using System; |
| 2 | +using System.Collections.Generic; |
| 3 | +using System.Diagnostics; |
| 4 | +using System.IO.Pipelines; |
| 5 | +using System.Linq; |
| 6 | +using System.Net; |
| 7 | +using System.Net.Sockets; |
| 8 | +using System.Text; |
| 9 | +using System.Threading; |
| 10 | +using System.Threading.Tasks; |
| 11 | +using AustinHarris.JsonRpc; |
| 12 | +using AustinHarris.JsonRpc.AspNetCore; |
| 13 | +using Microsoft.AspNetCore.Builder; |
| 14 | +using Microsoft.AspNetCore.Connections; |
| 15 | +using Microsoft.AspNetCore.Hosting; |
| 16 | +using Microsoft.Extensions.DependencyInjection; |
| 17 | +using Microsoft.Extensions.Logging; |
| 18 | +using StreamJsonRpc; |
| 19 | +using SjrMethod = StreamJsonRpc.JsonRpcMethodAttribute; |
| 20 | + |
| 21 | +namespace TestServer_Console; |
| 22 | + |
| 23 | +/// <summary> |
| 24 | +/// The same five requests through JSON-RPC.Net and through StreamJsonRpc (Microsoft's JSON-RPC library, the one |
| 25 | +/// behind Visual Studio and the language-server stack), on the same Kestrel TCP listener with the same |
| 26 | +/// pipelining client, so the only variable is the RPC library. Every request carries <c>"jsonrpc":"2.0"</c> |
| 27 | +/// because StreamJsonRpc requires it. StreamJsonRpc has no "document in, document out" call, so its in-process |
| 28 | +/// row runs over a pair of <see cref="Pipe"/>s, the closest thing it has to a direct call; a sequential |
| 29 | +/// proxy row shows what a typical <c>await proxy.AddAsync(1, 2)</c> costs end to end. |
| 30 | +/// </summary> |
| 31 | +internal static class CompareBenchmark |
| 32 | +{ |
| 33 | + private enum Framing { NewLine, Header } |
| 34 | + private enum Formatter { SystemTextJson, Newtonsoft } |
| 35 | + |
| 36 | + private static readonly string[] Requests = BenchmarkRunner.taskInputs.Select(t => "{\"jsonrpc\":\"2.0\"," + t.Substring(1)).ToArray(); |
| 37 | + private static readonly byte[] JsonRpcPrefix = Encoding.UTF8.GetBytes("{\"jsonrpc\":\"2.0\","); |
| 38 | + |
| 39 | + /// <summary>StreamJsonRpc target with the same five methods as <see cref="CalculatorService"/>.</summary> |
| 40 | + public class Target |
| 41 | + { |
| 42 | + [SjrMethod("add")] public double Add(double l, double r) => l + r; |
| 43 | + [SjrMethod("addInt")] public int AddInt(int l, int r) => l + r; |
| 44 | + [SjrMethod("NullableFloatToNullableFloat")] public float? NullableFloatToNullableFloat(float? a) => a; |
| 45 | + [SjrMethod("Test2")] public decimal? Test2(decimal x) => x; |
| 46 | + [SjrMethod("StringMe")] public string StringMe(string x) => x; |
| 47 | + } |
| 48 | + |
| 49 | + /// <summary>Client proxy for the sequential round-trip row.</summary> |
| 50 | + public interface ICalculator |
| 51 | + { |
| 52 | + [SjrMethod("add")] Task<double> AddAsync(double l, double r); |
| 53 | + [SjrMethod("addInt")] Task<int> AddIntAsync(int l, int r); |
| 54 | + [SjrMethod("NullableFloatToNullableFloat")] Task<float?> NullableFloatToNullableFloatAsync(float? a); |
| 55 | + [SjrMethod("Test2")] Task<decimal?> Test2Async(decimal x); |
| 56 | + [SjrMethod("StringMe")] Task<string> StringMeAsync(string x); |
| 57 | + } |
| 58 | + |
| 59 | + internal static async Task RunAsync(Action<string> print, double seconds = 3, int clients = 0, int pipeline = 256) |
| 60 | + { |
| 61 | + print ??= Console.WriteLine; |
| 62 | + if (clients <= 0) clients = Environment.ProcessorCount; |
| 63 | + |
| 64 | + var raw = Requests.Select(r => Encoding.UTF8.GetBytes(r)).ToArray(); |
| 65 | + var newLine = Requests.Select(r => Encoding.UTF8.GetBytes(r + "\n")).ToArray(); |
| 66 | + var header = Requests.Select(r => Encoding.UTF8.GetBytes("Content-Length: " + Encoding.UTF8.GetByteCount(r) + "\r\n\r\n" + r)).ToArray(); |
| 67 | + |
| 68 | + int oursPort = KestrelBenchmark.FreePort(), sjrNewLineStjPort = KestrelBenchmark.FreePort(), sjrHeaderStjPort = KestrelBenchmark.FreePort(), sjrHeaderNewtonsoftPort = KestrelBenchmark.FreePort(); |
| 69 | + var builder = WebApplication.CreateBuilder(); |
| 70 | + builder.Logging.ClearProviders(); |
| 71 | + builder.WebHost.ConfigureKestrel(k => |
| 72 | + { |
| 73 | + k.Listen(IPAddress.Loopback, oursPort, l => l.UseConnectionHandler<JsonRpcConnectionHandler>()); |
| 74 | + k.Listen(IPAddress.Loopback, sjrNewLineStjPort, l => l.Run(c => ServeStreamJsonRpc(c, Framing.NewLine, Formatter.SystemTextJson))); |
| 75 | + k.Listen(IPAddress.Loopback, sjrHeaderStjPort, l => l.Run(c => ServeStreamJsonRpc(c, Framing.Header, Formatter.SystemTextJson))); |
| 76 | + k.Listen(IPAddress.Loopback, sjrHeaderNewtonsoftPort, l => l.Run(c => ServeStreamJsonRpc(c, Framing.Header, Formatter.Newtonsoft))); |
| 77 | + }); |
| 78 | + builder.Services.AddJsonRpc(); |
| 79 | + var app = builder.Build(); |
| 80 | + await app.StartAsync(); |
| 81 | + |
| 82 | + try |
| 83 | + { |
| 84 | + print($"JSON-RPC.Net {typeof(JsonRpcProcessor).Assembly.GetName().Version} vs StreamJsonRpc {typeof(JsonRpc).Assembly.GetName().Version}; Kestrel TCP on loopback, {clients} clients, pipeline {pipeline}, {seconds:0.#} s per row\n"); |
| 85 | + var rows = new List<BenchmarkRunner.ChartRow>(); |
| 86 | + var session = Handler.DefaultSessionId(); |
| 87 | + |
| 88 | + // ---- in-process floors |
| 89 | + var memInputs = raw.Select(i => (ReadOnlyMemory<byte>)i).ToArray(); |
| 90 | + BenchmarkRunner.RunSync(session, memInputs, Config.Serializer, 1, 0.5, out _); |
| 91 | + var elapsed = BenchmarkRunner.RunSync(session, memInputs, Config.Serializer, 1, seconds, out long total); |
| 92 | + rows.Add(new BenchmarkRunner.ChartRow("JSON-RPC.Net in-process, 1 thread", total / elapsed, $"{total,12:N0} RPCs direct call, bytes in, bytes out")); |
| 93 | + print($" JSON-RPC.Net in-process done ({rows[^1].RpcPerSec:N0} RPC/s)"); |
| 94 | + |
| 95 | + PipeRun(newLine, 0.5, pipeline); |
| 96 | + var (count, secs) = PipeRun(newLine, seconds, pipeline); |
| 97 | + rows.Add(new BenchmarkRunner.ChartRow("StreamJsonRpc in-process, 1 client", count / secs, $"{count,12:N0} RPCs Pipe pair, newline framing, STJ formatter, {pipeline} pipelined")); |
| 98 | + print($" StreamJsonRpc in-process done ({count / secs:N0} RPC/s)"); |
| 99 | + |
| 100 | + (count, secs) = await ProxyRun(0.5); |
| 101 | + (count, secs) = await ProxyRun(seconds); |
| 102 | + rows.Add(new BenchmarkRunner.ChartRow("StreamJsonRpc proxy, sequential await", count / secs, $"{count,12:N0} RPCs {secs / count * 1e6:N1} us per round trip")); |
| 103 | + print($" StreamJsonRpc proxy done ({count / secs:N0} RPC/s)"); |
| 104 | + |
| 105 | + // ---- Kestrel TCP, same client for every row |
| 106 | + foreach (var (label, port, inputs, prefix) in new[] |
| 107 | + { |
| 108 | + ("JSON-RPC.Net TCP, raw documents", oursPort, raw, JsonRpcPrefix), |
| 109 | + ("StreamJsonRpc TCP, newline + STJ", sjrNewLineStjPort, newLine, JsonRpcPrefix), |
| 110 | + ("StreamJsonRpc TCP, Content-Length + STJ", sjrHeaderStjPort, header, (byte[])null), |
| 111 | + ("StreamJsonRpc TCP, Content-Length + Json.NET", sjrHeaderNewtonsoftPort, header, (byte[])null), |
| 112 | + }) |
| 113 | + { |
| 114 | + Probe(port, inputs, label); |
| 115 | + KestrelBenchmark.TcpRun(port, inputs, clients, 0.5, pipeline, prefix); |
| 116 | + (count, secs) = KestrelBenchmark.TcpRun(port, inputs, clients, seconds, pipeline, prefix); |
| 117 | + rows.Add(new BenchmarkRunner.ChartRow(label, count / secs, $"{count,12:N0} RPCs")); |
| 118 | + print($" {label} done ({count / secs:N0} RPC/s)"); |
| 119 | + } |
| 120 | + |
| 121 | + BenchmarkRunner.PrintBarChart("JSON-RPC.Net vs StreamJsonRpc - RPC/s", "Library / transport", rows); |
| 122 | + } |
| 123 | + finally |
| 124 | + { |
| 125 | + await app.StopAsync(); |
| 126 | + await app.DisposeAsync(); |
| 127 | + } |
| 128 | + } |
| 129 | + |
| 130 | + private static IJsonRpcMessageHandler CreateHandler(PipeWriter writer, PipeReader reader, Framing framing, Formatter formatter) |
| 131 | + { |
| 132 | + IJsonRpcMessageTextFormatter f = formatter == Formatter.SystemTextJson ? new SystemTextJsonFormatter() : new JsonMessageFormatter(); |
| 133 | + return framing == Framing.NewLine |
| 134 | + ? new NewLineDelimitedMessageHandler(writer, reader, f) |
| 135 | + : new HeaderDelimitedMessageHandler(writer, reader, f); |
| 136 | + } |
| 137 | + |
| 138 | + private static async Task ServeStreamJsonRpc(ConnectionContext connection, Framing framing, Formatter formatter) |
| 139 | + { |
| 140 | + using var rpc = new JsonRpc(CreateHandler(connection.Transport.Output, connection.Transport.Input, framing, formatter)); |
| 141 | + rpc.AddLocalRpcTarget(new Target()); |
| 142 | + rpc.StartListening(); |
| 143 | + try { await rpc.Completion; } |
| 144 | + catch (Exception) { /* the client hung up */ } |
| 145 | + } |
| 146 | + |
| 147 | + /// <summary>Sends each request once and checks every response is a result; a wrong method name or a formatter quirk would otherwise inflate a row.</summary> |
| 148 | + private static void Probe(int port, byte[][] inputs, string label) |
| 149 | + { |
| 150 | + using var socket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp) { NoDelay = true, ReceiveTimeout = 5000 }; |
| 151 | + socket.Connect(IPAddress.Loopback, port); |
| 152 | + foreach (var doc in inputs) socket.Send(doc); |
| 153 | + var counter = new KestrelBenchmark.ResponseCounter(null); |
| 154 | + var received = new List<byte>(); |
| 155 | + var buffer = new byte[16 * 1024]; |
| 156 | + int docs = 0; |
| 157 | + while (docs < inputs.Length) |
| 158 | + { |
| 159 | + int n = socket.Receive(buffer); |
| 160 | + if (n == 0) break; |
| 161 | + received.AddRange(buffer.AsSpan(0, n).ToArray()); |
| 162 | + docs += counter.Count(buffer.AsSpan(0, n)); |
| 163 | + } |
| 164 | + var text = Encoding.UTF8.GetString(received.ToArray()); |
| 165 | + if (docs != inputs.Length || text.Contains("\"error\"") || !text.Contains("\"result\"")) |
| 166 | + throw new InvalidOperationException($"{label}: unexpected responses:\n{text}"); |
| 167 | + } |
| 168 | + |
| 169 | + /// <summary> |
| 170 | + /// StreamJsonRpc served over an in-process <see cref="Pipe"/> pair, driven exactly like the TCP client: a ring |
| 171 | + /// of newline-framed requests, a cursor, and refills whenever the pipeline has room. |
| 172 | + /// </summary> |
| 173 | + private static (long count, double seconds) PipeRun(byte[][] inputs, double seconds, int pipeline) |
| 174 | + { |
| 175 | + var toServer = new Pipe(); |
| 176 | + var toClient = new Pipe(); |
| 177 | + using var rpc = new JsonRpc(CreateHandler(toClient.Writer, toServer.Reader, Framing.NewLine, Formatter.SystemTextJson)); |
| 178 | + rpc.AddLocalRpcTarget(new Target()); |
| 179 | + rpc.StartListening(); |
| 180 | + |
| 181 | + const int ringCount = 4096; |
| 182 | + var offsets = new int[ringCount + pipeline + 1]; |
| 183 | + var ring = KestrelBenchmark.BuildRing(inputs, ringCount, pipeline, offsets); |
| 184 | + var counter = new KestrelBenchmark.ResponseCounter(JsonRpcPrefix); |
| 185 | + int cursor = 0, inFlight = 0; |
| 186 | + long received = 0; |
| 187 | + var sw = Stopwatch.StartNew(); |
| 188 | + while (sw.Elapsed.TotalSeconds < seconds) |
| 189 | + { |
| 190 | + int free = pipeline - inFlight; |
| 191 | + if (free > 0) |
| 192 | + { |
| 193 | + toServer.Writer.WriteAsync(new ReadOnlyMemory<byte>(ring, offsets[cursor], offsets[cursor + free] - offsets[cursor])).AsTask().GetAwaiter().GetResult(); |
| 194 | + cursor += free; |
| 195 | + if (cursor >= ringCount) cursor -= ringCount; |
| 196 | + inFlight += free; |
| 197 | + } |
| 198 | + var result = toClient.Reader.ReadAsync().AsTask().GetAwaiter().GetResult(); |
| 199 | + foreach (var segment in result.Buffer) |
| 200 | + { |
| 201 | + int docs = counter.Count(segment.Span); |
| 202 | + received += docs; |
| 203 | + inFlight -= docs; |
| 204 | + } |
| 205 | + toClient.Reader.AdvanceTo(result.Buffer.End); |
| 206 | + if (result.IsCompleted) break; |
| 207 | + } |
| 208 | + // Drain so the server is idle before the next row. |
| 209 | + while (inFlight > 0) |
| 210 | + { |
| 211 | + var result = toClient.Reader.ReadAsync().AsTask().GetAwaiter().GetResult(); |
| 212 | + foreach (var segment in result.Buffer) |
| 213 | + { |
| 214 | + int docs = counter.Count(segment.Span); |
| 215 | + received += docs; |
| 216 | + inFlight -= docs; |
| 217 | + } |
| 218 | + toClient.Reader.AdvanceTo(result.Buffer.End); |
| 219 | + if (result.IsCompleted) break; |
| 220 | + } |
| 221 | + sw.Stop(); |
| 222 | + toServer.Writer.Complete(); |
| 223 | + return (received, sw.Elapsed.TotalSeconds); |
| 224 | + } |
| 225 | + |
| 226 | + /// <summary>A typed StreamJsonRpc proxy awaiting one call at a time over an in-process pipe pair: the usual way the library is used.</summary> |
| 227 | + private static async Task<(long count, double seconds)> ProxyRun(double seconds) |
| 228 | + { |
| 229 | + var toServer = new Pipe(); |
| 230 | + var toClient = new Pipe(); |
| 231 | + using var server = new JsonRpc(CreateHandler(toClient.Writer, toServer.Reader, Framing.Header, Formatter.SystemTextJson)); |
| 232 | + server.AddLocalRpcTarget(new Target()); |
| 233 | + server.StartListening(); |
| 234 | + using var client = new JsonRpc(CreateHandler(toServer.Writer, toClient.Reader, Framing.Header, Formatter.SystemTextJson)); |
| 235 | + var proxy = client.Attach<ICalculator>(); |
| 236 | + client.StartListening(); |
| 237 | + |
| 238 | + long n = 0; |
| 239 | + var sw = Stopwatch.StartNew(); |
| 240 | + while (sw.Elapsed.TotalSeconds < seconds) |
| 241 | + { |
| 242 | + if (await proxy.AddAsync(1, 2) != 3) throw new InvalidOperationException("add"); |
| 243 | + if (await proxy.AddIntAsync(1, 7) != 8) throw new InvalidOperationException("addInt"); |
| 244 | + if (await proxy.NullableFloatToNullableFloatAsync(1.23f) != 1.23f) throw new InvalidOperationException("NullableFloatToNullableFloat"); |
| 245 | + if (await proxy.Test2Async(3.456m) != 3.456m) throw new InvalidOperationException("Test2"); |
| 246 | + if (await proxy.StringMeAsync("Foo") != "Foo") throw new InvalidOperationException("StringMe"); |
| 247 | + n += 5; |
| 248 | + } |
| 249 | + sw.Stop(); |
| 250 | + toServer.Writer.Complete(); |
| 251 | + toClient.Writer.Complete(); |
| 252 | + return (n, sw.Elapsed.TotalSeconds); |
| 253 | + } |
| 254 | +} |
0 commit comments