o
    `j:                     @   s:  d Z ddlZddlZddlZddlmZmZm	Z
 ddlmZmZ ddlmZmZ ddlmZmZ ddlmZmZ ddlm Z!m"Z# d	Z$d
d e%dD Z&ej'j(Z)dd e%dD Z*dd e%dD Z+dd e%dD Z,ej-Z.e.j/Z0dd e%dD Z1dd e%dD Z2e#j3Z4e4j5Z6dd e%dD Z7dd e%dD Z8e Z9eddZ:edddZ;eddddZ<eddddde&dZ=ed ddddde&d!Z>ed"dddddde*g d#d$
Z?eddd%e+d&Z@e!dddd%e+d'ZAe!d(ddd%e+d'ZBe
de8d)ZCedde8d*ZDede,d+ZEed de,d,ZFed"ddde1d-ZGee2dd.ZHe#de2dd/ZIe#d(e7dd/ZJeCK ZLeEK ZMeHK ZNeGjKd"dZOeJjKd(dZPe9K e:jKddksJJ d0e;K e<jKddksYJ d1e=K e>jKd dkshJ d2e@K eAjKddkswJ d3eCK eDjKddksJ d4eEK eFjKd dksJ d5eHK eIjKddksJ d6d7d8 ZQd9d: ZRd;d<eQe9jKd=eQe:jKddfd>d?eQe;jKd@eQe<jKddfdAdBeQe=jKdCeQe>jKd dfdDdddEeQe?jKd"dfdFdGeQe@jKdHeQeAjKddfdIdddJeQeBjKd(dfgdKdLeRe
eLdMeReeLddfdNdOeReeNdPeRe#eNddfdQdddReRe#ePd(dfdSdTeReeMdUeReeMd dfdVdddWeReeOd"dfgdXZSdYdZ ZTd[d\ ZUd]d^ ZVeWd_kreX ZYi ZZd`da Z[eS\ D ])\Z]Z^e^D ]!Z_e_\Z`ZaZbZcZdeadur{e[eaeYeeaeb e[eceYeeced qdq^eZrdbejfvreVeZ dS dS dS dS )ca  Benchmark old (Struct-based) vs new (JSON/ApiMessage-based) protocol encode/decode.

Benchmarks focus on realistic client operations:
  - Encode: Request objects (what the client sends)
  - Decode: Response objects (what the client receives)

Usage:
    python kafka/benchmarks/protocol_old_vs_new.py [--fast] [--quiet]

    --fast      Use pyperf fast mode (fewer iterations, less stable)
    --quiet     Suppress per-benchmark warnings

Comparison tables are printed after all benchmarks complete.
    N)ApiVersionsRequest_v0ApiVersionsRequest_v3ApiVersionsResponse_v0)FetchRequest_v4FetchResponse_v4)ProduceRequest_v3ProduceResponse_v3)ApiVersionsRequestApiVersionsResponse)FetchRequestFetchResponse)ProduceRequestProduceResponses                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                   c                 C   &   g | ]}d | dd t dD fqS )test-topic-%dc                 S   s   g | ]	}||d  dfqS )      .0pr   r   b/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/benchmarks/protocol_old_vs_new.py
<listcomp><   s    <listcomp>.<listcomp>   ranger   tr   r   r   r   ;       r      c                 C   s.   g | ]}t jd d| dd tdD dqS )   r   c                 S   s&   g | ]}t d |d|d ddddqS )r!   r   r   r   )version	partitioncurrent_leader_epochfetch_offsetlast_fetched_epochlog_start_offsetpartition_max_bytes)_FetchPartitionr   r   r   r   r   F   s    r   r   r#   topic
partitions)NewFetchReq
FetchTopicr   r   r   r   r   r   E   s    c                 C   r   )r   c                 S   s   g | ]}|t fqS r   _RECORDS_BYTESr   r   r   r   r   Q   s    r   r   r   r   r   r   r   r   P   r   c                 C   r   )r   c                 S   s$   g | ]}|d d| d| dt fqS )r   r     Nr0   r   r   r   r   r   Z   s    r   r   r   r   r   r   r   r   Y   r   c                 C   ,   g | ]}t d d| dd tdD dqS )r!   r   c                 S   s.   g | ]}t d |dd| d| dg tdd	qS )r!   r   r   r2   r"   )	r#   partition_index
error_codehigh_watermarklast_stable_offsetr(   aborted_transactionsrecordspreferred_read_replica)_FetchRespPartr1   r   r   r   r   r   e   s    r   r   r+   )_FetchRespTopicr   r   r   r   r   r   d   s    c                 C   r   )r   c                 S   s   g | ]
}|d |d dfqS )r   r   r"   r   r   r   r   r   r   q   s    r   r   r   r   r   r   r   r   p   r   c                 C   r3   )	   r   c                 S   s,   g | ]}t d |d|d d|d g ddqS )r=   r   r   r"     N)r#   indexr5   base_offsetlog_append_time_msr(   record_errorserror_message)_ProduceRespPartr   r   r   r   r   |   s    r   r   )r#   namepartition_responses)_ProduceRespTopicr   r   r   r   r   r   {   s    c                 C   s   g | ]	}|d |d fqS )r   r   r   )r   ir   r   r   r      s    2   )r#   zkafka-pythonz3.0.0)client_software_nameclient_software_versionr   )r#   rJ   rK   r"   r>      i   )
replica_idmax_wait_ms	min_bytes	max_bytesisolation_leveltopics   )r#   rM   rN   rO   rP   rQ   rR   r!    )
r#   rN   rO   rP   rQ   
session_idsession_epochrR   forgotten_topics_datarack_idi0u  )transactional_idacks
timeout_ms
topic_data)r#   rY   rZ   r[   r\   r=   )r5   api_keys)r#   r5   r]   )throttle_time_ms	responses)r#   r^   r_   )r#   r^   r5   rU   r_   )r_   r^   )r#   r_   r^   z&ApiVersionsRequest v0 encode mismatch!z&ApiVersionsRequest v3 encode mismatch!z FetchRequest v4 encode mismatch!z"ProduceRequest v3 encode mismatch!z'ApiVersionsResponse v0 encode mismatch!z!FetchResponse v4 encode mismatch!z#ProduceResponse v3 encode mismatch!c                        fdd}|S )Nc                    s0   t  }t| D ]	} i  qt  | S N)pyperfperf_counterr   loopst0_argsfunckwargsr   r   bench  s   z_make_bench_func.<locals>.benchr   )rj   ri   rk   rl   r   rh   r   _make_bench_func     rm   c                    r`   )Nc                    s:   t  }t| D ]} jtfi  qt  | S ra   )rb   rc   r   decodeioBytesIOrd   clsencodedrk   r   r   rl   %  s   z _make_decode_func.<locals>.benchr   )rs   rt   rk   rl   r   rr   r   _make_decode_func$  rn   ru   zApiVersionsReq v0 (simple)encode_ApiVersionsReq_v0_oldencode_ApiVersionsReq_v0_newzApiVersionsReq v3 (flexible)encode_ApiVersionsReq_v3_oldencode_ApiVersionsReq_v3_newzFetchReq v4encode_FetchReq_v4_oldencode_FetchReq_v4_newzFetchReq v12 (flexible)encode_FetchReq_v12_newzProduceReq v3encode_ProduceReq_v3_oldencode_ProduceReq_v3_newzProduceReq v9 (flexible)encode_ProduceReq_v9_newzApiVersionsResp v0 (simple)decode_ApiVersionsResp_v0_olddecode_ApiVersionsResp_v0_newzProduceResp v3decode_ProduceResp_v3_olddecode_ProduceResp_v3_newzProduceResp v9 (flexible)decode_ProduceResp_v9_newzFetchResp v4decode_FetchResp_v4_olddecode_FetchResp_v4_newzFetchResp v12 (flexible)decode_FetchResp_v12_new)zencode (requests)zdecode (responses)c                 C   s@   | d }|dkrd| S |dkrd| S |dkrd| S d| S )zBFormat seconds into a human-readable string with appropriate unit.g    .Ar   z%.0f usd   
   z%.1f usz%.2f usr   )secondsusr   r   r   _format_time^  s   r   c           	      C   s   t   t d|   t   tdd |D }t|td}d|ddddf }d	d
