diff --git a/eng/code-quality-reports/src/main/resources/checkstyle/checkstyle-suppressions.xml b/eng/code-quality-reports/src/main/resources/checkstyle/checkstyle-suppressions.xml
index 725adce69d21..f57dcb91ab03 100755
--- a/eng/code-quality-reports/src/main/resources/checkstyle/checkstyle-suppressions.xml
+++ b/eng/code-quality-reports/src/main/resources/checkstyle/checkstyle-suppressions.xml
@@ -98,6 +98,7 @@
+
diff --git a/sdk/communication/azure-communication-callautomation/src/main/java/com/azure/communication/callautomation/models/CallConnectionProperties.java b/sdk/communication/azure-communication-callautomation/src/main/java/com/azure/communication/callautomation/models/CallConnectionProperties.java
index 618a42e18b68..79179db87b6b 100644
--- a/sdk/communication/azure-communication-callautomation/src/main/java/com/azure/communication/callautomation/models/CallConnectionProperties.java
+++ b/sdk/communication/azure-communication-callautomation/src/main/java/com/azure/communication/callautomation/models/CallConnectionProperties.java
@@ -139,8 +139,7 @@ public String getCallConnectionId() {
*
* @return the mediaSubscriptionId value.
*/
- public String getMediaSubscriptionId() {
+ public String getMediaSubscriptionId() {
return mediaSubscriptionId;
}
-
}
diff --git a/sdk/core/azure-core-tracing-opentelemetry/src/main/java/com/azure/core/tracing/opentelemetry/OpenTelemetryTracer.java b/sdk/core/azure-core-tracing-opentelemetry/src/main/java/com/azure/core/tracing/opentelemetry/OpenTelemetryTracer.java
index 4cf115d3b117..0587e5e32008 100644
--- a/sdk/core/azure-core-tracing-opentelemetry/src/main/java/com/azure/core/tracing/opentelemetry/OpenTelemetryTracer.java
+++ b/sdk/core/azure-core-tracing-opentelemetry/src/main/java/com/azure/core/tracing/opentelemetry/OpenTelemetryTracer.java
@@ -362,7 +362,7 @@ private SpanBuilder createSpanBuilder(String spanName,
SpanBuilder spanBuilder = tracer.spanBuilder(spanNameKey)
.setSpanKind(spanKind);
- io.opentelemetry.context.Context parentContext = getTraceContextOrDefault(context, io.opentelemetry.context.Context.current());
+ io.opentelemetry.context.Context parentContext = getTraceContextOrDefault(context, io.opentelemetry.context.Context.current());
// if remote parent is provided, it has higher priority
if (remoteParentContext != null) {
spanBuilder.setParent(parentContext.with(Span.wrap(remoteParentContext)));
diff --git a/sdk/servicebus/azure-messaging-servicebus/pom.xml b/sdk/servicebus/azure-messaging-servicebus/pom.xml
index 8e8af51e1933..6ae86d91f3e1 100644
--- a/sdk/servicebus/azure-messaging-servicebus/pom.xml
+++ b/sdk/servicebus/azure-messaging-servicebus/pom.xml
@@ -118,6 +118,27 @@
4.5.1
test
+
+
+ com.azure
+ azure-core-tracing-opentelemetry
+ 1.0.0-beta.28
+ test
+
+
+
+ io.opentelemetry
+ opentelemetry-api
+ 1.14.0
+ test
+
+
+
+ io.opentelemetry
+ opentelemetry-sdk
+ 1.14.0
+ test
+
diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/FluxTrace.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/FluxTrace.java
new file mode 100644
index 000000000000..0f7a49f0ca52
--- /dev/null
+++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/FluxTrace.java
@@ -0,0 +1,91 @@
+// Copyright (c) Microsoft Corporation. All rights reserved.
+// Licensed under the MIT License.
+
+package com.azure.messaging.servicebus;
+
+import com.azure.core.util.Context;
+import com.azure.messaging.servicebus.implementation.ServiceBusReceiverTracer;
+import org.reactivestreams.Subscription;
+import reactor.core.CoreSubscriber;
+import reactor.core.publisher.BaseSubscriber;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.FluxOperator;
+
+import java.util.Objects;
+
+/**
+ * Flux operator that traces receive and process calls
+ */
+final class FluxTrace extends FluxOperator {
+ static final String PROCESS_ERROR_KEY = "process-error";
+ private final ServiceBusReceiverTracer tracer;
+
+ FluxTrace(Flux extends ServiceBusMessageContext> upstream, ServiceBusReceiverTracer tracer) {
+ super(upstream);
+ this.tracer = tracer;
+ }
+
+ @Override
+ public void subscribe(CoreSubscriber super ServiceBusMessageContext> coreSubscriber) {
+ Objects.requireNonNull(coreSubscriber, "'coreSubscriber' cannot be null.");
+
+ source.subscribe(new TracingSubscriber(coreSubscriber, tracer));
+ }
+
+ private static class TracingSubscriber extends BaseSubscriber {
+
+ private final CoreSubscriber super ServiceBusMessageContext> downstream;
+ private final ServiceBusReceiverTracer tracer;
+ TracingSubscriber(CoreSubscriber super ServiceBusMessageContext> downstream, ServiceBusReceiverTracer tracer) {
+ this.downstream = downstream;
+ this.tracer = tracer;
+ }
+
+ @Override
+ public reactor.util.context.Context currentContext() {
+ return downstream.currentContext();
+ }
+
+ @Override
+ protected void hookOnSubscribe(Subscription subscription) {
+ downstream.onSubscribe(this);
+ }
+
+ @Override
+ protected void hookOnNext(ServiceBusMessageContext message) {
+ if (tracer == null || tracer.isSync()) {
+ downstream.onNext(message);
+ return;
+ }
+
+ Throwable exception = null;
+ Context span = tracer.startProcessSpan("ServiceBus.process", message.getMessage(), Context.NONE);
+ AutoCloseable scope = tracer.makeSpanCurrent(span);
+
+ try {
+ downstream.onNext(message);
+ } catch (Throwable t) {
+ exception = t;
+ } finally {
+ Context context = message.getMessage().getContext();
+ if (context != null) {
+ Object processorException = context.getData(PROCESS_ERROR_KEY).orElse(null);
+ if (processorException instanceof Throwable) {
+ exception = (Throwable) processorException;
+ }
+ }
+ tracer.endSpan(exception, span, scope);
+ }
+ }
+
+ @Override
+ protected void hookOnError(Throwable throwable) {
+ downstream.onError(throwable);
+ }
+
+ @Override
+ protected void hookOnComplete() {
+ downstream.onComplete();
+ }
+ }
+}
diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java
index a8449e77cd00..d705c78be934 100644
--- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java
+++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusClientBuilder.java
@@ -17,7 +17,6 @@
import com.azure.core.amqp.implementation.ReactorProvider;
import com.azure.core.amqp.implementation.StringUtil;
import com.azure.core.amqp.implementation.TokenManagerProvider;
-import com.azure.core.amqp.implementation.TracerProvider;
import com.azure.core.amqp.models.CbsAuthorizationType;
import com.azure.core.annotation.ServiceClientBuilder;
import com.azure.core.annotation.ServiceClientProtocol;
@@ -34,14 +33,15 @@
import com.azure.core.util.Configuration;
import com.azure.core.util.CoreUtils;
import com.azure.core.util.logging.ClientLogger;
-import com.azure.core.util.tracing.Tracer;
import com.azure.messaging.servicebus.implementation.MessageUtils;
import com.azure.messaging.servicebus.implementation.MessagingEntityType;
import com.azure.messaging.servicebus.implementation.ServiceBusAmqpConnection;
import com.azure.messaging.servicebus.implementation.ServiceBusConnectionProcessor;
import com.azure.messaging.servicebus.implementation.ServiceBusConstants;
+import com.azure.messaging.servicebus.implementation.ServiceBusReceiverTracer;
import com.azure.messaging.servicebus.implementation.ServiceBusReactorAmqpConnection;
import com.azure.messaging.servicebus.implementation.ServiceBusSharedKeyCredential;
+import com.azure.messaging.servicebus.implementation.ServiceBusTracer;
import com.azure.messaging.servicebus.implementation.models.ServiceBusProcessorClientOptions;
import com.azure.messaging.servicebus.models.ServiceBusReceiveMode;
import com.azure.messaging.servicebus.models.SubQueue;
@@ -59,7 +59,6 @@
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
-import java.util.ServiceLoader;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
@@ -212,7 +211,6 @@ public final class ServiceBusClientBuilder implements
private final Object connectionLock = new Object();
private final MessageSerializer messageSerializer = new ServiceBusMessageSerializer();
- private final TracerProvider tracerProvider = new TracerProvider(ServiceLoader.load(Tracer.class));
private ClientOptions clientOptions;
private Configuration configuration;
@@ -963,8 +961,9 @@ public ServiceBusSenderAsyncClient buildAsyncClient() {
clientIdentifier = UUID.randomUUID().toString();
}
+ final ServiceBusSenderTracer tracer = new ServiceBusSenderTracer(ServiceBusTracer.getDefaultTracer(), connectionProcessor.getFullyQualifiedNamespace(), entityName);
return new ServiceBusSenderAsyncClient(entityName, entityType, connectionProcessor, retryOptions,
- tracerProvider, messageSerializer, ServiceBusClientBuilder.this::onClientClose, null, clientIdentifier);
+ tracer, messageSerializer, ServiceBusClientBuilder.this::onClientClose, null, clientIdentifier);
}
/**
@@ -1047,8 +1046,7 @@ public final class ServiceBusSessionProcessorClientBuilder {
private ServiceBusSessionProcessorClientBuilder() {
sessionReceiverClientBuilder = new ServiceBusSessionReceiverClientBuilder();
processorClientOptions = new ServiceBusProcessorClientOptions()
- .setMaxConcurrentCalls(1)
- .setTracerProvider(tracerProvider);
+ .setMaxConcurrentCalls(1);
sessionReceiverClientBuilder.maxConcurrentSessions(1);
}
@@ -1444,11 +1442,11 @@ ServiceBusReceiverAsyncClient buildAsyncClientForProcessor() {
}
final ServiceBusSessionManager sessionManager = new ServiceBusSessionManager(entityPath, entityType,
- connectionProcessor, tracerProvider, messageSerializer, receiverOptions, clientIdentifier);
-
+ connectionProcessor, messageSerializer, receiverOptions, clientIdentifier);
+ final ServiceBusReceiverTracer tracer = new ServiceBusReceiverTracer(ServiceBusTracer.getDefaultTracer(), connectionProcessor.getFullyQualifiedNamespace(), entityPath, false);
return new ServiceBusReceiverAsyncClient(connectionProcessor.getFullyQualifiedNamespace(), entityPath,
entityType, receiverOptions, connectionProcessor, ServiceBusConstants.OPERATION_TIMEOUT,
- tracerProvider, messageSerializer, ServiceBusClientBuilder.this::onClientClose, sessionManager);
+ tracer, messageSerializer, ServiceBusClientBuilder.this::onClientClose, sessionManager);
}
/**
@@ -1466,7 +1464,7 @@ ServiceBusReceiverAsyncClient buildAsyncClientForProcessor() {
* queueName()} or {@link #topicName(String) topicName()}, respectively.
*/
public ServiceBusSessionReceiverAsyncClient buildAsyncClient() {
- return buildAsyncClient(true);
+ return buildAsyncClient(true, false);
}
/**
@@ -1484,12 +1482,12 @@ public ServiceBusSessionReceiverAsyncClient buildAsyncClient() {
*/
public ServiceBusSessionReceiverClient buildClient() {
final boolean isPrefetchDisabled = prefetchCount == 0;
- return new ServiceBusSessionReceiverClient(buildAsyncClient(false),
+ return new ServiceBusSessionReceiverClient(buildAsyncClient(false, true),
isPrefetchDisabled,
MessageUtils.getTotalTimeout(retryOptions));
}
- private ServiceBusSessionReceiverAsyncClient buildAsyncClient(boolean isAutoCompleteAllowed) {
+ private ServiceBusSessionReceiverAsyncClient buildAsyncClient(boolean isAutoCompleteAllowed, boolean syncConsumer) {
final MessagingEntityType entityType = validateEntityPaths(connectionStringEntityName, topicName,
queueName);
final String entityPath = getEntityPath(entityType, queueName, topicName, subscriptionName,
@@ -1520,8 +1518,10 @@ private ServiceBusSessionReceiverAsyncClient buildAsyncClient(boolean isAutoComp
clientIdentifier = UUID.randomUUID().toString();
}
+ final ServiceBusReceiverTracer tracer = new ServiceBusReceiverTracer(ServiceBusTracer.getDefaultTracer(),
+ connectionProcessor.getFullyQualifiedNamespace(), entityPath, syncConsumer);
return new ServiceBusSessionReceiverAsyncClient(connectionProcessor.getFullyQualifiedNamespace(),
- entityPath, entityType, receiverOptions, connectionProcessor, tracerProvider, messageSerializer,
+ entityPath, entityType, receiverOptions, connectionProcessor, tracer, messageSerializer,
ServiceBusClientBuilder.this::onClientClose, clientIdentifier);
}
}
@@ -1582,8 +1582,7 @@ public final class ServiceBusProcessorClientBuilder {
private ServiceBusProcessorClientBuilder() {
serviceBusReceiverClientBuilder = new ServiceBusReceiverClientBuilder();
processorClientOptions = new ServiceBusProcessorClientOptions()
- .setMaxConcurrentCalls(1)
- .setTracerProvider(tracerProvider);
+ .setMaxConcurrentCalls(1);
}
/**
@@ -1908,7 +1907,7 @@ public ServiceBusReceiverClientBuilder topicName(String topicName) {
* queueName()} or {@link #topicName(String) topicName()}, respectively.
*/
public ServiceBusReceiverAsyncClient buildAsyncClient() {
- return buildAsyncClient(true);
+ return buildAsyncClient(true, false);
}
/**
@@ -1926,12 +1925,12 @@ public ServiceBusReceiverAsyncClient buildAsyncClient() {
*/
public ServiceBusReceiverClient buildClient() {
final boolean isPrefetchDisabled = prefetchCount == 0;
- return new ServiceBusReceiverClient(buildAsyncClient(false),
+ return new ServiceBusReceiverClient(buildAsyncClient(false, true),
isPrefetchDisabled,
MessageUtils.getTotalTimeout(retryOptions));
}
- ServiceBusReceiverAsyncClient buildAsyncClient(boolean isAutoCompleteAllowed) {
+ ServiceBusReceiverAsyncClient buildAsyncClient(boolean isAutoCompleteAllowed, boolean syncConsumer) {
final MessagingEntityType entityType = validateEntityPaths(connectionStringEntityName, topicName,
queueName);
final String entityPath = getEntityPath(entityType, queueName, topicName, subscriptionName,
@@ -1962,9 +1961,11 @@ ServiceBusReceiverAsyncClient buildAsyncClient(boolean isAutoCompleteAllowed) {
clientIdentifier = UUID.randomUUID().toString();
}
+ final ServiceBusReceiverTracer tracer = new ServiceBusReceiverTracer(ServiceBusTracer.getDefaultTracer(),
+ connectionProcessor.getFullyQualifiedNamespace(), entityPath, syncConsumer);
return new ServiceBusReceiverAsyncClient(connectionProcessor.getFullyQualifiedNamespace(), entityPath,
entityType, receiverOptions, connectionProcessor, ServiceBusConstants.OPERATION_TIMEOUT,
- tracerProvider, messageSerializer, ServiceBusClientBuilder.this::onClientClose, clientIdentifier);
+ tracer, messageSerializer, ServiceBusClientBuilder.this::onClientClose, clientIdentifier);
}
}
diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusMessageBatch.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusMessageBatch.java
index 3308d1aac920..6a8ea03376fe 100644
--- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusMessageBatch.java
+++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusMessageBatch.java
@@ -7,7 +7,6 @@
import com.azure.core.amqp.exception.AmqpException;
import com.azure.core.amqp.implementation.ErrorContextProvider;
import com.azure.core.amqp.implementation.MessageSerializer;
-import com.azure.core.amqp.implementation.TracerProvider;
import com.azure.core.util.logging.ClientLogger;
import org.apache.qpid.proton.message.Message;
@@ -17,7 +16,6 @@
import java.util.Locale;
import java.util.Objects;
-import static com.azure.messaging.servicebus.implementation.MessageUtils.traceMessageSpan;
/**
* A class for aggregating {@link ServiceBusMessage messages} into a single, size-limited, batch. It is treated as a
@@ -31,11 +29,11 @@ public final class ServiceBusMessageBatch {
private final List serviceBusMessageList;
private final byte[] eventBytes;
private int sizeInBytes;
- private final TracerProvider tracerProvider;
+ private final ServiceBusSenderTracer tracer;
private final String entityPath;
private final String hostname;
- ServiceBusMessageBatch(int maxMessageSize, ErrorContextProvider contextProvider, TracerProvider tracerProvider,
+ ServiceBusMessageBatch(int maxMessageSize, ErrorContextProvider contextProvider, ServiceBusSenderTracer tracer,
MessageSerializer serializer, String entityPath, String hostname) {
this.maxMessageSize = maxMessageSize;
this.contextProvider = contextProvider;
@@ -43,7 +41,7 @@ public final class ServiceBusMessageBatch {
this.serviceBusMessageList = new ArrayList<>();
this.sizeInBytes = (maxMessageSize / 65536) * 1024; // reserve 1KB for every 64KB
this.eventBytes = new byte[maxMessageSize];
- this.tracerProvider = tracerProvider;
+ this.tracer = tracer;
this.entityPath = entityPath;
this.hostname = hostname;
}
@@ -94,15 +92,11 @@ public boolean tryAddMessage(final ServiceBusMessage serviceBusMessage) {
if (serviceBusMessage == null) {
throw LOGGER.logExceptionAsWarning(new NullPointerException("'serviceBusMessage' cannot be null"));
}
- ServiceBusMessage serviceBusMessageUpdated =
- tracerProvider.isEnabled()
- ? traceMessageSpan(serviceBusMessage, serviceBusMessage.getContext(), hostname, entityPath,
- tracerProvider)
- : serviceBusMessage;
+ tracer.createMessageSpan(serviceBusMessage);
final int size;
try {
- size = getSize(serviceBusMessageUpdated, serviceBusMessageList.isEmpty());
+ size = getSize(serviceBusMessage, serviceBusMessageList.isEmpty());
} catch (BufferOverflowException exception) {
final RuntimeException ex = new ServiceBusException(
new AmqpException(false, AmqpErrorCondition.LINK_PAYLOAD_SIZE_EXCEEDED,
@@ -117,7 +111,7 @@ public boolean tryAddMessage(final ServiceBusMessage serviceBusMessage) {
}
this.sizeInBytes += size;
- this.serviceBusMessageList.add(serviceBusMessageUpdated);
+ this.serviceBusMessageList.add(serviceBusMessage);
return true;
}
diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusProcessorClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusProcessorClient.java
index 55d2d83e7f55..99212d649516 100644
--- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusProcessorClient.java
+++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusProcessorClient.java
@@ -3,39 +3,22 @@
package com.azure.messaging.servicebus;
-import com.azure.core.amqp.implementation.TracerProvider;
-import com.azure.core.util.Context;
import com.azure.core.util.logging.ClientLogger;
-import com.azure.core.util.tracing.ProcessKind;
import com.azure.messaging.servicebus.ServiceBusClientBuilder.ServiceBusProcessorClientBuilder;
import com.azure.messaging.servicebus.ServiceBusClientBuilder.ServiceBusSessionProcessorClientBuilder;
import com.azure.messaging.servicebus.implementation.models.ServiceBusProcessorClientOptions;
import org.reactivestreams.Subscription;
import reactor.core.CoreSubscriber;
import reactor.core.Disposable;
-import reactor.core.publisher.Signal;
import reactor.core.scheduler.Schedulers;
-
-import java.util.Locale;
import java.util.Map;
import java.util.Objects;
-import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
-import static com.azure.core.util.tracing.Tracer.AZ_TRACING_NAMESPACE_KEY;
-import static com.azure.core.util.tracing.Tracer.DIAGNOSTIC_ID_KEY;
-import static com.azure.core.util.tracing.Tracer.ENTITY_PATH_KEY;
-import static com.azure.core.util.tracing.Tracer.HOST_NAME_KEY;
-import static com.azure.core.util.tracing.Tracer.MESSAGE_ENQUEUED_TIME;
-import static com.azure.core.util.tracing.Tracer.SCOPE_KEY;
-import static com.azure.core.util.tracing.Tracer.SPAN_CONTEXT_KEY;
-import static com.azure.messaging.servicebus.implementation.ServiceBusConstants.AZ_TRACING_NAMESPACE_VALUE;
-import static com.azure.messaging.servicebus.implementation.ServiceBusConstants.AZ_TRACING_SERVICE_NAME;
-
/**
* The processor client for processing Service Bus messages. {@link ServiceBusProcessorClient} provides a push-based
* mechanism that invokes the message processing callback when a message is received or the error handler when an error
@@ -133,7 +116,6 @@ public final class ServiceBusProcessorClient implements AutoCloseable {
private final Map receiverSubscriptions = new ConcurrentHashMap<>();
private final AtomicReference asyncClient = new AtomicReference<>();
private final AtomicBoolean isRunning = new AtomicBoolean();
- private final TracerProvider tracerProvider;
private final String queueName;
private final String topicName;
private final String subscriptionName;
@@ -160,9 +142,9 @@ public final class ServiceBusProcessorClient implements AutoCloseable {
this.processMessage = Objects.requireNonNull(processMessage, "'processMessage' cannot be null");
this.processError = Objects.requireNonNull(processError, "'processError' cannot be null");
this.processorOptions = Objects.requireNonNull(processorOptions, "'processorOptions' cannot be null");
+
this.asyncClient.set(sessionReceiverBuilder.buildAsyncClientForProcessor());
this.receiverBuilder = null;
- this.tracerProvider = processorOptions.getTracerProvider();
this.queueName = queueName;
this.topicName = topicName;
this.subscriptionName = subscriptionName;
@@ -189,7 +171,6 @@ public final class ServiceBusProcessorClient implements AutoCloseable {
this.processorOptions = Objects.requireNonNull(processorOptions, "'processorOptions' cannot be null");
this.asyncClient.set(receiverBuilder.buildAsyncClient());
this.sessionReceiverBuilder = null;
- this.tracerProvider = processorOptions.getTracerProvider();
this.queueName = queueName;
this.topicName = topicName;
this.subscriptionName = subscriptionName;
@@ -345,22 +326,15 @@ public void onNext(ServiceBusMessageContext serviceBusMessageContext) {
if (serviceBusMessageContext.hasError()) {
handleError(serviceBusMessageContext.getThrowable());
} else {
- Context processSpanContext = null;
try {
ServiceBusReceivedMessageContext serviceBusReceivedMessageContext =
new ServiceBusReceivedMessageContext(receiverClient, serviceBusMessageContext);
- processSpanContext =
- startProcessTracingSpan(serviceBusMessageContext.getMessage(),
- receiverClient.getEntityPath(), receiverClient.getFullyQualifiedNamespace());
- if (processSpanContext.getData(SPAN_CONTEXT_KEY).isPresent()) {
- serviceBusMessageContext.getMessage().addContext(SPAN_CONTEXT_KEY, processSpanContext);
- }
processMessage.accept(serviceBusReceivedMessageContext);
- endProcessTracingSpan(processSpanContext, Signal.complete());
} catch (Exception ex) {
+ serviceBusMessageContext.getMessage().addContext(FluxTrace.PROCESS_ERROR_KEY, ex);
handleError(new ServiceBusException(ex, ServiceBusErrorSource.USER_CALLBACK));
- endProcessTracingSpan(processSpanContext, Signal.error(ex));
+
if (!processorOptions.isDisableAutoComplete()) {
LOGGER.warning("Error when processing message. Abandoning message.", ex);
abandonMessage(serviceBusMessageContext, receiverClient);
@@ -407,54 +381,6 @@ public void onComplete() {
}
}
- private void endProcessTracingSpan(Context processSpanContext, Signal signal) {
- if (processSpanContext == null) {
- return;
- }
-
- Optional