diff --git a/Assets/Tests/Editor/JsonRpcHeartbeatTests.cs b/Assets/Tests/Editor/JsonRpcHeartbeatTests.cs index 303671bbdd..37073f24b9 100644 --- a/Assets/Tests/Editor/JsonRpcHeartbeatTests.cs +++ b/Assets/Tests/Editor/JsonRpcHeartbeatTests.cs @@ -55,8 +55,9 @@ public async Task SendHeartbeatsAsync_WhenRunning_WritesFramesUntilCancelled() // without leaving background work behind. int writtenFrameCount = 0; using CancellationTokenSource cancellationSource = new(); + UnityCliLoopBridgeHeartbeatService heartbeatService = new(); - Task heartbeatTask = UnityCliLoopBridgeServer.SendHeartbeatsAsync( + Task heartbeatTask = heartbeatService.SendHeartbeatsAsync( () => "{}", _ => { @@ -79,8 +80,9 @@ public async Task SendHeartbeatsAsync_WhenWriteThrowsIOException_StopsWithoutFau // Tests that a broken connection ends the heartbeat loop silently; teardown is // owned by the read loop, not the heartbeat writer. using CancellationTokenSource cancellationSource = new(); + UnityCliLoopBridgeHeartbeatService heartbeatService = new(); - Task heartbeatTask = UnityCliLoopBridgeServer.SendHeartbeatsAsync( + Task heartbeatTask = heartbeatService.SendHeartbeatsAsync( () => "{}", _ => throw new System.IO.IOException("broken pipe"), TimeSpan.FromMilliseconds(1), diff --git a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs new file mode 100644 index 0000000000..f765870192 --- /dev/null +++ b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs @@ -0,0 +1,56 @@ +using System; +using System.Threading; +using System.Threading.Tasks; + +namespace io.github.hatayama.UnityCliLoop.Infrastructure +{ + /// + /// Cancels accepted project IPC requests when their client connection disappears. + /// + internal sealed class UnityCliLoopBridgeClientDisconnectMonitor + { + private const int ClientDisconnectMonitorPollMilliseconds = 100; + + /// + /// Monitors an accepted client connection and cancels the request token source when the client disconnects. + /// + internal async Task MonitorClientDisconnectAsync( + BridgeClientConnection client, + CancellationTokenSource requestCancellationTokenSource) + { + while (!requestCancellationTokenSource.IsCancellationRequested) + { + if (!client.IsConnected) + { + requestCancellationTokenSource.Cancel(); + return; + } + + try + { + await Task.Delay(ClientDisconnectMonitorPollMilliseconds, requestCancellationTokenSource.Token); + } + catch (OperationCanceledException) + { + // Cancellation is the normal stop signal from StopClientDisconnectMonitorAsync. + // Without the token the delay always ran to completion, adding one poll + // interval of tail latency to every request teardown and server shutdown. + return; + } + } + } + + internal async Task StopClientDisconnectMonitorAsync( + Task clientDisconnectMonitorTask, + CancellationTokenSource requestCancellationTokenSource) + { + if (clientDisconnectMonitorTask == null) + { + return; + } + + requestCancellationTokenSource.Cancel(); + await clientDisconnectMonitorTask; + } + } +} diff --git a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs.meta b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs.meta new file mode 100644 index 0000000000..315569730c --- /dev/null +++ b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeClientDisconnectMonitor.cs.meta @@ -0,0 +1,11 @@ +fileFormatVersion: 2 +guid: 3331ed1250b74839a192452fa1298bd8 +MonoImporter: + externalObjects: {} + serializedVersion: 2 + defaultReferences: [] + executionOrder: 0 + icon: {instanceID: 0} + userData: + assetBundleName: + assetBundleVariant: diff --git a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs new file mode 100644 index 0000000000..068e6997da --- /dev/null +++ b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs @@ -0,0 +1,58 @@ +using System; +using System.IO; +using System.Threading; +using System.Threading.Tasks; + +namespace io.github.hatayama.UnityCliLoop.Infrastructure +{ + /// + /// Sends and stops project IPC heartbeat frames for an accepted client request. + /// + internal sealed class UnityCliLoopBridgeHeartbeatService + { + /// + /// Sends heartbeat frames at the given interval until cancelled. Write failures end + /// the loop silently because the connection teardown is owned by the read loop. + /// + internal async Task SendHeartbeatsAsync( + Func createHeartbeatJson, + Func writeFrameAsync, + TimeSpan interval, + CancellationToken ct) + { + while (true) + { + try + { + await Task.Delay(interval, ct); + await writeFrameAsync(createHeartbeatJson()); + } + catch (OperationCanceledException) + { + return; + } + catch (IOException) + { + return; + } + catch (ObjectDisposedException) + { + return; + } + } + } + + internal async Task StopHeartbeatsAsync( + Task heartbeatTask, + CancellationTokenSource heartbeatCancellationSource) + { + if (heartbeatTask == null) + { + return; + } + + heartbeatCancellationSource?.Cancel(); + await heartbeatTask; + } + } +} diff --git a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs.meta b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs.meta new file mode 100644 index 0000000000..b17b27fa62 --- /dev/null +++ b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeHeartbeatService.cs.meta @@ -0,0 +1,11 @@ +fileFormatVersion: 2 +guid: 535bfd8a9aaa4ffbb3544b4aa611fb80 +MonoImporter: + externalObjects: {} + serializedVersion: 2 + defaultReferences: [] + executionOrder: 0 + icon: {instanceID: 0} + userData: + assetBundleName: + assetBundleVariant: diff --git a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeServer.cs b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeServer.cs index a6455970e0..a57d343120 100644 --- a/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeServer.cs +++ b/Packages/src/Editor/Infrastructure/UnityCliLoopBridgeServer.cs @@ -35,7 +35,12 @@ internal UnityCliLoopBridgeServerInstanceFactory(IDomainReloadDetectionService d public IUnityCliLoopServerInstance Create() { - UnityCliLoopBridgeServer server = new(_domainReloadDetectionService); + UnityCliLoopBridgeHeartbeatService heartbeatService = new(); + UnityCliLoopBridgeClientDisconnectMonitor clientDisconnectMonitor = new(); + UnityCliLoopBridgeServer server = new( + _domainReloadDetectionService, + heartbeatService, + clientDisconnectMonitor); server.ServerLoopExited += NotifyServerLoopExited; return server; @@ -57,6 +62,8 @@ public class UnityCliLoopBridgeServer : IUnityCliLoopServerInstance // Subscribers must marshal to main thread before accessing Unity APIs. public event Action ServerLoopExited; private readonly IDomainReloadDetectionService _domainReloadDetectionService; + private readonly UnityCliLoopBridgeHeartbeatService _heartbeatService; + private readonly UnityCliLoopBridgeClientDisconnectMonitor _clientDisconnectMonitor; // HResult error codes for normal disconnection detection private static readonly HashSet NormalDisconnectionHResults = new() @@ -79,14 +86,22 @@ public class UnityCliLoopBridgeServer : IUnityCliLoopServerInstance private readonly ConcurrentDictionary _clientStreams = new(); private readonly ConcurrentDictionary _clientTasks = new(); private int _nextClientTaskId; - private const int ClientDisconnectMonitorPollMilliseconds = 100; - internal UnityCliLoopBridgeServer(IDomainReloadDetectionService domainReloadDetectionService) + internal UnityCliLoopBridgeServer( + IDomainReloadDetectionService domainReloadDetectionService, + UnityCliLoopBridgeHeartbeatService heartbeatService, + UnityCliLoopBridgeClientDisconnectMonitor clientDisconnectMonitor) { System.Diagnostics.Debug.Assert(domainReloadDetectionService != null, "domainReloadDetectionService must not be null"); + System.Diagnostics.Debug.Assert(heartbeatService != null, "heartbeatService must not be null"); + System.Diagnostics.Debug.Assert(clientDisconnectMonitor != null, "clientDisconnectMonitor must not be null"); _domainReloadDetectionService = domainReloadDetectionService ?? throw new ArgumentNullException(nameof(domainReloadDetectionService)); + _heartbeatService = heartbeatService + ?? throw new ArgumentNullException(nameof(heartbeatService)); + _clientDisconnectMonitor = clientDisconnectMonitor + ?? throw new ArgumentNullException(nameof(clientDisconnectMonitor)); } /// @@ -616,7 +631,7 @@ await WriteJsonResponseLockedAsync( // server token an in-flight write could ignore StopHeartbeatsAsync // and stall the final response behind a slow client. CancellationToken heartbeatToken = heartbeatCancellationSource.Token; - heartbeatTask = SendHeartbeatsAsync( + heartbeatTask = _heartbeatService.SendHeartbeatsAsync( createHeartbeatJson, heartbeatJson => WriteJsonResponseLockedAsync( stream, streamWriteLock, heartbeatJson, heartbeatToken), @@ -630,72 +645,29 @@ await WriteJsonResponseLockedAsync( } clientDisconnectMonitorTask = - MonitorClientDisconnectAsync(client, requestCancellationTokenSource); + _clientDisconnectMonitor.MonitorClientDisconnectAsync( + client, + requestCancellationTokenSource); }); // Stop heartbeats before the final response so no heartbeat frame can be // queued after the response the CLI stops reading at. - await StopHeartbeatsAsync(heartbeatTask, heartbeatCancellationSource); + await _heartbeatService.StopHeartbeatsAsync(heartbeatTask, heartbeatCancellationSource); heartbeatTask = null; await WriteJsonResponseLockedAsync(stream, streamWriteLock, responseJson, serverCancellationToken); } finally { - await StopHeartbeatsAsync(heartbeatTask, heartbeatCancellationSource); + await _heartbeatService.StopHeartbeatsAsync(heartbeatTask, heartbeatCancellationSource); heartbeatCancellationSource?.Dispose(); - await StopClientDisconnectMonitorAsync( + await _clientDisconnectMonitor.StopClientDisconnectMonitorAsync( clientDisconnectMonitorTask, requestCancellationTokenSource); } } } - /// - /// Sends heartbeat frames at the given interval until cancelled. Write failures end - /// the loop silently because the connection teardown is owned by the read loop. - /// - internal static async Task SendHeartbeatsAsync( - Func createHeartbeatJson, - Func writeFrameAsync, - TimeSpan interval, - CancellationToken ct) - { - while (true) - { - try - { - await Task.Delay(interval, ct); - await writeFrameAsync(createHeartbeatJson()); - } - catch (OperationCanceledException) - { - return; - } - catch (IOException) - { - return; - } - catch (ObjectDisposedException) - { - return; - } - } - } - - private static async Task StopHeartbeatsAsync( - Task heartbeatTask, - CancellationTokenSource heartbeatCancellationSource) - { - if (heartbeatTask == null) - { - return; - } - - heartbeatCancellationSource?.Cancel(); - await heartbeatTask; - } - private async Task WriteJsonResponseLockedAsync( Stream stream, SemaphoreSlim streamWriteLock, @@ -716,45 +688,6 @@ private async Task WriteJsonResponseLockedAsync( } } - private static async Task MonitorClientDisconnectAsync( - BridgeClientConnection client, - CancellationTokenSource requestCancellationTokenSource) - { - while (!requestCancellationTokenSource.IsCancellationRequested) - { - if (!client.IsConnected) - { - requestCancellationTokenSource.Cancel(); - return; - } - - try - { - await Task.Delay(ClientDisconnectMonitorPollMilliseconds, requestCancellationTokenSource.Token); - } - catch (OperationCanceledException) - { - // Cancellation is the normal stop signal from StopClientDisconnectMonitorAsync. - // Without the token the delay always ran to completion, adding one poll - // interval of tail latency to every request teardown and server shutdown. - return; - } - } - } - - private static async Task StopClientDisconnectMonitorAsync( - Task clientDisconnectMonitorTask, - CancellationTokenSource requestCancellationTokenSource) - { - if (clientDisconnectMonitorTask == null) - { - return; - } - - requestCancellationTokenSource.Cancel(); - await clientDisconnectMonitorTask; - } - /// /// Determines if the given exception represents a normal client disconnection. ///