The stream's element is now the enum InotifyEvent, and the struct that describes a change to a watched item is FileSystemEvent. A queue overflow was an event with descriptor -1 and an empty path that every consumer had to know about; as a case, the compiler makes them handle it. The enum is also where failed watches of a growing tree will be reported, since no call site can catch them.
220 lines
7.7 KiB
Swift
220 lines
7.7 KiB
Swift
import Dispatch
|
|
import CInotify
|
|
import SystemPackage
|
|
|
|
public actor Inotify {
|
|
private let fd: CInt
|
|
private var exclusions = ExclusionList()
|
|
private var watches = InotifyWatchManager()
|
|
private nonisolated(unsafe) let eventReader: any DispatchSourceRead
|
|
private nonisolated let eventStream: AsyncStream<RawInotifyEvent>
|
|
private nonisolated let continuation: AsyncStream<RawInotifyEvent>.Continuation
|
|
public nonisolated var events: AsyncCompactMapSequence<AsyncStream<RawInotifyEvent>, InotifyEvent> {
|
|
self.eventStream.compactMap(self.transform(_:))
|
|
}
|
|
|
|
/// Creates an inotify instance.
|
|
///
|
|
/// Events are read from the kernel as soon as they arrive and buffered
|
|
/// until they are consumed from ``events``.
|
|
///
|
|
/// - Parameter bufferingPolicy: How events are kept while no consumer is
|
|
/// reading ``events``. The default `.unbounded` keeps every event, so a
|
|
/// burst of changes is never lost; a bounded policy trades memory for
|
|
/// dropped events.
|
|
public init(bufferingPolicy: AsyncStream<RawInotifyEvent>.Continuation.BufferingPolicy = .unbounded) throws {
|
|
self.fd = inotify_init1(CInt(IN_NONBLOCK | IN_CLOEXEC))
|
|
guard self.fd >= 0 else {
|
|
throw InotifyError.initFailed(errno: cinotify_get_errno())
|
|
}
|
|
(self.eventReader, self.eventStream, self.continuation) = Self.createEventReader(
|
|
forFileDescriptor: fd,
|
|
bufferingPolicy: bufferingPolicy
|
|
)
|
|
}
|
|
|
|
/// Whether an item with this name is skipped, by an excluded name or
|
|
/// an excluded pattern.
|
|
public func isExcluded(_ name: String) -> Bool {
|
|
self.exclusions.excludes(name)
|
|
}
|
|
|
|
public func exclude(name: String) {
|
|
self.exclusions.add(name: name)
|
|
}
|
|
|
|
public func exclude(names: String...) {
|
|
self.exclude(names: names)
|
|
}
|
|
|
|
public func exclude(names: [String]) {
|
|
for name in names {
|
|
self.exclusions.add(name: name)
|
|
}
|
|
}
|
|
|
|
/// Excludes every item whose name matches a shell pattern such as
|
|
/// `*.tmp` or `@*`, with the same effect as an excluded name.
|
|
///
|
|
/// The pattern is matched against the item's own name, not its path,
|
|
/// as the shell matches file names: `*` and `?` stand for any
|
|
/// characters and `[…]` for a set of characters. A leading dot needs
|
|
/// no special treatment, so `.*` excludes hidden items.
|
|
public func exclude(pattern: String) {
|
|
self.exclusions.add(pattern: pattern)
|
|
}
|
|
|
|
public func exclude(patterns: String...) {
|
|
self.exclude(patterns: patterns)
|
|
}
|
|
|
|
public func exclude(patterns: [String]) {
|
|
for pattern in patterns {
|
|
self.exclusions.add(pattern: pattern)
|
|
}
|
|
}
|
|
|
|
@discardableResult
|
|
public func addWatch(path: String, mask: InotifyEventMask) throws -> CInt {
|
|
let wd = inotify_add_watch(self.fd, path, mask.rawValue)
|
|
guard wd >= 0 else {
|
|
throw InotifyError.addWatchFailed(path: path, errno: cinotify_get_errno())
|
|
}
|
|
watches.add(path, withId: wd, mask: mask)
|
|
return wd
|
|
}
|
|
|
|
@discardableResult
|
|
public func addRecursiveWatch(forDirectory path: String, mask: InotifyEventMask) async throws -> [CInt] {
|
|
let directoryPaths = try await DirectoryResolver.resolve([path], excluding: self.exclusions)
|
|
var result: [CInt] = []
|
|
for path in directoryPaths {
|
|
let wd = try self.addWatch(path: path.string, mask: mask)
|
|
result.append(wd)
|
|
}
|
|
return result
|
|
}
|
|
|
|
@discardableResult
|
|
public func addWatchWithAutomaticSubtreeWatching(forDirectory path: String, mask: InotifyEventMask) async throws -> [CInt] {
|
|
let wds = try await self.addRecursiveWatch(forDirectory: path, mask: mask)
|
|
watches.enableAutomaticSubtreeWatching(forIds: wds)
|
|
return wds
|
|
}
|
|
|
|
public func removeWatch(_ wd: CInt) throws {
|
|
guard inotify_rm_watch(self.fd, wd) == 0 else {
|
|
throw InotifyError.removeWatchFailed(watchDescriptor: wd, errno: cinotify_get_errno())
|
|
}
|
|
watches.remove(forId: wd)
|
|
}
|
|
|
|
deinit {
|
|
// The file descriptor is closed by the reader's cancel handler once
|
|
// libdispatch has unregistered it. Closing it here would leave a
|
|
// registration behind that a later instance reusing the descriptor
|
|
// number could inherit, silently losing its events.
|
|
self.eventReader.cancel()
|
|
}
|
|
|
|
private func transform(_ rawEvent: RawInotifyEvent) async -> InotifyEvent? {
|
|
if rawEvent.mask.contains(.queueOverflow) {
|
|
return .queueOverflow
|
|
}
|
|
guard let path = self.watches.path(forId: rawEvent.watchDescriptor) else { return nil }
|
|
guard !self.exclusions.excludes(rawEvent.name) else { return nil }
|
|
let event = FileSystemEvent(from: rawEvent, inDirectory: path)
|
|
self.forgetWatchInCaseTheKernelRemovedIt(event)
|
|
self.removeWatchesInCaseADirectoryLeftTheTree(event)
|
|
await self.addWatchInCaseOfAutomaticSubtreeWatching(event)
|
|
return .fileSystem(event)
|
|
}
|
|
|
|
/// The kernel reports `IN_IGNORED` once a watch is gone, whether it was
|
|
/// removed explicitly or because its item was deleted or unmounted.
|
|
/// Forgetting it keeps a reused descriptor number from mapping to a
|
|
/// stale path.
|
|
private func forgetWatchInCaseTheKernelRemovedIt(_ event: FileSystemEvent) {
|
|
guard event.mask.contains(.ignored) else { return }
|
|
self.watches.remove(forId: event.watchDescriptor)
|
|
}
|
|
|
|
/// A directory moved out of a watched tree keeps its kernel watches,
|
|
/// which would then report events under the old path. Those watches
|
|
/// are removed instead.
|
|
private func removeWatchesInCaseADirectoryLeftTheTree(_ event: FileSystemEvent) {
|
|
guard event.mask.contains(.movedFrom), event.mask.contains(.isDir) else { return }
|
|
for wd in self.watches.descriptors(under: event.path.string) {
|
|
inotify_rm_watch(self.fd, wd)
|
|
self.watches.remove(forId: wd)
|
|
}
|
|
}
|
|
|
|
private func addWatchInCaseOfAutomaticSubtreeWatching(_ event: FileSystemEvent) async {
|
|
guard !event.synthesized,
|
|
watches.isAutomaticSubtreeWatching(event.watchDescriptor),
|
|
event.mask.contains(.isDir),
|
|
let kind = Self.subtreeTrigger(in: event.mask) else {
|
|
return
|
|
}
|
|
|
|
guard let mask = self.watches.mask(forId: event.watchDescriptor) else { return }
|
|
guard let wds = try? await self.addWatchWithAutomaticSubtreeWatching(forDirectory: event.path.string, mask: mask) else { return }
|
|
await self.synthesizeEvents(forContentOfWatches: wds, kind: kind, cookie: event.cookie)
|
|
}
|
|
|
|
private static func subtreeTrigger(in mask: InotifyEventMask) -> InotifyEventMask? {
|
|
if mask.contains(.create) { return .create }
|
|
if mask.contains(.movedTo) { return .movedTo }
|
|
return nil
|
|
}
|
|
|
|
/// Items that already exist when a directory becomes watched never
|
|
/// produce kernel events, so they are reported as if they had just
|
|
/// appeared, marked as synthesized.
|
|
private func synthesizeEvents(forContentOfWatches wds: [CInt], kind: InotifyEventMask, cookie: UInt32) async {
|
|
for wd in wds {
|
|
guard let directory = self.watches.path(forId: wd) else { continue }
|
|
guard let entries = try? await DirectoryResolver.entries(of: FilePath(directory), excluding: self.exclusions) else { continue }
|
|
for entry in entries {
|
|
let mask: InotifyEventMask = entry.isDirectory ? [kind, .isDir] : kind
|
|
self.continuation.yield(RawInotifyEvent(
|
|
watchDescriptor: wd,
|
|
mask: mask,
|
|
cookie: cookie,
|
|
name: entry.name,
|
|
synthesized: true
|
|
))
|
|
}
|
|
}
|
|
}
|
|
|
|
private static func createEventReader(
|
|
forFileDescriptor fd: CInt,
|
|
bufferingPolicy: AsyncStream<RawInotifyEvent>.Continuation.BufferingPolicy
|
|
) -> (any DispatchSourceRead, AsyncStream<RawInotifyEvent>, AsyncStream<RawInotifyEvent>.Continuation) {
|
|
let (stream, continuation) = AsyncStream<RawInotifyEvent>.makeStream(
|
|
of: RawInotifyEvent.self,
|
|
bufferingPolicy: bufferingPolicy
|
|
)
|
|
|
|
let reader = DispatchSource.makeReadSource(
|
|
fileDescriptor: fd,
|
|
queue: DispatchQueue(label: "Inotify.read", qos: .utility)
|
|
)
|
|
|
|
reader.setEventHandler {
|
|
for rawEvent in InotifyEventParser.parse(fromFileDescriptor: fd) {
|
|
continuation.yield(rawEvent)
|
|
}
|
|
}
|
|
reader.setCancelHandler {
|
|
cinotify_deinit(fd)
|
|
continuation.finish()
|
|
}
|
|
reader.activate()
|
|
|
|
return (reader, stream, continuation)
|
|
}
|
|
}
|