Merge branch 'dev' into nucleic/trunk

This commit is contained in:
2026-07-04 20:15:01 -07:00
2 changed files with 643 additions and 460 deletions
@@ -0,0 +1,459 @@
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<Void, Never>?
private var connectTask: Task<Void, Never>?
private var retryTask: Task<Void, Never>?
private var lanConnectTimeout: Task<Void, Never>?
private var reconnectAttempts = 0
private var seenSeq: Set<UInt64> = []
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)
}
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
}
}
@@ -101,8 +101,12 @@ final class RemoteStore: ObservableObject {
/// UX_IOS §11.5, and the host dedupes/`alreadyResolved`s a lost race anyway).
func respondFromNotification(_ id: ApprovalID, allow: Bool) {
let decision: Decision = allow ? .allow(updatedInput: nil) : .deny(reason: nil)
if connectivity.isLive {
send(.approvalRespond(id, decision))
// The approval may belong to any connected Mac (mesh P3) and the notification carries no
// hostID — broadcast to every live connection; the owning host resolves it, the rest see an
// unknown/`alreadyResolved` approval and no-op.
let live = connections.values.filter { $0.connectivity.isLive }
if !live.isEmpty {
for conn in live { conn.send(.approvalRespond(id, decision)) }
} else {
pendingNotificationDecision = (id, decision, Date())
reconnect()
@@ -114,9 +118,11 @@ final class RemoteStore: ObservableObject {
private func flushPendingNotificationDecision() {
guard let (id, decision, at) = pendingNotificationDecision else { return }
let live = connections.values.filter { $0.connectivity.isLive }
guard !live.isEmpty else { return } // wait until something's connected
pendingNotificationDecision = nil
guard Date().timeIntervalSince(at) < 30 else { return } // stale — require the app
send(.approvalRespond(id, decision))
for conn in live { conn.send(.approvalRespond(id, decision)) }
}
private static let lastOpenedKey = "nucleic.lastOpenedAt"
@@ -153,22 +159,19 @@ final class RemoteStore: ObservableObject {
private let identity = IdentityStore.loadOrCreateIdentity()
let discovery = LANDiscovery()
private var client: SyncClient?
private var eventTask: Task<Void, Never>?
/// In-flight async connection setup (tailnet node start + dial); cancelled on teardown.
private var connectTask: Task<Void, Never>?
/// The pending reconnect backoff timer; cancelled on teardown so a stale retry can't
/// tear down a newer in-flight attempt.
private var retryTask: Task<Void, Never>?
/// The transport the current connection attempt uses (drives the "Connected · …" chip).
private var activeTransport: SyncTransportHint = .lan
/// The phone's embedded Tailscale node state, for Settings (nil = not running).
/// One live connection per paired Mac (mesh P3 multiplexer). All connected at once; the flat
/// `sessions`/`connectivity`/`capabilities`/… above mirror the *active* one (`activeHostID`), so
/// switching hosts is instant (no reconnect). Empty in demo mode (demo seeds state directly).
private var connections: [String: HostConnection] = [:]
/// The active connection — the Mac whose projection the flat state reflects + whose transcript
/// is open.
private var activeConnection: HostConnection? { activeHostID.flatMap { connections[$0] } }
/// The phone's embedded Tailscale node state, for Settings (nil = not running). Shared across all
/// connections (one embedded node, many dials).
@Published private(set) var tailnetStatus: String?
/// The Tailscale interactive-login URL while the node waits for a browser login (a
/// first tailnet start with no auth key). Auto-opened; Settings shows a re-open button.
@Published private(set) var tailnetLoginURL: URL?
private var seenSeq: Set<UInt64> = []
private var reconnectAttempts = 0
/// Offline demo mode: seeds mock state and simulates the agent locally so every surface
/// renders — and the core loops (send, approve, start chat, to-dos) actually respond —
@@ -205,8 +208,9 @@ final class RemoteStore: ObservableObject {
}
}
/// Switch which paired Mac is active. Live: re-point the (single) connection via the existing
/// reconnect machinery. Demo: swap the mock host's projection, preserving in-demo edits.
/// Switch which paired Mac is active. Live: every Mac is already connected, so this is instant —
/// close the current transcript, re-point at the target connection, and mirror its state. Demo:
/// swap the mock host's projection, preserving in-demo edits.
func switchHost(to id: String) {
guard id != activeHostID else { return }
if demoMode {
@@ -217,17 +221,19 @@ final class RemoteStore: ObservableObject {
guard let next = demoHosts.first(where: { $0.id == id }) else { return }
applyDemoHost(next)
} else {
// Close the transcript open on the old host.
activeConnection?.openSessionID = nil
openSessionID = nil; openEvents = []; openApprovals = []; openDiff = nil; diffLoading = false
activeHostID = id
reconnect()
// The target should already be connected; (re)dial only if it isn't.
if let host = IdentityStore.pairedHost(id: id) {
let conn = connection(for: host)
if !conn.connectivity.isLive { conn.reconnect(to: host) }
}
mirrorActive()
}
}
/// The paired Mac the live connection targets — the explicitly-active one, else the most-recent.
private func activeHost() -> PairedHost? {
if let id = activeHostID, let host = IdentityStore.pairedHost(id: id) { return host }
return IdentityStore.pairedHosts().last
}
/// Mock hosts for the multi-host demo. Each carries its own name + sessions + dashboard;
/// switching swaps the flat active-host projection. Empty outside demo.
private struct DemoHost {
@@ -405,252 +411,162 @@ final class RemoteStore: ObservableObject {
""")
}
/// One dialable way to reach the Mac. A connect builds an ordered candidate list — LAN
/// first (cheapest when reachable), then the tailnet (SYNC §3.2's LAN-then-fallback
/// ordering) — and `attempt` walks it until one carries a session.
private enum TransportAttempt {
case lan(NWEndpoint)
case tailnet(host: String, port: UInt16)
}
// MARK: - Connections (mesh P3 multiplexer)
/// The in-progress connect: remaining candidates plus everything needed to start a
/// client on whichever one succeeds. Cleared on `.ready` (chain done) and by teardown.
private struct ConnectPlan {
var remaining: [TransportAttempt]
let hostStaticKey: Data
let mode: SyncClient.Mode
let deviceID: String
let pairingPayload: PairingPayload?
}
private var connectPlan: ConnectPlan?
/// Kills a LAN attempt whose TCP connect just hangs (stale IP hint) so the chain can
/// move on — NWConnection's own timeout is far too slow for a fallback decision.
private var lanConnectTimeout: Task<Void, Never>?
/// Pair from a scanned QR (SYNC §4.2): try the QR's transports in order (LAN hint or
/// Bonjour first, then the Mac's tailnet IP via the phone's embedded node), run XXpsk0,
/// and on success pin the host key for future IK reconnects.
/// Pair a newly-scanned Mac, make it the active host, and start its connection — while any
/// existing connections keep running.
func pair(with payload: PairingPayload) {
teardown()
connectivity = .connecting
hostName = payload.hostName
activeHostID = payload.hostStaticKey.fingerprintHex
guard let hint = payload.transportHint else {
connectivity = .failed("This pairing code needs a newer version of Nucleic Remote.")
return
}
guard hint != .relay else {
connectivity = .failed("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 {
connectivity = .failed(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()
let id = payload.hostStaticKey.fingerprintHex
connections[id]?.teardown()
let conn = makeConnection(hostID: id, hostName: payload.hostName)
connections[id] = conn
activeHostID = id
mirrorActive()
conn.pair(with: payload)
}
/// Reconnect to the already-paired host using IK against the pinned static key: LAN
/// when reachable, else the pairing's tailnet hint — so a phone that leaves the Mac's
/// Wi‑Fi rolls over to the tailnet and rolls back when it returns.
/// Get or build the connection for a paired Mac.
private func connection(for host: PairedHost) -> HostConnection {
if let existing = connections[host.fingerprint] { return existing }
let conn = makeConnection(hostID: host.fingerprint, hostName: host.hostName)
connections[host.fingerprint] = conn
return conn
}
private func makeConnection(hostID: String, hostName: String) -> HostConnection {
HostConnection(
hostID: hostID, hostName: hostName, identity: identity, discovery: discovery,
callbacks: callbacks(for: hostID))
}
/// The callbacks one `HostConnection` uses to drive shared/aggregate state. `hostID` is captured
/// so a change updates the flat projection only for the active host, while badges, notifications,
/// and the single open transcript are handled globally across hosts.
private func callbacks(for hostID: String) -> HostConnection.Callbacks {
var cb = HostConnection.Callbacks()
cb.didUpdate = { [weak self] in
guard let self else { return }
if hostID == self.activeHostID { self.mirrorActive() }
self.refreshAggregate()
self.flushPendingNotificationDecision()
}
cb.openSnapshot = { [weak self] snap in
guard let self, hostID == self.activeHostID, snap.summary.sessionID == self.openSessionID else { return }
self.openEvents = snap.recentEvents
self.openApprovals = snap.pendingApprovals
}
cb.openEvents = { [weak self] batch in
guard let self, hostID == self.activeHostID, batch.sessionID == self.openSessionID else { return }
self.openEvents.append(contentsOf: batch.events)
}
cb.openDiff = { [weak self] diff in
guard let self, hostID == self.activeHostID, diff.sessionID == self.openSessionID else { return }
self.openDiff = diff
self.diffLoading = false
}
cb.approvalRequested = { [weak self] req, title in
guard let self else { return }
if hostID == self.activeHostID, req.sessionID == self.openSessionID,
!self.openApprovals.contains(where: { $0.id == req.id }) {
self.openApprovals.append(req)
}
// Post for any host; the router suppresses the banner if the user is on this session.
NotificationRouter.shared.postApproval(req, sessionTitle: title)
}
cb.approvalResolved = { [weak self] resolved in
guard let self else { return }
self.openApprovals.removeAll { $0.id == resolved.id }
NotificationRouter.shared.withdrawApproval(resolved.id)
}
cb.sessionBecameWaiting = { [weak self] summary in
guard let self, !self.isActive else { return }
NotificationRouter.shared.postSessionUpdate(summary)
}
cb.wireError = { [weak self] error in
self?.showError(error.message, sessionID: error.sessionID)
}
cb.didPair = { [weak self] host in
guard let self else { return }
IdentityStore.savePairedHost(host)
self.activeHostID = host.fingerprint
}
cb.tailnetStatus = { [weak self] status, loginURL in
self?.tailnetStatus = status
self?.tailnetLoginURL = loginURL
}
return cb
}
/// Mirror the active connection's projection onto the flat @Published state the UI binds to.
private func mirrorActive() {
guard let c = activeConnection else { return }
connectivity = c.connectivity
hostName = c.hostName
sessions = c.sessions
capabilities = c.capabilities
grantedScope = c.grantedScope
modelCatalog = c.modelCatalog
dashboard = c.dashboard
meshPeers = c.meshPeers
}
/// Every connected Mac's live (non-archived) sessions — the aggregate the app-icon badge and
/// Live Activity summarize across hosts.
private var allLiveSessions: [WireSessionSummary] {
if demoMode { return demoHosts.flatMap(\.sessions).filter { !$0.archived } }
return connections.values.flatMap(\.sessions).filter { !$0.archived }
}
/// Refresh cross-host aggregates: the app-icon badge (needs-you across *all* Macs) + the Live
/// Activity summary.
private func refreshAggregate() {
let needs = allLiveSessions.filter { $0.status.needsYou($0.disposition) }.count
NotificationRouter.shared.updateBadge(needs)
LiveActivityManager.shared.sync(hostName: hostName, sessions: allLiveSessions)
}
/// Tear down every live connection (leaving the paired registry intact) — for entering demo.
private func teardownAll() {
for conn in connections.values { conn.teardown() }
connections.removeAll()
}
/// Connect to every paired Mac at once (mesh P3 multiplexer): each `HostConnection` runs its own
/// IK reconnect (LAN→tailnet, with backoff), so all your Macs are live simultaneously and the
/// switcher flips between them instantly. Idempotent — an already-live connection is left alone;
/// an offline one (re)dials. Connections for since-unpaired Macs are dropped. The launch +
/// "Reconnect" path.
func reconnect() {
guard let host = activeHost() else { connectivity = .unpaired; return }
teardown()
activeHostID = host.fingerprint
connectivity = reconnectAttempts == 0 ? .connecting : .reconnecting
hostName = host.hostName
guard host.transportHint != .relay else {
connectivity = .failed("Relay connections aren't supported yet.")
return
let hosts = IdentityStore.pairedHosts()
guard !hosts.isEmpty else { connectivity = .unpaired; return }
let paired = Set(hosts.map(\.fingerprint))
for (id, conn) in connections where !paired.contains(id) { conn.teardown(); connections[id] = nil }
for host in hosts {
let conn = connection(for: host)
if !conn.connectivity.isLive { conn.reconnect(to: host) }
}
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
scheduleRetry()
return
if activeHostID == nil || connections[activeHostID!] == nil {
activeHostID = hosts.last?.fingerprint
}
connectPlan = ConnectPlan(
remaining: candidates, hostStaticKey: host.hostStaticKey,
mode: .reconnect, deviceID: host.deviceID, pairingPayload: nil)
_ = tryNextCandidate()
mirrorActive()
refreshAggregate()
}
/// LAN first (explicit hint, else a Bonjour match), tailnet second when the pairing
/// carries one and this build can dial it.
private func buildCandidates(
fingerprint: String?, lanHost: String?, lanPort: UInt16?,
tailnet: (host: String?, port: UInt16?)?
) -> [TransportAttempt] {
var candidates: [TransportAttempt] = []
if let endpoint = resolveEndpoint(fingerprint: 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
}
/// Pop and dial the next candidate. False when the plan is exhausted (or gone) — the
/// caller then applies its terminal failure handling.
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)
// Close the channel if TCP isn't up within the window (a stale IP hint would
// otherwise hang the chain on NWConnection's slow timeout); the finished stream
// then advances to the next candidate. The timer checks readiness itself — it's
// bound to exactly this channel, so a stale timer can never hit a later attempt.
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)
}
}
}
}
/// The tailnet is always the last candidate, so its failure ends the chain: terminal
/// for pairing and for anything retrying can't fix; otherwise offline + backoff.
private func tailnetAttemptFailed(_ error: Error, isPairing: Bool) {
if tryNextCandidate() { return }
if isPairing {
connectivity = .failed(error.localizedDescription)
return
}
switch error {
case TailnetError.notBuiltIn, TailnetError.notConfigured:
connectivity = .failed(error.localizedDescription)
default:
connectivity = .hostOffline
scheduleRetry()
}
}
/// Create the `SyncClient` on an established channel and start consuming its events —
/// the tail of every connect path, LAN or tailnet, pair or reconnect.
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)
}
/// Bring the phone's embedded Tailscale node up (first run needs the auth key from
/// Settings ▸ Tailscale; afterwards the on-disk state carries the registration) and dial
/// the Mac's tailnet address.
private func tailnetChannel(host: String, port: UInt16) async throws -> FDFrameChannel {
guard TailnetSupport.isBuiltIn else { throw TailnetError.notBuiltIn }
let config = Self.phoneTailnetConfig()
tailnetStatus = "Starting…"
// Mirror node status while the start is in flight. A first start with no auth key
// goes through the interactive browser login; pairing usually runs from the scanner
// sheet — not Settings — so take the user straight to the approval page.
let watcher = Task { [weak self] in
for await status in await TailnetNode.shared.statusStream() {
guard let self, !Task.isCancelled else { break }
self.tailnetStatus = status.label
if case .needsLogin(let url) = status, let loginURL = URL(string: url) {
if self.tailnetLoginURL != loginURL {
self.tailnetLoginURL = loginURL
// Eject to Safari only for user-initiated pairing — the user is
// actively watching. A routine reconnect that suddenly needs a
// login (node revoked, state wiped) must not yank them out of the
// app; Settings ▸ Tailscale carries the login link instead.
if self.connectPlan?.pairingPayload != nil {
UIApplication.shared.open(loginURL, options: [:], completionHandler: nil)
}
}
} else {
self.tailnetLoginURL = nil
}
}
}
defer {
watcher.cancel()
tailnetLoginURL = nil
}
do {
try await TailnetNode.shared.ensureRunning(config: config)
} catch TailnetError.timedOut(let message) {
// A login/auth timeout won't fix itself — retrying would just block per lap.
// Rethrow as .notConfigured so reconnect() treats it as terminal, not offline.
tailnetStatus = await TailnetNode.shared.status.label
throw TailnetError.notConfigured(message)
} catch {
tailnetStatus = await TailnetNode.shared.status.label
throw error
}
tailnetStatus = await TailnetNode.shared.status.label
return try await TailnetNode.shared.dial(host: host, port: port)
}
/// The phone's embedded-node config. State lives in this app's sandboxed Application
/// Support (no cross-channel collision — each channel is its own app container).
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)
}
func unpair() {
teardown()
// Forget the *active* Mac specifically (the switcher may have made it not the last).
if let id = activeHostID { IdentityStore.removePairedHost(id: id) } else { IdentityStore.clearPairedHost() }
// Forget + drop the active Mac's connection specifically (the switcher may have made it not
// the last). Other Macs' connections keep running.
if let id = activeHostID {
connections[id]?.teardown()
connections[id] = nil
IdentityStore.removePairedHost(id: id)
} else {
IdentityStore.clearPairedHost()
}
activeHostID = nil
// Mesh P3: if other Macs remain paired, switch to one rather than going fully unpaired.
if IdentityStore.pairedHosts().last != nil {
// Mesh P3: if other Macs remain paired, switch to one (already connected) rather than unpairing.
if let next = IdentityStore.pairedHosts().last {
activeHostID = next.fingerprint
reconnect()
return
}
@@ -667,10 +583,9 @@ final class RemoteStore: ObservableObject {
/// flag so it survives relaunch (a reviewer may relaunch), and seed the mock world. `isPaired`
/// then returns true, so `RootView` shows the full TabView.
func enterDemo() {
teardown()
teardownAll()
demoMode = true
UserDefaults.standard.set(true, forKey: Self.demoModeKey)
reconnectAttempts = 0
seedDemo()
}
@@ -707,9 +622,11 @@ final class RemoteStore: ObservableObject {
openApprovals = []
openDiff = nil
diffLoading = false
seenSeq.removeAll()
markOpened(sessionID)
if demoMode { seedDemoTranscript(sessionID); return }
// The open transcript belongs to the active host; tell its connection so it forwards the
// snapshot/events (and dedupes them), then subscribe.
activeConnection?.openSessionID = sessionID
send(.subscribe(Subscribe(sessionID: sessionID, sinceSeq: nil, verbosity: .full)))
}
@@ -793,6 +710,7 @@ final class RemoteStore: ObservableObject {
send(.unsubscribe(target))
markOpened(target) // everything up to now has been seen
guard openSessionID == target else { return }
activeConnection?.openSessionID = nil
openSessionID = nil
openEvents = []
openApprovals = []
@@ -863,10 +781,12 @@ final class RemoteStore: ObservableObject {
// MARK: - Plumbing
/// Route an intent to the active host's connection (mesh P3). Every current intent — subscribe,
/// send-input, approvals, favorite/archive/delete, start-chat — targets a session or project on
/// the Mac whose list is on screen, i.e. the active one.
private func send(_ msg: ClientMsg) {
if demoMode { demoHandle(msg); return }
guard let client else { return }
Task { await client.send(msg) }
activeConnection?.send(msg)
}
// MARK: - Demo simulator (offline, interactive)
@@ -1052,30 +972,10 @@ final class RemoteStore: ObservableObject {
todos: transform(dashboard.todos), usage: dashboard.usage, statusFeeds: dashboard.statusFeeds)
}
private func makeChannel(_ endpoint: NWEndpoint) -> NWFrameChannel {
let channel = NWFrameChannel(endpoint: endpoint)
// Surface a transport-level failure with an actionable message. Without this, a refused
// connection or a denied local-network permission finished the frame stream and only
// showed up as a generic "handshake closed" — hiding the real (usually fixable) cause.
channel.onFailed = { [weak self] error in
Task { @MainActor in self?.handleTransportFailure(error) }
}
return channel
}
/// Map an `NWConnection` failure to a connectivity state the pairing/onboarding UI can act on.
/// Only meaningful while we're still establishing the link; once connected, a drop is the
/// normal `.closed` → reconnect path's job. And only when the connect chain has nothing
/// left to try — a failed LAN probe about to fall back to the tailnet is routine, not news.
private func handleTransportFailure(_ error: String) {
guard !connectivity.isLive else { return }
guard connectPlan?.remaining.isEmpty != false else { return }
connectivity = .failed(Self.friendlyTransportError(error))
}
/// 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."
@@ -1086,150 +986,6 @@ final class RemoteStore: ObservableObject {
return "Couldn't connect to your Mac — make sure it's on the same Wi‑Fi and remote access is on."
}
private func resolveEndpoint(fingerprint: String?, lanHost: String?, lanPort: UInt16?) -> NWEndpoint? {
discovery.endpoint(forFingerprint: fingerprint, lanHost: lanHost, lanPort: lanPort)
}
private func consume(_ client: SyncClient, pairingPayload: PairingPayload?) {
eventTask = Task { [weak self] in
let stream = await client.start()
for await event in stream {
// A replaced client's tail events (`.failed` is always chased by `.closed`)
// must not leak into the new attempt — they'd advance the candidate chain
// or schedule retries against a connection that no longer exists.
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 // the chain found its transport
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() {
IdentityStore.savePairedHost(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))
}
send(.listSessions)
send(.listDashboard)
if let id = openSessionID { send(.subscribe(Subscribe(sessionID: id, sinceSeq: nil, verbosity: .full))) }
flushPendingNotificationDecision()
case .sessionList(let list):
sessions = list
NotificationRouter.shared.updateBadge(needsYouCount)
LiveActivityManager.shared.sync(hostName: hostName, sessions: sessions)
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)
}
// Notify on the transition into "waiting on you" / "finished" — only when the
// app isn't foreground-active (in-app, the list's washes and badges carry it).
let becameWaiting = summary.status == .awaitingInput
&& previous?.status != .awaitingInput
if becameWaiting, !isActive, !summary.archived {
NotificationRouter.shared.postSessionUpdate(summary)
}
NotificationRouter.shared.updateBadge(needsYouCount)
LiveActivityManager.shared.sync(hostName: hostName, sessions: sessions)
case .dashboard(let snapshot):
dashboard = snapshot
case .snapshot(let snapshot):
guard snapshot.summary.sessionID == openSessionID else { break }
seenSeq = Set(snapshot.recentEvents.map(\.seq))
openEvents = snapshot.recentEvents
openApprovals = snapshot.pendingApprovals
case .events(let batch):
guard batch.sessionID == openSessionID else { break }
for e in batch.events where !seenSeq.contains(e.seq) {
seenSeq.insert(e.seq)
openEvents.append(e)
}
case .approvalRequested(let req):
if req.sessionID == openSessionID, !openApprovals.contains(where: { $0.id == req.id }) {
openApprovals.append(req)
}
// Always post; the router suppresses the banner when the user is already
// looking at this session, and resolution (any device) withdraws it.
let title = sessions.first { $0.sessionID == req.sessionID }?.title ?? "Approval"
NotificationRouter.shared.postApproval(req, sessionTitle: title)
case .approvalResolved(let resolved):
openApprovals.removeAll { $0.id == resolved.id }
NotificationRouter.shared.withdrawApproval(resolved.id)
case .sessionDiff(let diff):
guard diff.sessionID == openSessionID else { break }
openDiff = diff
diffLoading = false
case .peerList(let peers):
// Mesh P4: the host's known peers, for the (future) mesh device list and transfer
// picker. Recorded now so multi-host UI can render it; single-host builds ignore it.
meshPeers = peers
case .transferAccept, .transferReject, .transferReady, .transferCommitted, .transferChunkAck:
// Session-transfer replies (mesh P5) only reach a *source* Mac driving a transfer; a
// phone is never a transfer source, so these are inert here.
break
case .wireError(let error):
// A channel mismatch means the host and this remote were built from incompatible
// release channels. Surface it as a persistent failure with a clear message rather
// than a transient bubble; the follow-on `.closed` still schedules a backoff retry,
// so the connection recovers on its own once either side is updated to a matching
// channel.
if error.code == .channelMismatch {
connectivity = .failed(error.message)
break
}
// Not fatal — surface as a transient bubble (the Mac's last-error overlay).
// Losing an approval race isn't an error worth interrupting for; the card
// collapses on the matching `approvalResolved`.
guard error.code != .alreadyResolved else { break }
// While the connect chain is still resolving a transport, a pre-ready error
// (e.g. "handshake closed" from a LAN probe about to fall back to tailnet) is
// routine, not news — the follow-on `.closed` advances the chain silently.
guard connectPlan == nil else { break }
showError(error.message, sessionID: error.sessionID)
case .failed(let message):
// Another candidate may still carry the session (e.g. LAN died → tailnet).
if tryNextCandidate() { break }
connectPlan = nil
// The channel's onFailed may have just surfaced a friendlier transport-level
// cause (denied Local Network, connection refused) — don't clobber it.
if case .failed = connectivity {} else { connectivity = .failed(message) }
scheduleRetry()
case .closed:
if !connectivity.isLive, tryNextCandidate() { break }
if connectivity.isLive {
connectivity = .reconnecting
} else if connectPlan?.pairingPayload != nil {
// A pairing chain died silently (every candidate closed pre-welcome).
// There's no retry loop before a pairing succeeds, so without a terminal
// state this would sit on "Connecting…" forever. Keep a friendlier
// transport-level failure if one was already surfaced.
if case .failed = connectivity {} else {
connectivity = .failed("Couldn't connect to your Mac — check that it's reachable, then scan again.")
}
}
connectPlan = nil
scheduleRetry()
}
}
/// Show a transient error bubble, replacing any current one; auto-dismisses after 6s
/// (matching the Mac's last-error overlay cadence).
private func showError(_ message: String, sessionID: SessionID?) {
@@ -1248,38 +1004,6 @@ final class RemoteStore: ObservableObject {
lastError = nil
}
private func scheduleRetry() {
guard isPaired else { return }
reconnectAttempts += 1
let delay = min(Double(reconnectAttempts) * 1.5, 10)
retryTask?.cancel()
retryTask = Task { [weak self] in
// `try?` swallows the sleep's CancellationError, so check explicitly.
try? await Task.sleep(for: .seconds(delay))
guard !Task.isCancelled, let self, !self.connectivity.isLive else { return }
self.reconnect()
}
}
private func teardown() {
retryTask?.cancel()
retryTask = nil
connectTask?.cancel()
connectTask = nil
connectPlan = nil
teardownClient()
}
/// Drop just the current client/channel — what moving to the next transport candidate
/// needs, without discarding the rest of the plan.
private func teardownClient() {
lanConnectTimeout?.cancel()
lanConnectTimeout = nil
eventTask?.cancel()
eventTask = nil
if let client { Task { await client.disconnect() } }
client = nil
}
}
extension Data {