fix(bedrock): replay thinking as reasoningText; resend once without redacted blocks on encrypted-content rejection (#115865)

This commit is contained in:
Yagna Vudathu
2026-09-19 03:44:11 -04:00
committed by Teknium
parent b873a3a1d9
commit 8437310aed
2 changed files with 212 additions and 6 deletions

View File

@@ -536,6 +536,62 @@ def recover_from_cache_point_rejection(exc: BaseException, kwargs: Dict[str, Any
return retry_kwargs
# --- Encrypted-content redacted-reasoning suppression ---
# Redacted thinking blobs are sealed to the issuing model/flow; replaying them after a model switch or
# across regions fails with an encrypted-content ValidationException. Drop the redacted blocks and
# resend once (mirrors the cachePoint self-heal above: same-object return means no retry can help).
_REDACTED_REASONING_REJECTION_PATTERN = re.compile(
r"ValidationException.*(?:redacted|encrypt)", re.IGNORECASE | re.DOTALL,
)
def _without_redacted_reasoning(blocks):
"""``blocks`` minus reasoningContent entries carrying redactedContent, or None if not a list /
nothing removed. Pure reasoningText blocks are kept."""
if not isinstance(blocks, list):
return None
cleaned = [b for b in blocks if not (
isinstance(b, dict) and isinstance(b.get("reasoningContent"), dict)
and "redactedContent" in b["reasoningContent"]
)]
return None if len(cleaned) == len(blocks) else cleaned
def strip_redacted_reasoning(kwargs):
"""Copy of Converse kwargs with redacted reasoning blocks removed; the SAME object
back when nothing was stripped (callers use identity to decide a retry cannot help).
Turns left with empty content are dropped (they carried no replayable signal)."""
messages = kwargs.get("messages")
if not isinstance(messages, list):
return kwargs
cleaned_contents = [
_without_redacted_reasoning(msg.get("content") if isinstance(msg, dict) else None)
for msg in messages
]
if all(content is None for content in cleaned_contents):
return kwargs
return {**kwargs, "messages": [
{**msg, "content": content} for msg, content in zip(messages, cleaned_contents)
if content is None or len(content) > 0 or not isinstance(msg, dict)
]}
def recover_from_redacted_reasoning_rejection(exc, kwargs):
"""Return retry kwargs with redacted reasoning blocks stripped, or None when the error
was not an encrypted-content rejection / nothing redacted remained (caller re-raises)."""
if not _REDACTED_REASONING_REJECTION_PATTERN.search(str(exc)):
return None
retry_kwargs = strip_redacted_reasoning(kwargs)
if retry_kwargs is kwargs:
return None
logger.warning(
"bedrock: %s rejected replayed redacted reasoning (encrypted content is sealed to the issuing "
"model/flow) - stripping redacted blocks and resending once.",
str(kwargs.get("modelId", "")) or "model",
)
return retry_kwargs
# One optional regional/global inference-profile prefix, then the Claude model family.
_ANTHROPIC_BEDROCK_MODEL_RE = re.compile(
r"^(?:(?:global|us|eu|apac|ap|au|jp|ca|sa|me|af)\.)?anthropic\.claude", re.IGNORECASE,
@@ -632,15 +688,15 @@ def _replay_ordered_blocks(ordered_blocks: List) -> List[Dict]:
reasoning = block["reasoningContent"]
if not isinstance(reasoning, dict):
continue
replay = {"text": reasoning["text"]} if isinstance(reasoning.get("text"), str) else {}
# ReasoningContentBlock is a tagged union: reasoningText and redactedContent must go
# out as separate blocks (#115865). Undecodable redacted entries are skipped alone.
if isinstance(reasoning.get("text"), str):
content_blocks.append({"reasoningContent": {"reasoningText": {"text": reasoning["text"]}}})
encoded = reasoning.get("redactedContentBase64")
if isinstance(encoded, str) and encoded:
redacted = _decode_redacted(encoded)
if redacted is None:
continue
replay["redactedContent"] = redacted
if replay:
content_blocks.append({"reasoningContent": replay})
if redacted is not None:
content_blocks.append({"reasoningContent": {"redactedContent": redacted}})
elif "toolUse" in block and isinstance(block["toolUse"], dict):
tu = block["toolUse"]
content_blocks.append(_tool_use_block(tu.get("toolUseId", ""), tu.get("name", ""), tu.get("input", {})))
@@ -955,6 +1011,9 @@ def call_converse(
retry_kwargs = recover_from_cache_point_rejection(exc, kwargs)
if retry_kwargs is not None:
return normalize_converse_response(client.converse(**retry_kwargs))
redacted_retry_kwargs = recover_from_redacted_reasoning_rejection(exc, kwargs)
if redacted_retry_kwargs is not None:
return normalize_converse_response(client.converse(**redacted_retry_kwargs))
if is_stale_connection_error(exc):
logger.warning(
"bedrock: stale-connection error on converse(region=%s, model=%s): "
@@ -1238,6 +1297,11 @@ def call_converse_stream(
return normalize_converse_stream_events(
client.converse_stream(**retry_kwargs)
)
redacted_retry_kwargs = recover_from_redacted_reasoning_rejection(exc, kwargs)
if redacted_retry_kwargs is not None:
return normalize_converse_stream_events(
client.converse_stream(**redacted_retry_kwargs)
)
if is_streaming_access_denied_error(exc):
# IAM allows bedrock:InvokeModel but not
# InvokeModelWithResponseStream — permanent for this session.

View File

@@ -1537,3 +1537,145 @@ class TestBearerTokenRoutesToConverse:
runtime = self._resolve(monkeypatch, bearer=False)
assert runtime["api_mode"] == "anthropic_messages"
assert runtime.get("bedrock_anthropic") is True
# ---------------------------------------------------------------------------
# Reasoning-text replay + redacted-reasoning resend-once (#115865)
# ---------------------------------------------------------------------------
REDACTED_REJECTION = (
"An error occurred (ValidationException) when calling the Converse "
"operation: The provided redacted thinking block contains encrypted "
"content that is not valid for this model, please reformat your input "
"and try again."
)
class TestReasoningTextReplay:
"""Converse ReasoningContentBlock members are reasoningText/redactedContent
only — replaying captured thinking as {"text": ...} dies client-side with
ParamValidationError (#115865)."""
def test_ordered_replay_emits_reasoningText(self):
from agent.bedrock_adapter import _replay_ordered_blocks
blocks = _replay_ordered_blocks([{
"reasoningContent": {
"text": "let me think",
"redactedContentBase64": "cjE=",
},
}])
assert blocks == [
{"reasoningContent": {"reasoningText": {"text": "let me think"}}},
{"reasoningContent": {"redactedContent": b"r1"}},
]
def test_replayed_blocks_pass_botocore_converse_validation(self):
import botocore.session
from botocore.validate import validate_parameters
from agent.bedrock_adapter import _replay_ordered_blocks
blocks = _replay_ordered_blocks([{
"reasoningContent": {
"text": "let me think",
"redactedContentBase64": "cjE=",
},
}])
model = botocore.session.get_session().get_service_model("bedrock-runtime")
shape = model.operation_model("Converse").input_shape
validate_parameters(
{"modelId": "test-model",
"messages": [{"role": "assistant", "content": blocks}]},
shape,
) # raises ParamValidationError while replay emits {"text": ...}
def _redacted_history():
return [
{"role": "user", "content": "go"},
{
"role": "assistant",
"content": None,
"tool_calls": [{
"id": "call_1", "type": "function",
"function": {"name": "read_file", "arguments": "{}"},
}],
"reasoning_details": [{
"type": "redacted_thinking", "data": "cjE=",
}],
"bedrock_content_blocks": [
{"reasoningContent": {
"text": "let me think", "redactedContentBase64": "cjE=",
}},
{"toolUse": {
"toolUseId": "call_1", "name": "read_file", "input": {},
}},
],
},
{"role": "tool", "tool_call_id": "call_1", "content": "ok"},
]
def _ok_response():
return {
"output": {"message": {"role": "assistant",
"content": [{"text": "done"}]}},
"stopReason": "end_turn",
"usage": {"inputTokens": 1, "outputTokens": 1},
}
def _assert_second_attempt_stripped(call_args_list):
assert len(call_args_list) == 2
resent = call_args_list[1].kwargs["messages"]
assistant = next(m for m in resent if m["role"] == "assistant")
for block in assistant["content"]:
assert "redactedContent" not in block.get("reasoningContent", {})
class TestRedactedReasoningResendOnce:
"""Redacted thinking blobs are sealed to the issuing model/flow — replaying
them elsewhere fails with an encrypted-content ValidationException. Drop
the redacted blocks and resend once (#115865)."""
def test_call_converse_strips_redacted_blocks_and_resends_once(self):
from agent.bedrock_adapter import call_converse
client = MagicMock()
client.converse.side_effect = [Exception(REDACTED_REJECTION), _ok_response()]
with patch("agent.bedrock_adapter._get_bedrock_runtime_client",
return_value=client):
response = call_converse(
region="us-east-1", model="test-model", messages=_redacted_history(),
)
assert response.choices[0].message.content == "done"
_assert_second_attempt_stripped(client.converse.call_args_list)
def test_call_converse_reraises_when_nothing_redacted_to_strip(self):
from agent.bedrock_adapter import call_converse
client = MagicMock()
client.converse.side_effect = Exception(REDACTED_REJECTION)
with patch("agent.bedrock_adapter._get_bedrock_runtime_client",
return_value=client):
with pytest.raises(Exception, match="ValidationException"):
call_converse(
region="us-east-1", model="test-model",
messages=[{"role": "user", "content": "hi"}],
)
assert client.converse.call_count == 1
def test_call_converse_stream_strips_redacted_blocks_and_resends_once(self):
from agent.bedrock_adapter import call_converse_stream
client = MagicMock()
client.converse_stream.side_effect = [Exception(REDACTED_REJECTION), {"stream": [
{"messageStart": {"role": "assistant"}},
{"contentBlockStart": {"contentBlockIndex": 0, "start": {}}},
{"contentBlockDelta": {"contentBlockIndex": 0, "delta": {"text": "done"}}},
{"contentBlockStop": {"contentBlockIndex": 0}},
{"messageStop": {"stopReason": "end_turn"}},
{"metadata": {"usage": {"inputTokens": 1, "outputTokens": 1}}},
]}]
with patch("agent.bedrock_adapter._get_bedrock_runtime_client",
return_value=client):
response = call_converse_stream(
region="us-east-1", model="test-model", messages=_redacted_history(),
)
assert response.choices[0].message.content == "done"
_assert_second_attempt_stripped(client.converse_stream.call_args_list)