]> git.deb.at Git - rhonda/impostor.git/commitdiff
Use channel for recorder to fix lagg issue caused by locking
authorAeonLucid <aeonlucid@outlook.com>
Sun, 1 Nov 2020 22:21:18 +0000 (23:21 +0100)
committerAeonLucid <aeonlucid@outlook.com>
Sun, 1 Nov 2020 22:21:18 +0000 (23:21 +0100)
src/Impostor.Hazel/Udp/UdpClientConnection.cs
src/Impostor.Server/Program.cs
src/Impostor.Server/Recorder/PacketRecorder.cs

index 4d80ddd78b89df1afe349a9ec839abb2a50862b8..5125ebe8034522a2e4c8b6b08d73ba81fe34c3c0 100644 (file)
@@ -26,10 +26,7 @@ namespace Impostor.Hazel.Udp
 
         private readonly Timer _reliablePacketTimer;
         private readonly SemaphoreSlim _connectWaitLock;
-        private readonly ArrayPool<byte> _pool;
-        private readonly Channel<byte[]> _channel;
         private Task _listenTask;
-        private Task _handleTask;
 
         /// <summary>
         ///     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<byte>.Shared;
-            _channel = Channel.CreateUnbounded<byte[]>(new UnboundedChannelOptions
-            {
-                SingleReader = true,
-                SingleWriter = true
-            });
         }
 
         ~UdpClientConnection()
index 234ee01306d68db26de0449cd44ce64ba29c58ef..a58dfa084a826fad634e93d72bcd24cd14189460 100644 (file)
@@ -161,6 +161,7 @@ namespace Impostor.Server
                             });
 
                             services.AddSingleton<PacketRecorder>();
+                            services.AddHostedService(sp => sp.GetRequiredService<PacketRecorder>());
                             services.AddSingleton<IClientFactory, ClientFactory<ClientRecorder>>();
                         }
                         else
index 322af4f57cb62d514d5db94805bc6b300713abfb..0ac8a60efd198be1642b2729d779870f6dc664dc 100644 (file)
@@ -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
     /// <summary>
     ///     Records all packets received in <see cref="ClientRecorder.HandleMessageAsync"/>.
     /// </summary>
-    internal class PacketRecorder : IDisposable
+    internal class PacketRecorder : BackgroundService
     {
+        private readonly string _path;
         private readonly ILogger<PacketRecorder> _logger;
         private readonly ObjectPool<PacketSerializationContext> _pool;
-        private readonly SemaphoreSlim _writerLock;
-        private readonly FileStream _writer;
+        private readonly Channel<byte[]> _channel;
 
         public PacketRecorder(ILogger<PacketRecorder> logger, IOptions<DebugConfig> options, ObjectPool<PacketSerializationContext> 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<byte[]>(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());
         }
     }
 }