fix(delegation): late-attached children take the parent's soft/hard stop kind; dedupe the fallback replay
_attach_child now mirrors a pending parent stop with the same split interrupt() uses for its own fan-out (hard -> hard_interrupt, soft -> interrupt), so a redirect is not turned into a cancel on a child that was attached late. _restore_parent_cancellation collapses to re-attaching the rejected unit's children: the replay is the attach step's job now. Test fixture: _Batch gained origin_session_history_delivery on main after the salvaged PR was written. Co-authored-by: illidan <noequal666@gmail.com>
This commit is contained in:
@@ -101,7 +101,7 @@ def _batch(parent, *children):
|
||||
creds={"model": children[0].model}, context=None, top_role="leaf", max_children=len(children),
|
||||
live_deleg_id=None, live_writers=[], live_paths=[], origin_wake_sid="",
|
||||
origin_ui_session_id="", origin_owner_transport=None,
|
||||
origin_owner_session_record=None, overall_start=time.monotonic(),
|
||||
origin_owner_session_record=None, origin_session_history_delivery=False, overall_start=time.monotonic(),
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -69,8 +69,16 @@ def _attach_child(parent_agent: Any, child: Any) -> None:
|
||||
so the whole spawn tree dies with its parent."""
|
||||
if hasattr(parent_agent, "_active_children"):
|
||||
_with_children_lock(parent_agent, "append", child)
|
||||
if getattr(parent_agent, "_interrupt_requested", False) is True:
|
||||
_signal_child_stop(child, getattr(parent_agent, "_interrupt_message", None) or "parent agent interrupted")
|
||||
if getattr(parent_agent, "_interrupt_requested", False) is not True:
|
||||
return
|
||||
# Same soft/hard split as ``interrupt()``'s own fan-out: a hard stop cancels, a soft one redirects.
|
||||
message = getattr(parent_agent, "_interrupt_message", None)
|
||||
hard = getattr(parent_agent, "_hard_interrupt_requested", None)
|
||||
if hard is None or hard.is_set():
|
||||
_signal_child_stop(child, message or "parent agent interrupted")
|
||||
else:
|
||||
with _quiet("Failed to propagate interrupt to late child: %s"):
|
||||
child.interrupt(message)
|
||||
|
||||
def _detach_child(parent_agent: Any, child: Any) -> None:
|
||||
"""Remove the child from parent interrupt propagation (no-op if absent)."""
|
||||
|
||||
@@ -375,19 +375,10 @@ def _dispatch_unit(unit: _Batch, unit_id: Optional[str], slot_key: Optional[str]
|
||||
)
|
||||
|
||||
def _restore_parent_cancellation(unit: _Batch) -> None:
|
||||
# Rejected children remain owned by the parent. Attach before replaying a
|
||||
# cancellation that may have arrived while async admission had them detached.
|
||||
parent = unit.parent_agent
|
||||
"""Rejected children stay owned by the parent: re-attach them (``_attach_child`` replays a stop that
|
||||
arrived while async admission had them detached)."""
|
||||
for _, _, child in unit.children:
|
||||
_attach_child(parent, child)
|
||||
if getattr(parent, "_interrupt_requested", False) is True:
|
||||
hard_stop = getattr(parent, "_hard_interrupt_requested", None)
|
||||
for _, _, child in unit.children:
|
||||
if hard_stop is not None and hard_stop.is_set():
|
||||
_signal_child_stop(child, getattr(parent, "_interrupt_message", None))
|
||||
else:
|
||||
with _quiet("Failed to propagate interrupt to fallback child: %s"):
|
||||
child.interrupt(getattr(parent, "_interrupt_message", None))
|
||||
_attach_child(unit.parent_agent, child)
|
||||
|
||||
def _dispatch_background(batch: _Batch) -> str:
|
||||
"""Dispatch the call as independent async units (see ``_units_of``) and return the tool result JSON. Every unit
|
||||
|
||||
Reference in New Issue
Block a user