mod common;
use common::*;
use pylon_core::export::export_schema;
use pylon_core::query;
use pylon_core::schema::{PropertyDescriptor, SchemaDescriptor, TypeDescriptor};
use pylon_pgcon::ExtensionOids;
use pylon_value::DecodedValue;
fn ty(name: &str, module: &str, properties: Vec<PropertyDescriptor>) -> TypeDescriptor {
TypeDescriptor {
name: name.into(),
module: module.into(),
table: name.into(),
abstract_: false,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties,
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}
}
fn int_prop(name: &str) -> PropertyDescriptor {
let mut p = text_prop(name);
p.pg_type = "int8".into();
p
}
fn job_schema(module: &str) -> SchemaDescriptor {
let job = ty(
"Job",
module,
vec![id_prop(), text_prop("status"), int_prop("priority")],
);
SchemaDescriptor {
types: vec![job],
..Default::default()
}
}
async fn bootstrap(pool: &pylon_pgcon::PgPool, sd: &SchemaDescriptor) {
pool.batch_execute(&pylon_core::stdlib::export_stdlib()).await.unwrap();
pool.batch_execute(&export_schema(sd).unwrap()).await.unwrap();
}
async fn exec(pool: &pylon_pgcon::PgPool, sd: &SchemaDescriptor, pyql: &str) {
let compiled = query::compile(pyql, sd).unwrap();
pool.execute_typed(&compiled.sql, &[]).await.unwrap();
}
fn field(row: &DecodedValue, i: usize) -> &DecodedValue {
match row {
DecodedValue::Composite(fields) => fields.get(i).unwrap_or(&DecodedValue::Null),
other => panic!("expected a Composite-shaped row, got {other:?}"),
}
}
fn as_i64(v: &DecodedValue) -> i64 {
match v {
DecodedValue::I64(n) => *n,
other => panic!("expected I64, got {other:?}"),
}
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn skip_locked_lets_a_second_transaction_claim_a_different_row() {
let module = unique_module("live_lock_skip");
let sd = job_schema(&module);
let pool = test_pool().await;
bootstrap(&pool, &sd).await;
exec(
&pool,
&sd,
&format!("insert {module}::Job {{ status := 'pending', priority := 1 }}"),
)
.await;
exec(
&pool,
&sd,
&format!("insert {module}::Job {{ status := 'pending', priority := 2 }}"),
)
.await;
let pick_sql = query::compile(
&format!(
"select {module}::Job {{ priority }} filter .status = 'pending' \
order by .priority asc limit 1 for update skip locked"
),
&sd,
)
.unwrap()
.sql
.clone();
let tx1 = pool.begin_default().await.unwrap();
let rows1 = tx1
.query_typed(&pick_sql, &[], &ExtensionOids::default())
.await
.unwrap();
assert_eq!(rows1.len(), 1);
assert_eq!(
as_i64(field(&rows1[0], 2)),
1,
"tx1 should have claimed the lowest-priority pending job"
);
let tx2 = pool.begin_default().await.unwrap();
let rows2 = tx2
.query_typed(&pick_sql, &[], &ExtensionOids::default())
.await
.unwrap();
assert_eq!(rows2.len(), 1);
assert_eq!(
as_i64(field(&rows2[0], 2)),
2,
"tx2 must skip tx1's locked row and claim the other one"
);
let tx3 = pool.begin_default().await.unwrap();
let rows3 = tx3
.query_typed(&pick_sql, &[], &ExtensionOids::default())
.await
.unwrap();
assert!(rows3.is_empty(), "both pending jobs are already locked, got {rows3:?}");
tx3.commit().await.unwrap();
tx1.commit().await.unwrap();
tx2.commit().await.unwrap();
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn nowait_fails_immediately_instead_of_blocking_on_a_locked_row() {
let module = unique_module("live_lock_nowait");
let sd = job_schema(&module);
let pool = test_pool().await;
bootstrap(&pool, &sd).await;
exec(
&pool,
&sd,
&format!("insert {module}::Job {{ status := 'pending', priority := 1 }}"),
)
.await;
let claim_sql = query::compile(
&format!("select {module}::Job {{ priority }} filter .status = 'pending' for update"),
&sd,
)
.unwrap()
.sql
.clone();
let claim_nowait_sql = query::compile(
&format!("select {module}::Job {{ priority }} filter .status = 'pending' for update nowait"),
&sd,
)
.unwrap()
.sql
.clone();
let tx1 = pool.begin_default().await.unwrap();
let rows1 = tx1
.query_typed(&claim_sql, &[], &ExtensionOids::default())
.await
.unwrap();
assert_eq!(rows1.len(), 1, "tx1 should have locked the only pending job");
let tx2 = pool.begin_default().await.unwrap();
let err = tx2
.query_typed(&claim_nowait_sql, &[], &ExtensionOids::default())
.await
.unwrap_err();
assert_eq!(
err.sqlstate(),
Some(&tokio_postgres::error::SqlState::LOCK_NOT_AVAILABLE),
"expected NOWAIT's own lock_not_available error, got: {err:?}",
);
tx1.commit().await.unwrap();
tx2.rollback().await.unwrap();
}