o
    `jM                     @  s0  d Z ddlmZ ddlmZ ddlmZ ddlZddlm	Z	m
Z
mZmZmZmZ ddlmZ ddlmZmZmZ ddlmZ dd	lmZ dd
lmZ ddlmZ e	rXddlmZ e e!Z"G dd dee#eZ$G dd deZ%G dd deZ&G dd deZ'G dd deZ(G dd deZ)G dd dZ*dS )aK  Hanging-transaction tooling mixin for KafkaAdminClient (KIP-664).

Exposes four wire APIs (ListTransactions, DescribeTransactions,
DescribeProducers, WriteTxnMarkers in admin abort mode) plus the
``find_hanging_transactions`` convenience that ties them together,
mirroring the Java tool's ``kafka-transactions.sh --find-hanging``.
    )annotations)defaultdict)EnumN)TYPE_CHECKINGDictList
NamedTupleOptionalSet)DescribeProducersRequestDescribeTransactionsRequestListTransactionsRequest)CoordinatorType)WriteTxnMarkersRequest)TopicPartition)
EnumHelper)KafkaConnectionManagerc                   @  s4   e Zd ZdZdZdZdZdZdZdZ	dZ
d	Zd
ZdS )TransactionStatez]Broker-reported transaction states (DescribeTransactions /
    ListTransactions wire values).EmptyOngoingPrepareCommitPrepareAbortCompleteCommitCompleteAbortDeadPrepareEpochFenceUnknownN)__name__
__module____qualname____doc__EMPTYONGOINGPREPARE_COMMITPREPARE_ABORTCOMPLETE_COMMITCOMPLETE_ABORTDEADPREPARE_EPOCH_FENCEUNKNOWN r*   r*   W/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/admin/_transactions.pyr      s    r   c                   @  s*   e Zd ZU dZded< ded< ded< dS )	TransactionListingz)One row from a ListTransactions response.strtransactional_idintproducer_idr   stateNr   r   r   r    __annotations__r*   r*   r*   r+   r,   -   s
   
 r,   c                   @  sJ   e Zd ZU dZded< ded< ded< ded< ded< ded	< d
ed< dS )TransactionDescriptionzhOne transactional id's state as returned by DescribeTransactions,
    plus the coordinator that owns it.r/   coordinator_idr   r1   r0   producer_epochtransaction_timeout_mstransaction_start_time_mszSet[TopicPartition]topic_partitionsNr2   r*   r*   r*   r+   r4   4   s   
 r4   c                   @  sB   e Zd ZU dZded< ded< ded< ded< ded< ded< d	S )
ProducerStatez.One ActiveProducer row from DescribeProducers.r/   r0   r6   last_sequencelast_timestampcoordinator_epoch current_transaction_start_offsetNr2   r*   r*   r*   r+   r:   @   s   
 r:   c                   @  s   e Zd ZU ded< dS )PartitionProducerStatezList[ProducerState]active_producersN)r   r   r   r3   r*   r*   r*   r+   r?   J   s   
 r?   c                   @  s6   e Zd ZU dZded< ded< ded< dZded< d	S )
AbortTransactionSpeczInputs for ``abort_transaction``. ``coordinator_epoch=-1`` is the
    sentinel used by the Java admin tool to bypass the epoch check; the
    partition leader still validates ``producer_id``/``producer_epoch``
    against current state.r   topic_partitionr/   r0   r6   r=   N)r   r   r   r    r3   r=   r*   r*   r*   r+   rA   N   s   
 rA   c                   @  s   e Zd ZU dZded< ded< e			d,ddZed	d
 Z			d-ddZ			d-ddZ	edd Z
edd Zdd Zdd Zedd Zedd Zd.ddZd.ddZedd  Zed!d" Zd#d$ Zd%d& Z		'd/d(d)Z		'd/d*d+ZdS )0TransactionsAdminMixinz4Mixin providing KIP-664 hanging-transaction tooling.z'KafkaConnectionManager'_managerdictconfigNc                 C  s   ddi}| r
t | ng |d< |rt |ng |d< |d ur+t||d< t|d d|d< |d ur<||d< t|d d|d< td	i |S )
Nmin_versionr   state_filtersproducer_id_filtersduration_filter   transactional_id_pattern   r*   )listr/   maxr   )rI   rJ   duration_filter_msrM   kwargsr*   r*   r+   _list_transactions_request`   s   z1TransactionsAdminMixin._list_transactions_requestc                 C  sV   t | j}|t jur|d| g }| jD ]}|t|j|j	t
|jd q|S )Nz2ListTransactionsRequest failed with response '{}'.)r.   r0   r1   )Errorsfor_code
error_codeNoErrorformattransaction_statesappendr,   r.   r0   r   transaction_state)response
error_typelistingstxnr*   r*   r+   #_list_transactions_process_responseo   s   


