Commit e0bf6830ae
Verified · cmc
Layout: unified · split
Sources/OrgWorkspace/FSEventsWatcher.swift added +78
| @@ -0,0 +1,78 @@ | ||
| 1 | #if os(macOS) | |
| 2 | import CoreServices | |
| 3 | import Foundation | |
| 4 | ||
| 5 | public enum WatchEvent: Sendable, Equatable { | |
| 6 | case changed(URL) | |
| 7 | /// Events were dropped or coalesced for this folder, or a root moved: rescan it. | |
| 8 | case rescan(URL) | |
| 9 | } | |
| 10 | ||
| 11 | /// FSEvents with file-level events. Delivers batches on a private queue. | |
| 12 | public final class FSEventsWatcher: @unchecked Sendable { | |
| 13 | private let paths: [String] | |
| 14 | private let latency: CFTimeInterval | |
| 15 | private let handler: @Sendable ([WatchEvent]) -> Void | |
| 16 | private let queue = DispatchQueue(label: "orgstar.fsevents") | |
| 17 | private var stream: FSEventStreamRef? | |
| 18 | ||
| 19 | public init(roots: [URL], latency: CFTimeInterval = 0.3, handler: @escaping @Sendable ([WatchEvent]) -> Void) { | |
| 20 | paths = roots.map(\.path) | |
| 21 | self.latency = latency | |
| 22 | self.handler = handler | |
| 23 | } | |
| 24 | ||
| 25 | deinit { | |
| 26 | stop() | |
| 27 | } | |
| 28 | ||
| 29 | public func start() { | |
| 30 | guard stream == nil else { return } | |
| 31 | var context = FSEventStreamContext( | |
| 32 | version: 0, info: Unmanaged.passUnretained(self).toOpaque(), retain: nil, release: nil, copyDescription: nil | |
| 33 | ) | |
| 34 | let flags = FSEventStreamCreateFlags( | |
| 35 | kFSEventStreamCreateFlagFileEvents | kFSEventStreamCreateFlagUseCFTypes | |
| 36 | | kFSEventStreamCreateFlagNoDefer | kFSEventStreamCreateFlagWatchRoot | |
| 37 | ) | |
| 38 | let callback: FSEventStreamCallback = { _, info, count, eventPaths, eventFlags, _ in | |
| 39 | guard let info else { return } | |
| 40 | let watcher = Unmanaged<FSEventsWatcher>.fromOpaque(info).takeUnretainedValue() | |
| 41 | let paths = unsafeBitCast(eventPaths, to: NSArray.self) as? [String] ?? [] | |
| 42 | watcher.deliver(paths: paths, flags: Array(UnsafeBufferPointer(start: eventFlags, count: count))) | |
| 43 | } | |
| 44 | guard let stream = FSEventStreamCreate( | |
| 45 | nil, callback, &context, paths as CFArray, FSEventStreamEventId(kFSEventStreamEventIdSinceNow), latency, flags | |
| 46 | ) else { return } | |
| 47 | FSEventStreamSetDispatchQueue(stream, queue) | |
| 48 | FSEventStreamStart(stream) | |
| 49 | self.stream = stream | |
| 50 | } | |
| 51 | ||
| 52 | public func stop() { | |
| 53 | guard let stream else { return } | |
| 54 | FSEventStreamStop(stream) | |
| 55 | FSEventStreamInvalidate(stream) | |
| 56 | FSEventStreamRelease(stream) | |
| 57 | self.stream = nil | |
| 58 | } | |
| 59 | ||
| 60 | private func deliver(paths eventPaths: [String], flags: [FSEventStreamEventFlags]) { | |
| 61 | let rescanFlags = FSEventStreamEventFlags( | |
| 62 | kFSEventStreamEventFlagMustScanSubDirs | kFSEventStreamEventFlagUserDropped | |
| 63 | | kFSEventStreamEventFlagKernelDropped | kFSEventStreamEventFlagRootChanged | |
| 64 | ) | |
| 65 | var events: [WatchEvent] = [] | |
| 66 | for (path, flag) in zip(eventPaths, flags) { | |
| 67 | let url = URL(fileURLWithPath: path).standardizedFileURL | |
| 68 | if flag & rescanFlags != 0 { | |
| 69 | let root = paths.first { path.hasPrefix($0) } ?? path | |
| 70 | events.append(.rescan(URL(fileURLWithPath: root).standardizedFileURL)) | |
| 71 | } else { | |
| 72 | events.append(.changed(url)) | |
| 73 | } | |
| 74 | } | |
| 75 | if !events.isEmpty { handler(events) } | |
| 76 | } | |
| 77 | } | |
| 78 | #endif | |
Tests/OrgWorkspaceTests/WatcherTests.swift added +38
| @@ -0,0 +1,38 @@ | ||
| 1 | import Foundation | |
| 2 | import OrgIndex | |
| 3 | import Testing | |
| 4 | @testable import OrgWorkspace | |
| 5 | ||
| 6 | struct WatcherTests { | |
| 7 | @Test func reportsAWrittenFile() throws { | |
| 8 | let folder = try Folder() | |
| 9 | let received = Received() | |
| 10 | let watcher = FSEventsWatcher(roots: [folder.url.resolvingSymlinksInPath()], latency: 0.05) { received.add($0) } | |
| 11 | watcher.start() | |
| 12 | defer { watcher.stop() } | |
| 13 | Thread.sleep(forTimeInterval: 0.2) | |
| 14 | try folder.write("a.org", "* a\n") | |
| 15 | #expect(received.wait(for: "a.org", timeout: 5)) | |
| 16 | } | |
| 17 | } | |
| 18 | ||
| 19 | final class Received: @unchecked Sendable { | |
| 20 | private let lock = NSLock() | |
| 21 | private var events: [WatchEvent] = [] | |
| 22 | ||
| 23 | func add(_ batch: [WatchEvent]) { | |
| 24 | lock.withLock { events += batch } | |
| 25 | } | |
| 26 | ||
| 27 | func wait(for name: String, timeout: TimeInterval) -> Bool { | |
| 28 | let deadline = Date(timeIntervalSinceNow: timeout) | |
| 29 | while Date() < deadline { | |
| 30 | let found = lock.withLock { | |
| 31 | events.contains { if case .changed(let url) = $0 { return url.lastPathComponent == name } else { return false } } | |
| 32 | } | |
| 33 | if found { return true } | |
| 34 | Thread.sleep(forTimeInterval: 0.05) | |
| 35 | } | |
| 36 | return false | |
| 37 | } | |
| 38 | } | |