From 1e20cf5faf8c6012b7d2204ee04825cf4f94941b Mon Sep 17 00:00:00 2001 From: Ferenc Viasz-Kadi Date: Tue, 15 Sep 2026 17:20:29 +0200 Subject: [PATCH 1/3] Error-erasing StorageSequence and demand test Replace the AsyncStream> wrapper with a small ErrorErasingSequence that erases the failure type, simplifying StorageSequence internals and iterator typing (uses any AsyncIteratorProtocol/any AsyncSequence). This removes the Result wrapping and stream machinery while preserving behavior. Add a unit test (DemandProbe + DemandTrackingSequence) that verifies StorageSequence only requests upstream elements when the consumer advances (demand-driven). --- Sources/FeatherStorage/StorageSequence.swift | 61 +++++++++++-------- .../StorageSequenceTestSuite.swift | 58 ++++++++++++++++++ 2 files changed, 92 insertions(+), 27 deletions(-) diff --git a/Sources/FeatherStorage/StorageSequence.swift b/Sources/FeatherStorage/StorageSequence.swift index fa85794..ddeacb0 100644 --- a/Sources/FeatherStorage/StorageSequence.swift +++ b/Sources/FeatherStorage/StorageSequence.swift @@ -8,13 +8,13 @@ import NIOCore /// A type-erased async sequence of storage byte buffers. public struct StorageSequence: Sendable, AsyncSequence { - typealias BaseAsyncSequence = AsyncStream> - /// An async iterator over a `StorageSequence`. public struct AsyncIterator: AsyncIteratorProtocol { - private var base: BaseAsyncSequence.AsyncIterator + private var base: any AsyncIteratorProtocol - init(base: BaseAsyncSequence.AsyncIterator) { + init( + base: any AsyncIteratorProtocol + ) { self.base = base } @@ -25,7 +25,7 @@ public struct StorageSequence: Sendable, AsyncSequence { /// - Throws: Any error emitted by the underlying async sequence. @concurrent public mutating func next() async throws -> ByteBuffer? { - try await self.base.next(isolation: nil)?.get() + try await base.next(isolation: nil) } #else /// Returns the next available byte buffer from the sequence. @@ -33,7 +33,7 @@ public struct StorageSequence: Sendable, AsyncSequence { /// - Returns: The next `ByteBuffer`, or `nil` when the sequence is finished. /// - Throws: Any error emitted by the underlying async sequence. public mutating func next() async throws -> ByteBuffer? { - try await self.base.next()?.get() + try await base.next() } #endif @@ -45,11 +45,35 @@ public struct StorageSequence: Sendable, AsyncSequence { public mutating func next( isolation actor: isolated (any Actor)? ) async throws -> Element? { - try await self.base.next(isolation: actor)?.get() + try await base.next(isolation: actor) + } + } + + private struct ErrorErasingSequence: + AsyncSequence, + Sendable + where Base.Element == ByteBuffer { + typealias Element = ByteBuffer + typealias Failure = any Error + + struct AsyncIterator: AsyncIteratorProtocol { + var base: Base.AsyncIterator + + mutating func next( + isolation actor: isolated (any Actor)? + ) async throws(any Error) -> ByteBuffer? { + try await base.next(isolation: actor) + } + } + + let base: Base + + func makeAsyncIterator() -> AsyncIterator { + .init(base: base.makeAsyncIterator()) } } - private let makeIteratorCallback: @Sendable () -> BaseAsyncSequence + private let base: any AsyncSequence & Sendable /// Optional known byte length of the sequence. public let length: UInt64? @@ -64,24 +88,7 @@ public struct StorageSequence: Sendable, AsyncSequence { length: UInt64? = nil ) where S.Element == ByteBuffer { self.length = length - self.makeIteratorCallback = { - BaseAsyncSequence { continuation in - let task = Task { - do { - for try await element in asyncSequence { - continuation.yield(.success(element)) - } - } - catch { - continuation.yield(.failure(error)) - } - continuation.finish() - } - continuation.onTermination = { _ in - task.cancel() - } - } - } + self.base = ErrorErasingSequence(base: asyncSequence) } /// Creates a type-erased storage sequence from a byte buffer. @@ -106,6 +113,6 @@ public struct StorageSequence: Sendable, AsyncSequence { /// /// - Returns: A new `AsyncIterator` instance. public func makeAsyncIterator() -> AsyncIterator { - AsyncIterator(base: makeIteratorCallback().makeAsyncIterator()) + AsyncIterator(base: base.makeAsyncIterator()) } } diff --git a/Tests/FeatherStorageTests/StorageSequenceTestSuite.swift b/Tests/FeatherStorageTests/StorageSequenceTestSuite.swift index d5b0c6f..7a05d97 100644 --- a/Tests/FeatherStorageTests/StorageSequenceTestSuite.swift +++ b/Tests/FeatherStorageTests/StorageSequenceTestSuite.swift @@ -16,6 +16,41 @@ struct StorageSequenceTestSuite { case failed } + private actor DemandProbe { + private(set) var requestCount = 0 + + func recordRequest() { + requestCount += 1 + } + } + + private struct DemandTrackingSequence: AsyncSequence, Sendable { + typealias Element = ByteBuffer + + struct AsyncIterator: AsyncIteratorProtocol { + let probe: DemandProbe + var remainingCount: Int + + mutating func next( + isolation actor: isolated (any Actor)? + ) async -> ByteBuffer? { + guard remainingCount > 0 else { + return nil + } + remainingCount -= 1 + await probe.recordRequest() + return ByteBuffer(bytes: [UInt8(remainingCount)]) + } + } + + let probe: DemandProbe + let count: Int + + func makeAsyncIterator() -> AsyncIterator { + .init(probe: probe, remainingCount: count) + } + } + @Test func initFromAsyncSequencePreservesElementsAndLength() async throws { let allocator = ByteBufferAllocator() @@ -93,6 +128,29 @@ struct StorageSequenceTestSuite { } } + @Test + func requestsUpstreamElementsOnlyWhenConsumerAdvances() async throws { + let probe = DemandProbe() + let sequence = StorageSequence( + asyncSequence: DemandTrackingSequence(probe: probe, count: 3) + ) + var iterator = sequence.makeAsyncIterator() + + #expect(await probe.requestCount == 0) + + _ = try await iterator.next() + for _ in 0..<10 { + await Task.yield() + } + #expect(await probe.requestCount == 1) + + _ = try await iterator.next() + for _ in 0..<10 { + await Task.yield() + } + #expect(await probe.requestCount == 2) + } + private static func makeBuffer( _ bytes: [UInt8], allocator: ByteBufferAllocator From 7ca2d98005adaa78a7466fb750bcdcbfd5fde7ce Mon Sep 17 00:00:00 2001 From: Ferenc Viasz-Kadi Date: Tue, 15 Sep 2026 17:35:04 +0200 Subject: [PATCH 2/3] Fix StorageSequence iterator reuse StorageSequence was retaining a single base async sequence, which could reuse iterator state across calls. This change creates a fresh iterator per makeAsyncIterator invocation while preserving the existing error-erased sequence behavior. --- Sources/FeatherStorage/StorageSequence.swift | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/Sources/FeatherStorage/StorageSequence.swift b/Sources/FeatherStorage/StorageSequence.swift index ddeacb0..b4221d3 100644 --- a/Sources/FeatherStorage/StorageSequence.swift +++ b/Sources/FeatherStorage/StorageSequence.swift @@ -73,7 +73,7 @@ public struct StorageSequence: Sendable, AsyncSequence { } } - private let base: any AsyncSequence & Sendable + private let makeIterator: @Sendable () -> AsyncIterator /// Optional known byte length of the sequence. public let length: UInt64? @@ -88,7 +88,12 @@ public struct StorageSequence: Sendable, AsyncSequence { length: UInt64? = nil ) where S.Element == ByteBuffer { self.length = length - self.base = ErrorErasingSequence(base: asyncSequence) + self.makeIterator = { + AsyncIterator( + base: ErrorErasingSequence(base: asyncSequence) + .makeAsyncIterator() + ) + } } /// Creates a type-erased storage sequence from a byte buffer. @@ -113,6 +118,6 @@ public struct StorageSequence: Sendable, AsyncSequence { /// /// - Returns: A new `AsyncIterator` instance. public func makeAsyncIterator() -> AsyncIterator { - AsyncIterator(base: base.makeAsyncIterator()) + makeIterator() } } From ab6ff373a2605556576ad17e242805c38e2cb319 Mon Sep 17 00:00:00 2001 From: Ferenc Viasz-Kadi Date: Wed, 16 Sep 2026 11:49:32 +0200 Subject: [PATCH 3/3] Rename failure-erasing async sequence This change renames the internal StorageSequence wrapper from ErrorErasingSequence to FailureErasingAsyncSequence to better reflect its purpose. The sequence still erases failure types while preserving byte-buffer output, keeping the StorageSequence API behavior unchanged. --- Sources/FeatherStorage/StorageSequence.swift | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Sources/FeatherStorage/StorageSequence.swift b/Sources/FeatherStorage/StorageSequence.swift index b4221d3..e97041f 100644 --- a/Sources/FeatherStorage/StorageSequence.swift +++ b/Sources/FeatherStorage/StorageSequence.swift @@ -49,7 +49,7 @@ public struct StorageSequence: Sendable, AsyncSequence { } } - private struct ErrorErasingSequence: + private struct FailureErasingAsyncSequence: AsyncSequence, Sendable where Base.Element == ByteBuffer { @@ -90,7 +90,7 @@ public struct StorageSequence: Sendable, AsyncSequence { self.length = length self.makeIterator = { AsyncIterator( - base: ErrorErasingSequence(base: asyncSequence) + base: FailureErasingAsyncSequence(base: asyncSequence) .makeAsyncIterator() ) }