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/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) } } } 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) 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 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]))