+
    Q(i,                    P  a  0 t $ R t^ RIHt ^ RIt^ RIt^ RIt^ RIt^ RIH	t	 ^ RI
H
t
Ht ^ RIHtHtHt ^ RIt]'       d   ^ RIHt ^ RIHt ^ RIHt ]P0                  ! ]4      tR	tR
tRt^t^tR R lt R R lt!R R lt"]#! 4       t$R]%R&   / t&R]%R&   R R lt'R R lt(R R lt)R# )a  Distributed notification queue for background task events (SEP-1686).

Enables distributed Docket workers to send MCP notifications to clients
without holding session references. Workers push to a Redis queue,
the MCP server process subscribes and forwards to the client's session.

Pattern: Fire-and-forward with retry
- One queue per session_id
- LPUSH/BRPOP for reliable ordered delivery
- Retry up to 3 times on delivery failure, then discard
- TTL-based expiration for stale messages

Note: Docket's execution.subscribe() handles task state/progress events via
Redis Pub/Sub. This module handles elicitation-specific notifications that
require reliable delivery (input_required prompts, cancel signals).
)annotationsN)suppress)datetimetimezone)TYPE_CHECKINGAnycast)Docket)ServerSession)FastMCPz"fastmcp:notifications:{session_id}z)fastmcp:notifications:{session_id}:activei,  c               (    V ^8  d   QhRRRRRRRR/# )	   
session_idstrnotificationdict[str, Any]docketr	   returnNone )formats   "h/Users/agent/.openclaw/workspace/venv/lib/python3.14/site-packages/fastmcp/server/tasks/notifications.py__annotate__r   0   s0     : :: : : 
	:    c           
     
  "   VP                  \        P                  V R7      4      p\        P                  ! RVR^ R\
        P                  ! \        P                  4      P                  4       /4      pVP                  4       ;_uu_4       GRj  xL
 pVP                  W44      G Rj  xL
  VP                  V\        4      G Rj  xL
  RRR4      GRj  xL
  R#  LM L6 L L  + GRj  xL 
 '       g   i     R# ; i5i)ab  Push notification to session's queue (called from Docket worker).

Used for elicitation-specific notifications (input_required, cancel)
that need reliable delivery across distributed processes.

Args:
    session_id: Target session's identifier
    notification: MCP notification dict (method, params, _meta)
    docket: Docket instance for Redis access
r   r   attemptenqueued_atN)keyNOTIFICATION_QUEUE_KEYr   jsondumpsr   nowr   utc	isoformatredislpushexpireNOTIFICATION_TTL_SECONDS)r   r   r   r   messager%   s   &&&   r   push_notificationr*   0   s      **+22j2I
JCjjLq8<<5??A	
G ||~~~kk#'''ll3 8999 ~~'9 ~~~sr   BDCDC&)C *C&C"C&DC$D C&"C&$D&D 	,C/-
D 	8D 	:	Dc          
     ,    V ^8  d   QhRRRRRRRRR	R
/# r   r   r   sessionr
   r   r	   fastmcpr   r   r   r   )r   s   "r   r   r   L   sA     V# V#V#V# V# 	V#
 
V#r   c           
       "   VP                  \        P                  V R7      4      pVP                  \        P                  V R7      4      p\        P                  RV 4         VP                  4       ;_uu_4       GRj  xL
 pVP                  VR\        ^,          R7      G Rj  xL
  \        \        VP                  V.\        R7      4      G Rj  xL
 pV'       g    RRR4      GRj  xL
  K  Vw  r\        P                  ! V	4      p
V
R,          pV
P                  R^ 4      p \        WWV4      G Rj  xL
  \        P                  R	V V^,           4       RRR4      GRj  xL
  EK   L L L L L=  \          d   pT\"        ^,
          8  dn   T^,           T
