You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()
        }
    }
}

问题排查与解决方案

核心问题分析

  1. NSXPC连接未复用:每次调用start/stopStreamingData都会创建新的NSXPC连接,原连接被释放后,后续回调无法传递到主应用。
  2. XPC回调生命周期管理不当:Helper中的updateHandler未被正确保留,或主应用回调闭包因连接失效被提前释放。
  3. 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
    
相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 07:12:41