From ec4be25813d655a2665b7e0acf55ca60b7536042 Mon Sep 17 00:00:00 2001 From: probe Date: Mon, 5 Oct 2026 15:13:34 +0800 Subject: [PATCH] feat(webui): run the WebUI's own scheduled tasks, and record how that ends MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The WebUI is a loopback process the user starts. Neither cron engine in the tree can serve it: the v1 path is pinned off for v2-compat hosts and its table is empty after the v2 migration moved the data out, and the v2 `CronService` is created only for an Electron owner and is not carried on the host contract in any case. Measured on a real WebUI host; ADR 0012 carries the full chain and the decision. So the WebUI runs its own: a store at `/webui/scheduled-tasks.sqlite` with its own table, an in-process tick beside the heartbeat in `WebuiService`, and a `WebuiScheduledTaskPort` whose six methods name operations rather than an engine. It never selects from the runtime's cron tables — that would be a second scheduler with none of v2's concurrency guards writing into the same agent queue as the desktop client. Being alive is the execution guarantee. Slots missed while the process was down are **not** replayed; they are counted (`missed_count`, `last_missed_at_ms`) so the surface can say so rather than imply the task never existed. `better-sqlite3` is declared here for the first time: the WebUI did not persist anything before, and it had been reaching the module through the workspace's hoisted `node_modules` without declaring it. ADR 0012 states the retirement condition as two independently checkable gates, and the probe test reports on them on every run — quietly while both are shut, and with a notice when the runtime moves. It never fails a build: a gate that goes red for something that is not a defect only teaches people to ignore red. Retirement is the adapter, the table migration, and tearing down this tick loop in the same change; two schedulers over one queue is what the decision exists to prevent. --- ...uled-tasks-until-the-runtime-offers-one.md | 118 ++++ packages/webui/package.json | 1 + packages/webui/src/server/host.ts | 57 ++ packages/webui/src/server/operation/names.ts | 6 + .../server/operation/operation-handlers.ts | 16 + .../webui/src/server/operation/operations.ts | 13 + .../src/server/operation/scheduled-task.ts | 372 +++++++++++ packages/webui/src/server/port.ts | 147 +++- .../src/server/scheduled-task-scheduler.ts | 466 +++++++++++++ .../webui/src/server/scheduled-task-store.ts | 482 ++++++++++++++ packages/webui/src/server/service.ts | 149 +++++ .../unit/webui-host-shape-invariant.test.ts | 45 ++ ...ebui-scheduled-task-upstream-probe.test.ts | 624 +++++++++++++++++ .../test/unit/webui-scheduled-task.test.ts | 629 ++++++++++++++++++ .../webui/test/unit/webui-service.test.ts | 79 +++ packages/webui/tsconfig.test.json | 1 + pnpm-lock.yaml | 3 + release/public-source.json | 6 + test/vitest-suites.json | 2 + 19 files changed, 3215 insertions(+), 1 deletion(-) create mode 100644 docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md create mode 100644 packages/webui/src/server/operation/scheduled-task.ts create mode 100644 packages/webui/src/server/scheduled-task-scheduler.ts create mode 100644 packages/webui/src/server/scheduled-task-store.ts create mode 100644 packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts create mode 100644 packages/webui/test/unit/webui-scheduled-task.test.ts diff --git a/docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md b/docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md new file mode 100644 index 000000000..2845520f8 --- /dev/null +++ b/docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md @@ -0,0 +1,118 @@ +# The WebUI runs its own scheduled tasks until the runtime offers one + +The WebUI ships a scheduled-task surface backed by its own store and its own +in-process tick loop. It does not wait for the runtime's cron capability, which +is not reachable from a loopback WebUI host, and it does not read the runtime's +cron tables. This is a **temporary** arrangement with a written retirement +condition below, not a second permanent product. + +## Why the runtime's cron is not reachable from here + +Measured on `webui` at `2f064db` (2026-10-05), against a real WebUI host: + +- **The v1 path is off and its data is gone.** `local-runtime-v2/src/compat/v1/runtime.ts` + pins `cronConsumerEnabled: false` for every v2-compat host, and + `local-runtime/src/cron/api.ts` returns early from `ensureStarted` in that + case, so its registry is never filled. The table it reads, + `local_runtime_crons`, holds **0 rows**; the migration + `local-runtime-v2/src/infra/db/migrations/cron/migration-0002-copy-legacy-cron-data.ts` + moved the data out. +- **The v2 path is not created for this host.** `enableCron` comes from + `ownsElectronRuntimeCapabilities(runtimeOwnerKind)` + (`local-runtime-v2/src/application/agent/runtime-browser-use-composition.ts`), + and the WebUI declares `runtimeOwnerKind: "tui"` + (`packages/webui/src/server/assembly.ts`), so no `CronService` is built. +- **The service would not be visible even if it were built.** + `CreatedLocalRuntimeHost` carries only `application?` and `cliService?` + (`local-runtime-v2/src/local/host-contract.ts`). The one place the services + object is spread whole is the `cliService`'s own `options` + (`local-runtime-v2/src/runtime.ts`), and there the `cron` slot is + `undefined` for the same reason. + +An earlier attempt to reach it by declaring the WebUI an Electron owner made +things worse and was reverted: `isV2RuntimeOwner` requires +`capabilities.electronHost` for that kind, and without it the whole v2 service +group — including the `cliService` every other WebUI feature depends on — goes +away. A closed loop over the reachable object graph (depth 5) found no +`CronService` instance anywhere the WebUI can already hold. + +## What the WebUI owns instead + +- Its own file, `/webui/scheduled-tasks.sqlite`, and its own table. + It never selects from `local_runtime_crons` or + `local_runtime_v2_cron_definitions`. Reading the latter from here would be the + worst available combination: a second scheduler with none of v2's + concurrency guards, writing turns into the same agent queue as the desktop + client. +- An in-process tick loop in `WebuiService`, beside the existing heartbeat. The + runtime itself has no timer, so a runtime not attached to a service never + fires. +- A port, `WebuiScheduledTaskPort`, whose six methods — `listScheduledTasks`, + `createScheduledTask`, `updateScheduledTask`, `deleteScheduledTask`, + `triggerScheduledTaskNow`, `getScheduledTaskCapability` — name the + operations and not the engine. That is what makes the retirement below a + one-adapter change. + +## Considered Options + +- **Use the v1 cron path** — rejected: the switch is pinned off, and its table + is empty, so it would display nothing and could not be made to. +- **Read v2's tables directly** — rejected for the reason above: it is the one + combination with two executors and no claim. See + `local-runtime-v2/src/service/cron/adapters/run.repository.ts`, where + `claimExecution` is an atomic conditional update and + `insertPendingScheduled` converges on a trigger id. All of that protection + lives in the repository, and bypassing it is the whole cost. +- **Open the runtime's gate and wait** — the runtime's own comment marks cron + as an Electron-only capability absent on embedded hosts, so this is a + deliberate design position, not an oversight. The WebUI is a loopback process + the user starts; asking it to become the execution role is an upstream + architecture decision, and it was raised with the upstream maintainer rather + than decided here. +- **Keep the panel and show an empty list** — rejected: an empty list reads as + "you have no scheduled tasks", which is a statement about the user when it is + actually a statement about the host. + +## Retirement condition + +**Retire this as soon as the runtime hands a WebUI host a cron capability the +host can actually reach.** Concretely, both of these become true: + +1. `services.cron` is created for a `tui` + `cliEmbedded` host — today + `enableCron` is gated on `runtimeOwnerKind === "electron" | undefined`; and +2. the created service is reachable from the host, i.e. carried on + `CreatedLocalRuntimeHost` (or on the `cliService` options) rather than + sealed inside `createLocalRuntimeHostV2`'s locals. + +The retirement itself: + +1. Implement `WebuiScheduledTaskPort` against the runtime's `CronService` and + delete the in-process store, the tick loop and the SQLite dependency. Only + the adapter changes; the port, the operations and the panel do not. +2. Migrate `webui_scheduled_task` into the runtime's table. Tasks created before + the switch are otherwise stranded, so the migration is part of the change and + not a follow-up. +3. **Never run both.** The own loop must be torn down in the same change that + arms the runtime's scheduler. Two schedulers over one queue is precisely the + failure this decision exists to avoid. + +While the runtime keeps the capability closed, this decision stands as written. +It is not a claim that the WebUI is the right owner of scheduled execution; it +is that a panel which cannot see a real task list is worse than one that owns a +small, honest one. + +## Consequences + +- **A scheduled task in the WebUI and a scheduled task in the desktop client + are two different things that cannot see each other.** This is the price, and + it is paid knowingly. +- **The WebUI process being alive is the execution guarantee.** The package is + started on demand (`mcode-webui`, loopback port 8787); when nobody is running + it, nothing fires. Slots missed while it was down are **not** replayed — they + are counted (`missed_count`, `last_missed_at_ms`) so the surface can say so + rather than imply the task never existed. +- **The scope is fixed and the scheduler is dumb.** Once and fixed interval, + no cron expressions, no timezone and no active hours. A cron expression is a + new capability, not a gap in this one. +- **The WebUI takes a first dependency it did not have.** It persists now: + `better-sqlite3` is declared in `packages/webui/package.json` for this reason. diff --git a/packages/webui/package.json b/packages/webui/package.json index 6debf7c59..e9edb4211 100644 --- a/packages/webui/package.json +++ b/packages/webui/package.json @@ -30,6 +30,7 @@ "@mavis/shared": "workspace:^", "@xterm/addon-fit": "0.11.0", "@xterm/xterm": "6.0.0", + "better-sqlite3": "12.11.1", "highlight.js": "10.7.3", "katex": "0.18.7", "lottie-web": "^5.13.0", diff --git a/packages/webui/src/server/host.ts b/packages/webui/src/server/host.ts index d5496cee4..796a8b9a5 100644 --- a/packages/webui/src/server/host.ts +++ b/packages/webui/src/server/host.ts @@ -19,6 +19,7 @@ import { extractWorkspaceArchiveDirectory, readWorkspaceArchiveListing, } from "./workspace-archive.js"; +import type { WebuiScheduledTaskRuntime } from "./scheduled-task-scheduler.js"; import type { WebuiHarnessPort, WebuiSessionListRequest, @@ -339,6 +340,14 @@ export interface WebuiRuntimeHostHandle { * the runtime host; the WebUI only needs the structural shape to forward. */ readonly cliService?: WebuiRuntimeCliService; + /** + * Scheduled tasks, when the host carries a runtime for them. Optional on + * purpose: a host that predates this surface still type-checks, and every + * scheduled-task method below fails closed with one clear message instead of + * crashing on an undefined call. The WebUI's own service supplies its local + * runtime rather than routing through here — see `service.ts`. + */ + readonly scheduledTasks?: WebuiScheduledTaskRuntime; } /** @@ -444,6 +453,38 @@ export function createHarnessPortFromHost( async clearGoal(request) { return { success: await requireCliService(host).clearGoal(request.sessionId) }; }, + // Scheduled tasks. The capability probe answers instead of throwing, so a + // client can ask "can this host do scheduled tasks at all?" and render the + // reason; the other five are operations, and an operation that cannot be + // served reports the same reason as a `harness_error`. + async listScheduledTasks(request) { + return requireScheduledTasks(host).listScheduledTasks(request); + }, + async createScheduledTask(request) { + return requireScheduledTasks(host).createScheduledTask(request); + }, + async updateScheduledTask(request) { + return requireScheduledTasks(host).updateScheduledTask(request); + }, + async deleteScheduledTask(request) { + return requireScheduledTasks(host).deleteScheduledTask(request); + }, + async triggerScheduledTaskNow(request) { + return requireScheduledTasks(host).triggerScheduledTaskNow(request); + }, + async getScheduledTaskCapability() { + return host.scheduledTasks + ? host.scheduledTasks.getScheduledTaskCapability() + : { + available: false, + // No implementation answered, which is a different fact from "our + // own implementation is unavailable". See + // `WebuiScheduledTaskCapabilitySource`. + source: "none" as const, + reason: + "scheduled tasks are not available: this host exposes no scheduled-task runtime", + }; + }, async listWorkspaceFileTree(request) { const tree = await requireCliService(host).listWorkspaceFileTree!(request) as readonly WebuiWorkspaceFile[]; // The runtime reports names and shape but no file facts, while the port @@ -754,6 +795,22 @@ export function createHarnessPortFromHost( * helper so the failure message is the same as it was before the batch-C * seam work. */ +/** + * Resolve the host's scheduled-task runtime, or fail closed. The message is + * the one the capability probe reports, so a client that asked first and a + * client that called blind see the same reason. + */ +function requireScheduledTasks( + host: WebuiRuntimeHostHandle, +): WebuiScheduledTaskRuntime { + const runtime = host.scheduledTasks; + if (!runtime) + throw new Error( + "scheduled tasks are not available: this host exposes no scheduled-task runtime", + ); + return runtime; +} + function requireCliService(host: WebuiRuntimeHostHandle): WebuiRuntimeCliService { if (!host.cliService) throw new Error("runtime host does not expose the CLI service"); diff --git a/packages/webui/src/server/operation/names.ts b/packages/webui/src/server/operation/names.ts index af111ec1b..ebf2194fe 100644 --- a/packages/webui/src/server/operation/names.ts +++ b/packages/webui/src/server/operation/names.ts @@ -88,3 +88,9 @@ export const REVEAL_MODEL_PROVIDER_API_KEY_OPERATION_NAME = "revealModelProvider export const START_CODEX_OAUTH_LOGIN_OPERATION_NAME = "startCodexOAuthLogin" as const; export const CANCEL_CODEX_OAUTH_LOGIN_OPERATION_NAME = "cancelCodexOAuthLogin" as const; export const REFRESH_MODELS_OPERATION_NAME = "refreshModels" as const; +export const LIST_SCHEDULED_TASKS_OPERATION_NAME = "listScheduledTasks" as const; +export const CREATE_SCHEDULED_TASK_OPERATION_NAME = "createScheduledTask" as const; +export const UPDATE_SCHEDULED_TASK_OPERATION_NAME = "updateScheduledTask" as const; +export const DELETE_SCHEDULED_TASK_OPERATION_NAME = "deleteScheduledTask" as const; +export const TRIGGER_SCHEDULED_TASK_OPERATION_NAME = "triggerScheduledTaskNow" as const; +export const GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME = "getScheduledTaskCapability" as const; diff --git a/packages/webui/src/server/operation/operation-handlers.ts b/packages/webui/src/server/operation/operation-handlers.ts index ed845f1da..9115cdceb 100644 --- a/packages/webui/src/server/operation/operation-handlers.ts +++ b/packages/webui/src/server/operation/operation-handlers.ts @@ -107,6 +107,12 @@ export type WebuiOperationPort = Pick< | "refreshModels" | "requestCompaction" | "invalidateAuth" + | "listScheduledTasks" + | "createScheduledTask" + | "updateScheduledTask" + | "deleteScheduledTask" + | "triggerScheduledTaskNow" + | "getScheduledTaskCapability" >; export type WebuiOperationHandlers = { @@ -278,6 +284,16 @@ export function createOperationHandlers( createGoal: async (_context, body) => ({ body: await port.createGoal(body) }), patchGoal: async (_context, body) => ({ body: await port.patchGoal(body) }), clearGoal: async (_context, body) => ({ body: await port.clearGoal(body) }), + // Scheduled tasks are WebUI-owned, but they reach the client through the + // same registry as everything else, so the panel needs no transport of its + // own. The port methods are the seam; `WebuiService` supplies its own local + // runtime behind them. + listScheduledTasks: async (_context, body) => ({ body: await port.listScheduledTasks(body) }), + createScheduledTask: async (_context, body) => ({ body: await port.createScheduledTask(body) }), + updateScheduledTask: async (_context, body) => ({ body: await port.updateScheduledTask(body) }), + deleteScheduledTask: async (_context, body) => ({ body: await port.deleteScheduledTask(body) }), + triggerScheduledTaskNow: async (_context, body) => ({ body: await port.triggerScheduledTaskNow(body) }), + getScheduledTaskCapability: async () => ({ body: await port.getScheduledTaskCapability() }), listSessions: async (_context, body) => ({ body: await port.listSessions(body) }), listVisibleProjects: async (_context, body) => { if (!port.listVisibleProjects) throw new Error("runtime host does not expose project listing"); diff --git a/packages/webui/src/server/operation/operations.ts b/packages/webui/src/server/operation/operations.ts index 5679d2e17..819a62caf 100644 --- a/packages/webui/src/server/operation/operations.ts +++ b/packages/webui/src/server/operation/operations.ts @@ -7,6 +7,7 @@ import { watchEventsOperation, listPendingPermissionsOperation, getPendingQuesti import { abortSessionOperation, listQueueMessagesOperation, deleteQueueItemOperation, listModelsOperation, listSkillsOperation, selectModelOperation, getSessionUsageOperation, getUsageQuotaOperation, getAccountStatusOperation } from "./queue.js"; import { pluginManagementOperation } from "./plugin-management.js"; import { getPermissionModeOperation, setPermissionModeOperation } from "./permission-mode.js"; +import { listScheduledTasksOperation, createScheduledTaskOperation, updateScheduledTaskOperation, deleteScheduledTaskOperation, triggerScheduledTaskNowOperation, getScheduledTaskCapabilityOperation } from "./scheduled-task.js"; import { archiveSessionOperation, deleteSessionOperation, updateSessionOperation, getSessionForkOptionsOperation, forkSessionOperation, listUserModelProvidersOperation, createUserModelProviderOperation, updateUserModelProviderOperation, deleteUserModelProviderOperation, testUserModelProviderOperation, testUserModelOperation, discoverUserModelsCandidateOperation, saveUserModelProviderCandidateOperation, listProviderPresetsOperation, getMiniMaxApiKeyStatusOperation, upsertMiniMaxApiKeyOperation, getCodexOAuthStatusOperation, getMiniMaxModelSourceOperation, setMiniMaxModelSourceOperation, testUserModelCandidateOperation, revealModelProviderApiKeyOperation, startCodexOAuthLoginOperation, cancelCodexOAuthLoginOperation, refreshModelsOperation, runCommandOperation, getSigninPanelOperation, claimSigninOperation, signOutOperation } from "./provider.js"; export { versionOperation, listSessionsOperation, listVisibleProjectsOperation, getSessionTreeOperation, createSessionOperation, getSessionOperation, getActiveTurnOperation } from "./session.js"; export { listWorkspaceFileTreeOperation, browseWorkspaceDirsOperation, readWorkspaceFileOperation, getWorkspaceEnvironmentOperation, mutateWorkspaceGitOperation, getWorkspaceReviewSummaryOperation, listWorkspaceReviewFileDiffsOperation, getWorkspaceReviewFileContentOperation, searchWorkspaceReviewDiffsOperation, readCanvasOperation, applyCanvasOperation, readWorkspaceArchiveOperation, extractWorkspaceArchiveOperation, createTerminalOperation, listTerminalsOperation, writeTerminalOperation, resizeTerminalOperation, disposeTerminalOperation, watchTerminalOperation } from "./workspace.js"; @@ -17,6 +18,7 @@ export { watchEventsOperation, listPendingPermissionsOperation, getPendingQuesti export { abortSessionOperation, listQueueMessagesOperation, deleteQueueItemOperation, listModelsOperation, listSkillsOperation, selectModelOperation, getSessionUsageOperation, getUsageQuotaOperation, getAccountStatusOperation } from "./queue.js"; export { pluginManagementOperation } from "./plugin-management.js"; export { getPermissionModeOperation, setPermissionModeOperation } from "./permission-mode.js"; +export { listScheduledTasksOperation, createScheduledTaskOperation, updateScheduledTaskOperation, deleteScheduledTaskOperation, triggerScheduledTaskNowOperation, getScheduledTaskCapabilityOperation } from "./scheduled-task.js"; export { archiveSessionOperation, deleteSessionOperation, updateSessionOperation, getSessionForkOptionsOperation, forkSessionOperation, listUserModelProvidersOperation, createUserModelProviderOperation, updateUserModelProviderOperation, deleteUserModelProviderOperation, testUserModelProviderOperation, testUserModelOperation, discoverUserModelsCandidateOperation, saveUserModelProviderCandidateOperation, listProviderPresetsOperation, getMiniMaxApiKeyStatusOperation, upsertMiniMaxApiKeyOperation, getCodexOAuthStatusOperation, getMiniMaxModelSourceOperation, setMiniMaxModelSourceOperation, testUserModelCandidateOperation, revealModelProviderApiKeyOperation, startCodexOAuthLoginOperation, cancelCodexOAuthLoginOperation, refreshModelsOperation, runCommandOperation, getSigninPanelOperation, claimSigninOperation, signOutOperation } from "./provider.js"; import { createOperationHandlers, type WebuiOperationPort } from "./operation-handlers.js"; import type { @@ -109,6 +111,17 @@ export function createOperationRegistry( registerOperation(registry, { operation: selectModelOperation, handle: handlers.selectModel }); registerOperation(registry, { operation: listSkillsOperation, handle: handlers.listSkills }); registerOperation(registry, { operation: pluginManagementOperation, handle: handlers.pluginManagement }); + // The WebUI's own scheduled tasks. Registered unconditionally for the same + // reason the other eighteen are: the port type requires the methods, so the + // registry contents must not depend on which host fills them. A host that + // cannot serve them throws inside the handler and the dispatcher reports it + // as `harness_error` with the host's own message. + registerOperation(registry, { operation: listScheduledTasksOperation, handle: handlers.listScheduledTasks }); + registerOperation(registry, { operation: createScheduledTaskOperation, handle: handlers.createScheduledTask }); + registerOperation(registry, { operation: updateScheduledTaskOperation, handle: handlers.updateScheduledTask }); + registerOperation(registry, { operation: deleteScheduledTaskOperation, handle: handlers.deleteScheduledTask }); + registerOperation(registry, { operation: triggerScheduledTaskNowOperation, handle: handlers.triggerScheduledTaskNow }); + registerOperation(registry, { operation: getScheduledTaskCapabilityOperation, handle: handlers.getScheduledTaskCapability }); registerOperation(registry, { operation: getPermissionModeOperation, handle: handlers.getPermissionMode }); registerOperation(registry, { operation: setPermissionModeOperation, handle: handlers.setPermissionMode }); registerOperation(registry, { operation: getSessionUsageOperation, handle: handlers.getSessionUsage }); diff --git a/packages/webui/src/server/operation/scheduled-task.ts b/packages/webui/src/server/operation/scheduled-task.ts new file mode 100644 index 000000000..3c439ebda --- /dev/null +++ b/packages/webui/src/server/operation/scheduled-task.ts @@ -0,0 +1,372 @@ +// Scheduled tasks (`定时任务`) — the six WebUI operations. +// +// Validation is structural only: non-empty names and prompts, a recognised +// session target and schedule kind, and a schedule that is actually +// expressible. Nothing here parses a schedule expression or reaches for a +// scheduler, because the port this surface is defined against is deliberately +// engine-free — a host may satisfy it with an interval timer, a calendar, or +// something not written yet, and this file must not decide which. +// +// Two asymmetries are deliberate and load-bearing: +// * `deleteScheduledTask` is idempotent and reports presence in `success`. +// * `updateScheduledTask` has no such field and raises for an unknown task. +// A panel retries a delete after an ambiguous disconnect; it never retries an +// edit, so an edit that silently reported "gone" would lose the user's change. + +import { WebuiErrorCode } from "../envelope.js"; +import { + invalidBody, + requireNonEmptyString, + requireRecord, +} from "./operation-contract.js"; +import type { + ValidationFailure, + WebuiOperation, + WebuiOperationValidation, +} from "./operation-contract.js"; +import type { + WebuiCreateScheduledTaskRequest, + WebuiListScheduledTasksRequest, + WebuiScheduledTask, + WebuiScheduledTaskCapability, + WebuiScheduledTaskListResult, + WebuiScheduledTaskTriggerRequest, + WebuiScheduledTaskTriggerResult, + WebuiUpdateScheduledTaskRequest, +} from "../port.js"; +import { + CREATE_SCHEDULED_TASK_OPERATION_NAME, + DELETE_SCHEDULED_TASK_OPERATION_NAME, + GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME, + LIST_SCHEDULED_TASKS_OPERATION_NAME, + TRIGGER_SCHEDULED_TASK_OPERATION_NAME, + UPDATE_SCHEDULED_TASK_OPERATION_NAME, +} from "./names.js"; + +const SESSION_TARGETS = new Set(["new", "existing"]); +const SCHEDULE_KINDS = new Set(["once", "interval"]); +const MAX_TASK_NAME_LENGTH = 120; +const MAX_PROMPT_LENGTH = 20_000; +const MIN_INTERVAL_MS = 1_000; + +/** + * Whether a helper returned a validation failure rather than a value. + * + * A `typeof value === "object"` test is wrong here: every helper in this file + * can legitimately return `null` — an absent timestamp, an absent interval — + * and `typeof null` is `"object"`, so that test would hand a `null` back to + * the dispatcher as if it were a failure object. The discriminant is the + * `ok: false` marker `invalidBody` produces. + */ +function isValidationFailure(value: unknown): value is ValidationFailure { + return ( + typeof value === "object" && + value !== null && + "ok" in value && + (value as { readonly ok: unknown }).ok === false + ); +} + +/** Tolerant timestamp check: finite, and not so large it overflows a double. */ +function isTimestamp(value: unknown): value is number { + return typeof value === "number" && Number.isFinite(value); +} + +function readOptionalTimestamp( + operation: string, + key: string, + value: unknown, +): number | null | ValidationFailure { + if (value === undefined || value === null) return null; + return isTimestamp(value) + ? value + : invalidBody(`${operation} body requires ${key} to be a timestamp or null`); +} + +function readOptionalInterval( + operation: string, + value: unknown, +): number | null | ValidationFailure { + if (value === undefined || value === null) return null; + if (typeof value !== "number" || !Number.isFinite(value) || value < MIN_INTERVAL_MS) + return invalidBody( + `${operation} body requires intervalMs to be at least ${MIN_INTERVAL_MS}ms`, + ); + return value; +} + +function requireTaskId( + operation: string, + body: Record, +): string | ValidationFailure { + return requireNonEmptyString(operation, body, "taskId"); +} + +export const listScheduledTasksOperation: WebuiOperation< + WebuiListScheduledTasksRequest, + WebuiScheduledTaskListResult +> = { + name: LIST_SCHEDULED_TASKS_OPERATION_NAME, + validate: (body) => { + if (body === undefined) return { ok: true, body: {} }; + const record = requireRecord(LIST_SCHEDULED_TASKS_OPERATION_NAME, body); + if (!record.ok) return record as WebuiOperationValidation; + const agentName = record.body.agentName; + if (agentName !== undefined && (typeof agentName !== "string" || !agentName.trim())) + return invalidBody( + `${LIST_SCHEDULED_TASKS_OPERATION_NAME} body requires agentName to be a non-empty string`, + ); + return { + ok: true, + body: agentName === undefined ? {} : { agentName: agentName as string }, + }; + }, +}; + +export const createScheduledTaskOperation: WebuiOperation< + WebuiCreateScheduledTaskRequest, + WebuiScheduledTask +> = { + name: CREATE_SCHEDULED_TASK_OPERATION_NAME, + validate: (body) => { + const record = requireRecord(CREATE_SCHEDULED_TASK_OPERATION_NAME, body); + if (!record.ok) + return record as WebuiOperationValidation; + const value = record.body; + const name = requireNonEmptyString( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "name", + ); + if (isValidationFailure(name)) return name; + if (name.length > MAX_TASK_NAME_LENGTH) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires a name of at most ${MAX_TASK_NAME_LENGTH} characters`, + ); + const agentName = requireNonEmptyString( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "agentName", + ); + if (isValidationFailure(agentName)) return agentName; + const prompt = requireNonEmptyString( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "prompt", + ); + if (isValidationFailure(prompt)) return prompt; + if (prompt.length > MAX_PROMPT_LENGTH) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires a prompt of at most ${MAX_PROMPT_LENGTH} characters`, + ); + if (typeof value.sessionTarget !== "string" || !SESSION_TARGETS.has(value.sessionTarget)) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires sessionTarget to be "new" or "existing"`, + ); + const sessionTarget = value.sessionTarget as "new" | "existing"; + if (sessionTarget === "existing") { + const sessionId = requireNonEmptyString( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "sessionId", + ); + if (isValidationFailure(sessionId)) return sessionId; + } + if (typeof value.scheduleKind !== "string" || !SCHEDULE_KINDS.has(value.scheduleKind)) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires scheduleKind to be "once" or "interval"`, + ); + const scheduleKind = value.scheduleKind as "once" | "interval"; + const runAtMs = readOptionalTimestamp( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + "runAtMs", + value.runAtMs, + ); + if (isValidationFailure(runAtMs)) return runAtMs; + const intervalMs = readOptionalInterval( + CREATE_SCHEDULED_TASK_OPERATION_NAME, + value.intervalMs, + ); + if (isValidationFailure(intervalMs)) return intervalMs; + // An expressible schedule is the whole point of the row: a recurring task + // needs an interval, a one-shot task needs a moment, and neither may be + // created with the other's missing half. + if (scheduleKind === "interval" && intervalMs === null) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires intervalMs for an interval schedule`, + ); + if (scheduleKind === "once" && runAtMs === null) + return invalidBody( + `${CREATE_SCHEDULED_TASK_OPERATION_NAME} body requires runAtMs for a one-shot schedule`, + ); + return { + ok: true, + body: { + name, + agentName, + prompt, + sessionTarget, + ...(sessionTarget === "existing" + ? { sessionId: value.sessionId as string } + : {}), + scheduleKind, + runAtMs, + intervalMs, + }, + }; + }, +}; + +export const updateScheduledTaskOperation: WebuiOperation< + WebuiUpdateScheduledTaskRequest, + WebuiScheduledTask +> = { + name: UPDATE_SCHEDULED_TASK_OPERATION_NAME, + validate: (body) => { + const record = requireRecord(UPDATE_SCHEDULED_TASK_OPERATION_NAME, body); + if (!record.ok) + return record as WebuiOperationValidation; + const value = record.body; + const taskId = requireTaskId(UPDATE_SCHEDULED_TASK_OPERATION_NAME, value); + if (isValidationFailure(taskId)) return taskId; + const patch: Record = { taskId }; + if (value.name !== undefined) { + const name = requireNonEmptyString( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "name", + ); + if (isValidationFailure(name)) return name; + if (name.length > MAX_TASK_NAME_LENGTH) + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires a name of at most ${MAX_TASK_NAME_LENGTH} characters`, + ); + patch.name = name; + } + if (value.prompt !== undefined) { + const prompt = requireNonEmptyString( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "prompt", + ); + if (isValidationFailure(prompt)) return prompt; + if (prompt.length > MAX_PROMPT_LENGTH) + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires a prompt of at most ${MAX_PROMPT_LENGTH} characters`, + ); + patch.prompt = prompt; + } + if (value.agentName !== undefined) { + const agentName = requireNonEmptyString( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + value, + "agentName", + ); + if (isValidationFailure(agentName)) return agentName; + patch.agentName = agentName; + } + if (value.sessionTarget !== undefined) { + if (typeof value.sessionTarget !== "string" || !SESSION_TARGETS.has(value.sessionTarget)) + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires sessionTarget to be "new" or "existing"`, + ); + patch.sessionTarget = value.sessionTarget; + } + if (value.sessionId !== undefined) { + if (value.sessionId !== null && (typeof value.sessionId !== "string" || !value.sessionId.trim())) + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires sessionId to be a non-empty string or null`, + ); + patch.sessionId = value.sessionId; + } + if (value.scheduleKind !== undefined) { + if (typeof value.scheduleKind !== "string" || !SCHEDULE_KINDS.has(value.scheduleKind)) + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires scheduleKind to be "once" or "interval"`, + ); + patch.scheduleKind = value.scheduleKind; + } + if (value.runAtMs !== undefined) { + const runAtMs = readOptionalTimestamp( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + "runAtMs", + value.runAtMs, + ); + if (isValidationFailure(runAtMs)) return runAtMs; + patch.runAtMs = runAtMs; + } + if (value.intervalMs !== undefined) { + const intervalMs = readOptionalInterval( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + value.intervalMs, + ); + if (isValidationFailure(intervalMs)) return intervalMs; + patch.intervalMs = intervalMs; + } + if (value.nextRunAtMs !== undefined) { + const nextRunAtMs = readOptionalTimestamp( + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + "nextRunAtMs", + value.nextRunAtMs, + ); + if (isValidationFailure(nextRunAtMs)) return nextRunAtMs; + patch.nextRunAtMs = nextRunAtMs; + } + if (value.enabled !== undefined) { + if (typeof value.enabled !== "boolean") + return invalidBody( + `${UPDATE_SCHEDULED_TASK_OPERATION_NAME} body requires enabled to be a boolean`, + ); + patch.enabled = value.enabled; + } + // An update that changes nothing is a client bug worth surfacing rather + // than a silent no-op row write. + if (Object.keys(patch).length === 1) + return invalidBody(`${UPDATE_SCHEDULED_TASK_OPERATION_NAME} requires a patch`); + return { ok: true, body: patch as unknown as WebuiUpdateScheduledTaskRequest }; + }, +}; + +export const deleteScheduledTaskOperation: WebuiOperation< + { readonly taskId: string }, + { readonly success: boolean } +> = { + name: DELETE_SCHEDULED_TASK_OPERATION_NAME, + validate: (body) => { + const record = requireRecord(DELETE_SCHEDULED_TASK_OPERATION_NAME, body); + if (!record.ok) return record as WebuiOperationValidation<{ readonly taskId: string }>; + const taskId = requireTaskId(DELETE_SCHEDULED_TASK_OPERATION_NAME, record.body); + if (isValidationFailure(taskId)) return taskId; + return { ok: true, body: { taskId } }; + }, +}; + +export const triggerScheduledTaskNowOperation: WebuiOperation< + WebuiScheduledTaskTriggerRequest, + WebuiScheduledTaskTriggerResult +> = { + name: TRIGGER_SCHEDULED_TASK_OPERATION_NAME, + validate: (body) => { + const record = requireRecord(TRIGGER_SCHEDULED_TASK_OPERATION_NAME, body); + if (!record.ok) + return record as WebuiOperationValidation; + const taskId = requireTaskId(TRIGGER_SCHEDULED_TASK_OPERATION_NAME, record.body); + if (isValidationFailure(taskId)) return taskId; + return { ok: true, body: { taskId } }; + }, +}; + +export const getScheduledTaskCapabilityOperation: WebuiOperation< + undefined, + WebuiScheduledTaskCapability +> = { + name: GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME, + validate: (body) => + body === undefined + ? { ok: true, body: undefined } + : { + ok: false, + code: WebuiErrorCode.invalidBody, + message: `${GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME} does not accept a body`, + }, +}; diff --git a/packages/webui/src/server/port.ts b/packages/webui/src/server/port.ts index 4d90fc58e..8b1f9aa8d 100644 --- a/packages/webui/src/server/port.ts +++ b/packages/webui/src/server/port.ts @@ -435,6 +435,151 @@ export type WebuiGoalWaitReason = | "verification" | "unknown"; +/** + * WebUI scheduled tasks (`定时任务`) — the semantic port. + * + * Every name here is about the *task*, not about the machinery under it: no + * cron expression, no store, no scheduler, no table. That is deliberate. The + * desktop client has its own scheduled-task surface over a different engine + * with a different concurrency model, and binding this port to either one's + * internals is what made the two impossible to keep apart. The Wire contract + * can therefore be implemented twice — once per host — without either + * implementation importing the other. + * + * The two surfaces are independent by decision: a task created here is + * invisible to the desktop client and vice versa. Reads and writes go to the + * WebUI's own store, and a run only happens while a WebUI host is running. + */ +export type WebuiScheduledTaskSessionTarget = "new" | "existing"; +export type WebuiScheduledTaskScheduleKind = "once" | "interval"; +/** + * `missed` is not a failure: nothing was attempted. It means the slot came due + * while no WebUI host was running, which is a normal outcome of an on-demand + * process and is counted rather than replayed. + */ +export type WebuiScheduledTaskRunStatus = "succeeded" | "failed" | "missed"; + +export interface WebuiScheduledTask { + readonly taskId: string; + readonly name: string; + readonly agentName: string; + /** `new` opens a fresh session per run; `existing` reuses `sessionId`. */ + readonly sessionTarget: WebuiScheduledTaskSessionTarget; + readonly sessionId: string | null; + readonly prompt: string; + readonly scheduleKind: WebuiScheduledTaskScheduleKind; + /** First slot on the grid; the anchor a recurring schedule counts from. */ + readonly runAtMs: number | null; + readonly intervalMs: number | null; + readonly enabled: boolean; + readonly lastRunAtMs: number | null; + readonly lastStatus: WebuiScheduledTaskRunStatus | null; + readonly lastError: string | null; + readonly lastResult: string | null; + /** Session the last run actually used, when it opened one. */ + readonly lastSessionId: string | null; + readonly nextRunAtMs: number | null; + /** Slots lost while no host was running. Surfaced, never silently dropped. */ + readonly missedCount: number; + /** Due time of the most recent lost slot. */ + readonly lastMissedAtMs: number | null; + readonly createdAtMs: number; + readonly updatedAtMs: number; +} + +export interface WebuiCreateScheduledTaskRequest { + readonly name: string; + readonly agentName: string; + readonly sessionTarget: WebuiScheduledTaskSessionTarget; + readonly sessionId?: string | null; + readonly prompt: string; + readonly scheduleKind: WebuiScheduledTaskScheduleKind; + readonly runAtMs?: number | null; + readonly intervalMs?: number | null; +} + +export interface WebuiUpdateScheduledTaskRequest { + readonly taskId: string; + readonly name?: string; + readonly prompt?: string; + readonly agentName?: string; + readonly sessionTarget?: WebuiScheduledTaskSessionTarget; + readonly sessionId?: string | null; + readonly scheduleKind?: WebuiScheduledTaskScheduleKind; + readonly runAtMs?: number | null; + readonly intervalMs?: number | null; + readonly enabled?: boolean; + readonly nextRunAtMs?: number | null; +} + +export interface WebuiListScheduledTasksRequest { + readonly agentName?: string; +} + +export interface WebuiScheduledTaskListResult { + readonly tasks: readonly WebuiScheduledTask[]; + readonly total: number; +} + +export interface WebuiScheduledTaskTriggerRequest { + readonly taskId: string; +} + +export interface WebuiScheduledTaskTriggerResult { + readonly taskId: string; + readonly started: boolean; + readonly status: "succeeded" | "failed" | "already_running"; + readonly error?: string; +} + +/** + * Which implementation answered a capability question. + * + * `available` alone cannot tell "the WebUI's own scheduler is in place" from + * "the adapter was swapped for the runtime's cron", and those two states carry + * opposite obligations: the first is the arrangement ADR 0012 sanctions, the + * second is the retirement it describes. The value names the implementation + * rather than the outcome, so it is meaningful when `available` is false too: + * `"webui-own"` means the WebUI's own scheduler answered (including "it could + * not open its database"), `"none"` means no scheduler is wired at all, and + * `"runtime-cron"` is the post-retirement value. + */ +export type WebuiScheduledTaskCapabilitySource = + | "webui-own" + | "runtime-cron" + | "none"; + +/** + * Whether this host can serve the surface at all, and why not when it cannot. + * Asked before the panel renders, so an unavailable backing store is a state + * the client can show rather than an error on first interaction. + */ +export interface WebuiScheduledTaskCapability { + readonly available: boolean; + readonly source: WebuiScheduledTaskCapabilitySource; + readonly reason?: string; +} + +export interface WebuiScheduledTaskPort { + listScheduledTasks( + request?: WebuiListScheduledTasksRequest, + ): Promise; + createScheduledTask( + request: WebuiCreateScheduledTaskRequest, + ): Promise; + updateScheduledTask( + request: WebuiUpdateScheduledTaskRequest, + ): Promise; + /** Idempotent: an absent task reports `success: false`, not an error. */ + deleteScheduledTask( + request: { readonly taskId: string }, + ): Promise<{ readonly success: boolean }>; + triggerScheduledTaskNow( + request: WebuiScheduledTaskTriggerRequest, + ): Promise; + getScheduledTaskCapability(): Promise; +} + export interface WebuiGoal { readonly goalId: string; readonly sessionId: string; @@ -904,7 +1049,7 @@ export type WebuiRunCommandResult = | { readonly handled: true; readonly output: string; readonly data?: unknown } | { readonly handled: true; readonly output?: undefined; readonly data: unknown }; -export interface WebuiHarnessPort { +export interface WebuiHarnessPort extends WebuiScheduledTaskPort { version(): WebuiVersionInfo; listVisibleProjects?(request: { readonly limit?: number }): Promise; listSessions(request: WebuiSessionListRequest): Promise; diff --git a/packages/webui/src/server/scheduled-task-scheduler.ts b/packages/webui/src/server/scheduled-task-scheduler.ts new file mode 100644 index 000000000..6b6e46fb0 --- /dev/null +++ b/packages/webui/src/server/scheduled-task-scheduler.ts @@ -0,0 +1,466 @@ +// The WebUI's own scheduled-task runtime: due-time accounting, delivery +// through the host's `sendMessage`, and the product boundary that defines +// what happens to a run the process was not alive for. +// +// ── The product boundary, stated once ────────────────────────────────────── +// The WebUI is an on-demand process. A scheduled task runs **only while the +// WebUI host is running**. A slot that comes due while it is down is not +// replayed when it comes back — that is the decision, not a gap in the +// implementation, and the tests exist so it cannot quietly become a catch-up +// loop. +// +// Not replaying anything at all would lose the run silently, which is a worse +// product than losing it loudly. So a missed slot is recorded: `missedCount` +// and `lastMissedAtMs` are persisted, and the panel can say "missed N times +// while the WebUI was closed" with the timestamp of the slot that was lost. +// +// ── How "missed" is decided, and why it is decidable ─────────────────────── +// The loop runs every `tickIntervalMs` for the life of the process, so a slot +// that came due while the process was alive is at most one interval late. A +// slot older than `missGraceMs` therefore cannot have been observed on time by +// a live process: it came due during a shutdown. Comparing due time against +// the wall clock — rather than carrying an "am I catching up" flag — keeps the +// rule stateless, so it holds for the very first tick after a cold start, which +// is precisely the case that matters. + +import { + ScheduledTaskStore, + openScheduledTaskDatabase, + type ScheduledTaskCreateInput, + type ScheduledTaskDatabase, + type ScheduledTaskOutcome, + type ScheduledTaskPatch, + type ScheduledTaskRecord, +} from "./scheduled-task-store.js"; +import type { + WebuiCreateScheduledTaskRequest, + WebuiScheduledTask, + WebuiScheduledTaskCapability, + WebuiScheduledTaskListResult, + WebuiScheduledTaskTriggerResult, + WebuiUpdateScheduledTaskRequest, +} from "./port.js"; + +/** The delivery seam: the host's `sendMessage`, narrowed to what a tick needs. */ +export type ScheduledTaskSendMessage = (request: { + readonly id: string; + readonly content?: string; +}) => Promise<{ + readonly ok: true; + readonly source: + | AsyncIterable + | Iterable; +} | { + readonly ok: false; + readonly status: number; + readonly body?: { readonly message?: string; readonly detail?: string }; +}>; + +export interface ScheduledTaskTickResult { + readonly firedTaskIds: readonly string[]; + readonly missedTaskIds: readonly string[]; + readonly failedTaskIds: readonly string[]; +} + +export interface WebuiScheduledTaskRuntimeOptions { + /** Pre-built store. Takes precedence over `databaseFile`; tests inject one. */ + readonly store?: ScheduledTaskStore; + /** Path to this surface's own database file. Created if absent. */ + readonly databaseFile?: string; + /** Seam for the native-addon load, so an unavailable host is reportable. */ + readonly openDatabase?: (file: string) => ScheduledTaskDatabase; + readonly sendMessage: ScheduledTaskSendMessage; + /** + * Opens a session for a `new`-target run. The port's own + * `WebuiCreateSessionResult.sessionId` is optional, so this seam returns it + * as-is and the scheduler refuses to deliver a turn without one — a missing + * session is a failed run, not a message sent to `undefined`. + */ + readonly createSession: (request: { + readonly name: string; + }) => Promise<{ readonly sessionId?: string }>; + readonly now?: () => number; + /** Loop period the service uses; also the default lateness allowance. */ + readonly tickIntervalMs?: number; + /** + * How late a due slot may be and still count as observed on time. Must stay + * comfortably above `tickIntervalMs` or a merely-late slot is misreported as + * a missed one. + */ + readonly missGraceMs?: number; +} + +const DEFAULT_TICK_INTERVAL_MS = 30_000; +const DEFAULT_MISS_GRACE_MS = 5 * 60_000; + +function toWire(task: ScheduledTaskRecord): WebuiScheduledTask { + return { + taskId: task.taskId, + name: task.name, + agentName: task.agentName, + sessionTarget: task.sessionTarget, + sessionId: task.sessionId, + prompt: task.prompt, + scheduleKind: task.scheduleKind, + runAtMs: task.runAtMs, + intervalMs: task.intervalMs, + enabled: task.enabled, + lastRunAtMs: task.lastRunAtMs, + lastStatus: task.lastStatus, + lastError: task.lastError, + lastResult: task.lastResult, + lastSessionId: task.lastSessionId, + nextRunAtMs: task.nextRunAtMs, + missedCount: task.missedCount, + lastMissedAtMs: task.lastMissedAtMs, + createdAtMs: task.createdAtMs, + updatedAtMs: task.updatedAtMs, + }; +} + +type ScheduledGrid = Pick< + ScheduledTaskRecord, + "scheduleKind" | "intervalMs" | "runAtMs" +>; + +/** The grid's period, or `null` when the task has no recurring schedule. */ +function gridPeriod(task: ScheduledGrid, fallbackMs: number): number | null { + if (task.scheduleKind === "once") return null; + const interval = task.intervalMs ?? 0; + return interval > 0 ? interval : null; +} + +/** + * The next slot strictly after `fromMs`. + * + * Strictly after, not at-or-after: a run that consumed slot S must not be + * handed S back, or one slot would be consumed twice. + */ +function nextSlotAfter(task: ScheduledGrid, fromMs: number): number | null { + const interval = gridPeriod(task, fromMs); + if (interval === null) return null; + const anchor = task.runAtMs ?? fromMs; + if (fromMs < anchor) return anchor; + return anchor + (Math.floor((fromMs - anchor) / interval) + 1) * interval; +} + +/** + * The first slot at or after `fromMs` — the one currently open. + * + * This is the question a missed run asks, and it is deliberately not + * `nextSlotAfter`: after a downtime the slot under the clock right now is the + * one the user is waiting for, so arming the slot *after* it would push an + * hourly task another hour into the future purely because the process was off. + */ +function firstSlotFrom(task: ScheduledGrid, fromMs: number): number | null { + const interval = gridPeriod(task, fromMs); + if (interval === null) return null; + const anchor = task.runAtMs ?? fromMs; + if (fromMs <= anchor) return anchor; + return anchor + Math.ceil((fromMs - anchor) / interval) * interval; +} + +export class WebuiScheduledTaskRuntime { + private readonly now: () => number; + private readonly missGraceMs: number; + private readonly inFlight = new Set(); + private store: ScheduledTaskStore | undefined; + /** Why the store could not be opened, when it could not. */ + private unavailableReason: string | undefined; + private disposed = false; + + constructor(private readonly options: WebuiScheduledTaskRuntimeOptions) { + this.now = options.now ?? (() => Date.now()); + this.missGraceMs = options.missGraceMs ?? DEFAULT_MISS_GRACE_MS; + if (options.store) { + this.store = options.store; + return; + } + const file = options.databaseFile; + if (!file) { + this.unavailableReason = + "scheduled tasks have no store: neither a store nor a database file was configured"; + return; + } + try { + const open = options.openDatabase ?? openScheduledTaskDatabase; + this.store = new ScheduledTaskStore(open(file)); + } catch (error) { + this.unavailableReason = + error instanceof Error ? error.message : String(error); + } + } + + /** + * Whether this host can serve the surface at all, and why not when it + * cannot. A client asks this before rendering the panel, so an unavailable + * native addon is a state the UI can show rather than a 500 on first click. + */ + async getScheduledTaskCapability(): Promise { + // `source` is this class's own claim, not the outcome: even when the store + // failed to open the implementation is still the WebUI's own scheduler, + // and the tripwire that watches for the ADR 0012 retirement reads this + // field to tell the two apart. + if (this.disposed) + return { + available: false, + source: "webui-own", + reason: "scheduled tasks are disposed", + }; + if (!this.store) + return { + available: false, + source: "webui-own", + reason: this.unavailableReason ?? "scheduled task store is unavailable", + }; + return { available: true, source: "webui-own" }; + } + + async listScheduledTasks( + request: { readonly agentName?: string } = {}, + ): Promise { + const store = this.requireStore(); + const tasks = store + .list() + .filter((task) => !request.agentName || task.agentName === request.agentName) + .map(toWire); + return { tasks, total: tasks.length }; + } + + async createScheduledTask( + request: WebuiCreateScheduledTaskRequest, + ): Promise { + const store = this.requireStore(); + const nowMs = this.now(); + const input: ScheduledTaskCreateInput = { + name: request.name, + agentName: request.agentName, + sessionTarget: request.sessionTarget, + sessionId: + request.sessionTarget === "existing" ? request.sessionId ?? null : null, + prompt: request.prompt, + scheduleKind: request.scheduleKind, + runAtMs: request.runAtMs ?? null, + intervalMs: + request.scheduleKind === "interval" ? request.intervalMs ?? null : null, + nowMs, + }; + return toWire(store.create(input)); + } + + async updateScheduledTask( + request: WebuiUpdateScheduledTaskRequest, + ): Promise { + const store = this.requireStore(); + const nowMs = this.now(); + const patch: ScheduledTaskPatch & { nowMs: number } = { nowMs }; + if (request.name !== undefined) patch.name = request.name; + if (request.prompt !== undefined) patch.prompt = request.prompt; + if (request.agentName !== undefined) patch.agentName = request.agentName; + if (request.sessionTarget !== undefined) patch.sessionTarget = request.sessionTarget; + if (request.sessionId !== undefined) patch.sessionId = request.sessionId; + if (request.scheduleKind !== undefined) patch.scheduleKind = request.scheduleKind; + if (request.runAtMs !== undefined) patch.runAtMs = request.runAtMs; + if (request.intervalMs !== undefined) patch.intervalMs = request.intervalMs; + if (request.enabled !== undefined) patch.enabled = request.enabled; + // Re-enabling a task with no future slot would leave it permanently + // invisible-but-armed, so the first slot is recomputed from the schedule. + if (request.enabled === true && request.nextRunAtMs === undefined) { + const current = store.get(request.taskId); + if (current && current.nextRunAtMs === null) { + patch.nextRunAtMs = + firstSlotFrom( + { + scheduleKind: request.scheduleKind ?? current.scheduleKind, + intervalMs: request.intervalMs ?? current.intervalMs, + runAtMs: request.runAtMs ?? current.runAtMs, + }, + nowMs, + ); + } + } + const updated = store.update(request.taskId, patch); + if (!updated) throw new Error(`scheduled task not found: ${request.taskId}`); + return toWire(updated); + } + + /** Idempotent: deleting an absent task reports absence, not an error. */ + async deleteScheduledTask( + request: { readonly taskId: string }, + ): Promise<{ readonly success: boolean }> { + const store = this.requireStore(); + return { success: store.remove(request.taskId) }; + } + + /** + * Runs a task now, without touching its schedule. A manual run does not + * consume the next automatic slot, so pressing the button cannot silently + * skip a scheduled run. + */ + async triggerScheduledTaskNow( + request: { readonly taskId: string }, + ): Promise { + const store = this.requireStore(); + const task = store.get(request.taskId); + if (!task) throw new Error(`scheduled task not found: ${request.taskId}`); + if (this.inFlight.has(task.taskId)) + return { taskId: task.taskId, started: false, status: "already_running" }; + const outcome = await this.execute(task, this.now(), false); + return { + taskId: task.taskId, + started: true, + status: outcome.status, + ...(outcome.error ? { error: outcome.error } : {}), + }; + } + + /** + * One pass of the schedule. Called by the service's timer, which is the only + * thing that drives it — this method has no timer of its own, so a runtime + * that is constructed but never attached to a service never fires anything. + */ + async tick(): Promise { + const store = this.store; + if (this.disposed || !store) + return { firedTaskIds: [], missedTaskIds: [], failedTaskIds: [] }; + const nowMs = this.now(); + const firedTaskIds: string[] = []; + const missedTaskIds: string[] = []; + const failedTaskIds: string[] = []; + for (const task of store.dueTasks(nowMs)) { + if (this.inFlight.has(task.taskId)) continue; + const dueAtMs = task.nextRunAtMs ?? nowMs; + if (nowMs - dueAtMs > this.missGraceMs) { + // Came due during a shutdown. Counted, not replayed — see the header. + store.recordMissed(task.taskId, { + missedAtMs: dueAtMs, + nextRunAtMs: firstSlotFrom(task, nowMs), + // A one-shot task that was missed has nothing left to do; leaving it + // armed would mean a slot that can only ever be missed again. + enabled: task.scheduleKind === "interval", + nowMs, + }); + missedTaskIds.push(task.taskId); + continue; + } + this.inFlight.add(task.taskId); + try { + const outcome = await this.execute(task, dueAtMs); + if (outcome.status === "failed") failedTaskIds.push(task.taskId); + else firedTaskIds.push(task.taskId); + } finally { + this.inFlight.delete(task.taskId); + } + } + return { firedTaskIds, missedTaskIds, failedTaskIds }; + } + + dispose(): void { + if (this.disposed) return; + this.disposed = true; + this.inFlight.clear(); + try { + this.store?.close(); + } catch { + // The database is already gone; a failure to close must not mask the + // shutdown that is already in progress. + } + this.store = undefined; + } + + private requireStore(): ScheduledTaskStore { + const store = this.store; + if (this.disposed) throw new Error("scheduled tasks are disposed"); + if (!store) + throw new Error( + this.unavailableReason ?? "scheduled task store is unavailable", + ); + return store; + } + + /** + * Delivers one run and records the result. Nothing here is allowed to throw: + * a scheduler that dies on the first failed turn would turn a transient agent + * error into a permanently dead surface, so every failure becomes a recorded + * `failed` status instead. + */ + private async execute( + task: ScheduledTaskRecord, + ranAtMs: number, + advanceSchedule = true, + ): Promise { + const recurring = task.scheduleKind === "interval"; + // A manual run reports its outcome but leaves the automatic schedule + // alone: pressing the button must not consume the next scheduled slot. + const nextRunAtMs = advanceSchedule + ? nextSlotAfter(task, ranAtMs) + : task.nextRunAtMs; + const enabled = advanceSchedule ? recurring : task.enabled; + try { + let sessionId = task.sessionId; + if (task.sessionTarget === "new" || !sessionId) { + // The agent name belongs on session creation; a session's agent is + // fixed once it exists, which is why an existing target needs none. + const created = await this.options.createSession({ + name: task.agentName, + }); + if (!created.sessionId) + throw new Error( + `the runtime created no session for agent ${task.agentName}`, + ); + sessionId = created.sessionId; + } + const result = await this.options.sendMessage({ + id: sessionId, + content: task.prompt, + }); + if (!result.ok) { + const message = result.body?.message ?? `status ${result.status}`; + return this.failed(task, ranAtMs, nextRunAtMs, enabled, message); + } + // The turn runs in the runtime and reports through the stream. Draining + // it is what turns "the prompt was accepted" into "the prompt finished"; + // an abandoned iterator would leave the outcome unknown forever. + let frames = 0; + for await (const _frame of result.source) frames += 1; + const outcome: ScheduledTaskOutcome = { + ranAtMs, + status: "succeeded", + result: `${frames} stream frame(s)`, + sessionId, + nextRunAtMs, + enabled, + }; + this.requireStore().recordOutcome(task.taskId, outcome); + return outcome; + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + return this.failed(task, ranAtMs, nextRunAtMs, enabled, message); + } + } + + private failed( + task: ScheduledTaskRecord, + ranAtMs: number, + nextRunAtMs: number | null, + enabled: boolean, + message: string, + ): ScheduledTaskOutcome { + // A recurring task keeps its schedule after a failure — a transient agent + // error must not disarm it — but the failure is visible on the row. + this.requireStore().recordOutcome(task.taskId, { + ranAtMs, + status: "failed", + error: message, + nextRunAtMs, + enabled, + }); + return { + ranAtMs, + status: "failed", + error: message, + nextRunAtMs, + enabled, + }; + } +} diff --git a/packages/webui/src/server/scheduled-task-store.ts b/packages/webui/src/server/scheduled-task-store.ts new file mode 100644 index 000000000..faeff7585 --- /dev/null +++ b/packages/webui/src/server/scheduled-task-store.ts @@ -0,0 +1,482 @@ +// WebUI-owned scheduled-task storage. +// +// ── Why this file exists at all ──────────────────────────────────────────── +// The WebUI could have borrowed the desktop client's scheduled-task tables. It +// deliberately does not, and the reasons are load-bearing rather than taste: +// +// * The v1 cron path is hard-wired off for webui +// (`local-runtime-v2/src/compat/v1/runtime.ts`), and the table it reads is +// empty in practice — the data moved to v2 during the migration. +// * `local-runtime-v2`'s `CronService` is Electron-only (`enableCron` in +// `services.ts`); a loopback WebUI process never gets that handle. +// * v2's live rows live in `local_runtime_v2_cron_definitions`. Reading them +// from here would be the worst possible combination: a second scheduler +// with no access to v2's concurrency guards, writing turns into the same +// agent queue as the desktop client. +// +// So: our own file, our own table, our own loop. The consequence is a product +// boundary, not a bug — **a scheduled task in the WebUI and a scheduled task +// in the desktop client are two different things that cannot see each other**, +// and turning the WebUI off stops its tasks from running without replaying +// them. `scheduled-task-scheduler.ts` owns that boundary's runtime half. +// +// ── Storage shape ────────────────────────────────────────────────────────── +// This is the first persistence the WebUI server owns: until now it only read +// credentials and quota out of `dataDir`. The database is a separate file +// under `/webui/`, never the runtime's own database, so "does not +// touch the v2 tables" is enforced by the filesystem rather than by +// discipline at every call site. +// +// Migration is driven by SQLite's own `user_version` pragma: it is per-file, +// transactional with the DDL it guards, and needs no bookkeeping table. Every +// step is `IF NOT EXISTS`, so re-running the migration on an already-migrated +// file is a no-op and an interrupted upgrade leaves the next start to finish +// the job. +// +// `better-sqlite3` is loaded through `createRequire` and lazily, the same way +// `terminal.ts` loads `node-pty`: it is a native addon, and a static import +// would hand esbuild a `.node` binding to bundle into the server artifact. The +// shipped standalone CLI already declares it as a runtime dependency +// (`release/webui-npm/package.json`) and `scripts/package-webui-npm.mjs` +// externalises it; in a source checkout it resolves through the workspace's +// hoisted `node_modules`. A host where it cannot load is not broken — the +// capability probe reports it and the surface fails closed with the reason. + +import { mkdirSync } from "node:fs"; +import path from "node:path"; +import { createRequire } from "node:module"; +import { randomUUID } from "node:crypto"; + +const require = createRequire(import.meta.url); + +export const SCHEDULED_TASK_SCHEMA_VERSION = 1; + +/** The table is prefixed because the file is ours alone, and names outlive files. */ +const TABLE = "webui_scheduled_task"; + +const CREATE_TABLE_SQL = ` +CREATE TABLE IF NOT EXISTS ${TABLE} ( + task_id TEXT PRIMARY KEY, + name TEXT NOT NULL, + agent_name TEXT NOT NULL, + session_target TEXT NOT NULL CHECK (session_target IN ('new', 'existing')), + session_id TEXT, + prompt TEXT NOT NULL, + schedule_kind TEXT NOT NULL CHECK (schedule_kind IN ('once', 'interval')), + run_at_ms INTEGER, + interval_ms INTEGER, + enabled INTEGER NOT NULL DEFAULT 1 CHECK (enabled IN (0, 1)), + last_run_at_ms INTEGER, + last_status TEXT CHECK (last_status IS NULL OR last_status IN ('succeeded', 'failed', 'missed')), + last_error TEXT, + last_result TEXT, + last_session_id TEXT, + next_run_at_ms INTEGER, + missed_count INTEGER NOT NULL DEFAULT 0, + last_missed_at_ms INTEGER, + created_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL +); + +CREATE INDEX IF NOT EXISTS ${TABLE}_due + ON ${TABLE} (enabled, next_run_at_ms); +`; + +/** The row shape, with SQLite's integers already narrowed back to booleans. */ +export interface ScheduledTaskRecord { + readonly taskId: string; + readonly name: string; + readonly agentName: string; + readonly sessionTarget: "new" | "existing"; + readonly sessionId: string | null; + readonly prompt: string; + readonly scheduleKind: "once" | "interval"; + readonly runAtMs: number | null; + readonly intervalMs: number | null; + readonly enabled: boolean; + readonly lastRunAtMs: number | null; + readonly lastStatus: "succeeded" | "failed" | "missed" | null; + readonly lastError: string | null; + readonly lastResult: string | null; + readonly lastSessionId: string | null; + readonly nextRunAtMs: number | null; + readonly missedCount: number; + readonly lastMissedAtMs: number | null; + readonly createdAtMs: number; + readonly updatedAtMs: number; +} + +export interface ScheduledTaskCreateInput { + readonly name: string; + readonly agentName: string; + readonly sessionTarget: "new" | "existing"; + readonly sessionId?: string | null; + readonly prompt: string; + readonly scheduleKind: "once" | "interval"; + /** First slot on the grid. Defaults to "now" for a one-shot task. */ + readonly runAtMs?: number | null; + /** Recurring period. Ignored for a one-shot task. */ + readonly intervalMs?: number | null; + readonly nowMs: number; +} + +/** + * The fields a caller may change. Mutable because the runtime assembles a + * patch key by key; the record it is derived from is readonly on purpose, + * since a row handed out for reading must not be writable in place. + */ +export type ScheduledTaskPatch = Partial< + { + -readonly [Key in keyof Pick< + ScheduledTaskRecord, + | "name" + | "prompt" + | "agentName" + | "sessionTarget" + | "sessionId" + | "scheduleKind" + | "runAtMs" + | "intervalMs" + | "enabled" + | "nextRunAtMs" + >]: Pick< + ScheduledTaskRecord, + Key + >[Key]; + } +>; + +export interface ScheduledTaskOutcome { + readonly ranAtMs: number; + readonly status: "succeeded" | "failed"; + readonly error?: string | undefined; + readonly result?: string | undefined; + readonly sessionId?: string | undefined; + readonly nextRunAtMs: number | null; + readonly enabled: boolean; +} + +/** The subset of `better-sqlite3` this module uses, declared structurally. */ +export interface ScheduledTaskDatabase { + prepare(sql: string): { + run(...params: readonly unknown[]): unknown; + all(...params: readonly unknown[]): readonly Record[]; + get(...params: readonly unknown[]): Record | undefined; + }; + exec(sql: string): unknown; + pragma(source: string): unknown; + close(): void; +} + +interface SqliteStatement { + run(...params: readonly unknown[]): unknown; + all(...params: readonly unknown[]): readonly Record[]; + get(...params: readonly unknown[]): Record | undefined; +} + +interface SqliteDatabase { + prepare(sql: string): SqliteStatement; + exec(sql: string): unknown; + pragma(source: string): { readonly user_version: number } | readonly number[]; + close(): void; +} + +type SqliteConstructor = new ( + file: string, + options?: Record, +) => SqliteDatabase; + +let cachedConstructor: SqliteConstructor | undefined; + +/** + * Resolves the native addon on demand. Kept separate from `openDatabase` so a + * host that never touches scheduled tasks never pays for the load, and so the + * failure is a value the capability probe can report rather than a throw that + * takes the process down at import time. + */ +function loadSqliteConstructor(): SqliteConstructor { + if (cachedConstructor) return cachedConstructor; + const loaded = require("better-sqlite3") as SqliteConstructor; + cachedConstructor = loaded; + return loaded; +} + +/** + * Opens (creating if needed) the WebUI's scheduled-task database and brings it + * to {@link SCHEDULED_TASK_SCHEMA_VERSION}. + * + * WAL keeps a panel read from blocking a tick's write, and `foreign_keys` is + * left off deliberately: this file has exactly one table and no references. + */ +export function openScheduledTaskDatabase(file: string): ScheduledTaskDatabase { + if (file !== ":memory:") mkdirSync(path.dirname(file), { recursive: true }); + const Database = loadSqliteConstructor(); + const database = new Database(file); + database.pragma("journal_mode = WAL"); + const store = new ScheduledTaskStore(database as unknown as ScheduledTaskDatabase); + store.migrate(); + return database as unknown as ScheduledTaskDatabase; +} + +function readUserVersion(database: ScheduledTaskDatabase): number { + // better-sqlite3 answers a named pragma with an array of rows — pragma + // values are query results like any other — so `user_version` arrives as + // `[{ user_version: n }]`. Both that and a bare scalar are accepted so the + // store does not depend on which spelling the driver picked. + const result = database.pragma("user_version") as unknown; + if (Array.isArray(result)) { + const row = result[0]; + if (row && typeof row === "object" && "user_version" in row) + return Number((row as { readonly user_version: unknown }).user_version ?? 0); + return Number(row ?? 0); + } + if (result && typeof result === "object" && "user_version" in result) + return Number((result as { readonly user_version: unknown }).user_version ?? 0); + return Number(result ?? 0); +} + +function toText(value: unknown): string { + return typeof value === "string" ? value : String(value ?? ""); +} + +function toNullableText(value: unknown): string | null { + return typeof value === "string" && value !== "" ? value : null; +} + +function toNumber(value: unknown): number { + return typeof value === "number" ? value : Number(value ?? 0); +} + +function toNullableNumber(value: unknown): number | null { + return value === null || value === undefined ? null : toNumber(value); +} + +function toRecord(row: Record): ScheduledTaskRecord { + return { + taskId: toText(row.task_id), + name: toText(row.name), + agentName: toText(row.agent_name), + sessionTarget: toText(row.session_target) === "new" ? "new" : "existing", + sessionId: toNullableText(row.session_id), + prompt: toText(row.prompt), + scheduleKind: toText(row.schedule_kind) === "once" ? "once" : "interval", + runAtMs: toNullableNumber(row.run_at_ms), + intervalMs: toNullableNumber(row.interval_ms), + enabled: toNumber(row.enabled) === 1, + lastRunAtMs: toNullableNumber(row.last_run_at_ms), + lastStatus: + row.last_status === "succeeded" || + row.last_status === "failed" || + row.last_status === "missed" + ? row.last_status + : null, + lastError: toNullableText(row.last_error), + lastResult: toNullableText(row.last_result), + lastSessionId: toNullableText(row.last_session_id), + nextRunAtMs: toNullableNumber(row.next_run_at_ms), + missedCount: toNumber(row.missed_count), + lastMissedAtMs: toNullableNumber(row.last_missed_at_ms), + createdAtMs: toNumber(row.created_at_ms), + updatedAtMs: toNumber(row.updated_at_ms), + }; +} + +const COLUMNS = + "task_id, name, agent_name, session_target, session_id, prompt, " + + "schedule_kind, run_at_ms, interval_ms, enabled, last_run_at_ms, last_status, " + + "last_error, last_result, last_session_id, next_run_at_ms, missed_count, " + + "last_missed_at_ms, created_at_ms, updated_at_ms"; + +/** + * All reads and writes for the WebUI's scheduled tasks. + * + * Every mutation is a single statement so the store never needs a + * multi-statement transaction to stay consistent — a tick and a panel edit can + * interleave, and each of them leaves the row valid on its own. + */ +export class ScheduledTaskStore { + constructor(private readonly database: ScheduledTaskDatabase) {} + + /** + * Brings the file to the current schema version. + * + * Steps run in order and each is individually idempotent, so an interrupted + * upgrade resumes on the next start instead of needing a repair path. The + * version is written after its DDL, which SQLite commits together with it. + */ + migrate(): void { + const current = readUserVersion(this.database); + if (current >= SCHEDULED_TASK_SCHEMA_VERSION) return; + this.database.exec(CREATE_TABLE_SQL); + this.database.pragma(`user_version = ${SCHEDULED_TASK_SCHEMA_VERSION}`); + } + + schemaVersion(): number { + return readUserVersion(this.database); + } + + list(): ScheduledTaskRecord[] { + return this.database + .prepare(`SELECT ${COLUMNS} FROM ${TABLE} ORDER BY created_at_ms ASC, task_id ASC`) + .all() + .map(toRecord); + } + + get(taskId: string): ScheduledTaskRecord | undefined { + const row = this.database + .prepare(`SELECT ${COLUMNS} FROM ${TABLE} WHERE task_id = ?`) + .get(taskId); + return row ? toRecord(row) : undefined; + } + + /** + * Enabled tasks whose due time has arrived. The scheduler decides what to do + * with each one; this method only reports candidates, so the miss accounting + * stays in one place. + */ + dueTasks(nowMs: number): ScheduledTaskRecord[] { + return this.database + .prepare( + `SELECT ${COLUMNS} FROM ${TABLE} + WHERE enabled = 1 AND next_run_at_ms IS NOT NULL AND next_run_at_ms <= ? + ORDER BY next_run_at_ms ASC, task_id ASC`, + ) + .all(nowMs) + .map(toRecord); + } + + /** + * Inserts a task and returns the stored row. The identity is minted here + * rather than by the caller: it is a storage detail, and a caller that could + * supply it could also collide with an existing row. + */ + create(input: ScheduledTaskCreateInput): ScheduledTaskRecord { + const taskId = `sched-${randomUUID()}`; + // A one-shot task with no explicit moment runs at the moment it was made; + // a recurring task always carries its interval, and one with no explicit + // anchor starts counting from now. + const runAtMs = + input.runAtMs ?? (input.scheduleKind === "once" ? input.nowMs : null); + this.database + .prepare( + `INSERT INTO ${TABLE} ( + task_id, name, agent_name, session_target, session_id, prompt, + schedule_kind, run_at_ms, interval_ms, enabled, next_run_at_ms, + created_at_ms, updated_at_ms + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)`, + ) + .run( + taskId, + input.name, + input.agentName, + input.sessionTarget, + input.sessionId ?? null, + input.prompt, + input.scheduleKind, + runAtMs, + input.intervalMs ?? null, + runAtMs, + input.nowMs, + input.nowMs, + ); + const created = this.get(taskId); + if (!created) throw new Error("scheduled task row vanished after insert"); + return created; + } + + /** Returns the updated row, or `undefined` when the task does not exist. */ + update(taskId: string, patch: ScheduledTaskPatch & { readonly nowMs: number }): ScheduledTaskRecord | undefined { + const assignments: string[] = []; + const values: unknown[] = []; + const columns: Readonly> = { + name: "name", + prompt: "prompt", + agentName: "agent_name", + sessionTarget: "session_target", + sessionId: "session_id", + scheduleKind: "schedule_kind", + runAtMs: "run_at_ms", + intervalMs: "interval_ms", + enabled: "enabled", + nextRunAtMs: "next_run_at_ms", + }; + for (const [key, column] of Object.entries(columns)) { + if (!(key in patch)) continue; + assignments.push(`${column} = ?`); + const value = (patch as Record)[key]; + values.push(key === "enabled" ? (value ? 1 : 0) : value); + } + if (!assignments.length) return this.get(taskId); + assignments.push("updated_at_ms = ?"); + values.push(patch.nowMs, taskId); + this.database + .prepare(`UPDATE ${TABLE} SET ${assignments.join(", ")} WHERE task_id = ?`) + .run(...values); + return this.get(taskId); + } + + /** Idempotent: reports whether a row was actually removed. */ + remove(taskId: string): boolean { + const result = this.database + .prepare(`DELETE FROM ${TABLE} WHERE task_id = ?`) + .run(taskId) as { readonly changes?: number }; + return Number(result?.changes ?? 0) > 0; + } + + recordOutcome(taskId: string, outcome: ScheduledTaskOutcome): void { + this.database + .prepare( + `UPDATE ${TABLE} SET + last_run_at_ms = ?, last_status = ?, last_error = ?, last_result = ?, + last_session_id = COALESCE(?, last_session_id), next_run_at_ms = ?, + enabled = ?, updated_at_ms = ? + WHERE task_id = ?`, + ) + .run( + outcome.ranAtMs, + outcome.status, + outcome.error ?? null, + outcome.result ?? null, + outcome.sessionId ?? null, + outcome.nextRunAtMs, + outcome.enabled ? 1 : 0, + outcome.ranAtMs, + taskId, + ); + } + + /** + * Records a run that came due while the WebUI was not running. `missedAtMs` + * is the due time, not the moment the miss was noticed, so the panel can say + * which slot was lost. + */ + recordMissed( + taskId: string, + missed: { + readonly missedAtMs: number; + readonly nextRunAtMs: number | null; + readonly enabled: boolean; + readonly nowMs: number; + }, + ): void { + this.database + .prepare( + `UPDATE ${TABLE} SET + missed_count = missed_count + 1, last_missed_at_ms = ?, + last_status = 'missed', last_error = NULL, next_run_at_ms = ?, + enabled = ?, updated_at_ms = ? + WHERE task_id = ?`, + ) + .run( + missed.missedAtMs, + missed.nextRunAtMs, + missed.enabled ? 1 : 0, + missed.nowMs, + taskId, + ); + } + + close(): void { + this.database.close(); + } +} diff --git a/packages/webui/src/server/service.ts b/packages/webui/src/server/service.ts index 12652fa95..6af74028b 100644 --- a/packages/webui/src/server/service.ts +++ b/packages/webui/src/server/service.ts @@ -25,6 +25,7 @@ import { createOperationRegistry, type WebuiOperationRegistryEntry, } from "./operation/operations.js"; +import { WebuiScheduledTaskRuntime } from "./scheduled-task-scheduler.js"; import { createWebuiCredential, credentialMatches, @@ -46,6 +47,15 @@ import { webuiSessionTransferFileName } from "./session-transfer.js"; export const WEBUI_MAX_MESSAGE_BYTES = 256 * 1024; export const WEBUI_WEBSOCKET_HEARTBEAT_INTERVAL_MS = 15_000; +/** + * How often the scheduled-task loop looks for due work. Deliberately coarse: + * a task is a prompt sent to an agent, and the loop's real job is to notice + * due slots — not to hit them to the millisecond. The runtime treats anything + * later than its grace window as "the process was not running", so this value + * also sets how much lateness is tolerated before a slot counts as missed. + */ +export const WEBUI_SCHEDULED_TASK_TICK_MS = 30_000; + /** * Truthy env-var spellings that turn the development mode on. Anything that * is not on this list is ignored, so `WEBUI_DEV=0` / `WEBUI_DEV=false` / @@ -80,8 +90,32 @@ export interface WebuiServiceOptions { readonly maxMessageBytes?: number; /** Override the liveness sweep interval; tests use a short interval. */ readonly webSocketHeartbeatIntervalMs?: number; + /** + * Override the scheduled-task tick interval. Independent of + * `scheduledTasks` on purpose: a caller that wants a faster loop but the + * service's own runtime (or no runtime at all) should not have to construct + * one just to change a period. Defaults to 30s. + */ + readonly scheduledTaskTickIntervalMs?: number; /** Optional credential override; tests supply one to assert its shape. */ readonly credential?: WebuiCredential; + /** + * The WebUI's own scheduled-task runtime. The service is what drives it: + * a tick timer starts beside the heartbeat in `start()` and is cleared in + * `close()`, so the schedule lives exactly as long as the process does. + * + * The service also answers the six scheduled-task port methods from this + * runtime rather than from the harness port. That is the point of the + * surface — it is WebUI-owned state, not harness state — and it means a + * host that ships no scheduled-task runtime still serves the panel. + * Omit the option and the six methods fall through to the harness port, + * which fails closed with its own reason. + */ + readonly scheduledTasks?: { + readonly runtime: WebuiScheduledTaskRuntime; + /** Tick period. Defaults to 30s; the runtime is told so it can size its grace window. */ + readonly tickIntervalMs?: number; + }; /** Optional server factory; tests inject an HTTP server without listening. */ readonly httpServerFactory?: () => Server; /** @@ -136,6 +170,9 @@ export class WebuiService { private readonly connectionSignals = new Map(); private readonly connectionAlive = new WeakMap(); private heartbeatTimer: ReturnType | undefined; + private scheduledTaskTimer: ReturnType | undefined; + private readonly scheduledTasks: WebuiScheduledTaskRuntime | undefined; + private readonly scheduledTaskIntervalMs: number; private accepting = true; private startedPromise: Promise | undefined; private bound: { info: WebuiServiceInfo } | undefined; @@ -152,6 +189,12 @@ export class WebuiService { this.webSocketHeartbeatIntervalMs = options.webSocketHeartbeatIntervalMs ?? WEBUI_WEBSOCKET_HEARTBEAT_INTERVAL_MS; this.credential = options.credential ?? createWebuiCredential(); + this.scheduledTaskIntervalMs = + options.scheduledTaskTickIntervalMs ?? + options.scheduledTasks?.tickIntervalMs ?? + WEBUI_SCHEDULED_TASK_TICK_MS; + this.scheduledTasks = + options.scheduledTasks?.runtime ?? this.buildScheduledTaskRuntime(); this.protocolVersion = options.protocolVersion ?? WEBUI_PROTOCOL_VERSION; // The option is the source of truth for tests; the env var is the // convenience for the dev preview launcher. The option must win @@ -191,6 +234,33 @@ export class WebuiService { createGoal: (request) => this.port.createGoal(request), patchGoal: (request) => this.port.patchGoal(request), clearGoal: (request) => this.port.clearGoal(request), + // Scheduled tasks come from the service's own runtime when it has one, + // and from the harness port otherwise. Both halves are the same port + // contract, so the registry and the wire are identical either way. + listScheduledTasks: (request) => + this.scheduledTasks + ? this.scheduledTasks.listScheduledTasks(request) + : this.port.listScheduledTasks(request), + createScheduledTask: (request) => + this.scheduledTasks + ? this.scheduledTasks.createScheduledTask(request) + : this.port.createScheduledTask(request), + updateScheduledTask: (request) => + this.scheduledTasks + ? this.scheduledTasks.updateScheduledTask(request) + : this.port.updateScheduledTask(request), + deleteScheduledTask: (request) => + this.scheduledTasks + ? this.scheduledTasks.deleteScheduledTask(request) + : this.port.deleteScheduledTask(request), + triggerScheduledTaskNow: (request) => + this.scheduledTasks + ? this.scheduledTasks.triggerScheduledTaskNow(request) + : this.port.triggerScheduledTaskNow(request), + getScheduledTaskCapability: () => + this.scheduledTasks + ? this.scheduledTasks.getScheduledTaskCapability() + : this.port.getScheduledTaskCapability(), listWorkspaceFileTree: (request) => this.port.listWorkspaceFileTree(request), readWorkspaceFile: (request) => this.port.readWorkspaceFile(request), getWorkspaceEnvironment: (request) => this.port.getWorkspaceEnvironment(request), @@ -631,6 +701,7 @@ export class WebuiService { boundUrl: `ws://${this.host}:${tcpPort}`, }; this.bound = { info }; + this.startScheduledTaskLoop(); resolve(info); } catch (error) { reject(error); @@ -648,6 +719,77 @@ export class WebuiService { return this.bound.info; } + /** + * Build the WebUI's own scheduled-task runtime from the `dataDir` the port + * already reports, rather than adding another injection point for it. + * + * `version()` is where a host states its data directory, and the service + * already reads it for the version frame — so the store lands in + * `/webui/`, beside the credentials the assembly reads from the + * same place, with no change to the assembly or to either launcher. A host + * that reports no `dataDir` (every test double) gets no runtime, and the + * capability probe says so instead of the surface failing at first click. + * + * Delivery goes through the port's own `sendMessage` and `createSession`, + * which is the same path the panel's own send button takes. The scheduler + * therefore has no privileged route into the harness that a user-initiated + * turn does not also have. + */ + private buildScheduledTaskRuntime(): WebuiScheduledTaskRuntime | undefined { + const dataDir = this.port.version().dataDir; + // Only an absolute path *in this platform's terms*. A host can report a + // data directory for another platform — `session-transfer-route.test.ts` + // reports `C:/data` — and on POSIX `path.join` would treat that as a + // relative path and quietly create the tree inside the working directory. + // A directory that is not ours to write is a host we do not persist for. + if (!dataDir || !path.isAbsolute(dataDir)) return undefined; + return new WebuiScheduledTaskRuntime({ + databaseFile: path.join(dataDir, "webui", "scheduled-tasks.sqlite"), + tickIntervalMs: this.scheduledTaskIntervalMs, + createSession: (request) => this.port.createSession({ name: request.name }), + sendMessage: async (request) => { + const result = await this.port.sendMessage(request); + // The runtime drains the stream to learn the outcome; the WebUI's own + // frame projection is a rendering concern and has no place here. + return result.ok + ? { + ok: true as const, + source: result.source as + | AsyncIterable + | Iterable, + } + : { + ok: false as const, + status: result.status, + body: result.body, + }; + }, + }); + } + + /** + * Start the scheduled-task loop, one instance per process, beside the + * WebSocket heartbeat and on the same terms. + * + * `unref` matters here for the same reason it does on the heartbeat: the + * timer must never be the reason a process stays alive. A run in progress is + * allowed to finish — the tick's promise is not awaited by the timer — and + * the store is released by the runtime's own `dispose`. + * + * A tick that throws is swallowed on purpose. The loop is a background + * courtesy to the user, and a transient store error must not become an + * unhandled rejection that takes the service down; the next tick retries, and + * a task that cannot be delivered records its own failure. + */ + private startScheduledTaskLoop(): void { + const runtime = this.scheduledTasks; + if (!runtime || this.scheduledTaskTimer !== undefined) return; + this.scheduledTaskTimer = setInterval(() => { + void runtime.tick().catch(() => undefined); + }, this.scheduledTaskIntervalMs); + this.scheduledTaskTimer.unref?.(); + } + /** * Stop accepting new operations, drain every connection, then close * the harness port. The order is the one step 13 of the assembly @@ -662,6 +804,13 @@ export class WebuiService { clearInterval(this.heartbeatTimer); this.heartbeatTimer = undefined; } + if (this.scheduledTaskTimer !== undefined) { + clearInterval(this.scheduledTaskTimer); + this.scheduledTaskTimer = undefined; + } + // The store is this process's own file handle; leaving it open would keep + // the WAL alive after the last operation is gone. + this.scheduledTasks?.dispose(); // Force-terminate every connection before the server closes; otherwise // `wsServer.close()` waits for the client to ack the close handshake // and can hang for the duration of the platform TCP timeout. diff --git a/packages/webui/test/unit/webui-host-shape-invariant.test.ts b/packages/webui/test/unit/webui-host-shape-invariant.test.ts index 18ad55e5d..1853728d6 100644 --- a/packages/webui/test/unit/webui-host-shape-invariant.test.ts +++ b/packages/webui/test/unit/webui-host-shape-invariant.test.ts @@ -67,6 +67,30 @@ import { createHarnessPortFromHost } from "../../src/server/host.js"; * the matching `WebuiHarnessPort` field. Forgetting a member fails the * compile instead of silently returning `undefined`. */ +/** Structural zero for the scheduled-task surface, per this file's convention. */ +const INVARIANT_SCHEDULED_TASK = { + taskId: "invariant", + name: "invariant", + agentName: "invariant", + sessionTarget: "existing" as const, + sessionId: "invariant", + prompt: "invariant", + scheduleKind: "once" as const, + runAtMs: 0, + intervalMs: null, + enabled: false, + lastRunAtMs: null, + lastStatus: null, + lastError: null, + lastResult: null, + lastSessionId: null, + nextRunAtMs: null, + missedCount: 0, + lastMissedAtMs: null, + createdAtMs: 0, + updatedAtMs: 0, +}; + class FullPort implements WebuiHarnessPort { version() { return { version: "invariant-test", protocolVersion: 1 }; @@ -384,6 +408,27 @@ class FullPort implements WebuiHarnessPort { // `/compact` slash command reaches it through `runWebuiCommand`. return { success: true as const }; } + // Added with the WebUI scheduled-task surface. The invariant port exists to + // turn a forgotten port member into a compile error, so these land here + // rather than being left to an optional hook. + async listScheduledTasks() { + return { tasks: [{ ...INVARIANT_SCHEDULED_TASK }], total: 1 }; + } + async createScheduledTask() { + return { ...INVARIANT_SCHEDULED_TASK }; + } + async updateScheduledTask() { + return { ...INVARIANT_SCHEDULED_TASK }; + } + async deleteScheduledTask() { + return { success: false }; + } + async triggerScheduledTaskNow() { + return { taskId: "invariant", started: false, status: "already_running" as const }; + } + async getScheduledTaskCapability() { + return { available: false as const, source: "none" as const, reason: "invariant port" }; + } async close() { // no-op } diff --git a/packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts b/packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts new file mode 100644 index 000000000..21ae722fc --- /dev/null +++ b/packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts @@ -0,0 +1,624 @@ +// Retirement-condition probe for the WebUI's own scheduled tasks (ADR 0012). +// +// ── What this file is, and what it is not ────────────────────────────────── +// ADR 0012 ("The WebUI runs its own scheduled tasks until the runtime offers +// one") owns the retirement decision and the work that follows it. **This file +// only reports.** It does not gate anything: no state — not a closed upstream, +// not an open upstream, not a probe that can no longer read the upstream — is +// allowed to turn `test:webui` red, and this file must never fail because of an +// upstream change. A door that goes red for something that is not a defect only +// teaches people to ignore red. +// +// So the two retirement conditions in ADR 0012 are read as source text and +// printed. What the reader gets: +// +// gate1=closed gate2=closed one quiet status line, no warning +// gate1=open gate2=closed a warning: upstream moved, still not usable +// gate1=closed gate2=open a warning: upstream moved, still not usable +// gate1=open gate2=open a warning: retirement condition is met +// probe=unconfirmed a warning: this probe cannot read the +// upstream, and that is NOT a retirement +// signal and NOT a reason to delete +// anything +// +// "Unknown" is never collapsed into "closed". An unreadable probe is reported +// as unreadable, and its message says in words that it is not a retirement +// signal, so nobody reads the yellow line and tears down the WebUI's +// scheduler on the strength of it. +// +// ── Why source text and not a real host ──────────────────────────────────── +// Booting a real `createLocalRuntimeHostV2` takes 30–100s, opens sqlite and +// starts schedulers, which is not a unit test. The two conditions are, however, +// fully decided by expressions in the upstream package's own source, and this +// repository already reads other packages' sources from its checks +// (`scripts/source-inventory.mjs`, `webui-boundary-check.test.ts`). So the +// probe reads four files under `packages/local-runtime-v2` — read-only — and +// answers from what it found. +// +// The two gates are judged independently and only then combined, so a report +// can name which door opened. They can open separately, and gate 2 before +// gate 1 is entirely plausible: upstream can expose the slot on the host +// contract long before it builds the service for an embedded host. +// +// ── The gates, as ADR 0012 states them ───────────────────────────────────── +// gate 1 (creation) `services.cron` is created for a `tui` + +// `cliEmbedded` host — today `enableCron` is +// `ownsElectronRuntimeCapabilities(runtimeOwnerKind)`, +// and that predicate accepts only `undefined` and +// `'electron'`. Evidence: +// `local-runtime-v2/src/services.ts` (the assignment +// and the creation site that consumes it) and +// `application/agent/runtime-browser-use-composition.ts` +// (the predicate body). +// gate 2 (reachability) the created service is reachable from the host — +// carried on `CreatedLocalRuntimeHost` or on the +// `cliService` options, which is where +// `runtime.ts` spreads the whole services object. +// Evidence: `local/host-contract.ts` and +// `local/cli-service.ts`. Note that +// `RuntimeServices` already declares an optional `cron`; +// that is the value being spread, and it is `undefined` +// for this host because gate 1 is closed, so it is not +// itself evidence of an open gate 2. +// +// The only assertions in this file are the probe's own behaviour against +// literal fixture sources, which no upstream change can influence. + +import { describe, expect, it } from "vitest"; +import { mkdtempSync, readFileSync, rmSync } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; + +import { + openScheduledTaskDatabase, + ScheduledTaskStore, +} from "../../src/server/scheduled-task-store.js"; +import { WebuiScheduledTaskRuntime } from "../../src/server/scheduled-task-scheduler.js"; + +const here = path.dirname(fileURLToPath(import.meta.url)); +const repoRoot = path.resolve(here, "../../../.."); + +const ADR_0012 = + "docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md"; + +/** The upstream source text the two gates are read from. */ +interface UpstreamSources { + readonly services?: string; + readonly composition?: string; + readonly hostContract?: string; + readonly cliService?: string; +} + +const UPSTREAM_FILES = { + services: "packages/local-runtime-v2/src/services.ts", + composition: + "packages/local-runtime-v2/src/application/agent/runtime-browser-use-composition.ts", + hostContract: "packages/local-runtime-v2/src/local/host-contract.ts", + cliService: "packages/local-runtime-v2/src/local/cli-service.ts", +} as const; + +function readUpstreamSources(rootDir: string): UpstreamSources { + const read = (relative: string): string | undefined => { + try { + return readFileSync(path.join(rootDir, relative), "utf8"); + } catch { + // A moved or deleted file is a state to report, not a crash. + return undefined; + } + }; + return { + services: read(UPSTREAM_FILES.services), + composition: read(UPSTREAM_FILES.composition), + hostContract: read(UPSTREAM_FILES.hostContract), + cliService: read(UPSTREAM_FILES.cliService), + }; +} + +/** One gate's answer. `unconfirmed` is never read as `closed`. */ +interface Gate { + readonly state: "closed" | "open" | "unconfirmed"; + readonly detail: string; +} + +const unconfirmed = (detail: string): Gate => ({ state: "unconfirmed", detail }); +const closed = (detail: string): Gate => ({ state: "closed", detail }); +const opened = (detail: string): Gate => ({ state: "open", detail }); + +/** The `{ ... }` body of a declaration, or undefined if it cannot be found. */ +function braceBody( + source: string, + declaration: string, +): string | undefined { + const declarationStart = source.indexOf(declaration); + if (declarationStart < 0) return undefined; + const open = source.indexOf("{", declarationStart); + if (open < 0) return undefined; + let depth = 0; + for (let index = open; index < source.length; index += 1) { + const character = source[index]; + if (character === "{") depth += 1; + else if (character === "}") { + depth -= 1; + if (depth === 0) return source.slice(open + 1, index); + } + } + return undefined; +} + +function memberNames(body: string): string[] { + return [ + ...body.matchAll( + /^[ \t]*(?:readonly[ \t]+)?([A-Za-z_$][A-Za-z0-9_$]*)[ \t]*\??[ \t]*[:?]/gmu, + ), + ].map((match) => match[1] ?? ""); +} + +/** + * Gate 1: does a `tui` + `cliEmbedded` host get a `services.cron`? + * + * Read from the `enableCron` assignment and from the creation site that + * consumes it. Both halves matter: an `enableCron` that nothing consumes would + * mean the service is built unconditionally, and reading only the assignment + * would call that closed. + */ +function probeGate1(sources: UpstreamSources): Gate { + const { services, composition } = sources; + if (services === undefined) + return unconfirmed(`could not read ${UPSTREAM_FILES.services}`); + if (composition === undefined) + return unconfirmed(`could not read ${UPSTREAM_FILES.composition}`); + + const assignments = [ + ...services.matchAll(/^[ \t]*enableCron[ \t]*:[ \t]*(.*)$/gmu), + ].map((match) => (match[1] ?? "").trim().replace(/,$/u, "")); + if (assignments.length === 0) + return unconfirmed( + `no \`enableCron:\` assignment found in ${UPSTREAM_FILES.services}; this probe does not know where cron is gated now`, + ); + if (assignments.length > 1) + return unconfirmed( + `${assignments.length} \`enableCron:\` assignments found in ${UPSTREAM_FILES.services}; this probe expects exactly one`, + ); + const value = assignments[0] ?? ""; + + if (value === "true") + return opened( + `\`enableCron: true\` in ${UPSTREAM_FILES.services} — the flag is no longer gated on the host kind, so a tui + cliEmbedded host gets a cron service`, + ); + if (value === "false") + return closed( + `\`enableCron: false\` in ${UPSTREAM_FILES.services} — cron is pinned off for every host kind`, + ); + if (!value.includes("ownsElectronRuntimeCapabilities(")) + return unconfirmed( + `\`enableCron: ${value}\` in ${UPSTREAM_FILES.services} is not the electron-only predicate this probe reads, so the gate can no longer be classified`, + ); + + const predicate = braceBody( + composition, + "function ownsElectronRuntimeCapabilities", + ); + if (predicate === undefined) + return unconfirmed( + `ownsElectronRuntimeCapabilities has no readable body in ${UPSTREAM_FILES.composition}`, + ); + // Quoted literals only: an incidental mention of a host kind in a comment + // must not be able to move a gate. + if (/(['"])tui\1/u.test(predicate)) + return opened( + `ownsElectronRuntimeCapabilities() accepts 'tui' (${UPSTREAM_FILES.composition}), so a tui + cliEmbedded host is created a cron service`, + ); + if (!/(['"])electron\1/u.test(predicate)) + return unconfirmed( + `ownsElectronRuntimeCapabilities() mentions neither 'tui' nor 'electron' (${UPSTREAM_FILES.composition}); this probe cannot classify it`, + ); + if (!/input\.enableCron\s*\?/u.test(services)) + return unconfirmed( + `\`enableCron\` is assigned in ${UPSTREAM_FILES.services} but no creation site consumes \`input.enableCron\`, so this probe cannot tell whether the service is still built conditionally`, + ); + return closed( + `\`enableCron\` is ownsElectronRuntimeCapabilities(runtimeOwnerKind), whose body accepts only 'electron' or undefined, and \`input.enableCron ? initializeRuntimeCron(...)\` still gates the creation site`, + ); +} + +/** + * Gate 2: could a host reach the created service? + * + * Two surfaces, because ADR 0012 names two: the host contract itself, and the + * `cliService` options, which is where `runtime.ts` spreads the services + * object whole. Either one carrying cron is enough. + */ +function probeGate2(sources: UpstreamSources): Gate { + const { hostContract, cliService } = sources; + if (hostContract === undefined) + return unconfirmed(`could not read ${UPSTREAM_FILES.hostContract}`); + if (cliService === undefined) + return unconfirmed(`could not read ${UPSTREAM_FILES.cliService}`); + + const hostBody = braceBody(hostContract, "interface CreatedLocalRuntimeHost"); + if (hostBody === undefined) + return unconfirmed( + `no \`interface CreatedLocalRuntimeHost\` body found in ${UPSTREAM_FILES.hostContract}`, + ); + const cliBody = braceBody(cliService, "interface CliServiceOptions"); + if (cliBody === undefined) + return unconfirmed( + `no \`interface CliServiceOptions\` body found in ${UPSTREAM_FILES.cliService}`, + ); + + const hostMembers = memberNames(hostBody); + const cliMembers = memberNames(cliBody); + if (hostMembers.length === 0 || cliMembers.length === 0) + return unconfirmed( + "a host surface parsed to zero members, so this probe cannot classify it", + ); + + const onHost = hostMembers.filter((name) => /cron/iu.test(name)); + if (onHost.length > 0) + return opened( + `CreatedLocalRuntimeHost carries ${onHost.join(", ")} in ${UPSTREAM_FILES.hostContract}`, + ); + const onCliService = cliMembers.filter((name) => /cron/iu.test(name)); + if (onCliService.length > 0) + return opened( + `CliServiceOptions carries ${onCliService.join(", ")} in ${UPSTREAM_FILES.cliService}, and \`runtime.ts\` spreads the whole services object into it`, + ); + return closed( + `CreatedLocalRuntimeHost carries only ${hostMembers.join(", ")} and CliServiceOptions only ${cliMembers.join(", ")} — nothing the WebUI can hold reaches a cron service`, + ); +} + +type ProbeStatus = + | "closed" + | "gate1-open" + | "gate2-open" + | "both-open" + | "unconfirmed"; + +interface ProbeReport { + readonly status: ProbeStatus; + readonly gate1: Gate; + readonly gate2: Gate; + /** The one line printed on every run. Never a warning. */ + readonly summary: string; + /** Printed only when the state deserves a look. Never a failure. */ + readonly notice?: string; +} + +function assess(sources: UpstreamSources, implementation: string): ProbeReport { + const gate1 = probeGate1(sources); + const gate2 = probeGate2(sources); + const summary = [ + "webui scheduled tasks:", + `gate1=${gate1.state}`, + `gate2=${gate2.state}`, + `probe=${gate1.state === "unconfirmed" || gate2.state === "unconfirmed" ? "unconfirmed" : "ok"}`, + `implementation=${implementation}`, + ].join(" "); + + if (gate1.state === "unconfirmed" || gate2.state === "unconfirmed") { + const unreadable = [ + ...(gate1.state === "unconfirmed" ? [`gate 1: ${gate1.detail}`] : []), + ...(gate2.state === "unconfirmed" ? [`gate 2: ${gate2.detail}`] : []), + ].join("\n "); + return { + status: "unconfirmed", + gate1, + gate2, + summary, + // Deliberately unlike every other notice, and it leads with the + // negation: a reader who only skims must not take this line as licence + // to remove the WebUI's scheduler. + notice: [ + "webui scheduled tasks: PROBE COULD NOT CONFIRM THE UPSTREAM STATE — THIS IS NOT A RETIREMENT SIGNAL", + ` ${unreadable}`, + " An unreadable probe is not a closed gate and not an open one. Do NOT retire,", + " replace or disable the WebUI's own scheduler on the basis of this line.", + ` Update the probe in ${path.relative(repoRoot, fileURLToPath(import.meta.url))} to read where the upstream gates now live.`, + ].join("\n"), + }; + } + + if (gate1.state === "open" && gate2.state === "open") + return { + status: "both-open", + gate1, + gate2, + summary, + notice: [ + `webui scheduled tasks: ADR 0012's RETIREMENT CONDITION IS NOW MET — both gates are open.`, + ` gate 1 (creation) ${gate1.detail}`, + ` gate 2 (reachability) ${gate2.detail}`, + ` Read ${ADR_0012} and do the retirement it specifies:`, + " 1. Re-implement WebuiScheduledTaskPort against the runtime's CronService. The adapter is the only thing that changes.", + " 2. Migrate the webui_scheduled_task table into the runtime's table, in the same change — not as a follow-up.", + " 3. Tear down the WebUI's own tick loop in that same change. Two schedulers over one queue is exactly what the decision exists to prevent.", + ].join("\n"), + }; + + if (gate1.state === "closed" && gate2.state === "closed") + // The quiet state: one status line, no warning, nothing to decide. + return { status: "closed", gate1, gate2, summary }; + + const openGate = gate1.state === "open" ? gate1 : gate2; + const closedGate = gate1.state === "open" ? gate2 : gate1; + const which = + gate1.state === "open" + ? [ + "the runtime now builds a cron service for this host kind (gate 1 is open)", + "the host contract still does not carry it (gate 2 is closed)", + ] + : [ + "the host contract now exposes a cron slot (gate 2 is open)", + "no cron service is built for this host kind (gate 1 is closed)", + ]; + return { + status: gate1.state === "open" ? "gate1-open" : "gate2-open", + gate1, + gate2, + summary, + notice: [ + "webui scheduled tasks: UPSTREAM HAS MOVED, AND IT IS STILL NOT A RETIREMENT SIGNAL", + ` ${which[0]}, but ${which[1]}.`, + ` open: ${openGate.detail}`, + ` closed: ${closedGate.detail}`, + " The service is therefore not usable by a WebUI host, so the surface stays as it is.", + ` ADR 0012 still stands: keep the WebUI's own scheduler and its own table. Re-read this when both gates are open.`, + ].join("\n"), + }; +} + +/** + * The implementation currently wired, read from the object that answers — + * never asserted on, and never able to fail the run. `"unknown"` when it + * cannot be constructed, because reporting must not throw. + */ +async function currentImplementation(): Promise { + const directory = mkdtempSync(path.join(os.tmpdir(), "webui-scheduled-task-probe-")); + let runtime: WebuiScheduledTaskRuntime | undefined; + try { + const store = new ScheduledTaskStore( + openScheduledTaskDatabase(path.join(directory, "scheduled-tasks.sqlite")), + ); + runtime = new WebuiScheduledTaskRuntime({ + store, + sendMessage: async () => ({ ok: true, source: [] }), + createSession: async () => ({ sessionId: "unused" }), + }); + return (await runtime.getScheduledTaskCapability()).source; + } catch { + return "unknown"; + } finally { + try { + runtime?.dispose(); + } finally { + rmSync(directory, { recursive: true, force: true }); + } + } +} + +/** Prints the report. Returns nothing and throws nothing, by construction. */ +function report(report: ProbeReport): void { + // vitest's default reporter swallows console output from passing tests, + // which would make this reporter invisible in CI. Write to the real stream. + process.stdout.write(`${report.summary} +`); + if (report.notice) process.stderr.write(`${report.notice} +`); +} + +describe("ADR 0012 retirement-condition probe (reports; never fails)", () => { + it("reports the upstream state for a tui + cliEmbedded host", async () => { + // The live reading. Whatever it says, this test passes: upstream openness + // is reported, never enforced. See the header for why that is the rule. + report(assess(readUpstreamSources(repoRoot), await currentImplementation())); + expect(true).toBe(true); + }); + + it("names the upstream source files it reads", () => { + // The probe is only honest if the reader knows what it is reading. + for (const relative of Object.values(UPSTREAM_FILES)) + expect(relative.startsWith("packages/local-runtime-v2/")).toBe(true); + }); +}); + +/** + * The probe's own behaviour, against literal sources. These are the cases the + * live reading cannot exercise on demand: the fixture text here cannot change + * because upstream changed, so an assertion can be honest about what the + * probe would report. + */ +describe("the probe itself", () => { + const real = readUpstreamSources(repoRoot); + // The fixtures below start from the real files and move one thing each, so + // a fixture is never a hand-written file that could be wrong about the + // upstream's shape in some unrelated way. + const requireSource = (): UpstreamSources => { + if ( + real.services === undefined || + real.composition === undefined || + real.hostContract === undefined || + real.cliService === undefined + ) + throw new Error("the upstream sources this fixture builds on are unreadable"); + return real; + }; + + it("reads the live sources as both gates closed", () => { + const report = assess(requireSource(), "webui-own"); + expect(report.gate1.state).toBe("closed"); + expect(report.gate2.state).toBe("closed"); + expect(report.status).toBe("closed"); + // A closed upstream is the quiet state: a status line and no warning. + expect(report.notice).toBeUndefined(); + expect(report.summary).toContain("gate1=closed gate2=closed probe=ok"); + }); + + it("reports gate 1 open on its own as partial, and keeps gate 2 closed", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + services: (sources.services ?? "").replace( + "enableCron: ownsElectronRuntimeCapabilities(options.runtimeOwnerKind)", + "enableCron: true", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("open"); + expect(report.gate2.state).toBe("closed"); + expect(report.status).toBe("gate1-open"); + expect(report.notice).toContain("STILL NOT A RETIREMENT SIGNAL"); + expect(report.notice).toContain("keep the WebUI's own scheduler"); + }); + + it("reports a predicate that accepts tui as gate 1 open", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + composition: (sources.composition ?? "").replace( + "runtimeOwnerKind === undefined || runtimeOwnerKind === 'electron'", + "runtimeOwnerKind === 'electron' || runtimeOwnerKind === 'tui'", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("open"); + expect(report.status).toBe("gate1-open"); + }); + + it("reports gate 2 open on its own as partial, and keeps gate 1 closed", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + hostContract: (sources.hostContract ?? "").replace( + " cliService?: import('./cli-service.js').CliService;", + " cliService?: import('./cli-service.js').CliService;\n cron?: unknown;", + ), + }, + "webui-own", + ); + expect(report.gate2.state).toBe("open"); + expect(report.gate1.state).toBe("closed"); + expect(report.status).toBe("gate2-open"); + expect(report.notice).toContain("STILL NOT A RETIREMENT SIGNAL"); + }); + + it("reports a cron slot on the cliService options as gate 2 open", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + cliService: (sources.cliService ?? "").replace( + " readonly management: CliManagementApplication;", + " readonly management: CliManagementApplication;\n readonly cron?: unknown;", + ), + }, + "webui-own", + ); + expect(report.gate2.state).toBe("open"); + }); + + it("reports both gates open as the retirement condition being met", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + services: (sources.services ?? "").replace( + "enableCron: ownsElectronRuntimeCapabilities(options.runtimeOwnerKind)", + "enableCron: true", + ), + hostContract: (sources.hostContract ?? "").replace( + " cliService?: import('./cli-service.js').CliService;", + " cliService?: import('./cli-service.js').CliService;\n cron?: unknown;", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("open"); + expect(report.gate2.state).toBe("open"); + expect(report.status).toBe("both-open"); + expect(report.notice).toContain(ADR_0012); + expect(report.notice).toContain("CronService"); + expect(report.notice).toContain("webui_scheduled_task"); + // The step that is easiest to skip and worst to get wrong. + expect(report.notice).toContain("Tear down the WebUI's own tick loop"); + }); + + it("reports an unmatchable enableCron as unconfirmed, not as closed", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + services: (sources.services ?? "").replace( + "enableCron: ownsElectronRuntimeCapabilities(options.runtimeOwnerKind)", + "enableCron: resolveCronEnablement(options)", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("unconfirmed"); + expect(report.status).toBe("unconfirmed"); + expect(report.summary).toContain("probe=unconfirmed"); + // The message must not be mistakable for the retirement notice. + expect(report.notice).toContain("NOT A RETIREMENT SIGNAL"); + expect(report.notice).not.toContain("RETIREMENT CONDITION IS NOW MET"); + expect(report.notice).toContain("Do NOT retire"); + }); + + it("reports a renamed predicate as unconfirmed, not as closed", () => { + const sources = requireSource(); + const report = assess( + { + ...sources, + composition: (sources.composition ?? "").replace( + "function ownsElectronRuntimeCapabilities", + "function ownsDesktopRuntimeCapabilities", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("unconfirmed"); + expect(report.notice).toContain("COULD NOT CONFIRM"); + expect(report.notice).toContain("Do NOT retire"); + }); + + it("reports a flag nothing consumes as unconfirmed rather than closed", () => { + // The blind spot this guard exists for: an `enableCron` that is assigned + // but never read means the service is built unconditionally, and reading + // only the assignment would have called that closed. + const sources = requireSource(); + const report = assess( + { + ...sources, + services: (sources.services ?? "").replace( + "cron = input.enableCron", + "cron = ALWAYS_ON", + ), + }, + "webui-own", + ); + expect(report.gate1.state).toBe("unconfirmed"); + }); + + it("reports a missing upstream file as unconfirmed, not as closed", () => { + const sources = requireSource(); + const report = assess({ ...sources, hostContract: undefined }, "webui-own"); + expect(report.gate2.state).toBe("unconfirmed"); + expect(report.status).toBe("unconfirmed"); + expect(report.notice).toContain(UPSTREAM_FILES.hostContract); + }); + + it("keeps reporting the implementation without asserting on it", async () => { + // Informational by design: the moment this becomes a pass/fail + // assertion, swapping the adapter turns the suite red, which is the thing + // this file must not do. + expect(["webui-own", "runtime-cron", "none", "unknown"]).toContain( + await currentImplementation(), + ); + }); +}); diff --git a/packages/webui/test/unit/webui-scheduled-task.test.ts b/packages/webui/test/unit/webui-scheduled-task.test.ts new file mode 100644 index 000000000..ed375f159 --- /dev/null +++ b/packages/webui/test/unit/webui-scheduled-task.test.ts @@ -0,0 +1,629 @@ +// Server-side scheduled tasks (`定时任务`) for the WebUI process. +// +// The scope here is the WebUI's OWN scheduler: its own SQLite file, its own +// table, its own tick loop. It is deliberately not the desktop client's +// scheduled-task engine, and the two never see each other — see the header of +// `src/server/scheduled-task-scheduler.ts` for why reading the v2 table was +// rejected rather than merely avoided. +// +// What this file pins: +// 1. The store's schema migration is idempotent and non-destructive across +// restarts (the WebUI is an on-demand process; this file is the first +// persistence the service owns). +// 2. list / create / update / delete. +// 3. A due task fires exactly once, through the host's `sendMessage`. +// 4. A failed delivery is recorded as a failure, never swallowed. +// 5. THE product boundary: a run that came due while webui was NOT running +// is not replayed, but it is counted and timestamped so the panel can say +// "missed N times" instead of silently losing it. +// 6. The service owns the loop (beside the heartbeat) and the six operations +// are wired into the operation registry. + +import { describe, expect, it } from "vitest"; +import { existsSync, mkdtempSync, rmSync } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import WebSocket from "ws"; + +import { + SCHEDULED_TASK_SCHEMA_VERSION, + openScheduledTaskDatabase, + ScheduledTaskStore, +} from "../../src/server/scheduled-task-store.js"; +import { + WebuiScheduledTaskRuntime, + type ScheduledTaskSendMessage, +} from "../../src/server/scheduled-task-scheduler.js"; +import { createHarnessPortFromHost } from "../../src/server/host.js"; +import { + createOperationRegistry, + createScheduledTaskOperation, + deleteScheduledTaskOperation, + getScheduledTaskCapabilityOperation, + listScheduledTasksOperation, + updateScheduledTaskOperation, +} from "../../src/server/operation/operations.js"; +import { WebuiService } from "../../src/server/service.js"; +import { WEBUI_PROTOCOL_VERSION } from "../../src/server/envelope.js"; +import { + LIST_SCHEDULED_TASKS_OPERATION_NAME, + CREATE_SCHEDULED_TASK_OPERATION_NAME, + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + DELETE_SCHEDULED_TASK_OPERATION_NAME, + TRIGGER_SCHEDULED_TASK_OPERATION_NAME, + GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME, +} from "../../src/server/operation/names.js"; + +/** One request, one response frame — the client's actual path. */ +function requestOnce(ws: WebSocket, request: unknown): Promise { + return new Promise((resolve, reject) => { + const onMessage = (raw: import("ws").RawData) => { + ws.off("message", onMessage); + ws.off("error", onError); + try { + resolve(JSON.parse(raw.toString("utf8"))); + } catch (error) { + reject(error); + } + }; + const onError = (error: Error) => { + ws.off("message", onMessage); + reject(error); + }; + ws.on("message", onMessage); + ws.once("error", onError); + ws.send(JSON.stringify(request)); + }); +} + +/** A host that carries no scheduled-task runtime: the fail-closed path. */ +function hostWithoutScheduledTasks(): Parameters[0] { + return { apiHost: { close: async () => undefined } }; +} + +function temporaryDirectory(): string { + return mkdtempSync(path.join(os.tmpdir(), "webui-scheduled-task-")); +} + +interface Harness { + readonly store: ScheduledTaskStore; + readonly runtime: WebuiScheduledTaskRuntime; + readonly sent: Array<{ readonly sessionId: string; readonly prompt: string }>; + readonly created: Array<{ readonly name: string }>; + setNow(value: number): void; + setDelivery(failure: Error | undefined): void; + close(): void; +} + +function createHarness( + options: { + readonly tickIntervalMs?: number; + readonly missGraceMs?: number; + /** Use the wall clock, for the case where the service owns the timer. */ + readonly realClock?: boolean; + } = {}, +): Harness { + const directory = temporaryDirectory(); + let now = 1_700_000_000_000; + let failure: Error | undefined; + const sent: Array<{ sessionId: string; prompt: string }> = []; + const created: Array<{ name: string }> = []; + const store = new ScheduledTaskStore( + openScheduledTaskDatabase(path.join(directory, "scheduled-tasks.sqlite")), + ); + const sendMessage: ScheduledTaskSendMessage = async (request) => { + sent.push({ sessionId: request.id, prompt: request.content ?? "" }); + if (failure) throw failure; + return { ok: true, source: [{ eventJson: "{}" }] }; + }; + const runtime = new WebuiScheduledTaskRuntime({ + store, + now: options.realClock ? () => Date.now() : () => now, + tickIntervalMs: options.tickIntervalMs ?? 30, + missGraceMs: options.missGraceMs ?? 60_000, + createSession: async (request) => { + created.push({ name: request.name }); + if (failure) throw failure; + return { sessionId: `session-${created.length}` }; + }, + sendMessage, + }); + return { + store, + runtime, + sent, + created, + setNow: (value) => { + now = value; + }, + setDelivery: (value) => { + failure = value; + }, + close: () => { + runtime.dispose(); + rmSync(directory, { recursive: true, force: true }); + }, + }; +} + +describe("scheduled task store", () => { + it("migrates on first open and is idempotent across restarts", () => { + const directory = temporaryDirectory(); + const file = path.join(directory, "scheduled-tasks.sqlite"); + try { + const first = openScheduledTaskDatabase(file); + expect(new ScheduledTaskStore(first).schemaVersion()).toBe(SCHEDULED_TASK_SCHEMA_VERSION); + const store = new ScheduledTaskStore(first); + const created = store.create({ + name: "nightly", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-1", + prompt: "summarise the day", + scheduleKind: "interval", + intervalMs: 60_000, + nowMs: 1_000, + }); + store.close(); + + // A second start must neither re-run destructively nor lose the row. + const second = openScheduledTaskDatabase(file); + expect(new ScheduledTaskStore(second).schemaVersion()).toBe(SCHEDULED_TASK_SCHEMA_VERSION); + const reopened = new ScheduledTaskStore(second); + expect(reopened.list().map((task) => task.taskId)).toEqual([ + created.taskId, + ]); + // And migrating an already-migrated file is a no-op, not an error. + expect(() => reopened.migrate()).not.toThrow(); + expect(reopened.list()).toHaveLength(1); + reopened.close(); + } finally { + rmSync(directory, { recursive: true, force: true }); + } + }); + + it("supports list, create, update and delete", () => { + const harness = createHarness(); + try { + expect(harness.store.list()).toEqual([]); + const created = harness.store.create({ + name: "hourly", + agentName: "main", + sessionTarget: "new", + sessionId: null, + prompt: "check the build", + scheduleKind: "interval", + intervalMs: 3_600_000, + runAtMs: 5_000, + nowMs: 1_000, + }); + expect(created.enabled).toBe(true); + expect(created.nextRunAtMs).toBe(5_000); + expect(harness.store.list()).toHaveLength(1); + + const updated = harness.store.update(created.taskId, { + prompt: "check the release build", + enabled: false, + nowMs: 2_000, + }); + expect(updated?.prompt).toBe("check the release build"); + expect(updated?.enabled).toBe(false); + expect(harness.store.get(created.taskId)?.updatedAtMs).toBe(2_000); + + expect(harness.store.remove(created.taskId)).toBe(true); + // Deleting twice reports absence rather than throwing, so a panel that + // retries a delete does not have to distinguish the two cases. + expect(harness.store.remove(created.taskId)).toBe(false); + expect(harness.store.list()).toEqual([]); + } finally { + harness.close(); + } + }); +}); + +describe("scheduled task scheduler", () => { + it("fires a due task exactly once through the host's sendMessage", async () => { + const harness = createHarness(); + try { + const task = harness.store.create({ + name: "interval", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-7", + prompt: "run the report", + scheduleKind: "interval", + intervalMs: 1_000, + runAtMs: 1_500, + nowMs: 1_000, + }); + harness.setNow(1_500); + const fired = await harness.runtime.tick(); + expect(fired.firedTaskIds).toEqual([task.taskId]); + expect(harness.sent).toEqual([ + { sessionId: "session-7", prompt: "run the report" }, + ]); + const afterFirst = harness.store.get(task.taskId); + expect(afterFirst?.lastStatus).toBe("succeeded"); + expect(afterFirst?.lastRunAtMs).toBe(1_500); + // The interval is anchored on the due time, not on "now + interval", + // so a late tick does not drift the schedule forward. + expect(afterFirst?.nextRunAtMs).toBe(2_500); + + // A tick that finds nothing due must not fire again. + const second = await harness.runtime.tick(); + expect(second.firedTaskIds).toEqual([]); + expect(harness.sent).toHaveLength(1); + + harness.setNow(2_500); + await harness.runtime.tick(); + expect(harness.sent).toHaveLength(2); + } finally { + harness.close(); + } + }); + + it("records a failure instead of swallowing it", async () => { + const harness = createHarness(); + try { + const task = harness.store.create({ + name: "failing", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-7", + prompt: "run the report", + scheduleKind: "once", + runAtMs: 1_500, + nowMs: 1_000, + }); + harness.setDelivery(new Error("agent is busy")); + harness.setNow(1_500); + const fired = await harness.runtime.tick(); + expect(fired.failedTaskIds).toEqual([task.taskId]); + const failed = harness.store.get(task.taskId); + expect(failed?.lastStatus).toBe("failed"); + expect(failed?.lastError).toContain("agent is busy"); + // A one-shot task is finished either way: a failure must not leave it + // armed to fire forever, and must not leave it looking pending. + expect(failed?.enabled).toBe(false); + expect(failed?.nextRunAtMs).toBeNull(); + } finally { + harness.close(); + } + }); + + it("does not replay runs missed while webui was closed, but counts them", async () => { + const harness = createHarness(); + try { + // The task was due 3 hours ago. webui was not running then, by design: + // an on-demand process does not catch up on the work it was not alive + // to do. The product boundary says "do not replay" — this test exists so + // that boundary can never be quietly turned into a catch-up loop. + const task = harness.store.create({ + name: "while-asleep", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-7", + prompt: "run the report", + scheduleKind: "interval", + intervalMs: 3_600_000, + runAtMs: 1_000, + nowMs: 1_000, + }); + harness.setNow(1_000 + 3 * 3_600_000); + const outcome = await harness.runtime.tick(); + + expect(outcome.firedTaskIds).toEqual([]); + expect(harness.sent).toEqual([]); + expect(outcome.missedTaskIds).toEqual([task.taskId]); + const missed = harness.store.get(task.taskId); + expect(missed?.missedCount).toBe(1); + expect(missed?.lastMissedAtMs).toBe(1_000); + expect(missed?.lastRunAtMs).toBeNull(); + // A missed run is not a failure: nothing was attempted. + expect(missed?.lastStatus).toBe("missed"); + // The next occurrence is the next slot on the grid, not a replay of the + // slot that was lost, so one downtime cannot cascade into a burst. Here + // `now` sits exactly on a slot boundary, so that boundary is what the + // task is armed for — the slot that was lost is the earlier one, and the + // current slot still gets its turn on the next tick. + expect(missed?.nextRunAtMs).toBe(1_000 + 3 * 3_600_000); + harness.setNow(1_000 + 3 * 3_600_000 + 1_000); + await harness.runtime.tick(); + expect(harness.sent).toHaveLength(1); + } finally { + harness.close(); + } + }); + + it("creates a session for a new-session target and records which one", async () => { + const harness = createHarness(); + try { + const task = harness.store.create({ + name: "new-session", + agentName: "main", + sessionTarget: "new", + sessionId: null, + prompt: "start a fresh report", + scheduleKind: "once", + runAtMs: 1_500, + nowMs: 1_000, + }); + harness.setNow(1_500); + await harness.runtime.tick(); + expect(harness.created).toEqual([{ name: "main" }]); + expect(harness.sent).toEqual([ + { sessionId: "session-1", prompt: "start a fresh report" }, + ]); + expect(harness.store.get(task.taskId)?.lastSessionId).toBe("session-1"); + } finally { + harness.close(); + } + }); + + it("runs a task on demand without waiting for its schedule", async () => { + const harness = createHarness(); + try { + const task = harness.store.create({ + name: "manual", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-7", + prompt: "run the report", + scheduleKind: "interval", + intervalMs: 3_600_000, + runAtMs: 1_000 + 10 * 3_600_000, + nowMs: 1_000, + }); + const result = await harness.runtime.triggerScheduledTaskNow({ + taskId: task.taskId, + }); + expect(result.started).toBe(true); + expect(harness.sent).toHaveLength(1); + // A manual run must not consume the scheduled slot: the next automatic + // run stays where it was. + expect(harness.store.get(task.taskId)?.nextRunAtMs).toBe( + 1_000 + 10 * 3_600_000, + ); + } finally { + harness.close(); + } + }); + + it("reports the capability as unavailable when the store cannot be opened", async () => { + const runtime = new WebuiScheduledTaskRuntime({ + store: undefined, + databaseFile: path.join(temporaryDirectory(), "nested", "store.sqlite"), + openDatabase: () => { + throw new Error("sqlite native module is missing"); + }, + sendMessage: async () => ({ ok: true, source: [] }), + createSession: async () => ({ sessionId: "unused" }), + }); + try { + const capability = await runtime.getScheduledTaskCapability(); + expect(capability.available).toBe(false); + expect(capability.reason).toContain("sqlite native module is missing"); + await expect(runtime.listScheduledTasks()).rejects.toThrow( + /sqlite native module is missing/u, + ); + } finally { + runtime.dispose(); + } + }); +}); + +describe("scheduled task service wiring", () => { + it("rejects an inexpressible schedule instead of crashing on it", async () => { + // Regression: the validators return `null` for an absent optional field, + // and `typeof null === "object"` made that look like a failure object. The + // dispatcher then read `.ok` off a `null` and the whole message handler + // threw, taking the request down instead of answering it. + const harness = createHarness(); + try { + // Driven through the typed descriptors, so the payload each one accepts + // is checked at compile time too. + const invalidCases = [ + // A recurring task with no interval, and a one-shot with no moment: + // both are missing the half that makes them expressible. + [createScheduledTaskOperation, { name: "a", agentName: "main", sessionTarget: "new", prompt: "p", scheduleKind: "interval" }], + [createScheduledTaskOperation, { name: "a", agentName: "main", sessionTarget: "new", prompt: "p", scheduleKind: "once" }], + [createScheduledTaskOperation, { name: "", agentName: "main", sessionTarget: "new", prompt: "p", scheduleKind: "once", runAtMs: 1 }], + [createScheduledTaskOperation, { name: "a", agentName: "main", sessionTarget: "existing", prompt: "p", scheduleKind: "once", runAtMs: 1 }], + [createScheduledTaskOperation, { name: "a", agentName: "main", sessionTarget: "new", prompt: "p", scheduleKind: "interval", intervalMs: 1 }], + [updateScheduledTaskOperation, { taskId: "x" }], + [deleteScheduledTaskOperation, {}], + [listScheduledTasksOperation, { agentName: "" }], + [getScheduledTaskCapabilityOperation, { unexpected: true }], + ] as const; + for (const [operation, body] of invalidCases) { + const validated = operation.validate(body); + expect( + (validated as { readonly ok?: boolean })?.ok, + `${operation.name} must answer, not throw`, + ).toBe(false); + } + + // The other half of the same regression: an explicit `null` for an + // optional field is a *value*, not a failure, so the create must be + // accepted and the null must survive into the stored row. + const validated = createScheduledTaskOperation.validate({ + name: "a", + agentName: "main", + sessionTarget: "new", + prompt: "p", + scheduleKind: "interval", + intervalMs: 1_000, + runAtMs: null, + }); + expect(validated.ok).toBe(true); + if (!validated.ok) throw new Error("expected the create to validate"); + const created = await harness.runtime.createScheduledTask(validated.body); + expect(created.runAtMs).toBeNull(); + expect(created.intervalMs).toBe(1_000); + } finally { + harness.close(); + } + }); + + it("fails closed on a host with no scheduled-task runtime", async () => { + const port = createHarnessPortFromHost(hostWithoutScheduledTasks()); + await expect(port.listScheduledTasks()).rejects.toThrow( + /scheduled tasks are not available/u, + ); + expect(await port.getScheduledTaskCapability()).toEqual({ + available: false, + source: "none", + reason: expect.stringContaining("scheduled tasks are not available"), + }); + }); + + it("builds its own runtime from the dataDir the port reports", async () => { + // The production path: no caller passes a runtime, and the surface is + // still live, because `version().dataDir` is the only thing it needs. A + // service that only worked with an injected runtime would be unreachable + // in the shipped CLI, where nothing constructs one. + const dataDir = temporaryDirectory(); + const sent: Array<{ readonly id: string; readonly content?: string }> = []; + const port = createHarnessPortFromHost({ + apiHost: { close: async () => undefined }, + dataDir, + cliService: { + createSession: async (request: { readonly name: string }) => ({ + sessionId: `session-for-${request.name}`, + }), + sendMessage: async (request: { readonly id: string; readonly content?: string }) => { + sent.push({ id: request.id, ...(request.content === undefined ? {} : { content: request.content }) }); + return { ok: true as const, source: [{ eventJson: "{}" }] }; + }, + } as never, + }); + const service = new WebuiService({ + port, + dev: true, + scheduledTaskTickIntervalMs: 20, + }); + try { + // The runtime belongs to the service, not to the port, so it is reached + // the way a browser reaches it: over the wire. + const info = await service.start(); + const ws = new WebSocket(info.boundUrl); + ws.on("error", () => undefined); + await new Promise((resolve, reject) => { + ws.once("open", resolve); + ws.once("error", reject); + }); + const call = async (operation: string, body: unknown) => { + const response = (await requestOnce(ws, { + protocolVersion: WEBUI_PROTOCOL_VERSION, + kind: "request", + requestId: `req-${operation}`, + operation, + body, + })) as { readonly kind: string; readonly body?: unknown }; + expect(response.kind, `${operation} must not be an error frame`).toBe( + "response", + ); + return response.body; + }; + + expect(await call(GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME, undefined)).toEqual({ + available: true, + // The wire shape names the implementation, so the surface a client + // can reach is itself evidence of which scheduler is in place. + source: "webui-own", + }); + + const created = (await call(CREATE_SCHEDULED_TASK_OPERATION_NAME, { + name: "own-runtime", + agentName: "main", + sessionTarget: "new", + prompt: "run the report", + scheduleKind: "once", + runAtMs: Date.now(), + })) as { readonly taskId: string; readonly enabled: boolean }; + expect(created.enabled).toBe(true); + // The store is this surface's own file, under the host's dataDir. + expect( + existsSync(path.join(dataDir, "webui", "scheduled-tasks.sqlite")), + ).toBe(true); + + const deadline = Date.now() + 4_000; + while (sent.length === 0 && Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 20)); + } + expect(sent).toEqual([ + { id: "session-for-main", content: "run the report" }, + ]); + const listed = (await call(LIST_SCHEDULED_TASKS_OPERATION_NAME, {})) as { + readonly tasks: ReadonlyArray<{ readonly lastStatus: string | null }>; + }; + expect(listed.tasks[0]?.lastStatus).toBe("succeeded"); + ws.close(); + } finally { + await service.close(); + rmSync(dataDir, { recursive: true, force: true }); + } + }); + + it("reports the capability as unavailable when the host has no dataDir", async () => { + const port = createHarnessPortFromHost(hostWithoutScheduledTasks()); + expect(await port.getScheduledTaskCapability()).toEqual({ + available: false, + source: "none", + reason: expect.stringContaining("scheduled tasks are not available"), + }); + }); + + it("registers all six operations and runs the loop beside the heartbeat", async () => { + const harness = createHarness({ tickIntervalMs: 20, realClock: true }); + const service = new WebuiService({ + port: createHarnessPortFromHost(hostWithoutScheduledTasks()), + dev: true, + scheduledTasks: { runtime: harness.runtime, tickIntervalMs: 20 }, + }); + try { + const registry = createOperationRegistry( + createHarnessPortFromHost(hostWithoutScheduledTasks()), + ); + for (const name of [ + LIST_SCHEDULED_TASKS_OPERATION_NAME, + CREATE_SCHEDULED_TASK_OPERATION_NAME, + UPDATE_SCHEDULED_TASK_OPERATION_NAME, + DELETE_SCHEDULED_TASK_OPERATION_NAME, + TRIGGER_SCHEDULED_TASK_OPERATION_NAME, + GET_SCHEDULED_TASK_CAPABILITY_OPERATION_NAME, + ]) { + expect(registry.has(name), `expected ${name} in the registry`).toBe( + true, + ); + } + + const task = harness.store.create({ + name: "service-driven", + agentName: "main", + sessionTarget: "existing", + sessionId: "session-9", + prompt: "run the report", + scheduleKind: "once", + runAtMs: Date.now(), + nowMs: Date.now(), + }); + const info = await service.start(); + expect(info.host).toBe("127.0.0.1"); + + // The tick loop lives in the service, not in the runtime: nothing calls + // `tick()` here, so a fired task proves the service is driving it. + const deadline = Date.now() + 3_000; + while (harness.sent.length === 0 && Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + expect(harness.sent).toEqual([ + { sessionId: "session-9", prompt: "run the report" }, + ]); + expect(harness.store.get(task.taskId)?.lastStatus).toBe("succeeded"); + } finally { + await service.close(); + harness.close(); + } + }); +}); diff --git a/packages/webui/test/unit/webui-service.test.ts b/packages/webui/test/unit/webui-service.test.ts index a428d38c8..e619899f2 100644 --- a/packages/webui/test/unit/webui-service.test.ts +++ b/packages/webui/test/unit/webui-service.test.ts @@ -64,11 +64,20 @@ import type { WebuiImportSessionTransferRequest, WebuiImportSessionTransferResult, WebuiSessionTransferFile, + WebuiScheduledTask, + WebuiScheduledTaskTriggerResult, } from "../../src/server/port.js"; import { createWebuiTransport } from "../../src/client/transport.js"; import { WebuiTerminalManager } from "../../src/server/terminal.js"; import { getWorkspaceReviewSummaryOperation, listWorkspaceReviewFileDiffsOperation, getWorkspaceReviewFileContentOperation, searchWorkspaceReviewDiffsOperation } from "../../src/server/operation/workspace.js"; +/** + * The one message a host without a scheduled-task runtime reports. Kept as a + * constant so the port methods and the capability probe cannot drift apart. + */ +const SCHEDULED_TASKS_UNAVAILABLE = + "scheduled tasks are not available: this host exposes no scheduled-task runtime"; + type CloseEvent = [number, Buffer]; // `once` from `node:events` is overloaded and not generic, so // `ReturnType>` does not type-check; the concrete @@ -270,6 +279,32 @@ class ScriptedHarnessPort implements WebuiHarnessPort { return { success: true }; } + // WebUI scheduled tasks. This double stands in for a host, and a host with + // no scheduled-task runtime is exactly the fail-closed case, so the scripted + // port reports the capability as unavailable rather than pretending. + async listScheduledTasks() { + return { tasks: [], total: 0 }; + } + async createScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + } + async updateScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + } + async deleteScheduledTask(): Promise<{ readonly success: boolean }> { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + } + async triggerScheduledTaskNow(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + } + async getScheduledTaskCapability() { + return { + available: false as const, + source: "none" as const, + reason: SCHEDULED_TASKS_UNAVAILABLE, + }; + } + async invalidateAuth(): Promise { // Test fixture: nothing to invalidate. } @@ -3777,6 +3812,28 @@ describe("WebUI shutdown order (criterion 7)", () => { panel: { scene: 0, days: [] }, }; }, + async listScheduledTasks() { + return { tasks: [], total: 0 }; + }, + async createScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async updateScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async deleteScheduledTask(): Promise<{ readonly success: boolean }> { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async triggerScheduledTaskNow(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async getScheduledTaskCapability() { + return { + available: false as const, + source: "none" as const, + reason: SCHEDULED_TASKS_UNAVAILABLE, + }; + }, async close() { // The service awaits wsServer.close() and httpServer.close() // before calling port.close(), so by the time we land here the @@ -4155,6 +4212,28 @@ describe("WebUI shutdown order (criterion 7)", () => { panel: { scene: 0, days: [] }, }; }, + async listScheduledTasks() { + return { tasks: [], total: 0 }; + }, + async createScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async updateScheduledTask(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async deleteScheduledTask(): Promise<{ readonly success: boolean }> { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async triggerScheduledTaskNow(): Promise { + throw new Error(SCHEDULED_TASKS_UNAVAILABLE); + }, + async getScheduledTaskCapability() { + return { + available: false as const, + source: "none" as const, + reason: SCHEDULED_TASKS_UNAVAILABLE, + }; + }, async close() { await closeGate; }, diff --git a/packages/webui/tsconfig.test.json b/packages/webui/tsconfig.test.json index 770927efd..bc12a4eef 100644 --- a/packages/webui/tsconfig.test.json +++ b/packages/webui/tsconfig.test.json @@ -6,6 +6,7 @@ "include": [ "test/unit/webui-service.test.ts", "test/unit/webui-host-shape-invariant.test.ts", + "test/unit/webui-scheduled-task-upstream-probe.test.ts", "src/server/**/*.ts", "src/shared/**/*.ts", "src/client/contracts.ts", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 2a0c93a6a..5ce8f63cd 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -936,6 +936,9 @@ importers: '@xterm/xterm': specifier: 6.0.0 version: 6.0.0 + better-sqlite3: + specifier: 12.11.1 + version: 12.11.1 highlight.js: specifier: 10.7.3 version: 10.7.3 diff --git a/release/public-source.json b/release/public-source.json index 69f90f06c..4e87b3ebf 100644 --- a/release/public-source.json +++ b/release/public-source.json @@ -53,6 +53,7 @@ "docs/adr/0009-webui-reuses-the-desktop-visual-language.md", "docs/adr/0010-webui-ships-an-esbuild-artifact-with-vite-as-a-development-server.md", "docs/adr/0011-shared-event-corpus-in-local-runtime-v2.md", + "docs/adr/0012-the-webui-runs-its-own-scheduled-tasks-until-the-runtime-offers-one.md", "docs/agents/domain.md", "docs/agents/issue-tracker.md", "docs/agents/triage-labels.md", @@ -3645,6 +3646,7 @@ "packages/webui/src/server/operation/provider.ts", "packages/webui/src/server/operation/questionnaire.ts", "packages/webui/src/server/operation/queue.ts", + "packages/webui/src/server/operation/scheduled-task.ts", "packages/webui/src/server/operation/session.ts", "packages/webui/src/server/operation/workspace.ts", "packages/webui/src/server/port.ts", @@ -3654,6 +3656,8 @@ "packages/webui/src/server/projections/permissions.ts", "packages/webui/src/server/projections/usage.ts", "packages/webui/src/server/runtime-environment.ts", + "packages/webui/src/server/scheduled-task-scheduler.ts", + "packages/webui/src/server/scheduled-task-store.ts", "packages/webui/src/server/service.ts", "packages/webui/src/server/session-transfer.ts", "packages/webui/src/server/terminal.ts", @@ -3724,6 +3728,8 @@ "packages/webui/test/unit/webui-plan-mode.test.ts", "packages/webui/test/unit/webui-round3-acceptance.test.tsx", "packages/webui/test/unit/webui-runtime-environment.test.ts", + "packages/webui/test/unit/webui-scheduled-task.test.ts", + "packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts", "packages/webui/test/unit/webui-service.test.ts", "packages/webui/test/unit/webui-shell.test.ts", "packages/webui/test/unit/webui-stream.test.ts", diff --git a/test/vitest-suites.json b/test/vitest-suites.json index 0dbf30cb7..19b399622 100644 --- a/test/vitest-suites.json +++ b/test/vitest-suites.json @@ -265,6 +265,8 @@ "packages/webui/test/unit/webui-w0-css-structure.test.ts", "packages/webui/test/unit/webui-w2-effect-reducer.test.ts", "packages/webui/test/unit/webui-host-shape-invariant.test.ts", + "packages/webui/test/unit/webui-scheduled-task.test.ts", + "packages/webui/test/unit/webui-scheduled-task-upstream-probe.test.ts", "packages/webui/test/unit/composer-intent.test.ts", "packages/webui/test/unit/transcript-shape.test.ts", "packages/webui/test/unit/slash-classification.test.ts",