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Zd dlZzd dlmZ	 W n e
y/   dZ	Y nw ddlmZ d dlmZ eeZG dd deZG dd	 d	ZdS )
    N)Session   )SaslMechanism)KafkaConfigurationErrorc                   @   sD   e 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 )SaslMechanismAwsMskIamc                 K   sX   t d u rtd|dddkrtdd|vrtd|d | _d | _d| _d| _d S )	Nz+AWS_MSK_IAM requires the "botocore" packagesecurity_protocol SASL_SSLzAWS_MSK_IAM requires SASL_SSLhostz'AWS_MSK_IAM requires host configurationF)BotoSessionRuntimeErrorgetr   r
   _auth_is_done_is_authenticated)selfconfig r   P/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/net/sasl/msk.py__init__   s   

zSaslMechanismAwsMskIam.__init__c                 C   sD   t  }|  }|dstdt| j|j|j|d|j	dS )NregionzJUnable to determine region for AWS MSK cluster. Is AWS_DEFAULT_REGION set?)r
   
access_key
secret_keyr   token)
r   get_credentialsget_frozen_credentialsget_config_variabler   AwsMskIamClientr
   r   r   r   )r   sessioncredentialsr   r   r   _build_client$   s   
z$SaslMechanismAwsMskIam._build_clientc                 C   s   |   }td|j | S )Nz'Generating auth token for MSK scope: %s)r    logdebug_scopefirst_message)r   clientr   r   r   
auth_bytes1   s   z!SaslMechanismAwsMskIam.auth_bytesc                 C   s    d| _ |dk| _|d| _d S )NT    utf-8)r   r   decoder   )r   r&   r   r   r   receive6   s   
zSaslMechanismAwsMskIam.receivec                 C      | j S N)r   r   r   r   r   is_done;      zSaslMechanismAwsMskIam.is_donec                 C   r+   r,   )r   r-   r   r   r   is_authenticated>   r/   z'SaslMechanismAwsMskIam.is_authenticatedc                 C   s   | j stdd| jf S )NzNot authenticated yet!z'Authenticated via SASL / AWS_MSK_IAM %s)r0   r   r   r-   r   r   r   auth_detailsA   s   z#SaslMechanismAwsMskIam.auth_detailsN)
__name__
__module____qualname__r   r    r&   r*   r.   r0   r1   r   r   r   r   r      s    r   c                   @   s   e Zd Zejej d ZdddZedd Z	edd Z
ed	d
 Zedd Zedd Zedd Zedd Zedd Zdd Zdd Zdd ZdS )r   z-._~Nc                 C   s~   d| _ d| _tj| _d|fg| _d| _d| _d| j| _	t
j
 }|d| _|d| _|| _|| _|| _|| _|| _d	S )
aM  
        Arguments:
            host (str): The hostname of the broker.
            access_key (str): An AWS_ACCESS_KEY_ID.
            secret_key (str): An AWS_SECRET_ACCESS_KEY.
            region (str): An AWS_REGION.
            token (Optional[str]): An AWS_SESSION_TOKEN if using temporary
                credentials.
        zAWS4-HMAC-SHA256900r
   
2020_10_22zkafka-clusterz
{}:Connectz%Y%m%dz%Y%m%dT%H%M%SZN)	algorithmexpireshashlibsha256hashfuncheadersversionserviceformatactiondatetimeutcnowstrftime	datestamp	timestampr
   r   r   r   r   )r   r
   r   r   r   r   nowr   r   r   r   J   s    


zAwsMskIamClient.__init__c                 C   
   d | S )Nz{0.access_key}/{0._scope}r?   r-   r   r   r   _credentiali      
zAwsMskIamClient._credentialc                 C   rG   )Nz1{0.datestamp}/{0.region}/{0.service}/aws4_requestrH   r-   r   r   r   r#   m   rJ   zAwsMskIamClient._scopec                 C   s   d tdd | jD S )z
        Returns (str):
            An alphabetically sorted, semicolon-delimited list of lowercase
            request header names.
        ;c                 s   s    | ]	\}}|  V  qd S r,   )lower).0k_r   r   r   	<genexpr>x   s    z2AwsMskIamClient._signed_headers.<locals>.<genexpr>)joinsortedr<   r-   r   r   r   _signed_headersq   s   zAwsMskIamClient._signed_headersc                 C   s   d tdj | jd S )z
        Returns (str):
            A newline-delited list of header names and values.
            Header names are lowercased.
        
