Skip to content

Commit ef24fd0

Browse files
authored
Merge pull request #8 from ApodexAI/refactor/llm-runtime
refactor(runtime): extract shared LLM call layer
2 parents 8cf57e3 + 50fa295 commit ef24fd0

28 files changed

Lines changed: 6518 additions & 18 deletions

‎README.md‎

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ Version `0.1.x` contains the converged foundation layer:
2121
- loop configuration, lifecycle contexts, observer protocol, intervention
2222
merging, and observer dispatch helpers.
2323
- streamed tool-call recovery checks for missing required arguments.
24+
- LLM binding, response normalization, streaming assembly/watchdogs, retry
25+
classification, runaway recovery, and physical-call orchestration.
2426

2527
The initial extraction is based on the already-merged integration branches:
2628

@@ -30,10 +32,10 @@ The initial extraction is based on the already-merged integration branches:
3032
Those revisions are provenance, not runtime dependencies. AgentCore tests and
3133
builds without either product checkout.
3234

33-
The LLM call runtime and agent loop remain in the products until tool
34-
execution, model profiles, retry classification, execution-context storage,
35-
and runtime hooks have product-neutral boundaries. Moving those files before
36-
that boundary exists would only hide product coupling inside this package.
35+
The agent loop remains in the products until tool execution, model profiles,
36+
tool-call parsing, and execution-context storage have converged boundaries.
37+
The LLM runtime accepts the three remaining product decisions through explicit
38+
hooks: wall-deadline lookup, provider-chain state, and sticky-session policy.
3739

3840
## Repository boundary
3941

@@ -109,11 +111,10 @@ edit both products' core copies, that is evidence it belongs here.
109111

110112
1. **Foundation** (this version): messages, token estimation, compaction,
111113
context budget, trimming.
112-
2. **Runtime contracts** (in progress): LLM protocols and loop types are now
113-
shared; errors, execution-context storage, retry classification, and
114-
explicit product hooks remain.
115-
3. **LLM runtime:** binding, calls, streaming, response normalization, runaway
116-
recovery, and the public `llm_client` facade.
114+
2. **Runtime contracts** (complete): LLM protocols, loop types, errors, retry
115+
classification, and explicit product hooks are shared.
116+
3. **LLM runtime** (complete): binding, calls, streaming, response
117+
normalization, runaway recovery, and the public `llm_client` facade.
117118
4. **Agent loop:** model/tool parsing, tool execution, and `agent_loop`.
118119
5. Remove product compatibility facades after downstream imports have moved to
119120
`agent_core`.

‎agent_core/__init__.py‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
"""Product-neutral building blocks for Apodex agent runtimes."""
22

3+
from agent_core.errors import AgentCoreError
34
from agent_core.llm import LLMClient, LLMResponse, StreamDelta
45
from agent_core.messages import (
56
Message,
@@ -11,6 +12,7 @@
1112
)
1213

