Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
961c64e422
|
||
|
|
c22bd73759
|
||
|
|
0ac52b23d4
|
||
|
|
88ce35930c
|
||
|
|
01ccd1817d
|
||
|
|
b80d71e95e
|
@@ -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
|
||||
@@ -214,6 +215,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=""
|
||||
|
||||
@@ -36,7 +36,7 @@ from pyinotify import ProcessEvent
|
||||
from pyinotifyd._install import install, uninstall
|
||||
from pyinotifyd.scheduler import TaskScheduler, Cancel
|
||||
|
||||
__version__ = "0.0.6"
|
||||
__version__ = "0.0.7"
|
||||
|
||||
|
||||
def setLoglevel(loglevel, logname=None):
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
#################################
|
||||
|
||||
+20
-12
@@ -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]
|
||||
|
||||
Reference in New Issue
Block a user