From 1a922c927c3d199a37397c556e575879717e92ae Mon Sep 17 00:00:00 2001 From: ahmad-ajmal Date: Tue, 22 Sep 2026 11:14:36 +0100 Subject: [PATCH 1/3] cleanup: remove observation logic deadcode --- agent_core/__init__.py | 3 +- agent_core/core/action/__init__.py | 3 +- agent_core/core/action/action.py | 10 --- agent_core/core/action/observe.py | 120 ------------------------- agent_core/core/impl/action/manager.py | 68 -------------- agent_core/core/protocols/action.py | 2 +- app/action/__init__.py | 2 - 7 files changed, 3 insertions(+), 205 deletions(-) delete mode 100644 agent_core/core/action/observe.py diff --git a/agent_core/__init__.py b/agent_core/__init__.py index a7b399f8d..90b6e4948 100644 --- a/agent_core/__init__.py +++ b/agent_core/__init__.py @@ -64,7 +64,7 @@ BYTEPLUS_MAX_INPUT_TOKENS, GeminiCacheManager, ) -from agent_core.core.action import Action, Observe +from agent_core.core.action import Action from agent_core.core.event_stream import Event, EventRecord from agent_core.decorators import ( profile, @@ -268,7 +268,6 @@ "GeminiCacheManager", # Action framework "Action", - "Observe", "ActionRegistry", "ActionMetadata", "RegisteredAction", diff --git a/agent_core/core/action/__init__.py b/agent_core/core/action/__init__.py index 71cf29e40..235558e2a 100644 --- a/agent_core/core/action/__init__.py +++ b/agent_core/core/action/__init__.py @@ -2,6 +2,5 @@ """Action framework for defining and executing agent actions.""" from agent_core.core.action.action import Action -from agent_core.core.action.observe import Observe -__all__ = ["Action", "Observe"] +__all__ = ["Action"] diff --git a/agent_core/core/action/action.py b/agent_core/core/action/action.py index c797b5211..5f82e61a5 100644 --- a/agent_core/core/action/action.py +++ b/agent_core/core/action/action.py @@ -9,8 +9,6 @@ import datetime from typing import Optional, List, Dict, Any -from agent_core.core.action.observe import Observe - class Action: """ @@ -30,7 +28,6 @@ class Action: input_schema: Schema describing expected inputs output_schema: Schema describing expected outputs sub_actions: Child actions for divisible actions - observer: Optional observation step for validation platforms: List of supported platforms platform_overrides: Platform-specific code/schema overrides requirements: List of pip packages required @@ -53,7 +50,6 @@ def __init__( input_schema: Optional[dict] = None, output_schema: Optional[dict] = None, sub_actions: Optional[List["Action"]] = None, - observer: Optional[Observe] = None, last_use: bool = None, default: bool = False, platforms: List[str] = ["windows", "linux", "darwin"], @@ -87,7 +83,6 @@ def __init__( output_schema: Schema describing expected outputs in the same format as input_schema. sub_actions: Child actions to run when action_type is "divisible". - observer: Optional observation step to validate outputs after execution. last_use: Timestamp or marker for last usage, used for analytics. default: Whether this action should be offered as a default choice in routing flows. @@ -120,7 +115,6 @@ def __init__( self.output_schema = output_schema or {} self.sub_actions = sub_actions or [] - self.observer = observer self.created_at = datetime.datetime.utcnow().isoformat() self.updated_at = self.created_at self.last_use = last_use @@ -163,7 +157,6 @@ def to_dict(self) -> Dict[str, Any]: "input_schema": self.input_schema, "output_schema": self.output_schema, "subActions": [sub_action.to_dict() for sub_action in self.sub_actions], - "observer": self.observer.to_dict() if self.observer else None, "createdAt": self.created_at, "updatedAt": self.updated_at, "lastUse": self.last_use, @@ -189,8 +182,6 @@ def from_dict(cls, data: Dict[str, Any]) -> "Action": Action instance """ sub_actions = [cls.from_dict(sub) for sub in data.get("subActions", [])] - observer_data = data.get("observer") - observer = Observe.from_dict(observer_data) if observer_data else None # Fallback logic for older fields if input_schema/output_schema not present input_schema = data.get("input_schema") or data.get("input") or {} @@ -210,7 +201,6 @@ def from_dict(cls, data: Dict[str, Any]) -> "Action": input_schema=input_schema, output_schema=output_schema, sub_actions=sub_actions, - observer=observer, default=data.get("default", False), platforms=data.get("platforms", ["windows", "linux", "darwin"]), platform_overrides=data.get("platform_overrides", {}), diff --git a/agent_core/core/action/observe.py b/agent_core/core/action/observe.py deleted file mode 100644 index 10661abf5..000000000 --- a/agent_core/core/action/observe.py +++ /dev/null @@ -1,120 +0,0 @@ -# -*- coding: utf-8 -*- -""" -Observation class for action validation. - -Although implemented, this observe is not USED at all yet. -This observe will be fired after action completion, if implemented. -This causes the action to have an immediate validation step followed up. - -For example, action that creates a folder path will be followed by an observation step -to make sure the folder path is created successfully. This creates another -layer of validation. -""" - -from typing import Optional, Dict, Any - - -class Observe: - """ - Defines how to confirm that an action completed successfully in the real world. - - Observation logic is defined as Python code that executes repeatedly until - success or timeout is reached. - - Attributes: - name: Identifier for this observation (e.g., "check_file_created") - description: Human-readable description of what is being observed - code: Python code to confirm action success - retry_interval_sec: Seconds between retry attempts (default: 3) - max_retries: Maximum number of retry attempts (default: 20) - max_total_time_sec: Maximum total time allowed regardless of retries (default: 60) - wait_to_observe_sec: Optional delay before first observation or between observations - input_schema: Schema for observation inputs - success: Final result of observation (True/False/None) - message: Optional output message from observation - """ - - def __init__( - self, - name: str, - description: Optional[str] = None, - code: Optional[str] = None, - retry_interval_sec: int = 3, - max_retries: int = 20, - max_total_time_sec: int = 60, - wait_to_observe_sec: Optional[int] = None, - input_schema: Optional[dict] = None, - success: Optional[bool] = None, - message: Optional[str] = None, - ): - """ - Initialize an Observe instance. - - Args: - name: Unique identifier for this observation - description: Human-readable description - code: Python code to execute for validation - retry_interval_sec: Seconds between retries (default: 3) - max_retries: Maximum retry attempts (default: 20) - max_total_time_sec: Maximum total time in seconds (default: 60) - wait_to_observe_sec: Optional initial wait time - input_schema: Dictionary describing expected inputs - success: Whether observation succeeded (set after execution) - message: Optional message from observation - """ - self.name = name - self.description = description - self.code = code - - self.retry_interval_sec = retry_interval_sec - self.max_retries = max_retries - self.max_total_time_sec = max_total_time_sec - self.wait_to_observe_sec = wait_to_observe_sec - - self.input_schema = input_schema or {} - self.success = success - self.message = message - - def to_dict(self) -> Dict[str, Any]: - """ - Convert Observe to dictionary format for serialization. - - Returns: - Dictionary representation of the observation - """ - return { - "name": self.name, - "description": self.description, - "code": self.code, - "retry_interval_sec": self.retry_interval_sec, - "max_retries": self.max_retries, - "max_total_time_sec": self.max_total_time_sec, - "wait_to_observe_sec": self.wait_to_observe_sec, - "input_schema": self.input_schema, - "success": self.success, - "message": self.message, - } - - @classmethod - def from_dict(cls, data: Dict[str, Any]) -> "Observe": - """ - Create an Observe instance from a dictionary. - - Args: - data: Dictionary containing observation data - - Returns: - Observe instance - """ - return cls( - name=data["name"], - description=data.get("description"), - code=data.get("code"), - retry_interval_sec=data.get("retry_interval_sec", 3), - max_retries=data.get("max_retries", 20), - max_total_time_sec=data.get("max_total_time_sec", 600), - wait_to_observe_sec=data.get("wait_to_observe_sec"), - input_schema=data.get("input_schema") or {}, - success=data.get("success"), - message=data.get("message"), - ) diff --git a/agent_core/core/impl/action/manager.py b/agent_core/core/impl/action/manager.py index 7f350bd4f..529b88b41 100644 --- a/agent_core/core/impl/action/manager.py +++ b/agent_core/core/impl/action/manager.py @@ -9,13 +9,10 @@ from datetime import datetime import platform -import time import json import asyncio import nest_asyncio from typing import Optional, List, Dict, Any, Callable, Tuple -import io -import sys import re import uuid @@ -404,21 +401,6 @@ async def execute_action( f"[OUTPUT DATA] Completed execute_atomic_action: {outputs}" ) - # Observation step - if action.observer: - obs_result = await self.run_observe_step(action, outputs) - if not obs_result["success"]: - status = "error" - outputs["observation"] = { - "success": False, - "message": obs_result.get("message"), - } - else: - outputs["observation"] = { - "success": True, - "message": obs_result.get("message"), - } - else: logger.debug(f"Executing divisible action: {action.name}") try: @@ -855,56 +837,6 @@ async def execute_divisible_action(self, action, input_data, parent_id) -> Dict: ) return results - @profile("action_manager_run_observe_step", OperationCategory.ACTION_EXECUTION) - async def run_observe_step( - self, action: Action, action_output: Dict - ) -> Dict[str, Any]: - """ - Executes the observation code with retries, to confirm action outcome. - """ - observe = action.observer - if not observe or not observe.code: - return {"success": True, "message": "No observation step."} - - input_json = json.dumps(action_output) - python_script = f"""import json;output = {input_json};{observe.code}""" - - attempt = 0 - start_time = time.time() - while ( - attempt < observe.max_retries - and (time.time() - start_time) < observe.max_total_time_sec - ): - stdout_buf = io.StringIO() - stderr_buf = io.StringIO() - - sys.stdout = stdout_buf - sys.stderr = stderr_buf - local_env = {} - - try: - exec(python_script, {}, local_env) - sys.stdout = sys.__stdout__ - sys.stderr = sys.__stderr__ - - success = local_env.get("success", None) - message = local_env.get("message", "") - - if success is True: - return {"success": True, "message": message} - elif success is False: - return {"success": False, "message": message} - - except Exception as e: - sys.stdout = sys.__stdout__ - sys.stderr = sys.__stderr__ - logger.warning(f"[OBSERVE] Error during observation: {e}") - - await asyncio.sleep(observe.retry_interval_sec) - attempt += 1 - - return {"success": False, "message": "Observation failed or timed out."} - @staticmethod def _extract_base64_to_files(data: dict, action_name: str) -> dict: """ diff --git a/agent_core/core/protocols/action.py b/agent_core/core/protocols/action.py index 33b8b50aa..b303551db 100644 --- a/agent_core/core/protocols/action.py +++ b/agent_core/core/protocols/action.py @@ -189,7 +189,7 @@ class ActionManagerProtocol(Protocol): Protocol for action orchestration. This defines the minimal interface for managing action execution - lifecycles, including observation steps and history logging. + lifecycles and history logging. """ async def execute_action( diff --git a/app/action/__init__.py b/app/action/__init__.py index 629d3f27e..326a374c3 100644 --- a/app/action/__init__.py +++ b/app/action/__init__.py @@ -8,7 +8,6 @@ # Re-export from agent_core from agent_core import ( Action, - Observe, ActionExecutor, ActionLibrary, ActionRouter, @@ -21,7 +20,6 @@ __all__ = [ # From agent_core "Action", - "Observe", "ActionExecutor", "ActionLibrary", "ActionRouter", From f17090f7c45004bba143af4c5d42d4c9389de450 Mon Sep 17 00:00:00 2001 From: ahmad-ajmal Date: Tue, 22 Sep 2026 11:21:05 +0100 Subject: [PATCH 2/3] fix: stdout/strerr --- installer/api.py | 58 ++++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 54 insertions(+), 4 deletions(-) diff --git a/installer/api.py b/installer/api.py index e9f28f376..ac804381b 100644 --- a/installer/api.py +++ b/installer/api.py @@ -154,10 +154,15 @@ def _dispatch(self, label: str, fn: Callable[[], None]) -> dict: self._push_log(f"\n[{label}] Already running, ignoring click.\n") return {"started": False, "reason": "busy"} + stdout, stderr = _install_routed_streams() + def target() -> None: - saved_stdout, saved_stderr = sys.stdout, sys.stderr - sys.stdout = _BridgeWriter(self) - sys.stderr = _BridgeWriter(self) + # Route only THIS thread's prints into the log panel. The process + # streams are never swapped here, so the bridge thread polling + # get_state and any other thread keep writing to the real streams. + bridge = _BridgeWriter(self) + stdout.route_current_thread(bridge) + stderr.route_current_thread(bridge) try: self._push_log(f"\n━━━ {label} ━━━\n") fn() @@ -165,7 +170,8 @@ def target() -> None: except Exception as exc: self._push_log(f"\n[{label}] ERROR: {exc!r}\n") finally: - sys.stdout, sys.stderr = saved_stdout, saved_stderr + stdout.route_current_thread(None) + stderr.route_current_thread(None) self._push_event("workerDone", {"label": label}) self._worker = threading.Thread(target=target, daemon=True) @@ -285,3 +291,47 @@ def flush(self) -> None: def isatty(self) -> bool: return False + + +class _ThreadRoutedStream: + """Process stream that sends writes from a routed thread to that + thread's sink, and everything else to the original stream. + + Installed once for the life of the process (see _install_routed_streams) + instead of swapping sys.stdout/sys.stderr per action: a per-action swap + is process-wide, so every other thread's output would be captured too.""" + + def __init__(self, fallback) -> None: + self._fallback = fallback + self._local = threading.local() + + def route_current_thread(self, sink) -> None: + self._local.sink = sink + + def _target(self): + return getattr(self._local, "sink", None) or self._fallback + + def write(self, text: str) -> int: + return self._target().write(text) + + def flush(self) -> None: + self._target().flush() + + def isatty(self) -> bool: + return self._target().isatty() + + def __getattr__(self, name: str): + # encoding, fileno, buffer, ... come from the real stream. + return getattr(self._fallback, name) + + +_routed_streams_lock = threading.Lock() + + +def _install_routed_streams() -> tuple[_ThreadRoutedStream, _ThreadRoutedStream]: + with _routed_streams_lock: + if not isinstance(sys.stdout, _ThreadRoutedStream): + sys.stdout = _ThreadRoutedStream(sys.stdout) + if not isinstance(sys.stderr, _ThreadRoutedStream): + sys.stderr = _ThreadRoutedStream(sys.stderr) + return sys.stdout, sys.stderr From 1179f312fd57b874a55c3035d27bf86da4d7bcfb Mon Sep 17 00:00:00 2001 From: Korivi Date: Wed, 23 Sep 2026 09:26:02 +0900 Subject: [PATCH 3/3] fix(installer): survive absent process streams when routing stdout _install_routed_streams wrapped sys.stdout/sys.stderr unconditionally. The installer is frozen with console=False (packaging/CraftBotInstaller.spec), so under pythonw both are None -- the wrapper then held None as its fallback and every write from a non-routed thread raised "AttributeError: 'NoneType' object has no attribute 'write'". CPython's print() is a silent no-op when sys.stdout is None, so before this the same stray library print did nothing. Wrapping turned a no-op into a crash in the shipped, windowed build, and sys.stdout.encoding -- read by a lot of ordinary code -- crashed with it. A _NullStream stand-in restores the no-op and answers the attributes a stream is expected to have, so the per-thread routing this commit's parent introduced keeps working with or without a console. Reproduced by setting sys.stdout = None and calling print() through the wrapper: AttributeError before, silent no-op after, with worker-thread routing unchanged. Tests live under tests/ so `pytest` actually collects them (pytest.ini sets testpaths = tests). --- installer/api.py | 33 +++++++- tests/test_installer_stream_routing.py | 109 +++++++++++++++++++++++++ 2 files changed, 141 insertions(+), 1 deletion(-) create mode 100644 tests/test_installer_stream_routing.py diff --git a/installer/api.py b/installer/api.py index ac804381b..93d9349a1 100644 --- a/installer/api.py +++ b/installer/api.py @@ -293,6 +293,37 @@ def isatty(self) -> bool: return False +class _NullStream: + """Stand-in for a process stream that does not exist. + + The installer is frozen with console=False (packaging/CraftBotInstaller.spec), + so under pythonw sys.stdout and sys.stderr are None. CPython's print() + treats a None stream as a silent no-op, so anything WRAPPING it has to do + the same -- wrapping None directly turns every stray library print on a + non-routed thread into AttributeError: 'NoneType' has no attribute 'write', + and sys.stdout.encoding into a crash on a very common attribute read.""" + + def write(self, text: str) -> int: + return len(text) + + def flush(self) -> None: + pass + + def isatty(self) -> bool: + return False + + @property + def encoding(self) -> str: + return "utf-8" + + @property + def errors(self) -> str: + return "replace" + + def fileno(self) -> int: + raise OSError("stream has no file descriptor") + + class _ThreadRoutedStream: """Process stream that sends writes from a routed thread to that thread's sink, and everything else to the original stream. @@ -302,7 +333,7 @@ class _ThreadRoutedStream: is process-wide, so every other thread's output would be captured too.""" def __init__(self, fallback) -> None: - self._fallback = fallback + self._fallback = fallback if fallback is not None else _NullStream() self._local = threading.local() def route_current_thread(self, sink) -> None: diff --git a/tests/test_installer_stream_routing.py b/tests/test_installer_stream_routing.py new file mode 100644 index 000000000..bb193970f --- /dev/null +++ b/tests/test_installer_stream_routing.py @@ -0,0 +1,109 @@ +"""The installer's per-thread stdout routing. + +Two things have to hold at once, and they pull in opposite directions: + + * a worker thread's output goes to the wizard's log panel and nowhere else + (the whole point of routing per thread rather than swapping the process + streams, which would capture every other thread too); + * nothing crashes when the process has no streams at all. The installer is + frozen with console=False, so under pythonw sys.stdout and sys.stderr are + None. print() is a documented no-op in that case, and a wrapper that does + not preserve it turns a harmless library print into AttributeError. +""" + +import io +import sys +import threading + +import pytest + +from installer.api import _install_routed_streams, _NullStream, _ThreadRoutedStream + + +@pytest.fixture +def restore_streams(): + saved = sys.stdout, sys.stderr + yield + sys.stdout, sys.stderr = saved + + +def test_wrapping_absent_streams_keeps_print_a_no_op(restore_streams): + """The frozen, windowed installer: sys.stdout is None.""" + sys.stdout = None + sys.stderr = None + + out, err = _install_routed_streams() + + print("a library logging from a non-routed thread") # must not raise + print("and on stderr", file=sys.stderr) + sys.stdout.flush() + assert sys.stdout.encoding == "utf-8" + assert sys.stdout.isatty() is False + assert isinstance(out, _ThreadRoutedStream) + + +def test_absent_streams_still_route_to_a_worker_sink(restore_streams): + sys.stdout = None + sys.stderr = None + _install_routed_streams() + + sink = io.StringIO() + sys.stdout.route_current_thread(sink) + print("worker line") + sys.stdout.route_current_thread(None) + + assert sink.getvalue() == "worker line\n" + + +def test_only_the_routed_thread_is_captured(restore_streams): + """A per-action swap of sys.stdout is process-wide: the bridge thread + polling get_state, and every other thread, would be captured with it.""" + console = io.StringIO() + sink = io.StringIO() + sys.stdout = _ThreadRoutedStream(console) + + def worker(): + sys.stdout.route_current_thread(sink) + print("from the worker") + sys.stdout.route_current_thread(None) + + thread = threading.Thread(target=worker) + thread.start() + thread.join() + print("from the main thread") + + assert sink.getvalue() == "from the worker\n" + assert console.getvalue() == "from the main thread\n" + + +def test_routing_is_cleared_even_though_the_stream_persists(restore_streams): + console = io.StringIO() + sink = io.StringIO() + sys.stdout = _ThreadRoutedStream(console) + + sys.stdout.route_current_thread(sink) + sys.stdout.route_current_thread(None) + print("after the worker finished") + + assert sink.getvalue() == "" + assert console.getvalue() == "after the worker finished\n" + + +def test_install_is_idempotent(restore_streams): + """Called once per action; it must not wrap a wrapper each time.""" + console = io.StringIO() + sys.stdout = console + sys.stderr = console + + first_out, _ = _install_routed_streams() + second_out, _ = _install_routed_streams() + + assert first_out is second_out + assert first_out._fallback is console + + +def test_null_stream_reports_no_descriptor(): + """Callers that probe fileno() (subprocess wiring, isatty shims) get the + same OSError a detached stream raises, not a wrong file descriptor.""" + with pytest.raises(OSError): + _NullStream().fileno()