Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -16,22 +16,17 @@

package com.palantir.conjure.java.client.jaxrs;

import static org.hamcrest.Matchers.contains;
import static org.hamcrest.Matchers.hasSize;
import static org.hamcrest.Matchers.is;
import static org.hamcrest.Matchers.not;
import static org.junit.Assert.assertThat;
import static org.assertj.core.api.Assertions.assertThat;

import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
import com.palantir.conjure.java.okhttp.HostMetricsRegistry;
import com.palantir.tracing.Tracer;
import com.palantir.tracing.api.OpenSpan;
import com.palantir.tracing.api.Span;
import com.palantir.tracing.api.SpanType;
import com.palantir.tracing.api.TraceHttpHeaders;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
Expand Down Expand Up @@ -61,25 +56,40 @@ public void before() {
public void testClientIsInstrumentedWithTracer() throws InterruptedException {
server.enqueue(new MockResponse().setBody("\"server\""));
OpenSpan parentTrace = Tracer.startSpan("");
List<Map.Entry<SpanType, String>> observedSpans = Lists.newArrayList();
Tracer.subscribe(TracerTest.class.getName(),
span -> observedSpans.add(Maps.immutableEntry(span.type(), span.getOperation())));
List<Span> observedSpans = Lists.newArrayList();
Tracer.subscribe(TracerTest.class.getName(), observedSpans::add);

String traceId = Tracer.getTraceId();
service.param("somevalue");

Tracer.unsubscribe(TracerTest.class.getName());
assertThat(observedSpans, contains(
Maps.immutableEntry(SpanType.LOCAL, "OkHttp: acquire-limiter-enqueue"),
Maps.immutableEntry(SpanType.LOCAL, "OkHttp: acquire-limiter-run"),
Maps.immutableEntry(SpanType.LOCAL, "OkHttp: execute-enqueue"),
Maps.immutableEntry(SpanType.CLIENT_OUTGOING, "OkHttp: GET /{param}"),
Maps.immutableEntry(SpanType.LOCAL, "OkHttp: execute-run"),
Maps.immutableEntry(SpanType.LOCAL, "OkHttp: dispatcher")));
Span executeRunSpan = observedSpans.stream()
.filter(s -> s.getOperation().equals("OkHttp: execute-run"))
.findFirst()
.get();

assertThat(observedSpans).allSatisfy(span -> {
if (span.getOperation().equals("OkHttp: GET /{param}")) {
assertThat(span.type()).isEqualTo(SpanType.CLIENT_OUTGOING);
assertThat(span.getParentSpanId().get()).isEqualTo(executeRunSpan.getSpanId());
} else {
assertThat(span.type()).isEqualTo(SpanType.LOCAL);
assertThat(span.getParentSpanId().get()).isEqualTo(parentTrace.getSpanId());
}
});
assertThat(observedSpans)
.extracting(Span::getOperation)
.containsExactly(
"OkHttp: acquire-limiter-enqueue",
"OkHttp: acquire-limiter-run",
"ignored-span",
"OkHttp: execute-enqueue",
"OkHttp: GET /{param}",
"OkHttp: execute-run");

RecordedRequest request = server.takeRequest();
assertThat(request.getHeader(TraceHttpHeaders.TRACE_ID), is(traceId));
assertThat(request.getHeader(TraceHttpHeaders.SPAN_ID), is(not(parentTrace.getSpanId())));
assertThat(request.getHeader(TraceHttpHeaders.TRACE_ID)).isEqualTo(traceId);
assertThat(request.getHeader(TraceHttpHeaders.SPAN_ID)).isNotEqualTo(parentTrace.getSpanId());
}

@Test
Expand All @@ -89,7 +99,7 @@ public void testLimiterAcquisitionMultiThread() {
addTraceSubscriber(observedTraceIds);
runTwoRequestsInParallel();
removeTraceSubscriber();
assertThat(observedTraceIds, hasSize(2));
assertThat(observedTraceIds).containsExactly("first", "second");
}

private void runTwoRequestsInParallel() {
Expand Down Expand Up @@ -130,7 +140,7 @@ private void reduceConcurrencyLimitTo1() {
server.enqueue(new MockResponse().setResponseCode(429));
});
server.enqueue(new MockResponse().setBody("\"server\""));
assertThat(service.string(), is("server"));
assertThat(service.string()).isEqualTo("server");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
/*
* (c) Copyright 2019 Palantir Technologies Inc. All rights reserved.
*
* 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 com.palantir.conjure.java.okhttp;

import com.palantir.conjure.java.client.config.ImmutablesStyle;
import com.palantir.tracing.AsyncTracer;
import org.immutables.value.Value;

@Value.Modifiable
@ImmutablesStyle
public interface AsyncTracerTag {
AsyncTracer asyncTracer();
AsyncTracerTag setAsyncTracer(AsyncTracer asyncTracer);

static AsyncTracerTag create() {
return ModifiableAsyncTracerTag.create();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -16,19 +16,18 @@

package com.palantir.conjure.java.okhttp;

import com.palantir.tracing.AsyncTracer;
import java.io.IOException;
import okhttp3.Interceptor;
import okhttp3.Response;

public final class DispatcherTraceTerminatingInterceptor implements Interceptor {
@Override
public Response intercept(Chain chain) throws IOException {
AsyncTracer tracerTag = chain.request().tag(AsyncTracer.class);
if (tracerTag == null) {
AsyncTracerTag tracerTag = chain.request().tag(AsyncTracerTag.class);
if (tracerTag == null || tracerTag.asyncTracer() == null) {
return chain.proceed(chain.request());
}

return tracerTag.withTrace(() -> chain.proceed(chain.request()));
return tracerTag.asyncTracer().withTrace(() -> chain.proceed(chain.request()));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -81,8 +81,7 @@ public final class OkHttpClients {
* </li>
* </ol>
*/
private static final ExecutorService executionExecutor =
Tracers.wrap("OkHttp: dispatcher", Executors.newCachedThreadPool(executionThreads));
private static final ExecutorService executionExecutor = Executors.newCachedThreadPool(executionThreads);

/** Shared dispatcher with static executor service. */
private static final Dispatcher dispatcher;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import com.palantir.logsafe.exceptions.SafeIllegalStateException;
import com.palantir.logsafe.exceptions.SafeIoException;
import com.palantir.tracing.AsyncTracer;
import com.palantir.tracing.Tracer;
import java.io.IOException;
import java.io.InterruptedIOException;
import java.net.SocketTimeoutException;
Expand Down Expand Up @@ -165,8 +166,15 @@ public void enqueue(Callback callback) {
Futures.addCallback(limiterListener, new FutureCallback<Limiter.Listener>() {
@Override
public void onSuccess(Limiter.Listener listener) {
tracer.withTrace(() -> null);
enqueueInternal(callback);
tracer.withTrace(() -> {

@iamdanfox iamdanfox Apr 18, 2019

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think before merging this it would be nice to come up with a solution that doesn't require us to include the acquire-limiter-run and ignore span hacks

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That was my thought as well - this pr is here to point to intended behaviour

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was the hack here needed?
It seems like leaving the old behavior would suffice to leave the acquire-limiter-run span empty and the other spans with the appropriate parents:

tracer.withTrace(() -> null);
enqueueInternal(callback);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It won't work - the callback is executed on different thread and different trace. You need to add spans to the original trace

// terminate acquire-limiter-run span
Tracer.fastCompleteSpan();
request().tag(AsyncTracerTag.class).setAsyncTracer(new AsyncTracer("OkHttp: execute"));
enqueueInternal(callback);
// Need to recreate a span to make sure withTrace will not close parent spans
Tracer.startSpan("ignored-span");
return null;
});
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
import com.palantir.logsafe.SafeArg;
import com.palantir.logsafe.exceptions.SafeIllegalStateException;
import com.palantir.logsafe.exceptions.SafeRuntimeException;
import com.palantir.tracing.AsyncTracer;
import java.util.Optional;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService;
Expand Down Expand Up @@ -102,7 +101,7 @@ private Request createNewRequest(Request request) {
return request.newBuilder()
.url(getNewRequestUrl(request.url()))
.tag(ConcurrencyLimiterListener.class, ConcurrencyLimiterListener.create())
.tag(AsyncTracer.class, new AsyncTracer("OkHttp: execute"))
.tag(AsyncTracerTag.class, AsyncTracerTag.create())
.build();
}

Expand Down