macOS主应用仅接收Privileged Helper首次更新的问题求助
macOS Privileged Helper 流式输出仅首次传递到主应用问题排查
问题背景
在macOS环境中开发主应用与Privileged Helper交互系统,需求是让Privileged Helper每秒以交互模式执行root命令,并将输出流式返回给主应用。实际运行时,主应用仅能接收首次更新的数据,后续助手在macOS Console中显示每秒正常输出,但无法传递到主应用。以下是简化代码及两次优化尝试:
主应用代码
ContentView.swift
import SwiftUI struct ContentView: View { @StateObject private var viewModel = ContentViewModel() var body: some View { VStack { ScrollView { Text(viewModel.scriptOutput) } Button("Stop Streaming") { if viewModel.isStreaming { viewModel.stopStreaming() } } } .padding() .onAppear { viewModel.startStreaming() } } } class ContentViewModel: ObservableObject { @Published var scriptOutput = "" @Published var isStreaming = false func startStreaming() { Task { do { try await ExecutionService.startStreamingData { [weak self] output in DispatchQueue.main.async { print(output) self?.scriptOutput = output } } await MainActor.run { self.isStreaming = true } } catch { await MainActor.run { self.scriptOutput = error.localizedDescription } } } } func stopStreaming() { Task { do { try await ExecutionService.stopStreamingData() await MainActor.run { self.isStreaming = false } } catch { await MainActor.run { self.scriptOutput = error.localizedDescription } } } } }
ExecutionService.swift
enum ExecutionService { // MARK: Execute static func startStreamingData(updateHandler: @escaping (String) -> Void) async throws { let helper = try await HelperRemoteProvider.remote() helper.startStreamingData { output in print("Received: \(output)") DispatchQueue.main.async { updateHandler(output) } } } static func stopStreamingData() async throws { let helper = try await HelperRemoteProvider.remote() helper.stopStreamingData() } }
HelperRemoteProvider.swift
import Foundation import ServiceManagement // MARK: - HelperRemoteProvider /// Provide a `HelperProtocol` object to request the helper. enum HelperRemoteProvider { // MARK: Computed private static var isHelperInstalled: Bool { FileManager.default.fileExists(atPath: HelperConstants.helperPath) } } // MARK: - Remote extension HelperRemoteProvider { static func remote() async throws -> some HelperProtocol { let connection = try connection() return try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<any HelperProtocol, Error>) in let continuationResume = ContinuationResume() let helper = connection.remoteObjectProxyWithErrorHandler { error in guard continuationResume.shouldResume() else { return } // 1st error to arrive, it will be the one thrown continuation.resume(throwing: error) } if let unwrappedHelper = helper as? HelperProtocol { continuation.resume(returning: unwrappedHelper) } else { if continuationResume.shouldResume() { // 1st error to arrive, it will be the one thrown let error = EchoError.helperConnection("Unable to get a valid 'HelperProtocol' object for an unknown reason") continuation.resume(throwing: error) } } } } } // MARK: - Install helper extension HelperRemoteProvider { /// Install the Helper in the privileged helper tools folder and load the daemon private static func installHelper() throws { // try to get a valid empty authorization var authRef: AuthorizationRef? try AuthorizationCreate(nil, nil, [.preAuthorize], &authRef).checkError("AuthorizationCreate") defer { if let authRef { AuthorizationFree(authRef, []) } } // create an AuthorizationItem to specify we want to bless a privileged Helper let authStatus = kSMRightBlessPrivilegedHelper.withCString { authorizationString in var authItem = AuthorizationItem(name: authorizationString, valueLength: 0, value: nil, flags: 0) return withUnsafeMutablePointer(to: &authItem) { pointer in var authRights = AuthorizationRights(count: 1, items: pointer) let flags: AuthorizationFlags = [.interactionAllowed, .extendRights, .preAuthorize] return AuthorizationCreate(&authRights, nil, flags, &authRef) } } guard authStatus == errAuthorizationSuccess else { throw EchoError.helperInstallation("Unable to get a valid loading authorization reference to load Helper daemon") } var blessErrorPointer: Unmanaged<CFError>? let wasBlessed = SMJobBless(kSMDomainSystemLaunchd, HelperConstants.domain as CFString, authRef, &blessErrorPointer) guard !wasBlessed else { return } // throw error since authorization was not blessed let blessError: Error = if let blessErrorPointer { blessErrorPointer.takeRetainedValue() as Error } else { EchoError.unknown } throw EchoError.helperInstallation("Error while installing the Helper: \(blessError.localizedDescription)") } } // MARK: - Connection extension HelperRemoteProvider { static private func connection() throws -> NSXPCConnection { if !isHelperInstalled { try installHelper() } return createConnection() } private static func createConnection() -> NSXPCConnection { let connection = NSXPCConnection(machServiceName: HelperConstants.domain, options: .privileged) connection.remoteObjectInterface = NSXPCInterface(with: HelperProtocol.self) connection.exportedInterface = NSXPCInterface(with: RemoteApplicationProtocol.self) connection.exportedObject = self connection.invalidationHandler = { if isHelperInstalled { print("Unable to connect to Helper although it is installed") } else { print("Helper is not installed") } } connection.resume() return connection } } // MARK: - ContinuationResume extension HelperRemoteProvider { /// Helper class to safely access a boolean when using a continuation to get the remote. private final class ContinuationResume: @unchecked Sendable { // MARK: Properties private let unfairLockPointer: UnsafeMutablePointer<os_unfair_lock_s> private var alreadyResumed = false // MARK: Computed /// `true` if the continuation should resume. func shouldResume() -> Bool { os_unfair_lock_lock(unfairLockPointer) defer { os_unfair_lock_unlock(unfairLockPointer) } if alreadyResumed { return false } else { alreadyResumed = true return true } } // MARK: Init init() { unfairLockPointer = UnsafeMutablePointer<os_unfair_lock_s>.allocate(capacity: 1) unfairLockPointer.initialize(to: os_unfair_lock()) } deinit { unfairLockPointer.deallocate() } } }
助手代码
HelperProtocol.swift
@objc public protocol HelperProtocol { @objc func startStreamingData(updateHandler: @escaping (String) -> Void) @objc func stopStreamingData() }
Helper.swift
import Foundation // MARK: - Helper final class Helper: NSObject { let listener: NSXPCListener private var streamingProcess: Process? private var updateHandler: ((String) -> Void)? override init() { listener = NSXPCListener(machServiceName: HelperConstants.domain) super.init() listener.delegate = self } } // MARK: - HelperProtocol extension Helper: HelperProtocol { func startStreamingData(updateHandler: @escaping (String) -> Void) { self.updateHandler = updateHandler streamingProcess = ExecutionService.streamData { [weak self] output in NSLog("Received output in Helper") self?.updateHandler?(output) } } func stopStreamingData() { streamingProcess?.terminate() streamingProcess = nil updateHandler = nil } } // MARK: - Run extension Helper { func run() { // start listening on new connections listener.resume() // prevent the terminal application to exit RunLoop.current.run() } } // MARK: - NSXPCListenerDelegate extension Helper: NSXPCListenerDelegate { func listener(_ listener: NSXPCListener, shouldAcceptNewConnection newConnection: NSXPCConnection) -> Bool { do { try ConnectionIdentityService.checkConnectionIsValid(connection: newConnection) } catch { NSLog("Connection \(newConnection) has not been validated. \(error.localizedDescription)") return false } newConnection.exportedInterface = NSXPCInterface(with: HelperProtocol.self) newConnection.remoteObjectInterface = NSXPCInterface(with: RemoteApplicationProtocol.self) newConnection.exportedObject = self newConnection.resume() return true } }
ExecutionService.swift
enum ExecutionService { static func streamData(updateHandler: @escaping (String) -> Void) -> Process { let process = Process() process.executableURL = URL(fileURLWithPath: "/path/to/root/command") process.arguments = [ "-i", "1000" // update every 1 second ] let outputPipe = Pipe() process.standardOutput = outputPipe process.standardError = outputPipe let outHandle = outputPipe.fileHandleForReading outHandle.readabilityHandler = { fileHandle in autoreleasepool { let data = fileHandle.availableData if data.count > 0 { if let output = String(data: data, encoding: .utf8) { NSLog("🟢 Received output: %@", output) updateHandler(output) } else { NSLog("🔴 Failed to convert output to string") } } else { NSLog("🔴 Received empty data from data, ending read") fileHandle.readabilityHandler = nil } } } do { try process.run() NSLog("🟠 process started") } catch { NSLog("🔴 Error starting process: %@", error.localizedDescription) } return process } }
更新尝试
Update 1:改用AsyncThrowingStream处理流式数据
主应用ExecutionService.swift和ContentView.swift代码更新后,问题仍未解决:
// ExecutionService.swift enum ExecutionService { static func startStreamingData() async throws -> AsyncThrowingStream<String, Error> { let helper = try await HelperRemoteProvider.remote() return AsyncThrowingStream { continuation in helper.startStreamingData { output in continuation.yield(output) } continuation.onTermination = { @Sendable _ in Task { await helper.stopStreamingData() } } } } static func stopStreamingData() async throws { let helper = try await HelperRemoteProvider.remote() helper.stopStreamingData() } }
// ContentView.swift中ContentViewModel class ContentViewModel: ObservableObject { @Published var scriptOutput = "" @Published var isStreaming = false private var streamingTask: Task<Void, Never>? func startStreaming() { streamingTask = Task { do { let stream = try await ExecutionService.startStreamingData() for try await output in stream { await MainActor.run { self.scriptOutput = output self.isStreaming = true } } } catch { await MainActor.run { self.scriptOutput = error.localizedDescription } } } } func stopStreaming() { Task { do { try await ExecutionService.stopStreamingData() streamingTask?.cancel() await MainActor.run { self.isStreaming = false } } catch { await MainActor.run { self.scriptOutput = error.localizedDescription } } } } }
Update 2:自定义AsyncSequence实现数据序列
新增Data.swift代码后,问题仍未解决:
import Foundation struct DataSequence: AsyncSequence { typealias Element = String let helper: HelperProtocol func makeAsyncIterator() -> AsyncIterator { return AsyncIterator(helper: helper) } class AsyncIterator: AsyncIteratorProtocol { private let helper: HelperProtocol private var continuations: [CheckedContinuation<String?, Error>] = [] private var isStreaming = false private var isFinished = false init(helper: HelperProtocol) { self.helper = helper } func next() async throws -> String? { print("NEXT") if isFinished { return nil } if !isStreaming { print("isStreaming") isStreaming = true helper.startStreamingData { [weak self] output in guard let self = self else { return } print("my output: \(output)") if let continuation = self.continuations.first { self.continuations.removeFirst() continuation.resume(returning: output) } } } return try await withCheckedThrowingContinuation { continuation in continuations.append(continuation) } } func cancel() { isFinished = true helper.stopStreamingData() for continuation in continuations { continuation.resume(returning: nil) } continuations.removeAll() } } }
问题排查与解决方案
核心问题分析
- NSXPC连接未复用:每次调用
start/stopStreamingData都会创建新的NSXPC连接,原连接被释放后,后续回调无法传递到主应用。 - XPC回调生命周期管理不当:Helper中的
updateHandler未被正确保留,或主应用回调闭包因连接失效被提前释放。 - Process输出处理不完整:Pipe的
readabilityHandler未处理数据缓冲,可能导致部分输出被截断或未及时传递。
修复方案
1. 复用NSXPC连接
修改HelperRemoteProvider,保留全局连接实例,避免重复创建:
// HelperRemoteProvider.swift private static var sharedConnection: NSXPCConnection? static private func connection() throws -> NSXPCConnection { if let connection = sharedConnection, connection.isValid { return connection } if !isHelperInstalled { try installHelper() } let newConnection = createConnection() sharedConnection = newConnection return newConnection } private static func createConnection() -> NSXPCConnection { let connection = NSXPCConnection(machServiceName: HelperConstants.domain, options: .privileged) connection.remoteObjectInterface = NSXPCInterface(with: HelperProtocol.self) connection.exportedInterface = NSXPCInterface(with: RemoteApplicationProtocol.self) connection.exportedObject = self connection.invalidationHandler = { sharedConnection = nil // 连接失效时清空实例 if isHelperInstalled { print("Unable to connect to Helper although it is installed") } else { print("Helper is not installed") } } connection.resume() return connection }
2. 确保XPC回调正确性
更新协议和Helper中的回调处理,保证闭包被正确持有:
// HelperProtocol.swift @objc public protocol HelperProtocol { @objc func startStreamingData(updateHandler: @escaping @Sendable (String) -> Void) @objc func stopStreamingData() } // Helper.swift extension Helper: HelperProtocol { func startStreamingData(updateHandler: @escaping @Sendable (String) -> Void) { stopStreamingData() // 先终止之前的流 self.updateHandler = updateHandler streamingProcess = ExecutionService.streamData { [weak self] output in guard let self = self else { return } NSLog("Received output in Helper") DispatchQueue.main.async { // 在XPC队列安全调用回调 self.updateHandler?(output) } } } }
3. 优化Process输出处理
处理数据缓冲,按行分割输出,避免截断:
// ExecutionService.swift static func streamData(updateHandler: @escaping @Sendable (String) -> Void) -> Process { let process = Process() process.executableURL = URL(fileURLWithPath: "/path/to/root/command") process.arguments = ["-i", "1000"] let outputPipe = Pipe() process.standardOutput = outputPipe process.standardError = outputPipe
相关产品推荐
相关产品推荐

