Files
nucleic-remote-ios/NucleicRemote/NucleicRemote/Models/HostConnection.swift
T

907 lines
47 KiB
Swift
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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.
/// Changing which session is open resets the per-session cursor + dedup set (they only make
/// sense within one transcript).
var openSessionID: SessionID? {
didSet {
guard openSessionID != oldValue else { return }
openMaxSeq = nil
seenSeq = []
transcriptFetchRequestID = nil
didStartFullFetch = false
}
}
/// The highest event `seq` delivered for the open session — the warm-resubscribe cursor. On a
/// reconnect we re-subscribe `sinceSeq: openMaxSeq` so the host replays only what we missed
/// instead of cold-resetting to the 200-event tail (which would truncate a long transcript).
private var openMaxSeq: UInt64?
/// The in-flight full-history fetch's request id (mesh full-transcript sync). On open we pull the
/// *whole* transcript — the cold `subscribe` only returns a 200-event tail — and merge its chunks
/// into the transcript on screen; replies are matched on this id so a late batch from a previous
/// open (the session changed underneath us) is dropped. `didStartFullFetch` guards against issuing
/// it twice — both `open` and the reconnect nudge call `fetchFullTranscript`, but only the first
/// that finds a live, capable connection actually sends.
private var transcriptFetchRequestID: String?
private var didStartFullFetch = false
/// The in-flight *background* transcript prefetch (one at a time): its request id, target
/// session, and accumulated chunks. Independent of the open-session fetch above — replies
/// are matched on this id, so the two flows never mix, and prefetching keeps working while a
/// different session is open on screen. Driven by RemoteStore, which warms the offline cache
/// with the result so a never-before-opened session still opens at its end instantly.
private var prefetchRequestID: String?
private var prefetchSessionID: SessionID?
private var prefetchEvents: [AgentEvent] = []
// 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 }
/// Backfilled history for the open session (the full-transcript fetch). Merged into the
/// transcript by seq, not appended — these events precede the tail already on screen.
var openBackfill: (EventBatch) -> Void = { _ in }
/// A background transcript prefetch finished — everything it fetched (empty when the
/// host had nothing / the fetch died with the connection). RemoteStore merges it into
/// the offline cache and starts the next queued prefetch either way.
var transcriptPrefetched: (SessionID, [AgentEvent]) -> 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 mesh roster changed (mesh "join"): a Mac was learned or revoked via gossip — the
/// registry was updated, so (re)connect to every paired Mac and drop any that left.
var meshRosterChanged: () -> Void = {}
/// The embedded Tailscale node's status changed (for Settings).
var tailnetStatus: (String?, URL?) -> Void = { _, _ in }
/// A join code this host minted at our request ("add a device to this mesh") — the
/// `nucleic://pair?d=…` string to show as a QR / copyable code, or nil if it couldn't.
var pairingCodeReceived: (String?) -> Void = { _ in }
/// This host forwarded a Mac-pair confirm — a Mac is joining via a code we shared and
/// wants the user's allow/deny. The phone can approve it (`respondMacPair`).
var macPairRequested: (WireMacPairRequest) -> Void = { _ in }
/// The forwarded Mac-pair (by deviceID) was answered/withdrawn — dismiss the prompt.
var macPairResolved: (String) -> Void = { _ in }
/// The outcome of a `createProject` this phone sent (CLOUD_RUNTIME §4.3), correlated by
/// requestID — settle the Add Project sheet (the project row rides the dashboard push).
var projectCreated: (WireProjectCreated) -> Void = { _ in }
}
private let callbacks: Callbacks
// MARK: Shared dependencies
private let identity: DeviceIdentity
private let discovery: LANDiscovery
private let pathMonitor: NetworkPathMonitor
// 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 revalidateTask: Task<Void, Never>?
private var lanConnectTimeout: Task<Void, Never>?
/// In-flight LAN-reachability probe (the tailnet/relay → LAN return leg). Nil unless a probe is
/// running; at most one at a time so repeated path-change events don't stack.
private var lanUpgradeProbe: Task<Void, Never>?
private var reconnectAttempts = 0
private var seenSeq: Set<UInt64> = []
private enum TransportAttempt {
case lan(NWEndpoint)
case tailnet(host: String, port: UInt16)
/// Nucleic Private Relay (mesh P2) — always the last candidate: works from anywhere,
/// but a direct path beats a brokered one when both exist.
case relay(base: URL, membershipToken: String)
}
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, pathMonitor: NetworkPathMonitor,
callbacks: Callbacks
) {
self.hostID = hostID
self.hostName = hostName
self.identity = identity
self.discovery = discovery
self.pathMonitor = pathMonitor
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
}
let candidates = buildCandidates(
fingerprint: payload.hostStaticKey.fingerprintHex,
lanHost: payload.lanHost, lanPort: payload.lanPort,
tailnet: hint == .tailnet ? (payload.tailnetHost, payload.tailnetPort) : nil,
relay: (payload.relayMembershipToken, payload.relayURL))
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()
let candidates = buildCandidates(
fingerprint: host.fingerprint,
lanHost: host.lanHost, lanPort: host.lanPort,
tailnet: host.transportHint == .tailnet ? (host.tailnetHost, host.tailnetPort) : nil,
relay: (host.relayMembershipToken, host.relayURL))
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()
}
/// Translate a gossiped mesh member (mesh "join") into a pinnable/dialable `PairedHost`. The
/// member carries the host static key + last-known addresses — everything IK reconnect needs.
/// No relay token yet (the phone never directly paired this Mac); it's adopted from the host's
/// `relayMembership` push after the first LAN/tailnet contact, so first contact needs one of
/// those paths.
static func pairedHost(from member: MeshMember) -> PairedHost {
let (lanHost, lanPort) = splitHostPort(member.addresses?.lanHint)
let tailnetHost = member.addresses?.tailnet
let hasTailnet = !(tailnetHost?.isEmpty ?? true)
return PairedHost(
deviceID: IdentityStore.deviceID(),
hostName: member.label,
hostStaticKey: member.staticPublicKey,
fingerprint: member.staticPublicKey.fingerprintHex,
lanHost: lanHost, lanPort: lanPort,
transport: (hasTailnet ? SyncTransportHint.tailnet : .lan).rawValue,
tailnetHost: hasTailnet ? tailnetHost : nil,
// The sync port on a tailnet is fixed protocol-wide (SyncTransportSetting.tailnetPort);
// a gossiped address carries only the IP.
tailnetPort: hasTailnet ? 43_753 : nil,
relayRoomID: member.addresses?.relayRoomID,
relayMembershipToken: nil, relayURL: nil)
}
/// Whether the dialable endpoints of a gossiped record differ from the one on file — the
/// signal that a known Mac must be re-dialed (its LAN port, tailnet IP, or relay room moved).
/// Ignores relay *credentials* (token/URL), which `mergePairedHost` preserves and which don't
/// change where the Mac is reached.
private static func dialableAddressChanged(from existing: PairedHost, to incoming: PairedHost) -> Bool {
existing.lanHost != incoming.lanHost
|| existing.lanPort != incoming.lanPort
|| existing.tailnetHost != incoming.tailnetHost
|| existing.tailnetPort != incoming.tailnetPort
|| existing.relayRoomID != incoming.relayRoomID
}
/// Split a "host:port" hint (last-colon split so a bracketed IPv6 host survives).
private static func splitHostPort(_ hint: String?) -> (String?, UInt16?) {
guard let hint, let colon = hint.lastIndex(of: ":"),
let port = UInt16(hint[hint.index(after: colon)...]) else { return (nil, nil) }
var host = String(hint[..<colon])
if host.hasPrefix("["), host.hasSuffix("]") { host = String(host.dropFirst().dropLast()) }
return (host.isEmpty ? nil : host, port)
}
private func buildCandidates(
fingerprint: String?, lanHost: String?, lanPort: UInt16?,
tailnet: (host: String?, port: UInt16?)?,
relay: (membershipToken: String?, url: String?)? = nil
) -> [TransportAttempt] {
var candidates: [TransportAttempt] = []
// Only offer LAN when the device actually has a LAN-capable path (Wi-Fi/wired). On cellular
// the pinned `lanHost:lanPort` is unreachable, so including it here would just burn the
// connect timeout before falling through — skip it and dial tailnet/relay immediately.
if pathMonitor.canUseLAN,
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))
}
if let relay, let token = relay.membershipToken, !token.isEmpty {
candidates.append(.relay(base: RelayAPI.baseURL(relay.url), membershipToken: token))
}
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.asyncAttemptFailed(error, isPairing: plan.pairingPayload != nil)
}
}
case .relay(let base, let membershipToken):
activeTransport = .relay
connectTask = Task { [weak self] in
guard let self else { return }
do {
let channel = try await RelayFrameChannel.dial(
base: base, membershipToken: membershipToken)
guard !Task.isCancelled else { channel.close(); return }
// A room with no live host swallows frames silently until presence says
// otherwise — bound the handshake like the LAN path does.
self.lanConnectTimeout?.cancel()
self.lanConnectTimeout = Task { [weak self, weak channel] in
try? await Task.sleep(for: .seconds(10))
guard !Task.isCancelled, let self, !self.connectivity.isLive else { return }
channel?.close()
}
self.startClient(
channel: channel, hostStaticKey: plan.hostStaticKey,
mode: plan.mode, deviceID: plan.deviceID, pairingPayload: plan.pairingPayload)
} catch {
guard !Task.isCancelled else { return }
self.asyncAttemptFailed(error, isPairing: plan.pairingPayload != nil)
}
}
}
}
/// A candidate that dials asynchronously (tailnet, relay) failed before producing a
/// channel — fall through to the next candidate or schedule a retry.
private func asyncAttemptFailed(_ 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,
// Our own bundle id is the exact APNs topic the relay must address; per-channel
// TestFlight builds are suffixed (…`.canary`), so a hardcoded topic would `BadTopic`.
pushTopic: Bundle.main.bundleIdentifier,
releaseChannel: BuildInfo.current.channel.releaseChannel,
// Mesh "join": advertise roster gossip so a host pushes its group view — the phone
// then auto-learns and connects to every Mac in the mesh, not just the one it scanned.
clientCaps: WireClientCapabilities(mesh: 1, canSyncRoster: true))
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,
relayRoomID: payload.relayRoomID,
relayMembershipToken: payload.relayMembershipToken,
relayURL: payload.relayURL))
}
callbacks.didUpdate()
send(.listSessions)
send(.listDashboard)
// Warm-resubscribe from what we already have, so a reconnect on a long transcript replays
// only the gap rather than snapping back to the host's 200-event tail.
if let id = openSessionID {
send(.subscribe(Subscribe(sessionID: id, sinceSeq: openMaxSeq, verbosity: .full)))
// Also pull full history — the subscribe above only replays the tail (or, on a warm
// reconnect, the missed gap). No-op after the first open of this session, so an
// ordinary reconnect keeps the transcript on screen instead of refetching it.
fetchFullTranscript(id)
}
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 }
// A snapshot may be a fresh open (tail window) or a warm-resubscribe delta; either way
// union its seqs into the dedup set and advance the cursor. RemoteStore merges the events
// into the transcript rather than replacing, so an existing transcript isn't truncated.
seenSeq.formUnion(snapshot.recentEvents.map(\.seq))
if let maxSeq = snapshot.recentEvents.map(\.seq).max() { openMaxSeq = max(openMaxSeq ?? 0, maxSeq) }
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 let maxSeq = batch.events.map(\.seq).max() { openMaxSeq = max(openMaxSeq ?? 0, maxSeq) }
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 .meshRoster(let push):
// Mesh "join": the host gossiped its group view. Auto-learn every Mac member into the
// paired-hosts registry (so the multiplexer connects to all of them), and drop any the
// group revoked. Phone members are ignored — a phone doesn't dial other phones.
var changed = false
for member in push.members where member.kind == .mac || member.capabilities.canHost {
// A mac deviceID is `sha256(staticPublicKey)`; reject a forged (id, key) pair.
guard DeviceIdentity.hostID(ofStaticKey: member.staticPublicKey) == member.deviceID
else { continue }
let fingerprint = member.staticPublicKey.fingerprintHex
guard fingerprint != self.hostID else { continue } // that's the Mac we're on
let incoming = Self.pairedHost(from: member)
let existing = IdentityStore.pairedHost(id: fingerprint)
// Re-dial not only when a Mac is brand-new, but also when a known Mac's dialable
// address changed — a Mac's LAN port is OS-assigned (`.any`), so it lands on a
// fresh one every relaunch and re-gossips it. `mergePairedHost` updates the
// registry in place, but a connection already retrying the dead old address won't
// pick that up on its own (its retry loop reuses the address it was last handed),
// so without forcing a fresh reconnect here the phone never reconnects to it.
if existing == nil || Self.dialableAddressChanged(from: existing!, to: incoming) {
changed = true
}
IdentityStore.mergePairedHost(incoming)
}
for tombstone in push.tombstones {
// fingerprint = first 16 hex of the hostID (sha256 prefix), the registry key.
let fingerprint = String(tombstone.deviceID.prefix(16))
if IdentityStore.pairedHost(id: fingerprint) != nil {
IdentityStore.removePairedHost(id: fingerprint)
changed = true
}
}
if changed { callbacks.meshRosterChanged() }
case .transferAccept, .transferReject, .transferReady, .transferCommitted, .transferChunkAck:
// Session-transfer replies (mesh P5) only reach a *source* Mac; a phone is never one.
break
case .transcriptChunk(let chunk):
// A slice of a *background prefetch*: accumulate it whole — it never touches the
// open transcript's cursor/dedup state, which belongs to the on-screen session.
if chunk.requestID == prefetchRequestID {
if chunk.sessionID == prefetchSessionID { prefetchEvents.append(contentsOf: chunk.events) }
break
}
// One ordered slice of the full-history fetch we started on open. Match on request id so
// a stale batch from a previous open (session switched underneath us) is ignored.
guard chunk.requestID == transcriptFetchRequestID, chunk.sessionID == openSessionID else { break }
// De-dup against the tail/live stream and advance the warm-resubscribe cursor exactly like
// a snapshot/events batch, then hand the fresh (older) events to RemoteStore, which merges
// them into the transcript by seq rather than appending.
let fresh = chunk.events.filter { !seenSeq.contains($0.seq) }
seenSeq.formUnion(chunk.events.map(\.seq))
if let maxSeq = chunk.events.map(\.seq).max() { openMaxSeq = max(openMaxSeq ?? 0, maxSeq) }
if !fresh.isEmpty { callbacks.openBackfill(EventBatch(sessionID: chunk.sessionID, events: fresh)) }
case .transcriptFetchComplete(let done):
if done.requestID == prefetchRequestID { finishPrefetch(delivering: true); break }
// Terminal success — every batch shipped. Clear the in-flight id; the cursor/dedup state
// is already advanced by the chunks above.
if done.requestID == transcriptFetchRequestID { transcriptFetchRequestID = nil }
case .transcriptUnavailable(let un):
if un.requestID == prefetchRequestID { finishPrefetch(delivering: false); break }
// The host holds no transcript for this session (or can't serve the format). The cold tail
// still shows; there's just no deeper history to add. Clear the in-flight id.
if un.requestID == transcriptFetchRequestID { transcriptFetchRequestID = nil }
case .toolSummaries:
// Collapsed Bash summary lines the owner pushes to peer *Macs* (mesh session sync).
// The phone renders its own deterministic command summaries, so it ignores these.
break
case .relayMembership(let membership):
// The host issued/refreshed this device's relay credential (mesh P2). Persist it
// in place (no reordering — this can arrive from a non-active Mac) so the relay
// is a dial candidate on the next reconnect; also patch the in-memory pin so an
// imminent retry uses the fresh token without a registry round-trip.
IdentityStore.updatePairedHost(id: hostID) {
$0.relayRoomID = membership.roomID
$0.relayMembershipToken = membership.token
$0.relayURL = membership.url
}
if pinnedHost?.fingerprint == hostID {
pinnedHost?.relayRoomID = membership.roomID
pinnedHost?.relayMembershipToken = membership.token
pinnedHost?.relayURL = membership.url
}
case .pairingCode(let qr):
// This Mac minted a join code we asked for ("add a device to this mesh") — hand it up
// to RemoteStore for the QR/copy sheet.
callbacks.pairingCodeReceived(qr)
case .macPairRequested(let req):
// A Mac is joining via a code we shared and needs allow/deny — forward it up so the
// phone can approve the join.
callbacks.macPairRequested(req)
case .macPairResolved(let deviceID, _):
// Answered on this Mac or another device (or timed out) — dismiss the prompt.
callbacks.macPairResolved(deviceID)
case .projectCreated(let outcome):
// The host settled a createProject we sent — hand it up so the Add Project sheet
// resolves (success or failure). Correlation by requestID happens in RemoteStore.
callbacks.projectCreated(outcome)
case .intelligenceRequest, .credentialNeeded, .credentialUpdate:
// Antimatter runner verbs (docs/ANTIMATTER_RUNNER.md §5–6): a runner host delegating
// intelligence work or asking for / mirroring sealed credentials. Inert here until the
// phone-side executor/vault land — and a host only sends these to clients that
// advertised the matching `WireClientCapabilities`, which this app doesn't yet.
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) }
}
/// Pull the open session's *full* transcript via mesh full-transcript sync and merge it into the
/// transcript on screen. The cold `subscribe` only returns the host's 200-event tail, so without
/// this the phone shows nothing from before it connected. `afterSeq: 0` asks for the whole history
/// (overlap with the tail is deduped on arrival). Idempotent per open — guarded by
/// `didStartFullFetch` — and a no-op when the host predates `canSyncTranscripts` (the tail still
/// shows) or the connection isn't live yet (the reconnect handler retries once it is).
func fetchFullTranscript(_ sessionID: SessionID) {
guard sessionID == openSessionID, !didStartFullFetch,
capabilities.canSyncTranscripts, client != nil else { return }
didStartFullFetch = true
let requestID = UUID().uuidString
transcriptFetchRequestID = requestID
send(.fetchTranscript(TranscriptFetch(requestID: requestID, sessionID: sessionID, afterSeq: 0)))
}
/// Whether a background transcript prefetch can start now: live, the host can serve
/// transcript fetches, and no prefetch is already in flight (they run one at a time).
var canPrefetchTranscript: Bool {
connectivity.isLive && capabilities.canSyncTranscripts && client != nil && prefetchRequestID == nil
}
/// Fetch a session's transcript in the background — no subscribe, no open — so the offline
/// cache is warm before the user ever opens the session. `afterSeq` bounds the pull to what
/// the cache is missing. The result (all chunks, merged) arrives via
/// `callbacks.transcriptPrefetched`; returns false when it can't start (caller retries on a
/// later session-list update).
@discardableResult
func prefetchTranscript(_ sessionID: SessionID, afterSeq: UInt64) -> Bool {
guard canPrefetchTranscript else { return false }
let requestID = UUID().uuidString
prefetchRequestID = requestID
prefetchSessionID = sessionID
prefetchEvents = []
send(.fetchTranscript(TranscriptFetch(requestID: requestID, sessionID: sessionID, afterSeq: afterSeq)))
return true
}
/// Settle the in-flight prefetch: hand what arrived to RemoteStore (nothing on
/// unavailable/disconnect) and clear the slot so the next queued prefetch can start.
private func finishPrefetch(delivering: Bool) {
guard let sessionID = prefetchSessionID else { return }
let events = delivering ? prefetchEvents : []
prefetchRequestID = nil
prefetchSessionID = nil
prefetchEvents = []
callbacks.transcriptPrefetched(sessionID, events)
}
/// The device's network path changed (Wi-Fi ⇄ cellular, joined/left a network). React now rather
/// than waiting for a zombie LAN socket to time out or the reconnect backoff to elapse — this is
/// what makes the transport switch feel immediate. Only reconnect-managed hosts (a pinned host)
/// re-plan here; a host mid-pairing runs its own one-shot candidate chain.
func networkPathChanged() {
guard let host = pinnedHost else { return }
// Using (or dialing) LAN but LAN is no longer reachable → the socket is dead in all but name.
// Drop it and re-plan now; `buildCandidates` skips LAN and dials tailnet/relay.
if activeTransport == .lan, !pathMonitor.canUseLAN {
reconnectAttempts = 0
reconnect(to: host)
return
}
// Offline/reconnecting and the network just came back or changed → retry immediately instead
// of sitting out the remaining backoff.
if !connectivity.isLive, pathMonitor.isSatisfied {
reconnectAttempts = 0
reconnect(to: host)
return
}
// Live over a fallback transport (tailnet/relay) and a LAN path just reappeared (rejoined the
// Mac's Wi-Fi) → the return leg of the switch above. The tunnel survives the interface change
// on its own, so nothing else would ever re-plan back down to the direct path — probe LAN and
// swap to it if the Mac actually answers here.
maybeUpgradeToLAN()
}
/// Probe whether the pinned Mac is reachable over LAN right now and, if so, switch back to the
/// direct connection — the tailnet/relay → LAN return leg. Only meaningful while live over a
/// non-LAN transport with a LAN-capable path available; a bare TCP connect to the pinned LAN
/// endpoint that reaches `.ready` is the "the Mac is on this network" signal. The probe runs
/// entirely off the live connection, so joining a *foreign* Wi-Fi (no Mac here) costs one short,
/// silent connect attempt and leaves the working tunnel untouched.
private func maybeUpgradeToLAN() {
guard let host = pinnedHost, connectivity.isLive, activeTransport != .lan,
pathMonitor.canUseLAN, lanUpgradeProbe == nil,
let endpoint = discovery.endpoint(
forFingerprint: host.fingerprint, lanHost: host.lanHost, lanPort: host.lanPort)
else { return }
lanUpgradeProbe = Task { [weak self] in
let reachable = await Self.probeReachable(endpoint)
guard let self else { return }
self.lanUpgradeProbe = nil
// Re-check the world didn't move under us during the probe (still live, still not on LAN,
// LAN still available) before tearing a working connection down.
guard !Task.isCancelled, reachable, self.connectivity.isLive,
self.activeTransport != .lan, self.pathMonitor.canUseLAN,
let host = self.pinnedHost
else { return }
self.reconnectAttempts = 0
self.reconnect(to: host) // rebuilds the candidate chain LAN-first
}
}
/// Bare TCP reachability check for `maybeUpgradeToLAN`: does `endpoint` accept a connection within
/// `timeoutSeconds`? Tears the probe socket down regardless — it only tests reachability, never
/// carries traffic. A short timeout so a foreign Wi-Fi (Mac absent) doesn't stall the check.
private static func probeReachable(_ endpoint: NWEndpoint, timeoutSeconds: Double = 3) async -> Bool {
let params = NWParameters.tcp
params.includePeerToPeer = true
let connection = NWConnection(to: endpoint, using: params)
let queue = DispatchQueue(label: "nucleic.remote.lanprobe")
let once = ProbeOnce()
return await withTaskCancellationHandler {
await withCheckedContinuation { (cont: CheckedContinuation<Bool, Never>) in
connection.stateUpdateHandler = { state in
switch state {
case .ready:
if once.take() { cont.resume(returning: true) }
connection.cancel()
case .failed, .cancelled:
if once.take() { cont.resume(returning: false) }
default:
break
}
}
queue.asyncAfter(deadline: .now() + timeoutSeconds) {
if once.take() { cont.resume(returning: false) }
connection.cancel()
}
connection.start(queue: queue)
}
} onCancel: {
connection.cancel()
}
}
/// Foreground liveness check — the fast-reconnect path. iOS suspends the app when it's
/// backgrounded, which silently kills the socket; but no `.closed`/`.failed` reaches us while
/// we're frozen, so this connection can wake still marked `.connected`. Trusting that stale
/// state means nothing reconnects until the 55s keepalive deadline (or ~15s TCP keepalive)
/// finally notices — the several-second (or worse) stall on return. Instead: if we're already
/// not live, re-dial now with no backoff; if we *look* live, actively ping and re-dial the
/// instant the link fails to answer within a tight deadline. A pinned host only (a host
/// mid-pairing runs its own one-shot chain and has no `pinnedHost`).
func revalidate() {
guard let host = pinnedHost else { return }
guard connectivity.isLive, let client else {
reconnectAttempts = 0
reconnect(to: host)
return
}
// Already live — but if we're on a fallback transport and a LAN path exists now, try to move
// back to the direct connection (foregrounding on the Mac's Wi-Fi after being away on cellular,
// where no path-change event fired while suspended).
maybeUpgradeToLAN()
revalidateTask?.cancel()
revalidateTask = Task { [weak self] in
// A generous bound on a pong round-trip (LAN <10ms, relay <200ms) — only paid in full
// when the link is actually dead; a live one answers and exits in ~1 RTT. Kept well
// under the follow-on re-dial so the whole dead-link reconnect stays sub-second.
let alive = await client.probeAlive(within: .milliseconds(600))
guard !Task.isCancelled, let self, self.client === client,
self.connectivity.isLive else { return }
if !alive {
self.reconnectAttempts = 0
self.reconnect(to: host)
}
}
}
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
revalidateTask?.cancel(); revalidateTask = nil
connectTask?.cancel(); connectTask = nil
lanUpgradeProbe?.cancel(); lanUpgradeProbe = nil
connectPlan = nil
teardownClient()
}
private func teardownClient() {
lanConnectTimeout?.cancel(); lanConnectTimeout = nil
eventTask?.cancel(); eventTask = nil
if let client { Task { await client.disconnect() } }
client = nil
// A dropped connection loses the in-flight prefetch's remaining chunks — settle it empty
// so RemoteStore's queue isn't left waiting on a completion that will never arrive.
finishPrefetch(delivering: false)
}
}
/// A thread-safe "fire once" latch: the `NWConnection` reachability probe resolves its continuation
/// from a state callback *or* a timeout, on the same queue but from different closures — `take()`
/// returns true to exactly one caller so the continuation is resumed exactly once.
private final class ProbeOnce: @unchecked Sendable {
private let lock = NSLock()
private var done = false
func take() -> Bool {
lock.lock(); defer { lock.unlock() }
if done { return false }
done = true
return true
}
}