Commit c02bcd9
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 9a2926d commit c02bcd9
13 files changed
Lines changed: 2322 additions & 1 deletion
File tree
- apps/dev-playground
- client/src
- routes
- server
- packages/appkit-ui/src/js/sse
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
5 | 5 | | |
6 | 6 | | |
7 | 7 | | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
8 | 11 | | |
9 | 12 | | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
| 23 | + | |
23 | 24 | | |
24 | 25 | | |
25 | 26 | | |
| |||
81 | 82 | | |
82 | 83 | | |
83 | 84 | | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
84 | 90 | | |
85 | 91 | | |
86 | 92 | | |
| |||
113 | 119 | | |
114 | 120 | | |
115 | 121 | | |
| 122 | + | |
116 | 123 | | |
117 | 124 | | |
118 | 125 | | |
| |||
131 | 138 | | |
132 | 139 | | |
133 | 140 | | |
| 141 | + | |
134 | 142 | | |
135 | 143 | | |
136 | 144 | | |
| |||
150 | 158 | | |
151 | 159 | | |
152 | 160 | | |
| 161 | + | |
153 | 162 | | |
154 | 163 | | |
155 | 164 | | |
| |||
170 | 179 | | |
171 | 180 | | |
172 | 181 | | |
| 182 | + | |
173 | 183 | | |
174 | 184 | | |
175 | 185 | | |
| |||
188 | 198 | | |
189 | 199 | | |
190 | 200 | | |
| 201 | + | |
191 | 202 | | |
192 | 203 | | |
193 | 204 | | |
| |||
206 | 217 | | |
207 | 218 | | |
208 | 219 | | |
| 220 | + | |
209 | 221 | | |
210 | 222 | | |
211 | 223 | | |
| |||
225 | 237 | | |
226 | 238 | | |
227 | 239 | | |
| 240 | + | |
228 | 241 | | |
229 | 242 | | |
230 | 243 | | |
| |||
317 | 330 | | |
318 | 331 | | |
319 | 332 | | |
| 333 | + | |
| 334 | + | |
| 335 | + | |
| 336 | + | |
| 337 | + | |
| 338 | + | |
| 339 | + | |
320 | 340 | | |
321 | 341 | | |
322 | 342 | | |
| |||
361 | 381 | | |
362 | 382 | | |
363 | 383 | | |
| 384 | + | |
364 | 385 | | |
365 | 386 | | |
366 | 387 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
64 | 64 | | |
65 | 65 | | |
66 | 66 | | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
67 | 75 | | |
68 | 76 | | |
69 | 77 | | |
| |||
0 commit comments