Я реализовал большую часть датчика, но по какой-то причине при запуске наблюдателя WatchDog -- по сути, это просто потоковая обработка. Thread — процесс зависает.
Код: Выделить всё
"""StackStorm sensor to to watch for file changes"""
import json
from typing import TypedDict, Callable, Any
from logging import Logger
import re
from watchdog.observers import Observer
from watchdog.observers.api import BaseObserver
from watchdog.events import FileSystemEventHandler, FileSystemEvent
import eventlet
from st2reactor.sensor.base import Sensor
import time
class _FileWatchTriggerParams(TypedDict):
watch_directory: str
filename_regex: str
recursive: bool
class _FileWatchTrigger(TypedDict):
uid: str
id: str
ref: str
name: str
pack: str
type: str
parameters: _FileWatchTriggerParams
class _FileWatcherHandler(FileSystemEventHandler):
dispatch_trigger: Callable[[str, Any], None]
file_regex: re.Pattern | None
def __init__(
self,
dispatch_trigger: Callable[[str, Any], None],
trigger_ref: str,
file_regex: str,
logger: Logger,
):
self.dispatch_trigger = dispatch_trigger
self.trigger_ref = trigger_ref
self.logger = logger
if file_regex and file_regex != ".*":
self.file_regex = re.compile(file_regex)
else:
self.file_regex = None
def on_any_event(self, event: FileSystemEvent) -> None:
self.logger.info(
"%s event occurred on file %s", event.event_type.title(), event.src_path
)
if event.is_directory:
return
self.logger.info(
"Trigger %s Evaluating event on file %s", self.trigger_ref, event.src_path
)
if self.file_regex is not None:
if event.dest_path and not self.file_regex.match(event.dest_path):
return
if not self.file_regex.match(event.src_path):
return
payload = {
"event_type": event.event_type,
"src_path": event.src_path,
"dest_path": event.dest_path,
"watch_path": self.file_regex,
}
self.logger.info(
"Trigger %s dispatching event %s",
self.trigger_ref,
json.dumps(payload, indent=2),
)
self.dispatch_trigger(self.trigger_ref, payload)
class FileWatchSensor(Sensor):
"""Sensor allowing for configurable file watchers using the Watchdog library"""
_logger: Logger
_reload_needed: bool # whether the observer needs to be reloaded.
_triggers: dict[str, _FileWatchTrigger] # the triggers that are being watched
_observer: BaseObserver
def __init__(self, sensor_service, config=None):
super().__init__(sensor_service=sensor_service, config=config)
self._logger = self.sensor_service.get_logger(type(self).__qualname__)
self._reload_needed = False
self._observer = None
self._triggers = {}
def setup(self):
pass
def run(self):
while True:
try:
self._logger.info("Polling...")
if self._reload_needed:
self.reload_observer()
eventlet.sleep(1)
time.sleep(1)
except Exception: # pylint: disable=broad-exception-caught
self._logger.exception(
"Unexpected exception while running %s", type(self).__name__
)
def cleanup(self):
if self._observer:
self._observer.stop()
self._observer.join()
def add_trigger(self, trigger: _FileWatchTrigger):
try:
self._logger.info(
"Adding watch on dir [%s] for file [%s]",
trigger["parameters"].get("watch_directory", "No Watch Directory"),
trigger["parameters"].get("filename_regex", "No File Pattern"),
)
self._triggers[trigger["ref"]] = trigger
self._reload_needed = True
except Exception: # pylint: disable=broad-exception-caught
self._logger.exception(
"Unexpected exception while adding trigger: %s",
json.dumps(trigger, indent=2),
)
def update_trigger(self, trigger: _FileWatchTrigger):
try:
if trigger["ref"] not in self._triggers:
self._triggers[trigger["ref"]] = trigger
else:
self._triggers[trigger["ref"]].update(trigger)
except Exception: # pylint: disable=broad-exception-caught
self._logger.exception(
"Unexpected exception while updating trigger: %s",
json.dumps(trigger, indent=2),
)
def remove_trigger(self, trigger: _FileWatchTrigger):
try:
if trigger["ref"] in self._triggers:
removed_trigger = self._triggers.pop(trigger["ref"])
self._logger.info(
"Removing trigger watch %s on dir [%s] for file [%s]",
trigger["ref"],
removed_trigger["parameters"].get(
"watch_directory", "No Watch Directory"
),
removed_trigger["parameters"].get("filename_regex", ""),
)
self._reload_needed = True
except Exception: # pylint: disable=broad-exception-caught
self._logger.exception(
"Unexcpected exception while removing trigger: %s",
json.dumps(trigger, indent=2),
)
def reload_observer(self):
"""Build a new observer with the current set of watches"""
# The observer objects are just threads.
# Threads can only be started once.
# Schedules cannot be modified on started threads.
old_observer = None
if self._observer:
old_observer = self._observer
self._observer = Observer()
watches = []
for trigger in self._triggers.values():
watch_dir = trigger["parameters"].get("watch_directory")
if not watch_dir:
self._logger.error(
"Trigger missing `watch_directory` parameter: %s for ",
json.dumps(trigger, indent=2),
)
continue
pattern = trigger["parameters"].get("filename_regex")
event_handler = _FileWatcherHandler(
self.sensor_service.dispatch,
trigger["ref"],
pattern,
self._logger,
)
recursive = trigger["parameters"].get("recursive", False)
self._observer.schedule(
event_handler,
watch_dir,
recursive,
)
watches.append(
(watch_dir, "recursive" if recursive else "standard", pattern or ".*")
)
if old_observer is not None:
self._logger.info("Stopping old observer")
old_observer.stop()
old_observer.join()
self._logger.info("Old observer stopped")
self._logger.info("Starting observer with %d watches", len(watches))
self._observer.start() # Main process gets stuck right here
self._logger.info("Observer thread started.")
self._reload_needed = False
self._logger.info("Watching files: %s", json.dumps(watches, indent=2))
При реализации этого я обнаружил, что после запуска запланированные часы наблюдателя нельзя изменить. Как и в текущем состоянии, попытка сделать это приведет к блокировке основного потока на неопределенный срок.
Я запускаю это в StackStorm 3.8.0 с использованием Python 3.11.3
Вот логи, которые генерирует датчик при попытке запуска:
Код: Выделить всё
2024-07-31 21:58:51 2024-08-01 02:58:51,369 INFO [-] Sensor ava_core.FileWatchSensor updated. Reloading sensor.
2024-07-31 21:58:52 2024-08-01 02:58:52,387 INFO [-] Sensor ava_core.FileWatchSensor reloaded.
2024-07-31 21:58:53 2024-08-01 02:58:53,200 INFO [-] No config found for sensor "FileWatchSensor"
2024-07-31 21:58:53 2024-08-01 02:58:53,201 INFO [-] Watcher started
2024-07-31 21:58:53 2024-08-01 02:58:53,201 INFO [-] Running sensor initialization code
2024-07-31 21:58:53 2024-08-01 02:58:53,201 INFO [-] Running sensor in passive mode
2024-07-31 21:58:53 2024-08-01 02:58:53,201 INFO [-] Polling...
2024-07-31 21:58:53 2024-08-01 02:58:53,231 INFO [-] Adding watch on dir [/appl/appcloud] for file [.*\\.txt]
2024-07-31 21:58:53 2024-08-01 02:58:53,237 INFO [-] Connected to amqp://guest:**@rabbitmq:5672//
2024-07-31 21:58:55 2024-08-01 02:58:55,204 INFO [-] Polling...
2024-07-31 21:58:55 2024-08-01 02:58:55,205 INFO [-] Starting observer with 1 watches
Подробнее здесь: https://stackoverflow.com/questions/788 ... fter-start