Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
05eb46a0bd
|
||
|
|
53eb8f48ed
|
||
|
|
5556902e54
|
||
|
|
5a019f1352
|
||
|
|
56b19120f4
|
||
|
|
546b8c4a2d
|
||
|
|
8d8e08d94f
|
||
|
|
0ea2dadea4
|
||
|
|
006d631682
|
||
|
|
669129d919
|
@@ -165,21 +165,23 @@ pyinotifyd = Pyinotifyd(
|
||||
```
|
||||
|
||||
### Watches
|
||||
A watch connects the *path* to an *event_map*. Automatically add a watch on each sub-directories in *path* if *rec* is set to True. If *auto_add* is True, a watch will be added automatically on newly created sub-directories in *path*.
|
||||
A watch connects the *path* to an *event_map*. Automatically add a watch on each sub-directories in *path* if *rec* is set to True. If *auto_add* is True, a watch will be added automatically on newly created sub-directories in *path*. All events for paths matching one of the regular expressions in *exclude_filter* are ignored. If the value of *exclude_filter* is a string, it is assumed to be a path to a file from which the list of regular expressions will be read.
|
||||
```python
|
||||
# Add a watch directly to Pyinotifyd.
|
||||
pyinotifyd.add_watch(
|
||||
path="/src_path",
|
||||
event_map=event_map,
|
||||
rec=False,
|
||||
auto_add=False)
|
||||
auto_add=False,
|
||||
exclude_filter=["^/src_path/subpath$"])
|
||||
|
||||
# Or instantiate and add it
|
||||
w = Watch(
|
||||
path="/src_path",
|
||||
event_map=event_map,
|
||||
rec=False,
|
||||
auto_add=False)
|
||||
auto_add=False,
|
||||
exclude_filter=["^/src_path/subpath$"])
|
||||
|
||||
pyinotifyd.add_watch(watch=w)
|
||||
```
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
# Copyright 2020 Gentoo Authors
|
||||
# Distributed under the terms of the GNU General Public License v2
|
||||
|
||||
EAPI=7
|
||||
PYTHON_COMPAT=( python3_{8..10} )
|
||||
DISTUTILS_USE_SETUPTOOLS=rdepend
|
||||
EAPI=8
|
||||
PYTHON_COMPAT=( python3_{13..14} )
|
||||
DISTUTILS_USE_PEP517=setuptools
|
||||
|
||||
SCM=""
|
||||
if [ "${PV#9999}" != "${PV}" ] ; then
|
||||
|
||||
Executable → Regular
+55
-31
@@ -31,12 +31,12 @@ import pyinotify
|
||||
import signal
|
||||
import sys
|
||||
|
||||
from pyinotify import ProcessEvent
|
||||
from pyinotify import ProcessEvent, ExcludeFilter
|
||||
|
||||
from pyinotifyd._install import install, uninstall
|
||||
from pyinotifyd.scheduler import TaskScheduler, Cancel
|
||||
|
||||
__version__ = "0.0.7"
|
||||
__version__ = "1.0.0"
|
||||
|
||||
|
||||
def setLoglevel(loglevel, logname=None):
|
||||
@@ -56,16 +56,15 @@ def enableSyslog(loglevel=None, address="/dev/log", logname=None):
|
||||
|
||||
|
||||
class _SchedulerList:
|
||||
def __init__(self, schedulers=[], loop=None):
|
||||
def __init__(self, schedulers=[]):
|
||||
if not isinstance(schedulers, list):
|
||||
schedulers = [schedulers]
|
||||
|
||||
self._schedulers = schedulers
|
||||
self._loop = (loop or asyncio.get_event_loop())
|
||||
|
||||
def process_event(self, event):
|
||||
for scheduler in self._schedulers:
|
||||
self._loop.create_task(scheduler.process_event(event))
|
||||
asyncio.create_task(scheduler.process_event(event))
|
||||
|
||||
def schedulers(self):
|
||||
return self._schedulers
|
||||
@@ -76,10 +75,10 @@ class EventMap(ProcessEvent):
|
||||
**pyinotify.EventsCodes.OP_FLAGS,
|
||||
**pyinotify.EventsCodes.EVENT_FLAGS}
|
||||
|
||||
def my_init(self, event_map=None, default_sched=None, loop=None,
|
||||
def my_init(self, event_map=None, default_sched=None, exclude_filter=None,
|
||||
logname="eventmap"):
|
||||
self._map = {}
|
||||
self._loop = (loop or asyncio.get_event_loop())
|
||||
self._exclude_filter = None
|
||||
|
||||
if default_sched is not None:
|
||||
for flag in EventMap.flags:
|
||||
@@ -91,6 +90,7 @@ class EventMap(ProcessEvent):
|
||||
for flag, schedulers in event_map.items():
|
||||
self.set_scheduler(flag, schedulers)
|
||||
|
||||
self.set_exclude_filter(exclude_filter)
|
||||
self._log = logging.getLogger((logname or __name__))
|
||||
|
||||
def set_scheduler(self, flag, schedulers):
|
||||
@@ -106,27 +106,43 @@ class EventMap(ProcessEvent):
|
||||
isinstance(scheduler, Cancel):
|
||||
instances.append(scheduler)
|
||||
else:
|
||||
instances.append(
|
||||
TaskScheduler(scheduler, loop=self._loop))
|
||||
instances.append(TaskScheduler(scheduler))
|
||||
|
||||
self._map[flag] = _SchedulerList(instances, loop=self._loop)
|
||||
self._map[flag] = _SchedulerList(instances)
|
||||
|
||||
elif flag in self._map:
|
||||
del self._map[flag]
|
||||
|
||||
def set_exclude_filter(self, exclude_filter):
|
||||
if exclude_filter is None:
|
||||
self._exclude_filter = None
|
||||
return
|
||||
|
||||
if not isinstance(exclude_filter, ExcludeFilter):
|
||||
self._exclude_filter = ExcludeFilter(exclude_filter)
|
||||
else:
|
||||
self._exclude_filter = exclude_filter
|
||||
|
||||
def process_default(self, event):
|
||||
msg = "received event"
|
||||
attrs = ""
|
||||
for attr in [
|
||||
"dir", "mask", "maskname", "pathname", "src_pathname", "wd"]:
|
||||
value = getattr(event, attr, None)
|
||||
if attr == "mask":
|
||||
value = hex(value)
|
||||
if value:
|
||||
msg += f", {attr}={value}"
|
||||
attrs += f", {attr}={value}"
|
||||
|
||||
self._log.debug(msg)
|
||||
self._log.debug(f"received event{attrs}")
|
||||
maskname = event.maskname.split("|")[0]
|
||||
if maskname in self._map:
|
||||
|
||||
if maskname not in self._map:
|
||||
return
|
||||
|
||||
if self._exclude_filter and self._exclude_filter(event.pathname):
|
||||
self._log.debug(f"pathname {event.pathname} is excluded")
|
||||
return
|
||||
|
||||
self._map[maskname].process_event(event)
|
||||
|
||||
def schedulers(self):
|
||||
@@ -139,23 +155,32 @@ class EventMap(ProcessEvent):
|
||||
|
||||
|
||||
class Watch:
|
||||
def __init__(self, path, event_map=None, default_sched=None, rec=False,
|
||||
auto_add=False, logname="watch", loop=None):
|
||||
assert isinstance(path, str), \
|
||||
f"path: expected {type('')}, got {type(path)}"
|
||||
def __init__(self, path, event_map=None, default_sched=None,
|
||||
rec=False, auto_add=False, exclude_filter=None,
|
||||
logname="watch"):
|
||||
assert (isinstance(path, str) or isinstance(path, list)), \
|
||||
f"path: expected {type('')} or {type([])}, got {type(path)}"
|
||||
|
||||
if isinstance(event_map, EventMap):
|
||||
self._event_map = event_map
|
||||
else:
|
||||
self._event_map = EventMap(
|
||||
event_map=event_map, default_sched=default_sched)
|
||||
event_map=event_map, default_sched=default_sched,
|
||||
exclude_filter=exclude_filter)
|
||||
|
||||
assert isinstance(rec, bool), \
|
||||
f"rec: expected {type(bool)}, got {type(rec)}"
|
||||
assert isinstance(auto_add, bool), \
|
||||
f"auto_add: expected {type(bool)}, got {type(auto_add)}"
|
||||
|
||||
self._exclude_filter = None
|
||||
if exclude_filter:
|
||||
if not isinstance(exclude_filter, ExcludeFilter):
|
||||
self._exclude_filter = ExcludeFilter(exclude_filter)
|
||||
else:
|
||||
self._exclude_filter = exclude_filter
|
||||
|
||||
logname = (logname or __name__)
|
||||
self._loop = loop
|
||||
|
||||
self._path = path
|
||||
self._rec = rec
|
||||
@@ -171,14 +196,16 @@ class Watch:
|
||||
def event_map(self):
|
||||
return self._event_map
|
||||
|
||||
def start(self, loop=None):
|
||||
loop = (loop or self._loop)
|
||||
def start(self):
|
||||
self._watch_manager.add_watch(self._path, pyinotify.ALL_EVENTS,
|
||||
rec=self._rec, auto_add=self._auto_add,
|
||||
exclude_filter=self._exclude_filter,
|
||||
do_glob=True)
|
||||
|
||||
self._notifier = pyinotify.AsyncioNotifier(
|
||||
self._watch_manager, loop, default_proc_fun=self._event_map)
|
||||
self._watch_manager,
|
||||
asyncio.get_event_loop(),
|
||||
default_proc_fun=self._event_map)
|
||||
|
||||
def stop(self):
|
||||
self._notifier.stop()
|
||||
@@ -189,12 +216,10 @@ class Watch:
|
||||
class Pyinotifyd:
|
||||
name = "pyinotifyd"
|
||||
|
||||
def __init__(self, watches=[], shutdown_timeout=30, logname="daemon",
|
||||
loop=None):
|
||||
def __init__(self, watches=[], shutdown_timeout=30, logname="daemon"):
|
||||
self.set_watches(watches)
|
||||
self.set_shutdown_timeout(shutdown_timeout)
|
||||
logname = (logname or __name__)
|
||||
self._loop = (loop or asyncio.get_event_loop())
|
||||
|
||||
self._log = logging.getLogger(logname)
|
||||
|
||||
@@ -246,9 +271,7 @@ class Pyinotifyd:
|
||||
schedulers.extend(w.event_map().schedulers())
|
||||
return list(set(schedulers))
|
||||
|
||||
def start(self, loop=None):
|
||||
loop = (loop or self._loop)
|
||||
|
||||
def start(self):
|
||||
if len(self._watches) == 0:
|
||||
self._log.warning(
|
||||
"no watches configured, the daemon will not do anything")
|
||||
@@ -256,7 +279,7 @@ class Pyinotifyd:
|
||||
for watch in self._watches:
|
||||
self._log.info(
|
||||
f"start listening for inotify events on '{watch.path()}'")
|
||||
watch.start(loop)
|
||||
watch.start()
|
||||
|
||||
def pause(self):
|
||||
for scheduler in self.schedulers():
|
||||
@@ -435,7 +458,7 @@ def main():
|
||||
f"%(asctime)s - {name}/%(name)s - %(levelname)s - %(message)s")
|
||||
ch.setFormatter(formatter)
|
||||
|
||||
loop = asyncio.get_event_loop()
|
||||
loop = asyncio.new_event_loop()
|
||||
loop.add_signal_handler(
|
||||
signal.SIGTERM, lambda: loop.create_task(
|
||||
daemon.shutdown("SIGTERM")))
|
||||
@@ -445,6 +468,7 @@ def main():
|
||||
loop.add_signal_handler(
|
||||
signal.SIGHUP, lambda: loop.create_task(
|
||||
daemon.reload("SIGHUP", args.config, args.debug)))
|
||||
asyncio.set_event_loop(loop)
|
||||
|
||||
daemon.start()
|
||||
loop.run_forever()
|
||||
|
||||
@@ -37,7 +37,7 @@
|
||||
#move_rule = Rule(
|
||||
# action="move",
|
||||
# src_re="^/src_path/(?P<path>.*).to_move",
|
||||
# dst_re="/dst_path/\g<path>.moved",
|
||||
# dst_re="/dst_path/\\g<path>.moved",
|
||||
# auto_create=True,
|
||||
# filemode=0o755,
|
||||
# dirmode=0o644,
|
||||
@@ -88,7 +88,8 @@
|
||||
# path="/watched/directory",
|
||||
# event_map = event_map,
|
||||
# rec=True,
|
||||
# auto_add=True)
|
||||
# auto_add=True,
|
||||
# exclude_filter=["^/watched/directory/subpath$"])
|
||||
|
||||
|
||||
################
|
||||
|
||||
Executable → Regular
+7
-9
@@ -52,7 +52,7 @@ class TaskScheduler:
|
||||
self.cancelable = cancelable
|
||||
|
||||
def __init__(self, job, files=True, dirs=False, delay=0, logname="sched",
|
||||
loop=None, global_vars={}, singlejob=False):
|
||||
global_vars={}, singlejob=False):
|
||||
assert iscoroutinefunction(job), \
|
||||
f"job: expected coroutine, got {type(job)}"
|
||||
assert isinstance(files, bool), \
|
||||
@@ -69,7 +69,6 @@ class TaskScheduler:
|
||||
self._dirs = dirs
|
||||
self._delay = delay
|
||||
self._log = logging.getLogger((logname or __name__))
|
||||
self._loop = (loop or asyncio.get_event_loop())
|
||||
self._globals = global_vars
|
||||
self._singlejob = singlejob
|
||||
self._tasks = {}
|
||||
@@ -91,8 +90,7 @@ class TaskScheduler:
|
||||
self._log.info(
|
||||
f"wait {timeout} seconds for {len(pending)} "
|
||||
f"remaining task(s) to complete")
|
||||
done, pending = await asyncio.wait([*pending], timeout=timeout,
|
||||
loop=self._loop)
|
||||
done, pending = await asyncio.wait([*pending], timeout=timeout)
|
||||
if pending:
|
||||
self._log.warning(
|
||||
f"shutdown timeout exceeded, "
|
||||
@@ -100,7 +98,7 @@ class TaskScheduler:
|
||||
for task in pending:
|
||||
task.cancel()
|
||||
try:
|
||||
await asyncio.gather(*pending, loop=self._loop)
|
||||
await asyncio.gather(*pending)
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
else:
|
||||
@@ -115,7 +113,7 @@ class TaskScheduler:
|
||||
"id": task_state.id})
|
||||
|
||||
if self._delay > 0:
|
||||
task_state.task = self._loop.create_task(
|
||||
task_state.task = asyncio.create_task(
|
||||
asyncio.sleep(self._delay))
|
||||
try:
|
||||
if restart:
|
||||
@@ -134,11 +132,11 @@ class TaskScheduler:
|
||||
local_vars = {"self": self,
|
||||
"event": event,
|
||||
"task_id": task_state.id}
|
||||
task_state.task = self._loop.create_task(
|
||||
task_state.task = asyncio.create_task(
|
||||
eval("self._job(event, task_id)", self._globals, local_vars))
|
||||
|
||||
else:
|
||||
task_state.task = self._loop.create_task(
|
||||
task_state.task = asyncio.create_task(
|
||||
self._job(event, task_state.id))
|
||||
|
||||
try:
|
||||
@@ -244,7 +242,7 @@ class ShellScheduler(TaskScheduler):
|
||||
|
||||
logger.info(f"execute shell command, cmd={cmd}")
|
||||
try:
|
||||
proc = await asyncio.create_subprocess_shell(cmd, loop=self._loop)
|
||||
proc = await asyncio.create_subprocess_shell(cmd)
|
||||
await proc.communicate()
|
||||
except Exception as e:
|
||||
logger.error(e)
|
||||
|
||||
Reference in New Issue
Block a user