Преглед на файлове

feat: add generic event execution kernel

Problem: event resolution and execution were coupled to ToolRegistry and EventAgent, making direct tools and future capabilities duplicate orchestration. Risk: compatibility adapters preserve current APIs while deterministic-first resolution changes when EventAgent LLM is invoked; tests cover fallback boundaries.
zhenyu.hu преди 2 седмици
родител
ревизия
d157e9706d

+ 45 - 65
src/agent_lab/application/event_agent.py

@@ -1,10 +1,15 @@
-import asyncio
 import json
 from collections.abc import AsyncIterator, Iterable, Sequence
 from dataclasses import dataclass
 from typing import Any, Protocol
 
 from agent_lab.application.contracts import AgentParams
+from agent_lab.application.events import (
+    EventDefinition,
+    EventExecutionContext,
+    EventKernel,
+    EventRequest,
+)
 from agent_lab.application.tools import (
     ToolExecutionContext,
     ToolRegistry,
@@ -49,6 +54,12 @@ class EventAgent:
         self.registry = registry or build_default_tool_registry()
         self.chat_client = chat_client
         self.params = params or AgentParams()
+        self.kernel = EventKernel(
+            self.registry.event_registry,
+            argument_fallback=(
+                self._resolve_arguments_with_llm if chat_client is not None else None
+            ),
+        )
         self._raw_chunks_by_event: dict[str, list[dict[str, Any]]] = {}
 
     async def handle(
@@ -59,48 +70,16 @@ class EventAgent:
         extra_body: dict[str, Any] | None = None,
     ) -> ChatMessage:
         self._raw_chunks_by_event[event.id] = []
-        if event.name not in self.enabled_tools:
-            return self._tool_reply(
-                event,
-                {
-                    "tool": event.name,
-                    "error": "tool disabled",
-                },
-            )
-
-        try:
-            if self.chat_client is None:
-                payload = self.registry.handle(
-                    event,
-                    context=ToolExecutionContext(
-                        history=history,
-                        system_prompt=system_prompt,
-                        extra_body=extra_body or {},
-                    ),
-                )
-            else:
-                resolved_event = await self._resolve_event_with_llm(
-                    event,
-                    history,
-                    system_prompt,
-                    extra_body,
-                )
-                if resolved_event is None:
-                    payload = self.registry.handle(
-                        event,
-                        context=ToolExecutionContext(
-                            history=history,
-                            system_prompt=system_prompt,
-                            extra_body=extra_body or {},
-                        ),
-                    )
-                else:
-                    payload = self.registry.execute(resolved_event)
-        except Exception as exc:
-            payload = {
-                "tool": event.name,
-                "error": f"tool handler failed: {exc}",
-            }
+        result = await self.kernel.execute(
+            self.registry.event_request(event),
+            enabled_names=self.enabled_tools,
+            context=ToolExecutionContext(
+                history=history,
+                system_prompt=system_prompt,
+                extra_body=extra_body or {},
+            ),
+        )
+        payload = self.registry.tool_payload(result)
         return self._tool_reply(event, payload)
 
     async def handle_many(
@@ -112,17 +91,17 @@ class EventAgent:
     ) -> list[ChatMessage]:
         for event in events:
             self._raw_chunks_by_event[event.id] = []
-        return await asyncio.gather(
-            *[
-                self.handle(
+        replies: list[ChatMessage] = []
+        for event in events:
+            replies.append(
+                await self.handle(
                     event,
                     history=history,
                     system_prompt=system_prompt,
                     extra_body=extra_body,
                 )
-                for event in events
-            ]
-        )
+            )
+        return replies
 
     def raw_model_chunks(
         self,
@@ -137,23 +116,29 @@ class EventAgent:
             for event in events
         ]
 
-    async def _resolve_event_with_llm(
+    async def _resolve_arguments_with_llm(
         self,
-        event: ToolCallEvent,
-        history: Sequence[ChatMessage],
-        system_prompt: str,
-        extra_body: dict[str, Any] | None,
-    ) -> ToolCallEvent | None:
+        definition: EventDefinition,
+        request: EventRequest,
+        context: EventExecutionContext,
+    ) -> dict[str, Any] | None:
         assert self.chat_client is not None
-        tool = self.registry.tool_schema(event.name)
+        tool = self.registry.tool_schema(request.name)
         if tool is None:
-            return event
+            return None
 
-        messages = self._build_argument_messages(event, history, system_prompt)
+        event = self.registry._tool_event(request)
+        messages = self._build_argument_messages(
+            event,
+            context.history,
+            context.system_prompt,
+        )
         params = self.params.model_copy(
             update={
                 "extra_body": (
-                    extra_body if extra_body is not None else self.params.extra_body
+                    context.extra_body
+                    if context.extra_body is not None
+                    else self.params.extra_body
                 ),
             }
         )
@@ -173,12 +158,7 @@ class EventAgent:
 
             if item.kind == "provider_tool_call" and item.event is not None:
                 self._raw_chunks_by_event[event.id] = raw_chunks
-                return event.model_copy(
-                    update={
-                        "arguments": item.event.arguments,
-                        "raw_arguments": item.event.raw_arguments,
-                    }
-                )
+                return dict(item.event.arguments)
 
         self._raw_chunks_by_event[event.id] = raw_chunks
         return None

+ 32 - 0
src/agent_lab/application/events/__init__.py

@@ -0,0 +1,32 @@
+from agent_lab.application.events.kernel import ArgumentFallback, EventKernel
+from agent_lab.application.events.models import (
+    ConfirmationPolicy,
+    EventArgumentResolver,
+    EventDefinition,
+    EventExecutionContext,
+    EventHandler,
+    EventRequest,
+    EventResult,
+    EventSource,
+    EventStatus,
+    ResultPolicy,
+    RiskLevel,
+)
+from agent_lab.application.events.registry import EventRegistry
+
+__all__ = [
+    "ArgumentFallback",
+    "ConfirmationPolicy",
+    "EventArgumentResolver",
+    "EventDefinition",
+    "EventExecutionContext",
+    "EventHandler",
+    "EventKernel",
+    "EventRegistry",
+    "EventRequest",
+    "EventResult",
+    "EventSource",
+    "EventStatus",
+    "ResultPolicy",
+    "RiskLevel",
+]

+ 247 - 0
src/agent_lab/application/events/kernel.py

@@ -0,0 +1,247 @@
+from __future__ import annotations
+
+import inspect
+from collections.abc import Awaitable, Callable, Iterable
+from dataclasses import replace
+from typing import Any
+
+from agent_lab.application.events.models import (
+    EventDefinition,
+    EventExecutionContext,
+    EventRequest,
+    EventResult,
+    EventSource,
+    EventStatus,
+)
+from agent_lab.application.events.registry import EventRegistry
+
+
+ArgumentFallback = Callable[
+    [EventDefinition, EventRequest, EventExecutionContext],
+    Awaitable[dict[str, Any] | None],
+]
+
+
+class EventKernel:
+    def __init__(
+        self,
+        registry: EventRegistry,
+        argument_fallback: ArgumentFallback | None = None,
+    ) -> None:
+        self.registry = registry
+        self.argument_fallback = argument_fallback
+
+    async def execute(
+        self,
+        request: EventRequest,
+        *,
+        enabled_names: Iterable[str] | None = None,
+        context: EventExecutionContext | None = None,
+    ) -> EventResult:
+        definition, early_result = self._lookup(request, enabled_names)
+        if early_result is not None:
+            return early_result
+        assert definition is not None
+        resolved_context = context or EventExecutionContext()
+        arguments = self._deterministic_arguments(
+            definition,
+            request,
+            resolved_context,
+        )
+        validation_error = self._validation_error(definition, arguments)
+        used_fallback = False
+        if (
+            validation_error is not None
+            and request.source is not EventSource.PROVIDER_RESOLVED
+            and definition.fallback_allowed
+            and self.argument_fallback is not None
+        ):
+            used_fallback = True
+            fallback_request = replace(request, arguments=arguments)
+            fallback_arguments = await self.argument_fallback(
+                definition,
+                fallback_request,
+                resolved_context,
+            )
+            if fallback_arguments is not None:
+                arguments = fallback_arguments
+            validation_error = self._validation_error(definition, arguments)
+
+        if validation_error is not None:
+            return self._result(
+                definition,
+                request,
+                EventStatus.INVALID_ARGUMENTS,
+                arguments=arguments,
+                error=validation_error,
+                used_fallback=used_fallback,
+            )
+
+        resolved_request = replace(request, arguments=arguments)
+        try:
+            payload = definition.handler(resolved_request)
+            if inspect.isawaitable(payload):
+                payload = await payload
+        except Exception as exc:
+            return self._result(
+                definition,
+                request,
+                EventStatus.HANDLER_ERROR,
+                arguments=arguments,
+                error=f"event handler failed: {exc}",
+                used_fallback=used_fallback,
+            )
+        return self._result(
+            definition,
+            request,
+            EventStatus.SUCCESS,
+            arguments=arguments,
+            payload=dict(payload),
+            used_fallback=used_fallback,
+        )
+
+    def execute_sync(
+        self,
+        request: EventRequest,
+        *,
+        enabled_names: Iterable[str] | None = None,
+        context: EventExecutionContext | None = None,
+    ) -> EventResult:
+        definition, early_result = self._lookup(request, enabled_names)
+        if early_result is not None:
+            return early_result
+        assert definition is not None
+        arguments = self._deterministic_arguments(
+            definition,
+            request,
+            context or EventExecutionContext(),
+        )
+        validation_error = self._validation_error(definition, arguments)
+        if validation_error is not None:
+            return self._result(
+                definition,
+                request,
+                EventStatus.INVALID_ARGUMENTS,
+                arguments=arguments,
+                error=validation_error,
+            )
+        try:
+            payload = definition.handler(replace(request, arguments=arguments))
+            if inspect.isawaitable(payload):
+                raise TypeError("async event handlers require EventKernel.execute")
+        except Exception as exc:
+            return self._result(
+                definition,
+                request,
+                EventStatus.HANDLER_ERROR,
+                arguments=arguments,
+                error=f"event handler failed: {exc}",
+            )
+        return self._result(
+            definition,
+            request,
+            EventStatus.SUCCESS,
+            arguments=arguments,
+            payload=dict(payload),
+        )
+
+    def _lookup(
+        self,
+        request: EventRequest,
+        enabled_names: Iterable[str] | None,
+    ) -> tuple[EventDefinition | None, EventResult | None]:
+        definition = self.registry.definition(request.name)
+        if definition is None:
+            return None, EventResult(
+                event_id=request.id,
+                event_name=request.name,
+                status=EventStatus.UNKNOWN,
+                source=request.source,
+                arguments=dict(request.arguments),
+                error="unknown event",
+            )
+        if enabled_names is not None and request.name not in set(enabled_names):
+            return definition, self._result(
+                definition,
+                request,
+                EventStatus.DISABLED,
+                arguments=dict(request.arguments),
+                error="event disabled",
+            )
+        return definition, None
+
+    def _deterministic_arguments(
+        self,
+        definition: EventDefinition,
+        request: EventRequest,
+        context: EventExecutionContext,
+    ) -> dict[str, Any]:
+        if request.source is EventSource.PROVIDER_RESOLVED:
+            return dict(request.arguments)
+        if definition.resolver is None:
+            return dict(request.arguments)
+        return dict(definition.resolver(request, context))
+
+    def _validation_error(
+        self,
+        definition: EventDefinition,
+        arguments: dict[str, Any],
+    ) -> str | None:
+        required = definition.parameters.get("required", [])
+        missing = [name for name in required if name not in arguments]
+        if missing:
+            return f"missing required arguments: {', '.join(missing)}"
+
+        properties = definition.parameters.get("properties", {})
+        for name, value in arguments.items():
+            expected_name = properties.get(name, {}).get("type")
+            if expected_name is not None and not self._matches_json_type(
+                value,
+                expected_name,
+            ):
+                return f"invalid argument type for {name}: expected {expected_name}"
+        return None
+
+    def _matches_json_type(self, value: Any, expected_name: str) -> bool:
+        if expected_name == "number":
+            return not isinstance(value, bool) and isinstance(value, (int, float))
+        if expected_name == "integer":
+            return not isinstance(value, bool) and isinstance(value, int)
+        python_types = {
+            "string": str,
+            "boolean": bool,
+            "object": dict,
+            "array": list,
+        }
+        expected_type = python_types.get(expected_name)
+        return expected_type is None or isinstance(value, expected_type)
+
+    def _result(
+        self,
+        definition: EventDefinition,
+        request: EventRequest,
+        status: EventStatus,
+        *,
+        arguments: dict[str, Any],
+        payload: dict[str, Any] | None = None,
+        error: str | None = None,
+        used_fallback: bool = False,
+    ) -> EventResult:
+        return EventResult(
+            event_id=request.id,
+            event_name=request.name,
+            status=status,
+            source=request.source,
+            arguments=dict(arguments),
+            payload=payload or {},
+            error=error,
+            used_fallback=used_fallback,
+            result_policy=definition.result_policy,
+            confirmation_policy=definition.confirmation_policy,
+            risk_level=definition.risk_level,
+            idempotency_key_fields=definition.idempotency_key_fields,
+            concurrency_class=definition.concurrency_class,
+            conflict_keys=definition.conflict_keys,
+            timeout_seconds=definition.timeout_seconds,
+            terminal=definition.terminal,
+        )

+ 102 - 0
src/agent_lab/application/events/models.py

@@ -0,0 +1,102 @@
+from __future__ import annotations
+
+from collections.abc import Awaitable, Callable, Sequence
+from dataclasses import dataclass, field
+from enum import StrEnum
+from typing import Any
+
+from agent_lab.domain.messages import ChatMessage
+
+
+class EventSource(StrEnum):
+    TEXT_EVENT = "text_event"
+    PROVIDER_RESOLVED = "provider_resolved"
+
+
+class EventStatus(StrEnum):
+    SUCCESS = "success"
+    UNKNOWN = "unknown"
+    DISABLED = "disabled"
+    INVALID_ARGUMENTS = "invalid_arguments"
+    HANDLER_ERROR = "handler_error"
+
+
+class ResultPolicy(StrEnum):
+    TERMINATE = "terminate"
+    SILENT_SUCCESS = "silent_success"
+    TEMPLATE_FOLLOW_UP = "template_follow_up"
+    LLM_FOLLOW_UP = "llm_follow_up"
+
+
+class ConfirmationPolicy(StrEnum):
+    NONE = "none"
+    REQUIRED = "required"
+
+
+class RiskLevel(StrEnum):
+    LOW = "low"
+    MEDIUM = "medium"
+    HIGH = "high"
+
+
+@dataclass(frozen=True)
+class EventExecutionContext:
+    history: Sequence[ChatMessage] = ()
+    system_prompt: str = ""
+    extra_body: dict[str, Any] | None = None
+
+
+@dataclass(frozen=True)
+class EventRequest:
+    id: str
+    name: str
+    arguments: dict[str, Any] = field(default_factory=dict)
+    source: EventSource = EventSource.TEXT_EVENT
+    raw_arguments: str = "{}"
+
+
+EventHandlerResult = dict[str, Any] | Awaitable[dict[str, Any]]
+EventHandler = Callable[[EventRequest], EventHandlerResult]
+EventArgumentResolver = Callable[
+    [EventRequest, EventExecutionContext],
+    dict[str, Any],
+]
+
+
+@dataclass(frozen=True)
+class EventDefinition:
+    name: str
+    description: str
+    parameters: dict[str, Any]
+    handler: EventHandler
+    resolver: EventArgumentResolver | None = None
+    schema_version: str = "1"
+    fallback_allowed: bool = True
+    result_policy: ResultPolicy = ResultPolicy.LLM_FOLLOW_UP
+    confirmation_policy: ConfirmationPolicy = ConfirmationPolicy.NONE
+    risk_level: RiskLevel = RiskLevel.LOW
+    idempotency_key_fields: tuple[str, ...] = ()
+    concurrency_class: str | None = None
+    conflict_keys: tuple[str, ...] = ()
+    timeout_seconds: float | None = None
+    terminal: bool = False
+
+
+@dataclass(frozen=True)
+class EventResult:
+    event_id: str
+    event_name: str
+    status: EventStatus
+    source: EventSource
+    arguments: dict[str, Any] = field(default_factory=dict)
+    payload: dict[str, Any] = field(default_factory=dict)
+    error: str | None = None
+    used_fallback: bool = False
+    result_policy: ResultPolicy = ResultPolicy.LLM_FOLLOW_UP
+    confirmation_policy: ConfirmationPolicy = ConfirmationPolicy.NONE
+    risk_level: RiskLevel = RiskLevel.LOW
+    idempotency_key_fields: tuple[str, ...] = ()
+    concurrency_class: str | None = None
+    conflict_keys: tuple[str, ...] = ()
+    timeout_seconds: float | None = None
+    terminal: bool = False

+ 46 - 0
src/agent_lab/application/events/registry.py

@@ -0,0 +1,46 @@
+from __future__ import annotations
+
+from collections.abc import Iterable
+from copy import deepcopy
+
+from agent_lab.application.events.models import EventDefinition
+
+
+class EventRegistry:
+    def __init__(self, definitions: Iterable[EventDefinition] = ()) -> None:
+        self._definitions: dict[str, EventDefinition] = {}
+        for definition in definitions:
+            self.register(definition)
+
+    def register(self, definition: EventDefinition) -> None:
+        if definition.name in self._definitions:
+            raise ValueError(f"duplicate event definition: {definition.name}")
+        self._definitions[definition.name] = definition
+
+    def definition(self, name: str) -> EventDefinition | None:
+        return self._definitions.get(name)
+
+    def catalog(self, enabled_names: Iterable[str] | None = None) -> list[dict]:
+        enabled = set(enabled_names) if enabled_names is not None else None
+        return [
+            {
+                "name": definition.name,
+                "description": definition.description,
+                "parameters": deepcopy(definition.parameters),
+            }
+            for definition in self._definitions.values()
+            if enabled is None or definition.name in enabled
+        ]
+
+    def tool_schema(self, name: str) -> dict | None:
+        definition = self.definition(name)
+        if definition is None:
+            return None
+        return {
+            "type": "function",
+            "function": {
+                "name": definition.name,
+                "description": definition.description,
+                "parameters": deepcopy(definition.parameters),
+            },
+        }

+ 106 - 57
src/agent_lab/application/tools.py

@@ -3,15 +3,24 @@ from copy import deepcopy
 from dataclasses import dataclass
 from typing import Any
 
+from agent_lab.application.events import (
+    ConfirmationPolicy,
+    EventDefinition,
+    EventExecutionContext,
+    EventKernel,
+    EventRegistry,
+    EventRequest,
+    EventResult,
+    EventSource,
+    EventStatus,
+    ResultPolicy,
+    RiskLevel,
+)
 from agent_lab.domain.events import EVENT_BLOCK_END, EVENT_BLOCK_START, ToolCallEvent
 from agent_lab.domain.messages import ChatMessage
 
 
-@dataclass(frozen=True)
-class ToolExecutionContext:
-    history: Sequence[ChatMessage]
-    system_prompt: str = ""
-    extra_body: dict[str, Any] | None = None
+ToolExecutionContext = EventExecutionContext
 
 
 ToolHandler = Callable[[ToolCallEvent], dict[str, Any]]
@@ -28,21 +37,32 @@ class ToolDefinition:
     parameters: dict[str, Any]
     handler: ToolHandler
     argument_resolver: ToolArgumentResolver | None = None
+    schema_version: str = "1"
+    fallback_allowed: bool = True
+    result_policy: ResultPolicy = ResultPolicy.LLM_FOLLOW_UP
+    confirmation_policy: ConfirmationPolicy = ConfirmationPolicy.NONE
+    risk_level: RiskLevel = RiskLevel.LOW
+    idempotency_key_fields: tuple[str, ...] = ()
+    concurrency_class: str | None = None
+    conflict_keys: tuple[str, ...] = ()
+    timeout_seconds: float | None = None
+    terminal: bool = False
 
 
 class ToolRegistry:
     def __init__(self, definitions: Iterable[ToolDefinition]) -> None:
-        self._definitions = {definition.name: definition for definition in definitions}
+        self._definitions: dict[str, ToolDefinition] = {}
+        event_definitions: list[EventDefinition] = []
+        for definition in definitions:
+            if definition.name in self._definitions:
+                raise ValueError(f"duplicate event definition: {definition.name}")
+            self._definitions[definition.name] = definition
+            event_definitions.append(self._to_event_definition(definition))
+        self.event_registry = EventRegistry(event_definitions)
+        self.kernel = EventKernel(self.event_registry)
 
     def available_tools(self) -> list[dict[str, Any]]:
-        return [
-            {
-                "name": definition.name,
-                "description": definition.description,
-                "parameters": deepcopy(definition.parameters),
-            }
-            for definition in self._definitions.values()
-        ]
+        return self.event_registry.catalog()
 
     def chat_event_system_message(self, enabled_names: Iterable[str]) -> str:
         enabled = set(enabled_names)
@@ -71,61 +91,90 @@ class ToolRegistry:
         )
 
     def tool_schema(self, name: str) -> dict[str, Any] | None:
-        definition = self._definitions.get(name)
-        if definition is None:
-            return None
-        return {
-            "type": "function",
-            "function": {
-                "name": definition.name,
-                "description": definition.description,
-                "parameters": deepcopy(definition.parameters),
-            },
-        }
+        return self.event_registry.tool_schema(name)
 
     def handle(
         self,
         event: ToolCallEvent,
         context: ToolExecutionContext | None = None,
     ) -> dict[str, Any]:
-        definition = self._definitions.get(event.name)
-        if definition is None:
-            return {
-                "tool": event.name,
-                "error": "unknown tool",
-            }
-        resolved_context = context or ToolExecutionContext(history=())
-        resolved_event = event.model_copy(
-            update={
-                "arguments": self._resolve_arguments(
-                    definition,
-                    event,
-                    resolved_context,
-                ),
-            }
+        result = self.kernel.execute_sync(
+            self.event_request(event),
+            context=context or ToolExecutionContext(history=()),
         )
-        return definition.handler(resolved_event)
+        return self.tool_payload(result)
 
     def execute(self, event: ToolCallEvent) -> dict[str, Any]:
-        definition = self._definitions.get(event.name)
-        if definition is None:
-            return {
-                "tool": event.name,
-                "error": "unknown tool",
-            }
-        return definition.handler(event)
+        result = self.kernel.execute_sync(
+            self.event_request(event, source=EventSource.PROVIDER_RESOLVED)
+        )
+        return self.tool_payload(result)
 
-    def _resolve_arguments(
+    def event_request(
         self,
-        definition: ToolDefinition,
         event: ToolCallEvent,
-        context: ToolExecutionContext,
-    ) -> dict[str, Any]:
-        if definition.argument_resolver is not None:
-            return definition.argument_resolver(event, context)
-        if context.history:
-            return {}
-        return deepcopy(event.arguments)
+        source: EventSource = EventSource.TEXT_EVENT,
+    ) -> EventRequest:
+        return EventRequest(
+            id=event.id,
+            name=event.name,
+            arguments=deepcopy(event.arguments),
+            raw_arguments=event.raw_arguments,
+            source=source,
+        )
+
+    def tool_payload(self, result: EventResult) -> dict[str, Any]:
+        if result.status is EventStatus.SUCCESS:
+            return result.payload
+        error = result.error or result.status.value
+        if result.status is EventStatus.UNKNOWN:
+            error = "unknown tool"
+        elif result.status is EventStatus.DISABLED:
+            error = "tool disabled"
+        elif result.status is EventStatus.HANDLER_ERROR:
+            error = error.replace("event handler failed:", "tool handler failed:", 1)
+        return {"tool": result.event_name, "error": error}
+
+    def _to_event_definition(self, definition: ToolDefinition) -> EventDefinition:
+        def resolver(
+            request: EventRequest,
+            context: EventExecutionContext,
+        ) -> dict[str, Any]:
+            event = self._tool_event(request)
+            if definition.argument_resolver is not None:
+                return definition.argument_resolver(event, context)
+            if context.history:
+                return {}
+            return deepcopy(request.arguments)
+
+        def handler(request: EventRequest) -> dict[str, Any]:
+            return definition.handler(self._tool_event(request))
+
+        return EventDefinition(
+            name=definition.name,
+            description=definition.description,
+            parameters=deepcopy(definition.parameters),
+            handler=handler,
+            resolver=resolver,
+            schema_version=definition.schema_version,
+            fallback_allowed=definition.fallback_allowed,
+            result_policy=definition.result_policy,
+            confirmation_policy=definition.confirmation_policy,
+            risk_level=definition.risk_level,
+            idempotency_key_fields=definition.idempotency_key_fields,
+            concurrency_class=definition.concurrency_class,
+            conflict_keys=definition.conflict_keys,
+            timeout_seconds=definition.timeout_seconds,
+            terminal=definition.terminal,
+        )
+
+    def _tool_event(self, request: EventRequest) -> ToolCallEvent:
+        return ToolCallEvent(
+            id=request.id,
+            name=request.name,
+            arguments=deepcopy(request.arguments),
+            raw_arguments=request.raw_arguments,
+        )
 
 
 def build_default_tool_registry() -> ToolRegistry:

+ 2 - 14
tests/test_debug_runtime.py

@@ -817,14 +817,7 @@ async def test_runtime_audit_includes_model_params_prompts_results_and_usage():
     assert event_response["details"]["replies"][0]["role"] == "tool"
     assert '"tool": "handoff_note"' in event_response["details"]["replies"][0]["content"]
     assert event_response["details"]["raw_model_chunks"][0]["event_name"] == "handoff_note"
-    assert event_response["details"]["raw_model_chunks"][0]["chunks"][0]["choices"][0]["delta"] == {
-        "tool_calls": [
-            {
-                "index": 0,
-                "function": {"name": "handoff_note"},
-            }
-        ]
-    }
+    assert event_response["details"]["raw_model_chunks"][0]["chunks"] == []
 
 
 @pytest.mark.asyncio
@@ -896,12 +889,7 @@ async def test_runtime_event_agent_history_excludes_chat_agent_system_context():
         message["role"] == "system"
         for message in event_request["details"]["history"]
     )
-    assert client.event_agent_messages
-    assert [
-        (message.role, message.content)
-        for message in client.event_agent_messages[0]
-        if message.content in {"ChatAgent root prompt.", "ChatAgent pre system."}
-    ] == []
+    assert client.event_agent_messages == []
 
 
 @pytest.mark.asyncio

+ 101 - 44
tests/test_event_agent.py

@@ -95,7 +95,7 @@ class WrongSourceToolCallChatClient:
 
 
 @pytest.mark.asyncio
-async def test_event_agent_resolves_tool_arguments_with_llm_tool_call():
+async def test_event_agent_uses_deterministic_resolver_before_llm_fallback():
     chat_client = ToolCallingChatClient({"message": "LLM generated handoff"})
     agent = EventAgent(
         enabled_tools=["handoff_note"],
@@ -120,56 +120,20 @@ async def test_event_agent_resolves_tool_arguments_with_llm_tool_call():
     assert reply.name == "handoff_note"
     assert json.loads(reply.content) == {
         "tool": "handoff_note",
-        "message": "LLM generated handoff",
-    }
-    assert chat_client.calls[0]["tools"] == [
-        {
-            "type": "function",
-            "function": {
-                "name": "handoff_note",
-                "description": "Send a note to the event agent.",
-                "parameters": {
-                    "type": "object",
-                    "properties": {
-                        "message": {"type": "string"},
-                    },
-                    "required": ["message"],
-                },
-            },
-        }
-    ]
-    assert chat_client.calls[0]["tool_choice"] == {
-        "type": "function",
-        "function": {"name": "handoff_note"},
+        "message": "I need the event agent.",
     }
-    assert chat_client.calls[0]["params"].model == "event-model"
+    assert chat_client.calls == []
     assert agent.raw_model_chunks([event]) == [
         {
             "event_id": "call_1",
             "event_name": "handoff_note",
-            "chunks": [
-                {
-                    "choices": [
-                        {
-                            "delta": {
-                                "tool_calls": [
-                                    {
-                                        "index": 0,
-                                        "function": {"name": "handoff_note"},
-                                    }
-                                ]
-                            },
-                            "finish_reason": None,
-                        }
-                    ]
-                }
-            ],
+            "chunks": [],
         }
     ]
 
 
 @pytest.mark.asyncio
-async def test_event_agent_falls_back_to_context_arguments_when_llm_returns_no_tool_call():
+async def test_event_agent_deterministic_mock_resolver_skips_llm():
     chat_client = NoToolCallChatClient()
     agent = EventAgent(
         enabled_tools=["mock_search"],
@@ -193,7 +157,52 @@ async def test_event_agent_falls_back_to_context_arguments_when_llm_returns_no_t
     assert payload["tool"] == "mock_search"
     assert payload["query"] == "Need a search for latency docs"
     assert "event agent did not return arguments" not in reply.content
-    assert chat_client.calls[0]["tools"][0]["function"]["name"] == "mock_search"
+    assert chat_client.calls == []
+
+
+@pytest.mark.asyncio
+async def test_event_agent_calls_llm_once_for_incomplete_deterministic_arguments():
+    chat_client = ToolCallingChatClient({"message": "LLM generated handoff"})
+    registry = ToolRegistry(
+        [
+            ToolDefinition(
+                name="ambiguous_handoff",
+                description="Resolve an ambiguous handoff.",
+                parameters={
+                    "type": "object",
+                    "properties": {"message": {"type": "string"}},
+                    "required": ["message"],
+                },
+                handler=lambda event: {
+                    "tool": event.name,
+                    "message": event.arguments["message"],
+                },
+                argument_resolver=lambda event, context: {},
+            )
+        ]
+    )
+    event = ToolCallEvent(
+        id="call_1",
+        name="ambiguous_handoff",
+        arguments={},
+        raw_arguments="{}",
+    )
+
+    reply = await EventAgent(
+        enabled_tools=["ambiguous_handoff"],
+        registry=registry,
+        chat_client=chat_client,
+    ).handle(event, history=[ChatMessage(role="user", content="ambiguous")])
+
+    assert json.loads(reply.content) == {
+        "tool": "ambiguous_handoff",
+        "message": "LLM generated handoff",
+    }
+    assert len(chat_client.calls) == 1
+    assert chat_client.calls[0]["tool_choice"] == {
+        "type": "function",
+        "function": {"name": "ambiguous_handoff"},
+    }
 
 
 @pytest.mark.asyncio
@@ -226,8 +235,27 @@ async def test_event_agent_ignores_non_provider_tool_call_sources(item: StreamIt
         arguments={},
         raw_arguments="{}",
     )
+    registry = ToolRegistry(
+        [
+            ToolDefinition(
+                name="mock_search",
+                description="Search mock external knowledge for the current turn.",
+                parameters={
+                    "type": "object",
+                    "properties": {"query": {"type": "string"}},
+                    "required": ["query"],
+                },
+                handler=lambda resolved: {
+                    "tool": resolved.name,
+                    "query": resolved.arguments["query"],
+                },
+                argument_resolver=lambda event, context: {},
+            )
+        ]
+    )
     agent = EventAgent(
         enabled_tools=["mock_search"],
+        registry=registry,
         chat_client=WrongSourceToolCallChatClient(item),
     )
 
@@ -238,7 +266,10 @@ async def test_event_agent_ignores_non_provider_tool_call_sources(item: StreamIt
 
     payload = json.loads(reply.content)
     assert payload["tool"] == "mock_search"
-    assert payload["query"] == "fallback query"
+    assert payload == {
+        "tool": "mock_search",
+        "error": "missing required arguments: query",
+    }
 
 
 @pytest.mark.asyncio
@@ -332,7 +363,14 @@ async def test_event_agent_llm_receives_history_and_agent_config_context():
             ToolDefinition(
                 name="handoff_note",
                 description="Send a note to the event agent.",
-                parameters={"type": "object"},
+                parameters={
+                    "type": "object",
+                    "properties": {
+                        "message": {"type": "string"},
+                        "thinking": {"type": "string"},
+                    },
+                    "required": ["message", "thinking"],
+                },
                 handler=lambda event: {
                     "tool": event.name,
                     "message": event.arguments["message"],
@@ -374,6 +412,24 @@ async def test_event_agent_llm_receives_history_and_agent_config_context():
 @pytest.mark.asyncio
 async def test_event_agent_projects_complete_tool_round_to_visible_history():
     chat_client = ToolCallingChatClient({"message": "projected history"})
+    registry = ToolRegistry(
+        [
+            ToolDefinition(
+                name="handoff_note",
+                description="Send a note to the event agent.",
+                parameters={
+                    "type": "object",
+                    "properties": {"message": {"type": "string"}},
+                    "required": ["message"],
+                },
+                handler=lambda resolved: {
+                    "tool": resolved.name,
+                    "message": resolved.arguments["message"],
+                },
+                argument_resolver=lambda event, context: {},
+            )
+        ]
+    )
     event = ToolCallEvent(
         id="call_current",
         name="handoff_note",
@@ -389,6 +445,7 @@ async def test_event_agent_projects_complete_tool_round_to_visible_history():
 
     await EventAgent(
         enabled_tools=["handoff_note"],
+        registry=registry,
         chat_client=chat_client,
     ).handle(
         event,

+ 401 - 0
tests/test_event_kernel.py

@@ -0,0 +1,401 @@
+from __future__ import annotations
+
+from typing import Any
+
+import pytest
+
+from agent_lab.application.events import (
+    ConfirmationPolicy,
+    EventDefinition,
+    EventExecutionContext,
+    EventKernel,
+    EventRegistry,
+    EventRequest,
+    EventSource,
+    EventStatus,
+    ResultPolicy,
+    RiskLevel,
+)
+from agent_lab.application.tools import (
+    ToolDefinition,
+    ToolExecutionContext,
+    ToolRegistry,
+)
+from agent_lab.domain.events import ToolCallEvent
+from agent_lab.domain.messages import ChatMessage
+
+
+def _definition(
+    name: str = "example.lookup",
+    **overrides: Any,
+) -> EventDefinition:
+    values: dict[str, Any] = {
+        "name": name,
+        "description": "Look up an example value.",
+        "parameters": {
+            "type": "object",
+            "properties": {"query": {"type": "string"}},
+            "required": ["query"],
+        },
+        "handler": lambda request: {
+            "event": request.name,
+            "query": request.arguments["query"],
+        },
+    }
+    values.update(overrides)
+    return EventDefinition(**values)
+
+
+def test_registry_registers_flat_definitions_and_filters_enabled_catalog():
+    registry = EventRegistry(
+        [_definition("example.lookup"), _definition("device.inspect")]
+    )
+
+    assert [item["name"] for item in registry.catalog()] == [
+        "example.lookup",
+        "device.inspect",
+    ]
+    assert registry.catalog(["device.inspect"]) == [
+        {
+            "name": "device.inspect",
+            "description": "Look up an example value.",
+            "parameters": {
+                "type": "object",
+                "properties": {"query": {"type": "string"}},
+                "required": ["query"],
+            },
+        }
+    ]
+    assert registry.tool_schema("example.lookup")["function"]["name"] == (
+        "example.lookup"
+    )
+    assert registry.tool_schema("missing") is None
+
+
+def test_registry_rejects_duplicate_definition_names():
+    with pytest.raises(ValueError, match="duplicate event definition: example.lookup"):
+        EventRegistry([_definition(), _definition()])
+
+
+def test_tool_registry_public_api_remains_compatible():
+    registry = ToolRegistry(
+        [
+            ToolDefinition(
+                name="compat.lookup",
+                description="Look up compatibility data.",
+                parameters={
+                    "type": "object",
+                    "properties": {"query": {"type": "string"}},
+                    "required": ["query"],
+                },
+                handler=lambda event: {
+                    "tool": event.name,
+                    "query": event.arguments["query"],
+                },
+                argument_resolver=lambda event, context: {
+                    "query": context.history[-1].content
+                },
+            )
+        ]
+    )
+    event = ToolCallEvent(
+        id="call-1",
+        name="compat.lookup",
+        arguments={"query": "provider"},
+        raw_arguments='{"query":"provider"}',
+    )
+
+    assert registry.available_tools() == [
+        {
+            "name": "compat.lookup",
+            "description": "Look up compatibility data.",
+            "parameters": {
+                "type": "object",
+                "properties": {"query": {"type": "string"}},
+                "required": ["query"],
+            },
+        }
+    ]
+    assert "- compat.lookup: Look up compatibility data." in (
+        registry.chat_event_system_message(["compat.lookup"])
+    )
+    assert registry.tool_schema("compat.lookup")["function"]["name"] == (
+        "compat.lookup"
+    )
+    assert registry.handle(
+        event,
+        ToolExecutionContext(history=[ChatMessage(role="user", content="history")]),
+    ) == {"tool": "compat.lookup", "query": "history"}
+    assert registry.execute(event) == {"tool": "compat.lookup", "query": "provider"}
+
+
+@pytest.mark.asyncio
+async def test_kernel_executes_complete_deterministic_arguments_without_fallback():
+    fallback_calls: list[str] = []
+
+    async def fallback(*args: Any) -> dict[str, Any]:
+        fallback_calls.append("called")
+        return {"query": "fallback"}
+
+    registry = EventRegistry(
+        [_definition(resolver=lambda request, context: {"query": "deterministic"})]
+    )
+
+    result = await EventKernel(registry, argument_fallback=fallback).execute(
+        EventRequest(id="event-1", name="example.lookup"),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.SUCCESS
+    assert result.arguments == {"query": "deterministic"}
+    assert result.payload == {"event": "example.lookup", "query": "deterministic"}
+    assert result.used_fallback is False
+    assert fallback_calls == []
+
+
+@pytest.mark.asyncio
+async def test_kernel_calls_fallback_once_when_required_arguments_are_incomplete():
+    fallback_calls: list[dict[str, Any]] = []
+
+    async def fallback(
+        definition: EventDefinition,
+        request: EventRequest,
+        context: EventExecutionContext,
+    ) -> dict[str, Any]:
+        fallback_calls.append(dict(request.arguments))
+        return {"query": "resolved once"}
+
+    registry = EventRegistry([_definition(resolver=lambda request, context: {})])
+
+    result = await EventKernel(registry, argument_fallback=fallback).execute(
+        EventRequest(id="event-1", name="example.lookup"),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.SUCCESS
+    assert result.arguments == {"query": "resolved once"}
+    assert result.used_fallback is True
+    assert fallback_calls == [{}]
+
+
+@pytest.mark.asyncio
+async def test_kernel_does_not_fallback_when_definition_disallows_it():
+    fallback_calls = 0
+
+    async def fallback(*args: Any) -> dict[str, Any]:
+        nonlocal fallback_calls
+        fallback_calls += 1
+        return {"query": "not allowed"}
+
+    registry = EventRegistry(
+        [
+            _definition(
+                resolver=lambda request, context: {},
+                fallback_allowed=False,
+            )
+        ]
+    )
+
+    result = await EventKernel(registry, argument_fallback=fallback).execute(
+        EventRequest(id="event-1", name="example.lookup"),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.INVALID_ARGUMENTS
+    assert result.error == "missing required arguments: query"
+    assert result.used_fallback is False
+    assert fallback_calls == 0
+
+
+@pytest.mark.asyncio
+async def test_provider_resolved_arguments_are_not_rewritten_or_fallen_back():
+    resolver_calls = 0
+    fallback_calls = 0
+
+    def resolver(*args: Any) -> dict[str, Any]:
+        nonlocal resolver_calls
+        resolver_calls += 1
+        return {"query": "rewritten"}
+
+    async def fallback(*args: Any) -> dict[str, Any]:
+        nonlocal fallback_calls
+        fallback_calls += 1
+        return {"query": "fallback"}
+
+    registry = EventRegistry([_definition(resolver=resolver)])
+
+    result = await EventKernel(registry, argument_fallback=fallback).execute(
+        EventRequest(
+            id="event-1",
+            name="example.lookup",
+            arguments={"query": "provider value"},
+            source=EventSource.PROVIDER_RESOLVED,
+        ),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.SUCCESS
+    assert result.arguments == {"query": "provider value"}
+    assert resolver_calls == 0
+    assert fallback_calls == 0
+
+
+@pytest.mark.asyncio
+async def test_provider_resolved_missing_arguments_return_invalid_without_fallback():
+    fallback_calls = 0
+
+    async def fallback(*args: Any) -> dict[str, Any]:
+        nonlocal fallback_calls
+        fallback_calls += 1
+        return {"query": "fallback"}
+
+    result = await EventKernel(
+        EventRegistry([_definition()]), argument_fallback=fallback
+    ).execute(
+        EventRequest(
+            id="event-1",
+            name="example.lookup",
+            arguments={},
+            source=EventSource.PROVIDER_RESOLVED,
+        ),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.INVALID_ARGUMENTS
+    assert result.error == "missing required arguments: query"
+    assert fallback_calls == 0
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("event_request", "enabled_names", "expected_status", "expected_error"),
+    [
+        (
+            EventRequest(id="event-1", name="missing"),
+            ["missing"],
+            EventStatus.UNKNOWN,
+            "unknown event",
+        ),
+        (
+            EventRequest(id="event-1", name="example.lookup"),
+            [],
+            EventStatus.DISABLED,
+            "event disabled",
+        ),
+        (
+            EventRequest(
+                id="event-1",
+                name="example.lookup",
+                arguments={"query": 42},
+                source=EventSource.PROVIDER_RESOLVED,
+            ),
+            ["example.lookup"],
+            EventStatus.INVALID_ARGUMENTS,
+            "invalid argument type for query: expected string",
+        ),
+    ],
+)
+async def test_kernel_normalizes_lookup_and_validation_failures(
+    event_request: EventRequest,
+    enabled_names: list[str],
+    expected_status: EventStatus,
+    expected_error: str,
+):
+    result = await EventKernel(EventRegistry([_definition()])).execute(
+        event_request,
+        enabled_names=enabled_names,
+    )
+
+    assert result.status is expected_status
+    assert result.error == expected_error
+
+
+@pytest.mark.asyncio
+async def test_kernel_normalizes_handler_exceptions():
+    def fail(request: EventRequest) -> dict[str, Any]:
+        raise RuntimeError("boom")
+
+    result = await EventKernel(
+        EventRegistry([_definition(handler=fail)])
+    ).execute(
+        EventRequest(
+            id="event-1",
+            name="example.lookup",
+            arguments={"query": "value"},
+            source=EventSource.PROVIDER_RESOLVED,
+        ),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.HANDLER_ERROR
+    assert result.error == "event handler failed: boom"
+
+
+@pytest.mark.asyncio
+async def test_kernel_rejects_boolean_for_json_number_arguments():
+    definition = _definition(
+        parameters={
+            "type": "object",
+            "properties": {"query": {"type": "number"}},
+            "required": ["query"],
+        }
+    )
+
+    result = await EventKernel(EventRegistry([definition])).execute(
+        EventRequest(
+            id="event-1",
+            name="example.lookup",
+            arguments={"query": True},
+            source=EventSource.PROVIDER_RESOLVED,
+        ),
+        enabled_names=["example.lookup"],
+    )
+
+    assert result.status is EventStatus.INVALID_ARGUMENTS
+    assert result.error == "invalid argument type for query: expected number"
+
+
+@pytest.mark.asyncio
+async def test_definition_metadata_survives_registration_and_result_creation():
+    definition = _definition(
+        result_policy=ResultPolicy.TEMPLATE_FOLLOW_UP,
+        confirmation_policy=ConfirmationPolicy.REQUIRED,
+        risk_level=RiskLevel.HIGH,
+        idempotency_key_fields=("session_id", "event_id"),
+        concurrency_class="device-write",
+        conflict_keys=("device",),
+        timeout_seconds=1.5,
+        terminal=True,
+    )
+    registry = EventRegistry([definition])
+
+    result = await EventKernel(registry).execute(
+        EventRequest(
+            id="event-1",
+            name="example.lookup",
+            arguments={"query": "value"},
+            source=EventSource.PROVIDER_RESOLVED,
+        ),
+        enabled_names=["example.lookup"],
+    )
+
+    registered = registry.definition("example.lookup")
+    assert registered is definition
+    assert result.result_policy is ResultPolicy.TEMPLATE_FOLLOW_UP
+    assert result.confirmation_policy is ConfirmationPolicy.REQUIRED
+    assert result.risk_level is RiskLevel.HIGH
+    assert result.idempotency_key_fields == ("session_id", "event_id")
+    assert result.concurrency_class == "device-write"
+    assert result.conflict_keys == ("device",)
+    assert result.timeout_seconds == 1.5
+    assert result.terminal is True
+
+
+def test_kernel_source_has_no_builtin_event_name_branches():
+    from pathlib import Path
+
+    source = Path("src/agent_lab/application/events/kernel.py").read_text()
+
+    assert "handoff_note" not in source
+    assert "mock_search" not in source
+    assert "mock_ticket" not in source