§
    ±½jË)  ã                  óÆ   — d Z ddlmZ ddlZddlZddlZddlmZmZm	Z	m
Z
mZmZmZmZ ddlmZ eegef         Z G d„ d¦  «        Z G d„ d	¦  «        Z G d
„ d¦  «        ZdS )z¥
Distributed Chat Cluster Router with Sharded Ring Buffers for 1,000,000 CCU.
Supports Hybrid Redis Cluster Pub/Sub bridging and In-Memory high-throughput sharding.
é    )ÚannotationsN)ÚAnyÚCallableÚ	CoroutineÚDictÚListÚOptionalÚSetÚTuple)ÚChatMessageDTOc                  ó\   — e Zd ZdZddd„Zd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 )$ÚShardedChannelRegistryz‡
    Sharded subscriber registry that minimizes lock contention when
    managing up to 1,000,000 concurrent client subscriptions.
    é@   Ú
num_shardsÚintÚreturnÚNonec                óÈ   — || _         d„ t          |¦  «        D ¦   «         | _        d„ t          |¦  «        D ¦   «         | _        d„ t          |¦  «        D ¦   «         | _        d S )Nc                ó   — g | ]}i ‘ŒS © r   ©Ú.0Ú_s     ú8C:\Projects\FreeExile\server\chat\chat_cluster_router.pyú
<listcomp>z3ShardedChannelRegistry.__init__.<locals>.<listcomp>   ó   € ÐFeÐFeÐFeÈaÀrÐFeÐFeÐFeó    c                ó   — g | ]}i ‘ŒS r   r   r   s     r   r   z3ShardedChannelRegistry.__init__.<locals>.<listcomp>   r   r   c                ó   — g | ]}i ‘ŒS r   r   r   s     r   r   z3ShardedChannelRegistry.__init__.<locals>.<listcomp>   s   € ÐGfÐGfÐGfÈqÈÐGfÐGfÐGfr   )r   ÚrangeÚshardsÚcached_syncÚcached_async)Úselfr   s     r   Ú__init__zShardedChannelRegistry.__init__   sn   € Ø$ˆŒàFeÐFeÕSXÐYcÑSdÔSdÐFeÑFeÔFeˆŒàFeÐFeÕSXÐYcÑSdÔSdÐFeÑFeÔFeˆÔØGfÐGfÕTYÐZdÑTeÔTeÐGfÑGfÔGfˆÔÐÐr   Úchannel_keyÚstrc                ó0   — t          |¦  «        | j        z  S ©N)Úhashr   ©r$   r&   s     r   Ú_get_shard_indexz'ShardedChannelRegistry._get_shard_index!   s   € Ý�KÑ Ô  4¤?Ñ2Ð2r   Ú	shard_idxc                ó  — g }g }| j         |         |                              ¦   «         D ]A}t          j        |¦  «        r|                     |¦  «         Œ,|                     |¦  «         ŒB|| j        |         |<   || j        |         |<   d S r)   )r!   ÚvaluesÚinspectÚiscoroutinefunctionÚappendr"   r#   )r$   r-   r&   Ú	sync_listÚ
async_listÚcbs         r   Ú_rebuild_cachez%ShardedChannelRegistry._rebuild_cache$   sž   € Ø.0ˆ	Ø/1ˆ
Ø”+˜iÔ(¨Ô5×<Ò<Ñ>Ô>ð 	%ð 	%ˆBÝÔ*¨2Ñ.Ô.ð %Ø×!Ò! "Ñ%Ô%Ð%Ð%à× Ò  Ñ$Ô$Ð$Ð$Ø3<ˆÔ˜Ô# KÑ0Ø4>ˆÔ˜)Ô$ [Ñ1Ð1Ð1r   Ú	client_idÚcallbackÚSubscriberCallbackc                óž   — |                       |¦  «        }| j        |         }||vri ||<   |||         |<   |                      ||¦  «         d S r)   )r,   r!   r6   )r$   r&   r7   r8   r-   Úshards         r   Ú	subscribez ShardedChannelRegistry.subscribe/   sa   € Ø×)Ò)¨+Ñ6Ô6ˆ	Ø”˜IÔ&ˆØ˜eÐ#Ð#Ø!#ˆE�+ÑØ(0ˆˆkÔ˜9Ñ%Ø×Ò˜I {Ñ3Ô3Ð3Ð3Ð3r   c                óX  — |                       |¦  «        }| j        |         }||vrd S ||                              |d ¦  «         ||         sG||= | j        |                              |d ¦  «         | j        |                              |d ¦  «         d S |                      ||¦  «         d S r)   )r,   r!   Úpopr"   r#   r6   )r$   r&   r7   r-   r;   s        r   Úunsubscribez"ShardedChannelRegistry.unsubscribe7   s¿   € Ø×)Ò)¨+Ñ6Ô6ˆ	Ø”˜IÔ&ˆØ˜eÐ#Ð#ØˆFØˆkÔ×Ò˜y¨$Ñ/Ô/Ð/Ø�[Ô!ð 	8Ø�kÐ"ØÔ˜YÔ'×+Ò+¨K¸Ñ>Ô>Ð>ØÔ˜iÔ(×,Ò,¨[¸$Ñ?Ô?Ð?Ð?Ð?à×Ò 	¨;Ñ7Ô7Ð7Ð7Ð7r   ú9Tuple[List[SubscriberCallback], List[SubscriberCallback]]c                ó°   — |                       |¦  «        }| j        |                              |g ¦  «        | j        |                              |g ¦  «        fS r)   ©r,   r"   Úgetr#   )r$   r&   r-   s      r   Úget_partitioned_callbacksz0ShardedChannelRegistry.get_partitioned_callbacksD   sU   € Ø×)Ò)¨+Ñ6Ô6ˆ	àÔ˜YÔ'×+Ò+¨K¸Ñ<Ô<ØÔ˜iÔ(×,Ò,¨[¸"Ñ=Ô=ð
ð 	
r   úList[SubscriberCallback]c                óº   — |                       |¦  «        }| j        |                              |g ¦  «        }| j        |                              |g ¦  «        }||z   S r)   rB   )r$   r&   r-   Úsync_cbsÚ	async_cbss        r   Úget_callbacksz$ShardedChannelRegistry.get_callbacksK   sZ   € Ø×)Ò)¨+Ñ6Ô6ˆ	ØÔ# IÔ.×2Ò2°;ÀÑCÔCˆØÔ% iÔ0×4Ò4°[À"ÑEÔEˆ	Ø˜)Ñ#Ð#r   ú$List[Tuple[int, SubscriberCallback]]c                óº   — |                       |¦  «        }| j        |         }|                     |¦  «        }|r!t          |                     ¦   «         ¦  «        ng S r)   )r,   r!   rC   ÚlistÚitems)r$   r&   r-   r;   Úsubscriberss        r   Úget_subscribersz&ShardedChannelRegistry.get_subscribersQ   sV   € Ø×)Ò)¨+Ñ6Ô6ˆ	Ø”˜IÔ&ˆØ—i’i Ñ,Ô,ˆØ,7Ð?�t�K×%Ò%Ñ'Ô'Ñ(Ô(Ð(¸RÐ?r   c                óp   — d}| j         D ]+}|                     ¦   «         D ]}|t          |¦  «        z  }ŒŒ,|S )Nr   )r!   r/   Úlen)r$   Úcountr;   Úsubss       r   Útotal_subscribers_countz.ShardedChannelRegistry.total_subscribers_countW   sK   € ØˆØ”[ð 	#ð 	#ˆEØŸš™œð #ð #�Ø�˜T™œÑ"��ð#àˆr   N)r   )r   r   r   r   )r&   r'   r   r   )r-   r   r&   r'   r   r   )r&   r'   r7   r   r8   r9   r   r   )r&   r'   r7   r   r   r   )r&   r'   r   r@   )r&   r'   r   rE   )r&   r'   r   rJ   )r   r   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r%   r,   r6   r<   r?   rD   rI   rO   rT   r   r   r   r   r      sà   € € € € € ðð ð
gð gð gð gð gð3ð 3ð 3ð 3ð	?ð 	?ð 	?ð 	?ð4ð 4ð 4ð 4ð8ð 8ð 8ð 8ð
ð 
ð 
ð 
ð$ð $ð $ð $ð@ð @ð @ð @ðð ð ð ð ð r   r   c                  óP   — e Zd ZdZ	 	 ddd	„Zedd„¦   «         Zdd„Zdd„Zdd„Z	dS )ÚRedisShardedPubSubBridgezå
    Redis 7 Sharded Pub/Sub bridge interface (SPUBLISH / SSUBSCRIBE).
    Routes messages using Redis cluster hash tags {...} to confine pub/sub
    traffic to the slot-owning node, preventing cluster-wide broadcast storms.
    NÚpublish_handlerú8Optional[Callable[[str, str], Coroutine[Any, Any, int]]]Úsubscribe_handlerúHOptional[Callable[[str, SubscriberCallback], Coroutine[Any, Any, None]]]r   r   c                ó0   — || _         || _        g | _        d S r)   )Ú_publish_handlerÚ_subscribe_handlerÚpublished_messages)r$   r[   r]   s      r   r%   z!RedisShardedPubSubBridge.__init__f   s"   € ð
 !0ˆÔØ"3ˆÔØ9;ˆÔÐÐr   r&   r'   c                óf   — |                       d¦  «        r|                      d¦  «        r| S d| › d�S )zª
        Wraps channel key into Redis Cluster hash tag format.
        e.g. 'world' -> '{world}', 'zone:barrow' -> '{zone:barrow}', 'guild:101' -> '{guild:101}'.
        Ú{Ú})Ú
