Skip to content

Commit 972f8d5

Browse files
GordonBeemingclaudegitbutler-client
authored
Self-heal indexing so a stalled watcher can't silently stop it (#12)
* Self-heal indexing so a stalled watcher can't silently stop it New screenshots are indexed only via the live FSEvents->OCR pipeline, and rescanAll() fires solely from the manual Rescan Now button. When that live path stalls - a long-lived FSEvent stream on an iCloud Drive folder quietly stops delivering, or one file hangs OCR forever - indexing dies until the app is relaunched. It stayed dead for 8 days in the wild. Two complementary fixes: - OCR can no longer hang. recognize(at:) now wraps the decode and Vision call in a withTimeout watchdog (60s default) that throws OCRError.timedOut, so a single bad file can't wedge the drain queue on an await that never resumes. - The indexer runs a maintenance tick every 5 minutes that re-arms the FSEvents stream and rescans. It runs independently of the drain queue, so even a wedged queue recovers, and a dead watcher gets a fresh stream. Co-authored-by: Claude <noreply@anthropic.com> Co-authored-by: GitButler <gitbutler@gitbutler.com> * Address PR review: make the OCR watchdog actually free the caller The first cut of withTimeout raced the work against a sleeper inside a withThrowingTaskGroup. A structured group can't return from its scope until every child finishes, so when the OCR work hangs in a non-cancellable Vision/ImageIO call the group hangs with it - the watchdog never fires and the indexer can still wedge on the exact bad-file case it was meant to fix. Codex and Copilot both flagged this. Rework withTimeout to run the work in an unstructured task the caller never awaits directly, suspending instead on a continuation that whichever of {work, deadline} settles first resumes via a one-shot gate. A hung work task is orphaned rather than awaited, so the caller is always freed at the deadline. Added a regression test with a Thread.sleep operation that ignores cooperative cancellation. Also from review: - Guard periodicMaintenance with isMaintenanceRunning so a scan that outruns the 5-minute interval can't stack a second concurrent scan. - guard let self in the maintenance loop so the task exits if the actor is deallocated instead of looping forever. - Fix a comment typo (worked -> worker task). Co-authored-by: Claude <noreply@anthropic.com> Co-authored-by: GitButler <gitbutler@gitbutler.com> --------- Co-authored-by: Claude <noreply@anthropic.com> Co-authored-by: GitButler <gitbutler@gitbutler.com>
1 parent c852a35 commit 972f8d5

3 files changed

Lines changed: 186 additions & 10 deletions

File tree

‎Sources/VistaCore/Indexer.swift‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,22 @@ public actor Indexer {
5050
private var pendingIndex: Int = 0
5151
private var workTask: Task<Void, Never>?
5252
private var eventTask: Task<Void, Never>?
53+
private var maintenanceTask: Task<Void, Never>?
54+
// Guards against overlapping maintenance ticks. `periodicMaintenance()`
55+
// awaits `rescanAll()`, and the actor suspends across that await; if a scan
56+
// runs longer than `maintenanceInterval` (large backlog, slow disk/iCloud),
57+
// the next tick would otherwise start a second concurrent scan. The flag
58+
// makes a still-running scan swallow the next tick instead.
59+
private var isMaintenanceRunning = false
60+
61+
// How often the self-heal loop re-arms the watcher and rescans. New
62+
// files normally index instantly via FSEvents; this is the safety net
63+
// for when that live path silently stops — most notably long-lived
64+
// FSEvent streams on iCloud Drive folders, which have been observed to
65+
// stop delivering while the app keeps running. 5 minutes keeps the
66+
// worst-case staleness small without meaningful cost (a rescan with a
67+
// populated DB is just a fingerprint walk).
68+
private static let maintenanceInterval: Duration = .seconds(300)
5369

5470
private var progressContinuation: AsyncStream<Progress>.Continuation?
5571
// nonisolated because the stream is immutable after init and AsyncStream
@@ -98,16 +114,43 @@ public actor Indexer {
98114
}
99115

100116
await emitWatchingProgress()
117+
118+
maintenanceTask = Task { [weak self] in
119+
while !Task.isCancelled {
120+
try? await Task.sleep(for: Self.maintenanceInterval)
121+
// Bail if the actor was deallocated — `self?` would silently
122+
// no-op the call while the loop kept sleeping forever, leaking
123+
// the task.
124+
guard let self, !Task.isCancelled else { return }
125+
await self.periodicMaintenance()
126+
}
127+
}
101128
}
102129

103130
public func stop() {
104131
watcher.stop()
105132
eventTask?.cancel()
106133
workTask?.cancel()
134+
maintenanceTask?.cancel()
107135
progressContinuation?.finish()
108136
accessContinuation?.finish()
109137
}
110138

