From 6523d829d8451452d665297cee96ca11d4267c19 Mon Sep 17 00:00:00 2001 From: Dimitri Bouniol Date: Thu, 24 Sep 2026 03:15:54 -0700 Subject: [PATCH 1/4] Fixed an issue where the backpressure stream task didn't explicitly handle all thrown errors --- .../AsyncThrowingBackpressureStream.swift | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/Sources/CodableDatastore/Persistence/Disk Persistence/AsyncThrowingBackpressureStream.swift b/Sources/CodableDatastore/Persistence/Disk Persistence/AsyncThrowingBackpressureStream.swift index 3b092c8..ebda064 100644 --- a/Sources/CodableDatastore/Persistence/Disk Persistence/AsyncThrowingBackpressureStream.swift +++ b/Sources/CodableDatastore/Persistence/Disk Persistence/AsyncThrowingBackpressureStream.swift @@ -128,9 +128,11 @@ struct AsyncThrowingBackpressureStream: Sendable { } } - fileprivate func finish(throwing error: (any Error)? = nil) async throws { - try await withCheckedThrowingContinuation { continuation in - guard let stateMachine else { continuation.resume(throwing: CancellationError()) + fileprivate func finish(throwing error: (any Error)? = nil) async { + /// Do nothing if the continuation fails, as the stream was already consumed or cancelled. + try? await withCheckedThrowingContinuation { continuation in + guard let stateMachine else { + continuation.resume(throwing: CancellationError()) return } Task(name: "CodableDatastore.AsyncThrowingBackpressureStream.Continuation.finish(throwing:)") { @@ -153,9 +155,9 @@ struct AsyncThrowingBackpressureStream: Sendable { Task(name: "CodableDatastore.AsyncThrowingBackpressureStream.init(provider:)") { do { try await provider(continuation) - try await continuation.finish() + await continuation.finish() } catch { - try await continuation.finish(throwing: error) + await continuation.finish(throwing: error) } } } From dfb64a9ec3672a4ac6b5e398ab8d6a68c0a786b0 Mon Sep 17 00:00:00 2001 From: Dimitri Bouniol Date: Thu, 24 Sep 2026 03:20:19 -0700 Subject: [PATCH 2/4] Fixed warnings using async sequences by marking the related methods as concurrent --- Sources/CodableDatastore/Datastore/AsyncInstances.swift | 9 +++++++++ .../Disk Persistence/Datastore/DatastorePage.swift | 3 +++ 2 files changed, 12 insertions(+) diff --git a/Sources/CodableDatastore/Datastore/AsyncInstances.swift b/Sources/CodableDatastore/Datastore/AsyncInstances.swift index 4988d21..50b2bc4 100644 --- a/Sources/CodableDatastore/Datastore/AsyncInstances.swift +++ b/Sources/CodableDatastore/Datastore/AsyncInstances.swift @@ -13,6 +13,9 @@ extension AsyncInstances { /// Returns the first instance of the sequence, if it exists. /// /// - Returns: The first instance of the sequence, or `nil` if the sequence is empty. + #if compiler(>=6.2) + @concurrent + #endif public var firstInstance: Element? { get async throws { try await first { _ in true } } } @@ -24,6 +27,9 @@ extension AsyncInstances { /// - Warning: This method is only safe to use from sequences vended by a Datastore ranged ``Datastore/load(range:order:)-(IndexRangeExpression,_)`` operation as they guarantee that the returned sequence won't stall due to unavailable instances. Do not use it when collecting observations as there is no guarantee observations will be returned! /// - Parameter collectionLimit: The maximum amount of entries to collect. Specify `.infinity` to _questionably_ collect all instances. /// - Returns: An array of instances up to the collection limit. + #if compiler(>=6.2) + @concurrent + #endif public func collectInstances(upTo collectionLimit: Int) async throws -> [Element] { var instances: [Element] = [] for try await instance in self { @@ -42,6 +48,9 @@ extension AsyncInstances { /// - Warning: This method is only safe to use from sequences vended by a Datastore ranged ``Datastore/load(range:order:)-(IndexRangeExpression,_)`` operation as they guarantee that the returned sequence won't stall due to unavailable instances. Do not use it when collecting observations as there is no guarantee observations will be returned! /// - Parameter collectionLimit: The maximum amount of entries to collect. Specify `.infinity` to _questionably_ collect all instances. /// - Returns: An array of instances up to the collection limit. + #if compiler(>=6.2) + @concurrent + #endif public func collectInstances(upTo collectionLimit: AsyncInstancesLimit) async throws -> [Element] { try await collectInstances(upTo: .max) } diff --git a/Sources/CodableDatastore/Persistence/Disk Persistence/Datastore/DatastorePage.swift b/Sources/CodableDatastore/Persistence/Disk Persistence/Datastore/DatastorePage.swift index b08b509..68eab26 100644 --- a/Sources/CodableDatastore/Persistence/Disk Persistence/Datastore/DatastorePage.swift +++ b/Sources/CodableDatastore/Persistence/Disk Persistence/Datastore/DatastorePage.swift @@ -224,6 +224,9 @@ actor MultiplexedAsyncSequence: AsyncSequence wh } extension RangeReplaceableCollection where Self: Sendable { + #if compiler(>=6.2) + @concurrent + #endif init(_ sequence: sending S) async throws where S.Element == Element { self = try await sequence.reduce(into: Self.init()) { @Sendable partialResult, element in partialResult.append(element) From 1f5c7ac9ae9f2987aaf0c89470cdeaf338ded77b Mon Sep 17 00:00:00 2001 From: Dimitri Bouniol Date: Thu, 24 Sep 2026 03:21:16 -0700 Subject: [PATCH 3/4] Fixed a warning surrounding collated writes, and added a note to revisit it --- .../Persistence/Disk Persistence/Transaction/Transaction.swift | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/Sources/CodableDatastore/Persistence/Disk Persistence/Transaction/Transaction.swift b/Sources/CodableDatastore/Persistence/Disk Persistence/Transaction/Transaction.swift index 6add351..154813a 100644 --- a/Sources/CodableDatastore/Persistence/Disk Persistence/Transaction/Transaction.swift +++ b/Sources/CodableDatastore/Persistence/Disk Persistence/Transaction/Transaction.swift @@ -108,7 +108,8 @@ extension DiskPersistence { if options.contains(.collateWrites) { /// If we are skipping immediate writes, kick off persistence in a separate task. Task(name: "CodableDatastore.DiskPersistence.Transaction.run()") { - try await self.persist() + // TODO: Rethink collated writes here… + try? await self.persist() } } else { /// If we don't care to collate our writes, go ahead and wait for the persistence to stick From bbf14ecadf134a55f217fec9a57e43a4b460a024 Mon Sep 17 00:00:00 2001 From: Dimitri Bouniol Date: Thu, 24 Sep 2026 03:36:54 -0700 Subject: [PATCH 4/4] Fixed tests that created unstructured tasks that threw errors --- .../DiskPersistenceDatastoreIndexTests.swift | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Tests/CodableDatastoreTests/DiskPersistenceDatastoreIndexTests.swift b/Tests/CodableDatastoreTests/DiskPersistenceDatastoreIndexTests.swift index a129a34..e7ec55e 100644 --- a/Tests/CodableDatastoreTests/DiskPersistenceDatastoreIndexTests.swift +++ b/Tests/CodableDatastoreTests/DiskPersistenceDatastoreIndexTests.swift @@ -423,7 +423,7 @@ final class DiskPersistenceDatastoreIndexTests: XCTestCase, @unchecked Sendable let exp = expectation(description: "Finished") Task { [pageInfos, pageLookup] in for _ in 0..<1000 { - _ = try await index.pageIndex(for: UInt64.random(in: 0..<1000000), in: pageInfos, requiresCompleteEntries: false) { pageID in + _ = try? await index.pageIndex(for: UInt64.random(in: 0..<1000000), in: pageInfos, requiresCompleteEntries: false) { pageID in pageLookup[pageID]! } comparator: { lhs, rhs in lhs.sortOrder(comparedTo: try UInt64(bigEndianBytes: rhs.headers[0]))