startswithÚendswith)r&   s    r   Úto_sharded_channelz+RedisShardedPubSubBridge.to_sharded_channelo   sF   € ð ×!Ò! #Ñ&Ô&ð 	¨;×+?Ò+?ÀÑ+DÔ+Dð 	ØÐØ#�KÐ#Ð#Ð#Ð#r   ÚmessageúChatMessageDTO | strr   c              ƒ  óî   K  — |                       |¦  «        }t          |t          ¦  «        r|n|j        }| j                             ||f¦  «         | j        �|                      ||¦  «        ƒ d{V —†S dS )z2Dispatches SPUBLISH command to Redis 7 shard slot.Né   )rh   Ú
isinstancer'   Úraw_contentrb   r2   r`   )r$   r&   ri   Ú
sharded_chÚpayloads        r   Úspublishz!RedisShardedPubSubBridge.spublishy   s†   è è € à×,Ò,¨[Ñ9Ô9ˆ
Ý'¨µÑ5Ô5ÐN�'�'¸7Ô;NˆØÔ×&Ò&¨
°GÐ'<Ñ=Ô=Ð=ØÔ Ð,Ø×.Ò.¨z¸7ÑCÔCÐCÐCÐCÐCÐCÐCÐCØˆqr   r8   r9   c              ƒ  ó~   K  — |                       |¦  «        }| j        �|                      ||¦  «        ƒ d{V —† dS dS )z2Registers SSUBSCRIBE on the Redis 7 shard channel.N)rh   ra   )r$   r&   r8   ro   s       r   Ú
ssubscribez#RedisShardedPubSubBridge.ssubscribe‚   sX   è è € à×,Ò,¨[Ñ9Ô9ˆ
ØÔ"Ð.Ø×)Ò)¨*°hÑ?Ô?Ð?Ð?Ð?Ð?Ð?Ð?Ð?Ð?Ð?ð /Ð.r   c              ƒ  ó
   K  — dS )z Removes SSUBSCRIBE registration.Nr   r+   s     r   Úsunsubscribez%RedisShardedPubSubBridge.sunsubscribeˆ   s   è è € àˆr   )NN)r[   r\   r]   r^   r   r   )r&   r'   r   r'   )r&   r'   ri   rj   r   r   )r&   r'   r8   r9   r   r   )r&   r'   r   r   )
