
    epj;{                       d Z ddlmZ ddl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
mZmZ ddlmZ ddlmZmZ ddlmZ ddlmZmZ d	Zd
ZdZdZdZdZdZdZdZ eeeeeh          Z dZ!dZ"dZ#dZ$dZ%dZ&dZ'dZ(dZ)dZ*dZ+dZ,dZ-dZ.dyd"Z/dzd$Z0dd%d%d%d&d'd{d3Z1d|d7Z2d}d;Z3d~d>Z4dd@Z5ddAZ6ddCZ7ddDZ8dzdEZ9dzdFZ:ddHZ;	 	 dddMZ<dddPZ=dddSZ>dddUZ?ddWZ@ddYZA	 dd&dZdd_ZBddd`ZCddaZDdddbZEdzdcZF G dd de          ZGdfZHdgZIdydhZJ G di dj          ZK G dk dl          ZL eL            ZM G dm dn          ZNddpZOddqZPdddsZQdddvZRddxZSdS )u\  
A2A protocol helpers — Agent Card construction, JSON-RPC framing, task store,
and disk-backed conversation persistence.

Wire shape follows A2A Protocol v1.0 (JSON-RPC 2.0 binding over HTTP):
  - Agent Card served at GET /.well-known/agent-card.json (canonical v1.0; legacy agent.json also answers)
  - Tasks via POST {jsonrpc:"2.0", method:"message/send", params:{...}}
  - Streaming via ``message/stream`` → SSE; events are StreamResponse objects
    discriminated by member presence (``statusUpdate`` / ``artifactUpdate``),
    stream closure signals the terminal state (no ``final`` field in v1.0)
  - Task states / message roles are v1.0 SCREAMING_SNAKE_CASE enums
  - Parts are the v1.0 unified shape ({"text": ..., "mediaType": ...}),
    discriminated by member presence (no ``kind`` field)
  - Push notification configs carry ``configId`` + ``createdAt`` and can be
    passed inline in ``message/send`` via configuration.taskPushNotificationConfig

We deliberately implement the subset of A2A needed for text task exchange with
stdlib only (no a2a-sdk). ``extract_text`` stays tolerant of v0.3 peers.
    )annotationsN)OrderedDictdefaultdictdeque)Future)datetimetimezone)Path)AnyOptionalz1.0TASK_STATE_SUBMITTEDTASK_STATE_WORKINGTASK_STATE_INPUT_REQUIREDTASK_STATE_AUTH_REQUIREDTASK_STATE_COMPLETEDTASK_STATE_FAILEDTASK_STATE_CANCELEDTASK_STATE_REJECTED	ROLE_USER
ROLE_AGENTz[INPUT_REQUIRED]iDiiiiii΂i͂î      returnintc                     	 t          t          j        dt          t                                        } t          dt          | t                              S # t          t          f$ r
 t          cY S w xY w)NA2A_MAX_PINGPONG_TURNS   )
r   osgetenvstr_DEFAULT_MAX_PINGPONGmaxmin_HARD_MAX_PINGPONG
ValueError	TypeError)vs    D/home/thesage/.hermes/hermes-agent/plugins/platforms/a2a/protocol.pymax_pingpong_turnsr)   N   sq    %	2C8M4N4NOOPP1c!/00111	" % % %$$$$%s   AA A43A4r    c                 z    t          j        t          j                                      d          dd         dz   S )z=ISO 8601 UTC timestamp with millisecond precision (A2A v1.0).z%Y-%m-%dT%H:%M:%S.%fNZ)r   nowr	   utcstrftime     r(   now_isor2   V   s1    <%%../EFFssKcQQr1   F )skills	streamingpush_notificationsauth_requiredtenantnameurldescriptionr4   Optional[list[dict]]r5   boolr6   r7   r8   dictc                    |dt           d}|r||d<   | ||dt          j        dd          t          j        dd          p|d	|g||d
d
ddgdg|pg d
}	|rddddi|	d<   dg ig|	d<   |	S )zConstruct an A2A v1.0 Agent Card document.

    ``tenant`` is the optional v1.0 multi-tenancy routing key advertised on
    AgentInterface. When present, clients MUST echo it in request params.
    JSONRPC)r:   protocolBindingprotocolVersionr8   z1.0.0A2A_PROVIDER_ORGzHermes AgentA2A_PROVIDER_URLr3   )organizationr:   F)r5   pushNotificationsstateTransitionHistoryextendedAgentCard
