nvrsion: Add: Implement Activity Logic, fix network connection logic for LAN→Tailnet→LAN switching in iOS, fix Live Activity Session Load, and investigate Live Activity Session Load.
Nucleic-Promote: 1 Co-authored-by: Nucleic <[email protected]>
This commit is contained in:
@@ -49,8 +49,20 @@ final class LiveActivityManager {
|
|||||||
/// is backgrounded and its sync socket is suspended (UX_IOS §5.3).
|
/// is backgrounded and its sync socket is suspended (UX_IOS §5.3).
|
||||||
var onPushToken: ((_ token: String, _ activityID: String) -> Void)?
|
var onPushToken: ((_ token: String, _ activityID: String) -> Void)?
|
||||||
var onActivityEnded: ((_ activityID: String) -> Void)?
|
var onActivityEnded: ((_ activityID: String) -> Void)?
|
||||||
|
/// Set by `RemoteStore` to ship this device's **push-to-start** token to the paired Macs. Unlike
|
||||||
|
/// `onPushToken` (a per-activity update token that exists only once an Activity does), this token
|
||||||
|
/// is device-scoped and exists before any Activity — it's what lets a Mac *create* the glance
|
||||||
|
/// over APNs when work starts while the app is closed, so it appears without the user opening the
|
||||||
|
/// app first (iOS 17.2+, UX_IOS §5.3).
|
||||||
|
var onPushToStartToken: ((_ token: String) -> Void)?
|
||||||
/// Streams the activity's per-activity APNS update token (it can rotate); cancelled on end.
|
/// Streams the activity's per-activity APNS update token (it can rotate); cancelled on end.
|
||||||
private var tokenObservation: Task<Void, Never>?
|
private var tokenObservation: Task<Void, Never>?
|
||||||
|
/// Streams this device's push-to-start token (it can rotate). Device-scoped, so — unlike
|
||||||
|
/// `tokenObservation` — it lives for the whole process and is never cancelled on `end()`.
|
||||||
|
private var pushToStartObservation: Task<Void, Never>?
|
||||||
|
/// Observes Activities that appear without us creating them — i.e. ones a Mac push-started while
|
||||||
|
/// the app was closed — so we adopt them and forward their update token. Also process-lived.
|
||||||
|
private var activityAdoptionObservation: Task<Void, Never>?
|
||||||
|
|
||||||
/// Reconcile the Activity with the current session set.
|
/// Reconcile the Activity with the current session set.
|
||||||
func sync(hostName: String, sessions: [WireSessionSummary]) {
|
func sync(hostName: String, sessions: [WireSessionSummary]) {
|
||||||
@@ -177,6 +189,41 @@ final class LiveActivityManager {
|
|||||||
enqueue(state)
|
enqueue(state)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Start the process-lived observers that make push-to-start work: the device's push-to-start
|
||||||
|
/// token (forwarded to the Macs so they can create the glance over APNs) and adoption of any
|
||||||
|
/// Activity a Mac push-started while the app was closed. Idempotent — safe to call on every
|
||||||
|
/// bridge setup; the guards keep a single observer each.
|
||||||
|
func beginPushToStartObservation() {
|
||||||
|
guard ActivityAuthorizationInfo().areActivitiesEnabled else { return }
|
||||||
|
if pushToStartObservation == nil {
|
||||||
|
pushToStartObservation = Task { [weak self] in
|
||||||
|
for await data in Activity<NucleicSessionAttributes>.pushToStartTokenUpdates {
|
||||||
|
let hex = data.map { String(format: "%02x", $0) }.joined()
|
||||||
|
self?.onPushToStartToken?(hex)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if activityAdoptionObservation == nil {
|
||||||
|
activityAdoptionObservation = Task { [weak self] in
|
||||||
|
for await activity in Activity<NucleicSessionAttributes>.activityUpdates {
|
||||||
|
self?.adopt(activity)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Adopt an Activity we didn't create locally — almost always one a Mac push-started while the
|
||||||
|
/// app was closed. Taking ownership wires up its update-token stream (`observePushToken`) so the
|
||||||
|
/// Macs can keep the glance fresh the normal way, and collapses any strays to one. If we already
|
||||||
|
/// track an Activity, leave it: the local `Activity.request` path and `endStrays` already keep a
|
||||||
|
/// single glance, and re-adopting would just churn the token observation.
|
||||||
|
private func adopt(_ activity: Activity<NucleicSessionAttributes>) {
|
||||||
|
guard self.activity == nil else { return }
|
||||||
|
self.activity = activity
|
||||||
|
observePushToken(activity)
|
||||||
|
endStrays(keeping: activity.id)
|
||||||
|
}
|
||||||
|
|
||||||
/// Forward the activity's APNS update token (and its rotations) to `RemoteStore`.
|
/// Forward the activity's APNS update token (and its rotations) to `RemoteStore`.
|
||||||
private func observePushToken(_ activity: Activity<NucleicSessionAttributes>) {
|
private func observePushToken(_ activity: Activity<NucleicSessionAttributes>) {
|
||||||
tokenObservation?.cancel()
|
tokenObservation?.cancel()
|
||||||
|
|||||||
@@ -117,6 +117,9 @@ final class HostConnection {
|
|||||||
private var retryTask: Task<Void, Never>?
|
private var retryTask: Task<Void, Never>?
|
||||||
private var revalidateTask: Task<Void, Never>?
|
private var revalidateTask: Task<Void, Never>?
|
||||||
private var lanConnectTimeout: 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 reconnectAttempts = 0
|
||||||
private var seenSeq: Set<UInt64> = []
|
private var seenSeq: Set<UInt64> = []
|
||||||
|
|
||||||
@@ -686,6 +689,72 @@ final class HostConnection {
|
|||||||
if !connectivity.isLive, pathMonitor.isSatisfied {
|
if !connectivity.isLive, pathMonitor.isSatisfied {
|
||||||
reconnectAttempts = 0
|
reconnectAttempts = 0
|
||||||
reconnect(to: host)
|
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()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -704,6 +773,10 @@ final class HostConnection {
|
|||||||
reconnect(to: host)
|
reconnect(to: host)
|
||||||
return
|
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?.cancel()
|
||||||
revalidateTask = Task { [weak self] in
|
revalidateTask = Task { [weak self] in
|
||||||
// A generous bound on a pong round-trip (LAN <10ms, relay <200ms) — only paid in full
|
// A generous bound on a pong round-trip (LAN <10ms, relay <200ms) — only paid in full
|
||||||
@@ -735,6 +808,7 @@ final class HostConnection {
|
|||||||
retryTask?.cancel(); retryTask = nil
|
retryTask?.cancel(); retryTask = nil
|
||||||
revalidateTask?.cancel(); revalidateTask = nil
|
revalidateTask?.cancel(); revalidateTask = nil
|
||||||
connectTask?.cancel(); connectTask = nil
|
connectTask?.cancel(); connectTask = nil
|
||||||
|
lanUpgradeProbe?.cancel(); lanUpgradeProbe = nil
|
||||||
connectPlan = nil
|
connectPlan = nil
|
||||||
teardownClient()
|
teardownClient()
|
||||||
}
|
}
|
||||||
@@ -746,3 +820,17 @@ final class HostConnection {
|
|||||||
client = nil
|
client = nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -558,13 +558,20 @@ final class RemoteStore: ObservableObject {
|
|||||||
var cb = HostConnection.Callbacks()
|
var cb = HostConnection.Callbacks()
|
||||||
cb.didUpdate = { [weak self] in
|
cb.didUpdate = { [weak self] in
|
||||||
guard let self else { return }
|
guard let self else { return }
|
||||||
|
// Complete a session opened before its owning Mac was live — a Live Activity /
|
||||||
|
// notification cold-launch deep link opens the detail while the app is still dialing, so
|
||||||
|
// `open()` couldn't bind it to a connection that didn't exist yet. Runs before the merge
|
||||||
|
// so the aggregate picks up the freshly-bound host's connectivity/scope this pass.
|
||||||
|
self.bindOpenSessionIfNeeded()
|
||||||
// Any host's change re-merges the aggregate the flat state binds to (mesh P3).
|
// Any host's change re-merges the aggregate the flat state binds to (mesh P3).
|
||||||
self.rebuildAggregate()
|
self.rebuildAggregate()
|
||||||
self.refreshAggregate()
|
self.refreshAggregate()
|
||||||
self.flushPendingNotificationDecision()
|
self.flushPendingNotificationDecision()
|
||||||
// A host that just connected (or reconnected) needs the current Live Activity token
|
// A host that just connected (or reconnected) needs the current Live Activity token
|
||||||
// so it can push while the phone is away.
|
// so it can push while the phone is away — and the push-to-start token so it can create
|
||||||
|
// the glance cold when work starts before the app is opened.
|
||||||
self.syncLiveActivityRegistration()
|
self.syncLiveActivityRegistration()
|
||||||
|
self.syncPushToStartRegistration()
|
||||||
}
|
}
|
||||||
cb.openSnapshot = { [weak self] snap in
|
cb.openSnapshot = { [weak self] snap in
|
||||||
guard let self, hostID == self.openSessionHostID, snap.summary.sessionID == self.openSessionID else { return }
|
guard let self, hostID == self.openSessionHostID, snap.summary.sessionID == self.openSessionID else { return }
|
||||||
@@ -922,6 +929,26 @@ final class RemoteStore: ObservableObject {
|
|||||||
connection(owningSession: sessionID)?.fetchFullTranscript(sessionID)
|
connection(owningSession: sessionID)?.fetchFullTranscript(sessionID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Bind the open session to its owning Mac once that Mac is live — the deferred half of `open()`
|
||||||
|
/// for a session opened before any connection existed (a Live Activity / notification cold-launch
|
||||||
|
/// deep link). Because `open()` ran while the app was still dialing, `connection(owningSession:)`
|
||||||
|
/// found nothing: the host binding stayed nil, no `subscribe` went out, and the transcript
|
||||||
|
/// callbacks — all gated on `openSessionHostID` — dropped everything the host later sent, leaving
|
||||||
|
/// the detail stuck on "Disconnected" with an empty transcript and a dead composer. Once a host
|
||||||
|
/// connects and lists its sessions, `didUpdate` calls this: it finds the owner and finishes the
|
||||||
|
/// subscription so the transcript loads and connectivity/scope follow the bound host. Idempotent
|
||||||
|
/// and cheap — a no-op on the common path where `open()` already bound a live connection.
|
||||||
|
private func bindOpenSessionIfNeeded() {
|
||||||
|
guard !demoMode, let sessionID = openSessionID,
|
||||||
|
let conn = connection(owningSession: sessionID), conn.connectivity.isLive else { return }
|
||||||
|
// Already bound to this live host with its subscription in place — nothing to redo.
|
||||||
|
guard openSessionHostID != conn.hostID || conn.openSessionID != sessionID else { return }
|
||||||
|
openSessionHostID = conn.hostID
|
||||||
|
conn.openSessionID = sessionID
|
||||||
|
conn.send(.subscribe(Subscribe(sessionID: sessionID, sinceSeq: nil, verbosity: .full)))
|
||||||
|
conn.fetchFullTranscript(sessionID)
|
||||||
|
}
|
||||||
|
|
||||||
/// Ask the host for the open session's full patch (Diff tab). No-op when the host
|
/// Ask the host for the open session's full patch (Diff tab). No-op when the host
|
||||||
/// doesn't advertise the capability — the view falls back to the diffstat summary.
|
/// doesn't advertise the capability — the view falls back to the diffstat summary.
|
||||||
func fetchDiff(_ sessionID: SessionID) {
|
func fetchDiff(_ sessionID: SessionID) {
|
||||||
@@ -1184,7 +1211,7 @@ final class RemoteStore: ObservableObject {
|
|||||||
// `requestPairingCode`/`cancelPairingCode` are sent straight to the chosen host by
|
// `requestPairingCode`/`cancelPairingCode` are sent straight to the chosen host by
|
||||||
// `requestPairingCode()`/`cancelPairingCode()`, not through this owner-routing switch.
|
// `requestPairingCode()`/`cancelPairingCode()`, not through this owner-routing switch.
|
||||||
case .hello, .ping, .listPeers, .addressUpdate, .meshRoster,
|
case .hello, .ping, .listPeers, .addressUpdate, .meshRoster,
|
||||||
.registerLiveActivity, .endLiveActivity,
|
.registerLiveActivity, .endLiveActivity, .registerPushToStartToken,
|
||||||
.transferOffer, .transferChunk, .transferCommit, .transferCancel, .fetchTranscript,
|
.transferOffer, .transferChunk, .transferCommit, .transferCancel, .fetchTranscript,
|
||||||
.requestPairingCode, .cancelPairingCode, .respondMacPair:
|
.requestPairingCode, .cancelPairingCode, .respondMacPair:
|
||||||
break
|
break
|
||||||
@@ -1201,6 +1228,12 @@ final class RemoteStore: ObservableObject {
|
|||||||
/// Host ids that already have the current token (re-sent to a host that (re)connects, and
|
/// Host ids that already have the current token (re-sent to a host that (re)connects, and
|
||||||
/// re-sent to everyone when the token rotates).
|
/// re-sent to everyone when the token rotates).
|
||||||
private var liveActivitySentTo: Set<String> = []
|
private var liveActivitySentTo: Set<String> = []
|
||||||
|
/// This device's push-to-start token (iOS 17.2+), shipped to every capable Mac so it can create
|
||||||
|
/// the Live Activity over APNs when work starts before the app is opened. Device-scoped, so —
|
||||||
|
/// unlike `liveActivityReg` — it persists across activities and isn't cleared when one ends.
|
||||||
|
private var pushToStartToken: String?
|
||||||
|
/// Host ids that already have the current push-to-start token (mirrors `liveActivitySentTo`).
|
||||||
|
private var pushToStartSentTo: Set<String> = []
|
||||||
|
|
||||||
private func setupLiveActivityBridge() {
|
private func setupLiveActivityBridge() {
|
||||||
LiveActivityManager.shared.onPushToken = { [weak self] token, activityID in
|
LiveActivityManager.shared.onPushToken = { [weak self] token, activityID in
|
||||||
@@ -1209,6 +1242,15 @@ final class RemoteStore: ObservableObject {
|
|||||||
self.liveActivitySentTo.removeAll() // a fresh token must reach every host again
|
self.liveActivitySentTo.removeAll() // a fresh token must reach every host again
|
||||||
self.syncLiveActivityRegistration()
|
self.syncLiveActivityRegistration()
|
||||||
}
|
}
|
||||||
|
LiveActivityManager.shared.onPushToStartToken = { [weak self] token in
|
||||||
|
guard let self, self.pushToStartToken != token else { return }
|
||||||
|
self.pushToStartToken = token
|
||||||
|
self.pushToStartSentTo.removeAll() // a fresh token must reach every host again
|
||||||
|
self.syncPushToStartRegistration()
|
||||||
|
}
|
||||||
|
// Start observing the push-to-start token (and adopting any push-started activity) now — it's
|
||||||
|
// device-scoped and must be captured even before any Activity exists.
|
||||||
|
LiveActivityManager.shared.beginPushToStartObservation()
|
||||||
LiveActivityManager.shared.onActivityEnded = { [weak self] activityID in
|
LiveActivityManager.shared.onActivityEnded = { [weak self] activityID in
|
||||||
guard let self else { return }
|
guard let self else { return }
|
||||||
self.liveActivityReg = nil
|
self.liveActivityReg = nil
|
||||||
@@ -1234,6 +1276,23 @@ final class RemoteStore: ObservableObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Send the current push-to-start token to every live, capable host that hasn't got it yet — the
|
||||||
|
/// mirror of `syncLiveActivityRegistration` for the device-scoped token that lets a host create
|
||||||
|
/// the glance over APNs when work starts before the app is opened (iOS 17.2+).
|
||||||
|
private func syncPushToStartRegistration() {
|
||||||
|
guard let token = pushToStartToken else { return }
|
||||||
|
// A host that dropped should re-register when it returns.
|
||||||
|
pushToStartSentTo = pushToStartSentTo.filter { connections[$0]?.connectivity.isLive == true }
|
||||||
|
for (id, conn) in connections {
|
||||||
|
// Gate on the dedicated push-to-start bit, not `canPushLiveActivity`: an older host can
|
||||||
|
// advertise the latter yet throw on the unknown `registerPushToStartToken` tag.
|
||||||
|
guard conn.connectivity.isLive, conn.capabilities.canPushToStartLiveActivity,
|
||||||
|
!pushToStartSentTo.contains(id) else { continue }
|
||||||
|
conn.send(.registerPushToStartToken(token))
|
||||||
|
pushToStartSentTo.insert(id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// MARK: - Demo simulator (offline, interactive)
|
// MARK: - Demo simulator (offline, interactive)
|
||||||
//
|
//
|
||||||
// In demo mode there's no host, so writes can't go over the wire. Instead they mutate the
|
// In demo mode there's no host, so writes can't go over the wire. Instead they mutate the
|
||||||
@@ -1291,7 +1350,7 @@ final class RemoteStore: ObservableObject {
|
|||||||
.addressUpdate, .meshRoster,
|
.addressUpdate, .meshRoster,
|
||||||
// Live Activity push registration is a real-connection concern (there's no host to
|
// Live Activity push registration is a real-connection concern (there's no host to
|
||||||
// push in demo), so it's inert here.
|
// push in demo), so it's inert here.
|
||||||
.registerLiveActivity, .endLiveActivity,
|
.registerLiveActivity, .endLiveActivity, .registerPushToStartToken,
|
||||||
// Session transfer (mesh P5) is a Mac↔Mac flow — the phone never originates these,
|
// Session transfer (mesh P5) is a Mac↔Mac flow — the phone never originates these,
|
||||||
// and demo has no peer Macs, so they're inert here.
|
// and demo has no peer Macs, so they're inert here.
|
||||||
.transferOffer, .transferChunk, .transferCommit, .transferCancel,
|
.transferOffer, .transferChunk, .transferCommit, .transferCancel,
|
||||||
|
|||||||
Reference in New Issue
Block a user