From 0f05547f42c2ccd53722a5c6adf4992e6750aaba Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Sat, 12 Sep 2026 06:31:16 -0400 Subject: [PATCH] feat(rest): add policy management APIs --- crates/paimon/src/api/api_request.rs | 151 ++++- crates/paimon/src/api/api_response.rs | 38 +- crates/paimon/src/api/management.rs | 536 +++++++++++++++++- crates/paimon/src/api/mod.rs | 11 +- crates/paimon/src/api/resource_paths.rs | 28 + crates/paimon/src/api/rest_api.rs | 91 ++- .../paimon/src/catalog/rest/rest_catalog.rs | 40 +- crates/paimon/tests/mock_server.rs | 227 +++++++- crates/paimon/tests/rest_api_test.rs | 228 +++++++- crates/paimon/tests/rest_catalog_test.rs | 42 +- 10 files changed, 1362 insertions(+), 30 deletions(-) diff --git a/crates/paimon/src/api/api_request.rs b/crates/paimon/src/api/api_request.rs index 42ac80ada..7cc1febd4 100644 --- a/crates/paimon/src/api/api_request.rs +++ b/crates/paimon/src/api/api_request.rs @@ -22,7 +22,10 @@ use serde::{Deserialize, Deserializer, Serialize}; use std::collections::HashMap; -use crate::api::management::{PermissionAccess, PermissionAssignment, PermissionResource}; +use crate::api::management::{ + bad_request, is_blank, ColumnMask, DataPolicy, PermissionAccess, PermissionAssignment, + PermissionResource, PolicyType, RowFilter, +}; use crate::{ catalog::{Function, FunctionDefinition, Identifier, ViewSchema}, spec::{DataField, PartitionStatistics, Schema, SchemaChange}, @@ -353,6 +356,62 @@ impl RevokePermissionRequest { } } +/// Body of `POST .../tables/{table}/policies`: the policy without its resource, which the +/// path already names. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct PolicyRequest { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub row_filter: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub column_mask: Option, + pub principal: String, +} + +impl From<&DataPolicy> for PolicyRequest { + fn from(policy: &DataPolicy) -> Self { + Self { + row_filter: policy.row_filter().cloned(), + column_mask: policy.column_mask().cloned(), + principal: policy.principal().to_string(), + } + } +} + +/// Body of `POST .../tables/{table}/policies/drop`: a policy identity (Java `DropPolicyRequest`). +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DropPolicyRequest { + #[serde(rename = "type")] + pub policy_type: PolicyType, + pub principal: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub column: Option, +} + +impl DropPolicyRequest { + pub fn new( + policy_type: PolicyType, + principal: &str, + column: Option<&str>, + ) -> crate::Result { + PermissionAssignment::validate_principal(principal)?; + let column = column.filter(|column| !is_blank(column)); + match (policy_type, column) { + (PolicyType::RowFilter, Some(_)) => { + Err(bad_request("ROW_FILTER identity cannot contain a column.")) + } + (PolicyType::ColumnMasking, None) => Err(bad_request( + "column is required for COLUMN_MASKING identity.", + )), + (_, column) => Ok(Self { + policy_type, + principal: principal.to_string(), + column: column.map(str::to_string), + }), + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -434,6 +493,96 @@ mod tests { ); } + #[test] + fn test_policy_request_requires_a_definition_and_carries_no_resource() { + let policy = DataPolicy::new_column_mask( + PermissionResource::table("sales", "orders"), + ColumnMask::new("email", "{}").unwrap(), + "analyst", + ) + .unwrap(); + assert_eq!( + serde_json::to_value(PolicyRequest::from(&policy)).unwrap(), + serde_json::json!({"columnMask": {"onColumn": "email", "transform": "{}"}, "principal": "analyst"}) + ); + // A policy with neither definition, or with both, cannot be deserialized at all, so it + // can never reach a request body. Java's `@JsonCreator` constructor rejects the same two. + for body in [ + r#"{"resource":{"type":"TABLE","database":"sales","table":"orders"},"principal":"analyst"}"#, + r#"{"resource":{"type":"TABLE","database":"sales","table":"orders"},"rowFilter":{"predicate":"{}"},"columnMask":{"onColumn":"email","transform":"{}"},"principal":"analyst"}"#, + ] { + let error = serde_json::from_str::(body).unwrap_err(); + assert!(error.to_string().contains("exactly one"), "{error}"); + } + } + + #[test] + fn test_drop_policy_request_round_trip_and_identity_rules() { + let request = + DropPolicyRequest::new(PolicyType::ColumnMasking, "analyst", Some("email")).unwrap(); + let json = serde_json::to_string(&request).unwrap(); + assert_eq!( + json, + r#"{"type":"COLUMN_MASKING","principal":"analyst","column":"email"}"# + ); + assert_eq!( + serde_json::from_str::(&json).unwrap(), + request + ); + assert_eq!( + serde_json::to_string( + &DropPolicyRequest::new(PolicyType::RowFilter, "analyst", Some(" ")).unwrap() + ) + .unwrap(), + r#"{"type":"ROW_FILTER","principal":"analyst"}"# + ); + let message = |result: crate::Result| result.unwrap_err().to_string(); + assert!(message(DropPolicyRequest::new( + PolicyType::ColumnMasking, + "analyst", + None + )) + .contains("column is required")); + assert!(message(DropPolicyRequest::new( + PolicyType::RowFilter, + "analyst", + Some("email") + )) + .contains("cannot contain a column")); + assert!( + message(DropPolicyRequest::new(PolicyType::RowFilter, " ", None)) + .contains("principal cannot be empty") + ); + + // Blank here is Java `String.trim()`, not Rust's Unicode `trim()`: NUL is blank and a + // non-breaking space is not. The two disagree in opposite directions. + assert_eq!( + DropPolicyRequest::new(PolicyType::ColumnMasking, "analyst", Some("\u{a0}")) + .unwrap() + .column + .as_deref(), + Some("\u{a0}") + ); + assert!(message(DropPolicyRequest::new( + PolicyType::RowFilter, + "analyst", + Some("\u{a0}") + )) + .contains("cannot contain a column")); + assert!( + DropPolicyRequest::new(PolicyType::RowFilter, "analyst", Some("\0")) + .unwrap() + .column + .is_none() + ); + assert!(message(DropPolicyRequest::new( + PolicyType::ColumnMasking, + "analyst", + Some("\0") + )) + .contains("column is required")); + } + #[test] fn test_create_partitions_request_serialization() { let req = CreatePartitionsRequest::new( diff --git a/crates/paimon/src/api/api_response.rs b/crates/paimon/src/api/api_response.rs index 20d1b7b77..be95d9894 100644 --- a/crates/paimon/src/api/api_response.rs +++ b/crates/paimon/src/api/api_response.rs @@ -22,7 +22,7 @@ use serde::{Deserialize, Deserializer, Serialize}; use std::collections::HashMap; -use crate::api::management::PermissionAssignment; +use crate::api::management::{DataPolicy, PermissionAssignment}; use crate::catalog::{Function, FunctionDefinition, ViewSchema}; use crate::spec::{DataField, Schema, Snapshot}; @@ -41,6 +41,9 @@ pub struct ErrorResponse { } impl ErrorResponse { + /// `resource_type` of a 404/409 about a policy (Java `ErrorResponse.RESOURCE_TYPE_POLICY`). + pub const RESOURCE_TYPE_POLICY: &'static str = "POLICY"; + /// Create a new ErrorResponse. pub fn new( resource_type: Option, @@ -523,6 +526,25 @@ impl ListPermissionsResponse { } } +/// Response of `GET .../tables/{table}/policies`. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct ListPoliciesResponse { + #[serde(default)] + pub policies: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub next_page_token: Option, +} + +impl ListPoliciesResponse { + pub fn new(policies: Vec, next_page_token: Option) -> Self { + Self { + policies, + next_page_token, + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -681,6 +703,20 @@ mod tests { ); } + #[test] + fn test_list_policies_response_deserialization() { + let response: ListPoliciesResponse = serde_json::from_str( + r#"{"policies":[{"resource":{"type":"TABLE","database":"sales","table":"orders"},"rowFilter":{"predicate":"{}"},"principal":"analyst"}],"nextPageToken":"next"}"#, + ) + .unwrap(); + assert_eq!(response.policies.len(), 1); + assert_eq!(response.policies[0].row_filter().unwrap().predicate(), "{}"); + assert_eq!(response.next_page_token.as_deref(), Some("next")); + let last: ListPoliciesResponse = serde_json::from_str("{}").unwrap(); + assert!(last.policies.is_empty()); + assert_eq!(ErrorResponse::RESOURCE_TYPE_POLICY, "POLICY"); + } + #[test] fn test_error_response_serialization() { let resp = ErrorResponse::new( diff --git a/crates/paimon/src/api/management.rs b/crates/paimon/src/api/management.rs index 188ba370f..874fa35c2 100644 --- a/crates/paimon/src/api/management.rs +++ b/crates/paimon/src/api/management.rs @@ -15,13 +15,13 @@ // specific language governing permissions and limitations // under the License. -//! REST management API model: permission assignments and (later) data policies. +//! REST management API model: permission assignments and data policies. //! -//! Mirrors Java `org.apache.paimon.management`, whose request-side constructors store a -//! corrected value rather than only checking it; `Deserialize` does neither, so a server may -//! list values a client could not have sent. The send paths therefore run `canonicalized` on -//! the types whose Java constructor corrects, `validate` on the types whose Java constructor -//! only checks, and send what comes back. +//! Mirrors Java `org.apache.paimon.management`. The three policy types deserialize through +//! their constructors, as Java does with `@JsonCreator`, so they cannot hold a value a client +//! could not have sent. The permission types follow Java the other way: its `@JsonCreator` +//! there is the non-validating constructor, so `Deserialize` is derived and the send paths run +//! `canonicalized` or `validate` instead. use std::fmt; use std::str::FromStr; @@ -216,6 +216,15 @@ impl PermissionResource { ) } + pub fn validate_policy_attachment(&self) -> Result<()> { + if self.resource_type != ResourceType::Table { + return Err(bad_request( + "Policies can currently be attached only to TABLE resources.", + )); + } + Ok(()) + } + pub fn new( resource_type: ResourceType, database: Option<&str>, @@ -670,6 +679,322 @@ impl ListPermissionsRequest { } } +/// The kind of a data policy (Java `PolicyType`). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum PolicyType { + RowFilter, + ColumnMasking, +} + +impl PolicyType { + pub fn as_str(&self) -> &'static str { + match self { + PolicyType::RowFilter => "ROW_FILTER", + PolicyType::ColumnMasking => "COLUMN_MASKING", + } + } +} + +impl fmt::Display for PolicyType { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.as_str()) + } +} + +impl FromStr for PolicyType { + type Err = Error; + + /// Case-insensitive, like Java `PolicyType.fromString`. + fn from_str(value: &str) -> Result { + [PolicyType::RowFilter, PolicyType::ColumnMasking] + .into_iter() + .find(|policy_type| policy_type.as_str() == value.to_uppercase()) + .ok_or_else(|| bad_request(format!("Unknown policy type '{value}'."))) + } +} + +impl Serialize for PolicyType { + fn serialize(&self, serializer: S) -> std::result::Result { + serializer.serialize_str(self.as_str()) + } +} + +impl<'de> Deserialize<'de> for PolicyType { + fn deserialize>(deserializer: D) -> std::result::Result { + String::deserialize(deserializer)? + .parse() + .map_err(serde::de::Error::custom) + } +} + +/// Largest serialized predicate or transform a policy may carry, in UTF-8 bytes. +const MAX_POLICY_PAYLOAD_BYTES: usize = 60 * 1024; + +fn check_payload(value: &str, field: &str) -> Result<()> { + if is_blank(value) { + return Err(bad_request(format!("{field} cannot be empty."))); + } + if value.len() > MAX_POLICY_PAYLOAD_BYTES { + return Err(bad_request(format!( + "{field} must not exceed {MAX_POLICY_PAYLOAD_BYTES} UTF-8 bytes." + ))); + } + Ok(()) +} + +/// A serialized Paimon `Predicate` applied to every scan of the table (Java `RowFilter`). +/// Same JSON as one entry of `AuthTableQueryResponse::filter`. +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct RowFilter { + predicate: String, +} + +impl RowFilter { + pub const MAX_PREDICATE_BYTES: usize = MAX_POLICY_PAYLOAD_BYTES; + + pub fn new(predicate: impl Into) -> Result { + let predicate = predicate.into(); + check_payload(&predicate, "predicate")?; + Ok(Self { predicate }) + } + + pub fn predicate(&self) -> &str { + &self.predicate + } +} + +/// A serialized Paimon `Transform` whose result replaces `on_column` (Java `ColumnMask`). +/// Same JSON as one value of `AuthTableQueryResponse::column_masking`. +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct ColumnMask { + on_column: String, + transform: String, +} + +impl ColumnMask { + pub const MAX_TRANSFORM_BYTES: usize = MAX_POLICY_PAYLOAD_BYTES; + + pub fn new(on_column: impl Into, transform: impl Into) -> Result { + let on_column = on_column.into(); + if is_blank(&on_column) { + return Err(bad_request("onColumn cannot be empty.")); + } + let transform = transform.into(); + check_payload(&transform, "transform")?; + Ok(Self { + on_column, + transform, + }) + } + + pub fn on_column(&self) -> &str { + &self.on_column + } + + pub fn transform(&self) -> &str { + &self.transform + } +} + +/// One row filter or one column mask attached to a table for one principal +/// (Java `DataPolicy`). A row filter is identified by `(table, principal)`, a column mask by +/// `(table, principal, on_column)`. +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct DataPolicy { + resource: PermissionResource, + #[serde(skip_serializing_if = "Option::is_none")] + row_filter: Option, + #[serde(skip_serializing_if = "Option::is_none")] + column_mask: Option, + principal: String, +} + +impl DataPolicy { + /// Java `DataPolicy.rowFilter`. The `new_` prefix is forced: `row_filter` is the getter. + pub fn new_row_filter( + resource: PermissionResource, + row_filter: RowFilter, + principal: &str, + ) -> Result { + Self::new(resource, Some(row_filter), None, principal) + } + + /// Java `DataPolicy.columnMask`. + pub fn new_column_mask( + resource: PermissionResource, + column_mask: ColumnMask, + principal: &str, + ) -> Result { + Self::new(resource, None, Some(column_mask), principal) + } + + pub fn new( + resource: PermissionResource, + row_filter: Option, + column_mask: Option, + principal: &str, + ) -> Result { + Self { + resource, + row_filter, + column_mask, + principal: principal.to_string(), + } + .canonicalized() + } + + /// A policy corrects nothing of its own; only its resource can need canonicalizing. + fn canonicalized(mut self) -> Result { + self.resource.validate_policy_attachment()?; + self.resource = self.resource.canonicalized()?; + PermissionAssignment::validate_principal(&self.principal)?; + if self.row_filter.is_none() == self.column_mask.is_none() { + return Err(bad_request( + "A policy must contain exactly one of rowFilter and columnMask.", + )); + } + Ok(self) + } + + pub fn policy_type(&self) -> PolicyType { + if self.row_filter.is_some() { + PolicyType::RowFilter + } else { + PolicyType::ColumnMasking + } + } + + pub fn resource(&self) -> &PermissionResource { + &self.resource + } + + pub fn row_filter(&self) -> Option<&RowFilter> { + self.row_filter.as_ref() + } + + pub fn column_mask(&self) -> Option<&ColumnMask> { + self.column_mask.as_ref() + } + + pub fn principal(&self) -> &str { + &self.principal + } +} + +// The three policy types deserialize through their constructors, so a value arriving from the +// wire has passed the same checks as one built locally. Java gets this from `@JsonCreator`. + +impl<'de> Deserialize<'de> for RowFilter { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + #[derive(Deserialize)] + struct Repr { + predicate: String, + } + let repr = Repr::deserialize(deserializer)?; + Self::new(repr.predicate).map_err(serde::de::Error::custom) + } +} + +impl<'de> Deserialize<'de> for ColumnMask { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + #[derive(Deserialize)] + #[serde(rename_all = "camelCase")] + struct Repr { + on_column: String, + transform: String, + } + let repr = Repr::deserialize(deserializer)?; + Self::new(repr.on_column, repr.transform).map_err(serde::de::Error::custom) + } +} + +impl<'de> Deserialize<'de> for DataPolicy { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + #[derive(Deserialize)] + #[serde(rename_all = "camelCase")] + struct Repr { + resource: PermissionResource, + #[serde(default)] + row_filter: Option, + #[serde(default)] + column_mask: Option, + principal: String, + } + let repr = Repr::deserialize(deserializer)?; + Self::new( + repr.resource, + repr.row_filter, + repr.column_mask, + &repr.principal, + ) + .map_err(serde::de::Error::custom) + } +} + +/// Filters for `GET .../tables/{table}/policies` (Java `ListPoliciesRequest`). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ListPoliciesRequest { + /// The table the policies are attached to; must be a `TABLE` resource. + pub resource: PermissionResource, + pub policy_type: Option, + pub principal: Option, + /// Only masks on this column; requires `policy_type == Some(PolicyType::ColumnMasking)`. + pub column: Option, + pub max_results: Option, + pub page_token: Option, +} + +impl ListPoliciesRequest { + pub fn new(resource: PermissionResource) -> Self { + Self { + resource, + policy_type: None, + principal: None, + column: None, + max_results: None, + page_token: None, + } + } + + /// The validated query string as `(name, value)` pairs, in the order Java sends them. + pub fn query_params(&self) -> Result> { + self.resource.validate_policy_attachment()?; + let mut params = Vec::new(); + if let Some(policy_type) = self.policy_type { + params.push(("type", policy_type.to_string())); + } + if let Some(principal) = self.principal.as_deref().filter(|value| !is_blank(value)) { + PermissionAssignment::validate_principal(principal)?; + params.push(("principal", principal.to_string())); + } + if let Some(column) = self.column.as_deref().filter(|value| !is_blank(value)) { + if self.policy_type != Some(PolicyType::ColumnMasking) { + return Err(bad_request("column filter requires type COLUMN_MASKING.")); + } + params.push(("column", column.to_string())); + } + if let Some(max_results) = self.max_results { + validate_max_results(max_results)?; + params.push(("maxResults", max_results.to_string())); + } + if let Some(page_token) = self.page_token.as_deref().filter(|value| !value.is_empty()) { + params.push(("pageToken", page_token.to_string())); + } + Ok(params) + } +} + #[cfg(test)] mod tests { use super::*; @@ -1295,4 +1620,203 @@ mod tests { ] ); } + + const PREDICATE_JSON: &str = r#"{"kind":"LEAF","transform":{"name":"FIELD_REF","fieldRef":{"index":0,"name":"region","type":"STRING"}},"function":"EQUAL","literals":["APAC"]}"#; + const TRANSFORM_JSON: &str = + r#"{"name":"CONCAT","inputs":[{"index":0,"name":"region","type":"STRING"},"****"]}"#; + + #[test] + fn test_policy_type_wire_names() { + assert_eq!( + serde_json::to_string(&PolicyType::RowFilter).unwrap(), + r#""ROW_FILTER""# + ); + assert_eq!( + serde_json::from_str::(r#""column_masking""#).unwrap(), + PolicyType::ColumnMasking + ); + assert!("MASK".parse::().is_err()); + assert_eq!(PolicyType::ColumnMasking.to_string(), "COLUMN_MASKING"); + } + + #[test] + fn test_policies_round_trip_java_wire_json() { + let mask = DataPolicy::new_column_mask( + table_resource(), + ColumnMask::new("email", TRANSFORM_JSON).unwrap(), + "analyst", + ) + .unwrap(); + let json = serde_json::to_string(&mask).unwrap(); + assert_eq!( + json, + format!( + r#"{{"resource":{{"type":"TABLE","database":"sales","table":"orders"}},"columnMask":{{"onColumn":"email","transform":{}}},"principal":"analyst"}}"#, + serde_json::to_string(TRANSFORM_JSON).unwrap() + ) + ); + let round_trip: DataPolicy = serde_json::from_str(&json).unwrap(); + assert_eq!(round_trip, mask); + assert_eq!(round_trip.policy_type(), PolicyType::ColumnMasking); + assert_eq!(round_trip.column_mask().unwrap().on_column(), "email"); + assert_eq!( + round_trip.column_mask().unwrap().transform(), + TRANSFORM_JSON + ); + assert!(round_trip.row_filter().is_none()); + + let filter = DataPolicy::new_row_filter( + table_resource(), + RowFilter::new(PREDICATE_JSON).unwrap(), + "analyst", + ) + .unwrap(); + let round_trip: DataPolicy = + serde_json::from_str(&serde_json::to_string(&filter).unwrap()).unwrap(); + assert_eq!(round_trip.policy_type(), PolicyType::RowFilter); + assert_eq!(round_trip.row_filter().unwrap().predicate(), PREDICATE_JSON); + assert!(round_trip.column_mask().is_none()); + assert_eq!(round_trip.principal(), "analyst"); + assert_eq!(round_trip.resource(), &table_resource()); + } + + #[test] + fn test_policy_validation_and_payload_bounds() { + let message = |error: Error| error.to_string(); + assert!(message(RowFilter::new(" ").unwrap_err()).contains("predicate cannot be empty")); + assert!(message(ColumnMask::new("email", " ").unwrap_err()) + .contains("transform cannot be empty")); + assert!(message(ColumnMask::new(" ", TRANSFORM_JSON).unwrap_err()) + .contains("onColumn cannot be empty")); + assert!(message( + DataPolicy::new_column_mask( + PermissionResource::catalog(), + ColumnMask::new("email", TRANSFORM_JSON).unwrap(), + "analyst" + ) + .unwrap_err() + ) + .contains("only to TABLE")); + assert!(message( + DataPolicy::new_row_filter( + table_resource(), + RowFilter::new(PREDICATE_JSON).unwrap(), + " " + ) + .unwrap_err() + ) + .contains("principal cannot be empty")); + assert!(message( + RowFilter::new("p".repeat(RowFilter::MAX_PREDICATE_BYTES + 1)).unwrap_err() + ) + .contains("UTF-8 bytes")); + assert!(RowFilter::new("p".repeat(RowFilter::MAX_PREDICATE_BYTES)).is_ok()); + assert!(message( + ColumnMask::new("email", "t".repeat(ColumnMask::MAX_TRANSFORM_BYTES + 1)).unwrap_err() + ) + .contains("UTF-8 bytes")); + assert!(PermissionResource::column("sales", "orders") + .validate_policy_attachment() + .is_err()); + assert!(table_resource().validate_policy_attachment().is_ok()); + } + + #[test] + fn test_policy_deserialization_runs_the_same_checks_as_the_constructors() { + let rejected = |json: &str| { + serde_json::from_str::(json) + .unwrap_err() + .to_string() + }; + let policy = |definition: &str| -> String { + rejected(&format!( + r#"{{"resource":{{"type":"TABLE","database":"sales","table":"orders"}},{definition}"principal":"analyst"}}"# + )) + }; + assert!(serde_json::from_str::(r#"{"predicate":" "}"#) + .unwrap_err() + .to_string() + .contains("predicate cannot be empty")); + assert!(serde_json::from_value::( + serde_json::json!({"predicate": "p".repeat(RowFilter::MAX_PREDICATE_BYTES + 1)}) + ) + .unwrap_err() + .to_string() + .contains("UTF-8 bytes")); + assert!( + serde_json::from_str::(r#"{"onColumn":" ","transform":" "}"#) + .unwrap_err() + .to_string() + .contains("onColumn cannot be empty") + ); + assert!(policy(r#""rowFilter":{"predicate":" "},"#).contains("predicate cannot be empty")); + assert!( + policy(r#""columnMask":{"onColumn":"email","transform":" "},"#) + .contains("transform cannot be empty") + ); + for definition in [ + "", + r#""rowFilter":{"predicate":"{}"},"columnMask":{"onColumn":"email","transform":"{}"},"#, + ] { + assert!(policy(definition).contains("exactly one")); + } + assert!(rejected( + r#"{"resource":{"type":"TABLE","database":"sales","table":"orders"},"rowFilter":{"predicate":"{}"},"principal":" "}"# + ) + .contains("principal cannot be empty")); + assert!(rejected( + r#"{"resource":{"type":"TABLE","database":"sales","table":""},"rowFilter":{"predicate":"{}"},"principal":"analyst"}"# + ) + .contains("table is required for TABLE")); + + // A blank locator that canonicalizing can drop is corrected, not rejected, on the way in. + let corrected: DataPolicy = serde_json::from_str( + r#"{"resource":{"type":"TABLE","database":"sales","table":"orders","view":""},"rowFilter":{"predicate":"{}"},"principal":"analyst"}"#, + ) + .unwrap(); + assert_eq!(corrected.resource(), &table_resource()); + } + + #[test] + fn test_list_policies_request_query_params() { + let mut request = ListPoliciesRequest::new(table_resource()); + request.policy_type = Some(PolicyType::ColumnMasking); + request.principal = Some("analyst".to_string()); + request.column = Some("email".to_string()); + request.max_results = Some(25); + request.page_token = Some(" \t".to_string()); + assert_eq!( + request.query_params().unwrap(), + vec![ + ("type", "COLUMN_MASKING".to_string()), + ("principal", "analyst".to_string()), + ("column", "email".to_string()), + ("maxResults", "25".to_string()), + ("pageToken", " \t".to_string()), + ] + ); + assert!(ListPoliciesRequest::new(table_resource()) + .query_params() + .unwrap() + .is_empty()); + let mut request = ListPoliciesRequest::new(table_resource()); + request.column = Some("email".to_string()); + assert!(request + .query_params() + .unwrap_err() + .to_string() + .contains("COLUMN_MASKING")); + let mut request = ListPoliciesRequest::new(table_resource()); + request.max_results = Some(1001); + assert!(request + .query_params() + .unwrap_err() + .to_string() + .contains("1000")); + assert!(ListPoliciesRequest::new(PermissionResource::catalog()) + .query_params() + .unwrap_err() + .to_string() + .contains("only to TABLE")); + } } diff --git a/crates/paimon/src/api/mod.rs b/crates/paimon/src/api/mod.rs index 888567416..f7c0e0786 100644 --- a/crates/paimon/src/api/mod.rs +++ b/crates/paimon/src/api/mod.rs @@ -34,8 +34,8 @@ mod api_response; pub use api_request::{ AlterDatabaseRequest, AlterTableRequest, AuthTableQueryRequest, CreateDatabaseRequest, CreateFunctionRequest, CreatePartitionsRequest, CreateTableRequest, CreateTagRequest, - CreateViewRequest, DropPartitionsRequest, ListPartitionsByFilterRequest, - ListPartitionsByNamesRequest, RenameTableRequest, RevokePermissionRequest, + CreateViewRequest, DropPartitionsRequest, DropPolicyRequest, ListPartitionsByFilterRequest, + ListPartitionsByNamesRequest, PolicyRequest, RenameTableRequest, RevokePermissionRequest, }; // Re-export response types @@ -43,13 +43,14 @@ pub use api_response::{ AuditRESTResponse, AuthTableQueryResponse, ConfigResponse, ErrorResponse, GetDatabaseResponse, GetFunctionResponse, GetTableResponse, GetTableTokenResponse, GetTagResponse, GetViewResponse, ListDatabasesResponse, ListFunctionsResponse, ListPartitionsResponse, ListPermissionsResponse, - ListTablesResponse, ListViewsResponse, PagedList, + ListPoliciesResponse, ListTablesResponse, ListViewsResponse, PagedList, }; // Re-export management types pub use management::{ - ListPermissionsRequest, PermissionAccess, PermissionAssignment, PermissionColumns, - PermissionResource, ResourceType, + ColumnMask, DataPolicy, ListPermissionsRequest, ListPoliciesRequest, PermissionAccess, + PermissionAssignment, PermissionColumns, PermissionResource, PolicyType, ResourceType, + RowFilter, }; // Re-export error types diff --git a/crates/paimon/src/api/resource_paths.rs b/crates/paimon/src/api/resource_paths.rs index bc911db28..470b2f9b0 100644 --- a/crates/paimon/src/api/resource_paths.rs +++ b/crates/paimon/src/api/resource_paths.rs @@ -37,6 +37,7 @@ impl ResourcePaths { const VIEWS: &'static str = "views"; const FUNCTIONS: &'static str = "functions"; const PERMISSIONS: &'static str = "permissions"; + const POLICIES: &'static str = "policies"; /// Create a new ResourcePaths with the given prefix. pub fn new(prefix: &str) -> Self { @@ -274,6 +275,20 @@ impl ResourcePaths { pub fn revoke_permission(&self) -> String { format!("{}/revoke", self.permissions()) } + + /// Get the policy collection nested below its table (`.../tables/{table}/policies`). + pub fn policies(&self, database_name: &str, table_name: &str) -> String { + format!( + "{}/{}", + self.table(database_name, table_name), + Self::POLICIES + ) + } + + /// Get the action endpoint that drops one policy from its table. + pub fn drop_policy(&self, database_name: &str, table_name: &str) -> String { + format!("{}/drop", self.policies(database_name, table_name)) + } } #[cfg(test)] @@ -396,4 +411,17 @@ mod tests { assert_eq!(paths.grant_permission(), "/v1/catalog/permissions/grant"); assert_eq!(paths.revoke_permission(), "/v1/catalog/permissions/revoke"); } + + #[test] + fn test_policy_paths_nest_under_the_table_and_encode_names() { + let paths = ResourcePaths::new("catalog"); + assert_eq!( + paths.policies("sales db", "orders/all"), + "/v1/catalog/databases/sales+db/tables/orders%2Fall/policies" + ); + assert_eq!( + paths.drop_policy("sales db", "orders/all"), + "/v1/catalog/databases/sales+db/tables/orders%2Fall/policies/drop" + ); + } } diff --git a/crates/paimon/src/api/rest_api.rs b/crates/paimon/src/api/rest_api.rs index cac1c33f8..6fe0c8fe3 100644 --- a/crates/paimon/src/api/rest_api.rs +++ b/crates/paimon/src/api/rest_api.rs @@ -31,18 +31,22 @@ use crate::Result; use super::api_request::{ AlterDatabaseRequest, AlterTableRequest, AuthTableQueryRequest, CreateDatabaseRequest, CreateFunctionRequest, CreatePartitionsRequest, CreateTableRequest, CreateTagRequest, - CreateViewRequest, DropPartitionsRequest, ListPartitionsByFilterRequest, - ListPartitionsByNamesRequest, RenameTableRequest, RevokePermissionRequest, + CreateViewRequest, DropPartitionsRequest, DropPolicyRequest, ListPartitionsByFilterRequest, + ListPartitionsByNamesRequest, PolicyRequest, RenameTableRequest, RevokePermissionRequest, }; use super::api_response::{ - AuthTableQueryResponse, ConfigResponse, GetDatabaseResponse, GetFunctionResponse, - GetTableResponse, GetTagResponse, GetViewResponse, ListDatabasesResponse, - ListFunctionsResponse, ListPartitionsResponse, ListPermissionsResponse, ListTablesResponse, - ListViewsResponse, PagedList, + AuthTableQueryResponse, ConfigResponse, ErrorResponse, GetDatabaseResponse, + GetFunctionResponse, GetTableResponse, GetTagResponse, GetViewResponse, ListDatabasesResponse, + ListFunctionsResponse, ListPartitionsResponse, ListPermissionsResponse, ListPoliciesResponse, + ListTablesResponse, ListViewsResponse, PagedList, }; use super::auth::{AuthProviderFactory, RESTAuthFunction}; -use super::management::{ListPermissionsRequest, PermissionAssignment, PermissionResource}; +use super::management::{ + bad_request, DataPolicy, ListPermissionsRequest, ListPoliciesRequest, PermissionAssignment, + PermissionResource, PolicyType, +}; use super::resource_paths::ResourcePaths; +use super::rest_error::RestError; use super::rest_util::RESTUtil; /// Validate that a string is not empty after trimming. @@ -76,6 +80,18 @@ fn validate_non_empty_multi(values: &[(&str, &str)]) -> Result<()> { Ok(()) } +/// The canonical `(database, table)` of a policy resource, shared by all three calls. +fn policy_table(resource: PermissionResource) -> Result<(String, String)> { + resource.validate_policy_attachment()?; + let resource = resource.canonicalized()?; + match (resource.database_name(), resource.table_name()) { + (Some(database), Some(table)) => Ok((database.to_string(), table.to_string())), + _ => Err(bad_request( + "A TABLE resource needs both database and table.", + )), + } +} + /// REST API wrapper for Paimon catalog operations. /// /// This struct provides methods for database and table CRUD operations @@ -937,6 +953,67 @@ impl RESTApi { Ok(()) } + // ==================== Policy Management ==================== + // + // Experimental REST management API (Java `RESTPolicyManagement`). + + pub async fn list_policies_paged( + &self, + request: &ListPoliciesRequest, + ) -> Result> { + let params = request.query_params()?; + let (database, table) = policy_table(request.resource.clone())?; + let path = self.resource_paths.policies(&database, &table); + let response: ListPoliciesResponse = if params.is_empty() { + self.client.get(&path, None::<&[(&str, &str)]>).await? + } else { + self.client.get(&path, Some(¶ms)).await? + }; + Ok(PagedList::new(response.policies, response.next_page_token)) + } + + /// Create a policy on its table. A policy with the same identity surfaces as + /// `RestError::AlreadyExists` whose `resource_type` is `POLICY`. + pub async fn create_policy(&self, policy: &DataPolicy) -> Result<()> { + let (database, table) = policy_table(policy.resource().clone())?; + let request = PolicyRequest::from(policy); + let path = self.resource_paths.policies(&database, &table); + let _resp: serde_json::Value = self.client.post(&path, &request).await?; + Ok(()) + } + + /// Drop one policy by identity. `ignore_if_not_exists` swallows a 404 on the policy only: + /// a 404 on the table still surfaces. + pub async fn drop_policy( + &self, + resource: &PermissionResource, + policy_type: PolicyType, + principal: &str, + column: Option<&str>, + ignore_if_not_exists: bool, + ) -> Result<()> { + let (database, table) = policy_table(resource.clone())?; + let request = DropPolicyRequest::new(policy_type, principal, column)?; + let path = self.resource_paths.drop_policy(&database, &table); + match self + .client + .post::(&path, &request) + .await + { + Ok(_) => Ok(()), + Err(crate::Error::RestApi { + source: + RestError::NoSuchResource { + resource_type: Some(resource_type), + .. + }, + }) if ignore_if_not_exists && resource_type == ErrorResponse::RESOURCE_TYPE_POLICY => { + Ok(()) + } + Err(error) => Err(error), + } + } + // ==================== Commit Operations ==================== /// Commit a snapshot for a table. diff --git a/crates/paimon/src/catalog/rest/rest_catalog.rs b/crates/paimon/src/catalog/rest/rest_catalog.rs index fb4f47de1..724dbcdbb 100644 --- a/crates/paimon/src/catalog/rest/rest_catalog.rs +++ b/crates/paimon/src/catalog/rest/rest_catalog.rs @@ -25,10 +25,12 @@ use std::sync::Arc; use async_trait::async_trait; -use crate::api::management::{ListPermissionsRequest, PermissionAssignment, PermissionResource}; use crate::api::rest_api::RESTApi; use crate::api::rest_error::RestError; -use crate::api::{GetTagResponse, PagedList}; +use crate::api::{ + DataPolicy, GetTagResponse, ListPermissionsRequest, ListPoliciesRequest, PagedList, + PermissionAssignment, PermissionResource, PolicyType, +}; use crate::catalog::{ list_partitions_from_file_system, Catalog, Database, Identifier, DB_LOCATION_PROP, }; @@ -154,6 +156,40 @@ impl RESTCatalog { .revoke_permission(resource, access, principal) .await } + + // ======================= policy management ============================== + // + // Java `RESTCatalog.policyManagement()`. + + pub async fn list_policies_paged( + &self, + request: &ListPoliciesRequest, + ) -> Result> { + self.api.list_policies_paged(request).await + } + + pub async fn create_policy(&self, policy: &DataPolicy) -> Result<()> { + self.api.create_policy(policy).await + } + + pub async fn drop_policy( + &self, + resource: &PermissionResource, + policy_type: PolicyType, + principal: &str, + column: Option<&str>, + ignore_if_not_exists: bool, + ) -> Result<()> { + self.api + .drop_policy( + resource, + policy_type, + principal, + column, + ignore_if_not_exists, + ) + .await + } } // ============================================================================ diff --git a/crates/paimon/tests/mock_server.rs b/crates/paimon/tests/mock_server.rs index 929e16500..c8a64f9cd 100644 --- a/crates/paimon/tests/mock_server.rs +++ b/crates/paimon/tests/mock_server.rs @@ -36,12 +36,12 @@ use tokio::task::JoinHandle; use paimon::api::{ AlterDatabaseRequest, AlterTableRequest, AuditRESTResponse, ConfigResponse, CreateFunctionRequest, CreatePartitionsRequest, CreateTagRequest, CreateViewRequest, - DropPartitionsRequest, ErrorResponse, GetDatabaseResponse, GetFunctionResponse, - GetTableResponse, GetTagResponse, GetViewResponse, ListDatabasesResponse, + DataPolicy, DropPartitionsRequest, DropPolicyRequest, ErrorResponse, GetDatabaseResponse, + GetFunctionResponse, GetTableResponse, GetTagResponse, GetViewResponse, ListDatabasesResponse, ListFunctionsResponse, ListPartitionsByFilterRequest, ListPartitionsByNamesRequest, - ListPartitionsResponse, ListPermissionsResponse, ListTablesResponse, ListViewsResponse, - PermissionAssignment, PermissionResource, RenameTableRequest, ResourcePaths, ResourceType, - RevokePermissionRequest, + ListPartitionsResponse, ListPermissionsResponse, ListPoliciesResponse, ListTablesResponse, + ListViewsResponse, PermissionAssignment, PermissionResource, PolicyRequest, PolicyType, + RenameTableRequest, ResourcePaths, ResourceType, RevokePermissionRequest, }; use paimon::catalog::{Function, Identifier}; use paimon::spec::{CommitKind, Partition, Snapshot}; @@ -78,6 +78,13 @@ struct MockState { grant_permission_bodies: Vec, revoke_permission_bodies: Vec, grant_permission_error_status: Option, + /// Policies per `"{db}.{table}"`. + policies: HashMap>, + list_policies_queries: Vec>, + create_policy_bodies: Vec, + drop_policy_bodies: Vec, + create_policy_error: Option, + drop_policy_error: Option, /// ECS metadata role name (for token loader testing) ecs_role_name: Option, /// ECS metadata token (for token loader testing) @@ -1366,6 +1373,173 @@ impl RESTServer { StatusCode::OK.into_response() } + // ==================== Policy management ==================== + + fn policy_identity(policy: &DataPolicy) -> (PolicyType, String, Option) { + ( + policy.policy_type(), + policy.principal().to_string(), + policy + .column_mask() + .map(|mask| mask.on_column().to_string()), + ) + } + + fn policy_error(error: ErrorResponse) -> axum::response::Response { + let status = StatusCode::from_u16(error.code.unwrap_or(500) as u16) + .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); + (status, Json(error)).into_response() + } + + fn table_not_found(table: &str) -> axum::response::Response { + Self::policy_error(ErrorResponse::new( + Some("TABLE".to_string()), + Some(table.to_string()), + Some("Table not found".to_string()), + Some(404), + )) + } + + /// Handle GET .../tables/{table}/policies - policies on the table, filtered by identity parts. + pub async fn list_policies( + Path((db, table)): Path<(String, String)>, + Query(params): Query>, + Extension(state): Extension>, + ) -> impl IntoResponse { + let mut inner = state.inner.lock().unwrap(); + inner.list_policies_queries.push(params.clone()); + let key = format!("{db}.{table}"); + if !inner.tables.contains_key(&key) { + return Self::table_not_found(&table); + } + let matching: Vec = inner + .policies + .get(&key) + .cloned() + .unwrap_or_default() + .into_iter() + .filter(|policy| { + params + .get("type") + .is_none_or(|t| t == policy.policy_type().as_str()) + }) + .filter(|policy| { + params + .get("principal") + .is_none_or(|p| p == policy.principal()) + }) + .filter(|policy| { + params.get("column").is_none_or(|column| { + policy + .column_mask() + .is_some_and(|mask| mask.on_column() == column) + }) + }) + .collect(); + let page_size = params + .get("maxResults") + .and_then(|value| value.parse().ok()); + let (policies, next_page_token) = paginate(matching, ¶ms, page_size); + ( + StatusCode::OK, + Json(ListPoliciesResponse::new(policies, next_page_token)), + ) + .into_response() + } + + /// Handle POST .../tables/{table}/policies - 404 TABLE for an unknown table, 409 POLICY + /// for a second policy with the same identity. + pub async fn create_policy( + Path((db, table)): Path<(String, String)>, + Extension(state): Extension>, + Json(body): Json, + ) -> impl IntoResponse { + let mut inner = state.inner.lock().unwrap(); + inner.create_policy_bodies.push(body.clone()); + if let Some(error) = inner.create_policy_error.clone() { + return Self::policy_error(error); + } + let request: PolicyRequest = match serde_json::from_value(body) { + Ok(request) => request, + Err(error) => return Self::bad_request(error.to_string()), + }; + let key = format!("{db}.{table}"); + if !inner.tables.contains_key(&key) { + return Self::table_not_found(&table); + } + let resource = PermissionResource::table(db, table); + let policy = match (request.row_filter, request.column_mask) { + (Some(row_filter), None) => { + DataPolicy::new_row_filter(resource, row_filter, &request.principal) + } + (None, Some(column_mask)) => { + DataPolicy::new_column_mask(resource, column_mask, &request.principal) + } + _ => return Self::bad_request("exactly one of rowFilter and columnMask".to_string()), + }; + let policy = match policy { + Ok(policy) => policy, + Err(error) => return Self::bad_request(error.to_string()), + }; + let existing = inner.policies.entry(key).or_default(); + if existing + .iter() + .any(|candidate| Self::policy_identity(candidate) == Self::policy_identity(&policy)) + { + let (policy_type, principal, column) = Self::policy_identity(&policy); + let name = match column { + Some(column) => format!("{policy_type}:{principal}:{column}"), + None => format!("{policy_type}:{principal}"), + }; + return Self::policy_error(ErrorResponse::new( + Some(ErrorResponse::RESOURCE_TYPE_POLICY.to_string()), + Some(name), + Some("Policy already exists.".to_string()), + Some(409), + )); + } + existing.push(policy); + StatusCode::OK.into_response() + } + + /// Handle POST .../tables/{table}/policies/drop - 404 POLICY when the identity is absent. + pub async fn drop_policy( + Path((db, table)): Path<(String, String)>, + Extension(state): Extension>, + Json(body): Json, + ) -> impl IntoResponse { + let mut inner = state.inner.lock().unwrap(); + inner.drop_policy_bodies.push(body.clone()); + if let Some(error) = inner.drop_policy_error.clone() { + return Self::policy_error(error); + } + let request: DropPolicyRequest = match serde_json::from_value(body) { + Ok(request) => request, + Err(error) => return Self::bad_request(error.to_string()), + }; + let key = format!("{db}.{table}"); + if !inner.tables.contains_key(&key) { + return Self::table_not_found(&table); + } + let identity = ( + request.policy_type, + request.principal.clone(), + request.column.clone(), + ); + let policies = inner.policies.entry(key).or_default(); + let before = policies.len(); + policies.retain(|policy| Self::policy_identity(policy) != identity); + if policies.len() == before { + return Self::policy_error(ErrorResponse::new( + Some(ErrorResponse::RESOURCE_TYPE_POLICY.to_string()), + Some(format!("{}:{}", request.policy_type, request.principal)), + Some("Policy does not exist.".to_string()), + Some(404), + )); + } + StatusCode::OK.into_response() + } + /// Handle POST /rename-table - rename a table. pub async fn rename_table( Extension(state): Extension>, @@ -1782,6 +1956,41 @@ impl RESTServer { self.inner.lock().unwrap().grant_permission_error_status = status; } + /// Every query string received by `GET .../tables/{table}/policies`. + pub fn list_policies_queries(&self) -> Vec> { + self.inner.lock().unwrap().list_policies_queries.clone() + } + + /// Raw JSON bodies received by `POST .../tables/{table}/policies`. + pub fn create_policy_bodies(&self) -> Vec { + self.inner.lock().unwrap().create_policy_bodies.clone() + } + + /// Raw JSON bodies received by `POST .../tables/{table}/policies/drop`. + pub fn drop_policy_bodies(&self) -> Vec { + self.inner.lock().unwrap().drop_policy_bodies.clone() + } + + pub fn table_policies(&self, database: &str, table: &str) -> Vec { + self.inner + .lock() + .unwrap() + .policies + .get(&format!("{database}.{table}")) + .cloned() + .unwrap_or_default() + } + + /// Make every create-policy call answer with `error` (status from its `code`). + pub fn set_create_policy_error(&self, error: Option) { + self.inner.lock().unwrap().create_policy_error = error; + } + + /// Make every drop-policy call answer with `error` (status from its `code`). + pub fn set_drop_policy_error(&self, error: Option) { + self.inner.lock().unwrap().drop_policy_error = error; + } + /// Return all create-partitions calls received by the server. pub fn create_partitions_calls(&self) -> Vec<(String, String, CreatePartitionsRequest)> { self.inner.lock().unwrap().create_partitions_calls.clone() @@ -1965,6 +2174,14 @@ pub async fn start_mock_server( &format!("{prefix}/permissions/revoke"), post(RESTServer::revoke_permission), ) + .route( + &format!("{prefix}/databases/:db/tables/:table/policies"), + get(RESTServer::list_policies).post(RESTServer::create_policy), + ) + .route( + &format!("{prefix}/databases/:db/tables/:table/policies/drop"), + post(RESTServer::drop_policy), + ) // ECS metadata endpoints (for token loader testing) .route( "/ram/security-credentials/", diff --git a/crates/paimon/tests/rest_api_test.rs b/crates/paimon/tests/rest_api_test.rs index f8ef468b2..9ca3d15ae 100644 --- a/crates/paimon/tests/rest_api_test.rs +++ b/crates/paimon/tests/rest_api_test.rs @@ -26,8 +26,9 @@ use axum::http::StatusCode; use paimon::api::auth::{DLFECSTokenLoader, DLFToken, DLFTokenLoader}; use paimon::api::rest_api::RESTApi; use paimon::api::{ - ConfigResponse, CreatePartitionsRequest, DropPartitionsRequest, ListPermissionsRequest, - PermissionAssignment, PermissionColumns, PermissionResource, RestError, + ColumnMask, ConfigResponse, CreatePartitionsRequest, DataPolicy, DropPartitionsRequest, + ErrorResponse, ListPermissionsRequest, ListPoliciesRequest, PermissionAssignment, + PermissionColumns, PermissionResource, PolicyType, RestError, RowFilter, }; use paimon::catalog::{Function, FunctionDefinition, Identifier, ViewSchema}; use paimon::common::Options; @@ -1239,3 +1240,226 @@ async fn test_ecs_loader_token() { .to_string() .contains("Failed to parse token Expiration")); } + +// ==================== Policy Management Tests ==================== + +const PREDICATE_JSON: &str = r#"{"kind":"LEAF","transform":{"name":"FIELD_REF","fieldRef":{"index":0,"name":"region","type":"STRING"}},"function":"EQUAL","literals":["APAC"]}"#; +const TRANSFORM_JSON: &str = + r#"{"name":"CONCAT","inputs":[{"index":0,"name":"region","type":"STRING"},"****"]}"#; + +fn managed_table() -> PermissionResource { + PermissionResource::table("default", "managed_table") +} + +fn mask_policy() -> DataPolicy { + DataPolicy::new_column_mask( + managed_table(), + ColumnMask::new("email", TRANSFORM_JSON).unwrap(), + "analyst", + ) + .unwrap() +} + +fn policy_error(resource_type: &str, code: i32) -> ErrorResponse { + ErrorResponse::new( + Some(resource_type.to_string()), + Some("orders".to_string()), + Some(format!("{resource_type} error")), + Some(code), + ) +} + +#[tokio::test] +async fn test_policies_round_trip_through_the_table_nested_endpoints() { + let (ctx, _identifier) = setup_partition_api().await; + let mask = mask_policy(); + let filter = DataPolicy::new_row_filter( + managed_table(), + RowFilter::new(PREDICATE_JSON).unwrap(), + "analyst", + ) + .unwrap(); + ctx.api.create_policy(&mask).await.unwrap(); + ctx.api.create_policy(&filter).await.unwrap(); + + let create = &ctx.server.create_policy_bodies()[0]; + assert_eq!( + *create, + json!({"columnMask": {"onColumn": "email", "transform": TRANSFORM_JSON}, "principal": "analyst"}) + ); + assert!(create.get("resource").is_none()); + + let mut request = ListPoliciesRequest::new(managed_table()); + request.policy_type = Some(PolicyType::ColumnMasking); + request.principal = Some("analyst".to_string()); + request.column = Some("email".to_string()); + request.max_results = Some(25); + request.page_token = Some("0".to_string()); + let page = ctx.api.list_policies_paged(&request).await.unwrap(); + assert_eq!(page.elements, vec![mask]); + assert_eq!(page.next_page_token, None); + assert_eq!( + ctx.server.list_policies_queries(), + vec![HashMap::from([ + ("type".to_string(), "COLUMN_MASKING".to_string()), + ("principal".to_string(), "analyst".to_string()), + ("column".to_string(), "email".to_string()), + ("maxResults".to_string(), "25".to_string()), + ("pageToken".to_string(), "0".to_string()), + ])] + ); + + ctx.api + .drop_policy( + &managed_table(), + PolicyType::ColumnMasking, + "analyst", + Some("email"), + false, + ) + .await + .unwrap(); + assert_eq!( + ctx.server.drop_policy_bodies(), + vec![json!({"type": "COLUMN_MASKING", "principal": "analyst", "column": "email"})] + ); + let page = ctx + .api + .list_policies_paged(&ListPoliciesRequest::new(managed_table())) + .await + .unwrap(); + assert_eq!(page.elements, vec![filter]); +} + +#[tokio::test] +async fn test_create_policy_surfaces_conflicts_with_their_resource_type() { + let (ctx, _identifier) = setup_partition_api().await; + ctx.api.create_policy(&mask_policy()).await.unwrap(); + let error = ctx.api.create_policy(&mask_policy()).await.unwrap_err(); + assert!( + matches!(&error, paimon::Error::RestApi { source: RestError::AlreadyExists { resource_type: Some(t), .. } } if t == "POLICY"), + "{error:?}" + ); + ctx.server + .set_create_policy_error(Some(policy_error("TABLE", 409))); + let error = ctx.api.create_policy(&mask_policy()).await.unwrap_err(); + assert!( + matches!(&error, paimon::Error::RestApi { source: RestError::AlreadyExists { resource_type: Some(t), .. } } if t == "TABLE"), + "{error:?}" + ); + // An unknown table is a 404 on the table, not on the policy. + ctx.server.set_create_policy_error(None); + let error = ctx + .api + .create_policy( + &DataPolicy::new_column_mask( + PermissionResource::table("default", "missing"), + ColumnMask::new("email", "{}").unwrap(), + "analyst", + ) + .unwrap(), + ) + .await + .unwrap_err(); + assert!( + matches!(&error, paimon::Error::RestApi { source: RestError::NoSuchResource { resource_type: Some(t), .. } } if t == "TABLE"), + "{error:?}" + ); +} + +#[tokio::test] +async fn test_drop_policy_if_exists_ignores_only_a_missing_policy() { + let (ctx, _identifier) = setup_partition_api().await; + ctx.api + .drop_policy( + &managed_table(), + PolicyType::ColumnMasking, + "analyst", + Some("email"), + true, + ) + .await + .unwrap(); + let error = ctx + .api + .drop_policy( + &managed_table(), + PolicyType::ColumnMasking, + "analyst", + Some("email"), + false, + ) + .await + .unwrap_err(); + assert!( + matches!(&error, paimon::Error::RestApi { source: RestError::NoSuchResource { resource_type: Some(t), .. } } if t == "POLICY"), + "{error:?}" + ); + ctx.server + .set_drop_policy_error(Some(policy_error("TABLE", 404))); + let error = ctx + .api + .drop_policy( + &managed_table(), + PolicyType::ColumnMasking, + "analyst", + Some("email"), + true, + ) + .await + .unwrap_err(); + assert!( + matches!(&error, paimon::Error::RestApi { source: RestError::NoSuchResource { resource_type: Some(t), .. } } if t == "TABLE"), + "{error:?}" + ); + assert_eq!(ctx.server.drop_policy_bodies().len(), 3); +} + +#[tokio::test] +async fn test_invalid_policy_requests_never_reach_the_server() { + let (ctx, _identifier) = setup_partition_api().await; + let error = ctx + .api + .drop_policy( + &managed_table(), + PolicyType::RowFilter, + "analyst", + Some("email"), + false, + ) + .await + .unwrap_err(); + assert!( + error.to_string().contains("cannot contain a column"), + "{error}" + ); + let error = ctx + .api + .list_policies_paged(&ListPoliciesRequest::new(PermissionResource::catalog())) + .await + .unwrap_err(); + assert!(error.to_string().contains("only to TABLE"), "{error}"); + // A blank locator is not a TABLE-only violation, so only the resource rules catch it. + let blank_table = PermissionResource::table("default", ""); + let error = ctx + .api + .list_policies_paged(&ListPoliciesRequest::new(blank_table.clone())) + .await + .unwrap_err(); + assert!( + error.to_string().contains("table is required for TABLE"), + "{error}" + ); + let error = ctx + .api + .drop_policy(&blank_table, PolicyType::RowFilter, "analyst", None, false) + .await + .unwrap_err(); + assert!( + error.to_string().contains("table is required for TABLE"), + "{error}" + ); + assert!(ctx.server.create_policy_bodies().is_empty()); + assert!(ctx.server.drop_policy_bodies().is_empty()); + assert!(ctx.server.list_policies_queries().is_empty()); +} diff --git a/crates/paimon/tests/rest_catalog_test.rs b/crates/paimon/tests/rest_catalog_test.rs index 3d7805e0a..7deb412f7 100644 --- a/crates/paimon/tests/rest_catalog_test.rs +++ b/crates/paimon/tests/rest_catalog_test.rs @@ -28,7 +28,10 @@ use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as Arr use axum::http::StatusCode; use futures::TryStreamExt; use paimon::api::ConfigResponse; -use paimon::api::{ListPermissionsRequest, PermissionAssignment, PermissionResource}; +use paimon::api::{ + DataPolicy, ListPermissionsRequest, ListPoliciesRequest, PermissionAssignment, + PermissionResource, PolicyType, RowFilter, +}; use paimon::catalog::{Catalog, Function, FunctionDefinition, Identifier, RESTCatalog, ViewSchema}; use paimon::common::Options; use paimon::spec::{ @@ -2539,3 +2542,40 @@ async fn test_rest_catalog_manages_permissions_end_to_end() { .unwrap(); assert!(page.elements.is_empty()); } + +#[tokio::test] +async fn test_rest_catalog_manages_policies_end_to_end() { + let ctx = setup_catalog(vec!["default"]).await; + let identifier = Identifier::new("default", "orders"); + ctx.catalog + .create_table(&identifier, test_schema(), false) + .await + .unwrap(); + let resource = PermissionResource::table("default", "orders"); + let predicate = r#"{"kind":"LEAF","transform":{"name":"FIELD_REF","fieldRef":{"index":0,"name":"id","type":"BIGINT"}},"function":"EQUAL","literals":[1]}"#; + let policy = DataPolicy::new_row_filter( + resource.clone(), + RowFilter::new(predicate).unwrap(), + "analyst", + ) + .unwrap(); + + ctx.catalog.create_policy(&policy).await.unwrap(); + let page = ctx + .catalog + .list_policies_paged(&ListPoliciesRequest::new(resource.clone())) + .await + .unwrap(); + assert_eq!(page.elements, vec![policy]); + + ctx.catalog + .drop_policy(&resource, PolicyType::RowFilter, "analyst", None, false) + .await + .unwrap(); + let page = ctx + .catalog + .list_policies_paged(&ListPoliciesRequest::new(resource)) + .await + .unwrap(); + assert!(page.elements.is_empty()); +}