Skip to content

Commit 6f1aa5d

Browse files
adityaanikamolegz
authored andcommitted
fix(context): route KafkaNull payloads through conversion for non-Message types
convertInputIfNecessary returned the raw, unconverted Message as soon as it saw a KafkaNull payload, regardless of what the target function declared as its input type. For a function bound to a concrete type such as Consumer<MyType>, that Message was then cast to the declared type, throwing "class GenericMessage cannot be cast to class MyType" -- an error that gives no indication a Kafka tombstone was involved. The shortcut also ran before any MessageConverter or MessageConverterHelper was consulted, so the hook added in gh-1168 could not see these messages at all. Guard the shortcut with isInputTypeMessage(), the same predicate the class already uses elsewhere to test whether the declared input type is itself Message-compatible. A function genuinely declared to accept Message<?> still receives the raw message exactly as before; anything else now follows the normal conversion path. Fixes gh-1448 Resolves #1456 Signed-off-by: adityaanikam <adityanikam9502@gmail.com>
1 parent 389404c commit 6f1aa5d

3 files changed

Lines changed: 111 additions & 2 deletions

File tree

‎spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/catalog/SimpleFunctionRegistry.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1214,7 +1214,8 @@ else if (this.skipInputConversion) {
12141214
}
12151215
else if (input instanceof Message) {
12161216
input = this.filterOutHeaders((Message) input);
1217-
if (((Message) input).getPayload().getClass().getName().equals("org.springframework.kafka.support.KafkaNull")) {
1217+
if (this.isInputTypeMessage()
1218+
&& ((Message) input).getPayload().getClass().getName().equals("org.springframework.kafka.support.KafkaNull")) {
12181219
return input;
12191220
}
12201221

‎spring-cloud-function-context/src/test/java/org/springframework/cloud/function/context/catalog/SimpleFunctionRegistryTests.java‎

Lines changed: 73 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
import org.springframework.core.ResolvableType;
7070
import org.springframework.core.convert.ConversionService;
7171
import org.springframework.core.convert.support.DefaultConversionService;
72+
import org.springframework.kafka.support.KafkaNull;
7273
import org.springframework.lang.Nullable;
7374
import org.springframework.messaging.Message;
7475
import org.springframework.messaging.MessageHeaders;
@@ -157,6 +158,62 @@ public void concurrencyRegistrationTest() throws Exception {
157158
assertThat(c.size()).isEqualTo(1);
158159
}
159160

161+
@Test
162+
public void testKafkaNullWithConcreteConsumerTypeNoLongerReachesFunctionAsRawMessage() {
163+
// Regression test for #1448: a Kafka tombstone (KafkaNull payload) bound
164+
// to a Consumer<ConcreteType> used to bypass conversion entirely and
165+
// reach the function as a raw, unconverted Message, throwing an opaque
166+
// "GenericMessage cannot be cast to <Type>" ClassCastException that gave
167+
// no hint a Kafka tombstone was involved.
168+
//
169+
// With the fix, the message now flows through the same conversion path
170+
// as any other message. For a plain (non-generic) declared parameter
171+
// type such as this one, SmartCompositeMessageConverter's two-argument
172+
// fromMessage(Message, Class) overload does not consult a registered
173+
// MessageConverterHelper when every converter quietly returns null
174+
// (only when a converter throws), so the failure still surfaces as a
175+
// ClassCastException here rather than a MessageConversionException --
176+
// a pre-existing, orthogonal limitation of that overload, unrelated to
177+
// this fix. But the payload is now unwrapped before the cast, so the
178+
// exception names the real cause (KafkaNull) instead of the opaque
179+
// wrapping Message type -- a genuine diagnostic improvement.
180+
CompositeMessageConverter converter = new SmartCompositeMessageConverter(
181+
List.of(new ByteArrayMessageConverter()));
182+
183+
FunctionRegistration<ConsumePerson> registration = new FunctionRegistration<>(
184+
new ConsumePerson(), "consumePerson").type(ConsumePerson.class);
185+
SimpleFunctionRegistry catalog = new SimpleFunctionRegistry(this.conversionService, converter,
186+
new JacksonMapper(new ObjectMapper()));
187+
catalog.register(registration);
188+
FunctionInvocationWrapper lookedUpFunction = catalog.lookup("consumePerson");
189+
190+
Message<Object> kafkaNullMessage = MessageBuilder.withPayload((Object) KafkaNull.INSTANCE).build();
191+
192+
Assertions.assertThatThrownBy(() -> lookedUpFunction.apply(kafkaNullMessage))
193+
.isInstanceOf(ClassCastException.class)
194+
.hasMessageContaining("KafkaNull")
195+
.hasMessageNotContaining("GenericMessage");
196+
}
197+
198+
@Test
199+
public void testKafkaNullWithMessageTypedConsumerStillPassesThroughUnconverted() {
200+
// A function genuinely declared to accept Message<?>/KafkaNull must keep
201+
// working exactly as before -- only the mismatched-type case changes.
202+
ConsumeMessage function = new ConsumeMessage();
203+
FunctionRegistration<ConsumeMessage> registration = new FunctionRegistration<>(
204+
function, "consumeMessage").type(ConsumeMessage.class);
205+
SimpleFunctionRegistry catalog = new SimpleFunctionRegistry(this.conversionService, this.messageConverter,
206+
new JacksonMapper(new ObjectMapper()));
207+
catalog.register(registration);
208+
FunctionInvocationWrapper lookedUpFunction = catalog.lookup("consumeMessage");
209+
210+
Message<Object> kafkaNullMessage = MessageBuilder.withPayload((Object) KafkaNull.INSTANCE).build();
211+
lookedUpFunction.apply(kafkaNullMessage);
212+
213+
assertThat(function.received).isNotNull();
214+
assertThat(function.received.getPayload()).isSameAs(KafkaNull.INSTANCE);
215+
}
216+
160217
@Test
161218
public void testCachingOfFunction() {
162219
Echo function = new Echo();
@@ -820,9 +877,24 @@ public Object apply(Object t) {
820877

821878
}
822879

880+
private static final class ConsumePerson implements Consumer<Person> {
881+
@Override
882+
public void accept(Person person) {
883+
fail("function must not be invoked when conversion fails");
884+
}
885+
}
886+
887+
private static final class ConsumeMessage implements Consumer<Message<Object>> {
888+
private volatile Message<Object> received;
889+
890+
@Override
891+
public void accept(Message<Object> message) {
892+
this.received = message;
893+
}
894+
}
895+
823896
private static final class UpperCaseMessage
824897
implements Function<Message<String>, Message<String>> {
825-
826898
@Override
827899
public Message<String> apply(Message<String> t) {
828900
return MessageBuilder.withPayload(t.getPayload().toUpperCase(Locale.ROOT))
Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
/*
2+
* Copyright 2012-present the original author or authors.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* https://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
17+
package org.springframework.kafka.support;
18+
19+
/**
20+
* Minimal test double for {@code org.springframework.kafka.support.KafkaNull}.
21+
* {@code SimpleFunctionRegistry} detects a Kafka tombstone payload by comparing
22+
* {@code getClass().getName()} against this exact fully-qualified name (to avoid
23+
* a hard compile dependency on spring-kafka), so a class with the same name and
24+
* package is sufficient to exercise that code path in tests without pulling in
25+
* the real spring-kafka dependency.
26+
*
27+
* @author Aditya Nikam
28+
*/
29+
public final class KafkaNull {
30+
31+
public static final KafkaNull INSTANCE = new KafkaNull();
32+
33+
private KafkaNull() {
34+
}
35+
36+
}

0 commit comments

Comments
 (0)