feat: implement Layer 2 per-file transit, RAM session, mock buffer, and tests (v0.2.0)
This commit is contained in:
@@ -0,0 +1 @@
|
||||
"""Backend переиспользуемого модуля загрузки: upload_refs + session."""
|
||||
@@ -0,0 +1,38 @@
|
||||
"""In-memory хранилище сессий с TTL и лимитами.
|
||||
|
||||
Каждая операция — в отдельном файле, общее состояние — в state.py.
|
||||
Импорт как единый пакет:
|
||||
|
||||
from upload.backend.session import create_session, add_file, get_files
|
||||
"""
|
||||
|
||||
from .create_session import create_session
|
||||
from .add_file import add_file
|
||||
from .get_files import get_files, file_count
|
||||
from .store_result import store_result, get_result
|
||||
from .store_csv import store_csv, get_csv
|
||||
from .ttl import touch, pause_ttl, resume_ttl
|
||||
from .cancel import request_cancel, get_cancel_event
|
||||
from .cleanup import cleanup
|
||||
from .state import TTL_SECONDS, MAX_FILE_BYTES, MAX_SESSION_BYTES, configure
|
||||
|
||||
__all__ = [
|
||||
"create_session",
|
||||
"add_file",
|
||||
"get_files",
|
||||
"file_count",
|
||||
"store_result",
|
||||
"get_result",
|
||||
"store_csv",
|
||||
"get_csv",
|
||||
"touch",
|
||||
"pause_ttl",
|
||||
"resume_ttl",
|
||||
"request_cancel",
|
||||
"get_cancel_event",
|
||||
"cleanup",
|
||||
"configure",
|
||||
"TTL_SECONDS",
|
||||
"MAX_FILE_BYTES",
|
||||
"MAX_SESSION_BYTES",
|
||||
]
|
||||
@@ -0,0 +1,25 @@
|
||||
"""add_file — добавить файл в сессию (с проверкой суммарного лимита)."""
|
||||
|
||||
from . import state
|
||||
|
||||
|
||||
def add_file(sid: str, filename: str, content: bytes) -> bool:
|
||||
"""Добавить файл в сессию.
|
||||
|
||||
Args:
|
||||
sid: Идентификатор сессии
|
||||
filename: Имя файла
|
||||
content: Бинарное содержимое
|
||||
|
||||
Returns:
|
||||
True если добавлено; False если сессии нет или превышен лимит сессии.
|
||||
"""
|
||||
with state._lock:
|
||||
s = state._sessions.get(sid)
|
||||
if not s:
|
||||
return False
|
||||
total = sum(len(c) for _, c in s["files"])
|
||||
if total + len(content) > state.MAX_SESSION_BYTES:
|
||||
return False # превышен суммарный лимит сессии
|
||||
s["files"].append((filename, content))
|
||||
return True
|
||||
@@ -0,0 +1,27 @@
|
||||
"""request_cancel / get_cancel_event — мягкое прерывание обработки сессии."""
|
||||
|
||||
import threading
|
||||
from typing import Optional
|
||||
|
||||
from .state import _sessions, _lock
|
||||
|
||||
|
||||
def request_cancel(sid: str) -> bool:
|
||||
"""Запросить мягкое прерывание обработки сессии.
|
||||
|
||||
Returns:
|
||||
True если сессия существует и отмена запрошена, False если нет.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if not s:
|
||||
return False
|
||||
s["cancel"].set()
|
||||
return True
|
||||
|
||||
|
||||
def get_cancel_event(sid: str) -> Optional[threading.Event]:
|
||||
"""Получить событие отмены сессии (или None, если сессии нет)."""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
return s["cancel"] if s else None
|
||||
@@ -0,0 +1,11 @@
|
||||
"""cleanup — удалить сессию."""
|
||||
|
||||
from .state import _sessions, _lock
|
||||
|
||||
|
||||
def cleanup(sid: str):
|
||||
"""Удалить сессию и остановить её TTL-таймер."""
|
||||
with _lock:
|
||||
s = _sessions.pop(sid, None)
|
||||
if s and s.get("timer"):
|
||||
s["timer"].cancel()
|
||||
@@ -0,0 +1,23 @@
|
||||
"""create_session — создать новую сессию."""
|
||||
|
||||
import threading
|
||||
import uuid
|
||||
|
||||
from .state import _sessions, _lock, _start_timer
|
||||
|
||||
|
||||
def create_session() -> str:
|
||||
"""Создать новую сессию.
|
||||
|
||||
Returns:
|
||||
Уникальный идентификатор сессии (UUID).
|
||||
"""
|
||||
sid = uuid.uuid4().hex
|
||||
with _lock:
|
||||
_sessions[sid] = {
|
||||
"files": [],
|
||||
"result": None,
|
||||
"cancel": threading.Event(),
|
||||
"timer": _start_timer(sid),
|
||||
}
|
||||
return sid
|
||||
@@ -0,0 +1,23 @@
|
||||
"""get_files / file_count — чтение файлов сессии."""
|
||||
|
||||
from typing import List, Optional, Tuple
|
||||
|
||||
from .state import _sessions, _lock
|
||||
|
||||
|
||||
def get_files(sid: str) -> Optional[List[Tuple[str, bytes]]]:
|
||||
"""Получить все файлы сессии.
|
||||
|
||||
Returns:
|
||||
[(filename, content), ...] или None если сессия не найдена.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
return list(s["files"]) if s else None
|
||||
|
||||
|
||||
def file_count(sid: str) -> int:
|
||||
"""Количество файлов в сессии."""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
return len(s["files"]) if s else 0
|
||||
@@ -0,0 +1,45 @@
|
||||
"""Общее состояние сессий: хранилище, блокировка, константы, TTL-таймер.
|
||||
|
||||
Единая точка хранения состояния — все операции импортируют её.
|
||||
Дробить state.py на файл-на-переменную не нужно: это данные, а не функции.
|
||||
"""
|
||||
|
||||
import threading
|
||||
|
||||
# TTL сессии: 30 минут
|
||||
TTL_SECONDS = 30 * 60
|
||||
|
||||
# Максимальный объём одного файла и суммарный объём файлов в сессии (защита памяти)
|
||||
MAX_FILE_BYTES = 50 * 1024 * 1024 # 50 MB на один файл
|
||||
MAX_SESSION_BYTES = 500 * 1024 * 1024 # 500 MB суммарно на сессию
|
||||
|
||||
_sessions: dict = {}
|
||||
_lock = threading.Lock()
|
||||
|
||||
|
||||
def _start_timer(sid: str) -> threading.Timer:
|
||||
"""Запустить таймер автоочистки сессии через TTL."""
|
||||
|
||||
def _clean():
|
||||
with _lock:
|
||||
_sessions.pop(sid, None)
|
||||
|
||||
timer = threading.Timer(TTL_SECONDS, _clean)
|
||||
timer.daemon = True
|
||||
timer.start()
|
||||
return timer
|
||||
|
||||
|
||||
def configure(max_file_bytes: int = None, max_session_bytes: int = None,
|
||||
ttl_seconds: int = None):
|
||||
"""Переопределить лимиты/TTL из конфига приложения (глобально).
|
||||
|
||||
None — оставить текущее значение.
|
||||
"""
|
||||
global MAX_FILE_BYTES, MAX_SESSION_BYTES, TTL_SECONDS
|
||||
if max_file_bytes is not None:
|
||||
MAX_FILE_BYTES = max_file_bytes
|
||||
if max_session_bytes is not None:
|
||||
MAX_SESSION_BYTES = max_session_bytes
|
||||
if ttl_seconds is not None:
|
||||
TTL_SECONDS = ttl_seconds
|
||||
@@ -0,0 +1,30 @@
|
||||
"""store_csv / get_csv — сохранение и чтение CSV с таблицей замен."""
|
||||
|
||||
from typing import Optional
|
||||
|
||||
from .state import _sessions, _lock
|
||||
|
||||
|
||||
def store_csv(sid: str, csv_str: str) -> bool:
|
||||
"""Сохранить CSV с таблицей замен.
|
||||
|
||||
Returns:
|
||||
True если сохранено, False если сессии нет.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if not s:
|
||||
return False
|
||||
s["csv"] = csv_str
|
||||
return True
|
||||
|
||||
|
||||
def get_csv(sid: str) -> Optional[str]:
|
||||
"""Получить CSV с таблицей замен.
|
||||
|
||||
Returns:
|
||||
Строка CSV или None если нет.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
return s.get("csv") if s else None
|
||||
@@ -0,0 +1,30 @@
|
||||
"""store_result / get_result — сохранение и чтение результата обработки."""
|
||||
|
||||
from typing import Optional
|
||||
|
||||
from .state import _sessions, _lock
|
||||
|
||||
|
||||
def store_result(sid: str, zip_data: bytes) -> bool:
|
||||
"""Сохранить результат обработки (ZIP-архив).
|
||||
|
||||
Returns:
|
||||
True если сохранено, False если сессии нет.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if not s:
|
||||
return False
|
||||
s["result"] = zip_data
|
||||
return True
|
||||
|
||||
|
||||
def get_result(sid: str) -> Optional[bytes]:
|
||||
"""Получить результат обработки.
|
||||
|
||||
Returns:
|
||||
ZIP-архив или None если сессия не найдена/результат не готов.
|
||||
"""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
return s["result"] if s else None
|
||||
@@ -0,0 +1,34 @@
|
||||
"""touch / pause_ttl / resume_ttl — управление TTL-таймером сессии."""
|
||||
|
||||
from .state import _sessions, _lock, _start_timer
|
||||
|
||||
|
||||
def touch(sid: str):
|
||||
"""Продлить жизнь сессии: перезапустить TTL-таймер (если сессия существует)."""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if not s:
|
||||
return
|
||||
if s.get("timer"):
|
||||
s["timer"].cancel()
|
||||
s["timer"] = _start_timer(sid)
|
||||
|
||||
|
||||
def pause_ttl(sid: str):
|
||||
"""Приостановить TTL сессии (во время обработки): сессия живёт, пока идёт воркер."""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if s and s.get("timer"):
|
||||
s["timer"].cancel()
|
||||
s["timer"] = None
|
||||
|
||||
|
||||
def resume_ttl(sid: str):
|
||||
"""Возобновить TTL сессии (после завершения обработки): результат доступен ещё TTL."""
|
||||
with _lock:
|
||||
s = _sessions.get(sid)
|
||||
if not s:
|
||||
return
|
||||
if s.get("timer"):
|
||||
s["timer"].cancel()
|
||||
s["timer"] = _start_timer(sid)
|
||||
@@ -0,0 +1,7 @@
|
||||
"""upload_refs — Flask Blueprint для приёма ссылок и pull с ВМ-буфера."""
|
||||
|
||||
from .blueprint import create_upload_refs_blueprint
|
||||
from .safe_name import safe_name
|
||||
from .pull_file import pull_file
|
||||
|
||||
__all__ = ["create_upload_refs_blueprint", "safe_name", "pull_file"]
|
||||
@@ -0,0 +1,145 @@
|
||||
"""Переиспользуемый Blueprint слоя 2: POST /upload_refs (pull с ВМ-буфера в сессию).
|
||||
|
||||
Поведение:
|
||||
- safe_name (защита от path traversal)
|
||||
- SSRF-валидация url.startswith(vm_prefix)
|
||||
- лимит на один файл -> delete + skip
|
||||
- pull с ретраями через httpx stream
|
||||
- лимит сессии -> skip; отсутствие сессии -> 404
|
||||
- delete url с ВМ (best-effort)
|
||||
- поддержка пофайлового транзита и пачек
|
||||
- опциональный callback для интеграции/эмуляции Слоя 3 (on_file_received)
|
||||
"""
|
||||
|
||||
import httpx
|
||||
import logging
|
||||
from urllib.parse import urlsplit
|
||||
from flask import Blueprint, request, jsonify, current_app
|
||||
|
||||
from ..session import (create_session, add_file, get_files, file_count,
|
||||
MAX_FILE_BYTES, configure)
|
||||
from .config import PULL_RETRIES, PULL_RETRY_DELAY, VM_UPLOAD_PREFIX
|
||||
from .safe_name import safe_name
|
||||
from .pull_file import pull_file
|
||||
|
||||
log = logging.getLogger("upload.upload_refs")
|
||||
|
||||
|
||||
def create_upload_refs_blueprint(cfg: dict = None) -> Blueprint:
|
||||
"""Создать Blueprint с эндпоинтом upload_refs.
|
||||
|
||||
cfg (все ключи опциональны, есть дефолты):
|
||||
apiPrefix (str) — префикс Blueprint, по умолчанию "/api"
|
||||
vmUploadPrefix (str) — доверенный префикс ВМ-буфера (SSRF-валидация)
|
||||
maxFileBytes (int) — лимит на один файл
|
||||
maxSessionBytes (int) — суммарный лимит сессии (применяется к сессиям)
|
||||
ttlSeconds (int) — TTL сессии
|
||||
pullRetries (int) — ретраи pull
|
||||
pullRetryDelay (int) — пауза между ретраями (сек)
|
||||
pullTimeout (int) — таймаут одного GET pull
|
||||
onFileReceived (func) — опциональный callback (sid, name, content) для Слоя 3
|
||||
"""
|
||||
cfg = cfg or {}
|
||||
prefix = cfg.get("apiPrefix", "/api")
|
||||
vm_prefix = cfg.get("vmUploadPrefix", VM_UPLOAD_PREFIX)
|
||||
max_file_bytes = cfg.get("maxFileBytes", MAX_FILE_BYTES)
|
||||
pull_retries = cfg.get("pullRetries", PULL_RETRIES)
|
||||
pull_delay = cfg.get("pullRetryDelay", PULL_RETRY_DELAY)
|
||||
pull_timeout = cfg.get("pullTimeout", 120)
|
||||
on_file_received = cfg.get("onFileReceived")
|
||||
|
||||
# Применить лимиты сессии/TTL из конфига (глобально для всех сессий)
|
||||
configure(
|
||||
max_file_bytes=cfg.get("maxFileBytes"),
|
||||
max_session_bytes=cfg.get("maxSessionBytes"),
|
||||
ttl_seconds=cfg.get("ttlSeconds"),
|
||||
)
|
||||
|
||||
is_path_only_prefix = vm_prefix.startswith("/")
|
||||
|
||||
bp = Blueprint("upload_refs", __name__, url_prefix=prefix)
|
||||
|
||||
@bp.route("/upload_refs", methods=["POST"])
|
||||
def upload_refs():
|
||||
"""Принять ссылки на файлы (загружены на ВМ-буфер), забрать по egress.
|
||||
|
||||
Вход: JSON {"session": "...", "files": [{"name": str, "size": int, "url": str}]}.
|
||||
Каждый файл тянется исходящим GET с ВМ (egress не ограничен шлюзом),
|
||||
читается по частям (stream), кладётся в сессию. После успешного pull файл
|
||||
удаляется с ВМ (DELETE).
|
||||
"""
|
||||
data = request.get_json(silent=True) or {}
|
||||
sid = data.get("session") or create_session()
|
||||
refs = data.get("files") or []
|
||||
if not refs:
|
||||
log.warning("upload_refs: no files, sid=%s", sid)
|
||||
return jsonify({"ok": False, "error": "No files"}), 400
|
||||
added = 0
|
||||
transport = cfg.get("httpxTransport") or current_app.config.get("UPLOAD_HTTPX_TRANSPORT")
|
||||
try:
|
||||
with httpx.Client(transport=transport, timeout=pull_timeout, follow_redirects=True) as client:
|
||||
for ref in refs:
|
||||
name = safe_name(ref.get("name") or "")
|
||||
url = ref.get("url") or ""
|
||||
if not name or not url:
|
||||
continue
|
||||
# SSRF-защита: тянуть можно ТОЛЬКО с доверенного ВМ-буфера
|
||||
if is_path_only_prefix:
|
||||
if not urlsplit(url).path.startswith(vm_prefix):
|
||||
log.warning("upload_refs: unsafe URL path, skip sid=%s url=%r", sid, url)
|
||||
continue
|
||||
if not url.startswith(("http://", "https://")):
|
||||
url = request.host_url.rstrip("/") + ("/" if not url.startswith("/") else "") + url
|
||||
else:
|
||||
if not url.startswith(vm_prefix):
|
||||
log.warning("upload_refs: unsafe URL, skip sid=%s url=%r", sid, url)
|
||||
continue
|
||||
# Лимит на один файл: сверх лимита — пропускаем
|
||||
if (ref.get("size") or 0) > max_file_bytes:
|
||||
log.warning("upload_refs: file exceeds %dMB, skip sid=%s file=%r size=%s",
|
||||
max_file_bytes // (1024 * 1024), sid, name, ref.get("size"))
|
||||
try:
|
||||
client.delete(url)
|
||||
except Exception:
|
||||
pass
|
||||
continue
|
||||
# Pull с ретраями
|
||||
content = pull_file(client, url, pull_retries, pull_delay, sid=sid, name=name)
|
||||
log.info("upload_refs: pulled sid=%s file=%r size=%d", sid, name, len(content))
|
||||
if len(content) > max_file_bytes:
|
||||
log.warning("upload_refs: pulled file exceeds %dMB, skip sid=%s file=%r size=%d",
|
||||
max_file_bytes // (1024 * 1024), sid, name, len(content))
|
||||
try:
|
||||
client.delete(url)
|
||||
except Exception:
|
||||
pass
|
||||
continue
|
||||
if not add_file(sid, name, content):
|
||||
if get_files(sid) is None:
|
||||
log.warning("upload_refs: session not found, sid=%s file=%r", sid, name)
|
||||
return jsonify({"ok": False, "error": "Session not found"}), 404
|
||||
log.warning("upload_refs: session limit exceeded, skip sid=%s file=%r", sid, name)
|
||||
try:
|
||||
client.delete(url)
|
||||
except Exception:
|
||||
pass
|
||||
continue
|
||||
# Успешно добавлен в сессию — удаляем с ВМ-буфера
|
||||
try:
|
||||
client.delete(url)
|
||||
except Exception:
|
||||
pass
|
||||
# Если передан callback для Слоя 3 (эмуляция или реальный процессинг)
|
||||
if callable(on_file_received):
|
||||
try:
|
||||
on_file_received(sid, name, content)
|
||||
except Exception as cb_err:
|
||||
log.warning("upload_refs: on_file_received callback error: %r", cb_err)
|
||||
added += 1
|
||||
except Exception as e:
|
||||
log.error("upload_refs: pull error sid=%s: %r", sid, e)
|
||||
return jsonify({"ok": False, "error": "Pull failed: %s" % e}), 502
|
||||
log.info("upload_refs: done sid=%s added=%d total=%d", sid, added, file_count(sid))
|
||||
return jsonify({"ok": True, "session": sid, "count": file_count(sid), "added": added})
|
||||
|
||||
return bp
|
||||
@@ -0,0 +1,11 @@
|
||||
"""Параметры слоя 2 (бэк) по умолчанию.
|
||||
|
||||
Переопределяются из конфига приложения через create_upload_refs_blueprint(cfg).
|
||||
"""
|
||||
|
||||
# Ретраи pull из ВМ-буфера: защита от разовых DNS/сетевых сбоев (gaierror -5 и т.п.)
|
||||
PULL_RETRIES = 3
|
||||
PULL_RETRY_DELAY = 2 # секунды между попытками
|
||||
|
||||
# Доверенный префикс ВМ-буфера по умолчанию — валидация URL при pull (защита от SSRF)
|
||||
VM_UPLOAD_PREFIX = "https://contracts.kube5s.ru/drhider-upload/"
|
||||
@@ -0,0 +1,39 @@
|
||||
"""pull_file — вытащить файл с ВМ-буфера исходящим GET с ретраями."""
|
||||
|
||||
import logging
|
||||
import time
|
||||
|
||||
from .config import PULL_RETRIES, PULL_RETRY_DELAY
|
||||
|
||||
log = logging.getLogger("upload.upload_refs.pull")
|
||||
|
||||
|
||||
def pull_file(client, url: str, retries: int = PULL_RETRIES,
|
||||
delay: float = PULL_RETRY_DELAY, sid: str = None, name: str = None) -> bytes:
|
||||
"""GET url с ретраями; читает по частям (stream).
|
||||
|
||||
Args:
|
||||
client: httpx.Client
|
||||
url: URL файла на ВМ-буфере.
|
||||
retries: число попыток.
|
||||
delay: пауза между попытками (сек).
|
||||
sid/name: для логирования (опционально).
|
||||
|
||||
Returns:
|
||||
Содержимое файла (bytes).
|
||||
|
||||
Raises:
|
||||
Последнюю ошибку попытки, если все ретраи не удались.
|
||||
"""
|
||||
last_err = None
|
||||
for attempt in range(retries):
|
||||
try:
|
||||
with client.stream("GET", url) as resp:
|
||||
resp.raise_for_status()
|
||||
return b"".join(resp.iter_bytes())
|
||||
except Exception as e:
|
||||
last_err = e
|
||||
log.warning("pull: attempt %d/%d failed sid=%s file=%r: %r",
|
||||
attempt + 1, retries, sid, name, e)
|
||||
time.sleep(delay)
|
||||
raise last_err if last_err else RuntimeError("pull failed")
|
||||
@@ -0,0 +1,16 @@
|
||||
"""safe_name — санитизация имени файла (защита от path traversal)."""
|
||||
|
||||
|
||||
def safe_name(name: str) -> str:
|
||||
"""Санитизировать имя файла: защита от path traversal, сохраняя подпапки.
|
||||
|
||||
Запрещает '..' и абсолютные пути; нормализует слэши. Возвращает "" если
|
||||
имя пустое или небезопасное.
|
||||
"""
|
||||
if not name:
|
||||
return ""
|
||||
name = name.replace("\\", "/")
|
||||
parts = [p for p in name.split("/") if p and p != "."]
|
||||
if not parts or any(p == ".." for p in parts):
|
||||
return ""
|
||||
return "/".join(parts)
|
||||
Reference in New Issue
Block a user