Kaynağa Gözat

feat: add built-in event plugins

Problem:
- The default event registry only exposed legacy mock tools and lacked reusable built-in event boundaries.

Risk:
- Deterministic resolvers intentionally reject or defer ambiguous input, while default adapters remain in-memory and network-free.
zhenyu.hu 2 hafta önce
ebeveyn
işleme
083afdea8b

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

@@ -1,6 +1,18 @@
 from agent_lab.application.events.kernel import ArgumentFallback, EventKernel
+from agent_lab.application.events.builtin_plugins import (
+    CalendarSchedulePort,
+    DeviceVolumePort,
+    InMemoryCalendarScheduleAdapter,
+    InMemoryDeviceVolumeAdapter,
+    InMemorySessionTerminationAdapter,
+    InMemoryWebSearchAdapter,
+    SessionTerminationPort,
+    WebSearchPort,
+    build_builtin_event_definitions,
+)
 from agent_lab.application.events.models import (
     ConfirmationPolicy,
+    EventArgumentResolution,
     EventArgumentResolver,
     EventDefinition,
     EventExecutionContext,
@@ -17,7 +29,9 @@ from agent_lab.application.events.registry import EventRegistry
 
 __all__ = [
     "ArgumentFallback",
+    "CalendarSchedulePort",
     "ConfirmationPolicy",
+    "EventArgumentResolution",
     "EventArgumentResolver",
     "EventDefinition",
     "EventExecutionContext",
@@ -28,7 +42,15 @@ __all__ = [
     "EventResult",
     "EventSource",
     "EventStatus",
+    "DeviceVolumePort",
+    "InMemoryCalendarScheduleAdapter",
+    "InMemoryDeviceVolumeAdapter",
+    "InMemorySessionTerminationAdapter",
+    "InMemoryWebSearchAdapter",
     "ResolvedEventArguments",
     "ResultPolicy",
     "RiskLevel",
+    "SessionTerminationPort",
+    "WebSearchPort",
+    "build_builtin_event_definitions",
 ]

+ 560 - 0
src/agent_lab/application/events/builtin_plugins.py

@@ -0,0 +1,560 @@
+from __future__ import annotations
+
+import re
+from collections.abc import Awaitable, Callable
+from copy import deepcopy
+from datetime import datetime, timezone as datetime_timezone
+from typing import Any, Protocol
+
+from agent_lab.application.events.models import (
+    EventArgumentResolution,
+    EventDefinition,
+    EventExecutionContext,
+    EventRequest,
+    ResultPolicy,
+    RiskLevel,
+)
+
+
+_RFC3339_PATTERN = (
+    r"\d{4}-(?:0[1-9]|1[0-2])-(?:0[1-9]|[12]\d|3[01])"
+    r"T(?:[01]\d|2[0-3]):[0-5]\d:[0-5]\d"
+    r"(?:\.\d+)?(?:Z|[+-](?:0\d|1[0-4]):[0-5]\d)"
+)
+_RFC3339_RE = re.compile(_RFC3339_PATTERN)
+_TIMEZONE_PATTERN = r"[A-Za-z][A-Za-z0-9_+-]*(?:/[A-Za-z0-9_+-]+)+"
+
+
+PortResult = dict[str, Any] | Awaitable[dict[str, Any]]
+Clock = Callable[[], datetime | str]
+
+
+class SessionTerminationPort(Protocol):
+    def terminate(
+        self,
+        event_id: str,
+        *,
+        reason: str | None = None,
+    ) -> PortResult:
+        ...
+
+
+class DeviceVolumePort(Protocol):
+    def adjust(
+        self,
+        event_id: str,
+        *,
+        mode: str,
+        value: int | None = None,
+        delta: int | None = None,
+    ) -> PortResult:
+        ...
+
+
+class CalendarSchedulePort(Protocol):
+    def create(
+        self,
+        event_id: str,
+        *,
+        title: str,
+        start_at: str,
+        timezone: str,
+        recurrence: str | None = None,
+        reminder_minutes: int | None = None,
+    ) -> PortResult:
+        ...
+
+
+class WebSearchPort(Protocol):
+    def search(
+        self,
+        event_id: str,
+        *,
+        query: str,
+        max_results: int = 3,
+    ) -> PortResult:
+        ...
+
+
+class InMemorySessionTerminationAdapter:
+    def __init__(self) -> None:
+        self._results: dict[str, dict[str, Any]] = {}
+
+    def terminate(
+        self,
+        event_id: str,
+        *,
+        reason: str | None = None,
+    ) -> dict[str, Any]:
+        if event_id not in self._results:
+            self._results[event_id] = {
+                "tool": "session.terminate",
+                "status": "terminated",
+                "event_id": event_id,
+                "reason": reason,
+            }
+        return deepcopy(self._results[event_id])
+
+
+class InMemoryDeviceVolumeAdapter:
+    def __init__(self) -> None:
+        self._results: dict[str, dict[str, Any]] = {}
+
+    def adjust(
+        self,
+        event_id: str,
+        *,
+        mode: str,
+        value: int | None = None,
+        delta: int | None = None,
+    ) -> dict[str, Any]:
+        if event_id not in self._results:
+            payload: dict[str, Any] = {
+                "tool": "device.volume.adjust",
+                "status": "applied",
+                "event_id": event_id,
+                "mode": mode,
+            }
+            if value is not None:
+                payload["value"] = value
+            if delta is not None:
+                payload["delta"] = delta
+            self._results[event_id] = payload
+        return deepcopy(self._results[event_id])
+
+
+class InMemoryCalendarScheduleAdapter:
+    def __init__(self) -> None:
+        self._results: dict[str, dict[str, Any]] = {}
+
+    def create(
+        self,
+        event_id: str,
+        *,
+        title: str,
+        start_at: str,
+        timezone: str,
+        recurrence: str | None = None,
+        reminder_minutes: int | None = None,
+    ) -> dict[str, Any]:
+        if event_id not in self._results:
+            schedule: dict[str, Any] = {
+                "title": title,
+                "start_at": start_at,
+                "timezone": timezone,
+            }
+            if recurrence is not None:
+                schedule["recurrence"] = recurrence
+            if reminder_minutes is not None:
+                schedule["reminder_minutes"] = reminder_minutes
+            self._results[event_id] = {
+                "tool": "calendar.schedule.create",
+                "status": "created",
+                "event_id": event_id,
+                "schedule": schedule,
+            }
+        return deepcopy(self._results[event_id])
+
+
+class InMemoryWebSearchAdapter:
+    def __init__(self, clock: Clock | None = None) -> None:
+        self._clock = clock or _fixed_clock
+
+    def search(
+        self,
+        event_id: str,
+        *,
+        query: str,
+        max_results: int = 3,
+    ) -> dict[str, Any]:
+        del event_id
+        sources = [
+            {
+                "title": f"Mock source {index}",
+                "url": f"https://example.invalid/search/{index}",
+                "snippet": f"Deterministic mock result {index} for: {query}",
+            }
+            for index in range(1, max_results + 1)
+        ]
+        return {
+            "tool": "knowledge.web.search",
+            "query": query,
+            "sources": sources,
+            "retrieved_at": _format_retrieved_at(self._clock()),
+        }
+
+
+def build_builtin_event_definitions(
+    *,
+    session_termination_port: SessionTerminationPort | None = None,
+    device_volume_port: DeviceVolumePort | None = None,
+    calendar_schedule_port: CalendarSchedulePort | None = None,
+    web_search_port: WebSearchPort | None = None,
+    clock: Clock | None = None,
+) -> list[EventDefinition]:
+    session_port = session_termination_port or InMemorySessionTerminationAdapter()
+    volume_port = device_volume_port or InMemoryDeviceVolumeAdapter()
+    calendar_port = calendar_schedule_port or InMemoryCalendarScheduleAdapter()
+    search_port = web_search_port or InMemoryWebSearchAdapter(clock=clock)
+    return [
+        _session_terminate_definition(session_port),
+        _device_volume_definition(volume_port),
+        _calendar_schedule_definition(calendar_port),
+        _knowledge_search_definition(search_port),
+    ]
+
+
+def _session_terminate_definition(port: SessionTerminationPort) -> EventDefinition:
+    def handler(request: EventRequest) -> PortResult:
+        return port.terminate(
+            request.id,
+            reason=request.arguments.get("reason"),
+        )
+
+    return EventDefinition(
+        name="session.terminate",
+        description="Terminate the active dialogue session with an optional reason.",
+        parameters={
+            "type": "object",
+            "properties": {"reason": {"type": "string"}},
+            "additionalProperties": False,
+        },
+        handler=handler,
+        resolver=_resolve_session_terminate,
+        fallback_allowed=False,
+        result_policy=ResultPolicy.TERMINATE,
+        risk_level=RiskLevel.HIGH,
+        idempotency_key_fields=("event_id",),
+        concurrency_class="session-lifecycle",
+        conflict_keys=("session",),
+        timeout_seconds=5.0,
+        terminal=True,
+    )
+
+
+def _device_volume_definition(port: DeviceVolumePort) -> EventDefinition:
+    def handler(request: EventRequest) -> PortResult:
+        return port.adjust(
+            request.id,
+            mode=request.arguments["mode"],
+            value=request.arguments.get("value"),
+            delta=request.arguments.get("delta"),
+        )
+
+    mode_condition = lambda mode: {
+        "properties": {"mode": {"const": mode}},
+        "required": ["mode"],
+    }
+    return EventDefinition(
+        name="device.volume.adjust",
+        description="Adjust device volume absolutely, relatively, or by mute state.",
+        parameters={
+            "type": "object",
+            "properties": {
+                "mode": {
+                    "type": "string",
+                    "enum": ["absolute", "relative", "mute", "unmute"],
+                },
+                "value": {"type": "integer", "minimum": 0, "maximum": 100},
+                "delta": {
+                    "type": "integer",
+                    "minimum": -100,
+                    "maximum": 100,
+                    "not": {"const": 0},
+                },
+            },
+            "required": ["mode"],
+            "additionalProperties": False,
+            "allOf": [
+                {
+                    "if": mode_condition("absolute"),
+                    "then": {
+                        "required": ["value"],
+                        "properties": {"delta": False},
+                    },
+                },
+                {
+                    "if": mode_condition("relative"),
+                    "then": {
+                        "required": ["delta"],
+                        "properties": {"value": False},
+                    },
+                },
+                {
+                    "if": mode_condition("mute"),
+                    "then": {"properties": {"value": False, "delta": False}},
+                },
+                {
+                    "if": mode_condition("unmute"),
+                    "then": {"properties": {"value": False, "delta": False}},
+                },
+            ],
+        },
+        handler=handler,
+        resolver=_resolve_device_volume,
+        fallback_allowed=True,
+        result_policy=ResultPolicy.SILENT_SUCCESS,
+        risk_level=RiskLevel.MEDIUM,
+        idempotency_key_fields=("event_id",),
+        concurrency_class="device-volume",
+        conflict_keys=("device.volume",),
+        timeout_seconds=5.0,
+    )
+
+
+def _calendar_schedule_definition(port: CalendarSchedulePort) -> EventDefinition:
+    def handler(request: EventRequest) -> PortResult:
+        return port.create(
+            request.id,
+            title=request.arguments["title"],
+            start_at=request.arguments["start_at"],
+            timezone=request.arguments["timezone"],
+            recurrence=request.arguments.get("recurrence"),
+            reminder_minutes=request.arguments.get("reminder_minutes"),
+        )
+
+    return EventDefinition(
+        name="calendar.schedule.create",
+        description="Create a schedule entry from explicit RFC3339 date-time details.",
+        parameters={
+            "type": "object",
+            "properties": {
+                "title": {"type": "string", "minLength": 1},
+                "start_at": {
+                    "type": "string",
+                    "pattern": f"^{_RFC3339_PATTERN}$",
+                },
+                "timezone": {
+                    "type": "string",
+                    "pattern": f"^{_TIMEZONE_PATTERN}$",
+                },
+                "recurrence": {"type": "string", "minLength": 1},
+                "reminder_minutes": {"type": "integer", "minimum": 0},
+            },
+            "required": ["title", "start_at", "timezone"],
+            "additionalProperties": False,
+        },
+        handler=handler,
+        resolver=_resolve_calendar_schedule,
+        fallback_allowed=True,
+        result_policy=ResultPolicy.TEMPLATE_FOLLOW_UP,
+        risk_level=RiskLevel.MEDIUM,
+        idempotency_key_fields=("event_id",),
+        concurrency_class="schedule-write",
+        conflict_keys=("calendar.schedule",),
+        timeout_seconds=10.0,
+    )
+
+
+def _knowledge_search_definition(port: WebSearchPort) -> EventDefinition:
+    def handler(request: EventRequest) -> PortResult:
+        return port.search(
+            request.id,
+            query=request.arguments["query"],
+            max_results=request.arguments.get("max_results", 3),
+        )
+
+    return EventDefinition(
+        name="knowledge.web.search",
+        description="Search deterministic mock web sources for a focused query.",
+        parameters={
+            "type": "object",
+            "properties": {
+                "query": {"type": "string", "minLength": 1},
+                "max_results": {
+                    "type": "integer",
+                    "minimum": 1,
+                    "maximum": 5,
+                },
+            },
+            "required": ["query"],
+            "additionalProperties": False,
+        },
+        handler=handler,
+        resolver=_resolve_knowledge_search,
+        fallback_allowed=False,
+        result_policy=ResultPolicy.LLM_FOLLOW_UP,
+        risk_level=RiskLevel.LOW,
+        concurrency_class="read-only",
+        timeout_seconds=10.0,
+    )
+
+
+def _resolve_session_terminate(
+    request: EventRequest,
+    context: EventExecutionContext,
+) -> dict[str, Any]:
+    del context
+    return deepcopy(request.arguments)
+
+
+def _resolve_device_volume(
+    request: EventRequest,
+    context: EventExecutionContext,
+) -> dict[str, Any] | EventArgumentResolution:
+    if request.arguments:
+        arguments = deepcopy(request.arguments)
+        mode = arguments.get("mode")
+        if isinstance(mode, str):
+            arguments["mode"] = mode.strip().lower()
+            mode = arguments["mode"]
+        complete = mode in {"mute", "unmute"} or (
+            mode == "absolute" and "value" in arguments
+        ) or (mode == "relative" and "delta" in arguments)
+        if mode not in {"absolute", "relative", "mute", "unmute"}:
+            complete = True
+        return (
+            arguments
+            if complete
+            else EventArgumentResolution(arguments=arguments, complete=False)
+        )
+
+    content = _latest_user_content(context).lower()
+    unmute = bool(re.search(r"\bunmute\b|取消静音|解除静音|恢复声音", content))
+    mute = bool(re.search(r"\bmute\b|(?<!取消)(?<!解除)静音", content))
+    if mute and unmute:
+        return EventArgumentResolution(arguments={}, complete=False)
+    if mute or unmute:
+        arguments: dict[str, Any] = {"mode": "unmute" if unmute else "mute"}
+        number = re.search(r"\d+", content)
+        if number is not None:
+            arguments["value"] = int(number.group())
+        return arguments
+
+    absolute_matches = [
+        re.search(
+            r"音量.{0,6}?(?:调到|调至|设置为|设为|到)\s*(\d{1,3})",
+            content,
+        ),
+        re.search(
+            r"(?:set|change)\s+(?:the\s+)?volume\s+(?:to|at)\s*(\d{1,3})",
+            content,
+        ),
+        re.search(r"\bvolume\s*(?:to|at|=)\s*(\d{1,3})", content),
+    ]
+    positive_matches = [
+        re.search(r"音量.{0,4}?(?:增加|调高|提高|加)\s*(\d{1,3})", content),
+        re.search(
+            r"(?:increase|raise|turn\s+up)\s+(?:the\s+)?volume"
+            r"(?:\s+by)?\s*(\d{1,3})",
+            content,
+        ),
+    ]
+    negative_matches = [
+        re.search(r"音量.{0,4}?(?:降低|调低|减少|减)\s*(\d{1,3})", content),
+        re.search(
+            r"(?:decrease|lower|turn\s+down)\s+(?:the\s+)?volume"
+            r"(?:\s+by)?\s*(\d{1,3})",
+            content,
+        ),
+    ]
+    signed = re.search(r"(?:音量|\bvolume\b)\s*([+-])\s*(\d{1,3})", content)
+    absolute = next((match for match in absolute_matches if match), None)
+    positive = next((match for match in positive_matches if match), None)
+    negative = next((match for match in negative_matches if match), None)
+    matched_kinds = sum(item is not None for item in (absolute, positive, negative, signed))
+    if matched_kinds > 1:
+        return EventArgumentResolution(arguments={}, complete=False)
+    if absolute is not None:
+        return {"mode": "absolute", "value": int(absolute.group(1))}
+    if positive is not None:
+        return {"mode": "relative", "delta": int(positive.group(1))}
+    if negative is not None:
+        return {"mode": "relative", "delta": -int(negative.group(1))}
+    if signed is not None:
+        value = int(signed.group(2))
+        return {
+            "mode": "relative",
+            "delta": value if signed.group(1) == "+" else -value,
+        }
+    return EventArgumentResolution(arguments={}, complete=False)
+
+
+def _resolve_calendar_schedule(
+    request: EventRequest,
+    context: EventExecutionContext,
+) -> dict[str, Any] | EventArgumentResolution:
+    if request.arguments:
+        arguments = deepcopy(request.arguments)
+        complete = all(
+            key in arguments and arguments[key] not in (None, "")
+            for key in ("title", "start_at", "timezone")
+        )
+        return (
+            arguments
+            if complete
+            else EventArgumentResolution(arguments=arguments, complete=False)
+        )
+
+    content = _latest_user_content(context)
+    arguments: dict[str, Any] = {}
+    start_at = _RFC3339_RE.search(content)
+    timezone = re.search(
+        rf"(?:timezone|时区)\s*[::]?\s*({_TIMEZONE_PATTERN})",
+        content,
+        re.IGNORECASE,
+    )
+    english_title = re.search(
+        rf"\bschedule\s+[\"']([^\"']+)[\"']\s+at\s+{_RFC3339_PATTERN}",
+        content,
+        re.IGNORECASE,
+    )
+    chinese_title = re.search(r"标题\s*[::]\s*([^;;,\n]+)", content)
+    title = english_title or chinese_title
+    if title is not None:
+        arguments["title"] = title.group(1).strip()
+    if start_at is not None:
+        arguments["start_at"] = start_at.group()
+    if timezone is not None:
+        arguments["timezone"] = timezone.group(1)
+
+    recurrence = re.search(
+        r"(?:recurrence|重复)\s*[::]\s*([^;;,\n]+)", content, re.IGNORECASE
+    )
+    reminder = re.search(
+        r"(?:reminder_minutes|提醒分钟)\s*[::]\s*(\d+)",
+        content,
+        re.IGNORECASE,
+    )
+    if recurrence is not None:
+        arguments["recurrence"] = recurrence.group(1).strip()
+    if reminder is not None:
+        arguments["reminder_minutes"] = int(reminder.group(1))
+
+    if all(key in arguments for key in ("title", "start_at", "timezone")):
+        return arguments
+    return EventArgumentResolution(arguments=arguments, complete=False)
+
+
+def _resolve_knowledge_search(
+    request: EventRequest,
+    context: EventExecutionContext,
+) -> dict[str, Any] | EventArgumentResolution:
+    arguments = deepcopy(request.arguments)
+    query = arguments.get("query")
+    if not isinstance(query, str) or not query.strip():
+        query = _latest_user_content(context)
+    if query:
+        arguments["query"] = query.strip()
+        return arguments
+    return EventArgumentResolution(arguments=arguments, complete=False)
+
+
+def _latest_user_content(context: EventExecutionContext) -> str:
+    for message in reversed(context.history):
+        if message.role == "user" and message.content.strip():
+            return message.content.strip()
+    return ""
+
+
+def _fixed_clock() -> datetime:
+    return datetime(2026, 1, 1, tzinfo=datetime_timezone.utc)
+
+
+def _format_retrieved_at(value: datetime | str) -> str:
+    if isinstance(value, str):
+        return value
+    if value.tzinfo is None:
+        value = value.replace(tzinfo=datetime_timezone.utc)
+    rendered = value.astimezone(datetime_timezone.utc).isoformat(timespec="seconds")
+    return rendered.replace("+00:00", "Z")

+ 32 - 5
src/agent_lab/application/tools.py

@@ -6,6 +6,8 @@ from typing import Any
 
 from agent_lab.application.events import (
     ConfirmationPolicy,
+    CalendarSchedulePort,
+    DeviceVolumePort,
     EventDefinition,
     EventExecutionContext,
     EventKernel,
@@ -16,6 +18,9 @@ from agent_lab.application.events import (
     EventStatus,
     ResultPolicy,
     RiskLevel,
+    SessionTerminationPort,
+    WebSearchPort,
+    build_builtin_event_definitions,
 )
 from agent_lab.domain.events import EVENT_BLOCK_END, EVENT_BLOCK_START, ToolCallEvent
 from agent_lab.domain.messages import ChatMessage
@@ -54,14 +59,22 @@ class ToolDefinition:
 
 
 class ToolRegistry:
-    def __init__(self, definitions: Iterable[ToolDefinition]) -> None:
-        self._definitions: dict[str, ToolDefinition] = {}
+    def __init__(
+        self,
+        definitions: Iterable[ToolDefinition | EventDefinition],
+    ) -> None:
+        self._definitions: dict[str, EventDefinition] = {}
         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))
+            event_definition = (
+                definition
+                if isinstance(definition, EventDefinition)
+                else self._to_event_definition(definition)
+            )
+            self._definitions[definition.name] = event_definition
+            event_definitions.append(event_definition)
         self.event_registry = EventRegistry(event_definitions)
         self.kernel = EventKernel(self.event_registry)
 
@@ -232,7 +245,14 @@ class ToolRegistry:
         )
 
 
