Skip to content
Open
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
11 changes: 6 additions & 5 deletions agent_core/core/impl/action/context.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,9 @@
Scope rules:
- Set only by the internal executors (``_atomic_action_internal*``),
reset in a ``finally`` — never leaks across actions.
- Sync actions run in a thread pool where the caller's context does NOT
propagate, so the executor wraps the call and sets the var inside the
worker thread (see ``run_with_input_context``).
- Sync actions run in a thread pool inside a copy of the caller's
context; the executor wraps the call and sets the var inside the worker
thread (see ``run_with_input_context``).
- Sandboxed (subprocess) actions cannot see it at all — helpers must
treat a ``None`` value as "no context available".
"""
Expand All @@ -31,8 +31,9 @@ def run_with_input_context(
) -> dict:
"""Call a sync action with ``current_input_data`` set for its duration.

Used as the thread-pool target: the worker thread has its own context,
so the var must be set (and reset) inside the thread, not the caller.
Used as the thread-pool target, run inside a copy of the caller's
context: the var is set (and reset) inside the thread so it scopes to
this action alone.
"""
token = current_input_data.set(input_data)
try:
Expand Down
15 changes: 11 additions & 4 deletions agent_core/core/impl/action/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
"""

import asyncio
import contextvars
import importlib
import json
import os
Expand Down Expand Up @@ -634,14 +635,20 @@ async def _atomic_action_internal_async(
finally:
current_input_data.reset(ctx_token)
else:
# Sync function - run in thread pool to avoid blocking. The
# worker thread doesn't inherit this context, so the wrapper
# sets current_input_data inside the thread.
# Sync function - run in thread pool to avoid blocking. A pool
# thread does not inherit the caller's context, so the action
# runs inside a copy of it (as asyncio.to_thread does): the bound
# session, log context, and other context-scoped state follow the
# action into the thread, and the wrapper sets current_input_data
# there.
logger.debug(
f"[SYNC] Action '{action_name}' is sync, running in thread pool"
)
thread_future = THREAD_POOL.submit(
run_with_input_context, function_to_call, input_data
contextvars.copy_context().run,
run_with_input_context,
function_to_call,
input_data,
)
try:
execution_result = await asyncio.wrap_future(thread_future)
Expand Down
21 changes: 18 additions & 3 deletions agent_core/core/impl/llm/cache/gemini.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
import hashlib
import logging
import time
from typing import Any, Dict, TYPE_CHECKING
from typing import Any, Dict, Optional, TYPE_CHECKING

from .config import get_cache_config

Expand Down Expand Up @@ -78,21 +78,32 @@ def get_or_create_cache(
system_prompt: str,
user_prompt: str,
call_type: str,
temperature: float,
temperature: Optional[float],
max_tokens: int,
thinking_budget: Optional[int] = None,
thinking_level: Optional[str] = None,
) -> Dict[str, Any]:
"""Get response using explicit cache, creating cache if needed.

Args:
system_prompt: The system prompt to cache.
user_prompt: The user prompt for this request.
call_type: Type of LLM call (e.g., "reasoning", "action_selection").
temperature: Sampling temperature.
temperature: Sampling temperature (None: not sent).
max_tokens: Maximum output tokens.
thinking_budget: Reasoning token budget (Gemini 2.5), forwarded
to every generation call.
thinking_level: Reasoning level (Gemini 3.x), forwarded to every
generation call.

Returns:
Response dict with tokens_used, content, cached_tokens, etc.
"""
thinking: Dict[str, Any] = {
"thinking_budget": thinking_budget,
"thinking_level": thinking_level,
}

# Check if system prompt is large enough for explicit caching
# Gemini requires at least 1024 tokens; skip explicit cache if too small
estimated_tokens = self._estimate_tokens(system_prompt)
Expand All @@ -109,6 +120,7 @@ def get_or_create_cache(
temperature=temperature,
max_output_tokens=max_tokens,
json_mode=True,
**thinking,
)

cache_key = self._make_cache_key(system_prompt, call_type)
Expand All @@ -132,6 +144,7 @@ def get_or_create_cache(
temperature=temperature,
max_output_tokens=max_tokens,
json_mode=True,
**thinking,
)
except Exception as e:
logger.warning(
Expand Down Expand Up @@ -166,6 +179,7 @@ def get_or_create_cache(
temperature=temperature,
max_output_tokens=max_tokens,
json_mode=True,
**thinking,
)
except Exception as e:
logger.warning(
Expand All @@ -185,6 +199,7 @@ def get_or_create_cache(
temperature=temperature,
max_output_tokens=max_tokens,
json_mode=True,
**thinking,
)

def invalidate_cache(self, system_prompt: str, call_type: str) -> None:
Expand Down
Loading