o
    j+&                     @   s   d dl Z d dlZd dlZd dlZd dlZd dlmZ d dlmZm	Z	 d dl
mZ d dlmZ G dd dZG dd	 d	Z	 d dlZd dlZG d
d dejZG dd dZG dd deZG dd deZdS )    N)	UniTarget)
PacketizerStreamPacketizer)PacketizerSSL)UNITransportc                   @   s   e Zd Zd*dejdejdededef
ddZ	d	d
 Z
dd Zd+ddZdd Zdd Zdd Zd*ddZdd Zdd Zdd Zdd Zdd  Zd!d" Zd#d$ Zd%ejfd&d'Zd%ejfd(d)ZdS ),UniConnectionNreaderwriter
packetizerpeer_ip	peer_portc                 C   sV   || _ || _|| _|| _|| _d | _d| _t | _	t
 | _t | _| j  d S NF)r   r	   r
   r   r   packetizer_taskclosingasyncioEvent
closed_evtLock	read_lockread_resumeset)selfr   r	   r
   r   r    r   /root/aizidognhua/tmp/workspace/projects/ec89d86c-575f-41c9-af57-ac45cbdbf775/venv/lib/python3.10/site-packages/asysocks/unicomm/common/connection.py__init__   s   


zUniConnection.__init__c                    s   | S Nr   r   r   r   r   
__aenter__      zUniConnection.__aenter__c                    s   |   I d H  d S r   )close)r   exc_typeexctbr   r   r   	__aexit__   s   zUniConnection.__aexit__c                 C   s<   t | jdr| j||S |dkr| jd ur| j| jfS |S )Nget_extra_infopeername)hasattrr	   r$   r   r   )r   namedefaultr   r   r   r$   !   s
   zUniConnection.get_extra_infoc                 C   s
   | j  S r   )r
   get_peer_certificater   r   r   r   r)   *      
z"UniConnection.get_peer_certificatec                 C   s,   | j  }t| j tr|| j _ d S || _ d S r   )r
   flush_buffer
isinstancer   )r   r
   rem_datar   r   r   change_packetizer-   s   

zUniConnection.change_packetizerc                 O   s   | j j|i |S r   )r
   packetizer_control)r   argskwr   r   r   r/   4      z UniConnection.packetizer_controlc                    sV   |d u r| j }|d u rt }d|_tj|_t||| _ | j | j| j	I d H  d S r   )
r
   sslcreate_default_contextcheck_hostname	CERT_NONEverify_moder   do_handshaker   r	   )r   ssl_ctxr
   r   r   r   wrap_ssl7   s   zUniConnection.wrap_sslc                    s*   d| _ | jd ur| j  | j  d S NT)r   r	   r   r   r   r   r   r   r   r   A   s
   

zUniConnection.closec                    s   d S r   r   r   r   r   r   drainG   r   zUniConnection.drainc                    s>   | j |2 z3 d H W }| j| | j I d H  q6 d S r   )r
   data_outr	   writer<   )r   datapacketr   r   r   r>   J   s
   zUniConnection.writec                    s$   |   2 z	3 d H W }|  S 6 d S r   )read)r   r@   r   r   r   read_oneO   s   zUniConnection.read_onec              
   C  s   zQd }| j du r5| j|2 z3 d H W }|d u r n|V  q6 | j| jjI d H }|dkr0n| j du s	d }| j|2 z3 d H W }|d u rK W d S |V  q=6 W d S  tyh } z
d V  W Y d }~d S d }~ww )NF    )r   r
   data_inr   rA   buffer_size	Exception)r   r?   resulter   r   r   rA   S   s.   

zUniConnection.readc                    sN   t | jtstd	 | j| jjI d H }| j|I d H  |dkr&d S q)Nz<This function onaly available when StreamPacketizer is used!TrC   )r,   r
   r   rF   r   rA   rE   rD   r   r?   r   r   r   streami   s   zUniConnection.streamc              	      s\   | j   | j4 I d H  | j  I d H  W d   I d H  d S 1 I d H s'w   Y  d S r   )r   clearr   waitr   r   r   r   pause_readingt   s
   
.zUniConnection.pause_readingprotocolc              
      s"  d }zz^d }	 | j 4 I d H F | j|2 z3 d H W }|d u r&|   n|| q6 | j| jjI d H }|dkrK|  	 W d   I d H  nW d   I d H  n1 I d H s[w   Y  qW n tyz } zt	
  |}W Y d }~nd }~ww W || d S W || d S || w )NTrC   )r   r
   rD   eof_receiveddata_receivedr   rA   rE   rF   	traceback	print_excconnection_lost)r   rN   errr?   rG   rH   r   r   r   __transport_readery   s<   (z UniConnection.__transport_readerc                    s*   t | |}|| t| |}|S r   )r   connection_mader   create_task _UniConnection__transport_reader)r   rN   	transportxr   r   r   get_transport   s
   

zUniConnection.get_transport)NNr   )__name__
__module____qualname__r   StreamReaderStreamWriterr   strintr   r   r#   r$   r)   r.   r/   r:   r   r<   r>   rB   rA   rJ   rM   ProtocolrX   r[   r   r   r   r   r      s$    $
	

r   c                   @   s   e Zd Zdd ZdS )UniUDPConnectionc                 C   s   || _ || _|| _d S r   )socketr?   addr)r   re   r?   rf   r   r   r   r      s   
zUniUDPConnection.__init__N)r\   r]   r^   r   r   r   r   r   rd      s    rd   c                   @   sH   e Zd ZdZdd Zdd Zdd Zdd	 Zd
d Zdd Z	dd Z
dS )DatagramEndpointProtocolz8Datagram protocol for the endpoint high-level interface.c                 C   s
   || _ d S r   )	_endpoint)r   endpointr   r   r   r      r*   z!DatagramEndpointProtocol.__init__c                 C   s   || j _d S r   )rh   
