timsaucer commented on code in PR #24973:
URL: https://github.com/apache/datafusion/pull/24973#discussion_r4106906018


##########
datafusion/proto/src/physical_plan/mod.rs:
##########
@@ -2093,6 +2334,56 @@ impl ComposedPhysicalExtensionCodec {
             .encode(buf)
             .map_err(|e| internal_datafusion_err!("{e}"))
     }
+
+    /// Like [`Self::encode_protobuf`], but for the function hooks whose trait
+    /// default is `Ok(())` rather than an error.
+    ///
+    /// Those hooks treat an empty buffer as "no custom payload, encode by
+    /// name", and the decode side only consults the function registry when the
+    /// payload is absent. Wrapping an empty blob in a [`DataEncoderTuple`]
+    /// would make a by-name function look codec-encoded and strand it at
+    /// decode time, so a codec that writes nothing must leave `buf` untouched.
+    fn encode_protobuf_by_name_aware(
+        &self,
+        buf: &mut Vec<u8>,
+        mut encode: impl FnMut(&dyn PhysicalExtensionCodec, &mut Vec<u8>) -> 
Result<()>,
+    ) -> Result<()> {

Review Comment:
   A larger test
   
   
   ```
           /// Runs `f` against encode/decode contexts backed by a
           /// `ComposedPhysicalExtensionCodec` over `codecs`, with an empty
           /// function registry so every by-name decode takes the codec 
fallback.
           fn with_composed_ctx(
               codecs: Vec<Arc<dyn PhysicalExtensionCodec>>,
               f: impl FnOnce(&ExecutionPlanEncodeCtx<'_>, 
&ExecutionPlanDecodeCtx<'_>) -> Result<()>,
           ) -> Result<()> {
               let composed = ComposedPhysicalExtensionCodec::new(codecs);
               let converter = DefaultPhysicalProtoConverter {};
               let encoder = encode_ctx_over(&composed, &converter);
               let encode_ctx = ExecutionPlanEncodeCtx::new(&encoder);
   
               let task_ctx = TaskContext::default();
               let decode_context = PhysicalPlanDecodeContext::new(&task_ctx, 
&composed);
               let decoder = ConverterPlanDecoder {
                   ctx: &decode_context,
                   proto_converter: &converter,
               };
               let decode_ctx = ExecutionPlanDecodeCtx::new(&decoder);
   
               f(&encode_ctx, &decode_ctx)
           }
   
           /// Control: with the by-name codec at position 0, the registry-miss
           /// fallback (`try_decode_*(name, &[])`) reaches it through the
           /// composed codec.
           #[test]
           fn composed_by_name_fallback_reaches_codec_at_position_0() -> 
Result<()> {
               with_composed_ctx(vec![Arc::new(EmptyPayloadOnlyCodec)], |enc, 
dec| {
                   let payload = 
enc.encode_udf(&ScalarUDF::from(TestUdf::new()))?;
                   assert!(payload.is_none());
                   assert_eq!(dec.decode_udf("test_udf", 
payload.as_deref())?.name(), "test_udf");
   
                   let payload = 
enc.encode_udaf(&AggregateUDF::from(TestUdaf::new()))?;
                   assert!(payload.is_none());
                   assert_eq!(dec.decode_udaf("test_udaf", 
payload.as_deref())?.name(), "test_udaf");
                   Ok(())
               })
           }
   
           // The tests below place the by-name codec at position 1. Encoding
           // writes no payload, so decode misses the registry and calls the
           // composed codec with `&[]`. An empty buffer decodes as
           // `DataEncoderTuple { encoder_position: 0, blob: [] }`, so only 
codec 0
           // (which cannot resolve the function) is ever asked.
   
           #[test]
           fn composed_by_name_udf_fallback_reaches_later_codec() -> Result<()> 
{
               with_composed_ctx(
                   vec![
                       Arc::new(DefaultPhysicalExtensionCodec {}),
                       Arc::new(EmptyPayloadOnlyCodec),
                   ],
                   |enc, dec| {
                       let payload = 
enc.encode_udf(&ScalarUDF::from(TestUdf::new()))?;
                       assert!(payload.is_none(), "expected by-name encoding");
                       assert_eq!(dec.decode_udf("test_udf", 
payload.as_deref())?.name(), "test_udf");
                       Ok(())
                   },
               )
           }
   
           #[test]
           fn composed_by_name_udaf_fallback_reaches_later_codec() -> 
Result<()> {
               with_composed_ctx(
                   vec![
                       Arc::new(DefaultPhysicalExtensionCodec {}),
                       Arc::new(EmptyPayloadOnlyCodec),
                   ],
                   |enc, dec| {
                       let payload = 
enc.encode_udaf(&AggregateUDF::from(TestUdaf::new()))?;
                       assert!(payload.is_none(), "expected by-name encoding");
                       assert_eq!(dec.decode_udaf("test_udaf", 
payload.as_deref())?.name(), "test_udaf");
                       Ok(())
                   },
               )
           }
   
           #[test]
           fn composed_by_name_udwf_fallback_reaches_later_codec() -> 
Result<()> {
               with_composed_ctx(
                   vec![
                       Arc::new(DefaultPhysicalExtensionCodec {}),
                       Arc::new(EmptyPayloadOnlyCodec),
                   ],
                   |enc, dec| {
                       let payload = 
enc.encode_udwf(&WindowUDF::from(TestUdwf::new()))?;
                       assert!(payload.is_none(), "expected by-name encoding");
                       assert_eq!(dec.decode_udwf("test_udwf", 
payload.as_deref())?.name(), "test_udwf");
                       Ok(())
                   },
               )
           }
   ```



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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to