refactor(core): Synchronously emit the current signal value in toObservable (#49894)

As described in
https://github.com/angular/angular/discussions/49681#discussioncomment-5628930,
if an `Observable` created from a signal with `toObservable` is
subscribed to in a template, it will initially have `null` as the value.
Immediately after the template is done executing, effects are flushed
and this results in the `AsyncPipe` getting a new value before the
`checkNoChanges` pass, resulting in `ExpressionChanged` error.

```
template: '{{obs$ | async}}'
...
obs$ = toObservable(signal(0));
```

Instead, this commit updates the `toObservable` to synchronously emit
the initial value to the Observable stream.

Side note here: We don't exactly encourage this pattern. Instead of
using `AsyncPipe`, the template should just read signals.

PR Close #49894
This commit is contained in:
Andrew Scott
2023-04-17 12:14:18 -07:00
committed by Dylan Hunn
parent b4531f1d82
commit 5cdf6bc089
4 changed files with 88 additions and 28 deletions
+19 -16
View File
@@ -6,8 +6,8 @@
* found in the LICENSE file at https://angular.io/license
*/
import {assertInInjectionContext, effect, inject, Injector, Signal} from '@angular/core';
import {Observable} from 'rxjs';
import {assertInInjectionContext, DestroyRef, effect, EffectRef, inject, Injector, Signal, untracked} from '@angular/core';
import {Observable, ReplaySubject} from 'rxjs';
/**
* Options for `toObservable`.
@@ -38,20 +38,23 @@ export function toObservable<T>(
): Observable<T> {
!options?.injector && assertInInjectionContext(toObservable);
const injector = options?.injector ?? inject(Injector);
const subject = new ReplaySubject<T>(1);
// Creating a new `Observable` allows the creation of the effect to be lazy. This allows for all
// references to `source` to be dropped if the `Observable` is fully unsubscribed and thrown away.
return new Observable(observer => {
const watcher = effect(() => {
let value: T;
try {
value = source();
} catch (err) {
observer.error(err);
return;
}
observer.next(value);
}, {injector, manualCleanup: true, allowSignalWrites: true});
return () => watcher.destroy();
const watcher = effect(() => {
let value: T;
try {
value = source();
} catch (err) {
untracked(() => subject.error(err));
return;
}
untracked(() => subject.next(value));
}, {injector, manualCleanup: true});
injector.get(DestroyRef).onDestroy(() => {
watcher.destroy();
subject.complete();
});
return subject.asObservable();
}
@@ -6,14 +6,14 @@
* found in the LICENSE file at https://angular.io/license
*/
import {Component, computed, Injector, signal} from '@angular/core';
import {Component, computed, createEnvironmentInjector, EnvironmentInjector, Injector, Signal, signal} from '@angular/core';
import {toObservable} from '@angular/core/rxjs-interop';
import {ComponentFixture, TestBed} from '@angular/core/testing';
import {take, toArray} from 'rxjs/operators';
describe('toObservable()', () => {
let fixture!: ComponentFixture<unknown>;
let injector!: Injector;
let injector!: EnvironmentInjector;
@Component({
template: '',
@@ -24,7 +24,7 @@ describe('toObservable()', () => {
beforeEach(() => {
fixture = TestBed.createComponent(Cmp);
injector = TestBed.inject(Injector);
injector = TestBed.inject(EnvironmentInjector);
});
function flushEffects(): void {
@@ -81,7 +81,7 @@ describe('toObservable()', () => {
sub.unsubscribe();
});
it('should not monitor the signal if the Observable is never subscribed', () => {
it('monitors the signal even if the Observable is never subscribed', () => {
let counterRead = false;
const counter = computed(() => {
counterRead = true;
@@ -93,12 +93,12 @@ describe('toObservable()', () => {
// Simply creating the Observable shouldn't trigger a signal read.
expect(counterRead).toBeFalse();
// Nor should the signal be read after effects have run.
// The signal is read after effects have run.
flushEffects();
expect(counterRead).toBeFalse();
expect(counterRead).toBeTrue();
});
it('should not monitor the signal if the Observable has no active subscribers', () => {
it('should still monitor the signal if the Observable has no active subscribers', () => {
const counter = signal(0);
// Tracks how many reads of `counter()` there have been.
@@ -122,15 +122,52 @@ describe('toObservable()', () => {
flushEffects();
expect(readCount).toBe(2);
// Tear down the only subscription and hence the effect that's monitoring the signal.
// Tear down the only subscription.
sub.unsubscribe();
// Now, setting the signal shouldn't trigger any additional reads, as the Observable is no
// longer interested in its value.
// Now, setting the signal still triggers additional reads
counter.set(2);
flushEffects();
expect(readCount).toBe(3);
});
expect(readCount).toBe(2);
it('stops monitoring the signal once injector is destroyed', () => {
const counter = signal(0);
// Tracks how many reads of `counter()` there have been.
let readCount = 0;
const trackedCounter = computed(() => {
readCount++;
return counter();
});
const childInjector = createEnvironmentInjector([], injector);
toObservable(trackedCounter, {injector: childInjector});
expect(readCount).toBe(0);
flushEffects();
expect(readCount).toBe(1);
// Now, setting the signal shouldn't trigger any additional reads, as the Injector was destroyed
childInjector.destroy();
counter.set(2);
flushEffects();
expect(readCount).toBe(1);
});
it('does not track downstream signal reads in the effect', () => {
const counter = signal(0);
const emits = signal(0);
toObservable(counter, {injector}).subscribe(() => {
// Read emits. If we are still tracked in the effect, this will cause an infinite loop by
// triggering the effect again.
emits();
emits.update(v => v + 1);
});
flushEffects();
expect(emits()).toBe(1);
flushEffects();
expect(emits()).toBe(1);
});
});
+1
View File
@@ -25,6 +25,7 @@ ts_library(
"//packages/common",
"//packages/compiler",
"//packages/core",
"//packages/core/rxjs-interop",
"//packages/core/src/di/interface",
"//packages/core/src/interface",
"//packages/core/src/util",
@@ -6,7 +6,9 @@
* found in the LICENSE file at https://angular.io/license
*/
import {AsyncPipe} from '@angular/common';
import {AfterViewInit, Component, ContentChildren, createComponent, destroyPlatform, effect, EnvironmentInjector, inject, Injector, Input, NgZone, OnChanges, QueryList, signal, SimpleChanges, ViewChild} from '@angular/core';
import {toObservable} from '@angular/core/rxjs-interop';
import {TestBed} from '@angular/core/testing';
import {bootstrapApplication} from '@angular/platform-browser';
import {withBody} from '@angular/private/testing';
@@ -386,4 +388,21 @@ describe('effects', () => {
fixture.detectChanges();
expect(fixture.componentInstance.noOfCmpCreated).toBe(1);
});
it('should allow toObservable subscription in template (with async pipe)', () => {
@Component({
selector: 'test-cmp',
standalone: true,
imports: [AsyncPipe],
template: '{{counter$ | async}}',
})
class Cmp {
counter$ = toObservable(signal(0));
}
const fixture = TestBed.createComponent(Cmp);
expect(() => fixture.detectChanges(true)).not.toThrow();
fixture.detectChanges();
expect(fixture.nativeElement.textContent).toBe('0');
});
});