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
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
4 changes: 3 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -59,4 +59,6 @@ agent_file_system/ACTIONS.md
agent_bundle/
**/.craftbot/
app/data/.file_index/
.playwright-mcp
.playwright-mcp
# Sidecar Node runtime (install.py downloads it when the system Node is too old for Living UI)
runtime/
9 changes: 9 additions & 0 deletions agent_core/core/event_stream/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,12 @@ class Event:
uses it to keep the run's "Working…" indicator up across the
bubble instead of treating every agent bubble as a run-ending
reply. None/False for final replies and non-chat events.
question: For AGENT_MESSAGE events only: set when the message is a
question to the user with suggested responses (send_message with
suggested_responses). Shape:
``{"options": ["Yes", "No"], "allow_free_text": true}``. The UI
renders it as answer chips plus a pinned question box above the
chat composer. None for ordinary messages.
"""

message: str
Expand All @@ -157,6 +163,7 @@ class Event:
action_output: Optional[Dict[str, Any]] = None
platform: Optional[str] = None
continue_work: Optional[bool] = None
question: Optional[Dict[str, Any]] = None

def display_text(self) -> Optional[str]:
"""
Expand Down Expand Up @@ -189,6 +196,7 @@ def to_dict(self) -> Dict[str, Any]:
"action_output": self.action_output,
"platform": self.platform,
"continue_work": self.continue_work,
"question": self.question,
}

@classmethod
Expand Down Expand Up @@ -228,6 +236,7 @@ def from_dict(cls, data: Dict[str, Any]) -> "Event":
action_output=data.get("action_output"),
platform=data.get("platform"),
continue_work=data.get("continue_work"),
question=data.get("question"),
)

@property
Expand Down
41 changes: 41 additions & 0 deletions agent_core/core/impl/action/context.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""Execution-scoped context for in-process actions.

``current_input_data`` holds the full ``input_data`` dict of the action
currently executing in this context. It exists so cross-cutting helpers
deep inside an action's call tree (e.g. multi-account routing reading the
``account`` hint) can see routing keys without threading them through
every action function signature.

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``).
- Sandboxed (subprocess) actions cannot see it at all — helpers must
treat a ``None`` value as "no context available".
"""

from __future__ import annotations

from contextvars import ContextVar
from typing import Any, Callable, Dict, Optional

current_input_data: ContextVar[Optional[Dict[str, Any]]] = ContextVar(
"current_input_data", default=None
)


def run_with_input_context(
function_to_call: Callable[[dict], dict], input_data: dict
) -> 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.
"""
token = current_input_data.set(input_data)
try:
return function_to_call(input_data)
finally:
current_input_data.reset(token)
23 changes: 19 additions & 4 deletions agent_core/core/impl/action/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -571,7 +571,9 @@ def _atomic_action_internal(
"The action_code string did not define a callable Python function."
)

execution_result = function_to_call(input_data)
from agent_core.core.impl.action.context import run_with_input_context

execution_result = run_with_input_context(function_to_call, input_data)
return execution_result

except Exception as e:
Expand Down Expand Up @@ -618,16 +620,29 @@ async def _atomic_action_internal_async(
"The action_code string did not define a callable Python function."
)

from agent_core.core.impl.action.context import (
current_input_data,
run_with_input_context,
)

# Check if the function is async (coroutine function)
if inspect.iscoroutinefunction(function_to_call):
logger.debug(f"[ASYNC] Action '{action_name}' is async, awaiting directly")
execution_result = await function_to_call(input_data)
ctx_token = current_input_data.set(input_data)
try:
execution_result = await function_to_call(input_data)
finally:
current_input_data.reset(ctx_token)
else:
# Sync function - run in thread pool to avoid blocking
# 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.
logger.debug(
f"[SYNC] Action '{action_name}' is sync, running in thread pool"
)
thread_future = THREAD_POOL.submit(function_to_call, input_data)
thread_future = THREAD_POOL.submit(
run_with_input_context, function_to_call, input_data
)
try:
execution_result = await asyncio.wrap_future(thread_future)
except asyncio.CancelledError:
Expand Down
40 changes: 36 additions & 4 deletions agent_core/core/impl/action/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,41 @@ async def _compat_wait_for(fut, timeout):

nest_asyncio.apply()

# ============================================================================
# Second half of the nest_asyncio/3.14 shim: heal asyncio.current_task().
# nest_asyncio forces the PURE-PYTHON asyncio.Task class, whose tasks
# register in the Python-side registry (asyncio.tasks._py_current_task) —
# but asyncio.current_task stays bound to the C-accelerated registry, so it
# returns None inside EVERY task, on EVERY loop, process-wide. Everything
# built on `async with asyncio.timeout(...)` then dies with "Timeout
# (context manager) should be used inside a task" — most visibly the entire
# aiohttp CLIENT (every request enters a timeout context), which is what
# broke the external A2App adapter self-check on 2026-08-24 while the
# aiohttp SERVER (no timeout context on the request path) kept working.
# Rebinding current_task to the Python registry fixes timeout/aiohttp under
# both plain awaits and nested re-entry (verified on 3.14.7 + aiohttp
# 3.14.3). The wait_for replacement above stays: its explicit
# cancellation-wait semantics are load-bearing for force-stop (PR #410).
try:
import _asyncio as _compat_c_asyncio

if asyncio.Task is not getattr(_compat_c_asyncio, "Task", None) and hasattr(
asyncio.tasks, "_py_current_task"
):
asyncio.current_task = asyncio.tasks._py_current_task
asyncio.tasks.current_task = asyncio.tasks._py_current_task
try:
_compat_sys.stderr.write(
"[compat-shim] asyncio.current_task routed to the Python "
"task registry (action/manager)\n"
)
_compat_sys.stderr.flush()
except Exception:
pass
except Exception as _compat_ct_exc:
logger.warning(f"[compat-shim] current_task rebinding skipped: {_compat_ct_exc!r}")
# ============================================================================


def _to_pretty_json(value: Any) -> str:
"""Serialize a value to pretty-printed JSON for readable logs and event streams."""
Expand Down Expand Up @@ -247,10 +282,7 @@ async def execute_action(
# re-execute work the ledger shows as already completed (or as
# interrupted mid-flight, where the effect may have happened).
idem_key = None
# if getattr(action, "irreversible", False) and self._idempotency_guard:

# TODO: Temporary turning idempotency guard off.
if 1 == 0:
if getattr(action, "irreversible", False) and self._idempotency_guard:
try:
decision = self._idempotency_guard.begin(
action.name, input_data, session_id
Expand Down
1 change: 0 additions & 1 deletion agent_core/core/impl/config/watcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,6 @@ class ConfigWatcher:
- settings.json
- mcp_config.json
- skills_config.json
- external_comms_config.json

When a file changes, the appropriate reload callback is invoked.
"""
Expand Down
6 changes: 4 additions & 2 deletions agent_core/core/impl/context/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -457,8 +457,10 @@ def get_session_state(self, session_id: Optional[str] = None) -> str:
f"Session ID: {session.id}",
f"Session Type: {session.type}",
]
if session.title:
lines.append(f"Session Title: {session.title}")
# Session Title is intentionally omitted: it is auto-generated/
# updated a turn or two into a session, and this block sits in the
# cacheable prefix (ahead of the event stream), so a mutating title
# would break the KV-cache prefix every time it changed.
if getattr(session, "living_ui_project_id", None):
lines.append(f"Living UI Project: {session.living_ui_project_id}")
lines.append(f"Loaded Action Sets: {['core'] + list(session.action_sets)}")
Expand Down
Loading
Loading