Эх сурвалжийг харах

fix: preserve batch outcome and resource safety

Problem:
Batch replay lost deadline semantics, invalid canonical keys could inherit successful results, released scopes could be repopulated by late tasks, timed-out sync work released physical concurrency resources early, and logical duplicates leaked into policy and audit output.

Risk:
Worker threads cannot be force-killed after timeout. Conflict locks and parallel-slot leases remain held until the underlying kernel task actually completes, so production side-effect ports must remain idempotent and tolerate late completion.
zhenyu.hu 2 долоо хоног өмнө
parent
commit
171296e26c

+ 120 - 53
src/agent_lab/application/events/batch.py

@@ -1,6 +1,7 @@
 from __future__ import annotations
 
 import asyncio
+import inspect
 import json
 from collections import OrderedDict
 from collections.abc import Hashable, Iterable, Sequence
@@ -18,6 +19,9 @@ from agent_lab.application.events.models import (
 )
 
 
+_ReplayKey = tuple[Hashable, int, str, str]
+
+
 @dataclass(frozen=True)
 class _BatchItem:
     key: str
@@ -55,12 +59,11 @@ class EventBatchExecutor:
         self.max_replay_entries = max_replay_entries
         self._semaphore = asyncio.Semaphore(max_parallel_events)
         self._conflict_locks: dict[str, asyncio.Lock] = {}
-        self._replay_results: OrderedDict[
-            tuple[Hashable, str, str], EventResult
-        ] = OrderedDict()
-        self._inflight: dict[
-            tuple[Hashable, str, str], asyncio.Task[_ExecutionOutcome]
-        ] = {}
+        self._replay_results: OrderedDict[_ReplayKey, _ExecutionOutcome] = (
+            OrderedDict()
+        )
+        self._inflight: dict[_ReplayKey, asyncio.Task[_ExecutionOutcome]] = {}
+        self._scope_generations: dict[Hashable, int] = {}
         self._state_lock = asyncio.Lock()
 
     async def execute(
@@ -73,6 +76,7 @@ class EventBatchExecutor:
     ) -> EventBatchResult:
         if not requests:
             return EventBatchResult()
+        generation = self._scope_generations.get(scope, 0)
 
         items: list[_BatchItem] = []
         indexes_by_key: dict[str, list[int]] = {}
@@ -114,6 +118,7 @@ class EventBatchExecutor:
                     self._execute_replayable(
                         item,
                         scope=scope,
+                        generation=generation,
                         enabled_names=enabled_names,
                         context=context,
                         deadline=deadline,
@@ -140,11 +145,18 @@ class EventBatchExecutor:
                     result,
                     requests[index],
                     deduplicated_from=(
-                        primary_id if requests[index].id != primary_id else None
+                        primary_id if index != indexes[0] else None
                     ),
                 )
 
-        await self._cache_ordered_results(scope, requests, indexes_by_key, ordered)
+        await self._cache_ordered_results(
+            scope,
+            generation,
+            requests,
+            indexes_by_key,
+            outcomes_by_key,
+            ordered,
+        )
 
         return EventBatchResult(
             results=tuple(result for result in ordered if result is not None),
@@ -154,26 +166,32 @@ class EventBatchExecutor:
         )
 
     def release_scope(self, scope: Hashable) -> None:
+        generation = self._scope_generations.get(scope, 0)
+        self._scope_generations[scope] = generation + 1
         for key in tuple(self._replay_results):
-            if key[0] == scope:
+            if key[0] == scope and key[1] == generation:
                 self._replay_results.pop(key, None)
+        for key in tuple(self._inflight):
+            if key[0] == scope and key[1] == generation:
+                self._inflight.pop(key, None)
 
     async def _execute_replayable(
         self,
         item: _BatchItem,
         *,
         scope: Hashable,
+        generation: int,
         enabled_names: Iterable[str] | None,
         context: EventExecutionContext | None,
         deadline: float,
         count_deadline_exceeded: bool,
     ) -> tuple[_ExecutionOutcome, bool]:
-        replay_key = (scope, item.request.id, item.key)
+        replay_key = (scope, generation, item.request.id, item.key)
         async with self._state_lock:
             replayed = self._replay_results.get(replay_key)
             if replayed is not None:
                 self._replay_results.move_to_end(replay_key)
-                return self._outcome_from_cached(replayed), True
+                return replayed, True
             task = self._inflight.get(replay_key)
             reused = task is not None
             if task is None:
