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.
///