|
| 1 | +using System; |
| 2 | +using System.Buffers; |
| 3 | +using System.Threading.Tasks; |
| 4 | +using AustinHarris.JsonRpc.Serialization; |
| 5 | +using Microsoft.AspNetCore.Connections; |
| 6 | +using Microsoft.Extensions.Options; |
| 7 | + |
| 8 | +namespace AustinHarris.JsonRpc.AspNetCore |
| 9 | +{ |
| 10 | + /// <summary> |
| 11 | + /// JSON-RPC over a raw Kestrel connection (TCP, Unix socket, named pipe): clients send JSON documents back to |
| 12 | + /// back (optionally whitespace / newline separated) and receive responses in order. Wire it up with |
| 13 | + /// <c>kestrel.ListenLocalhost(port, l => l.UseConnectionHandler<JsonRpcConnectionHandler>())</c>. |
| 14 | + /// The connection's <see cref="ConnectionContext"/> is the RPC context for every call. |
| 15 | + /// </summary> |
| 16 | + public partial class JsonRpcConnectionHandler : ConnectionHandler |
| 17 | + { |
| 18 | + private readonly JsonRpcOptions _options; |
| 19 | + |
| 20 | + public JsonRpcConnectionHandler(IOptions<JsonRpcOptions> options) |
| 21 | + { |
| 22 | + _options = options?.Value ?? new JsonRpcOptions(); |
| 23 | + } |
| 24 | + |
| 25 | + /// <summary>Processes a connection using the hosting mode selected in options.</summary> |
| 26 | + public override Task OnConnectedAsync(ConnectionContext connection) |
| 27 | + { |
| 28 | + return _options.EnableAsyncMethods ? RunAsynchronousMethodsAsync(connection) : RunSynchronousMethodsAsync(connection); |
| 29 | + } |
| 30 | + |
| 31 | + private async Task RunSynchronousMethodsAsync(ConnectionContext connection) |
| 32 | + { |
| 33 | + var input = connection.Transport.Input; |
| 34 | + var output = connection.Transport.Output; |
| 35 | + string session = _options.SessionId ?? Handler.DefaultSessionId(); |
| 36 | + |
| 37 | + while (true) |
| 38 | + { |
| 39 | + var result = await input.ReadAsync(connection.ConnectionClosed).ConfigureAwait(false); |
| 40 | + var buffer = result.Buffer; |
| 41 | + bool wrote = false; |
| 42 | + |
| 43 | + while (JsonFramer.TryReadDocument(ref buffer, out var document)) |
| 44 | + { |
| 45 | + // the framer hands back a one-byte slice for anything that is not '{' or '[': drop it |
| 46 | + if (document.Length > 1 || document.First.Span[0] == (byte)'{' || document.First.Span[0] == (byte)'[') |
| 47 | + { |
| 48 | + if (document.Length > _options.MaxRequestBytes) |
| 49 | + { |
| 50 | + connection.Abort(new ConnectionAbortedException("JSON-RPC document exceeds MaxRequestBytes.")); |
| 51 | + return; |
| 52 | + } |
| 53 | + JsonRpcProcessor.Process(session, in document, output, connection, _options.Serializer); |
| 54 | + wrote = true; |
| 55 | + } |
| 56 | + } |
| 57 | + |
| 58 | + if (wrote) |
| 59 | + { |
| 60 | + var flush = await output.FlushAsync(connection.ConnectionClosed).ConfigureAwait(false); |
| 61 | + if (flush.IsCompleted || flush.IsCanceled) break; |
| 62 | + } |
| 63 | + |
| 64 | + if (result.IsCompleted || result.IsCanceled) |
| 65 | + { |
| 66 | + input.AdvanceTo(buffer.Start, buffer.End); |
| 67 | + break; |
| 68 | + } |
| 69 | + if (buffer.Length > _options.MaxRequestBytes) |
| 70 | + { |
| 71 | + connection.Abort(new ConnectionAbortedException("JSON-RPC document exceeds MaxRequestBytes.")); |
| 72 | + return; |
| 73 | + } |
| 74 | + // consumed up to the last complete document, examined everything |
| 75 | + input.AdvanceTo(buffer.Start, buffer.End); |
| 76 | + } |
| 77 | + } |
| 78 | + } |
| 79 | +} |
0 commit comments