@@ -192,7 +210,7 @@ class EventBatchExecutor:
 
     async def _run_and_store(
         self,
-        replay_key: tuple[Hashable, str, str],
+        replay_key: _ReplayKey,
         request: EventRequest,
         *,
         enabled_names: Iterable[str] | None,
@@ -213,47 +231,50 @@ class EventBatchExecutor:
             async with self._state_lock:
                 if self._inflight.get(replay_key) is asyncio.current_task():
                     self._inflight.pop(replay_key, None)
-                if outcome is not None:
-                    self._store_replay_locked(replay_key, outcome.result)
+                if (
+                    outcome is not None
+                    and self._scope_generations.get(replay_key[0], 0)
+                    == replay_key[1]
+                ):
+                    self._store_replay_locked(replay_key, outcome)
         assert outcome is not None
         return outcome
 
     async def _cache_ordered_results(
         self,
         scope: Hashable,
+        generation: int,
         requests: Sequence[EventRequest],
         indexes_by_key: dict[str, list[int]],
+        outcomes_by_key: dict[str, _ExecutionOutcome],
         ordered: list[EventResult | None],
     ) -> None:
         async with self._state_lock:
+            if self._scope_generations.get(scope, 0) != generation:
+                return
             for key, indexes in indexes_by_key.items():
+                outcome = outcomes_by_key[key]
                 for index in indexes:
                     result = ordered[index]
                     assert result is not None
                     self._store_replay_locked(
-                        (scope, requests[index].id, key),
-                        result,
+                        (scope, generation, requests[index].id, key),
+                        _ExecutionOutcome(
+                            result,
+                            batch_timed_out=outcome.batch_timed_out,
+                        ),
                     )
 
     def _store_replay_locked(
         self,
-        replay_key: tuple[Hashable, str, str],
-        result: EventResult,
+        replay_key: _ReplayKey,
+        outcome: _ExecutionOutcome,
     ) -> None:
-        self._replay_results[replay_key] = result
+        self._replay_results[replay_key] = outcome
         self._replay_results.move_to_end(replay_key)
         while len(self._replay_results) > self.max_replay_entries:
             self._replay_results.popitem(last=False)
 
-    def _outcome_from_cached(self, result: EventResult) -> _ExecutionOutcome:
-        return _ExecutionOutcome(
-            result,
-            batch_timed_out=(
-                result.status is EventStatus.TIMEOUT
-                and result.error == "batch deadline exceeded"
-            ),
-        )
-
     def _clone_result(
         self,
         result: EventResult,
@@ -262,13 +283,14 @@ class EventBatchExecutor:
         deduplicated_from: str | None,
     ) -> EventResult:
         payload = dict(result.payload)
-        if deduplicated_from is not None:
-            payload["deduplicated_from"] = deduplicated_from
+        if payload.get("event_id") == result.event_id:
+            payload["event_id"] = request.id
         return replace(
             result,
             event_id=request.id,
             raw_arguments=request.raw_arguments,
             payload=payload,
+            deduplicated_from=deduplicated_from,
         )
 
     async def _execute_one(
@@ -297,6 +319,30 @@ class EventBatchExecutor:
         timeout = remaining if event_timeout is None else min(remaining, event_timeout)
         acquired_locks: list[asyncio.Lock] = []
         semaphore_acquired = False
+        resources_released = False
+        kernel_task: asyncio.Task[EventResult] | None = None
+        sync_handler = (
+            definition is not None
+            and not inspect.iscoroutinefunction(definition.handler)
+        )
+
+        def release_resources() -> None:
+            nonlocal resources_released
+            if resources_released:
+                return
+            resources_released = True
+            if semaphore_acquired:
+                self._semaphore.release()
+            for lock in reversed(acquired_locks):
+                lock.release()
+
+        def release_sync_lease(task: asyncio.Task[EventResult]) -> None:
+            try:
+                task.exception()
+            except BaseException:
+                pass
+            release_resources()
+
         try:
             async with asyncio.timeout(timeout):
                 for key in sorted(set(definition.conflict_keys if definition else ())):
@@ -305,10 +351,17 @@ class EventBatchExecutor:
                     acquired_locks.append(lock)
                 await self._semaphore.acquire()
                 semaphore_acquired = True
-                result = await self.kernel.execute(
-                    request,
-                    enabled_names=enabled_names,
-                    context=context,
+                kernel_task = asyncio.create_task(
+                    self.kernel.execute(
+                        request,
+                        enabled_names=enabled_names,
+                        context=context,
+                    )
+                )
+                result = await (
+                    asyncio.shield(kernel_task)
+                    if sync_handler
+                    else kernel_task
                 )
                 return _ExecutionOutcome(result)
         except TimeoutError:
@@ -336,10 +389,10 @@ class EventBatchExecutor:
                 )
             )
         finally:
-            if semaphore_acquired:
-                self._semaphore.release()
-            for lock in reversed(acquired_locks):
-                lock.release()
+            if sync_handler and kernel_task is not None and not kernel_task.done():
+                kernel_task.add_done_callback(release_sync_lease)
+            else:
+                release_resources()
 
     def _failure_result(
         self,
@@ -376,32 +429,46 @@ class EventBatchExecutor:
         request: EventRequest,
         definition: EventDefinition | None,
     ) -> str:
-        arguments_value = self._safe_arguments(request.arguments)
-        key_kind = "original"
-        if definition is not None and definition.normalizer is not None:
+        arguments = self._strict_json_arguments(request.arguments)
+        key_kind = "original-json"
+        if arguments is None:
+            key_kind = "original-invalid"
+            arguments = request.raw_arguments
+        elif definition is not None and definition.normalizer is not None:
             try:
-                normalized = definition.normalizer(arguments_value)
+                normalized = definition.normalizer(
+                    json.loads(arguments)
+                )
                 if not isinstance(normalized, dict):
                     raise TypeError("normalizer returned non-object")
-                arguments_value = self._safe_arguments(normalized)
-                key_kind = "normalized"
+                normalized_arguments = self._strict_json_arguments(normalized)
+                if normalized_arguments is None:
+                    raise TypeError("normalizer returned non-JSON object")
+                arguments = normalized_arguments
+                key_kind = "normalized-json"
             except Exception:
-                arguments_value = self._safe_arguments(request.arguments)
+                key_kind = "normalizer-invalid"
+        return json.dumps(
+            [request.name, request.source.value, key_kind, arguments],
+            ensure_ascii=False,
+            separators=(",", ":"),
+        )
+
+    def _strict_json_arguments(self, arguments: dict[str, Any]) -> str | None:
         try:
-            arguments = json.dumps(
-                arguments_value,
+            serialized = json.dumps(
+                arguments,
                 ensure_ascii=False,
                 sort_keys=True,
                 separators=(",", ":"),
                 allow_nan=False,
             )
+            copied = json.loads(serialized)
         except (TypeError, ValueError):
-            arguments = request.raw_arguments
-        return json.dumps(
-            [request.name, request.source.value, key_kind, arguments],
-            ensure_ascii=False,
-            separators=(",", ":"),
-        )
+            return None
+        if not isinstance(copied, dict) or copied != arguments:
+            return None
+        return serialized
 
     def _safe_arguments(self, arguments: dict[str, Any]) -> dict[str, Any]:
         try:

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

@@ -121,6 +121,7 @@ class EventResult:
     conflict_keys: tuple[str, ...] = ()
     timeout_seconds: float | None = None
     terminal: bool = False
+    deduplicated_from: str | None = None
 
 
 @dataclass(frozen=True)

+ 27 - 9
src/agent_lab/application/runtime.py

@@ -15,6 +15,7 @@ from agent_lab.application.events import (
     EventBatchExecutor,
     EventBatchResult,
     EventKernel,
+    EventResult,
     EventSource,
 )
 from agent_lab.application.queues import RuntimeQueues
@@ -368,6 +369,10 @@ class DebugRuntime:
                     await queues.output.put(
                         {"type": "tool_result", "message": reply.model_dump()}
                     )
+                assert batch_result is not None
+                logical_result_count = len(
+                    self._logical_batch_results(batch_result)
+                )
                 if request.tool_invocation_mode == "chat_agent_tools":
                     await self._audit(
                         queues,
@@ -377,7 +382,7 @@ class DebugRuntime:
                         tool_invocation_mode=request.tool_invocation_mode,
                         event_source="provider_resolved",
                         events=self._event_snapshots(events),
-                        result_count=len(tool_replies),
+                        result_count=logical_result_count,
                     )
                 else:
                     await self._audit(
@@ -386,9 +391,8 @@ class DebugRuntime:
                         turn_started_at=turn_started_at,
                         round_index=round_index,
                         event_names=[event.name for event in events],
-                        result_count=len(tool_replies),
+                        result_count=logical_result_count,
                     )
-                assert batch_result is not None
                 await self._audit_batch_result(
                     queues,
                     batch_result,
@@ -769,6 +773,10 @@ class DebugRuntime:
                     await queues.output.put(
                         {"type": "tool_result", "message": reply.model_dump()}
                     )
+                assert batch_result is not None
+                logical_result_count = len(
+                    self._logical_batch_results(batch_result)
+                )
                 if request.tool_invocation_mode == "chat_agent_tools":
                     await self._audit(
                         queues,
@@ -780,7 +788,7 @@ class DebugRuntime:
                         tool_invocation_mode=request.tool_invocation_mode,
                         event_source="provider_resolved",
                         events=self._event_snapshots(events),
-                        result_count=len(tool_replies),
+                        result_count=logical_result_count,
                     )
                 else:
                     await self._audit(
@@ -791,9 +799,8 @@ class DebugRuntime:
                         turn_index=turn_index,
                         round_index=round_index,
                         event_names=[event.name for event in events],
-                        result_count=len(tool_replies),
+                        result_count=logical_result_count,
                     )
-                assert batch_result is not None
                 await self._audit_batch_result(
                     queues,
                     batch_result,
@@ -1100,6 +1107,7 @@ class DebugRuntime:
         turn_index: int | None = None,
         round_index: int | None = None,
     ) -> None:
+        results = self._logical_batch_results(batch)
         await self._audit(
             queues,
             "event_batch_results",
@@ -1116,10 +1124,10 @@ class DebugRuntime:
                     "result_policy": result.result_policy.value,
                     "terminal": result.terminal,
                 }
-                for result in batch.results
+                for result in results
             ],
             timeout_count=sum(
-                result.status.value == "timeout" for result in batch.results
+                result.status.value == "timeout" for result in results
             ),
             coalesced_count=batch.coalesced_count,
             replayed_count=batch.replayed_count,
@@ -1168,12 +1176,22 @@ class DebugRuntime:
                 round_index=round_index,
                 event_ids=[
                     result.event_id
-                    for result in batch.results
+                    for result in self._logical_batch_results(batch)
                     if result.terminal and result.status.value == "success"
                 ],
             )
         return decision
 
+    def _logical_batch_results(
+        self,
+        batch: EventBatchResult,
+    ) -> tuple[EventResult, ...]:
+        return tuple(
+            result
+            for result in batch.results
+            if result.deduplicated_from is None
+        )
+
     def _chat_messages_for_round(
         self,
         messages: list[ChatMessage],

+ 14 - 3
src/agent_lab/application/tools.py

@@ -233,7 +233,7 @@ class ToolRegistry:
         templates: list[str] = []
         llm_follow_up = False
         terminate = False
-        for result in batch.results:
+        for result in self._logical_results(batch):
             if result.status is not EventStatus.SUCCESS:
                 llm_follow_up = True
                 continue
@@ -275,7 +275,8 @@ class ToolRegistry:
         self,
         batch: EventBatchResult,
     ) -> ChatMessage | None:
-        if not batch.results:
+        results = self._logical_results(batch)
+        if not results:
             return None
         return ChatMessage(
             role="user",
@@ -283,12 +284,22 @@ class ToolRegistry:
                 "EventAgent results:\n"
                 + "\n".join(
                     json.dumps(self.tool_payload(result), ensure_ascii=False)
-                    for result in batch.results
+                    for result in results
                 )
             ),
             name="event_agent",
         )
 
