From f5ee895d3f8d22905c6a181a632d791b710770c3 Mon Sep 17 00:00:00 2001 From: Nucleic Date: Tue, 4 Aug 2026 02:46:53 -0700 Subject: [PATCH] Merge nucleic/warm-north-vole-pf37 into dev --- Nucleic.sln | 42 ++++ NucleicApp.Tests/Fakes.cs | 126 ++++++++++ NucleicApp.Tests/NucleicApp.Tests.csproj | 22 ++ NucleicApp.Tests/ProtocolInteropTests.cs | 84 +++++++ NucleicApp.Tests/RendererCoreTests.cs | 120 ++++++++++ NucleicApp/NucleicApp.csproj | 14 ++ NucleicApp/Services/HostConnection.cs | 223 ++++++++++++++++++ NucleicApp/Services/HostdLauncher.cs | 101 ++++++++ NucleicApp/Services/OnboardingService.cs | 69 ++++++ NucleicApp/Services/RendererDispatcher.cs | 27 +++ NucleicApp/Services/RendererStore.cs | 155 ++++++++++++ NucleicApp/Services/RendezvousFile.cs | 13 + NucleicProtocol.Interop/ClientIntents.cs | 80 +++++++ NucleicProtocol.Interop/NativeProtocol.cs | 111 +++++++++ .../NucleicProtocol.Interop.csproj | 11 + NucleicProtocol.Interop/ProtocolClient.cs | 139 +++++++++++ NucleicProtocol.Interop/ProtocolModels.cs | 137 +++++++++++ NucleicProtocol.Interop/ProtocolProjection.cs | 41 ++++ 18 files changed, 1515 insertions(+) create mode 100644 NucleicApp.Tests/Fakes.cs create mode 100644 NucleicApp.Tests/NucleicApp.Tests.csproj create mode 100644 NucleicApp.Tests/ProtocolInteropTests.cs create mode 100644 NucleicApp.Tests/RendererCoreTests.cs create mode 100644 NucleicApp/NucleicApp.csproj create mode 100644 NucleicApp/Services/HostConnection.cs create mode 100644 NucleicApp/Services/HostdLauncher.cs create mode 100644 NucleicApp/Services/OnboardingService.cs create mode 100644 NucleicApp/Services/RendererDispatcher.cs create mode 100644 NucleicApp/Services/RendererStore.cs create mode 100644 NucleicApp/Services/RendezvousFile.cs create mode 100644 NucleicProtocol.Interop/ClientIntents.cs create mode 100644 NucleicProtocol.Interop/NativeProtocol.cs create mode 100644 NucleicProtocol.Interop/NucleicProtocol.Interop.csproj create mode 100644 NucleicProtocol.Interop/ProtocolClient.cs create mode 100644 NucleicProtocol.Interop/ProtocolModels.cs create mode 100644 NucleicProtocol.Interop/ProtocolProjection.cs diff --git a/Nucleic.sln b/Nucleic.sln index 515c3c4..48bfd97 100644 --- a/Nucleic.sln +++ b/Nucleic.sln @@ -7,6 +7,12 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NucleicBroker", "NucleicBro EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NucleicBroker.Tests", "NucleicBroker.Tests\NucleicBroker.Tests.csproj", "{513C933F-89FF-4918-87AB-E9B70625BE4B}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NucleicProtocol.Interop", "NucleicProtocol.Interop\NucleicProtocol.Interop.csproj", "{A13FB10F-03E9-4F2E-A934-324820A3EBC1}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NucleicApp", "NucleicApp\NucleicApp.csproj", "{B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NucleicApp.Tests", "NucleicApp.Tests\NucleicApp.Tests.csproj", "{C37B2431-0F92-4B61-8038-18B7C0BD7BAA}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -41,6 +47,42 @@ Global {513C933F-89FF-4918-87AB-E9B70625BE4B}.Release|x64.Build.0 = Release|Any CPU {513C933F-89FF-4918-87AB-E9B70625BE4B}.Release|x86.ActiveCfg = Release|Any CPU {513C933F-89FF-4918-87AB-E9B70625BE4B}.Release|x86.Build.0 = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|Any CPU.Build.0 = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|x64.ActiveCfg = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|x64.Build.0 = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|x86.ActiveCfg = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Debug|x86.Build.0 = Debug|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|Any CPU.ActiveCfg = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|Any CPU.Build.0 = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|x64.ActiveCfg = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|x64.Build.0 = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|x86.ActiveCfg = Release|Any CPU + {A13FB10F-03E9-4F2E-A934-324820A3EBC1}.Release|x86.Build.0 = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|Any CPU.Build.0 = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|x64.ActiveCfg = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|x64.Build.0 = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|x86.ActiveCfg = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Debug|x86.Build.0 = Debug|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|Any CPU.ActiveCfg = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|Any CPU.Build.0 = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|x64.ActiveCfg = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|x64.Build.0 = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|x86.ActiveCfg = Release|Any CPU + {B2721DF8-DDA6-4B35-98E8-94C5CA52FB2B}.Release|x86.Build.0 = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|Any CPU.Build.0 = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|x64.ActiveCfg = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|x64.Build.0 = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|x86.ActiveCfg = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Debug|x86.Build.0 = Debug|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|Any CPU.ActiveCfg = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|Any CPU.Build.0 = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|x64.ActiveCfg = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|x64.Build.0 = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|x86.ActiveCfg = Release|Any CPU + {C37B2431-0F92-4B61-8038-18B7C0BD7BAA}.Release|x86.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE diff --git a/NucleicApp.Tests/Fakes.cs b/NucleicApp.Tests/Fakes.cs new file mode 100644 index 0000000..aa374f1 --- /dev/null +++ b/NucleicApp.Tests/Fakes.cs @@ -0,0 +1,126 @@ +using System.Runtime.InteropServices; +using NucleicProtocol.Interop; + +namespace NucleicApp.Tests; + +internal sealed class FakeNativeProtocol : INativeProtocol +{ + private NativeProtocolCallback? eventCallback; + private NativeProtocolCallback? stateCallback; + private nint context; + + public string? ConfigurationJson { get; private set; } + public string? RendezvousJson { get; private set; } + public List Intents { get; } = []; + public bool Closed { get; private set; } + public string ProjectionResult { get; set; } = "{\"splices\":[],\"version\":1}"; + + public nint ClientCreate(string configurationJson) + { + ConfigurationJson = configurationJson; + return 11; + } + + public void ClientSetCallbacks( + nint handle, NativeProtocolCallback? eventCallback, + NativeProtocolCallback? stateCallback, nint context) + { + this.eventCallback = eventCallback; + this.stateCallback = stateCallback; + this.context = context; + } + + public int ClientConnectLocal(nint handle, string rendezvousJson) + { + RendezvousJson = rendezvousJson; + return 0; + } + + public int ClientConnectRemote(nint handle, string endpointJson) => 0; + public int ClientPair(nint handle, string pairingCode) => 0; + public int ClientSendIntent(nint handle, string clientMessageJson) + { + Intents.Add(clientMessageJson); + return 0; + } + + public void ClientClose(nint handle) => Closed = true; + public nint ProjectionCreate() => 22; + public nint ProjectionApply(nint projection, string hostMessageJson) => + Marshal.StringToCoTaskMemUTF8(ProjectionResult); + public void ProjectionFree(nint projection) { } + public string CopyAndFreeString(nint value) + { + try { return Marshal.PtrToStringUTF8(value)!; } + finally { Marshal.FreeCoTaskMem(value); } + } + + public void EmitEvent(string json) => Emit(eventCallback, json); + public void EmitState(string json) => Emit(stateCallback, json); + + private void Emit(NativeProtocolCallback? callback, string json) + { + var pointer = Marshal.StringToCoTaskMemUTF8(json); + try { callback?.Invoke(context, pointer); } + finally { Marshal.FreeCoTaskMem(pointer); } + } +} + +internal sealed class FakeProtocolClient : IProtocolClient +{ + public event Action? MessageJsonReceived; + public event Action? StateJsonReceived; + public HostdRendezvous? Rendezvous { get; private set; } + public List Intents { get; } = []; + public Queue ConnectLocalResults { get; } = []; + public int ConnectLocalCount { get; private set; } + public bool Disposed { get; private set; } + + public ProtocolResult ConnectLocal(HostdRendezvous rendezvous) + { + ConnectLocalCount++; + Rendezvous = rendezvous; + return ConnectLocalResults.TryDequeue(out var result) ? result : ProtocolResult.Ok; + } + public ProtocolResult ConnectRemote(string endpointJson) => ProtocolResult.Ok; + public ProtocolResult Pair(string pairingCode) => ProtocolResult.Ok; + public ProtocolResult SendIntent(string clientMessageJson) + { + Intents.Add(clientMessageJson); + return ProtocolResult.Ok; + } + public void Dispose() => Disposed = true; + public void EmitMessage(string json) => MessageJsonReceived?.Invoke(json); + public void EmitState(string json) => StateJsonReceived?.Invoke(json); +} + +internal sealed class FakeProtocolClientFactory(FakeProtocolClient client) : IProtocolClientFactory +{ + public ProtocolClientConfiguration? Configuration { get; private set; } + public IProtocolClient Create(ProtocolClientConfiguration configuration) + { + Configuration = configuration; + return client; + } +} + +internal static class TestRepository +{ + public static string Root + { + get + { + var current = new DirectoryInfo(AppContext.BaseDirectory); + while (current is not null) + { + if (Directory.Exists(Path.Combine(current.FullName, "fixtures", "protocol-abi"))) + return current.FullName; + current = current.Parent; + } + throw new DirectoryNotFoundException("could not locate repository fixtures"); + } + } + + public static string Fixture(string side, string name) => File.ReadAllText( + Path.Combine(Root, "fixtures", "protocol-abi", "v1", side, name + ".json")); +} diff --git a/NucleicApp.Tests/NucleicApp.Tests.csproj b/NucleicApp.Tests/NucleicApp.Tests.csproj new file mode 100644 index 0000000..7a6ef31 --- /dev/null +++ b/NucleicApp.Tests/NucleicApp.Tests.csproj @@ -0,0 +1,22 @@ + + + + false + + + + + + + + + + + + + + + + + + diff --git a/NucleicApp.Tests/ProtocolInteropTests.cs b/NucleicApp.Tests/ProtocolInteropTests.cs new file mode 100644 index 0000000..4fb9233 --- /dev/null +++ b/NucleicApp.Tests/ProtocolInteropTests.cs @@ -0,0 +1,84 @@ +using System.Text.Json; +using NucleicProtocol.Interop; + +namespace NucleicApp.Tests; + +public sealed class ProtocolInteropTests +{ + [Fact] + public void EveryGoldenJsonFixtureHasTheFilenameDiscriminator() + { + var root = Path.Combine(TestRepository.Root, "fixtures", "protocol-abi", "v1"); + var clients = Directory.GetFiles(Path.Combine(root, "client"), "*.json"); + var hosts = Directory.GetFiles(Path.Combine(root, "host"), "*.json"); + + Assert.Equal(65, clients.Length); + Assert.Equal(52, hosts.Length); + foreach (var path in clients.Concat(hosts)) + { + var expected = Path.GetFileNameWithoutExtension(path); + var message = ProtocolMessage.Parse(File.ReadAllText(path)); + Assert.Equal(expected, message.Tag); + } + } + + [Fact] + public void ManagedClientPinsCallbackContextUntilNativeCloseReturns() + { + var native = new FakeNativeProtocol(); + var messages = new List(); + using (var client = new ProtocolClient(new ProtocolClientConfiguration + { + IdentityDirectory = @"C:\Users\me\AppData\Roaming\Nucleic-local\renderer", + Channel = "local", + }, native)) + { + client.MessageJsonReceived += messages.Add; + native.EmitEvent("{\"t\":\"pong\"}"); + var result = client.ConnectLocal(ValidRendezvous()); + Assert.Equal(ProtocolResult.Ok, result); + Assert.Contains("\"hostID\":\"host-1\"", native.RendezvousJson); + } + + Assert.True(native.Closed); + Assert.Single(messages, "{\"t\":\"pong\"}"); + } + + [Fact] + public void ProjectionCopiesThenFreesNativeString() + { + var native = new FakeNativeProtocol + { + ProjectionResult = "{\"sessionID\":\"s1\",\"splices\":[],\"version\":1}", + }; + using var projection = new ProtocolProjection(native); + var diff = projection.Apply("{\"t\":\"pong\"}"); + Assert.Equal("s1", diff.Root.GetProperty("sessionID").GetString()); + } + + [Fact] + public void IntentHelpersMatchFrozenWireShapes() + { + AssertJsonEqual(TestRepository.Fixture("client", "listSessions"), ClientIntents.ListSessions()); + AssertJsonEqual(TestRepository.Fixture("client", "unsubscribe"), ClientIntents.Unsubscribe("session-1")); + var expectedSend = "{\"input\":{\"parts\":[{\"text\":\"continue\",\"type\":\"text\"}]},\"sessionID\":\"session-1\",\"t\":\"sendInput\"}"; + AssertJsonEqual(expectedSend, ClientIntents.SendInput("session-1", "continue")); + } + + internal static HostdRendezvous ValidRendezvous() => new() + { + ProcessId = Environment.ProcessId, + Host = "127.0.0.1", + Port = 51900, + LocalPsk = Convert.ToBase64String(new byte[32]), + HostId = "host-1", + }; + + private static void AssertJsonEqual(string expected, string actual) + { + using var expectedDocument = JsonDocument.Parse(expected); + using var actualDocument = JsonDocument.Parse(actual); + Assert.True(JsonElement.DeepEquals(expectedDocument.RootElement, actualDocument.RootElement), + $"expected {expected}\nactual {actual}"); + } +} diff --git a/NucleicApp.Tests/RendererCoreTests.cs b/NucleicApp.Tests/RendererCoreTests.cs new file mode 100644 index 0000000..2c67742 --- /dev/null +++ b/NucleicApp.Tests/RendererCoreTests.cs @@ -0,0 +1,120 @@ +using System.Net; +using System.Net.Sockets; +using System.Text.Json; +using NucleicApp.Services; +using NucleicProtocol.Interop; + +namespace NucleicApp.Tests; + +public sealed class RendererCoreTests +{ + [Fact] + public void StoreProjectsSessionsApprovalsAndSettingsFromHostMessagesOnly() + { + using var store = new RendererStore(); + store.Apply(TestRepository.Fixture("host", "sessionList")); + store.Apply(TestRepository.Fixture("host", "approvalRequested")); + store.Apply(TestRepository.Fixture("host", "dashboard")); + store.Apply(TestRepository.Fixture("host", "settings")); + + var session = Assert.Single(store.Sessions); + Assert.Equal("session-1", session.Id); + Assert.Equal("awaitingApproval", session.Status); + Assert.Equal(1, session.PendingApprovalCount); + Assert.Equal("approval-1", Assert.Single(store.PendingApprovals).Id); + Assert.NotNull(store.Dashboard); + Assert.True(store.Settings!.Value.GetProperty("resolveApprovalsViaCloud").GetBoolean()); + + store.Apply(TestRepository.Fixture("host", "approvalResolved")); + Assert.Empty(store.PendingApprovals); + } + + [Fact] + public async Task HostConnectionReadsRendezvousAndRequestsInitialStateOnReady() + { + var directory = Directory.CreateTempSubdirectory("nucleic-renderer-"); + try + { + var rendezvousPath = Path.Combine(directory.FullName, "hostd.json"); + await File.WriteAllTextAsync(rendezvousPath, JsonSerializer.Serialize( + ProtocolInteropTests.ValidRendezvous())); + using var store = new RendererStore(); + var client = new FakeProtocolClient(); + var factory = new FakeProtocolClientFactory(client); + await using var connection = new HostConnection( + store, Path.Combine(directory.FullName, "identity"), factory, + InlineRendererDispatcher.Instance); + + await connection.StartLocalAsync(rendezvousPath); + Assert.Equal("host-1", client.Rendezvous?.HostId); + client.EmitState("{\"state\":\"ready\",\"hostID\":\"host-1\",\"welcome\":{}}"); + Assert.Equal(HostConnectionState.Connected, connection.State); + + var tags = client.Intents.Select(json => ProtocolMessage.Parse(json).Tag).ToArray(); + Assert.Equal(["listSessions", "listDashboard", "listPeers"], tags); + + client.EmitMessage(TestRepository.Fixture("host", "sessionUpdated")); + Assert.Equal("session-1", Assert.Single(store.Sessions).Id); + } + finally { directory.Delete(recursive: true); } + } + + [Fact] + public async Task HostConnectionKeepsRetryingAfterSynchronousAdmissionFailures() + { + var directory = Directory.CreateTempSubdirectory("nucleic-reconnect-"); + try + { + var rendezvousPath = Path.Combine(directory.FullName, "hostd.json"); + await File.WriteAllTextAsync(rendezvousPath, JsonSerializer.Serialize( + ProtocolInteropTests.ValidRendezvous())); + using var store = new RendererStore(); + var client = new FakeProtocolClient(); + client.ConnectLocalResults.Enqueue(ProtocolResult.InvalidState); + client.ConnectLocalResults.Enqueue(ProtocolResult.InvalidState); + client.ConnectLocalResults.Enqueue(ProtocolResult.Ok); + await using var connection = new HostConnection( + store, Path.Combine(directory.FullName, "identity"), + new FakeProtocolClientFactory(client), InlineRendererDispatcher.Instance); + + await connection.StartLocalAsync(rendezvousPath); + await WaitUntilAsync(() => client.ConnectLocalCount == 3, TimeSpan.FromSeconds(3)); + + Assert.Equal(3, client.ConnectLocalCount); + Assert.Equal(HostConnectionState.Reconnecting, connection.State); + } + finally { directory.Delete(recursive: true); } + } + + [Fact] + public async Task LauncherAcceptsOnlyALiveLoopbackRendezvous() + { + var directory = Directory.CreateTempSubdirectory("nucleic-launcher-"); + var listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + try + { + var record = ProtocolInteropTests.ValidRendezvous() with + { + Port = ((IPEndPoint)listener.LocalEndpoint).Port, + }; + var path = Path.Combine(directory.FullName, "hostd.json"); + await File.WriteAllTextAsync(path, JsonSerializer.Serialize(record)); + var launcher = new HostdLauncher(new HostdLaunchOptions( + Path.Combine(directory.FullName, "must-not-launch.exe"), path)); + var result = await launcher.EnsureRunningAsync(); + Assert.Equal(record.Port, result.Port); + } + finally + { + listener.Stop(); + directory.Delete(recursive: true); + } + } + + private static async Task WaitUntilAsync(Func condition, TimeSpan timeout) + { + using var deadline = new CancellationTokenSource(timeout); + while (!condition()) await Task.Delay(10, deadline.Token); + } +} diff --git a/NucleicApp/NucleicApp.csproj b/NucleicApp/NucleicApp.csproj new file mode 100644 index 0000000..ad814a7 --- /dev/null +++ b/NucleicApp/NucleicApp.csproj @@ -0,0 +1,14 @@ + + + + + NucleicApp + NucleicApp.Core + + + + + + + diff --git a/NucleicApp/Services/HostConnection.cs b/NucleicApp/Services/HostConnection.cs new file mode 100644 index 0000000..b9979c2 --- /dev/null +++ b/NucleicApp/Services/HostConnection.cs @@ -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? 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 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", + }; +} diff --git a/NucleicApp/Services/HostdLauncher.cs b/NucleicApp/Services/HostdLauncher.cs new file mode 100644 index 0000000..1c727f9 --- /dev/null +++ b/NucleicApp/Services/HostdLauncher.cs @@ -0,0 +1,101 @@ +using System.Diagnostics; +using System.Net.Sockets; +using NucleicProtocol.Interop; + +namespace NucleicApp.Services; + +public sealed record HostdLaunchOptions( + string ExecutablePath, + string RendezvousPath, + string? WorkingDirectory = null, + TimeSpan? StartupTimeout = null); + +/// Starts the channel-matched full-trust host and waits for a live owner-only rendezvous. +public sealed class HostdLauncher +{ + private readonly HostdLaunchOptions options; + + public HostdLauncher(HostdLaunchOptions options) => this.options = options; + + public async Task EnsureRunningAsync( + CancellationToken cancellationToken = default) + { + if (await TryHealthyRendezvousAsync(cancellationToken).ConfigureAwait(false) is { } healthy) + return healthy; + + Process? launched = null; + if (!HostMutexExists()) + { + if (!File.Exists(options.ExecutablePath)) + throw new FileNotFoundException("nucleic-hostd was not found", options.ExecutablePath); + launched = Process.Start(new ProcessStartInfo + { + FileName = options.ExecutablePath, + WorkingDirectory = options.WorkingDirectory + ?? Path.GetDirectoryName(options.ExecutablePath) ?? Environment.CurrentDirectory, + UseShellExecute = false, + CreateNoWindow = true, + }) ?? throw new InvalidOperationException("Windows refused to start nucleic-hostd"); + } + + var timeout = options.StartupTimeout ?? TimeSpan.FromSeconds(30); + using var deadline = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + deadline.CancelAfter(timeout); + while (true) + { + deadline.Token.ThrowIfCancellationRequested(); + if (launched is { HasExited: true }) + throw new InvalidOperationException($"nucleic-hostd exited during startup ({launched.ExitCode})"); + if (await TryHealthyRendezvousAsync(deadline.Token).ConfigureAwait(false) is { } ready) + return ready; + await Task.Delay(100, deadline.Token).ConfigureAwait(false); + } + } + + private async Task TryHealthyRendezvousAsync(CancellationToken cancellationToken) + { + HostdRendezvous record; + try { record = await RendezvousFile.ReadAsync(options.RendezvousPath, cancellationToken); } + catch (Exception error) when (error is IOException or UnauthorizedAccessException or System.Text.Json.JsonException) + { + return null; + } + + if (!ProcessIsAlive(record.ProcessId)) return null; + try + { + using var client = new TcpClient(); + using var probe = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + probe.CancelAfter(TimeSpan.FromSeconds(1)); + await client.ConnectAsync(record.Host, record.Port, probe.Token).ConfigureAwait(false); + return record; + } + catch (Exception error) when (error is SocketException or OperationCanceledException) + { + return null; + } + } + + private static bool ProcessIsAlive(int processId) + { + try { return !Process.GetProcessById(processId).HasExited; } + catch (ArgumentException) { return false; } + } + + private static bool HostMutexExists() + { + if (!OperatingSystem.IsWindows()) return false; + try + { + if (!Mutex.TryOpenExisting("Local\\nucleic-hostd" + BuildChannel.Suffix, out var mutex)) + return false; + mutex.Dispose(); + return true; + } + catch (UnauthorizedAccessException) + { + // Existence without access is still existence; hostd's own guard remains definitive. + return true; + } + } +} diff --git a/NucleicApp/Services/OnboardingService.cs b/NucleicApp/Services/OnboardingService.cs new file mode 100644 index 0000000..2af3720 --- /dev/null +++ b/NucleicApp/Services/OnboardingService.cs @@ -0,0 +1,69 @@ +namespace NucleicApp.Services; + +public enum OnboardingPhase +{ + Idle, + StartingHost, + Connecting, + Ready, + Failed, +} + +/// First-run coordinator for the item-11 core. Component installation and OAuth screens plug +/// into the same phase surface once their host intents land; host launch + authenticated local +/// sync are real today. +public sealed class OnboardingService +{ + private readonly HostdLauncher launcher; + private readonly HostConnection connection; + + public OnboardingPhase Phase { get; private set; } = OnboardingPhase.Idle; + public string? Failure { get; private set; } + public event Action? PhaseChanged; + + public OnboardingService(HostdLauncher launcher, HostConnection connection) + { + this.launcher = launcher; + this.connection = connection; + } + + public async Task StartAsync(string rendezvousPath, CancellationToken cancellationToken = default) + { + try + { + SetPhase(OnboardingPhase.StartingHost); + _ = await launcher.EnsureRunningAsync(cancellationToken).ConfigureAwait(false); + SetPhase(OnboardingPhase.Connecting); + + var ready = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + void Changed(HostConnectionState state) + { + if (state == HostConnectionState.Connected) ready.TrySetResult(); + else if (state == HostConnectionState.Failed) + ready.TrySetException(new InvalidOperationException( + connection.LastProtocolState?.Message ?? "local host connection failed")); + } + connection.StateChanged += Changed; + try + { + await connection.StartLocalAsync(rendezvousPath, cancellationToken).ConfigureAwait(false); + await ready.Task.WaitAsync(TimeSpan.FromSeconds(30), cancellationToken).ConfigureAwait(false); + } + finally { connection.StateChanged -= Changed; } + SetPhase(OnboardingPhase.Ready); + } + catch (Exception error) + { + Failure = error.Message; + SetPhase(OnboardingPhase.Failed); + throw; + } + } + + private void SetPhase(OnboardingPhase phase) + { + if (Phase == phase) return; + Phase = phase; + PhaseChanged?.Invoke(phase); + } +} diff --git a/NucleicApp/Services/RendererDispatcher.cs b/NucleicApp/Services/RendererDispatcher.cs new file mode 100644 index 0000000..55de052 --- /dev/null +++ b/NucleicApp/Services/RendererDispatcher.cs @@ -0,0 +1,27 @@ +namespace NucleicApp.Services; + +/// Isolates WinUI's DispatcherQueue from the renderer state machine and keeps tests portable. +public interface IRendererDispatcher +{ + void Post(Action action); +} + +public sealed class SynchronizationContextDispatcher : IRendererDispatcher +{ + private readonly SynchronizationContext context; + + public SynchronizationContextDispatcher(SynchronizationContext? context = null) + { + this.context = context ?? SynchronizationContext.Current + ?? throw new InvalidOperationException("renderer has no SynchronizationContext"); + } + + public void Post(Action action) => context.Post(static state => ((Action)state!).Invoke(), action); +} + +public sealed class InlineRendererDispatcher : IRendererDispatcher +{ + public static InlineRendererDispatcher Instance { get; } = new(); + private InlineRendererDispatcher() { } + public void Post(Action action) => action(); +} diff --git a/NucleicApp/Services/RendererStore.cs b/NucleicApp/Services/RendererStore.cs new file mode 100644 index 0000000..2852aa9 --- /dev/null +++ b/NucleicApp/Services/RendererStore.cs @@ -0,0 +1,155 @@ +using System.ComponentModel; +using System.Runtime.CompilerServices; +using System.Text.Json; +using NucleicProtocol.Interop; + +namespace NucleicApp.Services; + +/// Read-only projection of host authority (WINDOWS_PORT §7). Every mutation enters through a +/// HostMsg; user actions leave as ClientMsg intents through ``HostConnection``. +public sealed class RendererStore : INotifyPropertyChanged, IDisposable +{ + private readonly Dictionary sessions = new(StringComparer.Ordinal); + private readonly Dictionary approvals = new(StringComparer.Ordinal); + private readonly ITranscriptProjector? projector; + + private IReadOnlyList sessionSnapshot = []; + private IReadOnlyList approvalSnapshot = []; + private JsonElement? dashboard; + private JsonElement? settings; + private JsonElement? peers; + private JsonElement? intelligenceCatalog; + private string? lastError; + + public RendererStore(ITranscriptProjector? projector = null) => this.projector = projector; + + public IReadOnlyList Sessions => sessionSnapshot; + public IReadOnlyList PendingApprovals => approvalSnapshot; + public JsonElement? Dashboard => dashboard; + public JsonElement? Settings => settings; + public JsonElement? Peers => peers; + public JsonElement? IntelligenceCatalog => intelligenceCatalog; + public string? LastError => lastError; + + public event PropertyChangedEventHandler? PropertyChanged; + public event Action? ProjectionChanged; + public event Action? MessageApplied; + + public void Apply(string hostMessageJson) => Apply(ProtocolMessage.Parse(hostMessageJson)); + + public void Apply(ProtocolMessage message) + { + var root = message.Root; + switch (message.Tag) + { + case "sessionList": + sessions.Clear(); + if (root.TryGetProperty("input", out var list)) + foreach (var item in list.EnumerateArray()) UpsertSession(item); + PublishSessions(); + break; + case "sessionUpdated": + UpsertSession(root); + PublishSessions(); + break; + case "snapshot": + if (root.TryGetProperty("summary", out var summary)) + { + UpsertSession(summary); + PublishSessions(); + } + if (root.TryGetProperty("pendingApprovals", out var pending)) + { + foreach (var item in pending.EnumerateArray()) UpsertApproval(item); + PublishApprovals(); + } + break; + case "approvalRequested": + UpsertApproval(root); + PublishApprovals(); + break; + case "approvalResolved": + if (String(root, "id") is { } resolved) approvals.Remove(resolved); + PublishApprovals(); + break; + case "dashboard": + dashboard = root.Clone(); + Changed(nameof(Dashboard)); + break; + case "settings": + settings = root.GetProperty("settings").Clone(); + Changed(nameof(Settings)); + break; + case "peerList": + peers = root.GetProperty("input").Clone(); + Changed(nameof(Peers)); + break; + case "intelligenceCatalog": + intelligenceCatalog = root.GetProperty("intelligenceCatalog").Clone(); + Changed(nameof(IntelligenceCatalog)); + break; + case "error": + lastError = String(root, "message") ?? "The host reported an error."; + Changed(nameof(LastError)); + break; + } + + if (projector is not null) + { + var diff = projector.Apply(message.Json); + if (diff.Root.GetProperty("splices").GetArrayLength() > 0) + ProjectionChanged?.Invoke(diff); + } + MessageApplied?.Invoke(message); + } + + private void UpsertSession(JsonElement json) + { + var id = String(json, "sessionID") ?? throw new JsonException("session summary has no sessionID"); + sessions[id] = new SessionSummary( + id, String(json, "title") ?? "Untitled", String(json, "projectID"), + String(json, "projectName"), String(json, "status") ?? "unknown", + Int(json, "pendingApprovalCount"), json.Clone()); + } + + private void UpsertApproval(JsonElement json) + { + var id = String(json, "id") ?? throw new JsonException("approval has no id"); + approvals[id] = new ApprovalRequest( + id, String(json, "sessionID") ?? string.Empty, String(json, "title") ?? "Approval", + String(json, "toolName") ?? string.Empty, String(json, "risk") ?? string.Empty, + json.Clone()); + } + + private void PublishSessions() + { + sessionSnapshot = sessions.Values + .OrderByDescending(item => item.Raw.TryGetProperty("updatedAt", out var value) + ? value.GetDouble() : 0) + .ToArray(); + Changed(nameof(Sessions)); + } + + private void PublishApprovals() + { + approvalSnapshot = approvals.Values.ToArray(); + Changed(nameof(PendingApprovals)); + } + + private static string? String(JsonElement root, string name) => + root.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.String + ? value.GetString() : null; + private static int Int(JsonElement root, string name) => + root.TryGetProperty(name, out var value) && value.TryGetInt32(out var number) ? number : 0; + private void Changed([CallerMemberName] string? name = null) => + PropertyChanged?.Invoke(this, new PropertyChangedEventArgs(name)); + + public void Dispose() => projector?.Dispose(); +} + +public sealed record SessionSummary( + string Id, string Title, string? ProjectId, string? ProjectName, + string Status, int PendingApprovalCount, JsonElement Raw); + +public sealed record ApprovalRequest( + string Id, string SessionId, string Title, string ToolName, string Risk, JsonElement Raw); diff --git a/NucleicApp/Services/RendezvousFile.cs b/NucleicApp/Services/RendezvousFile.cs new file mode 100644 index 0000000..421a3d8 --- /dev/null +++ b/NucleicApp/Services/RendezvousFile.cs @@ -0,0 +1,13 @@ +using NucleicProtocol.Interop; + +namespace NucleicApp.Services; + +public static class RendezvousFile +{ + public static async Task ReadAsync( + string path, CancellationToken cancellationToken = default) + { + var json = await File.ReadAllTextAsync(path, cancellationToken).ConfigureAwait(false); + return ProtocolJson.DeserializeRendezvous(json); + } +} diff --git a/NucleicProtocol.Interop/ClientIntents.cs b/NucleicProtocol.Interop/ClientIntents.cs new file mode 100644 index 0000000..adc5e25 --- /dev/null +++ b/NucleicProtocol.Interop/ClientIntents.cs @@ -0,0 +1,80 @@ +using System.Text.Encodings.Web; +using System.Text.Json; + +namespace NucleicProtocol.Interop; + +/// Small typed constructors for the renderer's first intents. Swift remains the decoder and +/// source of truth; these helpers only keep hand-written UI JSON out of call sites. +public static class ClientIntents +{ + public static string ListSessions() => TagOnly("listSessions"); + public static string ListDashboard() => TagOnly("listDashboard"); + public static string ListPeers() => TagOnly("listPeers"); + public static string Ping() => TagOnly("ping"); + + public static string Subscribe(string sessionId, ulong? sinceSequence = null, string verbosity = "full") => + Write(writer => + { + writer.WriteString("sessionID", sessionId); + if (sinceSequence is { } sequence) writer.WriteNumber("sinceSeq", sequence); + writer.WriteString("t", "subscribe"); + writer.WriteString("verbosity", verbosity); + }); + + public static string Unsubscribe(string sessionId) => Write(writer => + { + writer.WriteString("sessionID", sessionId); + writer.WriteString("t", "unsubscribe"); + }); + + public static string ApprovalRespond( + string approvalId, string decisionType, JsonElement? updatedInput = null) => + Write(writer => + { + writer.WriteString("approvalID", approvalId); + writer.WritePropertyName("decision"); + writer.WriteStartObject(); + writer.WriteString("type", decisionType); + if (updatedInput is { } input) + { + writer.WritePropertyName("updatedInput"); + input.WriteTo(writer); + } + writer.WriteEndObject(); + writer.WriteString("t", "approvalRespond"); + }); + + public static string SendInput(string sessionId, string text) => Write(writer => + { + writer.WritePropertyName("input"); + writer.WriteStartObject(); + writer.WritePropertyName("parts"); + writer.WriteStartArray(); + writer.WriteStartObject(); + writer.WriteString("text", text); + writer.WriteString("type", "text"); + writer.WriteEndObject(); + writer.WriteEndArray(); + writer.WriteEndObject(); + writer.WriteString("sessionID", sessionId); + writer.WriteString("t", "sendInput"); + }); + + private static string TagOnly(string tag) => JsonSerializer.Serialize( + new TagOnlyIntent(tag), ProtocolJsonContext.Default.TagOnlyIntent); + + private static string Write(Action body) + { + using var stream = new MemoryStream(); + using (var writer = new Utf8JsonWriter(stream, new JsonWriterOptions + { + Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping, + })) + { + writer.WriteStartObject(); + body(writer); + writer.WriteEndObject(); + } + return System.Text.Encoding.UTF8.GetString(stream.GetBuffer(), 0, checked((int)stream.Length)); + } +} diff --git a/NucleicProtocol.Interop/NativeProtocol.cs b/NucleicProtocol.Interop/NativeProtocol.cs new file mode 100644 index 0000000..653254c --- /dev/null +++ b/NucleicProtocol.Interop/NativeProtocol.cs @@ -0,0 +1,111 @@ +using System.Runtime.InteropServices; + +namespace NucleicProtocol.Interop; + +public enum ProtocolResult +{ + Ok = 0, + InvalidArgument = -1, + InvalidState = -2, + DecodeError = -3, + Unsupported = -4, +} + +[UnmanagedFunctionPointer(CallingConvention.Cdecl)] +public delegate void NativeProtocolCallback(nint context, nint utf8Json); + +/// Injectable native surface: production delegates to the Swift DLL; tests use an in-memory +/// implementation and exercise callback lifetime without loading a platform binary. +public interface INativeProtocol +{ + nint ClientCreate(string configurationJson); + void ClientSetCallbacks( + nint handle, NativeProtocolCallback? eventCallback, + NativeProtocolCallback? stateCallback, nint context); + int ClientConnectLocal(nint handle, string rendezvousJson); + int ClientConnectRemote(nint handle, string endpointJson); + int ClientPair(nint handle, string pairingCode); + int ClientSendIntent(nint handle, string clientMessageJson); + void ClientClose(nint handle); + nint ProjectionCreate(); + nint ProjectionApply(nint projection, string hostMessageJson); + void ProjectionFree(nint projection); + string CopyAndFreeString(nint value); +} + +internal sealed class PInvokeNativeProtocol : INativeProtocol +{ + public static PInvokeNativeProtocol Instance { get; } = new(); + private PInvokeNativeProtocol() { } + + public nint ClientCreate(string json) => NativeMethods.np_client_create(json); + public void ClientSetCallbacks( + nint handle, NativeProtocolCallback? eventCallback, + NativeProtocolCallback? stateCallback, nint context) => + NativeMethods.np_client_set_callbacks(handle, eventCallback, stateCallback, context); + public int ClientConnectLocal(nint handle, string json) => + NativeMethods.np_client_connect_local(handle, json); + public int ClientConnectRemote(nint handle, string json) => + NativeMethods.np_client_connect_remote(handle, json); + public int ClientPair(nint handle, string code) => NativeMethods.np_client_pair(handle, code); + public int ClientSendIntent(nint handle, string json) => + NativeMethods.np_client_send_intent(handle, json); + public void ClientClose(nint handle) => NativeMethods.np_client_close(handle); + public nint ProjectionCreate() => NativeMethods.np_projection_create(); + public nint ProjectionApply(nint projection, string json) => + NativeMethods.np_projection_apply(projection, json); + public void ProjectionFree(nint projection) => NativeMethods.np_projection_free(projection); + + public string CopyAndFreeString(nint value) + { + if (value == 0) throw new InvalidOperationException("native protocol returned NULL"); + try { return Marshal.PtrToStringUTF8(value) ?? string.Empty; } + finally { NativeMethods.np_free(value); } + } + + private static class NativeMethods + { + private const string Dll = "NucleicProtocolC"; + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern nint np_client_create( + [MarshalAs(UnmanagedType.LPUTF8Str)] string configurationJson); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern void np_client_set_callbacks( + nint handle, NativeProtocolCallback? eventCallback, + NativeProtocolCallback? stateCallback, nint context); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern int np_client_connect_local( + nint handle, [MarshalAs(UnmanagedType.LPUTF8Str)] string rendezvousJson); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern int np_client_connect_remote( + nint handle, [MarshalAs(UnmanagedType.LPUTF8Str)] string endpointJson); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern int np_client_pair( + nint handle, [MarshalAs(UnmanagedType.LPUTF8Str)] string pairingCode); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern int np_client_send_intent( + nint handle, [MarshalAs(UnmanagedType.LPUTF8Str)] string clientMessageJson); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern void np_client_close(nint handle); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern nint np_projection_create(); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern nint np_projection_apply( + nint projection, [MarshalAs(UnmanagedType.LPUTF8Str)] string hostMessageJson); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern void np_projection_free(nint projection); + + [DllImport(Dll, CallingConvention = CallingConvention.Cdecl)] + internal static extern void np_free(nint value); + } +} diff --git a/NucleicProtocol.Interop/NucleicProtocol.Interop.csproj b/NucleicProtocol.Interop/NucleicProtocol.Interop.csproj new file mode 100644 index 0000000..2d7194d --- /dev/null +++ b/NucleicProtocol.Interop/NucleicProtocol.Interop.csproj @@ -0,0 +1,11 @@ + + + + + NucleicProtocol.Interop + NucleicProtocol.Interop + true + + + diff --git a/NucleicProtocol.Interop/ProtocolClient.cs b/NucleicProtocol.Interop/ProtocolClient.cs new file mode 100644 index 0000000..faeb7c2 --- /dev/null +++ b/NucleicProtocol.Interop/ProtocolClient.cs @@ -0,0 +1,139 @@ +using System.Runtime.InteropServices; +using System.Text.Json; + +namespace NucleicProtocol.Interop; + +/// Managed lifetime and callback barrier over one `np_handle`. +public interface IProtocolClient : IDisposable +{ + event Action? MessageJsonReceived; + event Action? StateJsonReceived; + ProtocolResult ConnectLocal(HostdRendezvous rendezvous); + ProtocolResult ConnectRemote(string endpointJson); + ProtocolResult Pair(string pairingCode); + ProtocolResult SendIntent(string clientMessageJson); +} + +public sealed class ProtocolClient : IProtocolClient +{ + private static readonly NativeProtocolCallback EventThunk = OnNativeEvent; + private static readonly NativeProtocolCallback StateThunk = OnNativeState; + + private readonly object gate = new(); + private readonly INativeProtocol native; + private nint handle; + private GCHandle callbackContext; + private bool disposed; + + public event Action? MessageJsonReceived; + public event Action? StateJsonReceived; + + public ProtocolClient( + ProtocolClientConfiguration configuration, INativeProtocol? native = null) + { + this.native = native ?? PInvokeNativeProtocol.Instance; + var json = JsonSerializer.Serialize( + configuration, ProtocolJsonContext.Default.ProtocolClientConfiguration); + handle = this.native.ClientCreate(json); + if (handle == 0) throw new InvalidOperationException("NucleicProtocolC rejected client configuration"); + + callbackContext = GCHandle.Alloc(this, GCHandleType.Normal); + try + { + this.native.ClientSetCallbacks( + handle, EventThunk, StateThunk, GCHandle.ToIntPtr(callbackContext)); + } + catch + { + this.native.ClientClose(handle); + handle = 0; + callbackContext.Free(); + throw; + } + } + + public ProtocolResult ConnectLocal(HostdRendezvous rendezvous) + { + rendezvous.Validate(); + return Invoke(handle => native.ClientConnectLocal( + handle, JsonSerializer.Serialize(rendezvous, ProtocolJsonContext.Default.HostdRendezvous))); + } + + public ProtocolResult ConnectRemote(string endpointJson) => + Invoke(handle => native.ClientConnectRemote(handle, endpointJson)); + + public ProtocolResult Pair(string pairingCode) => + Invoke(handle => native.ClientPair(handle, pairingCode)); + + public ProtocolResult SendIntent(string clientMessageJson) + { + // Fail malformed UI output before it reaches the ABI. The Swift decoder remains the + // authority for the actual ClientMsg shape and returns DecodeError for a valid-but-wrong + // document. + using (JsonDocument.Parse(clientMessageJson)) { } + return Invoke(handle => native.ClientSendIntent(handle, clientMessageJson)); + } + + private ProtocolResult Invoke(Func operation) + { + lock (gate) + { + ObjectDisposedException.ThrowIf(disposed, this); + return (ProtocolResult)operation(handle); + } + } + + public void Dispose() + { + nint closing; + lock (gate) + { + if (disposed) return; + disposed = true; + closing = handle; + handle = 0; + } + + // The Swift contract makes close a callback barrier. Only after it returns is it safe to + // free the GCHandle that native callbacks carry as their opaque context. + native.ClientClose(closing); + MessageJsonReceived = null; + StateJsonReceived = null; + if (callbackContext.IsAllocated) callbackContext.Free(); + } + + private static void OnNativeEvent(nint context, nint utf8Json) => + Dispatch(context, utf8Json, static (client, json) => client.MessageJsonReceived?.Invoke(json)); + + private static void OnNativeState(nint context, nint utf8Json) => + Dispatch(context, utf8Json, static (client, json) => client.StateJsonReceived?.Invoke(json)); + + private static void Dispatch( + nint context, nint utf8Json, Action deliver) + { + if (context == 0 || utf8Json == 0) return; + try + { + if (GCHandle.FromIntPtr(context).Target is ProtocolClient client) + { + var json = Marshal.PtrToStringUTF8(utf8Json); + if (json is not null) deliver(client, json); + } + } + catch + { + // Nothing may unwind through the C callback boundary. Renderer event subscribers own + // their diagnostics; a bad observer must not tear down Swift's callback thread. + } + } +} + +public interface IProtocolClientFactory +{ + IProtocolClient Create(ProtocolClientConfiguration configuration); +} + +public sealed class ProtocolClientFactory : IProtocolClientFactory +{ + public IProtocolClient Create(ProtocolClientConfiguration configuration) => new ProtocolClient(configuration); +} diff --git a/NucleicProtocol.Interop/ProtocolModels.cs b/NucleicProtocol.Interop/ProtocolModels.cs new file mode 100644 index 0000000..ffabff7 --- /dev/null +++ b/NucleicProtocol.Interop/ProtocolModels.cs @@ -0,0 +1,137 @@ +using System.Text.Json; +using System.Text.Json.Serialization; + +namespace NucleicProtocol.Interop; + +public sealed record ProtocolClientConfiguration +{ + [JsonPropertyName("identityDir")] + public required string IdentityDirectory { get; init; } + + [JsonPropertyName("hostKeyPins")] + public IReadOnlyDictionary? HostKeyPins { get; init; } + + [JsonPropertyName("deviceID")] + public string? DeviceId { get; init; } + + [JsonPropertyName("deviceLabel")] + public string DeviceLabel { get; init; } = "Nucleic for Windows"; + + [JsonPropertyName("scopeClaim")] + public string ScopeClaim { get; init; } = "control"; + + [JsonPropertyName("channel")] + public string? Channel { get; init; } +} + +/// The owner-only `run/hostd.json` written by nucleic-hostd. +public sealed record HostdRendezvous +{ + [JsonPropertyName("pid")] + public required int ProcessId { get; init; } + + [JsonPropertyName("host")] + public string Host { get; init; } = "127.0.0.1"; + + [JsonPropertyName("port")] + public required int Port { get; init; } + + [JsonPropertyName("localPSK")] + public required string LocalPsk { get; init; } + + /// Selects the DLL's persisted pin on renderer relaunch. A fresh identity has no matching + /// pin and falls back to the one-time local PSK. + [JsonPropertyName("hostID")] + public required string HostId { get; init; } + + public void Validate() + { + if (ProcessId <= 0) throw new JsonException("rendezvous pid must be positive"); + if (Host is not ("127.0.0.1" or "::1" or "localhost")) + throw new JsonException("renderer rendezvous must be loopback-only"); + if (Port is <= 0 or > 65_535) throw new JsonException("rendezvous port is invalid"); + if (string.IsNullOrWhiteSpace(HostId)) throw new JsonException("rendezvous hostID is missing"); + byte[] secret; + try { secret = Convert.FromBase64String(LocalPsk); } + catch (FormatException error) { throw new JsonException("rendezvous localPSK is invalid", error); } + if (secret.Length != 32) throw new JsonException("rendezvous localPSK must be 32 bytes"); + } +} + +/// A forward-compatible HostMsg view. The tagged enum remains owned by Swift; C# reads the +/// discriminator and keeps the complete JSON value so a newer field or message never gets lost. +public sealed record ProtocolMessage(string Tag, JsonElement Root, string Json) +{ + public static ProtocolMessage Parse(string json) + { + using var document = JsonDocument.Parse(json); + if (document.RootElement.ValueKind != JsonValueKind.Object || + !document.RootElement.TryGetProperty("t", out var tagValue) || + tagValue.ValueKind != JsonValueKind.String) + { + throw new JsonException("protocol message has no string 't' discriminator"); + } + return new ProtocolMessage(tagValue.GetString()!, document.RootElement.Clone(), json); + } +} + +public sealed record ProtocolState( + string State, string? Message, string? Transport, string? Peer, + string? HostId, JsonElement? Welcome, string Json) +{ + public static ProtocolState Parse(string json) + { + using var document = JsonDocument.Parse(json); + var root = document.RootElement; + if (!root.TryGetProperty("state", out var state) || state.ValueKind != JsonValueKind.String) + throw new JsonException("protocol state callback has no state"); + return new ProtocolState( + state.GetString()!, OptionalString(root, "message"), + OptionalString(root, "transport"), OptionalString(root, "peer"), + OptionalString(root, "hostID"), + root.TryGetProperty("welcome", out var welcome) ? welcome.Clone() : null, + json); + } + + private static string? OptionalString(JsonElement root, string name) => + root.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.String + ? value.GetString() : null; +} + +public sealed record ProjectionDiff(JsonElement Root, string Json) +{ + public static ProjectionDiff Parse(string json) + { + using var document = JsonDocument.Parse(json); + var root = document.RootElement; + if (root.ValueKind != JsonValueKind.Object || + !root.TryGetProperty("version", out var version) || version.GetInt32() != 1 || + !root.TryGetProperty("splices", out var splices) || + splices.ValueKind != JsonValueKind.Array) + { + throw new JsonException("invalid transcript projection result"); + } + return new ProjectionDiff(root.Clone(), json); + } +} + +internal sealed record TagOnlyIntent([property: JsonPropertyName("t")] string Tag); + +[JsonSourceGenerationOptions( + PropertyNamingPolicy = JsonKnownNamingPolicy.CamelCase, + DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull)] +[JsonSerializable(typeof(ProtocolClientConfiguration))] +[JsonSerializable(typeof(HostdRendezvous))] +[JsonSerializable(typeof(TagOnlyIntent))] +internal partial class ProtocolJsonContext : JsonSerializerContext; + +public static class ProtocolJson +{ + public static HostdRendezvous DeserializeRendezvous(string json) + { + var value = JsonSerializer.Deserialize(json, ProtocolJsonContext.Default.HostdRendezvous) + ?? throw new JsonException("empty hostd rendezvous"); + value.Validate(); + return value; + } +} diff --git a/NucleicProtocol.Interop/ProtocolProjection.cs b/NucleicProtocol.Interop/ProtocolProjection.cs new file mode 100644 index 0000000..8f26961 --- /dev/null +++ b/NucleicProtocol.Interop/ProtocolProjection.cs @@ -0,0 +1,41 @@ +namespace NucleicProtocol.Interop; + +/// Managed ownership wrapper for the shared Swift transcript projector. +public interface ITranscriptProjector : IDisposable +{ + ProjectionDiff Apply(string hostMessageJson); +} + +public sealed class ProtocolProjection : ITranscriptProjector +{ + private readonly object gate = new(); + private readonly INativeProtocol native; + private nint handle; + + public ProtocolProjection(INativeProtocol? native = null) + { + this.native = native ?? PInvokeNativeProtocol.Instance; + handle = this.native.ProjectionCreate(); + if (handle == 0) throw new InvalidOperationException("could not create transcript projection"); + } + + public ProjectionDiff Apply(string hostMessageJson) + { + lock (gate) + { + ObjectDisposedException.ThrowIf(handle == 0, this); + var value = native.ProjectionApply(handle, hostMessageJson); + return ProjectionDiff.Parse(native.CopyAndFreeString(value)); + } + } + + public void Dispose() + { + lock (gate) + { + if (handle == 0) return; + native.ProjectionFree(handle); + handle = 0; + } + } +}