mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-19 11:06:18 +00:00
* feat: add scheduled tasks MVP
* fix: harden scheduled task execution semantics
* feat(scheduled-tasks): preset-driven schedule form with timezone and live preview
Replace the raw cron input with a preset Select (hourly/daily/weekly/monthly/custom)
plus structured inputs (time picker, weekday toggles, day-of-month), datetime-local
for one-time tasks, a timezone selector defaulting to the browser timezone, and a
live human-readable preview. Reuses one ScheduledTaskScheduleInput for create and
edit; backend contract unchanged; zero new deps (pure Intl + DST-safe offset helpers).
* feat(scheduled-tasks): full-page i18n + recipe templates + E2E locale pin
Localize the rest of the scheduled-tasks page (filters, detail pane, actions,
edit form, run list, enum values) via t.scheduledTasks.* in en/zh. Add four
built-in recipe templates (GitHub Trending, news digest, issue triage, weekly
report) exposed as a chip row that pre-fills title + prompt + schedule. Pin
Playwright locale to en-US so E2E selectors stay stable against i18n. No backend
change, no new deps.
* fix(scheduled-tasks): idempotent 0003 migration, update head constants, future-date once test
Merge with main surfaced three CI failures:
- 0003_scheduled_tasks create_table collided with legacy test seeds that
build from full metadata; guard with inspector.has_table so the revision
no-ops when the table already exists (0004/0005 are already idempotent via
_helpers.py).
- persistence bootstrap concurrency/regression tests pinned HEAD to main's
0002_runs_token_usage; bump to the new head 0005_scheduled_task_thread_nullable.
- once-task router test used a fixed past run_at and tripped the
must-be-in-the-future validation; use a future date.
* address review: ok-check, 502 for trigger failure, mock fields, migration filename, doc fences
- fetchThreadScheduledTasks now checks response.ok like the other fetchers.
- trigger endpoint returns 502 (not 409) when dispatch fails outright, so
clients can distinguish a real conflict from a server-side failure.
- E2E mock normalizes scheduled-task objects with context_mode/last_thread_id
and nullable thread_id, matching the backend contract the UI renders against.
- Rename 0002_scheduled_tasks.py -> 0003_scheduled_tasks.py to match its
revision id (file was renamed in spirit already; filename now follows).
- CONFIGURATION.md: close the Tool Groups yaml fence and drop the stray fence
after the Scheduler notes so the sections render correctly.
* fix(scheduled-tasks): harden lease, poller, config, and frontend UX after review
* fix(scheduled-tasks): harden run lifecycle, overlap skip, non_interactive gating, and DST conversion after review
- defer a once task's terminal status to the run-completion hook; the task
stays running until the real outcome, and a startup sweep cancels once
tasks orphaned by a crash (launch-time 'completed' could stick forever)
- record interrupted runs as a distinct 'interrupted' run status with a
readable message; an interrupted once task ends 'cancelled', not 'failed'
- enforce overlap_policy=skip for fresh_thread_per_run via an active-run
pre-check (same-thread ConflictError can never fire across fresh threads)
- protect terminal run statuses from the late launch-path 'running' write
- honor context.non_interactive only for internally-authenticated callers;
arbitrary clients can no longer strip ask_clarification
- fix DST-stale timezone offset in zonedLocalToUtcIso by re-deriving the
offset at the resolved instant (once tasks fired an hour late around
spring-forward and the create->edit round-trip diverged)
- drop dead ScheduledTaskRunRepository.update_by_run_id; share one Gateway
API error helper between channels and scheduled-tasks frontends
* fix(scheduled-tasks): close review round-3 gaps in guards, concurrency, and API ergonomics
- scrub internal-only context keys (non_interactive) from the assembled run
config for non-internal callers: gating body.context alone left the same
key smuggle-able through the free-form body.config copied verbatim by
build_run_config
- guard update_after_launch with protect_terminal so the launch bookkeeping
write cannot clobber a once task already finalized by a fast-failing run's
completion hook (parent-row sibling of the run-row guard)
- reject a manual trigger while the task has an active run (409) instead of
launching a duplicate concurrent run on fresh_thread_per_run
- re-arm a terminal once task to enabled when PATCH pushes run_at into the
future; previously the endpoint returned 200 with a next_run_at that could
never be claimed
- make max_concurrent_runs a real global cap: each poll claims only into the
remaining budget of active (queued/running) scheduled runs
- paginate GET /scheduled-tasks/{id}/runs (limit<=200, offset) and push the
thread filter of /threads/{id}/scheduled-tasks into SQL
- stamp context.user_id on scheduler-launched runs, matching IM channels, so
user-scoped guardrail providers see the owning user
---------
Co-authored-by: Willem Jiang <willem.jiang@gmail.com>
165 lines
4.2 KiB
TypeScript
165 lines
4.2 KiB
TypeScript
import { throwGatewayApiError } from "@/core/api/errors";
|
|
import { fetch } from "@/core/api/fetcher";
|
|
import { getBackendBaseURL } from "@/core/config";
|
|
|
|
import type { ScheduledTask, ScheduledTaskRun } from "./types";
|
|
|
|
function scheduledTasksUrl(path: string): string {
|
|
return `${getBackendBaseURL()}/api/scheduled-tasks${path}`;
|
|
}
|
|
|
|
export async function fetchScheduledTasks(): Promise<ScheduledTask[]> {
|
|
const response = await fetch(scheduledTasksUrl(""));
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to load scheduled tasks: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function fetchThreadScheduledTasks(
|
|
threadId: string,
|
|
): Promise<ScheduledTask[]> {
|
|
const response = await fetch(
|
|
`${getBackendBaseURL()}/api/threads/${encodeURIComponent(threadId)}/scheduled-tasks`,
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to load thread scheduled tasks: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function fetchScheduledTaskRuns(
|
|
taskId: string,
|
|
): Promise<ScheduledTaskRun[]> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}/runs`),
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to load scheduled task runs: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export type ScheduledTaskPayload = {
|
|
context_mode: "fresh_thread_per_run" | "reuse_thread";
|
|
thread_id?: string | null;
|
|
title: string;
|
|
prompt: string;
|
|
schedule_type: "once" | "cron";
|
|
schedule_spec: Record<string, unknown>;
|
|
timezone: string;
|
|
};
|
|
|
|
export async function createScheduledTask(
|
|
payload: ScheduledTaskPayload,
|
|
): Promise<ScheduledTask> {
|
|
const response = await fetch(scheduledTasksUrl(""), {
|
|
method: "POST",
|
|
headers: { "Content-Type": "application/json" },
|
|
body: JSON.stringify(payload),
|
|
});
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to create scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function updateScheduledTask(
|
|
taskId: string,
|
|
payload: Partial<Omit<ScheduledTaskPayload, "thread_id" | "schedule_type">>,
|
|
): Promise<ScheduledTask> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}`),
|
|
{
|
|
method: "PATCH",
|
|
headers: { "Content-Type": "application/json" },
|
|
body: JSON.stringify(payload),
|
|
},
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to update scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function pauseScheduledTask(
|
|
taskId: string,
|
|
): Promise<ScheduledTask> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}/pause`),
|
|
{ method: "POST" },
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to pause scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function resumeScheduledTask(
|
|
taskId: string,
|
|
): Promise<ScheduledTask> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}/resume`),
|
|
{ method: "POST" },
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to resume scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function triggerScheduledTask(
|
|
taskId: string,
|
|
): Promise<{ id: string; triggered: boolean }> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}/trigger`),
|
|
{ method: "POST" },
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to trigger scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|
|
|
|
export async function deleteScheduledTask(
|
|
taskId: string,
|
|
): Promise<{ id: string; deleted: boolean }> {
|
|
const response = await fetch(
|
|
scheduledTasksUrl(`/${encodeURIComponent(taskId)}`),
|
|
{
|
|
method: "DELETE",
|
|
},
|
|
);
|
|
if (!response.ok) {
|
|
await throwGatewayApiError(
|
|
response,
|
|
`Failed to delete scheduled task: ${response.statusText}`,
|
|
);
|
|
}
|
|
return response.json();
|
|
}
|