import os import json import asyncio import tempfile from pathlib import Path import mimetypes import urllib.parse import re import httpx import aio_pika import pandas as pd import openpyxl from openpyxl.styles import Alignment, Font, PatternFill from supabase import create_client, Client # Настройки SUPABASE_URL = os.getenv("SUPABASE_URL", "") SUPABASE_KEY = os.getenv("SUPABASE_KEY", "") CLOUDAMQP_URL = os.getenv("CLOUDAMQP_URL", "") OPENROUTER_API_KEY = os.getenv("OPENROUTER_API_KEY", "") UNSTRUCTURED_API_KEY = os.getenv("UNSTRUCTURED_API_KEY", "") UNSTRUCTURED_API_URL = os.getenv( "UNSTRUCTURED_API_URL", "https://api.unstructuredapp.io/general/v0/general" ).strip() if UNSTRUCTURED_API_URL.startswith("https://api.unstructured.io"): UNSTRUCTURED_API_URL = UNSTRUCTURED_API_URL.replace( "https://api.unstructured.io", "https://api.unstructuredapp.io", 1 ) supabase: Client = create_client(SUPABASE_URL, SUPABASE_KEY) def _parse_quantity(value) -> float: if value is None: return 0.0 if isinstance(value, (int, float)): return float(value) raw = str(value).strip().replace(",", ".") if not raw: return 0.0 match = re.search(r"-?\d+(?:\.\d+)?", raw) if not match: return 0.0 try: return float(match.group(0)) except ValueError: return 0.0 def _parse_llm_json(raw_text: str) -> dict: try: data = json.loads(raw_text) if isinstance(data, dict): return data except json.JSONDecodeError: pass match = re.search(r"\{[\s\S]*\}", raw_text) if match: try: data = json.loads(match.group(0)) if isinstance(data, dict): return data except json.JSONDecodeError: pass raise RuntimeError("LLM вернул невалидный JSON") def _normalize_catalog_key(value: str) -> str: # убирает регистр/кавычки/тире и лишние пробелы, чтобы LLM-варианты имени находили совпадения в справочнике text = (value or "").strip().lower() text = text.replace("\u2013", "-").replace("\u2014", "-") text = re.sub(r"[\"'\u00ab\u00bb]", "", text) text = re.sub(r"\s+", " ", text) return text def _normalize_unrecognized_entry(entry, source: str) -> dict: if isinstance(entry, dict): return { "source": entry.get("source", source), "description": entry.get("description") or entry.get("name") or str(entry), "quantity": entry.get("quantity", ""), "raw_context": entry.get("raw_context", ""), } return { "source": source, "description": str(entry), "quantity": "", "raw_context": "", } TENDER_DATA_SECTIONS = ( ("НМЦК И ДОГОВОР", "nmck_contract"), ("ОБЪЕКТЫ И ПЛОЩАДИ", "sites_areas"), ("ПЕРСОНАЛ", "personnel"), ("ОБОРУДОВАНИЕ И ТЕХНИКА", "equipment"), ("РАСХОДНИКИ И МАТЕРИАЛЫ", "consumables"), ("ДОПОЛНИТЕЛЬНЫЕ РАБОТЫ", "additional_services"), ("РИСКИ И ТРЕБОВАНИЯ", "risks_requirements"), ("ПРОЧЕЕ", "other"), ) FACT_STATUSES = { "found": "Найдено", "missing": "Не найдено", "insufficient": "Недостаточно данных", "conflict": "Конфликт документов", "manual_review": "Требуется ручная проверка", } TENDER_DATA_PARAMETERS = { "nmck_contract": ( "НМЦК и валюта", "Разбивка начальной цены по видам работ", "Срок оказания услуг", "Лицензии", "Допуски и разрешения", "Страхование", "Арендные платежи исполнителя заказчику", ), "sites_areas": ( "Вид деятельности и назначение объекта", "Адреса объектов", "Удаленность объектов друг от друга", "Площади внутренних помещений", "Площади внешней территории", "Удаленность от остановок общественного транспорта", ), "personnel": ( "Вид и должности персонала", "Количество персонала", "Гражданство персонала", "Санитарные и медицинские книжки", "Работа по субботам и воскресеньям", "График и часы работы", "Дополнительные требования к персоналу", "Доставка персонала исполнителем", ), "equipment": ( "Оборудование для уборки и поставщик оборудования", "Механизированная уборка и спецтехника", ), "consumables": ( "Санузловые расходные материалы", "Характеристики санузловых расходников", "Количество санузловых расходников", "Количество сотрудников или посетителей для расчета расходников", "Прочие расходники и их количество", "Диспенсеры, дозаторы и электросушилки", "Материалы для ремонта и сторона предоставления", ), "additional_services": ( "Химчистка мебели, штор, тюля и ковров", "Мытье окон и витражей", "Высотные работы", "Вывоз снега", "Замена грязезащитных ковров", "Чистка кровли, сосулек и наледи", "Благоустройство и озеленение", "Дезинсекция и дератизация", "Мытье бассейна или фонтана", ), "risks_requirements": ( "Критичные условия и риски договора", "Требования, влияющие на стоимость", ), "other": ("Прочие значимые данные",), } # --- БЛОК 1: Парсеры --- def _safe_filename_from_url_or_response(url: str, response: httpx.Response) -> str: parsed = urllib.parse.urlparse(url) query = urllib.parse.parse_qs(parsed.query) filename_values = query.get("filename") if filename_values and filename_values[0].strip(): return Path(filename_values[0].strip()).name content_disposition = response.headers.get("content-disposition", "") match = re.search(r'filename\*?=(?:UTF-8\'\')?"?([^";]+)"?', content_disposition) if match: return Path(urllib.parse.unquote(match.group(1))).name suffix = Path(parsed.path).suffix if suffix: return "document" + suffix content_type = response.headers.get("content-type", "") ext = mimetypes.guess_extension(content_type.split(";")[0].strip()) or "" return "document" + ext def extract_text_from_excel(file_path: str) -> str: sheets = pd.read_excel(file_path, sheet_name=None, header=None, dtype=str) parts = [] for sheet_name, df in sheets.items(): df = df.fillna("") parts.append(f"\n\n# Лист: {sheet_name}") for row_num, row in enumerate(df.itertuples(index=False), start=1): values = [str(cell).strip() for cell in row if str(cell).strip()] if values: parts.append(f"Строка {row_num}: " + " | ".join(values)) text = "\n".join(parts).strip() if not text: raise RuntimeError("Excel-файл пустой или не удалось извлечь текст") return text def _elements_to_text(elements) -> str: parts = [] for element in elements: if not isinstance(element, dict): continue category = element.get("type") or element.get("category") or "Element" text = (element.get("text") or "").strip() metadata = element.get("metadata") or {} table_html = (metadata.get("text_as_html") or "").strip() if table_html: parts.append(f"\n\n## {category}\n{table_html}") elif text: parts.append(f"\n\n## {category}\n{text}") result = "".join(parts).strip() if not result: raise RuntimeError("Unstructured вернул пустой текст") return result async def call_smart_parser(file_path: str) -> str: if not UNSTRUCTURED_API_KEY: raise RuntimeError("UNSTRUCTURED_API_KEY is not set") filename = os.path.basename(file_path) or "document" mime_type = mimetypes.guess_type(filename)[0] or "application/octet-stream" headers = { "accept": "application/json", "unstructured-api-key": UNSTRUCTURED_API_KEY, } data = { "strategy": "hi_res", "languages": "rus", "skip_infer_table_types": "[]", "include_page_breaks": "true", } timeout_config = httpx.Timeout(900.0, read=900.0, connect=60.0) async with httpx.AsyncClient( timeout=timeout_config, follow_redirects=True ) as http_client: with open(file_path, "rb") as f: response = await http_client.post( UNSTRUCTURED_API_URL, headers=headers, data=data, files={"files": (filename, f, mime_type)}, ) try: response.raise_for_status() except httpx.HTTPStatusError as e: raise RuntimeError( f"Ошибка Unstructured API {response.status_code}: {response.text[:2000]}" ) from e elements = response.json() if not isinstance(elements, list): raise RuntimeError(f"Unstructured вернул неожиданный ответ: {elements}") return _elements_to_text(elements) # --- БЛОК 2: LLM (DeepSeek) + Интеграция с БД --- def get_pricebook_catalog_sync() -> list: catalog = [] try: p_res = ( supabase.table("personnel_rates") .select("raw_value") .eq("is_active", True) .execute() ) for r in p_res.data: if r.get("raw_value"): catalog.append(r.get("raw_value").split("|")[0].strip()) except: pass try: i_res = ( supabase.table("price_items").select("name").eq("is_active", True).execute() ) for r in i_res.data: if r.get("name"): catalog.append(r.get("name").strip()) except: pass result = sorted(list(set(catalog))) if not result: raise RuntimeError("Справочник пуст: не удалось выгрузить прайс-лист из БД") return result async def analyze_text_with_deepseek_analytics(text: str) -> dict: if not OPENROUTER_API_KEY: raise RuntimeError("OPENROUTER_API_KEY is not set") api_url = "https://openrouter.ai/api/v1/chat/completions" headers = { "Authorization": f"Bearer {OPENROUTER_API_KEY}", "Content-Type": "application/json", } system_prompt = """Ты — Главный аналитик тендерной документации и старший сметчик клининговой компании. Твоя задача — подготовить КРАТКИЙ УПРАВЛЕНЧЕСКИЙ АУДИТ тендера: понять, можно ли переходить к ручному расчету, какие условия влияют на стоимость и какие риски нужно снять до участия. НЕ рассчитывай цены, ФОТ, прибыль, налоги и итоговую стоимость. Не придумывай данные. Верни только строгий JSON по указанной схеме: { "tender_card": { "nmck": {"value": "", "status": "found|missing|insufficient|conflict|manual_review", "source_document": "", "raw_context": ""}, "service_term": {"value": "", "status": "", "source_document": "", "raw_context": ""}, "object_type": {"value": "", "status": "", "source_document": "", "raw_context": ""}, "addresses_and_areas": {"value": "", "status": "", "source_document": "", "raw_context": ""}, "personnel_summary": {"value": "", "status": "", "source_document": "", "raw_context": ""}, "work_schedule": {"value": "", "status": "", "source_document": "", "raw_context": ""} }, "cost_drivers": [ {"factor": "", "impact": "", "status": "", "source_document": "", "raw_context": ""} ], "risks": [ {"risk": "", "severity": "high|medium|low", "status": "", "source_document": "", "raw_context": "", "recommendation": ""} ], "conflicts": [ {"topic": "", "details": "", "source_document": "", "raw_context": ""} ], "critical_questions": [ {"question": "", "reason": "", "severity": "high|medium|low"} ], "overall_assessment": { "readiness": "ready|needs_clarification|high_risk", "summary": "" } } ПРАВИЛА: 1. Статусы: found — данные достаточны; missing — данных нет; insufficient — данные есть, но для решения недостаточны; conflict — документы противоречат; manual_review — нужна ручная трактовка. 2. У каждого найденного факта указывай source_document и дословную raw_context. Для missing источник и цитата могут быть пустыми. 3. При конфликте НЕ выбирай версию самостоятельно. Зафиксируй оба значения и источники в details/raw_context, а также добавь запись в conflicts. 4. К критичным вопросам относись строго: включай только то, что блокирует расчет, меняет цену или создает договорной риск. 5. В cost_drivers отрази персонал и графики, площади, работу по выходным, расходники, оборудование, спецтехнику, дополнительные работы, удаленность, лицензии и другие условия, влияющие на стоимость. 6. В risks отрази лицензии, допуски, медкнижки, гражданство, страхование, аренду, сторону предоставления оборудования/материалов, неясные объемы и рискованные договорные требования. 7. Если по карточке тендера данных нет — ставь status="missing", value="". НИКОГДА НЕ ДОДУМЫВАЙ. 8. В overall_assessment.readiness: ready — критичных пробелов нет; needs_clarification — есть вопросы до расчета; high_risk — есть существенные риски или конфликты. """ payload = { "model": "deepseek/deepseek-v4-pro", "messages": [ {"role": "system", "content": system_prompt}, {"role": "user", "content": f"Проанализируй текст и верни JSON:\n\n{text}"}, ], "reasoning": {"enabled": True}, "temperature": 0.1, "response_format": {"type": "json_object"}, } async with httpx.AsyncClient(timeout=240.0) as client: response = await client.post(api_url, headers=headers, json=payload) response.raise_for_status() return _parse_llm_json( response.json()["choices"][0]["message"] .get("content", "") .replace("```json", "") .replace("```", "") .strip() ) async def extract_tender_data_for_estimator(text: str) -> dict: if not OPENROUTER_API_KEY: raise RuntimeError("OPENROUTER_API_KEY is not set") api_url = "https://openrouter.ai/api/v1/chat/completions" headers = { "Authorization": f"Bearer {OPENROUTER_API_KEY}", "Content-Type": "application/json", } sections = ", ".join(section_key for _, section_key in TENDER_DATA_SECTIONS) system_prompt = f"""Ты — старший сметчик клининговой компании. Прочитай ОБЪЕДИНЕННЫЙ текст документов тендера и подготовь КАРТУ ИСХОДНЫХ ДАННЫХ для человека-сметчика. Твоя задача — только извлечь факты. НЕ рассчитывай цены, ФОТ, прибыль, налоги и не сопоставляй с прайс-листом. РАЗДЕЛЫ JSON: {sections}. Для каждого найденного факта верни объект: {{ "parameter": "точное название проверяемого параметра", "status": "found | missing | insufficient | conflict | manual_review", "value": "значение только из документов", "source_document": "имя документа после заголовка ДОКУМЕНТ", "source_location": "страница, лист, строка или раздел, если есть", "raw_context": "точная цитата из документа", "comment": "что должен проверить сметчик" }} ПРАВИЛА: 1. Не придумывай данные. Если значение отсутствует — status="missing", value="", а в comment сформулируй вопрос заказчику. 2. Если данные есть, но не позволяют считать (например, расходники без количества) — status="insufficient" и укажи чего не хватает. 3. Если документы противоречат друг другу — status="conflict"; сохрани оба значения и оба источника в value/raw_context. Не выбирай одно значение сам. 4. Если формулировка двусмысленна — status="manual_review" и объясни причину. 5. Для каждого статуса кроме missing обязательно заполняй source_document и raw_context дословной цитатой. 6. Все количественные значения сохраняй вместе с единицами измерения, периодичностью, площадью, графиком или стороной предоставления, когда они указаны. 7. Обязательно ищи: НМЦК и разбивку цены; срок, лицензии, допуски, страхование и аренду; типы объектов, адреса, площади; персонал, гражданство, медкнижки и графики; оборудование и кто его предоставляет; расходники и характеристики; все дополнительные работы. 8. Не создавай разделов вне списка. Возвращай только строгий JSON. Формат ответа: {{ "sections": {{ "personnel": [{{"parameter": "Количество персонала", "status": "found", "value": "8 уборщиков", "source_document": "ТЗ.pdf", "source_location": "стр. 12", "raw_context": "...", "comment": ""}}] }}, "critical_questions": ["вопрос заказчику, без которого нельзя рассчитывать"], "documents_summary": ["краткое замечание по документу"] }}""" payload = { "model": "deepseek/deepseek-v4-pro", "messages": [ {"role": "system", "content": system_prompt}, { "role": "user", "content": f"Вот объединенный текст всех документов тендера:\n\n{text}", }, ], "reasoning": {"enabled": True}, "temperature": 0.1, "response_format": {"type": "json_object"}, } async with httpx.AsyncClient(timeout=240.0) as client: response = await client.post(api_url, headers=headers, json=payload) response.raise_for_status() return _parse_llm_json( response.json()["choices"][0]["message"] .get("content", "") .replace("```json", "") .replace("```", "") .strip() ) # --- БЛОК 3: Обработка очередей --- async def download_file(url: str, temp_dir: str) -> str: if "yandex.ru" in url or "ya.disk" in url: async with httpx.AsyncClient(follow_redirects=True) as client: resp = await client.get( f"https://cloud-api.yandex.net/v1/disk/public/resources/download?public_key={urllib.parse.quote(url)}" ) resp.raise_for_status() url = resp.json().get("href", url) async with httpx.AsyncClient(follow_redirects=True) as client: resp = await client.get(url) resp.raise_for_status() filename = _safe_filename_from_url_or_response(url, resp) path = os.path.join(temp_dir, filename) with open(path, "wb") as f: f.write(resp.content) return path # --- БЛОК 2.5: МАГИЧЕСКИЙ КАЛЬКУЛЯТОР СМЕТЫ --- class TenderCalculator: def __init__(self, supabase_client): self.supabase = supabase_client self.price_book = {} self.normalized_lookup = {} self.cleaning_group_id = None self._load_reference_data() def _load_reference_data(self): groups = self.supabase.table("payroll_groups").select("id, code").execute() for g in groups.data: if g.get("code") == "CLEANING_STAFF": self.cleaning_group_id = g.get("id") if self.cleaning_group_id is None: raise RuntimeError("В payroll_groups отсутствует код CLEANING_STAFF.") personnel = ( self.supabase.table("personnel_rates") .select("raw_value, rate_min, payroll_group_id") .eq("is_active", True) .execute() ) for p in personnel.data: val = p.get("raw_value", "") if val: key = val.split("|")[0].strip() self.price_book[key] = { "type": "personnel", "price": float(p.get("rate_min") or 0), "group_id": p.get("payroll_group_id"), } self.normalized_lookup[_normalize_catalog_key(key)] = key items = ( self.supabase.table("price_items") .select("name, price_min, category_id") .eq("is_active", True) .execute() ) for i in items.data: name = i.get("name") if name: key = name.strip() self.price_book[key] = { "type": "item", "price": float(i.get("price_min") or 0), "category_id": i.get("category_id"), } self.normalized_lookup[_normalize_catalog_key(key)] = key def _find_price_entry(self, name: str): # сначала точное совпадение, затем поиск без учёта регистра/пробелов/тире if not name: return None, None if name in self.price_book: return self.price_book[name], name canonical_name = self.normalized_lookup.get(_normalize_catalog_key(name)) if canonical_name: return self.price_book[canonical_name], canonical_name return None, None def calculate(self, extracted_data: dict) -> dict: estimate = { "details": [], "totals": { "cleaning_payroll": 0, "other_payroll": 0, "materials_and_items": 0, "calculated_chemistry": 0, "grand_total": 0, }, "unrecognized": [ _normalize_unrecognized_entry(entry, source="llm_unrecognized") for entry in extracted_data.get("unrecognized_items", []) ], } for item in extracted_data.get("extracted_items", []): name, quantity = item.get("exact_name"), _parse_quantity( item.get("quantity", 0) ) if not name or quantity <= 0: continue ref_data, canonical_name = self._find_price_entry(name) if ref_data: price = ref_data["price"] total_cost = price * quantity estimate["details"].append( { "name": canonical_name or name, "quantity": quantity, "unit_price": price, "total_cost": total_cost, "type": ref_data["type"], "raw_context": item.get("raw_context", ""), } ) if ref_data["type"] == "personnel": if ref_data["group_id"] == self.cleaning_group_id: estimate["totals"]["cleaning_payroll"] += total_cost else: estimate["totals"]["other_payroll"] += total_cost else: estimate["totals"]["materials_and_items"] += total_cost else: estimate["unrecognized"].append( { "source": "no_price_match", "description": name, "quantity": quantity, "raw_context": item.get("raw_context", ""), } ) chemistry_cost = estimate["totals"]["cleaning_payroll"] * 0.10 estimate["totals"]["calculated_chemistry"] = chemistry_cost if chemistry_cost > 0: estimate["details"].append( { "name": "Бюджет химии и расходников (10% от ФОТ уборщиц)", "quantity": 1, "unit_price": chemistry_cost, "total_cost": chemistry_cost, "type": "calculated_rule", "raw_context": "Расчётное правило: не найдено в тексте документа", } ) estimate["totals"]["grand_total"] = sum( [ estimate["totals"]["cleaning_payroll"], estimate["totals"]["other_payroll"], estimate["totals"]["materials_and_items"], estimate["totals"]["calculated_chemistry"], ] ) return estimate # --- БЛОК 2.7: ГЕНЕРАТОР EXCEL --- def create_excel_report(estimate_data: dict, temp_dir: str, task_id: str) -> str: wb = openpyxl.Workbook() ws_summary = wb.active ws_summary.title = "СВОДНАЯ" headers_font, header_fill, bold_font = ( Font(bold=True, color="FFFFFF"), PatternFill("solid", fgColor="4F81BD"), Font(bold=True), ) ws_summary.append(["Статья расходов", "Сумма (руб.)"]) ws_summary.cell(row=1, column=1).font, ws_summary.cell(row=1, column=1).fill = ( headers_font, header_fill, ) ws_summary.cell(row=1, column=2).font, ws_summary.cell(row=1, column=2).fill = ( headers_font, header_fill, ) totals = estimate_data.get("totals", {}) ws_summary.append(["ФОТ (Уборочный персонал)", totals.get("cleaning_payroll", 0)]) ws_summary.append(["ФОТ (Прочий персонал)", totals.get("other_payroll", 0)]) ws_summary.append( ["Материалы, инвентарь, техника", totals.get("materials_and_items", 0)] ) ws_summary.append( ["Химия и расходники (по правилу 10%)", totals.get("calculated_chemistry", 0)] ) ws_summary.append(["", ""]) ws_summary.append(["ИТОГО РАСХОДЫ (Себестоимость):", "=SUM(B2:B5)"]) ws_summary.cell(row=7, column=1).font = bold_font ws_summary.cell(row=7, column=2).font = bold_font ws_summary.append(["Налоги УСН (6%):", "=B7*0.06"]) ws_summary.append(["Желаемая прибыль (10%):", "=B7*0.10"]) ws_summary.append(["", ""]) ws_summary.append(["ИТОГО СТОИМОСТЬ КЛИЕНТУ:", "=B7+B8+B9"]) ws_summary.cell(row=11, column=1).font, ws_summary.cell(row=11, column=2).font = ( Font(bold=True, color="FF0000", size=12), Font(bold=True, color="FF0000", size=12), ) ws_summary.column_dimensions["A"].width, ws_summary.column_dimensions["B"].width = ( 45, 20, ) ws_details = wb.create_sheet(title="ДЕТАЛИЗАЦИЯ") ws_details.append( [ "Категория", "Наименование", "Кол-во", "Цена за ед.", "Итоговая сумма", "Цитата из ТЗ", ] ) for col in range(1, 7): ( ws_details.cell(row=1, column=col).font, ws_details.cell(row=1, column=col).fill, ) = (headers_font, header_fill) type_labels = { "personnel": "Персонал", "item": "Материал/инвентарь", "calculated_rule": "Расчётное правило (не из ТЗ)", } row_idx = 2 for item in estimate_data.get("details", []): category = type_labels.get( item.get("type", "unknown"), item.get("type", "unknown") ) ws_details.append( [ category, item.get("name", ""), item.get("quantity", 0), item.get("unit_price", 0), f"=C{row_idx}*D{row_idx}", item.get("raw_context", ""), ] ) row_idx += 1 for idx, w in enumerate([20, 55, 10, 15, 15, 50], 1): ws_details.column_dimensions[openpyxl.utils.get_column_letter(idx)].width = w ws_unrecognized = wb.create_sheet(title="НЕРАСПОЗНАНО") ws_unrecognized.append(["Источник", "Описание", "Кол-во", "Цитата из ТЗ"]) for col in range(1, 5): ( ws_unrecognized.cell(row=1, column=col).font, ws_unrecognized.cell(row=1, column=col).fill, ) = (headers_font, PatternFill("solid", fgColor="C0504D")) source_labels = { "llm_unrecognized": "Не найдено в справочнике (по мнению ИИ)", "no_price_match": "Не сопоставлено с прайс-листом", } unrecognized_items = estimate_data.get("unrecognized", []) if not unrecognized_items: ws_unrecognized.append(["-", "Нераспознанных позиций нет", "-", "-"]) else: for entry in unrecognized_items: normalized = _normalize_unrecognized_entry(entry, source="unknown") ws_unrecognized.append( [ source_labels.get(normalized["source"], normalized["source"]), normalized["description"], normalized["quantity"], normalized["raw_context"], ] ) for idx, w in enumerate([35, 45, 10, 60], 1): ws_unrecognized.column_dimensions[ openpyxl.utils.get_column_letter(idx) ].width = w file_path = os.path.join(temp_dir, f"Smeta_Tender_{task_id[:8]}.xlsx") wb.save(file_path) return file_path def upload_excel_to_supabase(file_path: str, filename: str) -> str: with open(file_path, "rb") as f: supabase.storage.from_("tenders").upload( filename, f.read(), file_options={ "content-type": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet" }, ) public_url = supabase.storage.from_("tenders").get_public_url(filename) return ( public_url.get("publicURL") or public_url.get("public_url") or "" if isinstance(public_url, dict) else str(public_url) ) def _normalize_tender_fact(fact, parameter: str) -> dict: fact = fact if isinstance(fact, dict) else {} status = fact.get("status", "missing") if status not in FACT_STATUSES: status = "manual_review" return { "parameter": fact.get("parameter") or parameter, "status": status, "value": fact.get("value", ""), "source_document": fact.get("source_document", ""), "source_location": fact.get("source_location", ""), "raw_context": fact.get("raw_context", ""), "comment": fact.get("comment", ""), } def normalize_tender_data(raw_data: dict) -> dict: raw_sections = raw_data.get("sections", {}) normalized_sections = {} clarification_questions = list(raw_data.get("critical_questions", [])) for _, section_key in TENDER_DATA_SECTIONS: facts_by_parameter = {} for fact in raw_sections.get(section_key, []): normalized = _normalize_tender_fact(fact, "Прочие значимые данные") facts_by_parameter.setdefault(normalized["parameter"], []).append( normalized ) section_facts = [] for parameter in TENDER_DATA_PARAMETERS[section_key]: matched_facts = facts_by_parameter.pop(parameter, []) if matched_facts: section_facts.extend(matched_facts) else: section_facts.append( _normalize_tender_fact( { "status": "missing", "comment": f"Уточнить у заказчика: {parameter}.", }, parameter, ) ) for extra_facts in facts_by_parameter.values(): section_facts.extend(extra_facts) normalized_sections[section_key] = section_facts for facts in normalized_sections.values(): for fact in facts: if fact["status"] in {"missing", "insufficient", "conflict"}: question = fact["comment"] or fact["parameter"] if question not in clarification_questions: clarification_questions.append(question) return { "sections": normalized_sections, "critical_questions": clarification_questions, "documents_summary": raw_data.get("documents_summary", []), } def create_tender_data_report(data: dict, temp_dir: str, task_id: str) -> str: wb = openpyxl.Workbook() header_font = Font(bold=True, color="FFFFFF") header_fill = PatternFill("solid", fgColor="4F81BD") status_fills = { "found": PatternFill("solid", fgColor="C6EFCE"), "missing": PatternFill("solid", fgColor="FFC7CE"), "insufficient": PatternFill("solid", fgColor="FFEB9C"), "conflict": PatternFill("solid", fgColor="F4B183"), "manual_review": PatternFill("solid", fgColor="D9EAD3"), } columns = [ "Параметр", "Статус", "Найденные данные", "Документ-источник", "Страница / лист / строка", "Цитата из документа", "Комментарий для сметчика", "Решение сметчика", ] ws_summary = wb.active ws_summary.title = "СВОДКА" ws_summary.append(["Карта исходных данных тендера", ""]) ws_summary.cell(row=1, column=1).font = Font(bold=True, size=14) ws_summary.append(["Раздел", "Найдено", "Пробелы/уточнения", "Конфликты"]) for column in range(1, 5): ws_summary.cell(row=2, column=column).font = header_font ws_summary.cell(row=2, column=column).fill = header_fill row = 3 total_found = total_gaps = total_conflicts = 0 for section_name, section_key in TENDER_DATA_SECTIONS: facts = data["sections"][section_key] found = sum(fact["status"] == "found" for fact in facts) gaps = sum(fact["status"] in {"missing", "insufficient"} for fact in facts) conflicts = sum(fact["status"] == "conflict" for fact in facts) ws_summary.append([section_name, found, gaps, conflicts]) total_found += found total_gaps += gaps total_conflicts += conflicts row += 1 ws_summary.append(["ИТОГО", total_found, total_gaps, total_conflicts]) for column in range(1, 5): ws_summary.cell(row=row, column=column).font = Font(bold=True) ws_summary.append([]) ws_summary.append(["Документы", ""]) for document_summary in data.get("documents_summary", []): ws_summary.append([str(document_summary), ""]) for column, width in enumerate([42, 15, 20, 15], start=1): ws_summary.column_dimensions[openpyxl.utils.get_column_letter(column)].width = ( width ) ws_questions = wb.create_sheet(title="КРИТИЧНЫЕ УТОЧНЕНИЯ") ws_questions.append(["Вопрос или риск", "Действие сметчика"]) for column in range(1, 3): ws_questions.cell(row=1, column=column).font = header_font ws_questions.cell(row=1, column=column).fill = PatternFill( "solid", fgColor="C0504D" ) questions = data.get("critical_questions", []) if questions: for question in questions: ws_questions.append([str(question), "Запросить уточнение у заказчика"]) else: ws_questions.append(["Критичных уточнений не выявлено", ""]) ws_questions.column_dimensions["A"].width = 90 ws_questions.column_dimensions["B"].width = 38 ws_request = wb.create_sheet(title="ЗАПРОС ЗАКАЗЧИКУ") ws_request.append(["№", "Вопрос заказчику", "Причина"]) for column in range(1, 4): ws_request.cell(row=1, column=column).font = header_font ws_request.cell(row=1, column=column).fill = header_fill for index, question in enumerate(questions, start=1): ws_request.append([index, str(question), "Данные нужны для расчета сметы"]) if not questions: ws_request.append([1, "Уточнения не требуются", ""]) ws_request.column_dimensions["A"].width = 8 ws_request.column_dimensions["B"].width = 90 ws_request.column_dimensions["C"].width = 38 ws_sources = wb.create_sheet(title="ИСТОЧНИКИ") ws_sources.append(["Раздел", *columns[:-1]]) for column in range(1, 8): ws_sources.cell(row=1, column=column).font = header_font ws_sources.cell(row=1, column=column).fill = header_fill for section_name, section_key in TENDER_DATA_SECTIONS: ws_section = wb.create_sheet(title=section_name[:31]) ws_section.append(columns) for column in range(1, len(columns) + 1): ws_section.cell(row=1, column=column).font = header_font ws_section.cell(row=1, column=column).fill = header_fill for fact in data["sections"][section_key]: values = [ fact["parameter"], FACT_STATUSES[fact["status"]], fact["value"], fact["source_document"], fact["source_location"], fact["raw_context"], fact["comment"], "", ] ws_section.append(values) current_row = ws_section.max_row ws_section.cell(row=current_row, column=2).fill = status_fills[ fact["status"] ] ws_sources.append([section_name, *values[:-1]]) for column, width in enumerate([34, 24, 48, 28, 24, 65, 50, 30], start=1): ws_section.column_dimensions[ openpyxl.utils.get_column_letter(column) ].width = width ws_section.freeze_panes = "A2" ws_section.auto_filter.ref = ws_section.dimensions for worksheet_row in ws_section.iter_rows(min_row=2): for cell in worksheet_row: cell.alignment = Alignment(wrap_text=True, vertical="top") for column, width in enumerate([28, 34, 24, 48, 28, 24, 65], start=1): ws_sources.column_dimensions[openpyxl.utils.get_column_letter(column)].width = ( width ) ws_sources.freeze_panes = "A2" ws_sources.auto_filter.ref = ws_sources.dimensions for worksheet in (ws_summary, ws_questions, ws_request, ws_sources): for worksheet_row in worksheet.iter_rows(): for cell in worksheet_row: cell.alignment = Alignment(wrap_text=True, vertical="top") file_path = os.path.join(temp_dir, f"Karta_Dannyh_Tendera_{task_id[:8]}.xlsx") wb.save(file_path) return file_path # --- БЛОК 2.8: Финализация пакета после успешной обработки всех документов --- def _normalize_extracted_data(raw): if isinstance(raw, dict): return raw if isinstance(raw, str): try: parsed = json.loads(raw) if isinstance(parsed, dict): return parsed except json.JSONDecodeError: return {} return {} async def finalize_task_if_ready(task_id: str, mode: str): docs_resp = await asyncio.to_thread( lambda: supabase.table("documents") .select("id, status, document_name, extracted_data") .eq("task_id", task_id) .execute() ) docs = docs_resp.data or [] if not docs: raise RuntimeError(f"По task_id={task_id} не найдено документов") statuses = [doc.get("status", "pending") for doc in docs] if any(status in ("pending", "processing") for status in statuses): return if any(status == "error" for status in statuses): return extracted_by_doc = [] for doc in docs: doc_name = doc.get("document_name", "Unknown_Doc") parsed_data = _normalize_extracted_data(doc.get("extracted_data")) text = parsed_data.get("parsed_text", "") if not text: raise RuntimeError(f"У документа '{doc_name}' отсутствует parsed_text") extracted_by_doc.append((doc_name, text)) combined_text = "" for doc_name, text in extracted_by_doc: combined_text += f"\n\n{'='*30}\nДОКУМЕНТ: {doc_name}\n{'='*30}\n\n{text}" if mode == "analytics": final_result = await analyze_text_with_deepseek_analytics(combined_text) elif mode == "estimate": raw_tender_data = await extract_tender_data_for_estimator(combined_text) final_result = normalize_tender_data(raw_tender_data) final_result["report_type"] = "tender_data_map" final_result["report_note"] = ( "Карта исходных данных для ручного расчета сметчиком. " "Финансовый расчет не выполнялся." ) with tempfile.TemporaryDirectory() as temp_dir: excel_path = await asyncio.to_thread( create_tender_data_report, final_result, temp_dir, task_id ) final_result["excel_download_url"] = await asyncio.to_thread( upload_excel_to_supabase, excel_path, f"Karta_Dannyh_Tendera_{task_id[:8]}.xlsx", ) else: raise RuntimeError(f"Неизвестный режим: {mode}") payload = { "final_result": final_result, "stage": "finalized", } await asyncio.to_thread( lambda: supabase.table("documents") .update({"status": "completed", "extracted_data": payload}) .eq("task_id", task_id) .execute() ) # --- ИЗМЕНЕНИЕ: Воркер теперь обрабатывает документы поштучно и ждёт все документы перед расчетом --- async def process_message(message: aio_pika.IncomingMessage): async with message.process(requeue=True, reject_on_redelivered=True): body = json.loads(message.body.decode()) task_id = body.get("task_id") mode = body.get("mode", "analytics") if not task_id: raise RuntimeError("Некорректное сообщение: нет task_id") try: docs_resp = await asyncio.to_thread( lambda: supabase.table("documents") .select("id, document_url, document_name, status") .eq("task_id", task_id) .execute() ) documents = docs_resp.data or [] if not documents: raise RuntimeError(f"По task_id={task_id} не найдено документов") with tempfile.TemporaryDirectory() as temp_dir: for doc in documents: doc_id = doc.get("id") url = doc.get("document_url") doc_name = doc.get("document_name") or "Unknown_Doc" if not doc_id or not url: continue await asyncio.to_thread( lambda current_doc_id=doc_id: supabase.table("documents") .update( { "status": "processing", "extracted_data": {"stage": "downloading"}, } ) .eq("id", current_doc_id) .execute() ) file_path = await download_file(url, temp_dir) ext = Path(file_path).suffix.lower() if ext in [".xlsx", ".xls"]: text = extract_text_from_excel(file_path) else: text = await call_smart_parser(file_path) payload = { "stage": "parsed", "parsed_text": text, "source_file": os.path.basename(file_path), "document_name": doc_name, } await asyncio.to_thread( lambda current_doc_id=doc_id, current_payload=payload: supabase.table( "documents" ) .update( { "status": "completed", "extracted_data": current_payload, } ) .eq("id", current_doc_id) .execute() ) await finalize_task_if_ready(task_id, mode) except Exception as e: error_msg = str(e) if "ReadTimeout" in error_msg or "timeout" in error_msg.lower(): error_msg = "Пакет документов слишком тяжелый. Сервер парсинга не смог прочитать документы за 15 минут. Вырежи лишние страницы из ТЗ." await asyncio.to_thread( lambda: supabase.table("documents") .update({"status": "error", "extracted_data": {"error": error_msg}}) .eq("task_id", task_id) .execute() ) async def main(): conn = await aio_pika.connect_robust(CLOUDAMQP_URL) async with conn: channel = await conn.channel() await channel.set_qos(prefetch_count=1) queue = await channel.declare_queue("tender_tasks", durable=True) await queue.consume(process_message) await asyncio.Future() if __name__ == "__main__": asyncio.run(main())