Skip to content

Commit fe25056

Browse files
committed
impl: snowflake integration
1 parent 094a085 commit fe25056

14 files changed

Lines changed: 1621 additions & 2 deletions

File tree

Cargo.lock

Lines changed: 45 additions & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
[workspace]
2-
members = ["crates/binlog-protocol", "crates/mysql-binlog-source", "crates/kafka-source", "crates/kafka-producer", "crates/kafka-types", "crates/postgresql", "crates/postgresql-wal2json-source", "crates/postgresql-pgoutput-source", "crates/pgoutput-protocol", "crates/postgresql-trigger-source", "crates/csv-source", "crates/jsonl-source", "crates/file", "crates/sync-core", "crates/loadtest-generator", "crates/loadtest-populate", "crates/mysql-types", "crates/postgresql-types", "crates/mongodb-types", "crates/surreal2-types", "crates/surreal3-types", "crates/json-types", "crates/csv-types", "crates/neo4j-types", "crates/loadtest-populate-mysql", "crates/loadtest-populate-postgresql", "crates/loadtest-populate-mongodb", "crates/loadtest-populate-csv", "crates/loadtest-populate-jsonl", "crates/loadtest-populate-neo4j", "crates/loadtest-populate-kafka", "crates/loadtest-verify-surreal2", "crates/loadtest-verify-surreal3", "crates/surreal-sink", "crates/surreal2-sink", "crates/surreal3-sink", "crates/surreal-client", "crates/surreal2-client", "crates/surreal3-client", "crates/surreal-version", "crates/mongodb-changestream-source", "crates/neo4j-source", "crates/checkpoint", "crates/checkpoint-surreal3", "crates/mysql-trigger-source", "crates/loadtest-distributed", "crates/surrealdb-multi-sdk-demo", "crates/interleaved-snapshot"]
2+
members = ["crates/binlog-protocol", "crates/mysql-binlog-source", "crates/kafka-source", "crates/kafka-producer", "crates/kafka-types", "crates/postgresql", "crates/postgresql-wal2json-source", "crates/postgresql-pgoutput-source", "crates/pgoutput-protocol", "crates/postgresql-trigger-source", "crates/csv-source", "crates/jsonl-source", "crates/file", "crates/sync-core", "crates/loadtest-generator", "crates/loadtest-populate", "crates/mysql-types", "crates/postgresql-types", "crates/mongodb-types", "crates/surreal2-types", "crates/surreal3-types", "crates/json-types", "crates/csv-types", "crates/neo4j-types", "crates/loadtest-populate-mysql", "crates/loadtest-populate-postgresql", "crates/loadtest-populate-mongodb", "crates/loadtest-populate-csv", "crates/loadtest-populate-jsonl", "crates/loadtest-populate-neo4j", "crates/loadtest-populate-kafka", "crates/loadtest-verify-surreal2", "crates/loadtest-verify-surreal3", "crates/surreal-sink", "crates/surreal2-sink", "crates/surreal3-sink", "crates/surreal-client", "crates/surreal2-client", "crates/surreal3-client", "crates/surreal-version", "crates/mongodb-changestream-source", "crates/neo4j-source", "crates/checkpoint", "crates/checkpoint-surreal3", "crates/mysql-trigger-source", "crates/loadtest-distributed", "crates/surrealdb-multi-sdk-demo", "crates/interleaved-snapshot", "crates/snowflake-types", "crates/snowflake-source"]
33
resolver = "2"
44

55
[package]
@@ -142,6 +142,10 @@ surreal-sync-postgresql = { path = "crates/postgresql" }
142142
# Watermark-based interleaved snapshot full-sync framework
143143
surreal-sync-interleaved-snapshot = { path = "crates/interleaved-snapshot" }
144144

145+
# Snowflake ingestion source (SQL REST API v2, key-pair JWT auth)
146+
surreal-sync-snowflake-source = { path = "crates/snowflake-source" }
147+
snowflake-types = { path = "crates/snowflake-types" }
148+
145149
# Binlog protocol (Flavor for binlog test containers)
146150
binlog-protocol = { path = "crates/binlog-protocol" }
147151

crates/snowflake-source/Cargo.toml

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
[package]
2+
name = "surreal-sync-snowflake-source"
3+
version = "0.1.0"
4+
edition = "2021"
5+
description = "Snowflake ingestion source for surreal-sync (SQL REST API v2, key-pair JWT auth)"
6+
authors = ["surreal-sync contributors"]
7+
8+
[dependencies]
9+
# Core data model + sink trait
10+
sync-core = { path = "../sync-core" }
11+
surreal-sink = { path = "../surreal-sink" }
12+
snowflake-types = { path = "../snowflake-types" }
13+
14+
# HTTP client (rustls to match the process-wide aws-lc-rs provider; no OpenSSL)
15+
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }
16+
17+
# Snowflake key-pair JWT generation
18+
snowflake-jwt = "0.3"
19+
20+
# Async runtime (timers for async statement polling)
21+
tokio = { version = "1.49", features = ["time"] }
22+
23+
# Serialization
24+
serde = { version = "1.0", features = ["derive"] }
25+
serde_json = "1.0"
26+
27+
# Error handling / logging
28+
anyhow = "1.0"
29+
tracing = "0.1"
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
//! Table discovery for a Snowflake schema.
2+
3+
use anyhow::Result;
4+
5+
use crate::client::SnowflakeClient;
6+
7+
/// List base-table names in `database.schema` via `INFORMATION_SCHEMA.TABLES`.
8+
///
9+
/// Snowflake stores unquoted identifiers upper-cased, so `schema` is upper-cased
10+
/// before matching. Views and temporary tables are excluded (`TABLE_TYPE = 'BASE TABLE'`).
11+
pub async fn list_tables(
12+
client: &SnowflakeClient,
13+
database: &str,
14+
schema: &str,
15+
) -> Result<Vec<String>> {
16+
let db = database.to_ascii_uppercase();
17+
let schema_upper = schema.to_ascii_uppercase();
18+
let sql = format!(
19+
"SELECT TABLE_NAME FROM \"{db}\".INFORMATION_SCHEMA.TABLES \
20+
WHERE TABLE_SCHEMA = '{schema_upper}' AND TABLE_TYPE = 'BASE TABLE' \
21+
ORDER BY TABLE_NAME"
22+
);
23+
24+
let result = client.execute_query(&sql).await?;
25+
let tables = result
26+
.rows
27+
.iter()
28+
.filter_map(|row| row.first())
29+
.filter_map(|cell| cell.as_str())
30+
.map(|s| s.to_string())
31+
.collect();
32+
Ok(tables)
33+
}

0 commit comments

Comments
 (0)