rU   rV   rW   rX   r%   Ústaticmethodrh   rq   rs   ru   r   r   r   rZ   rZ   _   s¡   € € € € € ðð ð UYØfjð<ð <ð <ð <ð <ð ð$ð $ð $ñ „\ð$ðð ð ð ð@ð @ð @ð @ðð ð ð ð ð r   rZ   c                  óJ   — e Zd ZdZ	 	 	 ddd„Zd d„Zd!d„Zd"d„Zd#d„Zd$d„Z	dS )%ÚChatClusterRouterz¤
    Core distributed fan-out engine capable of horizontal cluster scaling.
    Combines local sharded non-blocking delivery with Redis Pub/Sub cluster bridge.
    r   éè  Nr   r   Úbatch_fanout_sizeÚredis_bridgeú"Optional[RedisShardedPubSubBridge]r   r   c                óh   — t          |¬¦  «        | _        || _        || _        g | _        d| _        d S )N)r   r   )r   Úregistryrz   r{   Ú_dispatch_latencies_msÚ_total_messages_routed)r$   r   rz   r{   s       r   r%   zChatClusterRouter.__init__“   s<   € õ /¸*ÐEÑEÔEˆŒØ!2ˆÔØ(ˆÔØ35ˆÔ#Ø+,ˆÔ#Ð#Ð#r   ÚbridgerZ   c                ó   — || _         dS )z4Attaches an external Redis 7 Sharded Pub/Sub bridge.N)r{   )r$   r�   s     r   Úattach_redis_bridgez%ChatClusterRouter.attach_redis_bridgeŸ   s   € à"ˆÔÐÐr   r5   r9   ri   r   ú/Tuple[bool, Optional[Coroutine[Any, Any, Any]]]c                ó¶   — 	 t          j        |¦  «        r ||¦  «        }d|fS  ||¦  «        }t          j        |¦  «        rd|fS dS # t          $ r Y dS w xY w)a1  
        Invokes subscriber callback safely:
        If cb is a coroutine function (async def), returns coroutine for batch gathering.
        If cb is a regular synchronous callable, invokes it directly without coroutine overhead.
        Eliminates unneeded coroutines and TypeError exceptions.
        T)TN)FN)r0   r1   ÚasyncioÚiscoroutineÚ	Exception)r$   r5   ri   ÚcoroÚress        r   Ú_safe_invokezChatClusterRouter._safe_invoke£   s‚   € ð		ÝÔ*¨2Ñ.Ô.ð "Ø�r˜'‘{”{�Ø˜T�zÐ!Ø�"�W‘+”+ˆCÝÔ" 3Ñ'Ô'ð !Ø˜S�yÐ Ø�:øÝð 	ð 	ð 	Ø�;�;ð	øøøs   ‚"A
 ¥"A
 Á

