a
    Ôi–  ã                   @   sL  U d Z ddlZddlZddlmZmZmZmZ ddlZddl	m	Z	 e 
e¡Zdaeej ed< e ¡ Zd!eeejdœdd	„Zdd
œdd„Zejd
œdd„Zd"eee eedœdd„Zeee edœdd„Zd#eee ee dœdd„Zd$eee ee dœdd„Zd%eee ee dœdd„Zd&eee ee dœdd„ZG dd „ d ƒZ dS )'zs
Async database layer using aiomysql.

Provides connection pooling, query helpers, and proper resource management.
é    N)ÚOptionalÚAnyÚTupleÚList)ÚconfigÚ_poolé
   )Ú	pool_sizeÚmax_overflowÚreturnc                 Ã   s  t 4 I dH šè tdur8t d¡ tW  d  ƒI dH  S | p@tj} zdtjtjtj	tj
tjtjd| ddtjdd�I dH at d| › d	tj› �¡ tW W  d  ƒI dH  S  tyÜ } zt d
|› �¡ ‚ W Y d}~n
d}~0 0 W d  ƒI dH  �q1 I dH �s0    Y  dS )zö
    Initialize the async database connection pool.

    Args:
        pool_size: Maximum pool size (default from config.MAX_CONNECTIONS)
        max_overflow: Additional connections beyond pool_size

    Returns:
        The connection pool
    NzPool already initializedé   Ti  r   )ÚhostÚportÚuserÚpasswordÚdbZminsizeÚmaxsizeÚ
autocommitZpool_recycleZechoÚconnect_timeoutz Database pool initialized: size=z, host=z$Failed to initialize database pool: )Ú
_pool_lockr   ÚloggerÚwarningr   ÚMAX_CONNECTIONSÚaiomysqlZcreate_poolÚDB_HOSTÚDB_PORTÚDB_USERÚDB_PASSWORDÚDB_NAMEÚDEBUGÚinfoÚ	ExceptionÚerror)r	   r
   Úe© r$   ú//var/www/lichun.app/lichun/ws/database_async.pyÚinitialize_pool   s0    

