From 4f6d3e7e640aa59e07f7a0b7deeadb493bebd2f4 Mon Sep 17 00:00:00 2001 From: Joseph Zaki <325440518+jz-krono@users.noreply.github.com> Date: Sun, 4 Oct 2026 05:39:28 +0900 Subject: [PATCH] fix: hold interrupted scheduled turns for explicit retry --- README.md | 2 + src/client/TaskActions.tsx | 12 ++--- src/client/WorkspaceDialog.tsx | 3 +- src/client/style.css | 3 +- src/server/runner.ts | 4 +- src/server/store.ts | 16 +++---- src/shared/types.ts | 8 +++- tests/controls.test.tsx | 13 ++++++ tests/runner.test.ts | 8 ++-- tests/store.test.ts | 81 ++++++++++++++++++++++++++++++---- 10 files changed, 120 insertions(+), 30 deletions(-) diff --git a/README.md b/README.md index 49dfde4..e788cb0 100644 --- a/README.md +++ b/README.md @@ -180,6 +180,8 @@ See [Setup](docs/SETUP.md) for configuration, Slack, calls, the browser service, | Automatic Learning | Per-Dot Learning containers, conversation evidence routing, and published-skill delivery; see [setup](docs/SETUP.md#automatic-learning) | | Deployment | Local Node setup and separate application/browser containers | +Scheduled tasks run in their original conversation. If a worker stops or its lease expires during a run, OpenDots marks that run **Interrupted** and waits for an explicit retry. Review its pages and computer actions, then use **Retry after review** when appropriate. Completed effects may already be present even when a run has no final result. + Local checks cover setup, persistence, permissions, SDK failure handling, and browser isolation. Automated tests use service fixtures. **Live Intelligence, model responses, and page-context chat were verified on September 29, 2026.** Live OpenBot computer browsing, file creation, shell verification, and file persistence across stop/start were also verified locally. Live Realtime speech, call controls, and receipt persistence were verified locally on September 30, 2026. Slack and spoken compute delegation still need connected-service verification. See [recording notes](docs/demos/README.md) for the demonstrated flows and limits. Automatic Learning routing and skill delivery are configured locally. Cloud schedules, eligible-thread counts, and published-skill delivery still need connected-service verification. Skills require review and publication in Intelligence; existing conversations without a container are not enrolled retroactively. diff --git a/src/client/TaskActions.tsx b/src/client/TaskActions.tsx index 7dca5cf..80bc1fa 100644 --- a/src/client/TaskActions.tsx +++ b/src/client/TaskActions.tsx @@ -27,11 +27,13 @@ export function TaskActions({ onClick={() => onAction('run')} > - {task.status === 'failed' - ? 'Retry task' - : task.status === 'paused' - ? 'Resume task' - : 'Run again'} + {task.status === 'interrupted' + ? 'Retry after review' + : task.status === 'failed' + ? 'Retry task' + : task.status === 'paused' + ? 'Resume task' + : 'Run again'} )} {task.status === 'completed' && !!task.intervalSeconds && ( diff --git a/src/client/WorkspaceDialog.tsx b/src/client/WorkspaceDialog.tsx index 234fa4d..2e00056 100644 --- a/src/client/WorkspaceDialog.tsx +++ b/src/client/WorkspaceDialog.tsx @@ -354,7 +354,8 @@ export function WorkspaceDialog({

Runs on the server in this same conversation, even with the tab - closed. Failed runs wait for manual retry. + closed. Failed or interrupted runs wait for manual retry. Review + completed work before retrying an interrupted run.

)} diff --git a/src/client/style.css b/src/client/style.css index 0b11b0b..88c9f61 100644 --- a/src/client/style.css +++ b/src/client/style.css @@ -780,7 +780,8 @@ h3 { color: #788bc3; background: #edf0f9; } -.status.failed { +.status.failed, +.status.interrupted { color: #b17d73; background: #faf0ed; } diff --git a/src/server/runner.ts b/src/server/runner.ts index d7ee26c..d56b705 100644 --- a/src/server/runner.ts +++ b/src/server/runner.ts @@ -26,9 +26,9 @@ export class Runner { this.timer = undefined; for (const task of this.store.tasks()) { if (this.active.has(task.id) && task.lease) - this.store.release( + this.store.interrupt( { ...task, lease: task.lease }, - 'Server stopping; queued for restart.', + 'Server stopped during this run. Review completed effects before retrying.', ); } this.abortAll(); diff --git a/src/server/store.ts b/src/server/store.ts index b361de9..6b0b601 100644 --- a/src/server/store.ts +++ b/src/server/store.ts @@ -74,8 +74,8 @@ export class Store { for (const task of running) { this.invalidate( task, - 'queued', - 'Run stopped because settings changed.', + 'interrupted', + 'Run interrupted because settings changed. Review completed effects before retrying.', ); } } @@ -141,9 +141,9 @@ export class Store { .run(now, reason, task.lease); this.db .prepare( - 'UPDATE tasks SET status=?, lease=NULL, leaseUntil=NULL, updatedAt=? WHERE id=?', + 'UPDATE tasks SET status=?, lease=NULL, leaseUntil=NULL, updatedAt=?, error=? WHERE id=?', ) - .run(status, now, task.id); + .run(status, now, status === 'interrupted' ? reason : null, task.id); this.event(task.id, task.lease, reason); } action(id: string, action: Action): Task | undefined { @@ -199,8 +199,8 @@ export class Store { for (const task of expired) this.invalidate( task, - 'queued', - 'Previous worker lease expired; safely retrying.', + 'interrupted', + 'Worker lease expired. Review completed effects before retrying.', ); const task = this.db .prepare( @@ -255,9 +255,9 @@ export class Store { return true; }); } - release(claim: Claim, reason: string) { + interrupt(claim: Claim, reason: string) { this.transaction(() => { - if (this.owns(claim)) this.invalidate(claim, 'queued', reason); + if (this.owns(claim)) this.invalidate(claim, 'interrupted', reason); }); } fail(claim: Claim, error: string) { diff --git a/src/shared/types.ts b/src/shared/types.ts index c87b548..af5a73c 100644 --- a/src/shared/types.ts +++ b/src/shared/types.ts @@ -1,5 +1,11 @@ export type Status = - 'queued' | 'running' | 'paused' | 'completed' | 'failed' | 'cancelled'; + | 'queued' + | 'running' + | 'paused' + | 'completed' + | 'failed' + | 'interrupted' + | 'cancelled'; export interface Settings { name: string; paused: boolean; diff --git a/tests/controls.test.tsx b/tests/controls.test.tsx index 04a9429..a5dc2b0 100644 --- a/tests/controls.test.tsx +++ b/tests/controls.test.tsx @@ -46,3 +46,16 @@ it('offers resume for a paused task without claiming to be running', () => { expect(html).toContain('Resume task'); expect(html).not.toContain('Pause schedule'); }); +it('labels an interrupted task for owner review before retry', () => { + const html = renderToStaticMarkup( + {}} + onSchedule={() => {}} + />, + ); + expect(html).toContain('Retry after review'); + expect(html).not.toContain('Pause task'); +}); diff --git a/tests/runner.test.ts b/tests/runner.test.ts index f55f507..d621ef8 100644 --- a/tests/runner.test.ts +++ b/tests/runner.test.ts @@ -34,7 +34,7 @@ it('aborts research when permissions are revoked outside the runner instance', a await tick; expect(requestSignal?.aborted).toBe(true); expect(request).toHaveBeenCalledOnce(); - expect(store.tasks()[0].status).toBe('queued'); + expect(store.tasks()[0].status).toBe('interrupted'); runner.stop(); store.close(); }); @@ -73,7 +73,7 @@ it('omits stored memories from research when memory permission is disabled', asy ); store.close(); }); -it('requeues active work on graceful shutdown instead of losing it', async () => { +it('holds active work for review on graceful shutdown', async () => { const store = new Store(':memory:'); const runner = new Runner(store, config); const fetch = vi.fn( @@ -92,7 +92,9 @@ it('requeues active work on graceful shutdown instead of losing it', async () => await vi.waitFor(() => expect(fetch).toHaveBeenCalledOnce()); runner.stop(); await pending; - expect(store.tasks()[0].status).toBe('queued'); + expect(store.tasks()[0].status).toBe('interrupted'); + expect(store.claim()).toBeNull(); + store.action(store.tasks()[0].id, 'run'); expect(store.claim()).toBeTruthy(); store.close(); }); diff --git a/tests/store.test.ts b/tests/store.test.ts index 9118685..54e4335 100644 --- a/tests/store.test.ts +++ b/tests/store.test.ts @@ -3,6 +3,8 @@ import { mkdtempSync, rmSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { Store } from '../src/server/store.js'; +import { WorkspaceStore } from '../src/server/workspace.js'; +import { pageAccess } from '../src/server/page-tools.js'; const resources: { store: Store; dir: string }[] = []; function fixture() { @@ -71,21 +73,82 @@ describe('durable task lifecycle', () => { store.finish(claim, { text: 'Late result', sources: [], sample: true }), ).toBe(false); store.updateSettings({ paused: false }); + expect(store.claim(Date.now())).toBeNull(); + expect(store.task(task.id)?.status).toBe('interrupted'); + store.action(task.id, 'run'); expect(store.claim(Date.now())?.id).toBe(task.id); }); - it('recovers expired work after restart without duplicate completion', () => { - const { store } = fixture(); + it('holds expired work until the owner retries it', () => { + const { store, path } = fixture(); store.createTask('Recover me'); const now = Date.now(); const old = store.claim(now)!; - const recovered = store.claim(now + 180_001)!; - expect(recovered.id).toBe(old.id); - expect(recovered.lease).not.toBe(old.lease); - expect(store.finish(old, { text: 'Old', sources: [], sample: true })).toBe( + const restarted = new Store(path); + try { + expect(restarted.claim(now + 180_001)).toBeNull(); + expect(restarted.task(old.id)?.status).toBe('interrupted'); + expect(restarted.detail(old.id)?.runs[0].status).toBe('interrupted'); + restarted.action(old.id, 'run'); + const recovered = restarted.claim(now + 180_001)!; + expect(recovered.id).toBe(old.id); + expect(recovered.lease).not.toBe(old.lease); + expect( + store.finish(old, { text: 'Old', sources: [], sample: true }), + ).toBe(false); + expect( + restarted.finish(recovered, { + text: 'New', + sources: [], + sample: true, + }), + ).toBe(true); + } finally { + restarted.close(); + } + }); + it('claims other queued work while an expired run waits for review', () => { + const { store } = fixture(); + const interrupted = store.createTask('Create the first page'); + const now = Date.now(); + const old = store.claim(now)!; + const next = store.createTask('Research another topic'); + + expect(store.claim(now + 180_001)?.id).toBe(next.id); + expect(store.task(interrupted.id)?.status).toBe('interrupted'); + expect(store.detail(interrupted.id)?.runs[0].status).toBe('interrupted'); + expect(store.finish(old, { text: 'Late', sources: [], sample: true })).toBe( false, ); - expect( - store.finish(recovered, { text: 'New', sources: [], sample: true }), - ).toBe(true); + }); + it('holds an expired run for review after a local page effect', () => { + const { store, path } = fixture(); + const workspace = new WorkspaceStore(path, 'synthetic-owner'); + try { + const dot = workspace.dots()[0]; + const thread = workspace.bindThread( + 'scheduled-thread', + dot.id, + 'Scheduled work', + ); + const task = store.createTask('Create the sample page once'); + workspace.bindTask(task.id, thread.id); + const now = Date.now(); + const old = store.claim(now)!; + const pages = pageAccess(workspace, dot.spaceId, thread.id, () => {}); + pages.create({ title: 'Sample', content: 'Synthetic fixture body' }); + + expect(store.claim(now + 180_001)).toBeNull(); + expect(store.task(task.id)?.status).toBe('interrupted'); + expect(store.detail(task.id)?.runs[0].status).toBe('interrupted'); + expect(workspace.pages.list(dot.spaceId)).toHaveLength(1); + expect( + store.finish(old, { text: 'Late', sources: [], sample: true }), + ).toBe(false); + + store.action(task.id, 'run'); + expect(store.claim(now + 180_002)?.id).toBe(task.id); + } finally { + workspace.close(); + } }); });