Skip to content

Commit

Permalink
Fix flaky spring webflux tests (open-telemetry#3150)
Browse files Browse the repository at this point in the history
  • Loading branch information
laurit authored and robododge committed Jun 17, 2021
1 parent b00ac92 commit 18cd1a7
Showing 1 changed file with 26 additions and 6 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import io.opentelemetry.api.trace.StatusCode;
import io.opentelemetry.context.Scope;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscription;
Expand Down Expand Up @@ -79,17 +80,18 @@ static void finishSpanIfPresent(io.opentelemetry.context.Context context, Throwa

private static void finishSpanIfPresentInAttributes(
Map<String, Object> attributes, Throwable throwable) {

io.opentelemetry.context.Context context =
(io.opentelemetry.context.Context) attributes.remove(CONTEXT_ATTRIBUTE);
finishSpanIfPresent(context, throwable);
}

public static class SpanFinishingSubscriber<T> implements CoreSubscriber<T> {
public static class SpanFinishingSubscriber<T> implements CoreSubscriber<T>, Subscription {

private final CoreSubscriber<? super T> subscriber;
private final io.opentelemetry.context.Context otelContext;
private final Context context;
private final AtomicBoolean completed = new AtomicBoolean();
private Subscription subscription;

public SpanFinishingSubscriber(
CoreSubscriber<? super T> subscriber, io.opentelemetry.context.Context otelContext) {
Expand All @@ -99,9 +101,10 @@ public SpanFinishingSubscriber(
}

@Override
public void onSubscribe(Subscription s) {
public void onSubscribe(Subscription subscription) {
this.subscription = subscription;
try (Scope scope = otelContext.makeCurrent()) {
subscriber.onSubscribe(s);
subscriber.onSubscribe(this);
}
}

Expand All @@ -114,19 +117,36 @@ public void onNext(T t) {

@Override
public void onError(Throwable t) {
finishSpanIfPresent(otelContext, t);
if (completed.compareAndSet(false, true)) {
finishSpanIfPresent(otelContext, t);
}
subscriber.onError(t);
}

@Override
public void onComplete() {
finishSpanIfPresent(otelContext, null);
if (completed.compareAndSet(false, true)) {
finishSpanIfPresent(otelContext, null);
}
subscriber.onComplete();
}

@Override
public Context currentContext() {
return context;
}

@Override
public void request(long n) {
subscription.request(n);
}

@Override
public void cancel() {
if (completed.compareAndSet(false, true)) {
finishSpanIfPresent(otelContext, null);
}
subscription.cancel();
}
}
}

0 comments on commit 18cd1a7

Please sign in to comment.