Merge pull request 'WP-5b: FSEvents ReturnWatcher' (#7) from wp5b/return-watcher-20260831 into feat/shotdeck-20260830
This commit was merged in pull request #7.
This commit is contained in:
@@ -0,0 +1,154 @@
|
||||
import Foundation
|
||||
import PDFKit
|
||||
import CoreServices // FSEventStream* APIs; system framework, no Package.swift change needed
|
||||
|
||||
public actor ReturnWatcher {
|
||||
private let ledger: ReturnLedger
|
||||
private var watchFolder: URL
|
||||
private var onChange: (@Sendable ([ReturnedDocument]) -> Void)?
|
||||
private var stream: FSEventStreamRef?
|
||||
private var bridge: FSEventBridge?
|
||||
private var pendingScanTask: Task<Void, Never>?
|
||||
private let eventQueue = DispatchQueue(label: "ai.flowmaster.shotdeck.returns.fsevents")
|
||||
|
||||
/// Watch folder is `paths.watchFolder`, which production constructs from
|
||||
/// `FolderSettings.resolve().watch`. This type never calls FolderSettings;
|
||||
/// `updateWatchFolder` is invoked by the UI layer only.
|
||||
public init(paths: AppSupportPaths, ledger: ReturnLedger) {
|
||||
self.ledger = ledger
|
||||
self.watchFolder = paths.watchFolder // never a literal "~/Downloads" here
|
||||
}
|
||||
|
||||
/// Starts watching paths.watchFolder for returned PDFs. Performs one immediate
|
||||
/// scanNow() before returning, then calls onChange after every subsequent debounced
|
||||
/// batch (even if that batch's result is empty — the caller decides what to do).
|
||||
public func start(onChange: @escaping @Sendable ([ReturnedDocument]) -> Void) async throws {
|
||||
self.onChange = onChange
|
||||
try startStream(on: watchFolder)
|
||||
let found = try await scanNow()
|
||||
onChange(found)
|
||||
}
|
||||
|
||||
/// Idempotent. Stops and releases the FSEventStream if one is running; safe to call
|
||||
/// when never started or already stopped. Cancels any pending debounced scan.
|
||||
public func stop() {
|
||||
// Idempotent: nil stream / already-stopped is a no-op; never started is the same.
|
||||
pendingScanTask?.cancel()
|
||||
pendingScanTask = nil
|
||||
if let stream {
|
||||
FSEventStreamStop(stream)
|
||||
FSEventStreamInvalidate(stream)
|
||||
FSEventStreamRelease(stream)
|
||||
}
|
||||
stream = nil
|
||||
bridge = nil
|
||||
}
|
||||
|
||||
/// Called by whoever owns the Settings "Choose..." folder action (WP-4c) after the user
|
||||
/// picks a new watch folder. If the watcher was running, stops the old FSEventStream,
|
||||
/// switches to the new folder, restarts, and performs one immediate scanNow (reporting
|
||||
/// through the same onChange callback given to start()). If the watcher was never
|
||||
/// started, only updates the stored folder for the next start() call.
|
||||
public func updateWatchFolder(_ url: URL) async throws {
|
||||
let wasRunning = stream != nil
|
||||
stop()
|
||||
watchFolder = url
|
||||
guard wasRunning else { return }
|
||||
try startStream(on: url)
|
||||
let found = try await scanNow()
|
||||
onChange?(found)
|
||||
}
|
||||
|
||||
/// Scans the watch folder once, immediately, without waiting for an event. Every
|
||||
/// recognized, stable, openable Shotdeck PDF present is (re-)inspected and (re-)recorded
|
||||
/// into the ledger; returns exactly the documents processed in this call.
|
||||
@discardableResult
|
||||
public func scanNow() async throws -> [ReturnedDocument] {
|
||||
let fm = FileManager.default
|
||||
let candidates = (try? fm.contentsOfDirectory(
|
||||
at: watchFolder, includingPropertiesForKeys: nil
|
||||
)) ?? []
|
||||
var results: [ReturnedDocument] = []
|
||||
for url in candidates.sorted(by: { $0.lastPathComponent < $1.lastPathComponent }) {
|
||||
let name = url.lastPathComponent
|
||||
// ".pdf.inprogress" already fails hasSuffix(".pdf") -> naturally skipped.
|
||||
guard name.hasSuffix(".pdf"), !name.hasPrefix(".") else { continue }
|
||||
guard await isStableAndReadable(url) else { continue } // leave for next event
|
||||
guard let document = PDFDocument(url: url),
|
||||
AnnotationInspector.isShotdeckDocument(document) else { continue }
|
||||
guard let inspected = try? AnnotationInspector.inspect(fileURL: url) else { continue }
|
||||
try await ledger.record(inspected)
|
||||
results.append(inspected)
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
/// Size-stable: two equal byte counts 250 ms apart AND PDFDocument opens;
|
||||
/// otherwise leave the file for the next event.
|
||||
private func isStableAndReadable(_ url: URL) async -> Bool {
|
||||
let fm = FileManager.default
|
||||
guard let size1 = try? fm.attributesOfItem(atPath: url.path)[.size] as? Int else { return false }
|
||||
try? await Task.sleep(for: .milliseconds(250))
|
||||
guard let size2 = try? fm.attributesOfItem(atPath: url.path)[.size] as? Int else { return false }
|
||||
guard size1 == size2, size1 > 0 else { return false }
|
||||
return PDFDocument(url: url) != nil
|
||||
}
|
||||
|
||||
private func startStream(on folder: URL) throws {
|
||||
let bridge = FSEventBridge { [weak self] in
|
||||
guard let self else { return }
|
||||
Task { await self.scheduleDebouncedScan() }
|
||||
}
|
||||
self.bridge = bridge
|
||||
var context = FSEventStreamContext()
|
||||
context.version = 0
|
||||
context.info = Unmanaged.passUnretained(bridge).toOpaque()
|
||||
context.retain = nil
|
||||
context.release = nil
|
||||
context.copyDescription = nil
|
||||
guard let stream = FSEventStreamCreate(
|
||||
kCFAllocatorDefault, shotdeckFSEventsCallback, &context,
|
||||
[folder.path] as CFArray, FSEventStreamEventId(kFSEventStreamEventIdSinceNow),
|
||||
0.0,
|
||||
FSEventStreamCreateFlags(kFSEventStreamCreateFlagFileEvents | kFSEventStreamCreateFlagNoDefer)
|
||||
) else {
|
||||
throw ShotdeckError.captureFailed(underlying: "could not create FSEventStream for \(folder.path)")
|
||||
}
|
||||
// Dispatch queue, not a run loop: this actor has no run loop of its own, and
|
||||
// FSEventStreamSetDispatchQueue is the modern replacement for
|
||||
// FSEventStreamScheduleWithRunLoop. One dedicated serial queue per watcher.
|
||||
FSEventStreamSetDispatchQueue(stream, eventQueue)
|
||||
guard FSEventStreamStart(stream) else {
|
||||
FSEventStreamInvalidate(stream)
|
||||
FSEventStreamRelease(stream)
|
||||
throw ShotdeckError.captureFailed(underlying: "FSEventStreamStart failed for \(folder.path)")
|
||||
}
|
||||
self.stream = stream
|
||||
}
|
||||
|
||||
private func scheduleDebouncedScan() async {
|
||||
pendingScanTask?.cancel()
|
||||
pendingScanTask = Task {
|
||||
try? await Task.sleep(for: .milliseconds(400)) // coalesce AirDrop's write+rename burst
|
||||
guard !Task.isCancelled else { return }
|
||||
guard let found = try? await self.scanNow() else { return }
|
||||
// One in-flight debounced callback may land after stop(); it is a
|
||||
// harmless read-only rescan (ledger upsert, no watch-folder mutation).
|
||||
self.onChange?(found)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Non-actor bridge because FSEventStreamCallback is a @convention(c) function pointer and
|
||||
/// cannot capture actor-isolated state directly; it hops back onto the actor via Task.
|
||||
private final class FSEventBridge: @unchecked Sendable {
|
||||
// @unchecked is safe: `notify` is a let, set once at init, never mutated after — the
|
||||
// type is immutable for its entire lifetime.
|
||||
let notify: @Sendable () -> Void
|
||||
init(notify: @escaping @Sendable () -> Void) { self.notify = notify }
|
||||
}
|
||||
|
||||
private let shotdeckFSEventsCallback: FSEventStreamCallback = { _, info, _, _, _, _ in
|
||||
guard let info else { return }
|
||||
Unmanaged<FSEventBridge>.fromOpaque(info).takeUnretainedValue().notify()
|
||||
}
|
||||
Reference in New Issue
Block a user