native macOS codings agent orchestrator prowl.onev.cat
Something went wrong. Try again.
Swift
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719import Foundationimport ProwlCLISharedimport Testing
@testable import Prowl
#if canImport(Darwin) import Darwin#elseif canImport(Glibc) import Glibc#endif
@MainActorstruct CLISocketServerTests { @Test func workflowPayloadEscapingFitsTheSocketFrame() async throws { let socketPath = temporarySocketPath(suffix: "workflow-large") let server = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try server.start() defer { server.stop() } let request = try JSONEncoder().encode( CommandEnvelope( output: .json, command: .workflow(.init(action: .deliver, body: String(repeating: "\u{1}", count: WorkflowSizeLimits.payload))) )) #expect(request.count > 32 * 1024 * 1024) #expect(request.count <= WorkflowSizeLimits.transportFrame) let response: Data = try await withCheckedThrowingContinuation { continuation in DispatchQueue.global(qos: .userInitiated).async { continuation.resume(with: Result { try Self.send(requestData: request, socketPath: socketPath) }) } } #expect(try JSONDecoder().decode(CommandResponse.self, from: response).command == "workflow") }
@Test func secondServerDoesNotReplaceReachableSocket() throws { let socketPath = temporarySocketPath(suffix: "reachable-owner") let first = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try first.start() defer { first.stop() }
#expect(canConnect(to: socketPath))
let second = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) #expect(throws: CLIServiceError.socketAlreadyOwned) { try second.start() }
#expect(canConnect(to: socketPath)) }
@Test func statusFollowsStartAndStop() throws { let socketPath = temporarySocketPath(suffix: "status") var published: [CLIServiceStatus] = [] let server = CLISocketServer( router: CLICommandRouter(), socketPath: socketPath, onStatusChanged: { published.append($0) }) #expect(server.status == .stopped)
try server.start() #expect(server.status == .listening(path: socketPath))
server.stop() #expect(server.status == .stopped) #expect(published == [.listening(path: socketPath), .stopped]) }
@Test func statusReportsAnAlreadyOwnedSocket() throws { let socketPath = temporarySocketPath(suffix: "status-owned") let first = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try first.start() defer { first.stop() }
let second = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) #expect(throws: CLIServiceError.socketAlreadyOwned) { try second.start() } #expect(second.status == .failed(.socketAlreadyOwned, path: socketPath)) }
@Test func ownerCanReplaceStaleSocketPath() throws { let socketPath = temporarySocketPath(suffix: "stale-owner") try createStaleSocket(at: socketPath) #expect(FileManager.default.fileExists(atPath: socketPath)) #expect(!canConnect(to: socketPath))
let server = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try server.start() defer { server.stop() }
#expect(canConnect(to: socketPath)) }
@Test func nonOwnerStopDoesNotRemoveOwnedSocketPath() throws { let socketPath = temporarySocketPath(suffix: "non-owner-stop") let first = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try first.start() defer { first.stop() }
let second = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) #expect(throws: CLIServiceError.socketAlreadyOwned) { try second.start() } second.stop()
#expect(canConnect(to: socketPath)) }
// `debugFileDescriptors` exists only in Debug builds, and so does this test — // `make bench` compiles this target under the Release configuration. #if DEBUG @Test func ownedDescriptorsAreClosedOnExec() throws { let socketPath = temporarySocketPath(suffix: "cloexec") let server = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try server.start() defer { server.stop() }
let descriptors = server.debugFileDescriptors #expect(isCloseOnExec(descriptors.server)) #expect(isCloseOnExec(descriptors.lock)) } #endif
@Test func socketFilesAreOwnerOnly() throws { let socketPath = temporarySocketPath(suffix: "permissions") let socketDirectory = (socketPath as NSString).deletingLastPathComponent let lockPath = "\(socketPath).lock" let server = CLISocketServer(router: CLICommandRouter(), socketPath: socketPath) try server.start() defer { server.stop() }
#expect(fileMode(at: socketDirectory) == 0o700) #expect(fileMode(at: socketPath) == 0o600) #expect(fileMode(at: lockPath) == 0o600) }
@Test func peerUIDMustMatchCurrentUser() { #expect(CLISocketServer.isAllowedPeerUID(501, currentUID: 501)) #expect(!CLISocketServer.isAllowedPeerUID(502, currentUID: 501)) }
@Test func descriptorDuplicationFailureRejectsConnectionBeforeRouting() async throws { let socketPath = temporarySocketPath(suffix: "monitor-dup-failure") let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) var recordedSignal: AgentSignal? let handler = AgentSignalCommandHandler( resolveCaller: { _ in pane }, recordSignal: { _, signal in recordedSignal = signal return .recorded(binding: .current) } ) let accepted = DispatchSemaphore(value: 0) let closed = DispatchSemaphore(value: 0) let peerClosed = DisconnectObservation() let rejected = DisconnectObservation() let server = CLISocketServer( router: CLICommandRouter(agentsSignalHandler: handler), socketPath: socketPath, onClientAccepted: { accepted.signal() if closed.wait(timeout: .now() + 10) == .success { peerClosed.signal() } }, onPeerMonitorUnavailable: { rejected.signal() }, duplicatePeerDescriptor: { _ in -1 } ) try server.start() defer { server.stop() } let requestData = try JSONEncoder().encode( CommandEnvelope( output: .json, command: .agentsSignal(AgentSignalInput(event: .turnEnded, detail: "must not route")) ) )
try await Task.detached { try Self.sendAndCloseAfterAcceptance( requestData: requestData, socketPath: socketPath, accepted: accepted, closed: closed ) }.value let closedBeforeRouting = await Task.detached { peerClosed.wait(timeout: .now() + 10) }.value #expect(closedBeforeRouting) let connectionRejected = await Task.detached { rejected.wait(timeout: .now() + 10) }.value
#expect(connectionRejected) #expect(recordedSignal == nil) }
@Test func socketRoundTripThreadsKernelPeerPIDIntoAgentSignalHandler() async throws { let socketPath = temporarySocketPath(suffix: "signal-context") let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) var recordedSignal: AgentSignal? let handler = AgentSignalCommandHandler( resolveCaller: { context in #expect(context.callerProcessID == getpid()) return pane }, recordSignal: { caller, signal in #expect(caller == pane) recordedSignal = signal return .recorded(binding: .current) }, now: { Date(timeIntervalSince1970: 1_000) } ) let server = CLISocketServer( router: CLICommandRouter(agentsSignalHandler: handler), socketPath: socketPath ) try server.start() defer { server.stop() } let envelope = CommandEnvelope( output: .json, command: .agentsSignal(AgentSignalInput(event: .turnEnded, detail: "socket result")) )
let requestData = try JSONEncoder().encode(envelope) let responseData = try await withCheckedThrowingContinuation { continuation in DispatchQueue.global(qos: .userInitiated).async { continuation.resume( with: Result { try Self.send(requestData: requestData, socketPath: socketPath) } ) } } let response = try JSONDecoder().decode(CommandResponse.self, from: responseData)
#expect(response.ok) #expect(response.command == "agents.signal") #expect(recordedSignal?.kind == .turnEnded) #expect(recordedSignal?.detail == "socket result") }
@Test func nativeHookRoundTripThreadsKernelPeerPIDWithoutExposingToken() async throws { let socketPath = temporarySocketPath(suffix: "hook-context") let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) let input = AgentNativeHookInput( runtime: .codex, token: "private-token", signal: AgentNativeHookSignal( event: .turnEnded, nativeEvent: "agent-turn-complete", cwd: "/tmp/project", sessionID: "thread-1" ) ) var recordedInput: AgentNativeHookInput? let handler = AgentNativeHookCommandHandler( resolveCaller: { context in #expect(context.callerProcessID == getpid()) #expect(context.callerProcessAncestry.first?.processID == getpid()) return pane }, recordHook: { caller, received in #expect(caller == pane) recordedInput = received return true } ) let server = CLISocketServer( router: CLICommandRouter(agentsHookHandler: handler), socketPath: socketPath ) try server.start() defer { server.stop() } let requestData = try JSONEncoder().encode( CommandEnvelope(output: .json, command: .agentsHook(input)) )
let responseData = try await withCheckedThrowingContinuation { continuation in DispatchQueue.global(qos: .userInitiated).async { continuation.resume( with: Result { try Self.send(requestData: requestData, socketPath: socketPath) } ) } } let response = try JSONDecoder().decode(CommandResponse.self, from: responseData)
#expect(response.ok) #expect(recordedInput == input) let responseText = try #require(String(bytes: responseData, encoding: .utf8)) #expect(!responseText.contains(input.token)) }
@Test func nativeHookRoutesAfterPeerClosesBeforeMainActorHandling() async throws { let socketPath = temporarySocketPath(suffix: "hook-short-lived-peer") let routeRecorded = DisconnectObservation() let pane = CallerPane(worktreeID: "wt", surfaceID: UUID()) let input = AgentNativeHookInput( runtime: .codex, token: "private-token", signal: AgentNativeHookSignal( event: .turnEnded, nativeEvent: "agent-turn-complete", cwd: "/tmp/project", sessionID: "thread-1" ) ) var recordedInput: AgentNativeHookInput? let handler = AgentNativeHookCommandHandler( resolveCaller: { context in #expect(context.callerProcessID == getpid()) #expect(context.callerProcessAncestry.first?.processID == getpid()) return pane }, recordHook: { _, received in recordedInput = received routeRecorded.signal() return true } ) let accepted = DispatchSemaphore(value: 0) let closed = DispatchSemaphore(value: 0) let peerClosed = DisconnectObservation() let server = CLISocketServer( router: CLICommandRouter(agentsHookHandler: handler), socketPath: socketPath, onClientAccepted: { accepted.signal() if closed.wait(timeout: .now() + 10) == .success { peerClosed.signal() } } ) try server.start() defer { server.stop() } let requestData = try JSONEncoder().encode( CommandEnvelope(output: .json, command: .agentsHook(input)) ) let sendTask = Task.detached { try Self.sendAndCloseAfterAcceptance( requestData: requestData, socketPath: socketPath, accepted: accepted, closed: closed ) }
try await sendTask.value let closedBeforeRouting = await Task.detached { peerClosed.wait(timeout: .now() + 10) }.value #expect(closedBeforeRouting) let didRoute = await Task.detached { routeRecorded.wait(timeout: .now() + 10) }.value #expect(didRoute) #expect(recordedInput == input) }
@Test func closingPeerCancelsInFlightWaitRequest() async throws { let socketPath = temporarySocketPath(suffix: "wait-peer-eof") let probe = CancellationProbe() let accepted = DispatchSemaphore(value: 0) let closed = DispatchSemaphore(value: 0) let peerClosed = DisconnectObservation() let server = CLISocketServer( router: CLICommandRouter(agentsWaitHandler: CancellationProbeHandler(probe: probe)), socketPath: socketPath, onClientAccepted: { accepted.signal() if closed.wait(timeout: .now() + 10) == .success { peerClosed.signal() } } ) try server.start() defer { server.stop() } let envelope = CommandEnvelope( output: .json, command: .agentsWait( AgentWaitInput(mode: .dispatch, dispatchID: "dispatch-peer-eof", timeoutSeconds: 600) ) ) let requestData = try JSONEncoder().encode(envelope)
try await Task.detached { try Self.sendAndCloseAfterAcceptance( requestData: requestData, socketPath: socketPath, accepted: accepted, closed: closed ) }.value
let closedBeforeRouting = await Task.detached { peerClosed.wait(timeout: .now() + 10) }.value #expect(closedBeforeRouting) let cancellationObserved = await Task.detached { probe.waitForCancellation(timeout: .now() + 10) }.value #expect(cancellationObserved) #expect(probe.wasCancelled) }
#if canImport(Darwin) @Test func disconnectMonitorActivatesDuringCreation() async throws { var descriptors: [Int32] = [-1, -1] #expect(socketpair(AF_UNIX, SOCK_STREAM, 0, &descriptors) == 0) defer { close(descriptors[0]) } let observation = DisconnectObservation() let monitor = try #require( CLIPeerDisconnectMonitor(fileDescriptor: descriptors[0]) { observation.signal() } )
close(descriptors[1]) let activatedDuringCreation = await Task.detached { observation.wait(timeout: .now() + 1) }.value
#expect(activatedDuringCreation) monitor.cancel() }
@Test func disconnectMonitorOutlivesOriginalDescriptor() async throws { var descriptors: [Int32] = [-1, -1] #expect(socketpair(AF_UNIX, SOCK_STREAM, 0, &descriptors) == 0) let observation = DisconnectObservation() let monitor = try #require( CLIPeerDisconnectMonitor(fileDescriptor: descriptors[0]) { observation.signal() } )
close(descriptors[0]) close(descriptors[1]) let disconnected = await Task.detached { observation.wait(timeout: .now() + 10) }.value
#expect(disconnected) monitor.cancel() }
@Test func acceptedSocketSuppressesSIGPIPEAndClosesOnExec() throws { var descriptors: [Int32] = [-1, -1] #expect(socketpair(AF_UNIX, SOCK_STREAM, 0, &descriptors) == 0) defer { close(descriptors[0]) close(descriptors[1]) }
#expect(CLISocketServer.configureAcceptedClient(descriptors[0])) var enabled: Int32 = 0 var length = socklen_t(MemoryLayout<Int32>.size) #expect(getsockopt(descriptors[0], SOL_SOCKET, SO_NOSIGPIPE, &enabled, &length) == 0) #expect(enabled == 1) #expect(fcntl(descriptors[0], F_GETFD) & FD_CLOEXEC != 0) }
@Test func kernelReportsPeerPIDForLocalSocket() throws { var descriptors: [Int32] = [-1, -1] #expect(socketpair(AF_UNIX, SOCK_STREAM, 0, &descriptors) == 0) defer { close(descriptors[0]) close(descriptors[1]) }
#expect(CLISocketServer.peerProcessID(descriptors[0]) == getpid()) } #endif
private func temporarySocketPath(suffix: String) -> String { URL(fileURLWithPath: "/tmp", isDirectory: true) .appending(path: "prowl-cli-tests-\(UUID().uuidString.prefix(8))", directoryHint: .isDirectory) .appending( path: "prowl-\(suffix)-\(UUID().uuidString.prefix(8)).sock", directoryHint: .notDirectory ) .path(percentEncoded: false) }
nonisolated private static func send(requestData: Data, socketPath: String) throws -> Data { let socketFD = socket(AF_UNIX, SOCK_STREAM, 0) guard socketFD >= 0 else { throw CLIServiceError.socketCreationFailed } defer { close(socketFD) }
let connected = withSocketAddress(socketPath) { address in withUnsafePointer(to: address) { pointer in pointer.withMemoryRebound(to: sockaddr.self, capacity: 1) { socketPointer in connect(socketFD, socketPointer, socklen_t(MemoryLayout<sockaddr_un>.size)) } } } guard connected == 0 else { throw TestSocketClientError.connectFailed }
var requestLength = UInt32(requestData.count).bigEndian try withUnsafeBytes(of: &requestLength) { try writeAll(socketFD, buffer: $0) } try requestData.withUnsafeBytes { try writeAll(socketFD, buffer: $0) }
let lengthData = try readExact(socketFD, count: 4) let responseLength = lengthData.withUnsafeBytes { UInt32(bigEndian: $0.load(as: UInt32.self)) } return try readExact(socketFD, count: Int(responseLength)) }
nonisolated private static func sendAndCloseAfterAcceptance( requestData: Data, socketPath: String, accepted: DispatchSemaphore, closed: DispatchSemaphore ) throws { let socketFD = socket(AF_UNIX, SOCK_STREAM, 0) guard socketFD >= 0 else { throw CLIServiceError.socketCreationFailed } defer { close(socketFD) closed.signal() } let connected = withSocketAddress(socketPath) { address in withUnsafePointer(to: address) { pointer in pointer.withMemoryRebound(to: sockaddr.self, capacity: 1) { socketPointer in connect(socketFD, socketPointer, socklen_t(MemoryLayout<sockaddr_un>.size)) } } } guard connected == 0 else { throw TestSocketClientError.connectFailed }
var requestLength = UInt32(requestData.count).bigEndian try withUnsafeBytes(of: &requestLength) { try writeAll(socketFD, buffer: $0) } try requestData.withUnsafeBytes { try writeAll(socketFD, buffer: $0) } guard accepted.wait(timeout: .now() + 10) == .success else { throw CLIServiceError.readFailed } }
nonisolated private static func writeAll(_ fileDescriptor: Int32, buffer: UnsafeRawBufferPointer) throws { var offset = 0 while offset < buffer.count { let written = Darwin.write( fileDescriptor, buffer.baseAddress!.advanced(by: offset), buffer.count - offset ) guard written > 0 else { throw CLIServiceError.writeFailed } offset += written } }
nonisolated private static func readExact(_ fileDescriptor: Int32, count: Int) throws -> Data { var data = Data(count: count) var offset = 0 while offset < count { let readCount = data.withUnsafeMutableBytes { buffer in Darwin.read(fileDescriptor, buffer.baseAddress!.advanced(by: offset), count - offset) } guard readCount > 0 else { throw CLIServiceError.readFailed } offset += readCount } return data }
nonisolated private static func withSocketAddress<Result>( _ socketPath: String, _ body: (sockaddr_un) throws -> Result ) rethrows -> Result { var address = sockaddr_un() address.sun_family = sa_family_t(AF_UNIX) let pathBytes = Array(socketPath.utf8) let maxLength = MemoryLayout.size(ofValue: address.sun_path) - 1 precondition(pathBytes.count <= maxLength) withUnsafeMutableBytes(of: &address.sun_path) { buffer in for index in pathBytes.indices { buffer[index] = pathBytes[index] } buffer[pathBytes.count] = 0 } return try body(address) }
private func fileMode(at path: String) -> mode_t? { var statValue = stat() guard stat(path, &statValue) == 0 else { return nil } return statValue.st_mode & mode_t(0o777) }
private func createStaleSocket(at socketPath: String) throws { unlink(socketPath) try FileManager.default.createDirectory( atPath: (socketPath as NSString).deletingLastPathComponent, withIntermediateDirectories: true ) let socketFD = socket(AF_UNIX, SOCK_STREAM, 0) guard socketFD >= 0 else { throw CLIServiceError.socketCreationFailed } defer { close(socketFD) } try bindSocket(socketFD, to: socketPath) }
private func canConnect(to socketPath: String) -> Bool { let socketFD = socket(AF_UNIX, SOCK_STREAM, 0) guard socketFD >= 0 else { return false } defer { close(socketFD) } return withSocketAddress(socketPath) { address in withUnsafePointer(to: address) { pointer in pointer.withMemoryRebound(to: sockaddr.self, capacity: 1) { socketPointer in connect(socketFD, socketPointer, socklen_t(MemoryLayout<sockaddr_un>.size)) } } } == 0 }
private func bindSocket(_ socketFD: Int32, to socketPath: String) throws { let result = withSocketAddress(socketPath) { address in withUnsafePointer(to: address) { pointer in pointer.withMemoryRebound(to: sockaddr.self, capacity: 1) { socketPointer in bind(socketFD, socketPointer, socklen_t(MemoryLayout<sockaddr_un>.size)) } } } guard result == 0 else { throw CLIServiceError.bindFailed } }
private func isCloseOnExec(_ fileDescriptor: Int32) -> Bool { let flags = fcntl(fileDescriptor, F_GETFD) return flags >= 0 && (flags & FD_CLOEXEC) == FD_CLOEXEC }
private func withSocketAddress<Result>( _ socketPath: String, _ body: (sockaddr_un) throws -> Result ) rethrows -> Result { var address = sockaddr_un() address.sun_family = sa_family_t(AF_UNIX) let pathBytes = Array(socketPath.utf8) let maxLength = MemoryLayout.size(ofValue: address.sun_path) - 1 precondition(pathBytes.count <= maxLength) let copyLength = min(pathBytes.count, maxLength) withUnsafeMutableBytes(of: &address.sun_path) { buffer in for index in 0..<copyLength { buffer[index] = pathBytes[index] } buffer[copyLength] = 0 } return try body(address) }}
private nonisolated final class DisconnectObservation: @unchecked Sendable { private let semaphore = DispatchSemaphore(value: 0)
func signal() { semaphore.signal() }
func wait(timeout: DispatchTime) -> Bool { semaphore.wait(timeout: timeout) == .success }}
private enum TestSocketClientError: Error { case connectFailed}
private struct CancellationProbeHandler: CommandHandler { let probe: CancellationProbe
func handle(envelope: CommandEnvelope) async -> CommandResponse { await probe.suspendUntilCancelled() return CommandResponse( ok: false, command: "agents.wait", schemaVersion: "prowl.cli.agents.wait.v1", error: CommandError(code: CLIErrorCode.timeout, message: "Cancelled.") ) }}
private nonisolated final class CancellationProbe: @unchecked Sendable { private let lock = NSLock() private let cancellationObserved = DispatchSemaphore(value: 0) private var cancellationContinuation: CheckedContinuation<Void, Never>? private var cancelled = false
var wasCancelled: Bool { lock.withLock { cancelled } }
func suspendUntilCancelled() async { await withTaskCancellationHandler { await withCheckedContinuation { continuation in let resumeImmediately = lock.withLock { () -> Bool in guard !cancelled else { return true } cancellationContinuation = continuation return false } if resumeImmediately { continuation.resume() } } } onCancel: { markCancelled() } }
nonisolated func waitForCancellation(timeout: DispatchTime) -> Bool { cancellationObserved.wait(timeout: timeout) == .success }
private func markCancelled() { let continuation = lock.withLock { cancelled = true defer { cancellationContinuation = nil } return cancellationContinuation } cancellationObserved.signal() continuation?.resume() }}