-def build_default_tool_registry() -> ToolRegistry:
+def build_default_tool_registry(
+    *,
+    session_termination_port: SessionTerminationPort | None = None,
+    device_volume_port: DeviceVolumePort | None = None,
+    calendar_schedule_port: CalendarSchedulePort | None = None,
+    web_search_port: WebSearchPort | None = None,
+    clock: Callable[[], Any] | None = None,
+) -> ToolRegistry:
     return ToolRegistry(
         [
             ToolDefinition(
@@ -274,6 +294,13 @@ def build_default_tool_registry() -> ToolRegistry:
                 handler=_handle_mock_ticket,
                 argument_resolver=_resolve_mock_ticket_arguments,
             ),
+            *build_builtin_event_definitions(
+                session_termination_port=session_termination_port,
+                device_volume_port=device_volume_port,
+                calendar_schedule_port=calendar_schedule_port,
+                web_search_port=web_search_port,
+                clock=clock,
+            ),
         ]
     )
 

+ 670 - 0
tests/test_builtin_event_plugins.py

@@ -0,0 +1,670 @@
+import json
+from datetime import datetime, timezone
+
+import pytest
+
+import agent_lab.application.events.builtin_plugins as builtin_plugins
+from agent_lab.application.events import (
+    EventArgumentResolution,
+    EventExecutionContext,
+    EventRequest,
+    ResultPolicy,
+    RiskLevel,
+)
+from agent_lab.application.tools import build_default_tool_registry
+from agent_lab.domain.events import ToolCallEvent
+from agent_lab.domain.messages import ChatMessage
+
+
+BUILTIN_EVENT_NAMES = [
+    "session.terminate",
+    "device.volume.adjust",
+    "calendar.schedule.create",
+    "knowledge.web.search",
+]
+
+
+def _event(event_id: str, name: str, arguments: dict):
+    return ToolCallEvent(
+        id=event_id,
+        name=name,
+        arguments=arguments,
+        raw_arguments=json.dumps(arguments, ensure_ascii=False),
+    )
+
+
+def test_default_registry_preserves_legacy_tools_and_adds_builtin_event_plugins():
+    registry = build_default_tool_registry()
+
+    assert [item["name"] for item in registry.available_tools()] == [
+        "handoff_note",
+        "mock_search",
+        "mock_ticket",
+        *BUILTIN_EVENT_NAMES,
+    ]
+
+
+@pytest.mark.parametrize(
+    (
+        "name",
+        "result_policy",
+        "risk_level",
+        "fallback_allowed",
+        "idempotency_key_fields",
+        "concurrency_class",
+        "conflict_keys",
+        "timeout_seconds",
+        "terminal",
+    ),
+    [
+        (
+            "session.terminate",
+            ResultPolicy.TERMINATE,
+            RiskLevel.HIGH,
+            False,
+            ("event_id",),
+            "session-lifecycle",
+            ("session",),
+            5.0,
+            True,
+        ),
+        (
+            "device.volume.adjust",
+            ResultPolicy.SILENT_SUCCESS,
+            RiskLevel.MEDIUM,
+            True,
+            ("event_id",),
+            "device-volume",
+            ("device.volume",),
+            5.0,
+            False,
+        ),
+        (
+            "calendar.schedule.create",
+            ResultPolicy.TEMPLATE_FOLLOW_UP,
+            RiskLevel.MEDIUM,
+            True,
+            ("event_id",),
+            "schedule-write",
+            ("calendar.schedule",),
+            10.0,
+            False,
+        ),
+        (
+            "knowledge.web.search",
+            ResultPolicy.LLM_FOLLOW_UP,
+            RiskLevel.LOW,
+            False,
+            (),
+            "read-only",
+            (),
+            10.0,
+            False,
+        ),
+    ],
+)
+def test_builtin_definition_metadata_is_owned_by_each_flat_plugin(
+    name,
+    result_policy,
+    risk_level,
+    fallback_allowed,
+    idempotency_key_fields,
+    concurrency_class,
+    conflict_keys,
+    timeout_seconds,
+    terminal,
+):
+    definition = build_default_tool_registry().event_registry.definition(name)
+
+    assert definition is not None
+    assert definition.result_policy is result_policy
+    assert definition.risk_level is risk_level
+    assert definition.fallback_allowed is fallback_allowed
+    assert definition.idempotency_key_fields == idempotency_key_fields
+    assert definition.concurrency_class == concurrency_class
+    assert definition.conflict_keys == conflict_keys
+    assert definition.timeout_seconds == timeout_seconds
+    assert definition.terminal is terminal
+
+
+def test_session_terminate_schema_accepts_optional_reason_only():
+    registry = build_default_tool_registry().event_registry
+
+    assert registry.iter_validation_errors("session.terminate", {}) == ()
+    assert registry.iter_validation_errors(
+        "session.terminate", {"reason": "user requested"}
+    ) == ()
+    assert registry.iter_validation_errors(
+        "session.terminate", {"reason": 1}
+    )
+    assert registry.iter_validation_errors(
+        "session.terminate", {"unexpected": True}
+    )
+
+
+@pytest.mark.parametrize(
+    "arguments",
+    [
+        {"mode": "absolute", "value": 0},
+        {"mode": "absolute", "value": 100},
+        {"mode": "relative", "delta": -100},
+        {"mode": "relative", "delta": 100},
+        {"mode": "mute"},
+        {"mode": "unmute"},
+    ],
+)
+def test_volume_schema_accepts_mode_specific_valid_arguments(arguments):
+    registry = build_default_tool_registry().event_registry
+
+    assert registry.iter_validation_errors("device.volume.adjust", arguments) == ()
+
+
+@pytest.mark.parametrize(
+    "arguments",
+    [
+        {"mode": "absolute"},
+        {"mode": "absolute", "value": -1},
+        {"mode": "absolute", "value": 101},
+        {"mode": "absolute", "value": 20, "delta": 5},
+        {"mode": "relative"},
+        {"mode": "relative", "delta": 0},
+        {"mode": "relative", "delta": 101},
+        {"mode": "relative", "delta": 5, "value": 20},
+        {"mode": "mute", "value": 0},
+        {"mode": "unmute", "delta": 1},
+    ],
+)
+def test_volume_schema_rejects_invalid_or_cross_mode_arguments(arguments):
+    registry = build_default_tool_registry().event_registry
+
+    assert registry.iter_validation_errors("device.volume.adjust", arguments)
+
+
+@pytest.mark.parametrize(
+    ("content", "expected"),
+    [
+        ("把音量调到 35", {"mode": "absolute", "value": 35}),
+        ("set volume to 42", {"mode": "absolute", "value": 42}),
+        ("音量增加 15", {"mode": "relative", "delta": 15}),
+        ("decrease volume by 20", {"mode": "relative", "delta": -20}),
+        ("请静音", {"mode": "mute"}),
+        ("please unmute", {"mode": "unmute"}),
+    ],
+)
+def test_volume_resolver_recognizes_clear_chinese_and_english(content, expected):
+    resolved = _resolve_text_event("device.volume.adjust", content)
+
+    assert resolved == expected
+
+
+@pytest.mark.parametrize("content", ["把音量调高", "静音还是取消静音"])
+def test_volume_resolver_marks_incomplete_or_ambiguous_requests_for_fallback(content):
+    resolved = _resolve_text_event("device.volume.adjust", content)
+
+    assert isinstance(resolved, EventArgumentResolution)
+    assert resolved.complete is False
+
+
+def test_schedule_schema_requires_explicit_rfc3339_and_timezone():
+    registry = build_default_tool_registry().event_registry
+    valid = {
+        "title": "Design review",
+        "start_at": "2026-07-14T09:30:00+08:00",
+        "timezone": "Asia/Shanghai",
+        "recurrence": "FREQ=WEEKLY",
+        "reminder_minutes": 15,
+    }
+
+    assert registry.iter_validation_errors("calendar.schedule.create", valid) == ()
+    assert registry.iter_validation_errors(
+        "calendar.schedule.create",
+        {**valid, "start_at": "tomorrow at nine"},
+    )
+    assert registry.iter_validation_errors(
+        "calendar.schedule.create",
+        {**valid, "timezone": "Shanghai"},
+    )
+    assert registry.iter_validation_errors(
+        "calendar.schedule.create",
+        {**valid, "reminder_minutes": -1},
+    )
+
+
+@pytest.mark.parametrize(
+    ("content", "expected"),
+    [
+        (
+            'Schedule "Design review" at 2026-07-14T09:30:00+08:00 '
+            "timezone Asia/Shanghai",
+            {
+                "title": "Design review",
+                "start_at": "2026-07-14T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+            },
+        ),
+        (
+            "标题:项目评审;开始时间:2026-07-14T09:30:00+08:00;"
+            "时区:Asia/Shanghai",
+            {
+                "title": "项目评审",
+                "start_at": "2026-07-14T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+            },
+        ),
+    ],
+)
+def test_schedule_resolver_accepts_only_explicit_datetime_details(content, expected):
+    resolved = _resolve_text_event("calendar.schedule.create", content)
+
+    assert resolved == expected
+
+
+def test_schedule_resolver_does_not_guess_relative_time():
+    resolved = _resolve_text_event(
+        "calendar.schedule.create", "Schedule standup tomorrow at nine"
+    )
+
+    assert isinstance(resolved, EventArgumentResolution)
+    assert resolved.complete is False
+
+
+def test_web_search_resolver_uses_latest_relevant_user_request():
+    registry = build_default_tool_registry()
+    definition = registry.event_registry.definition("knowledge.web.search")
+    assert definition is not None
+    assert definition.resolver is not None
+
+    resolved = definition.resolver(
+        EventRequest(id="search-1", name="knowledge.web.search"),
+        EventExecutionContext(
+            history=(
+                ChatMessage(role="user", content="old query"),
+                ChatMessage(role="assistant", content="assistant summary"),
+                ChatMessage(role="user", content="latest focused query"),
+                ChatMessage(role="assistant", content="I will search"),
+            )
+        ),
+    )
+
+    assert resolved == {"query": "latest focused query"}
+
+
+def test_web_search_schema_rejects_empty_query_and_out_of_range_result_limit():
+    registry = build_default_tool_registry().event_registry
+
+    assert registry.iter_validation_errors(
+        "knowledge.web.search", {"query": "focused", "max_results": 1}
+    ) == ()
+    assert registry.iter_validation_errors(
+        "knowledge.web.search", {"query": "focused", "max_results": 5}
+    ) == ()
+    assert registry.iter_validation_errors(
+        "knowledge.web.search", {"query": ""}
+    )
+    assert registry.iter_validation_errors(
+        "knowledge.web.search", {"query": "focused", "max_results": 0}
+    )
+    assert registry.iter_validation_errors(
+        "knowledge.web.search", {"query": "focused", "max_results": 6}
+    )
+
+
+def test_catalog_and_direct_provider_modes_share_the_same_builtin_schemas():
+    registry = build_default_tool_registry()
+    catalog = {item["name"]: item for item in registry.available_tools()}
+    direct = {
+        item["function"]["name"]: item["function"]
+        for item in registry.provider_tool_schemas(BUILTIN_EVENT_NAMES)
+    }
+
+    for name in BUILTIN_EVENT_NAMES:
+        assert direct[name]["description"] == catalog[name]["description"]
+        assert direct[name]["parameters"] == catalog[name]["parameters"]
+
+
+class RecordingSessionPort:
+    def __init__(self):
+        self.calls = []
+
+    def terminate(self, event_id, *, reason=None):
+        self.calls.append((event_id, reason))
+        return {"port": "session", "event_id": event_id, "reason": reason}
+
+
+class RecordingVolumePort:
+    def __init__(self):
+        self.calls = []
+
+    def adjust(self, event_id, *, mode, value=None, delta=None):
+        self.calls.append((event_id, mode, value, delta))
+        return {
+            "port": "volume",
+            "event_id": event_id,
+            "mode": mode,
+            "value": value,
+            "delta": delta,
+        }
+
+
+class RecordingCalendarPort:
+    def __init__(self):
+        self.calls = []
+
+    def create(
+        self,
+        event_id,
+        *,
+        title,
+        start_at,
+        timezone,
+        recurrence=None,
+        reminder_minutes=None,
+    ):
+        self.calls.append(
+            (
+                event_id,
+                title,
+                start_at,
+                timezone,
+                recurrence,
+                reminder_minutes,
+            )
+        )
+        return {"port": "calendar", "event_id": event_id, "title": title}
+
+
+class RecordingSearchPort:
+    def __init__(self):
+        self.calls = []
+
+    def search(self, event_id, *, query, max_results=3):
+        self.calls.append((event_id, query, max_results))
+        return {"port": "search", "query": query, "max_results": max_results}
+
+
+def test_injected_sync_ports_receive_validated_arguments_and_event_ids():
+    session = RecordingSessionPort()
+    volume = RecordingVolumePort()
+    calendar = RecordingCalendarPort()
+    search = RecordingSearchPort()
+    registry = build_default_tool_registry(
+        session_termination_port=session,
+        device_volume_port=volume,
+        calendar_schedule_port=calendar,
+        web_search_port=search,
+    )
+
+    assert registry.execute(
+        _event("session-1", "session.terminate", {"reason": "done"})
+    ) == {"port": "session", "event_id": "session-1", "reason": "done"}
+    assert registry.execute(
+        _event("volume-1", "device.volume.adjust", {"mode": "absolute", "value": 30})
+    )["port"] == "volume"
+    assert registry.execute(
+        _event(
+            "calendar-1",
+            "calendar.schedule.create",
+            {
+                "title": "Review",
+                "start_at": "2026-07-14T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+                "reminder_minutes": 10,
+            },
+        )
+    )["port"] == "calendar"
+    assert registry.execute(
+        _event(
+            "search-1",
+            "knowledge.web.search",
+            {"query": "agent kernels", "max_results": 2},
+        )
+    ) == {"port": "search", "query": "agent kernels", "max_results": 2}
+    assert session.calls == [("session-1", "done")]
+    assert volume.calls == [("volume-1", "absolute", 30, None)]
+    assert calendar.calls == [
+        (
+            "calendar-1",
+            "Review",
+            "2026-07-14T09:30:00+08:00",
+            "Asia/Shanghai",
+            None,
+            10,
+        )
+    ]
+    assert search.calls == [("search-1", "agent kernels", 2)]
+
+
+@pytest.mark.asyncio
+async def test_injected_async_ports_are_awaited_for_all_builtin_plugins():
+    class AsyncSessionPort(RecordingSessionPort):
+        async def terminate(self, event_id, *, reason=None):
+            return super().terminate(event_id, reason=reason)
+
+    class AsyncVolumePort(RecordingVolumePort):
+        async def adjust(self, event_id, *, mode, value=None, delta=None):
+            return super().adjust(
+                event_id, mode=mode, value=value, delta=delta
+            )
+
+    class AsyncCalendarPort(RecordingCalendarPort):
+        async def create(self, event_id, **arguments):
+            return super().create(event_id, **arguments)
+
+    class AsyncSearchPort(RecordingSearchPort):
+        async def search(self, event_id, *, query, max_results=3):
+            return super().search(
+                event_id, query=query, max_results=max_results
+            )
+
+    registry = build_default_tool_registry(
+        session_termination_port=AsyncSessionPort(),
+        device_volume_port=AsyncVolumePort(),
+        calendar_schedule_port=AsyncCalendarPort(),
+        web_search_port=AsyncSearchPort(),
+    )
+    events = [
+        _event("session-1", "session.terminate", {}),
+        _event(
+            "volume-1",
+            "device.volume.adjust",
+            {"mode": "relative", "delta": -10},
+        ),
+        _event(
+            "calendar-1",
+            "calendar.schedule.create",
+            {
+                "title": "Review",
+                "start_at": "2026-07-14T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+            },
+        ),
+        _event("search-1", "knowledge.web.search", {"query": "agent kernels"}),
+    ]
+
+    payloads = [await registry.execute_async(event) for event in events]
+
+    assert [payload["port"] for payload in payloads] == [
+        "session",
+        "volume",
+        "calendar",
+        "search",
+    ]
+
+
+@pytest.mark.parametrize(
+    "arguments",
+    [
+        {"mode": "absolute", "value": 101},
+        {"mode": "relative", "delta": 0},
+        {"mode": "mute", "value": 1},
+    ],
+)
+def test_invalid_volume_arguments_never_call_the_port(arguments):
+    volume = RecordingVolumePort()
+    registry = build_default_tool_registry(device_volume_port=volume)
+
+    payload = registry.execute(_event("volume-invalid", "device.volume.adjust", arguments))
+
+    assert payload["tool"] == "device.volume.adjust"
+    assert "invalid" in payload["error"]
+    assert volume.calls == []
+
+
+def test_ambiguous_volume_request_without_fallback_never_calls_the_port():
+    volume = RecordingVolumePort()
+    registry = build_default_tool_registry(device_volume_port=volume)
+
+    payload = registry.handle(
+        _event("volume-ambiguous", "device.volume.adjust", {}),
+        EventExecutionContext(
+            history=(ChatMessage(role="user", content="静音还是取消静音"),)
+        ),
+    )
+
+    assert payload == {
+        "tool": "device.volume.adjust",
+        "error": "missing required arguments: mode",
+    }
+    assert volume.calls == []
+
+
+@pytest.mark.parametrize(
+    ("port_name", "event"),
+    [
+        ("session_termination_port", _event("s", "session.terminate", {})),
+        (
+            "device_volume_port",
+            _event("v", "device.volume.adjust", {"mode": "mute"}),
+        ),
+        (
+            "calendar_schedule_port",
+            _event(
+                "c",
+                "calendar.schedule.create",
+                {
+                    "title": "Review",
+                    "start_at": "2026-07-14T09:30:00+08:00",
+                    "timezone": "Asia/Shanghai",
+                },
+            ),
+        ),
+        (
+            "web_search_port",
+            _event("w", "knowledge.web.search", {"query": "agent kernels"}),
+        ),
+    ],
+)
+def test_port_failures_are_normalized_by_the_generic_kernel(port_name, event):
+    class FailingPort:
+        def terminate(self, *args, **kwargs):
+            raise RuntimeError("port unavailable")
+
+        adjust = terminate
+        create = terminate
+        search = terminate
+
+    registry = build_default_tool_registry(**{port_name: FailingPort()})
+
+    payload = registry.execute(event)
+
+    assert payload == {
+        "tool": event.name,
+        "error": "tool handler failed: port unavailable",
+    }
+
+
+def test_default_stateful_adapters_are_idempotent_by_event_id():
+    registry = build_default_tool_registry()
+
+    first_session = registry.execute(
+        _event("same-session", "session.terminate", {"reason": "first"})
+    )
+    second_session = registry.execute(
+        _event("same-session", "session.terminate", {"reason": "second"})
+    )
+    first_volume = registry.execute(
+        _event(
+            "same-volume",
+            "device.volume.adjust",
+            {"mode": "absolute", "value": 20},
+        )
+    )
+    second_volume = registry.execute(
+        _event(
+            "same-volume",
+            "device.volume.adjust",
+            {"mode": "absolute", "value": 80},
+        )
+    )
+    first_schedule = registry.execute(
+        _event(
+            "same-schedule",
+            "calendar.schedule.create",
+            {
+                "title": "First",
+                "start_at": "2026-07-14T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+            },
+        )
+    )
+    second_schedule = registry.execute(
+        _event(
+            "same-schedule",
+            "calendar.schedule.create",
+            {
+                "title": "Second",
+                "start_at": "2026-07-15T09:30:00+08:00",
+                "timezone": "Asia/Shanghai",
+            },
+        )
+    )
+
+    assert second_session == first_session
+    assert first_session["reason"] == "first"
+    assert second_volume == first_volume
+    assert first_volume["value"] == 20
+    assert second_schedule == first_schedule
+    assert first_schedule["schedule"]["title"] == "First"
+
+
+def test_default_web_search_is_deterministic_compact_and_uses_injected_clock():
+    registry = build_default_tool_registry(
+        clock=lambda: datetime(2030, 1, 2, 3, 4, 5, tzinfo=timezone.utc)
+    )
+
+    payload = registry.execute(
+        _event(
+            "search-1",
+            "knowledge.web.search",
+            {"query": "agent kernels", "max_results": 2},
+        )
+    )
+
+    assert payload["tool"] == "knowledge.web.search"
+    assert payload["query"] == "agent kernels"
+    assert payload["retrieved_at"] == "2030-01-02T03:04:05Z"
+    assert len(payload["sources"]) == 2
+    assert all(source["url"].startswith("https://example.invalid/") for source in payload["sources"])
+    assert all(set(source) == {"title", "url", "snippet"} for source in payload["sources"])
+
+
+def test_builtin_port_protocols_and_default_adapter_types_are_public():
+    assert builtin_plugins.SessionTerminationPort
+    assert builtin_plugins.DeviceVolumePort
+    assert builtin_plugins.CalendarSchedulePort
+    assert builtin_plugins.WebSearchPort
+    assert builtin_plugins.InMemorySessionTerminationAdapter
+    assert builtin_plugins.InMemoryDeviceVolumeAdapter
+    assert builtin_plugins.InMemoryCalendarScheduleAdapter
+    assert builtin_plugins.InMemoryWebSearchAdapter
+
+
+def _resolve_text_event(name: str, content: str):
+    registry = build_default_tool_registry()
+    definition = registry.event_registry.definition(name)
+    assert definition is not None
+    assert definition.resolver is not None
+    return definition.resolver(
+        EventRequest(id="event-1", name=name),
+        EventExecutionContext(history=(ChatMessage(role="user", content=content),)),
+    )

+ 5 - 1
tests/test_websocket_api.py

@@ -197,7 +197,7 @@ def test_health_returns_ok():
     assert response.json() == {"status": "ok"}
 
 
-def test_api_tools_returns_handoff_note_metadata():
+def test_api_tools_returns_default_tool_catalog():
     app = create_app(runtime_factory=FakeRuntime)
     client = TestClient(app)
 
@@ -209,6 +209,10 @@ def test_api_tools_returns_handoff_note_metadata():
         "handoff_note",
         "mock_search",
         "mock_ticket",
+        "session.terminate",
+        "device.volume.adjust",
+        "calendar.schedule.create",
+        "knowledge.web.search",
     ]
     assert tools[0]["parameters"]["required"] == ["message"]