From 4f12b18d63253d6dfeb34da92bc6145159efd9c1 Mon Sep 17 00:00:00 2001 From: Andrew Blakeslee Moore Date: Sat, 4 Jul 2026 17:56:31 -0700 Subject: [PATCH] Mesh P3 (multiplexer, WIP): extract HostConnection engine MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Begins the simultaneous multiplexer. HostConnection.swift is the per-host connection engine lifted out of RemoteStore: one instance owns a single Mac's SyncClient, LAN→tailnet candidate chain, reconnect backoff, tailnet-node lifecycle, and event stream, and keeps that Mac's projection (connectivity, sessions, dashboard, capabilities, meshPeers). It talks back to its owner through a Callbacks struct so aggregate concerns — the badge, Live Activity, notifications, and the single open transcript — can be merged across hosts by the coordinator. Not yet wired: RemoteStore still runs its own inline single-connection machine. The next step swaps RemoteStore onto a `[HostID: HostConnection]` map (connect all paired Macs, aggregate their sessions, route intents by host) — the behavior-critical part, to be verified with two real Macs. Compiles (iOS Simulator BUILD SUCCEEDED); no behavior change yet (the engine is unused until the rewire). Also: RemoteStore.friendlyTransportError made non-private so the engine shares it. Co-Authored-By: Claude Opus 4.8 --- .../NucleicRemote/Models/HostConnection.swift | 460 ++++++++++++++++++ .../NucleicRemote/Models/RemoteStore.swift | 2 +- 2 files changed, 461 insertions(+), 1 deletion(-) create mode 100644 NucleicRemote/NucleicRemote/Models/HostConnection.swift diff --git a/NucleicRemote/NucleicRemote/Models/HostConnection.swift b/NucleicRemote/NucleicRemote/Models/HostConnection.swift new file mode 100644 index 0000000..c067131 --- /dev/null +++ b/NucleicRemote/NucleicRemote/Models/HostConnection.swift @@ -0,0 +1,460 @@ +import Foundation +import Network +import NucleicProtocol +import NucleicTailnet +#if canImport(UIKit) +import UIKit +#endif + +/// One phone→Mac connection and its projected state (mesh P3 multiplexer foundation). +/// +/// This is the per-host connection engine extracted from `RemoteStore`: it owns the `SyncClient`, +/// the LAN→tailnet candidate chain, reconnect backoff, and the event stream for a single paired +/// Mac, and it keeps that Mac's projection (connectivity, sessions, dashboard, capabilities, …). +/// `RemoteStore` will own one of these per paired host (`[HostID: HostConnection]`), aggregating +/// their state and routing intents to the right one; aggregate concerns (badge, Live Activity, +/// notifications, the single open transcript) are surfaced through `Callbacks` so `RemoteStore` +/// can merge across hosts. A phone in demo mode has no `HostConnection` — demo seeds state directly. +@MainActor +final class HostConnection { + /// The Mac this connection targets, by `HostID` (its static-key fingerprint). + let hostID: String + + // MARK: Projected state (this host's slice of what the phone shows) + + private(set) var hostName: String + private(set) var connectivity: RemoteStore.Connectivity = .connecting + private(set) var sessions: [WireSessionSummary] = [] + private(set) var capabilities = WireCapabilities(canModifyToolInput: false, allowAlwaysScopes: []) + private(set) var grantedScope: DeviceScope = .approve + private(set) var modelCatalog: WireModelCatalog = .empty + private(set) var dashboard = DashboardSnapshot.empty + private(set) var meshPeers: [PeerSummary] = [] + /// The live transport, for the "Connected · …" chip. + private(set) var activeTransport: SyncTransportHint = .lan + + /// The session RemoteStore currently has open on this host, if any — so snapshot/events are only + /// forwarded (and deduped) for the transcript on screen. Set by RemoteStore on open/close. + var openSessionID: SessionID? + + // MARK: Callbacks up to RemoteStore (aggregate concerns) + + /// How a `HostConnection` talks back to `RemoteStore`. All fire on the main actor. + struct Callbacks { + /// This host's projected state changed — refresh the aggregate list/dashboard/chip. + var didUpdate: () -> Void = {} + /// The open session's transcript snapshot arrived (only when `openSessionID` matches). + var openSnapshot: (SessionSnapshot) -> Void = { _ in } + /// New transcript events for the open session. + var openEvents: (EventBatch) -> Void = { _ in } + /// The open session's full diff arrived. + var openDiff: (WireSessionDiff) -> Void = { _ in } + /// An approval was requested (with the resolved session title) — post the notification and, + /// if it's the open session, add it to the approval strip. + var approvalRequested: (ApprovalRequest, String) -> Void = { _, _ in } + /// An approval resolved anywhere — withdraw the notification and clear the strip. + var approvalResolved: (ApprovalResolved) -> Void = { _ in } + /// A session became "waiting on you" (for a background notification). + var sessionBecameWaiting: (WireSessionSummary) -> Void = { _ in } + /// A non-fatal host error to surface as a transient bubble. + var wireError: (WireError) -> Void = { _ in } + /// A pairing handshake succeeded — persist the pinned host record. + var didPair: (PairedHost) -> Void = { _ in } + /// The embedded Tailscale node's status changed (for Settings). + var tailnetStatus: (String?, URL?) -> Void = { _, _ in } + } + private let callbacks: Callbacks + + // MARK: Shared dependencies + + private let identity: DeviceIdentity + private let discovery: LANDiscovery + + // MARK: Connection machine (moved from RemoteStore, one per host) + + private var client: SyncClient? + private var eventTask: Task? + private var connectTask: Task? + private var retryTask: Task? + private var lanConnectTimeout: Task? + private var reconnectAttempts = 0 + private var seenSeq: Set = [] + + private enum TransportAttempt { + case lan(NWEndpoint) + case tailnet(host: String, port: UInt16) + } + + private struct ConnectPlan { + var remaining: [TransportAttempt] + let hostStaticKey: Data + let mode: SyncClient.Mode + let deviceID: String + let pairingPayload: PairingPayload? + } + private var connectPlan: ConnectPlan? + + init( + hostID: String, hostName: String, + identity: DeviceIdentity, discovery: LANDiscovery, callbacks: Callbacks + ) { + self.hostID = hostID + self.hostName = hostName + self.identity = identity + self.discovery = discovery + self.callbacks = callbacks + } + + // MARK: - Connect (pair / reconnect) + + /// Pair from a scanned QR (SYNC §4.2): try the QR's transports in order, run XXpsk0, and on + /// success pin the host key for future IK reconnects. + func pair(with payload: PairingPayload) { + teardown() + connectivity = .connecting + hostName = payload.hostName + callbacks.didUpdate() + guard let hint = payload.transportHint else { + fail("This pairing code needs a newer version of Nucleic Remote.") + return + } + guard hint != .relay else { + fail("Relay connections aren't supported yet.") + return + } + let candidates = buildCandidates( + fingerprint: payload.hostStaticKey.fingerprintHex, + lanHost: payload.lanHost, lanPort: payload.lanPort, + tailnet: hint == .tailnet ? (payload.tailnetHost, payload.tailnetPort) : nil) + guard !candidates.isEmpty else { + fail(hint == .tailnet && !TailnetSupport.isBuiltIn + ? TailnetError.notBuiltIn.errorDescription ?? "Tailscale support isn't built in" + : "No Mac found on this network") + return + } + connectPlan = ConnectPlan( + remaining: candidates, hostStaticKey: payload.hostStaticKey, + mode: .pair(secret: payload.pairingSecret), deviceID: IdentityStore.deviceID(), + pairingPayload: payload) + _ = tryNextCandidate() + } + + /// Reconnect to the pinned host using IK: LAN when reachable, else the pairing's tailnet hint. + func reconnect(to host: PairedHost) { + teardown() + pinnedHost = host + connectivity = reconnectAttempts == 0 ? .connecting : .reconnecting + hostName = host.hostName + callbacks.didUpdate() + guard host.transportHint != .relay else { + fail("Relay connections aren't supported yet.") + return + } + let candidates = buildCandidates( + fingerprint: host.fingerprint, + lanHost: host.lanHost, lanPort: host.lanPort, + tailnet: host.transportHint == .tailnet ? (host.tailnetHost, host.tailnetPort) : nil) + guard !candidates.isEmpty else { + connectivity = .hostOffline + callbacks.didUpdate() + scheduleRetry(host: host) + return + } + connectPlan = ConnectPlan( + remaining: candidates, hostStaticKey: host.hostStaticKey, + mode: .reconnect, deviceID: host.deviceID, pairingPayload: nil) + _ = tryNextCandidate() + } + + /// The host this connection reconnects to (for its retry loop). Set on `reconnect(to:)`. + private var pinnedHost: PairedHost? + + private func fail(_ message: String) { + connectivity = .failed(message) + callbacks.didUpdate() + } + + private func buildCandidates( + fingerprint: String?, lanHost: String?, lanPort: UInt16?, + tailnet: (host: String?, port: UInt16?)? + ) -> [TransportAttempt] { + var candidates: [TransportAttempt] = [] + if let endpoint = discovery.endpoint(forFingerprint: fingerprint, lanHost: lanHost, lanPort: lanPort) { + candidates.append(.lan(endpoint)) + } + if let tailnet, let host = tailnet.host, let port = tailnet.port, TailnetSupport.isBuiltIn { + candidates.append(.tailnet(host: host, port: port)) + } + return candidates + } + + private func tryNextCandidate() -> Bool { + guard var plan = connectPlan, !plan.remaining.isEmpty else { return false } + let next = plan.remaining.removeFirst() + connectPlan = plan + attempt(next, plan: plan) + return true + } + + private func attempt(_ candidate: TransportAttempt, plan: ConnectPlan) { + teardownClient() + switch candidate { + case .lan(let endpoint): + activeTransport = .lan + let channel = makeChannel(endpoint) + lanConnectTimeout?.cancel() + lanConnectTimeout = Task { [weak channel] in + try? await Task.sleep(for: .seconds(4)) + guard !Task.isCancelled, let channel, !channel.isReady else { return } + channel.close() + } + startClient( + channel: channel, hostStaticKey: plan.hostStaticKey, + mode: plan.mode, deviceID: plan.deviceID, pairingPayload: plan.pairingPayload) + case .tailnet(let host, let port): + activeTransport = .tailnet + connectTask = Task { [weak self] in + guard let self else { return } + do { + let channel = try await self.tailnetChannel(host: host, port: port) + guard !Task.isCancelled else { channel.close(); return } + self.startClient( + channel: channel, hostStaticKey: plan.hostStaticKey, + mode: plan.mode, deviceID: plan.deviceID, pairingPayload: plan.pairingPayload) + } catch { + guard !Task.isCancelled else { return } + self.tailnetAttemptFailed(error, isPairing: plan.pairingPayload != nil) + } + } + } + } + + private func tailnetAttemptFailed(_ error: Error, isPairing: Bool) { + if tryNextCandidate() { return } + if isPairing { + fail(error.localizedDescription) + return + } + switch error { + case TailnetError.notBuiltIn, TailnetError.notConfigured: + fail(error.localizedDescription) + default: + connectivity = .hostOffline + callbacks.didUpdate() + scheduleRetry(host: pinnedHost) + } + } + + private func startClient( + channel: any FrameChannel, hostStaticKey: Data, mode: SyncClient.Mode, + deviceID: String, pairingPayload: PairingPayload? + ) { + let client = SyncClient( + channel: channel, identity: identity, hostStaticKey: hostStaticKey, + mode: mode, deviceID: deviceID, + deviceLabel: UIDevice.current.name, pushToken: PushRegistrar.shared.tokenHex, + releaseChannel: BuildInfo.current.channel.releaseChannel) + self.client = client + consume(client, pairingPayload: pairingPayload) + } + + private func tailnetChannel(host: String, port: UInt16) async throws -> FDFrameChannel { + guard TailnetSupport.isBuiltIn else { throw TailnetError.notBuiltIn } + let config = Self.phoneTailnetConfig() + callbacks.tailnetStatus("Starting…", nil) + let watcher = Task { [weak self] in + for await status in await TailnetNode.shared.statusStream() { + guard let self, !Task.isCancelled else { break } + var loginURL: URL? + if case .needsLogin(let url) = status, let u = URL(string: url) { + loginURL = u + // Eject to Safari only for user-initiated pairing (the user is watching). + if self.connectPlan?.pairingPayload != nil { + UIApplication.shared.open(u, options: [:], completionHandler: nil) + } + } + self.callbacks.tailnetStatus(status.label, loginURL) + } + } + defer { + watcher.cancel() + callbacks.tailnetStatus(nil, nil) + } + do { + try await TailnetNode.shared.ensureRunning(config: config) + } catch TailnetError.timedOut(let message) { + callbacks.tailnetStatus(await TailnetNode.shared.status.label, nil) + throw TailnetError.notConfigured(message) + } catch { + callbacks.tailnetStatus(await TailnetNode.shared.status.label, nil) + throw error + } + callbacks.tailnetStatus(await TailnetNode.shared.status.label, nil) + return try await TailnetNode.shared.dial(host: host, port: port) + } + + private static func phoneTailnetConfig() -> TailnetConfig { + let base = FileManager.default.urls(for: .applicationSupportDirectory, in: .userDomainMask)[0] + .appendingPathComponent("Nucleic", isDirectory: true) + .appendingPathComponent("tailnet", isDirectory: true) + return TailnetConfig( + hostName: TailnetConfig.nodeName(for: UIDevice.current.name), + stateDirectory: base, + authKey: TailnetAuthStore.loadAuthKey()) + } + + private func makeChannel(_ endpoint: NWEndpoint) -> NWFrameChannel { + let channel = NWFrameChannel(endpoint: endpoint) + channel.onFailed = { [weak self] error in + Task { @MainActor in self?.handleTransportFailure(error) } + } + return channel + } + + private func handleTransportFailure(_ error: String) { + guard !connectivity.isLive else { return } + guard connectPlan?.remaining.isEmpty != false else { return } + connectivity = .failed(RemoteStore.friendlyTransportError(error)) + callbacks.didUpdate() + } + + // MARK: - Event stream + + private func consume(_ client: SyncClient, pairingPayload: PairingPayload?) { + eventTask = Task { [weak self] in + let stream = await client.start() + for await event in stream { + guard let self, self.client === client else { break } + await self.handle(event, pairingPayload: pairingPayload) + } + } + } + + private func handle(_ event: SyncClient.Event, pairingPayload: PairingPayload?) async { + switch event { + case .connecting: + break + case .ready(let welcome): + reconnectAttempts = 0 + connectPlan = nil + connectivity = .connected(activeTransport) + hostName = welcome.host.hostName + capabilities = welcome.capabilities + grantedScope = welcome.grantedScope + modelCatalog = welcome.modelCatalog + if let payload = pairingPayload, let hostKey = await client?.hostKey() { + callbacks.didPair(PairedHost( + deviceID: IdentityStore.deviceID(), hostName: welcome.host.hostName, + hostStaticKey: hostKey, fingerprint: hostKey.fingerprintHex, + lanHost: payload.lanHost, lanPort: payload.lanPort, + transport: payload.transport, tailnetHost: payload.tailnetHost, + tailnetPort: payload.tailnetPort)) + } + callbacks.didUpdate() + send(.listSessions) + send(.listDashboard) + if let id = openSessionID { send(.subscribe(Subscribe(sessionID: id, sinceSeq: nil, verbosity: .full))) } + case .sessionList(let list): + sessions = list + callbacks.didUpdate() + case .sessionUpdated(let summary): + let previous: WireSessionSummary? + if let i = sessions.firstIndex(where: { $0.sessionID == summary.sessionID }) { + previous = sessions[i] + sessions[i] = summary + } else { + previous = nil + sessions.append(summary) + } + let becameWaiting = summary.status == .awaitingInput && previous?.status != .awaitingInput + if becameWaiting, !summary.archived { callbacks.sessionBecameWaiting(summary) } + callbacks.didUpdate() + case .dashboard(let snapshot): + dashboard = snapshot + callbacks.didUpdate() + case .snapshot(let snapshot): + guard snapshot.summary.sessionID == openSessionID else { break } + seenSeq = Set(snapshot.recentEvents.map(\.seq)) + callbacks.openSnapshot(snapshot) + case .events(let batch): + guard batch.sessionID == openSessionID else { break } + let fresh = batch.events.filter { !seenSeq.contains($0.seq) } + for e in fresh { seenSeq.insert(e.seq) } + if !fresh.isEmpty { callbacks.openEvents(EventBatch(sessionID: batch.sessionID, events: fresh)) } + case .approvalRequested(let req): + let title = sessions.first { $0.sessionID == req.sessionID }?.title ?? "Approval" + callbacks.approvalRequested(req, title) + case .approvalResolved(let resolved): + callbacks.approvalResolved(resolved) + case .sessionDiff(let diff): + guard diff.sessionID == openSessionID else { break } + callbacks.openDiff(diff) + case .peerList(let peers): + meshPeers = peers + callbacks.didUpdate() + case .transferAccept, .transferReject, .transferReady, .transferCommitted, .transferChunkAck: + // Session-transfer replies (mesh P5) only reach a *source* Mac; a phone is never one. + break + case .wireError(let error): + if error.code == .channelMismatch { + connectivity = .failed(error.message) + callbacks.didUpdate() + break + } + guard error.code != .alreadyResolved else { break } + guard connectPlan == nil else { break } + callbacks.wireError(error) + case .failed(let message): + if tryNextCandidate() { break } + connectPlan = nil + if case .failed = connectivity {} else { connectivity = .failed(message) } + callbacks.didUpdate() + scheduleRetry(host: pinnedHost) + case .closed: + if !connectivity.isLive, tryNextCandidate() { break } + if connectivity.isLive { + connectivity = .reconnecting + } else if connectPlan?.pairingPayload != nil { + if case .failed = connectivity {} else { + connectivity = .failed("Couldn't connect to your Mac — check that it's reachable, then scan again.") + } + } + connectPlan = nil + callbacks.didUpdate() + scheduleRetry(host: pinnedHost) + } + } + + // MARK: - Send / retry / teardown + + func send(_ msg: ClientMsg) { + guard let client else { return } + Task { await client.send(msg) } + } + + private func scheduleRetry(host: PairedHost?) { + guard let host else { return } + reconnectAttempts += 1 + let delay = min(Double(reconnectAttempts) * 1.5, 10) + retryTask?.cancel() + retryTask = Task { [weak self] in + try? await Task.sleep(for: .seconds(delay)) + guard !Task.isCancelled, let self, !self.connectivity.isLive else { return } + self.reconnect(to: host) + } + } + + func teardown() { + retryTask?.cancel(); retryTask = nil + connectTask?.cancel(); connectTask = nil + connectPlan = nil + teardownClient() + } + + private func teardownClient() { + lanConnectTimeout?.cancel(); lanConnectTimeout = nil + eventTask?.cancel(); eventTask = nil + if let client { Task { await client.disconnect() } } + client = nil + } +} diff --git a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift index 64b4baf..795dff7 100644 --- a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift +++ b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift @@ -1076,7 +1076,7 @@ final class RemoteStore: ObservableObject { /// Turn a raw `NWError` string into a short, fixable hint. The common onboarding failures are /// a denied Local Network permission (EPERM / -65555) and the Mac not listening (refused). - private static func friendlyTransportError(_ error: String) -> String { + static func friendlyTransportError(_ error: String) -> String { let lower = error.lowercased() if lower.contains("denied") || lower.contains("not permitted") || lower.contains("65555") { return "Can't reach the local network. In Settings ▸ Nucleic, allow Local Network access, then try again."