import asyncio
import logging
from app.core.postgres.erp_client import ErpDatabaseClient
from app.core.module_framework.event_bus import event_bus

logger = logging.getLogger("dscons.outbox_worker")

async def process_outbox_events():
    """Poll the erp_outbox_events table and publish them to the event bus."""
    db = ErpDatabaseClient()
    
    while True:
        try:
            with db.get_connection() as conn, conn.cursor() as cur:
                # Use SELECT FOR UPDATE SKIP LOCKED for concurrent worker safety
                cur.execute(
                    """
                    SELECT id, event_type, payload
                    FROM erp_outbox_events
                    WHERE status = 'pending'
                    ORDER BY created_at ASC
                    LIMIT 10
                    FOR UPDATE SKIP LOCKED
                    """
                )
                events = cur.fetchall()
                
                for row in events:
                    event_id = row["id"]
                    event_type = row["event_type"]
                    payload = row["payload"]
                    
                    try:
                        # Mark as processing
                        cur.execute(
                            "UPDATE erp_outbox_events SET status = 'processing' WHERE id = %s",
                            (event_id,)
                        )
                        conn.commit()
                        
                        # Publish to memory subscribers (wait if it's async to ensure DLQ logic could be handled, 
                        # but in this design publish is detached. If we want DLQ, event_bus.publish could return status)
                        await event_bus.publish(event_type, payload)
                        
                        # Mark as completed
                        with db.get_connection() as conn2, conn2.cursor() as cur2:
                            cur2.execute(
                                "UPDATE erp_outbox_events SET status = 'completed', processed_at = CURRENT_TIMESTAMP WHERE id = %s",
                                (event_id,)
                            )
                            conn2.commit()
                            
                    except Exception as e:
                        logger.error(f"[OUTBOX] Failed to process event {event_id}: {e}")
                        with db.get_connection() as conn3, conn3.cursor() as cur3:
                            cur3.execute(
                                "UPDATE erp_outbox_events SET status = 'failed', error_message = %s WHERE id = %s",
                                (str(e), event_id)
                            )
                            conn3.commit()
                            
            # Sleep if no events or just yield
            if not events:
                await asyncio.sleep(2.0)
            else:
                await asyncio.sleep(0.1)
                
        except Exception as e:
            logger.error(f"[OUTBOX] Worker error: {e}")
            await asyncio.sleep(5.0)

async def start_outbox_worker_loop():
    logger.info("Starting Transactional Outbox Worker...")
    asyncio.create_task(process_outbox_events())
