mjsax commented on code in PR #23159:
URL: https://github.com/apache/kafka/pull/23159#discussion_r3808809194


##########
streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchema.java:
##########
@@ -91,6 +92,9 @@ public CombinedKey<KRight, KLeft> fromBytes(final Bytes data, 
final Headers head
         final byte[] dataArray = data.get();
         final ByteBuffer dataBuffer = ByteBuffer.wrap(dataArray);
         final int foreignKeyLength = dataBuffer.getInt();
+        if (foreignKeyLength < 0 || foreignKeyLength > dataBuffer.remaining()) 
{
+            throw new BufferUnderflowException();

Review Comment:
   ```suggestion
               throw new SerializationException();
   ```



##########
streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java:
##########
@@ -184,4 +185,37 @@ public void 
shouldPassHeadersToUnderlyingDeserializersOnFromBytes() {
         verify(mockDeserializer).deserialize(PK_TOPIC, HEADERS, primaryKeyRaw);
         verify(mockDeserializer, never()).deserialize(PK_TOPIC, primaryKeyRaw);
     }
+
+    @Test
+    public void shouldThrowWhenForeignKeyLengthExceedsRemainingBytes() {
+        final CombinedKeySchema<String, Integer> cks = new CombinedKeySchema<>(
+                () -> FK_TOPIC, Serdes.String(),
+                () -> PK_TOPIC, Serdes.Integer()
+        );
+
+        final byte[] foreignKeyRaw = 
Serdes.String().serializer().serialize(FK_TOPIC, "foreignKey");
+        final byte[] primaryKeyRaw = 
Serdes.Integer().serializer().serialize(PK_TOPIC, 1);
+        final int fakeForeignKeyLength = foreignKeyRaw.length + 
primaryKeyRaw.length + 1; // one past what's actually there
+
+        final ByteBuffer buf = ByteBuffer.allocate(Integer.BYTES + 
foreignKeyRaw.length + primaryKeyRaw.length);
+        buf.putInt(fakeForeignKeyLength);
+        buf.put(foreignKeyRaw).put(primaryKeyRaw);
+        final Bytes corrupted = Bytes.wrap(buf.array());
+
+        assertThrows(BufferUnderflowException.class, () -> 
cks.fromBytes(corrupted, HEADERS));
+    }
+
+    @Test
+    public void shouldThrowWhenForeignKeyLengthIsNegative() {
+        final CombinedKeySchema<String, Integer> cks = new CombinedKeySchema<>(
+                () -> FK_TOPIC, Serdes.String(),
+                () -> PK_TOPIC, Serdes.Integer()

Review Comment:
   ```suggestion
               () -> FK_TOPIC, Serdes.String(),
               () -> PK_TOPIC, Serdes.Integer()
   ```



##########
streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java:
##########
@@ -184,4 +185,37 @@ public void 
shouldPassHeadersToUnderlyingDeserializersOnFromBytes() {
         verify(mockDeserializer).deserialize(PK_TOPIC, HEADERS, primaryKeyRaw);
         verify(mockDeserializer, never()).deserialize(PK_TOPIC, primaryKeyRaw);
     }
+
+    @Test
+    public void shouldThrowWhenForeignKeyLengthExceedsRemainingBytes() {
+        final CombinedKeySchema<String, Integer> cks = new CombinedKeySchema<>(
+                () -> FK_TOPIC, Serdes.String(),
+                () -> PK_TOPIC, Serdes.Integer()
+        );

Review Comment:
   ```suggestion
           final CombinedKeySchema<String, Integer> cks = new 
CombinedKeySchema<>(
               () -> FK_TOPIC, Serdes.String(),
               () -> PK_TOPIC, Serdes.Integer()
           );
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to