: AsyncSequence where P.Failure =
self.publisher = publisher
}
- public struct Iterator: NonThrowingAsyncIteratorProtocol {
+ public struct Iterator: AsyncIteratorProtocol, TypedAsyncIteratorProtocol {
public typealias Element = P.Output
+ public typealias Failure = P.Failure
@usableFromInline
internal let inner = AsyncSubscriber()
@usableFromInline
internal let reference:AnyCancellable
+ @_disfavoredOverload
@inlinable
- public mutating func next() async -> P.Output? {
- let result = await withTaskCancellationHandler(operation: inner.next) { [reference] in
+ public mutating func next() async throws(Never) -> P.Output? {
+ await next(isolation: nil)
+ }
+
+ @inlinable
+ public func next(isolation actor: isolated (any Actor)? = #isolation) async -> P.Output? {
+ let result: Result
? = await withTaskCancellationHandler { [inner] in
+ await inner.next(isolation: actor)
+ } onCancel: { [reference] in
reference.cancel()
}
+
switch result {
case .none:
return nil
diff --git a/Sources/Tetra/Combine/CompatAsyncThrowingPublisher.swift b/Sources/Tetra/Combine/CompatAsyncThrowingPublisher.swift
index 1bef09b..cbe037c 100644
--- a/Sources/Tetra/Combine/CompatAsyncThrowingPublisher.swift
+++ b/Sources/Tetra/Combine/CompatAsyncThrowingPublisher.swift
@@ -8,11 +8,12 @@
import Foundation
@preconcurrency import Combine
+public import BackPortAsyncSequence
-public struct CompatAsyncThrowingPublisher: AsyncTypedSequence {
+public struct CompatAsyncThrowingPublisher: AsyncSequence, TypedAsyncSequence {
public typealias AsyncIterator = Iterator
- public typealias Element = P.Output
+ public typealias Failure = P.Failure
public var publisher:P
@@ -21,19 +22,23 @@ public struct CompatAsyncThrowingPublisher: AsyncTypedSequence {
Iterator(source: publisher)
}
- public struct Iterator: AsyncIteratorProtocol {
+ public struct Iterator: AsyncIteratorProtocol, TypedAsyncIteratorProtocol {
public typealias Element = P.Output
+ public typealias Failure = P.Failure
@usableFromInline
internal let inner = AsyncSubscriber