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",