Skip to content
Merged
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
31 changes: 31 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -92,17 +92,48 @@ jobs:
CODE_SIGNING_REQUIRED=NO \
CODE_SIGN_IDENTITY=""

# `-resultBundlePath` so a failure names the test (GitHub #82). Two
# App-target runs have reported "1 failure" with no name and neither
# reproduced; a count is not a bug report.
- name: Run app-target tests (xcodebuild)
run: |
xcodebuild test \
-project InterlinedList.xcodeproj \
-scheme InterlinedList \
-destination "platform=macOS" \
-derivedDataPath build/DerivedData \
-resultBundlePath build/test-results/app.xcresult \
CODE_SIGNING_ALLOWED=NO \
CODE_SIGNING_REQUIRED=NO \
CODE_SIGN_IDENTITY=""

- name: Name the failing tests
if: failure()
run: |
xcrun xcresulttool get test-results tests \
--path build/test-results/app.xcresult --format json \
| python3 -c '
import json, sys
def walk(node):
for child in node.get("children") or []:
yield from walk(child)
if node.get("nodeType") == "Test Case" and node.get("result") == "Failed":
yield node.get("nodeIdentifier") or node.get("name")
doc = json.load(sys.stdin)
names = [n for root in doc.get("testNodes", []) for n in walk(root)]
print("::error::Failing tests: " + (", ".join(names) if names else "none listed"))
for n in names:
print("FAILED:", n)
'

- name: Upload the App-target result bundle
if: failure()
uses: actions/upload-artifact@v7
with:
name: app-test-results-${{ github.sha }}
path: build/test-results/app.xcresult
retention-days: 14

