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
2 changes: 0 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -171,8 +171,6 @@ let results = await withOrderedTaskGroup(of: Int.self) { group in
print(results) // 😁 [0, 2, 4, 6, 8, 10]
```

They are also used for async functions of `Sequence`.

## Requirements
Swift 5.10 or later.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ extension AsyncSequence where Element: Sendable {
priority: TaskPriority? = nil,
_ body: @escaping @Sendable (Element) async throws -> Void
) async rethrows {
try await withThrowingOrderedTaskGroup(of: Void.self) { group in
try await withThrowingTaskGroup(of: Void.self) { group in
var counter = 0
var asyncIterator = self.makeAsyncIterator()
while let element = try await asyncIterator.next() {
Expand Down
2 changes: 0 additions & 2 deletions Sources/AsyncOperations/Documentation.docc/Documentation.md
Original file line number Diff line number Diff line change
Expand Up @@ -129,5 +129,3 @@ let results = await withOrderedTaskGroup(of: Int.self) { group in
print(results) // 😁 [0, 2, 4, 6, 8, 10]
```

They are also used for async functions of `Sequence`.

28 changes: 21 additions & 7 deletions Sources/AsyncOperations/Sequence/InternalForEach.swift
Original file line number Diff line number Diff line change
@@ -1,24 +1,38 @@
extension Sequence where Element: Sendable {
func internalForEach<T: Sendable>(
group: inout ThrowingOrderedTaskGroup<T, any Error>,
group: inout ThrowingTaskGroup<(T, Int), any Error>,
numberOfConcurrentTasks: UInt,
priority: TaskPriority?,
taskOperation: @escaping @Sendable (Element) async throws -> T,
nextOperation: (T) -> ()
) async throws {
var currentIndex = 0
var results: [Int: T] = [:]

func doNextOperationIfNeeded() {
while let result = results[currentIndex] {
let index = currentIndex
nextOperation(result)
currentIndex += 1
results.removeValue(forKey: index)
}
}

for (index, element) in self.enumerated() {
if index >= numberOfConcurrentTasks {
if let value = try await group.next() {
nextOperation(value)
if numberOfConcurrentTasks <= index {
if let (value, index) = try await group.next() {
results[index] = value
doNextOperationIfNeeded()
}
}
group.addTask(priority: priority) {
try await taskOperation(element)
try await (taskOperation(element), index)
}
}

for try await value in group {
nextOperation(value)
for try await (value, index) in group {
results[index] = value
doNextOperationIfNeeded()
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ extension Sequence where Element: Sendable {
priority: TaskPriority? = nil,
_ transform: @escaping @Sendable (Element) async throws -> T?
) async rethrows -> [T] {
try await withThrowingOrderedTaskGroup(of: T?.self) { group in
try await withThrowingTaskGroup(of: (T?, Int).self) { group in
var values: [T] = []

try await internalForEach(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ extension Sequence where Element: Sendable {
priority: TaskPriority? = nil,
_ isIncluded: @escaping @Sendable (Element) async throws -> Bool
) async rethrows -> [Element] {
try await withThrowingOrderedTaskGroup(of: Element?.self) { group in
try await withThrowingTaskGroup(of: (Element?, Int).self) { group in
var values: [Element] = []

try await internalForEach(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ extension Sequence where Element: Sendable {
priority: TaskPriority? = nil,
_ transform: @escaping @Sendable (Element) async throws -> [T]
) async rethrows -> [T] {
try await withThrowingOrderedTaskGroup(of: [T].self) { group in
try await withThrowingTaskGroup(of: ([T], Int).self) { group in
var values: [T] = []

try await internalForEach(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ extension Sequence where Element: Sendable {
priority: TaskPriority? = nil,
_ body: @escaping @Sendable (Element) async throws -> Void
) async rethrows {
try await withThrowingOrderedTaskGroup(of: Void.self) { group in
try await withThrowingTaskGroup(of: (Void, Int).self) { group in
try await internalForEach(
group: &group,
numberOfConcurrentTasks: numberOfConcurrentTasks,
Expand Down
2 changes: 1 addition & 1 deletion Sources/AsyncOperations/Sequence/Sequence+AsyncMap.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ extension Sequence where Element: Sendable {
priority: TaskPriority? = nil,
_ transform: @escaping @Sendable (Element) async throws -> T
) async rethrows -> [T] {
try await withThrowingOrderedTaskGroup(of: T.self) { group in
try await withThrowingTaskGroup(of: (T, Int).self) { group in
var values: [T] = []

try await internalForEach(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,17 +5,21 @@ final class AsyncSequenceAsyncForEachTests: XCTestCase {
@MainActor
func testAsyncForEach() async throws {
var results: [Int] = []
var events: [ConcurrentTaskEvent] = []

let asyncSequence = AsyncStream { c in
(0..<5).forEach { c.yield($0) }
c.finish()
}

await asyncSequence.asyncForEach(numberOfConcurrentTasks: 3) { @MainActor number in
await Task.yield()
try await asyncSequence.asyncForEach(numberOfConcurrentTasks: 3) { @MainActor number in
events.append(.start)
try await Task.sleep(for: .milliseconds(100 * (5 - number)))
events.append(.end)
results.append(number)
}
XCTAssertEqual(results.count, 5)
XCTAssertEqual(Set(results), [0, 1, 2, 3, 4])
XCTAssertEqual(events, [.start, .start, .start, .end, .start, .end, .start, .end, .end, .end])
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,43 @@ final class SequenceAsyncForEachTests: XCTestCase {
@MainActor
func testAsyncForEach() async throws {
var results: [Int] = []
await [0, 1, 2, 3, 4].asyncForEach { @MainActor number in
await Task.yield()
try await [0, 1, 2, 3, 4].asyncForEach { @MainActor number in
try await Task.sleep(for: .milliseconds(100 * (5 - number)))
results.append(number)
}
XCTAssertEqual(results.count, 5)
XCTAssertEqual(Set(results), [0, 1, 2, 3, 4])
XCTAssertEqual(results, [0, 1, 2, 3, 4])
}

@MainActor
func testAsyncForEachConcurrently() async throws {
var events: [ConcurrentTaskEvent] = []
let numberOfElements = 10
let numberOfConcurrentTasks = 3
// A random offset to stagger iteration timings.
let randomOffsets: [Double] = [
470,
150,
420,
290,
670,
810,
930,
120,
540,
760,
]
try await (0..<numberOfElements).asyncForEach(numberOfConcurrentTasks: UInt(numberOfConcurrentTasks)) { @MainActor number in
events.append(.start)
let offsets = randomOffsets[number]
try await Task.sleep(for: .milliseconds(100 * Double(numberOfElements - number) + offsets))
events.append(.end)
}
print(events.map(\.rawValue))
XCTAssertEqual(Array(events.prefix(numberOfConcurrentTasks)), Array(repeating: .start, count: numberOfConcurrentTasks))
XCTAssertEqual(
events[numberOfConcurrentTasks..<(2 * numberOfElements - numberOfConcurrentTasks)].map { $0 },
Array(repeating: [ConcurrentTaskEvent.end, .start], count: numberOfElements - numberOfConcurrentTasks).flatMap { $0 }
)
XCTAssertEqual(Array(events.suffix(Int(numberOfConcurrentTasks))), Array(repeating: .end, count: Int(numberOfConcurrentTasks)))
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ final class SequenceAsyncMapTests: XCTestCase {
let difference = endTime.timeIntervalSince(startTime)

XCTAssertEqual(results, [0, 2, 4, 6, 8])
XCTAssertLessThan(difference, 4)
XCTAssertLessThan(difference, 2)
}

func testAsyncMapMultipleTasksWithNumberOfConcurrentTasks() async throws {
Expand Down
4 changes: 4 additions & 0 deletions Tests/AsyncOperationsTests/Utils/ConcurrentTaskEvent.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
enum ConcurrentTaskEvent: String {
case start
case end
}