Skip to content

Commit 0c54eeb

Browse files
committed
feat(dev-playground): durable-task example + typed SSE client helper
The reference implementation for plugin authors. A demo plugin covering both TaskFlow recovery patterns (manual via `ctx.previousEvents` and structural via `step()`), a typed frontend SSE consumer (`subscribeToTaskflowTask<TEvents>`), and the `connectSSE` parser extension that captures `event:` field names. Bottom-up teaching: server plugin → bridge → client helper → React UI. Server demo plugin (`apps/dev-playground/server/durable-task-example-plugin.ts`): - `count-to-n` task — manual recovery via `ctx.previousEvents`. Ticks once per `sleepMs`, emitting typed `tick` events. On recovery, scans the event log to find the last persisted tick and resumes from there. The pattern for "checkpoint is the last time I emitted X" with no expensive computation to memoize. - `pipeline-with-steps` task — automatic recovery via `step()`. Wraps each stage (extract → transform → load) with `step()`, which memoizes its result in the WAL the first time it runs. On recovery, completed stages return the cached value without re-executing. The pattern for stages that are expensive (LLM calls, large queries) and unsafe to replay. - Routes (mounted under `/api/durable-example`): - `POST /run`, `POST /run-pipeline` — start + bridge SSE via `executeTask`. - `POST /crash/:id` — `simulateCrash` (gated behind `NODE_ENV !== "production"`). - `POST /stop/:id` — cooperative `taskflow.stop({ reason })`. - `POST /nudge-recovery` — re-submits the original input so the same IK triggers stale-Running recovery (`engine.resume()` only applies to Suspended tasks, so the demo "nudges" the engine). - `GET /reattach/:id` — bridges an SSE stream onto an existing task by IK via `subscribe()` + `setupSseHeaders` + `writeSseFrame` directly (the `executeTask` path would derive a new IK). Performs an OBO ownership check via `asUser(req).reconnect(id, userId)` before subscribing (F11 fix). - Registers with `enableTestMode: true` so `simulateCrash` is available; the route handler additionally gates on `NODE_ENV` so a misconfigured production deployment can't crash live tasks. `apps/dev-playground/server/index.ts`: - Registers the demo plugin and enables test mode on the TaskFlow config (`taskflow: { engine: { enableTestMode: true } }`) so `simulateCrash` is callable from the demo route. Client React route (`apps/dev-playground/client/src/routes/durable-task.route.tsx`): - Exercises both tasks end-to-end: `POST /run` then opens an SSE stream via `subscribeToTaskflowTask<CountEvents>`. Renders `tick` / `recovered` events for `count-to-n`; renders `stage_started` / `stage_done` / `recovered` for `pipeline-with-steps` (which surfaces "from cache" on recovered stages). - Buttons for Stop, Crash, Nudge, and Reattach exercise the full cancellation / crash / recovery / re-attach loop. - Adds nav entries in `__root.tsx`, `index.tsx`, and the TanStack-generated `routeTree.gen.ts`. Typed client helper (`packages/appkit-ui/src/js/sse/subscribe-taskflow-task.ts`): - `subscribeToTaskflowTask<TEvents>(url, { onEvent, onComplete, onError, signal? })` — typed async API consuming the AppKit SSE bridge. Each `event: <name>` frame is dispatched to `onEvent[name]` with `payload` typed as `TEvents[name]`. - Terminal events (`completed`, `failed`, `cancelled`) resolve / reject the returned promise so plugins can `await` the durable run without an event handler. - `Last-Event-ID` reconnection: the helper tracks the highest seen `id:` frame and reattaches with that header on transient network failure. Tests assert the reconnect math is correct. - Includes tests for happy-path streaming, terminal events, abort via `AbortSignal`, and Last-Event-ID reconnect. `connect-sse` extension (`packages/appkit-ui/src/js/sse/connect-sse.ts`, `types.ts`, `index.ts`): - The generic SSE parser captures `event:` field names alongside `data:` payloads. `SSEMessage` gains `event?: string` so any AppKit SSE consumer can inspect the event name without re-parsing. Tests cover multi-line `data:` joining, CRLF normalisation, and comment-frame handling. - Export the new typed helper from `index.ts`. Gitignore: - `apps/dev-playground/.gitignore` adds `tasks.*` / `*.wal` patterns as a defensive belt-and-braces around the existing `.appkit/` exclusion. The demo plugin may configure storage at the playground root for diagnostics; the additional patterns keep `tasks.db` and the rotating WAL out of git regardless of `databasePath`. Verify: - `pnpm -r typecheck`, `pnpm build`, `pnpm test` (125 files, 2304 tests) all green. - `pnpm exec biome check` clean on touched files. - `pnpm exec knip` clean. Risk. Demo plugin is unauthenticated by design (it ships with the dev playground, not the SDK). `/crash/:id` returns 404 in production via `NODE_ENV` gate; `enableTestMode` flips on `simulateCrash`. The demo route handlers do not enforce auth — they assume the playground sits behind the Databricks Apps proxy. Document in deployment notes. Not in this PR. No production-plugin changes. No doc rewrite — that's PR 7. The `subscribeToTaskflowTask` helper currently requires plugin authors to redeclare `TEvents` client-side; a future follow-up (F26) would derive it from the registered `TaskHandle`. Stacked on: stack/taskflow/analytics-migration. Signed-off-by: ditadi <victordperd@gmail.com>
1 parent 561da92 commit 0c54eeb

