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,5 +1,9 @@
**Unreleased:**

- Breaking: use Swift Async Algorithms for `Sequence.async` and two- or three-input `zip`/`merge`, removing import ambiguities. Add the `AsyncAlgorithms` product dependency and import when migrating these APIs. Preserve variadic `zip`/`merge` and the explicit `AsyncLazySequence` constructor.
- Breaking: rename the buffered Date timer from `AsyncTimerSequence` to `AsyncBufferedTimerSequence` to avoid ambiguity with Apple's clock-based timer.
- SwiftPM: require Swift 5.8 or later for the Swift Async Algorithms test dependency.

- SwitchToLatest: finish cancelled collection while the latest channel or outer sequence remains open, and discard late producer results (https://github.com/sideeffect-io/AsyncExtensions/issues/53).
- Subjects: fix a deadlock when sending values or termination concurrently with consumer cancellation (https://github.com/sideeffect-io/AsyncExtensions/issues/52).
- Subjects: preserve a shared order for concurrent sends, keep current-value/replay delivery consistent with stored state, and queue termination after accepted values without locking during delivery (https://github.com/sideeffect-io/AsyncExtensions/issues/61). Concurrent sends may return while another sender drains their queued delivery.
Expand Down
13 changes: 11 additions & 2 deletions Package.resolved

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 9 additions & 3 deletions Package.swift
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
// swift-tools-version:5.5
// swift-tools-version:5.8
// The swift-tools-version declares the minimum version of Swift required to build this package.

import PackageDescription
Expand All @@ -16,7 +16,10 @@ let package = Package(
name: "AsyncExtensions",
targets: ["AsyncExtensions"]),
],
dependencies: [.package(url: "https://github.com/apple/swift-collections.git", .upToNextMajor(from: "1.0.3"))],
dependencies: [
.package(url: "https://github.com/apple/swift-async-algorithms.git", .upToNextMajor(from: "1.0.0")),
.package(url: "https://github.com/apple/swift-collections.git", .upToNextMajor(from: "1.0.3"))
],
targets: [
.target(
name: "AsyncExtensions",
Expand All @@ -32,7 +35,10 @@ let package = Package(
),
.testTarget(
name: "AsyncExtensionsTests",
dependencies: ["AsyncExtensions"],
dependencies: [
"AsyncExtensions",
.product(name: "AsyncAlgorithms", package: "swift-async-algorithms")
],
path: "Tests"),
]
)
54 changes: 45 additions & 9 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@

**AsyncExtensions** provides a collection of operators that intends to ease the creation and combination of `AsyncSequences`.

**AsyncExtensions** can be seen as a companion to Apple [swift-async-algorithms](https://github.com/apple/swift-async-algorithms). For now there is an overlap between both libraries, but when **swift-async-algorithms** becomes stable the overlapping operators while be deprecated in **AsyncExtensions**. Nevertheless **AsyncExtensions** will continue to provide the operators that the community needs and are not provided by Apple.
**AsyncExtensions** complements Apple [swift-async-algorithms](https://github.com/apple/swift-async-algorithms) with subjects, buffered channels, and operators that Apple does not provide.

## Adding AsyncExtensions as a Dependency

Expand All @@ -28,6 +28,46 @@ Include `"AsyncExtensions"` as a dependency for your executable target:

Finally, add `import AsyncExtensions` to your source code.

## Using Swift Async Algorithms alongside AsyncExtensions

SwiftPM requires Swift 5.8 or later. To use both libraries, declare both package dependencies and add the `AsyncAlgorithms` product to your target alongside `AsyncExtensions`:

```swift
.package(url: "https://github.com/apple/swift-async-algorithms.git", from: "1.0.0"),
```

```swift
.target(
name: "<target>",
dependencies: [
"AsyncExtensions",
.product(name: "AsyncAlgorithms", package: "swift-async-algorithms")
]
),
```

Import both modules in each source file that uses them:

```swift
import AsyncAlgorithms
import AsyncExtensions

let values = [1, 2, 3].async
let pairs = zip(values, ["a", "b", "c"].async)
let merged = merge(values, [4, 5, 6].async)
let rows = zip(values, values, values, values)
```

`Sequence.async`, two- and three-input `zip`/`merge`, and their corresponding sequence types now come from `AsyncAlgorithms`. This is a breaking API change: add that dependency and import when migrating these calls.

AsyncExtensions retains variadic `zip` and `merge`, including support for four or more inputs. Use `AsyncExtensions.zip(...)` or `AsyncExtensions.merge(...)` explicitly when you want the variadic implementations with two or three inputs; variadic `zip` produces arrays, while Apple's fixed overloads produce tuples.

Variadic inputs share one concrete sequence type. Use `eraseToAnyAsyncSequence()` when combining different sequence types through these variadic operators.

The `AsyncLazySequence(sequence)` constructor remains available when you only need AsyncExtensions. The `.async` extension is supplied exclusively by AsyncAlgorithms.

Rename uses of AsyncExtensions' `AsyncTimerSequence` to `AsyncBufferedTimerSequence`. It retains the buffered `Date` values and `DispatchTimeInterval` initializer. The unqualified `AsyncTimerSequence` name now refers to Apple's clock-based timer when both modules are imported.

## Features

### Channels
Expand Down Expand Up @@ -55,12 +95,8 @@ the value. The active sender drains pending deliveries before returning. A new c
or replay consumer receives the latest stored state, followed by subsequent sends.

### Combiners
* [`zip(_:_:)`](./Sources/Combiners/Zip/AsyncZip2Sequence.swift): Zips two `AsyncSequence` into an AsyncSequence of tuple of elements
* [`zip(_:_:_:)`](./Sources/Combiners/Zip/AsyncZip3Sequence.swift): Zips three `AsyncSequence` into an AsyncSequence of tuple of elements
* [`zip(_:)`](./Sources/Combiners/Zip/AsyncZipSequence.swift): Zips any async sequences into an array of elements
* [`merge(_:_:)`](./Sources/Combiners/Merge/AsyncMerge2Sequence.swift): Merges two `AsyncSequence` into an AsyncSequence of elements
* [`merge(_:_:_:)`](./Sources/Combiners/Merge/AsyncMerge3Sequence.swift): Merges three `AsyncSequence` into an AsyncSequence of elements
* [`merge(_:)`](./Sources/Combiners/Merge/AsyncMergeSequence.swift): Merges any `AsyncSequence` into an AsyncSequence of elements
* [`zip(_:)`](./Sources/Combiners/Zip/AsyncZipSequence.swift): Zips any number of async sequences into arrays of elements
* [`merge(_:)`](./Sources/Combiners/Merge/AsyncMergeSequence.swift): Merges any number of async sequences into one sequence
* [`withLatest(_:)`](./Sources/Combiners/WithLatestFrom/AsyncWithLatestFromSequence.swift): Combines elements from self with the last known element from an other `AsyncSequence`
* [`withLatest(_:_:)`](./Sources/Combiners/WithLatestFrom/AsyncWithLatestFrom2Sequence.swift): Combines elements from self with the last known elements from two other async sequences

Expand All @@ -69,8 +105,8 @@ or replay consumer receives the latest stored state, followed by subsequent send
* [AsyncFailSequence](./Sources/Creators/AsyncFailSequence.swift): Creates an `AsyncSequence` that immediately fails
* [AsyncJustSequence](./Sources/Creators/AsyncJustSequence.swift): Creates an `AsyncSequence` that emits an element an finishes
* [AsyncThrowingJustSequence](./Sources/Creators/AsyncThrowingJustSequence.swift): Creates an `AsyncSequence` that emits an elements and finishes bases on a throwing closure
* [AsyncLazySequence](./Sources/Creators/AsyncLazySequence.swift): Creates an `AsyncSequence` of the elements from the base sequence
* [AsyncTimerSequence](./Sources/Creators/AsyncTimerSequence.swift): Creates an `AsyncSequence` that emits a date value periodically
* [AsyncLazySequence](./Sources/Creators/AsyncLazySequence.swift): Creates an async sequence from an explicit synchronous sequence
* [AsyncBufferedTimerSequence](./Sources/Creators/AsyncBufferedTimerSequence.swift): Creates an `AsyncSequence` that buffers date values emitted periodically
* [AsyncStream Pipe](./Sources/Creators/AsyncStream+Pipe.swift): Creates an AsyncStream and returns a tuple standing for its inputs and outputs

### Operators
Expand Down
63 changes: 0 additions & 63 deletions Sources/Combiners/Merge/AsyncMerge2Sequence.swift

This file was deleted.

69 changes: 0 additions & 69 deletions Sources/Combiners/Merge/AsyncMerge3Sequence.swift

This file was deleted.

61 changes: 0 additions & 61 deletions Sources/Combiners/Merge/MergeStateMachine.swift
Original file line number Diff line number Diff line change
Expand Up @@ -29,67 +29,6 @@ struct MergeStateMachine<Element>: Sendable {
let state: ManagedCriticalState<State>
let task: Task<Void, Never>

init<Base1: AsyncSequence, Base2: AsyncSequence>(
_ base1: Base1,
_ base2: Base2
) where Base1.Element == Element, Base2.Element == Element {
self.state = ManagedCriticalState(State(buffer: .idle, basesToTerminate: 2))

let regulator1 = Regulator(base1, onNextRegulatedElement: { [state] in Self.onNextRegulatedElement($0, state: state) })
let regulator2 = Regulator(base2, onNextRegulatedElement: { [state] in Self.onNextRegulatedElement($0, state: state) })

self.requestNextRegulatedElements = {
regulator1.requestNextRegulatedElement()
regulator2.requestNextRegulatedElement()
}

self.task = Task {
await withTaskGroup(of: Void.self) { group in
group.addTask {
await regulator1.iterate()
}

group.addTask {
await regulator2.iterate()
}
}
}
}

init<Base1: AsyncSequence, Base2: AsyncSequence, Base3: AsyncSequence>(
_ base1: Base1,
_ base2: Base2,
_ base3: Base3
) where Base1.Element == Element, Base2.Element == Element, Base3.Element == Base1.Element {
self.state = ManagedCriticalState(State(buffer: .idle, basesToTerminate: 3))

let regulator1 = Regulator(base1, onNextRegulatedElement: { [state] in Self.onNextRegulatedElement($0, state: state) })
let regulator2 = Regulator(base2, onNextRegulatedElement: { [state] in Self.onNextRegulatedElement($0, state: state) })
let regulator3 = Regulator(base3, onNextRegulatedElement: { [state] in Self.onNextRegulatedElement($0, state: state) })

self.requestNextRegulatedElements = {
regulator1.requestNextRegulatedElement()
regulator2.requestNextRegulatedElement()
regulator3.requestNextRegulatedElement()
}

self.task = Task {
await withTaskGroup(of: Void.self) { group in
group.addTask {
await regulator1.iterate()
}

group.addTask {
await regulator2.iterate()
}

group.addTask {
await regulator3.iterate()
}
}
}
}

init<Base: AsyncSequence>(
_ bases: [Base]
) where Base.Element == Element {
Expand Down
Loading