rxjs riprova con ritardo e multicast

Sep 27 2020

Sto cercando di costruire una catena rxjs che alla fine chiama una funzione di libreria che restituisce una promessa (cioè, non posso cambiare nulla oltre quella chiamata):

  libraryCall(): Promise<SomeOutput> { /* opaque to me */ }

Detta libraryCallpromessa può essere risolta o rifiutata (come qualsiasi promessa).

Se la promessa si risolve, allora tutto deve funzionare come una semplice chiamata, nient'altro aggiunto.

Se la promessa viene rifiutata, deve riprovare (chiamare di libraryCallnuovo) indefinitamente (sì, a tempo indeterminato, poiché quella promessa restituisce un'informazione molto importante), dopo un certo tempo costante (1 s per esempio).

L'osservabile risultante verrà emesso solo una volta (dopo libraryCallil rifiuto delle promesse di N libraryCalle la risoluzione della promessa di 1 ).

Se un nuovo client si iscrive all'osservabile risultante mentre è in loop, riprovando per sempre, allora dovrebbe piggyback su quel loop. Una volta risolta una promessa, tutti i clienti in sospeso dovrebbero ottenere quel valore risolto.

Come ultimo requisito, se ogni client in sospeso annulla l'iscrizione, il ciclo di nuovi tentativi dovrebbe essere interrotto.

Cosa ho fatto finora?

    defer(() => libraryCall()).pipe(
      retryWhen(errors => errors.pipe(delay(1000))),
      shareReplay(0),
    );

Sta quasi funzionando. Ma ... dopo aver risolto una volta, qualsiasi nuova sottoscrizione emette immediatamente l'ultimo valore risolto. Ho bisogno di richiamare libraryCallin questi casi.

Sono abbastanza sicuro che abbia qualcosa a che fare con shareReplay(0), anche se ho impostato il suo buffer 0sperando che nulla venga mai memorizzato nel buffer (cioè, ho bisogno della parte "condivisione", non della parte "replay" shareReplay).

Come posso modificarlo (o riscriverlo)?

Risposte

1 AndreiGătej Sep 28 2020 at 07:49

Ma ... dopo aver risolto una volta, qualsiasi nuova sottoscrizione emette immediatamente l'ultimo valore risolto

Ciò accade probabilmente a causa di come ReplaySubjectgestisce il caso quando l' bufferSizeargomento è0 :

/* ... */

this.bufferSize = Math.max(1, bufferSize);

/* ... */

Quindi, anche se stai passando 0come bufferSize, verrà impostato bufferSizesu 1, il che dovrebbe spiegare perché quando un nuovo abbonato si registra, riceverà il valore bufferizzato, invece di invocare nuovamente la funzione di libreria.

Penso che un modo rapido per risolverlo sarebbe:

defer(() => libraryCall()).pipe(
  retryWhen(errors => errors.pipe(delay(1000))),
  shareReplay({
    bufferSize: 1,
    refCount: true,
  }),
  first(),
);

Utilizzando l' refCount: trueopzione, garantisce che quando non ci sono più abbonati attivi, si riattiverà alla sorgente osservabile, richiamando così la funzione, quando un nuovo abbonato si registrerà.

Penso che aiuterebbe vedere cosa succede nel codice sorgente.

Quando l' intero stream è iscritto :

if (!subject) {
  // subscribed for the first time

  subject = new ReplaySubject<T>(bufferSize, windowTime, scheduler);

  // adding the new subscriber to the `ReplaySubject`'s subscribers list
  // (the `ReplaySubject` extends `Subject`, so this is why it also has a list of subscribers)
  innerSub = subject.subscribe(subscriber);

  // subscribing to the source - this will cause the library function to be called
  subscription = source.subscribe({
    next(value) { subject!.next(value); },
    error(err) {
      const dest = subject;
      subscription = undefined;
      subject = undefined;
      dest!.error(err);
    },
    complete() {
      subscription = undefined;
      subject!.complete();
    },
  });

  // The following condition is needed because source can complete synchronously
  // upon subscription. When that happens `subscription` is first set to `undefined`
  // and right after is set to the "closed subscription" returned by `subscribe`
  if (subscription.closed) {
    subscription = undefined;
  }
} else {
  // subscribed for the second, third etc... time

  // when other subscribers register, they will be added to the `ReplaySubject`'s subscribers list
  // so that every time the source emits, each subscriber will get the same value
  innerSub = subject.subscribe(subscriber);
}

Utilizzando first(), dopo che un abbonato riceve un valore, completesi verificherà un evento, il che significa che l'abbonato verrà rimosso dall'elenco degli iscritti del soggetto.

Per quanto riguarda l' shareReplayoperatore, questo è ciò che accade quando un iscritto viene rimosso dalla lista (a causa di una notifica complete/ error):

subscriber.add(() => {
  refCount--;
  innerSub.unsubscribe();
  if (useRefCount && refCount === 0 && subscription) {
    subscription.unsubscribe();
    subscription = undefined;
    subject = undefined;
  }
});

Come puoi vedere, a causa di refCount === 0(non ci sono più iscritti nell'elenco) e useRefCount( refCount: true), dovresti ora ottenere i risultati attesi. Se if blockviene raggiunto, subjectdiventerà undefined, il che significa che quando un nuovo abbonato si iscrive allo stream, raggiungerà il if (!subject) { ... }blocco, quindi la sorgente verrà nuovamente iscritta.

Inoltre, direi che anche questo requisito sarà soddisfatto:

se ogni client in sospeso annulla la sottoscrizione, il ciclo di nuovi tentativi dovrebbe essere interrotto.


Una nota a margine, dal momento che non hai bisogno della parte di replay , penso che potresti sostituirla shareReplay()con share(). ( EDIT: ho notato il commento troppo tardi ).