Ë
    UV.jm/  ã                  óŽ  — 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mZ 	 ddlZdd	lmZmZmZmZ dd
lmZmZ dZdZddlmZ  ee«      ZdZ G d„ de«      Z G d„ dej>                  «      Z G d„ dej@                  «      Z  G d„ dejB                  «      Z! G d„ dejD                  «      Z"y# e$ r	 dZdxZZY Œ}w xY w)a	  confluent-kafka transport module for Kombu.

Kafka transport using confluent-kafka library.

**References**

- http://docs.confluent.io/current/clients/confluent-kafka-python

**Limitations**

The confluent-kafka transport does not support PyPy environment.

Features
========
* Type: Virtual
* Supports Direct: Yes
* Supports Topic: Yes
* Supports Fanout: No
* Supports Priority: No
* Supports TTL: No

Connection String
=================
Connection string has the following format:

.. code-block::

    confluentkafka://[USER:PASSWORD@]KAFKA_ADDRESS[:PORT]

Transport Options
=================
* ``connection_wait_time_seconds`` - Time in seconds to wait for connection
  to succeed. Default ``5``
* ``wait_time_seconds`` - Time in seconds to wait to receive messages.
  Default ``5``
* ``security_protocol`` - Protocol used to communicate with broker.
  Visit https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md for
  an explanation of valid values. Default ``plaintext``
* ``sasl_mechanism`` - SASL mechanism to use for authentication.
  Visit https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md for
  an explanation of valid values.
* ``num_partitions`` - Number of partitions to create. Default ``1``
* ``replication_factor`` - Replication factor of partitions. Default ``1``
* ``topic_config`` - Topic configuration. Must be a dict whose key-value pairs
  correspond with attributes in the
  http://kafka.apache.org/documentation.html#topicconfigs.
* ``kafka_common_config`` - Configuration applied to producer, consumer and
  admin client. Must be a dict whose key-value pairs correspond with attributes
  in the https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md.
* ``kafka_producer_config`` - Producer configuration. Must be a dict whose
  key-value pairs correspond with attributes in the
  https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md.
* ``kafka_consumer_config`` - Consumer configuration. Must be a dict whose
  key-value pairs correspond with attributes in the
  https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md.
* ``kafka_admin_config`` - Admin client configuration. Must be a dict whose
  key-value pairs correspond with attributes in the
  https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md.
é    )Úannotations)ÚEmpty)Úvirtual)Úcached_property)Ústr_to_bytes)ÚdumpsÚloadsN)ÚConsumerÚKafkaExceptionÚProducerÚTopicPartition)ÚAdminClientÚNewTopic© )Ú
get_loggeri„#  c                  ó   — e Zd ZdZdZy)ÚNoBrokersAvailablez(Kafka broker is not available exception.TN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú	retriabler   ó    ú`C:\xampp\htdocs\tradingbinance\backend\.venv\Lib\site-packages\kombu/transport/confluentkafka.pyr   r   Z   s
   „ Ù2à�Ir   r   c                  ó$   ‡ — e Zd ZdZdˆ fd„	Zˆ xZS )ÚMessagezMessage object.c                óV   •— |j                  d«      | _        t        ‰| �  |fd|i|¤Ž y )NÚtopicÚchannel)Úgetr   ÚsuperÚ__init__)ÚselfÚpayloadr   ÚkwargsÚ	__class__s       €r   r"   zMessage.__init__c   s*   ø€ Ø—[‘[ Ó)ˆŒ
Ü‰Ñ˜Ñ<¨'Ð<°VÓ<r   ©N)r   r   r   r   r"   Ú__classcell__©r&   s   @r   r   r   `   s   ø„ Ù÷=ñ =r   r   c                  óB   — e Zd ZdZi Zd„ Zd„ Zd„ Zd„ Zd„ Z	d
d„Z
dd	„Zy)ÚQoSzQuality of Service guarantees.c                ód   — | j                    xs" t        | j                  «      | j                   k  S )z�Return true if the channel can be consumed from.

        :returns: True, if this QoS object can accept a message.
        :rtype: bool
        ©Úprefetch_countÚlenÚ_not_yet_acked©r#   s    r   Úcan_consumezQoS.can_consumem   s3   € ð ×&Ñ&Ð&ò ¬#¨d×.AÑ.AÓ*BÀTß‰^ñ+ð 	r   c                ó`   — | j                   r"| j                   t        | j                  «      z
  S y)Né   r-   r1   s    r   Úcan_consume_max_estimatezQoS.can_consume_max_estimatev   s*   € Ø×ÒØ×&Ñ&¬¨T×-@Ñ-@Ó)AÑAÐAàr   c                ó"   — || j                   |<   y r'   ©r0   )r#   ÚmessageÚdelivery_tags      r   Úappendz
