DanielLeens commented on PR #11559:
URL: https://github.com/apache/seatunnel/pull/11559#issuecomment-5651183049

   Confirmed still present at the current head `4d75844b4a`, unchanged by my 
two recent commits — you're right that neither `779004c073` nor `4d75844b4a` 
touched `KafkaSourceReader.java`; this protection predates both, introduced by 
`f40ad4eaa1e1930ca9fa29003f632819345837e1` ("[Fix][Zeta] Fix dynamic lookup 
checkpoint and gate recovery"), which is round 10's fix per the thread history 
and is unaffected by my dev merge.
   
   `KafkaSourceReader.java` at `4d75844b4a`, class `KafkaGateObjectInputStream` 
(lines 418-441):
   
   ```java
   /** Restricts restored gate split payloads to the split classes owned by 
this reader. */
   private static final class KafkaGateObjectInputStream extends 
ObjectInputStream {
   
       private KafkaGateObjectInputStream(ByteArrayInputStream input) throws 
IOException {
           super(input);
       }
   
       @Override
       protected Class<?> resolveClass(ObjectStreamClass descriptor)
               throws IOException, ClassNotFoundException {
           String className = descriptor.getName();
           if (isAllowedClass(className)) {
               return super.resolveClass(descriptor);
           }
           throw new IOException("Rejected Kafka source gate split class: " + 
className);
       }
   
       private static boolean isAllowedClass(String className) {
           return className.equals("[B")
                   || className.equals("java.lang.String")
                   || className.equals("java.lang.Long")
                   || className.equals("java.lang.Integer")
                   || className.equals(KafkaSourceSplit.class.getName())
                   || className.equals(TopicPartition.class.getName())
                   || className.equals(
                           
org.apache.seatunnel.api.table.catalog.TablePath.class.getName());
       }
   }
   ```
   
   The `readObject()` call site is `deserializeSplit(byte[])` at line 401, 
which constructs an instance of this class rather than a plain 
`ObjectInputStream`:
   
   ```java
   private static KafkaSourceSplit deserializeSplit(byte[] serializedSplit)
           throws IOException, ClassNotFoundException {
       try (ObjectInputStream objectInputStream =
               new KafkaGateObjectInputStream(new 
ByteArrayInputStream(serializedSplit))) {
           Object split = objectInputStream.readObject();
           ...
   ```
   
   So this is a closed allowlist of seven exact class names (`resolveClass` 
throws `IOException` on anything else), not a package-prefix allowlist — the 
tightest of the three `*ObjectInputStream` sites in this PR. I re-verified this 
is exercised, not dead code, by tracing the caller chain back to the gate-state 
restore path.
   


-- 
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