207 lines
9.0 KiB
Swift
207 lines
9.0 KiB
Swift
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 { lhs, rhs in
|
||
let l = StatusStyle.sortRank(lhs), r = StatusStyle.sortRank(rhs)
|
||
return l == r ? lhs.updatedAt > rhs.updatedAt : l < r
|
||
}
|
||
|
||
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)
|
||
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
|
||
guard let activity else { return }
|
||
self.activity = nil
|
||
onActivityEnded?(activity.id) // let the Macs stop pushing to this token
|
||
Task {
|
||
await activity.end(
|
||
ActivityContent(state: activity.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)
|
||
}
|
||
}
|
||
}
|