"""spec_events — event sourcing: apply ops, reset contract.""" import json, uuid from db.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, )