+    def _logical_results(
+        self,
+        batch: EventBatchResult,
+    ) -> tuple[EventResult, ...]:
+        return tuple(
+            result
+            for result in batch.results
+            if result.deduplicated_from is None
+        )
+
     def _to_event_definition(self, definition: ToolDefinition) -> EventDefinition:
         def resolver(
             request: EventRequest,

+ 65 - 0
tests/test_debug_runtime.py

@@ -2641,6 +2641,71 @@ async def test_tool_only_terminate_emits_plugin_farewell_once(mode: str):
     assert outputs[-1] == {"type": "done"}
 
 
+@pytest.mark.asyncio
+@pytest.mark.parametrize("mode", ["dual_agent", "chat_agent_tools"])
+async def test_coalesced_terminate_has_one_farewell_and_one_audit_result(mode: str):
+    events = [
+        ToolCallEvent(
+            id=event_id,
+            name="session.terminate",
+            arguments={},
+            raw_arguments="{}",
+        )
+        for event_id in ("terminate-primary", "terminate-duplicate")
+    ]
+    items = [
+        (
+            StreamItem.text_event(event)
+            if mode == "dual_agent"
+            else StreamItem.provider_tool_call(event)
+        )
+        for event in events
+    ]
+    client = ScriptedChatClient([items])
+
+    outputs = await _collect_outputs(
+        DebugRuntime(client).run(
+            _policy_request(enabled_tools=["session.terminate"], mode=mode)
+        )
+    )
+
+    assert [
+        message["content"]
+        for message in outputs
+        if message["type"] == "message_delta"
+    ] == ["Goodbye."]
+    assert [
+        message["message"]["tool_call_id"]
+        for message in outputs
+        if message["type"] == "tool_result"
+    ] == [event.id for event in events]
+    batch_audit = next(
+        message
+        for message in outputs
+        if message.get("event") == "event_batch_results"
+    )
+    terminal_audit = next(
+        message
+        for message in outputs
+        if message.get("event") == "terminal_completed"
+    )
+    completion_audit = next(
+        message
+        for message in outputs
+        if message.get("event")
+        == (
+            "provider_tools_completed"
+            if mode == "chat_agent_tools"
+            else "event_agent_completed"
+        )
+    )
+    assert [
+        result["event_id"] for result in batch_audit["details"]["results"]
+    ] == ["terminate-primary"]
+    assert completion_audit["details"]["result_count"] == 1
+    assert terminal_audit["details"]["event_ids"] == ["terminate-primary"]
+
+
 @pytest.mark.asyncio
 async def test_tool_only_terminate_ends_reusable_session_after_plugin_farewell():
     client = ScriptedChatClient(

+ 520 - 6
tests/test_event_batch.py

@@ -1,6 +1,7 @@
 from __future__ import annotations
 
 import asyncio
+import json
 import threading
 from collections.abc import Callable
 from typing import Any
@@ -15,6 +16,7 @@ from agent_lab.application.events import (
     EventRequest,
     EventSource,
     EventStatus,
+    ResultPolicy,
 )
 from agent_lab.application.tools import ToolDefinition, ToolRegistry
 from agent_lab.domain.events import ToolCallEvent
@@ -250,7 +252,10 @@ async def test_batch_executor_coalesces_exact_duplicates_inside_one_batch():
 
     assert calls == 1
     assert batch.coalesced_count == 1
-    assert batch.results[0] == batch.results[1]
+    assert batch.results[0].event_id == batch.results[1].event_id == "same"
+    assert batch.results[0].payload == batch.results[1].payload
+    assert batch.results[0].deduplicated_from is None
+    assert batch.results[1].deduplicated_from == "same"
 
 
 @pytest.mark.asyncio
@@ -403,10 +408,64 @@ async def test_batch_executor_coalesces_normalized_provider_arguments_across_ids
         "float-id",
     ]
     assert batch.results[0].payload == {"value": 30}
-    assert batch.results[1].payload == {
-        "value": 30,
-        "deduplicated_from": "integer-id",
-    }
+    assert batch.results[1].payload == {"value": 30}
+    assert getattr(batch.results[0], "deduplicated_from", None) is None
+    assert getattr(batch.results[1], "deduplicated_from", None) == "integer-id"
+
+
+@pytest.mark.asyncio
+async def test_coalesced_clone_rewrites_only_correlated_payload_event_id():
+    correlated = _executor(
+        [
+            _definition(
+                "example.correlated",
+                lambda request: {"event_id": request.id, "value": "same"},
+            )
+        ]
+    )
+    correlated_batch = await correlated.execute(
+        [
+            EventRequest(
+                id="primary",
+                name="example.correlated",
+                arguments={},
+                raw_arguments='{"call":"primary"}',
+            ),
+            EventRequest(
+                id="duplicate",
+                name="example.correlated",
+                arguments={},
+                raw_arguments='{"call":"duplicate"}',
+            ),
+        ],
+        scope="correlated-payload",
+    )
+
+    primary, duplicate = correlated_batch.results
+    assert primary.event_id == "primary"
+    assert duplicate.event_id == "duplicate"
+    assert duplicate.raw_arguments == '{"call":"duplicate"}'
+    assert duplicate.payload == {"event_id": "duplicate", "value": "same"}
+    assert getattr(duplicate, "deduplicated_from", None) == "primary"
+
+    domain = _executor(
+        [
+            _definition(
+                "example.domain",
+                lambda request: {"event_id": "domain-object"},
+            )
+        ]
+    )
+    domain_batch = await domain.execute(
+        [
+            _request("primary", "example.domain"),
+            _request("duplicate", "example.domain"),
+        ],
+        scope="domain-payload",
+    )
+
+    assert domain_batch.results[1].payload == {"event_id": "domain-object"}
+    assert getattr(domain_batch.results[1], "deduplicated_from", None) == "primary"
 
 
 @pytest.mark.asyncio
@@ -417,7 +476,7 @@ async def test_batch_executor_tool_replies_keep_each_coalesced_call_id():
                 name="example.tool",
                 description="Execute a coalesced tool.",
                 parameters={"type": "object"},
-                handler=lambda event: {"handled_by": event.id},
+                handler=lambda event: {"event_id": event.id},
             )
         ]
     )
