907 lines
47 KiB
Swift
907 lines
47 KiB
Swift
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
|
||
}
|
||
}
|