(opts: RxResourceOptions<T, R>)
| 60 | */ |
| 61 | export function rxResource<T, R>(opts: RxResourceOptions<T, R>): ResourceRef<T | undefined>; |
| 62 | export 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); |
no test coverage detected