Supongamos que tengo una suscripción como esta:
mySubject$.subscribe(async inp => await complexFunction(inp))
Supongamos que mySubject$ es un Asunto que poseo. Necesito una manera de hacer esto de alguna manera:
await mySubject$.next(inp)
para que pueda esperar a que se complete la complexFunction .
Intenté crear una implementación de Sujeto personalizada para esto, pero no parece ser fácil. Tal vez me estoy perdiendo algo y en realidad es simple.
¿Cuáles son mis opciones aquí?
Actualizar:
Basado en la idea de que las API de rxjs deberían aceptar Promesas, probé el siguiente código. No funciona , así que de nuevo, me quedé sin ideas.
function rigObservable<T>(observable: Observable<T>) { let resolve; return { awaiter: new Promise<void>(res => resolve = res), observable: observable.pipe(tap(() => resolve())) }; } fdescribe('test', () => { it('should work', async () => { const subject = new Subject<string>(); const wrapper = rigObservable(subject); let testValue; wrapper.observable.subscribe(inp => of(inp).pipe(delay(1000), tap(t => testValue = t)).toPromise()); subject.next('text'); await wrapper.awaiter; expect(testValue).toEqual('text'); }); });Encontré una solución. No es muy simple, pero no es difícil de entender. A continuación no se muestra una implementación completa, debe agregar limpieza y controles, pero funciona.
function awaitNext(testFunc) { return async function() { Subject.prototype['saved_subscribe'] = Subject.prototype.subscribe; Subject.prototype['saved_next'] = Subject.prototype.next; Subject.prototype['promiseArray'] = []; Subject.prototype['subscribe'] = function (n, e, c) { let wrappedN, wrappedE, wrappedC; wrappedN = v => { const p = n(v); if (p instanceof Promise) this.promiseArray.push(p); } wrappedE = v => { const p = e(v); if (p instanceof Promise) this.promiseArray.push(p); } wrappedC = v => { const p = c(v); if (p instanceof Promise) this.promiseArray.push(p); } return this.saved_subscribe(wrappedN, wrappedE, wrappedC); } as any; Subject.prototype['next'] = function (v) { this.saved_next(v); return Promise.all(this.promiseArray).then(() => this.promiseArray = []); } return await testFunc(); } }Ejemplo de uso:
it('works', awaitNext(async () => { const subject = new Subject<string>(); let actualValue; subject.subscribe(v => of(v).pipe(delay(2000), tap(t => actualValue = t)).toPromise()); await subject.next('foo'); expect(actualValue).toEqual('foo'); }));No estoy exactamente seguro de tu intención, pero tal vez algo como esto funcione:
const source$ = mySubject$.pipe( mergeMap(inp => complexFunction(inp)) ); source$.subscribe( value => console.log('received value from complexFunction', value) ); source$ es un observable que emitirá el resultado de complexFunction() cada vez mySubject$ emita.