Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,25 @@
# Changelog

## Unreleased

### Added

- `execute_query` and `execute_query_batch` accept a `session_id` and keep
that session's connection when the statement or batch leaves an explicit
transaction open, so `BEGIN`, the changes, a verifying `SELECT` and
`COMMIT` can each be their own run instead of having to be one script.
With a `session_id` the replies carry `{ result, in_transaction }` and
`{ results, in_transaction }`; without one they stay the bare
`QueryResult` and array, so older hosts are unaffected.
- `release_session` RPC method, called when the owning editor tab closes.

### Fixed

- A batch that left a transaction open returned its connection to the pool
as-is. The pool recycles with `RecyclingMethod::Fast`, which resets
nothing, so the next borrower inherited the open transaction and its
locks. A connection is now always rolled back before it goes back.

## [1.0.0-rc.4] - 2026-09-17

### Added
Expand Down
3 changes: 2 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,8 @@ instead of reconnecting.
| `get_materialized_views` / `refresh_materialized_view` | Materialized view lifecycle |
| `get_routines` / `get_routine_parameters` / `get_routine_definition` | Function/procedure metadata |
| `get_triggers` / `get_trigger_definition` / `create_trigger` / `drop_trigger` | Trigger lifecycle |
| `execute_query` / `execute_query_batch` / `explain_query` | Query execution, multi-statement batches, and query plans |
| `execute_query` / `execute_query_batch` / `explain_query` | Query execution, multi-statement batches, and query plans. Both query methods accept an optional `session_id` that keeps the connection while an explicit transaction is open |
| `release_session` | Rolls back and releases the connection a session pinned for an open transaction |
| `insert_record` / `update_record` / `delete_record` | Row-level CRUD with type-aware value binding |
| `get_create_table_sql` / `get_add_column_sql` / `get_alter_column_sql` / `get_create_index_sql` / `get_create_foreign_key_sql` / `drop_index` / `drop_foreign_key` | DDL generation and execution |
| `save_blob_to_file` / `fetch_blob_as_data_url` | BLOB (`bytea`) export and preview |
Expand Down
4 changes: 4 additions & 0 deletions src/handlers/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,5 +40,9 @@ pub async fn test_connection(id: Value, params: &Value) -> Value {

/// Graceful shutdown — drain pools and exit.
pub async fn shutdown(id: Value, _params: &Value) -> Value {
// Any session still holding a connection has an open transaction. The
// pool recycles without resetting, so roll them back rather than let
// the server keep them open until it times the backend out.
crate::session::release_all().await;
ok_response(id, Value::Null)
}
263 changes: 243 additions & 20 deletions src/handlers/query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use crate::client;
use crate::extract::extract_value;
use crate::models::{inner_params, ConnectionParams};
use crate::rpc::{error_response, ok_response};
use crate::session;

pub async fn execute_query(id: Value, params: &Value) -> Value {
let conn_params = ConnectionParams::from_value(inner_params(params));
Expand All @@ -18,6 +19,37 @@ pub async fn execute_query(id: Value, params: &Value) -> Value {
.map(|v| v as u32);
let page = params.get("page").and_then(Value::as_u64).unwrap_or(1) as u32;
let schema = params.get("schema").and_then(Value::as_str);
let session_id = params.get("session_id").and_then(Value::as_str);

// One statement at a time is how a transaction is actually driven -
// BEGIN, the changes, a verifying SELECT, COMMIT - so this path needs
// the same pinning as a batch. Implemented as a one-statement batch so
// there is a single session code path.
if session_id.is_some() {
let queries = [query.to_string()];
return match run_batch_in_session(&conn_params, &queries, limit, page, schema, session_id)
.await
{
Ok((mut results, in_transaction)) => {
let statement = results.pop().unwrap_or(Value::Null);
match statement.get("error").and_then(Value::as_str) {
// Still a reply, not an RPC error, so the host learns the state a failed COMMIT left.
Some(error) => ok_response(
id,
json!({ "error": error, "in_transaction": in_transaction }),
),
None => {
let result = statement.get("result").cloned().unwrap_or(Value::Null);
ok_response(
id,
json!({ "result": result, "in_transaction": in_transaction }),
)
}
}
}
Err(e) => error_response(id, -32603, &e),
};
}

match exec_query(&conn_params, query, limit, page, schema).await {
Ok(result) => ok_response(id, result),
Expand All @@ -42,41 +74,87 @@ pub async fn execute_query_batch(id: Value, params: &Value) -> Value {
.map(|v| v as u32);
let page = params.get("page").and_then(Value::as_u64).unwrap_or(1) as u32;
let schema = params.get("schema").and_then(Value::as_str);
// The host sends the editor tab's id. A batch that leaves a transaction
// open keeps its connection under this key so the next batch from the
// same tab continues that transaction.
let session_id = params.get("session_id").and_then(Value::as_str);

match run_batch_in_session(&conn_params, &queries, limit, page, schema, session_id).await {
// Only a session-aware call gets the richer shape; a host that did
// not send a session_id still receives the bare array it expects.
Ok((results, in_transaction)) => match session_id {
Some(_) => ok_response(
id,
json!({ "results": results, "in_transaction": in_transaction }),
),
None => ok_response(id, json!(results)),
},
Err(e) => error_response(id, -32603, &e),
}
}

// Acquire ONE connection for the entire batch (session state must survive)
let pool = match client::build_pool_pub(&conn_params).await {
Ok(p) => p,
Err(e) => return error_response(id, -32603, &e),
/// Run `queries` on one connection, reusing the session's pinned one when
/// it left a transaction open and re-pinning it if one is still open after.
///
/// Returns one per-statement result object per query — the same shape
/// `execute_query_batch` replies with — plus whether the session is still
/// inside a transaction. Shared with `execute_query` so a single statement
/// and a batch take the same session path.
///
/// A batch that leaves a transaction open with no session to pin it to is
/// rolled back: the pool recycles with `RecyclingMethod::Fast`, so the next
/// borrower would otherwise inherit the transaction and its locks.
async fn run_batch_in_session(
conn_params: &ConnectionParams,
queries: &[String],
limit: Option<u32>,
page: u32,
schema: Option<&str>,
session_id: Option<&str>,
) -> Result<(Vec<Value>, bool), String> {
// Held for the whole run, so an overlapping call for the same session waits its turn.
let mut slot = match session_id {
Some(sid) => Some(session::lock(sid).await),
None => None,
};
let pg_client = match pool.get().await {
Ok(c) => c,
Err(e) => {
return error_response(
id,
-32603,
&format!("Connection failed: {}", client::format_pool_error(&e)),
)
let pinned = slot.as_mut().and_then(|s| s.take());
// A pinned connection only exists because its transaction is still open.
let mut in_transaction = pinned.is_some();
let reused = pinned.is_some();

// Acquire ONE connection for the entire batch (session state must survive)
let pg_client = match pinned {
Some(client) => client,
None => {
let pool = client::build_pool_pub(conn_params).await?;
pool.get()
.await
.map_err(|e| format!("Connection failed: {}", client::format_pool_error(&e)))?
}
};

if let Some(s) = schema {
// Applying it to a reused connection would run inside the open
// transaction and change what the rest of it sees.
if let Some(s) = schema.filter(|_| !reused) {
let set_path = format!("SET search_path TO \"{}\"", s.replace('"', "\"\""));
if let Err(e) = pg_client.batch_execute(&set_path).await {
return error_response(
id,
-32603,
&format!("Failed to set search_path: {}", client::format_pg_error(&e)),
);
return Err(format!(
"Failed to set search_path: {}",
client::format_pg_error(&e)
));
}
}

let mut results: Vec<Value> = Vec::new();

for query in &queries {
for query in queries {
let start = Instant::now();
let outcome = exec_query_on_client(&pg_client, query, limit, page).await;
let elapsed_ms = start.elapsed().as_secs_f64() * 1000.0;

in_transaction =
transaction_effect(query).in_transaction_after(outcome.is_ok(), in_transaction);

match outcome {
Ok(result) => results.push(json!({
"result": result,
Expand All @@ -91,7 +169,28 @@ pub async fn execute_query_batch(id: Value, params: &Value) -> Value {
}
}

ok_response(id, json!(results))
match slot.as_mut() {
Some(slot) if in_transaction => slot.pin(pg_client),
_ => {
if in_transaction {
session::rollback_and_release(pg_client).await;
}
}
}

Ok((results, in_transaction))
}

/// Roll back and release the connection pinned to a session. Called when
/// the owning editor tab closes.
pub async fn release_session(id: Value, params: &Value) -> Value {
match params.get("session_id").and_then(Value::as_str) {
Some(session_id) => {
session::release(session_id).await;
ok_response(id, json!({ "released": true }))
}
None => error_response(id, -32602, "session_id is required"),
}
}

pub async fn explain_query(id: Value, params: &Value) -> Value {
Expand Down Expand Up @@ -297,6 +396,130 @@ async fn exec_query_on_client(
}))
}

/// What a statement does to the surrounding transaction.
///
/// Decides whether the connection a batch ran on must be kept for the next
/// batch of the same session (an explicit transaction is still open) or may
/// go back to the pool.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TransactionEffect {
/// Opens an explicit transaction (`BEGIN`, `START TRANSACTION`).
Opens,
/// Closes the current transaction (`COMMIT`, `ROLLBACK`, `END`).
///
/// `ROLLBACK TO SAVEPOINT` does not close it and is classified as
/// [`TransactionEffect::None`].
Closes,
/// Ends the current transaction and opens a new one (`COMMIT AND CHAIN`).
Chains,
/// Leaves the transaction state as it was.
None,
}

impl TransactionEffect {
/// Whether a transaction is open after a statement with this effect ran.
pub fn in_transaction_after(self, succeeded: bool, before: bool) -> bool {
match (self, succeeded) {
(TransactionEffect::Opens | TransactionEffect::Chains, true) => true,
// PostgreSQL ends the transaction even when COMMIT itself fails, e.g. on a deferred constraint.
(TransactionEffect::Closes, _) | (TransactionEffect::Chains, false) => false,
_ => before,
}
}
}

/// Classify a statement's effect on the transaction state.
///
/// Only the leading keywords are inspected, so a `BEGIN` inside a string
/// literal or a later clause cannot be mistaken for transaction control. A
/// PL/pgSQL `BEGIN … END` body is not a concern: it arrives inside a `DO`
/// or `CREATE FUNCTION` statement, whose leading keyword is neither.
pub fn transaction_effect(query: &str) -> TransactionEffect {
// Six cover the longest form, `ROLLBACK TRANSACTION AND NO CHAIN`.
let mut words = leading_keywords(query, 6).into_iter();
let Some(first) = words.next() else {
return TransactionEffect::None;
};
let mut rest: Vec<String> = words.collect();
// `WORK` / `TRANSACTION` after a transaction-control verb is optional noise.
if matches!(first.as_str(), "COMMIT" | "END" | "ROLLBACK" | "ABORT")
&& matches!(
rest.first().map(String::as_str),
Some("WORK" | "TRANSACTION")
)
{
rest.remove(0);
}
let second = rest.first().map(String::as_str);
let chains = rest.windows(2).any(|p| p[0] == "AND" && p[1] == "CHAIN");

match first.as_str() {
// `BEGIN` alone, `BEGIN TRANSACTION`, `BEGIN ISOLATION LEVEL …`.
"BEGIN" => TransactionEffect::Opens,
"START" if second == Some("TRANSACTION") => TransactionEffect::Opens,
// Two-phase commit acts on a prepared transaction, not this session's.
"COMMIT" | "ROLLBACK" if second == Some("PREPARED") => TransactionEffect::None,
// `ROLLBACK [WORK] TO [SAVEPOINT] x` unwinds to a savepoint and leaves the transaction open.
"ROLLBACK" if second == Some("TO") => TransactionEffect::None,
"COMMIT" | "END" | "ROLLBACK" | "ABORT" if chains => TransactionEffect::Chains,
"COMMIT" | "END" | "ROLLBACK" | "ABORT" => TransactionEffect::Closes,
// `PREPARE TRANSACTION` dissociates the transaction from the session.
"PREPARE" if second == Some("TRANSACTION") => TransactionEffect::Closes,
_ => TransactionEffect::None,
}
}

/// The first `n` keywords of `query`, uppercased. Every `--` and (nested)
/// `/* */` comment is skipped, wherever it appears, and reading stops at a
/// string literal, dollar quote or `;`, so no data is read as a keyword.
fn leading_keywords(query: &str, n: usize) -> Vec<String> {
let mut words = Vec::new();
let mut word = String::new();
let mut chars = query.chars().peekable();
let mut comment_depth = 0usize;
while let Some(c) = chars.next() {
if comment_depth > 0 {
if c == '*' && chars.peek() == Some(&'/') {
chars.next();
comment_depth -= 1;
} else if c == '/' && chars.peek() == Some(&'*') {
chars.next();
comment_depth += 1;
}
continue;
}
if c.is_ascii_alphanumeric() || c == '_' {
word.push(c.to_ascii_uppercase());
continue;
}
if !word.is_empty() {
words.push(std::mem::take(&mut word));
if words.len() == n {
return words;
}
}
match (c, chars.peek()) {
('-', Some('-')) => {
for c in chars.by_ref() {
if c == '\n' {
break;
}
}
}
('/', Some('*')) => {
chars.next();
comment_depth = 1;
}
('\'' | '"' | '$' | ';', _) => return words,
_ => {}
}
}
if !word.is_empty() && words.len() < n {
words.push(word);
}
words
}

/// Strip leading SQL comments (`-- …` line comments and `/* … */` block
/// comments) and whitespace so the first statement keyword is at position 0.
/// Matches the builtin driver's `drivers/common/query.rs::strip_leading_sql_comments`
Expand Down
Loading
Loading