From e628cad72e175b17e4f35454edd06095f65b6062 Mon Sep 17 00:00:00 2001 From: nonoqing Date: Mon, 28 Sep 2026 11:10:19 +0800 Subject: [PATCH] fix: handle compacted tool specs and paused cron badges --- .../product_runtime/get_tool_spec_tool.rs | 2 +- .../product_runtime/loaded_spec_state.rs | 109 ++++++++++++++++- .../execution/agent-runtime/src/prompt.rs | 2 +- .../prompt_contracts.rs | 9 +- .../tool-contracts/src/deferred_tool.rs | 3 +- .../execution/tool-contracts/src/framework.rs | 11 +- .../tool-contracts/tests/tool_contracts.rs | 30 ++--- .../scheduled-jobs/cronJobCountsStore.test.ts | 114 ++++++++++++++++-- .../scheduled-jobs/cronJobCountsStore.ts | 30 +++-- 9 files changed, 266 insertions(+), 44 deletions(-) diff --git a/src/crates/assembly/core/src/agentic/tools/product_runtime/get_tool_spec_tool.rs b/src/crates/assembly/core/src/agentic/tools/product_runtime/get_tool_spec_tool.rs index 8537cadf86..0baf4d1cac 100644 --- a/src/crates/assembly/core/src/agentic/tools/product_runtime/get_tool_spec_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/product_runtime/get_tool_spec_tool.rs @@ -167,6 +167,6 @@ mod tests { assert!(result_for_assistant .as_deref() .unwrap_or_default() - .contains("already loaded in the current conversation")); + .contains("already loaded in the current context")); } } diff --git a/src/crates/assembly/core/src/agentic/tools/product_runtime/loaded_spec_state.rs b/src/crates/assembly/core/src/agentic/tools/product_runtime/loaded_spec_state.rs index f70d232eec..a940e06566 100644 --- a/src/crates/assembly/core/src/agentic/tools/product_runtime/loaded_spec_state.rs +++ b/src/crates/assembly/core/src/agentic/tools/product_runtime/loaded_spec_state.rs @@ -67,7 +67,12 @@ fn get_tool_spec_load_observation(message: &Message) -> Option openbitfun_agent_tools::LoadedDeferredToolSpec { @@ -77,6 +82,108 @@ mod tests { } } + #[test] + fn compaction_requires_reload_only_when_the_full_spec_leaves_the_context() { + let deferred_tools = vec!["Cron".to_string()]; + let input = json!({ "tool_name": "Cron" }); + let spec_result = Message::tool_result(ToolResult { + tool_id: "load-cron".to_string(), + tool_name: GET_TOOL_SPEC_TOOL_NAME.to_string(), + effective_tool_name: None, + result: json!({ + "tool_name": "Cron", + "catalog_generation": 42, + "description": "Manage scheduled jobs.", + "input_schema": { "type": "object", "properties": { "action": { "enum": ["list"] } } }, + }), + result_for_assistant: None, + is_error: false, + duration_ms: None, + image_attachments: None, + }); + let history = vec![ + Message::user("Check the scheduled jobs.".to_string()), + Message::assistant_with_tools( + String::new(), + vec![ToolCall { + tool_id: "load-cron".to_string(), + tool_name: GET_TOOL_SPEC_TOOL_NAME.to_string(), + arguments: input.clone(), + raw_arguments: None, + is_error: false, + parse_error: None, + recovered_from_truncation: false, + repair_kind: Default::default(), + }], + ), + spec_result.clone(), + ]; + let compressor = ContextCompressor::new(); + + // The same summary may be generated with or without a retained tool + // result. Only the actual retained result is an execution receipt. + for (recent_tokens, retains_spec) in [(0, false), (10_000, true)] { + let plan = compressor + .plan_compression("session", &history, 128_000, recent_tokens) + .unwrap() + .unwrap(); + let mut compressed = compressor + .compress_plan_with_contract( + "session", + plan, + None, + "The Cron definition was loaded with GetToolSpec earlier.".to_string(), + ) + .unwrap() + .messages; + let loaded = collect_product_loaded_deferred_tool_specs(&compressed, &deferred_tools); + let admission = validate_deferred_tool_usage( + "Cron", + &deferred_tools, + &loaded, + 42, + GET_TOOL_SPEC_TOOL_NAME, + ); + let names = loaded + .iter() + .map(|spec| spec.tool_name.clone()) + .collect::>(); + let reload_plan = resolve_get_tool_spec_execution_plan(&input, &names).unwrap(); + + if retains_spec { + assert_eq!(loaded, vec![loaded_spec("Cron")]); + assert!(admission.is_ok()); + assert!(matches!( + reload_plan, + GetToolSpecExecutionPlan::DuplicateLoad(_) + )); + } else { + assert!(loaded.is_empty()); + assert!(admission + .unwrap_err() + .to_string() + .contains("reload it even if the summary says it was loaded")); + assert!(matches!( + reload_plan, + GetToolSpecExecutionPlan::LoadDetail { tool_name: "Cron" } + )); + + // Reading the definition again restores normal admission. + compressed.push(spec_result.clone()); + let reloaded = + collect_product_loaded_deferred_tool_specs(&compressed, &deferred_tools); + validate_deferred_tool_usage( + "Cron", + &deferred_tools, + &reloaded, + 42, + GET_TOOL_SPEC_TOOL_NAME, + ) + .expect("a fresh spec must unlock Cron after compaction"); + } + } + } + #[test] fn product_loaded_spec_state_collects_visible_get_tool_spec_results() { let visible_get_tool_spec_result = Message::tool_result(ToolResult { diff --git a/src/crates/execution/agent-runtime/src/prompt.rs b/src/crates/execution/agent-runtime/src/prompt.rs index 8a04ba09f6..5f3f72ac51 100644 --- a/src/crates/execution/agent-runtime/src/prompt.rs +++ b/src/crates/execution/agent-runtime/src/prompt.rs @@ -16,7 +16,7 @@ const DIRECT_TOOL_LISTING_GUIDANCE: &str = r#"Their definitions are already avai Each entry below is a directly callable tool name."#; const DEFERRED_TOOL_LISTING_TITLE: &str = "## Deferred tools"; const DEFERRED_TOOL_LISTING_GUIDANCE: &str = r#"Their definitions are not loaded at the start of the conversation. -You must obtain the tool definition using GetToolSpec before you first invoke a deferred tool. Once its definition is available in the conversation, you can call it through CallDeferredTool. +Use GetToolSpec to read a deferred tool's full definition before invoking it through CallDeferredTool. Reuse it while the successful GetToolSpec result remains in the current context. If compaction or truncation removed that result, load it again; a summary or a past call does not keep the definition loaded. Each entry below is a deferred tool name with an optional short description."#; pub fn render_direct_tool_listing_body<'a>( diff --git a/src/crates/execution/agent-runtime/tests/agent_definition_contracts/prompt_contracts.rs b/src/crates/execution/agent-runtime/tests/agent_definition_contracts/prompt_contracts.rs index 4e9c76a0b1..421e88ba9b 100644 --- a/src/crates/execution/agent-runtime/tests/agent_definition_contracts/prompt_contracts.rs +++ b/src/crates/execution/agent-runtime/tests/agent_definition_contracts/prompt_contracts.rs @@ -70,9 +70,12 @@ fn tool_listing_sections_render_only_present_sections() { assert!(deferred_tool_listing .contains("Their definitions are not loaded at the start of the conversation.")); assert!(deferred_tool_listing - .contains("You must obtain the tool definition using GetToolSpec before you first invoke a deferred tool.")); + .contains("Use GetToolSpec to read a deferred tool's full definition before invoking it through CallDeferredTool.")); assert!(deferred_tool_listing.contains( - "Once its definition is available in the conversation, you can call it through CallDeferredTool." + "Reuse it while the successful GetToolSpec result remains in the current context." + )); + assert!(deferred_tool_listing.contains( + "If compaction or truncation removed that result, load it again; a summary or a past call does not keep the definition loaded." )); assert!(deferred_tool_listing .contains("Each entry below is a deferred tool name with an optional short description.")); @@ -80,7 +83,7 @@ fn tool_listing_sections_render_only_present_sections() { "## Direct tools\nTheir definitions are already available. You can call them directly.\nEach entry below is a directly callable tool name.\n\n\n- Read\n- GetToolSpec\n- CallDeferredTool\n" )); assert!(deferred_tool_listing.contains( - "## Deferred tools\nTheir definitions are not loaded at the start of the conversation.\nYou must obtain the tool definition using GetToolSpec before you first invoke a deferred tool. Once its definition is available in the conversation, you can call it through CallDeferredTool." + "## Deferred tools\nTheir definitions are not loaded at the start of the conversation.\nUse GetToolSpec" )); assert!(deferred_tool_listing.ends_with("Search: summary")); } diff --git a/src/crates/execution/tool-contracts/src/deferred_tool.rs b/src/crates/execution/tool-contracts/src/deferred_tool.rs index 0bdcdc62f5..1c2411e322 100644 --- a/src/crates/execution/tool-contracts/src/deferred_tool.rs +++ b/src/crates/execution/tool-contracts/src/deferred_tool.rs @@ -62,7 +62,7 @@ pub fn call_deferred_tool_input_schema() -> Value { "properties": { "tool_name": { "type": "string", - "description": "Exact deferred tool name previously loaded with GetToolSpec." + "description": "Exact deferred tool name whose full GetToolSpec result is still visible in the current context." }, "args": { "type": "object", @@ -80,6 +80,7 @@ pub fn call_deferred_tool_short_description() -> String { pub fn call_deferred_tool_description() -> String { r#"Call a deferred tool after reading its full schema with GetToolSpec. +The full GetToolSpec result must still be visible in the current context. If compaction or truncation removed it, reload it with GetToolSpec first; a summary or a past call is not a loaded definition. Pass the exact deferred tool name in tool_name and put only that tool's arguments inside args. The order is important. ALWAYS output tool_name first, then args."# .to_string() diff --git a/src/crates/execution/tool-contracts/src/framework.rs b/src/crates/execution/tool-contracts/src/framework.rs index 013ae9219e..c99677df60 100644 --- a/src/crates/execution/tool-contracts/src/framework.rs +++ b/src/crates/execution/tool-contracts/src/framework.rs @@ -91,7 +91,7 @@ impl fmt::Display for DeferredToolUsageError { get_tool_spec_tool_name, } => write!( formatter, - "Tool '{tool_name}' is deferred. Call {get_tool_spec_tool_name} first with {{\"tool_name\":\"{tool_name}\"}} to read its full usage instructions and input schema before invoking it." + "Tool '{tool_name}' has no loaded definition in the current context. Call {get_tool_spec_tool_name} with {{\"tool_name\":\"{tool_name}\"}} to read its full usage instructions and input schema before invoking it. If compaction removed an earlier definition, reload it even if the summary says it was loaded." ), Self::StaleSpec { tool_name, @@ -371,7 +371,8 @@ pub fn get_tool_spec_short_description() -> String { pub fn build_get_tool_spec_description() -> String { r#"Read the full schema before first calling a deferred tool through CallDeferredTool. -Do not call GetToolSpec again for a tool whose definition is already loaded in the current conversation."# +Do not call GetToolSpec again while its successful result with the full definition is still visible in the current context. +If compaction or truncation removed that result, call GetToolSpec again before using CallDeferredTool. A summary mentioning a previously loaded tool or a past successful call does not load its definition. Reload also when the runtime reports a stale definition."# .to_string() } @@ -439,7 +440,7 @@ pub fn validate_get_tool_spec_input(input: &Value) -> ValidationResult { pub fn build_get_tool_spec_duplicate_load_hint(tool_name: &str) -> String { format!( - "Tool '{}' is already loaded in the current conversation. Do not call GetToolSpec again for it. Use CallDeferredTool with tool_name '{}' and put the tool arguments inside args.", + "Tool '{}' is already loaded in the current context. Use CallDeferredTool with tool_name '{}' and put the tool arguments inside args. Reload with GetToolSpec only if its full definition leaves the context or the runtime reports it stale.", tool_name, tool_name ) } @@ -2688,6 +2689,10 @@ mod tests { assert!(description.contains("Read the full schema")); assert!(description.contains("Do not call GetToolSpec again")); + assert!(description.contains("full definition is still visible in the current context")); + assert!(description + .contains("If compaction or truncation removed that result, call GetToolSpec again")); + assert!(description.contains("does not load its definition")); } #[test] diff --git a/src/crates/execution/tool-contracts/tests/tool_contracts.rs b/src/crates/execution/tool-contracts/tests/tool_contracts.rs index 9da730ecf6..0c07534076 100644 --- a/src/crates/execution/tool-contracts/tests/tool_contracts.rs +++ b/src/crates/execution/tool-contracts/tests/tool_contracts.rs @@ -97,6 +97,8 @@ fn call_deferred_tool_contract_uses_nested_object_arguments() { assert!(call_deferred_tool_description() .contains("The order is important. ALWAYS output tool_name first, then args.")); + assert!(call_deferred_tool_description() + .contains("If compaction or truncation removed it, reload it with GetToolSpec first")); assert_eq!(schema["additionalProperties"], false); assert_eq!(schema["required"], json!(["tool_name", "args"])); assert_eq!(schema["properties"]["args"]["type"], "object"); @@ -1470,7 +1472,7 @@ fn deferred_tool_usage_gate_preserves_get_tool_spec_unlock_contract() { .expect_err("deferred tool should require GetToolSpec unlock"); assert_eq!( err.to_string(), - "Tool 'WebFetch' is deferred. Call GetToolSpec first with {\"tool_name\":\"WebFetch\"} to read its full usage instructions and input schema before invoking it." + "Tool 'WebFetch' has no loaded definition in the current context. Call GetToolSpec with {\"tool_name\":\"WebFetch\"} to read its full usage instructions and input schema before invoking it. If compaction removed an earlier definition, reload it even if the summary says it was loaded." ); let loaded_deferred_tool_specs = vec![LoadedDeferredToolSpec { @@ -2030,7 +2032,7 @@ fn get_tool_spec_contract_escapes_assistant_detail_for_xml_sections() { fn get_tool_spec_contract_preserves_duplicate_load_hint() { assert_eq!( build_get_tool_spec_duplicate_load_hint("WebFetch"), - "Tool 'WebFetch' is already loaded in the current conversation. Do not call GetToolSpec again for it. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args." + "Tool 'WebFetch' is already loaded in the current context. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args. Reload with GetToolSpec only if its full definition leaves the context or the runtime reports it stale." ); } @@ -2052,7 +2054,7 @@ fn get_tool_spec_contract_builds_duplicate_load_result() { assert_eq!( result_for_assistant.as_deref(), Some( - "Tool 'WebFetch' is already loaded in the current conversation. Do not call GetToolSpec again for it. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args." + "Tool 'WebFetch' is already loaded in the current context. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args. Reload with GetToolSpec only if its full definition leaves the context or the runtime reports it stale." ) ); assert_eq!(image_attachments, None); @@ -2127,7 +2129,7 @@ fn get_tool_spec_contract_plans_duplicate_load_without_core_context() { assert!(result_for_assistant .as_deref() .unwrap_or_default() - .contains("already loaded in the current conversation")); + .contains("already loaded in the current context")); assert_eq!(image_attachments, None); } @@ -3058,7 +3060,7 @@ async fn get_tool_spec_detail_resolver_preserves_contextual_detail_contract() { .expect("collapsed WebFetch detail"); assert_eq!(detail.tool_name, "WebFetch"); - assert_eq!(detail.description, "WebFetch description for agentic"); + assert_eq!(detail.description, "WebFetch description for Standard"); assert_eq!( detail.input_schema["properties"]["agent"]["const"], "Standard" @@ -3067,7 +3069,7 @@ async fn get_tool_spec_detail_resolver_preserves_contextual_detail_contract() { detail.to_value(), json!({ "tool_name": "WebFetch", - "description": "WebFetch description for agentic", + "description": "WebFetch description for Standard", "input_schema": { "type": "object", "properties": { @@ -3118,7 +3120,7 @@ async fn get_tool_spec_catalog_provider_preserves_runtime_catalog_contract() { .await .expect("provider-backed detail"); assert_eq!(detail.tool_name, "WebFetch"); - assert_eq!(detail.description, "WebFetch description for agentic"); + assert_eq!(detail.description, "WebFetch description for Standard"); } #[tokio::test] @@ -3153,7 +3155,7 @@ async fn get_tool_spec_provider_execution_returns_duplicate_result_without_detai assert!(result_for_assistant .as_deref() .unwrap_or_default() - .contains("already loaded in the current conversation")); + .contains("already loaded in the current context")); assert_eq!(image_attachments, None); } @@ -3189,15 +3191,15 @@ async fn get_tool_spec_provider_execution_returns_detail_result_from_provider() }; assert_eq!(data["tool_name"], "WebFetch"); - assert_eq!(data["description"], "WebFetch description for agentic"); + assert_eq!(data["description"], "WebFetch description for Standard"); assert_eq!( data["input_schema"]["properties"]["agent"]["const"], "Standard" ); let assistant = result_for_assistant.expect("assistant detail"); - assert!(assistant.contains("\nWebFetch description for agentic")); + assert!(assistant.contains("\nWebFetch description for Standard")); assert!(assistant.contains("\"agent\"")); - assert!(assistant.contains("\"agentic\"")); + assert!(assistant.contains("\"Standard\"")); assert_eq!(image_attachments, None); } @@ -3267,7 +3269,7 @@ async fn get_tool_spec_runtime_facade_owns_execution_path() { panic!("expected normal tool result"); }; assert_eq!(data["tool_name"], "WebFetch"); - assert_eq!(data["description"], "WebFetch description for agentic"); + assert_eq!(data["description"], "WebFetch description for Standard"); assert_eq!( data["input_schema"]["properties"]["agent"]["const"], "Standard" @@ -3306,7 +3308,7 @@ async fn get_tool_spec_runtime_facade_owns_tool_result_vector_adapter_shape() { assert_eq!(data["tool_name"], "WebFetch"); assert!(result_for_assistant .expect("assistant detail") - .contains("\nWebFetch description for agentic")); + .contains("\nWebFetch description for Standard")); assert_eq!(image_attachments, None); let duplicate_runtime = @@ -3339,7 +3341,7 @@ async fn get_tool_spec_runtime_facade_owns_tool_result_vector_adapter_shape() { assert_eq!( result_for_assistant.as_deref(), Some( - "Tool 'WebFetch' is already loaded in the current conversation. Do not call GetToolSpec again for it. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args." + "Tool 'WebFetch' is already loaded in the current context. Use CallDeferredTool with tool_name 'WebFetch' and put the tool arguments inside args. Reload with GetToolSpec only if its full definition leaves the context or the runtime reports it stale." ) ); assert!(image_attachments.is_none()); diff --git a/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.test.ts b/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.test.ts index 54ce6ee96c..5eb8e6fd81 100644 --- a/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.test.ts +++ b/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.test.ts @@ -1,12 +1,16 @@ -/** - * The nav badge is the only scheduled-job signal visible without opening a - * panel, so what it counts has to match what the panel lists. A job the agent - * created with "enabled": false used to disappear from the badge entirely. - */ +// @vitest-environment jsdom -import { describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import type { CronJob } from '@/infrastructure/api'; import { computeCronJobCounts } from './cronJobCountsStore'; +import { notifyScheduledJobsChanged, SCHEDULED_JOBS_CHANGED_EVENT } from './scheduledJobDraft'; + +const { listJobs, onJobsChanged } = vi.hoisted(() => ({ + listJobs: vi.fn(), + onJobsChanged: vi.fn(), +})); + +vi.mock('@/infrastructure/api', () => ({ cronAPI: { listJobs, onJobsChanged } })); function sessionJob(id: string, sessionId: string, workspaceId: string, enabled = true): CronJob { return { @@ -66,14 +70,25 @@ describe('computeCronJobCounts', () => { expect(counts.bySessionId.size).toBe(0); }); - it('keeps disabled jobs counted, matching the panel list', () => { + it('counts only enabled jobs when a session also has paused schedules', () => { const counts = computeCronJobCounts([ sessionJob('cron_e', 'session_1', 'ws_1', false), sessionJob('cron_f', 'session_1', 'ws_1'), ]); - expect(counts.byWorkspaceId.get('ws_1')).toBe(2); - expect(counts.bySessionId.get('session_1')).toBe(2); + expect(counts.byWorkspaceId.get('ws_1')).toBe(1); + expect(counts.bySessionId.get('session_1')).toBe(1); + }); + + it('clears clock counts when every schedule is paused', () => { + const pausedWorkspaceJob = { ...workspaceJob('cron_workspace', 'ws_1'), enabled: false }; + const counts = computeCronJobCounts([ + sessionJob('cron_e', 'session_1', 'ws_1', false), + pausedWorkspaceJob, + ]); + + expect(counts.byWorkspaceId.size).toBe(0); + expect(counts.bySessionId.size).toBe(0); }); it('ignores jobs whose workspace has no id', () => { @@ -86,3 +101,84 @@ describe('computeCronJobCounts', () => { expect(counts.bySessionId.get('session_1')).toBe(1); }); }); + +describe('scheduled-job badge updates', () => { + let store: typeof import('./cronJobCountsStore'); + let backendChanged: () => void; + let addEventListener: ReturnType>; + + beforeEach(async () => { + vi.resetModules(); + listJobs.mockReset().mockResolvedValue([]); + onJobsChanged.mockReset().mockImplementation((callback: () => void) => { + backendChanged = callback; + return () => {}; + }); + addEventListener = vi.spyOn(window, 'addEventListener'); + store = await import('./cronJobCountsStore'); + }); + + afterEach(() => { + for (const [type, listener, options] of addEventListener.mock.calls) { + window.removeEventListener(type, listener, options); + } + vi.restoreAllMocks(); + }); + + it('publishes pause, resume and deletion from local and backend changes', async () => { + const job = sessionJob('cron_a', 'session_1', 'ws_1'); + listJobs.mockResolvedValue([job]); + const notify = vi.fn(); + store.subscribeCronJobCounts(notify); + store.ensureCronJobCountsListener(); + store.ensureCronJobCountsListener(); + + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.get('session_1')).toBe(1)); + expect(onJobsChanged).toHaveBeenCalledTimes(1); + expect(addEventListener.mock.calls.filter(([type]) => type === SCHEDULED_JOBS_CHANGED_EVENT)).toHaveLength(1); + + listJobs.mockResolvedValue([{ ...job, enabled: false }]); + notifyScheduledJobsChanged('todos'); + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.size).toBe(0)); + + listJobs.mockResolvedValue([job]); + backendChanged(); + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.get('session_1')).toBe(1)); + + listJobs.mockResolvedValue([]); + backendChanged(); + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.size).toBe(0)); + expect(notify).toHaveBeenCalledTimes(4); + }); + + it('refreshes again when a pause arrives while an older read is in flight', async () => { + const job = sessionJob('cron_a', 'session_1', 'ws_1'); + listJobs.mockResolvedValueOnce([job]); + store.ensureCronJobCountsListener(); + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.get('session_1')).toBe(1)); + + let finishRead!: (jobs: CronJob[]) => void; + listJobs.mockReturnValueOnce(new Promise(resolve => { finishRead = resolve; })); + listJobs.mockResolvedValueOnce([{ ...job, enabled: false }]); + backendChanged(); + notifyScheduledJobsChanged('todos'); + backendChanged(); + expect(listJobs).toHaveBeenCalledTimes(2); + + finishRead([job]); + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.size).toBe(0)); + expect(listJobs).toHaveBeenCalledTimes(3); + }); + + it('drains a pending change even when the in-flight read fails', async () => { + let failRead!: (error: Error) => void; + listJobs.mockReturnValueOnce(new Promise((_resolve, reject) => { failRead = reject; })); + listJobs.mockResolvedValueOnce([sessionJob('cron_a', 'session_1', 'ws_1')]); + store.ensureCronJobCountsListener(); + backendChanged(); + failRead(new Error('Host temporarily unavailable')); + + await vi.waitFor(() => expect(store.getCronJobCountsSnapshot().bySessionId.get('session_1')).toBe(1)); + expect(listJobs).toHaveBeenCalledTimes(2); + }); +}); diff --git a/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.ts b/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.ts index 1d4c0eecf2..1fa8610130 100644 --- a/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.ts +++ b/src/web-ui/src/app/components/scheduled-jobs/cronJobCountsStore.ts @@ -17,9 +17,9 @@ import type { CronJob } from '@/infrastructure/api'; const log = createLogger('CronJobCountsStore'); export interface CronJobCountsSnapshot { - /** Count per workspace id. */ + /** Enabled job count per workspace id. */ byWorkspaceId: ReadonlyMap; - /** Count per session id (session-targeted jobs only). */ + /** Enabled job count per session id (session-targeted jobs only). */ bySessionId: ReadonlyMap; } @@ -31,6 +31,7 @@ const EMPTY_SNAPSHOT: CronJobCountsSnapshot = { let snapshot: CronJobCountsSnapshot = EMPTY_SNAPSHOT; const listeners = new Set<() => void>(); let loadInFlight = false; +let reloadPending = false; let listenerRegistered = false; function publish(): void { @@ -55,14 +56,14 @@ function countsEqual( } /** - * Counts per workspace and per session. Mirrors what the scheduled-jobs views - * list, disabled jobs included: a paused job the agent just created must still - * be visible in the nav. + * Navigation clocks indicate enabled schedules. Paused jobs remain available + * in the scheduled-job views, but do not keep a workspace/session clock active. */ export function computeCronJobCounts(jobs: CronJob[]): CronJobCountsSnapshot { const byWorkspaceId = new Map(); const bySessionId = new Map(); for (const job of jobs) { + if (!job.enabled) continue; const workspaceId = job.target.workspace.workspaceId; if (workspaceId) { byWorkspaceId.set(workspaceId, (byWorkspaceId.get(workspaceId) ?? 0) + 1); @@ -82,15 +83,22 @@ function applyJobs(jobs: CronJob[]): void { } async function reload(): Promise { + reloadPending = true; if (loadInFlight) return; loadInFlight = true; try { - const jobs = await cronAPI.listJobs({}); - applyJobs(jobs); - } catch (error) { - // Badges are a hint; keep the previous counts and let the next change - // signal retry. - log.warn('Failed to load scheduled job counts', { error }); + do { + reloadPending = false; + try { + const jobs = await cronAPI.listJobs({}); + // A change during this read may have made its result stale. Drain the + // queued refresh before publishing counts instead of losing the hint. + if (!reloadPending) applyJobs(jobs); + } catch (error) { + // Keep the previous counts and retry any change queued during the read. + log.warn('Failed to load scheduled job counts', { error }); + } + } while (reloadPending); } finally { loadInFlight = false; }