MCPcopy Create free account
hub / github.com/angular/angular / rxResource

Function rxResource

packages/core/rxjs-interop/src/rx_resource.ts:62–145  ·  view source on GitHub ↗
(opts: RxResourceOptions<T, R>)

Source from the content-addressed store, hash-verified

60 */
61export function rxResource<T, R>(opts: RxResourceOptions<T, R>): ResourceRef<T | undefined>;
62export function rxResource<T, R>(opts: RxResourceOptions<T, R>): ResourceRef<T | undefined> {
63 if (ngDevMode && !opts?.injector) {
64 assertInInjectionContext(rxResource);
65 }
66 return resource<T, R>({
67 ...opts,
68 loader: undefined,
69 stream: (params) => {
70 let sub: Subscription | undefined;
71
72 // `abort` can fire synchronously while the subscription is not initialized yet.
73 // Use this flag to unsubscribe immediately once `sub` exists.
74 let aborted = false;
75
76 // Start off stream as undefined.
77 const stream = signal<ResourceStreamItem<T>>({value: undefined as T});
78 const {resolve, promise} = promiseWithResolvers<Signal<ResourceStreamItem<T>>>();
79 let hasResolved = false;
80
81 function resolveOnce(): void {
82 if (!hasResolved) {
83 hasResolved = true;
84 resolve(stream);
85 }
86 }
87
88 // Track the abort listener so it can be removed if the Observable completes (as a memory
89 // optimization).
90 const onAbort = () => {
91 aborted = true;
92 sub?.unsubscribe();
93 // Remove the listener immediately since unsubscribe won't trigger the subscription's
94 // error/complete handlers. This ensures the promise resolves and PendingTask is released.
95 params.abortSignal.removeEventListener('abort', onAbort);
96 // Resolve the promise with the current stream state if it hasn't been resolved yet.
97 // This ensures the PendingTask created for this request is released.
98 resolveOnce();
99 };
100 params.abortSignal.addEventListener('abort', onAbort);
101
102 function send(value: ResourceStreamItem<T>): void {
103 stream.set(value);
104 resolveOnce();
105 }
106
107 const streamFn = opts.stream;
108 if (streamFn === undefined) {
109 throw new ɵRuntimeError(
110 ɵRuntimeErrorCode.MUST_PROVIDE_STREAM_OPTION,
111 ngDevMode && `Must provide \`stream\` option.`,
112 );
113 }
114
115 sub = streamFn(params).subscribe({
116 next: (value) => send({value}),
117 error: (error: unknown) => {
118 send({error: encapsulateResourceError(error)});
119 params.abortSignal.removeEventListener('abort', onAbort);

Callers 4

TestComponentClass · 0.90

Calls 10

assertInInjectionContextFunction · 0.90
resourceFunction · 0.90
signalFunction · 0.90
promiseWithResolversFunction · 0.90
encapsulateResourceErrorFunction · 0.90
sendFunction · 0.85
addEventListenerMethod · 0.65
subscribeMethod · 0.65
removeEventListenerMethod · 0.65
unsubscribeMethod · 0.65

Tested by

no test coverage detected