- name: Locate built .app
id: locate
run: |
Expand Down
47 changes: 13 additions & 34 deletions App/Composition/ComposerEventBus.swift
Original file line number Diff line number Diff line change
Expand Up @@ -50,48 +50,27 @@ enum ComposerEvent: Sendable, Equatable {
/// terminate the stream by cancelling the consuming task.
final class ComposerEventBus: Sendable {

private let storage = Storage()
/// Subscriber registry. Shared with the other three feature buses; see
/// `EventBusStorage` for why registration is synchronous (GitHub #82).
private let storage = EventBusStorage<ComposerEvent>()

init() {}

/// Returns an `AsyncStream` that yields every event posted after
/// subscription. The stream finishes when the consumer cancels.
/// subscription. The subscriber is registered before this returns, so an
/// immediately-following `post` is delivered. The stream finishes when the
/// consuming task is cancelled.
func events() -> AsyncStream<ComposerEvent> {
let id = UUID()
return AsyncStream { continuation in
Task { await self.storage.register(id: id, continuation: continuation) }
continuation.onTermination = { _ in
Task { await self.storage.unregister(id: id) }
}
}
storage.stream()
}

/// Publish an event to every active subscriber. Late subscribers
/// do not receive past events.
/// Publish an event to every active subscriber. Late subscribers do not
/// receive past events. Delivery is synchronous with the call.
func post(_ event: ComposerEvent) {
Task { await storage.broadcast(event) }
storage.broadcast(event)
}

// MARK: - Storage

/// Holds the live continuations keyed by registration UUID. An
/// actor because subscribers / publishers are not serialized to
/// any thread.
private actor Storage {
private var continuations: [UUID: AsyncStream<ComposerEvent>.Continuation] = [:]

func register(id: UUID, continuation: AsyncStream<ComposerEvent>.Continuation) {
continuations[id] = continuation
}

func unregister(id: UUID) {
continuations[id] = nil
}

func broadcast(_ event: ComposerEvent) {
for continuation in continuations.values {
continuation.yield(event)
}
}
}
/// Live subscriber count, for tests that need to assert a subscription
/// exists rather than wait for one.
var subscriberCount: Int { storage.subscriberCount }
}
46 changes: 13 additions & 33 deletions App/Composition/DirectMessagesEventBus.swift
Original file line number Diff line number Diff line change
Expand Up @@ -42,47 +42,27 @@ enum DirectMessagesEvent: Sendable, Equatable {
/// a subscription stream; terminate by cancelling the consuming task.
final class DirectMessagesEventBus: Sendable {

private let storage = Storage()
/// Subscriber registry. Shared with the other three feature buses; see
/// `EventBusStorage` for why registration is synchronous (GitHub #82).
private let storage = EventBusStorage<DirectMessagesEvent>()

init() {}

/// Returns an `AsyncStream` that yields every event posted after
/// subscription. The stream finishes when the consumer cancels.
/// subscription. The subscriber is registered before this returns, so an
/// immediately-following `post` is delivered. The stream finishes when the
/// consuming task is cancelled.
func events() -> AsyncStream<DirectMessagesEvent> {
let id = UUID()
return AsyncStream { continuation in
Task { await self.storage.register(id: id, continuation: continuation) }
continuation.onTermination = { _ in
Task { await self.storage.unregister(id: id) }
}
}
storage.stream()
}

/// Publish an event to every active subscriber. Late subscribers do
/// not receive past events.
/// Publish an event to every active subscriber. Late subscribers do not
/// receive past events. Delivery is synchronous with the call.
func post(_ event: DirectMessagesEvent) {
Task { await storage.broadcast(event) }
storage.broadcast(event)
}

// MARK: - Storage

/// Holds the live continuations keyed by registration UUID. An actor
/// because publishers and subscribers aren't serialized.
private actor Storage {
private var continuations: [UUID: AsyncStream<DirectMessagesEvent>.Continuation] = [:]

func register(id: UUID, continuation: AsyncStream<DirectMessagesEvent>.Continuation) {
continuations[id] = continuation
}

func unregister(id: UUID) {
continuations[id] = nil
}

func broadcast(_ event: DirectMessagesEvent) {
for continuation in continuations.values {
continuation.yield(event)
}
}
}
/// Live subscriber count, for tests that need to assert a subscription
/// exists rather than wait for one.
var subscriberCount: Int { storage.subscriberCount }
}
75 changes: 75 additions & 0 deletions App/Composition/EventBusStorage.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
// EventBusStorage
//
// The shared subscriber registry behind the four feature event buses
// (`ListsEventBus`, `ComposerEventBus`, `NotificationsEventBus`,
// `DirectMessagesEventBus`). Each of those was carrying its own private copy of
// the same actor; this is that copy, written once and made synchronous.
//
// Why synchronous registration matters (GitHub #82). The previous actor-backed
// version registered the continuation *inside a `Task`*:
//
// return AsyncStream { continuation in
// Task { await self.storage.register(id: id, continuation: continuation) }
// }
//
// so `events()` returned a stream that was not yet in the subscriber table. A
// caller that subscribed and then immediately performed a write could miss its
// own event, and no number of `Task.yield()`s closed the window — registration
// was waiting on actor scheduling, not on cooperative yields. The tests papered
// over it with a fixed 10 ms sleep ("give the subscription a beat to register"),
// which is exactly the pattern that passes on an idle machine and loses the race
// when the suite competes with a cold build. In the running app the same window
// showed up as an unread badge that occasionally did not move.
//
// Registering under a lock inside the `AsyncStream` build closure — which
// `AsyncStream` invokes synchronously during init — closes it: by the time
// `events()` returns, the subscriber is live.
//
// `Mutex` rather than an actor because every operation here is a short,
// non-suspending dictionary mutation and `yield` never blocks (`AsyncStream`'s
// default buffering policy is unbounded). An actor buys serialization this does
// not need and costs the synchrony this does.
//
// Per Decision 0003 this file lives in `App/Composition/` and imports no kit.

import Foundation
import Synchronization

final class EventBusStorage<Event: Sendable>: Sendable {

private let continuations = Mutex<[UUID: AsyncStream<Event>.Continuation]>([:])

init() {}

/// A stream that is **already registered** by the time it is returned.
///
/// The stream finishes when the consuming task is cancelled; termination
/// unregisters the continuation so a dropped subscriber does not leak.
func stream() -> AsyncStream<Event> {
let id = UUID()
return AsyncStream { continuation in
continuations.withLock { $0[id] = continuation }
continuation.onTermination = { [weak self] _ in
self?.continuations.withLock { $0[id] = nil }
}
}
}

/// Delivers `event` to every subscriber registered at the moment of the
/// call. Late subscribers do not receive past events.
///
/// The values are copied out under the lock and yielded outside it, so a
/// subscriber that reacts by subscribing or unsubscribing cannot deadlock.
func broadcast(_ event: Event) {
let live = continuations.withLock { Array($0.values) }
for continuation in live {
continuation.yield(event)
}
}

/// Live subscriber count. Exists so a test can assert that subscription
/// happened without waiting on a clock.
var subscriberCount: Int {
continuations.withLock { $0.count }
}
}
46 changes: 13 additions & 33 deletions App/Composition/ListsEventBus.swift
Original file line number Diff line number Diff line change
Expand Up @@ -63,47 +63,27 @@ enum ListsEvent: Sendable, Equatable {
/// subscription stream; terminate by cancelling the consuming task.
final class ListsEventBus: Sendable {

private let storage = Storage()
/// Subscriber registry. Shared with the other three feature buses; see
/// `EventBusStorage` for why registration is synchronous (GitHub #82).
private let storage = EventBusStorage<ListsEvent>()

init() {}

/// Returns an `AsyncStream` that yields every event posted after
/// subscription. The stream finishes when the consumer cancels.
/// subscription. The subscriber is registered before this returns, so an
/// immediately-following `post` is delivered. The stream finishes when the
/// consuming task is cancelled.
func events() -> AsyncStream<ListsEvent> {
let id = UUID()
return AsyncStream { continuation in
Task { await self.storage.register(id: id, continuation: continuation) }
continuation.onTermination = { _ in
Task { await self.storage.unregister(id: id) }
}
}
storage.stream()
}

/// Publish an event to every active subscriber. Late subscribers
/// do not receive past events.
/// Publish an event to every active subscriber. Late subscribers do not
/// receive past events. Delivery is synchronous with the call.
func post(_ event: ListsEvent) {
Task { await storage.broadcast(event) }
storage.broadcast(event)
}

// MARK: - Storage

/// Holds the live continuations keyed by registration UUID. An
/// actor because publishers and subscribers aren't serialized.
private actor Storage {
private var continuations: [UUID: AsyncStream<ListsEvent>.Continuation] = [:]

func register(id: UUID, continuation: AsyncStream<ListsEvent>.Continuation) {
continuations[id] = continuation
}

func unregister(id: UUID) {
continuations[id] = nil
}

func broadcast(_ event: ListsEvent) {
for continuation in continuations.values {
continuation.yield(event)
}
}
}
/// Live subscriber count, for tests that need to assert a subscription
/// exists rather than wait for one.
var subscriberCount: Int { storage.subscriberCount }
}
46 changes: 13 additions & 33 deletions App/Composition/NotificationsEventBus.swift
Original file line number Diff line number Diff line change
Expand Up @@ -50,47 +50,27 @@ enum NotificationsEvent: Sendable, Equatable {
/// cancelling the consuming task.
final class NotificationsEventBus: Sendable {

private let storage = Storage()
/// Subscriber registry. Shared with the other three feature buses; see
/// `EventBusStorage` for why registration is synchronous (GitHub #82).
private let storage = EventBusStorage<NotificationsEvent>()

init() {}

/// Returns an `AsyncStream` that yields every event posted after
/// subscription. The stream finishes when the consumer cancels.
/// subscription. The subscriber is registered before this returns, so an
/// immediately-following `post` is delivered. The stream finishes when the
/// consuming task is cancelled.
func events() -> AsyncStream<NotificationsEvent> {
let id = UUID()
return AsyncStream { continuation in
Task { await self.storage.register(id: id, continuation: continuation) }
continuation.onTermination = { _ in
Task { await self.storage.unregister(id: id) }
}
}
storage.stream()
}

/// Publish an event to every active subscriber. Late subscribers
/// do not receive past events.
/// Publish an event to every active subscriber. Late subscribers do not
/// receive past events. Delivery is synchronous with the call.
func post(_ event: NotificationsEvent) {
Task { await storage.broadcast(event) }
storage.broadcast(event)
}

// MARK: - Storage

/// Holds the live continuations keyed by registration UUID. An
/// actor because publishers and subscribers aren't serialized.
private actor Storage {
private var continuations: [UUID: AsyncStream<NotificationsEvent>.Continuation] = [:]

func register(id: UUID, continuation: AsyncStream<NotificationsEvent>.Continuation) {
continuations[id] = continuation
}

func unregister(id: UUID) {
continuations[id] = nil
}

func broadcast(_ event: NotificationsEvent) {
for continuation in continuations.values {
continuation.yield(event)
}
}
}
/// Live subscriber count, for tests that need to assert a subscription
/// exists rather than wait for one.
var subscriberCount: Int { storage.subscriberCount }
}
13 changes: 4 additions & 9 deletions AppTests/CurrentUserStoreTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -78,18 +78,13 @@ final class CurrentUserStoreTests: XCTestCase {

// MARK: - Helpers

/// Polls `condition` until it returns `true` or 2 s elapses. Used
/// when the assertion depends on a value arriving asynchronously
/// from a stream we don't directly own a continuation on.
/// Thin shim onto the shared `settle(until:)` helper (`Support/AsyncSettle.swift`).
/// This file's private polling loop was the pattern the rest of the suite
/// should have been using all along; it now lives in one place (GitHub #82).
private func waitFor(
condition: @MainActor () -> Bool,
timeout: Double = 2.0
) async throws {
let deadline = Date().addingTimeInterval(timeout)
while Date() < deadline {
if condition() { return }
try await Task.sleep(nanoseconds: 10_000_000) // 10ms
}
XCTFail("Condition did not become true within \(timeout)s")
await settle(until: condition, timeout: .seconds(timeout))
}
}
Loading