U
    ÌZjÚ"  ã                   @   s€   d Z ddlmZ ddlmZ ddlmZ ddlmZ ddlmZ ddlm	Z	 dd	gZ
G d
d	„ d	eƒZG dd„ de	ƒZdd„ ZdS )a  Horizontal sharding support.

Defines a rudimental 'horizontal sharding' system which allows a Session to
distribute queries and persistence operations across multiple databases.

For a usage example, see the :ref:`examples_sharding` example included in
the source distribution.

é   )Úevent)Úexc)Úinspect)Úutil)ÚQuery)ÚSessionÚShardedSessionÚShardedQueryc                       s$   e Zd Z‡ fdd„Zdd„ Z‡  ZS )r	   c                    s:   t t| ƒj||Ž | jj| _| jj| _| jj| _d | _d S ©N)Úsuperr	   Ú__init__ÚsessionÚ
id_chooserÚquery_chooserÚexecute_chooserZ	_shard_id)ÚselfÚargsÚkwargs©Ú	__class__© úb/var/www/html/TRUCKING_PROJECT/venv/lib/python3.8/site-packages/sqlalchemy/ext/horizontal_shard.pyr      s
    


zShardedQuery.__init__c                 C   s   | j |d�S )aÃ  Return a new query, limited to a single shard ID.

        All subsequent operations with the returned query will
        be against the single shard regardless of other state.

        The shard_id can be passed for a 2.0 style execution to the
        bind_arguments dictionary of :meth:`.Session.execute`::

            results = session.execute(
                stmt,
                bind_arguments={"shard_id": "my_shard"}
            )

        )Ú_sa_shard_id)Úexecution_options)r   Úshard_idr   r   r   Ú	set_shard$   s    zShardedQuery.set_shard)Ú__name__Ú
__module__Ú__qualname__r   r   Ú__classcell__r   r   r   r   r	      s   c                       sV   e Zd Zddef‡ fdd„	Zd‡ fdd„	Zdd„ Zddd	„Zdd
d„Zdd„ Z	‡  Z
S )r   Nc                    s®   |  dd¡‰ tt| ƒjf d|i|—Ž tj| dtdd� || _|| _ˆ rvt	 
dd¡ |rbt d	¡‚‡ fd
d„}|| _n|| _ˆ | _i | _|dk	rª|D ]}|  ||| ¡ q”dS )a¥  Construct a ShardedSession.

        :param shard_chooser: A callable which, passed a Mapper, a mapped
          instance, and possibly a SQL clause, returns a shard ID.  This id
          may be based off of the attributes present within the object, or on
          some round-robin scheme. If the scheme is based on a selection, it
          should set whatever state on the instance to mark it in the future as
          participating in that shard.

        :param id_chooser: A callable, passed a query and a tuple of identity
          values, which should return a list of shard ids where the ID might
          reside.  The databases will be queried in the order of this listing.

        :param execute_chooser: For a given :class:`.ORMExecuteState`,
          returns the list of shard_ids
          where the query should be issued.  Results from all shards returned
          will be combined together into a single listing.

          .. versionchanged:: 1.4  The ``execute_chooser`` parameter
             supersedes the ``query_chooser`` parameter.

        :param shards: A dictionary of string shard names
          to :class:`~sqlalchemy.engine.Engine` objects.

        r   NÚ	query_clsZdo_orm_executeT)ÚretvalzMThe ``query_choser`` parameter is deprecated; please use ``execute_chooser``.z1.4z>Can't pass query_chooser and execute_chooser at the same time.c                    s
   ˆ | j ƒS r
   )Z	statement©Úorm_context©r   r   r   r   n   s    z0ShardedSession.__init__.<locals>.execute_chooser)Úpopr   r   r   r   ÚlistenÚexecute_and_instancesÚshard_chooserr   r   Zwarn_deprecatedr   ÚArgumentErrorr   r   Ú_ShardedSession__bindsÚ
