#![allow(clippy::unwrap_used, clippy::panic, clippy::print_stderr)]
use std::sync::{Arc, LazyLock, Mutex};
use fraiseql_auth::audit::logger::{AuditEntry, AuditLogger, init_audit_logger};
use fraiseql_core::{
db::postgres::PostgresAdapter,
prelude::DatabaseAdapter as _,
schema::{
CompiledSchema, FieldType, QueryDefinition, SessionVariableMapping, SessionVariableSource,
TypeDefinition,
},
};
use fraiseql_server::server_config::{AdminSqlConfig, ServerConfig};
use fraiseql_test_support::try_database_url;
use serde_json::{Value, json};
mod common;
use crate::common::server_harness::TestServer;
const SCHEMA: &str = "p20_console";
const WRITE_TOKEN: &str = "p20-console-write-token-0123456789abcdef";
const READONLY_TOKEN: &str = "p20-console-readonly-token-0123456789abc";
const TENANT_A: &str = "11111111-1111-1111-1111-111111111111";
const TENANT_B: &str = "22222222-2222-2222-2222-222222222222";
const TENANT_VAR: &str = "app.tenant_id";
fn database_url_or_skip(test: &str) -> Option<String> {
let url = try_database_url();
if url.is_none() {
eprintln!("SKIP {test}: DATABASE_URL not set");
}
url
}
struct CapturingAuditLogger {
entries: Mutex<Vec<AuditEntry>>,
}
impl AuditLogger for CapturingAuditLogger {
fn log_entry(&self, entry: AuditEntry) {
self.entries.lock().unwrap().push(entry);
}
}
static LEDGER: LazyLock<Arc<CapturingAuditLogger>> = LazyLock::new(|| {
let logger = Arc::new(CapturingAuditLogger {
entries: Mutex::new(Vec::new()),
});
init_audit_logger(logger.clone());
logger
});
fn ledger() -> &'static Arc<CapturingAuditLogger> {
&LEDGER
}
fn console_entries() -> Vec<AuditEntry> {
ledger()
.entries
.lock()
.unwrap()
.iter()
.filter(|e| e.event_type.as_str() == "admin_sql_execution")
.cloned()
.collect()
}
async fn seed(adapter: &PostgresAdapter) {
let stmts = vec![
format!("DROP SCHEMA IF EXISTS {SCHEMA} CASCADE"),
format!("CREATE SCHEMA {SCHEMA}"),
format!(
"CREATE TABLE {SCHEMA}.tb_doc (
id bigint PRIMARY KEY,
tenant_id uuid NOT NULL,
title text NOT NULL
)"
),
format!(
"INSERT INTO {SCHEMA}.tb_doc (id, tenant_id, title) VALUES
(1,'{TENANT_A}','a-one'),
(2,'{TENANT_A}','a-two'),
(3,'{TENANT_A}','a-three'),
(4,'{TENANT_B}','b-one'),
(5,'{TENANT_B}','b-two')"
),
format!(
"CREATE VIEW {SCHEMA}.v_scoped_doc AS
SELECT id, tenant_id, title FROM {SCHEMA}.tb_doc
WHERE tenant_id::text = current_setting('{TENANT_VAR}', true)"
),
format!(
"CREATE VIEW {SCHEMA}.v_doc AS
SELECT id, tenant_id,
jsonb_build_object('id', id, 'title', title) AS data
FROM {SCHEMA}.tb_doc ORDER BY id"
),
];
for stmt in stmts {
let _: Vec<std::collections::HashMap<String, Value>> =
adapter.execute_raw_query(&stmt).await.expect("fixture setup");
}
}
fn schema() -> CompiledSchema {
let mut schema = CompiledSchema::new();
let mut doc = TypeDefinition::new("ConsoleDoc", format!("{SCHEMA}.v_doc"));
doc.fields = vec![
fraiseql_core::schema::FieldDefinition::new("id", FieldType::Int),
fraiseql_core::schema::FieldDefinition::new("title", FieldType::String),
];
schema.types.push(doc);
schema.queries.push(
QueryDefinition::new("docs", "ConsoleDoc")
.returning_list()
.with_sql_source(format!("{SCHEMA}.v_doc")),
);
schema.session_variables.variables.push(SessionVariableMapping {
name: TENANT_VAR.to_string(),
source: SessionVariableSource::Jwt {
claim: "tenant_id".to_string(),
},
});
schema.build_indexes();
schema
}
const fn console_config() -> AdminSqlConfig {
AdminSqlConfig {
enabled: true,
statement_timeout_ms: 2_000,
max_rows: 3,
allow_commit: true,
}
}
fn config(admin_sql: Option<AdminSqlConfig>) -> ServerConfig {
ServerConfig {
cors_enabled: false,
admin_api_enabled: true,
admin_token: Some(WRITE_TOKEN.to_string()),
admin_readonly_token: Some(READONLY_TOKEN.to_string()),
admin_sql,
..ServerConfig::default()
}
}
async fn run_sql(server: &TestServer, token: &str, body: Value) -> (reqwest::StatusCode, Value) {
let resp = reqwest::Client::new()
.post(format!("{}/api/v1/admin/sql", server.url))
.bearer_auth(token)
.json(&body)
.send()
.await
.expect("request");
let status = resp.status();
let text = resp.text().await.expect("body");
let parsed = serde_json::from_str(&text).unwrap_or(Value::String(text));
(status, parsed)
}
async fn ok_sql(server: &TestServer, token: &str, body: Value) -> Value {
let (status, parsed) = run_sql(server, token, body).await;
assert_eq!(status, 200, "expected success, got {status}: {parsed}");
parsed["data"].clone()
}
async fn table_count(adapter: &PostgresAdapter) -> i64 {
let rows: Vec<std::collections::HashMap<String, Value>> = adapter
.execute_raw_query(&format!("SELECT count(*)::bigint AS n FROM {SCHEMA}.tb_doc"))
.await
.expect("count");
rows[0]["n"].as_i64().expect("count is a number")
}
struct Rig {
server: TestServer,
adapter: Arc<PostgresAdapter>,
}
async fn boot(admin_sql: Option<AdminSqlConfig>) -> Option<Rig> {
let _ = ledger();
let url = try_database_url()?;
let adapter = Arc::new(PostgresAdapter::new(&url).await.expect("adapter"));
seed(&adapter).await;
let server =
Box::pin(TestServer::start_with_config(config(admin_sql), schema(), adapter.clone())).await;
Some(Rig { server, adapter })
}
async fn boot_console() -> Option<Rig> {
boot(Some(console_config())).await
}
#[tokio::test]
async fn rollback_is_the_default_and_the_write_does_not_persist() {
if database_url_or_skip("rollback_is_the_default_and_the_write_does_not_persist").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let before = table_count(&rig.adapter).await;
let data = ok_sql(
&rig.server,
WRITE_TOKEN,
json!({
"sql": format!(
"INSERT INTO {SCHEMA}.tb_doc (id, tenant_id, title) \
VALUES (900, '{TENANT_A}', 'preview-only') RETURNING id, title"
)
}),
)
.await;
assert_eq!(
data["rows"],
json!([[900, "preview-only"]]),
"the statement must actually have run: {data}"
);
assert_eq!(data["committed"], json!(false), "the default must not commit: {data}");
assert_eq!(table_count(&rig.adapter).await, before, "the row must not survive the request");
}
#[tokio::test]
async fn commit_opt_in_persists_the_write() {
if database_url_or_skip("commit_opt_in_persists_the_write").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let before = table_count(&rig.adapter).await;
let data = ok_sql(
&rig.server,
WRITE_TOKEN,
json!({
"sql": format!(
"INSERT INTO {SCHEMA}.tb_doc (id, tenant_id, title) \
VALUES (901, '{TENANT_B}', 'committed')"
),
"commit": true
}),
)
.await;
assert_eq!(data["committed"], json!(true), "the opt-in must be honoured: {data}");
assert_eq!(data["rows_affected"], json!(1), "one row inserted: {data}");
assert_eq!(
table_count(&rig.adapter).await,
before + 1,
"the committed row must survive the request"
);
}
#[tokio::test]
async fn a_readonly_token_cannot_update() {
if database_url_or_skip("a_readonly_token_cannot_update").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let (status, body) = run_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": format!("UPDATE {SCHEMA}.tb_doc SET title = 'hijacked'") }),
)
.await;
assert_eq!(status, 403, "a read-only token's write must be refused: {body}");
let hijacked: Vec<std::collections::HashMap<String, Value>> = rig
.adapter
.execute_raw_query(&format!(
"SELECT count(*)::bigint AS n FROM {SCHEMA}.tb_doc WHERE title = 'hijacked'"
))
.await
.unwrap();
assert_eq!(hijacked[0]["n"].as_i64(), Some(0), "no row may have been touched");
}
#[tokio::test]
async fn a_readonly_token_can_select() {
if database_url_or_skip("a_readonly_token_can_select").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let data = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": format!("SELECT title FROM {SCHEMA}.tb_doc WHERE id = 1") }),
)
.await;
assert_eq!(data["rows"], json!([["a-one"]]), "{data}");
assert_eq!(data["read_only"], json!(true), "the transaction ran READ ONLY: {data}");
}
#[tokio::test]
async fn a_readonly_token_cannot_ask_for_a_commit() {
if database_url_or_skip("a_readonly_token_cannot_ask_for_a_commit").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let (status, body) =
run_sql(&rig.server, READONLY_TOKEN, json!({ "sql": "SELECT 1", "commit": true })).await;
assert_eq!(status, 403, "{body}");
assert!(
body["error"].as_str().unwrap_or_default().contains("read-only"),
"the refusal must name the reason: {body}"
);
}
#[tokio::test]
async fn the_row_cap_truncates_and_says_so() {
if database_url_or_skip("the_row_cap_truncates_and_says_so").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let data = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": format!("SELECT id FROM {SCHEMA}.tb_doc ORDER BY id") }),
)
.await;
assert_eq!(data["rows"].as_array().map(Vec::len), Some(3), "capped at 3: {data}");
assert_eq!(data["truncated"], json!(true), "and it must say so: {data}");
let exact = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": format!("SELECT id FROM {SCHEMA}.tb_doc ORDER BY id LIMIT 3") }),
)
.await;
assert_eq!(exact["rows"].as_array().map(Vec::len), Some(3), "{exact}");
assert_eq!(exact["truncated"], json!(false), "3 of 3 is complete: {exact}");
}
#[tokio::test]
async fn a_request_may_lower_a_bound_but_never_raise_it() {
if database_url_or_skip("a_request_may_lower_a_bound_but_never_raise_it").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let raised = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": format!("SELECT id FROM {SCHEMA}.tb_doc ORDER BY id"),
"max_rows": 10_000,
"statement_timeout_ms": 600_000
}),
)
.await;
assert_eq!(raised["max_rows"], json!(3), "the ceiling wins: {raised}");
assert_eq!(raised["statement_timeout_ms"], json!(2000), "the ceiling wins: {raised}");
assert_eq!(raised["rows"].as_array().map(Vec::len), Some(3), "{raised}");
let lowered = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": format!("SELECT id FROM {SCHEMA}.tb_doc ORDER BY id"),
"max_rows": 1
}),
)
.await;
assert_eq!(lowered["max_rows"], json!(1), "a tighter request is honoured: {lowered}");
assert_eq!(lowered["rows"].as_array().map(Vec::len), Some(1), "{lowered}");
let (status, body) = run_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": "SELECT 1", "statement_timeout_ms": 0 }),
)
.await;
assert_eq!(status, 400, "a zero timeout must be refused: {body}");
}
#[tokio::test]
async fn the_statement_timeout_cancels_a_long_statement() {
if database_url_or_skip("the_statement_timeout_cancels_a_long_statement").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let (status, body) = run_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": "SELECT pg_sleep(30)", "statement_timeout_ms": 400 }),
)
.await;
assert_eq!(status, 408, "a cancelled statement is a timeout, not a server fault: {body}");
assert!(
body["error"].as_str().unwrap_or_default().contains("statement_timeout"),
"the refusal must name the bound that fired: {body}"
);
}
#[tokio::test]
async fn only_one_statement_is_accepted() {
if database_url_or_skip("only_one_statement_is_accepted").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let (status, body) = run_sql(
&rig.server,
WRITE_TOKEN,
json!({
"sql": format!("SELECT 1; DROP TABLE {SCHEMA}.tb_doc"),
"commit": true
}),
)
.await;
assert_ne!(status, 200, "a two-statement request must not succeed: {body}");
assert_eq!(table_count(&rig.adapter).await, 5, "the second statement must not have run");
let data = ok_sql(&rig.server, READONLY_TOKEN, json!({ "sql": "SELECT 'a;b' AS s" })).await;
assert_eq!(data["rows"], json!([["a;b"]]), "{data}");
}
#[tokio::test]
async fn impersonation_refuses_a_tenant_claim_that_contradicts_the_tenant() {
if database_url_or_skip("impersonation_refuses_a_tenant_claim_that_contradicts_the_tenant")
.is_none()
{
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let scoped = format!("SELECT title FROM {SCHEMA}.v_scoped_doc ORDER BY id");
let (status, body) = run_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": scoped,
"impersonate": {
"user_id": "operator-preview",
"tenant_id": TENANT_A,
"claims": { "tenant_id": TENANT_B }
}
}),
)
.await;
assert_ne!(status, 200, "two tenants for one preview must be refused: {body}");
assert!(body.to_string().contains("different tenants"), "{body}");
let agreeing = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": scoped,
"max_rows": 100,
"impersonate": {
"user_id": "operator-preview",
"tenant_id": TENANT_A,
"claims": { "tenant_id": TENANT_A }
}
}),
)
.await;
assert_eq!(agreeing["rows"], json!([["a-one"], ["a-two"], ["a-three"]]), "{agreeing}");
}
#[tokio::test]
async fn impersonation_applies_the_session_variables() {
if database_url_or_skip("impersonation_applies_the_session_variables").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let scoped = format!("SELECT title FROM {SCHEMA}.v_scoped_doc ORDER BY id");
let as_a = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": scoped,
"max_rows": 100,
"impersonate": { "user_id": "operator-preview", "tenant_id": TENANT_A }
}),
)
.await;
assert_eq!(
as_a["rows"],
json!([["a-one"], ["a-two"], ["a-three"]]),
"tenant A's rows only: {as_a}"
);
let as_b = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": scoped,
"max_rows": 100,
"impersonate": { "user_id": "operator-preview", "tenant_id": TENANT_B }
}),
)
.await;
assert_eq!(as_b["rows"], json!([["b-one"], ["b-two"]]), "tenant B's rows only: {as_b}");
let unscoped =
ok_sql(&rig.server, READONLY_TOKEN, json!({ "sql": scoped, "max_rows": 100 })).await;
assert_eq!(
unscoped["rows"],
json!([]),
"with no impersonation the variable is unset and the view admits nothing: {unscoped}"
);
let read_back = ok_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": format!("SELECT current_setting('{TENANT_VAR}', true) AS v"),
"impersonate": { "user_id": "operator-preview", "tenant_id": TENANT_A }
}),
)
.await;
assert_eq!(read_back["rows"], json!([[TENANT_A]]), "{read_back}");
}
#[tokio::test]
async fn a_reserved_namespace_claim_is_refused() {
if database_url_or_skip("a_reserved_namespace_claim_is_refused").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let (status, body) = run_sql(
&rig.server,
READONLY_TOKEN,
json!({
"sql": "SELECT 1",
"impersonate": {
"user_id": "operator-preview",
"claims": { "fraiseql.actor_type": "system_job" }
}
}),
)
.await;
assert_eq!(status, 400, "{body}");
assert!(
body["error"].as_str().unwrap_or_default().contains("fraiseql.actor_type"),
"the refusal must name the claim: {body}"
);
}
#[tokio::test]
async fn every_execution_lands_in_the_audit_ledger() {
if database_url_or_skip("every_execution_lands_in_the_audit_ledger").is_none() {
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let before = console_entries().len();
let _ = ok_sql(&rig.server, READONLY_TOKEN, json!({ "sql": "SELECT 42 AS answer" })).await;
let _ = run_sql(
&rig.server,
READONLY_TOKEN,
json!({ "sql": format!("UPDATE {SCHEMA}.tb_doc SET title = 'x'") }),
)
.await;
let _ =
run_sql(&rig.server, READONLY_TOKEN, json!({ "sql": "SELECT 1", "commit": true })).await;
let _ = ok_sql(
&rig.server,
WRITE_TOKEN,
json!({
"sql": format!(
"INSERT INTO {SCHEMA}.tb_doc (id, tenant_id, title) \
VALUES (902, '{TENANT_A}', 'audited')"
),
"commit": true
}),
)
.await;
let entries = console_entries();
assert_eq!(
entries.len(),
before + 4,
"all four executions are recorded, not just the ones that worked"
);
let recent = &entries[before..];
assert!(recent[0].success, "the successful read is recorded as a success");
assert!(!recent[1].success, "the database's refusal is recorded as a failure");
assert!(!recent[2].success, "the endpoint's own refusal is recorded too");
assert!(recent[3].success, "the committed write is recorded as a success");
let ctx = recent[0].context.clone().unwrap_or_default();
assert!(ctx.contains("SELECT 42 AS answer"), "the statement is named: {ctx}");
assert!(ctx.contains("sha256="), "and identified beyond the truncation point: {ctx}");
assert_eq!(
recent[0].subject.as_deref(),
Some("admin_readonly_token"),
"the credential that authenticated is named"
);
assert_eq!(
recent[3].subject.as_deref(),
Some("admin_token"),
"and it distinguishes the write credential"
);
assert!(
recent[3].context.clone().unwrap_or_default().contains("committed=true"),
"whether it persisted is the fact the audit is for: {:?}",
recent[3].context
);
assert_eq!(recent[0].secret_type.as_str(), "admin_token");
}
#[tokio::test]
async fn the_console_is_absent_when_there_is_no_section() {
if database_url_or_skip("the_console_is_absent_when_there_is_no_section").is_none() {
return;
}
let rig = Box::pin(boot(None)).await.unwrap();
let (status, _) = run_sql(&rig.server, WRITE_TOKEN, json!({ "sql": "SELECT 1" })).await;
assert_eq!(status, 404, "the console must not be mounted");
let resp = reqwest::Client::new()
.get(format!("{}/api/v1/admin/config", rig.server.url))
.bearer_auth(READONLY_TOKEN)
.send()
.await
.expect("request");
assert_eq!(resp.status(), 200, "the admin API itself is mounted");
}
#[tokio::test]
async fn the_console_is_absent_when_the_section_disables_it() {
if database_url_or_skip("the_console_is_absent_when_the_section_disables_it").is_none() {
return;
}
let rig = Box::pin(boot(Some(AdminSqlConfig {
enabled: false,
..console_config()
})))
.await
.unwrap();
let (status, body) = run_sql(&rig.server, WRITE_TOKEN, json!({ "sql": "SELECT 1" })).await;
assert_eq!(status, 404, "enabled = false must not mount the console: {body}");
}
#[tokio::test]
async fn the_console_refuses_an_unauthenticated_or_wrong_credential() {
if database_url_or_skip("the_console_refuses_an_unauthenticated_or_wrong_credential").is_none()
{
return;
}
let rig = Box::pin(boot_console()).await.unwrap();
let anonymous = reqwest::Client::new()
.post(format!("{}/api/v1/admin/sql", rig.server.url))
.json(&json!({ "sql": "SELECT 1" }))
.send()
.await
.expect("request");
assert_eq!(anonymous.status(), 401, "no Authorization header");
let (status, _) = run_sql(
&rig.server,
"not-the-token-but-long-enough-0123456789",
json!({ "sql": "SELECT 1" }),
)
.await;
assert_eq!(status, 403, "a token matching neither admin credential");
}