@@ -440,12 +499,53 @@ async def test_batch_executor_tool_replies_keep_each_coalesced_call_id():
     replies = registry.tool_replies(events, batch)
 
     assert [reply.tool_call_id for reply in replies] == ["first-id", "second-id"]
+    assert [json.loads(reply.content) for reply in replies] == [
+        {"event_id": "first-id"},
+        {"event_id": "second-id"},
+    ]
     assert [result.event_id for result in batch.results] == [
         "first-id",
         "second-id",
     ]
 
 
+@pytest.mark.asyncio
+async def test_batch_policy_and_compact_summary_skip_logical_duplicates():
+    registry = ToolRegistry(
+        [
+            ToolDefinition(
+                name="example.template",
+                description="Render one template.",
+                parameters={"type": "object"},
+                handler=lambda event: {"event_id": event.id},
+                result_policy=ResultPolicy.TEMPLATE_FOLLOW_UP,
+                result_message_factory=lambda result: "Completed once.",
+            )
+        ]
+    )
+    executor = _executor_type()(EventKernel(registry.event_registry))
+    events = [
+        ToolCallEvent(
+            id=event_id,
+            name="example.template",
+            arguments={},
+            raw_arguments="{}",
+        )
+        for event_id in ("primary", "duplicate")
+    ]
+    batch = await executor.execute(
+        [registry.event_request(event) for event in events],
+        scope="policy-dedupe",
+    )
+
+    decision = registry.batch_decision(batch)
+    summary = registry.compact_results_message(batch)
+
+    assert decision.template_messages == ("Completed once.",)
+    assert summary is not None
+    assert len(summary.content.splitlines()) == 2
+
+
 @pytest.mark.asyncio
 async def test_batch_key_falls_back_to_original_arguments_when_normalizer_fails():
     calls = 0