bind_shard)r   r(   r   r   Zshardsr    r   Úkr   r$   r   r   7   s6    "   ÿýÿzShardedSession.__init__c           	         sˆ   |dk	r&t t| ƒj||fd|i|—ŽS |  |¡}|r>| |¡}|  ||¡D ]4}t t| ƒj||f||dœ|—Ž}|dk	rJ|  S qJdS dS )a_  override the default :meth:`.Session._identity_lookup` method so
        that we search for a given non-token primary key identity across all
        possible identity tokens (e.g. shard ids).

        .. versionchanged:: 1.4  Moved :meth:`.Session._identity_lookup` from
           the :class:`_query.Query` object to the :class:`.Session`.

        NÚidentity_token)r-   Úlazy_loaded_from)r   r   Ú_identity_lookupÚqueryZ_set_lazyload_fromr   )	r   ÚmapperZprimary_key_identityr-   r.   ÚkwÚqr   Úobjr   r   r   r/   z   s2    
þýü


þüû
zShardedSession._identity_lookupc                 K   s^   |d k	r<t |ƒ}|jr0|jd }|d k	s,t‚|S |jr<|jS | j||f|Ž}|d k	rZ||_|S )Nr   )r   ÚkeyÚAssertionErrorr-   r(   )r   r1   Úinstancer2   ÚstateÚtokenr   r   r   r   Ú_choose_shard_and_assign£   s    
z'ShardedSession._choose_shard_and_assignc                 K   sJ   |dkr|   ||¡}|  ¡ r.|  ¡ j||d�S | j|||d�jf |ŽS dS )zaProvide a :class:`_engine.Connection` to use in the unit of work
        flush process.

        N)r   )r   r7   )r:   Zin_transactionZget_transactionÚ
connectionÚget_bindÚconnect)r   r1   r7   r   r   r   r   r   Úconnection_callable²   s      ÿþz"ShardedSession.connection_callablec                 K   s"   |d kr| j |||d�}| j| S )N)Úclause)r:   r*   )r   r1   r   r7   r?   r2   r   r   r   r<   Ä   s      ÿzShardedSession.get_bindc                 C   s   || j |< d S r
   )r*   )r   r   Úbindr   r   r   r+   Í   s    zShardedSession.bind_shard)NN)NNN)NNNN)r   r   r   r	   r   r/   r:   r>   r<   r+   r   r   r   r   r   r   6   s$   úG  û)     ÿ
       ÿ
	c           	         sî   ˆ j rˆ j }}d }n(ˆ js"ˆ jr2d }ˆ j }}nd  } }}ˆ j}‡ fdd„}|rf|jd k	rf|j}n0dˆ jkr|ˆ jd }ndˆ jkr’ˆ jd }nd }|d k	rª||||ƒS g }| 	ˆ ¡D ]}||||ƒ}| 
|¡ q¸|d j|dd … Ž S d S )Nc                    sf   t ˆ jƒ}t ˆ jƒ}| |d< ˆ jr8|d| i7 }||d< n ˆ jsDˆ jrX|d| i7 }||d< ˆ j||d�S )Nr   Ú_refresh_identity_tokenZ_sa_orm_load_optionsZ_sa_orm_update_options)Úbind_argumentsr   )ÚdictZlocal_execution_optionsrB   Ú	is_selectÚ	is_updateÚ	is_deleteZinvoke_statement)r   Úload_optionsÚupdate_optionsr   rB   r"   r   r   Úiter_for_shardÞ   s    


 ÿz-execute_and_instances.<locals>.iter_for_shardr   r   é    é   )rD   rG   rE   rF   Zupdate_delete_optionsr   rA   r   rB   r   ÚappendÚmerge)	r#   rG   Zactive_optionsrH   r   rI   r   ÚpartialZresult_r   r"   r   r'   Ñ   s.    


r'   N)Ú__doc__Ú r   r   r   r   Z	orm.queryr   Zorm.sessionr   Ú__all__r	   r   r'   r   r   r   r   Ú<module>   s   
 