text/plain)
r9   r;   r:   versionprovidersupportedInterfacescapabilitiesdefaultInputModesdefaultOutputModesr4   bearerhttp)typeschemesecuritySchemessecurity)PROTOCOL_VERSIONr   r   )
r9   r:   r;   r4   r5   r6   r7   r8   ifacecards
             r(   build_agent_cardrY   _   s    " $+ E
  ! h "I&8.II9/44;
 
 !&w"!3&+!&	
 
 +^+n,B% D(  ,v::#
 &rN+ZKr1   toolsets)'list[str] | dict[str, list[str]] | None'
list[dict]c           
        g }t          | t                    rft          |                                           D ]C}d | |         pg D             }|                    d| |d| d|g|dd         z   d           DnCt          t          | pg                     D ]$}|                    d| |d| d|gd           %|s|                    ddd	dgd           |S )
u  Derive A2A skill descriptors from the agent's toolsets.

    Accepts either a plain list of toolset names, or a mapping of toolset name
    → tool names (built from the live tool registry for dynamic Agent Cards —
    tool names become tags so peers can match tasks to us).
    c                ,    g | ]}t          |          S r0   )r    ).0ts     r(   
<listcomp>z(skills_from_toolsets.<locals>.<listcomp>   s    DDDQ#a&&DDDr1   ztoolset.zHermes 'z' capabilitiesN
   )idr9   r;   tagsgeneralz$General-purpose conversational agent)
isinstancer>   sortedkeysappendset)rZ   r4   ts_name
tool_namestss        r(   skills_from_toolsetsrn      sW    F(D!! hmmoo.. 	 	GDD8G+<+BDDDJMM***A'AAA 	JssO3	     	 X^,,-- 	 	BMM%oo<"<<<	       AK	
 
 	 	 	 Mr1   req_idr   resultc                    d| |dS )N2.0)jsonrpcrc   rp   r0   )ro   rp   s     r(   jsonrpc_resultrt      s    Ff===r1   codemessagec                    d| ||ddS )Nrr   )ru   rv   )rs   rc   errorr0   )ro   ru   rv   s      r(   jsonrpc_errorry      s    Fdw5W5WXXXr1   payloadc                    t          | t                    r.|                     d          r|                     d          rd| iS d| iS )zA2A v1.0 SendMessageResponse oneof wrapper.

    The JSON-RPC ``SendMessage`` result is not a bare Task/Message; it is a
    wrapper containing exactly one of ``task`` or ``message``. Legacy methods
    still return bare payloads for compatibility.
    statusrc   taskrv   rf   r>   get)rz   s    r(   send_message_responser      sR     '4   !W[[%:%: !w{{4?P?P !  wr1   c                    t          | t                    r`t          |                     d          t                    r| d         S t          |                     d          t                    r| d         S | S )zGReturn the Task/Message inside a v1.0 response, or pass legacy through.r}   rv   r~   )rp   s    r(   unwrap_send_message_responser      sj    &$ %fjj(($// 	"&>!fjj++T22 	%)$$Mr1   r}   c                
    d| iS )z'v1.0 StreamResponse with a task member.r}   r0   )r}   s    r(   stream_taskr      s    D>r1   c                
    d| iS )z*v1.0 StreamResponse with a message member.rv   r0   )rv   s    r(   stream_messager      s    wr1   c                 H    dt          j                    j        d d         z   S )Nztask-   uuiduuid4hexr0   r1   r(   new_task_idr      s    TZ\\%crc***r1   c                 H    dt          j                    j        d d         z   S )Nzctx-r   r   r0   r1   r(   new_context_idr      s    DJLL$SbS)))r1   textc                    | ddS )zDBuild a v1.0 text Part (member-presence discriminated, no ``kind``).rI   )r   	mediaTyper0   )r   s    r(   	text_partr      s    |444r1   application/octet-streamrawfilename
media_typec                :    d|i}|r||d<   | r| |d<   n|r||d<   |S )u   Build a v1.0 file Part.

    Either ``url`` (file reference) or ``raw`` (base64-encoded bytes) must be
    provided. Discrimination is by member presence — no ``kind`` field.
    r   r   r:   r   r0   )r:   r   r   r   parts        r(   	file_partr      sI     (4D $#Z
 U	 UKr1   application/jsondatac                    | |dS )z<Build a v1.0 data Part (structured data, no ``kind`` field).)r   r   r0   )r   r   s     r(   	data_partr      s    z222r1   role
context_idc                h    | t          |          gt          j                    j        d}|r||d<   |S )z2Build an A2A v1.0 Message with a single text Part.r   parts	messageId	contextId)r   r   r   r   )r   r   r   msgs       r(   text_messager     sD     D//"Z\\% C
  &%KJr1   r   c                L    | |t          j                    j        d}|r||d<   |S )zBBuild an A2A v1.0 Message with arbitrary Parts (text, file, data).r   r   r   )r   r   r   r   s       r(   message_with_partsr     s;     Z\\% C
  &%KJr1   message_or_paramsc                   |                      d|           }t          |t                    r|                     dg           ng }g }|D ]}t          |t                    s|                     d          }t          |t                    r|                    |           Y|                     d          dk    rDt          |                     d          t                    r|                    |d                    |                     d          }t          |t                    r|r|                     d          p|                     d          pd}|                     d	          p|                     d
          pd}|rd| dnd}	|                    |	 d| |rd| dndz              k|                     d          }
