Skip to content

Commit 5fb8b69

Browse files
mdrxyopen-swe
andauthored
fix(code): serialize transcript tail reconciliation (#6143)
Long transcripts no longer duplicate rows when new output arrives during history hydration. --- The bounded tail jump introduced by #6057 could overlap with scroll-triggered hydration. Both paths built widgets from the same stale visible range, so the second mount hit duplicate DOM IDs and could drop fresh output or desynchronize the transcript store. Serialize transcript store/DOM mutations across append, hydration, pruning, and clear operations. The tail jump now derives mounted IDs from the actual container and releases removed tool-group summaries before regrouping surviving rows. Made by [Open SWE](https://openswe.vercel.app/agents/708f22e9-c9ed-554d-858f-1c2090a9482b) Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
1 parent a1a0887 commit 5fb8b69

2 files changed

Lines changed: 145 additions & 10 deletions

File tree

‎libs/code/deepagents_code/app.py‎

Lines changed: 64 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -4302,6 +4302,9 @@ def __init__(
43024302
self._transcript_generation = 0
43034303
"""Invalidates hydration/pruning work when the transcript is cleared."""
43044304

4305+
self._transcript_mutation_lock = asyncio.Lock()
4306+
"""Serializes message-store and transcript DOM reconciliation."""
4307+
43054308
self._hydration_requests: set[Literal["above", "below"]] = set()
43064309
"""Coalesced transcript hydration directions awaiting one UI slice."""
43074310

@@ -9390,9 +9393,19 @@ async def _hydrate_messages(
93909393
) -> int:
93919394
"""Hydrate one contiguous batch at a mounted-window edge.
93929395

9393-
Args:
9394-
direction: Edge receiving stored messages.
9395-
count: Maximum messages to mount; defaults to `HYDRATE_BUFFER`.
9396+
Returns:
9397+
Number of messages mounted.
9398+
"""
9399+
async with self._transcript_mutation_lock:
9400+
return await self._hydrate_messages_unlocked(direction, count=count)
9401+
9402+
async def _hydrate_messages_unlocked(
9403+
self,
9404+
direction: Literal["above", "below"],
9405+
*,
9406+
count: int | None = None,
9407+
) -> int:
9408+
"""Hydrate one batch while transcript mutation is already serialized.
93969409

93979410
Returns:
93989411
Number of messages mounted.
@@ -19984,6 +19997,22 @@ async def _mount_message(
1998419997
| ReasoningMessage
1998519998
| ToolCallMessage
1998619999
| SkillMessage,
20000+
) -> bool:
20001+
"""Mount one message while transcript mutation is serialized.
20002+
20003+
Returns:
20004+
Whether the widget reached the screen.
20005+
"""
20006+
async with self._transcript_mutation_lock:
20007+
return await self._mount_message_unlocked(widget)
20008+
20009+
async def _mount_message_unlocked(
20010+
self,
20011+
widget: Static
20012+
| AssistantMessage
20013+
| ReasoningMessage
20014+
| ToolCallMessage
20015+
| SkillMessage,
1998720016
) -> bool:
1998820017
"""Mount a message widget to the messages area.
1998920018

@@ -20118,13 +20147,18 @@ async def _move_transcript_window_to_tail(
2011820147
if tail is None:
2011920148
while self._message_store.has_messages_below:
2012020149
before = self._message_store.get_visible_range()[1]
20121-
await self._hydrate_messages("below")
20150+
await self._hydrate_messages_unlocked("below")
2012220151
if self._message_store.get_visible_range()[1] == before:
2012320152
return False
2012420153
return True
2012520154

2012620155
generation = self._transcript_generation
20127-
mounted_ids = {data.id for data in self._message_store.get_visible_messages()}
20156+
mounted_ids = {
20157+
child.id
20158+
for child in messages_container.children
20159+
if child.id is not None
20160+
and self._message_store.get_message(child.id) is not None
20161+
}
2012820162
tail_ids = {data.id for data in tail}
2012920163
entries = [
2013020164
self._build_hydration_entry(data)
@@ -20163,6 +20197,9 @@ async def _move_transcript_window_to_tail(
2016320197
]
2016420198
await messages_container.remove_children(hydrated_nodes)
2016520199
return False
20200+
for child in obsolete:
20201+
if isinstance(child, ToolGroupSummary):
20202+
child._release_all_collapsible()
2016620203
await messages_container.remove_children(obsolete)
2016720204
self._schedule_message_height_measurements([data.id for data in moved])
2016820205
self._sync_transcript_spacers(messages_container)
@@ -20176,12 +20213,24 @@ async def _prune_messages(
2017620213
*,
2017720214
count: int | None = None,
2017820215
) -> int:
20179-
"""Prune a bounded batch from one mounted-window edge.
20216+
"""Prune one mounted-window edge while mutation is serialized.
2018020217

20181-
Args:
20182-
direction: Edge to remove messages from.
20183-
messages_container: Cached transcript container, when available.
20184-
count: Maximum messages to remove; defaults to the soft-limit excess.
20218+
Returns:
20219+
Number of messages removed.
20220+
"""
20221+
async with self._transcript_mutation_lock:
20222+
return await self._prune_messages_unlocked(
20223+
direction, messages_container, count=count
20224+
)
20225+
20226+
async def _prune_messages_unlocked(
20227+
self,
20228+
direction: Literal["above", "below"],
20229+
messages_container: Container | None = None,
20230+
*,
20231+
count: int | None = None,
20232+
) -> int:
20233+
"""Prune a bounded batch while transcript mutation is serialized.
2018520234

2018620235
Returns:
2018720236
Number of messages removed.
@@ -20535,6 +20584,11 @@ def on_user_message_expansion_changed(
2053520584
self._schedule_message_height_measurements([event.widget.id])
2053620585

2053720586
async def _clear_messages(self) -> None:
20587+
"""Clear the transcript while mutation is serialized."""
20588+
async with self._transcript_mutation_lock:
20589+
await self._clear_messages_unlocked()
20590+
20591+
async def _clear_messages_unlocked(self) -> None:
2053820592
"""Clear the messages area and message store."""
2053920593
self._transcript_generation += 1
2054020594
# Drop buffered `!` shell output so it never leaks across a thread

‎libs/code/tests/unit_tests/test_app.py‎

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11271,6 +11271,87 @@ async def count_hydration(
1127111271
with pytest.raises(NoMatches):
1127211272
app.query_one(f"#{message_id}", UserMessage)
1127311273

11274+
async def test_mount_message_waits_for_inflight_hydration(
11275+
self, monkeypatch: pytest.MonkeyPatch
11276+
) -> None:
11277+
"""Concurrent hydration and append keep one mounted row per store entry."""
11278+
from deepagents_code.tui.widgets.message_store import MessageData, MessageType
11279+
11280+
app = DeepAgentsApp()
11281+
monkeypatch.setattr(app._message_store, "WINDOW_SIZE", 3)
11282+
monkeypatch.setattr(app, "_schedule_transcript_prune", lambda *_args: None)
11283+
11284+
async with app.run_test() as pilot:
11285+
await pilot.pause()
11286+
history = [
11287+
MessageData(
11288+
type=MessageType.USER,
11289+
content=f"m{index}",
11290+
id=f"race-{index}",
11291+
)
11292+
for index in range(6)
11293+
]
11294+
app._message_store.bulk_load(history)
11295+
messages = app.query_one("#messages", Container)
11296+
tail_entries = [
11297+
app._build_hydration_entry(data)
11298+
for data in app._message_store.get_visible_messages()
11299+
]
11300+
assert await app._mount_hydration_batch(
11301+
messages,
11302+
tail_entries,
11303+
generation=app._transcript_generation,
11304+
)
11305+
await messages.remove_children(
11306+
[
11307+
child
11308+
for child in messages.children
11309+
if child.id
11310+
and any(
11311+
child.id.startswith(f"race-{index}") for index in range(3, 6)
11312+
)
11313+
]
11314+
)
11315+
app._message_store._visible_start = 0
11316+
app._message_store._visible_end = 2
11317+
head_entries = [
11318+
app._build_hydration_entry(data)
11319+
for data in app._message_store.get_visible_messages()
11320+
]
11321+
assert await app._mount_hydration_batch(
11322+
messages,
11323+
head_entries,
11324+
generation=app._transcript_generation,
11325+
)
11326+
11327+
hydration_started = asyncio.Event()
11328+
release_hydration = asyncio.Event()
11329+
mount_batch = app._mount_hydration_batch
11330+
11331+
async def pause_first_batch(*args: Any, **kwargs: Any) -> bool:
11332+
if not hydration_started.is_set():
11333+
hydration_started.set()
11334+
await release_hydration.wait()
11335+
return await mount_batch(*args, **kwargs)
11336+
11337+
monkeypatch.setattr(app, "_mount_hydration_batch", pause_first_batch)
11338+
hydrate = asyncio.create_task(app._hydrate_messages("below", count=3))
11339+
await hydration_started.wait()
11340+
append = asyncio.create_task(
11341+
app._mount_message(UserMessage("new", id="race-new"))
11342+
)
11343+
await asyncio.sleep(0)
11344+
assert not append.done()
11345+
11346+
release_hydration.set()
11347+
assert await hydrate == 3
11348+
assert await append is True
11349+
await pilot.pause()
11350+
11351+
assert app._message_store.get_visible_range() == (3, 7)
11352+
for message_id in ["race-3", "race-4", "race-5", "race-new"]:
11353+
assert len(app.query(f"#{message_id}")) == 1
11354+
1127411355
async def test_mount_message_hydrates_tail_blocked_by_protected_row(
1127511356
self, monkeypatch: pytest.MonkeyPatch
1127611357
) -> None:

0 commit comments

Comments
 (0)