Files

194 lines
9.5 KiB
Python

"""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/<cid> — 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")