TeraRuntimeBoundedTask.swift (6220B)
1 import Foundation 2 3 enum TeraRuntimeBoundedOutcome<Value: Sendable>: Sendable { 4 case completed(Result<Value, TeraRuntimeFailure>) 5 case timedOut 6 case cancelled 7 } 8 9 /// The wrapper has one immutable state owner; every mutable field is protected 10 /// by its lock. Task cancellation and continuation resumption occur outside it. 11 final class TeraRuntimeBoundedTask<Value: Sendable>: @unchecked Sendable { 12 typealias Outcome = TeraRuntimeBoundedOutcome<Value> 13 14 final class State: @unchecked Sendable { 15 private let lock = NSLock() 16 private var outcome: Outcome? 17 private var continuations: [UUID: CheckedContinuation<Outcome, Never>] = [:] 18 private var cancelledWaiters: Set<UUID> = [] 19 private var operationTask: Task<Void, Never>? 20 private var timeoutTask: Task<Void, Never>? 21 private var settledResult: Result<Value, TeraRuntimeFailure>? 22 private var settlementWaiters: [CheckedContinuation<Result<Value, TeraRuntimeFailure>, Never>] = [] 23 24 func install(operationTask: Task<Void, Never>, timeoutTask: Task<Void, Never>) { 25 let terminal = lock.withLock { () -> Outcome? in 26 if let outcome { 27 return outcome 28 } 29 self.operationTask = operationTask 30 self.timeoutTask = timeoutTask 31 return nil 32 } 33 guard let terminal else { return } 34 timeoutTask.cancel() 35 switch terminal { 36 case .cancelled, .timedOut: 37 operationTask.cancel() 38 case .completed: 39 break 40 } 41 } 42 43 func value(cancelsOperationWhenWaiterCancelled: Bool) async -> Outcome { 44 let waiterID = UUID() 45 return await withTaskCancellationHandler { 46 await withCheckedContinuation { requestedContinuation in 47 let immediate = lock.withLock { () -> Outcome? in 48 if let outcome { 49 return outcome 50 } 51 if cancelledWaiters.remove(waiterID) != nil { 52 return .cancelled 53 } 54 continuations[waiterID] = requestedContinuation 55 return nil 56 } 57 if let immediate { 58 requestedContinuation.resume(returning: immediate) 59 } 60 } 61 } onCancel: { 62 if cancelsOperationWhenWaiterCancelled { 63 cancel() 64 } else { 65 cancelWaiter(waiterID) 66 } 67 } 68 } 69 70 func cancel() { 71 guard resolve(.cancelled) else { return } 72 let tasks = lock.withLock { (operationTask, timeoutTask) } 73 tasks.0?.cancel() 74 tasks.1?.cancel() 75 } 76 77 func expire() { 78 guard resolve(.timedOut) else { return } 79 let taskToCancel: Task<Void, Never>? = lock.withLock { self.operationTask } 80 taskToCancel?.cancel() 81 } 82 83 func finishOperation(_ result: Result<Value, TeraRuntimeFailure>) { 84 let waiters = lock.withLock { 85 operationTask = nil 86 settledResult = result 87 let pending = settlementWaiters 88 settlementWaiters.removeAll() 89 return pending 90 } 91 for waiter in waiters { 92 waiter.resume(returning: result) 93 } 94 } 95 96 func settlement() -> Result<Value, TeraRuntimeFailure>? { 97 lock.withLock { settledResult } 98 } 99 100 func settle() async -> Result<Value, TeraRuntimeFailure> { 101 await withCheckedContinuation { continuation in 102 let immediate = lock.withLock { () -> Result<Value, TeraRuntimeFailure>? in 103 if let settledResult { 104 return settledResult 105 } 106 settlementWaiters.append(continuation) 107 return nil 108 } 109 if let immediate { 110 continuation.resume(returning: immediate) 111 } 112 } 113 } 114 115 private func cancelWaiter(_ waiterID: UUID) { 116 let continuation = lock.withLock { () -> CheckedContinuation<Outcome, Never>? in 117 guard outcome == nil else { return nil } 118 guard let continuation = continuations.removeValue(forKey: waiterID) else { 119 cancelledWaiters.insert(waiterID) 120 return nil 121 } 122 return continuation 123 } 124 continuation?.resume(returning: .cancelled) 125 } 126 127 @discardableResult 128 func resolve(_ requestedOutcome: Outcome) -> Bool { 129 var pendingContinuations: [CheckedContinuation<Outcome, Never>] = [] 130 var timeoutToCancel: Task<Void, Never>? 131 let accepted = lock.withLock { () -> Bool in 132 guard outcome == nil else { return false } 133 outcome = requestedOutcome 134 pendingContinuations = Array(continuations.values) 135 continuations.removeAll(keepingCapacity: false) 136 cancelledWaiters.removeAll(keepingCapacity: false) 137 timeoutToCancel = timeoutTask 138 return true 139 } 140 if accepted { 141 timeoutToCancel?.cancel() 142 for continuation in pendingContinuations { 143 continuation.resume(returning: requestedOutcome) 144 } 145 } 146 return accepted 147 } 148 } 149 150 private let state: State 151 152 init( 153 deadlineNanoseconds: UInt64, 154 operation: @escaping @Sendable () async -> Result<Value, TeraRuntimeFailure>, 155 onAbandonedResult: 156 @escaping @Sendable ( 157 Result<Value, TeraRuntimeFailure> 158 ) async -> Void = { _ in } 159 ) { 160 let state = State() 161 self.state = state 162 let operationTask = Task { [state] in 163 let result = await operation() 164 if !state.resolve(.completed(result)) { 165 await onAbandonedResult(result) 166 } 167 state.finishOperation(result) 168 } 169 170 let timeoutTask = Task { [state] in 171 do { 172 try await Task.sleep(nanoseconds: deadlineNanoseconds) 173 } catch { 174 return 175 } 176 state.expire() 177 } 178 state.install(operationTask: operationTask, timeoutTask: timeoutTask) 179 } 180 181 func value(cancelsOperationWhenWaiterCancelled: Bool = true) async -> Outcome { 182 await state.value( 183 cancelsOperationWhenWaiterCancelled: cancelsOperationWhenWaiterCancelled 184 ) 185 } 186 187 func cancel() { 188 state.cancel() 189 } 190 191 /// Only the shutdown owner awaits actual completion. A caller deadline or 192 /// cancellation resolves `value`, but cannot prove that the operation ended. 193 func settle() async -> Result<Value, TeraRuntimeFailure> { 194 await state.settle() 195 } 196 197 func settlement() -> Result<Value, TeraRuntimeFailure>? { 198 state.settlement() 199 } 200 }