QoS.append|   s   € Ø,3ˆ×Ñ˜LÒ)r   c                ó    — | j                   |   S r'   r7   )r#   r9   s     r   r    zQoS.get   s   € Ø×"Ñ" <Ñ0Ð0r   c                óÂ   — || j                   vry | j                   j                  |«      }| j                  j                  |j                  «      }|j                  «        y r'   )r0   Úpopr   Ú_get_consumerr   Úcommit)r#   r9   r8   Úconsumers       r   ÚackzQoS.ack‚   sK   € Ø˜t×2Ñ2Ñ2ØØ×%Ñ%×)Ñ)¨,Ó7ˆØ—<‘<×-Ñ-¨g¯m©mÓ<ˆØ�‰Õr   c                ó`  — |r›| j                   j                  |«      }| j                  j                  |j                  «      }|j                  «       D ]G  }t        |j                  |j                  «      }|j                  |g«      \  }|j                  |«       ŒI y| j                  |«       y)zÝReject a message by delivery tag.

        If requeue is True, then the last consumed message is reverted so
        it'll be refetched on the next attempt.
        If False, that message is consumed and ignored.
        N)r0   r=   r   r>   r   Ú
assignmentr   Ú	partitionÚ	committedÚseekrA   )r#   r9   Úrequeuer8   r@   rC   Útopic_partitionÚcommitted_offsets           r   Úrejectz
QoS.reject‰   s”   € ñ Ø×)Ñ)×-Ñ-¨lÓ;ˆGØ—|‘|×1Ñ1°'·-±-Ó@ˆHØ&×1Ñ1Ö3�
Ü"0°·±Ø1;×1EÑ1Eó#G�à%-×%7Ñ%7¸Ð8IÓ%JÑ"Ð!Ø—‘Ð.Õ/ñ	 4ð �H‰H�\Õ"r   Nc                 ó   — y r'   r   )r#   Ústderrs     r   Úrestore_unacked_oncezQoS.restore_unacked_once›   ó   € Ør   )Fr'   )r   r   r   r   r0   r2   r5   r:   r    rA   rJ   rM   r   r   r   r+   r+   h   s-   „ Ù(à€Nòòò4ò1òó#ô$r   r+   c                  óÜ   ‡ — e Zd ZdZeZeZdZdZdZˆ fd„Z	d„ Z
d„ Zd„ Zd„ Zd	„ Zd
„ Zd„ Zd„ Zd„ Zd„ Zed„ «       Zed„ «       Zed„ «       Zed„ «       Zed„ «       Zed„ «       Zˆ fd„Zˆ xZS )ÚChannelzKafka Channel.é   Nc                ój   •— t        ‰| �  |i |¤Ž i | _        i | _        | j	                  «       | _        y r'   )r!   r"   Ú_kafka_consumersÚ_kafka_producersÚ_openÚ_client)r#   Úargsr%   r&   s      €r   r"   zChannel.__init__©   s2   ø€ Ü‰Ñ˜$Ð) &Ò)à "ˆÔØ "ˆÔà—z‘z“|ˆ�r   c                ó8   — t        |«      j                  dd«      S )z>Need to sanitize the name, celery sometimes pushes in @ signs.Ú@Ú )ÚstrÚreplace)r#   Úqueues     r   Úsanitize_queue_namezChannel.sanitize_queue_name±   s   € ä�5‹z×!Ñ! # rÓ*Ð*r   c                óî   — | j                  |«      }| j                  j                  |d«      }|€Et        i | j                  ¥| j
                  j                  d«      xs i ¥«      }|| j                  |<   |S )z9Create/get a producer instance for the given topic/queue.NÚkafka_producer_config)r^   rT   r    r   Úcommon_configÚoptions)r#   r]   Úproducers      r   Ú_get_producerzChannel._get_producerµ   s�   € à×(Ñ(¨Ó/ˆØ×(Ñ(×,Ñ,¨U°DÓ9ˆØÐÜð !Ø×$Ñ$ð!à—<‘<×#Ñ#Ð$;Ó<ÒBÀð!ó ˆHð ,4ˆD×!Ñ! %Ñ(àˆr   c                ó   — | j                  |«      }| j                  j                  |d«      }|€^t        |› d�dddœ| j                  ¥| j
                  j                  d«      xs i ¥«      }|j                  |g«       || j                  |<   |S )z9Create/get a consumer instance for the given topic/queue.Nz-consumer-groupÚearliestF)zgroup.idzauto.offset.resetzenable.auto.commitÚkafka_consumer_config)r^   rS   r    r
   ra   rb   Ú	subscribe)r#   r]   r@   s      r   r>   zChannel._get_consumerÂ   s¥   € à×(Ñ(¨Ó/ˆØ×(Ñ(×,Ñ,¨U°DÓ9ˆØÐÜØ$˜g _Ð5Ø%/Ø&+ñ!ð ×$Ñ$ð	!ð
 —<‘<×#Ñ#Ð$;Ó<ÒBÀð!ó ˆHð ×Ñ ˜wÔ'Ø+3ˆD×!Ñ! %Ñ(àˆr   c                ó°   — | j                  |«      }| j                  |«      }|j                  |t        t	        |«      «      «       |j                  «        y)z!Put a message on the topic/queue.N)r^   rd   Úproducer   r   Úflush)r#   r]   r8   r%   rc   s        r   Ú_putzChannel._putÓ   sE   € à×(Ñ(¨Ó/ˆØ×%Ñ% eÓ,ˆØ×Ñ˜¤¬U°7«^Ó <Ô=Ø�‰Õr   c                ót  — | j                  |«      }| j                  |«      }d}	 |j                  | j                  «      }|s
