"""routes/parse.py — SSE-парсинг: потоковый разбор документов договора.""" import json as _json import time as _time import io as _io import zipfile as _zipfile import base64 as _b64 from flask import Response import db import parser as parser_mod import textify import mimeutil def parse(cid): """ GET /parse/ — Server-Sent Events поток. Для каждого файла договора: 1. Достаёт байты из БД (original_bytes или original_b64) 2. Парсит через parser.parse() 3. Сохраняет parsed_text + elements_json в БД 4. Отправляет SSE-сообщения: file_start → file_done → summary ZIP-файлы распаковываются, каждый внутренний файл парсится отдельно. """ def generate(): start_time = _time.time() total_bytes = 0 files_processed = 0 # ── Получить все файлы договора ───────────────────── supps, _ = db.query( "SELECT s.id, d.id as doc_id, d.filename, d.mime_type, " "LENGTH(d.original_bytes) as file_size " "FROM supplements s JOIN documents d ON s.document_id = d.id " "WHERE s.contract_id = %s", (cid,), ) if not supps: yield f"data: {_json.dumps({'type': 'error', 'message': 'Нет файлов'})}\n\n" return supp_rows = [dict(zip(supps["columns"], r)) for r in supps["rows"]] # ── Вспомогательная: собрать результат парсинга ───── def _file_done(name, elapsed, elements, errors, text): """ Сформировать SSE-сообщение file_done со статистикой: количество параграфов, таблиц, строк таблиц, заголовки таблиц. """ paragraphs = sum(1 for e in elements if e.get("type") == "paragraph") tables = sum(1 for e in elements if e.get("type") == "table") table_rows = sum(len(e.get("rows", [])) for e in elements if e.get("type") == "table") table_headers = [ e["rows"][0] if e.get("rows") else [] for e in elements if e.get("type") == "table" ] return { "type": "file_done", "name": name, "time_s": elapsed, "elements": len(elements), "paragraphs": paragraphs, "tables": tables, "table_rows": table_rows, "table_headers": table_headers, "errors": errors, "text_len": len(text) if text else 0, "text_preview": text[:200] if text else "", } # ── Основной цикл: каждый файл ────────────────────── for s in supp_rows: # Пропустить уже распарсенные existing, _ = db.query( "SELECT 1 FROM documents WHERE id = %s AND status IN ('parsed', 'expanded')", (s["doc_id"],), ) if existing and existing["rows"]: continue file_start = _time.time() file_bytes = s.get("file_size", 0) or 0 total_bytes += file_bytes # SSE: начало обработки файла yield f"data: {_json.dumps({'type': 'file_start', 'name': s['filename'], 'bytes': file_bytes})}\n\n" # ── Получить байты файла ──────────────────────── doc, _ = db.query_one( "SELECT original_bytes, original_b64 FROM documents WHERE id = %s", (s["doc_id"],), ) raw = doc.get("original_bytes") if doc else None b64 = doc.get("original_b64") if doc else None if raw: # Обычный файл: original_bytes file_data = bytes(raw) if isinstance(raw, memoryview) else raw elif b64: # Старый формат: base64 file_data = _b64.b64decode(b64) else: yield f"data: {_json.dumps({'type': 'file_error', 'name': s['filename'], 'error': 'Нет данных'})}\n\n" continue # ── ZIP: распаковать и парсить каждый файл ────── if s["mime_type"] == "application/zip": zip_names = [] try: with _zipfile.ZipFile(_io.BytesIO(file_data)) as zf: for zname in zf.namelist(): if zname.endswith("/"): continue # пропускаем папки zdata = zf.read(zname) zmime = mimeutil.guess_mime(zname) or "application/octet-stream" # Только PDF и Word внутри ZIP if zmime not in ( "application/pdf", "application/vnd.openxmlformats-officedocument.wordprocessingml.document", "application/msword", ): continue zip_names.append(zname) # Сохранить как отдельный документ db.execute( "INSERT INTO documents (filename, mime_type, original_bytes, status) VALUES (%s,%s,%s,'uploaded')", (zname, zmime, zdata), ) zdoc, _ = db.query_one( "SELECT id FROM documents WHERE filename = %s ORDER BY created_at DESC LIMIT 1", (zname,), ) zdoc_id = zdoc["id"] if zdoc else None if zdoc_id: db.execute( "INSERT INTO supplements (contract_id, document_id, type) VALUES (%s,%s,'initial')", (cid, str(zdoc_id)), ) # Парсинг внутреннего файла zf_start = _time.time() total_bytes += len(zdata) yield f"data: {_json.dumps({'type': 'file_start', 'name': zname, 'bytes': len(zdata)})}\n\n" try: zpr = parser_mod.parse(zdata, zmime) zelements = zpr.get("elements", []) ztext = textify.to_text(zelements) if zelements else "" db.execute( "UPDATE documents SET parsed_text = %s, elements_json = %s, status = 'parsed' WHERE id = %s", (ztext, _json.dumps(zelements, ensure_ascii=False), str(zdoc_id)), ) zelapsed = round(_time.time() - zf_start, 2) files_processed += 1 yield f"data: {_json.dumps(_file_done(zname, zelapsed, zelements, zpr.get('errors', []), ztext))}\n\n" except Exception as e: yield f"data: {_json.dumps({'type': 'file_error', 'name': zname, 'error': str(e)})}\n\n" except Exception as e: yield f"data: {_json.dumps({'type': 'file_error', 'name': s['filename'], 'error': 'ZIP: ' + str(e)})}\n\n" # Пометить ZIP как expanded db.execute("UPDATE documents SET status = 'expanded' WHERE id = %s", (s["doc_id"],)) yield f"data: {_json.dumps({'type': 'zip_expanded', 'name': s['filename'], 'count': len(zip_names)})}\n\n" else: # ── Обычный файл (PDF/DOCX) ───────────────── pr = parser_mod.parse(file_data, s["mime_type"]) elements = pr.get("elements", []) text = textify.to_text(elements) if elements else "" # Сохранить результат в БД db.execute( "UPDATE documents SET parsed_text = %s, elements_json = %s, status = 'parsed' WHERE id = %s", (text, _json.dumps(elements, ensure_ascii=False), str(s["doc_id"])), ) elapsed = round(_time.time() - file_start, 2) files_processed += 1 yield f"data: {_json.dumps(_file_done(s['filename'], elapsed, elements, pr.get('errors', []), text))}\n\n" # ── Итоговое сообщение ────────────────────────────── total_time = round(_time.time() - start_time, 2) yield f"data: {_json.dumps({'type': 'summary', 'total_time_s': total_time, 'total_bytes': total_bytes, 'files_processed': files_processed})}\n\n" return Response(generate(), mimetype="text/event-stream")