9 Commits
6 changed files with 83 additions and 84 deletions
+11 -6
View File
@@ -1,5 +1,5 @@
# pyinotifyd
A daemon to monitor filesystems events with inotify on Linux and run tasks like filesystem operations (copy, move or delete), a shell command or custom async python methods.
A daemon for monitoring filesystem events with inotify on Linux and run tasks like filesystem operations (copy, move or delete), a shell command or custom async python methods.
It is possible to schedule tasks with a delay, which can then be canceled again in case a canceling event occurs. A useful example for this is to run tasks only if a file has not changed within a certain amount of time.
@@ -49,7 +49,7 @@ The basic idea is to instantiate one or multiple schedulers and map specific ino
pyinotifyd has different schedulers to schedule tasks with an optional delay. The advantages of using a scheduler are consistent logging and the possibility to cancel delayed tasks. Furthermore, schedulers have the ability to differentiate between files and directories.
### TaskScheduler
Schedule a custom python method *job* with an optional *delay* in seconds. Skip scheduling of tasks for files and/or directories according to *files* and *dirs* arguments. If there already is a scheduled task, re-schedule it with *delay*. Use *logname* in log messages. All additional modules, functions and variables that are defined in the config file and are needed within the *job*, need to be passed as dictionary to the TaskManager through *global_vars*.
Schedule a custom python method *job* with an optional *delay* in seconds. Skip scheduling of tasks for files and/or directories according to *files* and *dirs* arguments. If there already is a scheduled task, re-schedule it with *delay*. Use *logname* in log messages. All additional modules, functions and variables that are defined in the config file and are needed within the *job*, need to be passed as dictionary to the TaskManager through *global_vars*. If you want to limit the scheduler to run only one job at a time, set *singlejob* to True.
All arguments except for *job* are optional.
```python
# Please note that pyinotifyd uses pythons asyncio for asynchronous task execution.
@@ -71,7 +71,8 @@ task_sched = TaskScheduler(
dirs=False,
delay=0,
logname="sched",
global_vars=globals())
global_vars=globals(),
singlejob=False)
```
### ShellScheduler
@@ -164,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)
```
@@ -214,6 +217,8 @@ enableSyslog(lglevel=INFO, name="daemon")
## Schedule python method for all events on files and directories
```python
import logging
async def custom_job(event, task_id):
logging.info(f"{task_id}: execute example task: {event}")
@@ -1,47 +0,0 @@
# 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
SCM=""
if [ "${PV#9999}" != "${PV}" ] ; then
SCM="git-r3"
EGIT_REPO_URI="https://github.com/spacefreak86/${PN}"
fi
inherit ${SCM} distutils-r1 systemd
DESCRIPTION="Monitore filesystems events and execute Python methods or Shell commands."
HOMEPAGE="https://github.com/spacefreak86/pymodmilter"
if [ "${PV#9999}" != "${PV}" ] ; then
SRC_URI=""
KEYWORDS=""
# Needed for tests
S="${WORKDIR}/${PN}"
EGIT_CHECKOUT_DIR="${S}"
else
SRC_URI="https://github.com/spacefreak86/${PN}/archive/${PV}.tar.gz -> ${P}.tar.gz"
KEYWORDS="amd64 x86"
fi
LICENSE="GPL-3"
SLOT="0"
IUSE="systemd"
RDEPEND="dev-python/pyinotify[${PYTHON_USEDEP}]"
python_install_all() {
distutils-r1_python_install_all
dodir /etc/${PN}
insinto /etc/${PN}
newins ${PN}/misc/config.py.default config.py
use systemd && systemd_dounit ${PN}/misc/${PN}.service
newinitd ${PN}/misc/openrc/${PN}.initd ${PN}
newconfd ${PN}/misc/openrc/${PN}.confd ${PN}
}
@@ -2,7 +2,7 @@
# Distributed under the terms of the GNU General Public License v2
EAPI=7
PYTHON_COMPAT=( python3_{7,8,9} )
PYTHON_COMPAT=( python3_{8..10} )
DISTUTILS_USE_SETUPTOOLS=rdepend
SCM=""
Executable → Regular
+43 -13
View File
@@ -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.6"
__version__ = "0.0.8"
def setLoglevel(loglevel, logname=None):
@@ -76,9 +76,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, loop=None,
logname="eventmap"):
self._map = {}
self._exclude_filter = None
self._loop = (loop or asyncio.get_event_loop())
if default_sched is not None:
@@ -91,6 +92,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):
@@ -114,20 +116,37 @@ class EventMap(ProcessEvent):
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:
self._map[maskname].process_event(event)
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):
schedulers = []
@@ -139,21 +158,31 @@ 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", loop=None):
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
@@ -175,6 +204,7 @@ class Watch:
loop = (loop or self._loop)
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(
+8 -5
View File
@@ -6,15 +6,16 @@
#import logging
#
#async def custom_job(event, task_id):
# asyncio.sleep(1)
# await asyncio.sleep(1)
# logging.info(f"{task_id}: execute example task: {event}")
#
#task_sched = TaskScheduler(
# job=custom_job,
# files=True,
# dirs=False,
# delay=10
# global_vars=globals())
# delay=10,
# global_vars=globals(),
# singlejob=False)
###########################
@@ -25,7 +26,8 @@
# cmd="/usr/local/bin/task.sh {maskname} {pathname} {src_pathname}",
# files=True,
# dirs=False,
# delay=10)
# delay=10,
# singlejob=False)
#################################
@@ -86,7 +88,8 @@
# path="/watched/directory",
# event_map = event_map,
# rec=True,
# auto_add=True)
# auto_add=True,
# exclude_filter=["^/watched/directory/subpath$"])
################
+20 -12
View File
@@ -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={}):
loop=None, global_vars={}, singlejob=False):
assert iscoroutinefunction(job), \
f"job: expected coroutine, got {type(job)}"
assert isinstance(files, bool), \
@@ -71,6 +71,7 @@ class TaskScheduler:
self._log = logging.getLogger((logname or __name__))
self._loop = (loop or asyncio.get_event_loop())
self._globals = global_vars
self._singlejob = singlejob
self._tasks = {}
self._pause = False
@@ -105,6 +106,9 @@ class TaskScheduler:
else:
self._log.info("all remainig tasks completed")
def taskindex(self, event):
return "singlejob" if self._singlejob else event.pathname
async def _run_job(self, event, task_state, restart=False):
logger = SchedulerLogger(self._log, {
"event": event,
@@ -112,7 +116,7 @@ class TaskScheduler:
if self._delay > 0:
task_state.task = self._loop.create_task(
asyncio.sleep(self._delay, loop=self._loop))
asyncio.sleep(self._delay))
try:
if restart:
prefix = "re-"
@@ -145,7 +149,8 @@ class TaskScheduler:
else:
logger.info("task finished")
finally:
del self._tasks[event.pathname]
task_index = self.taskindex(event)
del self._tasks[task_index]
async def process_event(self, event):
if not ((not event.dir and self._files) or
@@ -153,9 +158,13 @@ class TaskScheduler:
return
restart = False
task_index = self.taskindex(event)
try:
task_state = self._tasks[event.pathname]
task_state = self._tasks[task_index]
except KeyError:
task_state = TaskScheduler.TaskState()
self._tasks[task_index] = task_state
else:
logger = SchedulerLogger(self._log, {
"event": event,
"id": task_state.id})
@@ -171,16 +180,13 @@ class TaskScheduler:
logger.warning("skip event due to ongoing task")
return
except KeyError:
task_state = TaskScheduler.TaskState()
self._tasks[event.pathname] = task_state
if not self._pause:
await self._run_job(event, task_state, restart)
async def process_cancel_event(self, event):
try:
task_state = self._tasks[event.pathname]
task_index = self.taskindex(event)
task_state = self._tasks[task_index]
except KeyError:
return
@@ -192,7 +198,8 @@ class TaskScheduler:
task_state.task.cancel()
logger.info("scheduled task cancelled")
task_state.task = None
del self._tasks[event.pathname]
logger.info(f"{task_index}")
del self._tasks[task_index]
else:
logger.warning("skip event due to ongoing task")
@@ -285,7 +292,8 @@ class FileManagerRule:
class FileManagerScheduler(TaskScheduler):
def __init__(self, rules, job=None, *args, **kwargs):
super().__init__(*args, **kwargs, job=self._manager_job)
super().__init__(
*args, **kwargs, job=self._manager_job, singlejob=False)
if not isinstance(rules, list):
rules = [rules]