Исходный код soniks_client.jobs.file_surveillance

"""Наблюдатель файловой системы: отправка декодированных кадров во время прохода."""

import threading
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from pathlib import Path

from apscheduler.schedulers.base import BaseScheduler
from watchdog.events import FileSystemEvent, FileSystemEventHandler

from core.configs import logger


[документация] class CreateFileHandler(FileSystemEventHandler): """Наблюдатель за появлением декодированных кадров в директории наблюдения. Кадры копятся в течение ``OBSERVATION__BATCH_DELAY`` и уходят одной пачкой, не дожидаясь конца прохода. Файл берётся в работу по ``on_closed`` (запись завершена) и по ``on_created`` (файл только появился). Второе нужно для наблюдателей, не отдающих ``IN_CLOSE_WRITE``; чтобы не выгрузить кадр обрезанным, перед отправкой размер файла сверяется с тем, что был на момент постановки в очередь. """ def __init__( self, observation_id: str, prefix: str, send_callback: Callable[[list[Path], str], None], scheduler: BaseScheduler, batch_delay: float, ) -> None: super().__init__() self.observation_id = observation_id self.prefix = prefix self.send_callback = send_callback self.scheduler = scheduler self.batch_delay = batch_delay # Путь -> размер на момент постановки в очередь. Словарь заодно даёт # дедуп: файл, пришедший и с on_created, и с on_closed, уйдёт один раз. self.pending: dict[Path, int] = {} self.lock = threading.Lock() self.is_scheduled_job = False
[документация] def on_created(self, event: FileSystemEvent) -> None: """Взять в работу появившийся файл.""" self._enqueue(event)
[документация] def on_closed(self, event: FileSystemEvent) -> None: """Взять в работу файл, запись в который завершена.""" self._enqueue(event)
[документация] def send_batch(self, force: bool = False) -> None: """Отправить накопленную пачку кадров. Args: force: Проход окончен, файлы дописаны — отправлять не сверяя размер. """ with self.lock: pending = self.pending self.pending = {} self.is_scheduled_job = False current_batch = [] for file_path, size in pending.items(): current_size = self._get_size(file_path) if current_size is None: logger.warning( "Файл исчез до отправки [наблюдение: %s]: %s", self.observation_id, file_path, ) continue if not force and current_size != size: # Файл ещё дописывается: вернуть в очередь с новым # размером, следующий цикл увидит совпадение и отправит. self.pending[file_path] = current_size continue current_batch.append(file_path) if self.pending: self._schedule_send_batch() if current_batch: logger.debug( "Добавлено задание отправки данных [наблюдение: %s]. " "Количество файлов для отправки: %d", self.observation_id, len(current_batch), ) # Вне лока: это сетевой вызов. self.send_callback(current_batch, self.observation_id)
def _enqueue(self, event: FileSystemEvent) -> None: file_path = Path(event.src_path) if not file_path.is_file() or not file_path.name.startswith(self.prefix): return size = self._get_size(file_path) if size is None: return logger.debug( "Создан файл с данными [наблюдение: %s]: %s", self.observation_id, file_path, ) with self.lock: self.pending[file_path] = size if not self.is_scheduled_job: self._schedule_send_batch() def _schedule_send_batch(self) -> None: self.is_scheduled_job = True run_time = datetime.now(UTC) + timedelta(seconds=self.batch_delay) self.scheduler.add_job( func=self.send_batch, trigger="date", next_run_time=run_time, ) @staticmethod def _get_size(file_path: Path) -> int | None: try: return file_path.stat().st_size except OSError: return None