
    `gjp8                     
   d Z ddlZddlZddlZddlZddlZddlmZ ddlm	Z	m
Z
 ddlmZ e	rddlmZ  ej                  e      Z G d d      Z G d	 d
      Z ej(                  dej*                        Z ej(                  dej*                        Z ej(                  dej*                        Z ej(                  dej*                        Z ej(                  d      Z ej(                  d      Z ej(                  dej8                        Z ej(                  d      Z ej(                  d      Zde de fdZ! G d d      Z"de de fdZ# ej(                  d      Z$de de%fdZ&de de'e    fdZ(de'e    de fd Z)de de fd!Z*y)"zShared helper classes for gateway platform adapters.

Extracts common patterns that were duplicated across 5-7 adapters:
message deduplication, text batch aggregation, markdown stripping,
and thread participation tracking.
    N)Path)TYPE_CHECKINGDict)atomic_json_write)MessageEventc                   X    e Zd ZdZddedefdZdedefdZ	dedefdZ
dedd	fd
Zd Zy	)MessageDeduplicatora|  TTL-based message deduplication cache.

    Replaces the identical ``_seen_messages`` / ``_is_duplicate()`` pattern
    previously duplicated in discord, slack, dingtalk, wecom, weixin,
    mattermost, and feishu adapters.

    Usage::

        self._dedup = MessageDeduplicator()

        # In message handler:
        if self._dedup.is_duplicate(msg_id):
            return
    max_sizettl_secondsc                 .    i | _         || _        || _        y N)_seen	_max_size_ttl)selfr
   r   s      L/root/.hermes/venv/lib/python3.12/site-packages/gateway/platforms/helpers.py__init__zMessageDeduplicator.__init__+   s    ')
!	    msg_idreturnc                 r   |syt        j                          }|| j                  v r-|| j                  |   z
  | j                  k  ry| j                  |= || j                  |<   t        | j                        | j                  kD  r|| j                  z
  }| j                  j                         D ci c]  \  }}||kD  s|| c}}| _        t        | j                        | j                  kD  rDt        | j                  j                         d       | j                   d }t        |      | _        yc c}}w )z?Return True if *msg_id* was already seen within the TTL window.FTc                     | d   S )N    )items    r   <lambda>z2MessageDeduplicator.is_duplicate.<locals>.<lambda>D   s
    T!W r   )keyN)timer   r   lenr   itemssorteddict)r   r   nowcutoffkvnewests          r   is_duplicatez MessageDeduplicator.is_duplicate0   s    iikTZZTZZ''$))3

6" 

6tzz?T^^+499_F+/::+;+;+=L41aV!Q$LDJ4::/  JJ$$&, >>/"$ "&\
 Ms   0D3>D3c                     |sy| j                   j                  |      }|yt        j                         |z
  | j                  k  ry| j                   |= y)zBReturn whether *msg_id* is live in the cache without inserting it.FT)r   getr   r   )r   r   seen_ats      r   containszMessageDeduplicator.containsI   sK    **..(?99; 499,JJvr   Nc                 <    | j                   j                  |d       y)z<Release a claimed message ID after cancelled/failed handoff.N)r   pop)r   r   s     r   discardzMessageDeduplicator.discardU   s    

vt$r   c                 8    | j                   j                          y)zClear all tracked messages.N)r   clearr   s    r   r1   zMessageDeduplicator.clearY   s    

r   )i  i,  )__name__
__module____qualname____doc__intfloatr   strboolr(   r,   r/   r1   r   r   r   r	   r	      sX       %  
3 4 2
s 
t 
%c %d %r   r	   c                   f    e Zd ZdZdddddededefd	Zd
efdZddde	d
dfdZ
de	d
dfdZddZy)TextBatchAggregatora@  Aggregates rapid-fire text events into single messages.

    Replaces the ``_enqueue_text_event`` / ``_flush_text_batch`` pattern
    previously duplicated in telegram, discord, matrix, wecom, and feishu.

    Usage::

        self._text_batcher = TextBatchAggregator(
            handler=self._message_handler,
            batch_delay=0.6,
            split_threshold=1900,
        )

        # In message dispatch:
        if msg_type == MessageType.TEXT and self._text_batcher.is_enabled():
            self._text_batcher.enqueue(event, session_key)
            return
    g333333?g       @i  )batch_delaysplit_delaysplit_thresholdr=   r>   r?   c                X    || _         || _        || _        || _        i | _        i | _        y r   )_handler_batch_delay_split_delay_split_threshold_pending_pending_tasks)r   handlerr=   r>   r?   s        r   r   zTextBatchAggregator.__init__u   s2      '' /3579r   r   c                      | j                   dkD  S )z.Return True if batching is active (delay > 0).r   )rB   r2   s    r   
is_enabledzTextBatchAggregator.is_enabled   s      1$$r   eventr   r   Nc                    t        |j                  xs d      }| j                  j                  |      }|s||_        || j                  |<   n'|j                   d|j                   |_        ||_        | j
                  j                  |      }|r |j                         s|j                          t        j                  | j                  |            | j
                  |<   y)z+Add *event* to the pending batch for *key*. 
N)r   textrE   r*   _last_chunk_lenrF   donecancelasynciocreate_task_flush)r   rJ   r   	chunk_lenexistingpriors         r   enqueuezTextBatchAggregator.enqueue   s    

(b)	==$$S)$-E!!&DMM#'}}oR

|<HM'0H$ ##'',LLN#*#6#6t{{37G#HC r   c                 X  K   | j                   j                  |      }| j                  j                  |      }|rt        |dd      nd}|| j                  k\  r| j
                  n| j                  }t        j                  |       d{    | j                  j                  |d      }|r	 | j                  |       d{    | j                   j                  |      |u r| j                   j                  |d       yy7 w7 A# t        $ r t        j                  d|       Y `w xY ww)z/Wait then dispatch the batched event for *key*.rO   r   Nz<[TextBatchAggregator] Error dispatching batched event for %s)rF   r*   rE   getattrrD   rC   rB   rR   sleepr.   rA   	Exceptionlogger	exception)r   r   current_taskpendinglast_lendelayrJ   s          r   rT   zTextBatchAggregator._flush   s    **..s3--##C(=D77$5q9! &.1F1F%F!!DL]L]mmE"""!!#t,fmmE*** ""3'<7##C. 8 	#
 + f  !_adefsH   BD*	D
