SSE-обработка: интерактивный прогресс с таймерами, итоги по времени/байтам
This commit is contained in:
+81
-57
@@ -1,7 +1,7 @@
|
||||
"""app.py — Точка входа и сборка слоёв."""
|
||||
|
||||
from dotenv import load_dotenv
|
||||
from flask import Flask, render_template, request, redirect, url_for
|
||||
from flask import Flask, render_template, request, redirect, url_for, Response
|
||||
import json
|
||||
|
||||
import schema
|
||||
@@ -26,7 +26,7 @@ class ContractsApp:
|
||||
|
||||
def _add_routes(self):
|
||||
self.app.add_url_rule("/", "index", self._index, methods=["GET", "POST"])
|
||||
self.app.add_url_rule("/process/<cid>", "process", self._process, methods=["POST"])
|
||||
self.app.add_url_rule("/process/<cid>", "process", self._process, methods=["GET"])
|
||||
self.app.add_url_rule("/health", "health", self._health)
|
||||
self.app.register_blueprint(test_bp)
|
||||
self.app.register_blueprint(upload_bp)
|
||||
@@ -132,70 +132,94 @@ class ContractsApp:
|
||||
contract=contract, supplements=supp_list,
|
||||
all_rows=all_rows, unprocessed=unprocessed)
|
||||
|
||||
# ── Шаг 2: LLM-обработка ────────────────────────────────────
|
||||
# ── Шаг 2: LLM-обработка (SSE) ──────────────────────────────
|
||||
|
||||
def _process(self, cid):
|
||||
"""Запустить LLM-обработку для всех необработанных допников договора."""
|
||||
supps, _ = db.query(
|
||||
"SELECT s.id, d.id as doc_id, d.parsed_text, d.filename "
|
||||
"FROM supplements s JOIN documents d ON s.document_id=d.id "
|
||||
"WHERE s.contract_id=%s AND d.parsed_text IS NOT NULL",
|
||||
(cid,),
|
||||
)
|
||||
if not supps:
|
||||
return redirect("/?id=" + cid)
|
||||
"""SSE-поток: обработка допников с интерактивным прогрессом."""
|
||||
import time as time_mod
|
||||
|
||||
supp_rows = [dict(zip(supps["columns"], r)) for r in supps["rows"]]
|
||||
total_rows = 0
|
||||
def generate():
|
||||
start_time = time_mod.time()
|
||||
total_bytes = 0
|
||||
files_processed = 0
|
||||
total_rows = 0
|
||||
|
||||
for s in supp_rows:
|
||||
text = s["parsed_text"]
|
||||
if not text:
|
||||
continue
|
||||
|
||||
# Проверим, не обработан ли уже
|
||||
existing, _ = db.query("SELECT 1 FROM spec_rows WHERE supplement_id=%s LIMIT 1", (s["id"],))
|
||||
if existing and existing["rows"]:
|
||||
continue
|
||||
|
||||
ext = extractor.extract(text)
|
||||
if "error" in ext:
|
||||
db.execute("UPDATE documents SET status='extract_error', error_message=%s WHERE id=%s",
|
||||
(ext["error"], s["doc_id"]))
|
||||
continue
|
||||
|
||||
rows = ext.get("rows", [])
|
||||
total_rows += len(rows)
|
||||
|
||||
for row in rows:
|
||||
db.execute(
|
||||
"INSERT INTO spec_rows (supplement_id,row_num,name,price,qty,sum,date_start) VALUES (%s,%s,%s,%s,%s,%s,%s)",
|
||||
(s["id"], row.get("row_num"), row.get("name"),
|
||||
row.get("price"), row.get("qty"), row.get("sum"), row.get("date_start")),
|
||||
)
|
||||
|
||||
# Diff
|
||||
prev_res, _ = db.query(
|
||||
"SELECT sr.row_num,sr.name,sr.price,sr.qty,sr.sum,sr.date_start "
|
||||
"FROM spec_rows sr JOIN supplements sup ON sr.supplement_id=sup.id "
|
||||
"WHERE sup.contract_id=%s AND sup.id!=%s ORDER BY sr.row_num",
|
||||
(cid, s["id"]),
|
||||
supps, _ = db.query(
|
||||
"SELECT s.id, d.id as doc_id, d.parsed_text, d.filename, "
|
||||
"LENGTH(d.original_bytes) as file_size "
|
||||
"FROM supplements s JOIN documents d ON s.document_id=d.id "
|
||||
"WHERE s.contract_id=%s AND d.parsed_text IS NOT NULL",
|
||||
(cid,),
|
||||
)
|
||||
prev_rows = [dict(zip(prev_res["columns"], r)) for r in prev_res["rows"]] if prev_res else []
|
||||
diff_r = differ.diff(prev_rows, rows)
|
||||
if not supps:
|
||||
yield f"data: {json.dumps({'type': 'error', 'message': 'Нет файлов для обработки'})}\n\n"
|
||||
return
|
||||
|
||||
for ch in diff_r.get("changes", []):
|
||||
db.execute(
|
||||
"INSERT INTO spec_history (contract_id,supplement_id,row_num,change_type,old_values,new_values) "
|
||||
"VALUES (%s,%s,%s,%s,%s,%s)",
|
||||
(cid, s["id"], ch["row_num"], ch["change_type"],
|
||||
json.dumps(ch.get("old_values"), ensure_ascii=False, default=str) if ch.get("old_values") else None,
|
||||
json.dumps(ch.get("new_values"), ensure_ascii=False, default=str) if ch.get("new_values") else None),
|
||||
supp_rows = [dict(zip(supps["columns"], r)) for r in supps["rows"]]
|
||||
|
||||
for s in supp_rows:
|
||||
text = s["parsed_text"]
|
||||
if not text:
|
||||
continue
|
||||
|
||||
# Проверим, не обработан ли уже
|
||||
existing, _ = db.query("SELECT 1 FROM spec_rows WHERE supplement_id=%s LIMIT 1", (s["id"],))
|
||||
if existing and existing["rows"]:
|
||||
continue
|
||||
|
||||
file_start = time_mod.time()
|
||||
file_bytes = s.get("file_size", 0) or 0
|
||||
total_bytes += file_bytes
|
||||
|
||||
yield f"data: {json.dumps({'type': 'file_start', 'name': s['filename'], 'bytes': file_bytes})}\n\n"
|
||||
|
||||
ext = extractor.extract(text)
|
||||
if "error" in ext:
|
||||
db.execute("UPDATE documents SET status='extract_error', error_message=%s WHERE id=%s",
|
||||
(ext["error"], s["doc_id"]))
|
||||
yield f"data: {json.dumps({'type': 'file_error', 'name': s['filename'], 'error': ext['error']})}\n\n"
|
||||
continue
|
||||
|
||||
rows = ext.get("rows", [])
|
||||
total_rows += len(rows)
|
||||
|
||||
for row in rows:
|
||||
db.execute(
|
||||
"INSERT INTO spec_rows (supplement_id,row_num,name,price,qty,sum,date_start) VALUES (%s,%s,%s,%s,%s,%s,%s)",
|
||||
(s["id"], row.get("row_num"), row.get("name"),
|
||||
row.get("price"), row.get("qty"), row.get("sum"), row.get("date_start")),
|
||||
)
|
||||
|
||||
# Diff
|
||||
prev_res, _ = db.query(
|
||||
"SELECT sr.row_num,sr.name,sr.price,sr.qty,sr.sum,sr.date_start "
|
||||
"FROM spec_rows sr JOIN supplements sup ON sr.supplement_id=sup.id "
|
||||
"WHERE sup.contract_id=%s AND sup.id!=%s ORDER BY sr.row_num",
|
||||
(cid, s["id"]),
|
||||
)
|
||||
prev_rows = [dict(zip(prev_res["columns"], r)) for r in prev_res["rows"]] if prev_res else []
|
||||
diff_r = differ.diff(prev_rows, rows)
|
||||
|
||||
db.execute("UPDATE documents SET status='extracted' WHERE id=%s", (s["doc_id"],))
|
||||
for ch in diff_r.get("changes", []):
|
||||
db.execute(
|
||||
"INSERT INTO spec_history (contract_id,supplement_id,row_num,change_type,old_values,new_values) "
|
||||
"VALUES (%s,%s,%s,%s,%s,%s)",
|
||||
(cid, s["id"], ch["row_num"], ch["change_type"],
|
||||
json.dumps(ch.get("old_values"), ensure_ascii=False, default=str) if ch.get("old_values") else None,
|
||||
json.dumps(ch.get("new_values"), ensure_ascii=False, default=str) if ch.get("new_values") else None),
|
||||
)
|
||||
|
||||
return redirect("/?id=" + cid)
|
||||
db.execute("UPDATE documents SET status='extracted' WHERE id=%s", (s["doc_id"],))
|
||||
|
||||
elapsed = round(time_mod.time() - file_start, 1)
|
||||
files_processed += 1
|
||||
|
||||
yield f"data: {json.dumps({'type': 'file_done', 'name': s['filename'], 'rows': len(rows), 'time_s': elapsed})}\n\n"
|
||||
|
||||
total_time = round(time_mod.time() - start_time, 1)
|
||||
yield f"data: {json.dumps({'type': 'summary', 'total_time_s': total_time, 'total_bytes': total_bytes, 'total_rows': total_rows, 'files_processed': files_processed})}\n\n"
|
||||
|
||||
return Response(generate(), mimetype="text/event-stream")
|
||||
|
||||
def run(self):
|
||||
self.app.run(host="0.0.0.0", port=5000)
|
||||
|
||||
Reference in New Issue
Block a user