Files
nucleic-remote-ios/NucleicRemote/NucleicRemote/LiveActivityManager.swift
T

228 lines
10 KiB
Swift
Raw Normal View History

import ActivityKit
import Foundation
import NucleicProtocol
/// Owns the one aggregate session Live Activity (UX_IOS §5.3): started when work exists,
/// updated as sessions change, ended when everything is idle or the device unpairs. State
/// flows in from `RemoteStore` on every session-list change; the widget extension renders it
/// (`SessionLiveActivity`).
@MainActor
final class LiveActivityManager {
static let shared = LiveActivityManager()
private init() {}
/// How many sessions the detail rows show. The glance stays a glance — the counts and churn
/// still summarize everything, this just bounds the per-session list.
private static let maxLines = 3
private var activity: Activity<NucleicSessionAttributes>?
/// The last content we pushed. Updates that don't change it are skipped so we don't spend
/// ActivityKit's update budget on no-ops — `sync` fires on every host message (dashboard,
/// connectivity, pong, diff ticks…), most of which leave the aggregate identical. Burning the
/// budget on those is exactly what makes a *real* change land late and the glance read stale.
private var lastState: NucleicSessionAttributes.ContentState?
/// The newest state waiting to be applied, and the single task draining it. Coalescing to the
/// latest through one serial task means the newest data always wins — firing an unstructured
/// `Task` per `sync` let a later update lose a race to an earlier one and freeze the glance.
private var pendingState: NucleicSessionAttributes.ContentState?
private var updateTask: Task<Void, Never>?
/// Set by `RemoteStore` to ship the activity's APNS push token to the paired Macs (and to tell
/// them it ended). The Macs use the token to keep this glance fresh over APNs while the phone
/// is backgrounded and its sync socket is suspended (UX_IOS §5.3).
var onPushToken: ((_ token: String, _ activityID: String) -> Void)?
var onActivityEnded: ((_ activityID: String) -> Void)?
/// Streams the activity's per-activity APNS update token (it can rotate); cancelled on end.
private var tokenObservation: Task<Void, Never>?
/// Reconcile the Activity with the current session set.
func sync(hostName: String, sessions: [WireSessionSummary]) {
guard ActivityAuthorizationInfo().areActivitiesEnabled else { return }
let live = sessions.filter { !$0.archived }
let running = live.filter { $0.status == .running || $0.status == .provisioning }
let needsYou = live.filter { $0.status.needsYou($0.disposition) }
guard !running.isEmpty || !needsYou.isEmpty else {
end()
return
}
// Everything in flight or waiting on the user, attention-first (approvals, then waiting
// input, then running), freshest within a rank. This is both the detail-row source and
// the set the aggregate churn/approval totals sum over.
let active = live
.filter {
$0.status == .running || $0.status == .provisioning
|| $0.status.needsYou($0.disposition)
}
.sorted(by: StatusStyle.attentionThenRecency)
let approvals = active.reduce(0) { $0 + $1.pendingApprovalCount }
let files = active.reduce(0) { $0 + ($1.diffStat?.filesChanged ?? 0) }
let added = active.reduce(0) { $0 + ($1.diffStat?.added ?? 0) }
let removed = active.reduce(0) { $0 + ($1.diffStat?.removed ?? 0) }
let lines = active.prefix(Self.maxLines).map(Self.line(for:))
let state = NucleicSessionAttributes.ContentState(
runningCount: running.count,
needsYouCount: needsYou.count,
approvalCount: approvals,
filesChanged: files,
linesAdded: added,
linesRemoved: removed,
lines: Array(lines))
push(state, hostName: hostName)
}
/// Apply `state` to the Activity — deduped against the last push and coalesced through one
/// serial task so the newest state always wins.
private func push(_ state: NucleicSessionAttributes.ContentState, hostName: String) {
guard state != lastState else { return }
lastState = state
guard let activity else {
// Recover an Activity that survived an app relaunch before starting a new one; the
// recovered one still needs the fresh state, so fall through to the update pipeline.
if let existing = Activity<NucleicSessionAttributes>.activities.first {
activity = existing
observePushToken(existing)
2026-07-06 12:40:19 -07:00
// Dismiss any duplicates so only the adopted Activity renders. iOS stacks multiple
// Activities of one type in the Dynamic Island — an orphan (from a prior launch that
// was killed before `end()`) then shows *its* stale state in the collapsed pill while
// expanding reveals the fresh one, and `end()` on completion would leave it lingering
// on the last "needs attention" glance. Keep exactly one.
endStrays(keeping: existing.id)
enqueue(state)
} else {
// `pushType: .token` opts the activity into APNs updates — the Macs push new
// content-state to the token so the glance stays fresh while the phone is locked.
let started = try? Activity.request(
attributes: NucleicSessionAttributes(hostName: hostName),
content: ActivityContent(state: state, staleDate: nil),
pushType: .token)
activity = started
if let started { observePushToken(started) }
}
return
}
enqueue(state)
}
/// Forward the activity's APNS update token (and its rotations) to `RemoteStore`.
private func observePushToken(_ activity: Activity<NucleicSessionAttributes>) {
tokenObservation?.cancel()
tokenObservation = Task { [weak self] in
for await tokenData in activity.pushTokenUpdates {
let hex = tokenData.map { String(format: "%02x", $0) }.joined()
self?.onPushToken?(hex, activity.id)
}
}
}
/// Hand the newest state to the serial drainer (starting it if idle).
private func enqueue(_ state: NucleicSessionAttributes.ContentState) {
pendingState = state
guard updateTask == nil else { return } // the running drainer will pick this up
updateTask = Task { @MainActor [weak self] in
guard let self else { return }
while let next = self.pendingState {
self.pendingState = nil
await self.activity?.update(ActivityContent(state: next, staleDate: nil))
}
self.updateTask = nil
}
}
/// End the Activity (all idle, or unpaired).
func end() {
updateTask?.cancel()
updateTask = nil
tokenObservation?.cancel()
tokenObservation = nil
pendingState = nil
lastState = nil
2026-07-06 12:40:34 -07:00
let tracked = activity
self.activity = nil
2026-07-06 12:40:34 -07:00
if let tracked { onActivityEnded?(tracked.id) } // let the Macs stop pushing to this token
// Dismiss the tracked Activity *and any strays*: ending only the one we track would leave an
// orphan (from a prior launch killed before `end()`) rendering a stale "needs attention"
// glance in the Dynamic Island after the work it described has finished. Clear them all so
// "everything idle" can't get stuck on screen.
Task {
2026-07-06 12:40:34 -07:00
for activity in Activity<NucleicSessionAttributes>.activities {
await activity.end(
ActivityContent(state: activity.content.state, staleDate: nil),
dismissalPolicy: .immediate)
}
}
}
/// Dismiss every live Activity except `keep` — collapses accidental duplicates down to one so
/// iOS can't render a stale orphan alongside the tracked glance.
private func endStrays(keeping keep: String) {
for stray in Activity<NucleicSessionAttributes>.activities where stray.id != keep {
Task {
await stray.end(
ActivityContent(state: stray.content.state, staleDate: nil),
dismissalPolicy: .immediate)
}
}
}
// MARK: - Projection
/// Project a wire summary into the widget's self-contained row model.
private static func line(for s: WireSessionSummary) -> NucleicSessionAttributes.SessionLine {
NucleicSessionAttributes.SessionLine(
id: s.sessionID.rawValue,
title: s.title.isEmpty ? s.projectName : s.title,
project: s.projectName,
backend: backend(s.backend),
kind: kind(for: s),
detail: detail(for: s))
}
private static func kind(for s: WireSessionSummary) -> NucleicSessionAttributes.Kind {
switch s.status {
case .awaitingApproval: .approval
case .running: .running
case .provisioning: .provisioning
case .awaitingInput: s.disposition == .completed ? .done : .needsInput
case .idle: .idle
case .finished: .done
case .interrupted, .error: .error
}
}
private static func backend(_ id: BackendID) -> NucleicSessionAttributes.Backend {
switch id {
case .claudeCode: .claude
case .codex, .codexExec: .codex
case .grok: .grok
}
}
/// The compact right-aligned status for a row: the actionable ask wins (approvals, then
/// "waiting on you"), else the worktree churn, else what the agent is up to.
private static func detail(for s: WireSessionSummary) -> String {
if s.status == .awaitingApproval || s.pendingApprovalCount > 0 {
let n = max(s.pendingApprovalCount, 1)
return "\(n) to approve"
}
if s.status == .awaitingInput, s.disposition != .completed {
return "Waiting on you"
}
if let d = s.diffStat, d.filesChanged > 0 {
return "\(d.filesChanged) file\(d.filesChanged == 1 ? "" : "s") +\(d.added) −\(d.removed)"
}
switch s.status {
case .provisioning: return "Starting…"
case .running: return "Working…"
case .awaitingInput: return "Done"
default: return StatusStyle.label(s.status, disposition: s.disposition)
}
}
}