From 89cf135b569b2f7c943ddb1bd9d029d2b178d519 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 1 Apr 2019 11:20:34 -0400 Subject: [PATCH] GH-2874: Syslog - copy IpHeaders to message Resolves https://github.com/spring-projects/spring-integration/issues/2874 **cherry-pick to 5.1.x** # Conflicts: # spring-integration-syslog/src/test/java/org/springframework/integration/syslog/inbound/SyslogReceivingChannelAdapterTests.java --- .../integration/syslog/DefaultMessageConverter.java | 1 + .../integration/syslog/RFC5424MessageConverter.java | 3 ++- .../inbound/SyslogReceivingChannelAdapterTests.java | 9 ++++++++- 3 files changed, 11 insertions(+), 2 deletions(-) diff --git a/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/DefaultMessageConverter.java b/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/DefaultMessageConverter.java index c78d0c0a76..db162fc8d5 100644 --- a/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/DefaultMessageConverter.java +++ b/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/DefaultMessageConverter.java @@ -93,6 +93,7 @@ public class DefaultMessageConverter implements MessageConverter, BeanFactoryAwa } } return getMessageBuilderFactory().withPayload(this.asMap ? map : message.getPayload()) + .copyHeaders(message.getHeaders()) .copyHeaders(out) .build(); } diff --git a/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/RFC5424MessageConverter.java b/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/RFC5424MessageConverter.java index 7f6b4c9eb4..56ef63e258 100644 --- a/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/RFC5424MessageConverter.java +++ b/spring-integration-syslog/src/main/java/org/springframework/integration/syslog/RFC5424MessageConverter.java @@ -79,7 +79,8 @@ public class RFC5424MessageConverter extends DefaultMessageConverter { } AbstractIntegrationMessageBuilder builder = getMessageBuilderFactory().withPayload( - asMap() ? map : originalContent); + asMap() ? map : originalContent) + .copyHeaders(message.getHeaders()); if (!asMap() && isMap) { builder.copyHeaders(map); } diff --git a/spring-integration-syslog/src/test/java/org/springframework/integration/syslog/inbound/SyslogReceivingChannelAdapterTests.java b/spring-integration-syslog/src/test/java/org/springframework/integration/syslog/inbound/SyslogReceivingChannelAdapterTests.java index 14d49e8fa3..8e55c63f08 100644 --- a/spring-integration-syslog/src/test/java/org/springframework/integration/syslog/inbound/SyslogReceivingChannelAdapterTests.java +++ b/spring-integration-syslog/src/test/java/org/springframework/integration/syslog/inbound/SyslogReceivingChannelAdapterTests.java @@ -30,6 +30,7 @@ import java.net.DatagramPacket; import java.net.DatagramSocket; import java.net.InetSocketAddress; import java.net.Socket; +import java.nio.charset.StandardCharsets; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -44,6 +45,7 @@ import org.springframework.beans.factory.BeanFactory; import org.springframework.context.ApplicationEvent; import org.springframework.context.ApplicationEventPublisher; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.ip.IpHeaders; import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory; import org.springframework.integration.ip.tcp.connection.TcpNioServerConnectionFactory; import org.springframework.integration.ip.udp.UnicastReceivingChannelAdapter; @@ -77,7 +79,7 @@ public class SyslogReceivingChannelAdapterTests { UnicastReceivingChannelAdapter.class); TestingUtilities.waitListening(server, null); UdpSyslogReceivingChannelAdapter adapter = (UdpSyslogReceivingChannelAdapter) factory.getObject(); - byte[] buf = "<157>JUL 26 22:08:35 WEBERN TESTING[70729]: TEST SYSLOG MESSAGE".getBytes("UTF-8"); + byte[] buf = "<157>JUL 26 22:08:35 WEBERN TESTING[70729]: TEST SYSLOG MESSAGE".getBytes(StandardCharsets.UTF_8); DatagramPacket packet = new DatagramPacket(buf, buf.length, new InetSocketAddress("localhost", server.getPort())); DatagramSocket socket = new DatagramSocket(); @@ -86,6 +88,7 @@ public class SyslogReceivingChannelAdapterTests { Message message = outputChannel.receive(10000); assertNotNull(message); assertEquals("WEBERN", message.getHeaders().get("syslog_HOST")); + assertTrue(message.getHeaders().containsKey(IpHeaders.IP_ADDRESS)); adapter.stop(); } @@ -130,6 +133,7 @@ public class SyslogReceivingChannelAdapterTests { Message message = outputChannel.receive(10000); assertNotNull(message); assertEquals("WEBERN", message.getHeaders().get("syslog_HOST")); + assertTrue(message.getHeaders().containsKey(IpHeaders.IP_ADDRESS)); adapter.stop(); assertTrue(latch.await(10, TimeUnit.SECONDS)); } @@ -163,6 +167,7 @@ public class SyslogReceivingChannelAdapterTests { assertEquals("WEBERN", message.getHeaders().get("syslog_HOST")); assertEquals("<157>JUL 26 22:08:35 WEBERN TESTING[70729]: TEST SYSLOG MESSAGE", new String((byte[]) message.getPayload(), "UTF-8")); + assertTrue(message.getHeaders().containsKey(IpHeaders.IP_ADDRESS)); adapter.stop(); } @@ -212,6 +217,7 @@ public class SyslogReceivingChannelAdapterTests { Message> message = (Message>) outputChannel.receive(10000); assertNotNull(message); assertEquals("loggregator", message.getPayload().get("syslog_HOST")); + assertTrue(message.getHeaders().containsKey(IpHeaders.IP_ADDRESS)); adapter.stop(); assertTrue(latch.await(10, TimeUnit.SECONDS)); } @@ -245,6 +251,7 @@ public class SyslogReceivingChannelAdapterTests { Message> message = (Message>) outputChannel.receive(10000); assertNotNull(message); assertEquals("loggregator", message.getPayload().get("syslog_HOST")); + assertTrue(message.getHeaders().containsKey(IpHeaders.IP_ADDRESS)); adapter.stop(); }