t          |
t                    rt          |
                     d          t                    rg|
d         }|
                     d          pd}|
                     d
          pd}|rd| dnd}	|                    |	 d| |rd| dndz              $t          |                     d          t                    rw|                     d          pd}|                     d	          pd}|rd| dnd}	t          |d                    d}|                    |	 d| |rd| dndz              |                     d          }|x	 t          j        |dt                    }n&# t          t          f$ r t          |          }Y nw xY w|                     d	          pd}|                    d| d|            R|                     d          dk    r|                     d          j	 t          j        |d         dt                    }n,# t          t          f$ r t          |d                   }Y nw xY w|                    d|            d
                    |                                          S )aT  Pull concatenated text from an A2A Message / Task-result / params payload.

    v1.0 Parts carry a ``text`` member directly; v0.3 used ``kind: "text"``
    and some pre-0.3 peers used ``type``. All three shapes put the payload in
    ``part["text"]``, so presence of a string ``text`` member is the test.

    File and data Parts are rendered into the text stream so the agent sees
    them: file Parts with a URL include the URL and filename; data Parts
    include their JSON-serialised content. Raw (base64) file Parts are noted
    but not decoded (the agent can't act on binary inline).
    rv   r   r   kindr:   r   r9   r3   r   mimeTypez[file: ]z[file] z ()filefileWithUrir   z bytes base64-encodedr   NF)ensure_asciidefaultr   z[data (z)]
z[data]

)r   rf   r>   r    ri   lenjsondumpsr&   r%   joinstrip)r   r   r   chunksr   txtr:   fnamemtypelabelv03_fileuri	size_noter   rendereds                  r(   extract_textr     s    

	+<
=
=C$.sD$9$9ACGGGR   rEF 6 6$%% 	hhvc3 	MM#88Fv%%*TXXf5E5Es*K*K%MM$v,'''hhuooc3 	C 	HHZ((BDHHV,<,<BEHH[))GTXXj-A-AGRE*/=&e&&&&XEMMU**S**u.Lm5mmmm"MNNN88F##h%% 	*X\\-5P5PRU*V*V 	=)CLL((.BELL,,2E*/=&e&&&&XEMMU**S**u.Lm5mmmm"MNNNdhhuoos++ 	HHZ((.BEHH[))/RE*/=&e&&&&XEtE{++BBBIMMU00Y00U4RMMMMMPRSTTTxx%:dLLLz* % % %t99%HH[))?-?EMM9E99x99:::88Fv%%$((6*:*:*F-:d6lPSTTTz* - - -tF|,,-MM/X//00099V""$$$s$   #M   M#"M#	"O,,&PPparamsc                    |                      d          pi }d}t          |t                    r$t          |                     d          pd          }|p#t          |                      d          pd          S )zBv1.0 puts contextId inside the Message; tolerate legacy top-level.rv   r3   r   )r   rf   r>   r    )r   r   ctxs      r(   extract_context_idr   f  sq    
**Y


%2C
C#t .#''+&&,"--4#fjj--3444r1   
created_attask_idstate
agent_textr   c                   t                      }| |||dd}|rWt          t          ||          |d         d<   |t          k    r-t	          j                    j        t          |          gdg|d<   |S )u  Build an A2A v1.0 Task object for a message/send result.

    ``created_at`` is accepted for call-site compatibility but not serialized —
    the A2A v1.0 ``Task`` proto (``lf.a2a.v1.Task``) has no ``createdAt`` or
    ``lastModified`` field.  Strict ProtoJSON parsers (e.g. a2a-sdk 1.1.0)
    reject unknown fields, so we must not include them.  The spec's §5.6.1
    timestamp-format example mentions them but they are not in the proto.
    r   	timestamp)rc   r   r|   r|   rv   
