v0.0.32: фиксы по ревью — единый обфускатор, рекурсия ZIP, кириллица 1С, дедуп, SSE-дисконнект
Deploy drhider / validate (push) Canceled after 0s
Deploy drhider / validate (push) Canceled after 0s
This commit is contained in:
+93
-63
@@ -11,84 +11,130 @@ Five endpoints:
|
||||
|
||||
import io
|
||||
import json
|
||||
import queue
|
||||
import threading
|
||||
import zipfile
|
||||
import traceback
|
||||
from datetime import datetime, timedelta
|
||||
from flask import Blueprint, request, send_file, jsonify, Response, stream_with_context
|
||||
|
||||
from drhider import obfuscate_files, LLMClient
|
||||
from drhider.builder import build_zip, build_mapping_csv
|
||||
from session import (create_session, add_file, get_files, store_result,
|
||||
get_result, store_csv, get_csv, cleanup, file_count)
|
||||
|
||||
api_bp = Blueprint("api", __name__, url_prefix="/api")
|
||||
|
||||
|
||||
def _disconnect_exceptions():
|
||||
"""Исключения, означающие отключение клиента SSE."""
|
||||
return (GeneratorExit, BrokenPipeError, ConnectionResetError)
|
||||
|
||||
|
||||
@api_bp.route("/upload", methods=["POST"])
|
||||
def upload():
|
||||
"""Upload one file to session."""
|
||||
"""Upload files to session (один или несколько)."""
|
||||
sid = request.form.get("session", "")
|
||||
if not sid:
|
||||
sid = create_session()
|
||||
uploaded = request.files.getlist("files")
|
||||
if not uploaded:
|
||||
return jsonify({"ok": False, "error": "No file"}), 400
|
||||
f = uploaded[0]
|
||||
if not f.filename:
|
||||
return jsonify({"ok": False, "error": "No filename"}), 400
|
||||
ok = add_file(sid, f.filename, f.read())
|
||||
if not ok:
|
||||
return jsonify({"ok": False, "error": "Session not found"}), 404
|
||||
|
||||
added = 0
|
||||
had_unnamed = False
|
||||
for f in uploaded:
|
||||
if not f.filename:
|
||||
had_unnamed = True
|
||||
continue
|
||||
if not add_file(sid, f.filename, f.read()):
|
||||
return jsonify({"ok": False, "error": "Session not found"}), 404
|
||||
added += 1
|
||||
|
||||
if added == 0:
|
||||
err = "No filename" if had_unnamed else "No file"
|
||||
return jsonify({"ok": False, "error": err}), 400
|
||||
|
||||
return jsonify({"ok": True, "session": sid, "count": file_count(sid)})
|
||||
|
||||
|
||||
@api_bp.route("/process_stream/<sid>", methods=["GET"])
|
||||
def process_stream(sid):
|
||||
"""SSE: process all session files one by one, streaming per-file progress."""
|
||||
"""SSE: process all session files, streaming per-file progress.
|
||||
|
||||
Все файлы обрабатываются ЕДИНЫМ вызовом obfuscate_files (общий mapping,
|
||||
согласованные токены). Обработка идёт в отдельном потоке; прогресс
|
||||
передаётся через очередь. Разрыв соединения клиента корректно
|
||||
перехватывается и останавливает генератор.
|
||||
"""
|
||||
files = get_files(sid)
|
||||
if files is None:
|
||||
return jsonify({"ok": False, "error": "Session not found"}), 404
|
||||
if not files:
|
||||
return jsonify({"ok": False, "error": "No files"}), 400
|
||||
|
||||
total = len(files)
|
||||
all_files = [(fname, content, "") for fname, content in files]
|
||||
|
||||
def generate():
|
||||
llm = LLMClient()
|
||||
all_results = []
|
||||
all_mapping = {}
|
||||
for idx, (fname, content) in enumerate(files):
|
||||
# Отправляем: начали файл
|
||||
try:
|
||||
yield f"event: start\ndata: {json.dumps({'idx': idx, 'name': fname, 'total': total})}\n\n"
|
||||
except GeneratorExit:
|
||||
return # клиент отключился
|
||||
q = queue.Queue()
|
||||
cancel = threading.Event()
|
||||
|
||||
zip_data, _csv = obfuscate_files([(fname, content, "")], llm_client=llm)
|
||||
with zipfile.ZipFile(io.BytesIO(zip_data)) as zf:
|
||||
for name in zf.namelist():
|
||||
if name == "mapping.csv":
|
||||
csv_text = zf.read(name).decode("utf-8")
|
||||
lines = csv_text.strip().split("\n")
|
||||
if len(lines) > 1:
|
||||
for line in lines[1:]:
|
||||
parts = line.split(",", 2)
|
||||
if len(parts) >= 3:
|
||||
all_mapping[parts[1]] = parts[2]
|
||||
elif not name.endswith("/"):
|
||||
all_results.append((name, zf.read(name)))
|
||||
# Отправляем: закончили файл (если клиент ушёл — освобождаем поток)
|
||||
try:
|
||||
yield f"event: done\ndata: {json.dumps({'idx': idx, 'name': fname, 'total': total})}\n\n"
|
||||
except GeneratorExit:
|
||||
return # клиент отключился, НЕ обрабатываем остальные файлы
|
||||
def progress(phase, idx, total_, name):
|
||||
q.put(("progress", phase, idx, name, total_))
|
||||
|
||||
csv_str = build_mapping_csv(all_mapping) if all_mapping else ""
|
||||
final_zip = build_zip(all_results)
|
||||
store_result(sid, final_zip)
|
||||
if csv_str:
|
||||
store_csv(sid, csv_str)
|
||||
yield f"event: complete\ndata: {json.dumps({'total': len(all_results)})}\n\n"
|
||||
def worker():
|
||||
try:
|
||||
zip_data, csv_str = obfuscate_files(
|
||||
all_files, llm_client=llm, progress_cb=progress
|
||||
)
|
||||
q.put(("result", zip_data, csv_str))
|
||||
except Exception as e:
|
||||
q.put(("error", repr(e)))
|
||||
|
||||
threading.Thread(target=worker, daemon=True).start()
|
||||
|
||||
while True:
|
||||
try:
|
||||
evt = q.get(timeout=1)
|
||||
except queue.Empty:
|
||||
if cancel.is_set():
|
||||
return
|
||||
continue
|
||||
|
||||
kind = evt[0]
|
||||
|
||||
if kind == "progress":
|
||||
_, phase, idx, name, total_ = evt
|
||||
try:
|
||||
yield (
|
||||
f"event: {phase}\n"
|
||||
f"data: {json.dumps({'idx': idx, 'name': name, 'total': total_})}\n\n"
|
||||
)
|
||||
except _disconnect_exceptions():
|
||||
cancel.set()
|
||||
return
|
||||
|
||||
elif kind == "result":
|
||||
_, zip_data, csv_str = evt
|
||||
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"])
|
||||
try:
|
||||
yield f"event: complete\ndata: {json.dumps({'total': count})}\n\n"
|
||||
except _disconnect_exceptions():
|
||||
return
|
||||
return
|
||||
|
||||
elif kind == "error":
|
||||
_, msg = evt
|
||||
try:
|
||||
yield f"event: error\ndata: {json.dumps({'error': msg})}\n\n"
|
||||
except _disconnect_exceptions():
|
||||
return
|
||||
return
|
||||
|
||||
return Response(
|
||||
stream_with_context(generate()),
|
||||
@@ -99,7 +145,7 @@ def process_stream(sid):
|
||||
|
||||
@api_bp.route("/process/<sid>", methods=["POST"])
|
||||
def process(sid):
|
||||
"""Process all session files one by one -> ZIP."""
|
||||
"""Process all session files -> ZIP (единым вызовом, общий mapping)."""
|
||||
files = get_files(sid)
|
||||
if files is None:
|
||||
return jsonify({"ok": False, "error": "Session not found"}), 404
|
||||
@@ -107,28 +153,12 @@ def process(sid):
|
||||
return jsonify({"ok": False, "error": "No files"}), 400
|
||||
try:
|
||||
llm = LLMClient()
|
||||
all_results = []
|
||||
all_mapping = {}
|
||||
for fname, content in files:
|
||||
zip_data, _csv = obfuscate_files([(fname, content, "")], llm_client=llm)
|
||||
with zipfile.ZipFile(io.BytesIO(zip_data)) as zf:
|
||||
for name in zf.namelist():
|
||||
if name == "mapping.csv":
|
||||
csv_text = zf.read(name).decode("utf-8")
|
||||
lines = csv_text.strip().split("\n")
|
||||
if len(lines) > 1:
|
||||
for line in lines[1:]:
|
||||
parts = line.split(",", 2)
|
||||
if len(parts) >= 3:
|
||||
all_mapping[parts[1]] = parts[2]
|
||||
elif not name.endswith("/"):
|
||||
all_results.append((name, zf.read(name)))
|
||||
csv_str = build_mapping_csv(all_mapping) if all_mapping else ""
|
||||
final_zip = build_zip(all_results) # CSV — отдельно, не в ZIP
|
||||
store_result(sid, final_zip)
|
||||
all_files = [(fname, content, "") for fname, content in files]
|
||||
zip_data, csv_str = obfuscate_files(all_files, llm_client=llm)
|
||||
store_result(sid, zip_data)
|
||||
if csv_str:
|
||||
store_csv(sid, csv_str)
|
||||
return jsonify({"ok": True, "status": "done", "files": len(all_results)})
|
||||
return jsonify({"ok": True, "status": "done"})
|
||||
except Exception as e:
|
||||
traceback.print_exc()
|
||||
return jsonify({"ok": False, "error": str(e)}), 500
|
||||
|
||||
Reference in New Issue
Block a user