From 80dcd0c4163f75863cd82fd221295c938eaa6536 Mon Sep 17 00:00:00 2001 From: Jakob Stender Guldberg Date: Sun, 6 Sep 2026 07:45:26 +0200 Subject: [PATCH 1/3] fix(web): survive page reload without losing analysis or running LLM jobs The web UI kept repo, branch pair, analysis, and job ids only in React state, so a reload dropped everything while the server kept the jobs running. Persist the session and the active job list in sessionStorage (web only, per tab), re-run the analysis on load, and re-attach to stored job streams after probing that the server still knows them. A job is forgotten only on a terminal event; Chromium fires the EventSource error during page teardown, so cleaning up on error wiped the list before the new page loaded. Jobs are scoped to the repo they were started for, each stream closes only the EventSource it owns, and a stored session is ignored when the server's --repo changed. Refinement results read the current analysis through a ref instead of a stale closure, which also fixes cached refinements being dropped on the first analysis of a page; a live result wins over a cached one. --- crates/diffcore-tauri/ui/src/App.tsx | 206 +++++++++++++++--- .../ui/tests/e2e/session-restore.spec.ts | 47 ++++ 2 files changed, 221 insertions(+), 32 deletions(-) create mode 100644 crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts diff --git a/crates/diffcore-tauri/ui/src/App.tsx b/crates/diffcore-tauri/ui/src/App.tsx index 1c197bd..f42f5dd 100644 --- a/crates/diffcore-tauri/ui/src/App.tsx +++ b/crates/diffcore-tauri/ui/src/App.tsx @@ -92,6 +92,55 @@ function isPrUrl(value: string): boolean { return /^https?:\/\//i.test(value.trim()); } +const SESSION_KEY = "diffcore.session"; +const JOBS_KEY = "diffcore.jobs"; +type StoredSession = { + repoPath: string; + baseRef: string; + headRef: string | null; + defaultRepo: string | null; + analyzedRepo: string | null; + refsPinned: boolean; +}; +type StoredJob = AsyncLlmJobStart & { repoPath: string; baseRef: string; headRef: string | null }; + +// Desktop never reloads and its SSE URLs carry a per-launch port, so it opts out. +function loadStored(key: string): T | null { + if (IS_TAURI) return null; + try { + const raw = sessionStorage.getItem(key); + return raw ? (JSON.parse(raw) as T) : null; + } catch { + return null; + } +} + +function saveStored(key: string, value: unknown): void { + if (IS_TAURI) return; + try { + sessionStorage.setItem(key, JSON.stringify(value)); + } catch { + // Non-fatal: the session just won't survive a reload + } +} + +function loadJobs(): StoredJob[] { + return loadStored(JOBS_KEY) ?? []; +} + +function forgetJob(jobId: string): void { + saveStored(JOBS_KEY, loadJobs().filter((job) => job.job_id !== jobId)); +} + +async function jobStreamAlive(url: string): Promise { + try { + const res = await fetch(url, { signal: AbortSignal.timeout(5000) }); + void res.body?.cancel(); + return res.ok; + } catch { + return false; + } +} function TruncatedText({ text, @@ -239,15 +288,22 @@ export default function App() { const [activityError, setActivityError] = useState(null); const [activityViewMode, setActivityViewMode] = useState("stream"); const [inspectedActivityId, setInspectedActivityId] = useState(null); - const activitySourceRef = useRef(null); + const activitySources = useRef(new Set()); const activityLogRef = useRef(null); const [rightPanelTab, setRightPanelTab] = useState("annotations"); const [sourceFocusRequest, setSourceFocusRequest] = useState(null); // Repo and git state - const [repoPath, setRepoPath] = useState(HAS_BACKEND ? (DEFAULT_REPO ?? "") : "/demo/repo"); - const [baseRef, setBaseRef] = useState("main"); - const [headRef, setHeadRef] = useState(null); + const [session] = useState(() => loadStored(SESSION_KEY)); + const restored = + typeof session?.repoPath === "string" + && !!session.repoPath + && !isPrUrl(session.repoPath) + && session.defaultRepo === DEFAULT_REPO; + const initialRepo = HAS_BACKEND ? (DEFAULT_REPO ?? "") : "/demo/repo"; + const [repoPath, setRepoPath] = useState(restored ? session.repoPath : initialRepo); + const [baseRef, setBaseRef] = useState(restored ? session.baseRef : "main"); + const [headRef, setHeadRef] = useState(restored ? session.headRef : null); const [repoInfo, setRepoInfo] = useState(null); // Set after a PR/MR URL resolves; see the effect below runAnalysis. const [pendingAnalysis, setPendingAnalysis] = useState(false); @@ -259,6 +315,20 @@ export default function App() { * detached, so letting loadRepoInfo auto-detect would replace the PR's fork * point and tip with the checkout's default branch and a bare HEAD. */ const prRefs = useRef(false); + // Refs the user picked win over auto-detect for the restored repo only; + // auto-detected refs are re-detected so a reload still sees a new checkout. + const refsPinned = useRef(false); + const restoredRepo = useRef(restored && session.refsPinned ? session.repoPath : null); + useEffect(() => { + saveStored(SESSION_KEY, { + repoPath, + baseRef, + headRef, + defaultRepo: DEFAULT_REPO, + analyzedRepo: analysis ? lastBackendArgs.current.analyze?.repoPath ?? null : null, + refsPinned: refsPinned.current, + }); + }, [repoPath, baseRef, headRef, analysis]); const [branchDropdownOpen, setBranchDropdownOpen] = useState(false); const [headBranchDropdownOpen, setHeadBranchDropdownOpen] = useState(false); @@ -511,7 +581,10 @@ export default function App() { info = MOCK_REPO_INFO; } setRepoInfo(info); - if (!prRefs.current) { + const keepRefs = prRefs.current || restoredRepo.current === path; + if (restoredRepo.current !== path) restoredRepo.current = null; + if (!keepRefs) { + refsPinned.current = false; // Auto-set base ref to the detected default branch setBaseRef(info.default_branch); // Auto-set head ref to the current branch (what we're comparing FROM) @@ -568,10 +641,8 @@ export default function App() { }, [llmSettings]); const closeActivityStream = useCallback(() => { - if (activitySourceRef.current) { - activitySourceRef.current.close(); - activitySourceRef.current = null; - } + for (const source of activitySources.current) source.close(); + activitySources.current.clear(); }, []); useEffect(() => () => closeActivityStream(), [closeActivityStream]); @@ -609,21 +680,12 @@ export default function App() { [appendActivityEntry, closeActivityStream], ); - const runStreamingJob = useCallback( - async ( - command: string, - args: Record, - onComplete: (value: T) => void, - onJobId?: (jobId: string) => void, - ) => { - closeActivityStream(); + const attachJobStream = useCallback( + async (start: AsyncLlmJobStart, onComplete: (value: T) => void) => { setActivityViewMode("stream"); setInspectedActivityId(null); setActivityEntries([]); setActivityError(null); - - const start = await tauriInvoke(command, args); - onJobId?.(start.job_id); setActivityJob({ job_id: start.job_id, operation: start.operation, @@ -634,7 +696,11 @@ export default function App() { await new Promise((resolve, reject) => { const source = new EventSource(start.stream_url); - activitySourceRef.current = source; + activitySources.current.add(source); + const close = () => { + source.close(); + activitySources.current.delete(source); + }; source.addEventListener("job_started", (event) => { try { @@ -661,7 +727,8 @@ export default function App() { }); source.addEventListener("completed", (event) => { - closeActivityStream(); + forgetJob(start.job_id); + close(); try { const payload = JSON.parse((event as MessageEvent).data) as { result: T }; onComplete(payload.result); @@ -675,7 +742,8 @@ export default function App() { }); source.addEventListener("failed", (event) => { - closeActivityStream(); + forgetJob(start.job_id); + close(); try { const payload = JSON.parse((event as MessageEvent).data) as { error: string }; setActivityError(payload.error); @@ -696,14 +764,30 @@ export default function App() { }); source.onerror = () => { - closeActivityStream(); + close(); setActivityError("Activity stream disconnected"); setActivityJob(null); reject(new Error("Activity stream disconnected")); }; }); }, - [appendActivityEntry, closeActivityStream], + [appendActivityEntry], + ); + + const runStreamingJob = useCallback( + async ( + command: string, + args: Record, + onComplete: (value: T) => void, + onJobId?: (jobId: string) => void, + ) => { + closeActivityStream(); + const start = await tauriInvoke(command, args); + onJobId?.(start.job_id); + saveStored(JOBS_KEY, [...loadJobs(), { ...start, repoPath: repoPath.trim(), baseRef, headRef }]); + await attachJobStream(start, onComplete); + }, + [attachJobStream, closeActivityStream, repoPath, baseRef, headRef], ); const handleSelectFile = useCallback( @@ -969,6 +1053,7 @@ export default function App() { setRefinementModel(null); setRefinementHadChanges(null); setShowRefined(false); + refinementApplied.current = false; // Reset review tick-off state setReviewedGroupIds(new Set()); setDismissedEmptyGroupIds(new Set()); @@ -1028,7 +1113,7 @@ export default function App() { } finally { setLoading(false); } - }, [repoPath, baseRef, headRef, handleSelectGroup, closeActivityStream, describeGroups]); + }, [repoPath, baseRef, headRef, includeUncommitted, handleSelectGroup, closeActivityStream, describeGroups]); /** Analyze whatever is in the repository field — a local path, or a PR/MR URL * that we first clone and resolve to a base/head pair. */ @@ -1054,6 +1139,7 @@ export default function App() { return; } prRefs.current = true; + refsPinned.current = true; setRepoPath(resolved.path); setBaseRef(resolved.base); setHeadRef(resolved.head); @@ -1188,12 +1274,16 @@ export default function App() { } }, [selectedGroup, repoPath, baseRef, resolvedPrimaryModel, resolvedPrimaryProvider, runMockActivityJob, runStreamingJob]); + const refinementApplied = useRef(false); + const analysisRef = useRef(null); + analysisRef.current = analysis; const applyRefinementResult = useCallback((result: RefinementResult, opts?: { fromCache?: boolean }) => { - if (!analysis) return; + const current = analysisRef.current; + if (!current) return; + if (opts?.fromCache && refinementApplied.current) return; + refinementApplied.current = true; - if (!originalGroups) { - setOriginalGroups(analysis.groups); - } + setOriginalGroups((prev) => prev ?? current.groups); setRefinedGroups(result.refined_groups); setRefinementResponse(result.refinement_response); @@ -1238,7 +1328,7 @@ export default function App() { if (HAS_BACKEND && !opts?.fromCache) { tauriInvoke("store_refinement_cache", { result, repoPath: repoPath || null }).catch(() => {}); } - }, [analysis, originalGroups, handleSelectGroup, showToast, describeGroups, repoPath]); + }, [handleSelectGroup, showToast, describeGroups, repoPath]); /** Run LLM refinement pass on the current analysis groups. */ const runRefinement = useCallback(async () => { @@ -1298,6 +1388,56 @@ export default function App() { } }, []); + const reattachJobs = useCallback(async () => { + const jobs = loadJobs().slice(-1); + saveStored(JOBS_KEY, jobs); + const alive = await Promise.all(jobs.map((job) => jobStreamAlive(job.stream_url))); + jobs.forEach((job, i) => { + const stale = job.repoPath !== repoPath.trim() || job.baseRef !== baseRef || job.headRef !== headRef; + if (!alive[i] || stale) { + forgetJob(job.job_id); + return; + } + const isRefinement = job.operation === "refinement"; + if (isRefinement) { + refinementJobIdRef.current = job.job_id; + setRefining(true); + } else { + deepAnalyzingCount.current += 1; + setDeepAnalyzing(true); + } + attachJobStream(job, (result) => { + if (isRefinement) { + applyRefinementResult(result as RefinementResult); + } else { + const analysisResult = result as Pass2Response; + setDeepAnalyses((prev) => ({ ...prev, [analysisResult.group_id]: analysisResult })); + } + }) + .catch((e) => { + const message = String(e); + if (message.includes("Cancelled by user")) return; + setError(`${isRefinement ? "Refinement" : "Deep analysis"} failed: ${message}`); + }) + .finally(() => { + if (isRefinement) { + refinementJobIdRef.current = null; + setRefining(false); + } else { + deepAnalyzingCount.current = Math.max(0, deepAnalyzingCount.current - 1); + if (deepAnalyzingCount.current === 0) setDeepAnalyzing(false); + } + }); + }); + }, [attachJobStream, applyRefinementResult, repoPath, baseRef, headRef]); + + const restorePending = useRef(HAS_BACKEND && restored && session.analyzedRepo === session.repoPath.trim()); + useEffect(() => { + if (!restorePending.current || !llmSettings) return; + restorePending.current = false; + runAnalysis().then(reattachJobs); + }, [llmSettings, runAnalysis, reattachJobs]); + /** Toggle between original and refined groups. */ const toggleRefinedView = useCallback( (useRefined: boolean) => { @@ -1384,7 +1524,7 @@ export default function App() { setActivityEntries(entries); }, setError: (msg: string | null) => setError(msg), - clearAnalysis: () => { setAnalysis(null); setSelectedGroup(null); setSelectedFile(null); setFileDiff(null); setOverview(null); setDeepAnalyses({}); setOriginalGroups(null); setRefinedGroups(null); setRefinementResponse(null); setRefinementProvider(null); setRefinementModel(null); setRefinementHadChanges(null); setShowRefined(false); setReviewedGroupIds(new Set()); setComments([]); setCommentInput(null); setCommentText(""); setRightPanelTab("annotations"); setSourceFocusRequest(null); setActivityJob(null); setActivityEntries([]); setActivityError(null); setActivityViewMode("stream"); setInspectedActivityId(null); }, + clearAnalysis: () => { setAnalysis(null); setSelectedGroup(null); setSelectedFile(null); setFileDiff(null); setOverview(null); setDeepAnalyses({}); setOriginalGroups(null); setRefinedGroups(null); setRefinementResponse(null); setRefinementProvider(null); setRefinementModel(null); setRefinementHadChanges(null); setShowRefined(false); refinementApplied.current = false; setReviewedGroupIds(new Set()); setComments([]); setCommentInput(null); setCommentText(""); setRightPanelTab("annotations"); setSourceFocusRequest(null); setActivityJob(null); setActivityEntries([]); setActivityError(null); setActivityViewMode("stream"); setInspectedActivityId(null); }, openAiSetup: (step: OnboardingStep = "recommended") => openAiSetup(step), dismissAiSetup: () => dismissAiSetup(), getAiSetupState: () => ({ open: aiSetupOpen, step: aiSetupStep }), @@ -2527,11 +2667,13 @@ export default function App() { }, [handleSelectFile, handleSelectFileDebounced, handleSelectGroup, enterReplay, exitReplay, goToReplayStep, toggleGroupReviewed, copyFilePath, copyFlowPaths, openCommentInput, exportComments, closeTab]); const handleSelectBase = useCallback((branch: string) => { + refsPinned.current = true; setBaseRef(branch); setBranchDropdownOpen(false); }, []); const handleSelectHead = useCallback((branch: string) => { + refsPinned.current = true; setHeadBranchDropdownOpen(false); // If the picked branch is checked out in a different worktree, switch // the nav bar to that worktree's path. Otherwise the fallback diff diff --git a/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts b/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts new file mode 100644 index 0000000..1039d16 --- /dev/null +++ b/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts @@ -0,0 +1,47 @@ +import { test, expect, type Page } from "@playwright/test"; + +async function waitForDemoApp(page: Page) { + await expect(page.locator(".summary")).toBeVisible({ timeout: 10_000 }); + await expect(page.locator(".group-item.selected")).toBeVisible({ timeout: 10_000 }); +} + +test("repository field and branch selection survive a page reload", async ({ page }) => { + await page.goto("/"); + await waitForDemoApp(page); + + await page.locator(".repo-input").fill("/work/other-repo"); + await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("main"); + + await page.locator("[data-testid='base-branch-dropdown'] .branch-dropdown-trigger").click(); + await page.locator("[data-testid='base-branch-dropdown'] .branch-option-name", { hasText: "develop" }).click(); + await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("develop"); + await page.locator("[data-testid='head-branch-dropdown'] .branch-dropdown-trigger").click(); + await page.locator("[data-testid='head-branch-dropdown'] .branch-option-name", { hasText: "fix/login-bug" }).click(); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("fix/login-bug"); + + await page.reload(); + await waitForDemoApp(page); + + await expect(page.locator(".repo-input")).toHaveValue("/work/other-repo"); + await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("develop"); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("fix/login-bug"); + + await page.locator(".repo-input").fill("/work/third-repo"); + await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("main"); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("feature/user-auth"); +}); + +test("auto-detected branches are re-detected after a reload", async ({ page }) => { + await page.goto("/"); + await waitForDemoApp(page); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("feature/user-auth"); + + await page.evaluate(() => { + (window as { __TEST_API__: { setHeadRef: (ref: string) => void } }).__TEST_API__.setHeadRef("release/v2.0"); + }); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("release/v2.0"); + + await page.reload(); + await waitForDemoApp(page); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("feature/user-auth"); +}); From 04e82e7b308094ac830d0c9b84b2c91362c8220f Mon Sep 17 00:00:00 2001 From: Jakob Stender Guldberg Date: Mon, 21 Sep 2026 13:40:07 +0200 Subject: [PATCH 2/3] fix(web): keep pinned refs pinned across repeated reloads refsPinned was seeded false on mount, so the effect that persists the session rewrote the stored pin to false on the first reload. Refs survived one reload via restoredRepo, then loadRepoInfo auto-detected over them on the next. Seed the ref from the restored session. --- crates/diffcore-tauri/ui/src/App.tsx | 2 +- .../ui/tests/e2e/session-restore.spec.ts | 15 +++++++++------ 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/crates/diffcore-tauri/ui/src/App.tsx b/crates/diffcore-tauri/ui/src/App.tsx index f42f5dd..383341b 100644 --- a/crates/diffcore-tauri/ui/src/App.tsx +++ b/crates/diffcore-tauri/ui/src/App.tsx @@ -317,7 +317,7 @@ export default function App() { const prRefs = useRef(false); // Refs the user picked win over auto-detect for the restored repo only; // auto-detected refs are re-detected so a reload still sees a new checkout. - const refsPinned = useRef(false); + const refsPinned = useRef(restored && !!session.refsPinned); const restoredRepo = useRef(restored && session.refsPinned ? session.repoPath : null); useEffect(() => { saveStored(SESSION_KEY, { diff --git a/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts b/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts index 1039d16..7037f4a 100644 --- a/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts +++ b/crates/diffcore-tauri/ui/tests/e2e/session-restore.spec.ts @@ -19,12 +19,15 @@ test("repository field and branch selection survive a page reload", async ({ pag await page.locator("[data-testid='head-branch-dropdown'] .branch-option-name", { hasText: "fix/login-bug" }).click(); await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("fix/login-bug"); - await page.reload(); - await waitForDemoApp(page); - - await expect(page.locator(".repo-input")).toHaveValue("/work/other-repo"); - await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("develop"); - await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("fix/login-bug"); + // Twice: the pin itself has to survive a reload, not just the refs it produced. + for (let i = 0; i < 2; i++) { + await page.reload(); + await waitForDemoApp(page); + + await expect(page.locator(".repo-input")).toHaveValue("/work/other-repo"); + await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("develop"); + await expect(page.locator("[data-testid='head-branch-dropdown'] .branch-name")).toHaveText("fix/login-bug"); + } await page.locator(".repo-input").fill("/work/third-repo"); await expect(page.locator("[data-testid='base-branch-dropdown'] .branch-name")).toHaveText("main"); From 7c09fae81c6ec785cea77102274b42768d3b3019 Mon Sep 17 00:00:00 2001 From: Jakob Stender Guldberg Date: Mon, 21 Sep 2026 14:31:49 +0200 Subject: [PATCH 3/3] fix(tauri): evict finished LLM jobs so the activity map stays bounded MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ActivityManager.jobs only ever grew: create_job inserted, nothing removed. Every job's history stayed resident for the life of the process, including Completed payloads carrying whole RefinementResults, and jobStreamAlive kept reporting long-dead jobs as live. Stamp finished_at on the terminal event and sweep finished jobs past JOB_RETENTION when a new job starts — no timer task, and if nothing is created nothing grows. Live jobs (no terminal event yet) always survive. --- crates/diffcore-tauri/src/activity_stream.rs | 92 +++++++++++++++++++- 1 file changed, 91 insertions(+), 1 deletion(-) diff --git a/crates/diffcore-tauri/src/activity_stream.rs b/crates/diffcore-tauri/src/activity_stream.rs index 883551c..83e05c8 100644 --- a/crates/diffcore-tauri/src/activity_stream.rs +++ b/crates/diffcore-tauri/src/activity_stream.rs @@ -3,7 +3,7 @@ use std::convert::Infallible; use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4, TcpListener as StdTcpListener}; use std::sync::Arc; use std::thread; -use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use async_stream::stream; use axum::extract::{Path, State}; @@ -90,8 +90,15 @@ impl JobEvent { struct JobState { history: Vec, sender: broadcast::Sender, + /// Set once the job completes or fails; until then the job is live and kept. + finished_at: Option, } +/// How long a finished job stays subscribable. Long enough that a reloaded page +/// can still collect the result it was waiting on, short enough that completed +/// jobs — whose history holds the full result payload — do not accumulate. +const JOB_RETENTION: Duration = Duration::from_secs(30 * 60); + #[derive(Default)] pub struct ActivityManager { jobs: RwLock>, @@ -124,12 +131,15 @@ impl ActivityManager { timestamp_ms: timestamp_ms(), }; + self.evict_finished(JOB_RETENTION).await; + let mut jobs = self.jobs.write().await; jobs.insert( job_id.clone(), JobState { history: vec![started], sender, + finished_at: None, }, ); @@ -148,9 +158,22 @@ impl ActivityManager { Some((job.history.clone(), job.sender.subscribe())) } + /// Drop finished jobs older than `retention`. Called when a new job starts, + /// which keeps the map bounded without a timer task: nothing is created, so + /// nothing grows. + pub async fn evict_finished(&self, retention: Duration) { + self.jobs + .write() + .await + .retain(|_, job| job.finished_at.is_none_or(|at| at.elapsed() < retention)); + } + async fn push_event(&self, job_id: &str, event: JobEvent) { let mut jobs = self.jobs.write().await; if let Some(job) = jobs.get_mut(job_id) { + if event.is_terminal() { + job.finished_at = Some(Instant::now()); + } job.history.push(event.clone()); let _ = job.sender.send(event); } @@ -312,3 +335,70 @@ fn timestamp_ms() -> u64 { .map(|duration| duration.as_millis() as u64) .unwrap_or(0) } + +#[cfg(test)] +mod tests { + use super::*; + + async fn finished_job(manager: &Arc) -> String { + let handle = manager + .create_job("refinement", "codex", "default", "Refining") + .await; + handle + .complete("refinement", serde_json::json!({ "refined_groups": [] })) + .await; + handle.job_id().to_string() + } + + #[tokio::test] + async fn evicts_finished_jobs_past_retention_and_keeps_live_ones() { + let manager = Arc::new(ActivityManager::new()); + let done = finished_job(&manager).await; + let live = manager + .create_job("refinement", "codex", "default", "Still going") + .await + .job_id() + .to_string(); + + manager.evict_finished(Duration::ZERO).await; + + assert!( + manager.subscribe(&done).await.is_none(), + "finished job should be dropped once past retention" + ); + assert!( + manager.subscribe(&live).await.is_some(), + "a job that has not reached a terminal event must survive eviction" + ); + } + + #[tokio::test] + async fn keeps_finished_jobs_within_retention_so_a_reload_can_collect_them() { + let manager = Arc::new(ActivityManager::new()); + let done = finished_job(&manager).await; + + manager.evict_finished(JOB_RETENTION).await; + + let (history, _) = manager.subscribe(&done).await.expect("job still retained"); + assert!(history.iter().any(JobEvent::is_terminal)); + } + + #[tokio::test] + async fn starting_a_job_sweeps_the_map() { + let manager = Arc::new(ActivityManager::new()); + let done = finished_job(&manager).await; + manager + .jobs + .write() + .await + .get_mut(&done) + .unwrap() + .finished_at = Some(Instant::now() - JOB_RETENTION - Duration::from_secs(1)); + + let _ = manager + .create_job("refinement", "codex", "default", "Next") + .await; + + assert!(manager.subscribe(&done).await.is_none()); + } +}