use std::collections::HashMap;
use axum::{
Json,
extract::{Extension, State},
};
use fraiseql_core::{
db::{AdminSqlOutcome, AdminSqlRequest},
security::AuthenticatedUser,
types::UserId,
};
use serde::{Deserialize, Serialize};
use crate::{
middleware::{AdminCaller, AdminPrivilege},
routes::{
api::types::{ApiError, ApiResponse},
graphql::AppState,
},
server_config::AdminSqlConfig,
};
const AUDIT_SQL_PREVIEW_BYTES: usize = 1024;
#[derive(Clone)]
pub struct AdminSqlState {
pub app: AppState,
pub config: AdminSqlConfig,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub struct ImpersonateClaims {
pub user_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tenant_id: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub roles: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub scopes: Vec<String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub claims: HashMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub struct AdminSqlRequestBody {
pub sql: String,
#[serde(default)]
pub commit: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_rows: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub statement_timeout_ms: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub impersonate: Option<ImpersonateClaims>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct AdminSqlResponse {
pub columns: Vec<String>,
pub rows: Vec<Vec<serde_json::Value>>,
pub truncated: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub rows_affected: Option<u64>,
pub committed: bool,
pub read_only: bool,
pub statement_timeout_ms: u32,
pub max_rows: usize,
}
pub async fn admin_sql_handler(
State(state): State<AdminSqlState>,
caller: Option<Extension<AdminCaller>>,
Json(body): Json<AdminSqlRequestBody>,
) -> Result<Json<ApiResponse<AdminSqlResponse>>, ApiError> {
let Some(Extension(caller)) = caller else {
return Err(ApiError::new(
"Admin SQL console reached without an authenticated admin privilege",
"FORBIDDEN",
));
};
let privilege = caller.privilege;
let prepared = prepare(&state.config, privilege, &body);
let result = match prepared {
Ok(request) => execute(&state, request).await,
Err(refusal) => Err(refusal),
};
audit_execution(privilege, &caller.peer_ip, &body, result.as_ref());
result.map(|(applied, outcome)| {
ApiResponse::success(AdminSqlResponse {
columns: outcome.columns,
rows: outcome.rows,
truncated: outcome.truncated,
rows_affected: outcome.rows_affected,
committed: outcome.committed,
read_only: applied.read_only,
statement_timeout_ms: applied.statement_timeout_ms,
max_rows: applied.max_rows,
})
})
}
#[derive(Debug, Clone, Copy)]
struct AppliedBounds {
read_only: bool,
statement_timeout_ms: u32,
max_rows: usize,
}
fn prepare(
config: &AdminSqlConfig,
privilege: AdminPrivilege,
body: &AdminSqlRequestBody,
) -> Result<(AppliedBounds, AdminSqlRequestPlan), ApiError> {
if body.sql.trim().is_empty() {
return Err(ApiError::validation_error("sql must not be empty"));
}
let read_only = privilege == AdminPrivilege::ReadOnly;
if body.commit {
if !config.allow_commit {
return Err(ApiError::new(
"This server's SQL console is preview-only: [admin_sql] allow_commit = false, \
so a statement can be run and read but never committed",
"FORBIDDEN",
));
}
if read_only {
return Err(ApiError::new(
"A read-only admin token cannot commit. Its transaction runs READ ONLY, so a \
commit would persist nothing while reporting that it had; authenticate with \
admin_token to make a change",
"FORBIDDEN",
));
}
}
let statement_timeout_ms = match body.statement_timeout_ms {
Some(0) => {
return Err(ApiError::validation_error(
"statement_timeout_ms must be greater than 0 (PostgreSQL reads 0 as no timeout)",
));
},
Some(requested) => requested.min(config.statement_timeout_ms),
None => config.statement_timeout_ms,
};
let max_rows = match body.max_rows {
Some(0) => {
return Err(ApiError::validation_error("max_rows must be greater than 0"));
},
Some(requested) => requested.min(config.max_rows),
None => config.max_rows,
};
Ok((
AppliedBounds {
read_only,
statement_timeout_ms,
max_rows,
},
AdminSqlRequestPlan {
sql: body.sql.clone(),
commit: body.commit,
impersonate: body.impersonate.clone(),
},
))
}
struct AdminSqlRequestPlan {
sql: String,
commit: bool,
impersonate: Option<ImpersonateClaims>,
}
async fn execute(
state: &AdminSqlState,
(bounds, plan): (AppliedBounds, AdminSqlRequestPlan),
) -> Result<(AppliedBounds, AdminSqlOutcome), ApiError> {
let executor = state.app.executor();
let session_vars = match plan.impersonate {
Some(ref claims) => impersonated_session_vars(executor.schema(), claims)?,
None => Vec::new(),
};
let request = AdminSqlRequest {
sql: plan.sql,
read_only: bounds.read_only,
commit: plan.commit,
statement_timeout_ms: bounds.statement_timeout_ms,
max_rows: bounds.max_rows,
session_vars,
};
executor
.execute_admin_sql(&request)
.await
.map(|outcome| (bounds, outcome))
.map_err(classify_database_error)
}
const PG_READ_ONLY_SQL_TRANSACTION: &str = "25006";
const PG_QUERY_CANCELED: &str = "57014";
fn classify_database_error(e: fraiseql_core::error::FraiseQLError) -> ApiError {
match e {
fraiseql_core::error::FraiseQLError::Unsupported { message } => {
ApiError::new(format!("Unsupported: {message}"), "UNSUPPORTED_OPERATION")
},
fraiseql_core::error::FraiseQLError::Database {
ref message,
sql_state: Some(ref state),
} if state == PG_READ_ONLY_SQL_TRANSACTION => ApiError::new(
format!(
"Refused by the database: this statement writes, and a read-only admin token \
runs it in a READ ONLY transaction ({message})"
),
"FORBIDDEN",
),
fraiseql_core::error::FraiseQLError::Database {
ref message,
sql_state: Some(ref state),
} if state == PG_QUERY_CANCELED => ApiError::new(
format!("Statement cancelled by statement_timeout ({message})"),
"TIMEOUT",
),
other => ApiError::internal_error(other.to_string()),
}
}
fn impersonated_session_vars(
schema: &fraiseql_core::schema::CompiledSchema,
claims: &ImpersonateClaims,
) -> Result<Vec<(String, String)>, ApiError> {
const RESERVED_PREFIX: &str = "fraiseql.";
if let Some(reserved) = claims.claims.keys().find(|k| k.starts_with(RESERVED_PREFIX)) {
return Err(ApiError::validation_error(format!(
"impersonate.claims may not set '{reserved}': the `{RESERVED_PREFIX}` namespace is \
reserved for values the server derives, and is stripped from real tokens for the \
same reason"
)));
}
let tenant_claim = schema.tenant_claim();
let mut extra_claims = claims.claims.clone();
if let Some(ref tenant) = claims.tenant_id {
let tenant = serde_json::Value::String(tenant.clone());
if extra_claims.get(tenant_claim).is_some_and(|named| *named != tenant) {
return Err(ApiError::validation_error(format!(
"impersonate.tenant_id and impersonate.claims.{tenant_claim} name different \
tenants; `{tenant_claim}` is the schema's tenant claim"
)));
}
extra_claims.insert(tenant_claim.to_string(), tenant);
}
let user = AuthenticatedUser {
user_id: UserId::new(claims.user_id.clone()),
scopes: claims.scopes.clone(),
expires_at: chrono::Utc::now() + chrono::Duration::minutes(5),
email: None,
display_name: None,
extra_claims,
};
let mut context = crate::extractors::build_security_context(
&user,
format!("admin-sql-{}", uuid::Uuid::new_v4()),
Some(tenant_claim),
);
context.roles.clone_from(&claims.roles);
fraiseql_core::runtime::resolve_session_variables(schema, &context).map_err(|e| {
ApiError::validation_error(format!("impersonation could not be resolved: {e}"))
})
}
fn audit_execution(
privilege: AdminPrivilege,
peer_ip: &str,
body: &AdminSqlRequestBody,
result: Result<&(AppliedBounds, AdminSqlOutcome), &ApiError>,
) {
use fraiseql_auth::audit::logger::{AuditEntry, AuditEventType, SecretType, get_audit_logger};
use sha2::{Digest as _, Sha256};
let digest = hex::encode(Sha256::digest(body.sql.as_bytes()));
let preview: String = body.sql.chars().take(AUDIT_SQL_PREVIEW_BYTES).collect();
let privilege = match privilege {
AdminPrivilege::ReadOnly => "admin_readonly_token",
AdminPrivilege::ReadWrite => "admin_token",
};
let committed = result.is_ok_and(|(_, outcome)| outcome.committed);
let impersonating = body.impersonate.as_ref().map_or("none", |c| c.user_id.as_str());
get_audit_logger().log_entry(AuditEntry {
event_type: AuditEventType::AdminSqlExecution,
secret_type: SecretType::AdminToken,
subject: Some(privilege.to_string()),
operation: "admin_sql".to_string(),
success: result.is_ok(),
error_message: result.err().map(ToString::to_string),
context: Some(format!(
"peer_ip={peer_ip} sha256={digest} commit_requested={} committed={committed} \
impersonate={impersonating} sql={preview}",
body.commit
)),
chain_hash: None,
});
}