feat(auto-agents): archive entire task-tool subagent tree before DELETE root
Before this change the dispatcher archived only the top-level
wrapper session. The entire ``task``-tool subagent chain
(tier-dispatcher → estimator-implementation / tier-qwen-med →
task-implementor → git-isolator-util) was opaque the moment
the dispatcher's DELETE /session/{id} fired, so post-mortem
analysis of an implementer run was limited to whatever
live-API polling we'd done DURING the run. That's how the
recent optimization round had to work from two cherry-picked
live snapshots of task-implementor and git-isolator-util —
unreliable, only what happened to be active when polled.
Three changes:
1. tools/_opencode_worker.py
- New helpers: _walk_subagent_descendants (BFS over
GET /session keyed on parentID), _extract_subagent_agent_name
(parses OpenCode's "(@<agent> subagent)" title convention),
_ms_to_iso (epoch-ms to ISO-8601), and _archive_subagent_tree
(best-effort walk + fetch + write driver; never raises).
- _archive_session / _build_archive_payload gain optional
parent_session_id / subagent_title / subagent_depth kwargs.
When set, the filename includes a ``sub<depth>`` infix
(e.g. 2026-...__sub01__AUTO-IMP-PR-30__tier-dispatcher__ses_*.json)
so a directory listing groups every session from one
dispatcher cycle and reads top-down in BFS order.
- Archive schema bumped from v1 → v2. New fields are nullable;
v1 readers (the existing telemetry-console endpoints) treat
them as missing and remain forward-compatible.
- run_session_blocking's finally block calls
_archive_subagent_tree after the root archive write and
before the root DELETE. Both calls are wrapped in
try/except so a subagent-walk failure can never mask the
worker outcome or stop the dispatcher from cleaning up.
- The dispatcher's existing redact_values list (the Forgejo
PAT) propagates into every subagent archive too, so a
``git clone https://${PAT}@...`` in git-isolator-util's
bash history is masked the same way the wrapper's prompt is.
2. .opencode/telemetry/server.py
- _api_archived_sessions listing endpoint now surfaces the
three v2 fields (schema_version, parent_session_id,
subagent_title, subagent_depth) in each row payload so a
future UI render can nest subagents under their wrapper.
Additive — existing row keys are preserved.
3. tests/auto_agents/test_opencode_worker_observability.py
- 17 new tests across four classes:
- TestSubagentTitleExtraction (5): title parser edge cases
- TestWalkSubagentDescendants (6): BFS order, depth
annotation, transport-error / malformed-payload paths,
cycle safety
- TestArchiveSubagentTree (5): end-to-end orchestration
including a redaction-propagation test that asserts a
PAT inside a subagent's bash tool input is replaced
with <REDACTED>
- TestEndToEndSubagentArchive (1): drives the full
run_session_blocking lifecycle with a wired subagent
descendant and asserts BOTH archives land on disk
- Existing schema-version assertion updated to v2 + three
new ``None``-on-top-level field assertions.
- Two manually-wired archive tests (transport-error,
timeout) now wire GET /session so the walker doesn't emit
a spurious warning.
Total auto_agents suite: 1061 passed, 3 skipped (up from 1044).
This is the prerequisite for trustworthy quantification of the
upcoming default-flip of IMPLEMENTER_DISPATCHER_PREFETCH=1 and
IMPLEMENTER_DISPATCHER_PRECLONE=1. With the walker in place,
every cycle now leaves a complete trace on disk that a human
can read bottom-up months later.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -869,6 +869,17 @@ def _api_archived_sessions() -> dict[str, Any]:
|
|||||||
"wallclock_seconds": data.get("wallclock_seconds"),
|
"wallclock_seconds": data.get("wallclock_seconds"),
|
||||||
"message_count": len(data.get("messages") or []),
|
"message_count": len(data.get("messages") or []),
|
||||||
"turn_count": len(per_turn) if isinstance(per_turn, list) else 0,
|
"turn_count": len(per_turn) if isinstance(per_turn, list) else 0,
|
||||||
|
# v2 (2026-05-10) fields: present for subagent archives,
|
||||||
|
# ``None`` for top-level wrapper archives. The UI can use
|
||||||
|
# ``parent_session_id`` to group every session from one
|
||||||
|
# dispatcher cycle into a single tree row, and
|
||||||
|
# ``subagent_depth`` to nest the table view. Old (v1)
|
||||||
|
# archives lack these keys; ``.get`` returns None so the
|
||||||
|
# row schema stays uniform.
|
||||||
|
"schema_version": data.get("schema_version"),
|
||||||
|
"parent_session_id": data.get("parent_session_id"),
|
||||||
|
"subagent_title": data.get("subagent_title"),
|
||||||
|
"subagent_depth": data.get("subagent_depth"),
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
# Newest archive first. Use ``modified_at`` since it is ISO-formatted
|
# Newest archive first. Use ``modified_at`` since it is ISO-formatted
|
||||||
|
|||||||
@@ -5,6 +5,96 @@ The format follows [Keep a Changelog](https://keepachangelog.com/en/1.1.0/).
|
|||||||
|
|
||||||
## [Unreleased]
|
## [Unreleased]
|
||||||
|
|
||||||
|
### Added
|
||||||
|
|
||||||
|
- **Subagent session archive walker (2026-05-10).** Before this
|
||||||
|
change the dispatcher archived only the top-level wrapper
|
||||||
|
session — the entire ``task``-tool subagent chain
|
||||||
|
(``tier-dispatcher`` → ``estimator-implementation`` /
|
||||||
|
``tier-qwen-med`` → ``task-implementor`` → ``git-isolator-util``)
|
||||||
|
was opaque the moment the dispatcher's ``DELETE /session/{id}``
|
||||||
|
fired. The post-mortem analysis of an implementer run was
|
||||||
|
therefore limited to whatever live-API polling we'd done DURING
|
||||||
|
the run, which is unreliable and only catches the agents that
|
||||||
|
happened to be active when we polled.
|
||||||
|
|
||||||
|
The wrapper's ``run_session_blocking`` ``finally`` block now
|
||||||
|
walks ``GET /session`` BFS-from-the-root after the root archive
|
||||||
|
is written and BEFORE the root ``DELETE``. Every descendant
|
||||||
|
session — at any depth — has its full ``/session/{id}/message``
|
||||||
|
stream fetched and written to a per-descendant archive file.
|
||||||
|
The schema is the same as the wrapper archive (so the existing
|
||||||
|
telemetry-console endpoints render them) with three new
|
||||||
|
optional fields:
|
||||||
|
|
||||||
|
- ``parent_session_id`` — wrapper's ``session_id`` for direct
|
||||||
|
children; the immediate parent's id for deeper nodes.
|
||||||
|
- ``subagent_title`` — OpenCode's full session title (e.g.
|
||||||
|
``"Dispatch implementation to tier-dispatcher (@tier-dispatcher subagent)"``).
|
||||||
|
- ``subagent_depth`` — 1-based distance from the wrapper.
|
||||||
|
|
||||||
|
Filenames embed a ``sub<depth>`` infix
|
||||||
|
(``<started_at>__sub01__<tag>__<agent>__<session_id>.json``) so a
|
||||||
|
``ls -1 .dispatcher-logs/sessions`` listing groups every session
|
||||||
|
from one dispatcher cycle and reads top-down in BFS order. The
|
||||||
|
schema version is bumped to **v2**; v1 readers see ``None`` for
|
||||||
|
the three new fields and remain forward-compatible (no breaking
|
||||||
|
change for the telemetry UI). The redaction list passed by the
|
||||||
|
dispatcher (the Forgejo PAT) propagates into every subagent
|
||||||
|
archive too, so a ``git clone https://${PAT}@...`` command in
|
||||||
|
``git-isolator-util``'s bash history is masked the same way the
|
||||||
|
wrapper's prompt is.
|
||||||
|
|
||||||
|
Concretely:
|
||||||
|
|
||||||
|
- ``tools/_opencode_worker.py`` gains
|
||||||
|
``_walk_subagent_descendants(server_url, root_session_id)``
|
||||||
|
(BFS over ``GET /session`` keyed on ``parentID``, returns one
|
||||||
|
dict per descendant with an added ``_subagent_depth`` field),
|
||||||
|
``_extract_subagent_agent_name(title)`` (regex on the
|
||||||
|
OpenCode title convention with a ``unknown-subagent``
|
||||||
|
fallback), ``_ms_to_iso(ms)`` (epoch-ms → ISO-8601 helper),
|
||||||
|
and ``_archive_subagent_tree(...)`` (best-effort
|
||||||
|
walk + fetch + write driver; never raises, logs and
|
||||||
|
continues on per-descendant transport failure).
|
||||||
|
- ``_archive_session`` and ``_build_archive_payload`` gain
|
||||||
|
optional ``parent_session_id`` / ``subagent_title`` /
|
||||||
|
``subagent_depth`` kwargs. When set, the filename includes
|
||||||
|
the ``sub<depth>`` infix; when unset, the legacy filename
|
||||||
|
pattern is preserved byte-for-byte.
|
||||||
|
- ``run_session_blocking``'s ``finally`` block invokes
|
||||||
|
``_archive_subagent_tree`` after the root archive write and
|
||||||
|
before the root ``DELETE``. Both calls are wrapped in
|
||||||
|
``try/except BaseException`` so a subagent-walk failure can
|
||||||
|
never mask the worker outcome or stop the
|
||||||
|
dispatcher from cleaning up.
|
||||||
|
- ``.opencode/telemetry/server.py``'s ``_api_archived_sessions``
|
||||||
|
listing endpoint now surfaces the three v2 fields in the row
|
||||||
|
payload so a future UI render can nest subagents under their
|
||||||
|
wrapper. The existing row shape is preserved (additive only).
|
||||||
|
- **17 new tests** in ``test_opencode_worker_observability.py``:
|
||||||
|
``TestSubagentTitleExtraction`` (5) covers the title parser,
|
||||||
|
``TestWalkSubagentDescendants`` (6) covers BFS order, depth
|
||||||
|
annotation, transport-error / malformed-payload paths, and
|
||||||
|
cycle safety, ``TestArchiveSubagentTree`` (5) covers
|
||||||
|
end-to-end orchestration including a redaction-propagation
|
||||||
|
test that asserts a PAT inside a subagent's
|
||||||
|
``bash`` tool input is replaced with ``<REDACTED>``, and
|
||||||
|
``TestEndToEndSubagentArchive`` (1) drives the full
|
||||||
|
``run_session_blocking`` lifecycle with a wired subagent
|
||||||
|
descendant and asserts both archives land on disk.
|
||||||
|
Total auto_agents suite: **1061 passed, 3 skipped** (up from
|
||||||
|
1044).
|
||||||
|
|
||||||
|
Why this matters: the prior post-mortem analysis recommendations
|
||||||
|
(default-flip ``IMPLEMENTER_DISPATCHER_PREFETCH=1`` etc.) all
|
||||||
|
rested on live-poll snapshots of two specific subagents I
|
||||||
|
happened to capture mid-run. With the walker in place, every
|
||||||
|
cycle leaves a complete trace on disk that a human can read
|
||||||
|
bottom-up months later, and the telemetry console can render
|
||||||
|
the wrapper → tier-dispatcher → task-implementor →
|
||||||
|
git-isolator-util tree as one navigation context.
|
||||||
|
|
||||||
### Changed
|
### Changed
|
||||||
|
|
||||||
- **Tighten `implementation-worker.md` skill allowlist and Step 0
|
- **Tighten `implementation-worker.md` skill allowlist and Step 0
|
||||||
|
|||||||
@@ -115,6 +115,15 @@ def _wire_lifecycle_with_statuses(
|
|||||||
httpserver.expect_request(
|
httpserver.expect_request(
|
||||||
f"/session/{session_id}/message", method="GET"
|
f"/session/{session_id}/message", method="GET"
|
||||||
).respond_with_json(msg_payload)
|
).respond_with_json(msg_payload)
|
||||||
|
# Subagent descendant walk (added 2026-05-10): _walk_subagent_descendants
|
||||||
|
# issues GET /session in the run_session_blocking finally block to find
|
||||||
|
# any subagent sessions that need archiving. By default the test wires
|
||||||
|
# this to return ONLY the root session (no descendants), so existing
|
||||||
|
# tests are subagent-archive-clean. Tests that exercise the subagent
|
||||||
|
# path should ``httpserver.clear()`` and re-wire with descendants.
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
[_session_obj(session_id)]
|
||||||
|
)
|
||||||
httpserver.expect_request(
|
httpserver.expect_request(
|
||||||
f"/session/{session_id}", method="DELETE"
|
f"/session/{session_id}", method="DELETE"
|
||||||
).respond_with_data("", status=200)
|
).respond_with_data("", status=200)
|
||||||
@@ -566,7 +575,7 @@ def test_archive_completed_session_writes_payload_with_schema(
|
|||||||
assert len(files) == 1, f"expected one archive file, got {files}"
|
assert len(files) == 1, f"expected one archive file, got {files}"
|
||||||
payload = json.loads(files[0].read_text())
|
payload = json.loads(files[0].read_text())
|
||||||
|
|
||||||
assert payload["schema_version"] == 1
|
assert payload["schema_version"] == 2
|
||||||
assert payload["session_id"] == "s-obs-1"
|
assert payload["session_id"] == "s-obs-1"
|
||||||
assert payload["agent"] == "test-agent"
|
assert payload["agent"] == "test-agent"
|
||||||
assert payload["tag"] == "AUTO-OBS"
|
assert payload["tag"] == "AUTO-OBS"
|
||||||
@@ -581,6 +590,10 @@ def test_archive_completed_session_writes_payload_with_schema(
|
|||||||
assert payload["state_history"][0]["state"] in {"busy", "idle", "unknown"}
|
assert payload["state_history"][0]["state"] in {"busy", "idle", "unknown"}
|
||||||
assert isinstance(payload["messages"], list)
|
assert isinstance(payload["messages"], list)
|
||||||
assert payload["messages"][-1]["info"]["role"] == "assistant"
|
assert payload["messages"][-1]["info"]["role"] == "assistant"
|
||||||
|
# v2 subagent-marker fields default to None on a top-level archive.
|
||||||
|
assert payload["parent_session_id"] is None
|
||||||
|
assert payload["subagent_title"] is None
|
||||||
|
assert payload["subagent_depth"] is None
|
||||||
|
|
||||||
|
|
||||||
def test_archive_transport_error_path_still_writes_payload(
|
def test_archive_transport_error_path_still_writes_payload(
|
||||||
@@ -614,6 +627,13 @@ def test_archive_transport_error_path_still_writes_payload(
|
|||||||
}
|
}
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
# Subagent walker fires in finally{} on every terminal path including
|
||||||
|
# transport-error. Wire it to return ONLY the root (no descendants)
|
||||||
|
# so the warning-on-transport-error path isn't muddied with subagent
|
||||||
|
# walk noise.
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
[_session_obj()]
|
||||||
|
)
|
||||||
httpserver.expect_request(
|
httpserver.expect_request(
|
||||||
"/session/s-obs-1", method="DELETE"
|
"/session/s-obs-1", method="DELETE"
|
||||||
).respond_with_data("")
|
).respond_with_data("")
|
||||||
@@ -675,6 +695,10 @@ def test_archive_timeout_path_still_writes_payload(
|
|||||||
}
|
}
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
# Subagent walker fires on timeout too — return only the root.
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
[_session_obj()]
|
||||||
|
)
|
||||||
httpserver.expect_request(
|
httpserver.expect_request(
|
||||||
"/session/s-obs-1", method="DELETE"
|
"/session/s-obs-1", method="DELETE"
|
||||||
).respond_with_data("")
|
).respond_with_data("")
|
||||||
@@ -1072,3 +1096,499 @@ class TestRunSessionRedactsArchive:
|
|||||||
for rec in caplog.records
|
for rec in caplog.records
|
||||||
)
|
)
|
||||||
assert short in body
|
assert short in body
|
||||||
|
|
||||||
|
|
||||||
|
# ─── Subagent archive walker (2026-05-10) ───────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class TestSubagentTitleExtraction:
|
||||||
|
"""``_extract_subagent_agent_name`` parses the ``(@<agent> subagent)``
|
||||||
|
suffix OpenCode appends to subagent session titles. The archive
|
||||||
|
walker relies on this to populate the ``agent`` field of the
|
||||||
|
subagent archive filename, so the parsing has to be tolerant of
|
||||||
|
the OpenCode title conventions we've actually seen on live runs.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def test_extracts_basic_agent_name(self, mod):
|
||||||
|
title = "Dispatch implementation to tier-dispatcher (@tier-dispatcher subagent)"
|
||||||
|
assert mod._extract_subagent_agent_name(title) == "tier-dispatcher"
|
||||||
|
|
||||||
|
def test_handles_dotted_agent_name(self, mod):
|
||||||
|
title = "Foo (@some.dotted.agent subagent)"
|
||||||
|
assert mod._extract_subagent_agent_name(title) == "some.dotted.agent"
|
||||||
|
|
||||||
|
def test_handles_trailing_whitespace(self, mod):
|
||||||
|
title = "Foo (@bar subagent) "
|
||||||
|
assert mod._extract_subagent_agent_name(title) == "bar"
|
||||||
|
|
||||||
|
def test_returns_sentinel_for_non_subagent_title(self, mod):
|
||||||
|
# Top-level sessions don't carry the ``(@... subagent)`` suffix;
|
||||||
|
# the walker is only called with descendants, but defensive
|
||||||
|
# handling means a misclassified record still archives.
|
||||||
|
assert mod._extract_subagent_agent_name("[AUTO-IMP] foo") == "unknown-subagent"
|
||||||
|
|
||||||
|
def test_returns_sentinel_for_empty(self, mod):
|
||||||
|
assert mod._extract_subagent_agent_name("") == "unknown-subagent"
|
||||||
|
assert mod._extract_subagent_agent_name(None) == "unknown-subagent"
|
||||||
|
|
||||||
|
|
||||||
|
class TestWalkSubagentDescendants:
|
||||||
|
"""``_walk_subagent_descendants`` walks the ``/session`` payload
|
||||||
|
BFS from a root session id and returns every descendant. Coverage:
|
||||||
|
(a) BFS order, (b) depth annotation, (c) transport error path,
|
||||||
|
(d) malformed payload path, (e) cycle-safety (just in case).
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _wire_sessions(self, httpserver, sessions: list[dict]):
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
sessions
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_returns_empty_for_no_descendants(self, mod, httpserver):
|
||||||
|
self._wire_sessions(
|
||||||
|
httpserver,
|
||||||
|
[{"id": "root", "parentID": "", "title": "[AUTO] root"}],
|
||||||
|
)
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
assert out == []
|
||||||
|
|
||||||
|
def test_collects_two_level_chain_in_bfs_order(self, mod, httpserver):
|
||||||
|
# root -> child1 -> grandchild; root -> child2.
|
||||||
|
# BFS order: child1, child2, grandchild.
|
||||||
|
self._wire_sessions(
|
||||||
|
httpserver,
|
||||||
|
[
|
||||||
|
{"id": "root", "parentID": "", "title": "root", "time": {"created": 1}},
|
||||||
|
{
|
||||||
|
"id": "child1",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Call A (@a subagent)",
|
||||||
|
"time": {"created": 10},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "child2",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Call B (@b subagent)",
|
||||||
|
"time": {"created": 20},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "grandchild",
|
||||||
|
"parentID": "child1",
|
||||||
|
"title": "Inner (@inner subagent)",
|
||||||
|
"time": {"created": 30},
|
||||||
|
},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
ids = [s["id"] for s in out]
|
||||||
|
assert ids == ["child1", "child2", "grandchild"]
|
||||||
|
depths = [s["_subagent_depth"] for s in out]
|
||||||
|
assert depths == [1, 1, 2]
|
||||||
|
|
||||||
|
def test_orders_siblings_by_creation_time(self, mod, httpserver):
|
||||||
|
# Two siblings created out-of-order in the /session response;
|
||||||
|
# BFS within a depth level should still be creation-time
|
||||||
|
# ordered so the trace reads left-to-right by clock.
|
||||||
|
self._wire_sessions(
|
||||||
|
httpserver,
|
||||||
|
[
|
||||||
|
{"id": "root", "parentID": "", "title": "root"},
|
||||||
|
{
|
||||||
|
"id": "later",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Later (@x subagent)",
|
||||||
|
"time": {"created": 999},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "earlier",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Earlier (@x subagent)",
|
||||||
|
"time": {"created": 1},
|
||||||
|
},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
assert [s["id"] for s in out] == ["earlier", "later"]
|
||||||
|
|
||||||
|
def test_returns_empty_on_transport_error(self, mod, httpserver):
|
||||||
|
# No /session route wired; httpserver returns 500 for the
|
||||||
|
# unmatched request and _request raises HTTPError which the
|
||||||
|
# walker catches.
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
assert out == []
|
||||||
|
|
||||||
|
def test_returns_empty_on_non_list_payload(self, mod, httpserver):
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
{"unexpected": "shape"}
|
||||||
|
)
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
assert out == []
|
||||||
|
|
||||||
|
def test_handles_cycle_without_infinite_loop(self, mod, httpserver):
|
||||||
|
# Defensive: OpenCode should never report a cycle, but the
|
||||||
|
# walker should terminate even if it ever does. Each id is
|
||||||
|
# only visited once thanks to the ``seen`` set.
|
||||||
|
self._wire_sessions(
|
||||||
|
httpserver,
|
||||||
|
[
|
||||||
|
{"id": "root", "parentID": "loop-end", "title": "root"},
|
||||||
|
{
|
||||||
|
"id": "loop-mid",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Mid (@x subagent)",
|
||||||
|
"time": {"created": 1},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "loop-end",
|
||||||
|
"parentID": "loop-mid",
|
||||||
|
"title": "End (@x subagent)",
|
||||||
|
"time": {"created": 2},
|
||||||
|
},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
out = mod._walk_subagent_descendants(_server_url(httpserver), "root")
|
||||||
|
# Only ``loop-mid`` and ``loop-end`` get visited; the back-edge
|
||||||
|
# to ``root`` is blocked by the seen-set so we don't recurse.
|
||||||
|
ids = [s["id"] for s in out]
|
||||||
|
assert ids == ["loop-mid", "loop-end"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestArchiveSubagentTree:
|
||||||
|
"""``_archive_subagent_tree`` orchestrates: walk → fetch each
|
||||||
|
descendant's messages → write archive per descendant. Coverage:
|
||||||
|
(a) writes one file per descendant, (b) filename infix carries
|
||||||
|
``sub<depth>``, (c) payload v2 fields populated, (d) parent's tag
|
||||||
|
+ started_at inherited, (e) per-descendant fetch failure logs and
|
||||||
|
continues.
|
||||||
|
|
||||||
|
These tests call ``_archive_subagent_tree`` directly (i.e. they
|
||||||
|
bypass ``run_session_blocking`` and therefore bypass
|
||||||
|
``_resolve_archive_dir``). Because of that we have to create the
|
||||||
|
archive directory ourselves — the autouse fixture only sets the
|
||||||
|
env var; ``_resolve_archive_dir`` is what calls ``mkdir`` on it.
|
||||||
|
"""
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _ensure_archive_dir_exists(self, archive_dir):
|
||||||
|
archive_dir.mkdir(parents=True, exist_ok=True)
|
||||||
|
|
||||||
|
def _wire(self, httpserver, descendants: list[dict], messages_for: dict[str, list]):
|
||||||
|
# Root session is always returned by /session GET — the walker
|
||||||
|
# uses it as a no-op (excludes root from descendants).
|
||||||
|
all_sessions = [{"id": "root", "parentID": "", "title": "root"}] + descendants
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
all_sessions
|
||||||
|
)
|
||||||
|
for sid, msgs in messages_for.items():
|
||||||
|
httpserver.expect_request(
|
||||||
|
f"/session/{sid}/message", method="GET"
|
||||||
|
).respond_with_json(msgs)
|
||||||
|
|
||||||
|
def test_writes_one_archive_per_descendant(
|
||||||
|
self, mod, httpserver, archive_dir
|
||||||
|
):
|
||||||
|
descendants = [
|
||||||
|
{
|
||||||
|
"id": "child1",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Dispatch (@tier-dispatcher subagent)",
|
||||||
|
"time": {"created": 1_000_000_000_000, "updated": 1_000_000_005_000},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "child2",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Estimate (@estimator-implementation subagent)",
|
||||||
|
"time": {"created": 1_000_000_001_000, "updated": 1_000_000_003_500},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
messages_for = {
|
||||||
|
"child1": [
|
||||||
|
{
|
||||||
|
"info": {
|
||||||
|
"id": "m1",
|
||||||
|
"role": "assistant",
|
||||||
|
"time": {"created": 1, "completed": 2},
|
||||||
|
"tokens": {"input": 10, "output": 5},
|
||||||
|
},
|
||||||
|
"parts": [{"type": "text", "text": "child1 said hi"}],
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"child2": [
|
||||||
|
{
|
||||||
|
"info": {
|
||||||
|
"id": "m2",
|
||||||
|
"role": "assistant",
|
||||||
|
"time": {"created": 1, "completed": 2},
|
||||||
|
"tokens": {"input": 20, "output": 10},
|
||||||
|
},
|
||||||
|
"parts": [{"type": "text", "text": "child2 said hi"}],
|
||||||
|
}
|
||||||
|
],
|
||||||
|
}
|
||||||
|
self._wire(httpserver, descendants, messages_for)
|
||||||
|
|
||||||
|
paths = mod._archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
root_session_id="root",
|
||||||
|
root_tag="AUTO-IMP-PR-99",
|
||||||
|
root_started_at_iso="2026-05-10T00:00:00+00:00",
|
||||||
|
)
|
||||||
|
assert len(paths) == 2
|
||||||
|
# Files exist on disk and parse.
|
||||||
|
for p in paths:
|
||||||
|
assert p.exists()
|
||||||
|
payload = json.loads(p.read_text())
|
||||||
|
assert payload["schema_version"] == 2
|
||||||
|
assert payload["parent_session_id"] == "root"
|
||||||
|
assert payload["tag"] == "AUTO-IMP-PR-99"
|
||||||
|
assert payload["subagent_depth"] == 1
|
||||||
|
assert payload["status"] == "subagent"
|
||||||
|
assert payload["subagent_title"].startswith(("Dispatch", "Estimate"))
|
||||||
|
|
||||||
|
def test_filename_contains_sub_depth_infix(
|
||||||
|
self, mod, httpserver, archive_dir
|
||||||
|
):
|
||||||
|
descendants = [
|
||||||
|
{
|
||||||
|
"id": "child1",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "X (@agent-a subagent)",
|
||||||
|
"time": {"created": 1_000_000_000_000, "updated": 1_000_000_001_000},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
messages_for = {"child1": []}
|
||||||
|
self._wire(httpserver, descendants, messages_for)
|
||||||
|
paths = mod._archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
root_session_id="root",
|
||||||
|
root_tag="AUTO-IMP-PR-99",
|
||||||
|
root_started_at_iso="2026-05-10T00:00:00+00:00",
|
||||||
|
)
|
||||||
|
assert len(paths) == 1
|
||||||
|
# Filename layout: <started_at>__sub01__<tag>__<agent>__<sid>.json
|
||||||
|
name = paths[0].name
|
||||||
|
assert "__sub01__" in name
|
||||||
|
assert "__AUTO-IMP-PR-99__" in name
|
||||||
|
assert "__agent-a__" in name
|
||||||
|
assert name.endswith("__child1.json")
|
||||||
|
|
||||||
|
def test_continues_after_per_descendant_fetch_failure(
|
||||||
|
self, mod, httpserver, archive_dir, caplog
|
||||||
|
):
|
||||||
|
# Wire two descendants but only respond to one's /message GET
|
||||||
|
# — the other returns 500 and the walker must log + continue.
|
||||||
|
descendants = [
|
||||||
|
{
|
||||||
|
"id": "ok",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "OK (@ok-agent subagent)",
|
||||||
|
"time": {"created": 1_000_000_000_000, "updated": 1_000_000_001_000},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "broken",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Broken (@broken-agent subagent)",
|
||||||
|
"time": {"created": 1_000_000_002_000, "updated": 1_000_000_003_000},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
messages_for = {"ok": []} # ``broken`` deliberately omitted
|
||||||
|
self._wire(httpserver, descendants, messages_for)
|
||||||
|
|
||||||
|
import logging
|
||||||
|
caplog.set_level(logging.WARNING, logger="opencode_worker")
|
||||||
|
paths = mod._archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
root_session_id="root",
|
||||||
|
root_tag="AUTO-IMP-PR-99",
|
||||||
|
root_started_at_iso="2026-05-10T00:00:00+00:00",
|
||||||
|
)
|
||||||
|
# Only ``ok`` archives; ``broken`` is skipped with a warning.
|
||||||
|
assert len(paths) == 1
|
||||||
|
assert any(
|
||||||
|
"subagent archive: skipping broken" in rec.getMessage()
|
||||||
|
for rec in caplog.records
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_redact_values_propagate_to_subagent_archive(
|
||||||
|
self, mod, httpserver, archive_dir
|
||||||
|
):
|
||||||
|
# The PAT in a subagent's tool input must also be redacted —
|
||||||
|
# the dispatcher passes redact_values=[cfg.token] through.
|
||||||
|
secret = "FORGEJO-PAT-LONG-SECRET-VALUE-NEVER-COMMIT-3"
|
||||||
|
descendants = [
|
||||||
|
{
|
||||||
|
"id": "child1",
|
||||||
|
"parentID": "root",
|
||||||
|
"title": "Clone (@git-isolator-util subagent)",
|
||||||
|
"time": {"created": 1_000_000_000_000, "updated": 1_000_000_001_000},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
messages_for = {
|
||||||
|
"child1": [
|
||||||
|
{
|
||||||
|
"info": {
|
||||||
|
"id": "m1",
|
||||||
|
"role": "assistant",
|
||||||
|
"time": {"created": 1, "completed": 2},
|
||||||
|
"tokens": {"input": 10, "output": 0},
|
||||||
|
},
|
||||||
|
"parts": [
|
||||||
|
{
|
||||||
|
"type": "tool",
|
||||||
|
"tool": "bash",
|
||||||
|
"state": {
|
||||||
|
"status": "completed",
|
||||||
|
"input": {
|
||||||
|
"command": (
|
||||||
|
f"git clone https://{secret}@host/x.git /tmp/x"
|
||||||
|
)
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
],
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
self._wire(httpserver, descendants, messages_for)
|
||||||
|
paths = mod._archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
root_session_id="root",
|
||||||
|
root_tag="AUTO-IMP-PR-99",
|
||||||
|
root_started_at_iso="2026-05-10T00:00:00+00:00",
|
||||||
|
redact_values=[secret],
|
||||||
|
)
|
||||||
|
assert len(paths) == 1
|
||||||
|
body = paths[0].read_text()
|
||||||
|
assert secret not in body, "secret leaked into subagent archive"
|
||||||
|
assert "<REDACTED>" in body
|
||||||
|
|
||||||
|
def test_returns_empty_when_no_descendants_to_walk(
|
||||||
|
self, mod, httpserver, archive_dir
|
||||||
|
):
|
||||||
|
# Only the root in /session response → nothing to archive.
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
[{"id": "root", "parentID": "", "title": "root"}]
|
||||||
|
)
|
||||||
|
paths = mod._archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
root_session_id="root",
|
||||||
|
root_tag="AUTO-IMP-PR-99",
|
||||||
|
root_started_at_iso="2026-05-10T00:00:00+00:00",
|
||||||
|
)
|
||||||
|
assert paths == []
|
||||||
|
assert list(archive_dir.glob("*.json")) == []
|
||||||
|
|
||||||
|
|
||||||
|
class TestEndToEndSubagentArchive:
|
||||||
|
"""End-to-end: a real ``run_session_blocking`` call that emits a
|
||||||
|
subagent descendant in /session → both archives land on disk."""
|
||||||
|
|
||||||
|
def test_run_session_archives_root_and_subagents(
|
||||||
|
self, mod, httpserver, archive_dir
|
||||||
|
):
|
||||||
|
# We deliberately do NOT use ``_wire_lifecycle_with_statuses``
|
||||||
|
# here because that helper wires its own ``/session`` GET
|
||||||
|
# (returning only the root) and pytest-httpserver matches
|
||||||
|
# ``expect_request`` registrations in FIFO order — the helper's
|
||||||
|
# registration would win over our descendant-bearing override.
|
||||||
|
# Wire the full lifecycle by hand instead so the walker sees
|
||||||
|
# the subagent we care about.
|
||||||
|
httpserver.expect_request("/session", method="POST").respond_with_json(
|
||||||
|
_session_obj("s-obs-1")
|
||||||
|
)
|
||||||
|
httpserver.expect_request(
|
||||||
|
"/session/s-obs-1/prompt_async", method="POST"
|
||||||
|
).respond_with_json({})
|
||||||
|
# status sequence: busy then idle
|
||||||
|
httpserver.expect_oneshot_request(
|
||||||
|
"/session/status", method="GET"
|
||||||
|
).respond_with_json({"s-obs-1": {"type": "busy"}})
|
||||||
|
httpserver.expect_oneshot_request(
|
||||||
|
"/session/status", method="GET"
|
||||||
|
).respond_with_json({"s-obs-1": {"type": "idle"}})
|
||||||
|
httpserver.expect_request(
|
||||||
|
"/session/s-obs-1/message", method="GET"
|
||||||
|
).respond_with_json(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"info": {
|
||||||
|
"id": "m-final",
|
||||||
|
"role": "assistant",
|
||||||
|
"time": {"created": 1000, "completed": 5000},
|
||||||
|
"tokens": {"input": 100, "output": 20, "reasoning": 5},
|
||||||
|
},
|
||||||
|
"parts": [
|
||||||
|
{
|
||||||
|
"type": "text",
|
||||||
|
"text": '{"outcome":"resolved","files_touched":[]}',
|
||||||
|
}
|
||||||
|
],
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
# ``/session`` GET returns BOTH the root and one subagent
|
||||||
|
# descendant. The walker excludes the root and archives the
|
||||||
|
# subagent.
|
||||||
|
httpserver.expect_request("/session", method="GET").respond_with_json(
|
||||||
|
[
|
||||||
|
_session_obj("s-obs-1"),
|
||||||
|
{
|
||||||
|
"id": "ses-sub-1",
|
||||||
|
"parentID": "s-obs-1",
|
||||||
|
"title": "Dispatch (@tier-dispatcher subagent)",
|
||||||
|
"time": {
|
||||||
|
"created": 1_000_000_000_000,
|
||||||
|
"updated": 1_000_000_001_000,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
)
|
||||||
|
# The walker fetches the descendant's /message — wire it.
|
||||||
|
httpserver.expect_request(
|
||||||
|
"/session/ses-sub-1/message", method="GET"
|
||||||
|
).respond_with_json(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"info": {
|
||||||
|
"id": "sm1",
|
||||||
|
"role": "assistant",
|
||||||
|
"time": {"created": 0, "completed": 1},
|
||||||
|
"tokens": {"input": 5, "output": 1},
|
||||||
|
},
|
||||||
|
"parts": [{"type": "text", "text": "sub said done"}],
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
httpserver.expect_request(
|
||||||
|
"/session/s-obs-1", method="DELETE"
|
||||||
|
).respond_with_data("")
|
||||||
|
result = mod.run_session_blocking(
|
||||||
|
server_url=_server_url(httpserver),
|
||||||
|
agent="test-agent",
|
||||||
|
tag="AUTO-OBS",
|
||||||
|
prompt="x",
|
||||||
|
timeout_seconds=30,
|
||||||
|
poll_interval_seconds=0.01,
|
||||||
|
)
|
||||||
|
assert result.status == "completed"
|
||||||
|
files = sorted(archive_dir.glob("*.json"))
|
||||||
|
# One root + one subagent.
|
||||||
|
assert len(files) == 2
|
||||||
|
names = [p.name for p in files]
|
||||||
|
# Sorting puts ``2026-...`` (root, started_at first) before
|
||||||
|
# ``<iso of sub started_at>``. The depth-1 infix appears in
|
||||||
|
# the subagent's filename.
|
||||||
|
assert any("__sub01__" in n for n in names)
|
||||||
|
# Validate the subagent payload.
|
||||||
|
sub_file = next(p for p in files if "__sub01__" in p.name)
|
||||||
|
sub_payload = json.loads(sub_file.read_text())
|
||||||
|
assert sub_payload["schema_version"] == 2
|
||||||
|
assert sub_payload["parent_session_id"] == "s-obs-1"
|
||||||
|
assert sub_payload["agent"] == "tier-dispatcher"
|
||||||
|
assert sub_payload["tag"] == "AUTO-OBS"
|
||||||
|
assert sub_payload["subagent_depth"] == 1
|
||||||
|
|||||||
+285
-7
@@ -535,7 +535,29 @@ _ARCHIVE_DISABLED_ENV = "OPENCODE_WORKER_ARCHIVE_DISABLED"
|
|||||||
# Schema version for the archive payload. Bump when the on-disk JSON
|
# Schema version for the archive payload. Bump when the on-disk JSON
|
||||||
# shape changes; the telemetry console reads this to decide which
|
# shape changes; the telemetry console reads this to decide which
|
||||||
# fields to render.
|
# fields to render.
|
||||||
_ARCHIVE_SCHEMA_VERSION = 1
|
#
|
||||||
|
# v1 (initial) — top-level wrapper sessions only; carries
|
||||||
|
# session_id, agent, tag, status, started_at,
|
||||||
|
# archived_at, wallclock_seconds, per_turn,
|
||||||
|
# state_history, messages.
|
||||||
|
# v2 (2026-05-10) — adds optional ``parent_session_id``,
|
||||||
|
# ``subagent_title``, and ``subagent_depth`` fields so
|
||||||
|
# the dispatcher can archive the entire ``task`` tool
|
||||||
|
# subagent tree before deleting the root session. The
|
||||||
|
# first three fields are absent / null on top-level
|
||||||
|
# (wrapper) archives so a v1 reader treating them as
|
||||||
|
# missing is forward-compatible.
|
||||||
|
_ARCHIVE_SCHEMA_VERSION = 2
|
||||||
|
|
||||||
|
# Pattern that matches the trailing ``(@<agent> subagent)`` suffix
|
||||||
|
# OpenCode appends to a subagent session's title (e.g.
|
||||||
|
# ``"Implement PR fix #30 (@task-implementor subagent)"``). The capture
|
||||||
|
# group is the agent name and is what we use as the ``agent`` field in
|
||||||
|
# the subagent archive payload. If the title does not match — which
|
||||||
|
# would mean OpenCode changed its title convention or this is a
|
||||||
|
# manually-titled top-level session — we fall back to a sentinel string
|
||||||
|
# so the archive still lands on disk.
|
||||||
|
_SUBAGENT_TITLE_RE = re.compile(r"\(@([\w.-]+)\s+subagent\)\s*$")
|
||||||
|
|
||||||
|
|
||||||
def _resolve_archive_dir() -> Path | None:
|
def _resolve_archive_dir() -> Path | None:
|
||||||
@@ -590,16 +612,28 @@ def _build_archive_payload(
|
|||||||
wallclock_seconds: float,
|
wallclock_seconds: float,
|
||||||
messages: list[dict[str, Any]],
|
messages: list[dict[str, Any]],
|
||||||
state_history: list[dict[str, Any]] | None = None,
|
state_history: list[dict[str, Any]] | None = None,
|
||||||
|
parent_session_id: str | None = None,
|
||||||
|
subagent_title: str | None = None,
|
||||||
|
subagent_depth: int | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
"""Build the JSON-serialisable archive payload for a session.
|
"""Build the JSON-serialisable archive payload for a session.
|
||||||
|
|
||||||
Pulled out into a helper so :func:`_archive_session` and the
|
Pulled out into a helper so :func:`_archive_session` and the
|
||||||
test suite share the same shape; the on-disk schema is therefore
|
test suite share the same shape; the on-disk schema is therefore
|
||||||
derivable from this function alone.
|
derivable from this function alone.
|
||||||
|
|
||||||
|
``parent_session_id`` / ``subagent_title`` / ``subagent_depth`` are
|
||||||
|
populated only for subagent archives (sessions spawned by a parent's
|
||||||
|
``task`` tool call). For the top-level wrapper session they remain
|
||||||
|
``None`` and serialise as ``null`` so a downstream reader can tell
|
||||||
|
the two apart without checking the schema version.
|
||||||
"""
|
"""
|
||||||
return {
|
return {
|
||||||
"schema_version": _ARCHIVE_SCHEMA_VERSION,
|
"schema_version": _ARCHIVE_SCHEMA_VERSION,
|
||||||
"session_id": session_id,
|
"session_id": session_id,
|
||||||
|
"parent_session_id": parent_session_id,
|
||||||
|
"subagent_title": subagent_title,
|
||||||
|
"subagent_depth": subagent_depth,
|
||||||
"agent": agent,
|
"agent": agent,
|
||||||
"tag": tag,
|
"tag": tag,
|
||||||
"status": status,
|
"status": status,
|
||||||
@@ -667,6 +701,9 @@ def _archive_session(
|
|||||||
messages: list[dict[str, Any]],
|
messages: list[dict[str, Any]],
|
||||||
state_history: list[dict[str, Any]] | None = None,
|
state_history: list[dict[str, Any]] | None = None,
|
||||||
redact_values: Iterable[str] | None = None,
|
redact_values: Iterable[str] | None = None,
|
||||||
|
parent_session_id: str | None = None,
|
||||||
|
subagent_title: str | None = None,
|
||||||
|
subagent_depth: int | None = None,
|
||||||
) -> Path | None:
|
) -> Path | None:
|
||||||
"""Write the session archive JSON file. Returns the resolved path
|
"""Write the session archive JSON file. Returns the resolved path
|
||||||
on success, or ``None`` on any I/O / serialisation failure.
|
on success, or ``None`` on any I/O / serialisation failure.
|
||||||
@@ -677,13 +714,25 @@ def _archive_session(
|
|||||||
disk. The redaction runs on the serialised JSON so it catches the
|
disk. The redaction runs on the serialised JSON so it catches the
|
||||||
value in any nested location (prompt text, tool input, error
|
value in any nested location (prompt text, tool input, error
|
||||||
string). See :func:`_redact_secret_values` for the safety floor.
|
string). See :func:`_redact_secret_values` for the safety floor.
|
||||||
|
|
||||||
|
``parent_session_id`` / ``subagent_title`` / ``subagent_depth`` are
|
||||||
|
forwarded into the payload to support subagent archiving (see
|
||||||
|
:func:`_archive_subagent_tree`). For top-level (wrapper) archives
|
||||||
|
they should remain ``None``.
|
||||||
"""
|
"""
|
||||||
fname = (
|
fname_parts = [
|
||||||
f"{_safe_filename_component(started_at_iso)}__"
|
_safe_filename_component(started_at_iso),
|
||||||
f"{_safe_filename_component(tag)}__"
|
_safe_filename_component(tag),
|
||||||
f"{_safe_filename_component(agent)}__"
|
_safe_filename_component(agent),
|
||||||
f"{_safe_filename_component(session_id)}.json"
|
_safe_filename_component(session_id),
|
||||||
)
|
]
|
||||||
|
# Subagent archives get a ``sub<depth>`` infix so a directory listing
|
||||||
|
# sorts the wrapper first, then its subagents in BFS order
|
||||||
|
# (sub1 before sub2 before sub3...). This makes ``ls -1`` a usable
|
||||||
|
# post-mortem trace tool without needing the telemetry console.
|
||||||
|
if subagent_depth is not None:
|
||||||
|
fname_parts.insert(1, f"sub{subagent_depth:02d}")
|
||||||
|
fname = "__".join(fname_parts) + ".json"
|
||||||
path = archive_dir / fname
|
path = archive_dir / fname
|
||||||
payload = _build_archive_payload(
|
payload = _build_archive_payload(
|
||||||
session_id=session_id,
|
session_id=session_id,
|
||||||
@@ -694,6 +743,9 @@ def _archive_session(
|
|||||||
wallclock_seconds=wallclock_seconds,
|
wallclock_seconds=wallclock_seconds,
|
||||||
messages=messages,
|
messages=messages,
|
||||||
state_history=state_history,
|
state_history=state_history,
|
||||||
|
parent_session_id=parent_session_id,
|
||||||
|
subagent_title=subagent_title,
|
||||||
|
subagent_depth=subagent_depth,
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
# Render with a default=str so anything urllib hands back that
|
# Render with a default=str so anything urllib hands back that
|
||||||
@@ -708,6 +760,212 @@ def _archive_session(
|
|||||||
return path
|
return path
|
||||||
|
|
||||||
|
|
||||||
|
def _walk_subagent_descendants(
|
||||||
|
server_url: str,
|
||||||
|
root_session_id: str,
|
||||||
|
) -> list[dict[str, Any]]:
|
||||||
|
"""Walk OpenCode's ``GET /session`` response BFS starting from
|
||||||
|
``root_session_id`` and return every descendant session record.
|
||||||
|
|
||||||
|
The returned list excludes ``root_session_id`` itself and is sorted
|
||||||
|
by depth-then-creation-time, which matches the natural order a
|
||||||
|
human reads a trace ("the wrapper called tier-dispatcher first,
|
||||||
|
which called the estimator, which called…"). Each record is the
|
||||||
|
raw session dict from OpenCode with one added field:
|
||||||
|
``_subagent_depth`` (1-based distance from the root).
|
||||||
|
|
||||||
|
Returns ``[]`` on any transport error or unexpected response
|
||||||
|
shape — the caller's archive flow must continue regardless, so we
|
||||||
|
never raise here. A short warning is logged so a missed walk is
|
||||||
|
visible in the dispatcher log.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
all_sessions = _request("GET", f"{server_url}/session")
|
||||||
|
except _TRANSPORT_EXC as e:
|
||||||
|
logger.warning(
|
||||||
|
"subagent walk failed (transport): %s; archiving root only", e
|
||||||
|
)
|
||||||
|
return []
|
||||||
|
if not isinstance(all_sessions, list):
|
||||||
|
logger.warning(
|
||||||
|
"subagent walk: /session returned unexpected shape "
|
||||||
|
"(%s); archiving root only",
|
||||||
|
type(all_sessions).__name__,
|
||||||
|
)
|
||||||
|
return []
|
||||||
|
|
||||||
|
# parent_id -> [child_session_dict, ...]
|
||||||
|
children_by_parent: dict[str, list[dict[str, Any]]] = {}
|
||||||
|
for s in all_sessions:
|
||||||
|
if not isinstance(s, dict):
|
||||||
|
continue
|
||||||
|
pid = s.get("parentID") or s.get("parent_id") or ""
|
||||||
|
if not pid:
|
||||||
|
continue
|
||||||
|
children_by_parent.setdefault(pid, []).append(s)
|
||||||
|
# Sort children by creation time so the trace order is stable.
|
||||||
|
for kids in children_by_parent.values():
|
||||||
|
kids.sort(
|
||||||
|
key=lambda k: (k.get("time") or {}).get("created", 0) or 0
|
||||||
|
)
|
||||||
|
|
||||||
|
descendants: list[dict[str, Any]] = []
|
||||||
|
seen: set[str] = {root_session_id}
|
||||||
|
# BFS frontier: list of (session_id, depth).
|
||||||
|
frontier: list[tuple[str, int]] = [(root_session_id, 0)]
|
||||||
|
while frontier:
|
||||||
|
parent_id, parent_depth = frontier.pop(0)
|
||||||
|
for child in children_by_parent.get(parent_id, []):
|
||||||
|
cid = child.get("id")
|
||||||
|
if not isinstance(cid, str) or cid in seen:
|
||||||
|
continue
|
||||||
|
seen.add(cid)
|
||||||
|
depth = parent_depth + 1
|
||||||
|
child_with_depth = dict(child)
|
||||||
|
child_with_depth["_subagent_depth"] = depth
|
||||||
|
descendants.append(child_with_depth)
|
||||||
|
frontier.append((cid, depth))
|
||||||
|
return descendants
|
||||||
|
|
||||||
|
|
||||||
|
def _extract_subagent_agent_name(title: str | None) -> str:
|
||||||
|
"""Return the ``agent`` name from a subagent session title.
|
||||||
|
|
||||||
|
Subagent sessions in OpenCode are titled
|
||||||
|
``"<some description> (@<agent-name> subagent)"`` — see
|
||||||
|
:data:`_SUBAGENT_TITLE_RE`. Returns ``"unknown-subagent"`` if the
|
||||||
|
title is missing or does not match the convention, so the archive
|
||||||
|
still lands on disk with a useful (if generic) filename component.
|
||||||
|
"""
|
||||||
|
if not isinstance(title, str) or not title:
|
||||||
|
return "unknown-subagent"
|
||||||
|
m = _SUBAGENT_TITLE_RE.search(title)
|
||||||
|
if not m:
|
||||||
|
return "unknown-subagent"
|
||||||
|
return m.group(1)
|
||||||
|
|
||||||
|
|
||||||
|
def _ms_to_iso(ms: int | float | None) -> str | None:
|
||||||
|
"""Convert a Unix-epoch millisecond timestamp to an ISO-8601
|
||||||
|
UTC string, or return ``None`` if ``ms`` is not a usable number.
|
||||||
|
|
||||||
|
OpenCode's ``/session`` payload reports ``time.created`` /
|
||||||
|
``time.updated`` in milliseconds since the epoch. We carry that
|
||||||
|
forward into the archive as an ISO string so the on-disk format
|
||||||
|
matches the wrapper archive's ``started_at`` field.
|
||||||
|
"""
|
||||||
|
if ms is None:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
seconds = float(ms) / 1000.0
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
return None
|
||||||
|
return _dt.datetime.fromtimestamp(
|
||||||
|
seconds, tz=_dt.timezone.utc
|
||||||
|
).isoformat()
|
||||||
|
|
||||||
|
|
||||||
|
def _archive_subagent_tree(
|
||||||
|
*,
|
||||||
|
archive_dir: Path,
|
||||||
|
server_url: str,
|
||||||
|
root_session_id: str,
|
||||||
|
root_tag: str,
|
||||||
|
root_started_at_iso: str,
|
||||||
|
redact_values: Iterable[str] | None = None,
|
||||||
|
) -> list[Path]:
|
||||||
|
"""Walk every subagent descendant of ``root_session_id``, fetch
|
||||||
|
its message stream, and write an archive file per session.
|
||||||
|
|
||||||
|
Returns the list of archive paths that were written (in BFS order).
|
||||||
|
Always runs best-effort: a transport error on any single
|
||||||
|
subagent's message fetch logs a warning and continues to the next
|
||||||
|
descendant. The caller MUST invoke this BEFORE deleting the root
|
||||||
|
session — OpenCode may garbage-collect subagent sessions when the
|
||||||
|
root is removed.
|
||||||
|
|
||||||
|
``root_tag`` and ``root_started_at_iso`` are inherited so a
|
||||||
|
directory listing groups every session from one dispatcher cycle
|
||||||
|
together. The subagent's own per-session ``time.created`` is
|
||||||
|
captured under the new ``subagent_started_at`` slot for the
|
||||||
|
timeline view in the telemetry console.
|
||||||
|
"""
|
||||||
|
descendants = _walk_subagent_descendants(server_url, root_session_id)
|
||||||
|
if not descendants:
|
||||||
|
return []
|
||||||
|
paths: list[Path] = []
|
||||||
|
for desc in descendants:
|
||||||
|
sid = desc.get("id")
|
||||||
|
if not isinstance(sid, str) or not sid:
|
||||||
|
continue
|
||||||
|
title = desc.get("title") or ""
|
||||||
|
agent_name = _extract_subagent_agent_name(title)
|
||||||
|
depth = int(desc.get("_subagent_depth") or 1)
|
||||||
|
time_info = desc.get("time") or {}
|
||||||
|
created_ms = time_info.get("created")
|
||||||
|
updated_ms = time_info.get("updated")
|
||||||
|
sub_started_iso = _ms_to_iso(created_ms) or root_started_at_iso
|
||||||
|
sub_wallclock = 0.0
|
||||||
|
if isinstance(created_ms, (int, float)) and isinstance(
|
||||||
|
updated_ms, (int, float)
|
||||||
|
):
|
||||||
|
sub_wallclock = max(
|
||||||
|
0.0, (float(updated_ms) - float(created_ms)) / 1000.0
|
||||||
|
)
|
||||||
|
|
||||||
|
try:
|
||||||
|
messages = _request(
|
||||||
|
"GET", f"{server_url}/session/{sid}/message"
|
||||||
|
)
|
||||||
|
except _TRANSPORT_EXC as e:
|
||||||
|
logger.warning(
|
||||||
|
"subagent archive: skipping %s (%s) — transport: %s",
|
||||||
|
sid, agent_name, e,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
if not isinstance(messages, list):
|
||||||
|
logger.warning(
|
||||||
|
"subagent archive: skipping %s (%s) — unexpected "
|
||||||
|
"messages shape %s",
|
||||||
|
sid, agent_name, type(messages).__name__,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
|
||||||
|
# Inherit the parent's tag so a ``ls -1 .dispatcher-logs/sessions``
|
||||||
|
# listing groups every session from this dispatcher cycle
|
||||||
|
# together. The subagent's own creation time goes into
|
||||||
|
# ``subagent_started_at`` inside the payload.
|
||||||
|
try:
|
||||||
|
written = _archive_session(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
session_id=sid,
|
||||||
|
agent=agent_name,
|
||||||
|
tag=root_tag,
|
||||||
|
status="subagent",
|
||||||
|
started_at_iso=sub_started_iso,
|
||||||
|
wallclock_seconds=sub_wallclock,
|
||||||
|
messages=messages,
|
||||||
|
state_history=None,
|
||||||
|
redact_values=redact_values,
|
||||||
|
parent_session_id=desc.get("parentID") or root_session_id,
|
||||||
|
subagent_title=title,
|
||||||
|
subagent_depth=depth,
|
||||||
|
)
|
||||||
|
except Exception as e: # noqa: BLE001 — best-effort, never raise
|
||||||
|
logger.warning(
|
||||||
|
"subagent archive raised for %s (%s): %s",
|
||||||
|
sid, agent_name, e,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
if written is not None:
|
||||||
|
paths.append(written)
|
||||||
|
logger.info(
|
||||||
|
"subagent %s (depth=%d agent=%s) archived to %s",
|
||||||
|
sid, depth, agent_name, written,
|
||||||
|
)
|
||||||
|
return paths
|
||||||
|
|
||||||
|
|
||||||
def _summarize_per_turn(messages: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
def _summarize_per_turn(messages: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||||
"""Aggregate per-assistant-turn metrics out of a /session/{id}/message
|
"""Aggregate per-assistant-turn metrics out of a /session/{id}/message
|
||||||
payload.
|
payload.
|
||||||
@@ -1442,6 +1700,26 @@ def run_session_blocking(
|
|||||||
"session archive raised for %s: %s",
|
"session archive raised for %s: %s",
|
||||||
session_id, e,
|
session_id, e,
|
||||||
)
|
)
|
||||||
|
# Subagent tree archive — runs AFTER the root archive
|
||||||
|
# but BEFORE the DELETE below, so the dispatcher
|
||||||
|
# captures the complete ``task`` tool chain even if
|
||||||
|
# OpenCode garbage-collects subagent sessions on
|
||||||
|
# root deletion. Best-effort: any failure logs a
|
||||||
|
# warning and lets the DELETE proceed.
|
||||||
|
try:
|
||||||
|
_archive_subagent_tree(
|
||||||
|
archive_dir=archive_dir,
|
||||||
|
server_url=server_url,
|
||||||
|
root_session_id=session_id,
|
||||||
|
root_tag=tag,
|
||||||
|
root_started_at_iso=started_at_iso,
|
||||||
|
redact_values=redact_snapshot,
|
||||||
|
)
|
||||||
|
except Exception as e: # noqa: BLE001 — defensive belt-and-braces
|
||||||
|
logger.warning(
|
||||||
|
"subagent tree archive raised for root %s: %s",
|
||||||
|
session_id, e,
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
_request("DELETE", f"{server_url}/session/{session_id}")
|
_request("DELETE", f"{server_url}/session/{session_id}")
|
||||||
except _TRANSPORT_EXC as e:
|
except _TRANSPORT_EXC as e:
|
||||||
|
|||||||
Reference in New Issue
Block a user