
    `gj:                       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mZ ddl	m
Z
mZmZmZ ddlmZmZ  ej"                  e      ZddZddZdd	d
	 	 	 	 	 d dZ eh d      ZdZd!dZd"dZd#dZd$dZd%dZd	d	 	 	 	 	 	 	 	 	 	 	 d&dZ eh d      Z d'd(dZ!d'd)dZ"d*dZ#ddddddd	 	 	 	 	 d+dZ$d,d-dZ%d'd-dZ&g dZ'y).u  Codex API runtime — App Server and Responses-API streaming paths.

Extracted from :class:`AIAgent` to keep the agent loop file focused.
Each function takes the parent ``AIAgent`` as its first argument
(``agent``).  AIAgent keeps thin forwarder methods for backward
compatibility.

* ``run_codex_app_server_turn`` — drives one turn through the
  ``codex_app_server`` subprocess client (used when a Codex CLI install
  is the active provider).
* ``run_codex_stream`` — streams a Codex Responses API call (the
  ``codex_responses`` api_mode).
* ``run_codex_create_stream_fallback`` — recovery path when the
  Responses ``stream=True`` initial create fails.
    )annotationsN)SimpleNamespace)AnyCallableDictList)claim_stream_writerstream_writer_is_currentc                   t        | t              ryt        | t              rt        | d      S t        | t              rt        t        |       d      S t        | t
              r	 t        t        |       d      S y# t        $ r Y yw xY w)Nr   )
isinstanceboolintmaxfloatstr
ValueError)values    F/root/.hermes/venv/lib/python3.12/site-packages/agent/codex_runtime.py_coerce_usage_intr      sy    %%5!}%3u:q!!%	s5z1%%   		s   $A: :	BBc                
   | xj                   dz  c_         t        |dd      }t        |t              r|st        | dd      }|t        |dd      r|j	                  i        | j
                  rt| j                  rh	 | j                  s| j                          | j
                  j                  | j                  | j                  | j                  | j                  dd       i S i S d
dlm}m} t'        |j)                  d            }t'        |j)                  d            }t'        |j)                  d            }	t'        |j)                  d            }
t'        |j)                  d            } |||	|d
|
|      }|j*                  }|j,                  }|xs |j.                  }||||j0                  |j,                  |j2                  |j4                  |j6                  d}t        | dd      }|;	 |j	                  |       t        |dd      }t        |t8              r|d
kD  r||_        | xj<                  |z  c_        | xj>                  |z  c_        | xj@                  |z  c_         | xjB                  |j0                  z  c_!        | xjD                  |j,                  z  c_"        | xjF                  |j2                  z  c_#        | xjH                  |j4                  z  c_$        | xjJ                  |j6                  z  c_%         || j                  || j                  | j                  t        | dd            }|jL                  (| xjN                  tQ        |jL                        z  c_'        |jR                  | _*        |jV                  | _,        | j
                  r| j                  r	 | j                  s| j                          | j
                  j                  | j                  |j0                  |j,                  |j2                  |j4                  |j6                  |jL                  tQ        |jL                        nd|jR                  |jV                  | j                  | j                  |jR                  dk(  rdnd| j                  d       i |||jL                  tQ        |jL                        nd|jR                  |jV                  dS # t        $ r,}t        j                  d	| j                  |       Y d}~i S d}~ww xY w# t        $ r t        j                  dd       Y w xY w# t        $ r,}t        j                  d| j                  ||       Y d}~d}~ww xY w)a5  Translate Codex app-server token usage into Hermes accounting.

    Codex app-server reports usage via thread/tokenUsage/updated as:
    inputTokens, cachedInputTokens, outputTokens, reasoningOutputTokens,
    totalTokens.

    Hermes' canonical prompt bucket includes uncached input + cached input.
    The Codex app-server protocol does not currently expose cache-write tokens,
    so that bucket remains zero on this runtime.

    Even when Codex omits usage for a turn, Hermes should still count that turn
    as one API call for session/status accounting.
       token_usage_lastNcontext_compressor%awaiting_real_usage_after_compressionFsubscription_included)modelbilling_providerbilling_base_urlbilling_modeapi_call_countz=Codex app-server api-call persistence failed (session=%s): %sr   )CanonicalUsageestimate_usage_costinputTokenscachedInputTokensoutputTokensreasoningOutputTokenstotalTokens)input_tokensoutput_tokenscache_read_tokenscache_write_tokensreasoning_tokens	raw_usage)prompt_tokenscompletion_tokenstotal_tokensr(   r)   r*   r+   r,   model_context_windowz$codex app-server usage update failedTexc_infoapi_key )providerbase_urlr4   included)r(   r)   r*   r+   r,   estimated_cost_usdcost_statuscost_sourcer   r   r   r   r    zECodex app-server token persistence failed (session=%s, tokens=%d): %s)last_prompt_tokensr9   r:   r;   )-session_api_callsgetattrr   dictupdate_from_response_session_db
session_id_session_db_created_ensure_db_sessionupdate_token_countsr   r6   r7   	Exceptionloggerdebugagent.usage_pricingr!   r"   r   getr.   r)   r0   r(   r*   r+   r,   r   context_lengthsession_prompt_tokenssession_completion_tokenssession_total_tokenssession_input_tokenssession_output_tokenssession_cache_read_tokenssession_cache_write_tokenssession_reasoning_tokens
amount_usdsession_estimated_cost_usdr   statussession_cost_statussourcesession_cost_source)agentturnusage
compressorexcr!   r"   r(   r*   r)   r,   reported_totalcanonical_usager.   r/   r0   
usage_dictcontext_windowcost_results                      r   _record_codex_app_server_usagerd   .   s    
q D,d3EeT"%U$8$?
"
$KUS
 ++B/!1!100,,.!!55$$++%*^^%*^^!8#$ 6  	r	G$UYY}%=>L)%))4G*HI%eii&?@M(3J)KL&uyy'?@N$!#+)O $11M'55!A_%A%AL&.$'44(66,>>-@@+<<	J  4d;J	P++J7$T+A4HN.#.>A3E,:
) 
=0	##'88#	,.	/">">>	?#@#@@	##'H'HH#	$$(J(JJ$	""o&F&FF"%y"-K )((E+2H2H,II( + 2 2E + 2 2EU--	,,((*11  ,99-;;"1"C"C#2#E#E!0!A!A))5 $))?)?#@;?'..'..!&!&%%3 59=kk ! 2 0
+!!- $K$:$:;37"))")) A  S$$c  	X  	PLL?$LO	P\  	LLW  , 	sD   ;A%S 1:T *C%T. 	T!S??T T+*T+.	U#7"UU#F)approx_tokensforcec               0   |st        |dd      syt        |dd      xs d}t        |dd      xs d}t        j                  dt        | dd      xs d	|||       |s	 d
dlm} | j                  |       t        | dd      }|t        |dd
      dz   |_        |xs d
|_        t        t        |      dd      }t        |      r ||d       nt        |d      rd|_        t        |dd      sd|_        d
|_        d|_        d| _        	 t        | dd      rH| j#                  dt        | dd      xs dt        | dd      xs ddd|t        |dd
      nd
d||d       y# t        $ r Y w xY w# t        $ r t        j%                  dd       Y yw xY w)a+  Record a Codex-native context compaction boundary in Hermes state.

    The app-server owns the compacted thread context, so Hermes should not
    rewrite local transcript rows here; state.db records the boundary via the
    session event/usage counters while preserving the visible transcript.
    	compactedF	thread_idNr5   turn_idzKcodex app-server compaction observed: session=%s thread=%s turn=%s force=%srB   noner   )COMPACTION_STATUSr   compression_countr   record_completed_compaction)used_fallback$_verify_compaction_cleared_thresholdTr   event_callbackzsession:compressplatformcodex_app_server)rs   rB   old_session_idin_placerm   runtimeri   rj   z.event_callback error on codex session:compressr2   )r>   rG   infoagent.conversation_compressionrl   _emit_statusrF   rm   last_compression_rough_tokenstypecallablehasattrrp   r<   last_completion_tokensr   _last_compaction_in_placerr   rH   )	rZ   r[   re   rf   ri   rj   rl   r]   record_boundarys	            r   #_record_codex_app_server_compactionr      s    {E:k406BIdIt,2G
KKU|T*4f 	H01  4d;J'.+Q(
(
$ 4A3EA
0 ";T
 O$ Je<Z!GH>BJ;t/6,.J)01J-?CJ<&+E#V5*D1  " 'z4 @ FB")%t"D"J&( % "- *1"$7* 1!*&( c  		\  VEPTUVs%   E# AE2 #	E/.E/2 FF>   	webSearch
fileChangemcpToolCalldynamicToolCallcommandExecutionzhermes-toolsc                   | j                  d      xs d}|dk(  ry|dk(  ry|dk(  r=| j                  d      xs d	}| j                  d
      xs d}|t        k(  r|S d| d| S |dk(  r| j                  d
      xs dS |dk(  ry|xs dS )zSynthetic Hermes tool name for a codex item. Mirrors
    CodexEventProjector so the progress bubble and the projected
    tool_calls entry use the same identifier.r|   r5   r   exec_commandr   apply_patchr   servermcptoolunknownzmcp..r   dynamicr   
web_search)rJ   _INTERNAL_MCP_SERVER)item	item_typer   r   s       r   _codex_item_to_tool_namer   6  s      &BI&&L M!(#,uxx,9))KfXQtf%%%%xx,9,K!	!    c                   | j                  d      xs d}|dk(  r+| j                  d      xs d| j                  d      xs ddS |dk(  rqd| j                  d      xs g D cg c]P  }t        |t              r>|j                  d	      xs i j                  d      xs d
|j                  d      xs ddR c}iS |dv r+| j                  d      xs i }t        |t              r|S d|iS |dk(  rd| j                  d      xs diS i S c c}w )zArgs dict surfaced to tool_progress_callback("tool.started", ...).
    Mirrors the projector's _project_command / _project_file_change /
    _project_mcp_tool_call / _project_dynamic_tool_call shapes.r|   r5   r   commandcwd)r   r   r   changeskindupdatepath)r   r      r   r   	argumentsr   query)rJ   r   r?   )r   r   cargss       r   _codex_item_to_argsr   L  s%     &BI&&88I.4"xx,". 	.L  hhy)/R
 Jq$4G eeFm)r..v6B(UU6](b*
  	
 66xx$*!$-tFK3FFK'*0b11I
s   #AD	c                   | j                  d      xs d}|dk(  r| j                  d      xs d}|r|dd S dS |dk(  r| j                  d      xs g D cg c]4  }t        |t              r"|j                  d	      r|j                  d	      6 }}|syd
j                  |dd       }t	        |      dkD  r|dt	        |      dz
   dz  }|S |dv rC| j                  d      xs i }t        |t              r|sy	 t        j                  |d      dd S |dk(  r| j                  d      xs d}|r|dd S dS yc c}w # t        t        f$ r Y yw xY w)zShort human-readable preview for the tool.started bubble. Returns
    None when no useful preview is available (Hermes' UI tolerates None).r|   r5   r   r   Nx   r   r   r   ,    z, +z morer   r   Fensure_asciir   r   )	rJ   r   r?   joinlenjsondumps	TypeErrorr   )r   r   cmdr   pathspreviewr   r   s           r   _codex_item_to_previewr   b  ss     &BI&&hhy!'Rs4Cy)T)L )-))<)B ;1q$'AEE&M v ; ;))E"1I&u:>SZ!^,E22G66xx$*$%T	::d7== K!'R#uTc{--'; :& 		s   9D79D< <EEc                0   | j                  d      xs d}|dk(  rH| j                  d      xs d}| j                  d      }t        |duxr |dk7        }|rd| d	| }||fS |d
k(  r@| j                  d      xs d}t        | j                  d      xs g       }d| d| d|dvfS |dk(  re| j                  d      }|rdt        j                  |d      dd  dfS | j                  d      }|t        j                  |d      dd dfS ddfS |dk(  r| j                  d      xs g }	t        |	t              r8|	r6t        j                  |	d      dd t        | j                  dd             fS | j                  dd      }
d|
 t        |
       fS y) zReturn (result_text, is_error) for a completed codex tool item.
    Mirrors the projector's tool-result content so the bubble shows the
    same outcome string that ends up in the messages list.r|   r5   r   aggregatedOutputexitCodeNr   z[exit z]
r   rV   r   r   zapply_patch status=r   z
 change(s)>   appliedsuccess	completedr   errorz[error] Fr   i  Tresulti  r   contentItemsr   zsuccess=)r5   F)rJ   r   r   r   r   r   list)r   r   out	exit_codeis_errorrV   nr   r   content_itemsr   s              r   _codex_item_completion_payloadr     s     &BI&&hh)*0bHHZ(		-@)q.A9+S.CH}L (#0y#)r*!&A3j9==
 	
 M!!4::e%@$GHI  (# ! JJvE25D9
 	
')
 	

 %%06BmT*}

=u=etD)T233  ((9d+'#g%666r   c                z     i ddd	 fdd	 fdd
 fdd
 fdd	 fddfd}|S )u  Build an ``on_event`` callback that wires codex app-server JSON-RPC
    notifications into Hermes' gateway UI callbacks.

    Returns a single-argument callable suitable for
    ``CodexAppServerSession(on_event=...)``.

    Translation map:
      * ``item/started`` for tool-shaped items → ``tool_progress_callback(
        "tool.started", name, preview, args)``
      * ``item/completed`` for tool-shaped items → ``tool_progress_callback(
        "tool.completed", name, None, None, duration=..., is_error=...,
        result=...)``
      * ``item/agentMessage/delta`` → ``_fire_stream_delta(text)`` so chat
        adapters can render the assistant's reply as it streams.
      * ``item/reasoning/delta`` → ``_fire_reasoning_delta(text)``
      * ``item/completed`` for ``agentMessage`` →
        ``_emit_interim_assistant_message({"role": "assistant",
        "content": text})``. The gateway's ``already_streamed`` check
        dedupes against any text the stream-delta callback already
        rendered for the same message.

    All callback invocations are guarded — a buggy display callback must
    not tear down the codex turn loop. Errors are logged at DEBUG so the
    notification stream keeps flowing regardless.
    c                t   ddl m} | j                  d      xs d}| j                  d      xs d}|dk(  r	 |d|      S |dk(  r	 |d	|      S |d
k(  r9| j                  d      xs d}| j                  d      xs d} |d| d| |      S |dk(  r!| j                  d      xs d} |d| |      S  |||      S )zDeterministic tool_call id mirroring CodexEventProjector, so a
        live TUI tool card correlates with the same tool call after the
        session is resumed and history is projected.r   )_deterministic_call_ididr5   r|   r   execr   r   r   r   r   r   r   mcp____r   dyn_)&agent.transports.codex_event_projectorr   rJ   )r   namer   item_idr   r   r   s          r   _stable_call_idz;make_codex_app_server_event_bridge.<locals>._stable_call_id  s     	R((4.&BHHV$*	**)&'::$)-AA%XXh'05F88F#0yD)E&D6*BGLL))88F#0yD)D-AA%dG44r   c                   | j                  d      xs d}t        |       }t        |       }|r||t        j                         f|<   t        dd       }|	  |d|t        |       |       t        dd       }|	  | | |      ||       y y # t        $ r t        j                  d|d       Y Fw xY w# t        $ r t        j                  d	|d       Y y w xY w)
Nr   r5   tool_progress_callbackztool.startedz4tool_progress_callback raised on tool.started for %sTr2   tool_start_callbackz!tool_start_callback raised for %s)
rJ   r   r   time	monotonicr>   r   rF   rG   rH   )	r   r   r   r   cbstart_cbr   rZ   starteds	         r   _fire_tool_startedz>make_codex_app_server_event_bridge.<locals>._fire_tool_started  s    ((4.&B'-"4( $dDNN,<=GGU4d;>>4)?)EtL 5"7>t4dDA    J4    7  s$   B =B7 !B43B47!CCc           	     j   | j                  d      xs d}t        |       }j                  |d       }d }| j                  d      }t        |t        t
        f      r|dk\  r|dz  }n|t        j                         |d   z
  }t        |       \  }}t        dd       }|	  |d|d d |||	       t        dd       }	|	&||d   n
t        |       }
	  |	 | |      ||
|       y y # t        $ r t        j                  d
|d       Y Yw xY w# t        $ r t        j                  d|d       Y y w xY w)Nr   r5   
durationMsr   g     @@   r   ztool.completed)durationr   r   z6tool_progress_callback raised on tool.completed for %sTr2   tool_complete_callbackr   z$tool_complete_callback raised for %s)rJ   r   popr   r   r   r   r   r   r>   rF   rG   rH   r   )r   r   r   priorr   codex_msr   r   r   complete_cbr   r   rZ   r   s              r   _fire_tool_completedz@make_codex_app_server_event_bridge.<locals>._fire_tool_completed  s\   ((4.&B'-GT*
 88L)he-(a-&(H~~'%(2H9$?U4d;>#T4$xH e%=tD"$0586I$6ODOD$7tVL #  L4    :D4  s$   "C' D '!D
D!D21D2c                    | j                  d      xs | j                  d      xs d}t        |t              r|sy t        dd       }|y 	  ||       y # t        $ r t
        j                  dd       Y y w xY w)Ndeltatextr5   _fire_stream_deltaz_fire_stream_delta raisedTr2   rJ   r   r   r>   rF   rG   rH   paramsr   fnrZ   s      r   _fire_text_deltaz<make_codex_app_server_event_bridge.<locals>._fire_text_delta  sy    zz'">fjj&8>B$$DU0$7:	EtH 	ELL4tLD	E   A  A:9A:c                    | j                  d      xs | j                  d      xs d}t        |t              r|sy t        dd       }|y 	  ||       y # t        $ r t
        j                  dd       Y y w xY w)Nr   r   r5   _fire_reasoning_deltaz_fire_reasoning_delta raisedTr2   r   r   s      r   r   zAmake_codex_app_server_event_bridge.<locals>._fire_reasoning_delta+  sy    zz'">fjj&8>B$$DU3T::	HtH 	HLL7$LG	Hr   c                   | j                  d      xs d}t        |t              r|j                         sy t	        dd      sy t	        dd       }|y 	  |d|d       y # t
        $ r t        j                  dd	       Y y w xY w)
Nr   r5   show_commentaryT_emit_interim_assistant_message	assistant)rolecontentz&_emit_interim_assistant_message raisedr2   )rJ   r   r   stripr>   rF   rG   rH   )r   r   emitrZ   s      r   _fire_agent_message_completedzImake_codex_app_server_event_bridge.<locals>._fire_agent_message_completed7  s    xx%2$$DJJL u/6u?F<	+$78 	LL84  	s   A# # BBc                   t        | t              sy | j                  d      xs d}| j                  d      xs i }t        |t              si }|dk(  r	 |       y |dv r	 |       y |j                  d      }t        |t              sy |j                  d      xs d}|dk(  r|t        v r	 	|       y |d	k(  r |t        v r	 |       y |d
k(  r	 |       y y y )Nmethodr5   r   zitem/agentMessage/delta>   item/reasoning/deltaitem/reasoning/summaryDeltar   r|   zitem/startedzitem/completedagentMessage)r   r?   rJ   _CODEX_TOOL_ITEM_TYPES)
noter   r   r   r   r   r   r   r   r   s
        r   on_eventz4make_codex_app_server_event_bridge.<locals>.on_eventJ  s    $%(#)r(#)r&$'F..V$LL!&)zz&!$%HHV$*	^#	5K(Kt$%%22$T*n,-d3 - &r   )r   r?   r   r   returnr   )r   r?   r   None)r   r?   r   r   )r   r?   r   r    )	rZ   r   r   r   r   r   r   r   r   s	   ` @@@@@@@r   "make_codex_app_server_event_bridger    s<    : 35G5*8!F
