File size: 4,141 Bytes
92c4ae6 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 | """
Auto-Dev Event Hooks
Lightweight event bus for triggering learning loops from lifecycle events
in the episodic memory system. Decouples the core EpisodeService and
SandboxExecutor from the Auto-Dev engines.
Usage:
from core.auto_dev.event_hooks import event_bus
# Register a listener
@event_bus.on_task_fail
async def handle_failure(event: TaskEvent):
...
# Emit from EpisodeService
await event_bus.emit_task_fail(TaskEvent(...))
"""
import asyncio
import logging
from dataclasses import dataclass, field
from typing import Any, Callable, Coroutine
logger = logging.getLogger(__name__)
@dataclass
class TaskEvent:
"""Event payload for task lifecycle events."""
episode_id: str
agent_id: str
tenant_id: str
task_description: str = ""
error_trace: str | None = None
outcome: str = "" # "success", "failure", "partial"
metadata: dict[str, Any] = field(default_factory=dict)
@dataclass
class SkillExecutionEvent:
"""Event payload for skill execution events."""
execution_id: str
agent_id: str
tenant_id: str
skill_id: str
skill_name: str = ""
execution_seconds: float = 0.0
token_usage: int = 0
success: bool = False
output: str = ""
metadata: dict[str, Any] = field(default_factory=dict)
# Type alias for event handler functions
EventHandler = Callable[..., Coroutine[Any, Any, None]]
class EventBus:
"""
Simple in-process event bus for Auto-Dev lifecycle events.
Supports three event types:
- on_task_fail: triggered when an episode ends in failure
- on_task_success: triggered when an episode completes successfully
- on_skill_execution: triggered after a skill is executed in sandbox
"""
def __init__(self):
self._fail_handlers: list[EventHandler] = []
self._success_handlers: list[EventHandler] = []
self._skill_handlers: list[EventHandler] = []
# --- Decorator registration ---
def on_task_fail(self, handler: EventHandler) -> EventHandler:
"""Register a handler for task failure events."""
self._fail_handlers.append(handler)
return handler
def on_task_success(self, handler: EventHandler) -> EventHandler:
"""Register a handler for task success events."""
self._success_handlers.append(handler)
return handler
def on_skill_execution(self, handler: EventHandler) -> EventHandler:
"""Register a handler for skill execution events."""
self._skill_handlers.append(handler)
return handler
# --- Emission ---
async def emit_task_fail(self, event: TaskEvent) -> None:
"""Emit a task failure event to all registered handlers."""
await self._dispatch(self._fail_handlers, event, "on_task_fail")
async def emit_task_success(self, event: TaskEvent) -> None:
"""Emit a task success event to all registered handlers."""
await self._dispatch(self._success_handlers, event, "on_task_success")
async def emit_skill_execution(self, event: SkillExecutionEvent) -> None:
"""Emit a skill execution event to all registered handlers."""
await self._dispatch(self._skill_handlers, event, "on_skill_execution")
# --- Internal ---
async def _dispatch(
self, handlers: list[EventHandler], event: Any, event_name: str
) -> None:
"""Dispatch event to all handlers, catching exceptions to prevent cascade."""
if not handlers:
return
for handler in handlers:
try:
await handler(event)
except Exception as e:
logger.error(
f"Auto-Dev event handler error in {event_name} "
f"(handler={handler.__name__}): {e}",
exc_info=True,
)
def clear(self) -> None:
"""Remove all registered handlers. Useful for testing."""
self._fail_handlers.clear()
self._success_handlers.clear()
self._skill_handlers.clear()
# Global singleton — imported by EpisodeService and SandboxExecutor
event_bus = EventBus()
|