Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
52065a5
tracing demo project
Aug 29, 2019
ae4edb9
Scheduling backoffs
iamdanfox Aug 29, 2019
ce2a390
more tests
Aug 29, 2019
5039907
cleanup
Aug 29, 2019
e7182eb
test transformed futures
Aug 29, 2019
25e27ee
fix conccurent equality check
Aug 29, 2019
96c4a10
WIP?
iamdanfox Aug 29, 2019
0b1897a
Utility methods
iamdanfox Aug 29, 2019
4935fd8
Merge branch 'fo/async-tracing-demos' into dfox/cross-thread-tracing
iamdanfox Aug 29, 2019
c789947
Try it out with some demos
iamdanfox Aug 29, 2019
c118b8c
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
iamdanfox Aug 29, 2019
30d40e4
TracingDemos all work
iamdanfox Aug 29, 2019
d322009
World before rob
iamdanfox Aug 29, 2019
abb9ff4
Rename methods, javadoc, rob++
iamdanfox Aug 29, 2019
96ee230
Javadoc
iamdanfox Aug 29, 2019
4e2d1a7
Go away errorprone
iamdanfox Aug 29, 2019
43e6b2e
Merge remote-tracking branch 'origin/develop' into dfox/cross-thread-…
Aug 30, 2019
cd274c3
Thanks checkstyle
iamdanfox Aug 30, 2019
1685527
Smaller diff
iamdanfox Aug 30, 2019
6b40aa8
Add generated changelog entries
iamdanfox Aug 30, 2019
36a640c
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
b6366ec
change some names
Aug 30, 2019
b510442
add CloseableSpan, update test cases
Aug 30, 2019
00df1a2
revert unnecessary API expansion
Aug 30, 2019
d3e4725
Merge remote-tracking branch 'origin/ckozak/detached_tracing_api' int…
Aug 30, 2019
5e5911f
compile
Aug 30, 2019
3d4222a
check
Aug 30, 2019
fc93f5c
javadoc
iamdanfox Aug 30, 2019
572cc58
Delete pr-247.v2.yml
iamdanfox Aug 30, 2019
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
177 changes: 111 additions & 66 deletions tracing-demos/src/test/java/com/palantir/tracing/TracingDemos.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,10 +23,10 @@
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.PriorityBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.IntStream;
Expand All@@ -48,16 +48,17 @@ void handles_async_spans() throws Exception {
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
CountDownLatch countDownLatch = new CountDownLatch(numTasks);

try (CloseableTracer root = CloseableTracer.startSpan("root")) {
IntStream.range(0, numTasks).forEach(i -> {
// DetachedSpan detachedSpan = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
// detachedSpan.close();
IntStream.range(0, numTasks).forEach(i -> {
Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan crossThread = DetachedSpan.start("task-queue-time" + i);
executorService.submit(() -> {
try (CloseableSpan t = crossThread.completeAndStartChild("task" + i)) {
emit_nested_spans();
countDownLatch.countDown();
});
}
});
}
});

assertThat(countDownLatch.await(expectedDurationMillis + 1000, TimeUnit.MILLISECONDS)).isTrue();
}
Expand All@@ -67,29 +68,42 @@ void handles_async_spans() throws Exception {
void async_future() throws InterruptedException {
int numThreads = 2;
int numCallbacks = 10;
ExecutorService executorService = Tracers.wrap(Executors.newFixedThreadPool(numThreads));
ExecutorService executorService = Executors.newFixedThreadPool(numThreads);
final SettableFuture<Object> future = SettableFuture.create();
CountDownLatch latch = new CountDownLatch(numCallbacks);

IntStream.range(0, numCallbacks).forEach(i ->
try (CloseableTracer tracer = CloseableTracer.startSpan("I am a root span")) {
String traceId = Tracer.getTraceId();

IntStream.range(0, numCallbacks).forEach(i -> {

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");

Futures.addCallback(future, new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
sleep(10, "success" + i);
latch.countDown();
assertThat(Tracer.hasTraceId()).isFalse();
try (CloseableSpan tracer = span.completeAndStartChild("success" + i)) {
assertThat(Tracer.getTraceId()).isEqualTo(traceId);
sleep(10);
latch.countDown();
}
}

@Override
public void onFailure(Throwable throwable) {
Assertions.fail();
}
}, executorService));
}, executorService);
});

executorService.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
try (CloseableTracer root = CloseableTracer.startSpan("bbb")) {
executorService.submit(() -> {
future.set(null);
});
}
});

}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

