#!/usr/bin/env python3 """Contracts VM server — thin HTTP router. All logic in db/ and services/.""" from http.server import HTTPServer, BaseHTTPRequestHandler from socketserver import ThreadingMixIn, TCPServer from urllib.parse import urlparse, parse_qs import json, re, os # ── DB auto-seed ────────────────────────────────────────────────────────── from db import prompts as db_prompts from db.connection import DB_CONFIG db_prompts.seed_defaults() # ── DB modules ──────────────────────────────────────────────────────────── from db import supplements as db_supplements from db import documents as db_documents from db import spec_current as db_spec_current # ── Services ────────────────────────────────────────────────────────────── from services.upload import handle_upload from services.unzip import handle_unzip from services.process import run_pipeline from llm_prompt import build_prompt class ThreadingHTTPServer(ThreadingMixIn, HTTPServer): daemon_threads = True class Handler(BaseHTTPRequestHandler): def do_OPTIONS(self): self.send_response(200) self._send_cors() self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS") self.send_header("Access-Control-Allow-Headers", "Content-Type") self.end_headers() def do_GET(self): parsed = urlparse(self.path) if parsed.path == "/process-v2": self._handle_process_v2(parsed) elif parsed.path == "/health": self._json({"ok": True, "db": DB_CONFIG["dbname"]}) elif parsed.path == "/app.js" or parsed.path.startswith("/static/"): self._handle_app_js() elif parsed.path == "/api/supplements": self._handle_api_supplements(parsed) elif parsed.path == "/api/groups": self._handle_api_groups(parsed) elif parsed.path == "/api/batch-progress": self._handle_api_batch_progress(parsed) elif parsed.path.startswith("/api/documents/"): self._handle_api_document(parsed) else: self.send_error(404) def do_POST(self): if self.path == "/upload": self._handle_upload() elif self.path == "/unzip-upload": self._handle_unzip() elif self.path == "/llm-ops": self._handle_llm_ops() elif self.path == "/convert-doc": self._handle_convert_doc() elif self.path == "/api/classify-batch": self._handle_classify_batch() elif self.path == "/api/apply-groups": self._handle_apply_groups() elif self.path == "/api/cleanup": self._handle_cleanup() else: self.send_error(404) # ── /process-v2 (SSE) ───────────────────────────────────────────────── def _handle_process_v2(self, parsed): params = parse_qs(parsed.query) cid = params.get("contract_id", [None])[0] if not cid: self.send_error(400, "contract_id required") return 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() order_ids = params.get("order", [None])[0] or "" try: run_pipeline(cid, order_ids, self._sse, build_prompt) except Exception as e: self._sse({"type": "error", "message": str(e)}) def _sse(self, data): self.wfile.write(f"data: {json.dumps(data, ensure_ascii=False)}\n\n".encode()) self.wfile.flush() # ── /api/supplements ────────────────────────────────────────────────── def _handle_api_supplements(self, parsed): params = parse_qs(parsed.query) cid = params.get("contract_id", [None])[0] if not cid: self._json({"ok": False, "error": "contract_id required"}, 400) return rows = db_supplements.list_by_contract(cid) self._json({"ok": True, "supplements": rows}) # ── /api/documents/ ────────────────────────────────────────────── def _handle_api_document(self, parsed): doc_id = parsed.path.split("/")[-1] doc = db_documents.get(doc_id) if not doc: self._json({"ok": False, "error": "not found"}, 404) return self._json({ "ok": True, "id": doc["id"], "filename": doc["filename"], "status": doc["status"], "elements_json": doc.get("elements_json"), }) # ── /api/groups ?batch=X ───────────────────────────────────────────── def _handle_api_groups(self, parsed): params = parse_qs(parsed.query) batch_id = params.get("batch", [None])[0] if not batch_id: self._json({"ok": False, "error": "batch required"}, 400) return from services.grouping import group_documents result = group_documents(batch_id) self._json(result) # ── /api/batch-progress ?batch=X ───────────────────────────────────── def _handle_api_batch_progress(self, parsed): params = parse_qs(parsed.query) batch_id = params.get("batch", [None])[0] if not batch_id: self._json({"ok": False, "error": "batch required"}, 400) return counts = db_documents.count_by_status(batch_id) docs = db_documents.list_by_batch(batch_id) self._json({"ok": True, "counts": counts, "total": len(docs)}) # ── POST /classify-batch ───────────────────────────────────────────── def _handle_classify_batch(self): length = int(self.headers.get("Content-Length", 0)) body = json.loads(self.rfile.read(length)) batch_id = body.get("batch_id") if not batch_id: self._json({"ok": False, "error": "batch_id required"}, 400) return from services.classify import classify_batch result = classify_batch(batch_id) self._json(result) # ── POST /apply-groups ─────────────────────────────────────────────── def _handle_apply_groups(self): length = int(self.headers.get("Content-Length", 0)) body = json.loads(self.rfile.read(length)) batch_id = body.get("batch_id") groups = body.get("groups", []) if not batch_id: self._json({"ok": False, "error": "batch_id required"}, 400) return from services.grouping import apply_groups result = apply_groups(batch_id, groups) self._json(result) # ── POST /api/cleanup ──────────────────────────────────────────────── def _handle_cleanup(self): """Полная очистка всех данных (кроме prompts). Вызывается при загрузке страницы.""" from db.connection import execute execute("DELETE FROM spec_current") execute("DELETE FROM spec_events") execute("DELETE FROM supplements") execute("DELETE FROM upload_chunks") execute("DELETE FROM documents") execute("DELETE FROM contracts") self._json({"ok": True, "message": "all data cleaned"}) # ── /app.js (static) ─────────────────────────────────────────────────── def _handle_app_js(self): try: fname = self.path.lstrip("/").replace("static/", "") with open(os.path.join(os.path.dirname(__file__), fname), "rb") as f: js = f.read() self.send_response(200) self.send_header("Content-Type", "application/javascript; charset=utf-8") self.send_header("Content-Length", str(len(js))) self._send_cors() self.end_headers() self.wfile.write(js) except FileNotFoundError: self.send_error(404) # ── /upload ─────────────────────────────────────────────────────────── def _handle_upload(self): try: content_type = self.headers.get("Content-Type", "") content_length = int(self.headers.get("Content-Length", 0)) result = handle_upload(self.rfile, content_type, content_length) self._json(result) except Exception as e: self._json({"ok": False, "error": str(e)}, 500) # ── /unzip-upload ───────────────────────────────────────────────────── def _handle_unzip(self): try: content_length = int(self.headers.get("Content-Length", 0)) result = handle_unzip(self.rfile, content_length) self._json(result) except Exception as e: self._json({"ok": False, "error": str(e)}, 500) # ── /convert-doc ────────────────────────────────────────────────────── def _handle_convert_doc(self): """DOC → DOCX через LibreOffice.""" import subprocess, tempfile 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_files = [x for x in os.listdir(tmpdir) if x.endswith(".docx")] if docx_files: with open(os.path.join(tmpdir, docx_files[0]), "rb") as f: result = f.read() self.send_response(200) self.send_header("Content-Type", "application/vnd.openxmlformats-officedocument.wordprocessingml.document") self.send_header("Content-Length", str(len(result))) self._send_cors() self.end_headers() self.wfile.write(result) else: self.send_error(500, "conversion produced no output") except Exception as e: self.send_error(500, str(e)) finally: os.unlink(doc_path) for x in os.listdir(tmpdir): os.unlink(os.path.join(tmpdir, x)) os.rmdir(tmpdir) # ── /llm-ops (debug) ────────────────────────────────────────────────── def _handle_llm_ops(self): from services.llm import call_llm length = int(self.headers.get("Content-Length", 0)) body = json.loads(self.rfile.read(length)) try: result, prompt_id = call_llm( body.get("current_spec", []), body.get("doc_text", ""), build_prompt, ) self._json(result) except Exception as e: self._json({"error": str(e)}, 502) # ── Helpers ─────────────────────────────────────────────────────────── def _json(self, data, status=200): self.send_response(status) self._send_cors() self.send_header("Content-Type", "application/json; charset=utf-8") self.end_headers() self.wfile.write(json.dumps(data, ensure_ascii=False).encode()) def _send_cors(self): self.send_header("Access-Control-Allow-Origin", "*") def log_message(self, format, *args): pass if __name__ == "__main__": import os TCPServer.allow_reuse_address = True port = int(os.environ.get("PORT", "8766")) server = ThreadingHTTPServer(("0.0.0.0", port), Handler) print(f"Contracts VM server on :{port}, db={DB_CONFIG['dbname']}") try: server.serve_forever() except KeyboardInterrupt: server.shutdown()