unbridled-41 opened a new pull request, #4948:
URL: https://github.com/apache/rocketmq-dashboard/pull/4948
## Problem
A client that reconnects to a running agent run can lose the middle of its
transcript permanently.
`AiRunService.attach` replays the persisted events after the client's
cursor, but it reads **one page** of them and then raises the observer
watermark to the last row of that page. Rows past the first page are then
neither replayed (they sit behind the client's cursor, so the tail never sends
them) nor buffered (they were published before this observer existed), and
`AgentStreamSession.deliver` drops anything at or below the watermark anyway.
## Evidence
Base commit `1ef5d860`.
* `AiRunService.java:263-264` (base) — the whole replay is one read:
```java
List<RmqAiEvent> rows =
eventRepository.findByConversationIdAfterSeq(run.getConversationId(),
Math.max(0, afterSeq),
AiConversationService.DEFAULT_TIMELINE_LIMIT);
```
* `AiConversationService.java:105` (base) — `DEFAULT_TIMELINE_LIMIT = 200`;
`MybatisPlusAiEventRepository.java:59-63` caps the query at `min(limit, 500)`.
There is no loop, and `rows.size()` only feeds the log line at `:288`.
* `AgentStreamSession.java:152` (base) — `noteWatermark(lastRow.getSeq())`
is the dedup high-water mark; `:208` (base) drops live frames with `seq <=
maxSeqSeen`.
* `AgentRunRegistry.java:121-134` (base) — `attach` adds the observer to a
`CopyOnWriteArrayList`. The registry keeps **no backlog**, so a frame published
before the registration is fanned out to nobody.
* `AiRunService.java:275` (base) — the registration happens *after* the
replay loop, which contradicts the class javadoc at `:253-255` ("The observer
is registered before the replay is flushed ... nothing is lost in the gap") and
`docs/api-spec.md` §15.8.
Reachability: the web client passes its real cursor —
`useActiveRunAttach.ts:58` → `attach(conversationId, activeRun.id, lastSeq)` →
`web/src/api/ai.ts:407` sends `after=<lastSeq>`. A long tool-heavy run writes
two rows per tool call plus one per coalesced text/thinking block
(`AiEventSink`: `setType`/`setSeq` at `:392-395`, coalescing owned by
`AgentEventProjector.COALESCE_MAX_CHARS`), so a client that slept or lost its
connection for a few minutes comes back to more than 200 rows.
## Root cause
Two defects in the same method:
1. The replay reuses the *timeline page* constant as if it were a replay
bound. A page is a display unit; the replay's job is to close the gap between
the client's cursor and the live head, which is not a fixed size.
2. The observer is registered after the replay read instead of before it, so
a frame published during the read window reaches no observer at all and is then
too old for the tail to ever send.
## Fix
* `attach` now registers the observer **first** and detaches it if the
replay read throws, so a failed response cannot strand a half-replayed observer
in the registry.
* The replay moved into `replayInto`, which drains the backlog page by page
until a short page proves it is empty, advancing the cursor monotonically (also
across rows whose payload does not decode).
* Page size stays `DEFAULT_TIMELINE_LIMIT`, matching what the client itself
pages the timeline at.
## Scoring (AGENTS.md)
* PRIORITY **62** = impact 16 (a reconnecting client shows a permanently
truncated answer until the run finishes and the timeline is refetched) + reach
12 (the reconnect path of the AI console) + reproducibility 18 (deterministic
in a unit test; the trigger is a >200-row backlog) + maintenance value 16 (the
method now matches the invariant its own javadoc states).
* FIX_CONFIDENCE **85** — no interface change, the failure modes are the two
covered by the new tests.
## Tests
New tests in
`server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunServiceTest.java`:
1. `attachShouldReplayABacklogLongerThanOneTimelinePageTest` — a 200-row
first page plus a 201st row.
2. `attachShouldNotLoseAFramePublishedWhileTheReplayIsReadTest` — the run
publishes a frame from inside the replay read.
3. `attachShouldDetachTheObserverWhenTheReplayReadFailsTest` — the new
registration path must not leak an observer.
Pre-fix, with the test kept and only `AiRunService.java` restored to
`1ef5d860`:
```
[ERROR]
AiRunServiceTest.attachShouldReplayABacklogLongerThanOneTimelinePageTest -- <<<
FAILURE!
Expecting actual:
<200 replayed frames>
to contain:
"the newest block"
[ERROR]
AiRunServiceTest.attachShouldNotLoseAFramePublishedWhileTheReplayIsReadTest --
<<< FAILURE!
Expecting actual:
"event:agent"
to contain:
"published mid-replay"
```
(The third test cannot fail on the old code — it guards behaviour the fix
introduces. Noted so the numbers are not overread.)
Post-fix, from the branch:
```
mvn -o test
-Dtest='AiRunServiceTest,AiRunExecutorTest,AiConversationServiceTest,AiConversationVoAssemblerTest,AgentStreamSessionTest,AiEventSinkTest,AgentRunRegistryTest,AiEventCodecTest'
AiRunServiceTest 23/23, AiRunExecutorTest 18/18, AiConversationServiceTest
18/18,
AiConversationVoAssemblerTest 21/21, AgentRunRegistryTest 11/11,
AiEventProjectorTest 40/40,
AiEventContractTest 165/165, AiEventJsonSerializationTest 116/116,
AiTimelineRepositoryTest 35/35 ...
```
Full server suite on the branch: `mvn -o test` → **3153 tests, 0 failures,
17 errors**. The 17 errors are the environment baseline of this checkout, not a
regression: every report that fails contains `Failed to load
ApplicationContext` caused by `Communications link failure` (there is no MySQL
here).
## Risk
Low. The drain is bounded by the rows that actually exist after the caller's
cursor, each query is index-ranged and capped at 200 rows, and
`AgentStreamSession` already bounds what is buffered while a replay is in
flight (`MAX_BUFFERED_DURING_REPLAY`, `:98`), with the documented fallback that
the client re-reads the persisted timeline when the run finishes. The one
ordering change means a run that finishes *during* the replay may now have its
terminal frame buffered and dropped by `complete()`; the client still receives
`done` and refetches the timeline, and the terminal row is in the database, so
the transcript converges.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]