этап 2: переиспользуемый модуль upload (sink) вместо рукописного транспорта

- копирую модуль upload/ из drhider (слои 1-2)
- blueprint: параметр sink (drhider-сессия по умолчанию, сверка — DB+парсинг)
- upload_bp: contracts_upload_sink = _store_and_parse
- routes: регистрирую create_upload_refs_blueprint(cfg, sink=...)
- app.py: корень репо в sys.path (для import upload)
- History: план переиспользования + ревью Соннета
This commit is contained in:
“Naeel”
2026-08-26 08:07:43 +03:00
parent 55bd7a027d
commit db58a433fb
46 changed files with 1986 additions and 27 deletions
+1
View File
@@ -0,0 +1 @@
"""Backend переиспользуемого модуля загрузки: upload_refs + session."""
+38
View File
@@ -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",
]
+25
View File
@@ -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
+27
View File
@@ -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
+11
View File
@@ -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()
+23
View File
@@ -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
+23
View File
@@ -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
+45
View File
@@ -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
+30
View File
@@ -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
+30
View File
@@ -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
+34
View File
@@ -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)
+20
View File
@@ -0,0 +1,20 @@
"""Слой 2 (бэк): переиспользуемый Blueprint закачки через ВМ.
Использование:
from upload.backend.upload_refs import create_upload_refs_blueprint
app.register_blueprint(create_upload_refs_blueprint(cfg))
"""
from .blueprint import create_upload_refs_blueprint
from .safe_name import safe_name
from .pull_file import pull_file
from .config import PULL_RETRIES, PULL_RETRY_DELAY, VM_UPLOAD_PREFIX
__all__ = [
"create_upload_refs_blueprint",
"safe_name",
"pull_file",
"PULL_RETRIES",
"PULL_RETRY_DELAY",
"VM_UPLOAD_PREFIX",
]
+160
View File
@@ -0,0 +1,160 @@
"""Переиспользуемый Blueprint слоя 2: POST /upload_refs (pull с ВМ-буфера в сессию).
Поведение 1:1 с drhider v0.0.75 (site/routes/api_bp.py):
- _safe_name (path traversal)
- SSRF-валидация url.startswith(VM_UPLOAD_PREFIX)
- лимит на один файл -> delete+skip
- pull с ретраями
- лимит сессии -> skip; отсутствие сессии -> 404
- delete url с ВМ (best-effort)
"""
import httpx
import logging
from flask import Blueprint, request, jsonify
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, sink=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
sink (callable | None): если задан — вместо add_file в in-memory сессию
вызывается sink(name, content, **extra) на каждый вытянутый файл,
endpoint возвращает {"ok": true, "results": [...]}. extra — сквозные
поля тела запроса (batch_id, contract_id, zip_source). Если None —
прежнее поведение drhider (in-memory сессия).
"""
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)
# Применить лимиты сессии/TTL из конфига (глобально для всех сессий)
configure(
max_file_bytes=cfg.get("maxFileBytes"),
max_session_bytes=cfg.get("maxSessionBytes"),
ttl_seconds=cfg.get("ttlSeconds"),
)
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 файл
удаляется с ВМ (best-effort; TTL-чистка на ВМ тоже есть).
"""
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
extra = {k: data[k] for k in ("batch_id", "contract_id", "zip_source") if data.get(k) is not None}
results = [] if sink else None
added = 0
try:
with httpx.Client(timeout=pull_timeout, follow_redirects=True) as client:
for ref in refs:
name = safe_name(ref.get("name") or "")
url = ref.get("url")
if not name or not url:
continue
# SSRF-защита: тянуть можно ТОЛЬКО с доверенного ВМ-буфера
if not url.startswith(vm_prefix):
log.warning("upload_refs: unsafe URL, skip sid=%s url=%r", sid, url)
if sink:
results.append({"name": name, "ok": False, "error": "invalid url (SSRF guard)"})
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
if sink:
results.append({"name": name, "ok": False, "error": "file too large"})
continue
# Pull с ретраями: разовые DNS/сетевые сбои не роняют всю загрузку
try:
content = pull_file(client, url, pull_retries, pull_delay, sid=sid, name=name)
except Exception as e:
log.warning("upload_refs: pull failed sid=%s file=%r: %r", sid, name, e)
if sink:
results.append({"name": name, "ok": False, "error": "pull failed: %s" % e})
continue
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
if sink:
results.append({"name": name, "ok": False, "error": "file too large"})
continue
if sink:
# Отдать файл приложению через sink (вместо in-memory сессии)
try:
client.delete(url)
except Exception:
pass
try:
r = sink(name, content, **extra) or {}
results.append({"name": name, **r})
except Exception as e:
log.error("upload_refs: sink failed file=%r: %r", name, e)
results.append({"name": name, "ok": False, "error": str(e)})
continue
if not add_file(sid, name, content):
# Различить: сессия исчезла vs превышен суммарный лимит сессии
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
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
if sink:
return jsonify({"ok": True, "results": results})
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)})
return bp
+11
View File
@@ -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/"
+39
View File
@@ -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")
+16
View File
@@ -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)