t        «       ‚|j                  «       }|rt        j                  |«       t        «       ‚i t        |j                  «       «      ¥d|j                  «       i¥S # t        $ r Y Œuw xY w)z#Get a message from the topic/queue.Nr   )r^   r>   ÚpollÚwait_time_secondsÚStopIterationr   ÚerrorÚloggerr	   Úvaluer   )r#   r]   r%   r@   r8   rq   s         r   Ú_getzChannel._getÚ   s¤   € à×(Ñ(¨Ó/ˆØ×%Ñ% eÓ,ˆØˆð	Ø—m‘m D×$:Ñ$:Ó;ˆGñ Ü“'ˆMà—‘“ˆÙÜ�L‰L˜ÔÜ“'ˆMàC”%˜Ÿ™›Ó(ÐC¨'°7·=±=³?ÑCÐCøô ò 	Ùð	ús   ¦B+ Â+	B7Â6B7c                óÎ   — | j                  |«      }| j                  |   j                  «        | j                  j                  |«       | j                  j                  |g«       y)zDelete a queue/topic.N)r^   rS   Úcloser=   ÚclientÚdelete_topics)r#   r]   rW   r%   s       r   Ú_deletezChannel._deleteï   sQ   € à×(Ñ(¨Ó/ˆØ×Ñ˜eÑ$×*Ñ*Ô,Ø×Ñ×!Ñ! %Ô(Ø�‰×!Ñ! 5 'Õ*r   c                ó4  — | j                  |«      }| j                  j                  |d«      }|€yd}|j                  «       D ]R  }t	        ||j
                  «      }|j                  |«      \  }}|j                  |g«      \  }|||j                  z
  z  }ŒT |S )z6Get the number of pending messages in the topic/queue.Nr   )	r^   rS   r    rC   r   rD   Úget_watermark_offsetsrE   Úoffset)	r#   r]   r@   ÚsizerC   rH   Ú_Ú
end_offsetrI   s	            r   Ú_sizezChannel._sizeö   s¥   € à×(Ñ(¨Ó/ˆà×(Ñ(×,Ñ,¨U°DÓ9ˆØÐØàˆØ"×-Ñ-Ö/ˆJÜ,¨U°J×4HÑ4HÓIˆOØ&×<Ñ<¸_ÓM‰OˆQ�
Ø!)×!3Ñ!3°_Ð4EÓ!FÑÐØ�JÐ!1×!8Ñ!8Ñ8Ñ8‰Dð	 0ð
 ˆr   c           	     óh  — | j                  |«      }|| j                  j                  «       j                  v ryt	        || j
                  j                  dd«      | j
                  j                  dd«      | j
                  j                  di «      ¬«      }| j                  j                  |g¬«       y)z(Create a new topic if it does not exist.NÚnum_partitionsr4   Úreplication_factorÚtopic_config)r‚   rƒ   Úconfig)Ú
