Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
6842678
Add epoll constants
colemancda Aug 3, 2026
aeacfaa
Add kernel event C interop types
colemancda Aug 3, 2026
5774f91
Add epoll and kqueue syscalls
colemancda Aug 3, 2026
c02af4e
Add EventQueue protocol
colemancda Aug 3, 2026
3484947
Add PollEventQueue
colemancda Aug 3, 2026
a79a09e
Add KqueueEventQueue
colemancda Aug 3, 2026
5bc0400
Add EpollEventQueue
colemancda Aug 3, 2026
cb8810a
Use EventQueue in AsyncSocketManager
colemancda Aug 3, 2026
a177926
Add EventQueue tests
colemancda Aug 3, 2026
4548d28
Add syscall mocking driver
colemancda Aug 3, 2026
76b5dde
Fix mocking hooks in syscalls
colemancda Aug 3, 2026
6ca552a
Enable mocking for debug builds
colemancda Aug 3, 2026
d97c4ac
Add mocking tests
colemancda Aug 3, 2026
b8ded23
Discard stale socket state on file descriptor reuse
colemancda Aug 3, 2026
1abddea
Report half close as hangup in EpollEventQueue
colemancda Aug 3, 2026
2cec227
Revert "Report half close as hangup in EpollEventQueue"
colemancda Aug 3, 2026
aa482c4
Only close descriptors on behalf of their owner
colemancda Aug 3, 2026
8dcab42
Register write readiness only while a write is pending
colemancda Aug 3, 2026
0e83830
Avoid double close in EventQueue tests
colemancda Aug 3, 2026
3f21854
Update expected events for on demand write readiness
colemancda Aug 3, 2026
8bac213
Expose monitored events to tests
colemancda Aug 3, 2026
3e29526
Add idle CPU tests
colemancda Aug 3, 2026
0707518
Run tests serially in CI
colemancda Aug 3, 2026
8164074
Run tests on Linux
colemancda Aug 3, 2026
3df4606
Fix idle tests build on Glibc
colemancda Aug 3, 2026
9dd2494
Report every matrix job independently
colemancda Aug 3, 2026
e6b3125
Fix UDP test destination address
colemancda Aug 4, 2026
34e8c89
Limit job runtime
colemancda Aug 4, 2026
851af03
Bind test sockets to unused addresses
colemancda Aug 4, 2026
796efeb
Close idle test sockets before returning
colemancda Aug 4, 2026
ce8fffd
Build Android against Swift 6.3.3
colemancda Aug 4, 2026
dba199c
Listen before connecting in TCP test
colemancda Aug 4, 2026
f255b05
Close accepted socket and bound test runtime
colemancda Aug 4, 2026
b2898b4
Ignore coalesced connect notification in TCP test
colemancda Aug 4, 2026
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
12 changes: 10 additions & 2 deletions .github/workflows/swift.yml
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,9 @@ jobs:
macos:
name: macOS
runs-on: macos-15
timeout-minutes: 20
strategy:
fail-fast: false
matrix:
config: ["debug", "release"]
options: ["", "SWIFT_BUILD_DYNAMIC_LIBRARY=1"]
Expand All @@ -17,16 +19,20 @@ jobs:
- name: Build
run: ${{ matrix.options }} swift build -c ${{ matrix.config }}
- name: Test
run: ${{ matrix.options }} swift test -c ${{ matrix.config }}
# serially, sockets share the process file descriptor table and the idle
# benchmark measures CPU for the whole process
run: ${{ matrix.options }} swift test -c ${{ matrix.config }} --no-parallel

linux:
name: Linux
strategy:
fail-fast: false
matrix:
container: ["swift:6.0.3", "swift:6.1.2", "swift:6.2.3"]
config: ["debug", "release"]
options: ["", "SWIFT_BUILD_DYNAMIC_LIBRARY=1"]
runs-on: ubuntu-latest
timeout-minutes: 20
container: ${{ matrix.container }}-jammy
steps:
- name: Checkout
Expand All @@ -35,13 +41,15 @@ jobs:
run: swift --version
- name: Build
run: ${{ matrix.options }} swift build -c ${{ matrix.config }}
- name: Test
run: ${{ matrix.options }} swift test -c ${{ matrix.config }} --no-parallel

