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
20 changes: 10 additions & 10 deletions Sources/AsyncSubjects/AsyncCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -91,23 +91,23 @@ public final class AsyncCurrentValueSubject<Element>: AsyncSubject where Element
func handleNewConsumer() -> (iterator: AsyncBufferedChannel<Element>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncBufferedChannel<Element>()

let (terminalState, current) = self.state.withCriticalRegion { state -> (Termination?, Element) in
(state.terminalState, state.current)
}

if let terminalState = terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return (asyncBufferedChannel.makeAsyncIterator(), {})
}
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return nil
}

asyncBufferedChannel.send(current)
asyncBufferedChannel.send(state.current)

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
18 changes: 9 additions & 9 deletions Sources/AsyncSubjects/AsyncPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -75,21 +75,21 @@ public final class AsyncPassthroughSubject<Element: Sendable>: AsyncSubject {
func handleNewConsumer() -> (iterator: AsyncBufferedChannel<Element>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncBufferedChannel<Element>()

let terminalState = self.state.withCriticalRegion { state in
state.terminalState
}

if let terminalState = terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return (asyncBufferedChannel.makeAsyncIterator(), {})
}
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return nil
}

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
24 changes: 12 additions & 12 deletions Sources/AsyncSubjects/AsyncReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -75,25 +75,25 @@ public final class AsyncReplaySubject<Element>: AsyncSubject where Element: Send
func handleNewConsumer() -> (iterator: AsyncBufferedChannel<Element>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncBufferedChannel<Element>()

let (terminalState, elements) = self.state.withCriticalRegion { state -> (Termination?, [Element]) in
(state.terminalState, state.buffer)
}

if let terminalState = terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return (asyncBufferedChannel.makeAsyncIterator(), {})
}
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState, terminalState.isFinished {
asyncBufferedChannel.finish()
return nil
}

for element in elements {
asyncBufferedChannel.send(element)
}
for element in state.buffer {
asyncBufferedChannel.send(element)
}

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
28 changes: 14 additions & 14 deletions Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -97,28 +97,28 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: As
) -> (iterator: AsyncThrowingBufferedChannel<Element, Error>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncThrowingBufferedChannel<Element, Error>()

let (terminalState, current) = self.state.withCriticalRegion { state -> (Termination?, Element) in
(state.terminalState, state.current)
}

if let terminalState = terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
}
return nil
}
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

