Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Package.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -478,6 +478,12 @@ let package = Package(
"CAuditToken",
]
),
.testTarget(
name: "ContainerXPCTests",
dependencies: [
"ContainerXPC"
]
),
.target(
name: "ContainerOS",
dependencies: [
Expand Down
17 changes: 10 additions & 7 deletions Sources/ContainerXPC/XPCClient.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -111,15 +111,18 @@ extension XPCClient {
}

group.addTask {
try await withCheckedThrowingContinuation { cont in
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
do {
let message = try self.parseReply(reply)
cont.resume(returning: message)
} catch {
cont.resume(throwing: error)
let box = XPCReplyBox()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { cont in
box.store(cont)
xpc_connection_send_message_with_reply(self.connection, message.underlying, nil) { reply in
box.resume { try self.parseReply(reply) }
}
}
} onCancel: {
// XPC reply callbacks do not observe task cancellation, so
// resume the awaiting continuation when a timeout wins.
box.resume { throw CancellationError() }
}
}

Expand Down
54 changes: 54 additions & 0 deletions Sources/ContainerXPC/XPCReplyBox.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import Foundation

/// Resumes an XPC reply continuation once when either a reply or cancellation arrives.
final class XPCReplyBox: @unchecked Sendable {
private let lock = NSLock()
private var continuation: CheckedContinuation<XPCMessage, Error>?
private var result: Result<XPCMessage, Error>?

func store(_ continuation: CheckedContinuation<XPCMessage, Error>) {
lock.lock()
guard let result else {
self.continuation = continuation
lock.unlock()
return
}
lock.unlock()
continuation.resume(with: result)
}

func resume(_ body: () throws -> XPCMessage) {
let result = Result { try body() }
lock.lock()
guard self.result == nil else {
lock.unlock()
return
}
self.result = result
guard let continuation else {
lock.unlock()
return
}
self.continuation = nil
lock.unlock()
continuation.resume(with: result)
}
}
#endif
230 changes: 230 additions & 0 deletions Tests/ContainerXPCTests/XPCClientTests.swift
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,230 @@
//===----------------------------------------------------------------------===//
// Copyright © 2026 Apple Inc. and the container project authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//===----------------------------------------------------------------------===//

#if os(macOS)
import ContainerizationError
@preconcurrency import Foundation
import Testing

@testable import ContainerXPC

@Suite(.timeLimit(.minutes(1)), .serialized)
struct XPCClientTests {
@Test
func replyBoxDeliversCancellationThatPrecedesContinuationStorage() async throws {
let box = XPCReplyBox()
box.resume { throw CancellationError() }

do {
_ = try await withCheckedThrowingContinuation { continuation in
box.store(continuation)
}
Issue.record("expected the pending cancellation to be delivered")
} catch is CancellationError {
// Expected.
}
}

@Test
func responseTimeoutReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let clock = ContinuousClock()
let start = clock.now

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout for request to test.container.xpc/hang"))
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func callerCancellationReturnsWithinBound() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()
let request = Task {
try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .seconds(30)
)
}

try await Task.sleep(for: .milliseconds(50))
let clock = ContinuousClock()
let start = clock.now
request.cancel()

do {
_ = try await request.value
Issue.record("expected send to be cancelled")
} catch is CancellationError {
#expect(start.duration(to: clock.now) < .seconds(2))
}
}

@Test
func clientCanBeReusedAfterResponseTimeout() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}

@Test
func lateReplyAfterTimeoutIsIgnored() async throws {
let server = AnonymousXPCServer()
defer { server.close() }

let client = server.makeClient()

do {
_ = try await client.send(
XPCMessage(route: "hang"),
responseTimeout: .milliseconds(100)
)
Issue.record("expected send to time out")
} catch let error as ContainerizationError {
#expect(error.message.contains("XPC timeout"))
}

#expect(server.replyToPendingRequests())
try await Task.sleep(for: .milliseconds(50))

let response = try await client.send(
XPCMessage(route: "echo"),
responseTimeout: .seconds(1)
)
#expect(response.string(key: "result") == "ok")
}
}

private final class AnonymousXPCServer: @unchecked Sendable {
private let listener: xpc_connection_t
private let lock = NSLock()
private var connections = [xpc_connection_t]()
private var pendingRequests = [xpc_object_t]()

init() {
listener = xpc_connection_create(nil, nil)
xpc_connection_set_event_handler(listener) { [weak self] object in
switch xpc_get_type(object) {
case XPC_TYPE_CONNECTION:
self?.accept(connection: object)
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(listener)
}

func makeClient() -> XPCClient {
let endpoint = xpc_endpoint_create(listener)
let connection = xpc_connection_create_from_endpoint(endpoint)
return XPCClient(connection: connection, label: "test.container.xpc")
}

func close() {
xpc_connection_cancel(listener)
lock.withLock {
for connection in connections {
xpc_connection_cancel(connection)
}
connections.removeAll()
pendingRequests.removeAll()
}
}

func replyToPendingRequests() -> Bool {
let requests = lock.withLock {
let requests = pendingRequests
pendingRequests.removeAll()
return requests
}
guard let connection = lock.withLock({ connections.last }) else {
return false
}
for request in requests {
guard let reply = xpc_dictionary_create_reply(request) else {
continue
}
xpc_dictionary_set_string(reply, "result", "late")
xpc_connection_send_message(connection, reply)
}
return !requests.isEmpty
}

private func accept(connection: xpc_connection_t) {
lock.withLock {
connections.append(connection)
}
nonisolated(unsafe) let connection = connection
xpc_connection_set_event_handler(connection) { [weak self] object in
guard let self else {
return
}
switch xpc_get_type(object) {
case XPC_TYPE_DICTIONARY:
let message = XPCMessage(object: object)
switch message.string(key: XPCMessage.routeKey) {
case "echo":
guard let reply = xpc_dictionary_create_reply(object) else {
return
}
xpc_dictionary_set_string(reply, "result", "ok")
xpc_connection_send_message(connection, reply)
default:
self.lock.withLock {
self.pendingRequests.append(object)
}
}
case XPC_TYPE_ERROR:
break
default:
fatalError("unhandled xpc object type: \(xpc_get_type(object))")
}
}
xpc_connection_activate(connection)
}
}
#endif