Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 12 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,22 +40,30 @@ Provides streaming orchestration that delivers LLM responses token-by-token for

When a provider exposes the optional synchronous `request_budget` capability,
the loop checks the fully assembled request before dispatch. An oversized
request gets exactly one smaller, retention-aware context view; required
request gets up to two smaller, retention-aware context views; required
reminders and request-only injections are replayed without running hooks or
draining pending state again. If that view still cannot fit, the loop raises
locally and makes no SDK request. Providers without the capability retain the
existing request and dispatch behavior.

When `context.request_retention` advertises its optional `hard_fit` keyword,
that one provider-forced rebuild forwards `hard_fit=True`, allowing the context
to target the provider's requested budget directly. Older retention
each provider-forced rebuild forwards `hard_fit=True`, allowing the context to
target the provider's requested budget directly. The second rebuild runs only
when the provider requests a strictly smaller budget. Older retention
capabilities, uninspectable dynamic callables, and the generic context fallback
keep their existing `provider`/`retain_contents`/`token_budget` assembly; the
second preflight and provider's final payload guard remain the safety boundary.
final preflight and provider's final payload guard remain the safety boundary.
The existing `orchestrator:provider_budget` event exposes each preflight's
attempt, result, estimate, allowance, and requested context budget to mounted
observability consumers.

When a provider reports its effective output cap, an oversized compacted view
is also preflighted with progressively smaller response reserves: 50%, 40%,
30%, 20%, and 10% of the original cap, then 1,000 tokens. These are local
preflights: they do not resend a provider request or omit input. At caps below
10,000 tokens, a view-only system reminder asks the model to tell the user that
the session is degraded and to recommend a new session for substantial work.

## Configuration