@@ -490,6 +590,135 @@ async def test_batch_key_falls_back_to_original_arguments_when_normalizer_fails(
     assert calls == 1
 
 
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("invalid_arguments", "raw_arguments"),
+    [
+        ({"value": float("nan")}, '{"value":NaN}'),
+        ({"value": object()}, '{"value":"object"}'),
+    ],
+)
+async def test_invalid_json_arguments_do_not_coalesce_with_valid_empty_object(
+    invalid_arguments: dict[str, Any],
+    raw_arguments: str,
+):
+    calls = 0
+
+    def handler(request: EventRequest) -> dict[str, Any]:
+        nonlocal calls
+        calls += 1
+        return {"event_id": request.id}
+
+    executor = _executor([_definition("example.strict", handler)])
+    batch = await executor.execute(
+        [
+            EventRequest(
+                id="valid",
+                name="example.strict",
+                arguments={},
+                source=EventSource.PROVIDER_RESOLVED,
+                raw_arguments="{}",
+            ),
+            EventRequest(
+                id="invalid",
+                name="example.strict",
+                arguments=invalid_arguments,
+                source=EventSource.PROVIDER_RESOLVED,
+                raw_arguments=raw_arguments,
+            ),
+        ],
+        scope="strict-json-batch",
+    )
+
+    assert batch.coalesced_count == 0
+    assert [result.status for result in batch.results] == [
+        EventStatus.SUCCESS,
+        EventStatus.INVALID_ARGUMENTS,
+    ]
+    assert calls == 1
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize(
+    ("invalid_arguments", "raw_arguments"),
+    [
+        ({"value": float("nan")}, '{"value":NaN}'),
+        ({"value": object()}, '{"value":"object"}'),
+    ],
+)
+async def test_invalid_json_arguments_do_not_replay_valid_success(
+    invalid_arguments: dict[str, Any],
+    raw_arguments: str,
+):
+    executor = _executor(
+        [_definition("example.strict", lambda request: {"event_id": request.id})]
+    )
+    valid = EventRequest(
+        id="same-id",
+        name="example.strict",
+        arguments={},
+        source=EventSource.PROVIDER_RESOLVED,
+        raw_arguments="{}",
+    )
+    invalid = EventRequest(
+        id="same-id",
+        name="example.strict",
+        arguments=invalid_arguments,
+        source=EventSource.PROVIDER_RESOLVED,
+        raw_arguments=raw_arguments,
+    )
+
+    first = await executor.execute([valid], scope="strict-json-replay")
+    second = await executor.execute([invalid], scope="strict-json-replay")
+
+    assert first.results[0].status is EventStatus.SUCCESS
+    assert second.replayed_count == 0
+    assert second.results[0].status is EventStatus.INVALID_ARGUMENTS
+
+
+@pytest.mark.asyncio
+async def test_non_json_normalizer_output_has_distinct_canonical_key():
+    calls = 0
+
+    def normalize(arguments: dict[str, Any]) -> dict[str, Any]:
+        if arguments.get("invalid"):
+            return {"value": object()}
+        return {}
+
+    def handler(request: EventRequest) -> dict[str, Any]:
+        nonlocal calls
+        calls += 1
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [_definition("example.normalized", handler, normalizer=normalize)]
+    )
+    batch = await executor.execute(
+        [
+            EventRequest(
+                id="valid",
+                name="example.normalized",
+                arguments={},
+                source=EventSource.PROVIDER_RESOLVED,
+            ),
+            EventRequest(
+                id="invalid",
+                name="example.normalized",
+                arguments={"invalid": True},
+                source=EventSource.PROVIDER_RESOLVED,
+            ),
+        ],
+        scope="normalizer-invalid-json",
+    )
+
+    assert batch.coalesced_count == 0
+    assert [result.status for result in batch.results] == [
+        EventStatus.SUCCESS,
+        EventStatus.INVALID_ARGUMENTS,
+    ]
+    assert calls == 1
+
+
 @pytest.mark.asyncio
 async def test_concurrent_replay_waiters_share_execution_and_cancel_independently():
     calls = 0
