diff --git a/app/db/__init__.py b/app/db/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/db/connection.py b/app/db/connection.py new file mode 100644 index 0000000..6afefb4 --- /dev/null +++ b/app/db/connection.py @@ -0,0 +1,91 @@ +""" +Database connection — psycopg2 connection pool. + +Архитектурное решение (Opus): +- ThreadedConnectionPool(minconn=1, maxconn=10) — оптимально для ThreadingMixIn HTTP-сервера. + Каждый HTTP-запрос в отдельном потоке получает своё соединение из пула. +- RealDictCursor для query() — возвращает dict с lowercase ключами (консистентно с Python API). +- execute()/execute_returning() — для INSERT/UPDATE/DELETE с автокоммитом и откатом при ошибке. +- Пул создаётся лениво (при первом запросе) через get_pool(). +""" +import os +import psycopg2 +import psycopg2.pool +import psycopg2.extras + +# Глобальный пул соединений (singleton) +_pool = None + +# Параметры подключения из переменных окружения (systemd Environment) +DB_CONFIG = { + "host": os.getenv("DB_HOST", "127.0.0.1"), + "port": int(os.getenv("DB_PORT", "5432")), + "dbname": os.getenv("DB_NAME", "baza"), + "user": os.getenv("DB_USER", "super"), + "password": os.getenv("DB_PASS", ""), +} + + +def get_pool(): + """Возвращает глобальный пул соединений. Создаёт при первом вызове.""" + global _pool + if _pool is None: + _pool = psycopg2.pool.ThreadedConnectionPool( + minconn=1, maxconn=10, **DB_CONFIG + ) + return _pool + + +def query(sql, params=None): + """ + SELECT → list[dict] с lowercase ключами. + Используется всеми db/*.py модулями для чтения данных. + """ + pool = get_pool() + conn = pool.getconn() + try: + with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur: + cur.execute(sql, params) + rows = cur.fetchall() + return [{k.lower(): v for k, v in r.items()} for r in rows] + finally: + pool.putconn(conn) + + +def execute(sql, params=None): + """ + INSERT/UPDATE/DELETE → количество затронутых строк. + Автокоммит. При ошибке — rollback и проброс исключения. + """ + pool = get_pool() + conn = pool.getconn() + try: + with conn.cursor() as cur: + cur.execute(sql, params) + conn.commit() + return cur.rowcount + except Exception: + conn.rollback() + raise + finally: + pool.putconn(conn) + + +def execute_returning(sql, params=None): + """ + INSERT/UPDATE/DELETE с RETURNING → dict первой строки. + Используется для insert с автогенерацией UUID (gen_random_uuid()). + """ + pool = get_pool() + conn = pool.getconn() + try: + with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur: + cur.execute(sql, params) + conn.commit() + row = cur.fetchone() + return {k.lower(): v for k, v in row.items()} if row else None + except Exception: + conn.rollback() + raise + finally: + pool.putconn(conn) diff --git a/app/db/contracts.py b/app/db/contracts.py new file mode 100644 index 0000000..4d13216 --- /dev/null +++ b/app/db/contracts.py @@ -0,0 +1,25 @@ +"""Contracts CRUD.""" +from .connection import query, execute, execute_returning + + +def insert(number, client=""): + return execute_returning( + "INSERT INTO contracts (number, client) VALUES (%s, %s) RETURNING *", + (number, client), + ) + + +def get(contract_id): + rows = query("SELECT * FROM contracts WHERE id = %s", (contract_id,)) + return rows[0] if rows else None + + +def delete(contract_id): + return execute("DELETE FROM contracts WHERE id = %s", (contract_id,)) + + +def delete_orphaned(): + """Remove contracts with no supplements.""" + execute( + "DELETE FROM contracts WHERE id NOT IN (SELECT DISTINCT contract_id FROM supplements)" + ) diff --git a/app/db/documents.py b/app/db/documents.py new file mode 100644 index 0000000..03a3e85 --- /dev/null +++ b/app/db/documents.py @@ -0,0 +1,115 @@ +"""Documents CRUD.""" +from .connection import query, execute, execute_returning + + +def insert(filename, mime_type, original_bytes, status="uploaded", batch_id=None, zip_source=None, content_hash=None): + """Insert document, return row dict.""" + return execute_returning( + """INSERT INTO documents (filename, mime_type, original_bytes, status, batch_id, zip_source, content_hash) + VALUES (%s, %s, %s, %s, %s, %s, %s) RETURNING *""", + (filename, mime_type, original_bytes, status, batch_id, zip_source, content_hash), + ) + + +def get(doc_id): + rows = query("SELECT * FROM documents WHERE id = %s", (doc_id,)) + return rows[0] if rows else None + + +def get_by_hash(batch_id, content_hash): + """Find document by content hash within batch.""" + rows = query( + "SELECT id FROM documents WHERE batch_id = %s AND content_hash = %s LIMIT 1", + (batch_id, content_hash), + ) + return rows[0] if rows else None + + +def set_parsed(doc_id, elements_json): + """Update elements_json + status='parsed'.""" + import json + return execute( + "UPDATE documents SET elements_json = %s::jsonb, status = 'parsed' WHERE id = %s", + (json.dumps(elements_json, ensure_ascii=False), doc_id), + ) + + +def set_error(doc_id, error_message): + return execute( + "UPDATE documents SET status = 'error', error_message = %s WHERE id = %s", + (error_message, doc_id), + ) + + +def delete(doc_id): + return execute("DELETE FROM documents WHERE id = %s", (doc_id,)) + + +def set_classification(doc_id, doc_type, own_number, parent_number, doc_date, counterparty, classify_raw=None, classify_input=None): + """Store LLM classification results + raw response + input text.""" + return execute( + """UPDATE documents SET doc_type=%s, own_number=%s, parent_number=%s, + doc_date=%s, counterparty=%s, classify_status='classified', classify_raw=%s, classify_input=%s + WHERE id=%s""", + (doc_type, own_number, parent_number, doc_date, counterparty, classify_raw, classify_input, doc_id), + ) + + +def set_classify_failed(doc_id, error): + return execute( + "UPDATE documents SET classify_status='failed', error_message=%s WHERE id=%s", + (error, doc_id), + ) + + +def set_classify_garbage(doc_id, reason=""): + """Mark document as garbage (Stage 1-2 filter, no LLM call).""" + return execute( + "UPDATE documents SET doc_type='garbage', classify_status='garbage', error_message=%s WHERE id=%s", + (f"garbage: {reason}", doc_id), + ) + + +def list_pending(batch_id): + """Documents waiting for classification.""" + return query( + "SELECT id, filename, elements_json FROM documents WHERE batch_id=%s AND classify_status='pending'", + (batch_id,), + ) + + +def reset_classify_status(batch_id): + """Сбросить classify_status на 'pending' только для 'processing' (crash recovery). + Уже классифицированные ('classified', 'garbage', 'failed') НЕ трогаем.""" + return execute( + "UPDATE documents SET classify_status='pending', error_message=NULL WHERE batch_id=%s AND classify_status='processing'", + (batch_id,), + ) + + +def set_classify_processing(doc_id): + """Mark document as being processed (for crash recovery).""" + return execute( + "UPDATE documents SET classify_status='processing' WHERE id=%s", + (doc_id,), + ) + + +def list_by_batch(batch_id): + """All documents in a batch with classification fields.""" + return query( + """SELECT id, filename, status, doc_type, own_number, parent_number, + doc_date, counterparty, classify_status, error_message, classify_raw, classify_input, + zip_source + FROM documents WHERE batch_id=%s ORDER BY created_at""", + (batch_id,), + ) + + +def count_by_status(batch_id): + """Count documents by classify_status.""" + rows = query( + "SELECT classify_status, count(*) as cnt FROM documents WHERE batch_id=%s GROUP BY classify_status", + (batch_id,), + ) + return {r["classify_status"]: r["cnt"] for r in rows} diff --git a/app/db/prompts.py b/app/db/prompts.py new file mode 100644 index 0000000..cfbb3ef --- /dev/null +++ b/app/db/prompts.py @@ -0,0 +1,151 @@ +"""Prompts CRUD.""" +from .connection import query, execute, execute_returning + + +def _serialize(row): + """Convert datetime fields to strings for JSON serialization (non-mutating).""" + if row and row.get("created_at"): + row = dict(row) + row["created_at"] = str(row["created_at"]) + return row + + +def get_active(role): + rows = query( + "SELECT * FROM prompts WHERE role = %s AND is_active = true ORDER BY created_at DESC LIMIT 1", + (role,), + ) + return _serialize(rows[0]) if rows else None + + +def get(prompt_id): + rows = query("SELECT * FROM prompts WHERE id = %s", (prompt_id,)) + return rows[0] if rows else None + + +def seed_defaults(): + """Auto-seed default prompts if table is empty.""" + # Check each role separately — don't skip if one is missing + existing_roles = set(r["role"] for r in query("SELECT DISTINCT role FROM prompts")) + + if "extract" not in existing_roles: + extract_body = ( + "Ты — анализатор договоров облачного провайдера.\n\n" + "Ниже текст спецификации услуг из ПЕРВОГО документа (базовый договор).\n" + "Извлеки ВСЕ строки спецификации в JSON-массив.\n\n" + "Верни СТРОГО JSON без пояснений. Не используй markdown-блоки.\n\n" + "ФОРМАТ:\n" + '{\n "mode": "partial",\n "ops": [\n' + ' {"action": "ADD", "new_row": {"name": "полное наименование", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}, "comment": ""}\n' + " ]\n}\n\n" + "ПРАВИЛА:\n" + "1. Извлеки КАЖДУЮ строку таблицы спецификации как отдельную ADD-операцию.\n" + '2. Пропускай итоговые строки ("Итого...") и строки с подписями.\n' + "3. Если ячейка пустая — ставь null (не пиши 0).\n" + "4. price, qty, sum — ЧИСЛА, не строки.\n\n" + "ТЕКСТ ДОКУМЕНТА:\n---\n{doc_text}\n---" + ) + execute( + "INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)", + ("extract", "default-v1", extract_body, True, "Авто-создан из llm_prompt.py _build_initial"), + ) + + if "diff" not in existing_roles: + diff_body = ( + "Ты — анализатор допсоглашений к договорам облачного провайдера.\n\n" + "У тебя есть текущая спецификация услуг и текст нового допсоглашения (ДС).\n" + "Твоя задача — определить, какие изменения вносит ДС в текущую спецификацию.\n\n" + "Верни СТРОГО JSON без пояснений. Не используй markdown-блоки.\n\n" + "ФОРМАТ:\n" + '{\n "mode": "partial" | "full_replace",\n "ops": [\n' + ' {"action": "ADD", "new_row": {"name": "...", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}, "comment": "..."},\n' + ' {"action": "UPDATE", "target_id": "rN", "new_values": {"price": число}, "comment": "..."},\n' + ' {"action": "DELETE", "target_id": "rN", "comment": "..."},\n' + ' {"action": "UNRESOLVED", "new_values": {"name": "...", "price": число, ...}, "reason": "почему не смог сопоставить"}\n' + " ]\n}\n\n" + "ПРАВИЛА:\n" + '1. mode = "full_replace" — если в тексте есть фразы: «в следующей редакции», «заменить приложение», «излагается в следующей редакции». При full_replace — опиши ВСЕ новые строки как ADD.\n' + '2. mode = "partial" — если ДС меняет только отдельные строки.\n' + "3. Для UPDATE/DELETE — укажи target_id (r1, r2...) из списка текущей спецификации. НЕ придумывай новые id.\n" + "4. new_values в UPDATE — только ИЗМЕНЁННЫЕ поля (не все).\n" + "5. Если не можешь однозначно сопоставить строку — action: UNRESOLVED с reason.\n" + "6. Пропускай итоговые строки и подписи.\n" + "7. price, qty, sum — ЧИСЛА (не строки).\n\n" + "ТЕКУЩАЯ СПЕЦИФИКАЦИЯ:\n{spec_current}\n\n" + "ТЕКСТ ДОПСОГЛАШЕНИЯ:\n---\n{doc_text}\n---" + ) + execute( + "INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)", + ("diff", "default-v1", diff_body, True, "Авто-создан из llm_prompt.py _build_diff"), + ) + _ensure_classify_prompt() + + +def list_by_role(role): + """List all versions for a role, newest first.""" + rows = query( + "SELECT id, role, name, is_active, created_by, notes, created_at FROM prompts WHERE role=%s ORDER BY created_at DESC", + (role,), + ) + for r in rows: + if r.get("created_at"): + r["created_at"] = str(r["created_at"]) + return rows + + +def save_new_version(role, name, body, notes="", is_active=True): + """Save new prompt version. Deactivates all others for this role, inserts new one.""" + if is_active: + execute("UPDATE prompts SET is_active=false WHERE role=%s", (role,)) + return execute_returning( + """INSERT INTO prompts (role, name, body, is_active, notes) + VALUES (%s, %s, %s, %s, %s) RETURNING *""", + (role, name, body, is_active, notes), + ) + + +def activate(prompt_id): + """Activate a prompt version (deactivates others for same role).""" + row = query("SELECT role FROM prompts WHERE id=%s", (prompt_id,)) + if not row: + return False + role = row[0]["role"] + execute("UPDATE prompts SET is_active=false WHERE role=%s", (role,)) + execute("UPDATE prompts SET is_active=true WHERE id=%s", (prompt_id,)) + return True + + +def delete_prompt(prompt_id): + """Delete a prompt version (cannot delete active one).""" + row = query("SELECT is_active FROM prompts WHERE id=%s", (prompt_id,)) + if not row: + return False + if row[0]["is_active"]: + return False # cannot delete active + execute("DELETE FROM prompts WHERE id=%s", (prompt_id,)) + return True + + +def _ensure_classify_prompt(): + """Ensure classify prompt exists (idempotent).""" + existing = query("SELECT id FROM prompts WHERE role = 'classify' AND is_active = true LIMIT 1") + if existing: + return + body = ( + "Ты — система классификации договорных документов. " + "Проанализируй текст и верни СТРОГО ВАЛИДНЫЙ JSON ОДНОЙ СТРОКОЙ " + "(без переносов строк, без markdown, без лишних пробелов в начале/конце).\n\n" + "Поля:\n" + '- doc_type: "contract" (договор) / "supplement" (допсоглашение) / "specification" (спецификация/приложение) / "other"\n' + "- own_number: номер ЭТОГО документа, строка без лишних пробелов (или null)\n" + '- parent_number: номер родительского договора из фразы «к Договору №...» (или null)\n' + '- doc_date: дата в YYYY-MM-DD. «27 февраля 2026» → 2026-02-27 (или null)\n' + "- counterparty: полное название контрагента (Заказчик/Арендатор), без сокращений (или null)\n\n" + "Пример вывода:\n" + '{"doc_type":"supplement","own_number":"1","parent_number":"01300_2","doc_date":"2026-02-27","counterparty":"АО XXX003"}\n\n' + "ДОКУМЕНТ:\n---\n{header_text}\n---" + ) + execute( + "INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)", + ("classify", "default-v1", body, True, "Авто-создан: classify prompt v2 — LLM сам предложил"), + ) diff --git a/app/db/spec_current.py b/app/db/spec_current.py new file mode 100644 index 0000000..63e9de3 --- /dev/null +++ b/app/db/spec_current.py @@ -0,0 +1,19 @@ +"""spec_current — текущее состояние спецификации.""" +from .connection import query + + +def list_by_contract(contract_id): + """Return list of dicts with name_hash, name, price, qty, sum, date_start.""" + return query( + """SELECT name_hash, name, price, qty, sum, date_start + FROM spec_current WHERE contract_id = %s ORDER BY name""", + (contract_id,), + ) + + +def get_elements_json(document_id): + """Get elements_json for a document.""" + rows = query( + "SELECT elements_json FROM documents WHERE id = %s", (document_id,) + ) + return rows[0]["elements_json"] if rows and rows[0]["elements_json"] else None diff --git a/app/db/spec_events.py b/app/db/spec_events.py new file mode 100644 index 0000000..2eb10eb --- /dev/null +++ b/app/db/spec_events.py @@ -0,0 +1,202 @@ +"""spec_events — event sourcing: apply ops, reset contract.""" +import json, uuid +from .connection import query, execute, get_pool + + +def reset(contract_id): + """Clear spec_events + spec_current for contract.""" + execute("DELETE FROM spec_current WHERE contract_id = %s", (contract_id,)) + execute("DELETE FROM spec_events WHERE contract_id = %s", (contract_id,)) + + +def get_next_seq(contract_id): + """Get next sequence number with row lock to prevent race conditions.""" + pool = get_pool() + conn = pool.getconn() + try: + conn.autocommit = False + with conn.cursor() as cur: + cur.execute( + "SELECT seq FROM spec_events WHERE contract_id = %s ORDER BY seq DESC LIMIT 1 FOR UPDATE", + (contract_id,), + ) + row = cur.fetchone() + seq = (row[0] + 1) if row else 1 + conn.commit() + return seq + except Exception: + conn.rollback() + raise + finally: + conn.autocommit = True + pool.putconn(conn) + + +def apply_ops(contract_id, supplement_id, document_id, ops, prompt_id, raw_llm_response): + """Apply ADD/UPDATE/DELETE ops. Returns summary dict.""" + added = 0 + updated = 0 + deleted = 0 + seq = get_next_seq(contract_id) + + for op in ops: + action = op.get("action", "").upper() + + if action == "ADD": + nr = op.get("new_row", {}) + name = nr.get("name", "") + if not name: + # ADD without name → UNRESOLVED + _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "ADD with empty name") + seq += 1 + continue + name_hash = _hash(name, nr.get("date_start")) + execute( + """INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash, + new_values, comment, status, prompt_version, source_document_id, raw_llm_response) + VALUES (%s, %s, %s, 'ADD', %s, %s, %s, 'applied', %s, %s, %s)""", + ( + contract_id, supplement_id, seq, name_hash, + json.dumps(nr, ensure_ascii=False), + op.get("comment", ""), prompt_id, document_id, + json.dumps(raw_llm_response, ensure_ascii=False), + ), + ) + _upsert_spec_current(contract_id, name_hash, nr) + seq += 1 + added += 1 + + elif action == "UPDATE": + nv = op.get("new_values", {}) + th = op.get("target_hash", "") + if not th: + _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "UPDATE with empty target_hash") + seq += 1 + continue + execute( + """INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash, + new_values, comment, status, prompt_version, source_document_id, raw_llm_response) + VALUES (%s, %s, %s, 'UPDATE', %s, %s, %s, 'applied', %s, %s, %s)""", + ( + contract_id, supplement_id, seq, th, + json.dumps(nv, ensure_ascii=False), + op.get("comment", ""), prompt_id, document_id, + json.dumps(raw_llm_response, ensure_ascii=False), + ), + ) + _update_spec_current(contract_id, th, nv) + seq += 1 + updated += 1 + + elif action == "DELETE": + th = op.get("target_hash", "") + if not th: + _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "DELETE with empty target_hash") + seq += 1 + continue + execute( + """INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash, + new_values, comment, status, prompt_version, source_document_id, raw_llm_response) + VALUES (%s, %s, %s, 'DELETE', %s, %s, %s, 'applied', %s, %s, %s)""", + ( + contract_id, supplement_id, seq, th, + json.dumps({}), op.get("comment", ""), + prompt_id, document_id, + json.dumps(raw_llm_response, ensure_ascii=False), + ), + ) + execute("DELETE FROM spec_current WHERE contract_id = %s AND name_hash = %s", (contract_id, th)) + seq += 1 + deleted += 1 + + elif action == "UNRESOLVED": + # Log but don't apply + execute( + """INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash, + new_values, comment, status, prompt_version, source_document_id, raw_llm_response) + VALUES (%s, %s, %s, 'UNRESOLVED', %s, %s, %s, 'unresolved', %s, %s, %s)""", + ( + contract_id, supplement_id, seq, + op.get("target_hash", ""), + json.dumps(op.get("new_values", {}), ensure_ascii=False), + op.get("reason", op.get("comment", "")), + prompt_id, document_id, + json.dumps(raw_llm_response, ensure_ascii=False), + ), + ) + seq += 1 + + else: + # Unknown action — log as UNRESOLVED + _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, + f"unknown action: {action}") + + return {"added": added, "updated": updated, "deleted": deleted} + + +def _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, reason): + """Log an op as UNRESOLVED instead of silently ignoring it.""" + execute( + """INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash, + new_values, comment, status, prompt_version, source_document_id, raw_llm_response) + VALUES (%s, %s, %s, 'UNRESOLVED', %s, %s, %s, 'unresolved', %s, %s, %s)""", + ( + contract_id, supplement_id, seq, + op.get("target_hash", ""), + json.dumps(op.get("new_values", op.get("new_row", {})) or {}, ensure_ascii=False), + reason, + prompt_id, document_id, + json.dumps(raw_llm_response, ensure_ascii=False), + ), + ) + + +def _hash(name, date_start=None): + """Нормализованный хеш услуги. Включает нормализованный date_start чтобы различать периоды.""" + import hashlib + key = name.strip().lower() + if date_start: + from compare.metrics import normalize_date + nd = normalize_date(str(date_start)) + if nd: + key += "|" + nd + return hashlib.sha256(key.encode()).hexdigest()[:16] + + +def _upsert_spec_current(contract_id, name_hash, row): + """INSERT or UPDATE spec_current.""" + existing = query( + "SELECT id FROM spec_current WHERE contract_id = %s AND name_hash = %s", + (contract_id, name_hash), + ) + if existing: + execute( + """UPDATE spec_current SET name=%s, price=%s, qty=%s, sum=%s, date_start=%s, updated_at=now() + WHERE contract_id=%s AND name_hash=%s""", + (row.get("name"), row.get("price"), row.get("qty"), row.get("sum"), + row.get("date_start"), contract_id, name_hash), + ) + else: + execute( + """INSERT INTO spec_current (contract_id, name_hash, name, price, qty, sum, date_start) + VALUES (%s, %s, %s, %s, %s, %s, %s)""", + (contract_id, name_hash, row.get("name"), row.get("price"), + row.get("qty"), row.get("sum"), row.get("date_start")), + ) + + +def _update_spec_current(contract_id, name_hash, new_values): + """Update specific fields in spec_current.""" + sets = [] + params = [] + for field in ("name", "price", "qty", "sum", "date_start"): + if field in new_values: + sets.append(f"{field} = %s") + params.append(new_values[field]) + if sets: + sets.append("updated_at = now()") + params.extend([contract_id, name_hash]) + execute( + f"UPDATE spec_current SET {', '.join(sets)} WHERE contract_id = %s AND name_hash = %s", + params, + ) diff --git a/app/db/supplements.py b/app/db/supplements.py new file mode 100644 index 0000000..41fa40e --- /dev/null +++ b/app/db/supplements.py @@ -0,0 +1,59 @@ +"""Supplements CRUD.""" +from .connection import query, execute, execute_returning + + +def insert(contract_id, document_id, supp_type="additional"): + return execute_returning( + """INSERT INTO supplements (contract_id, document_id, type) + VALUES (%s, %s, %s) RETURNING *""", + (contract_id, document_id, supp_type), + ) + + +def list_by_contract(contract_id): + """Supplements with parsed documents, ordered by created_at.""" + return query( + """SELECT s.id, s.type, s.document_id, d.filename + FROM supplements s + JOIN documents d ON s.document_id = d.id + WHERE s.contract_id = %s AND d.elements_json IS NOT NULL + ORDER BY s.created_at""", + (contract_id,), + ) + + +def get(supp_id): + rows = query("SELECT * FROM supplements WHERE id = %s", (supp_id,)) + return rows[0] if rows else None + + +def delete_by_document(contract_id, filename): + """Delete ALL supplements+documents by contract+filename (cascade: spec_events first). + Handles duplicates from previously failed uploads.""" + rows = query( + """SELECT s.id as sid, s.document_id FROM supplements s + JOIN documents d ON d.id = s.document_id + WHERE s.contract_id = %s AND d.filename = %s""", + (contract_id, filename), + ) + if not rows: + return False + + for r in rows: + # 1. Delete spec_current rows referencing this supplement's events + execute( + """DELETE FROM spec_current WHERE contract_id = %s + AND last_event_id IN (SELECT id FROM spec_events WHERE supplement_id = %s)""", + (contract_id, r["sid"]), + ) + # 2. Delete spec_events referencing this supplement + execute("DELETE FROM spec_events WHERE supplement_id = %s", (r["sid"],)) + # 3. Delete supplement + execute("DELETE FROM supplements WHERE id = %s", (r["sid"],)) + # 4. Delete document + execute("DELETE FROM documents WHERE id = %s", (r["document_id"],)) + return True + + +def delete(supp_id): + return execute("DELETE FROM supplements WHERE id = %s", (supp_id,)) diff --git a/app/llm_prompt.py b/app/llm_prompt.py new file mode 100644 index 0000000..7b64351 --- /dev/null +++ b/app/llm_prompt.py @@ -0,0 +1,281 @@ +"""llm_prompt.py — Формирование промпта для LLM-анализа ДС. +Читает активный промпт из БД напрямую (db.prompts). При ошибке — fallback на хардкод.""" + +import os +from db import prompts as db_prompts + +# ── Fallback-промпты (если БД недоступна) ────────────────────── + +FALLBACK_EXTRACT = """Ты — анализатор договоров облачного провайдера и ЦОД (дата-центра). +Ты разбираешь спецификации услуг colocation, аренды стоек, питания, каналов связи и облачных ресурсов. + +ЗАДАЧА: ниже текст спецификации услуг из ПЕРВОГО документа (базовый договор). +Извлеки ВСЕ строки спецификации, каждую как отдельную ADD-операцию. + +Верни СТРОГО JSON без пояснений. Не используй markdown-блоки, не добавляй текст до или после JSON. + +ФОРМАТ: +{{ + "mode": "partial", + "ops": [ + {{"action": "ADD", "new_row": {{"name": "полное наименование", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}}, "comment": ""}} + ] +}} + +ДОМЕННЫЙ ГЛОССАРИЙ (для корректного разбора): +- Единицы измерения: + \u2022 кВт — мощность электропитания (номинальная/гарантированная). + \u2022 юнит, U — высота места в стойке (1U, 2U, 10U). + \u2022 шт. — счётные позиции (IP-адреса, кросс-соединения, порты). + \u2022 Мбит/с, Гбит/с — пропускная способность канала связи. + \u2022 ГБ, ТБ — объём диска/хранилища; vCPU — виртуальные ядра; RAM ГБ — память. +- Типичные услуги: + \u2022 «Стойко-место» / «Аренда стойко-места» / «Colocation» — размещение оборудования в стойке ЦОД. + \u2022 «Электропитание» / «Питание» — выделенная мощность в кВт. + \u2022 «IP-адрес» (IPv4/IPv6) — считается в шт. + \u2022 «Канал связи» / «Порт» / «Интернет» — пропускная способность. + \u2022 «Кросс-соединение» (cross-connect) — физическая коммутация, шт. + \u2022 «Облачные ресурсы» — vCPU, RAM, диск. +- Мощность и габариты часто входят В СОСТАВ названия услуги: + «Аренда стойко-места, в составе: Номинальная мощность – 10 кВт». Сохраняй такое название ЦЕЛИКОМ. + +ПРАВИЛА: +1. Извлеки КАЖДУЮ строку таблицы спецификации как отдельную ADD-операцию. +2. name — полное наименование услуги дословно, со всеми уточнениями (мощность, объём, кол-во в составе). Не сокращай. +3. Пропускай итоговые строки («Итого», «Всего», «НДС», «К оплате») и строки с подписями/реквизитами. +4. Если ячейка пустая или значение не указано — ставь null (НЕ пиши 0). +5. price, qty, sum — ЧИСЛА (без пробелов, без «руб.», точка как десятичный разделитель). «50 000,00 руб.» \u2192 50000. +6. date_start — дата начала оказания услуги в формате YYYY-MM-DD. Если в документе нет — null. +7. Не вычисляй и не «исправляй» суммы. Бери значения как в документе. + +ПРИМЕР: +Текст: «1. Аренда стойко-места, в составе: Номинальная мощность – 10 кВт — 1 шт. — 50 000,00 руб. — 50 000,00 руб. Дата начала: 01.01.2025 +2. IP-адрес IPv4 — 8 шт. — 300,00 руб. — 2 400,00 руб.» +Ответ: +{{ + "mode": "partial", + "ops": [ + {{"action": "ADD", "new_row": {{"name": "Аренда стойко-места, в составе: Номинальная мощность – 10 кВт", "price": 50000, "qty": 1, "sum": 50000, "date_start": "2025-01-01"}}, "comment": ""}}, + {{"action": "ADD", "new_row": {{"name": "IP-адрес IPv4", "price": 300, "qty": 8, "sum": 2400, "date_start": null}}, "comment": ""}} + ] +}} + +ТЕКСТ ДОКУМЕНТА: +--- +{doc_text} +---""" + +FALLBACK_DIFF = """Ты — анализатор допсоглашений (ДС) к договорам облачного провайдера и ЦОД. +У тебя есть ТЕКУЩАЯ спецификация услуг (с готовыми id строк) и текст нового ДС. +Задача — определить, какие изменения ДС вносит в текущую спецификацию. + +Верни СТРОГО JSON без пояснений. Не используй markdown-блоки, не добавляй текст до или после JSON. + +ФОРМАТ: +{{ + "mode": "partial" | "full_replace", + "ops": [ + {{"action": "ADD", "new_row": {{"name": "...", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}}, "comment": "..."}}, + {{"action": "UPDATE", "target_id": "rN", "new_values": {{"price": число}}, "comment": "..."}}, + {{"action": "DELETE", "target_id": "rN", "comment": "..."}}, + {{"action": "UNRESOLVED", "new_values": {{"name": "...", "price": число}}, "reason": "почему не смог сопоставить"}} + ] +}} + +ДОМЕННЫЙ ГЛОССАРИЙ: +- Единицы: кВт (мощность), юнит/U (высота в стойке), шт. (IP, кросс-соединения, порты), Мбит/с\u00b7Гбит/с (канал), ГБ\u00b7ТБ\u00b7vCPU (облако). +- Услуги: стойко-место / colocation (размещение в стойке); электропитание / питание (мощность кВт); IP-адрес IPv4/IPv6 (шт.); канал связи / порт (пропускная способность); кросс-соединение (шт.); облачные ресурсы (vCPU/RAM/диск). +- Мощность/объём часто ВНУТРИ названия услуги: «Аренда стойко-места, в составе: Номинальная мощность – 10 кВт». + Если ДС меняет мощность (10 кВт \u2192 15 кВт) — это UPDATE той же строки, причём меняется и name, и, как правило, price/sum. + +СОПОСТАВЛЕНИЕ СТРОК: +1. Для UPDATE/DELETE укажи target_id (r1, r2\u2026) ИЗ списка текущей спецификации ниже. НЕ придумывай новые id. +2. Сопоставляй по СМЫСЛУ услуги, а не по точному совпадению символов. «Аренда стойко-места» = «Размещение оборудования в стойке» = одна услуга. Различие тире/пробелов/кавычек игнорируй. +3. new_values в UPDATE — ТОЛЬКО изменённые поля (не дублируй неизменные). +4. Если ДС увеличивает количество той же услуги (было 8 IP, стало 12) — это UPDATE qty (и sum), а не новая ADD. + +РЕЖИМ mode: +5. mode = "full_replace" — если ДС полностью переиздаёт приложение/спецификацию. Признаки: «Приложение \u2026 излагается в следующей редакции», «изложить в новой редакции», «заменить приложение \u2116\u2026». При full_replace опиши ВСЕ строки новой редакции как ADD (UPDATE/DELETE не используй). +6. mode = "partial" — если ДС точечно меняет отдельные позиции (изменить цену, добавить/удалить услугу, изменить мощность/кол-во). +7. Если в тексте есть и фраза о новой редакции, и точечные правки — приоритет за «новой редакцией»: full_replace. + +EDGE-CASES: +8. UNRESOLVED — если ДС упоминает изменение услуги, которой НЕТ в текущей спецификации, ИЛИ название настолько отличается, что нельзя уверенно сопоставить с конкретным id. ВАЖНО: если СОМНЕВАЕШЬСЯ в сопоставлении — делай UNRESOLVED, а НЕ ADD. Лучше unresolved, чем ложный дубликат. В reason укажи причину. +9. Частичные данные: если в ДС нет цены/кол-ва/даты — ставь null для этих полей, не выдумывай. +10. Пропускай итоговые строки («Итого», «НДС», «К оплате») и подписи/реквизиты. +11. price, qty, sum — ЧИСЛА (без «руб.», без пробелов; «55 000,00» \u2192 55000). Не пересчитывай суммы сам — бери из ДС. +12. date_start — YYYY-MM-DD; используй дату вступления изменения в силу из ДС, если она указана. + +ПРИМЕР 1 (partial, UPDATE цены): +Текущая спецификация: +[id: r1] Аренда стойко-места, в составе: Номинальная мощность – 10 кВт | цена=50000 | объём=1 | сумма=50000 | начало=2025-01-01 +[id: r2] IP-адрес IPv4 | цена=300 | объём=8 | сумма=2400 | начало=2025-01-01 +Текст ДС: «С 01.03.2025 стоимость аренды стойко-места устанавливается в размере 55 000,00 руб. в месяц.» +Ответ: +{{ + "mode": "partial", + "ops": [ + {{"action": "UPDATE", "target_id": "r1", "new_values": {{"price": 55000, "sum": 55000, "date_start": "2025-03-01"}}, "comment": "Изменение стоимости аренды стойко-места"}} + ] +}} + +ПРИМЕР 2 (partial: ADD новая услуга + UPDATE количества + UPDATE мощности): +Текущая спецификация: +[id: r1] Аренда стойко-места, в составе: Номинальная мощность – 10 кВт | цена=50000 | объём=1 | сумма=50000 | начало=2025-01-01 +[id: r2] IP-адрес IPv4 | цена=300 | объём=8 | сумма=2400 | начало=2025-01-01 +Текст ДС: «С 01.04.2025: 1) увеличить номинальную мощность стойко-места до 15 кВт, стоимость — 70 000,00 руб.; +2) предоставить дополнительно 4 IP-адреса IPv4 (итого 12 шт., сумма 3 600,00 руб.); +3) предоставить услугу "Кросс-соединение" — 2 шт. по 1 500,00 руб., сумма 3 000,00 руб.» +Ответ: +{{ + "mode": "partial", + "ops": [ + {{"action": "UPDATE", "target_id": "r1", "new_values": {{"name": "Аренда стойко-места, в составе: Номинальная мощность – 15 кВт", "price": 70000, "sum": 70000, "date_start": "2025-04-01"}}, "comment": "Увеличение мощности 10\u219215 кВт"}}, + {{"action": "UPDATE", "target_id": "r2", "new_values": {{"qty": 12, "sum": 3600, "date_start": "2025-04-01"}}, "comment": "Увеличение количества IP-адресов 8\u219212"}}, + {{"action": "ADD", "new_row": {{"name": "Кросс-соединение", "price": 1500, "qty": 2, "sum": 3000, "date_start": "2025-04-01"}}, "comment": "Новая услуга"}} + ] +}} + +ПРИМЕР 3 (full_replace): +Текущая спецификация: +[id: r1] Аренда стойко-места, в составе: Номинальная мощность – 10 кВт | цена=50000 | объём=1 | сумма=50000 | начало=2025-01-01 +[id: r2] IP-адрес IPv4 | цена=300 | объём=8 | сумма=2400 | начало=2025-01-01 +Текст ДС: «Приложение №1 (Спецификация услуг) излагается в следующей редакции: +1. Аренда стойко-места, номинальная мощность 15 кВт — 1 шт. — 70 000,00 руб. +2. IP-адрес IPv4 — 12 шт. — 300,00 руб. — 3 600,00 руб. +3. Канал связи 1 Гбит/с — 1 шт. — 20 000,00 руб. Дата: 01.05.2025» +Ответ: +{{ + "mode": "full_replace", + "ops": [ + {{"action": "ADD", "new_row": {{"name": "Аренда стойко-места, номинальная мощность 15 кВт", "price": 70000, "qty": 1, "sum": 70000, "date_start": "2025-05-01"}}, "comment": "Новая редакция приложения"}}, + {{"action": "ADD", "new_row": {{"name": "IP-адрес IPv4", "price": 300, "qty": 12, "sum": 3600, "date_start": "2025-05-01"}}, "comment": "Новая редакция приложения"}}, + {{"action": "ADD", "new_row": {{"name": "Канал связи 1 Гбит/с", "price": 20000, "qty": 1, "sum": 20000, "date_start": "2025-05-01"}}, "comment": "Новая редакция приложения"}} + ] +}} + +ПРИМЕР 4 (UNRESOLVED): +Текущая спецификация: +[id: r1] Аренда стойко-места, в составе: Номинальная мощность – 10 кВт | цена=50000 | объём=1 | сумма=50000 | начало=2025-01-01 +Текст ДС: «Снизить стоимость услуги резервного копирования до 4 000,00 руб.» +Ответ: +{{ + "mode": "partial", + "ops": [ + {{"action": "UNRESOLVED", "new_values": {{"name": "Резервное копирование", "price": 4000}}, "reason": "В текущей спецификации нет услуги резервного копирования — не с чем сопоставить"}} + ] +}} + +ТЕКУЩАЯ СПЕЦИФИКАЦИЯ: +{spec_current} + +ТЕКСТ ДОПСОГЛАШЕНИЯ: +--- +{doc_text} +---""" + + +def _fetch_prompt(role: str) -> dict | None: + """Получить активный промпт из БД напрямую (а не через Lucee HTTP).""" + try: + row = db_prompts.get_active(role) + if row and row.get("body"): + return {"id": row.get("id", ""), "body": row["body"]} + except Exception: + pass + return None + + +def _build_spec_text(current_spec: list) -> str: + """Перечисление строк спецификации для подстановки в {spec_current}.""" + lines = [] + for i, r in enumerate(current_spec): + lines.append( + f"[id: r{i+1}] {r.get('name', '?')} | " + f"цена={r.get('price', '')} | объём={r.get('qty', '')} | " + f"сумма={r.get('sum', '')} | начало={r.get('date_start', '')}" + ) + return "\n".join(lines) + + +def build_prompt(current_spec: list, doc_text: str) -> tuple: + """ + Формирует промпт для LLM. + 1. Пробует получить активный промпт из БД (Lucee API). + 2. При неудаче — fallback на хардкод. + Возвращает (текст_промпта, prompt_id). + prompt_id — UUID версии промпта из БД, или "" если fallback. + """ + is_first = len(current_spec) == 0 + role = "extract" if is_first else "diff" + + db = _fetch_prompt(role) + if db: + template = db["body"] + prompt_id = db.get("id", "") + else: + template = FALLBACK_EXTRACT if is_first else FALLBACK_DIFF + prompt_id = "" + + # Подстановка плейсхолдеров + result = template.replace("{doc_text}", doc_text) + result = result.replace("{spec_current}", _build_spec_text(current_spec)) + + return result, prompt_id + + +def build_classify_prompt(header_text): + """Build classify prompt. Returns (prompt, prompt_id).""" + try: + from db import prompts as db_prompts + prompt = db_prompts.get_active("classify") + if prompt: + body = prompt["body"].replace("{header_text}", header_text) + return body, prompt.get("id", "") + except Exception: + pass # DB unavailable — use fallback + body = """Ты — классификатор договорных документов облачного провайдера НУБЕС. + +Ниже фрагмент текста документа. Определи: + +1. doc_type: + - "contract" — договор (заголовок «Договор», «Соглашение», преамбула с условиями) + - "supplement" — допсоглашение (ссылается на родительский договор, меняет условия) + - "specification" — спецификация / приложение с таблицей услуг (стойко-места, IP, каналы, питание) + - "other" — НЕ договорной документ: акт сверки, счёт, счёт-фактура, УПД, акт оказанных услуг, платёжное поручение, доверенность, письмо + +2. own_number — номер ЭТОГО документа (например «XXX001-03700», «МЭС-123/2024», «1» для допника). + Если номер не указан — null. + +3. parent_number — номер родительского договора (для supplement и specification). + Для doc_type="contract": ВСЕГДА null. + +4. doc_date — дата документа в формате YYYY-MM-DD. Если дата прописью — переведи в цифры. + Если нет даты — null. + +5. counterparty — название КОНТРАГЕНТА (Заказчика). + ВАЖНО: НУБЕС — всегда Исполнитель. НЕ возвращай НУБЕС как counterparty. + НУБЕС известен как: «НУБЕС», «ООО НУБЕС», «ООО "НУБЕС"», «Nubes». + counterparty — ВСЕГДА другая сторона (Заказчик/Покупатель/Абонент). + Если документ не содержит контрагента — null. + +Верни СТРОГО JSON без пояснений: +{"doc_type":"...","own_number":"...","parent_number":"...","doc_date":"...","counterparty":"..."} + +ПРИМЕР 1 (договор): +Текст: «Договор № XXX001-03700 от 15.03.2025. ООО "НУБЕС" (Исполнитель) и ЗАО "ТехноПлюс" (Заказчик)...» +Ответ: {"doc_type":"contract","own_number":"XXX001-03700","parent_number":null,"doc_date":"2025-03-15","counterparty":"ЗАО \"ТехноПлюс\""} + +ПРИМЕР 2 (допсоглашение): +Текст: «Допсоглашение №1 к Договору № XXX003-01300 от 05.06.2024...» +Ответ: {"doc_type":"supplement","own_number":"1","parent_number":"XXX003-01300","doc_date":"2024-06-05","counterparty":"АО XXX003"} + +ПРИМЕР 3 (мусор): +Текст: «Акт сверки взаимных расчётов за 1 квартал 2025 г. Стороны: НУБЕС и ООО Ромашка. Сальдо 150 000 руб.» +Ответ: {"doc_type":"other","own_number":null,"parent_number":null,"doc_date":"2025-03-31","counterparty":"ООО Ромашка"} + +ДОКУМЕНТ: +--- +{header_text} +---""".replace("{header_text}", header_text) + return body, ""