Skip to content
Draft
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
56 changes: 35 additions & 21 deletions MobiusCore/Source/EffectHandlers/EffectExecutor.swift
Original file line number Diff line number Diff line change
@@ -1,8 +1,6 @@
// Copyright Spotify AB.
// SPDX-License-Identifier: Apache-2.0

import Foundation

final class EffectExecutor<Effect, Event>: Connectable {
private let handleEffect: (Effect, EffectCallback<Event>) -> Disposable
private var output: Consumer<Event>?
Expand All @@ -20,22 +18,29 @@ final class EffectExecutor<Effect, Event>: Connectable {
}

func connect(_ consumer: @escaping Consumer<Event>) -> Connection<Effect> {
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) {
Expand All @@ -47,7 +52,7 @@ final class EffectExecutor<Effect, Event>: 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) }
)
Expand All @@ -63,18 +68,27 @@ final class EffectExecutor<Effect, Event>: 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)
}
}

Expand Down
19 changes: 13 additions & 6 deletions MobiusCore/Source/EffectHandlers/EffectRouter.swift
Original file line number Diff line number Diff line change
Expand Up @@ -148,19 +148,26 @@ private func compose<Input, Output>(

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
Expand Down
37 changes: 23 additions & 14 deletions MobiusCore/Source/EffectHandlers/ThreadSafeConnectable.swift
Original file line number Diff line number Diff line change
Expand Up @@ -14,24 +14,33 @@ final class ThreadSafeConnectable<Event, Effect>: Connectable {
self.connectable = AnyConnectable(connectable)
}

func connect(_ output: @escaping (Event) -> Void) -> Connection<Effect> {
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<Event>) -> Connection<Effect> {
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) {
Expand Down
33 changes: 21 additions & 12 deletions MobiusCore/Source/Lock.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Result>(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<Value> {
private let lock = DispatchQueue(label: "Mobius synchronized storage")
private let lock = Lock()
private var storage: Value

init(value: Value) {
Expand All @@ -26,21 +35,21 @@ final class Synchronized<Value> {

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)
}
}
Expand Down
48 changes: 47 additions & 1 deletion MobiusCore/Test/EffectHandlers/EffectRouterTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Effect, Event>()
.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<Effect, Event>()
.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")
Expand Down Expand Up @@ -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<Effect> {
Connection(
for event in eventsOnConnect {
consumer(event)
}
return Connection(
acceptClosure: { _ in
self.onEvent()
consumer(self.event)
Expand Down
Loading