Merge claude/keen-shamir-122122 into dev (mesh P3 multiplexer)
The iOS [HostID: HostConnection] multiplexer: connect to all paired Macs at once, mirror the active one, instant host switching, cross-host badge/Live Activity. Compile + demo-verified; the live multi-connection path needs a two-Mac test (may be rolled back). Resolves the RemoteStore.swift overlap with dev's Tailnet-node-name/auth-key change by taking the multiplexer version (which moved the tailnet config into HostConnection); dev's authKey removal is re-applied there. Co-Authored-By: Claude Opus 4.8 <[email protected]> # Conflicts: # ios/NucleicRemote/NucleicRemote/Models/RemoteStore.swift
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user