v0.0.63: TTL-фикс, прерывание с сохранением, трекинг файлов/чанков, ETA (global+per-file), UI-таблица 3 секции + кнопка Прервать
Deploy drhider / validate (push) Canceled after 0s

This commit is contained in:
“Naeel”
2026-08-24 14:59:03 +03:00
parent 0e5f3ec9a1
commit 25e4b46e76
7 changed files with 581 additions and 108 deletions
+178 -65
View File
@@ -13,6 +13,7 @@ import io
import json
import queue
import threading
import time
import zipfile
import traceback
import logging
@@ -23,7 +24,8 @@ from flask import Blueprint, request, send_file, jsonify, Response, stream_with_
from drhider import obfuscate_files, LLMClient
from session import (create_session, add_file, get_files, store_result,
get_result, store_csv, get_csv, cleanup, file_count,
MAX_FILE_BYTES)
MAX_FILE_BYTES, pause_ttl, resume_ttl,
request_cancel, get_cancel_event)
api_bp = Blueprint("api", __name__, url_prefix="/api")
log = logging.getLogger("routes.api_bp")
@@ -142,6 +144,20 @@ def session_files(sid):
})
@api_bp.route("/cancel/<sid>", methods=["POST"])
def cancel(sid):
"""Запросить мягкое прерывание обработки сессии.
Воркер останавливается на ближайшей границе файла (текущий добирается),
собирает частичный результат (готовые файлы + mapping) и шлёт SSE-событие `cancelled`.
"""
if request_cancel(sid):
log.info("cancel: requested sid=%s", sid)
return jsonify({"ok": True}), 200
log.warning("cancel: session not found sid=%s", sid)
return jsonify({"ok": False, "error": "Session not found"}), 404
@api_bp.route("/process_stream/<sid>", methods=["GET"])
def process_stream(sid):
"""SSE: process all session files, streaming per-file progress.
@@ -149,7 +165,7 @@ def process_stream(sid):
Все файлы обрабатываются ЕДИНЫМ вызовом obfuscate_files (общий mapping,
согласованные токены). Обработка идёт в отдельном потоке; прогресс
передаётся через очередь. Разрыв соединения клиента корректно
перехватывается и останавливает генератор.
перехватывается и останавливает генератор (и воркер — через cancel_event).
"""
files = get_files(sid)
if files is None:
@@ -159,30 +175,60 @@ def process_stream(sid):
all_files = [(fname, content, "") for fname, content in files]
log.info("process_stream: start sid=%s files=%d", sid, len(all_files))
# Сессия живёт, пока идёт обработка (TTL возобновляется в finally генератора)
pause_ttl(sid)
def generate():
llm = LLMClient()
q = queue.Queue()
cancel = threading.Event()
cancel = threading.Event() # локальный: разрыв клиента (стоп heartbeat)
cancel_event = get_cancel_event(sid) # из сессии: мягкая отмена (кнопка «Прервать»)
# Состояние для глобальной ETA (символы)
eta = {"total_chars": 0, "done_chars": 0, "cur_chars": 0, "cur_total": 0, "cur_done": 0}
_TOKENS_PER_CHAR = 8.0 # эмпирический коэффициент символов -> токенов LLM
def progress(phase, idx, total_, name, elapsed):
q.put(("progress", phase, idx, name, total_, elapsed))
# Состояние для per-file ETA (чанки)
fstate = {"t0": 0.0, "total": 0, "done": 0}
def file_progress(event, fname, **fields):
if event == "file_start":
fstate["t0"] = time.time()
fstate["total"] = fields.get("chunks", 0)
fstate["done"] = 0
elif event == "file_chunk":
fstate["done"] = fields.get("chunks_done", 0)
elapsed = time.time() - fstate["t0"]
rate = fstate["done"] / elapsed if elapsed > 0 else 0
rem = fstate["total"] - fstate["done"]
if rate > 0:
fields["eta_sec"] = max(0, round(rem / rate))
q.put(("file", event, fname, fields))
def worker():
log.info("worker: start sid=%s files=%d", sid, len(all_files))
t0 = datetime.utcnow()
try:
zip_data, csv_str = obfuscate_files(
all_files, llm_client=llm, progress_cb=progress
zip_data, csv_str, meta = obfuscate_files(
all_files, llm_client=llm, progress_cb=progress,
cancel_event=cancel_event, file_progress=file_progress,
)
stats = {
"tokens": llm.tokens_total,
"llm_sec": round(llm.llm_sec, 1),
}
dt = (datetime.utcnow() - t0).total_seconds()
log.info("worker: done sid=%s in %.1fs tokens=%d llm_sec=%.1f zip_len=%d",
sid, dt, llm.tokens_total, llm.llm_sec, len(zip_data))
q.put(("result", zip_data, csv_str, stats))
if meta and meta.get("cancelled"):
log.info("worker: cancelled sid=%s in %.1fs processed=%d/%d",
sid, dt, meta.get("processed", 0), meta.get("total", 0))
q.put(("cancelled", zip_data, csv_str, stats, meta))
else:
log.info("worker: done sid=%s in %.1fs tokens=%d llm_sec=%.1f zip_len=%d",
sid, dt, llm.tokens_total, llm.llm_sec, len(zip_data))
q.put(("result", zip_data, csv_str, stats))
except Exception as e:
log.error("worker: exception sid=%s: %r\n%s", sid, e, traceback.format_exc())
q.put(("error", repr(e)))
@@ -190,70 +236,133 @@ def process_stream(sid):
threading.Thread(target=worker, daemon=True).start()
log.debug("process_stream: worker thread started sid=%s", sid)
while True:
try:
evt = q.get(timeout=1)
except queue.Empty:
if cancel.is_set():
log.info("process_stream: cancelled sid=%s (generator exit)", sid)
return
# Heartbeat: живая статистика LLM (для таймера в UI)
try:
while True:
try:
yield (
f"event: llm\n"
f"data: {json.dumps({'active': llm.llm_active, 'elapsed': round(llm.llm_elapsed_now(), 1), 'tokens': llm.tokens_total})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect during heartbeat sid=%s err=%r", sid, e)
cancel.set()
return
continue
evt = q.get(timeout=1)
except queue.Empty:
if cancel.is_set():
log.info("process_stream: cancelled sid=%s (generator exit)", sid)
return
# Heartbeat: живая статистика LLM + глобальная ETA
tokens = llm.tokens_total
elapsed = llm.llm_elapsed_now()
eta_sec = None
if eta["total_chars"] > 0 and tokens > 0 and elapsed > 0:
est_total_tokens = eta["total_chars"] * _TOKENS_PER_CHAR
rate = tokens / elapsed
if rate > 0:
eta_sec = max(0, int((est_total_tokens - tokens) / rate))
done_chars = eta["done_chars"]
if eta["cur_total"] > 0:
done_chars += int(eta["cur_chars"] * (eta["cur_done"] / eta["cur_total"]))
try:
yield (
f"event: llm\n"
f"data: {json.dumps({'active': llm.llm_active, 'elapsed': round(elapsed, 1), 'tokens': tokens, 'eta_sec': eta_sec, 'done_chars': done_chars, 'total_chars': eta['total_chars']})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect during heartbeat sid=%s err=%r", sid, e)
cancel.set()
if cancel_event:
cancel_event.set()
return
continue
kind = evt[0]
kind = evt[0]
if kind == "progress":
_, phase, idx, name, total_, elapsed = evt
log.debug("process_stream: event=%s idx=%d name=%r elapsed=%s sid=%s",
phase, idx, name, elapsed, sid)
try:
yield (
f"event: {phase}\n"
f"data: {json.dumps({'idx': idx, 'name': name, 'total': total_, 'elapsed': elapsed})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on progress sid=%s phase=%s err=%r", sid, phase, e)
cancel.set()
if kind == "progress":
_, phase, idx, name, total_, elapsed = evt
log.debug("process_stream: event=%s idx=%d name=%r elapsed=%s sid=%s",
phase, idx, name, elapsed, sid)
try:
yield (
f"event: {phase}\n"
f"data: {json.dumps({'idx': idx, 'name': name, 'total': total_, 'elapsed': elapsed})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on progress sid=%s phase=%s err=%r", sid, phase, e)
cancel.set()
if cancel_event:
cancel_event.set()
return
elif kind == "file":
_, event, fname, fields = evt
# Обновляем состояние для глобальной ETA
if event == "extract_done":
eta["total_chars"] = fields.get("total_chars", 0)
elif event == "file_start":
eta["cur_chars"] = fields.get("chars", 0)
eta["cur_total"] = fields.get("chunks", 0)
eta["cur_done"] = 0
elif event == "file_chunk":
eta["cur_done"] = fields.get("chunks_done", 0)
elif event == "file_done":
eta["done_chars"] += eta["cur_chars"]
eta["cur_chars"] = 0
eta["cur_total"] = 0
eta["cur_done"] = 0
try:
yield (
f"event: {event}\n"
f"data: {json.dumps({'name': fname, **fields})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on file event sid=%s err=%r", sid, e)
cancel.set()
if cancel_event:
cancel_event.set()
return
elif kind == "result":
_, zip_data, csv_str, stats = evt
log.info("process_stream: result sid=%s, storing result", sid)
store_result(sid, zip_data)
if csv_str:
store_csv(sid, csv_str)
count = 0
with zipfile.ZipFile(io.BytesIO(zip_data)) as zf:
count = len([n for n in zf.namelist() if n != "mapping.csv"])
log.info("process_stream: complete sid=%s count=%d stats=%r", sid, count, stats)
try:
yield (
f"event: complete\n"
f"data: {json.dumps({'total': count, **stats})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on complete sid=%s err=%r", sid, e)
return
return
elif kind == "result":
_, zip_data, csv_str, stats = evt
log.info("process_stream: result sid=%s, storing result", sid)
store_result(sid, zip_data)
if csv_str:
store_csv(sid, csv_str)
count = 0
with zipfile.ZipFile(io.BytesIO(zip_data)) as zf:
count = len([n for n in zf.namelist() if n != "mapping.csv"])
log.info("process_stream: complete sid=%s count=%d stats=%r", sid, count, stats)
try:
yield (
f"event: complete\n"
f"data: {json.dumps({'total': count, **stats})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on complete sid=%s err=%r", sid, e)
elif kind == "cancelled":
_, zip_data, csv_str, stats, meta = evt
log.info("process_stream: cancelled sid=%s, storing partial result", sid)
store_result(sid, zip_data)
if csv_str:
store_csv(sid, csv_str)
try:
yield (
f"event: cancelled\n"
f"data: {json.dumps({'saved': meta.get('processed', 0), 'total': meta.get('total', 0), **stats})}\n\n"
)
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on cancelled sid=%s err=%r", sid, e)
return
return
return
elif kind == "error":
_, msg = evt
log.error("process_stream: error event sid=%s msg=%r", sid, msg)
try:
yield f"event: error\ndata: {json.dumps({'error': msg})}\n\n"
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on error sid=%s err=%r", sid, e)
elif kind == "error":
_, msg = evt
log.error("process_stream: error event sid=%s msg=%r", sid, msg)
try:
yield f"event: error\ndata: {json.dumps({'error': msg})}\n\n"
except _disconnect_exceptions() as e:
log.warning("process_stream: disconnect on error sid=%s err=%r", sid, e)
return
return
return
finally:
resume_ttl(sid)
log.debug("process_stream: generator exit sid=%s, TTL resumed", sid)
return Response(
stream_with_context(generate()),
@@ -273,7 +382,11 @@ def process(sid):
try:
llm = LLMClient()
all_files = [(fname, content, "") for fname, content in files]
zip_data, csv_str = obfuscate_files(all_files, llm_client=llm)
pause_ttl(sid)
try:
zip_data, csv_str, _ = obfuscate_files(all_files, llm_client=llm)
finally:
resume_ttl(sid)
store_result(sid, zip_data)
if csv_str:
store_csv(sid, csv_str)