From d1fff208764174d4b209a47e7e3f93b946236805 Mon Sep 17 00:00:00 2001 From: kmcbride Date: Wed, 15 Apr 2026 12:03:49 -0700 Subject: [PATCH 1/2] Refactor lock usage --- .../EffectHandlers/EffectExecutor.swift | 56 ++++++++++++------- .../ThreadSafeConnectable.swift | 37 +++++++----- MobiusCore/Source/Lock.swift | 33 +++++++---- .../EffectHandlers/EffectRouterTests.swift | 48 +++++++++++++++- 4 files changed, 126 insertions(+), 48 deletions(-) diff --git a/MobiusCore/Source/EffectHandlers/EffectExecutor.swift b/MobiusCore/Source/EffectHandlers/EffectExecutor.swift index 8aed9a1b..f627fd0a 100644 --- a/MobiusCore/Source/EffectHandlers/EffectExecutor.swift +++ b/MobiusCore/Source/EffectHandlers/EffectExecutor.swift @@ -1,8 +1,6 @@ // Copyright Spotify AB. // SPDX-License-Identifier: Apache-2.0 -import Foundation - final class EffectExecutor: Connectable { private let handleEffect: (Effect, EffectCallback) -> Disposable private var output: Consumer? @@ -20,22 +18,29 @@ final class EffectExecutor: Connectable { } func connect(_ consumer: @escaping Consumer) -> Connection { - return lock.synchronized { + let needsConnection = lock.synchronized { guard output == nil else { - MobiusHooks.errorHandler( - "Connection limit exceeded: The Connectable \(type(of: self)) is already connected. " + - "Unable to connect more than once", - #file, - #line - ) + return false } output = consumer - return Connection( - acceptClosure: handle, - disposeClosure: dispose + + return true + } + + guard needsConnection else { + MobiusHooks.errorHandler( + "Connection limit exceeded: The Connectable \(type(of: self)) is already connected. " + + "Unable to connect more than once", + #file, + #line ) } + + return Connection( + acceptClosure: handle, + disposeClosure: dispose + ) } func handle(_ effect: Effect) { @@ -47,7 +52,7 @@ final class EffectExecutor: Connectable { let callback = EffectCallback( // Any events produced as a result of handling the effect will be sent to this class's `output` consumer, // unless it has already been disposed. - onSend: { [weak self] event in self?.output?(event) }, + onSend: { [weak self] event in self?.dispatch(event: event) }, // Once an effect has been handled, remove the reference to its callback and disposable. onEnd: { [weak self] in self?.delete(id: id) } ) @@ -63,18 +68,27 @@ final class EffectExecutor: Connectable { } func dispose() { - lock.synchronized { - // Dispose any effects currently being handled. We also need to `end` their callbacks to remove the - // references we are keeping to them. - handlingEffects.values - .forEach { - $0.disposable.dispose() - $0.callback.end() - } + let handlingEffectStates = lock.synchronized { + let states = handlingEffects.values // Restore the state of this `Connectable` to its pre-connected state. handlingEffects = [:] output = nil + + return states + } + + // Dispose any effects currently being handled. We also need to `end` their callbacks to remove the + // references we are keeping to them. + handlingEffectStates.forEach { + $0.disposable.dispose() + $0.callback.end() + } + } + + private func dispatch(event: Event) { + if let output = lock.synchronized(closure: { output }) { + output(event) } } diff --git a/MobiusCore/Source/EffectHandlers/ThreadSafeConnectable.swift b/MobiusCore/Source/EffectHandlers/ThreadSafeConnectable.swift index 863f69ae..159673a6 100644 --- a/MobiusCore/Source/EffectHandlers/ThreadSafeConnectable.swift +++ b/MobiusCore/Source/EffectHandlers/ThreadSafeConnectable.swift @@ -14,24 +14,33 @@ final class ThreadSafeConnectable: Connectable { self.connectable = AnyConnectable(connectable) } - func connect(_ output: @escaping (Event) -> Void) -> Connection { - return lock.synchronized { - guard self.output == nil, connection == nil else { - MobiusHooks.errorHandler( - "Connection limit exceeded: The Connectable \(type(of: self)) is already connected. " + - "Unable to connect more than once", - #file, - #line - ) + func connect(_ consumer: @escaping Consumer) -> Connection { + let needsConnection = lock.synchronized { + guard output == nil else { + return false } - self.output = output - connection = connectable.connect(self.dispatch) - return Connection( - acceptClosure: accept, - disposeClosure: dispose + output = consumer + + return true + } + + guard needsConnection else { + MobiusHooks.errorHandler( + "Connection limit exceeded: The Connectable \(type(of: self)) is already connected. " + + "Unable to connect more than once", + #file, + #line ) } + + let innerConnection = connectable.connect(dispatch) + lock.synchronized { connection = innerConnection } + + return Connection( + acceptClosure: accept, + disposeClosure: dispose + ) } private func accept(_ effect: Effect) { diff --git a/MobiusCore/Source/Lock.swift b/MobiusCore/Source/Lock.swift index 4aae2620..04bd0945 100644 --- a/MobiusCore/Source/Lock.swift +++ b/MobiusCore/Source/Lock.swift @@ -3,21 +3,30 @@ import Foundation -struct Lock { - private let lock = NSRecursiveLock() +final class Lock { + private let lock: os_unfair_lock_t + init() { + lock = .allocate(capacity: 1) + lock.initialize(to: os_unfair_lock()) + } + + deinit { + lock.deinitialize(count: 1) + lock.deallocate() + } + + @discardableResult func synchronized(closure: () throws -> Result) rethrows -> Result { - lock.lock() - defer { - lock.unlock() - } + os_unfair_lock_lock(lock) + defer { os_unfair_lock_unlock(lock) } return try closure() } } final class Synchronized { - private let lock = DispatchQueue(label: "Mobius synchronized storage") + private let lock = Lock() private var storage: Value init(value: Value) { @@ -26,21 +35,21 @@ final class Synchronized { var value: Value { get { - return lock.sync { storage } + lock.synchronized { storage } } - set(newValue) { - lock.sync { self.storage = newValue } + set { + lock.synchronized { storage = newValue } } } func mutate(with closure: (inout Value) throws -> Void) rethrows { - try lock.sync { + try lock.synchronized { try closure(&storage) } } func read(in closure: (Value) throws -> Void) rethrows { - try lock.sync { + try lock.synchronized { try closure(storage) } } diff --git a/MobiusCore/Test/EffectHandlers/EffectRouterTests.swift b/MobiusCore/Test/EffectHandlers/EffectRouterTests.swift index 366b0203..3d92e7c6 100644 --- a/MobiusCore/Test/EffectHandlers/EffectRouterTests.swift +++ b/MobiusCore/Test/EffectHandlers/EffectRouterTests.swift @@ -187,6 +187,46 @@ class EffectRouterTests: QuickSpec { } } + context("Synchronous event emission during connect") { + it("should not deadlock when inner connectable emits synchronously during connect") { + var receivedEvents: [Event] = [] + + let connectable = TestConnectable( + dispatchEventsOnConnect: [.eventForEffect1] + ) + + let connection = EffectRouter() + .routeEffects(equalTo: .effect1).to(connectable) + .asConnectable + .connect { event in + receivedEvents.append(event) + } + + expect(receivedEvents).to(equal([.eventForEffect1])) + + connection.dispose() + } + + it("should forward multiple synchronous events emitted during connect") { + var receivedEvents: [Event] = [] + + let connectable = TestConnectable( + dispatchEventsOnConnect: [.eventForEffect1, .eventForEffect2] + ) + + let connection = EffectRouter() + .routeEffects(equalTo: .effect1).to(connectable) + .asConnectable + .connect { event in + receivedEvents.append(event) + } + + expect(receivedEvents).to(equal([.eventForEffect1, .eventForEffect2])) + + connection.dispose() + } + } + context("Running on different queues") { it("supports handling effects on a specified queue") { let testQueue = DispatchQueue(label: "test") @@ -215,21 +255,27 @@ class EffectRouterTests: QuickSpec { private class TestConnectable: Connectable { private let event: Event + private let eventsOnConnect: [Event] private let onDispose: () -> Void private let onEvent: () -> Void init( dispatchEvent event: Event = .eventForEffect1, + dispatchEventsOnConnect eventsOnConnect: [Event] = [], onDispose: @escaping () -> Void = {}, onEvent: @escaping () -> Void = {} ) { self.event = event + self.eventsOnConnect = eventsOnConnect self.onDispose = onDispose self.onEvent = onEvent } func connect(_ consumer: @escaping (Event) -> Void) -> Connection { - Connection( + for event in eventsOnConnect { + consumer(event) + } + return Connection( acceptClosure: { _ in self.onEvent() consumer(self.event) From 51acc63817e84b23de29ef30311ad2ac52bf9bfa Mon Sep 17 00:00:00 2001 From: kmcbride Date: Wed, 15 Apr 2026 12:21:29 -0700 Subject: [PATCH 2/2] Replace intermediate route handler array --- .../Source/EffectHandlers/EffectRouter.swift | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/MobiusCore/Source/EffectHandlers/EffectRouter.swift b/MobiusCore/Source/EffectHandlers/EffectRouter.swift index c22aa960..01a0bac2 100644 --- a/MobiusCore/Source/EffectHandlers/EffectRouter.swift +++ b/MobiusCore/Source/EffectHandlers/EffectRouter.swift @@ -148,19 +148,26 @@ private func compose( return Connection( acceptClosure: { effect in - let handlers = connectedRoutes - .compactMap { route in route.tryToHandle(effect) } + var handler: (() -> Void)? + var handlerCount = 0 - if let handleEffect = handlers.first, handlers.count == 1 { - handleEffect() - } else { + for route in connectedRoutes { + if let routeHandler = route.tryToHandle(effect) { + handler = routeHandler + handlerCount += 1 + } + } + + guard let handler, handlerCount == 1 else { MobiusHooks.errorHandler( - "Error: \(handlers.count) EffectHandlers could be found for effect: \(effect). " + + "Error: \(handlerCount) EffectHandlers could be found for effect: \(effect). " + "Exactly 1 is required.", #file, #line ) } + + handler() }, disposeClosure: { connectedRoutes