from __future__ import annotations import gc import json import threading import warnings from dataclasses import FrozenInstanceError from typing import Any import pytest from agent_lab.application.events import ( ConfirmationPolicy, EventDefinition, EventExecutionContext, EventKernel, EventRegistry, EventRequest, ResolvedEventArguments, EventSource, EventStatus, ResultPolicy, RiskLevel, ) from agent_lab.application.events.models import EventArgumentResolution 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_registry_rejects_invalid_draft_2020_12_schema(): definition = _definition(parameters={"type": 42}) with pytest.raises(ValueError, match="invalid event schema for example.lookup"): EventRegistry([definition]) def test_registry_does_not_expose_mutable_validator_instances(): registry = EventRegistry([_definition()]) assert not hasattr(registry, "validator") errors = list( registry.iter_validation_errors( "example.lookup", {"query": 42}, ) ) assert errors assert not hasattr(errors[0], "schema") assert not hasattr(errors[0], "instance") with pytest.raises(FrozenInstanceError): errors[0].message = "mutated" # type: ignore[misc] assert list( registry.iter_validation_errors( "example.lookup", {"query": 42}, ) ) 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, ) -> ResolvedEventArguments: fallback_calls.append(dict(request.arguments)) return ResolvedEventArguments( event_name=definition.name, arguments={"query": "resolved once"}, raw_arguments='{"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 @pytest.mark.parametrize( "arguments", [ {"query": 42}, {"query": "unsupported"}, {"query": "valid", "unexpected": True}, ], ) async def test_kernel_does_not_fallback_for_complete_invalid_arguments( arguments: dict[str, Any], ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "valid"}, raw_arguments='{"query":"valid"}', ) definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string", "enum": ["valid"]}}, "required": ["query"], "additionalProperties": False, }, resolver=lambda request, context: arguments, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio @pytest.mark.parametrize( ("arguments", "expected_error"), [ ( {"mode": "invalid"}, "invalid event arguments: 'invalid' is not one of ['valid']", ), ( {"mode": "valid", "count": "invalid"}, "invalid argument type for count: expected integer", ), ( {"mode": "valid", "unexpected": True}, "invalid event arguments: Additional properties are not allowed " "('unexpected' was unexpected)", ), ], ) async def test_kernel_does_not_fallback_for_required_plus_substantive_error( arguments: dict[str, Any], expected_error: str, ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters={ "type": "object", "properties": { "query": {"type": "string"}, "mode": {"enum": ["valid"]}, "count": {"type": "integer"}, }, "required": ["query"], "additionalProperties": False, }, resolver=lambda request, context: arguments, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name=definition.name), enabled_names=[definition.name], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == expected_error assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio @pytest.mark.parametrize( ("arguments", "properties", "expected_error"), [ ( {"mode": "invalid"}, {"mode": {"enum": ["valid"]}}, "invalid event arguments: 'invalid' is not one of ['valid']", ), ( {"count": "invalid"}, {"count": {"type": "integer"}}, "invalid argument type for count: expected integer", ), ( {"unexpected": True}, {}, "invalid event arguments: Additional properties are not allowed " "('unexpected' was unexpected)", ), ], ) async def test_structured_incomplete_does_not_fallback_over_substantive_error( arguments: dict[str, Any], properties: dict[str, Any], expected_error: str, ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters={ "type": "object", "properties": { "query": {"type": "string"}, **properties, }, "required": ["query"], "additionalProperties": False, }, resolver=lambda request, context: EventArgumentResolution( arguments=arguments, complete=False, ), ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name=definition.name), enabled_names=[definition.name], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == expected_error assert result.used_fallback is False assert fallback_calls == 0 def _discriminated_composed_parameters(composition: str) -> dict[str, Any]: return { "type": "object", "properties": { "choice": { composition: [ { "type": "object", "properties": { "kind": {"const": "a"}, "value": {"type": "string"}, }, "required": ["kind", "value"], }, { "type": "object", "properties": { "kind": {"const": "b"}, "count": {"type": "integer"}, }, "required": ["kind", "count"], }, ] } }, "required": ["choice"], } @pytest.mark.asyncio @pytest.mark.parametrize("composition", ["anyOf", "oneOf"]) async def test_kernel_falls_back_for_matching_composed_branch_missing_required( composition: str, ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"choice": {"kind": "a", "value": "resolved"}}, raw_arguments='{"choice":{"kind":"a","value":"resolved"}}', ) definition = _definition( parameters=_discriminated_composed_parameters(composition), resolver=lambda request, context: {"choice": {"kind": "a"}}, handler=lambda request: {"ok": True}, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.SUCCESS assert result.used_fallback is True assert fallback_calls == 1 @pytest.mark.asyncio @pytest.mark.parametrize("composition", ["anyOf", "oneOf"]) async def test_kernel_does_not_fallback_when_matching_composed_branch_is_invalid( composition: str, ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters=_discriminated_composed_parameters(composition), resolver=lambda request, context: { "choice": {"kind": "a", "value": 42} }, handler=lambda request: {"ok": True}, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio @pytest.mark.parametrize("composition", ["anyOf", "oneOf"]) @pytest.mark.parametrize("choice", [{}, {"kind": "other"}]) async def test_kernel_does_not_fallback_when_composed_branch_is_ambiguous( composition: str, choice: dict[str, Any], ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters=_discriminated_composed_parameters(composition), resolver=lambda request, context: {"choice": choice}, handler=lambda request: {"ok": True}, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio async def test_kernel_does_not_fallback_for_substantive_and_composed_missing_errors(): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") parameters = _discriminated_composed_parameters("anyOf") parameters["properties"]["mode"] = {"enum": ["valid"]} definition = _definition( parameters=parameters, resolver=lambda request, context: { "mode": "invalid", "choice": {"kind": "a"}, }, handler=lambda request: {"ok": True}, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name=definition.name), enabled_names=[definition.name], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == "invalid event arguments: 'invalid' is not one of ['valid']" assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio async def test_kernel_does_not_fallback_when_definition_disallows_it(): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "not allowed"}, raw_arguments='{"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_structured_resolution_can_mark_optional_arguments_incomplete(): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "resolved"}, raw_arguments='{"query":"resolved"}', ) definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string"}}, "additionalProperties": False, }, resolver=lambda request, context: EventArgumentResolution( arguments={}, complete=False, ), ) result = await EventKernel( EventRegistry([definition]), 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"} assert result.used_fallback is True assert fallback_calls == 1 @pytest.mark.asyncio async def test_structured_incomplete_optional_arguments_fail_without_fallback(): definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string"}}, }, resolver=lambda request, context: EventArgumentResolution( arguments={}, complete=False, ), handler=lambda request: {"ok": True}, ) result = await EventKernel(EventRegistry([definition])).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == "event arguments incomplete" @pytest.mark.asyncio async def test_plain_dict_resolution_remains_complete_for_optional_schema(): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string"}}, }, resolver=lambda request, context: {}, handler=lambda request: {"ok": True}, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.SUCCESS assert fallback_calls == 0 def test_sync_kernel_rejects_structured_incomplete_optional_arguments(): handler_calls = 0 def handler(request: EventRequest) -> dict[str, Any]: nonlocal handler_calls handler_calls += 1 return {"ok": True} definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string"}}, }, resolver=lambda request, context: EventArgumentResolution( arguments={}, complete=False, ), handler=handler, ) result = EventKernel(EventRegistry([definition])).execute_sync( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == "event arguments incomplete" assert handler_calls == 0 @pytest.mark.parametrize( "resolver", [ lambda request, context: EventArgumentResolution( arguments={}, complete=True, ), lambda request, context: {}, ], ) def test_sync_kernel_preserves_complete_compatible_resolvers(resolver: Any): handler_calls = 0 def handler(request: EventRequest) -> dict[str, Any]: nonlocal handler_calls handler_calls += 1 return {"ok": True} definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "string"}}, }, resolver=resolver, handler=handler, ) result = EventKernel(EventRegistry([definition])).execute_sync( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.SUCCESS assert result.payload == {"ok": True} assert handler_calls == 1 @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) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "fallback"}, raw_arguments='{"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) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "fallback"}, raw_arguments='{"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 @pytest.mark.parametrize("handler_kind", ["async", "sync", "failure"]) async def test_kernel_measures_only_actual_handler_latency(handler_kind: str): clock_values = iter([10.0, 10.025]) async def async_handler(request: EventRequest) -> dict[str, Any]: return {"query": request.arguments["query"]} def sync_handler(request: EventRequest) -> dict[str, Any]: return {"query": request.arguments["query"]} def failing_handler(request: EventRequest) -> dict[str, Any]: raise RuntimeError("boom") handler = { "async": async_handler, "sync": sync_handler, "failure": failing_handler, }[handler_kind] result = await EventKernel( EventRegistry([_definition(handler=handler)]), monotonic_clock=lambda: next(clock_values), ).execute( EventRequest( id="event-1", name="example.lookup", arguments={"query": "value"}, source=EventSource.PROVIDER_RESOLVED, ) ) assert result.tool_latency_ms == 25 if handler_kind == "failure": assert result.status is EventStatus.HANDLER_ERROR else: assert result.status is EventStatus.SUCCESS @pytest.mark.asyncio async def test_kernel_handler_latency_excludes_argument_fallback_time(): clock_values = iter([20.0, 20.012]) async def fallback(*args: Any) -> ResolvedEventArguments: return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "fallback"}, raw_arguments='{"query":"fallback"}', ) result = await EventKernel( EventRegistry([_definition()]), argument_fallback=fallback, monotonic_clock=lambda: next(clock_values), ).execute( EventRequest( id="event-1", name="example.lookup", arguments={}, raw_arguments="{}", ) ) assert result.used_fallback is True assert result.tool_latency_ms == 12 @pytest.mark.asyncio async def test_kernel_validates_complete_draft_2020_12_schema(): definition = _definition( parameters={ "type": "object", "properties": { "mode": {"enum": ["quick", "deep"]}, "target": {"type": ["string", "null"]}, "filters": { "type": "array", "items": { "type": "object", "properties": {"score": {"type": "number", "minimum": 0}}, "required": ["score"], "additionalProperties": False, }, }, }, "required": ["mode", "target", "filters"], "additionalProperties": False, }, handler=lambda request: {"event": request.name}, ) kernel = EventKernel(EventRegistry([definition])) valid = await kernel.execute( EventRequest( id="valid", name=definition.name, arguments={ "mode": "deep", "target": None, "filters": [{"score": 0.5}], }, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[definition.name], ) invalid = await kernel.execute( EventRequest( id="invalid", name=definition.name, arguments={ "mode": "other", "target": 7, "filters": [{"score": -1, "extra": True}], "unexpected": True, }, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[definition.name], ) assert valid.status is EventStatus.SUCCESS assert invalid.status is EventStatus.INVALID_ARGUMENTS assert invalid.error.startswith("invalid event arguments:") @pytest.mark.asyncio async def test_kernel_normalizes_validator_runtime_exceptions(): definition = _definition(parameters={"$ref": "urn:agent-lab:missing-schema"}) result = await EventKernel(EventRegistry([definition])).execute( EventRequest( id="event-1", name=definition.name, arguments={}, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[definition.name], ) assert result.status is EventStatus.DEFINITION_ERROR assert result.error.startswith("event argument validation failed:") @pytest.mark.asyncio async def test_tool_registry_maps_definition_errors_to_tool_compatibility_payload(): registry = ToolRegistry( [ ToolDefinition( name="broken.lookup", description="Broken lookup.", parameters={"$ref": "urn:agent-lab:missing-schema"}, handler=lambda event: {"tool": event.name}, ) ] ) payload = await registry.execute_async( ToolCallEvent( id="call-1", name="broken.lookup", arguments={}, raw_arguments="{}", ) ) assert payload["tool"] == "broken.lookup" assert payload["error"].startswith("tool definition validation failed:") @pytest.mark.asyncio @pytest.mark.parametrize("boundary", ["resolver", "fallback"]) async def test_kernel_normalizes_resolution_boundary_exceptions(boundary: str): def resolver(request: EventRequest, context: EventExecutionContext) -> dict[str, Any]: if boundary == "resolver": raise RuntimeError("resolver boom") return {} async def fallback(*args: Any) -> ResolvedEventArguments: raise RuntimeError("fallback boom") result = await EventKernel( EventRegistry([_definition(resolver=resolver)]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error == f"event argument {boundary} failed: {boundary} boom" @pytest.mark.asyncio @pytest.mark.parametrize("boundary", ["resolver", "fallback"]) async def test_kernel_normalizes_invalid_resolution_payloads(boundary: str): resolver = ( (lambda request, context: None) if boundary == "resolver" else (lambda request, context: {}) ) async def fallback(*args: Any) -> Any: return {"query": "legacy bare mapping"} result = await EventKernel( EventRegistry([_definition(resolver=resolver)]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error == f"event argument {boundary} returned invalid payload" @pytest.mark.asyncio async def test_kernel_normalizes_non_json_resolver_arguments(): result = await EventKernel( EventRegistry( [ _definition( resolver=lambda request, context: {"query": object()} ) ] ) ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error.startswith("event argument resolver failed to serialize:") class _ExplodingDeepcopyDict(dict[str, Any]): def __deepcopy__(self, memo: dict[int, Any]) -> dict[str, Any]: raise RuntimeError("deepcopy must not be used") @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) @pytest.mark.parametrize( ("source", "expected_status", "expected_error"), [ ( EventSource.TEXT_EVENT, EventStatus.RESOLUTION_ERROR, "event argument resolver failed to serialize: " "JSON round-trip changed payload", ), ( EventSource.PROVIDER_RESOLVED, EventStatus.INVALID_ARGUMENTS, "provider-resolved event arguments are not valid JSON", ), ], ) async def test_kernel_normalizes_argument_snapshot_failures_without_deepcopy( execution: str, source: EventSource, expected_status: EventStatus, expected_error: str, ): fallback_calls = 0 arguments = {"query": _ExplodingDeepcopyDict({"nested": "value"})} async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback should not run") definition = _definition( parameters={ "type": "object", "properties": {"query": {"type": "object"}}, "required": ["query"], }, resolver=lambda request, context: arguments, handler=lambda request: {"ok": True}, ) request = EventRequest( id="event-1", name="example.lookup", arguments=arguments if source is EventSource.PROVIDER_RESOLVED else {}, source=source, ) kernel = EventKernel(EventRegistry([definition]), argument_fallback=fallback) result = ( await kernel.execute(request, enabled_names=[definition.name]) if execution == "async" else kernel.execute_sync(request, enabled_names=[definition.name]) ) assert result.status is expected_status assert result.status is not EventStatus.DEFINITION_ERROR assert result.error == expected_error assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio async def test_kernel_normalizes_non_object_fallback_arguments(): async def fallback(*args: Any) -> ResolvedEventArguments: return ResolvedEventArguments( event_name="example.lookup", arguments=["not", "an", "object"], # type: ignore[arg-type] raw_arguments='["not","an","object"]', ) result = await EventKernel( EventRegistry([_definition(resolver=lambda request, context: {})]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error == "event argument fallback returned invalid payload" @pytest.mark.asyncio @pytest.mark.parametrize( ("resolved", "expected_error"), [ ( ResolvedEventArguments( event_name="another.event", arguments={"query": "value"}, raw_arguments='{"query":"value"}', ), "fallback returned tool another.event for example.lookup", ), ( ResolvedEventArguments( event_name="example.lookup", arguments={"query": "value"}, raw_arguments="not-json", ), "fallback raw arguments are not valid JSON", ), ( ResolvedEventArguments( event_name="example.lookup", arguments={"query": "parsed"}, raw_arguments='{"query":"raw"}', ), "fallback raw arguments do not match parsed arguments", ), ], ) async def test_kernel_rejects_inconsistent_structured_fallback( resolved: ResolvedEventArguments, expected_error: str, ): async def fallback(*args: Any) -> ResolvedEventArguments: return resolved result = await EventKernel( EventRegistry([_definition(resolver=lambda request, context: {})]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error == expected_error @pytest.mark.asyncio @pytest.mark.parametrize( ("arguments", "raw_arguments"), [ ({"query": 1}, '{"query":true}'), ({"query": 1.0}, '{"query":1}'), ({"query": {"nested": [1]}}, '{"query":{"nested":[true]}}'), ({"query": float("nan")}, '{"query":NaN}'), ({"query": float("inf")}, '{"query":Infinity}'), ], ) async def test_kernel_rejects_noncanonical_fallback_json( arguments: dict[str, Any], raw_arguments: str, ): async def fallback(*args: Any) -> ResolvedEventArguments: return ResolvedEventArguments( event_name="example.lookup", arguments=arguments, raw_arguments=raw_arguments, ) result = await EventKernel( EventRegistry([_definition(resolver=lambda request, context: {})]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.RESOLUTION_ERROR @pytest.mark.asyncio async def test_registry_schema_is_isolated_from_caller_mutation(): parameters = { "type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"], "additionalProperties": False, } registry = EventRegistry([_definition(parameters=parameters)]) parameters["properties"]["query"]["type"] = "integer" parameters["required"].clear() parameters["additionalProperties"] = True assert registry.catalog()[0]["parameters"] == { "type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"], "additionalProperties": False, } registered = registry.definition("example.lookup") assert registered is not None with pytest.raises(TypeError): registered.parameters["additionalProperties"] = True with pytest.raises(TypeError): registered.parameters["properties"]["query"]["type"] = "integer" with pytest.raises(AttributeError): registered.parameters["required"].append("unexpected") result = await EventKernel(registry).execute( EventRequest( id="event-1", name="example.lookup", arguments={"query": 42, "unexpected": True}, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=["example.lookup"], ) assert result.status is EventStatus.INVALID_ARGUMENTS @pytest.mark.asyncio async def test_kernel_passes_consistent_fallback_arguments_and_raw_json_to_handler(): captured: list[EventRequest] = [] async def fallback(*args: Any) -> ResolvedEventArguments: return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "resolved"}, raw_arguments='{"query":"resolved"}', ) result = await EventKernel( EventRegistry( [ _definition( resolver=lambda request, context: {}, handler=lambda request: captured.append(request) or {"ok": True}, ) ] ), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name="example.lookup"), enabled_names=["example.lookup"], ) assert result.status is EventStatus.SUCCESS assert captured[0].arguments == json.loads(captured[0].raw_arguments) assert result.raw_arguments == captured[0].raw_arguments @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) async def test_kernel_isolates_nested_handler_mutation_from_audit_values( execution: str, ): caller_arguments = {"nested": {"items": ["original"]}} handler_arguments: list[dict[str, Any]] = [] def handler(request: EventRequest) -> dict[str, Any]: handler_arguments.append(request.arguments) request.arguments["nested"]["items"].append("handler") return {"nested": request.arguments["nested"]} definition = _definition( parameters={ "type": "object", "properties": { "nested": { "type": "object", "properties": { "items": {"type": "array", "items": {"type": "string"}} }, "required": ["items"], } }, "required": ["nested"], }, handler=handler, ) request = EventRequest( id="event-1", name="example.lookup", arguments=caller_arguments, raw_arguments='{"nested":{"items":["original"]}}', source=EventSource.PROVIDER_RESOLVED, ) kernel = EventKernel(EventRegistry([definition])) result = ( await kernel.execute(request, enabled_names=["example.lookup"]) if execution == "async" else kernel.execute_sync(request, enabled_names=["example.lookup"]) ) assert result.status is EventStatus.SUCCESS assert caller_arguments == {"nested": {"items": ["original"]}} assert request.arguments == {"nested": {"items": ["original"]}} assert result.arguments == {"nested": {"items": ["original"]}} assert json.loads(result.raw_arguments) == result.arguments assert result.payload == {"nested": {"items": ["original", "handler"]}} handler_arguments[0]["nested"]["items"].append("later") assert result.payload == {"nested": {"items": ["original", "handler"]}} result.payload["nested"]["items"].append("result") assert result.arguments == {"nested": {"items": ["original"]}} assert caller_arguments == {"nested": {"items": ["original"]}} @pytest.mark.asyncio @pytest.mark.parametrize("payload", [None, "text", 1, ["item"]]) async def test_kernel_normalizes_invalid_handler_payloads(payload: Any): result = await EventKernel( EventRegistry([_definition(handler=lambda request: payload)]) ).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 returned non-object payload" @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) @pytest.mark.parametrize( "invalid_value", [object(), {"set-item"}, float("nan"), float("inf")], ) async def test_kernel_rejects_non_json_handler_dictionary_payloads( execution: str, invalid_value: Any, ): definition = _definition( handler=lambda request: {"invalid": invalid_value}, ) request = EventRequest( id="event-1", name="example.lookup", arguments={"query": "value"}, raw_arguments='{"query":"value"}', source=EventSource.PROVIDER_RESOLVED, ) kernel = EventKernel(EventRegistry([definition])) result = ( await kernel.execute(request, enabled_names=["example.lookup"]) if execution == "async" else kernel.execute_sync(request, enabled_names=["example.lookup"]) ) assert result.status is EventStatus.HANDLER_ERROR assert result.error == "event handler returned non-JSON payload" class _ExplodingItemsDict(dict[str, Any]): def items(self): raise RuntimeError("payload items failed") @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) async def test_kernel_normalizes_handler_payload_snapshot_exceptions( execution: str, ): definition = _definition( handler=lambda request: _ExplodingItemsDict(ok=True), ) request = EventRequest( id="event-1", name=definition.name, arguments={"query": "value"}, raw_arguments='{"query":"value"}', source=EventSource.PROVIDER_RESOLVED, ) kernel = EventKernel(EventRegistry([definition])) result = ( await kernel.execute(request, enabled_names=[definition.name]) if execution == "async" else kernel.execute_sync(request, enabled_names=[definition.name]) ) assert result.status is EventStatus.HANDLER_ERROR assert result.error == "event handler returned non-JSON payload" @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) async def test_kernel_does_not_swallow_base_exception_from_payload_snapshot( execution: str, ): class SnapshotAbort(BaseException): pass class AbortingItemsDict(dict[str, Any]): def items(self): raise SnapshotAbort definition = _definition( handler=lambda request: AbortingItemsDict(ok=True), ) request = EventRequest( id="event-1", name=definition.name, arguments={"query": "value"}, source=EventSource.PROVIDER_RESOLVED, ) kernel = EventKernel(EventRegistry([definition])) with pytest.raises(SnapshotAbort): if execution == "async": await kernel.execute(request, enabled_names=[definition.name]) else: kernel.execute_sync(request, enabled_names=[definition.name]) @pytest.mark.asyncio async def test_tool_registry_normalizes_handler_payload_snapshot_exceptions(): registry = ToolRegistry( [ ToolDefinition( name="bad_payload", description="Return a payload that fails during snapshot.", parameters={"type": "object"}, handler=lambda event: _ExplodingItemsDict(ok=True), ) ] ) event = ToolCallEvent( id="call-1", name="bad_payload", arguments={}, raw_arguments="{}", ) expected = { "tool": "bad_payload", "error": "event handler returned non-JSON payload", } assert registry.execute(event) == expected assert await registry.execute_async(event) == expected @pytest.mark.asyncio async def test_async_kernel_supports_async_handler(): async def handler(request: EventRequest) -> dict[str, Any]: return {"query": request.arguments["query"]} result = await EventKernel( EventRegistry([_definition(handler=handler)]) ).execute( EventRequest( id="event-1", name="example.lookup", arguments={"query": "async"}, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=["example.lookup"], ) assert result.status is EventStatus.SUCCESS assert result.payload == {"query": "async"} @pytest.mark.asyncio async def test_async_kernel_offloads_sync_handler_but_execute_sync_stays_inline(): caller_thread = threading.get_ident() handler_threads: list[int] = [] def handler(request: EventRequest) -> dict[str, Any]: handler_threads.append(threading.get_ident()) return {"query": request.arguments["query"]} kernel = EventKernel(EventRegistry([_definition(handler=handler)])) request = EventRequest( id="threaded", name="example.lookup", arguments={"query": "threaded"}, raw_arguments='{"query":"threaded"}', ) async_result = await kernel.execute(request) sync_result = kernel.execute_sync(request) assert async_result.status is EventStatus.SUCCESS assert sync_result.status is EventStatus.SUCCESS assert handler_threads[0] != caller_thread assert handler_threads[1] == caller_thread @pytest.mark.asyncio async def test_tool_registry_async_entry_points_support_async_handler(): async def handler(event: ToolCallEvent) -> dict[str, Any]: return {"tool": event.name, "query": event.arguments["query"]} registry = ToolRegistry( [ ToolDefinition( name="async.lookup", description="Async lookup.", parameters={ "type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"], }, handler=handler, ) ] ) event = ToolCallEvent( id="call-1", name="async.lookup", arguments={"query": "value"}, raw_arguments='{"query":"value"}', ) assert await registry.handle_async(event) == { "tool": "async.lookup", "query": "value", } assert await registry.execute_async(event) == { "tool": "async.lookup", "query": "value", } def test_tool_registry_sync_facade_rejects_async_handler_without_runtime_warning(): called = False async def handler(event: ToolCallEvent) -> dict[str, Any]: nonlocal called called = True return {"tool": event.name} registry = ToolRegistry( [ ToolDefinition( name="async.lookup", description="Async lookup.", parameters={"type": "object"}, handler=handler, ) ] ) event = ToolCallEvent( id="call-1", name="async.lookup", arguments={}, raw_arguments="{}", ) with warnings.catch_warnings(record=True) as captured: warnings.simplefilter("always") payload = registry.execute(event) gc.collect() assert payload == { "tool": "async.lookup", "error": "tool handler failed: async event handlers require execute_async", } assert called is False assert not [warning for warning in captured if issubclass(warning.category, RuntimeWarning)] @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 not 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 @pytest.mark.asyncio @pytest.mark.parametrize("name", ["alpha.one", "beta-two", "任意.事件"]) async def test_kernel_applies_identical_behavior_to_arbitrary_event_names(name: str): definition = _definition(name=name) result = await EventKernel(EventRegistry([definition])).execute( EventRequest( id="event-1", name=name, arguments={"query": "value"}, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[name], ) assert result.status is EventStatus.SUCCESS assert result.payload == {"event": name, "query": "value"} @pytest.mark.asyncio async def test_kernel_runs_normalizer_after_text_resolver_before_validation(): captured: list[EventRequest] = [] definition = _definition( resolver=lambda request, context: {"query": " TEXT Value "}, normalizer=lambda arguments: {"query": arguments["query"].strip().lower()}, handler=lambda request: captured.append(request) or {"ok": True}, ) result = await EventKernel(EventRegistry([definition])).execute( EventRequest(id="event-1", name=definition.name), enabled_names=[definition.name], ) assert result.status is EventStatus.SUCCESS assert result.arguments == {"query": "text value"} assert result.raw_arguments == '{"query":"text value"}' assert captured[0].arguments == {"query": "text value"} assert captured[0].raw_arguments == '{"query":"text value"}' @pytest.mark.asyncio async def test_kernel_normalizes_provider_arguments_and_preserves_original_raw_json(): captured: list[EventRequest] = [] original_raw = '{ "query": " PROVIDER Value " }' definition = _definition( resolver=lambda request, context: (_ for _ in ()).throw( AssertionError("provider must not run resolver") ), normalizer=lambda arguments: {"query": arguments["query"].strip().lower()}, handler=lambda request: captured.append(request) or {"ok": True}, ) result = await EventKernel(EventRegistry([definition])).execute( EventRequest( id="event-1", name=definition.name, arguments={"query": " PROVIDER Value "}, raw_arguments=original_raw, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[definition.name], ) assert result.status is EventStatus.SUCCESS assert result.arguments == {"query": "provider value"} assert result.raw_arguments == original_raw assert captured[0].arguments == {"query": "provider value"} assert captured[0].raw_arguments == original_raw @pytest.mark.asyncio async def test_provider_normalization_never_enables_argument_fallback(): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "fallback"}, raw_arguments='{"query":"fallback"}', ) definition = _definition(normalizer=lambda arguments: {}) original_raw = '{"query":"provider"}' result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback ).execute( EventRequest( id="event-1", name=definition.name, arguments={"query": "provider"}, raw_arguments=original_raw, source=EventSource.PROVIDER_RESOLVED, ), enabled_names=[definition.name], ) assert result.status is EventStatus.INVALID_ARGUMENTS assert result.error == "missing required arguments: query" assert result.raw_arguments == original_raw assert result.used_fallback is False assert fallback_calls == 0 @pytest.mark.asyncio @pytest.mark.parametrize("execution", ["async", "sync"]) @pytest.mark.parametrize( ("source", "expected_status", "expected_error"), [ ( EventSource.TEXT_EVENT, EventStatus.RESOLUTION_ERROR, "event argument normalizer failed", ), ( EventSource.PROVIDER_RESOLVED, EventStatus.INVALID_ARGUMENTS, "provider-resolved event arguments could not be normalized", ), ], ) @pytest.mark.parametrize("normalizer_kind", ["raises", "invalid_return"]) async def test_kernel_normalizes_normalizer_failures_without_leaking_details( execution: str, source: EventSource, expected_status: EventStatus, expected_error: str, normalizer_kind: str, ): fallback_calls = 0 async def fallback(*args: Any) -> ResolvedEventArguments: nonlocal fallback_calls fallback_calls += 1 raise AssertionError("fallback must not run") def normalizer(arguments: dict[str, Any]): if normalizer_kind == "raises": raise RuntimeError("secret adapter detail") return [arguments] definition = _definition( resolver=lambda request, context: {"query": "value"}, normalizer=normalizer, ) request = EventRequest( id="event-1", name=definition.name, arguments={"query": "value"} if source is EventSource.PROVIDER_RESOLVED else {}, raw_arguments='{"query":"value"}', source=source, ) kernel = EventKernel(EventRegistry([definition]), argument_fallback=fallback) result = ( await kernel.execute(request, enabled_names=[definition.name]) if execution == "async" else kernel.execute_sync(request, enabled_names=[definition.name]) ) assert result.status is expected_status assert result.error == expected_error assert "secret" not in (result.error or "") assert result.used_fallback is False assert fallback_calls == 0 def test_tool_definition_normalizer_remains_compatible_with_provider_execution(): registry = ToolRegistry( [ ToolDefinition( name="compat.normalize", description="Normalize compatibility data.", parameters={ "type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"], }, handler=lambda event: { "tool": event.name, "query": event.arguments["query"], }, normalizer=lambda arguments: { "query": arguments["query"].strip().lower() }, ) ] ) payload = registry.execute( ToolCallEvent( id="call-1", name="compat.normalize", arguments={"query": " PROVIDER "}, raw_arguments='{ "query": " PROVIDER " }', ) ) assert payload == {"tool": "compat.normalize", "query": "provider"} @pytest.mark.asyncio async def test_post_fallback_normalizer_failure_preserves_used_fallback_flag(): normalizer_calls = 0 def normalizer(arguments: dict[str, Any]) -> dict[str, Any]: nonlocal normalizer_calls normalizer_calls += 1 if "query" in arguments: raise RuntimeError("post-fallback failure") return arguments async def fallback(*args: Any) -> ResolvedEventArguments: return ResolvedEventArguments( event_name="example.lookup", arguments={"query": "fallback"}, raw_arguments='{"query":"fallback"}', ) definition = _definition( resolver=lambda request, context: {}, normalizer=normalizer, ) result = await EventKernel( EventRegistry([definition]), argument_fallback=fallback, ).execute( EventRequest(id="event-1", name=definition.name), enabled_names=[definition.name], ) assert result.status is EventStatus.RESOLUTION_ERROR assert result.error == "event argument normalizer failed" assert result.used_fallback is True assert normalizer_calls == 2 def test_event_definition_preserves_historical_positional_field_order(): definition = EventDefinition( "compat.event", "Compatibility event.", {"type": "object"}, lambda request: {"ok": True}, None, "legacy-schema", False, ) assert definition.resolver is None assert definition.schema_version == "legacy-schema" assert definition.fallback_allowed is False assert definition.normalizer is None def test_tool_definition_preserves_historical_positional_field_order(): definition = ToolDefinition( "compat.tool", "Compatibility tool.", {"type": "object"}, lambda event: {"ok": True}, None, "legacy-schema", False, ) assert definition.argument_resolver is None assert definition.schema_version == "legacy-schema" assert definition.fallback_allowed is False assert definition.normalizer is None