AÁAr&   r'   c              ƒ  óN  K  — t          j        ¦   «         }| j                             |¦  «        \  }}d}|D ]O}	  ||¦  «        }|�*t	          j        |¦  «        r|                     |¦  «         n|dz  }Œ@# t          $ r Y ŒLw xY w|r¤g }	|D ]j}	  ||¦  «        }
	 |
                     d¦  «         |	                     |
¦  «         n # t          $ r |dz  }Y nt          $ r Y nw xY wŒ[# t          $ r Y Œgw xY w|	r3t	          j
        |	ddiŽƒ d{V —†}|t          d„ |D ¦   «         ¦  «        z  }| j        �3	 | j                             ||¦  «        ƒ d{V —† n# t          $ r Y nw xY wt          j        ¦   «         |z
  dz  }|                      |¦  «         | xj        dz  c_        |S )zŸ
        Dispatches a chat message to all subscribers of a channel key.
        Executes fan-out in parallel batches to achieve ultra-low p99 latency.
        r   Nrl   Úreturn_exceptionsTc              3  óD   K  — | ]}t          |t          ¦  «        °d V — ŒdS )rl   N)rm   rˆ   )r   Úrs     r   ú	<genexpr>z9ChatClusterRouter.broadcast_to_channel.<locals>.<genexpr>æ   s1   è è € Ð&ZÐ&Z¨QÅÈAÍyÑAYÔAYÐ&Z qÐ&ZÐ&ZÐ&ZÐ&ZÐ&ZÐ&Zr   g     @�@)ÚtimeÚperf_counterr~   rD   r†   r‡   r2   rˆ   ÚsendÚStopIterationÚgatherÚsumr{   rq   Ú_record_latencyr€   )r$   r&   ri   Út_startrG   rH   Údelivered_countr5   rŠ   Úpending_coroutinesr‰   ÚresultsÚ
elapsed_mss                r   Úbroadcast_to_channelz&ChatClusterRouter.broadcast_to_channel¹   sl  è è € õ Ô#Ñ%Ô%ˆØ"œm×EÒEÀkÑRÔRÑˆ�)Øˆð ð 	ð 	ˆBðØ�b˜‘k”k�Ø�?¥wÔ':¸3Ñ'?Ô'?�?Ø×$Ò$ RÑ(Ô(Ð(Ð(à# qÑ(�OøøÝð ð ð Ø�ðøøøð ð 	[ØACÐØð ð �ðØ˜2˜g™;œ;�Dð	8àŸ	š	 $™œ˜ð +×1Ò1°$Ñ7Ô7Ð7Ð7øõ )ð -ð -ð -Ø'¨1Ñ,˜˜˜Ý$ð ð ð Ø˜ðøøøøøõ
 !ð ð ð Ø�Dðøøøð "ð [Ý '¤Ð0BÐ [ÐVZÐ [Ð [Ð[Ð[Ð[Ð[Ð[Ð[�Ø¥3Ð&ZÐ&Z°'Ð&ZÑ&ZÔ&ZÑ#ZÔ#ZÑZ�ð ÔÐ(ðØÔ'×0Ò0°¸gÑFÔFÐFÐFÐFÐFÐFÐFÐFÐFøÝð ð ð Ø�ðøøøõ Ô'Ñ)Ô)¨GÑ3°vÑ=ˆ