Expand All@@ -98,28 +112,42 @@ public void onFailure(Throwable throwable) {
void multi_producer_single_consumer() throws InterruptedException {
int numProducers = 2;
int numElem = 20;
PriorityBlockingQueue<String> work = new PriorityBlockingQueue<>();
ArrayBlockingQueue<QueuedWork> work = new ArrayBlockingQueue<QueuedWork>(numElem);

CountDownLatch submitLatch = new CountDownLatch(numElem);
CountDownLatch consumeLatch = new CountDownLatch(numElem);
ExecutorService producerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(numProducers));
ExecutorService consumerExecutorService = Tracers.wrap(Executors.newFixedThreadPool(1));
ExecutorService producerExecutorService = Executors.newFixedThreadPool(numProducers);
ExecutorService consumerExecutorService = Executors.newFixedThreadPool(1);

try (CloseableTracer submit = CloseableTracer.startSpan("submit")) {
IntStream.range(0, numElem).forEach(i -> {

Tracer.clearCurrentTrace(); // just pretending all these tasks are on a fresh request

DetachedSpan span = DetachedSpan.start("callback-pending" + i + " (cross thread span)");
producerExecutorService.submit(() -> {
try (CloseableTracer closeableTracer = CloseableTracer.startSpan("submit-work" + i)) {
work.add("work" + i);
submitLatch.countDown();
}
work.add(new QueuedWork() {
@Override
public String name() {
return "work" + i;
}

@Override
public DetachedSpan span() {
return span;
}
});
submitLatch.countDown();
});
});
assertThat(submitLatch.await(10, TimeUnit.SECONDS)).isTrue();

consumerExecutorService.submit(() -> {
for (int i = 0; i < numElem; i++) {
String poll = work.take();
sleep(10, "processing" + poll);
QueuedWork queuedWork = work.take();
try (CloseableSpan span = queuedWork.span().completeAndStartChild("consume" + queuedWork.name())) {
Thread.sleep(10);
}
consumeLatch.countDown();
}
return null;
Expand All@@ -129,35 +157,30 @@ void multi_producer_single_consumer() throws InterruptedException {
}

@Test
@TestTracing(snapshot = true, layout = LayoutStrategy.SPLIT_BY_TRACE)
@TestTracing(snapshot = true, layout = LayoutStrategy.CHRONOLOGICAL)
void backoffs_on_a_scheduled_executor() throws InterruptedException {
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
CountDownLatch latch = new CountDownLatch(1);

try (CloseableTracer t = CloseableTracer.startSpan("some-request")) {
executor.execute(() -> {
// first attempt at a network call
sleep(100, "first attempt");

executor.schedule(() -> {
// attempt number 2
sleep(100, "second attempt");

executor.schedule(() -> {
// attempt number 3
sleep(100, "final attempt");
DetachedSpan overall = DetachedSpan.start("overall request");
executor.execute(() -> {

latch.countDown();
}, 100, TimeUnit.MILLISECONDS);
try (CloseableTracer t = CloseableTracer.startSpan("first network call (pretending this fails)")) {
sleep(100);
}

sleep(200, "second tidying");
}, 100, TimeUnit.MILLISECONDS);
DetachedSpan backoff = overall.childDetachedSpan("backoff");
executor.schedule(() -> {
try (CloseableSpan attempt2 = backoff.completeAndStartChild("secondAttempt")) {
sleep(100);
overall.complete();
latch.countDown();

sleep(200, "first tidying");
});
}
}, 20, TimeUnit.MILLISECONDS);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();

MoreExecutors.shutdownAndAwaitTermination(executor, 1, TimeUnit.SECONDS);
}
Expand All@@ -167,54 +190,76 @@ void backoffs_on_a_scheduled_executor() throws InterruptedException {
@SuppressWarnings("CheckReturnValue")
void transformed_future() throws InterruptedException {
SettableFuture<Object> future = SettableFuture.create();
ScheduledExecutorService executor = Tracers.wrap(Executors.newScheduledThreadPool(2));
ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
CountDownLatch latch = new CountDownLatch(1);

DetachedSpan foo = DetachedSpan.start("foo");
FluentFuture.from(future)
.transform(result -> {
sleep(100, "first");
return result;
try (CloseableSpan t = foo.childSpan("first transform")) {
sleep(1000);
return result;
}
}, executor)
.transform(result -> {
sleep(100, "second");
latch.countDown();
return result;
try (CloseableSpan t = foo.childSpan("second transform")) {
sleep(1000);
latch.countDown();
return result;
}
}, executor)
.addCallback(new FutureCallback<Object>() {
@Override
public void onSuccess(@Nullable Object result) {
foo.complete();
}

@Override
public void onFailure(Throwable throwable) {
foo.complete();
}
}, executor);

executor.submit(() -> {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
future.set(null);
}
future.set(null);
});

assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
}

private static void sleep(int millis, String operation) {
try (CloseableTracer t = CloseableTracer.startSpan(operation)) {
private static void sleep(int millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
throw new RuntimeException("dont care", e);
Thread.currentThread().interrupt();
throw new RuntimeException(e);
}
}

private static void sleep(int millis) {
sleep(millis, "sleep " + millis);
private static void sleepSpan(int millis) {
try (CloseableTracer t = CloseableTracer.startSpan("sleep " + millis)) {
sleep(millis);
}
}

@SuppressWarnings("NestedTryDepth")
private static void emit_nested_spans() {
try (CloseableTracer root = CloseableTracer.startSpan("root")) {
try (CloseableTracer root = CloseableTracer.startSpan("emit_nested_spans")) {
try (CloseableTracer first = CloseableTracer.startSpan("first")) {
sleep(100);
sleepSpan(100);
try (CloseableTracer nested = CloseableTracer.startSpan("nested")) {
sleep(90);
sleepSpan(90);
}
sleep(10);
sleepSpan(10);
}
try (CloseableTracer second = CloseableTracer.startSpan("second")) {
sleep(100);
sleepSpan(100);
}
}
}

interface QueuedWork {
String name();
DetachedSpan span();
}
}
Loading