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
3 changes: 1 addition & 2 deletions agent_core/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -268,7 +268,6 @@
"GeminiCacheManager",
# Action framework
"Action",
"Observe",
"ActionRegistry",
"ActionMetadata",
"RegisteredAction",
Expand Down
3 changes: 1 addition & 2 deletions agent_core/core/action/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
10 changes: 0 additions & 10 deletions agent_core/core/action/action.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,6 @@
import datetime
from typing import Optional, List, Dict, Any

from agent_core.core.action.observe import Observe


class Action:
"""
Expand All @@ -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
Expand All @@ -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"],
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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 {}
Expand All @@ -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", {}),
Expand Down
120 changes: 0 additions & 120 deletions agent_core/core/action/observe.py

This file was deleted.

68 changes: 0 additions & 68 deletions agent_core/core/impl/action/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
"""
Expand Down
2 changes: 1 addition & 1 deletion agent_core/core/protocols/action.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
2 changes: 0 additions & 2 deletions app/action/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
# Re-export from agent_core
from agent_core import (
Action,
Observe,
ActionExecutor,
ActionLibrary,
ActionRouter,
Expand All @@ -21,7 +20,6 @@
__all__ = [
# From agent_core
"Action",
"Observe",
"ActionExecutor",
"ActionLibrary",
"ActionRouter",
Expand Down
Loading