z:TransactionsAdminMixin._list_transactions_process_responsec           
        s   |d ur| j jtdk rtd|d ur%| j jtdk r%td|d u r4dd | j j D }i }|D ]}| j||||d}| j j	||dI d H }	| 
|	||< q8|S )	NrL   zXduration_filter_ms requires broker support for ListTransactions v1+ (Apache Kafka 3.8+).rN   zhtransactional_id_pattern requires broker support for ListTransactions v2+ (Apache Kafka 4.1+, KIP-1152).c                 S  s   g | ]}|j qS r*   node_id).0brokerr*   r*   r+   
<listcomp>   s    zCTransactionsAdminMixin._async_list_transactions.<locals>.<listcomp>)rI   rJ   rQ   rM   ra   )rE   broker_version_dataapi_versionr   rT   UnsupportedVersionErrorclusterbrokersrS   sendr`   )
self
broker_idsrJ   rI   rQ   rM   results	broker_idrequestr\   r*   r*   r+   _async_list_transactions~   s0   z/TransactionsAdminMixin._async_list_transactionsc                 C  s   | j | j|||||S )a  List active transactions across all brokers (or a subset).

        Each broker hosts a slice of the ``__transaction_state`` topic,
        so a full listing requires sharding the request to every broker
        and concatenating the results.

        Keyword Arguments:
            broker_ids ([int], optional): Brokers to query. Default: every
                broker in the cluster metadata.
            producer_id_filters ([int], optional): Only return transactions
                whose ``producer_id`` is in this list.
            state_filters ([str], optional): Only return transactions whose
                broker-reported state matches. Accepts :class:`TransactionState`
                members or their string wire values.
            duration_filter_ms (int, optional): Only return transactions
                running longer than this. Requires broker >= 3.8
                (ListTransactions v1+).
            transactional_id_pattern (str, optional): Only return
                transactions whose transactional id matches this regex.
                Requires broker >= 4.1 (ListTransactions v2+, KIP-1152).

        Returns:
            dict: A dict mapping broker ``node_id`` to a list of
            :class:`TransactionListing`.
        )rE   runrq   )rl   rm   rJ   rI   rQ   rM   r*   r*   r+   list_transactions   s   z(TransactionsAdminMixin.list_transactionsc                 C  s   t t| dS )Ntransactional_ids)r   rO   rt   r*   r*   r+   _describe_transactions_request   s   z5TransactionsAdminMixin._describe_transactions_requestc              
   C  s   i }| j D ]B}t|j}|tjur|d|jt }|jD ]}|j	D ]}|
t|j| q%q t|t|j|j|j|j|j|d||j< q|S )Nz=DescribeTransactionsRequest failed for transactional id '{}'.)r5   r1   r0   r6   r7   r8   r9   )rY   rT   rU   rV   rW   rX   r.   settopics
partitionsaddr   topicr4   r   r[   r0   r6   r7   r8   )r\   r5   rn   r_   r]   r9   r{   	partitionr*   r*   r+   '_describe_transactions_process_response   s.   



	z>TransactionsAdminMixin._describe_transactions_process_responsec           
        s   t |}|s	i S | j|tjdI d H }tt }| D ]\}}|| | qi }| D ]\}}| |}| jj	||dI d H }	|
| |	| q.|S )N)key_typera   )rO   _find_coordinator_idsr   TRANSACTIONr   itemsrZ   rv   rE   rk   updater}   )
rl   ru   coordinator_idscoordinator_to_txn_idstxn_idcoord_idrn   txn_idsrp   r\   r*   r*   r+   _async_describe_transactions   s    
z3TransactionsAdminMixin._async_describe_transactionsc                 C     | j | j|S )a  Describe one or more transactions by transactional id.

        Each request is routed to the transaction coordinator that owns
        the transactional id (discovered via FindCoordinator with
        ``CoordinatorType.TRANSACTION``).

        Arguments:
            transactional_ids: Iterable of transactional id strings.

        Returns:
            dict: A dict mapping ``transactional_id`` (str) to
            :class:`TransactionDescription`.

        Raises:
            TransactionalIdNotFoundError: If a transactional id is unknown
                to its coordinator.
            BrokerResponseError: For any other per-id error.
        )rE   rr   r   )rl   ru   r*   r*   r+   describe_transactions   s   z,TransactionsAdminMixin.describe_transactionsc                   s&   t j  fdd|  D }t |dS )Nc                   s    g | ]\}} |t |d qS )namepartition_indexes)rO   )rc   r   partsTopicr*   r+   re     s    zFTransactionsAdminMixin._describe_producers_request.<locals>.<listcomp>)rx   )r   TopicRequestr   )partitions_by_topicrx   r*   r   r+   _describe_producers_request  s
   

z2TransactionsAdminMixin._describe_producers_requestc                 C  st   i }| j D ]2}|jD ],}t|j|j}t|j}|tjur'|d	||j
dd |jD }t|d||< q
q|S )Nz*DescribeProducersRequest failed for {}: {}c              
   S  s,   g | ]}t |j|j|j|j|j|jd qS ))r0   r6   r;   r<   r=   r>   )r:   r0   r6   r;   r<   r=   current_txn_start_offset)rc   pr*   r*   r+   re     s    zOTransactionsAdminMixin._describe_producers_process_response.<locals>.<listcomp>)r@   )rx   ry   r   r   partition_indexrT   rU   rV   rW   rX   error_messager@   r?   )r\   rn   r{   r|   tpr]   	producersr*   r*   r+   $_describe_producers_process_response
  s    




z;TransactionsAdminMixin._describe_producers_process_responsec                   s   t |}|s	i S |d ur4tt }|D ]}||j |j q| |}| jj||dI d H }| |S | 	|I d H }i }|
 D ].\}	}
tt }|
D ]}||j |j qL| |}| jj||	dI d H }|| | qB|S Nra   )rO   r   r{   rZ   r|   r   rE   rk   r    _async_get_leader_for_partitionsr   r   )rl   ry   ro   r   r   rp   r\   leader2partitionsrn   leader
leader_tpsr*   r*   r+   _async_describe_producers"  s*   


z0TransactionsAdminMixin._async_describe_producersc                 C     | j | j||S )a  Describe active producer state on a set of topic partitions.

        Arguments:
            partitions: Iterable of :class:`~kafka.TopicPartition`.

        Keyword Arguments:
            broker_id (int, optional): Replica to query. Default: the
                partition leader (discovered from cluster metadata).

        Returns:
            dict: A dict mapping :class:`~kafka.TopicPartition` to
            :class:`PartitionProducerState`.

        Raises:
            BrokerResponseError: For any per-partition error (e.g.
                ``NotLeaderOrFollowerError`` if the chosen broker is not
                a replica).
        )rE   rr   r   )rl   ry   ro   r*   r*   r+   describe_producers=  s   z)TransactionsAdminMixin.describe_producersc                 C  sD   t j}|j}|| j| jd|| jj| jjgdg| jd}t |gdS )NFr   )r0   r6   transaction_resultrx   r=   )markers)	r   WritableTxnMarkerWritableTxnMarkerTopicr0   r6   rB   r{   r|   r=   )specMarkerr   markerr*   r*   r+   _abort_transaction_requestT  s   z1TransactionsAdminMixin._abort_transaction_requestc                 C  sN   | j D ]!}|jD ]}|jD ]}t|j}|tjur"|d|jqqqd S )Nz%WriteTxnMarkers (abort) failed for {})	r   rx   ry   rT   rU   rV   rW   rX   rB   )r\   r   resultr{   r|   r]   r*   r*   r+   #_abort_transaction_process_responseb  s   




z:TransactionsAdminMixin._abort_transaction_process_responsec                   sR   |  |jgI d H }tt|}| |}| jj||dI d H }| || d S r   )r   rB   nextiterr   rE   rk   r   )rl   r   r   r   rp   r\   r*   r*   r+   _async_abort_transactionl  s   
z/TransactionsAdminMixin._async_abort_transactionc                 C  r   )av  Administratively abort an open transaction on a partition.

        Sends a WriteTxnMarkers request (with ``transaction_result=False``)
        to the partition leader. The leader validates ``producer_id`` /
        ``producer_epoch`` against current state before writing the
        abort marker. Pass ``coordinator_epoch=-1`` (the default) to
        signal an admin abort that bypasses the coordinator-epoch
        guard, matching the Java AdminClient behaviour.

        Arguments:
            spec (:class:`AbortTransactionSpec`): Target partition,
                producer id/epoch, and optional coordinator epoch.
        )rE   rr   r   )rl   r   r*   r*   r+   abort_transactions  s   z(TransactionsAdminMixin.abort_transaction頻 c                   s   |d }| j |dI d H }tdd | D }|sg S | |I d H }dd l}t| d }g }	| D ]7\}
}|jtj	tj
tjtjfv rIq7|jdk rOq7||j }||k rYq7|	|
|j|j|jj||jt|jd q7|	S )Ni )rm   c                 S  s   h | ]
}|D ]}|j qqS r*   )r.   )rc   txnstr*   r*   r+   	<setcomp>  s    zJTransactionsAdminMixin._async_find_hanging_transactions.<locals>.<setcomp>r   i  )r.   r0   r6   r1   age_msr5   r9   )rq   sortedvaluesr   timer/   r   r1   r   r!   r%   r&   r'   r8   rZ   r0   r6   valuer5   r9   )rl   rm   max_transaction_timeout_msthreshold_mslistings_by_brokerr   descriptionsr   now_mshangingr   descr   r*   r*   r+    _async_find_hanging_transactions  s<   


	z7TransactionsAdminMixin._async_find_hanging_transactionsc                 C  r   )a  Detect transactions whose age exceeds the broker timeout + 5min.

        Convenience wrapper that runs :meth:`list_transactions` against
        each broker, then :meth:`describe_transactions` to read
        ``transaction_start_time_ms``, and filters to transactions in an
        active state whose age exceeds the threshold. Mirrors
        ``kafka-transactions.sh --find-hanging``.

        Keyword Arguments:
            broker_ids ([int], optional): Brokers to query. Default: all.
            max_transaction_timeout_ms (int): Suspected-hang threshold.
                Default: 900000 (15 minutes -- Kafka's default
                ``transaction.max.timeout.ms``).

        Returns:
            list: One dict per suspected hanging transaction with keys
            ``transactional_id``, ``producer_id``, ``producer_epoch``,
            ``state``, ``age_ms``, ``coordinator_id``,
            ``topic_partitions``.
        )rE   rr   r   )rl   rm   r   r*   r*   r+   find_hanging_transactions  s   z0TransactionsAdminMixin.find_hanging_transactions)NNNN)NNNNN)N)Nr   )r   r   r   r    r3   staticmethodrS   r`   rq   rs   rv   r}   r   r   r   r   r   r   r   r   r   r   r   r   r*   r*   r*   r+   rD   Y   sV   
 


"







	
%rD   )+r    
__future__r   collectionsr   enumr   loggingtypingr   r   r   r   r	   r
   kafka.errorserrorsrT   !kafka.protocol.admin.transactionsr   r   r   kafka.protocol.metadatar   #kafka.protocol.producer.transactionr   kafka.structsr   
kafka.utilr   kafka.net.managerr   	getLoggerr   logr-   r   r,   r4   r:   r?   rA   rD   r*   r*   r*   r+   <module>   s,     