R&   \%        T4      T
R
&   TP'                  T\        P(                  ! T
4      4      G Rj  xL 
  \        P                  RT T^,           T4        Rp?L\        P+                  RT \"        T4        Rp?LRp?ii ; i L  + GRj  xL 
 '       g   i     EK  ; i  \,        P.                   d    \        P                  RT 4        R# \          dB   p\        P                  RY4       \,        P0                  ! ^4      G Rj  xL 
   Rp?EK`  Rp?ii ; i5i)aN  Subscribe to notification queue and forward to session.

Runs in the MCP server process. Bridges distributed workers to clients.

This loop:
1. Maintains a heartbeat (active subscriber marker for debugging)
2. Blocks on BRPOP waiting for notifications
3. Forwards notifications to the client's session
4. Retries failed deliveries, then discards (no dead-letter queue)

Args:
    session_id: Session identifier to subscribe to
    session: MCP ServerSession for sending notifications
    docket: Docket instance for Redis access
    fastmcp: FastMCP server instance (for elicitation relay)
r   z/Starting notification subscriber for session %sN1)ex)timeoutr   r   z1Delivered notification to session %s (attempt %d)
last_errorz5Requeued notification for session %s (attempt %d): %sz<Discarding notification for session %s after %d attempts: %sz0Notification subscriber cancelled for session %sz0Notification subscriber error for session %s: %s)r   r   r   NOTIFICATION_ACTIVE_KEYloggerdebugr%   setSUBSCRIBER_TIMEOUT_SECONDSr   r   brpopr    loadsget_send_mcp_notification	ExceptionMAX_DELIVERY_ATTEMPTSr   r&   r!   warningasyncioCancelledErrorsleep)r   r-   r   r.   	queue_key
active_keyr%   result_message_bytesr)   notification_dictr   
send_errores   &&&&           r   notification_subscriber_looprK   L   sB    , 

188J8OPI3::j:QRJ
LLBJO
:	#||~~~ii
C4NQR4RiSSS  $i[:TU    &~~ $* **]3$+N$;!!++i30J   LLK"!- &~~S && ! !6!::-4q[	*03J-#kk)TZZ5HIIIS&#aK&	  Z&1&	 7 &~~~b %% 	LLKZX 	#LLBJ --""""		#sE  A"K%I E2I #H1(E4)-H1E6H1#H1$I /E80I 4K65H1,E<<E:="E<I *H/+I /K2I 4H16H18I :E<<H,AH'G
$H'H1H'"H1'H,,H1/I 1I	7H:8
I	I	I 	KI +K:K=KK/K6J97K<KKKc               0    V ^8  d   QhRRRRRRRRR	R
RR/# )r   r-   r
   rH   r   r   r   r   r	   r.   r   r   r   r   )r   s   "r   r   r      sD     8> 8>8>%8> 8> 	8>
 8> 
8>r   c           
     h  "   VP                  RR4      pVR8w  d   \        RV 24      h\        P                  P                  P                  RRRVP                  R/ 4      RVP                  R4      /4      p\        P                  P                  V4      pV P                  V4      G Rj  xL
  VP                  R/ 4      pVP                  R4      R8X  d   VP                  R/ 4      p	V	P                  R	/ 4      p
V
P                  R
4      pV'       d   VP                  R4      pV'       g   \        P                  R4       R# ^ RI
Hp \        P                  ! V! WWV4      RVR,           2R7      p\        P                  V4       VP!                  \        P"                  4       R# R# R#  EL5i)a  Reconstruct MCP notification from dict and send to session.

For input_required notifications with elicitation metadata, also sends
a standard elicitation/create request to the client and relays the
response back to the worker via Redis.

Args:
    session: MCP ServerSession
    notification_dict: Notification as dict (method, params, _meta)
    session_id: Session identifier (for elicitation relay)
    docket: Docket instance (for notification delivery)
    fastmcp: FastMCP server instance (for elicitation relay)
methodznotifications/tasks/statusz0Unsupported notification method for subscriber: params_metaNstatusinput_requiredz$io.modelcontextprotocol/related-taskelicitationtaskIdz:input_required notification missing taskId, skipping relay)relay_elicitationzelicitation-relay-N   Nname)r;   
ValueErrormcptypesTaskStatusNotificationmodel_validateServerNotificationsend_notificationr5   r?    fastmcp.server.tasks.elicitationrU   r@   create_task_background_tasksaddadd_done_callbackdiscard)r-   rH   r   r   r.   rN   r   server_notificationrO   metarelated_taskrS   task_idrU   tasks   &&&&&          r   r<   r<      s    ( ""8-IJF--KF8TUU9933BB2'++Hb9&**73	
L ))66|D

#
#$7
888 ""8R0Fzz(// $$Wb1xx FK"&&}5jj*GP J&&!'wWU)'"+7D !!$'""#4#<#<= 	 0 9s    B)F2+F/,A(F2F2.BF2zset[asyncio.Task[None]]rc   z@dict[str, tuple[asyncio.Task[None], weakref.ref[ServerSession]]]_active_subscribersc          
     ,    V ^8  d   QhRRRRRRRRR	R
/# r,   r   )r   s   "r   r   r      sA     %O %O%O%O %O 	%O
 
%Or   c                2  "   V \         9   d   \         V ,          w  rEVP                  4       '       g   V! 4       e   R# VP                  4       '       gE   VP                  4        \        \        P
                  4      ;_uu_ 4        VG Rj  xL
  RRR4       \         V  \        P                  ! \        WW#4      RV R,           2R7      pV\        P                  ! V4      3\         V &   \        P                  RV 4       R#  Lv  + '       g   i     L|; i5i)ae  Start notification subscriber if not already running (idempotent).

Subscriber is created on first task submission and cleaned up on disconnect.
Safe to call multiple times for the same session.

Args:
    session_id: Session identifier
    session: MCP ServerSession
    docket: Docket instance
    fastmcp: FastMCP server instance (for elicitation relay)
Nznotification-subscriber-rV   rX   z.Started notification subscriber for session %s)rl   donecancelr   r@   rA   rb   rK   weakrefrefr5   r6   )r   r-   r   r.   rk   session_refs   &&&&  r   ensure_subscriber_runningrt      s     $ ((/
;yy{{{}8 yy{{KKM'0011

 2
+ $Z&J'
2'78D (,W[[-A&B
#
LLA:N  21s7   AD1DDDDA2DDD	Dc                    V ^8  d   QhRRRR/# )r   r   r   r   r   r   )r   s   "r   r   r     s     O Oc Od Or   c                `  "   V \         9  d   R# \         P                  V 4      w  rVP                  4       '       gE   VP                  4        \	        \
        P                  4      ;_uu_ 4        VG Rj  xL
  RRR4       \        P                  RV 4       R#  L$  + '       g   i     L*; i5i)zStop notification subscriber for a session.

Called when session disconnects. Pending messages remain in queue
for delivery if client reconnects (with TTL expiration).

Args:
    session_id: Session identifier
Nz.Stopped notification subscriber for session %s)	rl   popro   rp   r   r@   rA   r5   r6   )r   rk   rF   s   &  r   stop_subscriberrx     sr      ,,!%%j1GD99;;g,,--JJ .
LLA:N  .-s0   A,B..B4B5B9 B.BB+	&B.c                   V ^8  d   QhRR/# )r   r   intr   )r   s   "r   r   r   *  s     $ $c $r   c                      \        \        4      # )z2Get number of active subscribers (for monitoring).)lenrl   r   r   r   get_subscriber_countr}   *  s    "##r   )*__conditional_annotations____doc__
__future__r   r@   r    loggingrq   
contextlibr   r   r   typingr   r   r   	mcp.typesr[   r   r	   mcp.server.sessionr
   fastmcp.server.serverr   	getLogger__name__r5   r   r4   r(   r>   r8   r*   rK   r<   r7   rc   __annotations__rl   rt   rx   r}   )r~   s   @r   <module>r      s   " #      ' + + 0-			8	$ > E     :8V#r8>@ .1U * 2    
%OPO($r   