11 files changed

Lines changed: 2309 additions & 1 deletion

File tree

apps/dev-playground/client/src/lib/nav.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,13 @@ export const NAV_GROUPS: ReadonlyArray<NavGroup> = [
162162
"Resilient SSE streams: automatic Last-Event-ID tracking and reconnection.",
163163
icon: RadioIcon,
164164
},
165+
{
166+
to: "/durable-task",
167+
label: "Durable Task",
168+
description:
169+
"Crash-safe long-running tasks with typed SSE streaming and resume.",
170+
icon: ZapIcon,
171+
},
165172
],
166173
},
167174
];

apps/dev-playground/client/src/routeTree.gen.ts

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import { Route as LakebaseRouteRouteImport } from './routes/lakebase.route'
2121
import { Route as JobsRouteRouteImport } from './routes/jobs.route'
2222
import { Route as GenieRouteRouteImport } from './routes/genie.route'
2323
import { Route as FilesRouteRouteImport } from './routes/files.route'
24+
import { Route as DurableTaskRouteRouteImport } from './routes/durable-task.route'
2425
import { Route as DataVisualizationRouteRouteImport } from './routes/data-visualization.route'
2526
import { Route as ChartInferenceRouteRouteImport } from './routes/chart-inference.route'
2627
import { Route as ArrowAnalyticsRouteRouteImport } from './routes/arrow-analytics.route'
@@ -88,6 +89,11 @@ const FilesRouteRoute = FilesRouteRouteImport.update({
8889
path: '/files',
8990
getParentRoute: () => rootRouteImport,
9091
} as any)
92+
const DurableTaskRouteRoute = DurableTaskRouteRouteImport.update({
93+
id: '/durable-task',
94+
path: '/durable-task',
95+
getParentRoute: () => rootRouteImport,
96+
} as any)
9197
const DataVisualizationRouteRoute = DataVisualizationRouteRouteImport.update({
9298
id: '/data-visualization',
9399
path: '/data-visualization',
@@ -126,6 +132,7 @@ export interface FileRoutesByFullPath {
126132
'/arrow-analytics': typeof ArrowAnalyticsRouteRoute
127133
'/chart-inference': typeof ChartInferenceRouteRoute
128134
'/data-visualization': typeof DataVisualizationRouteRoute
135+
'/durable-task': typeof DurableTaskRouteRoute
129136
'/files': typeof FilesRouteRoute
130137
'/genie': typeof GenieRouteRoute
131138
'/jobs': typeof JobsRouteRoute
@@ -146,6 +153,7 @@ export interface FileRoutesByTo {
146153
'/arrow-analytics': typeof ArrowAnalyticsRouteRoute
147154
'/chart-inference': typeof ChartInferenceRouteRoute
148155
'/data-visualization': typeof DataVisualizationRouteRoute
156+
'/durable-task': typeof DurableTaskRouteRoute
149157
'/files': typeof FilesRouteRoute
150158
'/genie': typeof GenieRouteRoute
151159
'/jobs': typeof JobsRouteRoute
@@ -167,6 +175,7 @@ export interface FileRoutesById {
167175
'/arrow-analytics': typeof ArrowAnalyticsRouteRoute
168176
'/chart-inference': typeof ChartInferenceRouteRoute
169177
'/data-visualization': typeof DataVisualizationRouteRoute
178+
'/durable-task': typeof DurableTaskRouteRoute
170179
'/files': typeof FilesRouteRoute
171180
'/genie': typeof GenieRouteRoute
172181
'/jobs': typeof JobsRouteRoute
@@ -189,6 +198,7 @@ export interface FileRouteTypes {
189198
| '/arrow-analytics'
190199
| '/chart-inference'
191200
| '/data-visualization'
201+
| '/durable-task'
192202
| '/files'
193203
| '/genie'
194204
| '/jobs'
@@ -209,6 +219,7 @@ export interface FileRouteTypes {
209219
| '/arrow-analytics'
210220
| '/chart-inference'
211221
| '/data-visualization'
222+
| '/durable-task'
212223
| '/files'
213224
| '/genie'
214225
| '/jobs'
@@ -229,6 +240,7 @@ export interface FileRouteTypes {
229240
| '/arrow-analytics'
230241
| '/chart-inference'
231242
| '/data-visualization'
243+
| '/durable-task'
232244
| '/files'
233245
| '/genie'
234246
| '/jobs'
@@ -250,6 +262,7 @@ export interface RootRouteChildren {
250262
ArrowAnalyticsRouteRoute: typeof ArrowAnalyticsRouteRoute
251263
ChartInferenceRouteRoute: typeof ChartInferenceRouteRoute
252264
DataVisualizationRouteRoute: typeof DataVisualizationRouteRoute
265+
DurableTaskRouteRoute: typeof DurableTaskRouteRoute
253266
FilesRouteRoute: typeof FilesRouteRoute
254267
GenieRouteRoute: typeof GenieRouteRoute
255268
JobsRouteRoute: typeof JobsRouteRoute
@@ -350,6 +363,13 @@ declare module '@tanstack/react-router' {
350363
preLoaderRoute: typeof FilesRouteRouteImport
351364
parentRoute: typeof rootRouteImport
352365
}
366+
'/durable-task': {
367+
id: '/durable-task'
368+
path: '/durable-task'
369+
fullPath: '/durable-task'
370+
preLoaderRoute: typeof DurableTaskRouteRouteImport
371+
parentRoute: typeof rootRouteImport
372+
}
353373
'/data-visualization': {
354374
id: '/data-visualization'
355375
path: '/data-visualization'
@@ -402,6 +422,7 @@ const rootRouteChildren: RootRouteChildren = {
402422
ArrowAnalyticsRouteRoute: ArrowAnalyticsRouteRoute,
403423
ChartInferenceRouteRoute: ChartInferenceRouteRoute,
404424
DataVisualizationRouteRoute: DataVisualizationRouteRoute,
425+
DurableTaskRouteRoute: DurableTaskRouteRoute,
405426
FilesRouteRoute: FilesRouteRoute,
406427
GenieRouteRoute: GenieRouteRoute,
407428
JobsRouteRoute: JobsRouteRoute,

0 commit comments

Comments
 (0)