Ø×Ò˜ZÑ(Ô(Ð(ØÐ#Ô# qÑ(Ð#Ô#ØÐse   º<A7Á7
BÂBÂC(ÂCÂ1C(ÃC$ÃC(Ã	C$Ã!C(Ã#C$Ã$C(Ã(
C5Ã4C5Ä5!E Å
E$Å#E$Ú
latency_msÚfloatc                óž   — t          | j        ¦  «        dk    r| j                             d¦  «         | j                             |¦  «         d S )Nry   r   )rQ   r   r>   r2   )r$   rž   s     r   r—   z!ChatClusterRouter._record_latencyô   sL   € ÝˆtÔ*Ñ+Ô+¨tÒ3Ð3ØÔ'×+Ò+¨AÑ.Ô.Ð.ØÔ#×*Ò*¨:Ñ6Ô6Ð6Ð6Ð6r   úDict[str, float]c                óæ  — | j         sdddddœS t          | j         ¦  «        }t          |¦  «        }|t          |dz  ¦  «                 }|t	          t          |dz  ¦  «        |dz
  ¦  «                 }|t	          t          |dz  ¦  «        |dz
  ¦  «                 }t          |¦  «        |z  }t          |d¦  «        t          |d¦  «        t          |d¦  «        t          |d¦  «        | j        dœS )	z;Calculates p50, p95, p99 latency metrics for observability.g        )Úp50Úp95Úp99Úavgg      à?gffffffî?rl   g®Gáz®ï?é   )r£   r¤   r¥   r¦   Útotal_routed)r   ÚsortedrQ   r   Úminr–   Úroundr€   )r$   Úsorted_latenciesÚnr£   r¤   r¥   r¦   s          r   Úget_latency_statsz#ChatClusterRouter.get_latency_statsù   sí   € àÔ*ð 	DØ s°3¸sÐCÐCÐCå! $Ô"=Ñ>Ô>ÐÝÐ Ñ!Ô!ˆØ�s 1 t¡8™}œ}Ô-ˆØ�s¥3 q¨4¡x¡=¤=°!°a±%Ñ8Ô8Ô9ˆØ�s¥3 q¨4¡x¡=¤=°!°a±%Ñ8Ô8Ô9ˆÝÐ"Ñ#Ô# aÑ'ˆõ ˜˜a‘=”=Ý˜˜a‘=”=Ý˜˜a‘=”=Ý˜˜a‘=”=Ø Ô7ð
ð 
ð 	
r   )r   ry   N)r   r   rz   r   r{   r|   r   r   )r�   rZ   r   r   )r5   r9   ri   r   r   r„   )r&   r'   ri   r   r   r   )rž   rŸ   r   r   )r   r¡   )
rU   rV   rW   rX   r%   rƒ   r‹   r�   r—   r®   r   r   r   rx   rx   �   s§   € € € € € ðð ð Ø!%Ø;?ð	
-ð 
-ð 
-ð 
-ð 
-ð#ð #ð #ð #ðð ð ð ð,9ð 9ð 9ð 9ðv7ð 7ð 7ð 7ð

ð 
ð 
ð 
ð 
ð 
r   rx   )rX   Ú
__future__r   r†   r0   r‘   Útypingr   r   r   r   r   r	   r
   r   Úserver.chat.chat_typesr   r9   r   rZ   rx   r   r   r   ú<module>r²      s7  ððð ð
 #Ð "Ð "Ð "Ð "Ð "à €€€Ø €€€Ø €€€Ø MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ MÐ Mà 1Ð 1Ð 1Ð 1Ð 1Ð 1ð ˜~Ð.°Ð3Ô4Ð ðIð Ið Ið Ið Iñ Iô Ið IðX+ð +ð +ð +ð +ñ +ô +ð +ð\~
ð ~
ð ~
ð ~
ð ~
ñ ~
ô ~
ð ~
ð ~
ð ~
r   