Improved Defer Fulfillment (#5127)

* improved defer fullfillment

* applied also for async streams

* fix xcodeproject

* better documentation

* restict deferFulfillment  only for Publisher which can never fail
This commit is contained in:
Mauro
2026-02-20 18:21:35 +01:00
committed by GitHub
parent 5e5fb19f42
commit fda19a0273
4 changed files with 149 additions and 117 deletions

View File

@@ -7,6 +7,7 @@
//
import Combine
import Testing
struct DeferredFulfillment<T> {
let closure: () async throws -> T
@@ -18,225 +19,254 @@ struct DeferredFulfillment<T> {
}
struct DeferredFulfillmentError: Error {
enum Kind {
case noOutput
case unexpectedFulfillment
static func noOutput(message: String?, sourceLocation: SourceLocation) -> Self {
defer { Issue.record(Comment(rawValue: message ?? "No Output"), sourceLocation: sourceLocation) }
return .init()
}
let kind: Kind
let message: String?
static func noOutput(message: String?) -> Self {
.init(kind: .noOutput, message: message)
static func unexpectedFulfillment(message: String?, sourceLocation: SourceLocation) -> Self {
defer { Issue.record(Comment(rawValue: message ?? "Unexpected Fulfillment"), sourceLocation: sourceLocation) }
return .init()
}
static func unexpectedFulfillment(message: String?) -> Self {
.init(kind: .unexpectedFulfillment, message: message)
static var empty: Self {
.init()
}
}
/// Utility that assists in subscribing to a publisher and deferring the fulfilment and results until some other actions have been performed.
/// Test utility that assists in subscribing to a publisher and deferring the fulfilment and results until some other actions have been performed.
/// - Parameters:
/// - publisher: The publisher to wait on.
/// - timeout: A timeout after which we give up.
/// - message: An optional message to include in the error if the condition is never met.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - until: callback that evaluates outputs until some condition is reached
/// - Returns: The deferred fulfilment to be executed after some actions and that returns the result of the publisher.
func deferFulfillment<P: Publisher>(_ publisher: P,
timeout: Duration = .seconds(10),
message: String? = nil,
until condition: @escaping (P.Output) -> Bool) -> DeferredFulfillment<P.Output> {
var result: Result<P.Output, Error>?
var hasFulfilled = false
func deferFulfillment<P: Publisher<P.Output, Never>>(_ publisher: P,
timeout: Duration = .seconds(10),
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation,
until condition: @escaping (P.Output) -> Bool) -> DeferredFulfillment<P.Output> {
let (stream, continuation) = AsyncStream<P.Output>.makeStream()
let cancellable = publisher
.sink { completion in
switch completion {
case .failure(let error):
result = .failure(error)
hasFulfilled = true
case .finished:
break
}
.sink { _ in
continuation.finish()
} receiveValue: { value in
if condition(value), !hasFulfilled {
result = .success(value)
hasFulfilled = true
}
guard condition(value) else { return }
continuation.yield(value)
continuation.finish()
}
return DeferredFulfillment<P.Output> {
let startTime = ContinuousClock.now
return DeferredFulfillment {
defer { cancellable.cancel() }
while !hasFulfilled {
await Task.yield()
if ContinuousClock.now - startTime >= timeout {
break
return try await withThrowingTaskGroup(of: P.Output.self) { group in
group.addTask {
for await result in stream {
return result
}
guard !Task.isCancelled else {
// Required to avoid a double recording of the issue in the case where the task is cancelled due to timeout.
throw DeferredFulfillmentError.empty
}
throw DeferredFulfillmentError.noOutput(message: message, sourceLocation: sourceLocation)
}
group.addTask {
try await Task.sleep(for: timeout)
throw DeferredFulfillmentError.noOutput(message: message, sourceLocation: sourceLocation)
}
defer { group.cancelAll() }
return try #require(try await group.next())
}
cancellable.cancel()
guard let unwrappedResult = result else {
throw DeferredFulfillmentError.noOutput(message: message)
}
return try unwrappedResult.get()
}
}
/// Utility that assists in observing an async sequence, deferring the fulfilment and results until some condition has been met.
/// Test utility that assists in observing an async sequence, deferring the fulfilment and results until some condition has been met.
/// - Parameters:
/// - asyncSequence: The sequence to wait on.
/// - timeout: A timeout after which we give up.
/// - message: An optional message to include in the error if the condition is never met.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - until: callback that evaluates outputs until some condition is reached
/// - Returns: The deferred fulfilment to be executed after some actions and that returns the result of the sequence.
func deferFulfillment<Value>(_ asyncSequence: any AsyncSequence<Value, Never>,
timeout: Duration = .seconds(10),
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation,
until condition: @escaping (Value) -> Bool) -> DeferredFulfillment<Value> {
var result: Result<Value, Error>?
var hasFulfilled = false
let (stream, continuation) = AsyncStream<Value>.makeStream()
let task = Task {
for await value in asyncSequence {
if condition(value), !hasFulfilled {
result = .success(value)
hasFulfilled = true
}
for await value in asyncSequence where condition(value) {
continuation.yield(value)
continuation.finish()
return
}
continuation.finish()
}
return DeferredFulfillment<Value> {
let startTime = ContinuousClock.now
return DeferredFulfillment {
defer { task.cancel() }
while !hasFulfilled {
await Task.yield()
if ContinuousClock.now - startTime >= timeout {
break
return try await withThrowingTaskGroup(of: Value.self) { group in
group.addTask {
for await value in stream {
return value
}
guard !Task.isCancelled else {
// Required to avoid a double recording of the issue in the case where the task is cancelled due to timeout.
throw DeferredFulfillmentError.empty
}
throw DeferredFulfillmentError.noOutput(message: message, sourceLocation: sourceLocation)
}
group.addTask {
try await Task.sleep(for: timeout)
throw DeferredFulfillmentError.noOutput(message: message, sourceLocation: sourceLocation)
}
defer { group.cancelAll() }
return try #require(try await group.next())
}
task.cancel()
guard let unwrappedResult = result else {
throw DeferredFulfillmentError.noOutput(message: message)
}
return try unwrappedResult.get()
}
}
/// Utility that assists in subscribing to a publisher and deferring the fulfilment and results until some other actions have been performed.
/// Test utility that assists in subscribing to a publisher and deferring the fulfilment and results until some other actions have been performed.
/// - Parameters:
/// - publisher: The publisher to wait on.
/// - keyPath: the key path for the expected values
/// - transitionValues: the values through which the keypath needs to transition through
/// - timeout: A timeout after which we give up.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - Returns: The deferred fulfilment to be executed after some actions and that returns the result of the publisher.
func deferFulfillment<P: Publisher, K: KeyPath<P.Output, V>, V: Equatable>(_ publisher: P,
keyPath: K,
transitionValues: [V],
timeout: Duration = .seconds(10)) -> DeferredFulfillment<P.Output> {
func deferFulfillment<P: Publisher<P.Output, Never>, K: KeyPath<P.Output, V>, V: Equatable>(_ publisher: P,
keyPath: K,
transitionValues: [V],
timeout: Duration = .seconds(10),
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation) -> DeferredFulfillment<P.Output> {
var expectedOrder = transitionValues
return deferFulfillment(publisher, timeout: timeout) { value in
return deferFulfillment(publisher, timeout: timeout, message: message, sourceLocation: sourceLocation) { value in
let receivedValue = value[keyPath: keyPath]
if let index = expectedOrder.firstIndex(where: { $0 == receivedValue }), index == 0 {
expectedOrder.remove(at: index)
}
return expectedOrder.isEmpty
}
}
/// Utility that assists in subscribing to an async sequence and deferring the fulfilment and results until some other actions have been performed.
/// Test utility that assists in subscribing to an async sequence and deferring the fulfilment and results until some other actions have been performed.
/// - Parameters:
/// - asyncSequence: The sequence to wait on.
/// - transitionValues: the values through which the sequence needs to transition through
/// - timeout: A timeout after which we give up.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - Returns: The deferred fulfilment to be executed after some actions and that returns the result of the sequence.
func deferFulfillment<Value: Equatable>(_ asyncSequence: any AsyncSequence<Value, Never>,
transitionValues: [Value],
timeout: Duration = .seconds(10)) -> DeferredFulfillment<Value> {
timeout: Duration = .seconds(10),
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation) -> DeferredFulfillment<Value> {
var expectedOrder = transitionValues
return deferFulfillment(asyncSequence, timeout: timeout) { value in
return deferFulfillment(asyncSequence, timeout: timeout, message: message, sourceLocation: sourceLocation) { value in
if let index = expectedOrder.firstIndex(where: { $0 == value }), index == 0 {
expectedOrder.remove(at: index)
}
return expectedOrder.isEmpty
}
}
/// Utility that assists in subscribing to a publisher and deferring the failure for a particular value until some other actions have been performed.
/// Test utility that assists in subscribing to a publisher and deferring the failure for a particular value until some other actions have been performed.
/// - Parameters:
/// - publisher: The publisher to wait on.
/// - timeout: A timeout after which we give up.
/// - message: An optional message to include in the error if the condition is unexpectedly met.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - until: callback that evaluates outputs until some condition is reached
/// - Returns: The deferred fulfilment to be executed after some actions. The publisher's result is not returned from this fulfilment.
func deferFailure<P: Publisher>(_ publisher: P,
timeout: Duration,
message: String? = nil,
until condition: @escaping (P.Output) -> Bool) -> DeferredFulfillment<Void> where P.Failure == Never {
var hasFulfilled = false
func deferFailure<P: Publisher<P.Output, Never>>(_ publisher: P,
timeout: Duration,
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation,
until condition: @escaping (P.Output) -> Bool) -> DeferredFulfillment<Void> where P.Failure == Never {
let (stream, continuation) = AsyncStream<Void>.makeStream()
let cancellable = publisher
.sink { value in
if condition(value), !hasFulfilled {
hasFulfilled = true
}
guard condition(value) else { return }
continuation.yield(())
continuation.finish()
}
return DeferredFulfillment<Void> {
let startTime = ContinuousClock.now
return DeferredFulfillment {
defer { cancellable.cancel() }
while !hasFulfilled {
await Task.yield()
if ContinuousClock.now - startTime >= timeout {
break
try await withThrowingTaskGroup(of: Void.self) { group in
// If the condition fires before timeout, that's the unexpected failure.
group.addTask {
for await _ in stream {
throw DeferredFulfillmentError.unexpectedFulfillment(message: message, sourceLocation: sourceLocation)
}
// Stream finished without condition firing this shouldn't happen
// but is safe to treat as success.
}
}
cancellable.cancel()
// For deferFailure, if hasFulfilled is true, it means the condition was met (which is a failure)
if hasFulfilled {
throw DeferredFulfillmentError.unexpectedFulfillment(message: message)
// Timeout elapsing without the condition firing = success.
group.addTask {
try await Task.sleep(for: timeout)
}
defer { group.cancelAll() }
return try #require(try await group.next())
}
}
}
/// Utility that assists in subscribing to an async sequence and deferring the failure for a particular value until some other actions have been performed.
/// Test utility that assists in subscribing to an async sequence and deferring the failure for a particular value until some other actions have been performed.
/// - Parameters:
/// - asyncSequence: The sequence to wait on.
/// - timeout: A timeout after which we give up.
/// - message: An optional message to include in the error if the condition is unexpectedly met.
/// - sourceLocation: The source location to attach to any recorded issues.
/// - until: callback that evaluates outputs until some condition is reached
/// - Returns: The deferred fulfilment to be executed after some actions. The sequence's result is not returned from this fulfilment.
func deferFailure<Value>(_ asyncSequence: any AsyncSequence<Value, Never>,
timeout: Duration,
message: String? = nil,
sourceLocation: SourceLocation = #_sourceLocation,
until condition: @escaping (Value) -> Bool) -> DeferredFulfillment<Void> {
var hasFulfilled = false
let (stream, continuation) = AsyncStream<Void>.makeStream()
let task = Task {
for await value in asyncSequence {
if condition(value), !hasFulfilled {
hasFulfilled = true
}
for await value in asyncSequence where condition(value) {
continuation.yield(())
continuation.finish()
return
}
continuation.finish()
}
return DeferredFulfillment<Void> {
let startTime = ContinuousClock.now
return DeferredFulfillment {
defer { task.cancel() }
while !hasFulfilled {
await Task.yield()
if ContinuousClock.now - startTime >= timeout {
break
try await withThrowingTaskGroup(of: Void.self) { group in
// If the condition fires before timeout, that's the unexpected failure.
group.addTask {
for await _ in stream {
throw DeferredFulfillmentError.unexpectedFulfillment(message: message, sourceLocation: sourceLocation)
}
}
}
task.cancel()
// For deferFailure, if hasFulfilled is true, it means the condition was met (which is a failure)
if hasFulfilled {
throw DeferredFulfillmentError.unexpectedFulfillment(message: message)
// Timeout elapsing without the condition firing = success.
group.addTask {
try await Task.sleep(for: timeout)
}
defer { group.cancelAll() }
return try #require(try await group.next())
}
}
}