artifactIdr   	artifacts)r2   r   r   STATE_COMPLETEDr   r   r   r   )r   r   r   r   r   r-   r}   s          r(   
build_taskr   o  s      ))C!44 D
  $0Z$T$TXy!O##"jll.#J//0" " !D Kr1   c                j    |t                      d}|rt          t          ||          |d<   d| ||diS )z/v1.0 StreamResponse with a statusUpdate member.r   rv   statusUpdate)taskIdr   r|   )r2   r   r   )r   r   r   r   r|   s        r(   status_updater     sH    ',799EEF G(T:FFywZSYZZ[[r1   c                `    d| |t          j                    j        t          |          gddiS )z2v1.0 StreamResponse with an artifactUpdate member.artifactUpdater   )r   r   artifact)r   r   r   r   )r   r   r   s      r(   artifact_updater     sC     	#"jll.#D//* 
 
	 	r1   c                `    |t          ||           }n| }dt          j        |d           dS )uf  Encode one StreamResponse as a JSON-RPC-wrapped SSE data frame.

    A2A v1.0 §9.4 requires each SSE frame to be a full JSON-RPC response:
    ``{"jsonrpc":"2.0","id":<req_id>,"result":{StreamResponse}}``.  Emitting a
    bare StreamResponse (the REST binding shape) breaks JSON-RPC clients that
    expect the envelope, including the official a2a-sdk.
    Nzdata: Fr   z

)rt   r   r   )rz   ro   envelopes      r(   sse_datar     s@     !&'22BDJxe<<<BBBBr1   c                     dS )u(  SSE stream-closure marker — a comment, not a parseable data frame.

    A2A v1.0 signals terminal state by closing the stream.  Emitting
    ``data: {}`` causes JSON-RPC clients to try parsing an empty response and
    fail.  An SSE comment line (``: done``) is ignored by all SSE parsers.
    z: done

r0   r0   r1   r(   sse_doner     s	     <r1   c                  .    e Zd ZdZdZddZdd	Zdd
ZdS )TurnTrackeru   Counts inbound turns per context_id to stop infinite agent↔agent loops.

    A "turn" is one inbound message/send from a peer. When the count exceeds
    max_pingpong_turns(), the adapter rejects further messages for that context.
    i  r   Nonec                v    t          t                    | _        i | _        t	          j                    | _        d S N)r   r   _counts_timestamps	threadingLock_lockselfs    r(   __init__zTurnTracker.__init__  s,    '23'7'7-/^%%


r1   r   r    r   c                     j         5  t          j                     fd j                                        D             }|D ]8} j                            |d            j                            |d           9 j        |xx         dz  cc<    j        |<    j        |         cddd           S # 1 swxY w Y   dS )z;Increment and return the turn count; prunes stale contexts.c                6    g | ]\  }}|z
  j         k    |S r0   )_TTL)r_   cidrm   r-   r   s      r(   ra   z%TurnTracker.track.<locals>.<listcomp>  s-    YYYWS"C"HtyDXDXSDXDXDXr1   Nr   )r   timer   itemsr   pop)r   r   staler   r-   s   `   @r(   trackzTurnTracker.track  s   Z 	, 	,)++CYYYYY(8(>(>(@(@YYYE 0 0  d+++ $$S$////L$$$)$$$+.DZ(<
+	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	, 	,s   B B77B;>B;c                    | j         5  | j                            |d           | j                            |d           ddd           dS # 1 swxY w Y   dS )z<Reset turn count for a context (e.g. after explicit cancel).N)r   r   r   r   )r   r   s     r(   resetzTurnTracker.reset  s    Z 	3 	3LZ...  T222	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3 	3s   7AAANr   r   )r   r    r   r   )r   r    r   r   )__name__
__module____qualname____doc__r   r   r   r   r0   r1   r(   r   r     sa          D& & & &

, 
, 
, 
,3 3 3 3 3 3r1   r   <   g      N@c                     	 t          dt          t          j        dt	          t
                                                  S # t          t          f$ r
 t
          cY S w xY w)Nr   A2A_RATE_LIMIT)r"   r   r   r   r    _RATE_LIMIT_DEFAULTr%   r&   r0   r1   r(   _rate_limit_per_minuter    sa    #1c")$4c:M6N6NOOPPQQQ	" # # #""""#s   AA AAc                  "    e Zd ZdZd
dZddZd	S )RateLimiterzFSliding-window request limiter, one bucket per authenticated identity.r   r   c                h    t          t                    | _        t          j                    | _        d S r   )r   r   _bucketsr   r   r   r   s    r(   r   zRateLimiter.__init__  s$    1<U1C1C^%%


r1   identityr    r=   c                   | j         5  t                      }t          j                    }| j        |         }|r>||d         z
  t          k    r*|                                 |r||d         z
  t          k    *t          |          |k    r	 d d d            dS |                    |           	 d d d            dS # 1 swxY w Y   d S )Nr   FT)r   r  r   r	  _RATE_WINDOWpopleftr   ri   )r   r
  limitr-   buckets        r(   allowzRateLimiter.allow  s)   Z 		 		*,,E)++C]8,F !S6!9_|;;     !S6!9_|;;6{{e##		 		 		 		 		 		 		 		 MM#		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		s   BB;B;;B?B?Nr   )r
  r    r   r=   )r   r   r   r   r   r  r0   r1   r(   r  r    sB        PP& & & &
 
 
 
 
 
r1   r  c                  2    e Zd ZdZddZddZddZdd
ZdS )Metricsz#Simple counters for A2A operations.r   r   c                    d| _         d| _        d| _        d| _        d| _        d| _        d| _        d| _        d| _        t          j	                    | _
        t          d          | _        d S )Nr   d   )maxlen)inbound_totaloutbound_totalstreams_started	push_sentpush_failedtasks_completedtasks_failedanti_loop_triggersrate_limit_triggersr   _start_timer   
_latenciesr   s    r(   r   zMetrics.__init__  sm      "##$ 9;;(-S(9(9(9r1   secondsfloatc                :    | j                             |           d S r   )r   ri   )r   r!  s     r(   record_latencyzMetrics.record_latency!  s    w'''''r1   c                f    | j         sdS t          | j                   t          | j                   z  S )Ng        )r   sumr   r   s    r(   avg_latencyzMetrics.avg_latency$  s0     	34?##c$/&:&:::r1   dict[str, Any]c                   t          j                     | j        z
  }t          |d          | j        | j        | j        | j        | j        | j        | j	        | j
        | j        t          |                                 dz  d          dS )Nr   i  )uptime_secondsr  r  r  r  r  r  r  r  r  avg_latency_ms)r   r  roundr  r  r  r  r  r  r  r  r  r'  )r   uptimes     r(   snapshotzMetrics.snapshot)  s    t//#FA..!/"1#3+#3 -"&"9#'#;#D$4$4$6$6$=qAA
 
 	
r1   Nr   )r!  r"  r   r   )r   r"  )r   r(  )r   r   r   r   r   r$  r'  r.  r0   r1   r(   r  r    sj        --: : : :( ( ( (; ; ; ;

 
 
 
 
 
r1   r  c                      e Zd ZdZdZd6dZed7d8d            Z	 d7d9dZd:dZ		 d7d;dZ
ed<d            Z	 	 d=d>dZd7d?dZ	 	 d=d@dZdAdZd7dBdZdCdDd Zd7dEd"Z	 	 	 	 	 	 	 dFdGd*ZdHdId.Zd6d/ZedJdKd5            Zd0S )L	TaskStorea6  In-memory store of A2A tasks, kept after completion for tasks/get.

    Records carry the routed agent slug and tenant. All read/write helpers accept
    optional scope values and return not-found when the task exists but is not
    visible in that scope, satisfying the spec's authorization scoping rule.
    i  r   r   c                j    t                      | _        i | _        t          j                    | _        d S r   )r   _tasks	_watchersr   r   r   r   s    r(   r   zTaskStore.__init__K  s'    :E--24^%%


r1   r3   recr>   
agent_slugr    r8   r=   c                ~    |r|                      dd          |k    rdS |r|                      dd          |k    rdS dS )Nr5  r3   Fr8   Tr   )r4  r5  r8   s      r(   	_in_scopezTaskStore._in_scopeP  sQ     	#'',33zAA5 	cggh++v555tr1   r   r   peerc                    ||||pd|pdt           dt          j                    t                      ddd}| j        5  || j        |<   d d d            n# 1 swxY w Y   t          |          S )Nr3   )r   r   r9  r5  r8   r   replyr   created_isopush_urlpush_config_id)STATE_SUBMITTEDr   r2   r   r2  r>   )r   r   r   r9  r5  r8   r4  s          r(   createzTaskStore.createX  s     $$*l$)++"99 
 
 Z 	' 	'#&DK 	' 	' 	' 	' 	' 	' 	' 	' 	' 	' 	' 	' 	' 	' 	'Cyys   AAAr   c                    | j         5  | j                            |          }|r|d         t          vr||d<   d d d            d S # 1 swxY w Y   d S )Nr   )r   r2  r   TERMINAL_STATES)r   r   r   r4  s       r(   	set_statezTaskStore.set_statek  s    Z 	% 	%+//'**C %s7|?::$G	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	%s   1AA
A
r:   Optional[dict]c                F   | j         5  | j                            |          }|r|                     |||          s	 ddd           dS ||d<   dt	          j                    j        dd         z   |d<   |                     |          cddd           S # 1 swxY w Y   dS )zEAttach a push notification config; returns the stored config or None.Nr=  zcfg-   r>  )r   r2  r   r8  r   r   r   _push_config_view)r   r   r:   r5  r8   r4  s         r(   set_push_configzTaskStore.set_push_configq  s    Z 	/ 	/+//'**C dnnS*fEE 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ "C
O$*TZ\\-=crc-B$BC !))#..	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/s   5B
?BBBc                    |                      d          pd| d         |                      dd          d|                      d          pdidS )z9Build the JSON-RPC result for a push notification config.r>  r3   r   r<  r:   r=  )configIdr   	createdAtpushNotificationConfigr7  )r4  s    r(   rG  zTaskStore._push_config_view|  sY      0117R)n33',cggj.A.A.GR&H	
 
 	
r1   	config_idc                l   | j         5  | j                            |          }|r,|                     |||          r|                    d          s	 d d d            d S |r'|                    d          |k    r	 d d d            d S |                     |          cd d d            S # 1 swxY w Y   d S )Nr=  r>  r   r2  r   r8  rG  r   r   rM  r5  r8   r4  s         r(   get_push_configzTaskStore.get_push_config  s>   Z 	/ 	/+//'**C dnnS*fEE SWWU_M`M` 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/  SWW%566)CC	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ ))#..	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/s   A
B)B)B))B-0B-r\   c                   | j         5  | j                            |          }|r,|                     |||          r|                    d          sg cd d d            S |                     |          gcd d d            S # 1 swxY w Y   d S )Nr=  rO  r   r   r5  r8   r4  s        r(   list_push_configszTaskStore.list_push_configs  s    Z 	1 	1+//'**C dnnS*fEE SWWU_M`M` 	1 	1 	1 	1 	1 	1 	1 	1 **3//0		1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1 	1s   A
BBBBc                Z   | j         5  | j                            |          }|r,|                     |||          r|                    d          s	 d d d            dS |r'|                    d          |k    r	 d d d            dS d|d<   d|d<   	 d d d            dS # 1 swxY w Y   d S )Nr=  Fr>  r3   T)r   r2  r   r8  rP  s         r(   delete_push_configzTaskStore.delete_push_config  sD   Z 	 	+//'**C dnnS*fEE SWWU_M`M` 	 	 	 	 	 	 	 	  SWW%566)CC	 	 	 	 	 	 	 	 !C
O$&C !	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	s   A
B B B  B$'B$c                    | j         5  | j                            |          }|s	 d d d            dS |d         dc}|d<   |cd d d            S # 1 swxY w Y   d S )Nr3   r=  )r   r2  r   )r   r   r4  r:   s       r(   pop_push_urlzTaskStore.pop_push_url  s    Z 	 	+//'**C 	 	 	 	 	 	 	 	 $'z?B CZ	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	s   AAAAc                    | j         5  | j                            |          }|r|                     |||          s	 d d d            d S t	          |          cd d d            S # 1 swxY w Y   d S r   )r   r2  r   r8  r>   rS  s        r(   r   zTaskStore.get  s    Z 	 	+//'**C dnnS*fEE 	 	 	 	 	 	 	 	 99		 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	s   5A%
A%%A),A)r;  c                   g }| j         5  | j                            |          }|r|d         t          v r	 ddd           dS ||d<   ||d<   t	          j                    |d<   | j                            |g           }|                                  t          |          }ddd           n# 1 swxY w Y   |D ]-}|	                                s|
                    ||f           .|S )z2Transition a task to a terminal state. Idempotent.r   Nr;  completed_at)r   r2  r   rB  r   r3  r   _trim_lockedr>   done
set_result)r   r   r   r;  watchersr4  outfuts           r(   completezTaskStore.complete  sM   !#Z 		 		+//'**C #g,/99		 		 		 		 		 		 		 		 !CL CL"&)++C~))'266Hs))C		 		 		 		 		 		 		 		 		 		 		 		 		 		 		  	/ 	/C88:: /u~...
s   -B.AB..B25B2Optional[Future]c                   | j         5  | j                            |          }|r|                     |||          s	 d d d            d S t	                      }|d         t
          v r2|                    |d         |                    dd          f           n.| j                            |g           	                    |           |cd d d            S # 1 swxY w Y   d S )Nr   r;  r3   )
r   r2  r   r8  r   rB  r^  r3  
setdefaultri   )r   r   r5  r8   r4  ra  s         r(   watchzTaskStore.watch  s:   Z 		 		+//'**C dnnS*fEE 		 		 		 		 		 		 		 		 !((C7|..Gcgggr.B.BCDDDD))'266==cBBB		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		 		s   5C
A>CCC2   r   F	page_sizer   offset
with_totalc                    t          dt          t          |pd          d                    } j        5  d t	           j                                                  D             }ddd           n# 1 swxY w Y   sr fd|D             }rfd|D             }rfd|D             }t          |          }	||||z            }
||z   |	k     r||z   nd	}|r|
||	fS |
|fS )
zFiltered task page (newest first).

        Historical API returns ``(records, next_offset)``. v1.0 ListTasks needs
        ``totalSize``, so callers can opt into ``(records, next_offset, total)``.
        r   rg  r  c                ,    g | ]}t          |          S r0   )r>   )r_   rs     r(   ra   z"TaskStore.list.<locals>.<listcomp>  s    DDDDGGDDDr1   Nc                B    g | ]}                     |          |S r0   )r8  )r_   rm  r5  r   r8   s     r(   ra   z"TaskStore.list.<locals>.<listcomp>  s.    MMM!t~~aV'L'LMAMMMr1   c                ,    g | ]}|d          k    |S r   r0   )r_   rm  r   s     r(   ra   z"TaskStore.list.<locals>.<listcomp>  s'    EEE!q*'D'DA'D'D'Dr1   c                ,    g | ]}|d          k    |S r   r0   )r_   rm  r   s     r(   ra   z"TaskStore.list.<locals>.<listcomp>  s'    ;;;!qzU':':A':':':r1   r   )r"   r#   r   r   reversedr2  valuesr   )r   r   r   rh  ri  r5  r8   rj  recstotalpagenext_offsets   ```  ``     r(   listzTaskStore.list  s    3s9?33S99::	Z 	E 	EDDXdk.@.@.B.B%C%CDDDD	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	E 	N 	NMMMMMMtMMMD 	FEEEEtEEED 	<;;;;t;;;DD		F6I--.,2Y,>,F,Ffy((A 	,e++[  s   1A77A;>A;,  timeout_seconds	list[str]c                *   | j         5  t          j                    fd| j                                        D             }d d d            n# 1 swxY w Y   g }|D ]3}|                     |t
          d          r|                    |           4|S )Nc                V    g | ]%\  }}|d          t           vr|d         z
  k    #|&S )r   r   rB  )r_   tidr4  r-   r{  s      r(   ra   z*TaskStore.fail_orphans.<locals>.<listcomp>  sK        Sw<66#l++o== ===r1   u%   [task orphaned — no reply produced])r   r   r2  r   rb  STATE_FAILEDri   )r   r{  r   failedr  r-   s    `   @r(   fail_orphanszTaskStore.fail_orphans  s    Z 	 	)++C    $(K$5$5$7$7  E	 	 	 	 	 	 	 	 	 	 	 	 	 	 	  	# 	#C}}S,0WXX #c"""s   :AAAc                    d | j                                         D             }t          |          | j        z
  }|d t	          d|                   D ]}| j                             |d            d S )Nc                6    g | ]\  }}|d          t           v |S rr  r  )r_   r  r4  s      r(   ra   z*TaskStore._trim_locked.<locals>.<listcomp>  s*    ___HCs7|?^?^C?^?^?^r1   r   )r2  r   r   _MAX_TERMINALr"   r   )r   terminalexcessr  s       r(   r\  zTaskStore._trim_locked  sy    __(9(9(;(;___X!33OSF^^O, 	' 	'CKOOC&&&&	' 	'r1   NThistory_lengthOptional[int]include_artifactsc           
     .   t          | d         | d         | d         |                     dd          |                     dd                    }|s|                    dd	           |d
