From 15a611ded5d139eabaed1f35313f0fafe874ebd6 Mon Sep 17 00:00:00 2001 From: rahim Date: Mon, 12 Jan 2026 22:36:36 +1100 Subject: [PATCH] refactor(store): queue simplification (#302) --- .claude/plans/store-queue-simplification.md | 627 +++++++ packages/store/README.md | 172 +- packages/store/package.json | 4 - packages/store/src/core/errors.ts | 2 - packages/store/src/core/index.ts | 1 + packages/store/src/core/queue.ts | 270 +-- packages/store/src/core/request.ts | 5 +- packages/store/src/core/store.ts | 3 +- packages/store/src/core/task.ts | 77 + .../src/core/tests/integration/store.test.ts | 19 +- packages/store/src/core/tests/queue.test.ts | 1628 +++++++---------- .../store/src/core/tests/queue.types.test.ts | 71 +- packages/store/src/core/tests/slice.test.ts | 6 - .../store/src/core/tests/store.types.test.ts | 38 - packages/store/src/core/tests/task.test.ts | 240 +++ .../store/src/core/tests/task.types.test.ts | 76 + packages/store/src/dom/index.ts | 1 - packages/store/src/dom/schedulers.ts | 27 - .../store/src/dom/tests/schedulers.test.ts | 176 -- packages/store/src/dom/tsconfig.json | 10 - .../lit/controllers/mutation-controller.ts | 2 +- .../lit/controllers/optimistic-controller.ts | 2 +- .../store/src/react/hooks/use-mutation.ts | 2 +- .../store/src/react/hooks/use-optimistic.ts | 2 +- packages/store/tsdown.config.ts | 1 - packages/store/vitest.config.ts | 8 - tsconfig.json | 1 - 27 files changed, 1754 insertions(+), 1717 deletions(-) create mode 100644 .claude/plans/store-queue-simplification.md create mode 100644 packages/store/src/core/task.ts create mode 100644 packages/store/src/core/tests/task.test.ts create mode 100644 packages/store/src/core/tests/task.types.test.ts delete mode 100644 packages/store/src/dom/index.ts delete mode 100644 packages/store/src/dom/schedulers.ts delete mode 100644 packages/store/src/dom/tests/schedulers.test.ts delete mode 100644 packages/store/src/dom/tsconfig.json diff --git a/.claude/plans/store-queue-simplification.md b/.claude/plans/store-queue-simplification.md new file mode 100644 index 00000000..40b9064b --- /dev/null +++ b/.claude/plans/store-queue-simplification.md @@ -0,0 +1,627 @@ +# Store Queue Simplification Plan + +> **STATUS: COMPLETED** +> +> Final result: ~405 LOC total (328 queue.ts + 77 task.ts) down from ~524 LOC. +> Task types and guards extracted to separate `task.ts` module for better organization. + +## Overview + +Simplify the queue from ~524 LOC to ~200 LOC while preserving core value. + +**Goal:** Remove convenience features only used in tests/docs, keep what's essential for real-world async media operations. + +--- + +## Why the Queue Exists + +The store design: **all writes are async requests**. This isn't optional — it's the architecture. + +```ts +// Read path: sync state from target +store.state.paused; // Synced from video.paused + +// Write path: async requests to target +await store.request.play(); // Returns Promise, always +``` + +Media operations are inherently async and can fail: + +| Scenario | What Can Happen | +| ---------------------- | --------------------------------------------------- | +| **Chromecast/AirPlay** | 100-500ms+ network round-trip, device disconnect | +| **Source loading** | 1-5s+ load time, 404, DRM errors, codec unsupported | +| **Quality switching** | Segment fetching, ABR delays | +| **Network conditions** | Timeouts, flaky connections | +| **User behavior** | Impatient clicks, changing mind mid-operation | + +The queue handles these realities uniformly. + +--- + +## Real-World Scenarios + +### Chromecast Play + +```ts +// Command goes over network to Cast device +await store.request.play(); +``` + +**Without queue:** + +- How does UI show "Connecting..."? +- User taps 5 times impatiently → 5 network messages? +- Device disconnects mid-request → how to surface error? + +**With queue:** + +- `tasks['play'].status === 'pending'` → show loading +- Supersession → only 1 message sent +- `tasks['play'].status === 'error'` → show failure + +### Source Switching + +```ts +// User browses playlist quickly +store.request.setSource('a.mp4'); // Starts loading +store.request.setSource('b.mp4'); // User changed mind +store.request.setSource('c.mp4'); // Final choice +``` + +**Without queue:** + +- All three sources load simultaneously +- Wasted bandwidth, race conditions +- Which promise resolves? Which errors? + +**With queue:** + +- Supersession aborts a.mp4 and b.mp4 +- Only c.mp4 loads +- Clean promise semantics (superseded reject with SUPERSEDED) + +### Network Failure + +```ts +// Quality change needs segment fetch +await store.request.setQuality('1080p'); +// Network times out +``` + +**Without queue:** + +- Component needs try/catch + useState for error +- Every async operation repeats this boilerplate + +**With queue:** + +- `tasks['setQuality'].status === 'error'` +- `tasks['setQuality'].error` contains details +- `reset('setQuality')` + retry + +### Component Implementation Comparison + +**Without queue (manual):** + +```tsx +function SourceSelector() { + const [loading, setLoading] = useState(false); + const [error, setError] = useState(null); + + const handleChange = async (src) => { + setLoading(true); + setError(null); + try { + await loadSource(src); + } catch (e) { + setError(e); + } finally { + setLoading(false); + } + }; + + return ( + <> + source.mutate(e.target.value)} /> + {source.status === 'pending' && } + {source.status === 'error' && } + + ); +} +``` + +--- + +## Required Capabilities → Queue Features + +| Scenario | Required Capability | Queue Feature | +| ------------------------ | ---------------------- | ---------------------------------- | +| Cast/AirPlay round-trip | Loading indicator | `tasks[name].status === 'pending'` | +| Network/DRM failures | Error display | `tasks[name].status === 'error'` | +| User retry after failure | Error recovery | `reset()` + re-request | +| Impatient clicks | Prevent duplicate work | Supersession (same key) | +| User changes mind | Cancel in-flight | `abort()` + AbortSignal | +| Multiple operation types | Parallel execution | Key-based coordination | +| Async completion | Await result | Promise API | +| React/Lit integration | Observe changes | `subscribe()` | +| Debug production issues | Task visibility | `tasks` record | + +--- + +## What We Keep (And Why) + +### 1. Microtask Batching + +**Required for supersession to work.** + +```ts +// User drags volume slider rapidly +store.request.setVolume(0.3); +store.request.setVolume(0.5); +store.request.setVolume(0.7); +// Without batching: all three handlers execute (race condition) +// With batching: only 0.7 executes +``` + +Without batching, `target.volume = 0.3` runs before we know 0.5 is coming. + +### 2. Supersession (Same Key Cancels Previous) + +**Prevents race conditions on every slider, toggle, and rapid interaction.** + +```ts +// Seek bar scrubbing +store.request.seek(10); // User still dragging... +store.request.seek(15); // Previous aborted +store.request.seek(20); // Only this completes +``` + +```ts +// Play/pause rapid toggle +store.request.play(); // Queued +store.request.pause(); // play() superseded +store.request.play(); // pause() superseded — only final executes +``` + +Media Chrome lacks this — users implement manually or suffer bugs. + +### 3. AbortController Propagation + +**Graceful cancellation for long-running operations.** + +```ts +async setSource(src, { signal }) { + target.src = src; + await onEvent(target, 'loadeddata', { signal }); // Aborted if new source +} + +// User changes quality mid-load +store.request.setSource('720p.mp4'); // Loading... +store.request.setSource('1080p.mp4'); // 720p aborted, 1080p starts +``` + +- Prevents wasted bandwidth +- Prevents stale responses updating state +- Handlers can clean up resources on abort + +### 4. Lifecycle Tracking (pending/success/error) + +**Observability for UI and debugging.** + +```tsx +// Loading state +) => Promise; + } +``` + +### 3. Simplified Constructor + +```diff +- constructor(config: QueueConfig = {}) { +- this.#scheduler = config.scheduler ?? microtask; +- this.#onDispatch = tryCatch(config.onDispatch, logError); +- this.#onSettled = tryCatch(config.onSettled, logError); +- } ++ constructor() {} +``` + +--- + +## Final API Surface + +```ts +interface Queue { + // State + readonly tasks: TasksRecord; + readonly destroyed: boolean; + + // Core + enqueue(task: QueueTask): Promise; + + // Control + abort(name?: keyof Tasks): void; + reset(name?: keyof Tasks): void; + destroy(): void; + + // Observe + subscribe(listener: QueueListener): () => void; +} +``` + +--- + +## Migration Guide + +| Before | After | +| ------------------------- | ------------------------------------------- | +| `queue.isPending('play')` | `queue.tasks['play']?.status === 'pending'` | +| `queue.isSettled('play')` | `queue.tasks['play']?.status !== 'pending'` | +| `queue.isQueued('play')` | Remove (instant with microtask) | +| `queue.cancel('play')` | `queue.abort('play')` | +| `queue.flush('play')` | Remove (no configurable scheduling) | +| `queue.queued` | Remove (internal detail) | +| `schedule: delay(100)` | External: `debounce(request, 100)` | +| `onSettled: fn` | `queue.subscribe(tasks => ...)` | +| `createQueue({ ... })` | `createQueue()` | + +--- + +## Files to Change + +1. **`packages/store/src/core/queue.ts`** — Main simplification +2. **`packages/store/src/core/tests/queue.test.ts`** — Remove tests for removed features +3. **`packages/store/src/core/tests/queue.types.test.ts`** — Update type tests +4. **`packages/store/src/dom/schedulers.ts`** — Delete +5. **`packages/store/src/dom/index.ts`** — Delete (empty after scheduler removal) +6. **`packages/store/src/dom/tests/schedulers.test.ts`** — Delete +7. **`packages/store/src/core/index.ts`** — Remove scheduler exports +8. **`packages/store/README.md`** — Update docs + +--- + +## Type Changes + +Keep strong typing. Simplifications: + +```diff +- export type TaskScheduler = (flush: () => void) => (() => void) | void; + +- export interface QueueConfig { +- scheduler?: TaskScheduler; +- onDispatch?: ...; +- onSettled?: ...; +- } + +- export interface QueuedTaskId { ... } +- export type PublicQueuedRecord = { ... } +``` + +Keep all Task types (`PendingTask`, `SuccessTask`, `ErrorTask`, etc.) — used by hooks. + +--- + +## Estimated LOC + +| Section | Current | After | Actual | +| ----------- | -------- | -------- | -------- | +| Types | ~120 | ~80 | — | +| Schedulers | ~30 | 0 | 0 | +| Queue class | ~350 | ~120 | — | +| Factory | ~10 | ~5 | — | +| **Total** | **~524** | **~200** | **~405** | + +> **Note:** Final LOC higher than estimate because we kept more robust error handling, +> comprehensive type exports, and extracted task types to a separate module for better organization. +> The simplification goals were achieved — removed ~120 LOC of unnecessary features. + +--- + +## Verification + +After implementation: + +1. `pnpm -F @videojs/store test` — All 310 tests pass +2. `pnpm -F @videojs/store build` — Builds successfully +3. `pnpm typecheck` — No type errors +4. Hooks (`useMutation`, `useOptimistic`, `useTasks`) — Still work +5. Examples — Still work + +--- + +## Implementation Summary + +### Files Changed + +1. **`packages/store/src/core/queue.ts`** — Main simplification (queue-only, no task re-exports) +2. **`packages/store/src/core/task.ts`** — NEW: Extracted task types and type guards +3. **`packages/store/src/core/index.ts`** — Added `task.ts` export +4. **`packages/store/src/core/request.ts`** — Removed `schedule` field, import `TaskKey` from `task.ts` +5. **`packages/store/src/core/store.ts`** — Removed `schedule` from enqueue, import task types from `task.ts` +6. **`packages/store/src/core/errors.ts`** — Deprecated `REMOVED` error code +7. **`packages/store/src/core/tests/queue.test.ts`** — Removed tests for removed features +8. **`packages/store/src/core/tests/queue.types.test.ts`** — Updated type tests +9. **`packages/store/src/core/tests/task.test.ts`** — NEW: Task type guard tests +10. **`packages/store/src/core/tests/task.types.test.ts`** — NEW: Task type-level tests +11. **`packages/store/src/dom/`** — DELETED entire directory +12. **`packages/store/src/lit/controllers/*.ts`** — Updated imports (Task from `task.ts`) +13. **`packages/store/src/react/hooks/*.ts`** — Updated imports (Task from `task.ts`) +14. **`packages/store/README.md`** — Updated documentation +15. **`packages/store/package.json`** — Removed `./dom` export +16. **`packages/store/tsdown.config.ts`** — Removed `dom` entry +17. **`packages/store/vitest.config.ts`** — Removed `store/dom` test project +18. **`tsconfig.json`** (root) — Removed `packages/store/src/dom` reference + +### Features Removed + +- `cancel()` method +- `flush()` method +- `queued` getter +- `TaskScheduler` type +- `QueueConfig` interface +- `delay()`, `microtask` exports +- `schedule` param on `QueueTask` +- `onDispatch`, `onSettled` hooks +- DOM schedulers (`raf()`, `idle()`) + +--- + +## Bundle Size Analysis + +Tested with esbuild (minified + gzipped): + +| Scenario | Minified | Gzipped | +| --------------------------------------- | -------- | ---------- | +| **Minimal** (createStore + createSlice) | 6.2 KB | **2.3 KB** | +| **+ requests** (using queue) | 6.3 KB | 2.3 KB | +| **+ task guards** (isPendingTask, etc.) | 6.5 KB | 2.4 KB | +| **Full entry** (all exports) | 8.1 KB | 3.0 KB | + +### Tree-Shaking Results + +| Export | Tree-Shakes? | Notes | +| ------------------------------------------------------------------------------ | ------------ | ---------------------------------------- | +| Task guards (`isPendingTask`, `isSuccessTask`, `isErrorTask`, `isSettledTask`) | ✅ Yes | Removed when unused | +| `Queue` class | ❌ No | Store creates one internally (by design) | +| `StoreError` | ❌ No | Used by Queue for error handling | +| `State` class | ❌ No | Core dependency of Store | + +### What's Always Included + +``` +Queue class: ~159 lines +Store class: ~193 lines +State class: ~59 lines +StoreError: ~8 lines +Utils (@videojs/utils): isFunction, isNull, isUndefined, isObject, getSelectorKeys +``` + +### Notes + +- Queue is always present because the store architecture requires it — all writes are async requests +- Task types extracted to `task.ts` tree-shake properly when guards aren't used +- Savings of ~1.9 KB raw / ~700 bytes gzip when not importing everything diff --git a/packages/store/README.md b/packages/store/README.md index 95ce00a4..04583af5 100644 --- a/packages/store/README.md +++ b/packages/store/README.md @@ -319,17 +319,19 @@ request: { When a new request arrives with the same key: -- Queued request with that key is dropped +- Queued request with that key is superseded - Executing request with that key is aborted - New request takes over ```ts store.request.play(); // queued -store.request.pause(); // play dropped, pause queued -store.request.play(); // pause dropped, play queued +store.request.pause(); // play superseded, pause queued +store.request.play(); // pause superseded, play queued // only final play() executes ``` +This supersession happens automatically via microtask batching—all requests in the same synchronous batch are collected, and only the last request per key actually executes. + Dynamic keys for parallel execution: ```ts @@ -351,7 +353,7 @@ request: { ### Cancels Requests can cancel other in-flight requests by name. Cancellation happens immediately when the -request is enqueued, before guards or scheduling. +request is enqueued, before guards. ```ts request: { @@ -362,51 +364,6 @@ request: { } ``` -### Schedule - -Schedule controls _when_ a request executes. The schedule function receives a `flush` callback and -optionally returns a cancel function. Default schedule is microtask (executes at end of current tick). - -```ts -import { delay, microtask } from '@videojs/store'; -import { raf, idle } from '@videojs/store/dom'; - -request: { - // Microtask - default, executes at end of current tick - setVolume: { - schedule: microtask, - handler: (volume, { target }) => { target.volume = volume; }, - }, - - // Debounce 100ms - good for sliders - seek: { - schedule: delay(100), - handler: (time, { target }) => { target.currentTime = time; }, - }, - - // Sync with animation frame - updateOverlay: { - schedule: raf(), - handler: () => { ... }, - }, - - // Execute when browser is idle - preloadNext: { - schedule: idle(), - handler: () => { ... }, - }, - - // Custom schedule - custom: { - schedule: (flush) => { - const id = setTimeout(flush, 200); - return () => clearTimeout(id); // cancel function - }, - handler: () => { ... }, - }, -} -``` - ### Guards Guards gate request execution. A guard returns a `GuardResult`: @@ -478,7 +435,6 @@ All store errors include a `code` for programmatic handling: | `DETACHED` | Target detached | | `NO_TARGET` | No target attached | | `REJECTED` | Guard returned falsy | -| `REMOVED` | Task dequeued or cleared | | `SUPERSEDED` | Replaced by same-key request | | `TIMEOUT` | Guard timed out | @@ -523,51 +479,15 @@ try { ## Queue -A default queue is created automatically. Provide a custom queue for lifecycle hooks or custom -scheduling: - -```ts -import { createQueue } from '@videojs/store'; - -const store = createStore({ - slices: [ - /* ... */ - ], - queue: createQueue({ - // Default scheduler for requests without schedule - scheduler: (flush) => queueMicrotask(flush), - - // Lifecycle hooks - onDispatch: (request) => { - console.log('Started:', request.name); - }, - - onSettled: (task) => { - analytics.track(task.name, { - status: task.status, - duration: task.settledAt - task.startedAt, - }); - }, - }), -}); -``` +The queue manages request execution with automatic supersession and lifecycle tracking. A default queue is created automatically with the store. ### Queue API ```ts -const queue = store.queue; // accessed on the store +const queue = store.queue; -queue.queued; // tasks waiting to execute -queue.tasks; // task lifecycle map (pending/success/error) keyed by request name - -// Check task status by request name -queue.isPending('play'); // true if currently executing -queue.isQueued('seek'); // true if waiting to execute -queue.isSettled('seek'); // true if completed (success or error) - -// Cancel queued tasks (waiting to execute) -queue.cancel('seek'); // cancel specific request -queue.cancel(); // cancel all queued +// Task lifecycle map (pending/success/error) keyed by request name +queue.tasks; // Abort tasks (queued + executing) queue.abort('play'); // abort specific request @@ -577,9 +497,41 @@ queue.abort(); // abort all queue.reset('seek'); // clear specific request queue.reset(); // clear all settled -// Execute queued tasks immediately -queue.flush(); // execute all queued now -queue.flush('play'); // execute specific request now +// Subscribe to task changes +queue.subscribe((tasks) => { + const playTask = tasks.play; + if (playTask?.status === 'pending') { + console.log('Play in progress...'); + } +}); +``` + +### Task Lifecycle + +Each request creates a task that transitions through states: + +```ts +import { isErrorTask, isPendingTask, isSettledTask, isSuccessTask } from '@videojs/store'; + +const task = queue.tasks.play; + +// Type guards for status checking +if (isPendingTask(task)) { + console.log('In progress, started at:', task.startedAt); +} + +if (isSettledTask(task)) { + console.log('Duration:', task.settledAt - task.startedAt); +} + +if (isSuccessTask(task)) { + console.log('Result:', task.output); +} + +if (isErrorTask(task)) { + console.log('Failed:', task.error); + console.log('Was cancelled:', task.cancelled); +} ``` ### Direct Queue Usage @@ -587,18 +539,48 @@ queue.flush('play'); // execute specific request now You can use the queue directly without a store: ```ts +import { createQueue } from '@videojs/store'; + +const queue = createQueue(); + await queue.enqueue({ name: 'myTask', key: 'task-key', input: { some: 'data' }, - schedule: delay(100), - handler: async ({ signal }) => { + handler: async ({ input, signal }) => { // do work, check signal.aborted return result; }, }); ``` +### Observing Tasks + +Use `subscribe` to react to task changes—useful for loading states and error handling: + +```ts +queue.subscribe((tasks) => { + for (const [name, task] of Object.entries(tasks)) { + if (task?.status === 'error' && !task.cancelled) { + toast.error(`${name} failed: ${task.error}`); + } + } +}); + +// Analytics +queue.subscribe((tasks) => { + for (const task of Object.values(tasks)) { + if (task && task.status !== 'pending') { + analytics.track('request', { + name: task.name, + status: task.status, + duration: task.settledAt - task.startedAt, + }); + } + } +}); +``` + ## Advanced ### Custom State diff --git a/packages/store/package.json b/packages/store/package.json index e99af2c3..4202f7be 100644 --- a/packages/store/package.json +++ b/packages/store/package.json @@ -16,10 +16,6 @@ "types": "./dist/index.d.ts", "default": "./dist/index.js" }, - "./dom": { - "types": "./dist/dom.d.ts", - "default": "./dist/dom.js" - }, "./lit": { "types": "./dist/lit.d.ts", "default": "./dist/lit.js" diff --git a/packages/store/src/core/errors.ts b/packages/store/src/core/errors.ts index 9e9cff79..6a41b115 100644 --- a/packages/store/src/core/errors.ts +++ b/packages/store/src/core/errors.ts @@ -28,8 +28,6 @@ export type StoreErrorCode | 'NO_TARGET' /** Guard condition returned falsy - request preconditions not met. */ | 'REJECTED' - /** Task was removed from queue via `dequeue()` or `clear()`. */ - | 'REMOVED' /** Request was replaced by a newer request with the same key. */ | 'SUPERSEDED' /** Guard condition timed out waiting for a truthy result. */ diff --git a/packages/store/src/core/index.ts b/packages/store/src/core/index.ts index 076fa067..6c3d7ed7 100644 --- a/packages/store/src/core/index.ts +++ b/packages/store/src/core/index.ts @@ -6,3 +6,4 @@ export * from './request'; export * from './slice'; export * from './state'; export * from './store'; +export * from './task'; diff --git a/packages/store/src/core/queue.ts b/packages/store/src/core/queue.ts index 62dbc625..b448f0ce 100644 --- a/packages/store/src/core/queue.ts +++ b/packages/store/src/core/queue.ts @@ -1,7 +1,7 @@ import type { Request, RequestMeta } from './request'; +import type { ErrorTask, PendingTask, SuccessTask, Task, TaskContext, TaskKey } from './task'; -import { tryCatch } from '@videojs/utils/function'; -import { isFunction, isUndefined } from '@videojs/utils/predicate'; +import { isUndefined } from '@videojs/utils/predicate'; import { StoreError } from './errors'; @@ -9,17 +9,6 @@ import { StoreError } from './errors'; // Types // ---------------------------------------- -export type TaskKey = T & (string | symbol); - -export type EnsureTaskKey = T extends string | symbol ? T : never; - -/** - * A task scheduler controls when a task flushes. - * - * Returns an optional cancel function. - */ -export type TaskScheduler = (flush: () => void) => (() => void) | void; - export type TaskRecord = { [K in TaskKey]: Request; }; @@ -28,56 +17,11 @@ export type DefaultTaskRecord = Record>; export type EnsureTaskRecord = T extends TaskRecord ? T : never; -export interface TaskBase { - id: symbol; - name: string; - key: Key; - input: Input; - startedAt: number; - meta: RequestMeta | null; -} - -export interface PendingTask extends TaskBase { - status: 'pending'; - abort: AbortController; -} - -export interface SuccessTask extends TaskBase< - Key, - Input -> { - status: 'success'; - settledAt: number; - output: Output; -} - -export interface ErrorTask extends TaskBase { - status: 'error'; - settledAt: number; - error: unknown; - cancelled: boolean; -} - -export type Task - = | PendingTask - | SuccessTask - | ErrorTask; - -export type SettledTask - = | SuccessTask - | ErrorTask; - -export interface TaskContext { - input: Input; - signal: AbortSignal; -} - export interface QueueTask { name: string; key: Key; input?: Input; meta?: RequestMeta | null; - schedule?: TaskScheduler | undefined; handler: (ctx: TaskContext) => Promise; } @@ -87,30 +31,12 @@ interface QueuedTask) => Promise; resolve: (value: Output) => void; reject: (error: unknown) => void; - invalidate?: () => void; } -export interface QueueConfig { - /** Default scheduler when task has no schedule */ - scheduler?: TaskScheduler; - onDispatch?: (task: PendingTask, Tasks[K]['input']>) => void; - onSettled?: (task: SettledTask, Tasks[K]['input'], Tasks[K]['output']>) => void; -} - -export interface QueuedTaskId { - key: Key; - name: string; -} - -export type PublicQueuedRecord = { - readonly [K in keyof Tasks]?: QueuedTaskId>; -}; - -export type QueuedRecord = { +type QueuedRecord = { [K in keyof Tasks]?: QueuedTask>; }; @@ -120,66 +46,17 @@ export type TasksRecord = { export type QueueListener = (tasks: TasksRecord) => void; -// ---------------------------------------- -// Schedulers -// ---------------------------------------- - -/** - * Default scheduler, delay to next microtask. - * - * @see {@link https://developer.mozilla.org/en-US/docs/Web/API/HTML_DOM_API/Microtask_guide} - */ -export const microtask: TaskScheduler = (flush) => { - let cancelled = false; - - queueMicrotask(() => { - if (!cancelled) flush(); - }); - - return () => { - cancelled = true; - }; -}; - -/** - * Delay execution by ms. Resets on each new task. - * - * @param ms - Milliseconds to delay - * @see {@link https://developer.mozilla.org/en-US/docs/Web/API/WindowOrWorkerGlobalScope/setTimeout} - */ -export function delay(ms: number): TaskScheduler { - return (flush) => { - const id = setTimeout(flush, ms); - return () => clearTimeout(id); - }; -} - // ---------------------------------------- // Implementation // ---------------------------------------- export class Queue { - readonly #scheduler: TaskScheduler; - readonly #onDispatch: QueueConfig['onDispatch']; - readonly #onSettled: QueueConfig['onSettled']; readonly #subscribers = new Set>(); #queued: QueuedRecord = {}; #tasks: TasksRecord = {}; #destroyed = false; - - constructor(config: QueueConfig = {}) { - this.#scheduler = config.scheduler ?? microtask; - - // Wrap callbacks to catch errors and prevent breaking queue/scheduler - const logError = (e: unknown) => console.error('[vjs-queue]', e); - this.#onDispatch = tryCatch(config.onDispatch, logError); - this.#onSettled = tryCatch(config.onSettled, logError); - } - - get queued(): Readonly> { - return Object.freeze({ ...this.#queued }); - } + #flushScheduled = false; get tasks(): Readonly> { return Object.freeze({ ...this.#tasks }); @@ -189,26 +66,7 @@ export class Queue { return this.#destroyed; } - isPending(name: keyof Tasks): boolean { - return this.#tasks[name]?.status === 'pending'; - } - - isQueued(name: keyof Tasks): boolean { - // Note: #queued is keyed by key, but we need to search by name - for (const task of Object.values(this.#queued)) { - if (task?.name === name) return true; - } - return false; - } - - isSettled(name: keyof Tasks): boolean { - const task = this.#tasks[name]; - return task?.status === 'success' || task?.status === 'error'; - } - - /** - * Clear settled task(s). If name provided, clears that task. If no name, clears all settled. - */ + /** Clear settled task(s). If name provided, clears that task. If no name, clears all settled. */ reset(name?: keyof Tasks): void { if (!isUndefined(name)) { const task = this.#tasks[name]; @@ -257,7 +115,7 @@ export class Queue { enqueue( task: QueueTask, Tasks[K]['input'], Tasks[K]['output']>, ): Promise { - const { name, key, input, schedule, meta = null, handler } = task; + const { name, key, input, meta = null, handler } = task; if (this.#destroyed) { return Promise.reject(new StoreError('DESTROYED')); @@ -265,7 +123,6 @@ export class Queue { // Supersede any queued task with the same key const queued = this.#queued[key]; - queued?.invalidate?.(); queued?.reject(new StoreError('SUPERSEDED')); delete this.#queued[key]; @@ -278,99 +135,51 @@ export class Queue { } return new Promise((resolve, reject) => { - const task: QueuedTask = { + const queuedTask: QueuedTask = { id: Symbol('@videojs/task'), name, key, input, meta, - schedule, handler, resolve, reject, }; - this.#queued[key as keyof Tasks] = task; - - let flushed = false; - - try { - const scheduleFlush = schedule ?? this.#scheduler; - - const safeFlush = () => { - if (flushed) return; - flushed = true; - this.#flushKey(key); - }; - - const cancel = scheduleFlush(safeFlush); - - if (!flushed && isFunction(cancel)) { - task.invalidate = cancel; - } - } catch (err) { - if (!flushed) { - delete this.#queued[key]; - } - - reject(err); - } + this.#queued[key as keyof Tasks] = queuedTask; + this.#scheduleFlush(); }); } - /** - * Cancel queued task(s). If name provided, cancels that task. If no name, cancels all. - */ - cancel(name?: keyof Tasks): boolean { - if (!isUndefined(name)) { - // Find queued task by name (#queued is keyed by key) - for (const [key, queued] of Object.entries(this.#queued)) { - if (queued?.name === name) { - queued.invalidate?.(); - queued.reject(new StoreError('REMOVED')); - delete this.#queued[key as keyof Tasks]; - return true; - } - } - return false; - } + #scheduleFlush(): void { + if (this.#flushScheduled) return; - const hadQueued = Object.keys(this.#queued).length > 0; - for (const queued of Object.values(this.#queued)) { - queued.invalidate?.(); - queued.reject(new StoreError('REMOVED')); - } - - this.#queued = {}; - - return hadQueued; + this.#flushScheduled = true; + queueMicrotask(() => { + this.#flushScheduled = false; + this.#flushAll(); + }); } - async flush(name?: keyof Tasks): Promise { - if (!isUndefined(name)) { - // Find queued task by name and flush by its key - for (const [key, queued] of Object.entries(this.#queued)) { - if (queued?.name === name) { - await this.#flushKey(key as keyof Tasks); - return; - } - } - return; - } + #flushAll(): void { + if (this.#destroyed) return; - const keys = Reflect.ownKeys(this.#queued); - await Promise.allSettled(keys.map(k => this.#flushKey(k))); + const keys = Reflect.ownKeys(this.#queued) as (keyof Tasks)[]; + for (const key of keys) { + const task = this.#queued[key]; + if (task) { + delete this.#queued[key]; + this.#executeNow(task); + } + } } - /** - * Abort task(s). If name provided, aborts that task. If no name, aborts all. - */ + /** Abort task(s). If name provided, aborts that task. If no name, aborts all. */ abort(name?: keyof Tasks): void { if (!isUndefined(name)) { // Find and abort queued task by name for (const [key, queued] of Object.entries(this.#queued)) { if (queued?.name === name) { - queued.invalidate?.(); queued.reject(new StoreError('ABORTED')); delete this.#queued[key as keyof Tasks]; break; @@ -389,7 +198,6 @@ export class Queue { const error = new StoreError('ABORTED'); for (const queued of Object.values(this.#queued)) { - queued.invalidate?.(); queued.reject(error); } @@ -411,17 +219,6 @@ export class Queue { this.#tasks = {}; } - async #flushKey(key: keyof Tasks): Promise { - if (this.#destroyed) return; - - const task = this.#queued[key]; - if (!task) return; - - delete this.#queued[key]; - - await this.#executeNow(task); - } - async #executeNow( task: QueuedTask, Tasks[K]['input'], Tasks[K]['output']>, ): Promise { @@ -444,7 +241,6 @@ export class Queue { // Store tasks by name for controller access (different names can share same key) this.#tasks[name as keyof Tasks] = pendingTask; this.#notifySubscribers(); - this.#onDispatch?.(pendingTask); try { if (abort.signal.aborted) { @@ -471,8 +267,6 @@ export class Queue { this.#tasks[name as keyof Tasks] = successTask; this.#notifySubscribers(); } - - this.#onSettled?.(successTask); } catch (error) { reject(error); @@ -489,8 +283,6 @@ export class Queue { this.#tasks[name as keyof Tasks] = errorTask; this.#notifySubscribers(); } - - this.#onSettled?.(errorTask); } } } @@ -503,7 +295,7 @@ export class Queue { * Create a queue for managing task execution. * * - Same key = supersede previous (cancel queued, abort pending) - * - Tasks scheduled via schedule function (default: microtask) + * - Tasks batched via microtask for supersession * * @example * // Loose typing (default) @@ -516,8 +308,6 @@ export class Queue { * 'volume': Request; * }>(); */ -export function createQueue( - config: QueueConfig = {}, -): Queue { - return new Queue(config); +export function createQueue(): Queue { + return new Queue(); } diff --git a/packages/store/src/core/request.ts b/packages/store/src/core/request.ts index 643c34bf..23005936 100644 --- a/packages/store/src/core/request.ts +++ b/packages/store/src/core/request.ts @@ -1,6 +1,6 @@ import type { EventLike } from '@videojs/utils/events'; import type { Guard } from './guard'; -import type { TaskKey, TaskScheduler } from './queue'; +import type { TaskKey } from './task'; import { isFunction, isObject } from '@videojs/utils/predicate'; @@ -42,7 +42,6 @@ export type RequestHandler = ( export interface RequestConfig { key?: RequestKey; - schedule?: TaskScheduler; guard?: Guard | Guard[]; cancel?: RequestCancel; handler: RequestHandler; @@ -50,7 +49,6 @@ export interface RequestConfig { export interface ResolvedRequestConfig { key: RequestKey; - schedule?: TaskScheduler | undefined; guard: Guard[]; cancel?: RequestCancel | undefined; handler: RequestHandler; @@ -120,7 +118,6 @@ export function resolveRequests[] = AnySlice[ key, input, meta, - schedule: config.schedule, handler, }); } catch (error) { diff --git a/packages/store/src/core/task.ts b/packages/store/src/core/task.ts new file mode 100644 index 00000000..2bf07991 --- /dev/null +++ b/packages/store/src/core/task.ts @@ -0,0 +1,77 @@ +import type { RequestMeta } from './request'; + +// ---------------------------------------- +// Types +// ---------------------------------------- + +export type TaskKey = T & (string | symbol); + +export type EnsureTaskKey = T extends string | symbol ? T : never; + +export interface TaskBase { + id: symbol; + name: string; + key: Key; + input: Input; + startedAt: number; + meta: RequestMeta | null; +} + +export interface PendingTask extends TaskBase { + status: 'pending'; + abort: AbortController; +} + +export interface SuccessTask extends TaskBase< + Key, + Input +> { + status: 'success'; + settledAt: number; + output: Output; +} + +export interface ErrorTask extends TaskBase { + status: 'error'; + settledAt: number; + error: unknown; + cancelled: boolean; +} + +export type Task + = | PendingTask + | SuccessTask + | ErrorTask; + +export type SettledTask + = | SuccessTask + | ErrorTask; + +export interface TaskContext { + input: Input; + signal: AbortSignal; +} + +// ---------------------------------------- +// Type Guards +// ---------------------------------------- + +/** Check if task is pending (in-flight). */ +export function isPendingTask(task: Task | undefined): task is PendingTask { + return task?.status === 'pending'; +} + +/** Check if task is settled (success or error). */ +export function isSettledTask(task: Task | undefined): task is SettledTask { + return task?.status === 'success' || task?.status === 'error'; +} + +/** Check if task is a success. */ +export function isSuccessTask(task: Task | undefined): task is SuccessTask { + return task?.status === 'success'; +} + +/** Check if task is an error. */ +export function isErrorTask(task: Task | undefined): task is ErrorTask { + return task?.status === 'error'; +} diff --git a/packages/store/src/core/tests/integration/store.test.ts b/packages/store/src/core/tests/integration/store.test.ts index 96abc9a1..5098948f 100644 --- a/packages/store/src/core/tests/integration/store.test.ts +++ b/packages/store/src/core/tests/integration/store.test.ts @@ -1,6 +1,6 @@ -import { describe, expect, it, vi } from 'vitest'; +import { describe, expect, it } from 'vitest'; -import { createSlice, createStore, delay } from '../../index'; +import { createSlice, createStore } from '../../index'; describe('store lifecycle integration', () => { it('full lifecycle: create → attach → use → detach → destroy', async () => { @@ -54,9 +54,7 @@ describe('store lifecycle integration', () => { expect(store.destroyed).toBe(true); }); - it('request with guards and scheduling', async () => { - vi.useFakeTimers(); - + it('request with guards', async () => { let ready = false; const isReady = () => ready; @@ -65,8 +63,7 @@ describe('store lifecycle integration', () => { getSnapshot: () => ({}), subscribe: () => {}, request: { - delayedAction: { - schedule: delay(100), + guardedAction: { guard: [isReady], handler: (_, _ctx) => 'completed', }, @@ -81,17 +78,13 @@ describe('store lifecycle integration', () => { store.attach({}); // Test 1: Guard rejects when not ready - const failPromise = store.request.delayedAction().catch(e => e); - await vi.runAllTimersAsync(); + const failPromise = store.request.guardedAction().catch(e => e); await expect(failPromise).resolves.toMatchObject({ code: 'REJECTED' }); // Test 2: Guard passes when ready ready = true; - const successPromise = store.request.delayedAction(); - await vi.runAllTimersAsync(); + const successPromise = store.request.guardedAction(); await expect(successPromise).resolves.toBe('completed'); - - vi.useRealTimers(); }); }); diff --git a/packages/store/src/core/tests/queue.test.ts b/packages/store/src/core/tests/queue.test.ts index 0a4f0b92..ad2edbea 100644 --- a/packages/store/src/core/tests/queue.test.ts +++ b/packages/store/src/core/tests/queue.test.ts @@ -1,8 +1,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; -import { createQueue, delay } from '../queue'; +import { createQueue } from '../queue'; -describe('queue', () => { +describe('Queue', () => { beforeEach(() => { vi.useFakeTimers(); }); @@ -11,1106 +11,702 @@ describe('queue', () => { vi.useRealTimers(); }); - describe('delay scheduler', () => { - it('delays execution', async () => { - const handler = vi.fn(); - const schedule = delay(100); + describe('enqueue', () => { + it('executes task via microtask scheduler', async () => { + const queue = createQueue(); + const handler = vi.fn().mockResolvedValue('result'); - const cancel = schedule(handler); + const promise = queue.enqueue({ + name: 'test', + key: 'test-key', + handler, + }); expect(handler).not.toHaveBeenCalled(); - vi.advanceTimersByTime(99); - expect(handler).not.toHaveBeenCalled(); - vi.advanceTimersByTime(1); - expect(handler).toHaveBeenCalledOnce(); - - expect(cancel).toBeTypeOf('function'); + await vi.runAllTimersAsync(); + await expect(promise).resolves.toBe('result'); }); - it('cancel prevents execution', () => { - const handler = vi.fn(); - const schedule = delay(100); + it('supersedes queued task with same key', async () => { + vi.useRealTimers(); - const cancel = schedule(handler); - vi.advanceTimersByTime(50); - cancel!(); - vi.advanceTimersByTime(100); + const queue = createQueue(); + const first = vi.fn().mockResolvedValue('first'); + const second = vi.fn().mockResolvedValue('second'); - expect(handler).not.toHaveBeenCalled(); + const promise1 = queue.enqueue({ name: 'a', key: 'same', handler: first }); + const promise2 = queue.enqueue({ name: 'b', key: 'same', handler: second }); + + await expect(promise1).rejects.toMatchObject({ code: 'SUPERSEDED' }); + await expect(promise2).resolves.toBe('second'); + expect(first).not.toHaveBeenCalled(); + expect(second).toHaveBeenCalledOnce(); + }); + + it('aborts pending task with same key', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + let aborted = false; + + const longRunning = queue.enqueue({ + name: 'long', + key: 'shared', + handler: async ({ signal }) => { + await new Promise((resolve, reject) => { + const timeout = setTimeout(resolve, 1000); + signal.addEventListener('abort', () => { + clearTimeout(timeout); + aborted = true; + reject(new Error('aborted')); + }); + }); + }, + }); + + // Let first task start + await new Promise(r => setTimeout(r, 10)); + + const superseding = queue.enqueue({ + name: 'supersede', + key: 'shared', + handler: async () => 'new result', + }); + + await expect(longRunning).rejects.toThrow(); + await expect(superseding).resolves.toBe('new result'); + expect(aborted).toBe(true); + }); + + it('parallel execution with different keys', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const results: string[] = []; + + const task1 = queue.enqueue({ + name: 'task1', + key: 'key-a', + handler: async () => { + results.push('a-start'); + await new Promise(r => setTimeout(r, 20)); + results.push('a-end'); + return 'a'; + }, + }); + + const task2 = queue.enqueue({ + name: 'task2', + key: 'key-b', + handler: async () => { + results.push('b-start'); + await new Promise(r => setTimeout(r, 10)); + results.push('b-end'); + return 'b'; + }, + }); + + await Promise.all([task1, task2]); + + expect(results).toEqual(['a-start', 'b-start', 'b-end', 'a-end']); }); }); - describe('queue', () => { - describe('enqueue', () => { - it('executes task via default microtask scheduler', async () => { - const queue = createQueue(); - const handler = vi.fn().mockResolvedValue('result'); + describe('abort', () => { + it('abort(name) cancels queued and aborts pending', async () => { + vi.useRealTimers(); - const promise = queue.enqueue({ - name: 'test', - key: 'test-key', - handler, - }); + const queue = createQueue(); + let aborted = false; - expect(handler).not.toHaveBeenCalled(); - await vi.runAllTimersAsync(); - await expect(promise).resolves.toBe('result'); - }); - - it('executes with custom scheduler', async () => { - const queue = createQueue({ - scheduler: (flush) => { - const id = setTimeout(flush, 50); - return () => clearTimeout(id); - }, - }); - - const handler = vi.fn().mockResolvedValue('done'); - const promise = queue.enqueue({ - name: 'delayed', - key: 'key', - handler, - }); - - expect(handler).not.toHaveBeenCalled(); - vi.advanceTimersByTime(50); - await vi.runAllTimersAsync(); - await expect(promise).resolves.toBe('done'); - }); - - it('supersedes queued task with same key', async () => { - vi.useRealTimers(); // Use real timers for microtasks - - const queue = createQueue(); - const first = vi.fn().mockResolvedValue('first'); - const second = vi.fn().mockResolvedValue('second'); - - const promise1 = queue.enqueue({ name: 'a', key: 'same', handler: first }); - const promise2 = queue.enqueue({ name: 'b', key: 'same', handler: second }); - - await expect(promise1).rejects.toMatchObject({ code: 'SUPERSEDED' }); - await expect(promise2).resolves.toBe('second'); - expect(first).not.toHaveBeenCalled(); - expect(second).toHaveBeenCalledOnce(); - }); - - it('aborts pending task with same key', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - let aborted = false; - - const longRunning = queue.enqueue({ - name: 'long', - key: 'shared', - handler: async ({ signal }) => { - await new Promise((resolve, reject) => { - const timeout = setTimeout(resolve, 1000); - signal.addEventListener('abort', () => { - clearTimeout(timeout); - aborted = true; - reject(new Error('aborted')); - }); + const promise = queue.enqueue({ + name: 'test', + key: 'k', + handler: async ({ signal }) => { + await new Promise((_, reject) => { + signal.addEventListener('abort', () => { + aborted = true; + reject(signal.reason); }); - }, - }); - - // Let first task start - await new Promise(r => setTimeout(r, 10)); - - const superseding = queue.enqueue({ - name: 'supersede', - key: 'shared', - handler: async () => 'new result', - }); - - await expect(longRunning).rejects.toThrow(); - await expect(superseding).resolves.toBe('new result'); - expect(aborted).toBe(true); + setTimeout(() => {}, 1000); + }); + }, }); - it('parallel execution with different keys', async () => { - vi.useRealTimers(); + await new Promise(r => setTimeout(r, 10)); + queue.abort('test'); - const queue = createQueue(); - const results: string[] = []; + await expect(promise).rejects.toMatchObject({ code: 'ABORTED' }); + expect(aborted).toBe(true); + }); + }); - const task1 = queue.enqueue({ - name: 'task1', - key: 'key-a', - handler: async () => { - results.push('a-start'); - await new Promise(r => setTimeout(r, 20)); - results.push('a-end'); - return 'a'; - }, - }); + describe('destroy', () => { + it('rejects after destroy', async () => { + const queue = createQueue(); + queue.destroy(); - const task2 = queue.enqueue({ - name: 'task2', - key: 'key-b', - handler: async () => { - results.push('b-start'); - await new Promise(r => setTimeout(r, 10)); - results.push('b-end'); - return 'b'; - }, - }); - - await Promise.all([task1, task2]); - - expect(results).toEqual(['a-start', 'b-start', 'b-end', 'a-end']); + await expect(queue.enqueue({ name: 't', key: 'k', handler: vi.fn() })).rejects.toMatchObject({ + code: 'DESTROYED', }); }); - describe('cancel', () => { - it('cancel(name) removes specific queued task', async () => { - const queue = createQueue({ - scheduler: delay(100), - }); + it('aborts all pending on destroy', async () => { + vi.useRealTimers(); - const handler = vi.fn(); - const promise = queue.enqueue({ name: 'test', key: 'k', handler }); + const queue = createQueue(); + const aborted = vi.fn(); - expect(queue.cancel('test')).toBe(true); - expect(queue.cancel('test')).toBe(false); - - vi.advanceTimersByTime(100); - await expect(promise).rejects.toMatchObject({ code: 'REMOVED' }); - expect(handler).not.toHaveBeenCalled(); + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async ({ signal }) => { + signal.addEventListener('abort', aborted); + await new Promise(r => setTimeout(r, 100)); + }, }); - it('cancel() clears all queued tasks', async () => { - const queue = createQueue({ scheduler: delay(100) }); + await new Promise(r => setTimeout(r, 10)); + queue.destroy(); - const p1 = queue.enqueue({ name: 'a', key: 'a', handler: vi.fn() }); - const p2 = queue.enqueue({ name: 'b', key: 'b', handler: vi.fn() }); - - expect(queue.cancel()).toBe(true); - vi.advanceTimersByTime(100); - - await expect(p1).rejects.toMatchObject({ code: 'REMOVED' }); - await expect(p2).rejects.toMatchObject({ code: 'REMOVED' }); - expect(Reflect.ownKeys(queue.queued).length).toBe(0); - }); - - it('cancel() returns false when no queued tasks', () => { - const queue = createQueue(); - - expect(queue.cancel()).toBe(false); - }); + await expect(promise).rejects.toMatchObject({ code: 'ABORTED' }); + expect(aborted).toHaveBeenCalled(); + expect(queue.destroyed).toBe(true); }); - describe('flush', () => { - it('flush() executes all queued immediately', async () => { - vi.useRealTimers(); + it('clears all task references on destroy', async () => { + vi.useRealTimers(); - const queue = createQueue({ scheduler: delay(1000) }); - const handler = vi.fn().mockResolvedValue('ok'); + const queue = createQueue(); - const promise = queue.enqueue({ name: 't', key: 'k', handler }); + await queue.enqueue({ name: 'task', key: 'k', handler: async () => 'result' }); + expect(queue.tasks.task?.status).toBe('success'); - queue.flush(); - await expect(promise).resolves.toBe('ok'); + queue.destroy(); + + expect(Reflect.ownKeys(queue.tasks).length).toBe(0); + }); + }); + + describe('cleanup edge cases', () => { + it('explicitly removes superseded task from queue before adding new one', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + const promise1 = queue.enqueue({ + name: 'task1', + key: 'shared', + handler: vi.fn().mockResolvedValue('result1'), }); - it('flush(name) executes specific task', async () => { - vi.useRealTimers(); - - const queue = createQueue({ scheduler: delay(1000) }); - const handlerA = vi.fn().mockResolvedValue('a'); - const handlerB = vi.fn().mockResolvedValue('b'); - - queue.enqueue({ name: 'a', key: 'a', handler: handlerA }); - queue.enqueue({ name: 'b', key: 'b', handler: handlerB }); - - queue.flush('a'); - await new Promise(r => setTimeout(r, 10)); - - expect(handlerA).toHaveBeenCalled(); - expect(handlerB).not.toHaveBeenCalled(); + const promise2 = queue.enqueue({ + name: 'task2', + key: 'shared', + handler: vi.fn().mockResolvedValue('result2'), }); + + await expect(promise1).rejects.toMatchObject({ code: 'SUPERSEDED' }); + await expect(promise2).resolves.toBe('result2'); }); - describe('abort', () => { - it('abort(name) cancels queued and aborts pending', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - let aborted = false; - - const promise = queue.enqueue({ - name: 'test', - key: 'k', - handler: async ({ signal }) => { - await new Promise((_, reject) => { - signal.addEventListener('abort', () => { - aborted = true; - reject(signal.reason); - }); - setTimeout(() => {}, 1000); - }); - }, - }); - - await new Promise(r => setTimeout(r, 10)); - queue.abort('test'); - - await expect(promise).rejects.toMatchObject({ code: 'ABORTED' }); - expect(aborted).toBe(true); - }); - }); - - describe('lifecycle hooks', () => { - it('onDispatch called when task starts', async () => { - vi.useRealTimers(); - - const onDispatch = vi.fn(); - const queue = createQueue({ onDispatch }); - - await queue.enqueue({ - name: 'myTask', - key: 'k', - input: { value: 42 }, - handler: async () => 'result', - }); - - expect(onDispatch).toHaveBeenCalledWith( - expect.objectContaining({ - name: 'myTask', - key: 'k', - input: { value: 42 }, - }), - ); - }); - - it('onSettled called with success task', async () => { - vi.useRealTimers(); - - const onSettled = vi.fn(); - const queue = createQueue({ onSettled }); - - await queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => 'done', - }); - - expect(onSettled).toHaveBeenCalledWith( - expect.objectContaining({ name: 'task', status: 'success', output: 'done' }), - ); - }); - - it('onSettled called with error task', async () => { - vi.useRealTimers(); - - const onSettled = vi.fn(); - const queue = createQueue({ onSettled }); - - const promise = queue.enqueue({ - name: 'failing', - key: 'k', - handler: async () => { - throw new Error('oops'); - }, - }); - - await expect(promise).rejects.toThrow('oops'); - expect(onSettled).toHaveBeenCalledWith(expect.objectContaining({ status: 'error', cancelled: false })); - }); - - it('onSettled called with cancelled task', async () => { - vi.useRealTimers(); - - const onSettled = vi.fn(); - const queue = createQueue({ onSettled }); - - const promise = queue.enqueue({ - name: 'first', - key: 'k', - handler: async () => new Promise(r => setTimeout(r, 100)), - }); - - await new Promise(r => setTimeout(r, 10)); - - queue.enqueue({ - name: 'second', - key: 'k', - handler: async () => 'done', - }); - - await promise.catch(() => {}); - - expect(onSettled).toHaveBeenCalledWith( - expect.objectContaining({ name: 'first', status: 'error', cancelled: true }), - ); - }); - }); - - describe('destroy', () => { - it('rejects after destroy', async () => { - const queue = createQueue(); - queue.destroy(); - - await expect(queue.enqueue({ name: 't', key: 'k', handler: vi.fn() })).rejects.toMatchObject({ - code: 'DESTROYED', - }); - }); - - it('aborts all pending on destroy', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const aborted = vi.fn(); - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async ({ signal }) => { - signal.addEventListener('abort', aborted); - await new Promise(r => setTimeout(r, 100)); - }, - }); - - await new Promise(r => setTimeout(r, 10)); - queue.destroy(); - - await expect(promise).rejects.toMatchObject({ code: 'ABORTED' }); - expect(aborted).toHaveBeenCalled(); - expect(queue.destroyed).toBe(true); - }); - }); - - describe('cleanup edge cases', () => { - it('explicitly removes superseded task from queue before adding new one', async () => { - const queue = createQueue({ scheduler: delay(100) }); - - // Enqueue first task - const promise1 = queue.enqueue({ - name: 'task1', - key: 'shared', - handler: vi.fn().mockResolvedValue('result1'), - }); - - // Queue should have the first task (#queued is keyed by key 'shared') - expect(Reflect.ownKeys(queue.queued).length).toBe(1); - expect(queue.queued.shared?.name).toBe('task1'); - - // Immediately enqueue second task with same key (supersedes first) - const promise2 = queue.enqueue({ - name: 'task2', - key: 'shared', - handler: vi.fn().mockResolvedValue('result2'), - }); - - // First should be superseded - await expect(promise1).rejects.toMatchObject({ code: 'SUPERSEDED' }); - - // Queue should only have the second task (#queued keyed by key 'shared') - expect(Reflect.ownKeys(queue.queued).length).toBe(1); - expect(queue.queued.shared?.name).toBe('task2'); - - // Second should succeed - await vi.runAllTimersAsync(); - await expect(promise2).resolves.toBe('result2'); - }); - - it('allows pending tasks to self-cleanup after destroy', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const cleanupSpy = vi.fn(); - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async ({ signal }) => { - signal.addEventListener('abort', () => cleanupSpy('aborted')); - try { - await new Promise((_, reject) => { - signal.addEventListener('abort', () => reject(signal.reason)); - }); - } finally { - cleanupSpy('cleanup'); - } - }, - }); - - // Wait for task to start - await new Promise(r => setTimeout(r, 10)); - expect(queue.tasks.task?.status).toBe('pending'); - - // Destroy queue - queue.destroy(); - - // Task should self-cleanup - await promise.catch(() => {}); - expect(cleanupSpy).toHaveBeenCalledWith('aborted'); - expect(cleanupSpy).toHaveBeenCalledWith('cleanup'); - - // After destroy, all tasks are cleared for memory cleanup - expect(queue.tasks.task).toBeUndefined(); - }); - - it('handles scheduler error without double-cleanup when already flushed', async () => { - const queue = createQueue(); - let flushCalled = false; - - // Scheduler that flushes synchronously then throws - const faultyScheduler = (flush: () => void) => { - flush(); // Synchronous flush - flushCalled = true; - throw new Error('Scheduler error'); - }; - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - schedule: faultyScheduler, - handler: vi.fn().mockResolvedValue('result'), - }); - - // Task was flushed despite scheduler error - expect(flushCalled).toBe(true); - - // Promise should reject with scheduler error - await expect(promise).rejects.toThrow('Scheduler error'); - - // Task should not be in queued object (wasn't double-deleted) - expect(Reflect.ownKeys(queue.queued).length).toBe(0); - }); - - it('handles scheduler error with cleanup when not yet flushed', async () => { - const queue = createQueue(); - - // Scheduler that throws before flushing - const faultyScheduler = () => { - throw new Error('Scheduler error'); - }; - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - schedule: faultyScheduler, - handler: vi.fn().mockResolvedValue('result'), - }); - - // Promise should reject with scheduler error - await expect(promise).rejects.toThrow('Scheduler error'); - - // Task should be removed from queue - expect(Reflect.ownKeys(queue.queued).length).toBe(0); - }); - }); - - describe('isPending and isQueued', () => { - it('isPending returns true when task is executing', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - const promise = queue.enqueue({ - name: 'test', - key: 'test-key', - handler: async () => { - await new Promise(r => setTimeout(r, 50)); - return 'result'; - }, - }); - - await new Promise(r => setTimeout(r, 10)); - - expect(queue.isPending('test')).toBe(true); - expect(queue.isPending('other')).toBe(false); - - await promise; - - expect(queue.isPending('test')).toBe(false); - }); - - it('isQueued returns true when task is waiting to execute', async () => { - const queue = createQueue({ scheduler: delay(100) }); - - queue.enqueue({ - name: 'test', - key: 'test-key', - handler: vi.fn().mockResolvedValue('result'), - }); - - expect(queue.isQueued('test')).toBe(true); - expect(queue.isQueued('other')).toBe(false); - - vi.advanceTimersByTime(100); - await vi.runAllTimersAsync(); - - expect(queue.isQueued('test')).toBe(false); - }); - }); - - describe('subscribe', () => { - it('returns an unsubscribe function', () => { - const queue = createQueue(); - const listener = vi.fn(); - - const unsubscribe = queue.subscribe(listener); - - expect(unsubscribe).toBeTypeOf('function'); - }); - - it('notifies when task becomes pending', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const listener = vi.fn(); - - queue.subscribe(listener); - - const promise = queue.enqueue({ - name: 'test', - key: 'test-key', - handler: vi.fn().mockResolvedValue('result'), - }); - - await promise; - - // Called when pending (dispatch) and when settled - expect(listener).toHaveBeenCalledTimes(2); - }); - - it('notifies with tasks map on dispatch and settlement', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const listener = vi.fn(); - - queue.subscribe(listener); - - let resolveHandler: () => void; - const handlerPromise = new Promise((resolve) => { - resolveHandler = resolve; - }); - - const promise = queue.enqueue({ - name: 'test', - key: 'test-key', - handler: async () => { - await handlerPromise; - return 'result'; - }, - }); - - // Wait for dispatch - await new Promise(r => setTimeout(r, 10)); - - // First call should have pending task (keyed by name 'test') - expect(listener).toHaveBeenCalledTimes(1); - const pendingSnapshot = listener.mock.calls[0]![0] as Record; - expect(Reflect.ownKeys(pendingSnapshot).length).toBe(1); - expect(pendingSnapshot.test?.status).toBe('pending'); - - // Complete the handler - resolveHandler!(); - await promise; - - // Second call should have settled task (success) - expect(listener).toHaveBeenCalledTimes(2); - const settledSnapshot = listener.mock.calls[1]![0] as Record; - expect(Reflect.ownKeys(settledSnapshot).length).toBe(1); - expect(settledSnapshot.test?.status).toBe('success'); - }); - - it('unsubscribe stops notifications', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const listener = vi.fn(); - - const unsubscribe = queue.subscribe(listener); - unsubscribe(); - - await queue.enqueue({ - name: 'test', - key: 'test-key', - handler: vi.fn().mockResolvedValue('result'), - }); - - expect(listener).not.toHaveBeenCalled(); - }); - - it('supports multiple subscribers', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const listener1 = vi.fn(); - const listener2 = vi.fn(); - - queue.subscribe(listener1); - queue.subscribe(listener2); - - await queue.enqueue({ - name: 'test', - key: 'test-key', - handler: vi.fn().mockResolvedValue('result'), - }); - - expect(listener1).toHaveBeenCalledTimes(2); - expect(listener2).toHaveBeenCalledTimes(2); - }); - - it('catches and logs listener errors', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const consoleSpy = vi.spyOn(console, 'error').mockImplementation(() => {}); - const errorListener = vi.fn(() => { - throw new Error('Listener error'); - }); - const successListener = vi.fn(); - - queue.subscribe(errorListener); - queue.subscribe(successListener); - - await queue.enqueue({ - name: 'test', - key: 'test-key', - handler: vi.fn().mockResolvedValue('result'), - }); - - // Both listeners were called despite error - expect(errorListener).toHaveBeenCalled(); - expect(successListener).toHaveBeenCalled(); - - // Error was logged - expect(consoleSpy).toHaveBeenCalledWith('[vjs-queue]', expect.any(Error)); - - consoleSpy.mockRestore(); - }); - - it('provides strongly typed tasks object', async () => { - vi.useRealTimers(); - - // Use default queue - type safety is validated at compile time - const queue = createQueue(); - - queue.subscribe((tasks) => { - // Tasks is a frozen object - const task = tasks.playback; - if (task) { - expect(task.key).toBe('playback'); - expect(task.name).toBeDefined(); - } - }); - - await queue.enqueue({ - name: 'play', - key: 'playback', - handler: vi.fn().mockResolvedValue(undefined), - }); - }); - }); - - describe('task lifecycle', () => { - it('task starts as pending and transitions to success', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => { - await new Promise(r => setTimeout(r, 10)); - return 'result'; - }, - }); - - // Task should be pending (keyed by name 'task') - await new Promise(r => setTimeout(r, 5)); - const pendingTask = queue.tasks.task; - expect(pendingTask?.status).toBe('pending'); - expect(pendingTask?.name).toBe('task'); - - // Wait for completion - await promise; - - // Task should be success - const successTask = queue.tasks.task; - expect(successTask?.status).toBe('success'); - if (successTask?.status === 'success') { - expect(successTask.output).toBe('result'); - expect(successTask.settledAt).toBeGreaterThan(successTask.startedAt); - } - }); - - it('task starts as pending and transitions to error', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const error = new Error('test error'); - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => { - await new Promise(r => setTimeout(r, 10)); - throw error; - }, - }); - - // Task should be pending - await new Promise(r => setTimeout(r, 5)); - expect(queue.tasks.task?.status).toBe('pending'); - - // Wait for failure - await expect(promise).rejects.toThrow('test error'); - - // Task should be error - const errorTask = queue.tasks.task; - expect(errorTask?.status).toBe('error'); - if (errorTask?.status === 'error') { - expect(errorTask.error).toBe(error); - expect(errorTask.cancelled).toBe(false); - } - }); - - it('aborted task has cancelled flag set to true', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async ({ signal }) => { + it('allows pending tasks to self-cleanup after destroy', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const cleanupSpy = vi.fn(); + + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async ({ signal }) => { + signal.addEventListener('abort', () => cleanupSpy('aborted')); + try { await new Promise((_, reject) => { signal.addEventListener('abort', () => reject(signal.reason)); - setTimeout(() => {}, 1000); }); - }, - }); + } finally { + cleanupSpy('cleanup'); + } + }, + }); - // Wait for task to start - await new Promise(r => setTimeout(r, 10)); - expect(queue.tasks.task?.status).toBe('pending'); + await new Promise(r => setTimeout(r, 10)); + expect(queue.tasks.task?.status).toBe('pending'); - // Abort the task (by name) - queue.abort('task'); - await promise.catch(() => {}); + queue.destroy(); - // Task should be error with cancelled=true - const errorTask = queue.tasks.task; - expect(errorTask?.status).toBe('error'); - if (errorTask?.status === 'error') { - expect(errorTask.cancelled).toBe(true); + await promise.catch(() => {}); + expect(cleanupSpy).toHaveBeenCalledWith('aborted'); + expect(cleanupSpy).toHaveBeenCalledWith('cleanup'); + expect(queue.tasks.task).toBeUndefined(); + }); + }); + + describe('subscribe', () => { + it('returns an unsubscribe function', () => { + const queue = createQueue(); + const listener = vi.fn(); + + const unsubscribe = queue.subscribe(listener); + + expect(unsubscribe).toBeTypeOf('function'); + }); + + it('notifies when task becomes pending', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const listener = vi.fn(); + + queue.subscribe(listener); + + const promise = queue.enqueue({ + name: 'test', + key: 'test-key', + handler: vi.fn().mockResolvedValue('result'), + }); + + await promise; + + // Called when pending (dispatch) and when settled + expect(listener).toHaveBeenCalledTimes(2); + }); + + it('notifies with tasks map on dispatch and settlement', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const listener = vi.fn(); + + queue.subscribe(listener); + + let resolveHandler: () => void; + const handlerPromise = new Promise((resolve) => { + resolveHandler = resolve; + }); + + const promise = queue.enqueue({ + name: 'test', + key: 'test-key', + handler: async () => { + await handlerPromise; + return 'result'; + }, + }); + + // Wait for dispatch + await new Promise(r => setTimeout(r, 10)); + + expect(listener).toHaveBeenCalledTimes(1); + const pendingSnapshot = listener.mock.calls[0]![0] as Record; + expect(Reflect.ownKeys(pendingSnapshot).length).toBe(1); + expect(pendingSnapshot.test?.status).toBe('pending'); + + resolveHandler!(); + await promise; + + expect(listener).toHaveBeenCalledTimes(2); + const settledSnapshot = listener.mock.calls[1]![0] as Record; + expect(Reflect.ownKeys(settledSnapshot).length).toBe(1); + expect(settledSnapshot.test?.status).toBe('success'); + }); + + it('unsubscribe stops notifications', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const listener = vi.fn(); + + const unsubscribe = queue.subscribe(listener); + unsubscribe(); + + await queue.enqueue({ + name: 'test', + key: 'test-key', + handler: vi.fn().mockResolvedValue('result'), + }); + + expect(listener).not.toHaveBeenCalled(); + }); + + it('supports multiple subscribers', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const listener1 = vi.fn(); + const listener2 = vi.fn(); + + queue.subscribe(listener1); + queue.subscribe(listener2); + + await queue.enqueue({ + name: 'test', + key: 'test-key', + handler: vi.fn().mockResolvedValue('result'), + }); + + expect(listener1).toHaveBeenCalledTimes(2); + expect(listener2).toHaveBeenCalledTimes(2); + }); + + it('catches and logs listener errors', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const consoleSpy = vi.spyOn(console, 'error').mockImplementation(() => {}); + const errorListener = vi.fn(() => { + throw new Error('Listener error'); + }); + const successListener = vi.fn(); + + queue.subscribe(errorListener); + queue.subscribe(successListener); + + await queue.enqueue({ + name: 'test', + key: 'test-key', + handler: vi.fn().mockResolvedValue('result'), + }); + + expect(errorListener).toHaveBeenCalled(); + expect(successListener).toHaveBeenCalled(); + expect(consoleSpy).toHaveBeenCalledWith('[vjs-queue]', expect.any(Error)); + + consoleSpy.mockRestore(); + }); + + it('provides strongly typed tasks object', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + queue.subscribe((tasks) => { + const task = tasks.playback; + if (task) { + expect(task.key).toBe('playback'); + expect(task.name).toBeDefined(); } }); - it('new request replaces settled task', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - // First request - await queue.enqueue({ - name: 'first', - key: 'k', - handler: async () => 'first-result', - }); - - // Keyed by name 'first' - expect(queue.tasks.first?.status).toBe('success'); - if (queue.tasks.first?.status === 'success') { - expect(queue.tasks.first.output).toBe('first-result'); - } - - // Second request (different name, same key) - doesn't replace first since different name - await queue.enqueue({ - name: 'second', - key: 'k', - handler: async () => 'second-result', - }); - - // Both tasks exist at different keys (names) - expect(queue.tasks.second?.status).toBe('success'); - if (queue.tasks.second?.status === 'success') { - expect(queue.tasks.second.output).toBe('second-result'); - expect(queue.tasks.second.name).toBe('second'); - } + await queue.enqueue({ + name: 'play', + key: 'playback', + handler: vi.fn().mockResolvedValue(undefined), }); }); + }); - describe('reset', () => { - it('clears settled task', async () => { - vi.useRealTimers(); + describe('task lifecycle', () => { + it('task starts as pending and transitions to success', async () => { + vi.useRealTimers(); - const queue = createQueue(); + const queue = createQueue(); - await queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => 'result', - }); - - expect(queue.tasks.task?.status).toBe('success'); - - queue.reset('task'); - - expect(queue.tasks.task).toBeUndefined(); + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => { + await new Promise(r => setTimeout(r, 10)); + return 'result'; + }, }); - it('is no-op when task is pending', async () => { - vi.useRealTimers(); + await new Promise(r => setTimeout(r, 5)); + const pendingTask = queue.tasks.task; + expect(pendingTask?.status).toBe('pending'); + expect(pendingTask?.name).toBe('task'); - const queue = createQueue(); + await promise; - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => { - await new Promise(r => setTimeout(r, 50)); - return 'result'; - }, - }); - - // Wait for task to start - await new Promise(r => setTimeout(r, 10)); - expect(queue.tasks.task?.status).toBe('pending'); - - // Reset should be no-op - queue.reset('task'); - expect(queue.tasks.task?.status).toBe('pending'); - - await promise; - }); - - it('is no-op when task does not exist', () => { - const queue = createQueue(); - - // Should not throw - queue.reset('nonexistent'); - - expect(queue.tasks.nonexistent).toBeUndefined(); - }); - - it('notifies subscribers when reset clears a task', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const listener = vi.fn(); - - await queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => 'result', - }); - - queue.subscribe(listener); - - queue.reset('task'); - - expect(listener).toHaveBeenCalledTimes(1); - const snapshot = listener.mock.calls[0]![0] as Record; - expect(snapshot.task).toBeUndefined(); - }); - - it('does not notify subscribers when task does not exist', () => { - const queue = createQueue(); - const listener = vi.fn(); - - queue.subscribe(listener); - queue.reset('nonexistent'); - - expect(listener).not.toHaveBeenCalled(); - }); - - it('resets all settled tasks when no key provided', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - // Create multiple settled tasks - await queue.enqueue({ name: 'a', key: 'a', handler: async () => 'a-result' }); - await queue.enqueue({ name: 'b', key: 'b', handler: async () => 'b-result' }); - - expect(queue.tasks.a?.status).toBe('success'); - expect(queue.tasks.b?.status).toBe('success'); - - // Reset all - queue.reset(); - - expect(queue.tasks.a).toBeUndefined(); - expect(queue.tasks.b).toBeUndefined(); - }); - - it('preserves pending tasks when resetting all', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - - // Create a settled task - await queue.enqueue({ name: 'settled', key: 'settled', handler: async () => 'done' }); - - // Create a pending task - const pendingPromise = queue.enqueue({ - name: 'pending', - key: 'pending', - handler: async () => { - await new Promise(r => setTimeout(r, 100)); - return 'pending-done'; - }, - }); - - await new Promise(r => setTimeout(r, 10)); - expect(queue.tasks.settled?.status).toBe('success'); - expect(queue.tasks.pending?.status).toBe('pending'); - - // Reset all - should only clear settled - queue.reset(); - - expect(queue.tasks.settled).toBeUndefined(); - expect(queue.tasks.pending?.status).toBe('pending'); - - await pendingPromise; - }); + const successTask = queue.tasks.task; + expect(successTask?.status).toBe('success'); + if (successTask?.status === 'success') { + expect(successTask.output).toBe('result'); + expect(successTask.settledAt).toBeGreaterThan(successTask.startedAt); + } }); - describe('isSettled', () => { - it('returns true for success task', async () => { - vi.useRealTimers(); + it('task starts as pending and transitions to error', async () => { + vi.useRealTimers(); - const queue = createQueue(); - await queue.enqueue({ name: 'task', key: 'k', handler: async () => 'result' }); + const queue = createQueue(); + const error = new Error('test error'); - expect(queue.isSettled('task')).toBe(true); + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => { + await new Promise(r => setTimeout(r, 10)); + throw error; + }, }); - it('returns true for error task', async () => { - vi.useRealTimers(); + await new Promise(r => setTimeout(r, 5)); + expect(queue.tasks.task?.status).toBe('pending'); - const queue = createQueue(); - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => { - throw new Error('fail'); - }, - }); + await expect(promise).rejects.toThrow('test error'); - await promise.catch(() => {}); - - expect(queue.isSettled('task')).toBe(true); - }); - - it('returns false for pending task', async () => { - vi.useRealTimers(); - - const queue = createQueue(); - const promise = queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => { - await new Promise(r => setTimeout(r, 50)); - return 'result'; - }, - }); - - await new Promise(r => setTimeout(r, 10)); - - expect(queue.isSettled('task')).toBe(false); - - await promise; - }); - - it('returns false for non-existent task', () => { - const queue = createQueue(); - - expect(queue.isSettled('nonexistent')).toBe(false); - }); + const errorTask = queue.tasks.task; + expect(errorTask?.status).toBe('error'); + if (errorTask?.status === 'error') { + expect(errorTask.error).toBe(error); + expect(errorTask.cancelled).toBe(false); + } }); - describe('destroy cleanup', () => { - it('clears all task references on destroy', async () => { - vi.useRealTimers(); + it('aborted task has cancelled flag set to true', async () => { + vi.useRealTimers(); - const queue = createQueue(); + const queue = createQueue(); - await queue.enqueue({ name: 'task', key: 'k', handler: async () => 'result' }); - expect(queue.tasks.task?.status).toBe('success'); - - queue.destroy(); - - expect(Reflect.ownKeys(queue.tasks).length).toBe(0); + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async ({ signal }) => { + await new Promise((_, reject) => { + signal.addEventListener('abort', () => reject(signal.reason)); + setTimeout(() => {}, 1000); + }); + }, }); + + await new Promise(r => setTimeout(r, 10)); + expect(queue.tasks.task?.status).toBe('pending'); + + queue.abort('task'); + await promise.catch(() => {}); + + const errorTask = queue.tasks.task; + expect(errorTask?.status).toBe('error'); + if (errorTask?.status === 'error') { + expect(errorTask.cancelled).toBe(true); + } }); - describe('tasks getter', () => { - it('returns frozen snapshot', async () => { - vi.useRealTimers(); + it('new request replaces settled task', async () => { + vi.useRealTimers(); - const queue = createQueue(); - await queue.enqueue({ name: 'task', key: 'k', handler: async () => 'result' }); + const queue = createQueue(); - const tasks = queue.tasks; - - expect(Object.isFrozen(tasks)).toBe(true); + await queue.enqueue({ + name: 'first', + key: 'k', + handler: async () => 'first-result', }); - it('returns independent snapshots', async () => { - vi.useRealTimers(); + expect(queue.tasks.first?.status).toBe('success'); + if (queue.tasks.first?.status === 'success') { + expect(queue.tasks.first.output).toBe('first-result'); + } - const queue = createQueue(); - await queue.enqueue({ name: 'first', key: 'k', handler: async () => 'first' }); - - const snapshot1 = queue.tasks; - - await queue.enqueue({ name: 'second', key: 'k', handler: async () => 'second' }); - - const snapshot2 = queue.tasks; - - // Tasks keyed by name - first and second are different entries - expect(snapshot1).not.toBe(snapshot2); - if (snapshot1.first?.status === 'success' && snapshot2.second?.status === 'success') { - expect(snapshot1.first.output).toBe('first'); - expect(snapshot2.second.output).toBe('second'); - } + await queue.enqueue({ + name: 'second', + key: 'k', + handler: async () => 'second-result', }); + + expect(queue.tasks.second?.status).toBe('success'); + if (queue.tasks.second?.status === 'success') { + expect(queue.tasks.second.output).toBe('second-result'); + expect(queue.tasks.second.name).toBe('second'); + } + }); + }); + + describe('reset', () => { + it('clears settled task', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + await queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => 'result', + }); + + expect(queue.tasks.task?.status).toBe('success'); + + queue.reset('task'); + + expect(queue.tasks.task).toBeUndefined(); }); - describe('symbol keys', () => { - it('supports symbol names', async () => { - vi.useRealTimers(); + it('is no-op when task is pending', async () => { + vi.useRealTimers(); - const queue = createQueue(); - const name = Symbol('task'); + const queue = createQueue(); - await queue.enqueue({ - name: name as unknown as string, - key: name, - handler: async () => 'result', - }); - - const task = queue.tasks[name as unknown as string]; - expect(task?.status).toBe('success'); - if (task?.status === 'success') { - expect(task.output).toBe('result'); - } + const promise = queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => { + await new Promise(r => setTimeout(r, 50)); + return 'result'; + }, }); + + await new Promise(r => setTimeout(r, 10)); + expect(queue.tasks.task?.status).toBe('pending'); + + queue.reset('task'); + expect(queue.tasks.task?.status).toBe('pending'); + + await promise; }); - describe('meta propagation', () => { - it('meta defaults to null when not provided', async () => { - vi.useRealTimers(); + it('is no-op when task does not exist', () => { + const queue = createQueue(); - const queue = createQueue(); + queue.reset('nonexistent'); - await queue.enqueue({ - name: 'task', - key: 'k', - handler: async () => 'result', - }); + expect(queue.tasks.nonexistent).toBeUndefined(); + }); - expect(queue.tasks.task?.meta).toBeNull(); + it('notifies subscribers when reset clears a task', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const listener = vi.fn(); + + await queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => 'result', }); + + queue.subscribe(listener); + + queue.reset('task'); + + expect(listener).toHaveBeenCalledTimes(1); + const snapshot = listener.mock.calls[0]![0] as Record; + expect(snapshot.task).toBeUndefined(); + }); + + it('does not notify subscribers when task does not exist', () => { + const queue = createQueue(); + const listener = vi.fn(); + + queue.subscribe(listener); + queue.reset('nonexistent'); + + expect(listener).not.toHaveBeenCalled(); + }); + + it('resets all settled tasks when no key provided', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + await queue.enqueue({ name: 'a', key: 'a', handler: async () => 'a-result' }); + await queue.enqueue({ name: 'b', key: 'b', handler: async () => 'b-result' }); + + expect(queue.tasks.a?.status).toBe('success'); + expect(queue.tasks.b?.status).toBe('success'); + + queue.reset(); + + expect(queue.tasks.a).toBeUndefined(); + expect(queue.tasks.b).toBeUndefined(); + }); + + it('preserves pending tasks when resetting all', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + await queue.enqueue({ name: 'settled', key: 'settled', handler: async () => 'done' }); + + const pendingPromise = queue.enqueue({ + name: 'pending', + key: 'pending', + handler: async () => { + await new Promise(r => setTimeout(r, 100)); + return 'pending-done'; + }, + }); + + await new Promise(r => setTimeout(r, 10)); + expect(queue.tasks.settled?.status).toBe('success'); + expect(queue.tasks.pending?.status).toBe('pending'); + + queue.reset(); + + expect(queue.tasks.settled).toBeUndefined(); + expect(queue.tasks.pending?.status).toBe('pending'); + + await pendingPromise; + }); + }); + + describe('tasks getter', () => { + it('returns frozen snapshot', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + await queue.enqueue({ name: 'task', key: 'k', handler: async () => 'result' }); + + const tasks = queue.tasks; + + expect(Object.isFrozen(tasks)).toBe(true); + }); + + it('returns independent snapshots', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + await queue.enqueue({ name: 'first', key: 'k', handler: async () => 'first' }); + + const snapshot1 = queue.tasks; + + await queue.enqueue({ name: 'second', key: 'k', handler: async () => 'second' }); + + const snapshot2 = queue.tasks; + + expect(snapshot1).not.toBe(snapshot2); + if (snapshot1.first?.status === 'success' && snapshot2.second?.status === 'success') { + expect(snapshot1.first.output).toBe('first'); + expect(snapshot2.second.output).toBe('second'); + } + }); + }); + + describe('symbol keys', () => { + it('supports symbol names', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + const name = Symbol('task'); + + await queue.enqueue({ + name: name as unknown as string, + key: name, + handler: async () => 'result', + }); + + const task = queue.tasks[name as unknown as string]; + expect(task?.status).toBe('success'); + if (task?.status === 'success') { + expect(task.output).toBe('result'); + } + }); + }); + + describe('meta propagation', () => { + it('meta defaults to null when not provided', async () => { + vi.useRealTimers(); + + const queue = createQueue(); + + await queue.enqueue({ + name: 'task', + key: 'k', + handler: async () => 'result', + }); + + expect(queue.tasks.task?.meta).toBeNull(); }); }); }); diff --git a/packages/store/src/core/tests/queue.types.test.ts b/packages/store/src/core/tests/queue.types.test.ts index 32ceef0b..faafd438 100644 --- a/packages/store/src/core/tests/queue.types.test.ts +++ b/packages/store/src/core/tests/queue.types.test.ts @@ -1,42 +1,10 @@ -import type { ErrorTask, PendingTask, SuccessTask, Task, TasksRecord } from '../queue'; +import type { TasksRecord } from '../queue'; import { describe, expectTypeOf, it } from 'vitest'; + import { createQueue } from '../queue'; describe('queue types', () => { - describe('Task', () => { - it('is discriminated union of task states', () => { - const task: Task = {} as Task; - - if (task.status === 'pending') { - expectTypeOf(task).toExtend(); - expectTypeOf(task.abort).toExtend(); - } - - if (task.status === 'success') { - expectTypeOf(task).toExtend(); - expectTypeOf(task.output).toBeUnknown(); - expectTypeOf(task.settledAt).toBeNumber(); - } - - if (task.status === 'error') { - expectTypeOf(task).toExtend(); - expectTypeOf(task.error).toBeUnknown(); - expectTypeOf(task.cancelled).toBeBoolean(); - expectTypeOf(task.settledAt).toBeNumber(); - } - }); - - it('has common properties across all states', () => { - const task: Task = {} as Task; - - expectTypeOf(task.id).toEqualTypeOf(); - expectTypeOf(task.name).toEqualTypeOf(); - expectTypeOf(task.key).toExtend(); - expectTypeOf(task.startedAt).toBeNumber(); - }); - }); - describe('createQueue', () => { it('returns Queue with default task record', () => { const queue = createQueue(); @@ -47,27 +15,6 @@ describe('queue types', () => { }); describe('Queue methods', () => { - it('isPending takes name parameter', () => { - const queue = createQueue(); - - expectTypeOf(queue.isPending).toBeFunction(); - expectTypeOf(queue.isPending).returns.toBeBoolean(); - }); - - it('isQueued takes name parameter', () => { - const queue = createQueue(); - - expectTypeOf(queue.isQueued).toBeFunction(); - expectTypeOf(queue.isQueued).returns.toBeBoolean(); - }); - - it('isSettled takes name parameter', () => { - const queue = createQueue(); - - expectTypeOf(queue.isSettled).toBeFunction(); - expectTypeOf(queue.isSettled).returns.toBeBoolean(); - }); - it('reset takes optional name parameter', () => { const queue = createQueue(); @@ -75,13 +22,6 @@ describe('queue types', () => { expectTypeOf(queue.reset).returns.toBeVoid(); }); - it('cancel takes optional name parameter and returns boolean', () => { - const queue = createQueue(); - - expectTypeOf(queue.cancel).toBeFunction(); - expectTypeOf(queue.cancel).returns.toBeBoolean(); - }); - it('abort takes optional name parameter', () => { const queue = createQueue(); @@ -89,13 +29,6 @@ describe('queue types', () => { expectTypeOf(queue.abort).returns.toBeVoid(); }); - it('flush takes optional name parameter and returns promise', () => { - const queue = createQueue(); - - expectTypeOf(queue.flush).toBeFunction(); - expectTypeOf(queue.flush).returns.toExtend>(); - }); - it('subscribe takes listener and returns unsubscribe', () => { const queue = createQueue(); diff --git a/packages/store/src/core/tests/slice.test.ts b/packages/store/src/core/tests/slice.test.ts index 6fdef7a2..5abef461 100644 --- a/packages/store/src/core/tests/slice.test.ts +++ b/packages/store/src/core/tests/slice.test.ts @@ -49,10 +49,6 @@ describe('slice', () => { it('preserves full config options', () => { const guard = () => true; - const schedule = (flush: () => void) => { - setTimeout(flush, 100); - }; - const slice = createSlice({ initialState: {}, getSnapshot: () => ({}), @@ -61,7 +57,6 @@ describe('slice', () => { configured: { key: 'custom-key', guard: [guard], - schedule, handler: () => {}, }, }, @@ -69,7 +64,6 @@ describe('slice', () => { expect(slice.request.configured.key).toBe('custom-key'); expect(slice.request.configured.guard).toEqual([guard]); - expect(slice.request.configured.schedule).toBe(schedule); }); }); }); diff --git a/packages/store/src/core/tests/store.types.test.ts b/packages/store/src/core/tests/store.types.test.ts index a7c6ec5b..72fc17f4 100644 --- a/packages/store/src/core/tests/store.types.test.ts +++ b/packages/store/src/core/tests/store.types.test.ts @@ -188,28 +188,6 @@ describe('store types', () => { expectTypeOf(store.queue.tasks).toHaveProperty('pause'); }); - it('queue.isPending accepts request names', () => { - const store = createSingleSliceStore(); - - // Should compile - valid request names - store.queue.isPending('setVolume'); - store.queue.isPending('setMuted'); - }); - - it('queue.isQueued accepts request names', () => { - const store = createSingleSliceStore(); - - store.queue.isQueued('setVolume'); - store.queue.isQueued('setMuted'); - }); - - it('queue.isSettled accepts request names', () => { - const store = createSingleSliceStore(); - - store.queue.isSettled('setVolume'); - store.queue.isSettled('setMuted'); - }); - it('queue.reset accepts request names', () => { const store = createSingleSliceStore(); @@ -218,14 +196,6 @@ describe('store types', () => { store.queue.reset(); // all }); - it('queue.cancel accepts request names', () => { - const store = createSingleSliceStore(); - - store.queue.cancel('setVolume'); - store.queue.cancel('setMuted'); - store.queue.cancel(); // all - }); - it('queue.abort accepts request names', () => { const store = createSingleSliceStore(); @@ -234,14 +204,6 @@ describe('store types', () => { store.queue.abort(); // all }); - it('queue.flush accepts request names', () => { - const store = createSingleSliceStore(); - - store.queue.flush('setVolume'); - store.queue.flush('setMuted'); - store.queue.flush(); // all - }); - it('InferStoreTasks matches queue task record keys', () => { const _store = createSingleSliceStore(); type Tasks = InferStoreTasks; diff --git a/packages/store/src/core/tests/task.test.ts b/packages/store/src/core/tests/task.test.ts new file mode 100644 index 00000000..6f3733d5 --- /dev/null +++ b/packages/store/src/core/tests/task.test.ts @@ -0,0 +1,240 @@ +import type { ErrorTask, PendingTask, SuccessTask, Task } from '../task'; + +import { describe, expect, it } from 'vitest'; + +import { isErrorTask, isPendingTask, isSettledTask, isSuccessTask } from '../task'; + +describe('task', () => { + describe('isPendingTask', () => { + it('returns true for pending task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'pending', + abort: new AbortController(), + }; + + expect(isPendingTask(task)).toBe(true); + }); + + it('returns false for success task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'success', + settledAt: Date.now(), + output: 'result', + }; + + expect(isPendingTask(task)).toBe(false); + }); + + it('returns false for error task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'error', + settledAt: Date.now(), + error: new Error('test'), + cancelled: false, + }; + + expect(isPendingTask(task)).toBe(false); + }); + + it('returns false for undefined', () => { + expect(isPendingTask(undefined)).toBe(false); + }); + }); + + describe('isSettledTask', () => { + it('returns true for success task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'success', + settledAt: Date.now(), + output: 'result', + }; + + expect(isSettledTask(task)).toBe(true); + }); + + it('returns true for error task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'error', + settledAt: Date.now(), + error: new Error('test'), + cancelled: false, + }; + + expect(isSettledTask(task)).toBe(true); + }); + + it('returns false for pending task', () => { + const task: Task = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'pending', + abort: new AbortController(), + }; + + expect(isSettledTask(task)).toBe(false); + }); + + it('returns false for undefined', () => { + expect(isSettledTask(undefined)).toBe(false); + }); + }); + + describe('isSuccessTask', () => { + it('returns true for success task', () => { + const task: SuccessTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'success', + settledAt: Date.now(), + output: 'result', + }; + + expect(isSuccessTask(task)).toBe(true); + }); + + it('returns false for error task', () => { + const task: ErrorTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'error', + settledAt: Date.now(), + error: new Error('test'), + cancelled: false, + }; + + expect(isSuccessTask(task)).toBe(false); + }); + + it('returns false for pending task', () => { + const task: PendingTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'pending', + abort: new AbortController(), + }; + + expect(isSuccessTask(task)).toBe(false); + }); + + it('returns false for undefined', () => { + expect(isSuccessTask(undefined)).toBe(false); + }); + }); + + describe('isErrorTask', () => { + it('returns true for error task', () => { + const task: ErrorTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'error', + settledAt: Date.now(), + error: new Error('test'), + cancelled: false, + }; + + expect(isErrorTask(task)).toBe(true); + }); + + it('returns true for cancelled error task', () => { + const task: ErrorTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'error', + settledAt: Date.now(), + error: new Error('aborted'), + cancelled: true, + }; + + expect(isErrorTask(task)).toBe(true); + }); + + it('returns false for success task', () => { + const task: SuccessTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'success', + settledAt: Date.now(), + output: 'result', + }; + + expect(isErrorTask(task)).toBe(false); + }); + + it('returns false for pending task', () => { + const task: PendingTask = { + id: Symbol('task'), + name: 'test', + key: 'test', + input: undefined, + startedAt: Date.now(), + meta: null, + status: 'pending', + abort: new AbortController(), + }; + + expect(isErrorTask(task)).toBe(false); + }); + + it('returns false for undefined', () => { + expect(isErrorTask(undefined)).toBe(false); + }); + }); +}); diff --git a/packages/store/src/core/tests/task.types.test.ts b/packages/store/src/core/tests/task.types.test.ts new file mode 100644 index 00000000..90d1cc85 --- /dev/null +++ b/packages/store/src/core/tests/task.types.test.ts @@ -0,0 +1,76 @@ +import type { ErrorTask, PendingTask, SuccessTask, Task } from '../task'; + +import { describe, expectTypeOf, it } from 'vitest'; + +import { isErrorTask, isPendingTask, isSettledTask, isSuccessTask } from '../task'; + +describe('task types', () => { + describe('Task', () => { + it('is discriminated union of task states', () => { + const task: Task = {} as Task; + + if (task.status === 'pending') { + expectTypeOf(task).toExtend(); + expectTypeOf(task.abort).toExtend(); + } + + if (task.status === 'success') { + expectTypeOf(task).toExtend(); + expectTypeOf(task.output).toBeUnknown(); + expectTypeOf(task.settledAt).toBeNumber(); + } + + if (task.status === 'error') { + expectTypeOf(task).toExtend(); + expectTypeOf(task.error).toBeUnknown(); + expectTypeOf(task.cancelled).toBeBoolean(); + expectTypeOf(task.settledAt).toBeNumber(); + } + }); + + it('has common properties across all states', () => { + const task: Task = {} as Task; + + expectTypeOf(task.id).toEqualTypeOf(); + expectTypeOf(task.name).toEqualTypeOf(); + expectTypeOf(task.key).toExtend(); + expectTypeOf(task.startedAt).toBeNumber(); + }); + }); + + describe('type guards', () => { + it('isPendingTask narrows to PendingTask', () => { + const task: Task = {} as Task; + + if (isPendingTask(task)) { + expectTypeOf(task).toExtend(); + } + }); + + it('isSettledTask narrows to SuccessTask | ErrorTask', () => { + const task: Task = {} as Task; + + if (isSettledTask(task)) { + expectTypeOf(task.settledAt).toBeNumber(); + } + }); + + it('isSuccessTask narrows to SuccessTask', () => { + const task: Task = {} as Task; + + if (isSuccessTask(task)) { + expectTypeOf(task).toExtend(); + expectTypeOf(task.output).toBeUnknown(); + } + }); + + it('isErrorTask narrows to ErrorTask', () => { + const task: Task = {} as Task; + + if (isErrorTask(task)) { + expectTypeOf(task).toExtend(); + expectTypeOf(task.error).toBeUnknown(); + } + }); + }); +}); diff --git a/packages/store/src/dom/index.ts b/packages/store/src/dom/index.ts deleted file mode 100644 index 32617e9c..00000000 --- a/packages/store/src/dom/index.ts +++ /dev/null @@ -1 +0,0 @@ -export { idle, raf } from './schedulers'; diff --git a/packages/store/src/dom/schedulers.ts b/packages/store/src/dom/schedulers.ts deleted file mode 100644 index 0f15eb98..00000000 --- a/packages/store/src/dom/schedulers.ts +++ /dev/null @@ -1,27 +0,0 @@ -import type { TaskScheduler } from '../core/queue'; - -import { animationFrame, idleCallback } from '@videojs/utils/dom'; - -/** - * Scheduler using `requestAnimationFrame`. Ideal for UI updates. - * - * @example - * ```ts - * queue.enqueue({ key: 'ui', schedule: raf(), handler: async () => {} }); - * ``` - */ -export function raf(): TaskScheduler { - return flush => animationFrame(flush); -} - -/** - * Scheduler using `requestIdleCallback`. Ideal for background work. - * - * @example - * ```ts - * queue.enqueue({ key: 'bg', schedule: idle({ timeout: 2000 }), handler: async () => {} }); - * ``` - */ -export function idle(options?: IdleRequestOptions): TaskScheduler { - return flush => idleCallback(flush, options); -} diff --git a/packages/store/src/dom/tests/schedulers.test.ts b/packages/store/src/dom/tests/schedulers.test.ts deleted file mode 100644 index da7b82e1..00000000 --- a/packages/store/src/dom/tests/schedulers.test.ts +++ /dev/null @@ -1,176 +0,0 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; - -import { createQueue } from '../../core/queue'; -import { idle, raf } from '../schedulers'; - -describe('dom schedulers', () => { - beforeEach(() => { - vi.useFakeTimers(); - }); - - afterEach(() => { - vi.useRealTimers(); - }); - - describe('raf', () => { - it('creates a TaskScheduler', () => { - const scheduler = raf(); - expect(scheduler).toBeTypeOf('function'); - }); - - it('schedules flush on animation frame', async () => { - const flush = vi.fn(); - const scheduler = raf(); - - scheduler(flush); - - expect(flush).not.toHaveBeenCalled(); - await vi.runAllTimersAsync(); - expect(flush).toHaveBeenCalledOnce(); - }); - - it('returns cancel function', async () => { - const flush = vi.fn(); - const scheduler = raf(); - - const cancel = scheduler(flush); - - expect(cancel).toBeTypeOf('function'); - cancel!(); - await vi.runAllTimersAsync(); - expect(flush).not.toHaveBeenCalled(); - }); - - it('works with queue', async () => { - const queue = createQueue(); - const handler = vi.fn().mockResolvedValue('result'); - - const promise = queue.enqueue({ - name: 'raf-task', - key: 'raf', - schedule: raf(), - handler, - }); - - expect(handler).not.toHaveBeenCalled(); - await vi.runAllTimersAsync(); - await expect(promise).resolves.toBe('result'); - }); - }); - - describe('idle', () => { - it('creates a TaskScheduler', () => { - const scheduler = idle(); - expect(scheduler).toBeTypeOf('function'); - }); - - it('schedules flush when idle', async () => { - const flush = vi.fn(); - const scheduler = idle(); - - scheduler(flush); - - expect(flush).not.toHaveBeenCalled(); - await vi.runAllTimersAsync(); - expect(flush).toHaveBeenCalledOnce(); - }); - - it('returns cancel function', async () => { - const flush = vi.fn(); - const scheduler = idle(); - - const cancel = scheduler(flush); - - expect(cancel).toBeTypeOf('function'); - cancel!(); - await vi.runAllTimersAsync(); - expect(flush).not.toHaveBeenCalled(); - }); - - it('accepts options', async () => { - const flush = vi.fn(); - const scheduler = idle({ timeout: 1000 }); - - scheduler(flush); - - await vi.runAllTimersAsync(); - expect(flush).toHaveBeenCalledOnce(); - }); - - it('works with queue', async () => { - const queue = createQueue(); - const handler = vi.fn().mockResolvedValue('idle-result'); - - const promise = queue.enqueue({ - name: 'idle-task', - key: 'idle', - schedule: idle(), - handler, - }); - - expect(handler).not.toHaveBeenCalled(); - await vi.runAllTimersAsync(); - await expect(promise).resolves.toBe('idle-result'); - }); - }); - - describe('integration', () => { - it('different schedulers can coexist in same queue', async () => { - const queue = createQueue(); - const order: string[] = []; - - queue.enqueue({ - name: 'raf-task', - key: 'raf', - schedule: raf(), - handler: async () => { - order.push('raf'); - }, - }); - - queue.enqueue({ - name: 'idle-task', - key: 'idle', - schedule: idle(), - handler: async () => { - order.push('idle'); - }, - }); - - await vi.runAllTimersAsync(); - - expect(order).toContain('raf'); - expect(order).toContain('idle'); - }); - - it('superseding works with raf scheduler', async () => { - const queue = createQueue(); - const first = vi.fn().mockResolvedValue('first'); - const second = vi.fn().mockResolvedValue('second'); - - const promise1 = queue.enqueue({ - name: 'first', - key: 'shared', - schedule: raf(), - handler: first, - }); - - const promise2 = queue.enqueue({ - name: 'second', - key: 'shared', - schedule: raf(), - handler: second, - }); - - // Handle the rejection immediately to avoid unhandled rejection warning - promise1.catch(() => {}); - - await vi.runAllTimersAsync(); - - await expect(promise1).rejects.toThrow(); - await expect(promise2).resolves.toBe('second'); - expect(first).not.toHaveBeenCalled(); - expect(second).toHaveBeenCalledOnce(); - }); - }); -}); diff --git a/packages/store/src/dom/tsconfig.json b/packages/store/src/dom/tsconfig.json deleted file mode 100644 index 7525a101..00000000 --- a/packages/store/src/dom/tsconfig.json +++ /dev/null @@ -1,10 +0,0 @@ -{ - "extends": "../../../../tsconfig.base.json", - "compilerOptions": { - "composite": true, - "lib": ["ES2020", "DOM", "DOM.Iterable"], - "declarationDir": "../../types/dom" - }, - "references": [{ "path": "../.." }], - "include": ["./**/*.ts"] -} diff --git a/packages/store/src/lit/controllers/mutation-controller.ts b/packages/store/src/lit/controllers/mutation-controller.ts index b7580772..961bd025 100644 --- a/packages/store/src/lit/controllers/mutation-controller.ts +++ b/packages/store/src/lit/controllers/mutation-controller.ts @@ -1,7 +1,7 @@ import type { ReactiveController, ReactiveControllerHost } from '@lit/reactive-element'; import type { EnsureFunction } from '@videojs/utils/types'; -import type { Task } from '../../core/queue'; import type { AnyStore, InferStoreRequests } from '../../core/store'; +import type { Task } from '../../core/task'; import type { MutationResult } from '../../shared/types'; import type { StoreSource } from '../store-accessor'; diff --git a/packages/store/src/lit/controllers/optimistic-controller.ts b/packages/store/src/lit/controllers/optimistic-controller.ts index bc423262..dfc6a95b 100644 --- a/packages/store/src/lit/controllers/optimistic-controller.ts +++ b/packages/store/src/lit/controllers/optimistic-controller.ts @@ -1,7 +1,7 @@ import type { ReactiveController, ReactiveControllerHost } from '@lit/reactive-element'; import type { EnsureFunction } from '@videojs/utils/types'; -import type { Task } from '../../core/queue'; import type { AnyStore, InferStoreRequests, InferStoreState } from '../../core/store'; +import type { Task } from '../../core/task'; import type { OptimisticResult } from '../../shared/types'; import type { StoreSource } from '../store-accessor'; diff --git a/packages/store/src/react/hooks/use-mutation.ts b/packages/store/src/react/hooks/use-mutation.ts index bc47f4c5..e6cd0497 100644 --- a/packages/store/src/react/hooks/use-mutation.ts +++ b/packages/store/src/react/hooks/use-mutation.ts @@ -1,6 +1,6 @@ import type { EnsureFunction } from '@videojs/utils/types'; -import type { Task } from '../../core/queue'; import type { AnyStore, InferStoreRequests } from '../../core/store'; +import type { Task } from '../../core/task'; import type { MutationResult } from '../../shared/types'; import { useCallback, useRef, useSyncExternalStore } from 'react'; diff --git a/packages/store/src/react/hooks/use-optimistic.ts b/packages/store/src/react/hooks/use-optimistic.ts index 158d79d0..61df3ce5 100644 --- a/packages/store/src/react/hooks/use-optimistic.ts +++ b/packages/store/src/react/hooks/use-optimistic.ts @@ -1,6 +1,6 @@ import type { EnsureFunction } from '@videojs/utils/types'; -import type { Task } from '../../core/queue'; import type { AnyStore, InferStoreRequests, InferStoreState } from '../../core/store'; +import type { Task } from '../../core/task'; import type { OptimisticResult } from '../../shared/types'; import { useCallback, useReducer, useRef, useSyncExternalStore } from 'react'; diff --git a/packages/store/tsdown.config.ts b/packages/store/tsdown.config.ts index 3b13088f..01e13d70 100644 --- a/packages/store/tsdown.config.ts +++ b/packages/store/tsdown.config.ts @@ -3,7 +3,6 @@ import { defineConfig } from 'tsdown'; export default defineConfig({ entry: { index: './src/core/index.ts', - dom: './src/dom/index.ts', lit: './src/lit/index.ts', react: './src/react/index.ts', }, diff --git a/packages/store/vitest.config.ts b/packages/store/vitest.config.ts index 80163f1a..6b8f6f7b 100644 --- a/packages/store/vitest.config.ts +++ b/packages/store/vitest.config.ts @@ -11,14 +11,6 @@ export default defineConfig({ include: ['src/core/**/*.test.ts'], }, }, - { - extends: true, - test: { - name: 'store/dom', - include: ['src/dom/**/*.test.ts'], - environment: 'jsdom', - }, - }, { extends: true, test: { diff --git a/tsconfig.json b/tsconfig.json index 54b3c329..19606760 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -8,7 +8,6 @@ { "path": "packages/utils/src/dom" }, { "path": "packages/store" }, - { "path": "packages/store/src/dom" }, { "path": "packages/store/src/lit" }, { "path": "packages/store/src/react" },