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]

Reply via email to