1import { message as errMessage } from "@clo/lib/error.ts";
2import type { MutationClientFromConfig, Reactive } from "./client.ts";
3import type { MutationClientConfig } from "./client.ts";
4import type { Mutation, MutationEvent, RunOptions, SerializableValue } from "./types.ts";
5
6/**
7 * Argument to `defineBlocking`.
8 * @template Args - the parameters to the mutation
9 * @template Result - the result of the API call
10 * @template Config - global values and helpers from `MutationContext`
11 */
12export interface MutationOptions<
13 Args extends unknown[],
14 Result,
15 Auth extends boolean | string,
16 Config extends MutationClientConfig,
17> {
18 /**
19 * Stable identifier, unique per client; duplicate registration throws in
20 * production and replaces the previous registration in development, where
21 * hot reload re-runs `define`. Identifies the mutation in a stored
22 * `MutationAction` so a call attempted while signed out can be restored
23 * with `MutationClient.run`.
24 */
25 id: string;
26 /**
27 * Require a current user before the mutation can run. Passing one of the
28 * client's `authScopes` additionally requires that scope to read `true`.
29 */
30 auth?: Auth;
31 /** Additional guard for fine-grained auth or domain-specific availability. */
32 isAllowed?: Reactive<boolean>;
33 /**
34 * This function is only responsible for performing the underlying API call,
35 * syncronizing the optimistic state with reality. Throw on failure. A rest
36 * params type is used to allow type inference. Place this function first to
37 * ensure TypeScript correctly infers the argument type for the rest of the
38 * functions.
39 *
40 * In practice, optimistic context is never needed in this function, but it
41 * is provided as the `this` value if you truly desire it.
42 */
43 mutate: (this: Context<Config, Auth>, ...args: Args) => Promise<Result>;
44 /**
45 * Specifying the optimistic strategy is required. To disable, pass an empty
46 * function with a comment to document why it isn't needed.
47 */
48 optimistic: (
49 context: OptimisticContext<Args, Result, Config, Auth>,
50 ) => void;
51 /**
52 * Used in error messages and debug tools.
53 * Phrase it considering the template `Could not ${describe(...)}`
54 *
55 * Unlike every other function, user context is nullable: this also runs
56 * while signed out to build the `MutationAction` description given to the
57 * sign-in flow.
58 */
59 describe: string | ((context: Context<Config, false> & { args: Args }) => string);
60 /**
61 * Used in success messages.
62 * Phrase it as a complete success message: "Deleted Item"
63 */
64 describeResult:
65 | string
66 | ((context: Context<Config, Auth> & { args: Args; result: Result }) => string)
67 | null;
68 /**
69 * If the optimistic updator function is perfect, then this may be set to false.
70 * @default true
71 */
72 refetchOnSuccess?: boolean;
73 /**
74 * A key to associate related items. For example, returning a user ID. If
75 * specifying, then all mutations of the same key will evaluate in serial,
76 * but optimistic updates will apply instantly.
77 */
78 key?: (context: Context<Config, Auth> & { args: Args }) => string | string[];
79 /**
80 * Enable debouncing with "last call wins" behavior. When rapid calls arrive,
81 * the previous optimistic update is rolled back and the new one applied.
82 *
83 * All pending promises resolve with the final result.
84 */
85 debounceMs?: number;
86 /**
87 * When true, the first call executes immediately (leading edge), then
88 * subsequent rapid calls are debounced (trailing edge). After the debounce
89 * period ends, the next call executes immediately again.
90 *
91 * Requires `debounceMs` to be set.
92 */
93 debounceImmediate?: boolean;
94 /**
95 * Called before and after optimistic updates to detect no-op mutations.
96 * If the snapshots are equal (using deepEquals), the mutation is cancelled.
97 * Only `onSettled` callbacks fire, not `onSuccess` or global handlers.
98 */
99 snapshot?: (context: Context<Config, Auth> & { args: Args }) => unknown;
100}
101
102type Context<Config extends MutationClientConfig, Auth extends boolean | string> =
103 & Config["context"]
104 & (Auth extends false ? Partial<Config["userContext"]> : Config["userContext"]);
105
106export type OptimisticContext<
107 Args extends unknown[],
108 Result,
109 Config extends MutationClientConfig,
110 Auth extends boolean | string,
111> = Context<Config, Auth> & {
112 args: Args;
113 helpers: Config["optimisticHelpers"];
114 /** Add an event listener to roll back the update */
115 onRestore: (cb: () => void) => void;
116 /** Add an event listener to apply `Result` to the store. */
117 onSuccess: (cb: (result: Result) => void) => void;
118 /** Add an event listener to refetch data after mutation. */
119 onRefetch: (cb: () => Promise<void>) => void;
120};
121
122interface PendingDebouncedState<Args extends unknown[], Result, Context> {
123 /** Arguments from the most recent call */
124 args: Args;
125 /** Context captured when the most recent call was allowed */
126 context: Context;
127 /** Number of rollbacks the most recent call added */
128 rollbackCount: number;
129 /** All pending promises from all superseded calls */
130 pending: Array<{
131 resolve: (result: Result) => void;
132 reject: (error: unknown) => void;
133 }>;
134 /** Success callbacks from the most recent call */
135 onSuccess: Array<(result: Result) => void>;
136 /** Initial snapshot before any debounced calls (for no-op detection) */
137 initialSnapshot?: unknown;
138 /** Whether the last call wanted global handlers to be called */
139 shouldCallGlobalHandler: boolean;
140}
141
142/** Internal wrapper that preserves the pre-rollback description for reporting. */
143class MutationError extends Error {
144 constructor(
145 readonly error: unknown,
146 readonly description: string,
147 ) {
148 super(errMessage(error), { cause: error });
149 this.name = "MutationError";
150 }
151}
152
153function unwrapMutationError(caught: unknown) {
154 if (caught instanceof MutationError) {
155 return { error: caught.error, description: caught.description };
156 }
157 return { error: caught, description: null };
158}
159
160interface Channel<Args extends unknown[], Result, OptimisticHelpers, Context> {
161 listeners: Set<(update: MutationEvent<Result>) => void>;
162 status: "idle" | "waiting" | "mutating" | "refetching" | "skipped";
163 rollbacks: Array<() => void>;
164 refetches: Array<() => Promise<void>>;
165 queue: Array<Item<Args, Result, Context>>;
166 // Shared optimistic helpers instance for the channel
167 helpers: OptimisticHelpers | null;
168 // Debounce state (only used if debounce option is set)
169 debounceTimer: ReturnType<typeof setTimeout> | null;
170 pendingDebounced: PendingDebouncedState<Args, Result, Context> | null;
171 // Track when last debounced mutation executed (for debounceImmediate)
172 lastDebouncedExecutionTime: number | null;
173}
174
175interface Item<Args extends unknown[], Result, Context> {
176 args: Args;
177 context: Context;
178 rollbacks: number;
179 onSuccess: Array<(result: Result) => void>;
180 resolve: (result: Result) => void;
181 reject: (error: unknown) => void;
182}
183
184export class BlockingMutation<
185 Args extends unknown[],
186 Result,
187 Auth extends boolean | string,
188 Config extends MutationClientConfig,
189> implements Mutation<Args, Result> {
190 #options: MutationOptions<Args, Result, Auth, Config>;
191 #client: MutationClientFromConfig<Config>;
192 #channels: Map<
193 string,
194 Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>
195 > = new Map();
196 client: MutationClientFromConfig<Config>;
197
198 constructor(
199 client: MutationClientFromConfig<Config>,
200 options: MutationOptions<Args, Result, Auth, Config>,
201 ) {
202 this.#options = options;
203 this.#client = client;
204 this.client = client;
205 }
206
207 get id(): string {
208 return this.#options.id;
209 }
210
211 #context(): Context<Config, Auth> {
212 return {
213 ...this.#client.context,
214 ...(this.#client.userContext?.get() ?? {}),
215 } as Context<Config, Auth>;
216 }
217
218 #reportUnauthenticated(options: RunOptions<Result>, args: Args) {
219 const handler = options.onUnauthenticated ?? this.#client.handleUnauthenticated;
220 handler?.call(
221 this.#client,
222 // Reached only for `auth: true`, whose arguments the define-time guard
223 // constrains to be serializable.
224 { id: this.#options.id, description: this.describe(...args), args: args as SerializableValue[] },
225 this as unknown as Mutation<SerializableValue[], unknown>,
226 );
227 }
228
229 isUnauthenticated(): boolean {
230 return this.#options.auth === true && this.#client.userContext?.get() == null;
231 }
232
233 isAllowed(): boolean {
234 const { auth } = this.#options;
235 if (auth !== undefined && auth !== false && this.#client.userContext?.get() == null) return false;
236 if (typeof auth === "string" && !(this.#client.authScopes?.[auth]?.get() ?? false)) return false;
237 return this.#options.isAllowed?.get() ?? true;
238 }
239
240 subscribeAllowed(cb: () => void): () => void {
241 const { auth } = this.#options;
242 const unsubscribes = [
243 this.#options.isAllowed?.sub(cb),
244 typeof auth === "string" ? this.#client.authScopes?.[auth]?.sub(cb) : undefined,
245 ];
246 return () => unsubscribes.forEach((unsubscribe) => unsubscribe?.());
247 }
248
249 key(args: Args) {
250 const k = this.#options.key?.({ ...this.#context(), args })
251 ?? "shared";
252 return JSON.stringify(k);
253 }
254
255 #getOrPutChannel(key: string) {
256 let channel = this.#channels.get(key);
257 if (!channel) {
258 const rollbacks: Array<() => []> = [];
259 channel = {
260 listeners: new Set(),
261 status: "idle",
262 rollbacks,
263 refetches: [],
264 queue: [],
265 helpers: null,
266 debounceTimer: null,
267 pendingDebounced: null,
268 lastDebouncedExecutionTime: null,
269 };
270 this.#channels.set(key, channel);
271 }
272 return channel;
273 }
274
275 subscribe(
276 key: string,
277 cb: (update: MutationEvent<Result>) => void,
278 ): () => void {
279 const channel = this.#getOrPutChannel(key);
280 channel.listeners.add(cb);
281 return () => channel?.listeners.delete(cb);
282 }
283
284 #notify(
285 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
286 status: MutationEvent<Result>["status"],
287 result: Result | null = null,
288 error: unknown = null,
289 ) {
290 const event: MutationEvent<Result> = { status, result, error, debounced: this.#options.debounceMs !== undefined };
291 channel.listeners.forEach((cb) => cb(event));
292 }
293
294 #setIdle(
295 key: string,
296 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
297 ) {
298 // Check if there are pending debounced calls waiting
299 if (channel.pendingDebounced !== null) {
300 // Stay in waiting state
301 channel.status = "waiting";
302 this.#notify(channel, "waiting", null, null);
303 } else {
304 // Normal idle transition
305 channel.status = "idle";
306 // Discard any unconsumed refetch callbacks
307 channel.refetches = [];
308 this.#notify(channel, "idle", null, null);
309 // Clean up the channel if there are no listeners
310 if (channel.listeners.size === 0) {
311 // Clear any pending timers before deleting the channel
312 if (channel.debounceTimer !== null) {
313 clearTimeout(channel.debounceTimer);
314 channel.debounceTimer = null;
315 }
316 this.#channels.delete(key);
317 }
318 }
319 }
320
321 describe(...args: Args): string {
322 const { describe } = this.#options;
323 return typeof describe === "function"
324 ? describe({ ...this.#context(), args })
325 : describe;
326 }
327
328 describeResult(args: Args, result: Result): string | undefined {
329 const { describeResult } = this.#options;
330 if (describeResult === null) return undefined;
331 return typeof describeResult === "function"
332 ? describeResult({ ...this.#context(), args, result })
333 : describeResult;
334 }
335
336 /** Calling the mutation in a global scope. Errors are turned into UI toasts. */
337 run(...args: Args) {
338 this.runWithOptions(...args, {});
339 }
340
341 /**
342 * Using this in any situation is likely incorrect. Use {@linkcode runWithOptions}
343 * to handle success and error.
344 *
345 * By using this, you must handle the success and error conditions of the
346 * promise, or else the user will never see the result on screen. If the
347 * mutation has snapshots, be aware that cancelled mutations are implemented
348 * with promises that never resolve.
349 */
350 runAsHeadlessPromise(...array: [...Args, RunOptions<Result>]): Promise<Result> {
351 if (!this.#client.enabled) {
352 throw new Error(
353 "MutationClient was passed enabled: false. Are you trying to perform a mutation from SSR?",
354 );
355 }
356
357 const args = array.slice() as Args;
358 const options = args.pop() as RunOptions<Result>;
359 const { onSuccessUi: onSuccess, onSuccessData, onError, onSettled, onRestore } = options;
360
361 if (!this.isAllowed()) {
362 if (this.isUnauthenticated()) {
363 this.#reportUnauthenticated(options, args);
364 return Promise.reject(new Error("Mutation requires authentication."));
365 }
366 return Promise.reject(new Error("Mutation is not allowed."));
367 }
368
369 const promise = this.#runWithOptions(args, onRestore, true);
370 return promise.then((result) => {
371 // Call user handlers
372 onSuccess?.(result);
373 onSuccessData?.(result);
374 onSettled?.({ status: "success", result });
375 return result;
376 }).catch((caught: unknown) => {
377 const { error } = unwrapMutationError(caught);
378
379 onError?.(error);
380 onSettled?.({ status: "error", error });
381
382 throw error;
383 });
384 }
385
386 /** Calls the mutation with custom handlers that can suppress global handlers. */
387 runWithOptions(...array: [...Args, RunOptions<Result>]): void {
388 if (!this.#client.enabled) {
389 throw new Error(
390 "MutationClient was passed enabled: false. Are you trying to perform a mutation from SSR?",
391 );
392 }
393
394 const args = array.slice() as Args;
395 const options = args.pop() as RunOptions<Result>;
396 const { onSuccessUi, onSuccessData, onError, onSettled, onRestore } = options;
397 if (!this.isAllowed()) {
398 if (this.isUnauthenticated()) {
399 this.#reportUnauthenticated(options, args);
400 return;
401 }
402 const error = new Error("Mutation is not allowed.");
403 onError?.(error);
404 onSettled?.({ status: "error", error });
405 if (!onError) {
406 this.#client.reportError(formatFriendlyError(this.describe(...args), error), error);
407 }
408 return;
409 }
410
411 const suppressAll = this.#options.debounceMs !== undefined && !onSuccessUi && !onError;
412 const suppressGlobalSuccess = onSuccessUi !== undefined || suppressAll;
413 const suppressGlobalError = onError !== undefined || suppressAll;
414
415 const promise = this.#runWithOptions(args, onRestore, suppressAll);
416 promise.then((result) => {
417 // Call user handlers
418 onSuccessUi?.(result);
419 onSuccessData?.(result);
420 onSettled?.({ status: "success", result });
421
422 // Call global handler unless suppressed
423 if (!suppressGlobalSuccess) {
424 const message = this.describeResult(args, result);
425 if (message && this.#client.reportSuccess) {
426 this.#client.reportSuccess(message);
427 }
428 }
429 }).catch((caught: unknown) => {
430 const { error, description = this.describe(...args) } = unwrapMutationError(caught);
431
432 onError?.(error);
433 onSettled?.({ status: "error", error });
434
435 // Call global handler unless suppressed
436 if (!suppressGlobalError) {
437 this.#client.reportError(formatFriendlyError(description, error), error);
438 }
439 });
440 }
441
442 #runWithOptions(
443 args: Args,
444 userOnRestore: RunOptions<Result>["onRestore"],
445 suppressGlobalHandlers: boolean,
446 ): Promise<Result> {
447 if (!this.#client.enabled) {
448 throw new Error(
449 "MutationClient was passed enabled: false. Are you trying to perform a mutation from SSR?",
450 );
451 }
452 const key = this.key(args);
453 const channel = this.#getOrPutChannel(key);
454 const context = this.#context();
455
456 // Check if debouncing is enabled
457 if (this.#options.debounceMs !== undefined) {
458 // Check if we should execute immediately (leading edge)
459 const shouldExecuteImmediate = this.#options.debounceImmediate && (
460 channel.lastDebouncedExecutionTime === null
461 || Date.now() - channel.lastDebouncedExecutionTime >= this.#options.debounceMs
462 );
463
464 return this.#runDebouncedAndReturn(
465 args,
466 key,
467 channel,
468 context,
469 userOnRestore,
470 !!shouldExecuteImmediate,
471 suppressGlobalHandlers,
472 );
473 }
474
475 // Take snapshot before optimistic update (if snapshot function defined)
476 const beforeSnapshot = this.#options.snapshot
477 ? this.#options.snapshot.call(context, { ...context, args })
478 : undefined;
479
480 // Create shared optimistic helpers instance for the channel if it doesn't exist
481 if (channel.helpers === null) {
482 const onRefetch = (cb: () => Promise<void>) => {
483 channel.refetches.push(cb);
484 };
485
486 channel.helpers = this.#client.getOptimisticHelpers({
487 onRestore: (cb: () => void) => {
488 channel.rollbacks.push(cb);
489 },
490 onRefetch,
491 });
492 }
493
494 const onSuccess: Array<(result: Result) => void> = [];
495 let expired = false;
496 let rollbacks = 0;
497 const onRestore = (cb: () => void) => {
498 if (expired) {
499 throw new Error(
500 "Can only call onRestore from within the optimistic update function.",
501 );
502 }
503 channel.rollbacks.push(cb);
504 rollbacks += 1;
505 };
506
507 // Register user's onRestore callback if provided
508 if (userOnRestore) {
509 channel.rollbacks.push(userOnRestore);
510 rollbacks += 1;
511 }
512
513 try {
514 this.#options.optimistic({
515 ...context,
516 args,
517 helpers: channel.helpers,
518 onRestore,
519 onSuccess(cb) {
520 if (expired) {
521 throw new Error(
522 "Can only call onSuccess from within the optimistic update function.",
523 );
524 }
525 onSuccess.push(cb);
526 },
527 onRefetch(cb) {
528 if (expired) {
529 throw new Error(
530 "Can only call onRefetch from within the optimistic update function.",
531 );
532 }
533 channel.refetches.push(cb);
534 },
535 });
536 } catch (error) {
537 expired = true;
538 let next;
539 while (
540 next = channel.rollbacks.splice(channel.rollbacks.length - rollbacks, 1)[0]
541 ) {
542 next();
543 }
544 return Promise.reject(error);
545 }
546 expired = true;
547
548 // Take snapshot after optimistic update and check for no-op
549 if (beforeSnapshot !== undefined) {
550 const afterSnapshot = this.#options.snapshot!.call(
551 context,
552 { ...context, args },
553 );
554 const isNoOp = this.#client.deepEquals(beforeSnapshot, afterSnapshot);
555
556 if (isNoOp) {
557 let next;
558 while (
559 next = channel.rollbacks.splice(channel.rollbacks.length - rollbacks, 1)[0]
560 ) {
561 next();
562 }
563
564 // Notify listeners of skipped status
565 this.#notify(channel, "skipped");
566
567 return new Promise(() => {});
568 }
569 }
570
571 const { promise, resolve, reject } = Promise.withResolvers<Result>();
572 channel.queue.push({
573 args,
574 context,
575 rollbacks,
576 onSuccess,
577 resolve,
578 reject,
579 });
580
581 if (channel.status === "idle") {
582 this.#executeNext(key, channel);
583 }
584
585 return promise;
586 }
587
588 #executeNext(
589 key: string,
590 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
591 ) {
592 const item = channel.queue.shift();
593 if (!item) {
594 this.#setIdle(key, channel);
595 return;
596 }
597
598 const { args, context, onSuccess, resolve, reject } = item;
599 channel.status = "mutating";
600 this.#notify(channel, "mutating");
601
602 this.#options.mutate.call(context, ...args).then((result) => {
603 // remove rollbacks and apply optimistic success handlers
604 channel.rollbacks.splice(0, item.rollbacks);
605 onSuccess.forEach((cb) => cb(result));
606
607 if (this.#options.refetchOnSuccess !== false) {
608 channel.status = "refetching";
609 this.#notify(channel, "refetching", result);
610 // Call refetch and all refetch callbacks in parallel
611 const refetchCallbacks = channel.refetches.splice(0);
612 Promise.allSettled(refetchCallbacks.map((cb) => cb())).then(
613 (results) => {
614 // Report any errors from refetch or callbacks
615 results.forEach((result) => {
616 if (result.status === "rejected") {
617 const message = `Failed to refetch data: ${errMessage(result.reason)}`;
618 this.#client.reportError(message, result.reason);
619 }
620 });
621 },
622 ).finally(() => {
623 this.#executeNext(key, channel);
624 });
625 } else {
626 // Discard refetch callbacks if refetchOnSuccess is false
627 channel.refetches = [];
628 // Notify listeners with success status and result before moving to next
629 this.#notify(channel, "mutating", result);
630 this.#executeNext(key, channel);
631 }
632 resolve(result);
633 }, (error) => {
634 // Capture description BEFORE rollback so it sees optimistic state
635 const description = this.describe(...args);
636 const wrappedError = new MutationError(error, description);
637
638 // if an error happens, then every rollback is called in reverse order
639 let next;
640 while (next = channel.rollbacks.pop()) next();
641
642 // Cancel all remaining items in the channel
643 const remainingItems = channel.queue.splice(0);
644 remainingItems.forEach((queuedItem) => {
645 queuedItem.reject(wrappedError);
646 });
647
648 // Notify listeners of the error
649 this.#notify(channel, "mutating", null, error);
650
651 // Refetch to restore correct state
652 channel.status = "refetching";
653 this.#notify(channel, "refetching", null, error);
654 // Call refetch and all refetch callbacks in parallel
655 const refetchCallbacks = channel.refetches.splice(0);
656 Promise.allSettled(refetchCallbacks.map((cb) => cb())).then((results) => {
657 // Report any errors from refetch or callbacks
658 results.forEach((result) => {
659 if (result.status === "rejected") {
660 const message = `Failed to refetch data: ${errMessage(result.reason)}`;
661 this.#client.reportError(message, result.reason);
662 }
663 });
664 }).finally(() => {
665 this.#setIdle(key, channel);
666 });
667
668 reject(wrappedError);
669 });
670 }
671
672 #runDebouncedAndReturn(
673 args: Args,
674 key: string,
675 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
676 context: Context<Config, Auth>,
677 userOnRestore: (() => void) | undefined,
678 shouldExecuteImmediate: boolean,
679 shouldCallGlobalHandler: boolean,
680 ): Promise<Result> {
681 // Capture initial snapshot before first debounced call
682 const isFirstDebouncedCall = channel.pendingDebounced === null;
683 let initialSnapshot: unknown;
684 if (isFirstDebouncedCall && this.#options.snapshot) {
685 initialSnapshot = this.#options.snapshot.call(
686 context,
687 { ...context, args },
688 );
689 }
690
691 // If there's a pending debounced call, roll it back
692 if (channel.pendingDebounced) {
693 this.#rollbackPendingDebounced(channel);
694 }
695
696 // Create shared helpers if needed (same as current implementation)
697 if (channel.helpers === null) {
698 const onRefetch = (cb: () => Promise<void>) => {
699 channel.refetches.push(cb);
700 };
701 channel.helpers = this.#client.getOptimisticHelpers({
702 onRestore: (cb: () => void) => {
703 channel.rollbacks.push(cb);
704 },
705 onRefetch,
706 });
707 }
708
709 // Apply optimistic update (same logic as current runAndReturn)
710 const onSuccess: Array<(result: Result) => void> = [];
711 let expired = false;
712 let rollbacks = 0;
713 const onRestore = (cb: () => void) => {
714 if (expired) {
715 throw new Error(
716 "Can only call onRestore from within the optimistic update function.",
717 );
718 }
719 channel.rollbacks.push(cb);
720 rollbacks += 1;
721 };
722
723 // Register user's onRestore callback if provided
724 if (userOnRestore) {
725 channel.rollbacks.push(userOnRestore);
726 rollbacks += 1;
727 }
728
729 try {
730 this.#options.optimistic({
731 ...context,
732 args,
733 helpers: channel.helpers,
734 onRestore,
735 onSuccess(cb) {
736 if (expired) {
737 throw new Error(
738 "Can only call onSuccess from within the optimistic update function.",
739 );
740 }
741 onSuccess.push(cb);
742 },
743 onRefetch(cb) {
744 if (expired) {
745 throw new Error(
746 "Can only call onRefetch from within the optimistic update function.",
747 );
748 }
749 channel.refetches.push(cb);
750 },
751 });
752 } catch (error) {
753 expired = true;
754 // Roll back the rollbacks we just added
755 let next;
756 while (
757 next = channel.rollbacks.splice(channel.rollbacks.length - rollbacks, 1)[0]
758 ) {
759 next();
760 }
761 return Promise.reject(error);
762 }
763 expired = true;
764
765 // Check for no-op by comparing to initial snapshot
766 if (this.#options.snapshot) {
767 const currentSnapshot = this.#options.snapshot.call(
768 context,
769 { ...context, args },
770 );
771 const snapshotToCompare = isFirstDebouncedCall
772 ? initialSnapshot!
773 : channel.pendingDebounced?.initialSnapshot;
774
775 if (
776 snapshotToCompare !== undefined
777 && this.#client.deepEquals(snapshotToCompare, currentSnapshot)
778 ) {
779 // No-op detected - rollback optimistic update
780 let next;
781 while (
782 next = channel.rollbacks.splice(channel.rollbacks.length - rollbacks, 1)[0]
783 ) {
784 next();
785 }
786
787 // Clear debounce timer
788 if (channel.debounceTimer !== null) {
789 clearTimeout(channel.debounceTimer);
790 channel.debounceTimer = null;
791 }
792
793 // Clear pending debounced state
794 channel.pendingDebounced = null;
795
796 // Notify listeners of skipped status
797 this.#notify(channel, "skipped");
798
799 // Return resolved promise
800 return new Promise(() => {});
801 }
802 }
803
804 // Create promise for this call
805 const { promise, resolve, reject } = Promise.withResolvers<Result>();
806
807 // Store or update pending debounced state
808 if (channel.pendingDebounced === null) {
809 // First debounced call
810 channel.pendingDebounced = {
811 args,
812 context,
813 rollbackCount: rollbacks,
814 pending: [{ resolve, reject }],
815 onSuccess,
816 initialSnapshot: isFirstDebouncedCall ? initialSnapshot : undefined,
817 shouldCallGlobalHandler,
818 };
819
820 // Set status to waiting
821 channel.status = "waiting";
822 this.#notify(channel, "waiting");
823 } else {
824 // Subsequent debounced call - update state
825 channel.pendingDebounced.args = args;
826 channel.pendingDebounced.context = context;
827 channel.pendingDebounced.rollbackCount = rollbacks;
828 channel.pendingDebounced.pending.push({ resolve, reject });
829 channel.pendingDebounced.onSuccess = onSuccess;
830 channel.pendingDebounced.shouldCallGlobalHandler = shouldCallGlobalHandler;
831 // Keep the initial snapshot from the first call
832 // Status stays "waiting"
833 }
834
835 // Clear existing timer
836 if (channel.debounceTimer !== null) {
837 clearTimeout(channel.debounceTimer);
838 }
839
840 // Start new timer (0ms for immediate execution, debounceMs otherwise)
841 const delay = shouldExecuteImmediate ? 0 : this.#options.debounceMs!;
842 channel.debounceTimer = setTimeout(() => {
843 this.#enqueueDebouncedCall(key, channel);
844 }, delay);
845
846 return promise;
847 }
848
849 #rollbackPendingDebounced(
850 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
851 ) {
852 if (!channel.pendingDebounced) return;
853
854 const { rollbackCount } = channel.pendingDebounced;
855
856 // Roll back this call's optimistic updates (in reverse order)
857 // Remove from the end of the rollbacks array
858 for (let i = 0; i < rollbackCount; i++) {
859 const rollback = channel.rollbacks.pop();
860 if (rollback) rollback();
861 }
862
863 // Note: We do NOT reject the promises here
864 // They will all resolve when the final call completes
865 }
866
867 #enqueueDebouncedCall(
868 key: string,
869 channel: Channel<Args, Result, Config["optimisticHelpers"], Context<Config, Auth>>,
870 ) {
871 // Clear timer
872 channel.debounceTimer = null;
873
874 // Safety check
875 if (!channel.pendingDebounced) {
876 this.#setIdle(key, channel);
877 return;
878 }
879
880 const { args, context, rollbackCount, pending, onSuccess, shouldCallGlobalHandler } = channel.pendingDebounced;
881 channel.pendingDebounced = null;
882
883 // Track execution time for debounceImmediate
884 if (this.#options.debounceImmediate) {
885 channel.lastDebouncedExecutionTime = Date.now();
886 }
887
888 // Create wrapper resolve/reject that resolves ALL pending promises
889 const {
890 promise: wrapperPromise,
891 resolve: wrapperResolve,
892 reject: wrapperReject,
893 } = Promise.withResolvers<Result>();
894
895 // Resolve/reject pending promises and add global handler logic for execution-time checks
896 wrapperPromise.then(
897 (result) => {
898 // Resolve all pending promises
899 pending.forEach((p) => p.resolve(result));
900
901 // Call global handler if needed (based on whether component is watching success)
902 if (shouldCallGlobalHandler) {
903 const message = this.describeResult(args, result);
904 if (message && this.#client.reportSuccess) {
905 this.#client.reportSuccess(message);
906 }
907 }
908 },
909 (caught) => {
910 const { error, description = this.describe(...args) } = unwrapMutationError(caught);
911
912 // Reject all pending promises with unwrapped error
913 pending.forEach((p) => p.reject(error));
914
915 // Call global handler if needed (based on whether component is watching errors)
916 if (shouldCallGlobalHandler) {
917 this.#client.reportError(
918 formatFriendlyError(description, error),
919 error,
920 );
921 }
922 },
923 );
924
925 // Add to queue (same structure as regular blocking mutation)
926 channel.queue.push({
927 args,
928 context,
929 rollbacks: rollbackCount,
930 onSuccess,
931 resolve: wrapperResolve,
932 reject: wrapperReject,
933 });
934
935 // If queue was idle/waiting, start execution
936 if (channel.status === "idle" || channel.status === "waiting") {
937 this.#executeNext(key, channel);
938 }
939 // Otherwise, it will execute when the current item finishes
940 }
941}
942
943export function formatFriendlyError(
944 description: string | null,
945 error: unknown,
946) {
947 if (!description || !description[0]) return `Something went wrong: ${errMessage(error)}`;
948 description = description[0].toLowerCase() + description.slice(1);
949 return `Could not ${description}: ${errMessage(error)}`;
950}