android-arm:
name: Android
strategy:
fail-fast: false
matrix:
swift: ['6.2.3', 'nightly-6.3']
swift: ['6.3.3']
arch: ['aarch64', 'x86_64', 'armv7']
runs-on: macos-15
timeout-minutes: 30
Expand Down
6 changes: 6 additions & 0 deletions Package.swift
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ var package = Package(
name: "SystemPackage",
package: "swift-system"
),
],
swiftSettings: [
.define("ENABLE_MOCKING", .when(configuration: .debug))
]
),
.target(
Expand All @@ -53,6 +56,9 @@ var package = Package(
name: "Logging",
package: "swift-log"
)
],
swiftSettings: [
.define("ENABLE_MOCKING", .when(configuration: .debug))
]
)
]
Expand Down
202 changes: 137 additions & 65 deletions Sources/Socket/SocketManager/AsyncSocketManager.swift
Original file line number Diff line number Diff line change
Expand Up @@ -42,15 +42,38 @@ extension AsyncSocketConfiguration: SocketManagerConfiguration {
}
}

#if os(Linux) || os(Android)
internal typealias PlatformEventQueue = EpollEventQueue
#elseif canImport(Darwin)
internal typealias PlatformEventQueue = KqueueEventQueue
#else
internal typealias PlatformEventQueue = PollEventQueue
#endif