"D*-D DD <D*D D'$D*&D''D*c                     | j                   j                         D ]#  }|j                         r|j                          % | j                   j	                          | j
                  j	                          y)zCancel all pending flush tasks.N)rF   valuesrP   rQ   r1   rE   )r   tasks     r   
cancel_allzTextBatchAggregator.cancel_all   sV    ''..0 	D99;	 	!!#r   r   N)r3   r4   r5   r6   r8   r7   r   r:   rI   r9   rX   rT   rf   r   r   r   r<   r<   a   sw    . ! #: 	:
 : :%D %I^ I# I$ I"/ / /(r   r<   z\*\*(.+?)\*\*z	\*(.+?)\*z \b__(?![\s_])(.+?)(?<![\s_])__\bz\b_(?![\s_])(.+?)(?<![\s_])_\bz```[a-zA-Z0-9_+-]*\n?z`(.+?)`z
^#{1,6}\s+z\[([^\]]+)\]\([^\)]+\)z\n{3,}rN   r   c                    t         j                  d|       } t        j                  d|       } t        j                  d|       } t        j                  d|       } t
        j                  d|       } t        j                  d|       } t        j                  d|       } t        j                  d|       } t        j                  d|       } | j                         S )zStrip markdown formatting for plain-text platforms (SMS, iMessage, etc.).

    Replaces the identical ``_strip_markdown()`` functions previously
    duplicated in sms.py, bluebubbles.py, and feishu.py.
    z\1rL   

)_RE_BOLDsub_RE_ITALIC_STAR_RE_BOLD_UNDER_RE_ITALIC_UNDER_RE_CODE_BLOCK_RE_INLINE_CODE_RE_HEADING_RE_LINK_RE_MULTI_NEWLINEstrip)rN   s    r   strip_markdownru      s     <<t$Dud+DeT*Dt,Db$'Dud+D??2t$D<<t$D  .D::<r   c                   t    e Zd ZdZdZddedefdZdefdZ	de
e   fdZdd
Zdedd	fdZdedefdZddZy	)ThreadParticipationTrackera  Persistent tracking of threads the bot has participated in.

    Replaces the identical ``_load/_save_participated_threads`` +
    ``_mark_thread_participated`` pattern previously duplicated in
    discord.py and matrix.py.

    Usage::

        self._threads = ThreadParticipationTracker("discord")

        # Check membership:
        if thread_id in self._threads:
            ...

        # Mark participation:
        self._threads.mark(thread_id)
      platform_namemax_trackedc                     || _         || _        | j                         D ci c]  }t        |      d  c}| _        y c c}w r   )	_platform_max_tracked_loadr9   _threads)r   ry   rz   	thread_ids       r   r   z#ThreadParticipationTracker.__init__   s<    &'26**,*
