Skip to main content

iceberg/spec/
partition.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18/*!
19 * Partitioning
20 */
21use std::sync::Arc;
22
23use itertools::Itertools;
24use serde::{Deserialize, Serialize};
25use typed_builder::TypedBuilder;
26
27use super::transform::Transform;
28use super::{NestedField, PrimitiveType, Schema, SchemaRef, StructType};
29use crate::Result;
30use crate::error::invalid_data;
31use crate::spec::Struct;
32
33pub(crate) const UNPARTITIONED_LAST_ASSIGNED_ID: i32 = 999;
34pub(crate) const DEFAULT_PARTITION_SPEC_ID: i32 = 0;
35
36/// Partition fields capture the transform from table data to partition values.
37#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)]
38#[serde(rename_all = "kebab-case")]
39pub struct PartitionField {
40    /// A source column id from the table’s schema
41    pub source_id: i32,
42    /// A partition field id that is used to identify a partition field and is unique within a partition spec.
43    /// In v2 table metadata, it is unique across all partition specs.
44    pub field_id: i32,
45    /// A partition name.
46    pub name: String,
47    /// A transform that is applied to the source column to produce a partition value.
48    pub transform: Transform,
49}
50
51impl PartitionField {
52    /// To unbound partition field
53    pub fn into_unbound(self) -> UnboundPartitionField {
54        self.into()
55    }
56}
57
58/// Reference to [`PartitionSpec`].
59pub type PartitionSpecRef = Arc<PartitionSpec>;
60/// Partition spec that defines how to produce a tuple of partition values from a record.
61///
62/// A [`PartitionSpec`] is originally obtained by binding an [`UnboundPartitionSpec`] to a schema and is
63/// only guaranteed to be valid for that schema. The main difference between [`PartitionSpec`] and
64/// [`UnboundPartitionSpec`] is that the former has field ids assigned,
65/// while field ids are optional for [`UnboundPartitionSpec`].
66#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
67#[serde(rename_all = "kebab-case")]
68pub struct PartitionSpec {
69    /// Identifier for PartitionSpec
70    spec_id: i32,
71    /// Details of the partition spec
72    fields: Vec<PartitionField>,
73}
74
75impl PartitionSpec {
76    /// Create a new partition spec builder with the given schema.
77    pub fn builder(schema: impl Into<SchemaRef>) -> PartitionSpecBuilder {
78        PartitionSpecBuilder::new(schema)
79    }
80
81    /// Fields of the partition spec
82    pub fn fields(&self) -> &[PartitionField] {
83        &self.fields
84    }
85
86    /// Spec id of the partition spec
87    pub fn spec_id(&self) -> i32 {
88        self.spec_id
89    }
90
91    /// Get a new unpartitioned partition spec
92    pub fn unpartition_spec() -> Self {
93        Self {
94            spec_id: DEFAULT_PARTITION_SPEC_ID,
95            fields: vec![],
96        }
97    }
98
99    /// Returns if the partition spec is unpartitioned.
100    ///
101    /// A [`PartitionSpec`] is unpartitioned if it has no fields or all fields are [`Transform::Void`] transform.
102    pub fn is_unpartitioned(&self) -> bool {
103        self.fields.is_empty() || self.fields.iter().all(|f| f.transform == Transform::Void)
104    }
105
106    /// Returns the partition type of this partition spec.
107    /// If a source column is absent, preserves fixed transform result types and uses
108    /// unknown for result types that depend on the source type.
109    pub fn partition_type(&self, schema: &Schema) -> Result<StructType> {
110        let mut struct_fields = Vec::with_capacity(self.fields.len());
111        for partition_field in &self.fields {
112            let res_type = match schema.field_by_id(partition_field.source_id) {
113                Some(field) => partition_field.transform.result_type(&field.field_type)?,
114                // Historical specs may reference dropped source columns. Retain every
115                // field's position and any result type that is independent of its source.
116                None => match partition_field.transform {
117                    Transform::Bucket(_) | Transform::Year | Transform::Month | Transform::Hour => {
118                        PrimitiveType::Int.into()
119                    }
120                    Transform::Day => PrimitiveType::Date.into(),
121                    Transform::Unknown => PrimitiveType::String.into(),
122                    Transform::Identity | Transform::Truncate(_) | Transform::Void => {
123                        PrimitiveType::Unknown.into()
124                    }
125                },
126            };
127            struct_fields.push(
128                NestedField::optional(partition_field.field_id, &partition_field.name, res_type)
129                    .into(),
130            );
131        }
132        Ok(StructType::new(struct_fields))
133    }
134
135    /// Convert to unbound partition spec
136    pub fn into_unbound(self) -> UnboundPartitionSpec {
137        self.into()
138    }
139
140    /// Change the spec id of the partition spec
141    pub fn with_spec_id(self, spec_id: i32) -> Self {
142        Self { spec_id, ..self }
143    }
144
145    /// Check if this partition spec has sequential partition ids.
146    /// Sequential ids start from 1000 and increment by 1 for each field.
147    /// This is required for spec version 1
148    pub fn has_sequential_ids(&self) -> bool {
149        has_sequential_ids(self.fields.iter().map(|f| f.field_id))
150    }
151
152    /// Get the highest field id in the partition spec.
153    pub fn highest_field_id(&self) -> Option<i32> {
154        self.fields.iter().map(|f| f.field_id).max()
155    }
156
157    /// Check if this partition spec is compatible with another partition spec.
158    ///
159    /// Returns true if the partition spec is equal to the other spec with partition field ids ignored and
160    /// spec_id ignored. The following must be identical:
161    /// * The number of fields
162    /// * Field order
163    /// * Field names
164    /// * Source column ids
165    /// * Transforms
166    pub fn is_compatible_with(&self, other: &PartitionSpec) -> bool {
167        if self.fields.len() != other.fields.len() {
168            return false;
169        }
170
171        for (this_field, other_field) in self.fields.iter().zip(other.fields.iter()) {
172            if this_field.source_id != other_field.source_id
173                || this_field.name != other_field.name
174                || this_field.transform != other_field.transform
175            {
176                return false;
177            }
178        }
179
180        true
181    }
182
183    /// Returns partition path string containing partition type and partition
184    /// value as key-value pairs.
185    pub fn partition_to_path(&self, data: &Struct, schema: SchemaRef) -> String {
186        let partition_type = self.partition_type(&schema).unwrap();
187        let field_types = partition_type.fields();
188
189        self.fields
190            .iter()
191            .enumerate()
192            .map(|(i, field)| {
193                let value = data[i].as_ref();
194                form_urlencoded::Serializer::new(String::new())
195                    .append_pair(
196                        &field.name,
197                        &field
198                            .transform
199                            .to_human_string(&field_types[i].field_type, value),
200                    )
201                    .finish()
202            })
203            .join("/")
204    }
205}
206
207/// A partition key represents a specific partition in a table, containing the partition spec,
208/// schema, and the actual partition values.
209#[derive(Clone, Debug)]
210pub struct PartitionKey {
211    /// The partition spec that contains the partition fields.
212    spec: PartitionSpec,
213    /// The schema to which the partition spec is bound.
214    schema: SchemaRef,
215    /// Partition fields' values in struct.
216    data: Struct,
217}
218
219impl PartitionKey {
220    /// Creates a new partition key with the given spec, schema, and data.
221    pub fn new(spec: PartitionSpec, schema: SchemaRef, data: Struct) -> Self {
222        Self { spec, schema, data }
223    }
224
225    /// Creates a new partition key from another partition key, with a new data field.
226    pub fn copy_with_data(&self, data: Struct) -> Self {
227        Self {
228            spec: self.spec.clone(),
229            schema: self.schema.clone(),
230            data,
231        }
232    }
233
234    /// Generates a partition path based on the partition values.
235    pub fn to_path(&self) -> String {
236        self.spec.partition_to_path(&self.data, self.schema.clone())
237    }
238
239    /// Returns `true` if the partition key is absent (`None`)
240    /// or represents an unpartitioned spec.
241    pub fn is_effectively_none(partition_key: Option<&PartitionKey>) -> bool {
242        match partition_key {
243            None => true,
244            Some(pk) => pk.spec.is_unpartitioned(),
245        }
246    }
247
248    /// Returns the associated [`PartitionSpec`].
249    pub fn spec(&self) -> &PartitionSpec {
250        &self.spec
251    }
252
253    /// Returns the associated [`SchemaRef`].
254    pub fn schema(&self) -> &SchemaRef {
255        &self.schema
256    }
257
258    /// Returns the associated [`Struct`].
259    pub fn data(&self) -> &Struct {
260        &self.data
261    }
262}
263
264/// Reference to [`UnboundPartitionSpec`].
265pub type UnboundPartitionSpecRef = Arc<UnboundPartitionSpec>;
266/// Unbound partition field can be built without a schema and later bound to a schema.
267///
268/// The fields are private so that an instance is always well formed: in particular
269/// `source_ids` holds at least one id, which [`UnboundPartitionField::builder`] checks when the
270/// field is built.
271#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)]
272#[serde(
273    try_from = "self::_serde::UnboundPartitionFieldSerde",
274    into = "self::_serde::UnboundPartitionFieldSerde"
275)]
276#[builder(build_method(into = Result<UnboundPartitionField>))]
277pub struct UnboundPartitionField {
278    /// The source column ids from the table’s schema. A single-argument transform reads one
279    /// id; a v3 multi-argument transform reads several.
280    source_ids: Vec<i32>,
281    /// A partition field id that is used to identify a partition field and is unique within a partition spec.
282    /// In v2 table metadata, it is unique across all partition specs.
283    #[builder(default, setter(strip_option(fallback = field_id_opt)))]
284    field_id: Option<i32>,
285    /// A partition name.
286    #[builder(setter(into))]
287    name: String,
288    /// A transform that is applied to the source column to produce a partition value.
289    transform: Transform,
290}
291
292impl UnboundPartitionField {
293    /// The single source column id this field reads.
294    ///
295    /// Returns an error for a multi-argument field, which reads several columns and therefore
296    /// has no single source id. Use [`Self::source_ids`] to handle both shapes.
297    pub fn source_id(&self) -> Result<i32> {
298        match self.source_ids.as_slice() {
299            [source_id] => Ok(*source_id),
300            source_ids => Err(invalid_data!(
301                "Partition field '{}' reads {} source columns and has no single source id",
302                self.name,
303                source_ids.len()
304            )),
305        }
306    }
307
308    /// The source column ids this field reads, in order. Never empty.
309    pub fn source_ids(&self) -> &[i32] {
310        &self.source_ids
311    }
312
313    /// The partition field id, when one was assigned.
314    pub fn field_id(&self) -> Option<i32> {
315        self.field_id
316    }
317
318    /// The partition name.
319    pub fn name(&self) -> &str {
320        &self.name
321    }
322
323    /// The transform applied to the source columns to produce a partition value.
324    pub fn transform(&self) -> Transform {
325        self.transform
326    }
327
328    /// Return this field with the given partition field id assigned.
329    pub fn with_field_id(self, field_id: i32) -> Self {
330        Self {
331            field_id: Some(field_id),
332            ..self
333        }
334    }
335
336    fn validate(&self) -> Result<()> {
337        if self.source_ids.is_empty() {
338            return Err(invalid_data!("Empty source-ids is not allowed"));
339        }
340        Ok(())
341    }
342}
343
344impl From<UnboundPartitionField> for Result<UnboundPartitionField> {
345    fn from(field: UnboundPartitionField) -> Self {
346        field.validate()?;
347        Ok(field)
348    }
349}
350
351mod _serde {
352    use serde::{Deserialize, Serialize};
353
354    use super::UnboundPartitionField;
355    use crate::Error;
356    use crate::error::invalid_data;
357    use crate::spec::Transform;
358
359    /// Per the spec a single-argument field carries `source-id` and a multi-argument field
360    /// carries `source-ids`. Either spelling is read, but not both at once; the one that
361    /// matches the field is written.
362    #[derive(Serialize, Deserialize)]
363    #[serde(rename_all = "kebab-case")]
364    pub(super) struct UnboundPartitionFieldSerde {
365        #[serde(default, skip_serializing_if = "Option::is_none")]
366        source_id: Option<i32>,
367        #[serde(default, skip_serializing_if = "Option::is_none")]
368        source_ids: Option<Vec<i32>>,
369        #[serde(default, skip_serializing_if = "Option::is_none")]
370        field_id: Option<i32>,
371        name: String,
372        transform: Transform,
373    }
374
375    impl TryFrom<UnboundPartitionFieldSerde> for UnboundPartitionField {
376        type Error = Error;
377
378        fn try_from(value: UnboundPartitionFieldSerde) -> Result<Self, Error> {
379            let source_ids = match (value.source_id, value.source_ids) {
380                (Some(source_id), None) => vec![source_id],
381                (None, Some(source_ids)) => source_ids,
382                (Some(_), Some(_)) => {
383                    return Err(invalid_data!(
384                        "source-id and source-ids are mutually exclusive"
385                    ));
386                }
387                (None, None) => {
388                    return Err(invalid_data!(
389                        "Either `source-id` or `source-ids` must be present"
390                    ));
391                }
392            };
393
394            UnboundPartitionField::builder()
395                .source_ids(source_ids)
396                .field_id_opt(value.field_id)
397                .name(value.name)
398                .transform(value.transform)
399                .build()
400        }
401    }
402
403    impl From<UnboundPartitionField> for UnboundPartitionFieldSerde {
404        fn from(value: UnboundPartitionField) -> Self {
405            let (source_id, source_ids) = match value.source_ids.as_slice() {
406                [source_id] => (Some(*source_id), None),
407                _ => (None, Some(value.source_ids)),
408            };
409            Self {
410                source_id,
411                source_ids,
412                field_id: value.field_id,
413                name: value.name,
414                transform: value.transform,
415            }
416        }
417    }
418}
419
420/// Unbound partition spec can be built without a schema and later bound to a schema.
421/// They are used to transport schema information as part of the REST specification.
422/// The main difference to [`PartitionSpec`] is that the field ids are optional.
423#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Default)]
424#[serde(rename_all = "kebab-case")]
425pub struct UnboundPartitionSpec {
426    /// Identifier for PartitionSpec
427    #[serde(skip_serializing_if = "Option::is_none")]
428    pub(crate) spec_id: Option<i32>,
429    /// Details of the partition spec
430    pub(crate) fields: Vec<UnboundPartitionField>,
431}
432
433impl UnboundPartitionSpec {
434    /// Create unbound partition spec builder
435    pub fn builder() -> UnboundPartitionSpecBuilder {
436        UnboundPartitionSpecBuilder::default()
437    }
438
439    /// Bind this unbound partition spec to a schema.
440    pub fn bind(self, schema: impl Into<SchemaRef>) -> Result<PartitionSpec> {
441        PartitionSpecBuilder::new_from_unbound(self, schema)?.build()
442    }
443
444    /// Spec id of the partition spec
445    pub fn spec_id(&self) -> Option<i32> {
446        self.spec_id
447    }
448
449    /// Fields of the partition spec
450    pub fn fields(&self) -> &[UnboundPartitionField] {
451        &self.fields
452    }
453
454    /// Change the spec id of the partition spec
455    pub fn with_spec_id(self, spec_id: i32) -> Self {
456        Self {
457            spec_id: Some(spec_id),
458            ..self
459        }
460    }
461}
462
463fn has_sequential_ids(field_ids: impl Iterator<Item = i32>) -> bool {
464    for (index, field_id) in field_ids.enumerate() {
465        let expected_id = (UNPARTITIONED_LAST_ASSIGNED_ID as i64)
466            .checked_add(1)
467            .and_then(|id| id.checked_add(index as i64))
468            .unwrap_or(i64::MAX);
469
470        if field_id as i64 != expected_id {
471            return false;
472        }
473    }
474
475    true
476}
477
478impl From<PartitionField> for UnboundPartitionField {
479    fn from(field: PartitionField) -> Self {
480        UnboundPartitionField {
481            source_ids: vec![field.source_id],
482            field_id: Some(field.field_id),
483            name: field.name,
484            transform: field.transform,
485        }
486    }
487}
488
489impl From<PartitionSpec> for UnboundPartitionSpec {
490    fn from(spec: PartitionSpec) -> Self {
491        UnboundPartitionSpec {
492            spec_id: Some(spec.spec_id),
493            fields: spec.fields.into_iter().map(Into::into).collect(),
494        }
495    }
496}
497
498/// Create a new UnboundPartitionSpec
499#[derive(Debug, Default)]
500pub struct UnboundPartitionSpecBuilder {
501    spec_id: Option<i32>,
502    fields: Vec<UnboundPartitionField>,
503}
504
505impl UnboundPartitionSpecBuilder {
506    /// Create a new partition spec builder with the given schema.
507    pub fn new() -> Self {
508        Self {
509            spec_id: None,
510            fields: vec![],
511        }
512    }
513
514    /// Set the spec id for the partition spec.
515    pub fn with_spec_id(mut self, spec_id: i32) -> Self {
516        self.spec_id = Some(spec_id);
517        self
518    }
519
520    /// Add a new partition field to the partition spec from an unbound partition field.
521    pub fn add_partition_field(mut self, field: UnboundPartitionField) -> Result<Self> {
522        self.check_name_set_and_unique(&field.name)?;
523        self.check_for_redundant_partitions(&field.source_ids, &field.transform)?;
524        if let Some(partition_field_id) = field.field_id {
525            self.check_partition_id_unique(partition_field_id)?;
526        }
527        self.fields.push(field);
528        Ok(self)
529    }
530
531    /// Add multiple partition fields to the partition spec.
532    pub fn add_partition_fields(
533        self,
534        fields: impl IntoIterator<Item = UnboundPartitionField>,
535    ) -> Result<Self> {
536        let mut builder = self;
537        for field in fields {
538            builder = builder.add_partition_field(field)?;
539        }
540        Ok(builder)
541    }
542
543    /// Build the unbound partition spec.
544    pub fn build(self) -> UnboundPartitionSpec {
545        UnboundPartitionSpec {
546            spec_id: self.spec_id,
547            fields: self.fields,
548        }
549    }
550}
551
552/// Create valid partition specs for a given schema.
553#[derive(Debug)]
554pub struct PartitionSpecBuilder {
555    spec_id: Option<i32>,
556    last_assigned_field_id: i32,
557    fields: Vec<UnboundPartitionField>,
558    schema: SchemaRef,
559}
560
561impl PartitionSpecBuilder {
562    /// Create a new partition spec builder with the given schema.
563    pub fn new(schema: impl Into<SchemaRef>) -> Self {
564        Self {
565            spec_id: None,
566            fields: vec![],
567            last_assigned_field_id: UNPARTITIONED_LAST_ASSIGNED_ID,
568            schema: schema.into(),
569        }
570    }
571
572    /// Create a new partition spec builder from an existing unbound partition spec.
573    pub fn new_from_unbound(
574        unbound: UnboundPartitionSpec,
575        schema: impl Into<SchemaRef>,
576    ) -> Result<Self> {
577        let mut builder =
578            Self::new(schema).with_spec_id(unbound.spec_id.unwrap_or(DEFAULT_PARTITION_SPEC_ID));
579
580        for field in unbound.fields {
581            builder = builder.add_unbound_field(field)?;
582        }
583        Ok(builder)
584    }
585
586    /// Set the last assigned field id for the partition spec.
587    ///
588    /// Set this field when a new partition spec is created for an existing TableMetaData.
589    /// As `field_id` must be unique in V2 metadata, this should be set to
590    /// the highest field id used previously.
591    pub fn with_last_assigned_field_id(mut self, last_assigned_field_id: i32) -> Self {
592        self.last_assigned_field_id = last_assigned_field_id;
593        self
594    }
595
596    /// Set the spec id for the partition spec.
597    pub fn with_spec_id(mut self, spec_id: i32) -> Self {
598        self.spec_id = Some(spec_id);
599        self
600    }
601
602    /// Add a new partition field to the partition spec.
603    pub fn add_partition_field(
604        self,
605        source_name: impl AsRef<str>,
606        target_name: impl Into<String>,
607        transform: Transform,
608    ) -> Result<Self> {
609        let source_id = self
610            .schema
611            .field_by_name(source_name.as_ref())
612            .ok_or_else(|| {
613                invalid_data!(
614                    "Cannot find source column with name: {} in schema",
615                    source_name.as_ref()
616                )
617            })?
618            .id;
619        let field = UnboundPartitionField {
620            source_ids: vec![source_id],
621            field_id: None,
622            name: target_name.into(),
623            transform,
624        };
625
626        self.add_unbound_field(field)
627    }
628
629    /// Add a new partition field to the partition spec.
630    ///
631    /// If partition field id is set, it is used as the field id.
632    /// Otherwise, a new `field_id` is assigned.
633    pub fn add_unbound_field(mut self, field: UnboundPartitionField) -> Result<Self> {
634        self.check_name_set_and_unique(&field.name)?;
635        self.check_for_redundant_partitions(&field.source_ids, &field.transform)?;
636        Self::check_name_does_not_collide_with_schema(&field, &self.schema)?;
637        Self::check_transform_compatibility(&field, &self.schema)?;
638        if let Some(partition_field_id) = field.field_id {
639            self.check_partition_id_unique(partition_field_id)?;
640        }
641
642        // Non-fallible from here
643        self.fields.push(field);
644        Ok(self)
645    }
646
647    /// Wrapper around `with_unbound_fields` to add multiple partition fields.
648    pub fn add_unbound_fields(
649        self,
650        fields: impl IntoIterator<Item = UnboundPartitionField>,
651    ) -> Result<Self> {
652        let mut builder = self;
653        for field in fields {
654            builder = builder.add_unbound_field(field)?;
655        }
656        Ok(builder)
657    }
658
659    /// Build a bound partition spec with the given schema.
660    pub fn build(self) -> Result<PartitionSpec> {
661        let fields = Self::set_field_ids(self.fields, self.last_assigned_field_id)?;
662        Ok(PartitionSpec {
663            spec_id: self.spec_id.unwrap_or(DEFAULT_PARTITION_SPEC_ID),
664            fields,
665        })
666    }
667
668    fn set_field_ids(
669        fields: Vec<UnboundPartitionField>,
670        last_assigned_field_id: i32,
671    ) -> Result<Vec<PartitionField>> {
672        let mut last_assigned_field_id = last_assigned_field_id;
673        // Already assigned partition ids. If we see one of these during iteration,
674        // we skip it.
675        let assigned_ids = fields
676            .iter()
677            .filter_map(|f| f.field_id)
678            .collect::<std::collections::HashSet<_>>();
679
680        fn _check_add_1(prev: i32) -> Result<i32> {
681            prev.checked_add(1)
682                .ok_or_else(|| invalid_data!("Cannot assign more partition ids. Overflow."))
683        }
684
685        let mut bound_fields = Vec::with_capacity(fields.len());
686        for field in fields.into_iter() {
687            let partition_field_id = if let Some(partition_field_id) = field.field_id {
688                last_assigned_field_id = std::cmp::max(last_assigned_field_id, partition_field_id);
689                partition_field_id
690            } else {
691                last_assigned_field_id = _check_add_1(last_assigned_field_id)?;
692                while assigned_ids.contains(&last_assigned_field_id) {
693                    last_assigned_field_id = _check_add_1(last_assigned_field_id)?;
694                }
695                last_assigned_field_id
696            };
697
698            bound_fields.push(PartitionField {
699                source_id: field.source_id()?,
700                field_id: partition_field_id,
701                name: field.name,
702                transform: field.transform,
703            })
704        }
705
706        Ok(bound_fields)
707    }
708
709    /// Ensure that the partition name is unique among columns in the schema.
710    /// Duplicate names are allowed if:
711    /// 1. The column is sourced from the column with the same name.
712    /// 2. AND the transformation is identity
713    fn check_name_does_not_collide_with_schema(
714        field: &UnboundPartitionField,
715        schema: &Schema,
716    ) -> Result<()> {
717        match schema.field_by_name(field.name.as_str()) {
718            Some(schema_collision) => {
719                if field.transform == Transform::Identity {
720                    if schema_collision.id == field.source_id()? {
721                        Ok(())
722                    } else {
723                        Err(invalid_data!(
724                            "Cannot create identity partition sourced from different field in schema. Field name '{}' has id `{}` in schema but partition source id is `{}`",
725                            field.name,
726                            schema_collision.id,
727                            field.source_id()?
728                        ))
729                    }
730                } else {
731                    Err(invalid_data!(
732                        "Cannot create partition with name: '{}' that conflicts with schema field and is not an identity transform.",
733                        field.name
734                    ))
735                }
736            }
737            None => Ok(()),
738        }
739    }
740
741    /// Ensure that the transformation of the field is compatible with type of the field
742    /// in the schema. Implicitly also checks if the source field exists in the schema.
743    fn check_transform_compatibility(field: &UnboundPartitionField, schema: &Schema) -> Result<()> {
744        // A multi-argument field is rejected here on purpose: a bound `PartitionField` has
745        // no place for more than one source id yet.
746        let source_id = field.source_id()?;
747        let schema_field = schema.field_by_id(source_id).ok_or_else(|| {
748            invalid_data!("Cannot find partition source field with id `{source_id}` in schema")
749        })?;
750
751        if field.transform != Transform::Void {
752            if !schema_field.field_type.is_primitive() {
753                return Err(invalid_data!(
754                    "Cannot partition by non-primitive source field: '{}'.",
755                    schema_field.field_type
756                ));
757            }
758
759            if field
760                .transform
761                .result_type(&schema_field.field_type)
762                .is_err()
763            {
764                return Err(invalid_data!(
765                    "Invalid source type: '{}' for transform: '{}'.",
766                    schema_field.field_type,
767                    field.transform.dedup_name()
768                ));
769            }
770        }
771
772        Ok(())
773    }
774}
775
776/// Contains checks that are common to both PartitionSpecBuilder and UnboundPartitionSpecBuilder
777trait CorePartitionSpecValidator {
778    /// Ensure that the partition name is unique among the partition fields and is not empty.
779    fn check_name_set_and_unique(&self, name: &str) -> Result<()> {
780        if name.is_empty() {
781            return Err(invalid_data!("Cannot use empty partition name"));
782        }
783
784        if self.fields().iter().any(|f| f.name == name) {
785            return Err(invalid_data!(
786                "Cannot use partition name more than once: {name}"
787            ));
788        }
789        Ok(())
790    }
791
792    /// For a single source-column transformations must be unique.
793    fn check_for_redundant_partitions(
794        &self,
795        source_ids: &[i32],
796        transform: &Transform,
797    ) -> Result<()> {
798        let collision = self.fields().iter().find(|f| {
799            f.source_ids == source_ids && f.transform.dedup_name() == transform.dedup_name()
800        });
801
802        if let Some(collision) = collision {
803            Err(invalid_data!(
804                "Cannot add redundant partition with source ids `{:?}` and transform `{}`. A partition with the same source ids and transform already exists with name `{}`",
805                source_ids,
806                transform.dedup_name(),
807                collision.name
808            ))
809        } else {
810            Ok(())
811        }
812    }
813
814    /// Check field / partition_id unique within the partition spec if set
815    fn check_partition_id_unique(&self, field_id: i32) -> Result<()> {
816        if self.fields().iter().any(|f| f.field_id == Some(field_id)) {
817            return Err(invalid_data!(
818                "Cannot use field id more than once in one PartitionSpec: {field_id}"
819            ));
820        }
821
822        Ok(())
823    }
824
825    fn fields(&self) -> &Vec<UnboundPartitionField>;
826}
827
828impl CorePartitionSpecValidator for PartitionSpecBuilder {
829    fn fields(&self) -> &Vec<UnboundPartitionField> {
830        &self.fields
831    }
832}
833
834impl CorePartitionSpecValidator for UnboundPartitionSpecBuilder {
835    fn fields(&self) -> &Vec<UnboundPartitionField> {
836        &self.fields
837    }
838}
839
840#[cfg(test)]
841mod tests {
842    use super::*;
843    use crate::ErrorKind;
844    use crate::spec::{Literal, PrimitiveType, Type};
845
846    #[test]
847    fn test_partition_spec() {
848        let spec = r#"
849        {
850        "spec-id": 1,
851        "fields": [ {
852            "source-id": 4,
853            "field-id": 1000,
854            "name": "ts_day",
855            "transform": "day"
856            }, {
857            "source-id": 1,
858            "field-id": 1001,
859            "name": "id_bucket",
860            "transform": "bucket[16]"
861            }, {
862            "source-id": 2,
863            "field-id": 1002,
864            "name": "id_truncate",
865            "transform": "truncate[4]"
866            } ]
867        }
868        "#;
869
870        let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
871        assert_eq!(4, partition_spec.fields[0].source_id);
872        assert_eq!(1000, partition_spec.fields[0].field_id);
873        assert_eq!("ts_day", partition_spec.fields[0].name);
874        assert_eq!(Transform::Day, partition_spec.fields[0].transform);
875
876        assert_eq!(1, partition_spec.fields[1].source_id);
877        assert_eq!(1001, partition_spec.fields[1].field_id);
878        assert_eq!("id_bucket", partition_spec.fields[1].name);
879        assert_eq!(Transform::Bucket(16), partition_spec.fields[1].transform);
880
881        assert_eq!(2, partition_spec.fields[2].source_id);
882        assert_eq!(1002, partition_spec.fields[2].field_id);
883        assert_eq!("id_truncate", partition_spec.fields[2].name);
884        assert_eq!(Transform::Truncate(4), partition_spec.fields[2].transform);
885    }
886
887    #[test]
888    fn test_is_unpartitioned() {
889        let schema = Schema::builder()
890            .with_fields(vec![
891                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
892                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
893            ])
894            .build()
895            .unwrap();
896        let partition_spec = PartitionSpec::builder(schema.clone())
897            .with_spec_id(1)
898            .build()
899            .unwrap();
900        assert!(
901            partition_spec.is_unpartitioned(),
902            "Empty partition spec should be unpartitioned"
903        );
904
905        let partition_spec = PartitionSpec::builder(schema.clone())
906            .add_unbound_fields(vec![
907                UnboundPartitionField::builder()
908                    .source_ids(vec![1])
909                    .name("id".to_string())
910                    .transform(Transform::Identity)
911                    .build()
912                    .unwrap(),
913                UnboundPartitionField::builder()
914                    .source_ids(vec![2])
915                    .name("name_string".to_string())
916                    .transform(Transform::Void)
917                    .build()
918                    .unwrap(),
919            ])
920            .unwrap()
921            .with_spec_id(1)
922            .build()
923            .unwrap();
924        assert!(
925            !partition_spec.is_unpartitioned(),
926            "Partition spec with one non void transform should not be unpartitioned"
927        );
928
929        let partition_spec = PartitionSpec::builder(schema.clone())
930            .with_spec_id(1)
931            .add_unbound_fields(vec![
932                UnboundPartitionField::builder()
933                    .source_ids(vec![1])
934                    .name("id_void".to_string())
935                    .transform(Transform::Void)
936                    .build()
937                    .unwrap(),
938                UnboundPartitionField::builder()
939                    .source_ids(vec![2])
940                    .name("name_void".to_string())
941                    .transform(Transform::Void)
942                    .build()
943                    .unwrap(),
944            ])
945            .unwrap()
946            .build()
947            .unwrap();
948        assert!(
949            partition_spec.is_unpartitioned(),
950            "Partition spec with all void field should be unpartitioned"
951        );
952    }
953
954    #[test]
955    fn test_unbound_partition_spec() {
956        let spec = r#"
957		{
958		"spec-id": 1,
959		"fields": [ {
960			"source-id": 4,
961			"field-id": 1000,
962			"name": "ts_day",
963			"transform": "day"
964			}, {
965			"source-id": 1,
966			"field-id": 1001,
967			"name": "id_bucket",
968			"transform": "bucket[16]"
969			}, {
970			"source-id": 2,
971			"field-id": 1002,
972			"name": "id_truncate",
973			"transform": "truncate[4]"
974			} ]
975		}
976		"#;
977
978        let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap();
979        assert_eq!(Some(1), partition_spec.spec_id);
980
981        assert_eq!(4, partition_spec.fields[0].source_id().unwrap());
982        assert_eq!(Some(1000), partition_spec.fields[0].field_id);
983        assert_eq!("ts_day", partition_spec.fields[0].name);
984        assert_eq!(Transform::Day, partition_spec.fields[0].transform);
985
986        assert_eq!(1, partition_spec.fields[1].source_id().unwrap());
987        assert_eq!(Some(1001), partition_spec.fields[1].field_id);
988        assert_eq!("id_bucket", partition_spec.fields[1].name);
989        assert_eq!(Transform::Bucket(16), partition_spec.fields[1].transform);
990
991        assert_eq!(2, partition_spec.fields[2].source_id().unwrap());
992        assert_eq!(Some(1002), partition_spec.fields[2].field_id);
993        assert_eq!("id_truncate", partition_spec.fields[2].name);
994        assert_eq!(Transform::Truncate(4), partition_spec.fields[2].transform);
995
996        let spec = r#"
997		{
998		"fields": [ {
999			"source-id": 4,
1000			"name": "ts_day",
1001			"transform": "day"
1002			} ]
1003		}
1004		"#;
1005        let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap();
1006        assert_eq!(None, partition_spec.spec_id);
1007
1008        assert_eq!(4, partition_spec.fields[0].source_id().unwrap());
1009        assert_eq!(None, partition_spec.fields[0].field_id);
1010        assert_eq!("ts_day", partition_spec.fields[0].name);
1011        assert_eq!(Transform::Day, partition_spec.fields[0].transform);
1012    }
1013
1014    #[test]
1015    fn test_unbound_partition_spec_serialization_skips_none_fields() {
1016        let spec = UnboundPartitionSpec::builder()
1017            .add_partition_field(
1018                UnboundPartitionField::builder()
1019                    .source_ids(vec![4])
1020                    .name("ts_day")
1021                    .transform(Transform::Day)
1022                    .build()
1023                    .unwrap(),
1024            )
1025            .unwrap()
1026            .build();
1027
1028        let value = serde_json::to_value(&spec).unwrap();
1029        let object = value.as_object().unwrap();
1030        assert!(!object.contains_key("spec-id"));
1031        let field = object["fields"][0].as_object().unwrap();
1032        assert!(!field.contains_key("field-id"));
1033
1034        let value = serde_json::to_value(spec.with_spec_id(1)).unwrap();
1035        let object = value.as_object().unwrap();
1036        assert_eq!(Some(&serde_json::json!(1)), object.get("spec-id"));
1037
1038        // Explicit nulls must still deserialize to `None` for backwards
1039        // compatibility.
1040        let spec: UnboundPartitionSpec = serde_json::from_str(
1041            r#"{
1042                "spec-id": null,
1043                "fields": [
1044                    {"source-id": 4, "name": "ts_day", "transform": "day", "field-id": null}
1045                ]
1046            }"#,
1047        )
1048        .unwrap();
1049        assert_eq!(None, spec.spec_id);
1050        assert_eq!(None, spec.fields[0].field_id);
1051    }
1052
1053    #[test]
1054    fn test_new_unpartition() {
1055        let schema = Schema::builder()
1056            .with_fields(vec![
1057                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1058                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1059            ])
1060            .build()
1061            .unwrap();
1062        let partition_spec = PartitionSpec::builder(schema.clone())
1063            .with_spec_id(0)
1064            .build()
1065            .unwrap();
1066        let partition_type = partition_spec.partition_type(&schema).unwrap();
1067        assert_eq!(0, partition_type.fields().len());
1068
1069        let unpartition_spec = PartitionSpec::unpartition_spec();
1070        assert_eq!(partition_spec, unpartition_spec);
1071    }
1072
1073    #[test]
1074    fn test_partition_type() {
1075        let spec = r#"
1076            {
1077            "spec-id": 1,
1078            "fields": [ {
1079                "source-id": 4,
1080                "field-id": 1000,
1081                "name": "ts_day",
1082                "transform": "day"
1083                }, {
1084                "source-id": 1,
1085                "field-id": 1001,
1086                "name": "id_bucket",
1087                "transform": "bucket[16]"
1088                }, {
1089                "source-id": 2,
1090                "field-id": 1002,
1091                "name": "id_truncate",
1092                "transform": "truncate[4]"
1093                } ]
1094            }
1095            "#;
1096
1097        let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1098        let schema = Schema::builder()
1099            .with_fields(vec![
1100                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1101                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1102                NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1103                NestedField::required(4, "ts_day", Type::Primitive(PrimitiveType::Timestamp))
1104                    .into(),
1105                NestedField::required(5, "id_bucket", Type::Primitive(PrimitiveType::Int)).into(),
1106                NestedField::required(6, "id_truncate", Type::Primitive(PrimitiveType::Int)).into(),
1107            ])
1108            .build()
1109            .unwrap();
1110
1111        let partition_type = partition_spec.partition_type(&schema).unwrap();
1112        assert_eq!(3, partition_type.fields().len());
1113        assert_eq!(
1114            *partition_type.fields()[0],
1115            NestedField::optional(
1116                partition_spec.fields[0].field_id,
1117                &partition_spec.fields[0].name,
1118                Type::Primitive(PrimitiveType::Date)
1119            )
1120        );
1121        assert_eq!(
1122            *partition_type.fields()[1],
1123            NestedField::optional(
1124                partition_spec.fields[1].field_id,
1125                &partition_spec.fields[1].name,
1126                Type::Primitive(PrimitiveType::Int)
1127            )
1128        );
1129        assert_eq!(
1130            *partition_type.fields()[2],
1131            NestedField::optional(
1132                partition_spec.fields[2].field_id,
1133                &partition_spec.fields[2].name,
1134                Type::Primitive(PrimitiveType::String)
1135            )
1136        );
1137    }
1138
1139    #[test]
1140    fn test_partition_empty() {
1141        let spec = r#"
1142            {
1143            "spec-id": 1,
1144            "fields": []
1145            }
1146            "#;
1147
1148        let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1149        let schema = Schema::builder()
1150            .with_fields(vec![
1151                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1152                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1153                NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1154                NestedField::required(4, "ts_day", Type::Primitive(PrimitiveType::Timestamp))
1155                    .into(),
1156                NestedField::required(5, "id_bucket", Type::Primitive(PrimitiveType::Int)).into(),
1157                NestedField::required(6, "id_truncate", Type::Primitive(PrimitiveType::Int)).into(),
1158            ])
1159            .build()
1160            .unwrap();
1161
1162        let partition_type = partition_spec.partition_type(&schema).unwrap();
1163        assert_eq!(0, partition_type.fields().len());
1164    }
1165
1166    #[test]
1167    fn test_partition_type_with_dropped_source_column() {
1168        let spec = r#"
1169        {
1170        "spec-id": 1,
1171        "fields": [ {
1172            "source-id": 4,
1173            "field-id": 1000,
1174            "name": "ts_day",
1175            "transform": "day"
1176            }, {
1177            "source-id": 1,
1178            "field-id": 1001,
1179            "name": "id_bucket",
1180            "transform": "bucket[16]"
1181            }, {
1182            "source-id": 2,
1183            "field-id": 1002,
1184            "name": "id_truncate",
1185            "transform": "truncate[4]"
1186            } ]
1187        }
1188        "#;
1189
1190        let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap();
1191        let schema = Schema::builder()
1192            .with_fields(vec![
1193                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1194                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1195            ])
1196            .build()
1197            .unwrap();
1198
1199        assert_eq!(
1200            partition_spec.partition_type(&schema).unwrap(),
1201            StructType::new(vec![
1202                NestedField::optional(1000, "ts_day", PrimitiveType::Date.into()).into(),
1203                NestedField::optional(1001, "id_bucket", PrimitiveType::Int.into()).into(),
1204                NestedField::optional(1002, "id_truncate", PrimitiveType::String.into()).into(),
1205            ])
1206        );
1207
1208        // Dropping all sources must retain every partition field in its original position.
1209        let empty_schema = Schema::builder().build().unwrap();
1210        assert_eq!(
1211            partition_spec.partition_type(&empty_schema).unwrap(),
1212            StructType::new(vec![
1213                NestedField::optional(1000, "ts_day", PrimitiveType::Date.into()).into(),
1214                NestedField::optional(1001, "id_bucket", PrimitiveType::Int.into()).into(),
1215                NestedField::optional(1002, "id_truncate", PrimitiveType::Unknown.into()).into(),
1216            ])
1217        );
1218
1219        // Missing sources must not suppress validation of transforms on remaining sources.
1220        let incompatible_schema = Schema::builder()
1221            .with_fields(vec![
1222                NestedField::required(1, "id", PrimitiveType::Boolean.into()).into(),
1223            ])
1224            .build()
1225            .unwrap();
1226        let err = partition_spec
1227            .partition_type(&incompatible_schema)
1228            .unwrap_err();
1229        assert_eq!(err.kind(), ErrorKind::DataInvalid);
1230        assert_eq!(
1231            err.message(),
1232            "boolean is not a valid input type of bucket transform"
1233        );
1234    }
1235
1236    #[test]
1237    fn test_partition_type_without_source_type() {
1238        let schema = Schema::builder().build().unwrap();
1239        for (transform, expected_type) in [
1240            (Transform::Identity, PrimitiveType::Unknown),
1241            (Transform::Truncate(4), PrimitiveType::Unknown),
1242            (Transform::Void, PrimitiveType::Unknown),
1243            (Transform::Bucket(16), PrimitiveType::Int),
1244            (Transform::Year, PrimitiveType::Int),
1245            (Transform::Month, PrimitiveType::Int),
1246            (Transform::Day, PrimitiveType::Date),
1247            (Transform::Hour, PrimitiveType::Int),
1248            (Transform::Unknown, PrimitiveType::String),
1249        ] {
1250            let spec = PartitionSpec {
1251                spec_id: 0,
1252                fields: vec![PartitionField {
1253                    source_id: 1,
1254                    field_id: 1000,
1255                    name: "partition".to_string(),
1256                    transform,
1257                }],
1258            };
1259            assert_eq!(
1260                spec.partition_type(&schema).unwrap(),
1261                StructType::new(vec![
1262                    NestedField::optional(1000, "partition", expected_type.into()).into(),
1263                ])
1264            );
1265        }
1266    }
1267
1268    #[test]
1269    fn test_builder_disallow_duplicate_names() {
1270        UnboundPartitionSpec::builder()
1271            .add_partition_field(
1272                UnboundPartitionField::builder()
1273                    .source_ids(vec![1])
1274                    .name("ts_day")
1275                    .transform(Transform::Day)
1276                    .build()
1277                    .unwrap(),
1278            )
1279            .unwrap()
1280            .add_partition_field(
1281                UnboundPartitionField::builder()
1282                    .source_ids(vec![2])
1283                    .name("ts_day")
1284                    .transform(Transform::Day)
1285                    .build()
1286                    .unwrap(),
1287            )
1288            .unwrap_err();
1289    }
1290
1291    #[test]
1292    fn test_builder_disallow_duplicate_field_ids() {
1293        let schema = Schema::builder()
1294            .with_fields(vec![
1295                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1296                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1297            ])
1298            .build()
1299            .unwrap();
1300        PartitionSpec::builder(schema.clone())
1301            .add_unbound_field(UnboundPartitionField {
1302                source_ids: vec![1],
1303                field_id: Some(1000),
1304                name: "id".to_string(),
1305                transform: Transform::Identity,
1306            })
1307            .unwrap()
1308            .add_unbound_field(UnboundPartitionField {
1309                source_ids: vec![2],
1310                field_id: Some(1000),
1311                name: "id_bucket".to_string(),
1312                transform: Transform::Bucket(16),
1313            })
1314            .unwrap_err();
1315    }
1316
1317    #[test]
1318    fn test_builder_auto_assign_field_ids() {
1319        let schema = Schema::builder()
1320            .with_fields(vec![
1321                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1322                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1323                NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1324            ])
1325            .build()
1326            .unwrap();
1327        let spec = PartitionSpec::builder(schema.clone())
1328            .with_spec_id(1)
1329            .add_unbound_field(UnboundPartitionField {
1330                source_ids: vec![1],
1331                name: "id".to_string(),
1332                transform: Transform::Identity,
1333                field_id: Some(1012),
1334            })
1335            .unwrap()
1336            .add_unbound_field(UnboundPartitionField {
1337                source_ids: vec![2],
1338                name: "name_void".to_string(),
1339                transform: Transform::Void,
1340                field_id: None,
1341            })
1342            .unwrap()
1343            // Should keep its ID even if its lower
1344            .add_unbound_field(UnboundPartitionField {
1345                source_ids: vec![3],
1346                name: "year".to_string(),
1347                transform: Transform::Year,
1348                field_id: Some(1),
1349            })
1350            .unwrap()
1351            .build()
1352            .unwrap();
1353
1354        assert_eq!(1012, spec.fields[0].field_id);
1355        assert_eq!(1013, spec.fields[1].field_id);
1356        assert_eq!(1, spec.fields[2].field_id);
1357    }
1358
1359    #[test]
1360    fn test_builder_valid_schema() {
1361        let schema = Schema::builder()
1362            .with_fields(vec![
1363                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1364                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1365            ])
1366            .build()
1367            .unwrap();
1368
1369        PartitionSpec::builder(schema.clone())
1370            .with_spec_id(1)
1371            .build()
1372            .unwrap();
1373
1374        let spec = PartitionSpec::builder(schema.clone())
1375            .with_spec_id(1)
1376            .add_partition_field("id", "id_bucket[16]", Transform::Bucket(16))
1377            .unwrap()
1378            .build()
1379            .unwrap();
1380
1381        assert_eq!(spec, PartitionSpec {
1382            spec_id: 1,
1383            fields: vec![PartitionField {
1384                source_id: 1,
1385                field_id: 1000,
1386                name: "id_bucket[16]".to_string(),
1387                transform: Transform::Bucket(16),
1388            }],
1389        });
1390        assert_eq!(
1391            spec.partition_type(&schema).unwrap(),
1392            StructType::new(vec![
1393                NestedField::optional(1000, "id_bucket[16]", Type::Primitive(PrimitiveType::Int))
1394                    .into()
1395            ])
1396        )
1397    }
1398
1399    #[test]
1400    fn test_collision_with_schema_name() {
1401        let schema = Schema::builder()
1402            .with_fields(vec![
1403                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1404            ])
1405            .build()
1406            .unwrap();
1407
1408        PartitionSpec::builder(schema.clone())
1409            .with_spec_id(1)
1410            .build()
1411            .unwrap();
1412
1413        let err = PartitionSpec::builder(schema)
1414            .with_spec_id(1)
1415            .add_unbound_field(UnboundPartitionField {
1416                source_ids: vec![1],
1417                field_id: None,
1418                name: "id".to_string(),
1419                transform: Transform::Bucket(16),
1420            })
1421            .unwrap_err();
1422        assert!(err.message().contains("conflicts with schema"))
1423    }
1424
1425    #[test]
1426    fn test_builder_collision_is_ok_for_identity_transforms() {
1427        let schema = Schema::builder()
1428            .with_fields(vec![
1429                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1430                NestedField::required(2, "number", Type::Primitive(PrimitiveType::Int)).into(),
1431            ])
1432            .build()
1433            .unwrap();
1434
1435        PartitionSpec::builder(schema.clone())
1436            .with_spec_id(1)
1437            .build()
1438            .unwrap();
1439
1440        PartitionSpec::builder(schema.clone())
1441            .with_spec_id(1)
1442            .add_unbound_field(UnboundPartitionField {
1443                source_ids: vec![1],
1444                field_id: None,
1445                name: "id".to_string(),
1446                transform: Transform::Identity,
1447            })
1448            .unwrap()
1449            .build()
1450            .unwrap();
1451
1452        // Not OK for different source id
1453        PartitionSpec::builder(schema)
1454            .with_spec_id(1)
1455            .add_unbound_field(UnboundPartitionField {
1456                source_ids: vec![2],
1457                field_id: None,
1458                name: "id".to_string(),
1459                transform: Transform::Identity,
1460            })
1461            .unwrap_err();
1462    }
1463
1464    #[test]
1465    fn test_builder_all_source_ids_must_exist() {
1466        let schema = Schema::builder()
1467            .with_fields(vec![
1468                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1469                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1470                NestedField::required(3, "ts", Type::Primitive(PrimitiveType::Timestamp)).into(),
1471            ])
1472            .build()
1473            .unwrap();
1474
1475        // Valid
1476        PartitionSpec::builder(schema.clone())
1477            .with_spec_id(1)
1478            .add_unbound_fields(vec![
1479                UnboundPartitionField {
1480                    source_ids: vec![1],
1481                    field_id: None,
1482                    name: "id_bucket".to_string(),
1483                    transform: Transform::Bucket(16),
1484                },
1485                UnboundPartitionField {
1486                    source_ids: vec![2],
1487                    field_id: None,
1488                    name: "name".to_string(),
1489                    transform: Transform::Identity,
1490                },
1491            ])
1492            .unwrap()
1493            .build()
1494            .unwrap();
1495
1496        // Invalid
1497        PartitionSpec::builder(schema)
1498            .with_spec_id(1)
1499            .add_unbound_fields(vec![
1500                UnboundPartitionField {
1501                    source_ids: vec![1],
1502                    field_id: None,
1503                    name: "id_bucket".to_string(),
1504                    transform: Transform::Bucket(16),
1505                },
1506                UnboundPartitionField {
1507                    source_ids: vec![4],
1508                    field_id: None,
1509                    name: "name".to_string(),
1510                    transform: Transform::Identity,
1511                },
1512            ])
1513            .unwrap_err();
1514    }
1515
1516    #[test]
1517    fn test_builder_disallows_variant_source() {
1518        let schema = Schema::builder()
1519            .with_fields(vec![
1520                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1521                NestedField::optional(2, "v", Type::Variant(crate::spec::VariantType)).into(),
1522            ])
1523            .build()
1524            .unwrap();
1525
1526        let err = PartitionSpec::builder(schema)
1527            .with_spec_id(1)
1528            .add_unbound_fields(vec![UnboundPartitionField {
1529                source_ids: vec![2],
1530                field_id: None,
1531                name: "v_part".to_string(),
1532                transform: Transform::Identity,
1533            }])
1534            .expect_err("variant must not be allowed as a partition source");
1535
1536        assert_eq!(
1537            err.message(),
1538            "Cannot partition by non-primitive source field: 'variant'."
1539        );
1540    }
1541
1542    #[test]
1543    fn test_builder_disallows_redundant() {
1544        let err = UnboundPartitionSpec::builder()
1545            .with_spec_id(1)
1546            .add_partition_field(
1547                UnboundPartitionField::builder()
1548                    .source_ids(vec![1])
1549                    .name("id_bucket[16]")
1550                    .transform(Transform::Bucket(16))
1551                    .build()
1552                    .unwrap(),
1553            )
1554            .unwrap()
1555            .add_partition_field(
1556                UnboundPartitionField::builder()
1557                    .source_ids(vec![1])
1558                    .name("id_bucket_with_other_name")
1559                    .transform(Transform::Bucket(16))
1560                    .build()
1561                    .unwrap(),
1562            )
1563            .unwrap_err();
1564        assert!(err.message().contains("redundant partition"));
1565    }
1566
1567    #[test]
1568    fn test_builder_incompatible_transforms_disallowed() {
1569        let schema = Schema::builder()
1570            .with_fields(vec![
1571                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1572            ])
1573            .build()
1574            .unwrap();
1575
1576        PartitionSpec::builder(schema)
1577            .with_spec_id(1)
1578            .add_unbound_field(UnboundPartitionField {
1579                source_ids: vec![1],
1580                field_id: None,
1581                name: "id_year".to_string(),
1582                transform: Transform::Year,
1583            })
1584            .unwrap_err();
1585    }
1586
1587    #[test]
1588    fn test_build_unbound_specs_without_partition_id() {
1589        let spec = UnboundPartitionSpec::builder()
1590            .with_spec_id(1)
1591            .add_partition_fields(vec![UnboundPartitionField {
1592                source_ids: vec![1],
1593                field_id: None,
1594                name: "id_bucket[16]".to_string(),
1595                transform: Transform::Bucket(16),
1596            }])
1597            .unwrap()
1598            .build();
1599
1600        assert_eq!(spec, UnboundPartitionSpec {
1601            spec_id: Some(1),
1602            fields: vec![UnboundPartitionField {
1603                source_ids: vec![1],
1604                field_id: None,
1605                name: "id_bucket[16]".to_string(),
1606                transform: Transform::Bucket(16),
1607            }]
1608        });
1609    }
1610
1611    #[test]
1612    fn test_is_compatible_with() {
1613        let schema = Schema::builder()
1614            .with_fields(vec![
1615                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1616                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1617            ])
1618            .build()
1619            .unwrap();
1620
1621        let partition_spec_1 = PartitionSpec::builder(schema.clone())
1622            .with_spec_id(1)
1623            .add_unbound_field(UnboundPartitionField {
1624                source_ids: vec![1],
1625                field_id: None,
1626                name: "id_bucket".to_string(),
1627                transform: Transform::Bucket(16),
1628            })
1629            .unwrap()
1630            .build()
1631            .unwrap();
1632
1633        let partition_spec_2 = PartitionSpec::builder(schema)
1634            .with_spec_id(1)
1635            .add_unbound_field(UnboundPartitionField {
1636                source_ids: vec![1],
1637                field_id: None,
1638                name: "id_bucket".to_string(),
1639                transform: Transform::Bucket(16),
1640            })
1641            .unwrap()
1642            .build()
1643            .unwrap();
1644
1645        assert!(partition_spec_1.is_compatible_with(&partition_spec_2));
1646    }
1647
1648    #[test]
1649    fn test_not_compatible_with_transform_different() {
1650        let schema = Schema::builder()
1651            .with_fields(vec![
1652                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1653            ])
1654            .build()
1655            .unwrap();
1656
1657        let partition_spec_1 = PartitionSpec::builder(schema.clone())
1658            .with_spec_id(1)
1659            .add_unbound_field(UnboundPartitionField {
1660                source_ids: vec![1],
1661                field_id: None,
1662                name: "id_bucket".to_string(),
1663                transform: Transform::Bucket(16),
1664            })
1665            .unwrap()
1666            .build()
1667            .unwrap();
1668
1669        let partition_spec_2 = PartitionSpec::builder(schema)
1670            .with_spec_id(1)
1671            .add_unbound_field(UnboundPartitionField {
1672                source_ids: vec![1],
1673                field_id: None,
1674                name: "id_bucket".to_string(),
1675                transform: Transform::Bucket(32),
1676            })
1677            .unwrap()
1678            .build()
1679            .unwrap();
1680
1681        assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1682    }
1683
1684    #[test]
1685    fn test_not_compatible_with_source_id_different() {
1686        let schema = Schema::builder()
1687            .with_fields(vec![
1688                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1689                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1690            ])
1691            .build()
1692            .unwrap();
1693
1694        let partition_spec_1 = PartitionSpec::builder(schema.clone())
1695            .with_spec_id(1)
1696            .add_unbound_field(UnboundPartitionField {
1697                source_ids: vec![1],
1698                field_id: None,
1699                name: "id_bucket".to_string(),
1700                transform: Transform::Bucket(16),
1701            })
1702            .unwrap()
1703            .build()
1704            .unwrap();
1705
1706        let partition_spec_2 = PartitionSpec::builder(schema)
1707            .with_spec_id(1)
1708            .add_unbound_field(UnboundPartitionField {
1709                source_ids: vec![2],
1710                field_id: None,
1711                name: "id_bucket".to_string(),
1712                transform: Transform::Bucket(16),
1713            })
1714            .unwrap()
1715            .build()
1716            .unwrap();
1717
1718        assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1719    }
1720
1721    #[test]
1722    fn test_not_compatible_with_order_different() {
1723        let schema = Schema::builder()
1724            .with_fields(vec![
1725                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1726                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1727            ])
1728            .build()
1729            .unwrap();
1730
1731        let partition_spec_1 = PartitionSpec::builder(schema.clone())
1732            .with_spec_id(1)
1733            .add_unbound_field(UnboundPartitionField {
1734                source_ids: vec![1],
1735                field_id: None,
1736                name: "id_bucket".to_string(),
1737                transform: Transform::Bucket(16),
1738            })
1739            .unwrap()
1740            .add_unbound_field(UnboundPartitionField {
1741                source_ids: vec![2],
1742                field_id: None,
1743                name: "name".to_string(),
1744                transform: Transform::Identity,
1745            })
1746            .unwrap()
1747            .build()
1748            .unwrap();
1749
1750        let partition_spec_2 = PartitionSpec::builder(schema)
1751            .with_spec_id(1)
1752            .add_unbound_field(UnboundPartitionField {
1753                source_ids: vec![2],
1754                field_id: None,
1755                name: "name".to_string(),
1756                transform: Transform::Identity,
1757            })
1758            .unwrap()
1759            .add_unbound_field(UnboundPartitionField {
1760                source_ids: vec![1],
1761                field_id: None,
1762                name: "id_bucket".to_string(),
1763                transform: Transform::Bucket(16),
1764            })
1765            .unwrap()
1766            .build()
1767            .unwrap();
1768
1769        assert!(!partition_spec_1.is_compatible_with(&partition_spec_2));
1770    }
1771
1772    #[test]
1773    fn test_highest_field_id_unpartitioned() {
1774        let spec = PartitionSpec::builder(Schema::builder().with_fields(vec![]).build().unwrap())
1775            .with_spec_id(1)
1776            .build()
1777            .unwrap();
1778
1779        assert!(spec.highest_field_id().is_none());
1780    }
1781
1782    #[test]
1783    fn test_highest_field_id() {
1784        let schema = Schema::builder()
1785            .with_fields(vec![
1786                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1787                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1788            ])
1789            .build()
1790            .unwrap();
1791
1792        let spec = PartitionSpec::builder(schema)
1793            .with_spec_id(1)
1794            .add_unbound_field(UnboundPartitionField {
1795                source_ids: vec![1],
1796                field_id: Some(1001),
1797                name: "id".to_string(),
1798                transform: Transform::Identity,
1799            })
1800            .unwrap()
1801            .add_unbound_field(UnboundPartitionField {
1802                source_ids: vec![2],
1803                field_id: Some(1000),
1804                name: "name".to_string(),
1805                transform: Transform::Identity,
1806            })
1807            .unwrap()
1808            .build()
1809            .unwrap();
1810
1811        assert_eq!(Some(1001), spec.highest_field_id());
1812    }
1813
1814    #[test]
1815    fn test_has_sequential_ids() {
1816        let schema = Schema::builder()
1817            .with_fields(vec![
1818                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1819                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1820            ])
1821            .build()
1822            .unwrap();
1823
1824        let spec = PartitionSpec::builder(schema)
1825            .with_spec_id(1)
1826            .add_unbound_field(UnboundPartitionField {
1827                source_ids: vec![1],
1828                field_id: Some(1000),
1829                name: "id".to_string(),
1830                transform: Transform::Identity,
1831            })
1832            .unwrap()
1833            .add_unbound_field(UnboundPartitionField {
1834                source_ids: vec![2],
1835                field_id: Some(1001),
1836                name: "name".to_string(),
1837                transform: Transform::Identity,
1838            })
1839            .unwrap()
1840            .build()
1841            .unwrap();
1842
1843        assert_eq!(1000, spec.fields[0].field_id);
1844        assert_eq!(1001, spec.fields[1].field_id);
1845        assert!(spec.has_sequential_ids());
1846    }
1847
1848    #[test]
1849    fn test_sequential_ids_must_start_at_1000() {
1850        let schema = Schema::builder()
1851            .with_fields(vec![
1852                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1853                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1854            ])
1855            .build()
1856            .unwrap();
1857
1858        let spec = PartitionSpec::builder(schema)
1859            .with_spec_id(1)
1860            .add_unbound_field(UnboundPartitionField {
1861                source_ids: vec![1],
1862                field_id: Some(999),
1863                name: "id".to_string(),
1864                transform: Transform::Identity,
1865            })
1866            .unwrap()
1867            .add_unbound_field(UnboundPartitionField {
1868                source_ids: vec![2],
1869                field_id: Some(1000),
1870                name: "name".to_string(),
1871                transform: Transform::Identity,
1872            })
1873            .unwrap()
1874            .build()
1875            .unwrap();
1876
1877        assert_eq!(999, spec.fields[0].field_id);
1878        assert_eq!(1000, spec.fields[1].field_id);
1879        assert!(!spec.has_sequential_ids());
1880    }
1881
1882    #[test]
1883    fn test_sequential_ids_must_have_no_gaps() {
1884        let schema = Schema::builder()
1885            .with_fields(vec![
1886                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1887                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1888            ])
1889            .build()
1890            .unwrap();
1891
1892        let spec = PartitionSpec::builder(schema)
1893            .with_spec_id(1)
1894            .add_unbound_field(UnboundPartitionField {
1895                source_ids: vec![1],
1896                field_id: Some(1000),
1897                name: "id".to_string(),
1898                transform: Transform::Identity,
1899            })
1900            .unwrap()
1901            .add_unbound_field(UnboundPartitionField {
1902                source_ids: vec![2],
1903                field_id: Some(1002),
1904                name: "name".to_string(),
1905                transform: Transform::Identity,
1906            })
1907            .unwrap()
1908            .build()
1909            .unwrap();
1910
1911        assert_eq!(1000, spec.fields[0].field_id);
1912        assert_eq!(1002, spec.fields[1].field_id);
1913        assert!(!spec.has_sequential_ids());
1914    }
1915
1916    #[test]
1917    fn test_partition_to_path() {
1918        let schema = Schema::builder()
1919            .with_fields(vec![
1920                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1921                NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1922                NestedField::required(3, "timestamp", Type::Primitive(PrimitiveType::Timestamp))
1923                    .into(),
1924                NestedField::required(4, "empty", Type::Primitive(PrimitiveType::String)).into(),
1925            ])
1926            .build()
1927            .unwrap();
1928
1929        let spec = PartitionSpec::builder(schema.clone())
1930            .add_partition_field("id", "id", Transform::Identity)
1931            .unwrap()
1932            .add_partition_field("name", "name", Transform::Identity)
1933            .unwrap()
1934            .add_partition_field("timestamp", "ts_hour", Transform::Hour)
1935            .unwrap()
1936            .add_partition_field("empty", "empty_void", Transform::Void)
1937            .unwrap()
1938            .build()
1939            .unwrap();
1940
1941        let data = Struct::from_iter([
1942            Some(Literal::int(42)),
1943            Some(Literal::string("alice")),
1944            Some(Literal::int(1000)),
1945            Some(Literal::string("empty")),
1946        ]);
1947
1948        assert_eq!(
1949            spec.partition_to_path(&data, schema.into()),
1950            "id=42/name=alice/ts_hour=1970-02-11-16/empty_void=null"
1951        );
1952    }
1953
1954    #[test]
1955    fn test_partition_to_path_escaped_strings() {
1956        let schema = Schema::builder()
1957            .with_fields(vec![
1958                NestedField::required(1, "\"esc\"#1", Type::Primitive(PrimitiveType::String))
1959                    .into(),
1960                NestedField::required(2, "data", Type::Primitive(PrimitiveType::String)).into(),
1961            ])
1962            .build()
1963            .unwrap();
1964
1965        let spec = PartitionSpec::builder(schema.clone())
1966            .add_partition_field("\"esc\"#1", "\"esc\"#1", Transform::Identity)
1967            .unwrap()
1968            .build()
1969            .unwrap();
1970
1971        let data = Struct::from_iter([
1972            Some(Literal::string("a/b/c/d")),
1973            Some(Literal::string("val#1")),
1974        ]);
1975
1976        assert_eq!(
1977            spec.partition_to_path(&data, schema.into()),
1978            "%22esc%22%231=a%2Fb%2Fc%2Fd"
1979        );
1980    }
1981
1982    #[test]
1983    fn test_partition_to_path_escaped_field_name() {
1984        let schema = Schema::builder()
1985            .with_fields(vec![
1986                NestedField::required(1, "\"esc\"#1", Type::Primitive(PrimitiveType::String))
1987                    .into(),
1988                NestedField::required(2, "data", Type::Primitive(PrimitiveType::String)).into(),
1989            ])
1990            .build()
1991            .unwrap();
1992
1993        let spec = PartitionSpec::builder(schema.clone())
1994            .add_partition_field("data", "data", Transform::Identity)
1995            .unwrap()
1996            .add_partition_field("data", "data_truc_10", Transform::Truncate(10))
1997            .unwrap()
1998            .build()
1999            .unwrap();
2000
2001        let data = Struct::from_iter([
2002            Some(Literal::string("a/b/c/d")),
2003            Some(Literal::string("a/b/c/d")),
2004        ]);
2005
2006        assert_eq!(
2007            spec.partition_to_path(&data, schema.into()),
2008            "data=a%2Fb%2Fc%2Fd/data_truc_10=a%2Fb%2Fc%2Fd"
2009        );
2010    }
2011
2012    #[test]
2013    fn test_unbound_partition_field_reads_source_ids_only() {
2014        let field: UnboundPartitionField = serde_json::from_str(
2015            r#"{"source-ids": [1, 2], "name": "m", "transform": "bucket[4]"}"#,
2016        )
2017        .unwrap();
2018
2019        assert_eq!([1, 2], field.source_ids());
2020        // a multi-argument field has no single source id
2021        assert!(
2022            field
2023                .source_id()
2024                .unwrap_err()
2025                .to_string()
2026                .contains("has no single source id")
2027        );
2028
2029        let serialized = serde_json::to_value(&field).unwrap();
2030        assert_eq!(
2031            Some(&serde_json::json!([1, 2])),
2032            serialized.get("source-ids")
2033        );
2034        assert!(serialized.get("source-id").is_none());
2035    }
2036
2037    #[test]
2038    fn test_unbound_partition_field_single_source_id_round_trip() {
2039        let field: UnboundPartitionField =
2040            serde_json::from_str(r#"{"source-id": 1, "name": "m", "transform": "identity"}"#)
2041                .unwrap();
2042
2043        assert_eq!(1, field.source_id().unwrap());
2044        assert_eq!([1], field.source_ids());
2045
2046        // a single-argument field keeps writing source-id
2047        let serialized = serde_json::to_value(&field).unwrap();
2048        assert_eq!(Some(&serde_json::json!(1)), serialized.get("source-id"));
2049        assert!(serialized.get("source-ids").is_none());
2050    }
2051
2052    #[test]
2053    fn test_unbound_partition_field_rejects_malformed_source_ids() {
2054        for (input, expected) in [
2055            (
2056                r#"{"source-ids": [], "name": "m", "transform": "identity"}"#,
2057                "Empty source-ids is not allowed",
2058            ),
2059            (
2060                r#"{"name": "m", "transform": "identity"}"#,
2061                "Either `source-id` or `source-ids` must be present",
2062            ),
2063            (
2064                r#"{"source-id": 1, "source-ids": [1], "name": "m", "transform": "identity"}"#,
2065                "mutually exclusive",
2066            ),
2067        ] {
2068            let err = serde_json::from_str::<UnboundPartitionField>(input).unwrap_err();
2069            assert!(
2070                err.to_string().contains(expected),
2071                "unexpected error for {input}: {err}"
2072            );
2073        }
2074    }
2075
2076    #[test]
2077    fn test_binding_a_multi_argument_field_fails_loudly() {
2078        let schema = Schema::builder()
2079            .with_fields(vec![
2080                NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(),
2081                NestedField::required(2, "b", Type::Primitive(PrimitiveType::Int)).into(),
2082            ])
2083            .build()
2084            .unwrap();
2085
2086        let spec: UnboundPartitionSpec = serde_json::from_str(
2087            r#"{"spec-id": 1, "fields": [{"source-ids": [1, 2], "field-id": 1000, "name": "m", "transform": "bucket[4]"}]}"#,
2088        )
2089        .unwrap();
2090
2091        // A bound PartitionField still carries a single source id, so binding has to refuse
2092        // rather than quietly keep the first one.
2093        let err = spec.bind(schema).unwrap_err();
2094        assert!(
2095            err.to_string().contains("has no single source id"),
2096            "unexpected error: {err}"
2097        );
2098    }
2099
2100    #[test]
2101    fn test_unbound_partition_field_builder_rejects_empty_source_ids() {
2102        let err = UnboundPartitionField::builder()
2103            .source_ids(vec![])
2104            .name("m")
2105            .transform(Transform::Identity)
2106            .build()
2107            .unwrap_err();
2108        assert!(
2109            err.to_string().contains("Empty source-ids is not allowed"),
2110            "unexpected error: {err}"
2111        );
2112    }
2113}