packages/core/src/lib/signals/graph.ts
This is the source snapshot used to build these API details. View this revision on GitHub.
1 import type {
2 Signal,
3 SignalEquals,
4 SignalOptions,
5 WritableSignal,
6 } from './types.js';
7 import type { PibblRenderTransaction } from '../types.js';
8 import {
9 beginReactiveCollection,
10 collectReactiveDependency,
11 finishReactiveCollection,
12 linkReactiveEdgeToSource,
13 notifyReactiveEdges,
14 unlinkReactiveEdgeFromSource,
15 type ReactiveEdge,
16 type ReactiveNotificationCursor,
17 type SignalDependencyConsumer,
18 } from './edge.js';
19
20 const signalNodes = new WeakMap<object, SignalNode<unknown>>();
21 const signalValues = new WeakSet<object>();
22
23 type ComputedState = 'clean' | 'check' | 'evaluating';
24
25 interface ComputedSignalNode<T>
26 extends SignalNode<T>, SignalDependencyConsumer {
27 readonly kind: 'computed';
28 compute(): T;
29 state: ComputedState;
30 hasValue: boolean;
31 hasError: boolean;
32 error: unknown;
33 invalidate(): void;
34 validate(): void;
35 observedCount: number;
36 pendingValidation: boolean;
37 observe(): () => void;
38 adoptCandidate(
39 candidate: ComputedSignalNode<T>,
40 compute: () => T,
41 changed: boolean,
42 readers: ReadonlySet<SignalDependencyCollector>,
43 ): void;
44 replaceCompute(compute: () => T, invalidate: boolean): void;
45 }
46
47 /** @internal One isolated formula generation owned by a render transaction. */
48 export interface ComputedFormulaCandidate<T> {
49 readonly signal: Signal<T>;
50 readonly readers: Set<SignalDependencyCollector>;
51 changed: boolean;
52 evaluated: boolean;
53 value: T | undefined;
54 }
55
56 export interface SignalNode<T> {
57 value: T;
58 equals: SignalEquals<T>;
59 debugName: string | undefined;
60 version: number;
61 observerHead: ReactiveEdge | undefined;
62 observerTail: ReactiveEdge | undefined;
63 observerCount: number;
64 activeObserverCursor: ReactiveNotificationCursor | undefined;
65 observe?(): () => void;
66 validate?(): void;
67 }
68
69 type SignalRenderProvisionStatus = 'pending' | 'success' | 'failure';
70
71 interface SignalRenderProvisionEntry {
72 readonly transaction: PibblRenderTransaction;
73 value: unknown;
74 status: SignalRenderProvisionStatus;
75 }
76
77 interface SignalRenderProvisionState {
78 baseValue: unknown;
79 readonly entries: SignalRenderProvisionEntry[];
80 readonly entriesByTransaction: Map<
81 PibblRenderTransaction,
82 SignalRenderProvisionEntry
83 >;
84 }
85
86 export interface SignalDependencyCollector {
87 addDependency(node: SignalNode<unknown>): void;
88 }
89
90 export type SignalWriteOperation = 'direct' | 'staged';
91
92 export type SignalWritePolicy = (
93 node: SignalNode<unknown>,
94 operation: SignalWriteOperation,
95 ) => void;
96
97 /** @internal Authorization seam owned by the animation writer-lease registry. */
98 export interface SignalWriterLeasePolicy {
99 assertDirectWriteAllowed(node: SignalNode<unknown>): void;
100 captureStagedWrite(
101 node: SignalNode<unknown>,
102 writer: object,
103 ): unknown;
104 assertStagedWriteAllowed(
105 node: SignalNode<unknown>,
106 writer: object,
107 authorization: unknown,
108 ): void;
109 }
110
111 let activeSignalDependencyCollector: SignalDependencyCollector | undefined;
112 const writePolicyStack: SignalWritePolicy[] = [];
113 const computedStack: ComputedSignalNode<unknown>[] = [];
114 function createSignalNodeState<T>(
115 value: T,
116 equals: SignalEquals<T>,
117 debugName: string | undefined,
118 ): SignalNode<T> {
119 return {
120 value,
121 equals,
122 debugName,
123 version: 0,
124 observerHead: undefined,
125 observerTail: undefined,
126 observerCount: 0,
127 activeObserverCursor: undefined,
128 };
129 }
130
131 const commandAssertionNode = createSignalNodeState<unknown>(
132 undefined,
133 Object.is,
134 undefined,
135 );
136 const activeSignalRenderProvisions = new WeakMap<
137 SignalNode<unknown>,
138 SignalRenderProvisionState
139 >();
140 const pendingComputedValidations = new Set<ComputedSignalNode<unknown>>();
141 let flushingPendingComputedValidations = false;
142 let rendererConsumers = 0;
143 let rendererEdges = 0;
144 let createdRendererConsumers = 0;
145 let releasedRendererConsumers = 0;
146 let createdRendererEdges = 0;
147 let releasedRendererEdges = 0;
148 let signalCommitNotifications = 0;
149 let planComputedValidations = 0;
150 let computedSourceEdges = 0;
151 let createdComputedSourceEdges = 0;
152 let releasedComputedSourceEdges = 0;
153 let lazyRendererConsumers = 0;
154 let reusedRendererEdges = 0;
155 let plainInputResolutions = 0;
156 let signalInputResolutions = 0;
157 let shallowInputCopies = 0;
158 let receiverInputProvisions = 0;
159 let adoptedReceiverInputProvisions = 0;
160 let rolledBackReceiverInputProvisions = 0;
161 let batchedObserverInvalidations = 0;
162 let reactionConsumers = 0;
163 let reactionEdges = 0;
164 let reactionRuns = 0;
165 let reactionCleanups = 0;
166 let reactionFailures = 0;
167 let signalGraphPlanWorkNotifier: (() => void) | undefined;
168 let signalWriterLeasePolicy: SignalWriterLeasePolicy | undefined;
169
170 interface PublicSignalBatchEntry {
171 readonly value: unknown;
172 readonly computedHasError: boolean | undefined;
173 readonly computedError: unknown;
174 }
175
176 interface PublicSignalBatchState {
177 readonly entries: Map<SignalNode<unknown>, PublicSignalBatchEntry>;
178 readonly failures: unknown[];
179 }
180
181 let activePublicSignalBatch: PublicSignalBatchState | undefined;
182
183 export interface SignalGraphSnapshot {
184 readonly rendererConsumers: number;
185 readonly rendererEdges: number;
186 readonly computedSourceEdges: number;
187 readonly pendingComputedValidations: number;
188 readonly createdRendererConsumers: number;
189 readonly releasedRendererConsumers: number;
190 readonly createdRendererEdges: number;
191 readonly releasedRendererEdges: number;
192 readonly signalCommitNotifications: number;
193 readonly planComputedValidations: number;
194 readonly createdComputedSourceEdges: number;
195 readonly releasedComputedSourceEdges: number;
196 readonly lazyRendererConsumers: number;
197 readonly reusedRendererEdges: number;
198 readonly plainInputResolutions: number;
199 readonly signalInputResolutions: number;
200 readonly shallowInputCopies: number;
201 readonly receiverInputProvisions: number;
202 readonly adoptedReceiverInputProvisions: number;
203 readonly rolledBackReceiverInputProvisions: number;
204 readonly batchedObserverInvalidations: number;
205 readonly reactionConsumers: number;
206 readonly reactionEdges: number;
207 readonly reactionRuns: number;
208 readonly reactionCleanups: number;
209 readonly reactionFailures: number;
210 }
211
212 /** @internal Deterministic ownership and structural evidence for tests. */
213 export function signalGraphSnapshot(): SignalGraphSnapshot {
214 return Object.freeze({
215 rendererConsumers,
216 rendererEdges,
217 computedSourceEdges,
218 pendingComputedValidations: pendingComputedValidations.size,
219 createdRendererConsumers,
220 releasedRendererConsumers,
221 createdRendererEdges,
222 releasedRendererEdges,
223 signalCommitNotifications,
224 planComputedValidations,
225 createdComputedSourceEdges,
226 releasedComputedSourceEdges,
227 lazyRendererConsumers,
228 reusedRendererEdges,
229 plainInputResolutions,
230 signalInputResolutions,
231 shallowInputCopies,
232 receiverInputProvisions,
233 adoptedReceiverInputProvisions,
234 rolledBackReceiverInputProvisions,
235 batchedObserverInvalidations,
236 reactionConsumers,
237 reactionEdges,
238 reactionRuns,
239 reactionCleanups,
240 reactionFailures,
241 });
242 }
243
244 /** @internal Records one component consumer materialized by its first read. */
245 export function recordLazyRendererConsumerCreated(): void {
246 lazyRendererConsumers++;
247 }
248
249 /** @internal Records one render edge reused by recollection or restoration. */
250 export function recordRendererEdgeReused(restored = false): void {
251 if (restored) rendererEdges++;
252 reusedRendererEdges++;
253 }
254
255 /** @internal Records one declared plain input that took the identity path. */
256 export function recordPlainInputResolution(): void {
257 plainInputResolutions++;
258 }
259
260 /** @internal Records one declared signal-valued input read. */
261 export function recordSignalInputResolution(): void {
262 signalInputResolutions++;
263 }
264
265 /** @internal Records one copy-on-signal props or style object. */
266 export function recordShallowInputCopy(): void {
267 shallowInputCopies++;
268 }
269
270 /** @internal Records one pre-evaluation receiver input provision. */
271 export function recordReceiverInputProvisionCreated(): void {
272 receiverInputProvisions++;
273 }
274
275 /** @internal Records one provision adopted by its receiver generation. */
276 export function recordReceiverInputProvisionAdopted(): void {
277 receiverInputProvisions--;
278 adoptedReceiverInputProvisions++;
279 }
280
281 /** @internal Records one provision discarded without publication. */
282 export function recordReceiverInputProvisionRolledBack(): void {
283 receiverInputProvisions--;
284 rolledBackReceiverInputProvisions++;
285 }
286
287 /** @internal Adds one already-read source to the active dependency collector. */
288 export function trackSignalDependency(node: SignalNode<unknown>): void {
289 activeSignalDependencyCollector?.addDependency(node);
290 }
291
292 /** @internal Records one observer invalidation deduplicated by a public batch. */
293 export function recordBatchedObserverInvalidation(): void {
294 batchedObserverInvalidations++;
295 }
296
297 /** @internal Records one live component-owned reaction consumer. */
298 export function recordReactionConsumerCreated(): void {
299 reactionConsumers++;
300 }
301
302 /** @internal Records final release of one reaction consumer. */
303 export function recordReactionConsumerReleased(): void {
304 reactionConsumers--;
305 }
306
307 /** @internal Records one live reaction source edge. */
308 export function recordReactionEdgeCreated(): void {
309 reactionEdges++;
310 }
311
312 /** @internal Records release of one reaction source edge. */
313 export function recordReactionEdgeReleased(): void {
314 reactionEdges--;
315 }
316
317 /** @internal Records one reaction setup attempt. */
318 export function recordReactionRun(): void {
319 reactionRuns++;
320 }
321
322 /** @internal Records one exact cleanup generation execution. */
323 export function recordReactionCleanup(): void {
324 reactionCleanups++;
325 }
326
327 /** @internal Records one setup or cleanup failure. */
328 export function recordReactionFailure(): void {
329 reactionFailures++;
330 }
331
332 /** @internal Records one live renderer consumer for test diagnostics. */
333 export function recordRendererConsumerCreated(): void {
334 rendererConsumers++;
335 createdRendererConsumers++;
336 }
337
338 /** @internal Records release of one renderer consumer for test diagnostics. */
339 export function recordRendererConsumerReleased(): void {
340 rendererConsumers--;
341 releasedRendererConsumers++;
342 }
343
344 /** @internal Records one live renderer dependency edge for test diagnostics. */
345 export function recordRendererEdgeCreated(): void {
346 rendererEdges++;
347 createdRendererEdges++;
348 }
349
350 /** @internal Records release of one renderer dependency edge for tests. */
351 export function recordRendererEdgeReleased(): void {
352 rendererEdges--;
353 releasedRendererEdges++;
354 }
355
356 /**
357 * Runs a synchronous callback while coalescing signal notifications until the outer batch
358 * completes.
359 *
360 * @param callback - Synchronous work whose signal notifications are grouped.
361 * @returns The callback's return value.
362 *
363 * @see {@link signal}
364 * @see {@link computed}
365 * @see {@link WritableSignal}
366 */
367 export function batch<T>(callback: () => T): T {
368 if (activePublicSignalBatch !== undefined) return callback();
369
370 const state: PublicSignalBatchState = {
371 entries: new Map(),
372 failures: [],
373 };
374 activePublicSignalBatch = state;
375 let result!: T;
376 let callbackFailure: unknown;
377 let callbackFailed = false;
378 try {
379 result = callback();
380 } catch (error) {
381 callbackFailed = true;
382 callbackFailure = error;
383 }
384
385 const notificationFailures = flushPublicSignalBatch(state);
386 if (callbackFailed) {
387 if (notificationFailures.length === 0) throw callbackFailure;
388 throw new AggregateError(
389 [callbackFailure, ...notificationFailures],
390 'Pibbl signal batch callback and observer notifications failed.',
391 );
392 }
393 throwSignalObserverFailures(notificationFailures);
394 return result;
395 }
396
397 function recordPublicBatchEntry(node: SignalNode<unknown>): void {
398 const state = activePublicSignalBatch;
399 if (state === undefined || state.entries.has(node)) return;
400 const computedNode = isComputedSignalNode(node) ? node : undefined;
401 state.entries.set(node, {
402 value: node.value,
403 computedHasError: computedNode?.hasError,
404 computedError: computedNode?.error,
405 });
406 }
407
408 function flushPublicSignalBatch(state: PublicSignalBatchState): unknown[] {
409 try {
410 flushPendingComputedValidations();
411 } catch (error) {
412 if (
413 error instanceof AggregateError &&
414 error.message === 'Pibbl signal observer notifications failed.'
415 ) {
416 state.failures.push(...error.errors);
417 } else {
418 state.failures.push(error);
419 }
420 }
421
422 const seen = new Set<SignalDependencyConsumer>();
423 const consumers: Array<{
424 readonly consumer: SignalDependencyConsumer;
425 readonly source: SignalNode<unknown>;
426 }> = [];
427 for (const [source, entry] of state.entries) {
428 let unchanged = false;
429 try {
430 unchanged = publicBatchEntryEquals(source, entry);
431 } catch (error) {
432 state.failures.push(error);
433 }
434 if (unchanged) continue;
435 signalCommitNotifications++;
436 for (
437 let edge = source.observerHead;
438 edge !== undefined;
439 edge = edge.sourceNext
440 ) {
441 const consumer = edge.consumer;
442 if (consumer.propagateInvalidation === true) continue;
443 if (seen.has(consumer)) {
444 recordBatchedObserverInvalidation();
445 continue;
446 }
447 seen.add(consumer);
448 consumers.push({ consumer, source });
449 }
450 }
451
452 activePublicSignalBatch = undefined;
453 for (const { consumer, source } of consumers) {
454 try {
455 consumer.notify(source);
456 } catch (error) {
457 state.failures.push(error);
458 }
459 }
460 return state.failures;
461 }
462
463 function publicBatchEntryEquals(
464 node: SignalNode<unknown>,
465 entry: PublicSignalBatchEntry,
466 ): boolean {
467 if (isComputedSignalNode(node)) {
468 if (node.hasError !== entry.computedHasError) return false;
469 if (node.hasError && node.error !== entry.computedError) return false;
470 }
471 return node.equals(entry.value, node.value);
472 }
473
474 function computedWriteErrorLabel(node: SignalNode<unknown>): string {
475 return node.debugName === undefined ?
476 'Cannot write to a signal while evaluating a computed signal.' :
477 `Cannot write to a signal while evaluating computed signal "${node.debugName}".`;
478 }
479
480 function notifySignalObservers(node: SignalNode<unknown>): unknown[] {
481 const batchState = activePublicSignalBatch;
482 if (batchState !== undefined) {
483 batchState.failures.push(...notifyComputedInvalidationObservers(node));
484 return [];
485 }
486 signalCommitNotifications++;
487 return notifyReactiveEdges(node);
488 }
489
490 function notifyComputedInvalidationObservers(node: SignalNode<unknown>): unknown[] {
491 return notifyReactiveEdges(
492 node,
493 consumer => consumer.propagateInvalidation === true,
494 );
495 }
496
497 function notifyRendererInvalidationObservers(node: SignalNode<unknown>): unknown[] {
498 return notifyReactiveEdges(
499 node,
500 consumer => consumer.renderInvalidation === true,
501 );
502 }
503
504 function throwSignalObserverFailures(failures: readonly unknown[]): void {
505 if (failures.length > 0) {
506 throw new AggregateError(failures, 'Pibbl signal observer notifications failed.');
507 }
508 }
509
510 /**
511 * Creates a synchronous writable signal with configurable equality and an optional diagnostic
512 * name.
513 *
514 * @param initialValue - Value stored before the first write.
515 * @param options - Equality comparison and optional diagnostic name. See {@link SignalOptions}.
516 * @returns A writable signal holding the initial value. See {@link WritableSignal}.
517 *
518 * @see {@link SignalOptions}
519 * @see {@link WritableSignal}
520 */
521 export function signal<T>(
522 initialValue: T,
523 options: SignalOptions<T> = {},
524 ): WritableSignal<T> {
525 const equals = options.equals ?? Object.is;
526 const debugName = options.debugName;
527 const node = createSignalNodeState(initialValue, equals, debugName);
528
529 let readonlySignal: Signal<T> | undefined;
530 const read = (): T => {
531 activeSignalDependencyCollector?.addDependency(node as SignalNode<unknown>);
532 return node.value;
533 };
534 const assertWriteAllowed = (): void => {
535 writePolicyStack[writePolicyStack.length - 1]?.(
536 node as SignalNode<unknown>,
537 'direct',
538 );
539 signalWriterLeasePolicy?.assertDirectWriteAllowed(node as SignalNode<unknown>);
540 };
541 const commit = (nextValue: T): void => {
542 if (node.equals(node.value, nextValue)) {
543 return;
544 }
545
546 recordPublicBatchEntry(node as SignalNode<unknown>);
547 node.value = nextValue;
548 node.version++;
549 throwSignalObserverFailures(notifySignalObservers(node as SignalNode<unknown>));
550 };
551 const writable = {
552 get: read,
553 set(nextValue: T): void {
554 assertWriteAllowed();
555 commit(nextValue);
556 },
557 update(updateValue: (current: T) => T): void {
558 assertWriteAllowed();
559 commit(updateValue(node.value));
560 },
561 asReadonly(): Signal<T> {
562 if (readonlySignal === undefined) {
563 readonlySignal = { get: read } as Signal<T>;
564 signalValues.add(readonlySignal);
565 signalNodes.set(readonlySignal, node as SignalNode<unknown>);
566 }
567 return readonlySignal;
568 },
569 } as WritableSignal<T>;
570
571 signalValues.add(writable);
572 signalNodes.set(writable, node as SignalNode<unknown>);
573 return writable;
574 }
575
576 /**
577 * @internal Presents a provisional value to reads and computed dependencies
578 * during one render without publishing an ordinary signal notification.
579 * This is reserved for a not-yet-successfully-mounted owner whose render
580 * transaction owns acceptance or restoration of every overlapping preview.
581 */
582 export function provisionSignalValueForRender<T>(
583 value: WritableSignal<T>,
584 nextValue: T,
585 transaction: PibblRenderTransaction,
586 ): void {
587 const node = getSignalNode(value) as SignalNode<unknown>;
588 const publishComputedVersion = (next: unknown): void => {
589 node.value = next;
590 node.version++;
591 throwSignalObserverFailures(
592 notifyComputedInvalidationObservers(node),
593 );
594 };
595 let state = activeSignalRenderProvisions.get(node);
596 if (state === undefined) {
597 state = {
598 baseValue: node.value,
599 entries: [],
600 entriesByTransaction: new Map(),
601 };
602 activeSignalRenderProvisions.set(node, state);
603 }
604
605 const effectiveEntry = (): SignalRenderProvisionEntry | undefined => {
606 for (let index = state.entries.length - 1; index >= 0; index--) {
607 const entry = state.entries[index];
608 if (entry.status !== 'failure') return entry;
609 }
610 return undefined;
611 };
612 const publishEffectiveValue = (
613 forceEntry?: SignalRenderProvisionEntry,
614 ): boolean => {
615 const entry = effectiveEntry();
616 const effectiveValue = entry?.value ?? state.baseValue;
617 if (
618 (forceEntry !== undefined && entry === forceEntry) ||
619 !Object.is(node.value, effectiveValue)
620 ) {
621 publishComputedVersion(effectiveValue);
622 return true;
623 }
624 return false;
625 };
626
627 let entry = state.entriesByTransaction.get(transaction);
628 if (entry === undefined) {
629 entry = { transaction, value: nextValue, status: 'pending' };
630 state.entries.push(entry);
631 state.entriesByTransaction.set(transaction, entry);
632 const provisionState = state;
633 const provisionEntry = entry;
634 transaction.onFinish(success => {
635 if (provisionEntry.status !== 'pending') return;
636 provisionEntry.status = success ? 'success' : 'failure';
637
638 while (
639 provisionState.entries[0]?.status !== 'pending' &&
640 provisionState.entries.length > 0
641 ) {
642 const settled = provisionState.entries.shift()!;
643 provisionState.entriesByTransaction.delete(settled.transaction);
644 if (settled.status === 'success') {
645 provisionState.baseValue = settled.value;
646 }
647 }
648
649 const effectiveValueChanged = publishEffectiveValue();
650 if (!success && effectiveValueChanged) {
651 throwSignalObserverFailures(notifyRendererInvalidationObservers(node));
652 }
653 if (provisionState.entries.length === 0) {
654 activeSignalRenderProvisions.delete(node);
655 }
656 });
657 }
658 entry.value = nextValue;
659 publishEffectiveValue(entry);
660 }
661
662 /**
663 * Creates a lazy, cached readonly signal whose dependencies are tracked during evaluation.
664 *
665 * @param compute - Synchronous calculation whose signal reads become dependencies.
666 * @param options - Equality comparison and optional diagnostic name. See {@link SignalOptions}.
667 * @returns A readonly, lazy, cached signal that recomputes when its dependencies change. See
668 * {@link Signal} .
669 *
670 * @see {@link SignalOptions}
671 * @see {@link Signal}
672 */
673 export function computed<T>(
674 compute: () => T,
675 options: SignalOptions<T> = {},
676 ): Signal<T> {
677 const node = Object.assign(
678 createSignalNodeState(
679 undefined as T,
680 options.equals ?? Object.is,
681 options.debugName,
682 ),
683 {
684 kind: 'computed',
685 compute,
686 state: 'check',
687 hasValue: false,
688 hasError: false,
689 error: undefined,
690 dependencyHead: undefined,
691 dependencyTail: undefined,
692 dependencyFreeHead: undefined,
693 dependencyFreeCount: 0,
694 dependencyCount: 0,
695 dependencyHighWater: 0,
696 collectionCursor: undefined,
697 collectionGeneration: 0,
698 observedCount: 0,
699 pendingValidation: false,
700 propagateInvalidation: true,
701 addDependency(dependency: SignalNode<unknown>): void {
702 if (
703 !isComputedSignalNode(dependency) ||
704 dependency.state !== 'evaluating'
705 ) {
706 collectReactiveDependency(node, dependency);
707 }
708 },
709 notify: () => {},
710 observe: () => () => {},
711 invalidate: () => {},
712 validate: () => {},
713 adoptCandidate: () => {},
714 replaceCompute: () => {},
715 },
716 ) as ComputedSignalNode<T>;
717
718 const notifyObservers = (): unknown[] => {
719 return notifySignalObservers(node as SignalNode<unknown>);
720 };
721
722 const enqueueValidation = (): void => {
723 if (node.observedCount === 0 || node.pendingValidation) {
724 return;
725 }
726 node.pendingValidation = true;
727 const wasEmpty = pendingComputedValidations.size === 0;
728 pendingComputedValidations.add(node as ComputedSignalNode<unknown>);
729 if (wasEmpty) signalGraphPlanWorkNotifier?.();
730 };
731
732 const removePendingValidation = (): void => {
733 if (!node.pendingValidation) return;
734 node.pendingValidation = false;
735 pendingComputedValidations.delete(node as ComputedSignalNode<unknown>);
736 if (pendingComputedValidations.size === 0) signalGraphPlanWorkNotifier?.();
737 };
738
739 const invalidate = (): void => {
740 if (node.state !== 'clean') {
741 return;
742 }
743 node.state = 'check';
744 enqueueValidation();
745 throwSignalObserverFailures(
746 notifyComputedInvalidationObservers(node as SignalNode<unknown>),
747 );
748 };
749 node.notify = invalidate;
750
751 const acquireDependency = (edge: ReactiveEdge): void => {
752 if (edge.sourceLinked) return;
753 const dependency = edge.source!;
754 linkReactiveEdgeToSource(edge, dependency.observe?.());
755 computedSourceEdges++;
756 createdComputedSourceEdges++;
757 };
758
759 const releaseDependency = (edge: ReactiveEdge): void => {
760 if (!edge.sourceLinked) return;
761 unlinkReactiveEdgeFromSource(edge);
762 computedSourceEdges--;
763 releasedComputedSourceEdges++;
764 };
765
766 const acquireDependencies = (): void => {
767 for (
768 let edge = node.dependencyHead;
769 edge !== undefined;
770 edge = edge.consumerNext
771 ) {
772 acquireDependency(edge);
773 }
774 };
775
776 const releaseDependencies = (): void => {
777 for (
778 let edge = node.dependencyHead;
779 edge !== undefined;
780 edge = edge.consumerNext
781 ) {
782 releaseDependency(edge);
783 }
784 };
785
786 const finishDependencies = (): void => {
787 finishReactiveCollection(
788 node,
789 edge => {
790 if (node.observedCount > 0) acquireDependency(edge);
791 edge.observedVersion = edge.source!.version;
792 edge.collectionKind = undefined;
793 },
794 releaseDependency,
795 );
796 };
797
798 const cycleError = (): Error => {
799 const start = computedStack.indexOf(node as ComputedSignalNode<unknown>);
800 const cycle = [...computedStack.slice(start), node as ComputedSignalNode<unknown>]
801 .map(computedNode => computedNode.debugName ?? '<computed>')
802 .join(' -> ');
803 return new Error(`Pibbl computed cycle: ${cycle}`);
804 };
805
806 const evaluate = (): void => {
807 if (node.state === 'evaluating') {
808 throw cycleError();
809 }
810
811 beginReactiveCollection(node);
812 node.state = 'evaluating';
813 computedStack.push(node as ComputedSignalNode<unknown>);
814 let changed: boolean;
815 try {
816 const result = withSignalWriteForbidden(
817 computedWriteErrorLabel(node as SignalNode<unknown>),
818 () => {
819 const nextValue = withSignalTracking(node, node.compute);
820 return {
821 nextValue,
822 changed: !node.hasValue || node.hasError ||
823 !node.equals(node.value, nextValue),
824 };
825 },
826 );
827
828 changed = result.changed;
829 if (changed) {
830 recordPublicBatchEntry(node as SignalNode<unknown>);
831 node.value = result.nextValue;
832 node.version++;
833 }
834 node.hasValue = true;
835 node.hasError = false;
836 node.error = undefined;
837 } catch (error) {
838 changed = !node.hasError || node.error !== error;
839 if (changed) recordPublicBatchEntry(node as SignalNode<unknown>);
840 node.hasError = true;
841 node.error = error;
842 } finally {
843 computedStack.pop();
844 finishDependencies();
845 node.state = 'clean';
846 }
847 if (changed) {
848 if (node.hasError) node.version++;
849 throwSignalObserverFailures(notifyObservers());
850 }
851 };
852
853 const validate = (): void => {
854 if (node.state === 'evaluating') {
855 throw cycleError();
856 }
857
858 let dependencyChanged = false;
859 if (node.state === 'clean' && node.observedCount > 0) return;
860 for (
861 let edge = node.dependencyHead;
862 edge !== undefined;
863 edge = edge.consumerNext
864 ) {
865 const dependency = edge.source!;
866 dependency.validate?.();
867 if (dependency.version !== edge.observedVersion) {
868 dependencyChanged = true;
869 }
870 }
871 if (node.state === 'clean' && !dependencyChanged) return;
872 if (dependencyChanged || node.dependencyCount === 0 || !node.hasValue) {
873 evaluate();
874 } else {
875 node.state = 'clean';
876 }
877 };
878
879 const hasChangedDependencyVersion = (): boolean => {
880 for (
881 let edge = node.dependencyHead;
882 edge !== undefined;
883 edge = edge.consumerNext
884 ) {
885 if (edge.source!.version !== edge.observedVersion) return true;
886 }
887 return false;
888 };
889
890 node.invalidate = invalidate;
891 node.validate = validate;
892
893 node.replaceCompute = (nextCompute, shouldInvalidate): void => {
894 node.compute = nextCompute;
895 if (shouldInvalidate) invalidate();
896 };
897
898 node.adoptCandidate = (
899 candidate,
900 nextCompute,
901 changed,
902 readers,
903 ): void => {
904 if (!candidate.hasValue || candidate.hasError) {
905 throw new Error('Cannot adopt an unevaluated computed formula candidate.');
906 }
907 node.compute = nextCompute;
908 beginReactiveCollection(node);
909 for (
910 let edge = candidate.dependencyHead;
911 edge !== undefined;
912 edge = edge.consumerNext
913 ) {
914 collectReactiveDependency(node, edge.source!);
915 }
916 finishDependencies();
917 removePendingValidation();
918 node.state = 'clean';
919 node.hasValue = true;
920 node.hasError = false;
921 node.error = undefined;
922 if (!changed) return;
923
924 recordPublicBatchEntry(node as SignalNode<unknown>);
925 node.value = candidate.value;
926 node.version++;
927 signalCommitNotifications++;
928 throwSignalObserverFailures(notifyReactiveEdges(
929 node as SignalNode<unknown>,
930 consumer => !readers.has(consumer),
931 ));
932 };
933
934 node.observe = (): (() => void) => {
935 node.observedCount++;
936 if (node.observedCount === 1) {
937 acquireDependencies();
938 if (
939 node.state === 'clean' &&
940 hasChangedDependencyVersion()
941 ) {
942 node.state = 'check';
943 }
944 }
945 if (node.state === 'check') {
946 enqueueValidation();
947 }
948
949 let detached = false;
950 return () => {
951 if (detached) {
952 return;
953 }
954 detached = true;
955 node.observedCount--;
956 if (node.observedCount === 0) {
957 removePendingValidation();
958 releaseDependencies();
959 }
960 };
961 };
962
963 const readonly: Signal<T> = {
964 get(): T {
965 activeSignalDependencyCollector?.addDependency(
966 node as SignalNode<unknown>,
967 );
968 validate();
969 removePendingValidation();
970 if (node.hasError) {
971 throw node.error;
972 }
973 return node.value;
974 },
975 } as Signal<T>;
976
977 signalValues.add(readonly);
978 signalNodes.set(readonly, node as SignalNode<unknown>);
979 return readonly;
980 }
981
982 /** @internal Creates a nominal alias that retains one computed graph identity. */
983 export function createComputedSignalAlias<T>(
984 target: Signal<T>,
985 read: () => T,
986 ): Signal<T> {
987 const node = requireComputedSignalNode(target);
988 const alias = { get: read } as Signal<T>;
989 signalValues.add(alias);
990 signalNodes.set(alias, node as SignalNode<unknown>);
991 return alias;
992 }
993
994 /** @internal Creates one unobserved formula generation for transactional render. */
995 export function createComputedFormulaCandidate<T>(
996 target: Signal<T>,
997 compute: () => T,
998 ): ComputedFormulaCandidate<T> {
999 const node = requireComputedSignalNode(target);
1000 return {
1001 signal: computed(compute, {
1002 equals: node.equals,
1003 debugName: node.debugName,
1004 }),
1005 readers: new Set(),
1006 changed: false,
1007 evaluated: false,
1008 value: undefined,
1009 };
1010 }
1011
1012 /** @internal Reads a provisional formula while attributing the stable node. */
1013 export function readComputedFormulaCandidate<T>(
1014 target: Signal<T>,
1015 candidate: ComputedFormulaCandidate<T>,
1016 ): T {
1017 const node = requireComputedSignalNode(target);
1018 const reader = activeSignalDependencyCollector;
1019 reader?.addDependency(node as SignalNode<unknown>);
1020 if (reader !== undefined) candidate.readers.add(reader);
1021 if (!candidate.evaluated) {
1022 const value = untracked(() => candidate.signal.get());
1023 candidate.changed = !node.hasValue || node.hasError ||
1024 !node.equals(node.value, value);
1025 candidate.value = value;
1026 candidate.evaluated = true;
1027 }
1028 return candidate.value as T;
1029 }
1030
1031 /** @internal Publishes an evaluated formula generation after render success. */
1032 export function adoptComputedFormulaCandidate<T>(
1033 target: Signal<T>,
1034 candidate: ComputedFormulaCandidate<T>,
1035 compute: () => T,
1036 ): void {
1037 const node = requireComputedSignalNode(target);
1038 const candidateNode = requireComputedSignalNode(candidate.signal);
1039 node.adoptCandidate(
1040 candidateNode,
1041 compute,
1042 candidate.changed,
1043 candidate.readers,
1044 );
1045 }
1046
1047 /** @internal Refreshes a committed formula and optionally invalidates its cache. */
1048 export function replaceComputedFormula<T>(
1049 target: Signal<T>,
1050 compute: () => T,
1051 invalidate: boolean,
1052 ): void {
1053 requireComputedSignalNode(target).replaceCompute(compute, invalidate);
1054 }
1055
1056 function requireComputedSignalNode<T>(
1057 value: Signal<T>,
1058 ): ComputedSignalNode<T> {
1059 const node = getSignalNode(value);
1060 if (!isComputedSignalNode(node as SignalNode<unknown>)) {
1061 throw new TypeError('Expected a computed signal.');
1062 }
1063 return node as ComputedSignalNode<T>;
1064 }
1065
1066 /** @internal Called during the scheduler Plan phase. */
1067 export function flushPendingComputedValidations(): void {
1068 if (flushingPendingComputedValidations) {
1069 return;
1070 }
1071
1072 flushingPendingComputedValidations = true;
1073 const attempted = new Set<ComputedSignalNode<unknown>>();
1074 const failures: unknown[] = [];
1075 try {
1076 while (pendingComputedValidations.size > 0) {
1077 const node = Array.from(pendingComputedValidations)
1078 .find(candidate => !attempted.has(candidate));
1079 if (node === undefined) {
1080 break;
1081 }
1082 attempted.add(node);
1083 node.pendingValidation = false;
1084 pendingComputedValidations.delete(node);
1085 if (pendingComputedValidations.size === 0) signalGraphPlanWorkNotifier?.();
1086 if (node.observedCount > 0) {
1087 planComputedValidations++;
1088 try {
1089 node.validate();
1090 } catch (error) {
1091 if (
1092 error instanceof AggregateError &&
1093 error.message === 'Pibbl signal observer notifications failed.'
1094 ) {
1095 failures.push(...error.errors);
1096 } else {
1097 failures.push(error);
1098 }
1099 if (node.state === 'check') enqueueComputedValidation(node);
1100 }
1101 }
1102 }
1103 } finally {
1104 flushingPendingComputedValidations = false;
1105 }
1106 throwSignalObserverFailures(failures);
1107 }
1108
1109 /** @internal Whether the realm scheduler has graph Plan work to service. */
1110 export function hasPendingComputedValidations(): boolean {
1111 return pendingComputedValidations.size > 0;
1112 }
1113
1114 /** @internal Connects graph liveness to its host-agnostic scheduler owner. */
1115 export function setSignalGraphPlanWorkNotifier(notifier: () => void): void {
1116 signalGraphPlanWorkNotifier = notifier;
1117 if (pendingComputedValidations.size > 0) notifier();
1118 }
1119
1120 /**
1121 * Runs a callback without recording its signal reads as dependencies of the current consumer.
1122 *
1123 * @param callback - Synchronous work whose reads should not register dependencies.
1124 * @returns The callback's return value.
1125 *
1126 * @see {@link Signal}
1127 * @see {@link computed}
1128 */
1129 export function untracked<T>(callback: () => T): T {
1130 const previous = activeSignalDependencyCollector;
1131 activeSignalDependencyCollector = undefined;
1132 try {
1133 return callback();
1134 } finally {
1135 activeSignalDependencyCollector = previous;
1136 }
1137 }
1138
1139 /**
1140 * Checks the nominal Pibbl signal identity rather than accepting objects that merely have a get
1141 * method.
1142 *
1143 * @param value - Value to check for Pibbl's nominal signal brand.
1144 * @returns True if the value is a Pibbl signal, narrowing its type accordingly. See {@link Signal}.
1145 *
1146 * @see {@link Signal}
1147 */
1148 export function isSignal(value: unknown): value is Signal<unknown> {
1149 return typeof value === 'object' && value !== null && signalValues.has(value);
1150 }
1151
1152 export function getSignalNode<T>(value: Signal<T>): SignalNode<T> {
1153 const node =
1154 typeof value === 'object' && value !== null ? signalNodes.get(value) :
1155 undefined;
1156 if (node === undefined) {
1157 throw new TypeError('Expected a signal created by signal().');
1158 }
1159 return node as SignalNode<T>;
1160 }
1161
1162 interface StagedSignalValue {
1163 readonly node: SignalNode<unknown>;
1164 value: unknown;
1165 writer: object;
1166 authorization: unknown;
1167 provenance: object | undefined;
1168 previous: StagedSignalValue | undefined;
1169 }
1170
1171 interface SignalCommitBatchState {
1172 readonly staged: Map<SignalNode<unknown>, StagedSignalValue>;
1173 readonly lifecycles: Map<object, SignalCommitLifecycle>;
1174 }
1175
1176 /** @internal Couples a non-signal resource publication to a staged writer. */
1177 export interface SignalCommitLifecycle {
1178 /** Runs after all fallible commit validation and before signal observers. */
1179 apply(): void;
1180 /** Restores a completed apply if a later writer cannot publish. */
1181 rollback?(): void;
1182 /** Releases a candidate when its writer never reaches Commit. */
1183 abort(): void;
1184 }
1185
1186 const signalCommitBatchStates = new WeakMap<
1187 SignalCommitBatch,
1188 SignalCommitBatchState
1189 >();
1190
1191 export interface SignalCommitBatch {
1192 read<T>(signal: WritableSignal<T>): T;
1193 stage<T>(signal: WritableSignal<T>, value: T, writer: object): void;
1194 registerLifecycle(writer: object, lifecycle: SignalCommitLifecycle): void;
1195 discardWriter(writer: object): void;
1196 commit(onApplied?: () => void): void;
1197 clear(): void;
1198 }
1199
1200 export function createSignalCommitBatch(): SignalCommitBatch {
1201 const staged = new Map<SignalNode<unknown>, StagedSignalValue>();
1202
1203 const batch: SignalCommitBatch = {
1204 read<T>(value: WritableSignal<T>): T {
1205 const node = getSignalNode(value) as SignalNode<T>;
1206 const candidate = staged.get(node as SignalNode<unknown>);
1207 return candidate === undefined ? node.value : candidate.value as T;
1208 },
1209 stage<T>(value: WritableSignal<T>, nextValue: T, writer: object): void {
1210 stageSignalCommitBatchCandidate(
1211 staged,
1212 value,
1213 nextValue,
1214 writer,
1215 undefined,
1216 );
1217 },
1218 registerLifecycle(writer: object, lifecycle: SignalCommitLifecycle): void {
1219 const prior = state.lifecycles.get(writer);
1220 if (prior) abortLifecycle(prior);
1221 state.lifecycles.set(writer, lifecycle);
1222 },
1223 discardWriter(writer: object): void {
1224 for (const [node, candidate] of staged) {
1225 if (candidate.writer === writer) {
1226 staged.delete(node);
1227 }
1228 }
1229 const lifecycle = state.lifecycles.get(writer);
1230 state.lifecycles.delete(writer);
1231 if (lifecycle) abortLifecycle(lifecycle);
1232 },
1233 commit(onApplied?: () => void): void {
1234 const candidates = [...staged.values()];
1235 staged.clear();
1236 try {
1237 for (const candidate of candidates) {
1238 signalWriterLeasePolicy?.assertStagedWriteAllowed(
1239 candidate.node,
1240 candidate.writer,
1241 candidate.authorization,
1242 );
1243 }
1244
1245 const changed: StagedSignalValue[] = [];
1246 for (const candidate of candidates) {
1247 const node = candidate.node;
1248 if (!node.equals(node.value, candidate.value)) {
1249 changed.push(candidate);
1250 }
1251 }
1252
1253 const changedWriters = new Set(changed.map(candidate => candidate.writer));
1254 const applied: SignalCommitLifecycle[] = [];
1255 for (const candidate of candidates) {
1256 const lifecycle = state.lifecycles.get(candidate.writer);
1257 state.lifecycles.delete(candidate.writer);
1258 if (!lifecycle) continue;
1259 if (!changedWriters.has(candidate.writer)) {
1260 abortLifecycle(lifecycle);
1261 continue;
1262 }
1263 try {
1264 lifecycle.apply();
1265 applied.push(lifecycle);
1266 } catch (error) {
1267 abortLifecycle(lifecycle);
1268 for (const appliedLifecycle of applied.reverse()) {
1269 try { appliedLifecycle.rollback?.(); } catch {}
1270 abortLifecycle(appliedLifecycle);
1271 }
1272 throw error;
1273 }
1274 }
1275 for (const candidate of changed) {
1276 const node = candidate.node;
1277 recordPublicBatchEntry(node);
1278 node.value = candidate.value;
1279 node.version++;
1280 }
1281 onApplied?.();
1282 const failures: unknown[] = [];
1283 for (const candidate of changed) {
1284 failures.push(...notifySignalObservers(candidate.node));
1285 }
1286 for (const lifecycle of state.lifecycles.values()) abortLifecycle(lifecycle);
1287 state.lifecycles.clear();
1288 throwSignalObserverFailures(failures);
1289 } catch (error) {
1290 for (const lifecycle of state.lifecycles.values()) abortLifecycle(lifecycle);
1291 state.lifecycles.clear();
1292 throw error;
1293 }
1294 },
1295 clear(): void {
1296 staged.clear();
1297 for (const lifecycle of state.lifecycles.values()) abortLifecycle(lifecycle);
1298 state.lifecycles.clear();
1299 },
1300 };
1301 const state: SignalCommitBatchState = { staged, lifecycles: new Map() };
1302 signalCommitBatchStates.set(batch, state);
1303 return batch;
1304 }
1305
1306 function abortLifecycle(lifecycle: SignalCommitLifecycle): void {
1307 try {
1308 lifecycle.abort();
1309 } catch {
1310 // Abort is cleanup for an already failing or discarded transaction. It must
1311 // not replace the originating scheduler/commit failure.
1312 }
1313 }
1314
1315 /** @internal Stages a candidate with rollback provenance for one owner. */
1316 export function stageSignalCommitBatchWithProvenance<T>(
1317 batch: SignalCommitBatch,
1318 provenance: object,
1319 signal: WritableSignal<T>,
1320 value: T,
1321 writer: object,
1322 ): void {
1323 stageSignalCommitBatchCandidate(
1324 requireSignalCommitBatchState(batch).staged,
1325 signal,
1326 value,
1327 writer,
1328 provenance,
1329 );
1330 }
1331
1332 /** @internal Discards only candidates staged with one owner's provenance. */
1333 export function discardSignalCommitBatchProvenance(
1334 batch: SignalCommitBatch,
1335 provenance: object,
1336 ): void {
1337 discardSignalCommitBatchCandidates(
1338 requireSignalCommitBatchState(batch).staged,
1339 candidate => candidate.provenance === provenance,
1340 );
1341 }
1342
1343 /** @internal Applies friend discardWriter semantics inside one owner provenance. */
1344 export function discardSignalCommitBatchWriterForProvenance(
1345 batch: SignalCommitBatch,
1346 provenance: object,
1347 writer: object,
1348 ): void {
1349 discardSignalCommitBatchCandidates(
1350 requireSignalCommitBatchState(batch).staged,
1351 candidate =>
1352 candidate.provenance === provenance && candidate.writer === writer,
1353 );
1354 }
1355
1356 function stageSignalCommitBatchCandidate<T>(
1357 staged: Map<SignalNode<unknown>, StagedSignalValue>,
1358 value: WritableSignal<T>,
1359 nextValue: T,
1360 writer: object,
1361 provenance: object | undefined,
1362 ): void {
1363 const node = getSignalNode(value) as SignalNode<T>;
1364 const unknownNode = node as SignalNode<unknown>;
1365 writePolicyStack[writePolicyStack.length - 1]?.(unknownNode, 'staged');
1366 staged.set(unknownNode, {
1367 node: unknownNode,
1368 value: nextValue,
1369 writer,
1370 authorization: signalWriterLeasePolicy?.captureStagedWrite(
1371 unknownNode,
1372 writer,
1373 ),
1374 provenance,
1375 previous: staged.get(unknownNode),
1376 });
1377 }
1378
1379 function requireSignalCommitBatchState(
1380 batch: SignalCommitBatch,
1381 ): SignalCommitBatchState {
1382 const state = signalCommitBatchStates.get(batch);
1383 if (!state) throw new TypeError('Expected a Pibbl signal Commit batch.');
1384 return state;
1385 }
1386
1387 function discardSignalCommitBatchCandidates(
1388 staged: Map<SignalNode<unknown>, StagedSignalValue>,
1389 discard: (candidate: StagedSignalValue) => boolean,
1390 ): void {
1391 for (const [node, candidate] of staged) {
1392 const retained = discardStagedSignalCandidates(candidate, discard);
1393 if (retained === undefined) staged.delete(node);
1394 else staged.set(node, retained);
1395 }
1396 }
1397
1398 function discardStagedSignalCandidates(
1399 candidate: StagedSignalValue | undefined,
1400 discard: (candidate: StagedSignalValue) => boolean,
1401 ): StagedSignalValue | undefined {
1402 const retained: StagedSignalValue[] = [];
1403 let current = candidate;
1404 while (current !== undefined) {
1405 if (!discard(current)) retained.push(current);
1406 current = current.previous;
1407 }
1408 let previous: StagedSignalValue | undefined;
1409 for (let index = retained.length - 1; index >= 0; index--) {
1410 retained[index].previous = previous;
1411 previous = retained[index];
1412 }
1413 return previous;
1414 }
1415
1416 export function withSignalTracking<T>(
1417 collector: SignalDependencyCollector,
1418 callback: () => T,
1419 ): T {
1420 const previous = activeSignalDependencyCollector;
1421 activeSignalDependencyCollector = collector;
1422 try {
1423 return callback();
1424 } finally {
1425 activeSignalDependencyCollector = previous;
1426 }
1427 }
1428
1429 export function withSignalWritePolicy<T>(
1430 policy: SignalWritePolicy,
1431 callback: () => T,
1432 ): T {
1433 writePolicyStack.push(policy);
1434 try {
1435 return callback();
1436 } finally {
1437 writePolicyStack.pop();
1438 }
1439 }
1440
1441 export function withSignalWriteForbidden<T>(
1442 reason: string,
1443 callback: () => T,
1444 ): T {
1445 return withSignalWritePolicy(() => {
1446 throw new Error(reason);
1447 }, callback);
1448 }
1449
1450 export function withDirectSignalWriteForbidden<T>(
1451 reason: string,
1452 callback: () => T,
1453 ): T {
1454 return withSignalWritePolicy((_node, operation) => {
1455 if (operation === 'direct') throw new Error(reason);
1456 }, callback);
1457 }
1458
1459 /** @internal Checks the active direct-write policy without mutating a signal. */
1460 export function assertSignalCommandAllowed(): void {
1461 writePolicyStack[writePolicyStack.length - 1]?.(
1462 commandAssertionNode,
1463 'direct',
1464 );
1465 }
1466
1467 /** @internal Preflight a direct signal writer without changing values or versions. */
1468 export function assertDirectSignalWriteAllowed(value: Signal<unknown>): void {
1469 const node = getSignalNode(value);
1470 writePolicyStack[writePolicyStack.length - 1]?.(node, 'direct');
1471 signalWriterLeasePolicy?.assertDirectWriteAllowed(node);
1472 }
1473
1474 /** @internal Installs the singleton policy that validates active writer leases. */
1475 export function installSignalWriterLeasePolicy(
1476 policy: SignalWriterLeasePolicy,
1477 ): void {
1478 if (signalWriterLeasePolicy !== undefined && signalWriterLeasePolicy !== policy) {
1479 throw new Error('Pibbl signal writer-lease policy is already installed.');
1480 }
1481 signalWriterLeasePolicy = policy;
1482 }
1483
1484 function isComputedSignalNode(
1485 node: SignalNode<unknown>,
1486 ): node is ComputedSignalNode<unknown> {
1487 return 'kind' in node && node.kind === 'computed';
1488 }
1489
1490 function enqueueComputedValidation(node: ComputedSignalNode<unknown>): void {
1491 if (node.observedCount === 0 || node.pendingValidation) return;
1492 node.pendingValidation = true;
1493 const wasEmpty = pendingComputedValidations.size === 0;
1494 pendingComputedValidations.add(node);
1495 if (wasEmpty) signalGraphPlanWorkNotifier?.();
1496 }
1497
Documentation version
Section titled “Documentation version”Documentation built with @pibbl/core 0.0.2, revision 2dccb19. ALPHA — NOT FOR PRODUCTION USE.