õr&   ©r   c                	   Ã   sh   t 4 I dH šB tdur:t ¡  t ¡ I dH  dat d¡ W d  ƒI dH  qd1 I dH sZ0    Y  dS )z"Close the database connection poolNzDatabase pool closed)r   r   ÚcloseÚwait_closedr   r    r$   r$   r$   r%   Ú
close_poolB   s    r*   c                   Ã   s   t du rtdƒ‚t  ¡ S )a  
    Get a connection from the pool.

    Usage:
        async with get_connection() as conn:
            async with conn.cursor() as cursor:
                await cursor.execute("SELECT ...")

    Returns:
        Database connection (context manager)
    Nz<Database pool not initialized. Call initialize_pool() first.)r   ÚRuntimeErrorÚacquirer$   r$   r$   r%   Úget_connectionN   s    r-   T)ÚqueryÚparamsÚcommitr   c              
   Ã   s¸   t du rtdƒ‚t  ¡ 4 I dH š~}| ¡ 4 I dH šB}| | |¡I dH  |jW  d  ƒI dH  W  d  ƒI dH  S 1 I dH s€0    Y  W d  ƒI dH  q´1 I dH sª0    Y  dS )zþ
    Execute a query (INSERT, UPDATE, DELETE).

    Args:
        query: SQL query with %s placeholders
        params: Query parameters
        commit: Whether to commit (default True due to autocommit)

    Returns:
        Number of affected rows
    NúDatabase pool not initialized)r   r+   r,   ÚcursorÚexecuteÚrowcount)r.   r/   r0   Úconnr2   r$   r$   r%   Úexecute_query`   s    r6   )r.   Úparams_listr   c              
   Ã   s¸   t du rtdƒ‚t  ¡ 4 I dH š~}| ¡ 4 I dH šB}| | |¡I dH  |jW  d  ƒI dH  W  d  ƒI dH  S 1 I dH s€0    Y  W d  ƒI dH  q´1 I dH sª0    Y  dS )zí
    Execute a query multiple times with different parameters (batch insert).

    Args:
        query: SQL query with %s placeholders
        params_list: List of parameter tuples

    Returns:
        Total number of affected rows
    Nr1   )r   r+   r,   r2   Úexecutemanyr4   )r.   r7   r5   r2   r$   r$   r%   Úexecute_manyy   s    r9   )r.   r/   r   c              
   Ã   sÀ   t du rtdƒ‚t  ¡ 4 I dH š†}| ¡ 4 I dH šJ}| | |¡I dH  | ¡ I dH W  d  ƒI dH  W  d  ƒI dH  S 1 I dH sˆ0    Y  W d  ƒI dH  q¼1 I dH s²0    Y  dS )z·
    Fetch a single row.

    Args:
        query: SQL query with %s placeholders
        params: Query parameters

    Returns:
        Single row as tuple, or None if not found
    Nr1   )r   r+   r,   r2   r3   Úfetchone©r.   r/   r5   r2   r$   r$   r%   Ú	fetch_one�   s    r<   c              
   Ã   sÀ   t du rtdƒ‚t  ¡ 4 I dH š†}| ¡ 4 I dH šJ}| | |¡I dH  | ¡ I dH W  d  ƒI dH  W  d  ƒI dH  S 1 I dH sˆ0    Y  W d  ƒI dH  q¼1 I dH s²0    Y  dS )z 
    Fetch all rows.

    Args:
        query: SQL query with %s placeholders
        params: Query parameters

    Returns:
        List of rows as tuples
    Nr1   )r   r+   r,   r2   r3   Úfetchallr;   r$   r$   r%   Ú	fetch_all§   s    r>   c              
   Ã   sÄ   t du rtdƒ‚t  ¡ 4 I dH šŠ}| tj¡4 I dH šJ}| | |¡I dH  | ¡ I dH W  d  ƒI dH  W  d  ƒI dH  S 1 I dH sŒ0    Y  W d  ƒI dH  qÀ1 I dH s¶0    Y  dS )zÄ
    Fetch a single row as dictionary.

    Args:
        query: SQL query with %s placeholders
        params: Query parameters

    Returns:
        Single row as dict, or None if not found
    Nr1   )r   r+   r,   r2   r   Ú
DictCursorr3   r:   r;   r$   r$   r%   Úfetch_dict_one¾   s    r@   c              
   Ã   sÄ   t du rtdƒ‚t  ¡ 4 I dH šŠ}| tj¡4 I dH šJ}| | |¡I dH  | ¡ I dH W  d  ƒI dH  W  d  ƒI dH  S 1 I dH sŒ0    Y  W d  ƒI dH  qÀ1 I dH s¶0    Y  dS )z¯
    Fetch all rows as dictionaries.

    Args:
        query: SQL query with %s placeholders
        params: Query parameters

    Returns:
        List of rows as dicts
    Nr1   )r   r+   r,   r2   r   r?   r3   r=   r;   r$   r$   r%   Úfetch_dict_allÕ   s    rA   c                   @   s0   e Zd ZdZdd„ Zejdœdd„Zdd„ Zd	S )
ÚTransactiona<  
    Context manager for database transactions.

    Usage:
        async with Transaction() as conn:
            async with conn.cursor() as cursor:
                await cursor.execute("UPDATE ...")
                await cursor.execute("INSERT ...")
            # Auto-commit on success, rollback on exception
    c                 C   s
   d | _ d S )N)r5   ©Úselfr$   r$   r%   Ú__init__ø   s    zTransaction.__init__r'   c                 Ã   s6   t d u rtdƒ‚t  ¡ I d H | _| j ¡ I d H  | jS )Nr1   )r   r+   r,   r5   ÚbeginrC   r$   r$   r%   Ú
__aenter__û   s
    zTransaction.__aenter__c                 Ã   sX   |d ur,| j  ¡ I d H  t d|j› �¡ n| j  ¡ I d H  t | j ¡I d H  d | _ d S )NzTransaction rolled back due to )r5   Úrollbackr   r   Ú__name__r0   r   Úrelease)rD   Úexc_typeÚexc_valÚexc_tbr$   r$   r%   Ú	__aexit__  s    zTransaction.__aexit__N)	rI   Ú
__module__Ú__qualname__Ú__doc__rE   r   Ú
ConnectionrG   rN   r$   r$   r$   r%   rB   ì   s   rB   )Nr   )NT)N)N)N)N)!rQ   Úasyncior   Útypingr   r   r   r   Úloggingr   Ú	getLoggerrI   r   r   ZPoolÚ__annotations__ÚLockr   Úintr&   r*   rR   r-   ÚstrÚboolr6   r9   r<   r>   Údictr@   rA   rB   r$   r$   r$   r%   Ú<module>   sl   
  þý.  ýüý þý þý þý þý