```toml
Expand Down
228 changes: 200 additions & 28 deletions amplifier_module_loop_streaming/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3462,7 +3462,10 @@ async def request_messages(
)

def build_chat_request(
message_dicts: list[dict[str, Any]], *, tool_choice: str | None = None
message_dicts: list[dict[str, Any]],
*,
tool_choice: str | None = None,
max_output_tokens: int | None = None,
) -> ChatRequest:
tools_list = [_build_tool_spec(tool) for tool in tools.values()] if tools else None
kwargs: dict[str, Any] = {
Expand All @@ -3472,6 +3475,8 @@ def build_chat_request(
}
if tool_choice is not None:
kwargs["tool_choice"] = tool_choice
if max_output_tokens is not None:
kwargs["max_output_tokens"] = max_output_tokens
return ChatRequest(
**kwargs
)
Expand All @@ -3481,11 +3486,11 @@ async def check_request_budget(
base_messages: list[dict[str, Any]],
*,
attempt: int,
) -> int | None:
"""Return one requested smaller context budget, or ``None`` when it fits."""
) -> tuple[int | None, int | None]:
"""Return a smaller context budget and effective output cap when available."""
request_budget = getattr(provider, "request_budget", None)
if not budget_capable or not callable(request_budget):
return None
return None, None
context_estimate = sum(len(str(message)) // 4 for message in base_messages)
decision = request_budget(request, context_estimate=context_estimate)
required = (
Expand All @@ -3506,6 +3511,15 @@ async def check_request_budget(
estimated = decision["estimated_input_tokens"]
allowance = decision["input_limit_tokens"]
target = decision["context_token_budget"]
output_cap = decision.get("max_output_tokens")
if output_cap is not None and (
isinstance(output_cap, bool)
or not isinstance(output_cap, int)
or output_cap <= 0
):
raise ContextLengthError(
"Provider request_budget returned an invalid max_output_tokens"
)
fits = estimated <= allowance
await hooks.emit(
"orchestrator:provider_budget",
Expand All @@ -3519,12 +3533,70 @@ async def check_request_budget(
},
)
if fits:
return None
if target <= 0:
raise ContextLengthError(
"Provider request exceeds its input budget and cannot retain a smaller context"
return None, output_cap
return target, output_cap

def output_cap_candidates(original: int | None) -> list[int]:
"""Return the bounded lossless output-reserve ladder."""
if original is None or original <= 1_000:
return []
candidates: list[int] = []
for fraction in (0.50, 0.40, 0.30, 0.20, 0.10):
cap = max(1_000, int(original * fraction))
if cap < original and (not candidates or cap < candidates[-1]):
candidates.append(cap)
if candidates[-1:] != [1_000]:
candidates.append(1_000)
return candidates

def degraded_output_warning(cap: int) -> list[dict[str, Any]]:
"""Add a view-only warning only at the severe output-cap tier."""
if cap >= 10_000:
return []
return [
{
"role": "system",
"content": (
"<system-reminder source=\"orchestrator-context-degraded\">\n"
"This request is running with a severely reduced response budget "
"because the conversation is near the provider context limit. "
"Give the user a concise answer, explain that the session is in "
"a degraded state, and recommend starting a new session for "
"substantial further work.\n"
"</system-reminder>"
),
"metadata": {"ephemeral": True},
}
]

async def try_reduced_output(
message_dicts: list[dict[str, Any]],
base_messages: list[dict[str, Any]],
*,
original_output_cap: int | None,
tool_choice: str | None = None,
) -> ChatRequest | None:
"""Preflight unchanged input at bounded lower output reserves."""
for cap in output_cap_candidates(original_output_cap):
candidate_request = build_chat_request(
message_dicts + degraded_output_warning(cap),
tool_choice=tool_choice,
max_output_tokens=cap,
)
return target
smaller_budget, _ = await check_request_budget(
candidate_request, base_messages, attempt=1
)
if smaller_budget is None:
if cap < 10_000:
await hooks.emit(
"orchestrator:context_degradation",
{
"mode": "reduced_output",
"max_output_tokens": cap,
},
)
return candidate_request
return None

# --- Turn-start reminder assembly (reminder-redesign-spec.md,
# W1.2, Option D). Hoists iteration 1's provider:request emit to
Expand Down Expand Up @@ -4038,10 +4110,14 @@ async def exit_for_cancellation() -> None:
f"[ORCHESTRATOR] Tool names: {[t.name for t in tools.values()]}"
)

smaller_context_budget = await check_request_budget(
smaller_context_budget, original_output_cap = await check_request_budget(
chat_request, base_message_dicts, attempt=0
)
if smaller_context_budget is not None:
if smaller_context_budget <= 0:
raise ContextLengthError(
"Provider request exceeds its input budget and cannot retain a smaller context"
)
rebuilt_base_messages = list(
await request_messages(
retained_contents,
Expand All @@ -4060,15 +4136,57 @@ async def exit_for_cancellation() -> None:
else rebuilt_base_messages
)
rebuilt_request = build_chat_request(rebuilt_messages)
if (
await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=1
)
) is not None:
raise ContextLengthError(
"Provider request remains over budget after one context rebuild"
next_context_budget, _ = await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=1
)
if next_context_budget is not None:
reduced_output_request = await try_reduced_output(
rebuilt_messages,
rebuilt_base_messages,
original_output_cap=original_output_cap,
)
chat_request = rebuilt_request
if reduced_output_request is not None:
chat_request = reduced_output_request
else:
if next_context_budget <= 0:
raise ContextLengthError(
"Provider request exceeds its input budget and cannot retain "
"a smaller context"
)
if next_context_budget >= smaller_context_budget:
raise ContextLengthError(
"Provider request remains over budget and a second "
"context budget would not reduce it"
)
rebuilt_base_messages = list(
await request_messages(
retained_contents,
token_budget=next_context_budget,
hard_fit=True,
)
)
rebuilt_messages = (
_replay_request_overlays(
rebuilt_base_messages,
turn_start_view_block=replay_turn_start_block,
request_injection=replay_request_injection,
pending_injections=replay_pending_injections,
)
if self._ephemeral_injection_mode == "tail"
else rebuilt_base_messages
)
rebuilt_request = build_chat_request(rebuilt_messages)
if (
await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=2
)
)[0] is not None:
raise ContextLengthError(
"Provider request remains over budget after two context rebuilds"
)
chat_request = rebuilt_request
else:
chat_request = rebuilt_request

# Apply rate limit delay before provider call
await self._apply_rate_limit_delay(hooks, iteration)
Expand Down Expand Up @@ -4723,10 +4841,15 @@ async def exit_for_cancellation() -> None:
max_iter_chat_request = build_chat_request(
message_dicts, tool_choice="none"
)
smaller_context_budget = await check_request_budget(
smaller_context_budget, original_output_cap = await check_request_budget(
max_iter_chat_request, base_message_dicts, attempt=0
)
if smaller_context_budget is not None:
if smaller_context_budget <= 0:
raise ContextLengthError(
"Provider request exceeds its input budget and cannot retain "
"a smaller context"
)
rebuilt_base_messages = list(
await request_messages(
final_retained_contents,
Expand All @@ -4748,16 +4871,62 @@ async def exit_for_cancellation() -> None:
rebuilt_request = build_chat_request(
rebuilt_messages, tool_choice="none"
)
if (
await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=1
)
) is not None:
raise ContextLengthError(
"Provider finalization request remains over budget "
"after one context rebuild"
next_context_budget, _ = await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=1
)
if next_context_budget is not None:
reduced_output_request = await try_reduced_output(
rebuilt_messages,
rebuilt_base_messages,
original_output_cap=original_output_cap,
tool_choice="none",
)
max_iter_chat_request = rebuilt_request
if reduced_output_request is not None:
max_iter_chat_request = reduced_output_request
else:
if next_context_budget <= 0:
raise ContextLengthError(
"Provider request exceeds its input budget and cannot retain "
"a smaller context"
)
if next_context_budget >= smaller_context_budget:
raise ContextLengthError(
"Provider finalization request remains over budget and a "
"second context budget would not reduce it"
)
rebuilt_base_messages = list(
await request_messages(
final_retained_contents,
token_budget=next_context_budget,
hard_fit=True,
)
)
rebuilt_messages = (
_replay_request_overlays(
rebuilt_base_messages,
turn_start_view_block=None,
request_injection=final_replay_request_injection,
pending_injections=final_replay_pending_injections,
)
if self._ephemeral_injection_mode == "tail"
else list(rebuilt_base_messages)
)
rebuilt_messages.append(finalization_overlay)
rebuilt_request = build_chat_request(
rebuilt_messages, tool_choice="none"
)
if (
await check_request_budget(
rebuilt_request, rebuilt_base_messages, attempt=2
)
)[0] is not None:
raise ContextLengthError(
"Provider finalization request remains over budget "
"after two context rebuilds"
)
max_iter_chat_request = rebuilt_request
else:
max_iter_chat_request = rebuilt_request

kwargs = {}
if self.extended_thinking:
Expand Down Expand Up @@ -4840,6 +5009,9 @@ async def exit_for_cancellation() -> None:
)
raise
except ContextLengthError:
await close_finalization_tool_turn(
"The final response could not be generated because the context is too long."
)
raise
except LLMError as e:
await hooks.emit(
Expand Down
Loading
Loading