:)rQ   mapr<   r-   r   r   r   _canonical_headersz   s   z"AwsMskIamClient._canonical_headersc                 C   s*   |  d }ddd| j| j| j|fS )a'  
        Returns (str):
            An AWS Signature Version 4 canonical request in the format:
                <Method>

                <Path>

                <CanonicalQueryString>

                <CanonicalHeaders>

                <SignedHeaders>

                <HashedPayload>
        r'   rT   GET/)r;   	hexdigestrQ   _canonical_querystringrW   rS   )r   hashed_payloadr   r   r   _canonical_request   s   z"AwsMskIamClient._canonical_requestc                    s   g }| d jf | d jf | d jf | d jf | d jf  jr5| d jf | d jf d fd	d
|D S )za
        Returns (str):
            A '&'-separated list of URI-encoded key/value pairs.
        ActionzX-Amz-AlgorithmzX-Amz-Credentialz
X-Amz-DatezX-Amz-ExpireszX-Amz-Security-TokenzX-Amz-SignedHeaders&c                 3   s,    | ]\}}  |d    | V  qdS )=N)
_uriencode)rM   rN   vr-   r   r   rP      s   * z9AwsMskIamClient._canonical_querystring.<locals>.<genexpr>)	appendr@   r7   rI   rE   r8   r   rS   rQ   )r   paramsr   r-   r   r[      s   z&AwsMskIamClient._canonical_querystringc                 C   sF   |  d| j d| j}|  || j}|  || j}|  |d}|S )z
        Returns (bytes):
            An AWS Signature V4 signing key generated from the secret_key, date,
            region, service, and request type.
        AWS4r(   aws4_request)_hmacr   encoderD   r   r>   )r   keyr   r   r   _signing_key   s
   zAwsMskIamClient._signing_keyc                 C   s.   |  | jd }d| j| j| j|fS )z
        Returns (str):
            A string used to sign the AWS Signature V4 payload in the format:
                <Algorithm>

                <Timestamp>

                <Scope>

                <CanonicalRequestHash>
        r(   rT   )r;   r]   rh   rZ   rQ   r7   rE   r#   )r   canonical_request_hashr   r   r   _signing_str   s   
zAwsMskIamClient._signing_strc                 C   s   t jj|| jdS )a  
        Arguments:
            msg (str): A string to URI-encode.

        Returns (str):
            The URI-encoded version of the provided msg, following the encoding
            rules specified: https://github.com/aws/aws-msk-iam-auth#uriencode
        )safe)urllibparsequoteUNRESERVED_CHARS)r   msgr   r   r   ra      s   	zAwsMskIamClient._uriencodec                 C   s   t j||d| jd S )z
        Arguments:
            key (bytes): A key to use for the HMAC digest.
            msg (str): A value to include in the HMAC digest.
        Returns (bytes):
            An HMAC digest of the given key and msg.
        r(   	digestmod)hmacnewrh   r;   digest)r   ri   rr   r   r   r   rg      s   zAwsMskIamClient._hmacc                 C   sn   t j| j| jd| jd }| j| jd| j	| j
| j| j| j| j|d
}| jr-| j|d< tj|dddS )z
        Returns (bytes):
            An encoded JSON authentication payload that can be sent to the
            broker.
        r(   rs   zkafka-python)
r=   r
   z
user-agentr@   zx-amz-algorithmzx-amz-credentialz
x-amz-datezx-amz-signedheaderszx-amz-expireszx-amz-signaturezx-amz-security-token),rU   )
separators)ru   rv   rj   rl   rh   r;   rZ   r=   r
   r@   r7   rI   rE   rS   r8   r   jsondumps)r   	signaturerr   r   r   r   r$      s*   

zAwsMskIamClient.first_messager,   )r2   r3   r4   stringascii_lettersdigitsrq   r   propertyrI   r#   rS   rW   r]   r[   rj   rl   ra   rg   r$   r   r   r   r   r   G   s,    









r   )rA   r9   ru   rz   loggingr}   rn   botocore.sessionr   r   ImportErrorabcr   kafka.errorsr   	getLoggerr2   r!   r   r   r   r   r   r   <module>   s"    
0