Merge nucleic/humble-harbor-viper into dev
This commit is contained in:
@@ -396,7 +396,13 @@ final class HostConnection {
|
||||
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))
|
||||
// Also offer this phone as a low-tier AFM executor (ANTIMATTER_RUNNER §5) when the
|
||||
// OS has Foundation Models — a runner host's mesh queue may place background work
|
||||
// here; the interactive tiers stay on desktop-class devices by the queue's rules.
|
||||
clientCaps: WireClientCapabilities(
|
||||
mesh: 1, canSyncRoster: true,
|
||||
canProvideIntelligence: PhoneIntelligenceExecutor.isSupported,
|
||||
intelligenceProfile: PhoneIntelligenceExecutor.profile))
|
||||
self.client = client
|
||||
consume(client, pairingPayload: pairingPayload)
|
||||
}
|
||||
@@ -647,14 +653,22 @@ final class HostConnection {
|
||||
// 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,
|
||||
case .intelligenceRequest(let request):
|
||||
// A runner host delegated one AFM job here (docs/ANTIMATTER_RUNNER.md §5) — it only
|
||||
// ever sends these after this app advertised `canProvideIntelligence`. Generate off
|
||||
// the event stream (a model call takes seconds) and answer with the same id.
|
||||
Task { [weak self] in
|
||||
let result = await PhoneIntelligenceExecutor.execute(request)
|
||||
self?.send(.intelligenceResult(result))
|
||||
}
|
||||
case .credentialNeeded, .credentialUpdate,
|
||||
// The owner's runner-pool credential (item 4) — inert until the phone grows a
|
||||
// pool-management surface; Macs are the managers today.
|
||||
.runnerPoolCredential:
|
||||
// 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.
|
||||
// Antimatter runner credential verbs (docs/ANTIMATTER_RUNNER.md §6): a runner host
|
||||
// asking for / mirroring sealed credentials. Inert here until the phone-side vault
|
||||
// lands — 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 {
|
||||
|
||||
@@ -0,0 +1,135 @@
|
||||
import Foundation
|
||||
import NucleicProtocol
|
||||
#if canImport(FoundationModels)
|
||||
import FoundationModels
|
||||
#endif
|
||||
|
||||
/// The phone-side executor for delegated intelligence work (docs/ANTIMATTER_RUNNER.md §5,
|
||||
/// item 6): a runner host pushes `HostMsg.intelligenceRequest` at this device — the mesh AFM
|
||||
/// queue only ever sends it the lower tiers (`completion`/`background`), which don't need
|
||||
/// desktop tok/s — and this renders the shared `IntelligenceDelegate` template on the local
|
||||
/// Apple Foundation Models and answers `ClientMsg.intelligenceResult` with the same id.
|
||||
///
|
||||
/// Generations run one at a time through ``SerialGate`` (the phone's Neural Engine is a single
|
||||
/// resource, same reasoning as the Mac's `AFMRequestQueue`), and every failure — model off,
|
||||
/// unavailable, unknown kind — answers with `error` set so the runner falls back to heuristics
|
||||
/// immediately instead of waiting out its deadline.
|
||||
enum PhoneIntelligenceExecutor {
|
||||
/// Whether this device can execute at all (Foundation Models exist on this OS). Gates
|
||||
/// advertising `canProvideIntelligence`; live availability (Apple Intelligence enabled,
|
||||
/// model downloaded) is re-checked per request.
|
||||
static var isSupported: Bool {
|
||||
#if canImport(FoundationModels)
|
||||
if #available(iOS 26, *) { return true }
|
||||
#endif
|
||||
return false
|
||||
}
|
||||
|
||||
/// The worker profile advertised in the hello: mobile class (only the lower priority
|
||||
/// tiers land here) with the SoC name for the queue's power bias — an M-series iPad
|
||||
/// outranks an A-series iPhone.
|
||||
static var profile: IntelligenceWorkerProfile? {
|
||||
guard isSupported else { return nil }
|
||||
return IntelligenceWorkerProfile(deviceClass: .mobile, chip: chipName)
|
||||
}
|
||||
|
||||
/// The SoC brand string ("Apple A18 Pro", "Apple M4"), falling back to the hardware model
|
||||
/// identifier ("iPhone17,1" — unranked but still telling) when the sysctl is unreadable.
|
||||
static var chipName: String? {
|
||||
sysctlString("machdep.cpu.brand_string") ?? sysctlString("hw.machine")
|
||||
}
|
||||
|
||||
private static func sysctlString(_ name: String) -> String? {
|
||||
var size = 0
|
||||
guard sysctlbyname(name, nil, &size, nil, 0) == 0, size > 0 else { return nil }
|
||||
var buffer = [CChar](repeating: 0, count: size)
|
||||
guard sysctlbyname(name, &buffer, &size, nil, 0) == 0 else { return nil }
|
||||
let value = String(cString: buffer).trimmingCharacters(in: .whitespaces)
|
||||
return value.isEmpty ? nil : value
|
||||
}
|
||||
|
||||
/// Run one delegated request end to end. Never throws — every failure mode answers with
|
||||
/// `error` set, correlated by the request id.
|
||||
static func execute(_ request: WireIntelligenceRequest) async -> WireIntelligenceResult {
|
||||
guard let job = IntelligenceDelegate.job(for: request) else {
|
||||
return WireIntelligenceResult(
|
||||
id: request.id, error: "unsupported kind: \(request.kind.rawValue)")
|
||||
}
|
||||
#if canImport(FoundationModels)
|
||||
if #available(iOS 26, *) {
|
||||
let model = SystemLanguageModel.default
|
||||
guard model.isAvailable else {
|
||||
return WireIntelligenceResult(id: request.id, error: "model unavailable")
|
||||
}
|
||||
// Bound the total wait (queue + generation) by the request's deadline: the host
|
||||
// has already timed the job out by then, and a wedged `respond` must not park
|
||||
// every later `execute` behind the stall — the deadline path answers an error and
|
||||
// moves on. (An uncancellable stuck generation itself can't be reclaimed; it is
|
||||
// abandoned in the background.)
|
||||
let deadline = request.deadlineSeconds ?? 60
|
||||
let text = await raceAgainstDeadline(seconds: deadline) {
|
||||
await gate.run {
|
||||
try? await LanguageModelSession(model: model, instructions: job.instructions)
|
||||
.respond(to: job.prompt).content
|
||||
}
|
||||
}
|
||||
if let text, !text.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty {
|
||||
return WireIntelligenceResult(id: request.id, outputs: [text])
|
||||
}
|
||||
return WireIntelligenceResult(id: request.id, error: "generation failed")
|
||||
}
|
||||
#endif
|
||||
return WireIntelligenceResult(id: request.id, error: "model unavailable")
|
||||
}
|
||||
|
||||
private static let gate = SerialGate()
|
||||
|
||||
/// First-resume race between `operation` and a deadline: whichever finishes first answers;
|
||||
/// the loser's resume is dropped by the once-latch. A task group can't express this — its
|
||||
/// scope waits for ALL children, so an uncancellable straggler would still block the exit.
|
||||
private static func raceAgainstDeadline(
|
||||
seconds: Double, _ operation: @escaping @Sendable () async -> String?
|
||||
) async -> String? {
|
||||
let once = FirstResume()
|
||||
return await withCheckedContinuation { continuation in
|
||||
Task {
|
||||
let value = await operation()
|
||||
if once.take() { continuation.resume(returning: value) }
|
||||
}
|
||||
Task {
|
||||
try? await Task.sleep(for: .seconds(seconds))
|
||||
if once.take() { continuation.resume(returning: nil) }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A thread-safe "fire once" latch so the deadline race resumes its continuation exactly once.
|
||||
private final class FirstResume: @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
|
||||
}
|
||||
}
|
||||
|
||||
/// A minimal FIFO gate: chains each operation behind the previous one so delegated
|
||||
/// generations never contend for the Neural Engine (an actor alone doesn't serialize across
|
||||
/// its suspension points).
|
||||
private actor SerialGate {
|
||||
private var tail: Task<Void, Never>?
|
||||
|
||||
func run<T: Sendable>(_ operation: @escaping @Sendable () async -> T) async -> T {
|
||||
let previous = tail
|
||||
let task = Task<T, Never> {
|
||||
await previous?.value
|
||||
return await operation()
|
||||
}
|
||||
tail = Task { _ = await task.value }
|
||||
return await task.value
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user