%.C	ND *
 *
s   =r   c                 <    ddl m}  |       | j                   dz  S )Nr   )get_hermes_homez_threads.json)hermes_constantsr   r|   )r   r   s     r   _state_pathz&ThreadParticipationTracker._state_path   s    4 dnn%5]#CCCr   c                    | j                         }|j                         rR	 t        j                  |j	                  d            }t        |t              r|D cg c]  }t        |       c}S 	 g S g S c c}w # t        $ r Y g S w xY w)Nzutf-8)encoding)	r   existsjsonloads	read_text
isinstancelistr9   r\   )r   pathdatar   s       r   r~   z ThreadParticipationTracker._load   s    !;;=zz$..'."BCdD)<@AyC	NAA * 	r	 B 	s#   9A: A5-A: 5A: :	BBNc                     | j                         }t        | j                        }t        |      | j                  kD  r*|| j                   d  }t
        j                  |      | _        t        ||d        y )N)indent)r   r   r   r   r}   r"   fromkeysr   )r   r   thread_lists      r   _savez ThreadParticipationTracker._save  sc    !4==){d///%t'8'8&8&9:K MM+6DM$D9r   r   c                 `    || j                   vr d| j                   |<   | j                          yy)z-Mark *thread_id* as participated and persist.N)r   r   r   r   s     r   markzThreadParticipationTracker.mark  s*    DMM)'+DMM)$JJL *r   c                     || j                   v S r   )r   r   s     r   __contains__z'ThreadParticipationTracker.__contains__  s    DMM))r   c                 8    | j                   j                          y r   )r   r1   r2   s    r   r1   z ThreadParticipationTracker.clear  s    r   )rx   rg   )r3   r4   r5   r6   _MAX_TRACKEDr9   r7   r   r   r   r   r~   r   r   r:   r   r1   r   r   r   rw   rw      so    $ L
c 
 
DT D	tCy 	:c d *c *d *r   rw   phonec                 |    | syt        |       dk  rt        |       dkD  r| dd dz   | dd z   S dS | dd dz   | dd z   S )	zRedact a phone number for logging, preserving country code and last 4.

    Replaces the identical ``_redact_phone()`` functions in signal.py,
    sms.py, and bluebubbles.py.
    z<none>      N   z****)r   )r   s    r   redact_phoner     s\     
