Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 67 additions & 0 deletions goldens/public-api/core/index.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -1530,6 +1530,59 @@ export interface RendererType2 {
// @public
export function resolveForwardRef<T>(type: T): T;

// @public
export interface Resource<T> {
readonly error: Signal<unknown>;
Comment thread
hybrist marked this conversation as resolved.
Outdated
hasValue(): this is Resource<T> & {
value: Signal<T>;
};
readonly isLoading: Signal<boolean>;
reload(): boolean;
readonly status: Signal<ResourceStatus>;
readonly value: Signal<T | undefined>;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Q: Have you thought about exposing the request on the Resource, which triggered the resource?
It will add one more type parameter and bind it to the request type, but we would have one state object with all the needed information for displaying the resource.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We talked about it, but didn't really see a use case where that made sense to do. Given that the resource is initialized with the request function in the first place, the user already has both of them in the same place.

}

// @public
export function resource<T, R>(options: ResourceOptions<T, R>): ResourceRef<T>;

// @public
export type ResourceLoader<T, R> = (param: ResourceLoaderParams<R>) => PromiseLike<T>;

// @public
export interface ResourceLoaderParams<R> {
// (undocumented)
abortSignal: AbortSignal;
// (undocumented)
previous: {
status: ResourceStatus;
};
// (undocumented)
request: Exclude<NoInfer<R>, undefined>;
}

// @public
export interface ResourceOptions<T, R> {
equal?: ValueEqualityFn<T>;
injector?: Injector;
loader: ResourceLoader<T, R>;
request?: () => R;
}
Comment thread
pkozlowski-opensource marked this conversation as resolved.

// @public
export interface ResourceRef<T> extends WritableResource<T> {
destroy(): void;
}

// @public
export enum ResourceStatus {
Error = 1,
Idle = 0,
Loading = 2,
Local = 5,
Reloading = 3,
Resolved = 4
}

// @public
export function runInInjectionContext<ReturnT>(injector: Injector, fn: () => ReturnT): ReturnT;

Expand Down Expand Up @@ -1862,6 +1915,20 @@ export abstract class ViewRef extends ChangeDetectorRef {
abstract onDestroy(callback: Function): void;
}

// @public
export interface WritableResource<T> extends Resource<T> {
// (undocumented)
asReadonly(): Resource<T>;
// (undocumented)
hasValue(): this is WritableResource<T> & {
value: WritableSignal<T>;
};
set(value: T | undefined): void;
update(updater: (value: T | undefined) => T | undefined): void;
// (undocumented)
readonly value: WritableSignal<T | undefined>;
}

// @public
export interface WritableSignal<T> extends Signal<T> {
// (undocumented)
Expand Down
12 changes: 12 additions & 0 deletions goldens/public-api/core/rxjs-interop/index.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@ import { MonoTypeOperatorFunction } from 'rxjs';
import { Observable } from 'rxjs';
import { OutputOptions } from '@angular/core';
import { OutputRef } from '@angular/core';
import { ResourceLoaderParams } from '@angular/core';
import { ResourceOptions } from '@angular/core';
import { ResourceRef } from '@angular/core';
import { Signal } from '@angular/core';
import { Subscribable } from 'rxjs';
import { ValueEqualityFn } from '@angular/core/primitives/signals';
Expand All @@ -20,6 +23,15 @@ export function outputFromObservable<T>(observable: Observable<T>, opts?: Output
// @public
export function outputToObservable<T>(ref: OutputRef<T>): Observable<T>;

// @public
export function rxResource<T, R>(opts: RxResourceOptions<T, R>): ResourceRef<T>;

// @public
export interface RxResourceOptions<T, R> extends Omit<ResourceOptions<T, R>, 'loader'> {
// (undocumented)
loader: (params: ResourceLoaderParams<R>) => Observable<T>;
}

// @public
export function takeUntilDestroyed<T>(destroyRef?: DestroyRef): MonoTypeOperatorFunction<T>;

Expand Down
1 change: 1 addition & 0 deletions packages/core/rxjs-interop/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,4 @@ export {
toObservableMicrotask as ɵtoObservableMicrotask,
} from './to_observable';
export {toSignal, ToSignalOptions} from './to_signal';
export {RxResourceOptions, rxResource} from './rx_resource';
44 changes: 44 additions & 0 deletions packages/core/rxjs-interop/src/rx_resource.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
/**
* @license
* Copyright Google LLC All Rights Reserved.
*
* Use of this source code is governed by an MIT-style license that can be
* found in the LICENSE file at https://angular.dev/license
*/

import {
assertInInjectionContext,
ResourceOptions,
resource,
ResourceLoaderParams,
ResourceRef,
} from '@angular/core';
import {firstValueFrom, Observable, Subject} from 'rxjs';
import {takeUntil} from 'rxjs/operators';

/**
* Like `ResourceOptions` but uses an RxJS-based `loader`.
*
* @experimental
*/
export interface RxResourceOptions<T, R> extends Omit<ResourceOptions<T, R>, 'loader'> {
loader: (params: ResourceLoaderParams<R>) => Observable<T>;

@manfredsteyer manfredsteyer Oct 20, 2024 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Here, the loader type is "inline". The Promise-based loader type is a type of its own.

Perhaps we can do the same here. I think it would be beneficial for two reasons:

  • symmetry of code
  • easier for people writing generic helpers on top of rxResource

This could look as follows or similar:

type RxLoader = (params: ResourceLoaderParams<R>) => Observable<T>;

export interface RxResourceOptions<T, R> extends Omit<ResourceOptions<T, R>, 'loader'> {
  loader: RxLoader;
}

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hey @manfredsteyer!

We chatted about this. One of the reasons ResourceLoader is a type alias is exactly for composition - wrappers can e.g. add exponential backoff/retry.

I don't see the same need on the RxJS side, because RxJS already has a composition model of its own (operators), so extra behaviors would more likely be implemented through that mechanism.

I'm not opposed to it if there's a clear need, but we do try to keep the space of symbols intentional and not add aliases without a compelling reason.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @alxhub,

I see. Yes, this makes sense. In the case of my demo I wrote for trying out resource, I had a skipFirst and a debounce helper, wrapped around the loader.

I can totally use a respective RxJS operator. I think it's not that much about functionality but more about code-symmetry in user land code.

}

/**
* Like `resource` but uses an RxJS based `loader` which maps the request to an `Observable` of the
* resource's value. Like `firstValueFrom`, only the first emission of the Observable is considered.
*
* @experimental
*/
export function rxResource<T, R>(opts: RxResourceOptions<T, R>): ResourceRef<T> {
opts?.injector || assertInInjectionContext(rxResource);
return resource<T, R>({
...opts,
loader: (params) => {
const cancelled = new Subject<void>();
params.abortSignal.addEventListener('abort', () => cancelled.next());
return firstValueFrom(opts.loader(params).pipe(takeUntil(cancelled)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will this improve here so that all values of the observable are taken (all, until we switch over to the next request) or is the recommendation to go directly with RxJS when there is a real stream of data?

},
});
}
61 changes: 61 additions & 0 deletions packages/core/rxjs-interop/test/rx_resource_spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
/**
* @license
* Copyright Google LLC All Rights Reserved.
*
* Use of this source code is governed by an MIT-style license that can be
* found in the LICENSE file at https://angular.dev/license
*/

import {of, Observable} from 'rxjs';
import {TestBed} from '@angular/core/testing';
import {ApplicationRef, Injector, signal} from '@angular/core';
import {rxResource} from '@angular/core/rxjs-interop';

describe('rxResource()', () => {
it('should fetch data using an observable loader', async () => {
const injector = TestBed.inject(Injector);
const appRef = TestBed.inject(ApplicationRef);
const res = rxResource({
loader: () => of(1),
injector,
});
await appRef.whenStable();
expect(res.value()).toBe(1);
});

it('should cancel the fetch when a new request comes in', async () => {
const injector = TestBed.inject(Injector);
const appRef = TestBed.inject(ApplicationRef);
let unsub = false;
const request = signal(1);
const res = rxResource({
request,
loader: ({request}) =>
new Observable((sub) => {
if (request === 2) {
sub.next(true);
}
return () => {
if (request === 1) {
unsub = true;
}
};
}),
injector,
});

// Wait for the resource to reach loading state.
await waitFor(() => res.isLoading());

// Setting request = 2 should cancel request = 1
request.set(2);
await appRef.whenStable();
expect(unsub).toBe(true);
});
});

async function waitFor(fn: () => boolean): Promise<void> {
while (!fn()) {
await new Promise((resolve) => setTimeout(resolve, 1));
}
}
1 change: 1 addition & 0 deletions packages/core/src/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ export {ErrorHandler} from './error_handler';
export * from './core_private_export';
export * from './core_render3_private_export';
export * from './core_reactivity_export';
export * from './resource';
export {SecurityContext} from './sanitization/security';
export {Sanitizer} from './sanitization/sanitizer';
export {
Expand Down
8 changes: 8 additions & 0 deletions packages/core/src/pending_tasks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,10 @@ export class PendingTasksInternal implements OnDestroy {
return taskId;
}

has(taskId: number): boolean {
return this.pendingTasks.has(taskId);
}

remove(taskId: number): void {
this.pendingTasks.delete(taskId);
if (this.pendingTasks.size === 0 && this._hasPendingTasks) {
Expand Down Expand Up @@ -90,6 +94,10 @@ export class PendingTasks {
add(): () => void {
const taskId = this.internalPendingTasks.add();
return () => {
if (!this.internalPendingTasks.has(taskId)) {
// This pending task has already been cleared.
return;
}
// Notifying the scheduler will hold application stability open until the next tick.
this.scheduler.notify(NotificationSource.PendingTaskRemoved);
this.internalPendingTasks.remove(taskId);
Expand Down
Loading