Skip to content

Commit

Permalink
Browse files Browse the repository at this point in the history
  • Loading branch information
akarnokd authored Dec 20, 2016
1 parent 9b91d4e commit a902d4a
Show file tree
Hide file tree
Showing 4 changed files with 460 additions and 1 deletion.
21 changes: 21 additions & 0 deletions src/main/java/io/reactivex/Completable.java
Original file line number Diff line number Diff line change
Expand Up @@ -971,6 +971,27 @@ public final Throwable blockingGet(long timeout, TimeUnit unit) {
return observer.blockingGetError(timeout, unit);
}

/**
* Subscribes to this Completable only once, when the first CompletableObserver
* subscribes to the result Completable, caches its terminal event
* and relays/replays it to observers.
* <p>
* Note that this operator doesn't allow disposing the connection
* of the upstream source.
* <dl>
* <dt><b>Scheduler:</b></dt>
* <dd>{@code cache} does not operate by default on a particular {@link Scheduler}.</dd>
* </dl>
* @return the new Completable instance
* @since 2.0.4 - experimental
*/
@CheckReturnValue
@SchedulerSupport(SchedulerSupport.NONE)
@Experimental
public final Completable cache() {
return RxJavaPlugins.onAssembly(new CompletableCache(this));
}

/**
* Calls the given transformer function with this instance and returns the function's resulting
* Completable.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
/**
* Copyright 2016 Netflix, Inc.
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in
* compliance with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software distributed under the License is
* distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See
* the License for the specific language governing permissions and limitations under the License.
*/

package io.reactivex.internal.operators.completable;

import java.util.concurrent.atomic.*;

import io.reactivex.*;
import io.reactivex.annotations.Experimental;
import io.reactivex.disposables.Disposable;

/**
* Consume the upstream source exactly once and cache its terminal event.
*
* @since 2.0.4 - experimental
*/
@Experimental
public final class CompletableCache extends Completable implements CompletableObserver {

static final InnerCompletableCache[] EMPTY = new InnerCompletableCache[0];

static final InnerCompletableCache[] TERMINATED = new InnerCompletableCache[0];

final CompletableSource source;

final AtomicReference<InnerCompletableCache[]> observers;

final AtomicBoolean once;

Throwable error;

public CompletableCache(CompletableSource source) {
this.source = source;
this.observers = new AtomicReference<InnerCompletableCache[]>(EMPTY);
this.once = new AtomicBoolean();
}

@Override
protected void subscribeActual(CompletableObserver s) {
InnerCompletableCache inner = new InnerCompletableCache(s);
s.onSubscribe(inner);

if (add(inner)) {
if (inner.isDisposed()) {
remove(inner);
}

if (once.compareAndSet(false, true)) {
source.subscribe(this);
}
} else {
Throwable ex = error;
if (ex != null) {
s.onError(ex);
} else {
s.onComplete();
}
}
}

@Override
public void onSubscribe(Disposable d) {
// not used
}

@Override
public void onError(Throwable e) {
error = e;
for (InnerCompletableCache inner : observers.getAndSet(TERMINATED)) {
if (!inner.get()) {
inner.actual.onError(e);
}
}
}

@Override
public void onComplete() {
for (InnerCompletableCache inner : observers.getAndSet(TERMINATED)) {
if (!inner.get()) {
inner.actual.onComplete();
}
}
}

boolean add(InnerCompletableCache inner) {
for (;;) {
InnerCompletableCache[] a = observers.get();
if (a == TERMINATED) {
return false;
}
int n = a.length;
InnerCompletableCache[] b = new InnerCompletableCache[n + 1];
System.arraycopy(a, 0, b, 0, n);
b[n] = inner;
if (observers.compareAndSet(a, b)) {
return true;
}
}
}

void remove(InnerCompletableCache inner) {
for (;;) {
InnerCompletableCache[] a = observers.get();
int n = a.length;
if (n == 0) {
return;
}

int j = -1;

for (int i = 0; i < n; i++) {
if (a[i] == inner) {
j = i;
break;
}
}

if (j < 0) {
return;
}

InnerCompletableCache[] b;

if (n == 1) {
b = EMPTY;
} else {
b = new InnerCompletableCache[n - 1];
System.arraycopy(a, 0, b, 0, j);
System.arraycopy(a, j + 1, b, j, n - j - 1);
}

if (observers.compareAndSet(a, b)) {
break;
}
}
}

final class InnerCompletableCache
extends AtomicBoolean
implements Disposable {

private static final long serialVersionUID = 8943152917179642732L;

final CompletableObserver actual;

InnerCompletableCache(CompletableObserver actual) {
this.actual = actual;
}

@Override
public boolean isDisposed() {
return get();
}

@Override
public void dispose() {
if (compareAndSet(false, true)) {
remove(this);
}
}
}
}
2 changes: 1 addition & 1 deletion src/test/java/io/reactivex/JavadocWording.java
Original file line number Diff line number Diff line change
Expand Up @@ -662,7 +662,7 @@ public void completableDocRefersToCompletableTypes() throws Exception {
&& !m.signature.contains("TestObserver")) {

if (idx < 11 || !m.javadoc.substring(idx - 11, idx + 8).equals("CompletableObserver")) {
e.append("java.lang.RuntimeException: Maybe doc mentions Observer but not using Observable\r\n at io.reactivex.")
e.append("java.lang.RuntimeException: Completable doc mentions Observer but not using Observable\r\n at io.reactivex.")
.append("Completable (Completable.java:").append(m.javadocLine + lineNumber(m.javadoc, idx) - 1).append(")\r\n\r\n");
}
}
Expand Down
Loading

0 comments on commit a902d4a

Please sign in to comment.