_transport)r   rY   r   r   r   rV      s   z(DatagramEndpointProtocol.connection_madec                 C   s4   |d u sJ | j jd ur| j jd  | j   d S r   )rh   _write_ready_future
set_resultr   )r   r!   r   r   r   rS      s   z(DatagramEndpointProtocol.connection_lostc                 C   s   | j || d S r   )rh   feed_datagramr   r?   rf   r   r   r   datagram_received   r2   z*DatagramEndpointProtocol.datagram_receivedc                 C   s   d}t || d S )Nz Endpoint received an error: {!r})warningswarnformat)r   r!   msgr   r   r   error_received   s   z'DatagramEndpointProtocol.error_receivedc                 C   s*   | j jd u sJ | j jj}| | j _d S r   )rh   rk   rj   _loopcreate_future)r   loopr   r   r   pause_writing   s   
z&DatagramEndpointProtocol.pause_writingc                 C   s*   | j jd usJ | j jd  d | j _d S r   )rh   rk   rl   r   r   r   r   resume_writing   s   z'DatagramEndpointProtocol.resume_writingN)r\   r]   r^   __doc__r   rV   rS   ro   rt   rx   ry   r   r   r   r   rg      s    rg   c                   @   s   e Zd ZdZddefddZdd Zdd	 Zdd
dZdddZ	dddZ
dd Zdd Zdd Zdd Zedd Zedd ZdS )EndpointzHigh-level interface for UDP enpoints.
	Can either be local or remote.
	It is initialized with an optional queue size for the incoming datagrams.
	Ntargetc                 C   s4   || _ |d u r	d}t|| _d| _d | _d | _d S )Nr   F)r|   r   Queue_queue_closedrj   rk   )r   r|   
queue_sizer   r   r   r      s   
zEndpoint.__init__c                 C   s8   z| j ||f W d S  tjy   td Y d S w )NzEndpoint queue is full)r~   
put_nowaitr   	QueueFullrp   rq   rn   r   r   r   rm      s
   zEndpoint.feed_datagramc                 C   s>   | j rd S d| _ | j r| d d  | jr| j  d S d S r;   )r   r~   emptyrm   rj   r   r   r   r   r   r      s   
zEndpoint.closec                    s(   |du r| j  | j jf}| ||S )%Send a datagram to the given address.N)r|   get_ip_or_hostnameportsendrn   r   r   r   r>   	  s   zEndpoint.writeFc                 C  sZ   | j s+|  I dH \}}|du r|   dS ||f}|du r#|V  n|V  | j rdS dS nWait for an incoming datagram and return it with
		the corresponding address.
		This method is a coroutine.
		NT)r   receiver   )r   	with_addrr?   rf   rG   r   r   r   rA     s   zEndpoint.readc                    s<   |   I dH \}}|du r|   dS |du r||fS |S r   )r   r   )r   r   r?   rf   r   r   r   rB      s   zEndpoint.read_onec                 C   sJ   | j rtdt| jjdu rt| j||}dS | j|| dS )r   Enpoint is closedTN)r   IOErrorr   iscoroutinerj   sendtorW   )r   r?   rf   rZ   r   r   r   r   -  s
   zEndpoint.sendc                    sF   | j  r| jrtd| j  I dH \}}|du rtd||fS )r   r   N)r~   r   r   r   getrn   r   r   r   r   7  s   zEndpoint.receivec                 C   s$   | j rtd| j  |   dS )z Close the transport immediately.r   N)r   r   rj   abortr   r   r   r   r   r   C  s   
zEndpoint.abortc                    s    | j dur| j I dH  dS dS )z4Drain the transport buffer below the low-water mark.N)rk   r   r   r   r   r<   J  s   
zEndpoint.drainc                 C   s   | j d S )z-The endpoint address as a (host, port) tuple.re   )rj   r$   getsocknamer   r   r   r   addressQ  s   zEndpoint.addressc                 C   s   | j S )z0Indicates whether the endpoint is closed or not.)r   r   r   r   r   closedV  s   zEndpoint.closedr   )F)r\   r]   r^   rz   r   r   rm   r   r>   rA   rB   r   r   r   r<   propertyr   r   r   r   r   r   r{      s     




r{   c                   @   s   e Zd ZdZdS )LocalEndpointzyHigh-level interface for UDP local enpoints.
	It is initialized with an optional queue size for the incoming datagrams.
	N)r\   r]   r^   rz   r   r   r   r   r   \  s    r   c                       s,   e Zd ZdZ fddZ fddZ  ZS )RemoteEndpointzzHigh-level interface for UDP remote enpoints.
	It is initialized with an optional queue size for the incoming datagrams.
	c                    s   t  |d dS )z#Send a datagram to the remote host.N)superr   rI   	__class__r   r   r   h  s   zRemoteEndpoint.sendc                    s   t   I dH \}}|S )zU Wait for an incoming datagram from the remote host.
		This method is a coroutine.
		N)r   r   rn   r   r   r   r   l  s   zRemoteEndpoint.receive)r\   r]   r^   rz   r   r   __classcell__r   r   r   r   r   c  s    r   )r3   copyr   rQ   	ipaddressasysocks.unicomm.common.targetr   #asysocks.unicomm.common.packetizersr   r   'asysocks.unicomm.common.packetizers.sslr   !asysocks.unicomm.common.transportr   r   rd   rp   DatagramProtocolrg   r{   r   r   r   r   r   r   <module>   s&     
)x