From ca9cd5d818496bd4df181503d980f6ee18f3d83f Mon Sep 17 00:00:00 2001 From: AeonLucid Date: Sun, 1 Nov 2020 23:21:18 +0100 Subject: [PATCH] Use channel for recorder to fix lagg issue caused by locking --- src/Impostor.Hazel/Udp/UdpClientConnection.cs | 9 --- src/Impostor.Server/Program.cs | 1 + .../Recorder/PacketRecorder.cs | 75 ++++++++++--------- 3 files changed, 41 insertions(+), 44 deletions(-) diff --git a/src/Impostor.Hazel/Udp/UdpClientConnection.cs b/src/Impostor.Hazel/Udp/UdpClientConnection.cs index 4d80ddd..5125ebe 100644 --- a/src/Impostor.Hazel/Udp/UdpClientConnection.cs +++ b/src/Impostor.Hazel/Udp/UdpClientConnection.cs @@ -26,10 +26,7 @@ namespace Impostor.Hazel.Udp private readonly Timer _reliablePacketTimer; private readonly SemaphoreSlim _connectWaitLock; - private readonly ArrayPool _pool; - private readonly Channel _channel; private Task _listenTask; - private Task _handleTask; /// /// Creates a new UdpClientConnection. @@ -48,12 +45,6 @@ namespace Impostor.Hazel.Udp _reliablePacketTimer = new Timer(ManageReliablePacketsInternal, null, 100, Timeout.Infinite); _connectWaitLock = new SemaphoreSlim(1, 1); - _pool = ArrayPool.Shared; - _channel = Channel.CreateUnbounded(new UnboundedChannelOptions - { - SingleReader = true, - SingleWriter = true - }); } ~UdpClientConnection() diff --git a/src/Impostor.Server/Program.cs b/src/Impostor.Server/Program.cs index 234ee01..a58dfa0 100644 --- a/src/Impostor.Server/Program.cs +++ b/src/Impostor.Server/Program.cs @@ -161,6 +161,7 @@ namespace Impostor.Server }); services.AddSingleton(); + services.AddHostedService(sp => sp.GetRequiredService()); services.AddSingleton>(); } else diff --git a/src/Impostor.Server/Recorder/PacketRecorder.cs b/src/Impostor.Server/Recorder/PacketRecorder.cs index 322af4f..0ac8a60 100644 --- a/src/Impostor.Server/Recorder/PacketRecorder.cs +++ b/src/Impostor.Server/Recorder/PacketRecorder.cs @@ -2,11 +2,13 @@ using System.IO; using System.Net; using System.Threading; +using System.Threading.Channels; using System.Threading.Tasks; using Impostor.Api.Games; using Impostor.Api.Net.Messages; using Impostor.Server.Config; using Impostor.Server.Net; +using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Microsoft.Extensions.ObjectPool; using Microsoft.Extensions.Options; @@ -16,23 +18,52 @@ namespace Impostor.Server.Recorder /// /// Records all packets received in . /// - internal class PacketRecorder : IDisposable + internal class PacketRecorder : BackgroundService { + private readonly string _path; private readonly ILogger _logger; private readonly ObjectPool _pool; - private readonly SemaphoreSlim _writerLock; - private readonly FileStream _writer; + private readonly Channel _channel; public PacketRecorder(ILogger logger, IOptions options, ObjectPool pool) { var name = $"session_{DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()}.dat"; - var path = Path.Combine(options.Value.GameRecorderPath, name); + _path = Path.Combine(options.Value.GameRecorderPath, name); _logger = logger; - _logger.LogInformation("PacketRecorder is enabled, writing packets to {0}.", path); _pool = pool; - _writerLock = new SemaphoreSlim(1, 1); - _writer = File.Open(path, FileMode.CreateNew, FileAccess.Write, FileShare.Read); + + _channel = Channel.CreateUnbounded(new UnboundedChannelOptions + { + SingleReader = true, + SingleWriter = false, + }); + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + _logger.LogInformation("PacketRecorder is enabled, writing packets to {0}.", _path); + + var writer = File.Open(_path, FileMode.CreateNew, FileAccess.Write, FileShare.Read); + + // Handle messages. + while (!stoppingToken.IsCancellationRequested) + { + try + { + var result = await _channel.Reader.ReadAsync(stoppingToken); + + await writer.WriteAsync(result, stoppingToken); + await writer.FlushAsync(stoppingToken); + } + catch (TaskCanceledException) + { + break; + } + } + + // Clean up. + await writer.DisposeAsync(); } public async Task WriteConnectAsync(ClientRecorder client) @@ -163,35 +194,9 @@ namespace Impostor.Server.Recorder context.Stream.Position = length; } - private async Task WriteAsync(Stream data) - { - var hasLock = false; - - try - { - hasLock = await _writerLock.WaitAsync(TimeSpan.FromMinutes(1)); - - if (hasLock) - { - data.Position = 0; - - await data.CopyToAsync(_writer); - await _writer.FlushAsync(); - } - } - finally - { - if (hasLock) - { - _writerLock.Release(); - } - } - } - - public void Dispose() + private async Task WriteAsync(MemoryStream data) { - _writer.Dispose(); - _writerLock.Dispose(); + await _channel.Writer.WriteAsync(data.ToArray()); } } } -- 2.39.5