5zQ25e*q.uRay6!E"#J.LfL!9vbc
**r   z0^\s*\|?\s*:?-+:?\s*(?:\|\s*:?-+:?\s*){1,}\|?\s*$linec                 D    | j                         }t        |      xr d|v S )z:Return True if *line* could plausibly be a table data row.|)rt   r:   )r   strippeds     r   is_table_rowr   7  s     zz|H>-cXo-r   c                     | j                         }|j                  d      r|dd }|j                  d      r|dd }|j                  d      D cg c]  }|j                          c}S c c}w )z0Split a GFM table row into stripped cell values.r   r   N)rt   
startswithendswithsplit)r   r   cells      r   split_markdown_table_rowr   =  sb    zz|H3AB<CR=%-^^C%89TDJJL999s   A*table_blockc                 f   t        |       dk  rdj                  |       S t        | d         }t        |      dk  rdj                  |       S t        |       dkD  rt        | d         ng }t        |      t        |      dz   k(  }g }t        | dd d      D ]  \  }}t        |      }|r|r
|d   r|d   nd| }|dd }	nt	        d	 |D        d|       }|}	t        |	      t        |      k  r+|	j                  d
gt        |      t        |	      z
  z         n%t        |	      t        |      kD  r|	dt        |       }	g }
t        ||	      D ]$  \  }}|s||k(  r|
j                  d| d|        & d| dg|
}|j                  dj                  |             
 dj                  |      S )u4  Render a detected GFM table as bold-heading + bullet groups.

    Uses the same alignment logic as Telegram's renderer: for non-row-label
    tables, ``data_cells = cells`` (the full row) and the bullet whose value
    duplicates the heading is skipped.  This keeps header→value alignment
    correct.
       rM   r   r   r   N)startzRow c              3   &   K   | ]	  }|s|  y wr   r   ).0r   s     r   	<genexpr>z&_render_table_block.<locals>.<genexpr>d  s     ;TdD;s   rL   u   • z: z**ri   )r   joinr   	enumeratenextextendzipappend)r   headersfirst_data_rowhas_row_label_colrendered_groupsindexrowcellsheading
data_cellsbulletsheadervaluegroup_liness                 r   _render_table_blockr   G  s    ;!yy%%&{1~6G
7|ayy%% {a 	!Q0 
 N+s7|a/??!#OABq9 7
s(-"'E!HeAhD.GqrJ;U;tE7^LGJz?S\)rdc'lS_&DEF_s7|+#Nc'l3J *5 	5MFE$')9NNT&E734	5
 G9B'2'2tyy56+7. ;;''r   c                    d| vsd| vr| S | j                  d      }g }d}d}|t        |      k  r.||   }|j                         }|j                  d      r| }|j	                  |       |dz  }O|r|j	                  |       |dz  }hd|v r|dz   t        |      k  rt
        j                  ||dz            r|||dz      g}|dz   }|t        |      k  rDt        ||         r6|j	                  ||          |dz  }|t        |      k  rt        ||         r6|j	                  t        |             |}|j	                  |       |dz  }|t        |      k  r.dj                  |      S )	zuRewrite GFM pipe tables into bold-heading + bullet groups.

    Tables inside fenced code blocks are left alone.
    r   -rM   Fr   z```r   r   )
r   r   lstripr   r   TABLE_SEPARATOR_REmatchr   r   r   )	rN   linesoutin_fenceir   r   r   js	            r   convert_table_to_bulletsr   x  s{   
 $#T/JJtECH	A
c%j.Qx;;=u%#|HJJtFAJJtFA 4KAE
""((q1u6q1u.KAAc%j.\%(%;""58,Q c%j.\%(%; JJ*;78A

4	Q; c%j.> 99S>r   )+r6   rR   r   loggingrer   pathlibr   typingr   r   utilsr   gateway.platforms.baser   	getLoggerr3   r]   r	   r<   compileDOTALLrj   rl   rm   rn   ro   rp   	MULTILINErq   rr   rs   r9   ru   rw   r   r   r:   r   r   r   r   r   r   r   r   <module>r      s      	   & #3			8	$@ @LR Rp 2::&		2"**\2995?K2::?K 45"**Z(bjj52::/0BJJy)   *= =F
+ 
+ 
+,  RZZ7 
.s .t .:3 :49 :.(T#Y .(3 .(b+3 +3 +r   