asyncBufferedChannel.send(current)
asyncBufferedChannel.send(state.current)

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
26 changes: 13 additions & 13 deletions Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -83,26 +83,26 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: Asy
) -> (iterator: AsyncThrowingBufferedChannel<Element, Error>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncThrowingBufferedChannel<Element, Error>()

let terminalState = self.state.withCriticalRegion { state in
state.terminalState
}

if let terminalState = terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
}
return nil
}
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
32 changes: 16 additions & 16 deletions Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -80,30 +80,30 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: AsyncSub
) -> (iterator: AsyncThrowingBufferedChannel<Element, Error>.Iterator, unregister: @Sendable () -> Void) {
let asyncBufferedChannel = AsyncThrowingBufferedChannel<Element, Error>()

let (terminalState, elements) = self.state.withCriticalRegion { state -> (Termination?, [Element]) in
(state.terminalState, state.buffer)
}

if let terminalState = terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
let consumerId = self.state.withCriticalRegion { state -> Int? in
if let terminalState = state.terminalState {
switch terminalState {
case .finished:
asyncBufferedChannel.finish()
case .failure(let error):
asyncBufferedChannel.fail(error)
}
return nil
}
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

for element in elements {
asyncBufferedChannel.send(element)
}
for element in state.buffer {
asyncBufferedChannel.send(element)
}

let consumerId = self.state.withCriticalRegion { state -> Int in
state.ids += 1
state.channels[state.ids] = asyncBufferedChannel
return state.ids
}

guard let consumerId = consumerId else {
return (asyncBufferedChannel.makeAsyncIterator(), {})
}

let unregister = { @Sendable [state] in
state.withCriticalRegion { state in
state.channels[consumerId] = nil
Expand Down
28 changes: 28 additions & 0 deletions Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -204,4 +204,32 @@ final class AsyncCurrentValueSubjectTests: XCTestCase {
XCTAssertEqual(receivedElementsA, expectedElements)
XCTAssertEqual(receivedElementsB, expectedElements)
}

func test_subscription_racing_send_receives_sent_element() async {
for _ in 0..<10_000 {
let sut = AsyncCurrentValueSubject<Int>(0)
var iterator: AsyncCurrentValueSubject<Int>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(1) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.elements.last == 1 else {
return XCTFail("Expected to receive the sent element, received \(drained.elements)")
}
}
}

func test_subscription_racing_termination_is_terminated() async {
for _ in 0..<10_000 {
let sut = AsyncCurrentValueSubject<Int>(0)
var iterator: AsyncCurrentValueSubject<Int>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(.finished) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.isTerminated else {
return XCTFail("Expected the subscription to be terminated")
}
}
}
}
14 changes: 14 additions & 0 deletions Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -198,4 +198,18 @@ final class AsyncPassthroughSubjectTests: XCTestCase {
XCTAssertEqual(receivedElementsA, expectedElements)
XCTAssertEqual(receivedElementsB, expectedElements)
}

func test_subscription_racing_termination_is_terminated() async {
for _ in 0..<10_000 {
let sut = AsyncPassthroughSubject<Int>()
var iterator: AsyncPassthroughSubject<Int>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(.finished) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.isTerminated else {
return XCTFail("Expected the subscription to be terminated")
}
}
}
}
30 changes: 30 additions & 0 deletions Tests/AsyncSubjets/AsyncReplaySubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -227,4 +227,34 @@ final class AsyncReplaySubjectTests: XCTestCase {
XCTAssertEqual(receivedElementsA, expectedElements)
XCTAssertEqual(receivedElementsB, expectedElements)
}

func test_subscription_racing_send_receives_sent_element() async {
for _ in 0..<10_000 {
let sut = AsyncReplaySubject<Int>(bufferSize: 2)
sut.send(0)
var iterator: AsyncReplaySubject<Int>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(1) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.elements.last == 1 else {
return XCTFail("Expected to receive the sent element, received \(drained.elements)")
}
}
}

func test_subscription_racing_termination_is_terminated() async {
for _ in 0..<10_000 {
let sut = AsyncReplaySubject<Int>(bufferSize: 2)
sut.send(0)
var iterator: AsyncReplaySubject<Int>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(.finished) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.isTerminated else {
return XCTFail("Expected the subscription to be terminated")
}
}
}
}
28 changes: 28 additions & 0 deletions Tests/AsyncSubjets/AsyncThrowingCurrentValueSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -256,4 +256,32 @@ final class AsyncThrowingCurrentValueSubjectTests: XCTestCase {
XCTAssertEqual(receivedElementsA, expectedElements)
XCTAssertEqual(receivedElementsB, expectedElements)
}

func test_subscription_racing_send_receives_sent_element() async {
for _ in 0..<10_000 {
let sut = AsyncThrowingCurrentValueSubject<Int, Error>(0)
var iterator: AsyncThrowingCurrentValueSubject<Int, Error>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(1) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.elements.last == 1 else {
return XCTFail("Expected to receive the sent element, received \(drained.elements)")
}
}
}

func test_subscription_racing_termination_is_terminated() async {
for _ in 0..<10_000 {
let sut = AsyncThrowingCurrentValueSubject<Int, Error>(0)
var iterator: AsyncThrowingCurrentValueSubject<Int, Error>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(.failure(MockError(code: 1))) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.isTerminated else {
return XCTFail("Expected the subscription to be terminated")
}
}
}
}
14 changes: 14 additions & 0 deletions Tests/AsyncSubjets/AsyncThrowingPassthroughSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -261,4 +261,18 @@ final class AsyncThrowingPassthroughSubjectTests: XCTestCase {
XCTAssertEqual(receivedElementsA, expectedElements)
XCTAssertEqual(receivedElementsB, expectedElements)
}

func test_subscription_racing_termination_is_terminated() async {
for _ in 0..<10_000 {
let sut = AsyncThrowingPassthroughSubject<Int, Error>()
var iterator: AsyncThrowingPassthroughSubject<Int, Error>.Iterator?

race({ iterator = sut.makeAsyncIterator() }, { sut.send(.failure(MockError(code: 1))) })

let drained = await drainBufferedElements(of: iterator!)
guard drained.isTerminated else {
return XCTFail("Expected the subscription to be terminated")
}
}
}
}
Loading