From fa4c7fe3789e389f62326bba09fa50fadb11d9f1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Tue, 16 Jun 2026 14:34:00 +0400 Subject: [PATCH] =?UTF-8?q?SSE-=D0=BE=D0=B1=D1=80=D0=B0=D0=B1=D0=BE=D1=82?= =?UTF-8?q?=D0=BA=D0=B0:=20=D0=B8=D0=BD=D1=82=D0=B5=D1=80=D0=B0=D0=BA?= =?UTF-8?q?=D1=82=D0=B8=D0=B2=D0=BD=D1=8B=D0=B9=20=D0=BF=D1=80=D0=BE=D0=B3?= =?UTF-8?q?=D1=80=D0=B5=D1=81=D1=81=20=D1=81=20=D1=82=D0=B0=D0=B9=D0=BC?= =?UTF-8?q?=D0=B5=D1=80=D0=B0=D0=BC=D0=B8,=20=D0=B8=D1=82=D0=BE=D0=B3?= =?UTF-8?q?=D0=B8=20=D0=BF=D0=BE=20=D0=B2=D1=80=D0=B5=D0=BC=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8/=D0=B1=D0=B0=D0=B9=D1=82=D0=B0=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- site/app.py | 138 ++++++++++++++++++++++--------------- site/templates/upload.html | 89 +++++++++++++++++++++--- 2 files changed, 161 insertions(+), 66 deletions(-) diff --git a/site/app.py b/site/app.py index 3a90a89..5f5a8d5 100644 --- a/site/app.py +++ b/site/app.py @@ -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/", "process", self._process, methods=["POST"]) + self.app.add_url_rule("/process/", "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) diff --git a/site/templates/upload.html b/site/templates/upload.html index 7e755ed..35a75ac 100644 --- a/site/templates/upload.html +++ b/site/templates/upload.html @@ -76,12 +76,11 @@ {% endfor %} {% if unprocessed %} -
- -
-
+ {% endif %} @@ -118,11 +117,83 @@