From e021d410288314e3c937cf4b03f18e0dc7778b0e Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 10 Aug 2026 23:47:53 +0000 Subject: [PATCH] fix(claude-native): register nested subagents under their direct parent The Claude native forwarder always registered Task-tool subagents against the top-level session and only scanned the root subagents/ directory. Nested subagents therefore appeared as direct children of the main session in the web graph (and Agents list), flattening the true spawn chain. Walk nested subagents/ directories, resolve each spawn's toolUseId against the owning transcript, and persist each row's jsonl path so tail/cost forwarding follow nested files after restart. Co-authored-by: jan21deepak --- omnigent/claude_native_forwarder.py | 473 ++++++++++++++++++++------ tests/test_claude_native_forwarder.py | 228 +++++++++++++ 2 files changed, 603 insertions(+), 98 deletions(-) diff --git a/omnigent/claude_native_forwarder.py b/omnigent/claude_native_forwarder.py index def6cf67f4..6867bac5a2 100644 --- a/omnigent/claude_native_forwarder.py +++ b/omnigent/claude_native_forwarder.py @@ -515,6 +515,10 @@ class SubagentEntry: sub-agent — used to dedupe so we don't spam ``running`` or ``idle`` events on every tick when nothing changed. ``None`` means no status has been posted yet. + :param jsonl_path: Absolute path to this sub-agent's on-disk + ``agent-.jsonl`` transcript. Persisted so nested + sub-agents (under ``…/agent-/subagents/``) tail + from the correct file after a forwarder restart. """ subagent_id: str @@ -523,6 +527,7 @@ class SubagentEntry: seen_source_ids: tuple[str, ...] = () last_activity_ts: float | None = None last_status: str | None = None + jsonl_path: str | None = None @dataclass(frozen=True) @@ -1144,6 +1149,351 @@ def _subagents_dir_for_transcript(transcript_path: Path) -> Path: return transcript_path.parent / transcript_path.stem / "subagents" +_SPAWN_TOOL_NAMES = frozenset({"Agent", "Task"}) + + +def _subagent_jsonl_path( + entry: SubagentEntry, + *, + fallback_subagents_dir: Path, +) -> Path: + """ + Resolve the on-disk transcript path for a tracked sub-agent. + + :param entry: Forwarder cursor row for the sub-agent. + :param fallback_subagents_dir: Root ``subagents/`` directory used + when ``entry.jsonl_path`` was not persisted (legacy state). + :returns: Path to ``agent-.jsonl``. + """ + if entry.jsonl_path: + return Path(entry.jsonl_path) + return fallback_subagents_dir / f"agent-{entry.subagent_id}.jsonl" + + +def _spawn_tool_use_ids_in_jsonl(jsonl_path: Path) -> set[str]: + """ + Collect Task/Agent ``tool_use`` block ids from a Claude transcript. + + Used to map a nested sub-agent's ``toolUseId`` meta field back to + the Omnigent session that issued the spawn call. + + :param jsonl_path: Claude transcript JSONL to scan. + :returns: Tool-use ids for spawn calls found in the file. + """ + if not jsonl_path.is_file(): + return set() + spawn_ids: set[str] = set() + try: + lines = jsonl_path.read_text(encoding="utf-8").splitlines() + except OSError: + return spawn_ids + for line in lines: + if not line.strip(): + continue + try: + entry = json.loads(line) + except json.JSONDecodeError: + continue + if not isinstance(entry, dict): + continue + message = entry.get("message") + if not isinstance(message, dict): + continue + content = message.get("content") + if not isinstance(content, list): + continue + for block in content: + if not isinstance(block, dict) or block.get("type") != "tool_use": + continue + name = block.get("name") + tool_id = block.get("id") + if name in _SPAWN_TOOL_NAMES and isinstance(tool_id, str) and tool_id: + spawn_ids.add(tool_id) + return spawn_ids + + +def _spawn_tool_use_parent_map( + *, + root_session_id: str, + parent_transcript_path: Path, + root_subagents_dir: Path, + state: SubagentForwardState, +) -> dict[str, str]: + """ + Map spawn ``tool_use`` ids to the Omnigent session that owns them. + + Covers the root transcript plus every registered sub-agent + transcript so a nested sub-agent's ``toolUseId`` resolves to its + direct parent rather than the top-level session. + + :param root_session_id: Omnigent id for the main Claude session. + :param parent_transcript_path: Root transcript JSONL path. + :param root_subagents_dir: Root ``subagents/`` directory (legacy + jsonl-path fallback for entries without ``jsonl_path``). + :param state: Current sub-agent cursor map. + :returns: ``tool_use_id`` → ``conversation_id`` for spawn calls. + """ + mapping: dict[str, str] = {} + for tool_use_id in _spawn_tool_use_ids_in_jsonl(parent_transcript_path): + mapping[tool_use_id] = root_session_id + for entry in state.subagents.values(): + if not entry.child_conversation_id: + continue + jsonl_path = _subagent_jsonl_path(entry, fallback_subagents_dir=root_subagents_dir) + for tool_use_id in _spawn_tool_use_ids_in_jsonl(jsonl_path): + mapping[tool_use_id] = entry.child_conversation_id + return mapping + + +async def _register_subagents_in_dir( + *, + client: httpx.AsyncClient, + bridge_dir: Path, + subagents_dir: Path, + owner_transcript_path: Path, + default_parent_session_id: str, + spawn_parent_map: dict[str, str], + state: SubagentForwardState, + start_retry_tracker: _PostRetryTracker, +) -> SubagentForwardState: + """ + Mint Omnigent child rows for every unseen ``.meta.json`` in one + ``subagents/`` directory. + + Resolves each sub-agent's parent via its ``toolUseId`` (falling + back to ``default_parent_session_id`` when the spawn call lives in + ``owner_transcript_path``). Repeats until no new rows register so + nested spawns whose parent transcript was just registered in the + same poll can attach on a later pass. + + :param client: Omnigent HTTP client. + :param bridge_dir: Native Claude bridge directory. + :param subagents_dir: Directory holding ``agent-*.meta.json`` files. + :param owner_transcript_path: Transcript for the session that owns + this ``subagents/`` directory (root jsonl or a sub-agent jsonl). + :param default_parent_session_id: Omnigent parent when the spawn + ``tool_use`` appears in ``owner_transcript_path``. + :param spawn_parent_map: ``tool_use_id`` → parent conversation id. + :param state: Current sub-agent cursor map. + :param start_retry_tracker: Backoff tracker for start POST failures. + :returns: Updated cursor map including newly registered sub-agents. + """ + updated = state + owner_spawn_ids = _spawn_tool_use_ids_in_jsonl(owner_transcript_path) + meta_paths = await asyncio.to_thread(lambda: sorted(subagents_dir.glob(_SUBAGENT_META_GLOB))) + pending = [ + meta_path + for meta_path in meta_paths + if meta_path.stem.removeprefix("agent-").removesuffix(".meta") not in updated.subagents + ] + while pending: + progress = False + next_pending: list[Path] = [] + for meta_path in pending: + subagent_id = meta_path.stem.removeprefix("agent-").removesuffix(".meta") + if subagent_id in updated.subagents: + continue + retry_key = f"subagent_start:{subagent_id}" + if start_retry_tracker.retry_delay_s(retry_key) is not None: + next_pending.append(meta_path) + continue + meta = await asyncio.to_thread(_read_subagent_meta, meta_path) + if meta is None: + next_pending.append(meta_path) + continue + tool_use_id = meta["toolUseId"] + if tool_use_id in spawn_parent_map: + parent_session_id = spawn_parent_map[tool_use_id] + elif tool_use_id in owner_spawn_ids: + parent_session_id = default_parent_session_id + elif any( + other.stem.removeprefix("agent-").removesuffix(".meta") not in updated.subagents + for other in meta_paths + if other != meta_path + ): + # Another unseen sibling may register first and index + # this spawn call on a later pass (flat layout). + next_pending.append(meta_path) + continue + else: + # Spawn record not flushed yet — attach to the directory + # owner, matching the legacy single-level behavior. + parent_session_id = default_parent_session_id + jsonl_path = subagents_dir / f"agent-{subagent_id}.jsonl" + try: + child_id = await _post_external_subagent_start( + client, + parent_session_id=parent_session_id, + subagent_id=subagent_id, + agent_type=meta["agentType"], + description=meta["description"], + tool_use_id=meta["toolUseId"], + ) + except httpx.HTTPError as exc: + decision = start_retry_tracker.record_failure(retry_key, exc) + if decision.exhausted: + _logger.error( + "Dropping claude-native sub-agent after permanent HTTP failures; " + "parent_session=%s subagent_id=%s attempts=%s http_status=%s", + parent_session_id, + subagent_id, + decision.attempts, + _http_status_for_log(exc), + ) + append_dead_letter( + bridge_dir, + session_id=parent_session_id, + event_type="external_subagent_start", + payload={ + "subagent_id": subagent_id, + "agent_type": meta["agentType"], + "description": meta["description"], + "tool_use_id": meta["toolUseId"], + }, + reason="permanent HTTP failure after retries", + delivered_ambiguous=False, + http_status=_http_status_for_log(exc), + ) + updated = SubagentForwardState( + subagents={ + **updated.subagents, + subagent_id: SubagentEntry( + subagent_id=subagent_id, + child_conversation_id="", + jsonl_path=str(jsonl_path), + ), + } + ) + await _write_subagent_forward_state_async(bridge_dir, updated) + progress = True + continue + _logger.warning( + "Failed to register claude-native sub-agent; parent_session=%s " + "subagent_id=%s attempt=%s permanent=%s next_retry_s=%.3f " + "http_status=%s", + parent_session_id, + subagent_id, + decision.attempts, + decision.permanent, + decision.delay_s, + _http_status_for_log(exc), + exc_info=True, + ) + next_pending.append(meta_path) + continue + start_retry_tracker.clear(retry_key) + updated = SubagentForwardState( + subagents={ + **updated.subagents, + subagent_id: SubagentEntry( + subagent_id=subagent_id, + child_conversation_id=child_id, + jsonl_path=str(jsonl_path), + ), + } + ) + await _write_subagent_forward_state_async(bridge_dir, updated) + for tool_use_id in _spawn_tool_use_ids_in_jsonl(jsonl_path): + spawn_parent_map[tool_use_id] = child_id + progress = True + if not progress: + break + pending = next_pending + owner_spawn_ids = _spawn_tool_use_ids_in_jsonl(owner_transcript_path) + return updated + + +async def _discover_nested_subagents( + *, + client: httpx.AsyncClient, + bridge_dir: Path, + root_session_id: str, + parent_transcript_path: Path, + state: SubagentForwardState, + start_retry_tracker: _PostRetryTracker, +) -> SubagentForwardState: + """ + Walk the Claude sub-agent tree on disk and register every unseen row. + + Top-level sub-agents live under the root transcript's + ``subagents/`` directory; nested sub-agents live under + ``…/agent-/subagents/``. Each is registered against the + Omnigent session that issued its spawn ``tool_use``. + + :param client: Omnigent HTTP client. + :param bridge_dir: Native Claude bridge directory. + :param root_session_id: Omnigent id for the main Claude session. + :param parent_transcript_path: Root transcript JSONL path. + :param state: Current sub-agent cursor map. + :param start_retry_tracker: Backoff tracker for start POST failures. + :returns: Updated cursor map with nested sub-agents registered. + """ + root_subagents_dir = _subagents_dir_for_transcript(parent_transcript_path) + if not root_subagents_dir.is_dir(): + return state + + spawn_parent_map = _spawn_tool_use_parent_map( + root_session_id=root_session_id, + parent_transcript_path=parent_transcript_path, + root_subagents_dir=root_subagents_dir, + state=state, + ) + updated = await _register_subagents_in_dir( + client=client, + bridge_dir=bridge_dir, + subagents_dir=root_subagents_dir, + owner_transcript_path=parent_transcript_path, + default_parent_session_id=root_session_id, + spawn_parent_map=spawn_parent_map, + state=state, + start_retry_tracker=start_retry_tracker, + ) + + seen_dirs: set[Path] = {root_subagents_dir} + frontier: list[tuple[Path, Path, str]] = [] + for entry in updated.subagents.values(): + if not entry.child_conversation_id: + continue + jsonl_path = _subagent_jsonl_path(entry, fallback_subagents_dir=root_subagents_dir) + nested_dir = _subagents_dir_for_transcript(jsonl_path) + if nested_dir.is_dir() and nested_dir not in seen_dirs: + frontier.append((nested_dir, jsonl_path, entry.child_conversation_id)) + + while frontier: + nested_dir, owner_transcript_path, default_parent_session_id = frontier.pop() + if nested_dir in seen_dirs: + continue + seen_dirs.add(nested_dir) + spawn_parent_map = _spawn_tool_use_parent_map( + root_session_id=root_session_id, + parent_transcript_path=parent_transcript_path, + root_subagents_dir=root_subagents_dir, + state=updated, + ) + before_ids = set(updated.subagents) + updated = await _register_subagents_in_dir( + client=client, + bridge_dir=bridge_dir, + subagents_dir=nested_dir, + owner_transcript_path=owner_transcript_path, + default_parent_session_id=default_parent_session_id, + spawn_parent_map=spawn_parent_map, + state=updated, + start_retry_tracker=start_retry_tracker, + ) + for subagent_id in updated.subagents: + if subagent_id in before_ids: + continue + entry = updated.subagents[subagent_id] + if not entry.child_conversation_id: + continue + jsonl_path = _subagent_jsonl_path(entry, fallback_subagents_dir=root_subagents_dir) + child_nested_dir = _subagents_dir_for_transcript(jsonl_path) + if child_nested_dir.is_dir() and child_nested_dir not in seen_dirs: + frontier.append((child_nested_dir, jsonl_path, entry.child_conversation_id)) + return updated + + def _read_subagent_forward_state(bridge_dir: Path) -> SubagentForwardState: """ Read the sub-agent forwarder's durable cursor map. @@ -1174,6 +1524,7 @@ def _read_subagent_forward_state(bridge_dir: Path) -> SubagentForwardState: seen_source_ids = row.get("seen_source_ids", []) last_activity_ts = row.get("last_activity_ts") last_status = row.get("last_status") + jsonl_path = row.get("jsonl_path") # Empty string is a valid parked sentinel written by # ``_forward_available_subagents`` after the start POST exhausts # its permanent-failure budget. Preserving it across restarts is @@ -1190,6 +1541,8 @@ def _read_subagent_forward_state(bridge_dir: Path) -> SubagentForwardState: last_activity_ts = None if last_status is not None and not isinstance(last_status, str): last_status = None + if jsonl_path is not None and not isinstance(jsonl_path, str): + jsonl_path = None entries[subagent_id] = SubagentEntry( subagent_id=subagent_id, child_conversation_id=child_id, @@ -1197,6 +1550,7 @@ def _read_subagent_forward_state(bridge_dir: Path) -> SubagentForwardState: seen_source_ids=tuple(seen_source_ids), last_activity_ts=last_activity_ts, last_status=last_status, + jsonl_path=jsonl_path, ) return SubagentForwardState(subagents=entries) @@ -1218,6 +1572,7 @@ def _write_subagent_forward_state(bridge_dir: Path, state: SubagentForwardState) "seen_source_ids": list(entry.seen_source_ids), "last_activity_ts": entry.last_activity_ts, "last_status": entry.last_status, + "jsonl_path": entry.jsonl_path, } for entry in state.subagents.values() }, @@ -1407,103 +1762,16 @@ async def _forward_available_subagents( :returns: Updated state with new sub-agents registered and existing sub-agents' cursors advanced. """ - subagents_dir = _subagents_dir_for_transcript(transcript_path) - if not subagents_dir.is_dir(): - return state + root_subagents_dir = _subagents_dir_for_transcript(transcript_path) - # ── Register newly-appeared sub-agents ────────────── - # ``glob`` is sync; offload to a thread so we don't stat the - # filesystem on the event loop. - meta_paths = await asyncio.to_thread(lambda: sorted(subagents_dir.glob(_SUBAGENT_META_GLOB))) - updated = state - for meta_path in meta_paths: - # ``agent-.meta.json`` → ```` - subagent_id = meta_path.stem.removeprefix("agent-").removesuffix(".meta") - if subagent_id in updated.subagents: - continue - retry_key = f"subagent_start:{subagent_id}" - if start_retry_tracker.retry_delay_s(retry_key) is not None: - continue - meta = await asyncio.to_thread(_read_subagent_meta, meta_path) - if meta is None: - # File may be mid-write; try again on the next tick. - continue - try: - child_id = await _post_external_subagent_start( - client, - parent_session_id=parent_session_id, - subagent_id=subagent_id, - agent_type=meta["agentType"], - description=meta["description"], - tool_use_id=meta["toolUseId"], - ) - except httpx.HTTPError as exc: - decision = start_retry_tracker.record_failure(retry_key, exc) - if decision.exhausted: - _logger.error( - "Dropping claude-native sub-agent after permanent HTTP failures; " - "parent_session=%s subagent_id=%s attempts=%s http_status=%s", - parent_session_id, - subagent_id, - decision.attempts, - _http_status_for_log(exc), - ) - # Dead-letter the dropped payload for recovery (#1120; replay #1579). - append_dead_letter( - bridge_dir, - session_id=parent_session_id, - event_type="external_subagent_start", - payload={ - "subagent_id": subagent_id, - "agent_type": meta["agentType"], - "description": meta["description"], - "tool_use_id": meta["toolUseId"], - }, - reason="permanent HTTP failure after retries", - # Claude only dead-letters permanent 4xx (it retries - # transient failures forever), so the server proved it - # rejected the item: never ambiguous, never replayable (#1579). - delivered_ambiguous=False, - http_status=_http_status_for_log(exc), - ) - # Park this sub-agent: insert a sentinel entry so we - # don't keep retrying. ``child_conversation_id=""`` - # is filtered out by the tail / status loops below. - updated = SubagentForwardState( - subagents={ - **updated.subagents, - subagent_id: SubagentEntry( - subagent_id=subagent_id, - child_conversation_id="", - ), - } - ) - await _write_subagent_forward_state_async(bridge_dir, updated) - continue - _logger.warning( - "Failed to register claude-native sub-agent; parent_session=%s " - "subagent_id=%s attempt=%s permanent=%s next_retry_s=%.3f " - "http_status=%s", - parent_session_id, - subagent_id, - decision.attempts, - decision.permanent, - decision.delay_s, - _http_status_for_log(exc), - exc_info=True, - ) - continue - start_retry_tracker.clear(retry_key) - updated = SubagentForwardState( - subagents={ - **updated.subagents, - subagent_id: SubagentEntry( - subagent_id=subagent_id, - child_conversation_id=child_id, - ), - } - ) - await _write_subagent_forward_state_async(bridge_dir, updated) + updated = await _discover_nested_subagents( + client=client, + bridge_dir=bridge_dir, + root_session_id=parent_session_id, + parent_transcript_path=transcript_path, + state=state, + start_retry_tracker=start_retry_tracker, + ) # ── Tail each tracked sub-agent's transcript ──────── now = time.time() @@ -1511,7 +1779,7 @@ async def _forward_available_subagents( if not entry.child_conversation_id: # Parked after exhausted start retries — nothing to tail. continue - jsonl_path = subagents_dir / f"agent-{subagent_id}.jsonl" + jsonl_path = _subagent_jsonl_path(entry, fallback_subagents_dir=root_subagents_dir) if not jsonl_path.exists(): continue # Reuse the parent-transcript parser, but pass @@ -1592,6 +1860,7 @@ async def _forward_available_subagents( seen_source_ids=_bounded_seen_source_ids(seen_source_ids), last_activity_ts=new_entry.last_activity_ts, last_status=new_entry.last_status, + jsonl_path=entry.jsonl_path, ) updated = SubagentForwardState( subagents={**updated.subagents, subagent_id: new_entry} @@ -1621,6 +1890,7 @@ async def _forward_available_subagents( seen_source_ids=_bounded_seen_source_ids(seen_source_ids), last_activity_ts=new_entry.last_activity_ts, last_status=new_entry.last_status, + jsonl_path=entry.jsonl_path, ) updated = SubagentForwardState( subagents={**updated.subagents, subagent_id: new_entry} @@ -1657,6 +1927,7 @@ async def _forward_available_subagents( seen_source_ids=_bounded_seen_source_ids(seen_source_ids), last_activity_ts=now, last_status=new_entry.last_status, + jsonl_path=entry.jsonl_path, ) updated = SubagentForwardState(subagents={**updated.subagents, subagent_id: new_entry}) await _write_subagent_forward_state_async(bridge_dir, updated) @@ -1671,6 +1942,7 @@ async def _forward_available_subagents( seen_source_ids=_bounded_seen_source_ids(seen_source_ids), last_activity_ts=now if had_item else entry.last_activity_ts, last_status=entry.last_status, + jsonl_path=entry.jsonl_path, ) elif had_item: # Items DID flow but a later post failed — still record @@ -1684,6 +1956,7 @@ async def _forward_available_subagents( seen_source_ids=_bounded_seen_source_ids(seen_source_ids), last_activity_ts=now, last_status=entry.last_status, + jsonl_path=entry.jsonl_path, ) # Quiescence-based status. Sub-agent transcripts don't carry @@ -1731,6 +2004,7 @@ async def _forward_available_subagents( seen_source_ids=new_entry.seen_source_ids, last_activity_ts=new_entry.last_activity_ts, last_status=desired_status, + jsonl_path=new_entry.jsonl_path, ) if new_entry is not entry: @@ -1839,7 +2113,10 @@ def _session_cost_estimate( parent_transcript_path, include_sidechains=False, cache=cost_cache ) for entry in active_subagents: - jsonl_path = subagents_dir / f"agent-{entry.subagent_id}.jsonl" + jsonl_path = _subagent_jsonl_path( + entry, + fallback_subagents_dir=subagents_dir, + ) sub_cost = _transcript_cost_size_cached( jsonl_path, include_sidechains=True, cache=cost_cache ) diff --git a/tests/test_claude_native_forwarder.py b/tests/test_claude_native_forwarder.py index 4209990c13..19f8150702 100644 --- a/tests/test_claude_native_forwarder.py +++ b/tests/test_claude_native_forwarder.py @@ -4472,6 +4472,234 @@ def _seed_subagent_on_disk( return jsonl_path +def _seed_nested_subagent_on_disk( + *, + parent_subagent_jsonl: Path, + subagent_id: str, + agent_type: str, + description: str, + tool_use_id: str, + parent_spawn_tool_use_id: str, +) -> Path: + """ + Create a nested sub-agent under ``parent_subagent_jsonl/subagents/``. + + Also appends a spawn ``tool_use`` to the parent sub-agent transcript + so the forwarder can resolve the nested row's parent. + + :param parent_subagent_jsonl: Parent sub-agent ``agent-.jsonl``. + :param subagent_id: Nested Claude-side id. + :param agent_type: ``agentType`` for the meta file. + :param description: ``description`` for the meta file. + :param tool_use_id: Nested sub-agent's ``toolUseId``. + :param parent_spawn_tool_use_id: ``tool_use.id`` recorded in the + parent sub-agent transcript for the spawn call. + :returns: Path to the nested sub-agent's ``.jsonl``. + """ + parent_record = { + "type": "assistant", + "message": { + "role": "assistant", + "content": [ + { + "type": "tool_use", + "id": parent_spawn_tool_use_id, + "name": "Agent", + "input": {"description": description}, + } + ], + }, + } + if ( + parent_subagent_jsonl.exists() + and parent_subagent_jsonl.read_text(encoding="utf-8").strip() + ): + parent_subagent_jsonl.write_text( + parent_subagent_jsonl.read_text(encoding="utf-8").rstrip() + + "\n" + + json.dumps(parent_record) + + "\n", + encoding="utf-8", + ) + else: + parent_subagent_jsonl.write_text(json.dumps(parent_record) + "\n", encoding="utf-8") + + nested_transcript_path = ( + parent_subagent_jsonl.parent / parent_subagent_jsonl.stem / "subagents" + ) + nested_transcript_path.mkdir(parents=True, exist_ok=True) + meta_path = nested_transcript_path / f"agent-{subagent_id}.meta.json" + meta_path.write_text( + json.dumps( + { + "agentType": agent_type, + "description": description, + "toolUseId": tool_use_id, + } + ), + encoding="utf-8", + ) + jsonl_path = nested_transcript_path / f"agent-{subagent_id}.jsonl" + jsonl_path.write_text("", encoding="utf-8") + return jsonl_path + + +def test_spawn_tool_use_parent_map_links_nested_spawn_to_child_parent( + tmp_path: Path, +) -> None: + """Nested sub-agents resolve their spawn tool_use to the minted child row.""" + root_transcript = tmp_path / "session.jsonl" + root_transcript.write_text( + json.dumps( + { + "type": "assistant", + "message": { + "role": "assistant", + "content": [ + { + "type": "tool_use", + "id": "toolu_root_spawn", + "name": "Task", + "input": {"description": "parent task"}, + } + ], + }, + } + ) + + "\n", + encoding="utf-8", + ) + parent_jsonl = _seed_subagent_on_disk( + transcript_path=root_transcript, + subagent_id="parent_sa", + agent_type="Explore", + description="Parent sub-agent", + tool_use_id="toolu_root_spawn", + ) + _seed_nested_subagent_on_disk( + parent_subagent_jsonl=parent_jsonl, + subagent_id="nested_sa", + agent_type="Explore", + description="Nested sub-agent", + tool_use_id="toolu_nested_spawn", + parent_spawn_tool_use_id="toolu_nested_spawn", + ) + root_subagents_dir = forwarder._subagents_dir_for_transcript(root_transcript) + state = forwarder.SubagentForwardState( + subagents={ + "parent_sa": forwarder.SubagentEntry( + subagent_id="parent_sa", + child_conversation_id="conv_parent_sa", + jsonl_path=str(parent_jsonl), + ) + } + ) + mapping = forwarder._spawn_tool_use_parent_map( + root_session_id="conv_root", + parent_transcript_path=root_transcript, + root_subagents_dir=root_subagents_dir, + state=state, + ) + assert mapping["toolu_root_spawn"] == "conv_root" + assert mapping["toolu_nested_spawn"] == "conv_parent_sa" + + +async def test_subagent_watcher_registers_nested_subagent_under_parent_child( + tmp_path: Path, +) -> None: + """ + A nested ``agent-.meta.json`` under ``…/agent-/subagents/`` + registers via ``external_subagent_start`` against the parent sub-agent's + Omnigent child session, not the top-level session. + """ + bridge_dir = tmp_path / "bridge" + transcript_path = tmp_path / "session.jsonl" + transcript_path.write_text( + json.dumps( + { + "type": "assistant", + "message": { + "role": "assistant", + "content": [ + { + "type": "tool_use", + "id": "toolu_root_spawn", + "name": "Task", + "input": {"description": "parent task"}, + } + ], + }, + } + ) + + "\n", + encoding="utf-8", + ) + parent_jsonl = _seed_subagent_on_disk( + transcript_path=transcript_path, + subagent_id="parent_sa", + agent_type="Explore", + description="Parent sub-agent", + tool_use_id="toolu_root_spawn", + ) + _seed_nested_subagent_on_disk( + parent_subagent_jsonl=parent_jsonl, + subagent_id="nested_sa", + agent_type="Explore", + description="Nested sub-agent", + tool_use_id="toolu_nested_spawn", + parent_spawn_tool_use_id="toolu_nested_spawn", + ) + record_hook_event( + bridge_dir, + { + "hook_event_name": "SessionStart", + "session_id": "claude-session", + "transcript_path": str(transcript_path), + }, + ) + + start_requests: list[dict[str, Any]] = [] + + def response_for(body: dict[str, Any]) -> dict[str, Any]: + if body.get("type") == "external_subagent_start": + subagent_id = body["data"]["subagent_id"] + if subagent_id == "parent_sa": + return {"queued": False, "child_session_id": "conv_parent_sa"} + if subagent_id == "nested_sa": + return {"queued": False, "child_session_id": "conv_nested_sa"} + return {} + + server, _thread, base_url = _start_recording_server_with_responses(response_for) + task = asyncio.create_task( + forward_claude_transcript_to_session( + base_url=base_url, + headers={}, + session_id="conv_root", + bridge_dir=bridge_dir, + agent_name="claude-native-ui", + start_at_end=False, + poll_interval_s=0.01, + ) + ) + try: + for _ in range(40): + if len(start_requests) >= 2: + break + req = await _get_recorded_request(server, timeout_s=1.0) + if req["body"].get("type") == "external_subagent_start": + start_requests.append(req) + assert len(start_requests) >= 2, "forwarder did not register both sub-agents" + by_subagent = {req["body"]["data"]["subagent_id"]: req for req in start_requests} + assert by_subagent["parent_sa"]["path"] == "/v1/sessions/conv_root/events" + assert by_subagent["nested_sa"]["path"] == "/v1/sessions/conv_parent_sa/events" + finally: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + server.shutdown() + server.server_close() + + async def test_subagent_watcher_posts_external_subagent_start_for_new_meta( tmp_path: Path, ) -> None: