Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion CONTEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -2503,7 +2503,8 @@ this instance that the client never sent and so has no row of its own to anchor
those two lanes this read is the only thing that carries any of it (the wire tags an instance
on the four events of a *direct* turn and nothing else). Addressed by
`(session_key, agent, handle)` through a second live index, because the first one is keyed by
the record's directory - a task id no reader of a *conversation* ever sees.
the record's address - the conversation's node root plus a node id no reader of a
*conversation* ever sees.
_Avoid_: reading the absence of live rows as "the turn ended" - a transport with no per-step
visibility reports none for the whole of every turn.

Expand Down
6 changes: 4 additions & 2 deletions docs/specs/2026-09-18-desk-tasks-list-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,10 @@ stores it uncapped; a spawn's is the head of `.error.md`, or of `.out.md` when `
the files are read off the activity being collected for it in this process (the same
in-memory account `subagent.context` and `dag.node` serve a transcript from), since the record
on disk carries them only once the run finishes; what the lane has not reported yet stays null,
and a node that is not running takes nothing from that index (its key is a record id unique per
conversation only). A dag node that has finished while its run has not keeps the account the
and a node that is not running takes nothing from that index. The index is one per process and
keyed by the record's address (`history.spawn_live_key`: the conversation's node root plus the
id; `dag_store.node_live_key`: the run plus the node), so two conversations that named a spawn
alike never share an entry. A dag node that has finished while its run has not keeps the account the
runner set aside for it at its end (`activity.record_settled`) until the manifest is written,
so its usage does not vanish between the two.

Expand Down
24 changes: 14 additions & 10 deletions raven/agent/subagent/activity.py
Original file line number Diff line number Diff line change
Expand Up @@ -196,13 +196,17 @@ def as_meta(self) -> dict[str, Any]:
_current: ContextVar[RunActivity | None] = ContextVar("raven_subagent_activity", default=None)

_live: dict[str, RunActivity] = {}
"""Runs being collected right now, keyed by their record's call id.
"""Runs being collected right now, keyed by their record's address.

The disk record is written when the run finishes, so while it is in flight the
only account of it lives in the ``RunActivity`` being collected. This index is
what lets ``subagent.context`` serve that account to a panel watching the run,
instead of a prompt and nothing until the end. Entries live exactly as long as
their ``collecting`` block."""
what lets ``subagent.context`` and ``tasks.list`` serve that account to a panel
watching the run, instead of a prompt and nothing until the end. Entries live
exactly as long as their ``collecting`` block. The key names the record, not
the run's own id alone: a spawn's id is unique for one conversation only
(``history.spawn_live_key``), and a dag node's carries its run
(``dag_store.node_live_key``) -- the index is one per process, and two
conversations must not share an entry."""


_settled: dict[str, dict[str, Any]] = {}
Expand Down Expand Up @@ -235,12 +239,12 @@ def forget_settled(keys: Iterable[str]) -> None:
_live_instances: dict[tuple[str, str, str], RunActivity] = {}
"""The same activities, addressed the way a *conversation* reader has to ask.

``_live`` is keyed by the record's own directory, which is what the panel
watching one call already holds. A reader of an instance's conversation holds
``(session_key, agent, handle)`` and nothing else -- the record name is a task id
it never saw -- so the same run is indexed twice rather than having that reader
guess at a directory layout. Entries live exactly as long as their ``collecting``
block, as ``_live``'s do."""
``_live`` is keyed by the record's address -- a node root the panel watching
one call derives from its session, plus an id it holds. A reader of an
instance's conversation holds ``(session_key, agent, handle)`` and nothing else
-- the node id is one it never saw -- so the same run is indexed twice rather
than having that reader guess at a directory layout. Entries live exactly as
long as their ``collecting`` block, as ``_live``'s do."""


@contextmanager
Expand Down
7 changes: 4 additions & 3 deletions raven/agent/subagent/dag_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -410,9 +410,10 @@ def make_run_id() -> str:
def node_live_key(run_id: str, node_id: str) -> str:
"""The live-index key a node's activity is collected under.

Shared with the reader rather than spelled out on both sides: a spawn keys
its activity by the record directory's name, and a node has no such
directory, so the two namespaces are kept apart by this prefix.
Shared with the reader rather than spelled out on both sides. A spawn keys
its activity by its node root plus its id (``history.spawn_live_key``); a
node has no root of its own, so it keys by its run plus its id -- and the
two prefixes keep one process-wide index's namespaces apart.
"""
return f"dag:{run_id}:{node_id}"

Expand Down
15 changes: 15 additions & 0 deletions raven/agent/subagent/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,21 @@ def nodes_root(session_dir: Path) -> Path:
return session_history_root(session_dir) / _NODES_DIRNAME


def spawn_live_key(root: Path, node_id: str) -> str:
"""The live-index key a spawn's activity is collected under.

The record's own address -- the node root of its conversation plus its
id -- rather than the id alone: a node id is unique for one conversation
only, and the live index is one per process, so two conversations that
named a spawn alike would otherwise share an entry (and the first to
finish would drop the other's). Shared with the readers the way the dag
side shares ``node_live_key``; the writer holds the root as
``SpawnRecord.dir`` and a reader as the node files' root, which both
resolve from the same session directory.
"""
return f"spawn:{root}:{node_id}"


def node_file_in(root: Path, node_id: str, name: str) -> Path:
"""One artifact of one node, under a node root the caller already holds.

Expand Down
13 changes: 7 additions & 6 deletions raven/agent/subagent/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@
DirectTurnMeta,
NotAddressableError,
)
from raven.agent.subagent.history import SpawnRecord, session_history_root
from raven.agent.subagent.history import SpawnRecord, session_history_root, spawn_live_key
from raven.agent.subagent.instance_state import InstanceState, instance_state_path
from raven.agent.subagent.instances import get_registry, hold_handle, mint_handle
from raven.agent.subagent.mode_tiers import resolve_tier, turn_tier_in_force
Expand Down Expand Up @@ -1609,10 +1609,11 @@ async def _run_subagent_inner(
# publishing into nothing is a no-op. Every exit below therefore has the
# tool calls and token cost the run got as far as producing -- a failed
# run's are the ones worth keeping. Keyed into the live index by the
# record's own directory name, so `subagent.context` can serve the run
# while it is still in flight, and by instance so the conversation view
# can: a spawned call is a turn of the same instance a direct chat talks
# to, and watching it there is the same question.
# record's own address (`spawn_live_key`), so `subagent.context` and
# `tasks.list` can serve the run while it is still in flight, and by
# instance so the conversation view can: a spawned call is a turn of
# the same instance a direct chat talks to, and watching it there is
# the same question.
cancelled = False
# The record's own id, not its directory's name: the artifacts are a
# filename prefix in the shared node root now, so the directory names
Expand All @@ -1622,7 +1623,7 @@ async def _run_subagent_inner(
# queued behind a direct chat to the same instance is not that
# instance's turn yet, and registering it here took the slot from the
# turn that was (see ``activity.collecting``).
with activity.collecting(live_key=call_id, prompt=task) as did:
with activity.collecting(live_key=spawn_live_key(record.dir, call_id), prompt=task) as did:
try:
backend = self._resolve_backend(agent)
# The same message list a direct chat to this handle would carry.
Expand Down
7 changes: 4 additions & 3 deletions raven/rpc/methods/subagent.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@
from raven.agent.subagent import activity as run_activity
from raven.agent.subagent.dag_live import live_run_ids
from raven.agent.subagent.dag_store import REGISTRY_FILENAME
from raven.agent.subagent.history import dag_root, nodes_root, session_history_root
from raven.agent.subagent.history import dag_root, nodes_root, session_history_root, spawn_live_key
from raven.agent.subagent.instances import get_registry
from raven.agent.subagent.tool_vocabulary import normalize_row
from raven.config.loader import load_config
Expand Down Expand Up @@ -464,15 +464,16 @@ async def subagent_context(
# collected in this very process, so a watching panel reads that. Copied
# entry-by-entry because the collector republishes on every update from the
# agent, and a list mutated mid-iteration is a crash in a read-only path.
if not transcribed and (live := run_activity.live(files.node_id)) is not None:
live_key = spawn_live_key(files.root, files.node_id)
if not transcribed and (live := run_activity.live(live_key)) is not None:
stored.extend(
normalize_row(entry) for entry in list(live.transcript) if isinstance(entry, dict) and entry.get("role")
)
# The console tail, for the lane whose only in-flight account is its own
# output (a cli agent streams no transcript). Gone once the run finishes:
# the live index empties with the collecting block, and the record's answer
# takes over.
if (live_run := run_activity.live(files.node_id)) is not None and live_run.console:
if (live_run := run_activity.live(live_key)) is not None and live_run.console:
stored.append({"role": "console", "content": live_run.console})
if answer is not None:
answer_msg: dict[str, Any] = {"role": "assistant", "content": answer}
Expand Down
9 changes: 4 additions & 5 deletions raven/rpc/methods/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
from raven.agent.subagent import activity as run_activity
from raven.agent.subagent.activity import merge_file_change
from raven.agent.subagent.dag_store import REGISTRY_FILENAME, RUNNING, UNRECORDED, node_live_key
from raven.agent.subagent.history import dag_root, nodes_root, session_history_root
from raven.agent.subagent.history import dag_root, nodes_root, session_history_root, spawn_live_key
from raven.agent.subagent.instances import get_registry
from raven.rpc.methods.instances import _graph_of
from raven.rpc.methods.session import _safe_invoke_factory
Expand Down Expand Up @@ -124,9 +124,8 @@ def _overlay_live(node: dict[str, Any], live: Any) -> None:
nothing for the whole of its run. Only what the lane has said so far is
taken: a lane that has not spoken keeps its null.
"""
# A node that is not running has an account of its own on disk (or none),
# and the live index is keyed by a record id that is unique per conversation
# only: another conversation's run under the same id must not fill it.
# A node that is not running has an account of its own on disk (or none);
# an entry still in the index for it is a run this reader is not describing.
if live is None or node["status"] != "running":
return
for key in ("tokens_in", "tokens_out"):
Expand Down Expand Up @@ -288,7 +287,7 @@ def _spawn_row(files: "_NodeFiles", agent_loop_factory: "AgentLoopFactory | None
"prompt_template": None,
"files": _files_of(meta),
}
_overlay_live(node, run_activity.live(files.node_id))
_overlay_live(node, run_activity.live(spawn_live_key(files.root, files.node_id)))
task_summary = meta.get("task_summary") or meta.get("label") or _label_from_prompt(files) or None
return {
"id": files.node_id,
Expand Down
50 changes: 50 additions & 0 deletions tests/test_rpc_subagent_calls.py
Original file line number Diff line number Diff line change
Expand Up @@ -811,3 +811,53 @@ async def test_a_graph_node_s_own_order_is_pinned_not_incidental(workspace: Path
items = (await subagent_list({"session_id": SESSION}))["items"]

assert [i["node"] for i in items] == ["write", "survey"], "the id is the tie-break, descending like the stamp"


async def test_two_conversations_that_named_a_spawn_alike_each_read_their_own_live_transcript(workspace: Path) -> None:
"""A node id is unique for one conversation only and the live index is one
per process. Keyed by the record's address, two conversations running a
spawn under the same id each watch their own run, and the first to finish
takes nothing of the other's with it."""
from raven.agent.subagent import activity

other = "tui:other"

class _Holds:
"""Says which conversation it works for, then waits to be released."""

def __init__(self) -> None:
self.started = {key: asyncio.Event() for key in (SESSION, other)}
self.release = {key: asyncio.Event() for key in (SESSION, other)}

async def run(self, task: str, **kw: Any) -> str:
key = str(kw.get("session_key") or "")
activity.note_transcript([{"role": "assistant", "content": f"{key} working"}])
self.started[key].set()
await self.release[key].wait()
return f"{key} done"

hold = _Holds()
mgr = _manager(workspace)
mgr.registry.set_builtin_builder(lambda _row, _build: hold)
await mgr.spawn("the same first step", task_summary="step", session_key=SESSION, node_id="step1")
await mgr.spawn("the same first step", task_summary="step", session_key=other, node_id="step1")
await asyncio.wait_for(asyncio.gather(hold.started[SESSION].wait(), hold.started[other].wait()), 5)

said_a = json.dumps(await subagent_context({"id": "step1", "session_id": SESSION}))
said_b = json.dumps(await subagent_context({"id": "step1", "session_id": other}))
assert f"{SESSION} working" in said_a and f"{other} working" not in said_a
assert f"{other} working" in said_b and f"{SESSION} working" not in said_b

# The other conversation's run finishes first.
hold.release[other].set()
for _ in range(200):
if sum(1 for t in mgr._running_tasks.values() if not t.done()) <= 1:
break
await asyncio.sleep(0.01)
else:
pytest.fail("the other conversation's run never finished, so its exit was never tested")
still_a = json.dumps(await subagent_context({"id": "step1", "session_id": SESSION}))
assert f"{SESSION} working" in still_a, "its exit dropped only its own entry"

hold.release[SESSION].set()
await asyncio.gather(*mgr._running_tasks.values(), return_exceptions=True)
51 changes: 44 additions & 7 deletions tests/test_rpc_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@

import raven.home as raven_home_module
from raven.agent.subagent.dag_store import DagRunStore, index_guard
from raven.agent.subagent.history import SpawnRecord, dag_root, nodes_root, session_history_root
from raven.agent.subagent.history import SpawnRecord, dag_root, nodes_root, session_history_root, spawn_live_key
from raven.agent.subagent.prompt_backend import LocalFileBackend
from raven.rpc.methods import tasks as tasks_mod
from raven.rpc.methods.tasks import tasks_list
Expand Down Expand Up @@ -947,7 +947,7 @@ async def test_a_running_spawn_reads_usage_and_counts_off_the_live_activity(work
)
loop = _loop_stub(live_spawn_handles=frozenset({("Raven", "counting")}))

with activity_mod.collecting(live_key="counting"):
with activity_mod.collecting(live_key=spawn_live_key(nodes_root(session_dir), "counting")):
activity_mod.note_usage({"prompt_tokens": 40, "completion_tokens": 2})
activity_mod.note_usage({"prompt_tokens": 10, "completion_tokens": 3})
activity_mod.note_tool_call("exec")
Expand Down Expand Up @@ -981,7 +981,7 @@ async def test_a_lane_that_has_not_spoken_keeps_its_nulls_while_live(workspace:
)
loop = _loop_stub(live_spawn_handles=frozenset({("Raven", "quiet")}))

with activity_mod.collecting(live_key="quiet"):
with activity_mod.collecting(live_key=spawn_live_key(nodes_root(session_dir), "quiet")):
node = (await tasks_list({"session_key": SESSION}, agent_loop_factory=_factory(loop)))["tasks"][0]["nodes"][0]

assert node["tokens_in"] is None and node["tokens_out"] is None
Expand Down Expand Up @@ -1054,9 +1054,9 @@ async def test_a_dag_node_that_finished_mid_run_reads_the_account_the_runner_set


async def test_a_settled_spawn_ignores_a_live_activity_under_its_key(workspace: Path) -> None:
"""The live index is keyed by a record id that is unique per conversation
only: a finished spawn must not take the numbers of another run collected
under the same id."""
"""A settled row's account is the record's, and an entry still in the live
index under this record's address belongs to a run this row is not
describing -- so a finished spawn keeps its on-disk nulls."""
from raven.agent.subagent import activity as activity_mod

session_dir = _session_dir(workspace)
Expand All @@ -1070,10 +1070,47 @@ async def test_a_settled_spawn_ignores_a_live_activity_under_its_key(workspace:
)
record.finish(status="completed", output="done")

with activity_mod.collecting(live_key="shared"):
with activity_mod.collecting(live_key=spawn_live_key(nodes_root(session_dir), "shared")):
activity_mod.note_usage({"prompt_tokens": 40, "completion_tokens": 2})
activity_mod.note_tool_call("exec")
node = (await tasks_list({"session_key": SESSION}))["tasks"][0]["nodes"][0]

assert node["status"] == "completed"
assert node["tokens_in"] is None and node["tool_call_count"] is None and node["files"] == []


async def test_two_conversations_that_named_a_spawn_alike_each_read_their_own_live_run(workspace: Path) -> None:
"""A node id is unique for one conversation only and the live index is one
per process: keyed by the record's address, two conversations running a
spawn under the same id read their own account, and the first to finish
takes nothing of the other's with it."""
from raven.agent.subagent import activity as activity_mod
from raven.session.manager import SessionManager

other = "tui:other"
dir_a = _session_dir(workspace)
dir_b = SessionManager(workspace).session_dir(other)
for session_dir in (dir_a, dir_b):
await _claim_spawn_node(session_dir, "step1")
SpawnRecord.open(
session_dir,
task_id="step1",
task="the same first step",
meta=_spawn_meta(agent="Raven", handle="step1", task_summary="Step one"),
node_id="step1",
)
loop = _loop_stub(live_spawn_handles=frozenset({("Raven", "step1")}))

with activity_mod.collecting(live_key=spawn_live_key(nodes_root(dir_a), "step1")):
activity_mod.note_usage({"prompt_tokens": 10, "completion_tokens": 1})
with activity_mod.collecting(live_key=spawn_live_key(nodes_root(dir_b), "step1")):
activity_mod.note_usage({"prompt_tokens": 200, "completion_tokens": 2})
a = (await tasks_list({"session_key": SESSION}, agent_loop_factory=_factory(loop)))["tasks"][0]["nodes"][0]
b = (await tasks_list({"session_key": other}, agent_loop_factory=_factory(loop)))["tasks"][0]["nodes"][0]
assert (a["tokens_in"], a["tokens_out"]) == (10, 1)
assert (b["tokens_in"], b["tokens_out"]) == (200, 2)
# B finished first: A's entry is still its own.
a_after = (await tasks_list({"session_key": SESSION}, agent_loop_factory=_factory(loop)))["tasks"][0]["nodes"][
0
]
assert (a_after["tokens_in"], a_after["tokens_out"]) == (10, 1)
Loading