plusplusjiajia commented on code in PR #513:
URL: https://github.com/apache/paimon-rust/pull/513#discussion_r4178082795


##########
crates/paimon/src/table/table_read.rs:
##########
@@ -942,21 +943,124 @@ impl<'a> PaimonTableRead<'a> {
                      scan",
                 ));
             }
-            if !grant.is_unrestricted() {
+            // The server ruled on these columns only, whether or not it set 
rules.
+            if let Some(select) = grant.select() {
+                if let Some(outside) = self
+                    .read_type
+                    .iter()
+                    .map(|f| f.name())
+                    .chain(filter_columns.iter().map(String::as_str))
+                    .find(|name| !select.iter().any(|s| s == name))
+                {
+                    return Err(super::query_auth::unsupported(&format!(
+                        "'{outside}' is outside the columns this plan was 
authorized for; plan \
+                         with it"
+                    )));
+                }
+            }
+            if !grant.is_unrestricted() && restricted.is_none() {
+                restricted = Some(grant);
+            }
+        }
+        if let Some(grant) = restricted {
+            if data_splits
+                .iter()
+                .any(|s| s.query_auth_grant().is_none_or(|g| **g != **grant))
+            {
                 return Err(super::query_auth::unsupported(
-                    "this client cannot apply a row filter or column masking, 
so it refuses \
-                     rather than return unfiltered rows",
+                    "the splits were planned under different rules; re-plan 
the scan",
                 ));
             }
         }
-        Ok(())
+        Ok(restricted.cloned())
     }
 
     /// Returns an [`ArrowRecordBatchStream`].
     pub fn to_arrow(&self, data_splits: &[DataSplit]) -> 
crate::Result<ArrowRecordBatchStream> {
-        let has_primary_keys = !self.table.schema.primary_keys().is_empty();
         let core_options = self.table.schema.core_options();
-        self.ensure_authorized_by_splits(&core_options, data_splits)?;
+        match self.ensure_authorized_by_splits(&core_options, data_splits)? {
+            Some(grant) => self.read_restricted(data_splits, &grant),
+            None => self.read_splits(data_splits, &core_options),
+        }
+    }
+
+    /// Reads what the row filter needs, filters on stored values and projects
+    /// back (Java `doAuth`).
+    fn read_restricted(
+        &self,
+        data_splits: &[DataSplit],
+        grant: &super::query_auth::QueryAuthGrant,
+    ) -> crate::Result<ArrowRecordBatchStream> {
+        use super::query_auth::{filter_batch, unsupported};
+
+        let schema_fields = self.table.schema().fields().to_vec();
+        let rules = grant.rules();
+        let index_of = |field: &DataField| schema_fields.iter().position(|s| 
s.id() == field.id());
+        let needed = rules.filter_columns();
+        // A partly projected column would feed the filter a partial value 
(Java
+        // `validateReadType`).
+        for field in &self.read_type {
+            if let Some(index) = index_of(field) {
+                if needed.contains(&index) && field.data_type() != 
schema_fields[index].data_type()
+                {
+                    return Err(unsupported(&format!(
+                        "the server's row filter reads '{}', which the read 
projects only in part",
+                        field.name()
+                    )));
+                }
+            }
+        }
+
+        let mut physical = self.read_type.clone();
+        let mut needed: Vec<usize> = needed.into_iter().collect();
+        needed.sort_unstable();
+        for index in needed {
+            let field = &schema_fields[index];
+            if !physical.iter().any(|f| f.id() == field.id()) {
+                physical.push(field.clone());
+            }
+        }
+        let mut inner = self.clone();
+        inner.read_type = physical.clone();
+        inner.limit = None;
+        let stream = inner.read_splits(data_splits, 
&self.table.schema.core_options())?;

Review Comment:
   Agreed, and reproduced with your rows. DataFusion no longer pushes 
`variant_get` into a query-auth scan, so it runs above the scan on the admitted 
rows only, the way Spark keeps a throwing extraction. A restricted read also 
refuses a strict extraction in its read type, for callers other than 
DataFusion. Your probe is now a regression test.
   



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