Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions Sources/CodableDatastore/Datastore/AsyncInstances.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 } }
}
Expand All @@ -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<IdentifierType>,_)`` 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 {
Expand All @@ -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<IdentifierType>,_)`` 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)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -128,9 +128,11 @@ struct AsyncThrowingBackpressureStream<Element: Sendable>: 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:)") {
Expand All @@ -153,9 +155,9 @@ struct AsyncThrowingBackpressureStream<Element: Sendable>: 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)
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -224,6 +224,9 @@ actor MultiplexedAsyncSequence<Base: AsyncSequence & Sendable>: AsyncSequence wh
}

extension RangeReplaceableCollection where Self: Sendable {
#if compiler(>=6.2)
@concurrent
#endif
init<S: AsyncSequence>(_ sequence: sending S) async throws where S.Element == Element {
self = try await sequence.reduce(into: Self.init()) { @Sendable partialResult, element in
partialResult.append(element)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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]))
Expand Down
Loading