Skip to content
Open
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 @@ -28,6 +28,7 @@
import org.apache.commons.logging.LogFactory;
import org.jspecify.annotations.Nullable;

import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.batch.BatchingStrategy;
import org.springframework.amqp.rabbit.batch.SimpleBatchingStrategy;
Expand Down Expand Up @@ -195,7 +196,7 @@ public String getComponentType() {
resp.getEnvelope(), StandardCharsets.UTF_8.name());
messageProperties.setConsumerQueue(this.queue);
Map<String, @Nullable Object> headers = this.headerMapper.toHeadersFromRequest(messageProperties);
var amqpMessage = new org.springframework.amqp.core.Message(resp.getBody(), messageProperties);
var amqpMessage = new Message(resp.getBody(), messageProperties);
Object payload;
if (this.batchingStrategy.canDebatch(messageProperties)) {
List<Object> payloads = new ArrayList<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
import org.springframework.context.ApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.core.log.LogAccessor;
import org.springframework.expression.Expression;
import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint;
import org.springframework.integration.amqp.support.AmqpHeaderMapper;
import org.springframework.integration.channel.DirectChannel;
Expand Down Expand Up @@ -118,7 +119,7 @@ public void withHeaderMapperCustomHeaders() {
AmqpOutboundEndpoint endpoint = TestUtils.getPropertyValue(eventDrivenConsumer, "handler");
assertThat(TestUtils.<Object>getPropertyValue(endpoint, "defaultDeliveryMode")).isNotNull();
assertThat(TestUtils.<Boolean>getPropertyValue(endpoint, "lazyConnect")).isFalse();
assertThat(TestUtils.<org.springframework.expression.Expression>getPropertyValue(endpoint, "delayExpression")
assertThat(TestUtils.<Expression>getPropertyValue(endpoint, "delayExpression")
.getExpressionString()).isEqualTo("42");
assertThat(TestUtils.<Boolean>getPropertyValue(endpoint, "headersMappedLast")).isFalse();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.springframework.beans.factory.parsing.BeanDefinitionParsingException;
import org.springframework.context.ApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.expression.Expression;
import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint;
import org.springframework.integration.amqp.outbound.AsyncAmqpOutboundGateway;
import org.springframework.integration.channel.QueueChannel;
Expand Down Expand Up @@ -97,7 +98,7 @@ protected void checkGWProps(ApplicationContext context, Orderable gateway) {

assertThat(sendTimeout).isEqualTo(Long.valueOf(777));
assertThat(TestUtils.<Boolean>getPropertyValue(gateway, "lazyConnect")).isTrue();
assertThat(TestUtils.<org.springframework.expression.Expression>getPropertyValue(gateway, "delayExpression")
assertThat(TestUtils.<Expression>getPropertyValue(gateway, "delayExpression")
.getExpressionString()).isEqualTo("42");
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

package org.springframework.integration.amqp.dsl;

import com.rabbitmq.client.amqp.Environment;
import com.rabbitmq.client.amqp.impl.AmqpEnvironmentBuilder;
import org.junit.jupiter.api.Test;

Expand Down Expand Up @@ -94,7 +95,7 @@ void requestReplyOverAmqp() {
static class ContextConfiguration {

@Bean
com.rabbitmq.client.amqp.Environment environment() {
Environment environment() {
return new AmqpEnvironmentBuilder()
.connectionSettings()
.port(RabbitTestContainer.amqpPort())
Expand All @@ -103,7 +104,7 @@ com.rabbitmq.client.amqp.Environment environment() {
}

@Bean
AmqpConnectionFactory connectionFactory(com.rabbitmq.client.amqp.Environment environment) {
AmqpConnectionFactory connectionFactory(Environment environment) {
return new SingleAmqpConnectionFactory(environment);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.junit.jupiter.api.Test;

import org.springframework.amqp.core.Declarables;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageBuilder;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.rabbitmq.client.AmqpConnectionFactory;
Expand Down Expand Up @@ -72,7 +73,7 @@ void inboundGatewayExchange() {

@Test
void inboundGatewayExchangeWithAck() throws InterruptedException {
org.springframework.amqp.core.Message requestMessage =
Message requestMessage =
MessageBuilder.withBody("test data #2".getBytes())
.setMessageId("someMessageId")
.setContentType(MimeTypeUtils.TEXT_PLAIN_VALUE)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@
*
* @since 6.0
*/
public class CamelHeaderMapper implements HeaderMapper<org.apache.camel.Message> {
public class CamelHeaderMapper implements HeaderMapper<Message> {

private static final LogAccessor LOGGER = new LogAccessor(CamelHeaderMapper.class);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@

package org.springframework.integration.util;

import kotlin.coroutines.Continuation;
import kotlinx.coroutines.reactor.MonoKt;
import org.jspecify.annotations.Nullable;
import reactor.core.publisher.Mono;

Expand All @@ -38,7 +40,7 @@ public static boolean isContinuation(Object candidate) {
}

public static boolean isContinuationType(Class<?> candidate) {
return KotlinDetector.isKotlinPresent() && kotlin.coroutines.Continuation.class.isAssignableFrom(candidate);
return KotlinDetector.isKotlinPresent() && Continuation.class.isAssignableFrom(candidate);
}

@Nullable
Expand All @@ -47,8 +49,8 @@ public static <T> T monoAwaitSingleOrNull(Mono<? extends T> source, Object conti
Assert.state(isContinuation(continuation), () ->
"The 'continuation' must be an instance of 'kotlin.coroutines.Continuation', but it is: "
+ continuation.getClass());
return (T) kotlinx.coroutines.reactor.MonoKt.awaitSingleOrNull(
source, (kotlin.coroutines.Continuation<T>) continuation);
return (T) MonoKt.awaitSingleOrNull(
source, (Continuation<T>) continuation);
}

private CoroutinesUtils() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

package org.springframework.integration.grpc.outbound;

import java.util.ArrayList;
import java.util.List;

import io.grpc.ManagedChannel;
Expand Down Expand Up @@ -347,7 +348,7 @@ public void onCompleted() {
public StreamObserver<HelloRequest> helloToEveryOne(StreamObserver<HelloReply> responseObserver) {
return new StreamObserver<>() {

private final java.util.List<String> names = new java.util.ArrayList<>();
private final List<String> names = new ArrayList<>();

@Override
public void onNext(HelloRequest value) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import com.hazelcast.multimap.MultiMap;
import com.hazelcast.replicatedmap.ReplicatedMap;
import com.hazelcast.topic.ITopic;
import com.hazelcast.topic.Message;
import com.hazelcast.topic.MessageListener;

import org.springframework.integration.hazelcast.HazelcastIntegrationTestUser;
Expand Down Expand Up @@ -167,7 +168,7 @@ public static void testWriteToTopic(MessageChannel channel,
private int index = 1;

@Override
public void onMessage(com.hazelcast.topic.Message message) {
public void onMessage(Message message) {
HazelcastIntegrationTestUser user =
(HazelcastIntegrationTestUser) message.getMessageObject();
verifyHazelcastIntegrationTestUser(user, index);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ public void setOutboundPrefix(String outboundPrefix) {
}

@Override
public void fromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage) {
public void fromHeaders(MessageHeaders headers, Message jmsMessage) {
try {
populateCorrelationIdPropertyFromHeaders(headers, jmsMessage);
populateReplyToPropertyFromHeaders(headers, jmsMessage);
Expand All @@ -154,7 +154,7 @@ public void fromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage)
}
}

private void populateCorrelationIdPropertyFromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage) {
private void populateCorrelationIdPropertyFromHeaders(MessageHeaders headers, Message jmsMessage) {
Object jmsCorrelationId = headers.get(JmsHeaders.CORRELATION_ID);
if (jmsCorrelationId instanceof Number) {
jmsCorrelationId = jmsCorrelationId.toString();
Expand All @@ -169,7 +169,7 @@ private void populateCorrelationIdPropertyFromHeaders(MessageHeaders headers, ja
}
}

private void populateReplyToPropertyFromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage) {
private void populateReplyToPropertyFromHeaders(MessageHeaders headers, Message jmsMessage) {
Object jmsReplyTo = headers.get(JmsHeaders.REPLY_TO);
if (jmsReplyTo instanceof Destination destination) {
try {
Expand All @@ -181,7 +181,7 @@ private void populateReplyToPropertyFromHeaders(MessageHeaders headers, jakarta.
}
}

private void populateTypePropertyFromHeaders(MessageHeaders headers, jakarta.jms.Message jmsMessage) {
private void populateTypePropertyFromHeaders(MessageHeaders headers, Message jmsMessage) {
Object jmsType = headers.get(JmsHeaders.TYPE);
if (jmsType instanceof String jmsTypeStr) {
try {
Expand All @@ -193,7 +193,7 @@ private void populateTypePropertyFromHeaders(MessageHeaders headers, jakarta.jms
}
}

private void populateArbitraryHeaderToProperty(jakarta.jms.Message jmsMessage, String headerName, Object value)
private void populateArbitraryHeaderToProperty(Message jmsMessage, String headerName, Object value)
throws JMSException {

if (SUPPORTED_PROPERTY_TYPES.contains(value.getClass())) {
Expand Down Expand Up @@ -221,7 +221,7 @@ else if (IntegrationMessageHeaderAccessor.CORRELATION_ID.equals(headerName)) {
}

@Override
public Map<String, Object> toHeaders(jakarta.jms.Message jmsMessage) {
public Map<String, Object> toHeaders(Message jmsMessage) {
Map<String, Object> headers = new HashMap<>();
try {
mapMessageIdProperty(jmsMessage, headers);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@
import java.util.concurrent.atomic.AtomicBoolean;

import jakarta.jms.JMSException;
import jakarta.jms.ObjectMessage;
import jakarta.jms.TextMessage;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInfo;

Expand Down Expand Up @@ -240,7 +242,7 @@ private MessageHandler unwrapObjectMessageAndEchoHandler() {
template.setDefaultDestination((MessageChannel) message.getHeaders().getReplyChannel());
Message<?> origMessage = null;
try {
origMessage = (Message<?>) ((jakarta.jms.ObjectMessage) message.getPayload()).getObject();
origMessage = (Message<?>) ((ObjectMessage) message.getPayload()).getObject();
}
catch (JMSException e) {
fail("failed to deserialize message");
Expand All @@ -251,7 +253,7 @@ private MessageHandler unwrapObjectMessageAndEchoHandler() {

private MessageHandler unwrapTextMessageAndEchoHandler() {
return message -> {
assertThat(message.getPayload()).isInstanceOf(jakarta.jms.TextMessage.class);
assertThat(message.getPayload()).isInstanceOf(TextMessage.class);
MessagingTemplate template = new MessagingTemplate();
template.setDefaultDestination((MessageChannel) message.getHeaders().getReplyChannel());
String payload = null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ public Object fromMessage(Message message) throws JMSException, MessageConversio
return "converted-" + original;
}

public jakarta.jms.Message toMessage(Object object, Session session) throws JMSException, MessageConversionException {
public Message toMessage(Object object, Session session) throws JMSException, MessageConversionException {
return null;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import jakarta.jms.ConnectionFactory;
import jakarta.jms.Destination;
import jakarta.jms.Message;
import org.apache.activemq.artemis.core.config.Configuration;
import org.apache.activemq.artemis.core.config.impl.ConfigurationImpl;
import org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ;
Expand Down Expand Up @@ -100,9 +101,9 @@ public void testContainerWithDestBrokenConnection() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(amqConnectionFactory);
template.setReceiveTimeout(5000);
jakarta.jms.Message request = template.receive(requestQueue1);
Message request = template.receive(requestQueue1);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> jmsReply);
assertThat(latch2.await(10, TimeUnit.SECONDS)).isTrue();
assertThat(reply.get()).isNotNull();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import jakarta.jms.Destination;
import jakarta.jms.JMSException;
import jakarta.jms.Message;
import org.apache.activemq.artemis.jms.client.ActiveMQQueue;
import org.junit.jupiter.api.Test;

Expand Down Expand Up @@ -104,9 +105,9 @@ public void testContainerWithDest() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue1);
Message request = template.receive(requestQueue1);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> jmsReply);
assertThat(latch2.await(10, TimeUnit.SECONDS)).isTrue();
assertThat(reply.get()).isNotNull();
Expand Down Expand Up @@ -146,9 +147,9 @@ public void testContainerWithDestNoCorrelation() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue2);
Message request = template.receive(requestQueue2);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> {
jmsReply.setJMSCorrelationID(jmsReply.getJMSMessageID());
return jmsReply;
Expand Down Expand Up @@ -192,9 +193,9 @@ public void testContainerWithDestName() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue3);
Message request = template.receive(requestQueue3);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> jmsReply);
assertThat(latch2.await(10, TimeUnit.SECONDS)).isTrue();
assertThat(reply.get()).isNotNull();
Expand Down Expand Up @@ -234,9 +235,9 @@ public void testContainerWithDestNameNoCorrelation() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue4);
Message request = template.receive(requestQueue4);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> {
jmsReply.setJMSCorrelationID(jmsReply.getJMSMessageID());
return jmsReply;
Expand Down Expand Up @@ -280,9 +281,9 @@ public void testContainerWithTemporary() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue5);
Message request = template.receive(requestQueue5);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> jmsReply);
assertThat(latch2.await(10, TimeUnit.SECONDS)).isTrue();
assertThat(reply.get()).isNotNull();
Expand Down Expand Up @@ -321,9 +322,9 @@ public void testContainerWithTemporaryNoCorrelation() throws Exception {
JmsTemplate template = new JmsTemplate();
template.setConnectionFactory(connectionFactory);
template.setReceiveTimeout(10000);
jakarta.jms.Message request = template.receive(requestQueue6);
Message request = template.receive(requestQueue6);
assertThat(request).isNotNull();
final jakarta.jms.Message jmsReply = request;
final Message jmsReply = request;
template.send(request.getJMSReplyTo(), session -> {
jmsReply.setJMSCorrelationID(jmsReply.getJMSMessageID());
return jmsReply;
Expand Down Expand Up @@ -379,8 +380,8 @@ public void testLazyContainerWithDest() throws Exception {
}

private static void receiveAndSend(JmsTemplate template) {
jakarta.jms.Message request = template.receive(requestQueue7);
final jakarta.jms.Message jmsReply = request;
Message request = template.receive(requestQueue7);
final Message jmsReply = request;
try {
template.send(request.getJMSReplyTo(), session -> jmsReply);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

package org.springframework.integration.mail.inbound;

import jakarta.mail.MessagingException;
import org.jspecify.annotations.Nullable;

/**
Expand All @@ -29,6 +30,6 @@
*/
public interface MailReceiver {

Object @Nullable [] receive() throws jakarta.mail.MessagingException;
Object @Nullable [] receive() throws MessagingException;

}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import java.io.ByteArrayOutputStream;
import java.nio.charset.Charset;

import jakarta.mail.Message;
import jakarta.mail.Multipart;
import jakarta.mail.Part;

Expand Down Expand Up @@ -52,7 +53,7 @@ public void setCharset(String charset) {
}

@Override
protected AbstractIntegrationMessageBuilder<String> doTransform(jakarta.mail.Message mailMessage) {
protected AbstractIntegrationMessageBuilder<String> doTransform(Message mailMessage) {
try {
String payload;
Object content = mailMessage.getContent();
Expand Down
Loading