@@ -539,6 +768,92 @@ async def test_release_scope_removes_replay_entries_for_reuse():
     assert batch.replayed_count == 0
 
 
+@pytest.mark.asyncio
+async def test_release_scope_starts_new_generation_before_old_task_finishes():
+    calls = 0
+    started = [asyncio.Event(), asyncio.Event()]
+    releases = [asyncio.Event(), asyncio.Event()]
+
+    async def handler(request: EventRequest) -> dict[str, Any]:
+        nonlocal calls
+        call_index = calls
+        calls += 1
+        started[call_index].set()
+        await releases[call_index].wait()
+        return {"call": call_index + 1, "event_id": request.id}
+
+    executor = _executor([_definition("example.generated", handler)])
+    request = _request("same-id", "example.generated")
+    first = asyncio.create_task(
+        executor.execute([request], scope="generation-scope")
+    )
+    await asyncio.wait_for(started[0].wait(), timeout=0.2)
+
+    executor.release_scope("generation-scope")
+    second = asyncio.create_task(
+        executor.execute([request], scope="generation-scope")
+    )
+    await asyncio.sleep(0.02)
+    new_generation_started = started[1].is_set()
+
+    releases[0].set()
+    releases[1].set()
+    first_batch, second_batch = await asyncio.gather(first, second)
+
+    assert new_generation_started is True
+    assert first_batch.results[0].payload["call"] == 1
+    assert second_batch.results[0].payload["call"] == 2
+
+
+@pytest.mark.asyncio
+async def test_old_generation_completion_cannot_replace_new_inflight_or_cache():
+    calls = 0
+    started = [asyncio.Event(), asyncio.Event(), asyncio.Event()]
+    releases = [asyncio.Event(), asyncio.Event(), asyncio.Event()]
+
+    async def handler(request: EventRequest) -> dict[str, Any]:
+        nonlocal calls
+        call_index = calls
+        calls += 1
+        started[call_index].set()
+        await releases[call_index].wait()
+        return {"call": call_index + 1, "event_id": request.id}
+
+    executor = _executor([_definition("example.generated", handler)])
+    request = _request("same-id", "example.generated")
+    old = asyncio.create_task(executor.execute([request], scope="reused-scope"))
+    await asyncio.wait_for(started[0].wait(), timeout=0.2)
+
+    executor.release_scope("reused-scope")
+    current = asyncio.create_task(
+        executor.execute([request], scope="reused-scope")
+    )
+    await asyncio.sleep(0.02)
+    new_generation_started = started[1].is_set()
+    if not new_generation_started:
+        releases[0].set()
+        await asyncio.gather(old, current)
+    assert new_generation_started is True
+
+    releases[0].set()
+    old_batch = await asyncio.wait_for(old, timeout=0.2)
+    waiter = asyncio.create_task(
+        executor.execute([request], scope="reused-scope")
+    )
+    await asyncio.sleep(0.02)
+    assert calls == 2
+
+    releases[1].set()
+    current_batch, waiter_batch = await asyncio.gather(current, waiter)
+    replayed = await executor.execute([request], scope="reused-scope")
+
+    assert old_batch.results[0].payload["call"] == 1
+    assert current_batch.results[0].payload["call"] == 2
+    assert waiter_batch.results[0].payload["call"] == 2
+    assert replayed.replayed_count == 1
+    assert replayed.results[0].payload["call"] == 2
+
+
 @pytest.mark.asyncio
 async def test_replay_cache_stays_bounded_across_250_scopes():
     executor = _executor(
@@ -623,6 +938,49 @@ async def test_terminal_grace_is_bounded_by_terminal_event_timeout():
     assert batch.deadline_exceeded is False
 
 
+@pytest.mark.asyncio
+async def test_replayed_terminal_grace_timeout_preserves_deadline_flag():
+    async def blocked_terminal(request: EventRequest) -> dict[str, Any]:
+        await asyncio.Event().wait()
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [_definition("example.terminate", blocked_terminal, terminal=True)],
+        terminal_grace_seconds=0.01,
+    )
+    request = _request("terminal", "example.terminate")
+
+    first = await executor.execute([request], scope="terminal-replay")
+    replayed = await executor.execute([request], scope="terminal-replay")
+
+    assert first.results[0].status is EventStatus.TIMEOUT
+    assert first.results[0].error == "batch deadline exceeded"
+    assert first.deadline_exceeded is False
+    assert replayed.replayed_count == 1
+    assert replayed.deadline_exceeded is False
+
+
+@pytest.mark.asyncio
+async def test_replayed_sibling_batch_timeout_preserves_deadline_flag():
+    async def blocked(request: EventRequest) -> dict[str, Any]:
+        await asyncio.Event().wait()
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [_definition("example.blocked", blocked)],
+        batch_timeout_seconds=0.01,
+    )
+    request = _request("sibling", "example.blocked")
+
+    first = await executor.execute([request], scope="sibling-replay")
+    replayed = await executor.execute([request], scope="sibling-replay")
+
+    assert first.results[0].status is EventStatus.TIMEOUT
+    assert first.deadline_exceeded is True
+    assert replayed.replayed_count == 1
+    assert replayed.deadline_exceeded is True
+
+
 @pytest.mark.asyncio
 async def test_blocking_sync_handler_timeout_returns_before_thread_finishes():
     started = threading.Event()