/// Async Socket Manager
internal actor AsyncSocketManager: SocketManager {

// MARK: - Properties

fileprivate var state = ManagerState()


fileprivate var eventQueue: PlatformEventQueue?

/// Events every socket is registered for.
///
/// Write readiness is not included. A connected socket is almost always writable, so a
/// standing registration would make every wait return immediately for every idle socket.
/// It is added on demand while a write is pending, see ``addInterest(_:for:)``.
internal static let monitoredEvents: FileEvents = [
.read,
.readUrgent,
.error,
.hangup,
.invalidRequest
]

// MARK: - Initialization

static let shared = AsyncSocketManager()

private init() { }
Expand All @@ -60,9 +83,14 @@ internal actor AsyncSocketManager: SocketManager {
func add(
_ fileDescriptor: SocketDescriptor
) -> Socket.Event.Stream {
guard state.sockets.keys.contains(fileDescriptor) == false else {
fatalError("Another socket for file descriptor \(fileDescriptor) already exists.")
}
// The kernel only hands back a descriptor number once the previous owner is gone,
// so any existing entry belongs to a socket that was closed without notifying us.
// Kernel event queues drop closed descriptors silently, unlike `poll(2)` reporting `POLLNVAL`.
if state.sockets.keys.contains(fileDescriptor) {
log("Discard stale socket \(fileDescriptor)")
discard(fileDescriptor)
}
state.detached.remove(fileDescriptor)
log("Add socket \(fileDescriptor)")
// make sure its non blocking
do {
Expand All @@ -76,6 +104,18 @@ internal actor AsyncSocketManager: SocketManager {
log("Unable to set non blocking. \(error)")
assertionFailure("Unable to set non blocking. \(error)")
}
// register with kernel event queue
do {
if eventQueue == nil {
eventQueue = try PlatformEventQueue(maxEvents: 1024)
}
try eventQueue?.add(fileDescriptor, events: Self.monitoredEvents)
state.interests[fileDescriptor] = Self.monitoredEvents
}
catch {
log("Unable to register socket for events. \(error)")
assertionFailure("Unable to register socket for events. \(error)")
}
// append socket with events continuation
let eventStream = Socket.Event.Stream(bufferingPolicy: .bufferingNewest(1)) { continuation in
state.sockets[fileDescriptor] = SocketState(
Expand All @@ -90,12 +130,32 @@ internal actor AsyncSocketManager: SocketManager {
}

func remove(_ fileDescriptor: SocketDescriptor) {
if state.sockets[fileDescriptor] != nil {
log("Remove socket \(fileDescriptor)")
// deregister before closing, a closed descriptor cannot be deregistered by number
discard(fileDescriptor)
} else if state.detached.remove(fileDescriptor) == nil {
return // already closed by its owner
}
// close on behalf of the owner, including sockets left open by a hangup
try? fileDescriptor.close()
}

/// Deregister a file descriptor and tear down its state without closing it.
///
/// The manager never closes a descriptor it did not open. Closing one that the owner
/// still holds lets the kernel hand the same number to a new socket, and a later
/// ``remove(_:)`` would then close that unrelated socket instead. A descriptor kept
/// open this way is recorded in `detached` so its owner can still close it once.
private func discard(_ fileDescriptor: SocketDescriptor, detach: Bool = false) {
guard let socket = state.sockets[fileDescriptor] else {
return // could have been removed previously
return
}
log("Remove socket \(fileDescriptor)")
// close underlying socket
try? fileDescriptor.close()
if detach {
state.detached.insert(fileDescriptor)
}
try? eventQueue?.remove(fileDescriptor)
state.interests[fileDescriptor] = nil
// cancel all pending actions
Task(priority: .userInitiated) {
await socket.dequeueAll(Errno.connectionAbort)
Expand Down Expand Up @@ -228,7 +288,7 @@ private extension AsyncSocketManager {
// poll
let hasEvents = try poll(&tasks)
// stop monitoring if no sockets
if state.pollDescriptors.isEmpty {
if state.sockets.isEmpty {
state.isMonitoring = false
}
// wait for each task to complete
Expand All @@ -252,7 +312,28 @@ private extension AsyncSocketManager {
func contains(_ fileDescriptor: SocketDescriptor) -> Bool {
return state.sockets.keys.contains(fileDescriptor)
}


/// Subscribe to additional events for a registered socket.
func addInterest(_ events: FileEvents, for fileDescriptor: SocketDescriptor) {
guard let current = state.interests[fileDescriptor] else { return }
setInterest(current.union(events), for: fileDescriptor)
}

/// Stop monitoring events that no longer have a pending operation.
func removeInterest(_ events: FileEvents, for fileDescriptor: SocketDescriptor) {
guard let current = state.interests[fileDescriptor] else { return }
setInterest(current.subtracting(events).union(Self.monitoredEvents), for: fileDescriptor)
}

private func setInterest(_ events: FileEvents, for fileDescriptor: SocketDescriptor) {
// skip the syscall when the mask is unchanged
guard state.interests[fileDescriptor] != events else { return }
state.interests[fileDescriptor] = events
do { try eventQueue?.update(fileDescriptor, events: events) }
catch { log("Unable to update events for \(fileDescriptor). \(error)") }
}


nonisolated func wait(
for events: FileEvents,
fileDescriptor: SocketDescriptor
Expand All @@ -262,6 +343,8 @@ private extension AsyncSocketManager {
guard await socket.pendingEvents.contains(events) == false else {
return socket // execute immediately
}
// subscribe to events that are not monitored by default, like write readiness
await addInterest(events, for: fileDescriptor)
// store continuation to resume when event is polled
try await withThrowingContinuation(for: fileDescriptor) { (continuation: SocketContinuation<(), Swift.Error>) -> () in
// store pending continuation
Expand All @@ -285,78 +368,59 @@ private extension AsyncSocketManager {
/// Poll for events.
@discardableResult
func poll(_ tasks: inout [Task<Void, Never>]) throws -> Bool {
// build poll descriptor array
let sockets = state.sockets
.lazy
.sorted(by: { $0.key.rawValue < $1.key.rawValue })
state.pollDescriptors.removeAll(keepingCapacity: true)
state.pollDescriptors.reserveCapacity(sockets.count)
let events: FileEvents = [
.read,
.readUrgent,
.write,
.error,
.hangup,
.invalidRequest
]
for (fileDescriptor, _) in sockets {
let poll = SocketDescriptor.Poll(
socket: fileDescriptor,
events: events
)
state.pollDescriptors.append(poll)
}
assert(state.pollDescriptors.count == sockets.count)
// poll sockets
guard state.sockets.isEmpty == false else { return false }
var hasEvents = false
do {
try state.pollDescriptors.poll()
try eventQueue?.wait(timeout: 0) { buffer in
hasEvents = buffer.isEmpty == false
for readiness in buffer {
guard let socket = state.sockets[readiness.fileDescriptor] else {
continue // stale event for a removed descriptor
}
tasks.append(process(readiness, socket: socket))
}
}
}
catch {
log("Unable to poll for events. \(error.localizedDescription)")
throw error
}
// wait for concurrent handling
let hasEvents = state.pollDescriptors.contains(where: { $0.returnedEvents.isEmpty == false })
if hasEvents {
for poll in state.pollDescriptors {
guard let state = state.sockets[poll.socket] else {
preconditionFailure()
continue
}
let task = process(poll, socket: state)
tasks.append(task)
}
}
return hasEvents
}
func process(_ poll: SocketDescriptor.Poll, socket: AsyncSocketManager.SocketState) -> Task<Void, Never> {

func process(_ readiness: SocketReadiness, socket: AsyncSocketManager.SocketState) -> Task<Void, Never> {
Task(priority: state.configuration.monitorPriority) {
if poll.returnedEvents.contains(.read) {
if readiness.events.contains(.read) {
await socket.event(.read, notification: socket.isListening ? .connection : .read)
}
if poll.returnedEvents.contains(.write) {
if readiness.events.contains(.write) {
await socket.event(.write, notification: .write)
// unsubscribe once nothing is waiting to write, the socket stays
// writable and would otherwise report readiness on every wait
if await socket.isWaiting(for: .write) == false {
removeInterest(.write, for: readiness.fileDescriptor)
}
}
if poll.returnedEvents.contains(.invalidRequest) {
error(.badFileDescriptor, for: poll.socket)
if readiness.events.contains(.invalidRequest) {
error(.badFileDescriptor, for: readiness.fileDescriptor)
}
if poll.returnedEvents.contains(.error) {
error(.connectionReset, for: poll.socket)
if readiness.events.contains(.error) {
error(.connectionReset, for: readiness.fileDescriptor)
}
if poll.returnedEvents.contains(.hangup) {
hangup(poll.socket)
if readiness.events.contains(.hangup) {
hangup(readiness.fileDescriptor)
}
}
}

func error(_ error: Errno, for fileDescriptor: SocketDescriptor) {
state.sockets[fileDescriptor]?.continuation.yield(.error(error))
remove(fileDescriptor)
// stop monitoring but leave the descriptor open, see `discard(_:detach:)`
discard(fileDescriptor, detach: true)
}

func hangup(_ fileDescriptor: SocketDescriptor) {
remove(fileDescriptor)
discard(fileDescriptor, detach: true)
}
}

Expand Down Expand Up @@ -463,6 +527,10 @@ fileprivate extension AsyncSocketManager.SocketState {
}
}

func isWaiting(for event: FileEvents) -> Bool {
eventContinuation[event, default: []].isEmpty == false
}

func queue(_ event: FileEvents, _ continuation: SocketContinuation<(), Error>) {
guard pendingEvents.contains(event) == false else {
continuation.resume()
Expand Down Expand Up @@ -519,9 +587,13 @@ extension AsyncSocketManager {
var configuration = AsyncSocketConfiguration()

var sockets = [SocketDescriptor: SocketState]()

var pollDescriptors = [SocketDescriptor.Poll]()


/// Events each socket is currently registered for.
var interests = [SocketDescriptor: FileEvents]()

/// Sockets no longer monitored whose descriptor is still open.
var detached = Set<SocketDescriptor>()

var isMonitoring = false
}

Expand Down
Loading
Loading