new_topics)r^   rw   Úlist_topicsÚtopicsr   rb   r    Úcreate_topics)r#   r]   r%   r   s       r   Ú
_new_queuezChannel._new_queue  s”   € à×(Ñ(¨Ó/ˆØ�D—K‘K×+Ñ+Ó-×4Ñ4Ñ4ØäØØŸ<™<×+Ñ+Ð,<¸aÓ@Ø#Ÿ|™|×/Ñ/Ð0DÀaÓHØ—<‘<×#Ñ# N°BÓ7ô	
ˆð 	�‰×!Ñ!¨e¨WÐ!Õ5r   c                óp   — | j                  |«      }|| j                  j                  «       j                  v S )z Check if a topic already exists.)r^   rw   r‡   rˆ   )r#   r]   r%   s      r   Ú
_has_queuezChannel._has_queue  s0   € à×(Ñ(¨Ó/ˆØ˜Ÿ™×/Ñ/Ó1×8Ñ8Ð8Ð8r   c                óø   — t        i | j                  ¥| j                  j                  d«      xs i ¥«      }	 |j	                  | j
                  ¬«       |S # t        j                  $ r}t        |«      ‚d }~ww xY w)NÚkafka_admin_config)Útimeout)	r   ra   rb   r    r‡   ro   Úconfluent_kafkar   r   )r#   rw   Úes      r   rU   zChannel._open  s‚   € Üð 
Ø× Ñ ð
à�|‰|×ÑÐ 4Ó5Ò;¸ð
ó ˆð
	(à×Ñ t×'=Ñ'=ÐÔ>ð ˆøô ×-Ñ-ò 	(Ü$ QÓ'Ð'ûð	(ús   ¸A ÁA9Á)A4Á4A9c                ó\   — | j                   €| j                  «       | _         | j                   S r'   )rV   rU   r1   s    r   rw   zChannel.client'  s#   € à�<‰<ÐØŸ:™:›<ˆDŒLØ�|‰|Ðr   c                óB   — | j                   j                  j                  S r'   )Ú
connectionrw   Útransport_optionsr1   s    r   rb   zChannel.options-  s   € à�‰×%Ñ%×7Ñ7Ð7r   c                ó.   — | j                   j                  S r'   )r”   rw   r1   s    r   ÚconninfozChannel.conninfo1  s   € à�‰×%Ñ%Ð%r   c                óN   — | j                   j                  d| j                  «      S )Nro   )rb   r    Údefault_wait_time_secondsr1   s    r   ro   zChannel.wait_time_seconds5  s$   € à�|‰|×ÑØ ×!?Ñ!?ó
ð 	
r   c                óN   — | j                   j                  d| j                  «      S )NÚconnection_wait_time_seconds)rb   r    Ú$default_connection_wait_time_secondsr1   s    r   r›   z$Channel.connection_wait_time_seconds;  s%   € à�|‰|×ÑØ*Ø×5Ñ5ó
ð 	
r   c                óÎ  — | j                   j                  }d|j                  › dt        |j                  «      xs t
        › �i}| j                  j                  dd«      }|j                  «       dk7  rC|j                  ||j                  |j                  | j                  j                  d«      dœ«       |j                  | j                  j                  d«      xs i «       |S )Nzbootstrap.serversÚ:Úsecurity_protocolÚ	plaintextÚsasl_mechanism)zsecurity.protocolzsasl.usernamezsasl.passwordzsasl.mechanismÚkafka_common_config)r”   rw   ÚhostnameÚintÚportÚDEFAULT_PORTrb   r    ÚlowerÚupdateÚuseridÚpassword)r#   r—   r…   rŸ   s       r   ra   zChannel.common_configB  sÉ   € à—?‘?×)Ñ)ˆàØ×$Ñ$Ð% Q¤s¨8¯=©=Ó'9Ò'I¼\Ð&JÐKð
