From b4577f5ce9289fcea0d0c65dea17cee64ca97235 Mon Sep 17 00:00:00 2001 From: Nicholas Blumhardt Date: Thu, 27 Oct 2016 18:10:54 +1000 Subject: [PATCH 1/2] Work-in-progress on WebSocket streaming support --- example/SeqTail/Program.cs | 125 +++-------- .../SeqTail/Properties/launchSettings.json | 8 + example/SeqTail/project.json | 5 +- src/Seq.Api/Api/Client/SeqApiClient.cs | 30 ++- .../Api/ResourceGroups/EventsResourceGroup.cs | 43 ++++ src/Seq.Api/Api/Streams/ObservableStream.cs | 201 ++++++++++++++++++ src/Seq.Api/Seq.Api.xproj | 6 +- src/Seq.Api/project.json | 7 +- 8 files changed, 325 insertions(+), 100 deletions(-) create mode 100644 example/SeqTail/Properties/launchSettings.json create mode 100644 src/Seq.Api/Api/Streams/ObservableStream.cs diff --git a/example/SeqTail/Program.cs b/example/SeqTail/Program.cs index afc06e7..63ed136 100644 --- a/example/SeqTail/Program.cs +++ b/example/SeqTail/Program.cs @@ -1,11 +1,13 @@ using System; -using System.Collections.Generic; using System.Linq; -using System.Threading; using System.Threading.Tasks; using DocoptNet; using Seq.Api; -using Seq.Api.Model.Events; +using Serilog; +using System.Reactive.Linq; +using Serilog.Formatting.Compact.Reader; +using System.IO; +using Serilog.Events; namespace SeqTail { @@ -27,36 +29,26 @@ class Program static void Main(string[] args) { - var cts = new CancellationTokenSource(); + Log.Logger = new LoggerConfiguration() + .MinimumLevel.Verbose() + .WriteTo.LiterateConsole() + .CreateLogger(); - // ReSharper disable once MethodSupportsCancellation - var tail = Task.Run(async () => + try { - try - { - var arguments = new Docopt().Apply(Usage, args, version: "Seq Tail 0.2", exit: true); - - var server = arguments[""].ToString(); - var apiKey = Normalize(arguments["--apikey"]); - var filter = Normalize(arguments["--filter"]); - var window = arguments["--window"].AsInt; + var arguments = new Docopt().Apply(Usage, args, version: "Seq Tail 0.2", exit: true); - await Run(server, apiKey, filter, window, cts.Token); - } - catch (Exception ex) - { - Console.ForegroundColor = ConsoleColor.White; - Console.BackgroundColor = ConsoleColor.Red; - Console.WriteLine("seq-tail: {0}", ex); - Console.ResetColor(); - Environment.Exit(-1); - } - }); + var server = arguments[""].ToString(); + var apiKey = Normalize(arguments["--apikey"]); + var filter = Normalize(arguments["--filter"]); - Console.ReadKey(true); - cts.Cancel(); - // ReSharper disable once MethodSupportsCancellation - tail.Wait(); + Run(server, apiKey, filter).GetAwaiter().GetResult(); + } + catch (Exception ex) + { + Log.Fatal(ex, "Tailing aborted"); + Environment.Exit(-1); + } } static string Normalize(ValueObject v) @@ -66,10 +58,8 @@ static string Normalize(ValueObject v) return string.IsNullOrWhiteSpace(s) ? null : s; } - static async Task Run(string server, string apiKey, string filter, int window, CancellationToken cancel) + static async Task Run(string server, string apiKey, string filter) { - var startedAt = DateTime.UtcNow; - var connection = new SeqConnection(server, apiKey); string strict = null; @@ -79,68 +69,19 @@ static async Task Run(string server, string apiKey, string filter, int window, C strict = converted.StrictExpression; } - var result = await connection.Events.ListAsync(count: window, render: true, fromDateUtc: startedAt, filter: strict); - - // Since results may come late, we request an overlapping window and exclude - // events that have already been printed. If the last seen ID wasn't returned - // we assume the count was too small to cover the window. - var lastPrintedBatch = new HashSet(); - string lastReturnedId = null; - - while (!cancel.IsCancellationRequested) + using (var stream = await connection.Events.StreamDocumentsAsync(filter: strict)) { - if (result.Count == 0) - { - await Task.Delay(TimeSpan.FromSeconds(1)); - } - else + var subscription = stream.Subscribe(document => { - var noOverlap = result.All(e => e.Id != lastReturnedId); - - if (noOverlap && lastReturnedId != null) - Console.WriteLine(""); - - foreach (var eventEntity in ((IEnumerable)result).Reverse()) - { - if (lastPrintedBatch.Contains(eventEntity.Id)) - { - continue; - } - - lastReturnedId = eventEntity.Id; - - var exception = ""; - if (eventEntity.Exception != null) - exception = Environment.NewLine + eventEntity.Exception; - - var ts = DateTimeOffset.Parse(eventEntity.Timestamp).ToLocalTime(); - - var color = ConsoleColor.White; - switch (eventEntity.Level) - { - case "Verbose": - case "Debug": - color = ConsoleColor.Gray; - break; - case "Warning": - color = ConsoleColor.Yellow; - break; - case "Error": - case "Fatal": - color = ConsoleColor.Red; - break; - } - - Console.ForegroundColor = color; - Console.WriteLine("{0:G} [{1}] {2}{3}", ts, eventEntity.Level, eventEntity.RenderedMessage, exception); - Console.ResetColor(); - } - - lastPrintedBatch = new HashSet(result.Select(e => e.Id)); - } - - var fromDateUtc = lastReturnedId == null ? startedAt : DateTime.UtcNow.AddMinutes(-3); - result = await connection.Events.ListAsync(count: window, render: true, fromDateUtc: fromDateUtc, filter: strict); + var reader = new LogEventReader(new StringReader(document)); + LogEvent evt; + if (!reader.TryRead(out evt)) + throw new InvalidOperationException("Expected document to contain data."); + Log.Write(evt); + }); + + Console.ReadKey(true); + subscription.Dispose(); } } } diff --git a/example/SeqTail/Properties/launchSettings.json b/example/SeqTail/Properties/launchSettings.json new file mode 100644 index 0000000..33fcc88 --- /dev/null +++ b/example/SeqTail/Properties/launchSettings.json @@ -0,0 +1,8 @@ +{ + "profiles": { + "seq-tail": { + "executablePath": "C:\\Development\\nblumhardt\\seq-api\\example\\SeqTail\\bin\\Debug\\net46\\win7-x64\\seq-tail.exe", + "commandLineArgs": "http://localhost:5341 --apikey=fSCcBOyOfttZ0kiZFgq" + } + } +} \ No newline at end of file diff --git a/example/SeqTail/project.json b/example/SeqTail/project.json index 3a5ff21..00d5aab 100644 --- a/example/SeqTail/project.json +++ b/example/SeqTail/project.json @@ -7,7 +7,10 @@ "dependencies": { "Seq.Api": { "target": "project" }, - "docopt.net": "0.6.1.9" + "docopt.net": "0.6.1.9", + "Serilog.Formatting.Compact.Reader": "1.0.0-dev-00004", + "Serilog.Sinks.Literate": "2.0.0", + "System.Reactive": "3.0.0" }, "frameworks": { diff --git a/src/Seq.Api/Api/Client/SeqApiClient.cs b/src/Seq.Api/Api/Client/SeqApiClient.cs index 93051b5..73ddc24 100644 --- a/src/Seq.Api/Api/Client/SeqApiClient.cs +++ b/src/Seq.Api/Api/Client/SeqApiClient.cs @@ -13,6 +13,9 @@ using Seq.Api.Model.Root; using Seq.Api.Serialization; using Tavis.UriTemplates; +using System.Threading; +using Seq.Api.Streams; +using System.Net.WebSockets; namespace Seq.Api.Client { @@ -26,6 +29,7 @@ public class SeqApiClient : IDisposable const string SeqApiV3MediaType = "application/vnd.continuousit.seq.v3+json"; readonly HttpClient _httpClient; + readonly CookieContainer _cookies = new CookieContainer(); readonly JsonSerializer _serializer = JsonSerializer.Create( new JsonSerializerSettings { @@ -41,7 +45,7 @@ public SeqApiClient(string serverUrl, string apiKey = null) if (!string.IsNullOrEmpty(apiKey)) _apiKey = apiKey; - var handler = new HttpClientHandler { CookieContainer = new CookieContainer() }; + var handler = new HttpClientHandler { CookieContainer = _cookies }; var baseAddress = serverUrl; if (!baseAddress.EndsWith("/")) @@ -115,6 +119,30 @@ public async Task DeleteAsync(ILinked entity, string link, TEntity cont new StreamReader(stream).ReadToEnd(); } + public async Task> StreamAsync(ILinked entity, string link, IDictionary parameters = null) + { + return await WebSocketStreamAsync(entity, link, parameters, reader => _serializer.Deserialize(new JsonTextReader(reader))); + } + + public async Task> StreamTextAsync(ILinked entity, string link, IDictionary parameters = null) + { + return await WebSocketStreamAsync(entity, link, parameters, reader => reader.ReadToEnd()); + } + + async Task> WebSocketStreamAsync(ILinked entity, string link, IDictionary parameters, Func deserialize) + { + var linkUri = ResolveLink(entity, link, parameters); + + var socket = new ClientWebSocket(); + socket.Options.Cookies = _cookies; + if (_apiKey != null) + socket.Options.SetRequestHeader("X-Seq-ApiKey", _apiKey); + + await socket.ConnectAsync(new Uri(linkUri), CancellationToken.None); + + return new ObservableStream(socket, deserialize); + } + async Task HttpGetAsync(string url) { var request = new HttpRequestMessage(HttpMethod.Get, url); diff --git a/src/Seq.Api/Api/ResourceGroups/EventsResourceGroup.cs b/src/Seq.Api/Api/ResourceGroups/EventsResourceGroup.cs index de20e82..beb8ce5 100644 --- a/src/Seq.Api/Api/ResourceGroups/EventsResourceGroup.cs +++ b/src/Seq.Api/Api/ResourceGroups/EventsResourceGroup.cs @@ -5,6 +5,7 @@ using Seq.Api.Model.Events; using Seq.Api.Model.Shared; using Seq.Api.Model.Signals; +using Seq.Api.Streams; namespace Seq.Api.ResourceGroups { @@ -156,5 +157,47 @@ public async Task DeleteInSignalAsync( var body = signal ?? new SignalEntity(); return await GroupPostAsync("DeleteInSignal", body, parameters).ConfigureAwait(false); } + + /// + /// Connect to the live event stream. Dispose the resulting stream to disconnect. + /// + /// The type into which events should be deserialized. + /// If provided, a list of signal ids whose intersection will be filtered for the result. + /// A strict Seq filter expression to match (text expressions must be in double quotes). To + /// convert a "fuzzy" filter into a strict one the way the Seq UI does, use connection.Expressions.ToStrictAsync(). + /// An observable that will stream events from the server to subscribers. Events will be buffered server-side until the first + /// subscriber connects, ensure at least one subscription is made in order to avoid event loss. + public async Task> StreamAsync( + string[] intersectIds = null, + string filter = null) + { + var parameters = new Dictionary(); + if (intersectIds != null && intersectIds.Length > 0) { parameters.Add("intersectIds", string.Join(",", intersectIds)); } + if (filter != null) { parameters.Add("filter", filter); } + + var group = await LoadGroupAsync().ConfigureAwait(false); + return await Client.StreamAsync(group, "Stream", parameters).ConfigureAwait(false); + } + + /// + /// Retrieve a list of events that match a set of conditions. The complete result is buffered into memory, + /// so if a large result set is expected, use InSignalAsync() and lastReadEventId to page the results. + /// + /// If provided, a list of signal ids whose intersection will be filtered for the result. + /// A strict Seq filter expression to match (text expressions must be in double quotes). To + /// convert a "fuzzy" filter into a strict one the way the Seq UI does, use connection.Expressions.ToStrictAsync(). + /// An observable that will stream events from the server to subscribers. Events will be buffered server-side until the first + /// subscriber connects, ensure at least one subscription is made in order to avoid event loss. + public async Task> StreamDocumentsAsync( + string[] intersectIds = null, + string filter = null) + { + var parameters = new Dictionary(); + if (intersectIds != null && intersectIds.Length > 0) { parameters.Add("intersectIds", string.Join(",", intersectIds)); } + if (filter != null) { parameters.Add("filter", filter); } + + var group = await LoadGroupAsync().ConfigureAwait(false); + return await Client.StreamTextAsync(group, "Stream", parameters).ConfigureAwait(false); + } } } diff --git a/src/Seq.Api/Api/Streams/ObservableStream.cs b/src/Seq.Api/Api/Streams/ObservableStream.cs new file mode 100644 index 0000000..f1a430c --- /dev/null +++ b/src/Seq.Api/Api/Streams/ObservableStream.cs @@ -0,0 +1,201 @@ +// Copyright 2016 Datalust; based on code from Serilog.Sinks.Observable: +// +// Copyright 2013-2016 Serilog Contributors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Net.WebSockets; +using System.Threading.Tasks; +using System.IO; +using System.Text; + +namespace Seq.Api.Streams +{ + public sealed class ObservableStream : IObservable, IDisposable + { + // Uses memory barriers for non-blocking reads during Emit, and replaces the + // list of observers completely upon subscribe/unsubscribe. + // Makes the assumption that list iteration is not + // mutating - correct but not guaranteed by the BCL. + readonly object _syncRoot = new object(); + IList> _observers = new List>(); + bool _disposed; + Task _run; + + readonly ClientWebSocket _socket; + readonly Func _deserialize; + + sealed class Unsubscriber : IDisposable + { + readonly ObservableStream _sink; + readonly IObserver _observer; + + public Unsubscriber(ObservableStream sink, IObserver observer) + { + if (sink == null) throw new ArgumentNullException(nameof(sink)); + if (observer == null) throw new ArgumentNullException(nameof(observer)); + _sink = sink; + _observer = observer; + } + + public void Dispose() + { + _sink.Unsubscribe(_observer); + } + } + + internal ObservableStream(ClientWebSocket socket, Func deserialize) + { + if (socket == null) throw new ArgumentNullException(nameof(socket)); + if (deserialize == null) throw new ArgumentNullException(nameof(deserialize)); + + _deserialize = deserialize; + _socket = socket; + } + + public IDisposable Subscribe(IObserver observer) + { + if (observer == null) throw new ArgumentNullException(nameof(observer)); + + lock (_syncRoot) + { + if (_disposed) + throw new ObjectDisposedException(message: "The observable WebSocket is disposed.", innerException: null); + + var old = _observers; + var newObservers = _observers.Concat(new[] { observer }).ToList(); + while (old != Interlocked.Exchange(ref _observers, newObservers)) + { + old = _observers; + newObservers = _observers.Concat(new[] { observer }).ToList(); + } + + if (_run == null) + _run = Task.Run(Receive); + } + + return new Unsubscriber(this, observer); + } + + void Unsubscribe(IObserver observer) + { + if (observer == null) throw new ArgumentNullException(nameof(observer)); + + lock (_syncRoot) + { + if (_disposed) + throw new ObjectDisposedException(message: "The observable WebSocket is disposed.", innerException: null); + + var old = _observers; + var newObservers = _observers.Except(new[] { observer }).ToList(); + while (old != Interlocked.Exchange(ref _observers, newObservers)) + { + old = _observers; + newObservers = _observers.Except(new[] { observer }).ToList(); + } + } + } + + async Task Receive() + { + var buffer = new byte[16 * 1024]; + var current = new MemoryStream(); + var reader = new StreamReader(current, new UTF8Encoding(false)); + + while (_socket.State == WebSocketState.Open) + { + var received = await _socket.ReceiveAsync(new ArraySegment(buffer), CancellationToken.None); + if (received.MessageType == WebSocketMessageType.Close) + { + await _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None); + } + else + { + current.Write(buffer, 0, received.Count); + + if (received.EndOfMessage) + { + current.Position = 0; + var value = _deserialize(reader); + + current.SetLength(0); + reader.DiscardBufferedData(); + + Emit(value); + } + } + } + } + + void Emit(T value) + { + if (value == null) throw new ArgumentNullException(nameof(value)); + + Interlocked.MemoryBarrier(); + + IList exceptions = null; + + // Mutations are made by replacing _observers wholesale. + // ReSharper disable once InconsistentlySynchronizedField + foreach (var observer in _observers) + { + try + { + observer.OnNext(value); + } + catch (Exception ex) + { + if (exceptions == null) + exceptions = new List(); + exceptions.Add(ex); + } + } + + // TODO; these need to be caught and propagated. + if (exceptions != null) + throw new AggregateException("At least one observer failed to accept the event", exceptions); + } + + public void Dispose() + { + lock (_syncRoot) + { + if (_disposed) return; + _disposed = true; + + try + { + if (_socket.State == WebSocketState.Open) + _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None).ConfigureAwait(false).GetAwaiter().GetResult(); + } + catch { } + + foreach (var observer in _observers) + { + observer.OnCompleted(); + } + + Interlocked.Exchange(ref _observers, new List>()); + + if (_run != null) + _run.ConfigureAwait(false).GetAwaiter().GetResult(); + + _socket.Dispose(); + } + } + } +} diff --git a/src/Seq.Api/Seq.Api.xproj b/src/Seq.Api/Seq.Api.xproj index def013e..75ec30c 100644 --- a/src/Seq.Api/Seq.Api.xproj +++ b/src/Seq.Api/Seq.Api.xproj @@ -4,18 +4,16 @@ 14.0 $(MSBuildExtensionsPath32)\Microsoft\VisualStudio\v$(VisualStudioVersion) - d4e037de-9778-4e48-a4a7-e8c1751e637c - Superpower + Seq .\obj .\bin\ v4.5.2 - 2.0 - + \ No newline at end of file diff --git a/src/Seq.Api/project.json b/src/Seq.Api/project.json index 6435e19..20f7b91 100644 --- a/src/Seq.Api/project.json +++ b/src/Seq.Api/project.json @@ -12,7 +12,9 @@ }, "buildOptions": { - "warningsAsErrors": true + "warningsAsErrors": true, + "xmlDoc": true, + "nowarn": ["CS1591"] }, "configurations": { @@ -39,7 +41,8 @@ "netstandard1.3": { "dependencies": { "System.Collections.Concurrent": "4.0.12", - "System.Net.Http": "4.1.0" + "System.Net.Http": "4.1.0", + "System.Net.WebSockets.Client": "4.0.0" } }, "net4.5": { From 3cda2e7886344d6058dae99e4c2bfd367bac43b6 Mon Sep 17 00:00:00 2001 From: Nicholas Blumhardt Date: Fri, 28 Oct 2016 11:54:56 +1000 Subject: [PATCH 2/2] Update sample to use new LogEventReader features, all the basics working --- example/SeqTail/Program.cs | 29 ++-- example/SeqTail/project.json | 6 +- .../Streams/ObservableStream.Unsubscriber.cs | 40 +++++ src/Seq.Api/Api/Streams/ObservableStream.cs | 142 ++++++++++-------- 4 files changed, 138 insertions(+), 79 deletions(-) create mode 100644 src/Seq.Api/Api/Streams/ObservableStream.Unsubscriber.cs diff --git a/example/SeqTail/Program.cs b/example/SeqTail/Program.cs index 63ed136..fb9c208 100644 --- a/example/SeqTail/Program.cs +++ b/example/SeqTail/Program.cs @@ -6,8 +6,8 @@ using Serilog; using System.Reactive.Linq; using Serilog.Formatting.Compact.Reader; -using System.IO; -using Serilog.Events; +using System.Threading; +using Newtonsoft.Json.Linq; namespace SeqTail { @@ -42,7 +42,13 @@ static void Main(string[] args) var apiKey = Normalize(arguments["--apikey"]); var filter = Normalize(arguments["--filter"]); - Run(server, apiKey, filter).GetAwaiter().GetResult(); + var cancel = new CancellationTokenSource(); + Console.WriteLine("Tailing, press Ctrl+C to exit."); + Console.CancelKeyPress += (s,a) => cancel.Cancel(); + + var run = Task.Run(() => Run(server, apiKey, filter, cancel)); + + run.GetAwaiter().GetResult(); } catch (Exception ex) { @@ -58,7 +64,7 @@ static string Normalize(ValueObject v) return string.IsNullOrWhiteSpace(s) ? null : s; } - static async Task Run(string server, string apiKey, string filter) + static async Task Run(string server, string apiKey, string filter, CancellationTokenSource cancel) { var connection = new SeqConnection(server, apiKey); @@ -69,18 +75,13 @@ static async Task Run(string server, string apiKey, string filter) strict = converted.StrictExpression; } - using (var stream = await connection.Events.StreamDocumentsAsync(filter: strict)) + using (var stream = await connection.Events.StreamAsync(filter: strict)) { - var subscription = stream.Subscribe(document => - { - var reader = new LogEventReader(new StringReader(document)); - LogEvent evt; - if (!reader.TryRead(out evt)) - throw new InvalidOperationException("Expected document to contain data."); - Log.Write(evt); - }); + var subscription = stream + .Select(jObject => LogEventReader.ReadFromJObject(jObject)) + .Subscribe(evt => Log.Write(evt), () => cancel.Cancel()); - Console.ReadKey(true); + cancel.Token.WaitHandle.WaitOne(); subscription.Dispose(); } } diff --git a/example/SeqTail/project.json b/example/SeqTail/project.json index 00d5aab..f4197f9 100644 --- a/example/SeqTail/project.json +++ b/example/SeqTail/project.json @@ -6,11 +6,13 @@ }, "dependencies": { + "System.Threading.Tasks": "4.0.11", "Seq.Api": { "target": "project" }, "docopt.net": "0.6.1.9", - "Serilog.Formatting.Compact.Reader": "1.0.0-dev-00004", + "Serilog.Formatting.Compact.Reader": "1.0.0-dev-00008", "Serilog.Sinks.Literate": "2.0.0", - "System.Reactive": "3.0.0" + "System.Reactive": "3.0.0", + "Newtonsoft.Json": "9.0.1" }, "frameworks": { diff --git a/src/Seq.Api/Api/Streams/ObservableStream.Unsubscriber.cs b/src/Seq.Api/Api/Streams/ObservableStream.Unsubscriber.cs new file mode 100644 index 0000000..e4380e3 --- /dev/null +++ b/src/Seq.Api/Api/Streams/ObservableStream.Unsubscriber.cs @@ -0,0 +1,40 @@ +// Copyright 2016 Datalust; based on code from +// Serilog.Sinks.Observable, Copyright 2013-2016 Serilog Contributors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +using System; + +namespace Seq.Api.Streams +{ + public partial class ObservableStream + { + sealed class Unsubscriber : IDisposable + { + readonly ObservableStream _stream; + readonly IObserver _observer; + + public Unsubscriber(ObservableStream sink, IObserver observer) + { + if (sink == null) throw new ArgumentNullException(nameof(sink)); + if (observer == null) throw new ArgumentNullException(nameof(observer)); + _stream = sink; + _observer = observer; + } + + public void Dispose() + { + _stream.Unsubscribe(_observer); + } + } + } +} diff --git a/src/Seq.Api/Api/Streams/ObservableStream.cs b/src/Seq.Api/Api/Streams/ObservableStream.cs index f1a430c..5417fbc 100644 --- a/src/Seq.Api/Api/Streams/ObservableStream.cs +++ b/src/Seq.Api/Api/Streams/ObservableStream.cs @@ -1,6 +1,5 @@ -// Copyright 2016 Datalust; based on code from Serilog.Sinks.Observable: -// -// Copyright 2013-2016 Serilog Contributors +// Copyright 2016 Datalust; based on code from +// Serilog.Sinks.Observable, Copyright 2013-2016 Serilog Contributors // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -16,48 +15,27 @@ using System; using System.Collections.Generic; -using System.Linq; using System.Threading; using System.Net.WebSockets; using System.Threading.Tasks; using System.IO; using System.Text; +using System.Linq; namespace Seq.Api.Streams { - public sealed class ObservableStream : IObservable, IDisposable + // Some questionable synchronization in this one; probably better to take the Rx dependency + // and use the Rx support classes here. + public sealed partial class ObservableStream : IObservable, IDisposable { - // Uses memory barriers for non-blocking reads during Emit, and replaces the - // list of observers completely upon subscribe/unsubscribe. - // Makes the assumption that list iteration is not - // mutating - correct but not guaranteed by the BCL. readonly object _syncRoot = new object(); IList> _observers = new List>(); - bool _disposed; + bool _ended, _disposed; Task _run; readonly ClientWebSocket _socket; readonly Func _deserialize; - sealed class Unsubscriber : IDisposable - { - readonly ObservableStream _sink; - readonly IObserver _observer; - - public Unsubscriber(ObservableStream sink, IObserver observer) - { - if (sink == null) throw new ArgumentNullException(nameof(sink)); - if (observer == null) throw new ArgumentNullException(nameof(observer)); - _sink = sink; - _observer = observer; - } - - public void Dispose() - { - _sink.Unsubscribe(_observer); - } - } - internal ObservableStream(ClientWebSocket socket, Func deserialize) { if (socket == null) throw new ArgumentNullException(nameof(socket)); @@ -74,16 +52,16 @@ public IDisposable Subscribe(IObserver observer) lock (_syncRoot) { if (_disposed) - throw new ObjectDisposedException(message: "The observable WebSocket is disposed.", innerException: null); + throw new ObjectDisposedException(message: "The observable stream is disposed.", innerException: null); - var old = _observers; - var newObservers = _observers.Concat(new[] { observer }).ToList(); - while (old != Interlocked.Exchange(ref _observers, newObservers)) + if (_ended) { - old = _observers; - newObservers = _observers.Concat(new[] { observer }).ToList(); + observer.OnCompleted(); + return new Unsubscriber(this, observer); } + _observers = _observers.Concat(new[] { observer }).ToList(); + if (_run == null) _run = Task.Run(Receive); } @@ -98,15 +76,9 @@ void Unsubscribe(IObserver observer) lock (_syncRoot) { if (_disposed) - throw new ObjectDisposedException(message: "The observable WebSocket is disposed.", innerException: null); + throw new ObjectDisposedException(message: "The observable stream is disposed.", innerException: null); - var old = _observers; - var newObservers = _observers.Except(new[] { observer }).ToList(); - while (old != Interlocked.Exchange(ref _observers, newObservers)) - { - old = _observers; - newObservers = _observers.Except(new[] { observer }).ToList(); - } + _observers = _observers.Except(new[] { observer }).ToList(); } } @@ -122,6 +94,7 @@ async Task Receive() if (received.MessageType == WebSocketMessageType.Close) { await _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None); + End(); } else { @@ -145,13 +118,17 @@ void Emit(T value) { if (value == null) throw new ArgumentNullException(nameof(value)); - Interlocked.MemoryBarrier(); + IList> observers; + lock (_syncRoot) + { + if (_ended || _disposed) + return; - IList exceptions = null; + observers = _observers; + } - // Mutations are made by replacing _observers wholesale. - // ReSharper disable once InconsistentlySynchronizedField - foreach (var observer in _observers) + IList exceptions = null; + foreach (var observer in observers) { try { @@ -165,37 +142,76 @@ void Emit(T value) } } - // TODO; these need to be caught and propagated. if (exceptions != null) - throw new AggregateException("At least one observer failed to accept the event", exceptions); + OnError(exceptions); } - public void Dispose() + void End() { + IList> observers; lock (_syncRoot) { - if (_disposed) return; - _disposed = true; + if (_ended) + return; + observers = _observers; + _observers = new List>(); + _ended = true; + } + + IList exceptions = null; + foreach (var observer in observers) + { try { - if (_socket.State == WebSocketState.Open) - _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "", CancellationToken.None).ConfigureAwait(false).GetAwaiter().GetResult(); + observer.OnCompleted(); } - catch { } - - foreach (var observer in _observers) + catch (Exception ex) { - observer.OnCompleted(); + if (exceptions == null) + exceptions = new List(); + exceptions.Add(ex); } + } - Interlocked.Exchange(ref _observers, new List>()); + if (exceptions != null) + OnError(exceptions); + } - if (_run != null) - _run.ConfigureAwait(false).GetAwaiter().GetResult(); + void OnError(IList exceptions) + { + // This will hit TaskScheduler.UnobservedTaskException + throw new AggregateException("At least one observer failed to accept the event", exceptions); + } - _socket.Dispose(); + public void Dispose() + { + lock (_syncRoot) + { + if (_disposed) return; + _disposed = true; } + + try + { + if (_socket.State == WebSocketState.Open) + { + using (var timeout = new CancellationTokenSource()) + { + timeout.CancelAfter(30000); + _socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Close requested", timeout.Token) + .ConfigureAwait(false).GetAwaiter().GetResult(); + } + } + } + catch { } + + End(); + + if (_run != null) + _run.ConfigureAwait(false).GetAwaiter().GetResult(); + + _socket.Dispose(); } } }