@@ -42,6 +42,27 @@ final class CoreAIPipelinedEngine: InferenceEngine, Sendable {
4242 private let generationTask = Mutex < Task < Void , Never > ? > ( nil )
4343 let config : ModelConfig
4444
45+ // MARK: - Lifecycle
46+
47+ var isBusy : Bool { engineInUse. load ( ordering: . acquiring) }
48+
49+ func cancel( ) async throws {
50+ let task = generationTask. withLock { v in
51+ v? . cancel ( )
52+ return v
53+ }
54+ await task? . value
55+ await engine. computeStream. currentWorkCompleted ( )
56+ generationTask. withLock { $0 = nil }
57+ }
58+
59+ func reset( ) async throws {
60+ try await cancel ( )
61+ guard tryAcquireEngine ( ) else { return }
62+ defer { releaseEngine ( ) }
63+ engine. reset ( )
64+ }
65+
4566 init (
4667 config: ModelConfig ,
4768 preparedModel: PreparedModel ,
@@ -144,37 +165,6 @@ final class CoreAIPipelinedEngine: InferenceEngine, Sendable {
144165 return GenerationSequence ( base: base, stopReasonStore: stopReasonStore)
145166 }
146167
147- var isBusy : Bool { engineInUse. load ( ordering: . acquiring) }
148-
149- func cancel( ) async {
150- let task = generationTask. withLock { v in
151- v? . cancel ( )
152- return v
153- }
154- await task? . value
155- await engine. computeStream. currentWorkCompleted ( )
156- generationTask. withLock { $0 = nil }
157- }
158-
159- /// Wait for any in-flight generate() Task to return the engine.
160- private func drain( ) {
161- var attempts = 0
162- while engineInUse. load ( ordering: . acquiring) {
163- attempts += 1
164- if attempts > Self . drainTimeoutIterations {
165- fatalError ( " Engine not returned after drain() — tokenSequence Task stuck? " )
166- }
167- Thread . sleep ( forTimeInterval: 0.001 )
168- }
169- }
170-
171- func reset( ) {
172- drain ( )
173- guard tryAcquireEngine ( ) else { return }
174- defer { releaseEngine ( ) }
175- engine. reset ( )
176- }
177-
178168 func cleanup( ) async throws {
179169 let cleanupSpan = InstrumentsProfiler . beginCleanup ( engine: " CoreAI-Pipelined " )
180170 if tryAcquireEngine ( ) {
0 commit comments