Merge nucleic/warm-north-vole-pf37 into dev
This commit is contained in:
@@ -0,0 +1,223 @@
|
||||
using NucleicProtocol.Interop;
|
||||
|
||||
namespace NucleicApp.Services;
|
||||
|
||||
public enum HostConnectionState
|
||||
{
|
||||
Stopped,
|
||||
Connecting,
|
||||
Connected,
|
||||
Reconnecting,
|
||||
Failed,
|
||||
}
|
||||
|
||||
/// Owns one native client handle for one host. Native callbacks arrive on the DLL's dedicated
|
||||
/// thread and are marshalled onto the renderer dispatcher before state is projected.
|
||||
public sealed class HostConnection : IAsyncDisposable
|
||||
{
|
||||
private readonly object gate = new();
|
||||
private readonly RendererStore store;
|
||||
private readonly IProtocolClientFactory clientFactory;
|
||||
private readonly IRendererDispatcher dispatcher;
|
||||
private readonly string identityDirectory;
|
||||
private readonly string deviceLabel;
|
||||
private IProtocolClient? client;
|
||||
private CancellationTokenSource? lifetime;
|
||||
private Task? retryTask;
|
||||
private string? rendezvousPath;
|
||||
private int retryAttempt;
|
||||
|
||||
public HostConnectionState State { get; private set; } = HostConnectionState.Stopped;
|
||||
public ProtocolState? LastProtocolState { get; private set; }
|
||||
public event Action<HostConnectionState>? StateChanged;
|
||||
|
||||
public HostConnection(
|
||||
RendererStore store, string identityDirectory,
|
||||
IProtocolClientFactory? clientFactory = null,
|
||||
IRendererDispatcher? dispatcher = null,
|
||||
string deviceLabel = "Nucleic for Windows")
|
||||
{
|
||||
this.store = store;
|
||||
this.identityDirectory = identityDirectory;
|
||||
this.clientFactory = clientFactory ?? new ProtocolClientFactory();
|
||||
this.dispatcher = dispatcher ?? InlineRendererDispatcher.Instance;
|
||||
this.deviceLabel = deviceLabel;
|
||||
}
|
||||
|
||||
public async Task StartLocalAsync(string rendezvousPath, CancellationToken cancellationToken = default)
|
||||
{
|
||||
lock (gate)
|
||||
{
|
||||
if (lifetime is not null) throw new InvalidOperationException("connection is already started");
|
||||
var createdLifetime = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
|
||||
IProtocolClient? createdClient = null;
|
||||
try
|
||||
{
|
||||
createdClient = clientFactory.Create(new ProtocolClientConfiguration
|
||||
{
|
||||
IdentityDirectory = identityDirectory,
|
||||
DeviceLabel = deviceLabel,
|
||||
Channel = BuildChannel.Current,
|
||||
});
|
||||
createdClient.MessageJsonReceived += OnMessage;
|
||||
createdClient.StateJsonReceived += OnState;
|
||||
this.rendezvousPath = rendezvousPath;
|
||||
lifetime = createdLifetime;
|
||||
client = createdClient;
|
||||
}
|
||||
catch
|
||||
{
|
||||
createdClient?.Dispose();
|
||||
createdLifetime.Dispose();
|
||||
throw;
|
||||
}
|
||||
}
|
||||
if (await ConnectNowAsync(initial: true, lifetime.Token).ConfigureAwait(false)
|
||||
!= ProtocolResult.Ok)
|
||||
{
|
||||
ScheduleReconnect();
|
||||
}
|
||||
}
|
||||
|
||||
public ProtocolResult SendIntent(string json)
|
||||
{
|
||||
lock (gate) return client?.SendIntent(json) ?? ProtocolResult.InvalidState;
|
||||
}
|
||||
|
||||
private async Task<ProtocolResult> ConnectNowAsync(
|
||||
bool initial, CancellationToken cancellationToken)
|
||||
{
|
||||
IProtocolClient active;
|
||||
string path;
|
||||
lock (gate)
|
||||
{
|
||||
active = client ?? throw new InvalidOperationException("connection has no client");
|
||||
path = rendezvousPath ?? throw new InvalidOperationException("connection has no rendezvous");
|
||||
}
|
||||
SetState(initial ? HostConnectionState.Connecting : HostConnectionState.Reconnecting);
|
||||
var record = await RendezvousFile.ReadAsync(path, cancellationToken).ConfigureAwait(false);
|
||||
var result = active.ConnectLocal(record);
|
||||
if (result != ProtocolResult.Ok)
|
||||
SetState(HostConnectionState.Failed);
|
||||
return result;
|
||||
}
|
||||
|
||||
private void OnMessage(string json) => dispatcher.Post(() => store.Apply(json));
|
||||
|
||||
private void OnState(string json)
|
||||
{
|
||||
ProtocolState state;
|
||||
try { state = ProtocolState.Parse(json); }
|
||||
catch { return; }
|
||||
dispatcher.Post(() =>
|
||||
{
|
||||
LastProtocolState = state;
|
||||
switch (state.State)
|
||||
{
|
||||
case "connecting": SetState(HostConnectionState.Connecting); break;
|
||||
case "ready":
|
||||
retryAttempt = 0;
|
||||
SetState(HostConnectionState.Connected);
|
||||
_ = SendIntent(ClientIntents.ListSessions());
|
||||
_ = SendIntent(ClientIntents.ListDashboard());
|
||||
_ = SendIntent(ClientIntents.ListPeers());
|
||||
break;
|
||||
case "failed": SetState(HostConnectionState.Failed); ScheduleReconnect(); break;
|
||||
case "closed": ScheduleReconnect(); break;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void ScheduleReconnect()
|
||||
{
|
||||
lock (gate)
|
||||
{
|
||||
if (lifetime is null || lifetime.IsCancellationRequested || retryTask is { IsCompleted: false })
|
||||
return;
|
||||
var cancellationToken = lifetime.Token;
|
||||
var delay = TimeSpan.FromMilliseconds(Math.Min(30_000, 250 * Math.Pow(2, retryAttempt++)));
|
||||
retryTask = Task.Run(async () =>
|
||||
{
|
||||
while (!cancellationToken.IsCancellationRequested)
|
||||
{
|
||||
try
|
||||
{
|
||||
await Task.Delay(delay, cancellationToken).ConfigureAwait(false);
|
||||
if (await ConnectNowAsync(initial: false, cancellationToken)
|
||||
.ConfigureAwait(false) == ProtocolResult.Ok)
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
catch (OperationCanceledException) { return; }
|
||||
catch { SetState(HostConnectionState.Failed); }
|
||||
|
||||
// A synchronous admission failure (notably InvalidState while the native
|
||||
// close callback is just ahead of its task cleanup) must not strand the
|
||||
// connection. The old recursive scheduler observed its own retry task as
|
||||
// active and silently declined to schedule another attempt.
|
||||
delay = TimeSpan.FromMilliseconds(
|
||||
Math.Min(30_000, 250 * Math.Pow(2, retryAttempt++)));
|
||||
}
|
||||
}, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
private void SetState(HostConnectionState next)
|
||||
{
|
||||
if (State == next) return;
|
||||
State = next;
|
||||
StateChanged?.Invoke(next);
|
||||
}
|
||||
|
||||
public async ValueTask DisposeAsync()
|
||||
{
|
||||
CancellationTokenSource? cancellation;
|
||||
Task? retry;
|
||||
IProtocolClient? closing;
|
||||
lock (gate)
|
||||
{
|
||||
cancellation = lifetime;
|
||||
lifetime = null;
|
||||
retry = retryTask;
|
||||
retryTask = null;
|
||||
closing = client;
|
||||
client = null;
|
||||
}
|
||||
if (cancellation is not null) await cancellation.CancelAsync().ConfigureAwait(false);
|
||||
if (retry is not null) try { await retry.ConfigureAwait(false); } catch (OperationCanceledException) { }
|
||||
closing?.Dispose();
|
||||
cancellation?.Dispose();
|
||||
SetState(HostConnectionState.Stopped);
|
||||
}
|
||||
}
|
||||
|
||||
internal static class BuildChannel
|
||||
{
|
||||
public static string Current
|
||||
{
|
||||
get
|
||||
{
|
||||
#if NUCLEIC_STABLE
|
||||
return "release";
|
||||
#elif NUCLEIC_RC
|
||||
return "rc";
|
||||
#elif NUCLEIC_BETA
|
||||
return "beta";
|
||||
#elif NUCLEIC_CANARY
|
||||
return "canary";
|
||||
#else
|
||||
return "local";
|
||||
#endif
|
||||
}
|
||||
}
|
||||
|
||||
public static string Suffix => Current switch
|
||||
{
|
||||
"release" => string.Empty,
|
||||
"rc" => "-rc",
|
||||
"beta" => "-beta",
|
||||
"canary" => "-canary",
|
||||
_ => "-local",
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user