139+
/// Self-heal tick: re-arm the FSEvents stream (recovers a stream that
140+
/// silently stopped delivering) then rescan to catch anything the live
141+
/// path missed. Re-arm first so we're listening for new events before
142+
/// the rescan enumerates — anything that lands in the gap is still
143+
/// caught by the rescan, so there's no window where a file is lost.
144+
/// Skipped while paused, matching the intent of the pause control.
145+
private func periodicMaintenance() async {
146+
guard !isPaused, !isMaintenanceRunning else { return }
147+
isMaintenanceRunning = true
148+
defer { isMaintenanceRunning = false }
149+
watcher.stop()
150+
watcher.start(paths: watchedFolders)
151+
await rescanAll()
152+
}
153+
111154
public func setPaused(_ paused: Bool) {
112155
isPaused = paused
113156
if !paused {

‎Sources/VistaCore/OCRRecognizer.swift‎

Lines changed: 85 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -34,23 +34,38 @@ public final class OCRRecognizer: Sendable {
3434
self.languages = languages
3535
}
3636

37+
/// Upper bound on a single file's OCR. Vision has been observed to
38+
/// accept an image and then never fire its completion handler (nor
39+
/// throw), which would await forever; the ImageIO decode can likewise
40+
/// stall on a malformed iCloud file. A per-file timeout turns either
41+
/// hang into a thrown error the indexer can log and skip, so one bad
42+
/// file can't wedge the whole queue. 60s is far above real OCR cost
43+
/// (accurate mode on the largest screenshots is a few seconds) — it
44+
/// only ever trips on a genuine infinite hang.
45+
public static let defaultTimeout: Duration = .seconds(60)
46+
3747
/// Recognises text in the image at `url`. Returns a single joined
3848
/// string (lines separated by `\n`). Empty when OCR is off or when
3949
/// Vision found no text.
4050
///
41-
/// Throws on unreadable files or Vision failures — callers can decide
42-
/// whether to log and skip or surface to the user.
43-
public func recognize(at url: URL) async throws -> String {
51+
/// Throws on unreadable files, Vision failures, or when the work
52+
/// exceeds `timeout` — callers can decide whether to log and skip or
53+
/// surface to the user.
54+
public func recognize(at url: URL, timeout: Duration = OCRRecognizer.defaultTimeout) async throws -> String {
4455
guard level != .off else { return "" }
4556

46-
// Load once via ImageIO — avoids pulling the full decoded bitmap
47-
// into memory when Vision can stream from the CGImage source.
48-
guard let source = CGImageSourceCreateWithURL(url as CFURL, nil),
49-
let cgImage = CGImageSourceCreateImageAtIndex(source, 0, nil) else {
50-
throw OCRError.unreadableImage(url)
57+
// Both the decode and the Vision call run inside the timeout so a
58+
// stall in either one is bounded. The decode moves off the caller's
59+
// thread into the worker task as a side effect, which is harmless.
60+
return try await withTimeout(timeout, onTimeout: OCRError.timedOut(url)) {
61+
// Load once via ImageIO — avoids pulling the full decoded bitmap
62+
// into memory when Vision can stream from the CGImage source.
63+
guard let source = CGImageSourceCreateWithURL(url as CFURL, nil),
64+
let cgImage = CGImageSourceCreateImageAtIndex(source, 0, nil) else {
65+
throw OCRError.unreadableImage(url)
66+
}
67+
return try await self.recognize(cgImage: cgImage)
5168
}
52-
53-
return try await recognize(cgImage: cgImage)
5469
}
5570

5671
/// Recognises text in an already-decoded CGImage. Exposed so tests and
@@ -136,5 +151,65 @@ public final class OCRRecognizer: Sendable {
136151

137152
public enum OCRError: Error, Equatable {
138153
case unreadableImage(URL)
154+
case timedOut(URL)
155+
}
156+
}
157+
158+
/// Runs `operation`, throwing `onTimeout` if it hasn't finished within
159+
/// `timeout`. The caller is freed the instant either the work or the deadline
160+
/// settles, even when the work is stuck in a non-cancellable synchronous call
161+
/// (Vision's `perform`, an ImageIO decode) that ignores cooperative cancellation.
162+
///
163+
/// This deliberately does NOT use a `withThrowingTaskGroup`: a structured group
164+
/// cannot return from its scope until every child task has finished, so if the
165+
/// operation hangs the group hangs with it and the timeout never frees the
166+
/// caller — defeating the whole watchdog. Instead the work runs in an
167+
/// unstructured task that the caller never awaits directly. The caller suspends
168+
/// on a continuation that whichever of {work, deadline} finishes first resumes
169+
/// via a one-shot gate; a hung work task is simply orphaned (it finishes in the
170+
/// background and its late result is discarded). This bounds the *await*, not
171+
/// necessarily the CPU — the right trade for a watchdog.
172+
func withTimeout<T: Sendable>(
173+
_ timeout: Duration,
174+
onTimeout: @autoclosure @escaping @Sendable () -> Error,
175+
operation: @escaping @Sendable () async throws -> T
176+
) async throws -> T {
177+
return try await withCheckedThrowingContinuation { continuation in
178+
let gate = TimeoutGate(continuation)
179+
let deadline = Task {
180+
do {
181+
try await Task.sleep(for: timeout)
182+
gate.resume(with: .failure(onTimeout()))
183+
} catch {
184+
// Sleep cancelled because the work settled first — nothing to do.
185+
}
186+
}
187+
Task {
188+
do { gate.resume(with: .success(try await operation())) }
189+
catch { gate.resume(with: .failure(error)) }
190+
// Stop the watchdog as soon as the work settles so the common fast
191+
// path doesn't leave a task sleeping out the full timeout.
192+
deadline.cancel()
193+
}
194+
}
195+
}
196+
197+
/// One-shot resume gate for `withTimeout`: only the first of {work, deadline}
198+
/// resumes the caller; the loser's later attempt is dropped. The NSLock guards
199+
/// the single-resume invariant that `CheckedContinuation` asserts on. Declared
200+
/// at file scope because Swift can't nest a type inside a generic function.
201+
private final class TimeoutGate<T: Sendable>: @unchecked Sendable {
202+
private let lock = NSLock()
203+
private var resumed = false
204+
private let continuation: CheckedContinuation<T, Error>
205+
206+
init(_ continuation: CheckedContinuation<T, Error>) { self.continuation = continuation }
207+
208+
func resume(with result: Result<T, Error>) {
209+
lock.lock()
210+
let firstTime = !resumed
211+
resumed = true
212+
lock.unlock()
213+
if firstTime { continuation.resume(with: result) }
139214
}
140215
}

‎Tests/VistaCoreTests/OCRRecognizerTests.swift‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,64 @@ final class OCRRecognizerTests: XCTestCase {
3030
XCTAssertEqual(text, "")
3131
}
3232

33+
// MARK: - Timeout watchdog
34+
35+
private struct Boom: Error, Equatable {}
36+
37+
func testWithTimeoutReturnsFastResult() async throws {
38+
let value = try await withTimeout(.seconds(5), onTimeout: Boom()) {
39+
"done"
40+
}
41+
XCTAssertEqual(value, "done")
42+
}
43+
44+
func testWithTimeoutThrowsWhenWorkExceedsDeadline() async {
45+
do {
46+
_ = try await withTimeout(.milliseconds(20), onTimeout: Boom()) {
47+
// Far longer than the deadline; the watchdog should fire first.
48+
try await Task.sleep(for: .seconds(10))
49+
return "should not reach here"
50+
}
51+
XCTFail("expected the timeout to throw")
52+
} catch {
53+
XCTAssertEqual(error as? Boom, Boom())
54+
}
55+
}
56+
57+
func testWithTimeoutFreesCallerWhenWorkIgnoresCancellation() async {
58+
// The regression the codex/copilot review caught: a structured task
59+
// group won't return until every child finishes, so a work task stuck
60+
// in a non-cancellable synchronous call would hang the watchdog too.
61+
// Here the operation blocks the thread with Thread.sleep (which ignores
62+
// cooperative cancellation); the caller must still be freed at the
63+
// deadline rather than waiting the full 2s.
64+
let start = Date()
65+
do {
66+
_ = try await withTimeout(.milliseconds(50), onTimeout: Boom()) {
67+
Thread.sleep(forTimeInterval: 2)
68+
return "should not reach here"
69+
}
70+
XCTFail("expected the timeout to throw")
71+
} catch {
72+
XCTAssertEqual(error as? Boom, Boom())
73+
XCTAssertLessThan(Date().timeIntervalSince(start), 1.5,
74+
"caller should be freed at the deadline, not after the blocking work")
75+
}
76+
}
77+
78+
func testWithTimeoutPropagatesOperationError() async {
79+
do {
80+
_ = try await withTimeout(.seconds(5), onTimeout: Boom()) {
81+
throw OCRRecognizer.OCRError.unreadableImage(URL(fileURLWithPath: "/nope.png"))
82+
}
83+
XCTFail("expected the operation error to propagate")
84+
} catch let error as OCRRecognizer.OCRError {
85+
XCTAssertEqual(error, .unreadableImage(URL(fileURLWithPath: "/nope.png")))
86+
} catch {
87+
XCTFail("unexpected error: \(error)")
88+
}
89+
}
90+
3391
// MARK: - Fixture rendering
3492

3593
/// Draws black text on a white background. Large font so Vision has

0 commit comments

Comments
 (0)