268 lines
11 KiB
Python
Executable File
268 lines
11 KiB
Python
Executable File
#!/usr/bin/env python3
|
||
"""HTTP-сервер: конвертация .doc → .docx + LLM-анализ ДС (/llm-ops) + SSE v2 (/process-v2)"""
|
||
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||
from socketserver import ThreadingMixIn
|
||
|
||
class ThreadingHTTPServer(ThreadingMixIn, HTTPServer):
|
||
"""Многопоточный HTTP-сервер."""
|
||
daemon_threads = True
|
||
from urllib.parse import urlparse, parse_qs
|
||
import subprocess, tempfile, os, sys, json, time, re
|
||
|
||
import httpx
|
||
from llm_prompt import build_prompt
|
||
|
||
LLM_URL = "https://api.aillm.ru/v1/chat/completions"
|
||
LLM_KEY = os.environ.get("LLM_KEY", "")
|
||
LLM_MODEL = "gpt-oss-120b"
|
||
LUCEE_URL = "https://contractor.luceek8s.dev.nubes.ru"
|
||
|
||
|
||
def call_llm(current_spec, doc_text):
|
||
"""Вызов LLM API. Возвращает {mode, ops} или кидает исключение."""
|
||
prompt = build_prompt(current_spec, doc_text)
|
||
payload = {
|
||
"model": LLM_MODEL,
|
||
"messages": [{"role": "user", "content": prompt}],
|
||
"max_tokens": 8000,
|
||
"temperature": 0.1,
|
||
}
|
||
with httpx.Client(http2=True, timeout=120, verify=False) as client:
|
||
resp = client.post(LLM_URL, json=payload,
|
||
headers={"Authorization": f"Bearer {LLM_KEY}", "Content-Type": "application/json"})
|
||
resp.raise_for_status()
|
||
data = resp.json()
|
||
|
||
raw_text = data.get("choices", [{}])[0].get("message", {}).get("content", "")
|
||
json_text = raw_text
|
||
if "```json" in json_text:
|
||
json_text = json_text.split("```json")[1].split("```")[0]
|
||
elif "```" in json_text:
|
||
json_text = json_text.split("```")[1].split("```")[0]
|
||
return json.loads(json_text.strip())
|
||
|
||
|
||
class Handler(BaseHTTPRequestHandler):
|
||
def do_GET(self):
|
||
parsed = urlparse(self.path)
|
||
if parsed.path == "/process-v2":
|
||
self._handle_process_v2(parsed)
|
||
else:
|
||
self.send_error(404)
|
||
|
||
def do_OPTIONS(self):
|
||
self.send_response(200)
|
||
self.send_header("Access-Control-Allow-Origin", "*")
|
||
self.send_header("Access-Control-Allow-Methods", "GET, OPTIONS")
|
||
self.send_header("Access-Control-Allow-Headers", "Content-Type")
|
||
self.end_headers()
|
||
|
||
def do_POST(self):
|
||
if self.path == "/llm-ops":
|
||
self._handle_llm_ops()
|
||
else:
|
||
self._handle_doc_convert()
|
||
|
||
def _sse(self, data):
|
||
"""Отправить SSE-событие."""
|
||
self.wfile.write(f"data: {json.dumps(data, ensure_ascii=False)}\n\n".encode())
|
||
self.wfile.flush()
|
||
|
||
def _lucee_query(self, sql):
|
||
"""Запрос к Lucee API."""
|
||
resp = httpx.get(f"{LUCEE_URL}/api.cfm", params={"action": "query", "sql": sql}, timeout=30)
|
||
resp.raise_for_status()
|
||
d = resp.json()
|
||
if not d.get("OK"):
|
||
raise Exception(d.get("ERROR", "API error"))
|
||
cols = d["COLUMNS"]
|
||
rows = []
|
||
for r in d["ROWS"]:
|
||
row = {}
|
||
for i, c in enumerate(cols):
|
||
row[c.lower()] = r[i] if isinstance(r, list) else r[c]
|
||
rows.append(row)
|
||
return rows
|
||
|
||
def _handle_process_v2(self, parsed):
|
||
"""SSE-стриминг v2 пайплайна."""
|
||
params = parse_qs(parsed.query)
|
||
cid = params.get("contract_id", [None])[0]
|
||
if not cid:
|
||
self.send_error(400, "contract_id required")
|
||
return
|
||
# Validate UUID to prevent SQL injection
|
||
if not re.fullmatch(r'[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}', cid, re.I):
|
||
self.send_error(400, "invalid contract_id format")
|
||
return
|
||
|
||
self.send_response(200)
|
||
self.send_header("Content-Type", "text/event-stream; charset=utf-8")
|
||
self.send_header("Cache-Control", "no-cache")
|
||
self.send_header("Connection", "keep-alive")
|
||
self.end_headers()
|
||
self.wfile.write(b": ok\n\n")
|
||
self.wfile.flush()
|
||
|
||
t0 = time.time()
|
||
|
||
try:
|
||
# 1. Список supplements
|
||
supps = self._lucee_query(
|
||
f"SELECT s.id, s.type, d.filename FROM supplements s "
|
||
f"JOIN documents d ON s.document_id=d.id "
|
||
f"WHERE s.contract_id='{cid}' AND d.elements_json IS NOT NULL ORDER BY s.created_at"
|
||
)
|
||
if not supps:
|
||
self._sse({"type": "error", "message": "Нет распарсенных файлов"})
|
||
return
|
||
|
||
for s in supps:
|
||
sid = s["id"]
|
||
|
||
# Текущая спецификация
|
||
cur = self._lucee_query(
|
||
f"SELECT name_hash, name, price, qty, sum, date_start "
|
||
f"FROM spec_current WHERE contract_id='{cid}' ORDER BY name"
|
||
)
|
||
current_spec = []
|
||
for r in cur:
|
||
current_spec.append({
|
||
"hash": r["name_hash"],
|
||
"name": r["name"],
|
||
"price": float(r["price"]) if r.get("price") else None,
|
||
"qty": float(r["qty"]) if r.get("qty") else None,
|
||
"sum": float(r["sum"]) if r.get("sum") else None,
|
||
"date_start": r["date_start"],
|
||
})
|
||
|
||
# elements_json → текст
|
||
docs = self._lucee_query(
|
||
f"SELECT elements_json FROM documents d "
|
||
f"JOIN supplements s ON s.document_id=d.id WHERE s.id='{sid}'"
|
||
)
|
||
if not docs:
|
||
continue
|
||
ej = docs[0]["elements_json"]
|
||
if isinstance(ej, dict) and "Value" in ej:
|
||
ej = ej["Value"]
|
||
elements = json.loads(ej)
|
||
# Lucee serializeJSON → UPPERCASE keys, normalize to lowercase
|
||
elements = [{k.lower(): v for k, v in el.items()} for el in elements]
|
||
lines = []
|
||
for el in elements:
|
||
if el.get("type") == "paragraph":
|
||
prefix = f"[{el['style']}] " if el.get("style") else ""
|
||
lines.append(prefix + el.get("text", ""))
|
||
elif el.get("type") == "table":
|
||
rows = el.get("rows", [])
|
||
if rows:
|
||
ncols = len(rows[0])
|
||
lines.append(f"--- Таблица ({len(rows)}×{ncols}) ---")
|
||
for row in rows:
|
||
cells = [str(c).replace("\n", " ").replace("|", "\\|") for c in row[:ncols]]
|
||
lines.append("| " + " | ".join(cells) + " |")
|
||
lines.append("")
|
||
doc_text = "\n".join(lines)
|
||
|
||
self._sse({"type": "extract_start", "supplement_id": sid, "filename": s["filename"]})
|
||
|
||
# LLM (прямой вызов, не HTTP self-call)
|
||
import threading
|
||
t1 = time.time()
|
||
done = threading.Event()
|
||
result = [None]
|
||
error = [None]
|
||
|
||
def do_llm():
|
||
try:
|
||
result[0] = call_llm(current_spec, doc_text)
|
||
except Exception as e:
|
||
error[0] = str(e)
|
||
finally:
|
||
done.set()
|
||
|
||
t = threading.Thread(target=do_llm)
|
||
t.start()
|
||
|
||
# Keepalive каждые 15с пока LLM думает
|
||
while not done.wait(15):
|
||
self._sse({"type": "keepalive"})
|
||
|
||
t.join()
|
||
elapsed = round(time.time() - t1, 1)
|
||
|
||
if error[0]:
|
||
self._sse({"type": "extract_error", "supplement_id": sid, "filename": s["filename"],
|
||
"error": error[0], "time_s": elapsed})
|
||
continue
|
||
|
||
llm_result = result[0]
|
||
ops = llm_result.get("ops", [])
|
||
mode = llm_result.get("mode", "partial")
|
||
self._sse({"type": "llm_done", "supplement_id": sid, "filename": s["filename"],
|
||
"ops_count": len(ops), "mode": mode, "time_s": elapsed})
|
||
|
||
# Применить
|
||
apply_resp = httpx.post(
|
||
f"{LUCEE_URL}/apply_events.cfm?contract_id={cid}&mode={mode}",
|
||
json={"supplement_id": sid, "ops": ops}, timeout=30
|
||
)
|
||
if apply_resp.status_code == 200:
|
||
ar = apply_resp.json()
|
||
# Lucee serializeJSON → UPPERCASE keys
|
||
summary = {k.lower(): v for k, v in ar.get("SUMMARY", {}).items()}
|
||
self._sse({"type": "applied", "supplement_id": sid, "summary": summary, "ops": ops})
|
||
else:
|
||
self._sse({"type": "apply_error", "supplement_id": sid,
|
||
"error": f"apply HTTP {apply_resp.status_code}"})
|
||
|
||
total_time = round(time.time() - t0, 1)
|
||
self._sse({"type": "done", "total_time_s": total_time})
|
||
|
||
except Exception as e:
|
||
self._sse({"type": "error", "message": str(e)})
|
||
|
||
def _handle_llm_ops(self):
|
||
length = int(self.headers.get("Content-Length", 0))
|
||
body = json.loads(self.rfile.read(length))
|
||
try:
|
||
result = call_llm(body.get("current_spec", []), body.get("doc_text", ""))
|
||
self.send_response(200)
|
||
self.send_header("Content-Type", "application/json")
|
||
self.end_headers()
|
||
self.wfile.write(json.dumps(result, ensure_ascii=False).encode())
|
||
except Exception as e:
|
||
self.send_response(502)
|
||
self.send_header("Content-Type", "application/json")
|
||
self.end_headers()
|
||
self.wfile.write(json.dumps({"error": str(e)}, ensure_ascii=False).encode())
|
||
|
||
def _handle_doc_convert(self):
|
||
length = int(self.headers.get("Content-Length", 0))
|
||
data = self.rfile.read(length)
|
||
with tempfile.NamedTemporaryFile(suffix=".doc", delete=False) as f:
|
||
f.write(data); doc_path = f.name
|
||
tmpdir = tempfile.mkdtemp()
|
||
try:
|
||
subprocess.run(["libreoffice", "--headless", "--convert-to", "docx", "--outdir", tmpdir, doc_path],
|
||
timeout=30, capture_output=True)
|
||
docx = [x for x in os.listdir(tmpdir) if x.endswith(".docx")]
|
||
if docx:
|
||
with open(os.path.join(tmpdir, docx[0]), "rb") as f:
|
||
out = f.read()
|
||
self.send_response(200)
|
||
self.send_header("Content-Type", "application/vnd.openxmlformats-officedocument.wordprocessingml.document")
|
||
self.send_header("Content-Length", len(out))
|
||
self.end_headers()
|
||
self.wfile.write(out)
|
||
else:
|
||
self.send_error(500, "Conversion failed")
|
||
finally:
|
||
os.unlink(doc_path)
|
||
for x in os.listdir(tmpdir): os.unlink(os.path.join(tmpdir, x))
|
||
os.rmdir(tmpdir)
|
||
def log_message(self, *a): pass
|
||
|
||
ThreadingHTTPServer(("127.0.0.1", 8766), Handler).serve_forever()
|