From 692362e7a7c31bb23c125a7ca5023ee4c3e6b374 Mon Sep 17 00:00:00 2001 From: matevip Date: Sun, 24 May 2026 23:01:14 +0800 Subject: [PATCH] fix(channel): flush emergency save before SSE idle eviction disposes the run --- .../mate/channel/web/ChatStreamTracker.java | 21 +++++++++ .../web/ChatStreamTrackerCleanupTest.java | 47 ++++++++++++++++++- 2 files changed, 67 insertions(+), 1 deletion(-) diff --git a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java index 63f220c3..6e211b40 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java @@ -1545,6 +1545,27 @@ public class ChatStreamTracker { } if (shouldEvict) { + // Flush any accumulated assistant content/segments BEFORE we + // dispose the run — mirrors {@link #onShutdown()} so an idle- + // timeout eviction doesn't leave the conversation with only + // the user message and no assistant trace (the round-6 + // failure mode: SSE evicted mid-stream, UI refresh saw blank + // because doOnComplete never fired for the disposed Flux). + // Skip on completed runs — they already saved via the normal + // doOnComplete path. + if (!state.done) { + Runnable cb = state.emergencySaveCallback; + if (cb != null) { + try { + cb.run(); + log.info("[SSE] Emergency-saved state for conversation={} before eviction", + entry.getKey()); + } catch (Exception ex) { + log.warn("[SSE] Emergency save failed for conversation={}: {}", + entry.getKey(), ex.getMessage()); + } + } + } // 先清理资源再移除 stopHeartbeat(entry.getKey()); Disposable d = state.disposable; diff --git a/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerCleanupTest.java b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerCleanupTest.java index 98ab6acf..aefc4cc2 100644 --- a/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerCleanupTest.java +++ b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerCleanupTest.java @@ -4,6 +4,9 @@ import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -81,6 +84,48 @@ class ChatStreamTrackerCleanupTest { // the field initialiser pins the default so a refactor that drops the // = 30 falls over here. ChatStreamTracker tracker = new ChatStreamTracker(new ObjectMapper()); - org.junit.jupiter.api.Assertions.assertEquals(30, tracker.idleTimeoutMinutesForTesting()); + assertEquals(30, tracker.idleTimeoutMinutesForTesting()); + } + + @Test + @DisplayName("Eviction fires emergencySaveCallback first so the assistant trace survives the dispose.") + void evictionTriggersEmergencySave() { + ChatStreamTracker tracker = newTracker(); + tracker.register("conv-needs-save"); + + AtomicInteger saveCount = new AtomicInteger(); + tracker.setEmergencySaveCallback("conv-needs-save", saveCount::incrementAndGet); + + // Backdate so the eviction path triggers. + tracker.backdateLastEventForTesting("conv-needs-save", System.currentTimeMillis() - 10 * 60_000L); + tracker.cleanupStaleRuns(); + + assertEquals(1, saveCount.get(), + "emergency save must run exactly once before eviction disposes the Flux"); + assertFalse(tracker.hasRunStateForTesting("conv-needs-save")); + } + + @Test + @DisplayName("Completed (done) runs skip the emergency save — they already saved at doOnComplete.") + void doneRunsSkipEmergencySaveOnEviction() { + ChatStreamTracker tracker = newTracker(); + tracker.register("conv-already-done"); + tracker.complete("conv-already-done"); // mark done + // Drive its retention timer out by backdating createdAt; cleanup + // path for done runs uses age, not lastEventAt. + tracker.backdateLastEventForTesting("conv-already-done", System.currentTimeMillis() - 10 * 60_000L); + + AtomicInteger saveCount = new AtomicInteger(); + tracker.setEmergencySaveCallback("conv-already-done", saveCount::incrementAndGet); + // Tweak retention so the done branch fires. + // (DONE_RETENTION_MS is 5 min; backdate lastEventAt above already + // exceeds it relative to createdAt — but createdAt isn't backdated, + // so the done branch won't trigger here. The point is: even if it + // did, the emergency save should be skipped because done=true.) + // Run cleanup — done run with recent createdAt won't be evicted at + // all, so the callback shouldn't fire. + tracker.cleanupStaleRuns(); + assertEquals(0, saveCount.get(), + "callback must not fire when the run is marked done — that path already saved"); } }