From 2d165b3fd8a626fb54e007cb0099dd2654301750 Mon Sep 17 00:00:00 2001 From: Nipunn Koorapati Date: Wed, 16 Sep 2026 11:50:34 -0700 Subject: [PATCH] Drop prepared-statement-cache in postgres upon FEATURE_NOT_SUPPORTED Postgres returns this when a schema change to the database has made the prepared statement no longer valid. W/o this change, the query will repeatedly fail forever as it will try to use the incorrect cached prepared statement. --- sqlx-postgres/src/connection/executor.rs | 35 ++++-- sqlx-postgres/src/error.rs | 2 + tests/postgres/postgres.rs | 147 +++++++++++++++++++++++ 3 files changed, 176 insertions(+), 8 deletions(-) diff --git a/sqlx-postgres/src/connection/executor.rs b/sqlx-postgres/src/connection/executor.rs index 4841a25971..34870451be 100644 --- a/sqlx-postgres/src/connection/executor.rs +++ b/sqlx-postgres/src/connection/executor.rs @@ -1,4 +1,4 @@ -use crate::error::Error; +use crate::error::{error_codes, Error}; use crate::executor::{Execute, Executor}; use crate::io::{PortalId, StatementId}; use crate::logger::QueryLogger; @@ -20,6 +20,23 @@ use sqlx_core::sql_str::SqlStr; use sqlx_core::Either; use std::{pin::pin, sync::Arc}; +/// Detects PostgreSQL's `cached plan must not change result type` error. +/// +/// PostgreSQL: +/// pgJDBC: +fn is_cached_plan_error(error: &Error) -> bool { + error + .as_database_error() + .and_then(|error| error.try_downcast_ref::()) + .is_some_and(|error| { + error.code() == error_codes::FEATURE_NOT_SUPPORTED + && matches!( + error.routine(), + Some("RevalidateCachedQuery" | "RevalidateCachedPlan") + ) + }) +} + async fn prepare( conn: &mut PgConnection, sql: &str, @@ -318,9 +335,9 @@ impl PgConnection { self.invalidate_cached_statement(sql, clear_backend_cache) .await?; - // If we were in transaction mode we can't retry statement, - // so we can immediately return err - if is_in_tx { + // A changed result type is not retried automatically. Invalidating the + // statement makes the next execution heal without hiding this error. + if is_in_tx || clear_backend_cache { return Err(err); } @@ -543,15 +560,17 @@ impl<'c> Executor<'c> for &'c mut PgConnection { // transaction pooling mode // - `Some(true)` - if we should invalidate both backend and frontend caches fn check_stale_plan(error: &Error) -> Option { + if is_cached_plan_error(error) { + return Some(true); + } + let error = error .as_database_error()? .try_downcast_ref::()?; - match (error.code(), error.routine()) { - // "cached plan must not change result type" - ("0A000", Some("RevalidateCachedQuery")) => Some(true), + match error.code() { // DISCARD ALL / DEALLOCATE / pgbouncer - ("26000", _) => Some(false), + "26000" => Some(false), _ => None, } } diff --git a/sqlx-postgres/src/error.rs b/sqlx-postgres/src/error.rs index 7f787a4e5e..00d439b5d0 100644 --- a/sqlx-postgres/src/error.rs +++ b/sqlx-postgres/src/error.rs @@ -236,6 +236,8 @@ impl BackendMessage for PgDatabaseError { /// For reference: pub(crate) mod error_codes { + /// The requested feature is not supported. + pub const FEATURE_NOT_SUPPORTED: &str = "0A000"; /// Caused when a unique or primary key is violated. pub const UNIQUE_VIOLATION: &str = "23505"; /// Caused when a foreign key is violated. diff --git a/tests/postgres/postgres.rs b/tests/postgres/postgres.rs index 0e5f2829c9..f6b053831a 100644 --- a/tests/postgres/postgres.rs +++ b/tests/postgres/postgres.rs @@ -797,6 +797,153 @@ async fn it_caches_statements() -> anyhow::Result<()> { Ok(()) } +#[sqlx_macros::test] +async fn it_clears_cached_statements_after_schema_change() -> anyhow::Result<()> { + let mut conn = new::().await?; + + sqlx::raw_sql( + "CREATE TEMPORARY TABLE statement_cache_test (id INTEGER PRIMARY KEY, value TEXT);\ + INSERT INTO statement_cache_test VALUES (1, 'one')", + ) + .execute(&mut conn) + .await?; + + let query = "SELECT * FROM statement_cache_test WHERE id = $1"; + let row = sqlx::query(query).bind(1_i32).fetch_one(&mut conn).await?; + + assert_eq!(row.columns().len(), 2); + assert_eq!(conn.cached_statements_size(), 1); + + sqlx::raw_sql("ALTER TABLE statement_cache_test DROP COLUMN value") + .execute(&mut conn) + .await?; + + let mut transaction = conn.begin().await?; + let error = sqlx::query(query) + .bind(1_i32) + .fetch_one(&mut *transaction) + .await + .unwrap_err(); + + let error = error + .into_database_error() + .unwrap() + .downcast::(); + + // PostgreSQL reports an invalid cached plan as FEATURE_NOT_SUPPORTED (0A000). + assert_eq!(error.code(), "0A000"); + assert_eq!(error.routine(), Some("RevalidateCachedQuery")); + transaction.rollback().await?; + assert_eq!(conn.cached_statements_size(), 0); + + let row = sqlx::query(query).bind(1_i32).fetch_one(&mut conn).await?; + + assert_eq!(row.columns().len(), 1); + assert_eq!(row.get::("id"), 1); + assert_eq!(conn.cached_statements_size(), 1); + + Ok(()) +} + +#[sqlx_macros::test] +async fn it_does_not_retry_after_cached_statement_schema_change() -> anyhow::Result<()> { + let mut conn = new::().await?; + + sqlx::raw_sql( + "CREATE TEMPORARY TABLE statement_cache_retry_test \ + (id INTEGER PRIMARY KEY, value TEXT);\ + INSERT INTO statement_cache_retry_test VALUES (1, 'one')", + ) + .execute(&mut conn) + .await?; + + let query = "SELECT * FROM statement_cache_retry_test WHERE id = $1"; + let row = sqlx::query(query).bind(1_i32).fetch_one(&mut conn).await?; + + assert_eq!(row.columns().len(), 2); + assert_eq!(conn.cached_statements_size(), 1); + + sqlx::raw_sql("ALTER TABLE statement_cache_retry_test DROP COLUMN value") + .execute(&mut conn) + .await?; + + let error = sqlx::query(query) + .bind(1_i32) + .fetch_one(&mut conn) + .await + .unwrap_err() + .into_database_error() + .unwrap() + .downcast::(); + + assert_eq!(error.code(), "0A000"); + assert_eq!(error.routine(), Some("RevalidateCachedQuery")); + assert_eq!(conn.cached_statements_size(), 0); + + let row = sqlx::query(query).bind(1_i32).fetch_one(&mut conn).await?; + + assert_eq!(row.columns().len(), 1); + assert_eq!(row.get::("id"), 1); + assert_eq!(conn.cached_statements_size(), 1); + + Ok(()) +} + +#[sqlx_macros::test] +async fn it_keeps_cached_statements_after_unrelated_feature_not_supported() -> anyhow::Result<()> { + let mut conn = new::().await?; + + let cached_query = "SELECT $1::INTEGER"; + let value: i32 = sqlx::query_scalar(cached_query) + .bind(1_i32) + .fetch_one(&mut conn) + .await?; + + assert_eq!(value, 1); + assert_eq!(conn.cached_statements_size(), 1); + + sqlx::raw_sql( + r#" + CREATE FUNCTION pg_temp.raise_feature_not_supported(value INTEGER) + RETURNS INTEGER + LANGUAGE plpgsql + AS $$ + BEGIN + RAISE EXCEPTION 'unrelated unsupported feature' USING ERRCODE = '0A000'; + END + $$ + "#, + ) + .execute(&mut conn) + .await?; + + let error = sqlx::query("SELECT pg_temp.raise_feature_not_supported($1)") + .bind(1_i32) + .execute(&mut conn) + .await + .unwrap_err() + .into_database_error() + .unwrap() + .downcast::(); + + assert_eq!(error.code(), "0A000"); + assert!(!matches!( + error.routine(), + Some("RevalidateCachedQuery" | "RevalidateCachedPlan") + )); + assert_eq!(conn.cached_statements_size(), 2); + + let value: i32 = sqlx::query_scalar(cached_query) + .bind(2_i32) + .fetch_one(&mut conn) + .await?; + + assert_eq!(value, 2); + assert_eq!(conn.cached_statements_size(), 2); + + Ok(()) +} + #[sqlx_macros::test] async fn it_closes_statement_from_cache_issue_470() -> anyhow::Result<()> { sqlx_test::setup_if_needed();