E
H&4 44 Or   )should_review_memoryc          
        ddl m}m} t        | d      r| j                  eddlm} t        | dd      xs t         |             }		 ddl	m
}
  |
       }d}	 dd	lm}  |       } ||	| |||      t!        |             | _        	 | j                  j#                  |      }t        |dd      rBt        j)                  d|j*                         	 | j                  j'                          d| _        |j,                  r:|j/                  |j,                         t        | dd      	 | j1                  |       t        | dd      |j2                  z   | _        t7        | |       t9        | |      }d}d}| j:                  dkD  r0| j4                  | j:                  k\  rd| j<                  v r	d}d| _        |j>                  s,|j*                   	 | jA                  ||jB                  d|       |jB                  r.|j>                  s"|s|r	 | jE                  tG        |      ||       |jB                  |||j>                   xr |j*                  du |j>                  xs |j*                  du|j*                  d|jH                  |jJ                  d	|S # t        $ r d}Y Iw xY w# t        $ r t        j                  d
d       Y `w xY w# t        $ rg}t        j%                  d       	 | j                  j'                          n# t        $ r Y nw xY wd| _        d| d|dddt        |      dcY d}~S d}~ww xY w# t        $ r Y \w xY w# t        $ r t        j                  dd       Y 6w xY w# t        $ r t        j                  dd       Y w xY w# t        $ r t        j                  dd       Y w xY w) aN  Codex app-server runtime path. Hands the entire turn to a `codex
    app-server` subprocess and projects its events back into Hermes'
    messages list so memory/skill review keep working.

    Called from run_conversation() when agent.api_mode == "codex_app_server".
    Returns the same dict shape as the chat_completions path.
    r   )CodexAppServerSession_ServerRequestRouting_codex_sessionN)resolve_agent_cwdsession_cwd)_get_approval_callbackF)is_approval_bypass_activezLcodex app-server: approval-bypass lookup failed; keeping fail-closed defaultTr2   )auto_approve_execauto_approve_apply_patch)r   approval_callbackrequest_routingr   )
user_inputzcodex app-server turn failedzCodex app-server turn failed: z:. Fall back to default runtime with `/codex-runtime auto`.)final_responsemessages	api_callsr   partialr   should_retirez1codex app-server session retired (turn error: %s)rA   z/codex app-server projected-message flush failed_iters_since_skillr   skill_manage)original_user_messager  interruptedr  zexternal memory sync raised)messages_snapshotreview_memoryreview_skillszbackground review spawn raised)	r  r  r  r   r  r   agent_persistedcodex_thread_idcodex_turn_id)&)agent.transports.codex_app_server_sessionr  r  r~   r  agent.runtime_cwdr  r>   r   tools.terminal_toolr
  rF   tools.approvalr  rG   rH   r  run_turn	exceptionclosewarningr   projected_messagesextend_flush_messages_to_session_dbtool_iterationsr  r   rd   _skill_nudge_intervalvalid_tool_namesr  _sync_external_memory_for_turn
final_text_spawn_background_reviewr   ri   rj   )rZ   user_messager  r  effective_task_idr  r  r  r  r   r
  r  auto_approve_requestsr  r[   r^   usage_resultr  should_review_skillss                      r   run_codex_app_server_turnr6  g  s     5*+u/C/C/K7e]D1MS9J9L5M	%B 6 8 !&		@$=$?!  5/1"7)> 8> 

##,,,E6 t_e,?JJ	
	  &&(  $
 //0 5-.:33H=  	+Q/$2F2FF 
 (t41%>LI !##a'$$(C(CCe444##$  

 2	G00&;#!!	 1  	  !%9	J**"&x.22 +  //)))@djjD.@##=tzz'=  >>)* + O  	% $	%"  	LL.  	<  
78	  &&( 		# 1 6K L !X

 
	

B  		2  E!  T  	GLL6LF	G"  	JLL9DLI	Js   J J! K L; 'M M2 N JJ! KK	L8L3'LL3	LL3LL3-L83L8;	MM M/.M/2 NN N=<N=>   response.failedresponse.completedresponse.incompletec                p    t        | |d      }|"t        | t              r| j                  ||      }||S |S )zSField access that handles both attr-style (SDK objects) and dict (raw JSON) events.Nr>   r   r?   rJ   )eventr   defaultr   s       r   _event_fieldr>  g  s>    E4&E}E40		$(%5272r   c                p    t        | |d      }|"t        | t              r| j                  ||      }||S |S )zGField access for nested Response items (attr-style SDK object or dict).Nr;  )r   r   r=  r   s       r   _item_fieldr@  o  s>    D$%E}D$/w'%5272r   c                     ddl m} t         d      d
 fd} |d      }|t        |t              st	        |      }|xs dj                         xs d} || |d       |d      	      )a;  Raise a ``_StreamErrorEvent`` from a ``type=error`` SSE frame.

    The Responses spec puts the failure details at the top level of the
    frame (``{"type": "error", "code": ..., "message": ..., "param": ...}``),
    but the official OpenAI SDK and several OpenAI-compatible proxies wrap
    them in an HTTP-style nested envelope instead
    (``{"type": "error", "error": {"code": ..., "message": ..., "param": ...}}``).
    Read the top-level fields first, then fall back to the nested envelope so
    the error classifier sees the provider's real code/message (rate-limit vs
    context-overflow vs entitlement) rather than the generic placeholder.
    Port of anomalyco/opencode#36130.

    Imported lazily so this module stays importable from places that don't
    pull in ``run_agent`` (e.g. plugin code, doc tools).
    r   )_StreamErrorEventr   c                @    t        |       }|t        |       }|S N)r>  r@  )r   r   r<  nesteds     r   _error_fieldz)_raise_stream_error.<locals>._error_field  s*    UD)=V/-Er   messagezstream emitted error eventcodeparam)rH  rI  )r   r   r   r   )	run_agentrB  r>  r   r   r   )r<  rB  rF  raw_messagerG  rE  s   `    @r   _raise_stream_errorrL  w  sx      ,%)F y)Kz+s'C+&::AACcGcG
