use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
use pylon_client::{Client, DecodedValue, Isolation, Value};
use pylon_core::export::export_schema;
use pylon_core::schema::{
ChannelDescriptor, ChannelPayload, PropertyDescriptor, SchemaDescriptor, TriggerDescriptor, TypeDescriptor,
};
fn test_dsn() -> String {
std::env::var("PYLON_PGCON_TEST_DSN").expect("PYLON_PGCON_TEST_DSN must be set to run live-Postgres tests")
}
fn unique_module(prefix: &str) -> String {
static COUNTER: AtomicU64 = AtomicU64::new(0);
let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
let seq = COUNTER.fetch_add(1, Ordering::Relaxed);
format!("{prefix}_{nanos}_{seq}")
}
fn id_prop() -> PropertyDescriptor {
PropertyDescriptor {
name: "id".into(),
pg_type: "uuid".into(),
nullable: false,
default_sql: Some("gen_random_uuid()".into()),
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: true,
is_pk: true,
is_readonly: true,
rewrites: vec![],
tuple_members: None,
column_type: None,
}
}
fn text_prop(name: &str) -> PropertyDescriptor {
PropertyDescriptor {
name: name.into(),
pg_type: "text".into(),
nullable: false,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
}
}
fn float_prop(name: &str) -> PropertyDescriptor {
PropertyDescriptor {
name: name.into(),
pg_type: "float8".into(),
nullable: false,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
}
}
fn person_schema(module: &str) -> SchemaDescriptor {
SchemaDescriptor {
types: vec![TypeDescriptor {
name: "Person".into(),
module: module.into(),
table: "Person".into(),
abstract_: false,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![id_prop(), text_prop("name")],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}],
..Default::default()
}
}
async fn setup(schema: &SchemaDescriptor) -> Client {
let ddl = export_schema(schema).unwrap();
let pool = pylon_pgcon::PgPool::connect(&test_dsn(), 5).await.unwrap();
pool.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
pylon_core::migrate::ensure_internal_schema(&pool).await.unwrap();
pool.batch_execute(&pylon_core::stdlib::export_stdlib()).await.unwrap();
pool.batch_execute(&ddl).await.unwrap();
pylon_core::migrate::write_schema_snapshot(&pool, &serde_json::to_string(schema).unwrap())
.await
.unwrap();
connected(Client::builder(test_dsn()).max_pool_size(5).build().unwrap()).await
}
async fn connected(client: Client) -> Client {
client.ensure_connected().await.unwrap();
client
}
async fn setup_with_cache(schema: &SchemaDescriptor) -> Client {
let ddl = export_schema(schema).unwrap();
let pool = pylon_pgcon::PgPool::connect(&test_dsn(), 5).await.unwrap();
pool.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
pylon_core::migrate::ensure_internal_schema(&pool).await.unwrap();
pool.batch_execute(&pylon_core::stdlib::export_stdlib()).await.unwrap();
pool.batch_execute(&ddl).await.unwrap();
pylon_core::migrate::write_schema_snapshot(&pool, &serde_json::to_string(schema).unwrap())
.await
.unwrap();
let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
let cache_dir = std::env::temp_dir().join(format!("pylon-client-live-test-cache-{nanos}"));
connected(
Client::builder(test_dsn())
.max_pool_size(5)
.cache(cache_dir, 10)
.build()
.unwrap(),
)
.await
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn query_and_execute_round_trip() {
let module = unique_module("live_client_basic");
let client = setup(&person_schema(&module)).await;
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Alice".into()))],
)
.await
.unwrap();
let rows = client
.query::<Value, _>(&format!("select {module}::Person {{ name }}"), &[])
.await
.unwrap();
assert_eq!(rows.len(), 1);
let Value::Object(person) = &rows[0] else {
panic!("expected Object, got {:?}", rows[0])
};
assert_eq!(person.get("name"), Some(&Value::Str("Alice".into())));
assert_eq!(person.type_name(), Some(format!("{module}::Person").as_str()));
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn json_names_every_field_and_runs_the_query_once() {
let module = unique_module("live_client_json");
let client = setup(&person_schema(&module)).await;
let inserted = client
.query_single_json(
&format!("select (insert {module}::Person {{ name := 'Dora' }}) {{ name }}"),
&[],
)
.await
.unwrap();
assert_eq!(inserted.as_deref(), Some(r#"{"name": "Dora"}"#));
let count = client
.query_json(&format!("select count({module}::Person)"), &[])
.await
.unwrap();
assert_eq!(count, "[1]");
let limited = client
.query_single_json(&format!("select {module}::Person {{ name }} limit 1"), &[])
.await
.unwrap();
assert_eq!(limited.as_deref(), Some(r#"{"name": "Dora"}"#));
let all = client
.query_json(&format!("select {module}::Person {{ name }}"), &[])
.await
.unwrap();
assert_eq!(all, r#"[{"name": "Dora"}]"#);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn query_single_enforces_cardinality() {
let module = unique_module("live_client_single");
let client = setup(&person_schema(&module)).await;
assert_eq!(
client
.query_single::<Value, _>(&format!("select {module}::Person"), &[])
.await
.unwrap(),
None
);
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Bob".into()))],
)
.await
.unwrap();
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Carol".into()))],
)
.await
.unwrap();
let err = client
.query_single::<Value, _>(&format!("select {module}::Person"), &[])
.await
.unwrap_err();
assert!(
matches!(err, pylon_client::Error::ResultCardinality { got: 2 }),
"got: {err:?}"
);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn globals_fill_the_dunder_global_param_slot() {
let mut schema = person_schema(&unique_module("live_client_globals"));
schema.globals.push(pylon_core::schema::GlobalDescriptor {
name: "viewer_name".into(),
module: schema.types[0].module.clone(),
scalar_type: "text".into(),
required: false,
default_expr: None,
computed_expr: None,
});
let module = schema.types[0].module.clone();
let client = setup(&schema).await;
let authed = client.with_globals([(format!("{module}::viewer_name"), DecodedValue::Str("Dave".into()))]);
let rows = authed
.query::<Value, _>("select global viewer_name", &[])
.await
.unwrap();
assert_eq!(rows, vec![Value::Str("Dave".into())]);
let rows = client
.query::<Value, _>("select global viewer_name", &[])
.await
.unwrap();
assert_eq!(rows, vec![Value::Null]);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn transaction_commits_on_success() {
let module = unique_module("live_client_tx");
let client = setup(&person_schema(&module)).await;
client
.transaction(Isolation::Serializable, |tx| {
let module = module.clone();
Box::pin(async move {
tx.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Erin".into()))],
)
.await
})
})
.await
.unwrap();
let rows = client
.query::<Value, _>(&format!("select {module}::Person {{ name }}"), &[])
.await
.unwrap();
assert_eq!(rows.len(), 1);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn transaction_rolls_back_on_error_and_does_not_retry_non_retriable_errors() {
let module = unique_module("live_client_tx_rollback");
let client = setup(&person_schema(&module)).await;
let mut attempts = 0;
let result: Result<(), pylon_client::Error> = client
.transaction(Isolation::Serializable, |tx| {
attempts += 1;
let module = module.clone();
Box::pin(async move {
tx.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Frank".into()))],
)
.await?;
tx.execute(&format!("select {module}::Person.nonexistent_link"), &[])
.await
})
})
.await;
assert!(result.is_err());
assert_eq!(attempts, 1, "a non-retriable error must not be retried");
let rows = client
.query::<Value, _>(&format!("select {module}::Person"), &[])
.await
.unwrap();
assert!(rows.is_empty(), "the insert must have been rolled back, got: {rows:?}");
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn a_body_returning_rollback_sees_its_own_writes_and_leaves_nothing_behind() {
let module = unique_module("live_client_tx_deliberate_rollback");
let client = setup(&person_schema(&module)).await;
let mut attempts = 0;
let committed: Option<()> = client
.transaction_opt(Isolation::Serializable, |tx| {
attempts += 1;
let module = module.clone();
Box::pin(async move {
tx.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Ghost".into()))],
)
.await?;
let seen = tx
.query::<Value, _>(&format!("select {module}::Person {{ name }}"), &[])
.await?;
assert_eq!(seen.len(), 1);
Err(pylon_client::Error::Rollback)
})
})
.await
.unwrap();
assert!(committed.is_none());
assert_eq!(attempts, 1, "a deliberate rollback must not be retried");
let rows = client
.query::<Value, _>(&format!("select {module}::Person"), &[])
.await
.unwrap();
assert!(rows.is_empty(), "the insert must have been rolled back, got: {rows:?}");
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn rollback_still_propagates_as_an_error_through_plain_transaction() {
let module = unique_module("live_client_tx_rollback_err");
let client = setup(&person_schema(&module)).await;
let result: Result<(), pylon_client::Error> = client
.transaction(Isolation::Serializable, |tx| {
let module = module.clone();
Box::pin(async move {
tx.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Ghost".into()))],
)
.await?;
Err(pylon_client::Error::Rollback)
})
})
.await;
let err = result.expect_err("Rollback must reach the caller here");
assert!(err.is_rollback());
assert!(!err.is_retriable());
let rows = client
.query::<Value, _>(&format!("select {module}::Person"), &[])
.await
.unwrap();
assert!(rows.is_empty(), "the insert must have been rolled back, got: {rows:?}");
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn cached_query_serves_stale_data_until_something_else_invalidates_it() {
let module = unique_module("live_client_cache");
let client = setup_with_cache(&person_schema(&module)).await;
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("Alice".into()))],
)
.await
.unwrap();
let query = format!("select {module}::Person {{ name }}");
let first = client.query::<Value, _>(&query, &[]).await.unwrap();
assert_eq!(first.len(), 1);
let Value::Object(person) = &first[0] else {
panic!("expected Object")
};
assert_eq!(person.get("name"), Some(&Value::Str("Alice".into())));
client
.raw_connection()
.await
.unwrap()
.execute_typed(&format!("UPDATE \"{module}\".\"Person\" SET name = 'Mutated'"), &[])
.await
.unwrap();
let second = client.query::<Value, _>(&query, &[]).await.unwrap();
assert_eq!(
second, first,
"a read-through cache hit must return the stale cached value"
);
let stats = client.cache_stat().unwrap().expect("cache was configured");
assert!(stats.entry_count >= 1);
client.cache_clear().unwrap();
let stats_after_clear = client.cache_stat().unwrap().unwrap();
assert_eq!(stats_after_clear.entry_count, 0);
let third = client.query::<Value, _>(&query, &[]).await.unwrap();
let Value::Object(person) = &third[0] else {
panic!("expected Object")
};
assert_eq!(person.get("name"), Some(&Value::Str("Mutated".into())));
}
async fn recv_with_timeout(listener: &mut pylon_client::ChannelListener) -> pylon_client::Result<Value> {
tokio::time::timeout(std::time::Duration::from_secs(5), listener.recv())
.await
.expect("timed out waiting for a notification")
.expect("listener closed with no notification")
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn listen_decodes_a_scalar_channel_payload() {
let module = unique_module("live_client_listen_scalar");
let mut schema = person_schema(&module);
schema.types[0].triggers = vec![TriggerDescriptor {
on: 1, timing: "After".into(),
handler: "select notify(Pings, __new__.name)".into(),
}];
schema.channels = vec![ChannelDescriptor {
name: "Pings".into(),
module: module.clone(),
wire_name: format!("{module}__pings"),
payload: ChannelPayload::Scalar("text".into()),
description: None,
}];
let client = setup(&schema).await;
let mut listener = client.listen("Pings").await.unwrap();
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("gadget".into()))],
)
.await
.unwrap();
assert_eq!(
recv_with_timeout(&mut listener).await.unwrap(),
Value::Str("gadget".to_string())
);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn listen_decodes_a_type_channel_payload_as_the_rows_id() {
let module = unique_module("live_client_listen_type");
let mut schema = person_schema(&module);
schema.types[0].triggers = vec![TriggerDescriptor {
on: 1, timing: "After".into(),
handler: "select notify(PersonUpdates, __new__)".into(),
}];
schema.channels = vec![ChannelDescriptor {
name: "PersonUpdates".into(),
module: module.clone(),
wire_name: format!("{module}__person_updates"),
payload: ChannelPayload::Type(format!("{module}::Person")),
description: None,
}];
let client = setup(&schema).await;
let mut listener = client.listen("PersonUpdates").await.unwrap();
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name }}"),
&[("name", DecodedValue::Str("gadget".into()))],
)
.await
.unwrap();
let payload = recv_with_timeout(&mut listener).await.unwrap();
let Value::Uuid(id) = payload else {
panic!("expected Uuid, got {payload:?}")
};
let rows = client
.query::<Value, _>(&format!("select {module}::Person {{ id }}"), &[])
.await
.unwrap();
let Value::Object(person) = &rows[0] else {
panic!("expected Object")
};
assert_eq!(person.get("id"), Some(&Value::Uuid(id)));
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn listen_decodes_an_object_channel_payload() {
let module = unique_module("live_client_listen_object");
let mut schema = person_schema(&module);
schema.types[0].properties.push(float_prop("score"));
schema.types[0].triggers = vec![TriggerDescriptor {
on: 1, timing: "After".into(),
handler: "select notify(PersonReady, { name := __new__.name, score := __new__.score })".into(),
}];
schema.channels = vec![ChannelDescriptor {
name: "PersonReady".into(),
module: module.clone(),
wire_name: format!("{module}__person_ready"),
payload: ChannelPayload::Object(vec![("name".into(), "text".into()), ("score".into(), "float8".into())]),
description: None,
}];
let client = setup(&schema).await;
let mut listener = client.listen("PersonReady").await.unwrap();
client
.execute(
&format!("insert {module}::Person {{ name := <str>$name, score := <float64>$score }}"),
&[
("name", DecodedValue::Str("gadget".into())),
("score", DecodedValue::F64(0.75)),
],
)
.await
.unwrap();
let payload = recv_with_timeout(&mut listener).await.unwrap();
let Value::Object(obj) = payload else {
panic!("expected Object, got {payload:?}")
};
assert_eq!(obj.get("name"), Some(&Value::Str("gadget".to_string())));
assert_eq!(obj.get("score"), Some(&Value::Float64(0.75)));
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn listen_raises_on_a_malformed_payload() {
let module = unique_module("live_client_listen_malformed");
let mut schema = person_schema(&module);
schema.channels = vec![ChannelDescriptor {
name: "Ids".into(),
module: module.clone(),
wire_name: format!("{module}__ids"),
payload: ChannelPayload::Type(format!("{module}::Person")),
description: None,
}];
let client = setup(&schema).await;
let mut listener = client.listen("Ids").await.unwrap();
client
.raw_connection()
.await
.unwrap()
.execute_typed(&format!("SELECT pg_notify('{module}__ids', 'not-a-uuid')"), &[])
.await
.unwrap();
let err = recv_with_timeout(&mut listener)
.await
.expect_err("expected a malformed-payload error");
assert!(matches!(err, pylon_client::Error::MalformedPayload(_)), "got: {err:?}");
}