Skip to main content

rc_core/
catalog.rs

1//! Table catalog resources and operations, independent of HTTP and storage SDKs.
2
3use crate::{Error, Result};
4use async_trait::async_trait;
5use serde_json::Value;
6
7#[derive(Clone, Copy, Debug, PartialEq, Eq)]
8pub enum ResourceKind {
9    Warehouse,
10    Namespace,
11    Table,
12}
13
14#[derive(Clone, Debug, PartialEq, Eq)]
15pub struct CatalogTarget {
16    pub alias: String,
17    pub warehouse: String,
18    pub namespace: Vec<String>,
19    pub name: Option<String>,
20}
21
22impl CatalogTarget {
23    pub fn parse(path: &str, kind: ResourceKind) -> Result<Self> {
24        let parts: Vec<_> = path.split('/').collect();
25        let count = match kind {
26            ResourceKind::Warehouse => 2,
27            ResourceKind::Namespace => 3,
28            ResourceKind::Table => 4,
29        };
30        if parts.len() != count || parts[0].is_empty() {
31            return Err(Error::InvalidPath(
32                "Expected alias/warehouse[/namespace[/table]]".into(),
33            ));
34        }
35        validate_warehouse(parts[1])?;
36        let namespace = if count >= 3 {
37            if parts[2].len() > 512 {
38                return Err(Error::InvalidPath("Namespace exceeds 512 bytes".into()));
39            }
40            parts[2]
41                .split('.')
42                .map(|segment| {
43                    validate_segment(segment)?;
44                    Ok(segment.to_owned())
45                })
46                .collect::<Result<Vec<_>>>()?
47        } else {
48            Vec::new()
49        };
50        let name = if count == 4 {
51            validate_segment(parts[3])?;
52            Some(parts[3].to_owned())
53        } else {
54            None
55        };
56        Ok(Self {
57            alias: parts[0].into(),
58            warehouse: parts[1].into(),
59            namespace,
60            name,
61        })
62    }
63}
64
65/// Preserve existing bucket names while rejecting URL path syntax.
66pub fn validate_warehouse(value: &str) -> Result<()> {
67    if value.is_empty()
68        || value.len() > 255
69        || !value
70            .bytes()
71            .all(|b| b.is_ascii_alphanumeric() || b"._-".contains(&b))
72        || matches!(value, "." | "..")
73    {
74        return Err(Error::InvalidPath("Invalid warehouse bucket name".into()));
75    }
76    Ok(())
77}
78
79pub fn validate_segment(value: &str) -> Result<()> {
80    let boundary = |b: u8| b.is_ascii_lowercase() || b.is_ascii_digit();
81    if value.is_empty()
82        || value.len() > 64
83        || !value.bytes().all(|b| boundary(b) || b == b'_' || b == b'-')
84        || !value.bytes().next().is_some_and(boundary)
85        || !value.bytes().last().is_some_and(boundary)
86    {
87        return Err(Error::InvalidPath("Catalog names must be 1-64 lowercase ASCII letters, digits, '-' or '_', with alphanumeric boundaries".into()));
88    }
89    Ok(())
90}
91
92#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
93#[serde(rename_all = "snake_case")]
94pub enum CatalogOperation {
95    Config,
96    WarehouseShow,
97    WarehouseEnable,
98    NamespaceCreate,
99    NamespaceList,
100    NamespaceShow,
101    NamespaceExists,
102    NamespaceUpdate,
103    NamespaceRemove,
104    TableCreate,
105    TableRegister,
106    TableList,
107    TableShow,
108    TableExists,
109    TableRename,
110    TableRemove,
111    MetadataShow,
112    MetadataUpdate,
113    Commit,
114    RefList,
115    RefSet,
116    RefRemove,
117    ViewCreate,
118    ViewList,
119    ViewShow,
120    ViewExists,
121    ViewReplace,
122    ViewRemove,
123    MaintenancePlan,
124    MaintenanceRun,
125    MaintenanceConfigShow,
126    MaintenanceConfigSet,
127    MaintenanceJobShow,
128    SchedulerShow,
129    SchedulerRun,
130    WorkerRun,
131    JobHeartbeat,
132    JobQuarantine,
133    Diagnostics,
134    Export,
135    Import,
136    Recover,
137    Rollback,
138    ExternalShow,
139    ExternalSet,
140    ExternalSync,
141    MigrationStatus,
142    MigrationStart,
143    MigrationCancel,
144}
145
146impl CatalogOperation {
147    pub const fn is_list(self) -> bool {
148        matches!(self, Self::NamespaceList | Self::TableList | Self::ViewList)
149    }
150}
151
152#[derive(Clone, Debug)]
153pub struct CatalogRequest {
154    pub operation: CatalogOperation,
155    pub target: CatalogTarget,
156    pub body: Option<Value>,
157    /// Reference or maintenance job identifier, encoded as one path segment.
158    pub child: Option<String>,
159    pub page_size: u16,
160    pub page_token: Option<String>,
161    pub single_page: bool,
162    pub snapshots: Option<String>,
163}
164
165impl CatalogRequest {
166    pub fn new(operation: CatalogOperation, target: CatalogTarget) -> Self {
167        Self {
168            operation,
169            target,
170            body: None,
171            child: None,
172            page_size: 1000,
173            page_token: None,
174            single_page: false,
175            snapshots: None,
176        }
177    }
178}
179
180#[async_trait]
181pub trait TableCatalogApi: Send + Sync {
182    /// Execute a catalog operation. Writes are never automatically retried.
183    async fn catalog(&self, request: &CatalogRequest) -> Result<Value>;
184}