From c96568f66ca49d27beec4545bee9740b09d64018 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 01:45:40 -0700 Subject: [PATCH] perf(delegation): finished delegate children no longer pin their transcripts in the parent heap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A parent that fanned out 1,320 subagents over 13h reached 2.6 GB RSS (1.9 GB anonymous heap). Every closed child AIAgent stayed reachable and still owned a copy of its full message history. gc.get_referrers on a finished child (30-child fan-out bench, evals/fanout_resource_bench.py) showed two retainers: 1. bind_subagent_parent() stored the agent strongly in the `hermes_subagent_lifecycle_parent` ContextVar. Each child binds ITSELF for its own turn, and every asyncio Handle/Future scheduled during that turn (LSP reader loops, kernel pipe transports) snapshots the Context — 56 live Contexts held 14 finished children after the bench. The ContextVar now holds a weakref (non-weakrefable doubles fall back to a closure); get_active_subagent_parent() dereferences it. 2. AIAgent.close() cleared _session_messages but not the _db_flush_scan_prefix snapshot (a `messages[:]` shallow copy taken on every successful DB flush) nor _streamed_assistant_text_parts, so the agent — kept alive by (1) — retained every message dict. close() now drops both. The delegate_task result entry never carried `messages`; a pin test confirms the per-child result JSON is unchanged. Bench (30 children / 10 worktrees, ~100 KB final replies so retention is visible): post-fan-out live child AIAgents 14 -> 0; RSS after fan-out 636 MB -> 556 MB. With the harness' tiny default replies both runs sit at ~192-194 MB (the children's transcripts were never the dominant cost there; the leaked objects were). --- agent/subagent_lifecycle.py | 19 ++- run_agent.py | 8 + .../test_delegate_child_transcript_release.py | 145 ++++++++++++++++++ 3 files changed, 169 insertions(+), 3 deletions(-) create mode 100644 tests/tools/test_delegate_child_transcript_release.py diff --git a/agent/subagent_lifecycle.py b/agent/subagent_lifecycle.py index 319e85a784..f7ba0146c9 100644 --- a/agent/subagent_lifecycle.py +++ b/agent/subagent_lifecycle.py @@ -17,6 +17,7 @@ import math import secrets import threading import time +import weakref from contextlib import contextmanager from concurrent.futures import Future, TimeoutError from typing import Any, Callable, Mapping, Optional @@ -171,8 +172,19 @@ _ACTIVE_PARENT_AGENT: contextvars.ContextVar[Any] = contextvars.ContextVar( @contextmanager def bind_subagent_parent(parent_agent: Any): - """Bind the host-owned parent for the current agent turn.""" - token = _ACTIVE_PARENT_AGENT.set(parent_agent) + """Bind the host-owned parent for the current agent turn. + + Stored as a weakref: every asyncio Handle/Future scheduled from the turn + (LSP reader loops, kernel pipes, ...) snapshots the Context, and those + snapshots outlive the turn. A strong ref there pinned finished delegate + children — each of which binds itself here for its own turn — in the + parent process heap for the life of the background loop. + """ + try: + ref = weakref.ref(parent_agent) + except TypeError: + ref = lambda: parent_agent # noqa: E731 — non-weakrefable test doubles + token = _ACTIVE_PARENT_AGENT.set(ref) try: yield finally: @@ -181,7 +193,8 @@ def bind_subagent_parent(parent_agent: Any): def get_active_subagent_parent() -> Any: """Return the parent bound to this execution context, if any.""" - return _ACTIVE_PARENT_AGENT.get() + ref = _ACTIVE_PARENT_AGENT.get() + return ref() if ref is not None else None class SubagentLifecycleService: diff --git a/run_agent.py b/run_agent.py index 8796236a9e..a9eeb2ad9a 100644 --- a/run_agent.py +++ b/run_agent.py @@ -5176,6 +5176,14 @@ class AIAgent: # still holds the closed agent (e.g. a draining background task). try: self._session_messages = [] + # Shadow copies of the same transcript: the DB-flush settled-prefix + # snapshot (a shallow copy of the whole list, see + # _flush_session_to_db) and the streamed-text accumulator. On a + # closed delegate child these were the only remaining owners of + # every message dict, so a retained child kept its full history + # alive in the parent's heap. + self._db_flush_scan_prefix = None + self._streamed_assistant_text_parts = [] except Exception: pass diff --git a/tests/tools/test_delegate_child_transcript_release.py b/tests/tools/test_delegate_child_transcript_release.py new file mode 100644 index 0000000000..de4e274b5a --- /dev/null +++ b/tests/tools/test_delegate_child_transcript_release.py @@ -0,0 +1,145 @@ +"""Finished delegate children must not pin their transcripts in the parent heap. + +Profiled parent (1,320 children over 13h) reached 2.6 GB RSS: every closed +child AIAgent stayed reachable, and each still owned a shallow copy of its +full message list. Two retainers were proven with ``gc.get_referrers``: + +1. ``AIAgent.close()`` cleared ``_session_messages`` but not the + ``_db_flush_scan_prefix`` snapshot (``messages[:]``) or the streamed-text + accumulator, so every message dict stayed alive through the agent. +2. ``bind_subagent_parent`` stored the agent (each child binds ITSELF for + its own turn) strongly in a ContextVar; asyncio Handles/Futures scheduled + during the turn snapshot that Context and live as long as the background + LSP / kernel loops do, so the child object itself was never collected. +""" + +from __future__ import annotations + +import gc +import json +import weakref +from unittest.mock import MagicMock + +from agent.subagent_lifecycle import ( + _ACTIVE_PARENT_AGENT, + bind_subagent_parent, + get_active_subagent_parent, +) +from run_agent import AIAgent + + +def _bare_agent() -> AIAgent: + agent = AIAgent.__new__(AIAgent) + agent._active_children = [] + import threading + + agent._active_children_lock = threading.Lock() + agent._session_db = None + agent.session_id = "child-x" + return agent + + +def test_close_releases_transcript_shadow_copies(): + agent = _bare_agent() + + class Payload(str): # weakref-able stand-in for a message content string + pass + + payload = Payload("x" * 50_000) + big = {"role": "tool", "content": payload} + agent._session_messages = [big] + agent._db_flush_scan_prefix = agent._session_messages[:] + agent._streamed_assistant_text_parts = ["y" * 10_000] + probe = weakref.ref(payload) + + agent.close() + + assert agent._session_messages == [] + assert agent._db_flush_scan_prefix is None + assert agent._streamed_assistant_text_parts == [] + del big, payload + gc.collect() + assert probe() is None, "closed agent still owns its message dicts" + + +def test_bind_subagent_parent_does_not_pin_agent(): + agent = _bare_agent() + probe = weakref.ref(agent) + snapshots = [] + with bind_subagent_parent(agent): + assert get_active_subagent_parent() is agent + import contextvars + + # An asyncio Handle scheduled inside the turn keeps this snapshot. + snapshots.append(contextvars.copy_context()) + assert get_active_subagent_parent() is None + assert snapshots[0][_ACTIVE_PARENT_AGENT] is not agent + del agent + gc.collect() + assert probe() is None, "Context snapshot still pins the agent" + + +def test_bind_subagent_parent_accepts_non_weakrefable_doubles(): + class Slots: + __slots__ = () + + double = Slots() + with bind_subagent_parent(double): + assert get_active_subagent_parent() is double + + +def _fake_child(messages): + child = MagicMock() + child._credential_pool = None + child._delegate_role = "leaf" + child.session_estimated_cost_usd = 0.0123 + child.session_cost_status = "estimated" + child.session_id = "child-sess" + child.run_conversation.return_value = { + "final_response": "the summary", + "completed": True, + "interrupted": False, + "api_calls": 3, + "messages": messages, + } + return child + + +def test_run_single_child_result_json_unchanged_by_transcript_release(): + """Pin: the parent-visible result entry is byte-identical whether or not + the child released its transcript at close() (the entry never carried + ``messages``; only summary/tool_trace/tokens/cost derive from them).""" + from tests.tools.test_delegate import _make_mock_parent + from tools.delegate_tool import _run_single_child + + messages = [ + {"role": "user", "content": "goal"}, + { + "role": "assistant", + "content": None, + "tool_calls": [ + { + "id": "c1", + "type": "function", + "function": {"name": "read_file", "arguments": '{"path": "a.py"}'}, + } + ], + }, + {"role": "tool", "tool_call_id": "c1", "content": "z" * 5000}, + {"role": "assistant", "content": "the summary"}, + ] + results = [] + for _ in range(2): + child = _fake_child([dict(m) for m in messages]) + entry = _run_single_child( + task_index=0, goal="goal", child=child, parent_agent=_make_mock_parent() + ) + child.close.assert_called_once() + entry.pop("duration_seconds", None) + results.append(json.dumps(entry, sort_keys=True, default=str)) + assert results[0] == results[1] + parsed = json.loads(results[0]) + assert parsed["summary"] == "the summary" + assert parsed["tool_trace"][0]["tool"] == "read_file" + assert parsed["cost_usd"] == 0.0123 + assert "messages" not in parsed