GH-9792: Remove resource header in the StreamTransformer
Fixes: https://github.com/spring-projects/spring-integration/issues/9792 The `StreamTransformer` closes `IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE` header value, so this becomes unusable afterward. In addition, it may even cause some problems downstream when the message could be serialized for subsequent network interaction. * Add logic to the `StreamTransformer` to build a new message, but remove `IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE` header * Verify that header is removed in the `StreamingInboundTests`. Addition `StreamingInboundTests` clean up for better code style.
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2023 the original author or authors.
|
||||
* Copyright 2016-2025 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -97,37 +97,45 @@ public class StreamingInboundTests {
|
||||
}
|
||||
streamer.afterPropertiesSet();
|
||||
streamer.start();
|
||||
Message<byte[]> received = (Message<byte[]>) this.transformer.transform(streamer.receive());
|
||||
Message<InputStream> inputStreamMessage = streamer.receive();
|
||||
Message<byte[]> received = (Message<byte[]>) this.transformer.transform(inputStreamMessage);
|
||||
assertThat(received.getPayload()).isEqualTo("foo\nbar".getBytes());
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("foo");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "foo")
|
||||
.doesNotContainKey(IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE);
|
||||
String fileInfo = (String) received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO);
|
||||
assertThat(fileInfo).contains("remoteDirectory\":\"/foo");
|
||||
assertThat(fileInfo).contains("permissions\":\"-rw-rw-rw");
|
||||
assertThat(fileInfo).contains("size\":42");
|
||||
assertThat(fileInfo).contains("directory\":false");
|
||||
assertThat(fileInfo).contains("filename\":\"foo");
|
||||
assertThat(fileInfo).contains("modified\":42000");
|
||||
assertThat(fileInfo).contains("link\":false");
|
||||
assertThat(fileInfo)
|
||||
.contains("remoteDirectory\":\"/foo")
|
||||
.contains("permissions\":\"-rw-rw-rw")
|
||||
.contains("size\":42")
|
||||
.contains("directory\":false")
|
||||
.contains("filename\":\"foo")
|
||||
.contains("modified\":42000")
|
||||
.contains("link\":false");
|
||||
|
||||
// close after list, transform
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(received), times(2)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(inputStreamMessage), times(2)).close();
|
||||
|
||||
received = (Message<byte[]>) this.transformer.transform(streamer.receive());
|
||||
inputStreamMessage = streamer.receive();
|
||||
received = (Message<byte[]>) this.transformer.transform(inputStreamMessage);
|
||||
assertThat(received.getPayload()).isEqualTo("baz\nqux".getBytes());
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("bar");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "bar")
|
||||
.doesNotContainKey(IntegrationMessageHeaderAccessor.CLOSEABLE_RESOURCE);
|
||||
fileInfo = (String) received.getHeaders().get(FileHeaders.REMOTE_FILE_INFO);
|
||||
assertThat(fileInfo).contains("remoteDirectory\":\"/foo");
|
||||
assertThat(fileInfo).contains("permissions\":\"-rw-rw-rw");
|
||||
assertThat(fileInfo).contains("size\":42");
|
||||
assertThat(fileInfo).contains("directory\":false");
|
||||
assertThat(fileInfo).contains("filename\":\"bar");
|
||||
assertThat(fileInfo).contains("modified\":42000");
|
||||
assertThat(fileInfo).contains("link\":false");
|
||||
assertThat(fileInfo)
|
||||
.contains("remoteDirectory\":\"/foo")
|
||||
.contains("permissions\":\"-rw-rw-rw")
|
||||
.contains("size\":42")
|
||||
.contains("directory\":false")
|
||||
.contains("filename\":\"bar")
|
||||
.contains("modified\":42000")
|
||||
.contains("link\":false");
|
||||
|
||||
// close after transform
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(received), times(3)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(inputStreamMessage), times(3)).close();
|
||||
|
||||
verify(sessionFactory.getSession()).list("/foo");
|
||||
}
|
||||
@@ -142,21 +150,25 @@ public class StreamingInboundTests {
|
||||
streamer.setFilter(new AcceptOnceFileListFilter<>());
|
||||
streamer.afterPropertiesSet();
|
||||
streamer.start();
|
||||
Message<byte[]> received = (Message<byte[]>) this.transformer.transform(streamer.receive());
|
||||
Message<InputStream> inputStreamMessage = streamer.receive();
|
||||
Message<byte[]> received = (Message<byte[]>) this.transformer.transform(inputStreamMessage);
|
||||
assertThat(received.getPayload()).isEqualTo("foo\nbar".getBytes());
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("foo");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "foo");
|
||||
|
||||
// close after list, transform
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(received), times(2)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(inputStreamMessage), times(2)).close();
|
||||
|
||||
received = (Message<byte[]>) this.transformer.transform(streamer.receive());
|
||||
inputStreamMessage = streamer.receive();
|
||||
received = (Message<byte[]>) this.transformer.transform(inputStreamMessage);
|
||||
assertThat(received.getPayload()).isEqualTo("baz\nqux".getBytes());
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("bar");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "bar");
|
||||
|
||||
// close after transform
|
||||
verify(new IntegrationMessageHeaderAccessor(received).getCloseableResource(), times(3)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(inputStreamMessage), times(3)).close();
|
||||
|
||||
verify(sessionFactory.getSession()).list("/foo");
|
||||
}
|
||||
@@ -189,31 +201,35 @@ public class StreamingInboundTests {
|
||||
splitter.handleMessage(receivedStream);
|
||||
Message<?> received = out.receive(0);
|
||||
assertThat(received.getPayload()).isEqualTo("foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("foo");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "foo");
|
||||
received = out.receive(0);
|
||||
assertThat(received.getPayload()).isEqualTo("bar");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("foo");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "foo");
|
||||
assertThat(out.receive(0)).isNull();
|
||||
|
||||
// close by list, splitter
|
||||
verify(new IntegrationMessageHeaderAccessor(receivedStream).getCloseableResource(), times(3)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(receivedStream), times(3)).close();
|
||||
|
||||
receivedStream = streamer.receive();
|
||||
splitter.handleMessage(receivedStream);
|
||||
received = out.receive(0);
|
||||
assertThat(received.getPayload()).isEqualTo("baz");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("bar");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "bar");
|
||||
received = out.receive(0);
|
||||
assertThat(received.getPayload()).isEqualTo("qux");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_DIRECTORY)).isEqualTo("/foo");
|
||||
assertThat(received.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("bar");
|
||||
assertThat(received.getHeaders())
|
||||
.containsEntry(FileHeaders.REMOTE_DIRECTORY, "/foo")
|
||||
.containsEntry(FileHeaders.REMOTE_FILE, "bar");
|
||||
assertThat(out.receive(0)).isNull();
|
||||
|
||||
// close by splitter
|
||||
verify(new IntegrationMessageHeaderAccessor(receivedStream).getCloseableResource(), times(5)).close();
|
||||
verify(StaticMessageHeaderAccessor.getCloseableResource(receivedStream), times(5)).close();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
|
||||
Reference in New Issue
Block a user