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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
**Unreleased:**

- Subjects: fix a deadlock when sending values or termination concurrently with consumer cancellation (https://github.com/sideeffect-io/AsyncExtensions/issues/52).

**v0.5.2 - Oxygen:**

This version is a bug fix version.
Expand Down
19 changes: 11 additions & 8 deletions Sources/AsyncSubjects/AsyncCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -67,24 +67,27 @@ public final class AsyncCurrentValueSubject<Element>: AsyncSubject where Element
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.current = element
for channel in state.channels.values {
channel.send(element)
}
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the async sequences with a normal ending.
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()
for channel in channels {
channel.finish()
}
return channels
}
for channel in channels {
channel.finish()
}
}

Expand Down
19 changes: 11 additions & 8 deletions Sources/AsyncSubjects/AsyncPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -52,23 +52,26 @@ public final class AsyncPassthroughSubject<Element: Sendable>: AsyncSubject {
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
for channel in state.channels.values {
channel.send(element)
}
let channels = self.state.withCriticalRegion { state in
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the subject with a normal ending.
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()
for channel in channels {
channel.finish()
}
return channels
}
for channel in channels {
channel.finish()
}
}

Expand Down
19 changes: 11 additions & 8 deletions Sources/AsyncSubjects/AsyncReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -46,29 +46,32 @@ public final class AsyncReplaySubject<Element>: AsyncSubject where Element: Send
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
if state.buffer.count >= state.bufferSize && !state.buffer.isEmpty {
state.buffer.removeFirst()
}
state.buffer.append(element)
for channel in state.channels.values {
channel.send(element)
}
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the subject with a normal ending.
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()
state.buffer.removeAll()
state.bufferSize = 0
for channel in channels {
channel.finish()
}
return channels
}
for channel in channels {
channel.finish()
}
}

Expand Down
27 changes: 15 additions & 12 deletions Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -67,28 +67,31 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: As
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.current = element
for channel in state.channels.values {
channel.send(element)
}
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the subject with either a normal ending or an error.
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()
for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
return channels
}
for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
}
}
Expand Down
28 changes: 15 additions & 13 deletions Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -53,28 +53,30 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: Asy
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
for channel in state.channels.values {
channel.send(element)
}
let channels = self.state.withCriticalRegion { state in
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the subject with either a normal ending or an error.
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()

for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
return channels
}
for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
}
}
Expand Down
27 changes: 15 additions & 12 deletions Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -45,33 +45,36 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: AsyncSub
/// Sends a value to all consumers
/// - Parameter element: the value to send
public func send(_ element: Element) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
if state.buffer.count >= state.bufferSize && !state.buffer.isEmpty {
state.buffer.removeFirst()
}
state.buffer.append(element)
for channel in state.channels.values {
channel.send(element)
}
return Array(state.channels.values)
}
// Resuming a consumer must not hold the lock used by its cancellation handler.
for channel in channels {
channel.send(element)
}
}

/// Finishes the subject with either a normal ending or an error.
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
self.state.withCriticalRegion { state in
let channels = self.state.withCriticalRegion { state in
state.terminalState = termination
let channels = Array(state.channels.values)
state.channels.removeAll()
state.buffer.removeAll()
state.bufferSize = 0
for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
return channels
}
for channel in channels {
switch termination {
case .finished:
channel.finish()
case .failure(let error):
channel.fail(error)
}
}
}
Expand Down
Loading
Loading