from __future__ import annotations import uuid from dataclasses import dataclass from typing import Any from psycopg2.extras import Json, RealDictCursor from .auth import CurrentUser from .db import get_connection from .settings import get_settings ALLOWED_ANSWER_CODES = { "AUTHORIZED_LEGACY_PRODUCTION", "ODOO_START_OMITTED", "TEST_OR_SETUP", "RESIDUAL_MATERIAL", "MAINTENANCE", "NOT_PRODUCTION", "ODOO_DOUBLE_START_ONE_VALID", "ODOO_DOUBLE_START_BOTH_INVALID", "ODOO_CONSECUTIVE_ORDERS_NOT_CLOSED", "DATA_ISSUE", "OTHER", } MAPPED_CLASSIFICATION = { "AUTHORIZED_LEGACY_PRODUCTION": "AUTHORIZED_LEGACY_PRODUCTION", "ODOO_START_OMITTED": "ODOO_START_OMITTED", "TEST_OR_SETUP": "TEST_OR_SETUP", "RESIDUAL_MATERIAL": "RESIDUAL_MATERIAL", "MAINTENANCE": "OTHER", "NOT_PRODUCTION": "NOT_RELEVANT", "ODOO_DOUBLE_START_ONE_VALID": "ODOO_DOUBLE_START_ONE_VALID", "ODOO_DOUBLE_START_BOTH_INVALID": "ODOO_DOUBLE_START_BOTH_INVALID", "ODOO_CONSECUTIVE_ORDERS_NOT_CLOSED": "ODOO_CONSECUTIVE_ORDERS_NOT_CLOSED", "DATA_ISSUE": "DATA_ISSUE", "OTHER": "OTHER", } CONTEXT_SOURCE_TYPE = { "AUTHORIZED_LEGACY_PRODUCTION": "LEGACY_ERP", "TEST_OR_SETUP": "TEST_SETUP", "MAINTENANCE": "MAINTENANCE", } @dataclass(frozen=True) class ReviewInput: decision_group_key: str selected_case_ids: tuple[int, ...] answer_code: str notes: str selected_order_ref: str | None selected_workorder_id: int | None corrected_in_odoo: bool create_contexts: bool confirmation: bool class ReviewError(RuntimeError): pass def json_safe(value: Any) -> Any: if isinstance(value, dict): return { str(key): json_safe(item) for key, item in value.items() } if isinstance(value, (list, tuple, set)): return [ json_safe(item) for item in value ] if hasattr(value, "isoformat"): return value.isoformat() return value def fetch_decision( cursor: RealDictCursor, decision_group_key: str, ) -> dict[str, Any]: cursor.execute( """ SELECT * FROM mv_reports_ucepsa_prod .li_context_governance_decision_evidence_pack_v1 WHERE decision_group_key = %s """, (decision_group_key,), ) row = cursor.fetchone() if not row: raise ReviewError( "La decisión ya no existe." ) return dict(row) def fetch_cases_for_update( cursor: RealDictCursor, decision_group_key: str, ) -> list[dict[str, Any]]: cursor.execute( """ SELECT c.* FROM mv_reports_ucepsa_prod .li_context_decision_group_cases_v1 c JOIN mv_loss_intelligence .context_review_cases raw ON raw.case_id = c.case_id WHERE c.decision_group_key = %s ORDER BY c.case_id FOR UPDATE OF raw """, (decision_group_key,), ) return [ dict(row) for row in cursor.fetchall() ] def validate_input( decision: dict[str, Any], all_cases: list[dict[str, Any]], review: ReviewInput, ) -> tuple[ list[dict[str, Any]], list[dict[str, Any]], str, ]: if not review.confirmation: raise ReviewError( "Debe confirmar que la respuesta refleja la realidad de planta." ) if review.answer_code not in ALLOWED_ANSWER_CODES: raise ReviewError( "Código de respuesta no permitido." ) pending_cases = [ case for case in all_cases if case["case_review_status"] == "PENDING" ] if not pending_cases: raise ReviewError( "La decisión ya no tiene casos pendientes." ) selected_ids = set( review.selected_case_ids ) pending_ids = { int(case["case_id"]) for case in pending_cases } if not selected_ids: raise ReviewError( "Debe seleccionar al menos un intervalo o caso." ) if not selected_ids.issubset( pending_ids ): raise ReviewError( "La selección contiene casos que ya no están pendientes." ) selected_cases = [ case for case in pending_cases if int(case["case_id"]) in selected_ids ] family = str( decision["decision_family"] ) if family == "ODOO_ORDER_START_CONFLICT": allowed = { "ODOO_DOUBLE_START_ONE_VALID", "ODOO_DOUBLE_START_BOTH_INVALID", "ODOO_CONSECUTIVE_ORDERS_NOT_CLOSED", "DATA_ISSUE", "OTHER", } if review.answer_code not in allowed: raise ReviewError( "La respuesta no corresponde a un conflicto de arranque Odoo." ) if ( review.answer_code == "ODOO_DOUBLE_START_ONE_VALID" ): if ( not review.selected_order_ref or review.selected_workorder_id is None ): raise ReviewError( "Debe seleccionar la orden y el workorder válidos." ) valid_pairs = { ( str(item["production_order"]), int(item["odoo_workorder_id"]), ) for item in ( decision.get( "episode_evidence" ) or [] ) if item.get( "production_order" ) and item.get( "odoo_workorder_id" ) is not None } if ( review.selected_order_ref, review.selected_workorder_id, ) not in valid_pairs: raise ReviewError( "La orden seleccionada no pertenece a la evidencia." ) if not review.corrected_in_odoo: raise ReviewError( "Debe confirmar que la incoherencia ya fue corregida en Odoo." ) elif family == "UNCONFIRMED_ACTIVITY": allowed = { "AUTHORIZED_LEGACY_PRODUCTION", "ODOO_START_OMITTED", "TEST_OR_SETUP", "RESIDUAL_MATERIAL", "MAINTENANCE", "NOT_PRODUCTION", "OTHER", } if review.answer_code not in allowed: raise ReviewError( "La respuesta no corresponde a actividad sin contexto." ) if review.create_contexts: if ( review.answer_code not in CONTEXT_SOURCE_TYPE ): raise ReviewError( "Esta respuesta no permite crear contextos retrospectivos." ) for case in selected_cases: if ( case["case_status"] != "CLOSED" or case["ended_at"] is None ): raise ReviewError( "Solo pueden crearse contextos para intervalos terminados." ) if not case["machine_id"]: raise ReviewError( "El caso no tiene máquina asociada." ) mapped = MAPPED_CLASSIFICATION[ review.answer_code ] return ( pending_cases, selected_cases, mapped, ) def insert_contexts( cursor: RealDictCursor, decision: dict[str, Any], selected_cases: list[dict[str, Any]], review: ReviewInput, user: CurrentUser, request_id: uuid.UUID, ) -> list[int]: if not review.create_contexts: return [] source_type = CONTEXT_SOURCE_TYPE[ review.answer_code ] settings = get_settings() created_ids: list[int] = [] for case in selected_cases: case_id = int(case["case_id"]) source_ref = ( f"GOVERNANCE_UI_CASE:{case_id}:" f"{review.answer_code}" ) cursor.execute( """ SELECT context_session_id FROM mv_loss_intelligence .production_context_sessions WHERE source_ref = %s """, (source_ref,), ) existing = cursor.fetchone() if existing: created_ids.append( int( existing[ "context_session_id" ] ) ) continue cursor.execute( """ SELECT COUNT(*) AS overlap_count FROM ( SELECT 1 FROM mv_loss_intelligence .production_context_sessions s WHERE s.machine_id = %s AND s.started_at < %s AND COALESCE( s.ended_at, 'infinity'::timestamptz ) > %s UNION ALL SELECT 1 FROM mv_loss_intelligence .odoo_shopfloor_sessions o WHERE o.machine_id = %s AND o.started_at IS NOT NULL AND o.started_at < %s AND COALESCE( o.ended_at, 'infinity'::timestamptz ) > %s AND o.session_status IN ( 'OPEN', 'CLOSED' ) ) overlaps """, ( case["machine_id"], case["ended_at"], case["first_seen_at"], case["machine_id"], case["ended_at"], case["first_seen_at"], ), ) overlap_count = int( cursor.fetchone()["overlap_count"] ) if overlap_count > 0: raise ReviewError( f"El caso {case_id} se solapa con otro contexto autoritativo. " "No se ha creado ningún contexto." ) cursor.execute( """ INSERT INTO mv_loss_intelligence .production_context_sessions ( tenant, site, machine_id, source_type, source_ref, started_at, ended_at, status, authorization_status, confidence, official_eligible, authorized_by, authorized_at, closed_by, notes, evidence_json ) VALUES ( %s, %s, %s, %s, %s, %s, %s, 'CLOSED', 'AUTHORIZED', 'MEDIUM', false, %s, now(), %s, %s, %s ) RETURNING context_session_id """, ( settings.tenant, settings.site, case["machine_id"], source_type, source_ref, case["first_seen_at"], case["ended_at"], user.display_name, user.display_name, review.notes, Json({ "created_with": "governance_review_ui_v0314", "request_id": str(request_id), "decision_group_key": decision[ "decision_group_key" ], "case_id": case_id, "answer_code": review.answer_code, "retrospective": True, "deployment_mode": "SHADOW", }), ), ) created_ids.append( int( cursor.fetchone()[ "context_session_id" ] ) ) return created_ids def insert_case_and_episode_reviews( cursor: RealDictCursor, decision: dict[str, Any], selected_cases: list[dict[str, Any]], mapped_classification: str, review: ReviewInput, user: CurrentUser, request_id: uuid.UUID, ) -> tuple[list[int], list[int]]: case_review_ids: list[int] = [] episode_review_ids: list[int] = [] new_status = ( "DISMISSED" if mapped_classification == "NOT_RELEVANT" else "REVIEWED" ) for case in selected_cases: case_id = int(case["case_id"]) cursor.execute( """ INSERT INTO mv_loss_intelligence .context_review_case_reviews ( case_id, previous_values, new_review_status, resolution_classification, selected_order_ref, selected_workorder_id, reviewed_by, notes, evidence_json ) VALUES ( %s, %s, %s, %s, %s, %s, %s, %s, %s ) RETURNING review_id """, ( case_id, Json( json_safe({ "review_status": case[ "case_review_status" ], "resolution_classification": case[ "case_resolution_classification" ], "reviewed_by": case[ "case_reviewed_by" ], "reviewed_at": case[ "case_reviewed_at" ], "review_notes": case[ "case_review_notes" ], }) ), new_status, mapped_classification, review.selected_order_ref, review.selected_workorder_id, user.display_name, review.notes, Json({ "source": "governance_review_ui_v0314", "request_id": str(request_id), "decision_group_key": decision[ "decision_group_key" ], "answer_code": review.answer_code, }), ), ) case_review_ids.append( int( cursor.fetchone()[ "review_id" ] ) ) cursor.execute( """ SELECT h.* FROM mv_loss_intelligence .context_review_case_evidence e JOIN mv_reports_ucepsa_prod .li_shopfloor_context_incident_history_v1 h ON h.episode_id = e.episode_id WHERE e.case_id = %s ORDER BY h.episode_id """, (case_id,), ) episodes = [ dict(row) for row in cursor.fetchall() ] episode_classification = ( "DATA_ISSUE" if mapped_classification in { "ODOO_DOUBLE_START_ONE_VALID", "ODOO_DOUBLE_START_BOTH_INVALID", "ODOO_CONSECUTIVE_ORDERS_NOT_CLOSED", } else mapped_classification ) for episode in episodes: if ( episode["review_status"] != "PENDING" ): continue cursor.execute( """ INSERT INTO mv_loss_intelligence .shopfloor_context_incident_reviews ( episode_id, previous_values, new_review_status, review_classification, reviewed_by, notes, evidence_json ) VALUES ( %s, %s, %s, %s, %s, %s, %s ) RETURNING review_id """, ( episode["episode_id"], Json( json_safe({ "review_status": episode[ "review_status" ], "review_classification": episode[ "review_classification" ], "reviewed_by": episode[ "reviewed_by" ], "reviewed_at": episode[ "reviewed_at" ], "review_notes": episode[ "review_notes" ], }) ), new_status, episode_classification, user.display_name, review.notes, Json({ "source": "governance_review_ui_v0314", "request_id": str(request_id), "decision_group_key": decision[ "decision_group_key" ], "case_id": case_id, "answer_code": review.answer_code, }), ), ) episode_review_ids.append( int( cursor.fetchone()[ "review_id" ] ) ) return ( case_review_ids, episode_review_ids, ) def apply_review( review: ReviewInput, user: CurrentUser, source_ip: str | None, user_agent: str | None, ) -> dict[str, Any]: request_id = uuid.uuid4() settings = get_settings() with get_connection() as connection: connection.autocommit = False try: with connection.cursor( cursor_factory=RealDictCursor, ) as cursor: cursor.execute( """ SELECT pg_advisory_xact_lock( hashtext(%s) ) """, ( ( f"{settings.tenant}:" f"{settings.site}:" f"{review.decision_group_key}" ), ), ) decision = fetch_decision( cursor, review.decision_group_key, ) all_cases = fetch_cases_for_update( cursor, review.decision_group_key, ) ( pending_cases, selected_cases, mapped_classification, ) = validate_input( decision, all_cases, review, ) previous_state = { "decision": decision, "cases": all_cases, } created_context_ids = insert_contexts( cursor, decision, selected_cases, review, user, request_id, ) pending_ids = { int(case["case_id"]) for case in pending_cases } selected_ids = { int(case["case_id"]) for case in selected_cases } all_pending_selected = ( pending_ids == selected_ids ) decision_review_id = None if all_pending_selected: new_status = ( "DISMISSED" if mapped_classification == "NOT_RELEVANT" else "REVIEWED" ) cursor.execute( """ INSERT INTO mv_loss_intelligence .context_decision_group_reviews ( tenant, site, decision_group_key, previous_values, new_review_status, resolution_classification, selected_order_ref, selected_workorder_id, reviewed_by, notes, evidence_json ) VALUES ( %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s ) RETURNING review_id """, ( settings.tenant, settings.site, review.decision_group_key, Json( json_safe(decision) ), new_status, mapped_classification, review.selected_order_ref, review.selected_workorder_id, user.display_name, review.notes, Json({ "source": "governance_review_ui_v0314", "request_id": str(request_id), "answer_code": review.answer_code, "selected_case_ids": list( review.selected_case_ids ), }), ), ) decision_review_id = int( cursor.fetchone()[ "review_id" ] ) ( case_review_ids, episode_review_ids, ) = insert_case_and_episode_reviews( cursor, decision, selected_cases, mapped_classification, review, user, request_id, ) resulting_decision = fetch_decision( cursor, review.decision_group_key, ) resulting_state = { "decision": resulting_decision, "decision_review_id": decision_review_id, "case_review_ids": case_review_ids, "episode_review_ids": episode_review_ids, "created_context_session_ids": created_context_ids, } cursor.execute( """ INSERT INTO mv_loss_intelligence .context_governance_ui_actions ( tenant, site, request_id, decision_group_key, selected_case_ids, action_scope, answer_code, mapped_classification, selected_order_ref, selected_workorder_id, corrected_in_odoo, create_contexts, created_context_session_ids, actor_username, actor_display_name, actor_roles, notes, previous_state, resulting_state, status, source_ip, user_agent, official_eligible ) VALUES ( %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, 'APPLIED', %s, %s, false ) RETURNING ui_action_id """, ( settings.tenant, settings.site, str(request_id), review.decision_group_key, list(review.selected_case_ids), ( "ALL_PENDING_CASES" if all_pending_selected else "SELECTED_CASES" ), review.answer_code, mapped_classification, review.selected_order_ref, review.selected_workorder_id, review.corrected_in_odoo, review.create_contexts, created_context_ids, user.username, user.display_name, sorted(user.roles), review.notes, Json( json_safe( previous_state ) ), Json( json_safe( resulting_state ) ), source_ip, user_agent, ), ) ui_action_id = int( cursor.fetchone()[ "ui_action_id" ] ) connection.commit() return { "ui_action_id": ui_action_id, "request_id": str(request_id), "decision_review_id": decision_review_id, "case_review_ids": case_review_ids, "episode_review_ids": episode_review_ids, "created_context_session_ids": created_context_ids, "all_pending_cases_selected": all_pending_selected, } except Exception: connection.rollback() raise