Исходный код 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