@@ -660,3 +1018,159 @@ async def test_blocking_sync_handler_timeout_returns_before_thread_finishes():
     assert batch.results[0].status is EventStatus.TIMEOUT
     assert elapsed < 0.1
     assert finished.is_set()
+
+
+@pytest.mark.asyncio
+async def test_timed_out_sync_handler_keeps_conflict_lock_until_thread_finishes():
+    worker_started = threading.Event()
+    worker_release = threading.Event()
+    follower_started = asyncio.Event()
+
+    def blocking_handler(request: EventRequest) -> dict[str, Any]:
+        worker_started.set()
+        worker_release.wait(timeout=1)
+        return {"event_id": request.id}
+
+    async def follower_handler(request: EventRequest) -> dict[str, Any]:
+        follower_started.set()
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [
+            _definition(
+                "example.blocking",
+                blocking_handler,
+                conflict_keys=("shared",),
+                timeout_seconds=0.02,
+            ),
+            _definition(
+                "example.follower",
+                follower_handler,
+                conflict_keys=("shared",),
+                timeout_seconds=0.2,
+            ),
+        ]
+    )
+
+    timed_out = await executor.execute(
+        [_request("blocking", "example.blocking")],
+        scope="sync-conflict-blocking",
+    )
+    follower = asyncio.create_task(
+        executor.execute(
+            [_request("follower", "example.follower")],
+            scope="sync-conflict-follower",
+        )
+    )
+    await asyncio.sleep(0.03)
+    entered_before_worker_finished = follower_started.is_set()
+
+    worker_release.set()
+    follower_batch = await asyncio.wait_for(follower, timeout=0.3)
+
+    assert worker_started.is_set()
+    assert timed_out.results[0].status is EventStatus.TIMEOUT
+    assert entered_before_worker_finished is False
+    assert follower_batch.results[0].status is EventStatus.SUCCESS
+
+
+@pytest.mark.asyncio
+async def test_timed_out_sync_handler_keeps_parallel_slot_until_thread_finishes():
+    worker_started = threading.Event()
+    worker_release = threading.Event()
+    follower_started = asyncio.Event()
+
+    def blocking_handler(request: EventRequest) -> dict[str, Any]:
+        worker_started.set()
+        worker_release.wait(timeout=1)
+        return {"event_id": request.id}
+
+    async def follower_handler(request: EventRequest) -> dict[str, Any]:
+        follower_started.set()
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [
+            _definition(
+                "example.blocking",
+                blocking_handler,
+                timeout_seconds=0.02,
+            ),
+            _definition(
+                "example.follower",
+                follower_handler,
+                timeout_seconds=0.2,
+            ),
+        ],
+        max_parallel_events=1,
+    )
+
+    timed_out = await executor.execute(
+        [_request("blocking", "example.blocking")],
+        scope="sync-slot-blocking",
+    )
+    follower = asyncio.create_task(
+        executor.execute(
+            [_request("follower", "example.follower")],
+            scope="sync-slot-follower",
+        )
+    )
+    await asyncio.sleep(0.03)
+    entered_before_worker_finished = follower_started.is_set()
+
+    worker_release.set()
+    follower_batch = await asyncio.wait_for(follower, timeout=0.3)
+
+    assert worker_started.is_set()
+    assert timed_out.results[0].status is EventStatus.TIMEOUT
+    assert entered_before_worker_finished is False
+    assert follower_batch.results[0].status is EventStatus.SUCCESS
+
+
+@pytest.mark.asyncio
+async def test_timed_out_async_handler_releases_resources_after_cancellation():
+    cancelled = asyncio.Event()
+    follower_started = asyncio.Event()
+
+    async def blocked_handler(request: EventRequest) -> dict[str, Any]:
+        try:
+            await asyncio.Event().wait()
+        finally:
+            cancelled.set()
+        return {"event_id": request.id}
+
+    async def follower_handler(request: EventRequest) -> dict[str, Any]:
+        follower_started.set()
+        return {"event_id": request.id}
+
+    executor = _executor(
+        [
+            _definition(
+                "example.blocked",
+                blocked_handler,
+                conflict_keys=("shared",),
+                timeout_seconds=0.01,
+            ),
+            _definition(
+                "example.follower",
+                follower_handler,
+                conflict_keys=("shared",),
+                timeout_seconds=0.2,
+            ),
+        ],
+        max_parallel_events=1,
+    )
+
+    timed_out = await executor.execute(
+        [_request("blocked", "example.blocked")],
+        scope="async-release-blocked",
+    )
+    follower = await executor.execute(
+        [_request("follower", "example.follower")],
+        scope="async-release-follower",
+    )
+
+    assert timed_out.results[0].status is EventStatus.TIMEOUT
+    assert cancelled.is_set()
+    assert follower_started.is_set()
+    assert follower.results[0].status is EventStatus.SUCCESS