&!7# r   )on_text_deltaon_reasoning_deltaon_commentary_messageon_first_deltar   interrupt_checkc          
     
   g }g }	d}
d}d}g }d}d}d}d}d}d}| D ]  }|		  ||       |
 |       r n|t        |dd      }t        |t              sd}|d	k(  rt        |       |d
k(  rut        |d      }t        |dd      }|dk(  rEt        |dd      }t        |t              r|j                         j                         nd}|dk(  rg }nd}dt        |      v rd}
d|v s|dk(  rvt        |dd      }|r$|dk(  r|j                  |       |M|K	  ||       nA|r|dk(  r|8	  ||       n.|r,|	j                  |       |
s|sd}|	  |        |		  ||       Cd|v rd}
d|v r d|v rt        |dd      }|r|		  ||       m|dk(  rt        |d      }||j                  |       t        |dd      }t        |t              r|j                         j                         nd}|dk(  rs|qdj                  |      j                         }|sCt        |dg       }t        |t              r&dj                  d |D              j                         }|r		  ||       g }H|t        v sRd}t        |d      }|t!        |dd      }|!t        |t"              r|j%                  d      }t!        |dd      } | !t        |t"              r|j%                  d      } | }t!        |dd      }!|!!t        |t"              r|j%                  d      }!t        |!t              r|!}|d k(  r0t!        |d!d      }|!t        |t"              r|j%                  d!      }|d"k(  r0t!        |d	d      }|!t        |t"              r|j%                  d	      }|d#k(  r|xs d}n|d k(  r|xs d$}n|d"k(  r|xs d%} n |rt        |      }"n4|	r0|
s.dj                  |	      }#t'        dd&dt'        d'|#(      g)      g}"ng }"|s|"st)        d*      dj                  |	      }$t'        |"|$||||||+      }%|%S # t         t        f$ r  t        $ r t        j	                  dd       Y 1w xY w# t        $ r t        j	                  dd       Y &w xY w# t        $ r t        j	                  dd       Y Mw xY w# t        $ r t        j	                  dd       Y w xY w# t        $ r t        j	                  dd       Y w xY w# t        $ r t        j	                  dd       Y w xY w# t        $ r t        j	                  dd       Y w xY w),u  Consume a Codex Responses SSE event stream and return a final response.

    The returned object is a ``SimpleNamespace`` shaped like the SDK's typed
    ``Response`` for the fields downstream code actually reads:

    * ``output``: list of output items, assembled from ``response.output_item.done``.
      For tool-call turns this contains the function_call items; for plain-text
      turns it contains a synthesized ``message`` item built from streamed deltas
      if no message item was emitted directly.
    * ``output_text``: assembled text from ``response.output_text.delta`` deltas.
    * ``usage``: copied from the terminal event's ``response.usage`` (when present).
    * ``status``: ``completed`` / ``incomplete`` / ``failed`` (or ``completed`` if
      the stream ended without a terminal frame but produced content).
    * ``id``: ``response.id`` when present.
    * ``incomplete_details``: passed through for ``response.incomplete`` frames.
    * ``error``: passed through for ``response.failed`` frames.
    * ``model``: from kwargs (the wire model name is not authoritative).

    Critically, we never read ``response.output`` from the terminal event for
    content reconstruction — only ``usage``, ``status``, ``id``.  That field
    being ``null`` / ``[]`` / missing is fine.

    Callbacks:

    * ``on_text_delta(str)`` — fires per ``response.output_text.delta``, suppressed
      once a function_call event is seen (so tool-call turns don't bleed text
      into the chat).
    * ``on_reasoning_delta(str)`` — fires per ``response.reasoning.*.delta`` and
      ``phase=analysis`` message deltas. When no dedicated commentary callback
      is supplied, commentary also uses this legacy fallback.
    * ``on_commentary_message(str)`` — fires once per completed
      ``phase=commentary`` message, before any following tool item executes.
    * ``on_first_delta()`` — one-shot, fires on the first text delta only.
    * ``on_event(event)`` — fires for every event before any other processing.
      Used for watchdog activity, debug logging, anything wire-shape-agnostic.
    * ``interrupt_check()`` — returns True to break the loop early.
    FNr   z!Codex stream on_event hook raisedTr2   r|   r5   r   zresponse.output_item.addedr   rG  phase
