refactor(agent/codex*,bedrock_adapter): drop intra-function blank lines (AST-identical, -117 LOC)
This commit is contained in:
@@ -154,7 +154,6 @@ class BedrockOpenAISigV4Auth(httpx.Auth):
|
||||
import botocore.session
|
||||
from botocore.auth import SigV4Auth
|
||||
from botocore.awsrequest import AWSRequest
|
||||
|
||||
credentials = botocore.session.get_session().get_credentials()
|
||||
if credentials is None:
|
||||
raise RuntimeError(
|
||||
@@ -215,7 +214,6 @@ _STALE_LIB_MODULE_PREFIXES = ("urllib3.", "botocore.", "boto3.")
|
||||
def _stale_error_types() -> tuple:
|
||||
"""botocore + urllib3 transport-failure exception classes (best-effort import)."""
|
||||
import importlib
|
||||
|
||||
types: list = []
|
||||
for module, names in (
|
||||
("botocore.exceptions", ("ConnectionError", "HTTPClientError")),
|
||||
@@ -656,19 +654,16 @@ def _assistant_blocks(msg: Dict, content) -> List[Dict]:
|
||||
content_blocks = _replay_ordered_blocks(ordered_blocks)
|
||||
if content_blocks:
|
||||
return content_blocks
|
||||
|
||||
content_blocks = []
|
||||
for detail in (msg.get("reasoning_details") or []):
|
||||
if isinstance(detail, dict) and detail.get("type") == "redacted_thinking":
|
||||
redacted = _decode_redacted(detail.get("data") or detail.get("redactedContentBase64"))
|
||||
if redacted is not None:
|
||||
content_blocks.append({"reasoningContent": {"redactedContent": redacted}})
|
||||
|
||||
if isinstance(content, str) and content.strip():
|
||||
content_blocks.append({"text": content})
|
||||
elif isinstance(content, list):
|
||||
content_blocks.extend(_convert_content_to_converse(content))
|
||||
|
||||
for tc in (msg.get("tool_calls", []) or []):
|
||||
fn = tc.get("function", {})
|
||||
content_blocks.append(_tool_use_block(tc.get("id", ""), fn.get("name", ""), _parse_tool_args(fn.get("arguments", "{}"))))
|
||||
@@ -691,7 +686,6 @@ def convert_messages_to_converse(messages: List[Dict]) -> Tuple[Optional[List[Di
|
||||
converse_msgs[-1]["content"].extend(blocks)
|
||||
else:
|
||||
converse_msgs.append({"role": role, "content": blocks})
|
||||
|
||||
for msg in messages:
|
||||
role = msg.get("role", "")
|
||||
content = msg.get("content")
|
||||
@@ -707,7 +701,6 @@ def convert_messages_to_converse(messages: List[Dict]) -> Tuple[Optional[List[Di
|
||||
append_turn("assistant", _assistant_blocks(msg, content) or [dict(_PLACEHOLDER_BLOCK)])
|
||||
elif role == "user":
|
||||
append_turn("user", _convert_content_to_converse(content))
|
||||
|
||||
if converse_msgs and converse_msgs[0]["role"] != "user":
|
||||
converse_msgs.insert(0, {"role": "user", "content": [dict(_PLACEHOLDER_BLOCK)]})
|
||||
if converse_msgs and converse_msgs[-1]["role"] != "user":
|
||||
@@ -881,7 +874,6 @@ def stream_converse_with_callbacks(
|
||||
if encoded:
|
||||
parts.add_redacted(encoded)
|
||||
current_block({"reasoningContent": {}}).setdefault("reasoningContent", {})["redactedContentBase64"] = encoded
|
||||
|
||||
for event in event_stream.get("stream", []):
|
||||
if on_event is not None:
|
||||
try:
|
||||
@@ -890,7 +882,6 @@ def stream_converse_with_callbacks(
|
||||
pass
|
||||
if on_interrupt_check and on_interrupt_check():
|
||||
break
|
||||
|
||||
if "contentBlockStart" in event:
|
||||
start_event = event["contentBlockStart"]
|
||||
current_block_index = start_event.get("contentBlockIndex", len(stream_blocks))
|
||||
@@ -902,7 +893,6 @@ def stream_converse_with_callbacks(
|
||||
stream_blocks[current_block_index] = _tool_use_block(current_tool["toolUseId"], current_tool["name"], {})
|
||||
if on_tool_start:
|
||||
on_tool_start(current_tool["name"])
|
||||
|
||||
elif "contentBlockDelta" in event:
|
||||
delta = event["contentBlockDelta"].get("delta", {})
|
||||
if "text" in delta:
|
||||
@@ -916,7 +906,6 @@ def stream_converse_with_callbacks(
|
||||
current_tool["input_json"] += delta["toolUse"].get("input", "")
|
||||
elif "reasoningContent" in delta:
|
||||
on_reasoning(delta["reasoningContent"])
|
||||
|
||||
elif "contentBlockStop" in event:
|
||||
if current_tool is not None:
|
||||
input_dict = _parse_tool_args(current_tool["input_json"]) if current_tool["input_json"] else {}
|
||||
@@ -926,17 +915,14 @@ def stream_converse_with_callbacks(
|
||||
current_tool = None
|
||||
else:
|
||||
flush_text()
|
||||
|
||||
elif "messageStop" in event:
|
||||
stop_reason = event["messageStop"].get("stopReason", "end_turn")
|
||||
|
||||
elif "metadata" in event:
|
||||
meta_usage = event["metadata"].get("usage", {})
|
||||
usage_data = {
|
||||
key: meta_usage.get(key, 0)
|
||||
for key in ("inputTokens", "outputTokens", "cacheReadInputTokens", "cacheWriteInputTokens")
|
||||
}
|
||||
|
||||
flush_text()
|
||||
return parts.build([stream_blocks[i] for i in sorted(stream_blocks)], usage_data, stop_reason, "")
|
||||
|
||||
@@ -966,17 +952,13 @@ def build_converse_kwargs(
|
||||
|
||||
def cache_here(placement: str) -> bool:
|
||||
return cache_enabled and cache_point_allowed(model, placement)
|
||||
|
||||
inference_config: Dict[str, Any] = {}
|
||||
if max_tokens is not None:
|
||||
inference_config["maxTokens"] = max_tokens
|
||||
kwargs: Dict[str, Any] = {"modelId": model, "messages": converse_messages, "inferenceConfig": inference_config}
|
||||
|
||||
if system_prompt:
|
||||
kwargs["system"] = system_prompt + [dict(_CACHE_POINT)] if cache_here("system") else system_prompt
|
||||
|
||||
from agent.anthropic_adapter import _forbids_sampling_params
|
||||
|
||||
if not _forbids_sampling_params(model):
|
||||
if temperature is not None:
|
||||
inference_config["temperature"] = temperature
|
||||
@@ -984,7 +966,6 @@ def build_converse_kwargs(
|
||||
inference_config["topP"] = top_p
|
||||
if stop_sequences:
|
||||
inference_config["stopSequences"] = stop_sequences
|
||||
|
||||
converse_tools = convert_tools_to_converse(tools) if tools else []
|
||||
if converse_tools:
|
||||
# Non-tool-calling models (e.g. DeepSeek R1) reject toolConfig with a
|
||||
@@ -998,12 +979,10 @@ def build_converse_kwargs(
|
||||
"Model %s does not support tool calling — tools stripped. "
|
||||
"The agent will operate in text-only mode.", model
|
||||
)
|
||||
|
||||
if cache_here("messages") and len(converse_messages) >= 2:
|
||||
content = converse_messages[-2].get("content")
|
||||
if isinstance(content, list) and content:
|
||||
content.append(dict(_CACHE_POINT))
|
||||
|
||||
if guardrail_config:
|
||||
kwargs["guardrailConfig"] = guardrail_config
|
||||
if not inference_config:
|
||||
@@ -1100,7 +1079,6 @@ def _list_inference_profiles(client, filter_set: set, models: List[Dict[str, Any
|
||||
next_token = response.get("nextToken")
|
||||
if not next_token:
|
||||
break
|
||||
|
||||
seen_ids = {m["id"].lower() for m in models}
|
||||
for profile in profiles:
|
||||
profile_id = (profile.get("inferenceProfileId") or "").strip()
|
||||
@@ -1123,18 +1101,15 @@ def discover_bedrock_models(region: str, provider_filter: Optional[List[str]] =
|
||||
Returns [] when the client cannot be built.
|
||||
"""
|
||||
import time
|
||||
|
||||
cache_key = f"{region}:{','.join(sorted(provider_filter or []))}"
|
||||
cached = _discovery_cache.get(cache_key)
|
||||
if cached and (time.time() - cached["timestamp"]) < _DISCOVERY_CACHE_TTL_SECONDS:
|
||||
return cached["models"]
|
||||
|
||||
try:
|
||||
client = _get_bedrock_control_client(region)
|
||||
except Exception as e:
|
||||
logger.warning("Failed to create Bedrock client for model discovery: %s", e)
|
||||
return []
|
||||
|
||||
models: List[Dict[str, Any]] = []
|
||||
filter_set = {f.lower() for f in (provider_filter or [])}
|
||||
try:
|
||||
@@ -1145,7 +1120,6 @@ def discover_bedrock_models(region: str, provider_filter: Optional[List[str]] =
|
||||
_list_inference_profiles(client, filter_set, models)
|
||||
except Exception as e:
|
||||
logger.debug("Skipping inference profile discovery: %s", e)
|
||||
|
||||
models.sort(key=lambda m: (0 if m["id"].startswith("global.") else 1, m["name"].lower()))
|
||||
_discovery_cache[cache_key] = {"timestamp": time.time(), "models": models}
|
||||
return models
|
||||
@@ -1220,7 +1194,6 @@ def probe_bedrock_context_length(model_id: str, region: str) -> Optional[int]:
|
||||
except Exception as exc: # boto3 missing / credential resolution failure
|
||||
logger.debug("Bedrock context probe skipped for %s: %s", model_id, exc)
|
||||
return None
|
||||
|
||||
last_error = ""
|
||||
for tier_tokens in _BEDROCK_PROBE_TIERS:
|
||||
oversized = "data " * int(tier_tokens / _WORDS_PER_TOKEN)
|
||||
@@ -1242,7 +1215,6 @@ def probe_bedrock_context_length(model_id: str, region: str) -> Optional[int]:
|
||||
logger.info("Probed Bedrock context window for %s: %s tokens", model_id, f"{limit:,}")
|
||||
return limit
|
||||
# Opaque server error / auth / throttle at this tier — try the next.
|
||||
|
||||
logger.debug("Bedrock context probe for %s returned no parseable limit: %s", model_id, last_error[:200])
|
||||
return None
|
||||
|
||||
|
||||
@@ -43,7 +43,6 @@ def codex_cloudflare_headers(access_token: str, *, base_url: str = CODEX_AUX_BAS
|
||||
"""
|
||||
if is_official_codex_base_url(base_url):
|
||||
from hermes_cli import __version__
|
||||
|
||||
headers = {"User-Agent": f"HermesAgent/{__version__}", "originator": "hermes-agent"}
|
||||
else:
|
||||
headers = {"User-Agent": "codex_cli_rs/0.0.0 (Hermes Agent)", "originator": "codex_cli_rs"}
|
||||
|
||||
@@ -126,11 +126,9 @@ def _neutralize_harmony_tokens(text: str) -> str:
|
||||
"""Keep Harmony source readable without emitting reserved wire tokens."""
|
||||
if not text or "<" not in text or "|" not in text:
|
||||
return text
|
||||
|
||||
replacement = rf"<{_FULLWIDTH_PIPE}\1{_FULLWIDTH_PIPE}>"
|
||||
if not any(unicodedata.category(char) == "Cf" for char in text):
|
||||
return _HARMONY_CONTROL_TOKEN_RE.sub(replacement, text)
|
||||
|
||||
# The backend strips Unicode format controls (e.g. U+200B) before its
|
||||
# reserved-token check, so match on the visible text and rewrite the
|
||||
# original spans — any Cf-hidden variant is neutralized the same way.
|
||||
@@ -324,7 +322,6 @@ def _derive_responses_function_call_id(call_id: str, response_item_id: Optional[
|
||||
"""Build a valid Responses `function_call.id` (must start with `fc_`)."""
|
||||
if isinstance(response_item_id, str) and response_item_id.strip().startswith("fc_"):
|
||||
return response_item_id.strip()
|
||||
|
||||
source = (call_id or "").strip()
|
||||
sanitized = re.sub(r"[^A-Za-z0-9_-]", "", source)
|
||||
for candidate in (source, sanitized):
|
||||
@@ -494,7 +491,6 @@ def _tool_output_item(msg: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
call_id = raw_tool_call_id.strip()
|
||||
if not _nonblank(call_id):
|
||||
return None
|
||||
|
||||
# ``output`` may be a string or an ``input_text``/``input_image`` array.
|
||||
tool_content = msg.get("content")
|
||||
output_value: Any = (
|
||||
@@ -546,7 +542,6 @@ def _chat_messages_to_responses_input(
|
||||
def emit(new_items: List[Dict[str, Any]], msg: Dict[str, Any]) -> None:
|
||||
items.extend(new_items)
|
||||
item_sources.extend([msg] * len(new_items))
|
||||
|
||||
for msg in messages:
|
||||
if not isinstance(msg, dict):
|
||||
continue
|
||||
@@ -558,7 +553,6 @@ def _chat_messages_to_responses_input(
|
||||
continue
|
||||
if role not in {"user", "assistant"}:
|
||||
continue
|
||||
|
||||
content = msg.get("content", "")
|
||||
content_parts = _chat_content_to_responses_parts(content, role=role) # [] unless a list
|
||||
if isinstance(content, list):
|
||||
@@ -566,11 +560,9 @@ def _chat_messages_to_responses_input(
|
||||
content_text = "".join(p["text"] for p in content_parts if p["type"] == text_type)
|
||||
else:
|
||||
content_text = _str_or_empty(content)
|
||||
|
||||
if role == "user":
|
||||
emit([{"role": role, "content": content_parts or content_text}], msg)
|
||||
continue
|
||||
|
||||
reasoning_items = [] if not replay_encrypted_reasoning else _replay_reasoning_items(
|
||||
msg, seen_item_ids=seen_item_ids, current_issuer_kind=current_issuer_kind,
|
||||
native_compaction_eligible=native_compaction_eligible,
|
||||
@@ -578,7 +570,6 @@ def _chat_messages_to_responses_input(
|
||||
emit(reasoning_items, msg)
|
||||
message_items = _replay_message_items(msg, is_github_responses=is_github_responses)
|
||||
emit(message_items, msg)
|
||||
|
||||
if not message_items:
|
||||
if content_parts:
|
||||
emit([{"role": "assistant", "content": content_parts}], msg)
|
||||
@@ -587,9 +578,7 @@ def _chat_messages_to_responses_input(
|
||||
elif reasoning_items:
|
||||
# Every reasoning item needs a following item (else missing_following_item).
|
||||
emit([{"role": "assistant", "content": ""}], msg)
|
||||
|
||||
emit(_replay_tool_call_items(msg, start_index=len(items)), msg)
|
||||
|
||||
# Native server-side compaction renders nothing placed before a compaction
|
||||
# item, so pre-checkpoint history is dead upload weight and the user's
|
||||
# plaintext asks / merged local summaries silently vanish. Keep the newest
|
||||
@@ -597,9 +586,7 @@ def _chat_messages_to_responses_input(
|
||||
# messages within a token budget, leave the tail untouched.
|
||||
if not native_compaction_eligible:
|
||||
return items
|
||||
|
||||
from agent.native_compaction import prune_pre_checkpoint_items
|
||||
|
||||
return prune_pre_checkpoint_items(items, item_sources=item_sources)
|
||||
|
||||
|
||||
@@ -622,7 +609,6 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags:
|
||||
``https://evil.com/models.github.ai`` must not classify as GitHub.
|
||||
"""
|
||||
from utils import base_url_hostname
|
||||
|
||||
provider = getattr(agent, "provider", None)
|
||||
base_url = str(getattr(agent, "base_url", "") or "")
|
||||
hostname = str(getattr(agent, "_base_url_hostname", "") or "").lower() or base_url_hostname(base_url)
|
||||
@@ -630,7 +616,6 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags:
|
||||
|
||||
def _host_is(domain: str) -> bool:
|
||||
return hostname == domain or hostname.endswith("." + domain)
|
||||
|
||||
return ResponsesRouteFlags(
|
||||
is_codex_backend=provider == "openai-codex" or (_host_is("chatgpt.com") and "/backend-api/codex" in lower),
|
||||
is_xai_responses=provider in {"xai", "xai-oauth"} or hostname == "api.x.ai",
|
||||
@@ -654,17 +639,13 @@ def estimate_native_responses_preflight_tokens(
|
||||
"""
|
||||
if getattr(agent, "api_mode", None) != "codex_responses" or not isinstance(messages, list):
|
||||
return None
|
||||
|
||||
is_codex_backend, is_xai_responses, is_github_responses = classify_responses_route(agent)
|
||||
|
||||
from agent.native_compaction import native_compaction_context_management
|
||||
|
||||
if not native_compaction_context_management(
|
||||
agent, is_codex_backend=is_codex_backend, is_xai_responses=is_xai_responses,
|
||||
is_github_responses=is_github_responses,
|
||||
):
|
||||
return None
|
||||
|
||||
try:
|
||||
items = _chat_messages_to_responses_input(
|
||||
messages, is_xai_responses=is_xai_responses, is_github_responses=is_github_responses,
|
||||
@@ -680,9 +661,7 @@ def estimate_native_responses_preflight_tokens(
|
||||
return None
|
||||
if not isinstance(items, list):
|
||||
return None
|
||||
|
||||
from agent.model_metadata import estimate_request_tokens_rough
|
||||
|
||||
return estimate_request_tokens_rough(items, system_prompt=system_prompt or "", tools=tools)
|
||||
|
||||
|
||||
@@ -781,7 +760,6 @@ def _preflight_role_message(item: Dict[str, Any], idx: int, role: str, ctx: _Pre
|
||||
content = item.get("content", "")
|
||||
if not isinstance(content, list):
|
||||
return {"role": role, "content": ctx.sanitize_text(_str_or_empty(content))}
|
||||
|
||||
# Parts are already Responses-shaped; validate and re-type text for the role.
|
||||
# Unlike history conversion, empty text / empty image urls are kept, not dropped.
|
||||
text_type = _text_type_for(role)
|
||||
@@ -825,7 +803,6 @@ def _preflight_codex_input_items(
|
||||
) -> List[Dict[str, Any]]:
|
||||
if not isinstance(raw_items, list):
|
||||
raise ValueError("Codex Responses input must be a list of input items.")
|
||||
|
||||
ctx = _PreflightCtx(
|
||||
sanitize_text=_neutralize_harmony_tokens if sanitize_harmony_tokens else (lambda text: text),
|
||||
sanitize_harmony_tokens=sanitize_harmony_tokens,
|
||||
@@ -906,19 +883,15 @@ def _preflight_codex_api_kwargs(
|
||||
) -> Dict[str, Any]:
|
||||
if not isinstance(api_kwargs, dict):
|
||||
raise ValueError("Codex Responses request must be a dict.")
|
||||
|
||||
missing = [key for key in ("model", "instructions", "input") if key not in api_kwargs]
|
||||
if missing:
|
||||
raise ValueError(f"Codex Responses request missing required field(s): {', '.join(sorted(missing))}.")
|
||||
|
||||
model = api_kwargs.get("model")
|
||||
if not _nonblank(model):
|
||||
raise ValueError("Codex Responses request 'model' must be a non-empty string.")
|
||||
|
||||
instructions = _str_or_empty(api_kwargs.get("instructions")).strip() or DEFAULT_AGENT_IDENTITY
|
||||
if sanitize_harmony_tokens:
|
||||
instructions = _neutralize_harmony_tokens(instructions)
|
||||
|
||||
normalized: Dict[str, Any] = {
|
||||
"model": model.strip(),
|
||||
"instructions": instructions,
|
||||
@@ -929,7 +902,6 @@ def _preflight_codex_api_kwargs(
|
||||
),
|
||||
"store": False,
|
||||
}
|
||||
|
||||
tools = api_kwargs.get("tools")
|
||||
if tools is not None:
|
||||
if not isinstance(tools, list):
|
||||
@@ -938,15 +910,12 @@ def _preflight_codex_api_kwargs(
|
||||
if sanitize_harmony_tokens:
|
||||
normalized_tools = _neutralize_harmony_structure(normalized_tools)
|
||||
normalized["tools"] = normalized_tools
|
||||
|
||||
if api_kwargs.get("store", False) is not False:
|
||||
raise ValueError("Codex Responses contract requires 'store' to be false.")
|
||||
|
||||
for key, accept, coerce in _PREFLIGHT_OPTIONAL_FIELDS:
|
||||
value = api_kwargs.get(key)
|
||||
if accept(value):
|
||||
normalized[key] = coerce(value) if coerce else value
|
||||
|
||||
extra_headers = api_kwargs.get("extra_headers")
|
||||
if extra_headers is not None:
|
||||
if not isinstance(extra_headers, dict):
|
||||
@@ -956,7 +925,6 @@ def _preflight_codex_api_kwargs(
|
||||
normalized_headers = {key.strip(): str(value) for key, value in extra_headers.items() if value is not None}
|
||||
if normalized_headers:
|
||||
normalized["extra_headers"] = normalized_headers
|
||||
|
||||
extra_body = api_kwargs.get("extra_body")
|
||||
if extra_body is not None:
|
||||
if not isinstance(extra_body, dict):
|
||||
@@ -965,7 +933,6 @@ def _preflight_codex_api_kwargs(
|
||||
# the SDK serializes extra_body without per-field checks.
|
||||
if extra_body:
|
||||
normalized["extra_body"] = dict(extra_body)
|
||||
|
||||
allowed_keys = set(_PREFLIGHT_ALLOWED_KEYS)
|
||||
if allow_stream:
|
||||
stream = api_kwargs.get("stream")
|
||||
@@ -976,7 +943,6 @@ def _preflight_codex_api_kwargs(
|
||||
allowed_keys.add("stream")
|
||||
elif "stream" in api_kwargs:
|
||||
raise ValueError("Codex Responses stream flag is only allowed in fallback streaming requests.")
|
||||
|
||||
# Defense-in-depth slash-enum strip for xAI (rejects ``Qwen/Qwen3.5`` style
|
||||
# enum values). Gated on the model name because native Codex accepts slashes.
|
||||
is_xai_model = str(api_kwargs.get("model") or "").lower().startswith(("grok-", "x-ai/grok-"))
|
||||
@@ -986,7 +952,6 @@ def _preflight_codex_api_kwargs(
|
||||
normalized["tools"], _ = strip_slash_enum(normalized["tools"])
|
||||
except Exception:
|
||||
pass # Best-effort — the caller-level sanitization should have handled it
|
||||
|
||||
unexpected = sorted(key for key in api_kwargs if key not in allowed_keys)
|
||||
if unexpected:
|
||||
raise ValueError(f"Codex Responses request has unsupported field(s): {', '.join(unexpected)}.")
|
||||
@@ -1025,7 +990,6 @@ def _format_responses_error(error_obj: Any, response_status: str) -> str:
|
||||
def field(name: str) -> str:
|
||||
value = _field(error_obj, name)
|
||||
return str(value).strip() if isinstance(value, str) or value else ""
|
||||
|
||||
code_str, message_str = field("code"), field("message")
|
||||
if code_str and message_str:
|
||||
return f"{code_str}: {message_str}"
|
||||
@@ -1113,7 +1077,6 @@ class _OutputScan:
|
||||
if item_status in _INCOMPLETE_STATUSES and item_type not in _SERVER_SIDE_TOOL_CALL_TYPES:
|
||||
self.has_incomplete_items = True
|
||||
self.saw_streaming_or_item_incomplete = True
|
||||
|
||||
if item_type == "message":
|
||||
self._message(item, item_status)
|
||||
elif item_type == "reasoning":
|
||||
@@ -1170,7 +1133,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
|
||||
response_incomplete_content_filter = (
|
||||
response_status == "incomplete" and str(incomplete_reason or "").strip().lower() == "content_filter"
|
||||
)
|
||||
|
||||
output = getattr(response, "output", None)
|
||||
if not isinstance(output, list) or not output:
|
||||
# Codex can deliver the whole answer via stream events and return an
|
||||
@@ -1189,20 +1151,16 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
|
||||
else:
|
||||
raise RuntimeError("Responses API returned no output items")
|
||||
response.output = output
|
||||
|
||||
if response_status in {"failed", "cancelled"}:
|
||||
raise RuntimeError(_format_responses_error(getattr(response, "error", None), response_status))
|
||||
|
||||
scan = _OutputScan(response_status)
|
||||
scan.scan(output, issuer_kind)
|
||||
tool_calls, reasoning_parts = scan.tool_calls, scan.reasoning_parts
|
||||
|
||||
final_text = "\n".join(scan.content_parts).strip()
|
||||
if not final_text and hasattr(response, "output_text") and (scan.saw_final_answer_phase or not scan.saw_commentary_phase):
|
||||
out_text = getattr(response, "output_text", "")
|
||||
if isinstance(out_text, str):
|
||||
final_text = out_text.strip()
|
||||
|
||||
# Tool-call leak recovery: gpt-5.x sometimes emits the intended
|
||||
# ``function_call`` as plain Harmony text (``to=functions.foo {json}``) with
|
||||
# no structured item. Treat as incomplete so the continuation path
|
||||
@@ -1216,7 +1174,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
|
||||
"Leaked snippet: %r", final_text[:300],
|
||||
)
|
||||
final_text = ""
|
||||
|
||||
# Reasoning-channel answer salvage (xAI grok): grok-4.x sometimes puts the
|
||||
# final answer inside the reasoning item after its ``<response>`` delimiter.
|
||||
# Without salvage the reasoning-only rule marks the turn incomplete, and since
|
||||
@@ -1235,7 +1192,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
|
||||
final_text = salvaged
|
||||
reasoning_prefix = joined_reasoning[:marker].strip()
|
||||
reasoning_parts = [reasoning_prefix] if reasoning_prefix else []
|
||||
|
||||
assistant_message = SimpleNamespace(
|
||||
content=final_text,
|
||||
tool_calls=tool_calls,
|
||||
@@ -1245,7 +1201,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
|
||||
codex_reasoning_items=scan.reasoning_items_raw or None,
|
||||
codex_message_items=scan.message_items_raw or None,
|
||||
)
|
||||
|
||||
if tool_calls:
|
||||
finish_reason = "tool_calls"
|
||||
elif response_incomplete_content_filter:
|
||||
|
||||
@@ -41,7 +41,6 @@ def _codex_request_failure_details(error: BaseException) -> tuple[int | None, st
|
||||
exception_classes: list[str] = []
|
||||
current: BaseException | None = error
|
||||
seen: set[int] = set()
|
||||
|
||||
while current is not None and id(current) not in seen and len(seen) < 8:
|
||||
seen.add(id(current))
|
||||
exception_classes.append(type(current).__name__)
|
||||
@@ -102,7 +101,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]:
|
||||
no usage still counts as one API call for session/status accounting.
|
||||
"""
|
||||
agent.session_api_calls += 1
|
||||
|
||||
usage = getattr(turn, "token_usage_last", None)
|
||||
compressor = getattr(agent, "context_compressor", None)
|
||||
if not isinstance(usage, dict) or not usage:
|
||||
@@ -118,9 +116,7 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]:
|
||||
),
|
||||
)
|
||||
return {}
|
||||
|
||||
from agent.usage_pricing import CanonicalUsage, estimate_usage_cost
|
||||
|
||||
canonical_usage = CanonicalUsage(
|
||||
input_tokens=_coerce_usage_int(usage.get("inputTokens")),
|
||||
output_tokens=_coerce_usage_int(usage.get("outputTokens")),
|
||||
@@ -139,7 +135,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]:
|
||||
"prompt_tokens": prompt_tokens, "completion_tokens": canonical_usage.output_tokens, "total_tokens": total_tokens,
|
||||
**token_counts,
|
||||
}
|
||||
|
||||
if compressor is not None:
|
||||
try:
|
||||
compressor.update_from_response(usage_dict)
|
||||
@@ -148,10 +143,8 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]:
|
||||
compressor.context_length = context_window
|
||||
except Exception:
|
||||
logger.debug("codex app-server usage update failed", exc_info=True)
|
||||
|
||||
for key, value in usage_dict.items():
|
||||
setattr(agent, f"session_{key}", getattr(agent, f"session_{key}") + value)
|
||||
|
||||
cost_result = estimate_usage_cost(
|
||||
agent.model, canonical_usage,
|
||||
provider=agent.provider, base_url=agent.base_url, api_key=getattr(agent, "api_key", ""),
|
||||
@@ -161,7 +154,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]:
|
||||
agent.session_estimated_cost_usd += cost_usd
|
||||
agent.session_cost_status, agent.session_cost_source = cost_result.status, cost_result.source
|
||||
cost_fields = {"estimated_cost_usd": cost_usd, "cost_status": cost_result.status, "cost_source": cost_result.source}
|
||||
|
||||
_queue_token_counts(
|
||||
agent, "Codex app-server token persistence failed (session=%s, tokens=%d): %s", total_tokens,
|
||||
counts=lambda: dict(
|
||||
@@ -182,7 +174,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non
|
||||
"""
|
||||
if not force and not getattr(turn, "compacted", False):
|
||||
return False
|
||||
|
||||
thread_id = getattr(turn, "thread_id", None) or ""
|
||||
turn_id = getattr(turn, "turn_id", None) or ""
|
||||
logger.info(
|
||||
@@ -195,7 +186,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non
|
||||
agent._emit_status(COMPACTION_STATUS)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
compressor = getattr(agent, "context_compressor", None)
|
||||
if compressor is not None:
|
||||
compressor.compression_count = getattr(compressor, "compression_count", 0) + 1
|
||||
@@ -212,7 +202,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non
|
||||
compressor.last_prompt_tokens = -1
|
||||
compressor.last_completion_tokens = 0
|
||||
compressor.awaiting_real_usage_after_compression = True
|
||||
|
||||
# Provider-side context was rewritten; the usage anchor's transcript snapshot no longer matches.
|
||||
agent._usage_anchor = None
|
||||
agent._turn_base_usage_anchor = None
|
||||
@@ -333,7 +322,6 @@ def _codex_item_completion_payload(item: dict) -> tuple[str, bool]:
|
||||
def _stable_call_id(item: dict, name: str) -> str:
|
||||
"""Deterministic tool_call id mirroring CodexEventProjector (live TUI card correlates with projected history)."""
|
||||
from agent.transports.codex_event_projector import _deterministic_call_id
|
||||
|
||||
item_type = item.get("type") or ""
|
||||
tool = item.get("tool") or "unknown"
|
||||
if item_type == "mcpToolCall":
|
||||
@@ -414,7 +402,6 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]:
|
||||
(_fire_tool_completed if completed else _fire_tool_started)(item)
|
||||
elif completed and item_type == "agentMessage":
|
||||
_fire_agent_message_completed(item)
|
||||
|
||||
handlers: dict[str, Callable[[dict], None]] = {
|
||||
"item/agentMessage/delta": lambda p: _fire_delta(p, "_fire_stream_delta"),
|
||||
"item/reasoning/delta": lambda p: _fire_delta(p, "_fire_reasoning_delta"),
|
||||
@@ -428,7 +415,6 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]:
|
||||
if handler is not None:
|
||||
params = note.get("params") or {}
|
||||
handler(params if isinstance(params, dict) else {})
|
||||
|
||||
return on_event
|
||||
|
||||
|
||||
@@ -461,7 +447,6 @@ def _ensure_codex_session(agent) -> None:
|
||||
return
|
||||
from agent.runtime_cwd import resolve_agent_cwd
|
||||
from agent.transports.codex_app_server_session import CodexAppServerSession, _ServerRequestRouting
|
||||
|
||||
# Approval callback: Hermes' standard prompt flow when a CLI thread installed one.
|
||||
try:
|
||||
from tools.terminal_tool import _get_approval_callback
|
||||
@@ -478,7 +463,6 @@ def _ensure_codex_session(agent) -> None:
|
||||
auto_approve_requests = is_approval_bypass_active()
|
||||
except Exception:
|
||||
logger.debug("codex app-server: approval-bypass lookup failed; keeping fail-closed default", exc_info=True)
|
||||
|
||||
agent._codex_session = CodexAppServerSession(
|
||||
cwd=getattr(agent, "session_cwd", None) or str(resolve_agent_cwd()), approval_callback=approval_callback,
|
||||
request_routing=_ServerRequestRouting(auto_approve_exec=auto_approve_requests, auto_approve_apply_patch=auto_approve_requests),
|
||||
@@ -498,10 +482,8 @@ def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) ->
|
||||
if not turn.projected_messages:
|
||||
return
|
||||
from agent.message_metadata import append_message
|
||||
|
||||
for projected_message in turn.projected_messages:
|
||||
append_message(messages, projected_message)
|
||||
|
||||
if getattr(agent, "_session_db", None) is None:
|
||||
return
|
||||
try:
|
||||
@@ -528,14 +510,12 @@ def _finish_codex_turn(
|
||||
agent._iters_since_skill = getattr(agent, "_iters_since_skill", 0) + turn.tool_iterations
|
||||
_record_codex_app_server_compaction(agent, turn)
|
||||
usage_result = _record_codex_app_server_usage(agent, turn)
|
||||
|
||||
# Skill nudge check AFTER iters were incremented (same as chat_completions).
|
||||
should_review_skills = (
|
||||
0 < agent._skill_nudge_interval <= agent._iters_since_skill and "skill_manage" in agent.valid_tool_names
|
||||
)
|
||||
if should_review_skills:
|
||||
agent._iters_since_skill = 0
|
||||
|
||||
# External memory sync skipped on interrupt/error (no partial transcripts).
|
||||
if not turn.interrupted and turn.error is None:
|
||||
try:
|
||||
@@ -545,7 +525,6 @@ def _finish_codex_turn(
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("external memory sync raised", exc_info=True)
|
||||
|
||||
# Background review fork: only when a trigger tripped AND a real final response exists.
|
||||
if turn.final_text and not turn.interrupted and (should_review_memory or should_review_skills):
|
||||
try:
|
||||
@@ -554,7 +533,6 @@ def _finish_codex_turn(
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("background review spawn raised", exc_info=True)
|
||||
|
||||
return usage_result
|
||||
|
||||
|
||||
@@ -581,9 +559,7 @@ def run_codex_app_server_turn(
|
||||
"codex_app_server owns the authoritative thread and compacts it "
|
||||
"without a truthful pre-compaction transcript boundary"
|
||||
)
|
||||
|
||||
_ensure_codex_session(agent)
|
||||
|
||||
try:
|
||||
turn = agent._codex_session.run_turn(user_input=user_message)
|
||||
except Exception as exc:
|
||||
@@ -593,20 +569,16 @@ def run_codex_app_server_turn(
|
||||
_consume_user_interrupt(agent), messages, api_calls=0, completed=False, error=str(exc),
|
||||
final_response=f"Codex app-server turn failed: {exc}. Fall back to default runtime with `/codex-runtime auto`.",
|
||||
)
|
||||
|
||||
interrupt = _consume_user_interrupt(agent, turn.interrupted)
|
||||
|
||||
# Wedged client (deadline blown, watchdog tripped, OAuth refresh died,
|
||||
# subprocess exited): retire the session so the next turn respawns codex.
|
||||
if getattr(turn, "should_retire", False):
|
||||
logger.warning("codex app-server session retired (turn error: %s)", turn.error)
|
||||
_close_codex_session(agent)
|
||||
|
||||
_persist_projected_messages(agent, turn, messages)
|
||||
usage_result = _finish_codex_turn(
|
||||
agent, turn, messages, original_user_message=original_user_message, should_review_memory=should_review_memory,
|
||||
)
|
||||
|
||||
return _turn_result(
|
||||
interrupt, messages, api_calls=1, completed=not turn.interrupted and turn.error is None, error=turn.error,
|
||||
final_response=turn.final_text,
|
||||
@@ -659,13 +631,11 @@ def _raise_stream_error(event: Any) -> None:
|
||||
``run_agent`` is imported lazily to keep this module importable standalone.
|
||||
"""
|
||||
from run_agent import _StreamErrorEvent
|
||||
|
||||
nested = _event_field(event, "error")
|
||||
|
||||
def _error_field(name: str) -> Any:
|
||||
value = _event_field(event, name)
|
||||
return _event_field(nested, name) if value is None and nested is not None else value
|
||||
|
||||
raw_message = _error_field("message")
|
||||
message = (str(raw_message) if raw_message is not None else "stream emitted error event").strip() or "stream emitted error event"
|
||||
raise _StreamErrorEvent(message, code=_error_field("code"), param=_error_field("param"))
|
||||
@@ -879,7 +849,6 @@ class _CodexResponseAssembler:
|
||||
# executable; malformed non-empty JSON passes through untouched.
|
||||
arguments=(pending["arguments"] or "").strip() or "{}",
|
||||
)))
|
||||
|
||||
# output_index is optional and a partial ordering over mixed indexed/unindexed
|
||||
# entries is ill-defined: protocol order only when every entry has an index, else wire order.
|
||||
if all(entry[0] is not None for entry in indexed):
|
||||
@@ -898,17 +867,14 @@ class _CodexResponseAssembler:
|
||||
if not output and self.text_deltas and not self.has_tool_calls:
|
||||
content = [SimpleNamespace(type="output_text", text="".join(self.text_deltas))]
|
||||
output = [SimpleNamespace(type="message", role="assistant", status="completed", content=content)]
|
||||
|
||||
# Done items stay authoritative; settlement only fills the gap left by
|
||||
# backends that omit per-item done events on a successful completion.
|
||||
if self.pending_function_calls and self.saw_response_completed:
|
||||
output = self._settled_output()
|
||||
|
||||
# No terminal frame AND no usable content = truncated / rejected stream,
|
||||
# distinct from "completed with empty body" (what the SDK helper raised as RuntimeError).
|
||||
if not self.saw_terminal and not output:
|
||||
raise RuntimeError("Codex Responses stream did not emit a terminal response")
|
||||
|
||||
return SimpleNamespace(
|
||||
output=output, output_text="".join(self.text_deltas), usage=self.terminal_usage, status=self.terminal_status,
|
||||
id=self.terminal_response_id, model=self.model, incomplete_details=self.terminal_incomplete_details,
|
||||
@@ -1021,7 +987,6 @@ def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict:
|
||||
"""
|
||||
if os.environ.get("HERMES_CODEX_SDK_TRANSFORM", "").strip().lower() in {"1", "true", "yes", "on"}:
|
||||
return stream_kwargs
|
||||
|
||||
moved = {
|
||||
field: stream_kwargs[field]
|
||||
for field in _SDK_TRANSFORM_BYPASS_FIELDS
|
||||
@@ -1029,7 +994,6 @@ def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict:
|
||||
}
|
||||
if not moved:
|
||||
return stream_kwargs
|
||||
|
||||
bypassed = {key: value for key, value in stream_kwargs.items() if key not in moved}
|
||||
extra_body = bypassed.get("extra_body")
|
||||
merged = dict(extra_body) if isinstance(extra_body, dict) else {}
|
||||
@@ -1049,9 +1013,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta
|
||||
"""
|
||||
import httpx as _httpx
|
||||
from openai import APIConnectionError as _APIConnectionError
|
||||
|
||||
from agent import relay_llm
|
||||
|
||||
transport_errors = (_httpx.RemoteProtocolError, _httpx.ReadTimeout, _httpx.ConnectError, ConnectionError)
|
||||
active_client = client or agent._ensure_primary_openai_client(reason="codex_stream_direct")
|
||||
max_stream_retries = 1
|
||||
@@ -1144,7 +1106,6 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta
|
||||
# primary client, which is never reuse-cached and must not be force-shut.
|
||||
if client is not None:
|
||||
agent._abort_request_openai_client(active_client, reason="codex_stream_close_failed")
|
||||
|
||||
on_commentary_message = (
|
||||
_fenced(lambda text: agent._fire_streamed_codex_commentary(text))
|
||||
if getattr(agent, "interim_assistant_callback", None) is not None and getattr(agent, "show_commentary", True)
|
||||
@@ -1155,11 +1116,9 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta
|
||||
else "fallback" if int(getattr(agent, "_fallback_index", 0) or 0) > 0
|
||||
else "primary"
|
||||
)
|
||||
|
||||
for attempt in range(max_stream_retries + 1):
|
||||
if agent._interrupt_requested:
|
||||
raise InterruptedError("Agent interrupted before Codex stream retry")
|
||||
|
||||
intercepted_events: list = []
|
||||
writer_token["value"] = None
|
||||
event_stream = None
|
||||
@@ -1209,10 +1168,8 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta
|
||||
except _APIConnectionError as exc:
|
||||
_log_failure(exc)
|
||||
raise
|
||||
|
||||
if not agent._interrupt_requested:
|
||||
_drain_for_finalizer(event_stream)
|
||||
|
||||
if final.status in {"incomplete", "failed"}:
|
||||
logger.warning(
|
||||
"Codex Responses stream terminal status=%s "
|
||||
|
||||
Reference in New Issue
Block a user