ˆð !ŸL™L×,Ñ,Ð-@À+ÓNÐØ×"Ñ"Ó$¨Ò3Ø�M‰MØ%6Ø!)§¡Ø!)×!2Ñ!2Ø"&§,¡,×"2Ñ"2Ð3CÓ"Dñ	ô ð 	�‰�d—l‘l×&Ñ&Ð'<Ó=ÒCÀÔDØˆr   c                óœ   •— t         ‰| �  «        i | _        | j                  j	                  «       D ]  }|j                  «        Œ i | _        y r'   )r!   rv   rT   rS   Úvalues)r#   r@   r&   s     €r   rv   zChannel.closeU  sA   ø€ Ü‰‰ŒØ "ˆÔà×-Ñ-×4Ñ4Ö6ˆHØ�N‰NÕð 7ð !#ˆÕr   )r   r   r   r   r+   r   r™   rœ   rV   r"   r^   rd   r>   rl   rt   ry   r€   rŠ   rŒ   rU   Úpropertyrw   rb   r—   r   ro   r›   ra   rv   r(   r)   s   @r   rP   rP   Ÿ   sÛ   ø„ Ùà
€CØ€Gà !ÐØ+,Ð(Ø€Gô$ò+òòò"òDò*+òò 6ò9ò
ð ñó ðð
 ñ8ó ð8ð ñ&ó ð&ð ñ
ó ð
ð
 ñ
ó ð
ð ñó ð÷$#ð #r   rP   c                  ó\   ‡ — e Zd ZdZd	d
d„ZeZeZdZdZ	e
fZˆ fd„Zd„ Zˆ fd„Zˆ fd„Zˆ xZS )Ú	TransportzKafka Transport.c                 ó   — y r'   r   )r#   ÚuriÚinclude_passwordÚmasks       r   Úas_urizTransport.as_urib  rN   r   ÚkafkaÚconfluentkafkac                óH   •— t         €t        d«      ‚t        ‰| �  |fi |¤Ž y )Nz,The confluent-kafka library is not installed)r�   ÚImportErrorr!   r"   )r#   rw   r%   r&   s      €r   r"   zTransport.__init__p  s'   ø€ ÜÐ"ÜÐLÓMÐMÜ‰Ñ˜Ñ* 6Ó*r   c                ó"   — t         j                  S r'   )r�   Ú__version__r1   s    r   Údriver_versionzTransport.driver_versionu  s   € Ü×*Ñ*Ð*r   c                ó    •— t         ‰| �  «       S r'   )r!   Úestablish_connection)r#   r&   s    €r   r½   zTransport.establish_connectionx  s   ø€ Ü‰wÑ+Ó-Ð-r   c                ó"   •— t         ‰| �  |«      S r'   )r!   Úclose_connection)r#   r”   r&   s     €r   r¿   zTransport.close_connection{  s   ø€ Ü‰wÑ'¨
Ó3Ð3r   )Fz**)r±   r[   Úreturnr[   )r   r   r   r   r´   rP   r¦   Údefault_portÚdriver_typeÚdriver_namer   Úrecoverable_connection_errorsr"   r»   r½   r¿   r(   r)   s   @r   r¯   r¯   _  sG   ø„ Ùôð €Gà€Là€KØ"€Kð 	ð%Ð!ô+ò
+ô.÷4ð 4r   r¯   )#r   Ú
__future__r   r]   r   Úkombu.transportr   Úkombu.utilsr   Úkombu.utils.encodingr   Úkombu.utils.jsonr   r	   r�   r
   r   r   r   Úconfluent_kafka.adminr   r   ÚKAFKA_CONNECTION_ERRORSÚKAFKA_CHANNEL_ERRORSr¸   Ú	kombu.logr   r   rr   r¦   r   r   r+   rP   r¯   r   r   r   Ú<module>rÎ      sÉ   ðñ:õx #å å #Ý 'Ý -ß )ð8Û÷1ó 1ç;à ÐØÐõ !á	�HÓ	€à€ô˜ô ô=ˆg�o‰oô =ô4ˆ'�+‰+ô 4ôn}#ˆg�o‰oô }#ô@4�×!Ñ!õ 4øða ò 8Ø€OØ57Ð7ÐÒ2ð8ús   ªB6 Â6CÃC