commentaryfunction_callzoutput_text.deltazresponse.output_text.deltar   z&Codex stream on_reasoning_delta raisedanalysisz"Codex stream on_first_delta raisedz!Codex stream on_text_delta raised	reasoningzresponse.output_item.doner   c              3  p   K   | ].  }t        |d d      dk(  rt        t        |dd      xs d       0 yw)r|   r5   output_textr   N)r@  r   ).0parts     r   	<genexpr>z._consume_codex_event_stream.<locals>.<genexpr>@  s<      6$(#.tVR#@M#Q !$Kfb$A$GR H6s   46z)Codex stream on_commentary_message raisedresponser\   r   rV   r9  incomplete_detailsr7  r8  
incompletefailedr   rY  )r|   r   )r|   r   rV   r   z7Codex Responses stream did not emit a terminal response)outputrY  r\   rV   r   r   r^  r   )TimeoutErrorInterruptedErrorrF   rG   rH   r>  r   r   rL  r@  r   lowerappendr   r   _TERMINAL_EVENT_TYPESr>   r?   rJ   r   RuntimeError)&
event_iterr   rM  rN  rO  rP  r   rQ  collected_output_itemscollected_text_deltashas_tool_callsfirst_delta_firedactive_message_phasecommentary_text_deltasterminal_statusterminal_usageterminal_response_idterminal_incomplete_detailsterminal_errorsaw_terminalr<  
event_typer   r   rS  
delta_textreasoning_text	done_item
done_phasecommentary_textcontent_partsresp_objridrstatusra  	assembledassembled_textfinals&                                         r   _consume_codex_event_streamr    s/   ` )+')N'+(*&ON $'+NL V	Q &?+<!%4
*c*J  & 55v.D#D&"5II%#D'48@J5RU@Vu{{}':':'<\`$'<7-/*'+$#i.0!%*,
>Z0Z%eWb9J2lB&--j9 )05G5S^*:6  4
 B%1^*:6 %,,Z8%,,0))5b . 0 %0])*5 j(!N *$J)>)%"=N"4"@Z&~6 44$UF3I$&--i8(GTB
;EjRU;VZ--/557\`
-2G2S&(gg.D&E&K&K&MO*(3Iy"(M%mT:.0gg 6,96 / $eg	 ,
 '1/B .0*..L#E:6H#!(7D!A!)j4.H%-\\'%:Nhd3;:h#=",,t,C'*$!(Hd;?z(D'A&ll84Ggs+&-O!6629(DXZ^2_/2:z(TX?Y6>llCW6X3!22%,Xw%EN%-*Xt2L)1g)>11"1"@[44"1"A\00"1"=XmVv ,-	~GG12	!$-iHI	
   E
 	
 WW23N"6	E Lw !"23   Q @4PQ\ % ^%MX\]^ % ^%MX\]^ $- b &-Q\` ab
  ) ]"LL)LW[L\] ! ZLL!ITXLYZ.  ) "LL K)- ) s}   P)Q%R
R-S>S;T")/QQ RR R*)R*- SS S87S8; TT" UUc                >    ddl }|xs  j                  d      }d}g  _        d fd}d fd}d fd}	d fd	}
t        |dz         D ]m  } j                  rt        d
      t              }d|d<   	  |j                  j                  di |}t!               }|fd fd}	 t#        |d      r3t#        |d      s'|t%        |dd      }t'        |      r	  |        c S c S 	 t+        |j-                  d      ||t%         dd      t%         dd      r|	nd||
|      }|j.                  dv r`t        j1                  d|j.                  |j2                  |j4                  t7        d  j                  D               j                                |t%        |dd      }t'        |      r	  |        c S c S  y# |j                  |j                  |j                  t        f$ r>}||k  r3t        j                  d|dz   |dz    j                         |       Y d}~Ղ d}~ww xY w# t(        $ r Y c S w xY w# |j                  |j                  |j                  t        f$ rp}||k  ret        j                  d|dz   |dz    j                         |       Y d}~t%        |dd      }t'        |      sj	  |        t# t(        $ r Y w xY w d}~ww xY w# t(        $ r Y c S w xY w# t%        |dd      }t'        |      r	  |        w # t(        $ r Y w w xY ww xY w)u  Execute one streaming Responses API request and return the final response.

    Uses ``responses.create(stream=True)`` (low-level raw event iteration)
    rather than the high-level ``responses.stream(...)`` helper.  This makes
    us structurally immune to backend drift in the ``response.completed``
    payload shape — we never let the SDK reconstruct a typed object from
    the terminal event's ``output`` field.
    r   Ncodex_stream_direct)reasonr   c                ^    j                   j                  |        j                  |        y rD  )_codex_streamed_text_partsre  r   r   rZ   s    r   _on_text_deltaz(run_codex_stream.<locals>._on_text_delta  s%    ((//5  &r   c                (    j                  |        y rD  )r   r  s    r   _on_reasoning_deltaz-run_codex_stream.<locals>._on_reasoning_delta  s    ##D)r   c                (    j                  |        y rD  )_fire_streamed_codex_commentaryr  s    r   _on_commentary_messagez0run_codex_stream.<locals>._on_commentary_message  s    --d3r   c                Z    t        j                          _        j                  d       y )Nzreceiving stream response)r   _codex_stream_last_event_ts_touch_activity)r<  rZ   s    r   	_on_eventz#run_codex_stream.<locals>._on_event  s     ,0IIK)9:r   z+Agent interrupted before Codex stream retryTstreamzLCodex Responses stream connect failed (attempt %s/%s); retrying. %s error=%sc                    j                   ryt        |       s't        j                  dj	                  dd             yy)NTz~Codex streaming attempt superseded by a newer stream; stopping consumption to preserve the single-writer invariant (model=%s).r   r   F)_interrupt_requestedr
   rG   r'  rJ   )_tokrZ   
api_kwargss    r   _interrupt_or_supersededz2run_codex_stream.<locals>._interrupt_or_superseded  sB    ))+E48, NN7I6	 r   ra  __iter__r&  r   interim_assistant_callbackr   )r   rM  rN  rO  rP  r   rQ  z\Codex Responses stream transport failed mid-iteration (attempt %s/%s); retrying. %s error=%s>   r`  r_  zbCodex Responses stream terminal status=%s (incomplete_details=%s, error=%s, streamed_chars=%d). %sc              3  2   K   | ]  }t        |        y wrD  )r   )rZ  ps     r   r\  z#run_codex_stream.<locals>.<genexpr>
  s     I1AIs   )r   r   r   r   r<  r   r   r   r  )r   r   )httpx_ensure_primary_openai_clientr  ranger  rc  r?   	responsescreateRemoteProtocolErrorReadTimeoutConnectErrorConnectionErrorrG   rH   _client_log_contextr	   r~   r>   r}   rF   r  rJ   rV   r'  r^  r   sum)rZ   r  clientrP  _httpxactive_clientmax_stream_retriesr  r  r  r  attemptstream_kwargsevent_streamr^   _writer_tokenr  close_fnr  s   ``                 r   run_codex_streamr    sJ    _eAAI^A_M-/E$'*4;
 +a/0 [%%"#PQQZ("&h
	9=2299JMJL" ,E2*7 	3	 |X.w|Z7X#T |Wd;H!J "S3 $..1"0': $E+GNZ '/@$ G /
 "#1&$<8 ||77OLL%":":EKKI(H(HII--/ |Wd;H!J "o[ **F,>,>@S@SUde 	++baK!3a!7--/
 	b ! 5 ..0B0BFDWDWYhi 	//LLA!%7!%;113S	  |Wd;H!J  # 	4 ! 	 |Wd;H!J   "s   8F>)K*H*)>H;'A/K*/K>+H')2H"!H""H'*	H87H8;+K&2KK*7K  	KKKKK*	K'&K'*LLL	L	LL	Lc                    t        | ||      S )a  Backward-compatible alias for the unified event-driven path.

    Historically this was the fallback when the SDK's high-level
    ``responses.stream(...)`` helper raised on shape drift.  The primary
    path now does exactly what the fallback did, so this just forwards.
    Kept as a public symbol because tests and a small number of call sites
    still reference it by name.
    )r  )r  )rZ   r  r  s      r    run_codex_create_stream_fallbackr    s     E:f==r   )r6  r  r  r  r  )r   r   r   r   )r   zdict[str, Any])re   z
int | Nonerf   r   r   r   )r   r?   r   r   )r   r?   r   r?   )r   r?   r   r   )r   r?   r   ztuple[str, bool])r   zCallable[[dict], None])r1  r   r  r   r  zList[Dict[str, Any]]r2  r   r  r   r   zDict[str, Any]rD  )r<  r   r   r   r=  r   r   r   )r   r   r   r   r=  r   r   r   r  )rh  r   r   r   r   r   )NN)r  r?   r  r   )(__doc__
__future__r   r   loggingosr   typesr   typingr   r   r   r   agent.stream_single_writerr	   r
   	getLogger__name__rG   r   rd   r   	frozensetr   r   r   r   r   r   r  r6  rf  r>  r@  rL  r  r  r  __all__r  r   r   <module>r     sf    #   	  ! , , T			8	$Nj !%O 	O
 O 
ON #U  & ",,<)XxD "'_ _ 	_
 #_ _ _ _r " #  33"R || | |~zz	>r   