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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
}
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,31 +97,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
container1.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records1.add(record)
}
})

container2.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records2.add(record)
}
})

container3.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records3.add(record)
}
})

container4.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records4.add(record)
}
})
Expand Down Expand Up @@ -202,31 +198,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
container1.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records1.add(activeSpan())
}
})

container2.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records2.add(activeSpan())
}
})

container3.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records3.add(activeSpan())
}
})

container4.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records4.add(activeSpan())
}
})
Expand Down
Comment thread
amarziali marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
}
})
Expand Down Expand Up @@ -420,7 +419,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
Comment thread
amarziali marked this conversation as resolved.
}
})
Expand Down Expand Up @@ -471,10 +469,13 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
}
}

// sort a snapshot so the producer trace is deterministically first, regardless of write order
def sortedTraces = new ArrayList<>(TEST_WRITER)
sortedTraces.sort(SORT_TRACES_BY_ID)
def headers = received.headers()
headers.iterator().hasNext()
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${TEST_WRITER[0][2].traceId}"
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${TEST_WRITER[0][2].spanId}"
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${sortedTraces[0][2].traceId}"
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${sortedTraces[0][2].spanId}"

if (isDataStreamsEnabled()) {
StatsGroup first = TEST_DATA_STREAMS_WRITER.groups.find { it.parentHash == 0 }
Expand Down Expand Up @@ -553,7 +554,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
}
})
Expand Down Expand Up @@ -889,7 +889,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new BatchMessageListener<String, String>() {
@Override
void onMessage(List<ConsumerRecord<String, String>> consumerRecords) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
consumerRecords.each {
records.add(it)
}
Expand Down Expand Up @@ -1026,7 +1025,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
if (isDataStreamsEnabled()) {
// even if header propagation is disabled, we want data streams to work.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,31 +101,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
container1.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records1.add(record)
}
})

container2.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records2.add(record)
}
})

container3.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records3.add(record)
}
})

container4.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records4.add(record)
}
})
Expand Down Expand Up @@ -205,31 +201,27 @@ class KafkaClientCustomPropagationConfigTest extends InstrumentationSpecificatio
container1.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records1.add(activeSpan())
}
})

container2.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records2.add(activeSpan())
}
})

container3.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records3.add(activeSpan())
}
})

container4.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records4.add(activeSpan())
}
})
Expand Down
Comment thread
amarziali marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
}
})
Expand Down Expand Up @@ -349,7 +348,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
Comment thread
amarziali marked this conversation as resolved.
}
})
Expand Down Expand Up @@ -404,10 +402,13 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
}
}

// sort a snapshot so the producer trace is deterministically first, regardless of write order
def sortedTraces = new ArrayList<>(TEST_WRITER)
sortedTraces.sort(SORT_TRACES_BY_ID)
def headers = received.headers()
headers.iterator().hasNext()
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${TEST_WRITER[0][2].traceId}"
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${TEST_WRITER[0][2].spanId}"
new String(headers.headers("x-datadog-trace-id").iterator().next().value()) == "${sortedTraces[0][2].traceId}"
new String(headers.headers("x-datadog-parent-id").iterator().next().value()) == "${sortedTraces[0][2].spanId}"

if (isDataStreamsEnabled()) {
StatsGroup first = TEST_DATA_STREAMS_WRITER.groups.find { it.parentHash == 0 }
Expand Down Expand Up @@ -480,7 +481,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
}
})
Expand Down Expand Up @@ -824,7 +824,6 @@ abstract class KafkaClientTestBase extends VersionedNamingTestBase {
container.setupMessageListener(new MessageListener<String, String>() {
@Override
void onMessage(ConsumerRecord<String, String> record) {
TEST_WRITER.waitForTraces(1) // ensure consistent ordering of traces
records.add(record)
if (isDataStreamsEnabled()) {
// even if header propagation is disabled, we want data streams to work.
Expand Down