k    r|                    dd	           t          j        |          S )z2Render a stored record as an A2A v1.0 Task object.r   r   r   r;  r3   r<  r   r   Nr   history)r   r   r   copydeepcopy)r4  r  r  r}   s       r(   to_taskzTaskStore.to_task  s     	NLGGGR  ww}b11
 
 
 ! 	(HH[$'''QHHY%%%}T"""r1   r   )r3   r3   )r4  r>   r5  r    r8   r    r   r=   )r   r    r   r    r9  r    r5  r    r8   r    r   r>   )r   r    r   r    r   r   )
r   r    r:   r    r5  r    r8   r    r   rD  )r4  r>   r   r>   )r3   r3   r3   )
r   r    rM  r    r5  r    r8   r    r   rD  )r   r    r5  r    r8   r    r   r\   )
r   r    rM  r    r5  r    r8   r    r   r=   )r   r    r   r    )r   r    r5  r    r8   r    r   rD  r3   )r   r    r   r    r;  r    r   rD  )r   r    r5  r    r8   r    r   rc  )r3   r3   rg  r   r3   r3   F)r   r    r   r    rh  r   ri  r   r5  r    r8   r    rj  r=   )rz  )r{  r   r   r|  )NT)r4  r>   r  r  r  r=   r   r>   )r   r   r   r   r  r   staticmethodr8  r@  rC  rH  rG  rQ  rT  rV  rX  r   rb  rf  ry  r  r\  r  r0   r1   r(   r0  r0  A  s         M& & & &
     \ 46    &% % % % =?	/ 	/ 	/ 	/ 	/ 
 
 
 \
 >@<>/ / / / /1 1 1 1 1 AC?A
 
 
 
 
           $
 
 
 
 
  ! ! ! ! !>    ' ' ' ' # # # # \# # #r1   r0  r
   c                     	 ddl m}  t           |                       }n<# t          $ r/ t          t          j                            d                    }Y nw xY w|dz  S )Nr   )get_hermes_homez	~/.hermesa2a_conversations)hermes_constantsr  r
   	Exceptionr   path
expanduser)r  bases     r(   	_conv_dirr    sx    5444444OO%%&& 5 5 5BG&&{33445%%%s     6AAc                H    d                     d | pdD                       pdS )Nr3   c              3  J   K   | ]}|                                 s|d v |V  dS )z-_N)isalnum)r_   cs     r(   	<genexpr>z_safe_name.<locals>.<genexpr>!  s3      TT199;;T!t))1))))TTr1   r   )r   rp  s    r(   
_safe_namer     s.    77TTz6YTTTTTaXaar1   r   c                   	 t                      }|                    dd           t          j                    |||d}|t          |            dz                      dd          5 }|                    t          j        |d	          d
z              ddd           dS # 1 swxY w Y   dS # t          $ r Y dS w xY w)z=Append one message to the context's on-disk conversation log.T)parentsexist_ok)rm   r   r   r   .jsonlautf-8encodingFr   r   N)	r  mkdirr   r  openwriter   r   r  )r   r   r   r   dr4  fhs          r(   persist_messager  $  s)   KK	t,,,Y[[$QQZ
++333399#9PP 	ATVHHTZ%8884?@@@	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A 	A   s6   A'B0 )-B#B0 #B''B0 *B'+B0 0
B>=B>rg  r  c                   t                      t          |            dz  }|                                sg S g }	 |                    dd          5 }|D ]V}|                                }|s	 |                    t          j        |                     B# t          j        $ r Y Sw xY w	 ddd           n# 1 swxY w Y   n# t          $ r g cY S w xY w|| d         S )zBLoad the last *limit* messages for a context (empty list if none).r  rm  r  r  N)
r  r  existsr  r   ri   r   loadsJSONDecodeErrorr  )r   r  r  r`  r  lines         r(   load_conversationr  0  sR   ;;Jz22::::D;;== 	CYYsWY-- 	  zz|| JJtz$//0000+   H	 	 	 	 	 	 	 	 	 	 	 	 	 	 	    			vww<sY   C B9/'BB9B)&B9(B))B9-C 9B==C  B=C CCr|  c                     t                      } |                                 sg S t          d |                     d          D                       S )z;Return known context-ids that have persisted conversations.c              3  $   K   | ]}|j         V  d S r   )stem)r_   ps     r(   r  z%list_conversations.<locals>.<genexpr>J  s$      44Q!&444444r1   z*.jsonl)r  r  rg   glob)r  s    r(   list_conversationsr  E  sI    A88:: 	44!&&"3"3444444r1   )r   r   )r   r    )r9   r    r:   r    r;   r    r4   r<   r5   r=   r6   r=   r7   r=   r8   r    r   r>   )rZ   r[   r   r\   )ro   r   rp   r   r   r>   )ro   r   ru   r   rv   r    r   r>   )rz   r>   r   r>   )rp   r   r   r   )r}   r>   r   r>   )rv   r>   r   r>   )r   r    r   r>   )r3   r3   r3   r   )
r:   r    r   r    r   r    r   r    r   r>   )r   )r   r   r   r    r   r>   r  )r   r    r   r    r   r    r   r>   )r   r    r   r\   r   r    r   r>   )r   r>   r   r    )r   r>   r   r    )r   r    r   r    r   r    r   r    r   r    r   r>   )
r   r    r   r    r   r    r   r    r   r>   )r   r    r   r    r   r    r   r>   r   )rz   r>   ro   r   r   r    )r   r
   )r   r    r   r    )
r   r    r   r    r   r    r   r    r   r   )rg  )r   r    r  r   r   r\   )r   r|  )Tr   
__future__r   r   r  r   r   r   r   collectionsr   r   r   concurrent.futuresr   r   r	   pathlibr
   typingr   r   rV   r?  STATE_WORKINGSTATE_INPUT_REQUIREDSTATE_AUTH_REQUIREDr   r  STATE_CANCELEDSTATE_REJECTED	frozensetrB  r   r   INPUT_REQUIRED_MARKER	ERR_PARSEERR_INVALID_PARAMSERR_METHOD_NOT_FOUNDERR_TASK_NOT_FOUNDERR_TASK_NOT_CANCELABLEERR_PUSH_NOT_SUPPORTEDERR_UNAUTHORIZEDERR_RATE_LIMITEDERR_UNTRUSTED_PEERr!   r$   r)   r2   rY   rn   rt   ry   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r  r  r  r  r  metricsr0  r  r  r  r  r  r0   r1   r(   <module>r     s   ( # " " " " "   				       7 7 7 7 7 7 7 7 7 7 % % % % % % ' ' ' ' ' ' ' '                        )$2 0 ("&&)_lNN[\\ 	

 +  	            % % % %R R R R $($1 1 1 1 1 1h       N> > > >Y Y Y Y	  	  	  	       
       
+ + + +* * * *5 5 5 5
 =? :    "3 3 3 3 3
	 	 	 	 		 	 	 	 	F% F% F% F%R5 5 5 5 	      H\ \ \ \ \   C C C C C   3 3 3 3 3 3 3 3J  # # # #       4'
 '
 '
 '
 '
 '
 '
 '
T '))P# P# P# P# P# P# P# P#l& & & &b b b b	 	 	 	 	    *5 5 5 5 5 5r1   