1314
__all__ = [
15+
"AgentCoreError",
1416
"LLMClient",
1517
"LLMResponse",
1618
"Message",

‎agent_core/errors.py‎

Lines changed: 164 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,164 @@
1+
"""Exception hierarchy for AgentCore."""
2+
3+
from __future__ import annotations
4+
5+
from typing import Any
6+
7+
8+
class AgentCoreError(Exception):
9+
"""Base exception for all AgentCore errors."""
10+
11+
12+
# ── Kernel errors ───────────────────────────────────────────────────────────
13+
14+
15+
class KernelError(AgentCoreError):
16+
"""Errors originating from the OS kernel layer."""
17+
18+
19+
class TaskNotFoundError(KernelError):
20+
def __init__(self, task_id: str) -> None:
21+
super().__init__(f"Task not found: {task_id}")
22+
self.task_id = task_id
23+
24+
25+
class InvalidStateTransition(KernelError):
26+
def __init__(self, task_id: str, current: str, target: str) -> None:
27+
super().__init__(f"Invalid transition for {task_id}: {current} → {target}")
28+
29+
30+
class ServiceNotRegistered(KernelError):
31+
def __init__(self, service_type: type) -> None:
32+
super().__init__(f"Service not registered: {service_type.__name__}")
33+
34+
35+
class PermissionDenied(KernelError):
36+
def __init__(self, role: str, tool: str) -> None:
37+
super().__init__(f"Role '{role}' has no permission for tool '{tool}'")
38+
39+
# LLM request errors
40+
41+
class LLMError(AgentCoreError):
42+
"""Errors from the LLM/provider layer."""
43+
44+
45+
class LLMReasoningRunaway(LLMError):
46+
"""A live stream spent its semantic budget on reasoning-only output.
47+
48+
Unlike :class:`LLMStreamStalled`, the provider is healthy and actively
49+
emitting chunks. The failure is semantic: no non-whitespace visible text
50+
or tool-call delta appeared before the configured time/token guard fired.
51+
52+
``partial_response`` is intentionally carried separately from provider
53+
usage. Early stream cancellation often happens before the terminal usage
54+
chunk arrives, so its estimated reasoning tokens must never be presented
55+
as authoritative billing data.
56+
"""
57+
58+
def __init__(
59+
self,
60+
*,
61+
elapsed_s: float,
62+
estimated_tokens: int,
63+
trigger: str,
64+
partial_response: Any,
65+
) -> None:
66+
self.elapsed_s = float(elapsed_s)
67+
self.estimated_tokens = int(estimated_tokens)
68+
self.trigger = trigger
69+
self.partial_response = partial_response
70+
super().__init__(
71+
"reasoning-only stream exceeded "
72+
f"{trigger} guard (elapsed={self.elapsed_s:.1f}s, "
73+
f"estimated_tokens={self.estimated_tokens})",
74+
)
75+
76+
77+
class LLMStreamStalled(LLMError, TimeoutError):
78+
"""A streaming LLM call went silent mid-flight.
79+
80+
Subclasses ``asyncio.TimeoutError`` so every existing transient-
81+
timeout handler (retry/backoff in ``call_llm``, chain wrappers,
82+
classification) treats it identically without changes; carried
83+
fields make the distinct failure mode visible in logs and traces.
84+
"""
85+
86+
def __init__(
87+
self, stall_s: float, chunks_seen: int, elapsed_s: float,
88+
) -> None:
89+
self.stall_s = stall_s
90+
self.chunks_seen = chunks_seen
91+
self.elapsed_s = elapsed_s
92+
super().__init__(
93+
f"stream stalled: no chunks for {stall_s:.0f}s "
94+
f"(chunks_seen={chunks_seen}, elapsed={elapsed_s:.0f}s)",
95+
)
96+
97+
98+
class LLMDeadlineExceeded(LLMError, TimeoutError):
99+
"""An LLM attempt was stopped by an enclosing runtime deadline.
100+
101+
``reason`` is deliberately carried on the underlying exception as well as
102+
on :class:`LLMCallExhausted`. Some callers unwrap ``last_exc`` before
103+
handing it to a provider-chain policy; a dedicated type prevents that
104+
policy from mistaking an exhausted run budget for an ordinary transient
105+
provider timeout.
106+
"""
107+
108+
def __init__(self, reason: str, detail: str) -> None:
109+
self.reason = reason
110+
super().__init__(f"{reason}: {detail}")
111+
112+
113+
class LLMCallExhausted(LLMError, RuntimeError):
114+
"""Raised by ``call_llm`` when retries are exhausted or the error is
115+
structurally unrecoverable (4xx without proxy-wrap, or a chain-aware
116+
fallback signal like ``model_not_found``).
117+
118+
Wraps the last exception encountered so the caller (typically the
119+
product's agent loop) can surface it to a provider-chain wrapper for
120+
L1→L2→L3 rotation. Carries ``last_exc`` separately because
121+
``raise from`` is too opaque for chain-aware classification — a chain
122+
wrapper calls ``classify_error(last_exc)`` directly.
123+
124+
``last_exc`` must always agree with ``reason``: it is the exception that
125+
*caused this raise*, not merely the most recent failure seen. A deadline
126+
refusal therefore carries :class:`LLMDeadlineExceeded` even when earlier
127+
attempts failed for unrelated reasons. The wrapper's ``reason`` remains
128+
authoritative, while the underlying exception preserves the same reason
129+
if a caller unwraps it before classification.
130+
131+
``prior_exc`` is where that earlier, superseded failure goes: diagnostic
132+
context for logs and post-mortems, deliberately outside the field
133+
classification reads.
134+
"""
135+
136+
def __init__(
137+
self,
138+
last_exc: BaseException,
139+
reason: str,
140+
*,
141+
prior_exc: BaseException | None = None,
142+
) -> None:
143+
self.last_exc = last_exc
144+
self.reason = reason
145+
self.prior_exc = prior_exc
146+
detail = f"call_llm {reason}: {last_exc!r}"
147+
if prior_exc is not None and prior_exc is not last_exc:
148+
detail += f" (after {prior_exc!r})"
149+
super().__init__(detail)
150+
151+
152+
__all__ = [
153+
"AgentCoreError",
154+
"InvalidStateTransition",
155+
"KernelError",
156+
"LLMCallExhausted",
157+
"LLMDeadlineExceeded",
158+
"LLMError",
159+
"LLMReasoningRunaway",
160+
"LLMStreamStalled",
161+
"PermissionDenied",
162+
"ServiceNotRegistered",
163+
"TaskNotFoundError",
164+
]

‎agent_core/llm.py‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,9 +37,9 @@ class StreamDelta:
3737
usage: dict[str, int] = field(default_factory=dict[str, int])
3838
finish_reason: str = ""
3939
model: str = ""
40-
# Vendor label of the leg serving this stream, stamped by
41-
# ``LLMFallbackChain.stream`` (constant once the chain commits to an
42-
# entry — failover only fires before the first yield). The stream
40+
# Vendor label of the leg serving this stream, stamped by a product's
41+
# provider-chain wrapper (constant once the chain commits to an entry —
42+
# failover only fires before the first yield). The stream
4343
# assembler folds it into ``LLMResponse.response_metadata`` so per-call
4444
# billing attribution works for streamed calls too — without this the
4545
# streaming path had no channel for the provider and every billing

‎agent_core/messages.py‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,8 @@
2020
is preserved as the ``content`` value when the underlying client returns it;
2121
callers that need flat text use :func:`text_of`.
2222
23-
This module replaces ``langchain_core.messages`` (BaseMessage / SystemMessage /
24-
HumanMessage / AIMessage / ToolMessage). It is intentionally dependency-free.
23+
This module replaces the former framework-specific message classes with a
24+
small, dependency-free wire contract.
2525
"""
2626

2727
from __future__ import annotations
@@ -34,8 +34,8 @@
3434
class ToolCall(TypedDict):
3535
"""OpenAI-style tool_call payload — ``function.arguments`` is JSON-encoded.
3636
37-
Wire key order is fixed ``{type, id, function}`` to match the LangChain
38-
serializer the served checkpoints were aligned against; do not reorder.
37+
Wire key order is fixed ``{type, id, function}`` to match the serializer
38+
byte shape the served checkpoints were aligned against; do not reorder.
3939
"""
4040

4141
id: str
@@ -141,7 +141,7 @@ def for_wire(messages: list[Message]) -> list[Message]:
141141

142142
# Key insertion order: ``content`` first, then ``role``. Some served
143143
# checkpoints are sensitive to this byte shape (wire byte-equality with
144-
# LangChain's ``_convert_message_to_dict`` — see migration gotcha #2). Do
144+
# the legacy message serializer — see migration gotcha #2). Do
145145
# not reorder these dict literals.
146146

147147

‎agent_core/runtime/env.py‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
"""Shared env-variable prefix cascade.
2+
3+
Several modules independently re-walked the same
4+
``AGENT_CORE_ / MIROHARNESS_ / FRONTIER_AGENT_`` prefix order looking for a
5+
configured value — one copy per module, easy to drift if a prefix is ever
6+
added or reordered in only one of them. This is the single implementation
7+
they converge on: import :func:`first_configured` and let it own the order
8+
rather than passing a locally-spelled tuple back in.
9+
"""
10+
11+
from __future__ import annotations
12+
13+
import os
14+
15+
# The portable ``AGENT_CORE_`` spelling wins when multiple aliases are
16+
# configured, followed by the MiroHarness and FrontierAgent compatibility
17+
# names. Order matters:
18+
# callers rely on the first configured prefix winning.
19+
ENV_PREFIXES = ("AGENT_CORE_", "MIROHARNESS_", "FRONTIER_AGENT_")
20+
21+
22+
def first_configured(suffix: str, prefixes: tuple[str, ...] = ENV_PREFIXES) -> tuple[str, str] | None:
23+
"""Return the ``(name, value)`` of the first non-empty ``{prefix}{suffix}`` env var.
24+
25+
``None`` when none of the prefixed names are set (or all are blank).
26+
"""
27+
for prefix in prefixes:
28+
name = f"{prefix}{suffix}"
29+
raw = os.environ.get(name, "").strip()
30+
if raw:
31+
return name, raw
32+
return None
33+
34+
35+
__all__ = ["ENV_PREFIXES", "first_configured"]
Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
1+
"""Task-local overrides for one physical LLM request.
2+
3+
The runtime occasionally needs to change generation behaviour for one retry
4+
without mutating a cached/shared client. ``ContextVar`` keeps that override
5+
isolated across concurrent tasks and automatically restores the client's
6+
normal profile on exit.
7+
8+
Provider adapters opt in to the semantic override they understand. Today the
9+
OpenAI-compatible adapter maps it onto SGLang/Qwen
10+
``chat_template_kwargs`` and an explicitly configured ``reasoning_effort``.
11+
Unsupported adapters simply keep their normal request shape; the runtime's
12+
retry prompt and output cap remain the portable fallback.
13+
"""
14+
15+
from __future__ import annotations
16+
17+
from collections.abc import Generator
18+
from contextlib import contextmanager
19+
from contextvars import ContextVar
20+
from dataclasses import dataclass
21+
22+
23+
@dataclass(frozen=True)
24+
class ThinkingRetryOverride:
25+
"""Semantic thinking controls for a single retry attempt."""
26+
27+
mode: str = "reduced"
28+
thinking_budget: int | None = None
29+
reasoning_effort: str | None = None
30+
31+
@property
32+
def disabled(self) -> bool:
33+
return self.mode == "disabled"
34+
35+
36+
_THINKING_RETRY_OVERRIDE: ContextVar[ThinkingRetryOverride | None] = ContextVar(
37+
"agent_core_thinking_retry_override",
38+
default=None,
39+
)
40+
41+
42+
def current_thinking_retry_override() -> ThinkingRetryOverride | None:
43+
"""Return the override active for the current async task, if any."""
44+
45+
return _THINKING_RETRY_OVERRIDE.get()
46+
47+
48+
@contextmanager
49+
def thinking_retry_override(
50+
override: ThinkingRetryOverride | None,
51+
) -> Generator[None, None, None]:
52+
"""Apply ``override`` only inside this context and async task."""
53+
54+
if override is None:
55+
yield
56+
return
57+
token = _THINKING_RETRY_OVERRIDE.set(override)
58+
try:
59+
yield
60+
finally:
61+
_THINKING_RETRY_OVERRIDE.reset(token)
62+
63+
64+
__all__ = [
65+
"ThinkingRetryOverride",
66+
"current_thinking_retry_override",
67+
"thinking_retry_override",
68+
]

0 commit comments

Comments
 (0)