o
    `j                     @   sV   d dl Z d dlmZ d dlmZmZmZ d dlmZmZ d dl	m
Z
 G dd dZdS )    N)defaultdict)datetime	timedeltatimezone)
OffsetSpecOffsetTimestamp)TopicPartitionc                   @   sh   e Zd ZdZ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S )ResetGroupOffsetszreset-offsetsz,Reset committed offsets for a consumer groupc              	   C   s   |j ddtdd |j ddtddg d	d
 |jdd}|j ddtdd |j dtddd |j dtddd |j dtddd |j dtddd |j dddd d! d S )"Nz-gz
--group-idT)typerequiredz-pz--partitionappend
partitionszTOPIC:PARTITION pair (repeatable). Scopes the reset to these partitions. If omitted, the reset applies to every partition currently committed by the group.)r
   actiondestdefaulthelp)r   z-sz--specznSpec may be one of earliest, latest, max-timestamp, earliest-local, latest-tiered, or a millisecond timestamp.)r
   r   z--to-offset	to_offsetzVReset all in-scope partitions to this explicit offset (clamped to [earliest, latest]).)r
   r   r   z
--shift-byshift_byzShift each in-scope committed offset by N positions (may be negative); clamped to [earliest, latest]. Requires a current commit for each in-scope partition.z--by-durationby_durationzUReset to the offset at (now - DURATION). DURATION is ISO-8601, e.g. P7D, PT1H, PT30M.z--to-datetimeto_datetimez]Reset to the offset at the given ISO-8601 datetime (UTC assumed if no tz offset is provided).z--to-current
store_true
to_currentzRe-commit the current committed offsets, clamped to [earliest, latest]. Useful to heal out-of-range offsets. Requires a current commit for each in-scope partition.)r   r   r   )add_argumentstradd_mutually_exclusive_groupint)clsparsermode r   b/home/djax/ivt_ai_plugin/venv/lib/python3.10/site-packages/kafka/cli/admin/groups/reset_offsets.pyadd_arguments   s>   
zResetGroupOffsets.add_argumentsc           
      C   s   | |jg}||j d }|dvrtd|j d| d| ||}||j|}tt}| D ]\}}	|	d j|	d< |	||j	 |j
< q3t|S )Ngroup_state)EmptyDeadzGroup z is z, expecting Empty or Dead!error)describe_groupsgroup_idRuntimeError_build_targetsreset_group_offsetsr   dictitems__name__topic	partition)
r   clientargsgroupstatetargetsresultoutputtpresr   r   r    command2   s   zResetGroupOffsets.commandc                    s   j rfdd j D nd } jd ur1|d ur|nt| j} jfdd|D S  jd urM|d ur<|nt| j} fdd|D S  jr} j}t	t
ttj|  d |d url|nt| j}fdd|D S  jr j}t	t
| d |d ur|nt| j}fdd|D S | j|d ur|nt}fd	d|D }|rtd
| d jd urه fdd|D S  jrfdd|D S td)Nc                    s   g | ]}  |qS r   )	_parse_tp).0p)r   r   r    
<listcomp>C   s    z4ResetGroupOffsets._build_targets.<locals>.<listcomp>c                       i | ]}| qS r   r   r;   r7   )specr   r    
<dictcomp>H       z4ResetGroupOffsets._build_targets.<locals>.<dictcomp>c                    s   i | ]}| j qS r   )r   r?   )r1   r   r    rA   L   s    i  c                    r>   r   r   r?   tsr   r    rA   R   rB   c                    r>   r   r   r?   rC   r   r    rA   X   rB   c                    s   g | ]}| vr|qS r   r   r?   currentr   r    r=   \       zNo committed offset for zK; --shift-by / --to-current require a current commit per in-scope partitionc                    s   i | ]}|| j  j qS r   )offsetr   r?   )r1   rF   r   r    rA   c   s    c                    s   i | ]}| | j qS r   )rH   r?   rE   r   r    rA   e   rG   zNo reset mode selected)r   r@   listlist_group_offsetsr'   _parse_specr   r   _parse_durationr   r   r   nowr   utc	timestampr   _parse_datetime
ValueErrorr   r   )r   r0   r1   explicit_scopescopedeltadtmissingr   )r1   r   rF   r@   rD   r    r)   A   s<   

 

z ResetGroupOffsets._build_targetsc                 C   s   | dd\}}t|t|S )N:   )rsplitr   r   )r   entryr.   r/   r   r   r    r:   h   s   zResetGroupOffsets._parse_tpc                 C   s,   zt |W S  ty   tt| Y S w )N)r   
build_fromrQ   r   r   )r   spec_strr   r   r    rK   m   s
   zResetGroupOffsets._parse_specc                 C   s~   t d|}|rt| std| | \}}}}t|r$t|nd|r+t|nd|r2t|nd|r;t|dS ddS )Nz?^P(?:(\d+)D)?(?:T(?:(\d+)H)?(?:(\d+)M)?(?:(\d+(?:\.\d+)?)S)?)?$zInvalid ISO-8601 duration: r   )dayshoursminutesseconds)rematchanygroupsrQ   r   r   float)r   durationmr]   r^   minssecsr   r   r    rL   t   s   
z!ResetGroupOffsets._parse_durationc                 C   sH   zt |}W n ty   td| w |jd u r"|jtjd}|S )NzInvalid ISO-8601 datetime: )tzinfo)r   fromisoformatrQ   rj   replacer   rN   )r   dt_strrU   r   r   r    rP      s   
z!ResetGroupOffsets._parse_datetimeN)r-   
__module____qualname__COMMANDHELPclassmethodr!   r9   r)   r:   rK   rL   rP   r   r   r   r    r	   	   s"    
$

&


r	   )ra   collectionsr   r   r   r   kafka.adminr   r   kafka.structsr   r	   r   r   r   r    <module>   s    