|d  dddf }t | t | |D ]2\}}}|dur`|dkrL|| ntd}t d||t|t||f  q;t d||dt|df  q;dS )zPrint a markdown-style table.z### %sc                 s   s    | ]	}t |d  V  qdS )r   N)len)r   rr   r   r   	<genexpr>p  s    z_print_table.<locals>.<genexpr>Messagez| %-*s | %10s | %10s | %5s |OldNewRatioz|%s|%s|%s|%s|-   z------------z-------Nr   infz| %-*s | %10s | %10s | %5.1fx |zn/a)printmaxr   floatr   )	titlerowscol1headersepdescold_meannew_meanratior   r   r   _print_tablek  s,   

r   c                 C   s@  t   t d t d t d t D ]2\}}g }|D ] }|\}}}}}|| v r;|r/| |nd}	|||	| | f q|rEt| | qt   t d t   tD ]H}g }
t| D ]*}|\}}}}}|r|| v r|| v r| | }	| | }|
|	dkr||	 ntd qZ|
rt|
}t	|
}t d| ||f  qRt   dS )z9Print comparison tables from collected benchmark results.zF======================================================================z4Protocol Benchmark: Old (Struct) vs New (ApiMessage)Nz### Summaryr   r   z#- **%s**: %.1fx - %.1fx (new / old))
r   
BENCHMARKSitemsgetappendr   
capitalizer   minr   )resultscategory
bench_defsr   entryr   old_namerg   new_nameold_valratiosnew_valmin_rmax_rr   r   r   _print_summary  sF   

r   __main__c                 C   s4   |d urz	|  t| < W d S  ty   Y d S w d S ra   )meanr   	Exception)rE   rl   r   r   r   _record  s   r   z--worker)g__doc__rp   sysrb   kafka.protocol.old.api_versionsr   OldApiVersionsReq_v0r   OldApiVersionsReq_v3r   OldApiVersionsResp_v0kafka.protocol.old.fetchr   OldFetchReq_v4r   OldFetchResp_v4kafka.protocol.old.producer   OldProduceReq_v3r   OldProduceResp_v3$kafka.protocol.metadata.api_versionsr	   NewApiVersionsReqr
   NewApiVersionsRespkafka.protocol.consumer.fetchr   r.   r   NewFetchRespkafka.protocol.producer.producer   NewProduceReqr   NewProduceRespr1   r   _FETCH_REQ_TOPICS_V4r/   FetchPartitionr*   _FETCH_REQ_TOPICS_V12_PRODUCE_REQ_TOPICS_FETCH_RESP_TOPICS_V4FetchableTopicResponser<   PartitionDatar;   _FETCH_RESP_TOPICS_V12_PRODUCE_RESP_TOPICS_V3TopicProduceResponserG   PartitionProduceResponserD   _PRODUCE_RESP_TOPICS_V9_API_KEYS_DATAOLD_APIVER_REQ_V0NEW_APIVER_REQ_V0OLD_APIVER_REQ_V3NEW_APIVER_REQ_V3OLD_FETCH_REQ_V4NEW_FETCH_REQ_V4NEW_FETCH_REQ_V12OLD_PRODUCE_REQ_V3NEW_PRODUCE_REQ_V3NEW_PRODUCE_REQ_V9OLD_APIVER_RESP_V0NEW_APIVER_RESP_V0OLD_FETCH_RESP_V4NEW_FETCH_RESP_V4NEW_FETCH_RESP_V12OLD_PRODUCE_RESP_V3NEW_PRODUCE_RESP_V3NEW_PRODUCE_RESP_V9encodeENCODED_APIVER_RESP_V0ENCODED_FETCH_RESP_V4ENCODED_PRODUCE_RESP_V3ENCODED_FETCH_RESP_V12ENCODED_PRODUCE_RESP_V9rm   ru   r   r   r   r   __name__Runnerrunnerr   r   r   r   r   r   r   r   old_funcr   new_funcbench_time_funcargvr   r   r   r   <module>   s  
				
		






-
(
