Skip to content

Commit 66818c9

Browse files
BtXinclaude
andcommitted
refactor(validator): type-split ad-hoc policy, share DataFusion's parser, tidy validator
Second-round review follow-ups (#3#10): - #4 Encode the trust boundary in the type. Ad-hoc /query now validates against a dedicated AdhocSqlPolicy (access modes + denied schemas); the trusted pipeline path keeps a bare SqlValidatorConfig. Reaching for validate_sql on untrusted input can no longer silently skip the allowlist + schema denial. - #5 Add coverage proving the schema denial descends into indirect relations (set ops, scalar/IN/EXISTS subqueries) via visit_relations; document the residual (operator-defined views/federated aliases, which ad-hoc SQL cannot create) and its structural fix. See follow-up issue. - #6 Parse via datafusion::sql::sqlparser (DataFusion enables the visitor feature); drop the standalone sqlparser dep entirely so validator and engine share one parser by construction — a DF bump is now a compile break, not silent parse divergence. - #8 Single AUTH_SCHEMA const used at both the register and deny sites. - #9 check_denied_schemas reports the table via extract_table_name (quote-stripped), matching WriteNotAllowed. - #10 statement_keyword maps the AST variant to a &'static str instead of Display-rendering the whole statement on each rejection. - #7 AppState::new derives the policy once; removes the copy-pasted Arc::new(validator_config_from_sources(..)) across ~9 sites. - #3 Document the startup-snapshot invariant (no runtime access_mode writer) on the field and the builder. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent ca2c8a0 commit 66818c9

14 files changed

Lines changed: 230 additions & 186 deletions

Cargo.lock

Lines changed: 0 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: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,6 @@ rstest = "0.25.0"
6464
serde = { version = "1.0", features = ["derive"] }
6565
serde_json = "1.0"
6666
serde_yaml = "0.9"
67-
sqlparser = { version = "0.59", features = ["visitor"] }
6867
sqlx = { version = "0.8", default-features = false, features = ["runtime-tokio", "tls-rustls", "postgres"] }
6968
tempfile = "3.23.0"
7069
tokio = { version = "1.44.2", features = ["macros", "rt", "sync"] }

crates/server/src/auth/bridge.rs

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -178,8 +178,15 @@ impl TableProvider for AuthSessionsTable {
178178
// Registration helper
179179
// ---------------------------------------------------------------------------
180180

181+
/// Schema the auth virtual tables (`users`, `sessions`) register under.
182+
///
183+
/// Single source of truth for the reserved schema name: the ad-hoc `/query`
184+
/// policy denies exactly this schema (see `config::adhoc_policy_from_sources`),
185+
/// so the two must never drift apart.
186+
pub const AUTH_SCHEMA: &str = "auth";
187+
181188
/// Register `auth.users` and `auth.sessions` virtual tables into the given
182-
/// DataFusion `SessionContext` under a dedicated `auth` schema.
189+
/// DataFusion `SessionContext` under the [`AUTH_SCHEMA`] schema.
183190
pub fn register_auth_tables(
184191
ctx: &mut SessionContext,
185192
auth: Arc<BetterAuth<DieselSqliteAdapter>>,
@@ -204,12 +211,14 @@ pub fn register_auth_tables(
204211
.catalog("datafusion")
205212
.ok_or_else(|| anyhow::anyhow!("Default catalog 'datafusion' not found"))?;
206213

207-
if catalog.schema("auth").is_some() {
208-
return Err(anyhow::anyhow!("Auth schema 'auth' is already registered"));
214+
if catalog.schema(AUTH_SCHEMA).is_some() {
215+
return Err(anyhow::anyhow!(
216+
"Auth schema '{AUTH_SCHEMA}' is already registered"
217+
));
209218
}
210219

211220
catalog
212-
.register_schema("auth", schema)
221+
.register_schema(AUTH_SCHEMA, schema)
213222
.map_err(|e| anyhow::anyhow!("Failed to register auth schema: {}", e))?;
214223

215224
tracing::info!("Registered DataFusion tables: auth.users, auth.sessions");

crates/server/src/auth/routes.rs

Lines changed: 4 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -363,13 +363,12 @@ mod tests {
363363
fn make_no_auth_state() -> AppState {
364364
use crate::auth::layer::AuthLayer;
365365
use crate::config::{CliArgs, ServerConfig};
366-
use crate::metrics::PipelineMetrics;
367366
use crate::semantics::SemanticsRegistry;
368367
use crate::server::AppState;
369368
use datafusion::prelude::SessionContext;
370369
use skardi::engine::datafusion::DataFusionEngine;
371370
use std::path::PathBuf;
372-
use std::sync::{Arc, RwLock};
371+
use std::sync::Arc;
373372

374373
let config = ServerConfig {
375374
pipelines: Default::default(),
@@ -387,15 +386,7 @@ mod tests {
387386
};
388387
let session_ctx = Arc::new(SessionContext::new());
389388
let engine = Arc::new(DataFusionEngine::new_with_arc(session_ctx.clone()));
390-
AppState {
391-
config: Arc::new(RwLock::new(config)),
392-
engine,
393-
session_ctx,
394-
metrics: PipelineMetrics::new(),
395-
auth_layer: AuthLayer::None,
396-
jobs: None,
397-
validator_config: Arc::new(crate::config::validator_config_from_sources(&[])),
398-
}
389+
AppState::new(config, engine, session_ctx, AuthLayer::None, None)
399390
}
400391

401392
#[tokio::test]
@@ -408,13 +399,12 @@ mod tests {
408399
async fn make_better_auth_state() -> AppState {
409400
use crate::auth::layer::AuthLayer;
410401
use crate::config::{CliArgs, ServerConfig};
411-
use crate::metrics::PipelineMetrics;
412402
use crate::semantics::SemanticsRegistry;
413403
use crate::server::AppState;
414404
use datafusion::prelude::SessionContext;
415405
use skardi::engine::datafusion::DataFusionEngine;
416406
use std::path::PathBuf;
417-
use std::sync::{Arc, RwLock};
407+
use std::sync::Arc;
418408

419409
unsafe {
420410
std::env::set_var("AUTH_SECRET", "test-secret-that-is-at-least-32-characters!");
@@ -442,15 +432,7 @@ mod tests {
442432
};
443433
let session_ctx = Arc::new(SessionContext::new());
444434
let engine = Arc::new(DataFusionEngine::new_with_arc(session_ctx.clone()));
445-
AppState {
446-
config: Arc::new(RwLock::new(config)),
447-
engine,
448-
session_ctx,
449-
metrics: PipelineMetrics::new(),
450-
auth_layer: layer,
451-
jobs: None,
452-
validator_config: Arc::new(crate::config::validator_config_from_sources(&[])),
453-
}
435+
AppState::new(config, engine, session_ctx, layer, None)
454436
}
455437

456438
#[tokio::test]

crates/server/src/config.rs

Lines changed: 19 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ use skardi::sources::providers::redis::datasource::register_redis_tables;
1818
use skardi::sources::providers::seekdb::register_seekdb_tables;
1919
use skardi::sources::providers::sqlite::register_sqlite_tables;
2020
use skardi::sources::providers::sqlx::postgres::register_postgres_tables;
21-
use skardi::sources::sql_validator::{SqlValidatorConfig, validate_sql};
21+
use skardi::sources::sql_validator::{AdhocSqlPolicy, SqlValidatorConfig, validate_sql};
2222
use std::collections::HashMap;
2323
use std::path::Path;
2424
use std::path::PathBuf;
@@ -947,21 +947,31 @@ fn validate_schema_types(_schema: &HashMap<String, String>) -> Result<()> {
947947
Ok(())
948948
}
949949

950-
/// Build a SQL validator config mapping every data source name to its
951-
/// configured access mode. Used at config load (pipeline SQL) and at
952-
/// request time (`POST /query`).
950+
/// Build the access-mode map for every data source. Shared by both the
951+
/// trusted pipeline-load path and the untrusted `/query` policy below.
953952
pub fn validator_config_from_sources(data_sources: &[DataSource]) -> SqlValidatorConfig {
954-
// The `auth` schema (auth.users / auth.sessions, registered on the same
955-
// SessionContext) holds live bearer tokens; ad-hoc SQL must never reach
956-
// it. The denial only applies to `validate_single_sql`, so pipeline SQL
957-
// may still read auth tables.
958-
let mut validator_config = SqlValidatorConfig::new().with_denied_schema("auth");
953+
let mut validator_config = SqlValidatorConfig::new();
959954
for ds in data_sources {
960955
validator_config = validator_config.with_table(&ds.name, ds.access_mode);
961956
}
962957
validator_config
963958
}
964959

960+
/// Build the statement policy for the untrusted ad-hoc `/query` endpoint:
961+
/// the access-mode map plus the reserved [`AUTH_SCHEMA`] denial (auth.users /
962+
/// auth.sessions register on the same `SessionContext` and hold live bearer
963+
/// tokens, so ad-hoc SQL must never reach them). The denial is scoped to this
964+
/// policy, so operator-authored pipeline SQL may still read auth tables.
965+
///
966+
/// Callers snapshot this once at startup into `AppState`. That is correct only
967+
/// as long as nothing mutates a source's `access_mode` at runtime — there is
968+
/// no such writer today. If one is ever added, it must rebuild this policy (or
969+
/// the snapshot will serve a stale, potentially more-permissive gate).
970+
pub fn adhoc_policy_from_sources(data_sources: &[DataSource]) -> AdhocSqlPolicy {
971+
AdhocSqlPolicy::new(validator_config_from_sources(data_sources))
972+
.with_denied_schema(crate::auth::bridge::AUTH_SCHEMA)
973+
}
974+
965975
/// Validate pipeline SQL against data source access modes
966976
fn validate_pipeline_sql(
967977
pipeline_name: &str,

crates/server/src/pipeline_handlers.rs

Lines changed: 12 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -778,7 +778,6 @@ pub async fn execute_pipeline_by_name(
778778
mod tests {
779779
use super::*;
780780
use crate::config::{CliArgs, DataSource, DataSourceType, ServerConfig};
781-
use crate::metrics::PipelineMetrics;
782781
use crate::server::AppState;
783782
use arrow::array::{Int64Array, StringArray};
784783
use arrow::datatypes::{DataType, Field, Schema};
@@ -789,7 +788,7 @@ mod tests {
789788
use skardi::sources::AccessMode;
790789
use std::fs;
791790
use std::path::PathBuf;
792-
use std::sync::{Arc, RwLock};
791+
use std::sync::Arc;
793792
use tempfile::TempDir;
794793

795794
async fn create_test_pipeline_with_params() -> StandardPipeline {
@@ -847,18 +846,13 @@ spec:
847846
let session_ctx = Arc::new(SessionContext::new());
848847
let engine = Arc::new(DataFusionEngine::new_with_arc(session_ctx.clone()));
849848

850-
let validator_config = Arc::new(crate::config::validator_config_from_sources(
851-
&config.data_sources,
852-
));
853-
AppState {
854-
config: Arc::new(RwLock::new(config)),
849+
AppState::new(
850+
config,
855851
engine,
856852
session_ctx,
857-
metrics: PipelineMetrics::new(),
858-
auth_layer: crate::auth::layer::AuthLayer::None,
859-
jobs: None,
860-
validator_config,
861-
}
853+
crate::auth::layer::AuthLayer::None,
854+
None,
855+
)
862856
}
863857

864858
fn create_test_record_batch() -> RecordBatch {
@@ -941,18 +935,13 @@ spec:
941935
let session_ctx_arc = Arc::new(session_ctx);
942936
let engine = Arc::new(DataFusionEngine::new_with_arc(session_ctx_arc.clone()));
943937

944-
let validator_config = Arc::new(crate::config::validator_config_from_sources(
945-
&config.data_sources,
946-
));
947-
let app_state = AppState {
948-
config: Arc::new(RwLock::new(config)),
938+
let app_state = AppState::new(
939+
config,
949940
engine,
950-
session_ctx: session_ctx_arc,
951-
metrics: PipelineMetrics::new(),
952-
auth_layer: crate::auth::layer::AuthLayer::None,
953-
jobs: None,
954-
validator_config,
955-
};
941+
session_ctx_arc,
942+
crate::auth::layer::AuthLayer::None,
943+
None,
944+
);
956945

957946
let request = ExecuteRequest {
958947
parameters: {

crates/server/src/query_handlers.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ pub async fn execute_query(
6868
None => DEFAULT_MAX_ROWS,
6969
};
7070

71-
let statement_kind = match validate_single_sql(&request.sql, &app_state.validator_config) {
71+
let statement_kind = match validate_single_sql(&request.sql, &app_state.adhoc_policy) {
7272
Ok(kind) => kind,
7373
Err(e) => {
7474
tracing::info!("Rejected ad-hoc query: {}", e);

crates/server/src/server.rs

Lines changed: 34 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ use datafusion::prelude::SessionContext;
77
use skardi::engine::datafusion::DataFusionEngine;
88
use skardi::jobs::{JobExecutor, JobStore, SqliteJobStore};
99
use skardi::sources::DataSourceType;
10-
use skardi::sources::sql_validator::SqlValidatorConfig;
10+
use skardi::sources::sql_validator::AdhocSqlPolicy;
1111
use std::collections::HashMap;
1212
use std::path::PathBuf;
1313
use std::sync::{Arc, RwLock};
@@ -55,10 +55,38 @@ pub struct AppState {
5555
/// Jobs executor + run ledger. `None` when the server was started
5656
/// without `--jobs`, which disables every `/jobs/*` endpoint.
5757
pub jobs: Option<Arc<JobExecutor>>,
58-
/// Statement policy for the ad-hoc `/query` endpoint, built once at
59-
/// startup from the configured data sources (there is no runtime
60-
/// config writer, so a snapshot is sufficient).
61-
pub validator_config: Arc<SqlValidatorConfig>,
58+
/// Statement policy for the ad-hoc `/query` endpoint. Derived from the
59+
/// data sources' access modes once at startup by [`AppState::new`]; see
60+
/// [`crate::config::adhoc_policy_from_sources`] for why a snapshot is
61+
/// safe (no runtime writer mutates access modes).
62+
pub adhoc_policy: Arc<AdhocSqlPolicy>,
63+
}
64+
65+
impl AppState {
66+
/// Assemble the shared state, deriving the ad-hoc `/query` policy from the
67+
/// config's data sources. Centralizing construction here keeps the derived
68+
/// policy from drifting across the many call sites that build an
69+
/// `AppState` (handlers, tests, integration harnesses).
70+
pub fn new(
71+
config: ServerConfig,
72+
engine: Arc<DataFusionEngine>,
73+
session_ctx: Arc<SessionContext>,
74+
auth_layer: AuthLayer,
75+
jobs: Option<Arc<JobExecutor>>,
76+
) -> Self {
77+
let adhoc_policy = Arc::new(crate::config::adhoc_policy_from_sources(
78+
&config.data_sources,
79+
));
80+
Self {
81+
config: Arc::new(RwLock::new(config)),
82+
engine,
83+
session_ctx,
84+
metrics: PipelineMetrics::new(),
85+
auth_layer,
86+
jobs,
87+
adhoc_policy,
88+
}
89+
}
6290
}
6391

6492
/// Main server creation function - Primary public interface
@@ -221,20 +249,8 @@ pub async fn setup_app_state(config: ServerConfig) -> Result<AppState> {
221249
None
222250
};
223251

224-
let validator_config = Arc::new(crate::config::validator_config_from_sources(
225-
&config.data_sources,
226-
));
227-
228252
// Create shared application state with RwLock for runtime updates
229-
let app_state = AppState {
230-
config: Arc::new(RwLock::new(config)),
231-
engine,
232-
session_ctx: session_ctx_arc,
233-
metrics: PipelineMetrics::new(),
234-
auth_layer,
235-
jobs: jobs_bundle,
236-
validator_config,
237-
};
253+
let app_state = AppState::new(config, engine, session_ctx_arc, auth_layer, jobs_bundle);
238254

239255
tracing::info!("Application state setup completed successfully");
240256

crates/server/tests/jobs_http.rs

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,14 +15,13 @@ use skardi::jobs::{JobDefinition, JobExecutor, JobStore, SqliteJobStore};
1515
use skardi::sources::DataSourceType;
1616
use std::collections::HashMap;
1717
use std::io::Write;
18-
use std::sync::{Arc, RwLock};
18+
use std::sync::Arc;
1919
use std::time::Duration;
2020
use tempfile::TempDir;
2121
use tower::ServiceExt;
2222

2323
use skardi_server::auth::layer::AuthLayer;
2424
use skardi_server::config::{CliArgs, ServerConfig};
25-
use skardi_server::metrics::PipelineMetrics;
2625
use skardi_server::semantics::SemanticsRegistry;
2726
use skardi_server::server::{AppState, configure_routes};
2827

@@ -106,18 +105,7 @@ spec:
106105
port: 0,
107106
},
108107
};
109-
let validator_config = Arc::new(skardi_server::config::validator_config_from_sources(
110-
&config.data_sources,
111-
));
112-
let state = AppState {
113-
config: Arc::new(RwLock::new(config)),
114-
engine,
115-
session_ctx: Arc::clone(&ctx),
116-
metrics: PipelineMetrics::new(),
117-
auth_layer: AuthLayer::None,
118-
jobs: executor,
119-
validator_config,
120-
};
108+
let state = AppState::new(config, engine, Arc::clone(&ctx), AuthLayer::None, executor);
121109
(state, tmp)
122110
}
123111

crates/server/tests/pipelines_http.rs

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,12 @@ use serde_json::{Value, json};
1515
use skardi::pipeline::pipeline::{Pipeline, StandardPipeline};
1616
use std::collections::HashMap;
1717
use std::io::Write;
18-
use std::sync::{Arc, RwLock};
18+
use std::sync::Arc;
1919
use tempfile::TempDir;
2020
use tower::ServiceExt;
2121

2222
use skardi_server::auth::layer::AuthLayer;
2323
use skardi_server::config::{CliArgs, ServerConfig};
24-
use skardi_server::metrics::PipelineMetrics;
2524
use skardi_server::semantics::SemanticsRegistry;
2625
use skardi_server::server::{AppState, configure_routes};
2726

@@ -110,18 +109,7 @@ spec:
110109
port: 0,
111110
},
112111
};
113-
let validator_config = Arc::new(skardi_server::config::validator_config_from_sources(
114-
&config.data_sources,
115-
));
116-
let state = AppState {
117-
config: Arc::new(RwLock::new(config)),
118-
engine,
119-
session_ctx: Arc::clone(&ctx),
120-
metrics: PipelineMetrics::new(),
121-
auth_layer: AuthLayer::None,
122-
jobs: None,
123-
validator_config,
124-
};
112+
let state = AppState::new(config, engine, Arc::clone(&ctx), AuthLayer::None, None);
125113
(state, tmp)
126114
}
127115

0 commit comments

Comments
 (0)