Skip to content
Merged
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
10 changes: 7 additions & 3 deletions pgdog/src/admin/show_prepared_statements.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ impl Command for ShowPreparedStatements {
];
for (key, stmt) in statements.statements() {
let name = stmt.name();
let rewrite = statements.rewritten_parse(&name).ok_or(Error::Empty)?;
let rewrite = statements.rewritten_parse(&name);
let rewritten = statements.is_rewritten(&name);
let name_memory = statements
.names()
Expand All @@ -43,8 +43,12 @@ impl Command for ShowPreparedStatements {
let mut dr = DataRow::new();
dr.add(stmt.name())
.add(key.query()?)
.add(if rewritten {
rewrite.query().to_data_row_column()
.add(if let Some(rewrite) = rewrite {
if rewritten {
rewrite.query().to_data_row_column()
} else {
Data::null()
}
} else {
Data::null()
})
Expand Down
2 changes: 1 addition & 1 deletion pgdog/src/backend/pool/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -767,7 +767,7 @@ impl Cluster {
) -> Result<Vec<T>, crate::backend::Error> {
let shard = self
.shards
.get(round_robin::next() % self.shards.len().max(1))
.get(round_robin::next(self.shards.len().max(1)))
.ok_or(crate::backend::pool::Error::NoDatabases)?;
let mut server = shard.primary_or_replica(&Request::default()).await?;

Expand Down
3 changes: 1 addition & 2 deletions pgdog/src/backend/pool/connection/mirror/handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -201,7 +201,6 @@ mod tests {
use super::*;
use crate::backend::pool::ClusterMetrics;
use parking_lot::Mutex;
use pgdog_config::QueryParserEngine;
use std::sync::Arc;
use tokio::sync::mpsc::{Receiver, channel};

Expand Down Expand Up @@ -500,7 +499,7 @@ mod tests {

fn request_with_ast(query: &str) -> ClientRequest {
use crate::frontend::router::Ast;
let ast = Ast::new_record(query, QueryParserEngine::PgQueryProtobuf).unwrap();
let ast = Ast::new_record(query).unwrap();
ClientRequest {
ast: Some(ast),
..Default::default()
Expand Down
45 changes: 41 additions & 4 deletions pgdog/src/backend/prepared_statements.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,16 @@ use std::{
time::{Duration, Instant},
};

use crate::util::time::deadline;
use crate::{
frontend::{self, prepared_statements::GlobalCache},
net::{
Close, CloseComplete, FromBytes, Message, ParseComplete, Protocol, ProtocolMessage,
ToBytes,
messages::{ParameterDescription, RowDescription, parse::Parse},
},
state::State,
};
use crate::{net::ErrorResponse, util::time::deadline};
use parking_lot::RwLock;
use pgdog_stats::PreparedStatementsConfig;

Expand Down Expand Up @@ -105,6 +106,7 @@ pub struct PreparedStatements {
config: PreparedStatementsConfig,
memory_used: usize,
oids: Arc<Oids>,
server_state: State,
}

#[cfg(test)]
Expand All @@ -126,6 +128,7 @@ impl PreparedStatements {
config: PreparedStatementsConfig::default(),
memory_used: 0,
oids,
server_state: State::Idle,
}
}

Expand All @@ -135,6 +138,10 @@ impl PreparedStatements {
self.config = config;
}

pub(super) fn set_server_state(&mut self, state: State) {
self.server_state = state;
}

/// Current prepared statement settings.
pub fn config(&self) -> PreparedStatementsConfig {
self.config
Expand Down Expand Up @@ -272,6 +279,8 @@ impl PreparedStatements {

if !parse.anonymous() {
if self.contains(parse.name()) {
// TODO(lev): perform the same in errored transaction check
// as we do for PREPARE below.
self.state.add_simulated(ParseComplete.message()?);
return Ok(HandleResult::Drop);
} else {
Expand Down Expand Up @@ -302,15 +311,43 @@ impl PreparedStatements {
self.state.add('3');
}
}
ProtocolMessage::Prepare { name, .. } => {
if self.contains(name) {
ProtocolMessage::PrepareFromClient(prepare) => {
use crate::net::{CommandComplete, ReadyForQuery};
if self.contains(prepare.name()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think we need to use check_prepared at some extend - the ttl functionality and parses deduplication is not working for prepared rn

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good call!

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if self.server_state == State::TransactionError {
self.state
.add_simulated(ErrorResponse::in_failed_transaction().message()?);
} else {
self.state
.add_simulated(CommandComplete::from_str("PREPARE").message()?);
}

self.state.add_simulated(
if self.server_state == State::TransactionError {
ReadyForQuery::error()
} else {
ReadyForQuery::in_transaction(
self.server_state == State::IdleInTransaction,
)
}
.message()?,
);
return Ok(HandleResult::Drop);
} else {
self.parses.push_back(prepare.name().to_owned());
self.state.add(ExecutionCode::ReadyForQuery);
}
}
ProtocolMessage::EnsurePrepared(prepare) => {
if self.contains(prepare.name()) {
return Ok(HandleResult::Drop);
} else {
self.parses.push_back(name.clone());
self.parses.push_back(prepare.name().to_string());
self.state.add_ignore('C');

// Prepare turns into a Simple Query ('Q') so it expects a regular RFQ back.
self.state.add_ignore(ExecutionCode::ReadyForQuery);
self.parses.push_back(prepare.name().to_owned());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we pushing name twice in this block

return Ok(HandleResult::Forward);
}
}
Expand Down
149 changes: 137 additions & 12 deletions pgdog/src/backend/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -569,7 +569,7 @@ impl Server {
if let Some(message) = self.prepared_statements.state_mut().get_simulated() {
// INVARIANT: omni dedup in multi_shard relies on this being process-unique;
// never substitute a non-unique value here.
return Ok(message.backend(self.id));
break message.backend(self.id);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why this change? that seems like simulated prepared will now be processed and setting sync_prepared=true that will cause more work on check-in

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Because we need CommandComplete and ReadyForQuery to reset the server state back to Idle when we prepare a statement. Otherwise, it gets stuck in reading data. Preparing a statement using Query basically executes a separate query before the client's request, so the state management is a bit more complex.

We just push whatever message we "simulate" through the state manager though and it seems to work.

}
match self.stream_buffer.read(self.stream.as_mut().unwrap()).await {
Ok(message) => {
Expand Down Expand Up @@ -680,6 +680,9 @@ impl Server {
_ => (),
}

self.prepared_statements
.set_server_state(self.stats.get_state());

trace!("{:#?} <<< [{}]", message, self.addr());

Ok(message)
Expand Down Expand Up @@ -1326,16 +1329,18 @@ impl Drop for Server {
pub mod test {
use std::time::SystemTime;

use bytes::{BufMut, BytesMut};
use bytes::{BufMut, Bytes, BytesMut};
use pgdog_stats::PreparedStatementsConfig;
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::TcpListener,
};

use crate::{
backend::pool::token_cache::TokenCache, config::Memory, frontend::PreparedStatements,
net::*,
backend::pool::token_cache::TokenCache,
config::Memory,
frontend::{PreparedStatements, RewritePlan},
net::{Prepare, *},
};

use super::{Error, *};
Expand Down Expand Up @@ -2314,20 +2319,20 @@ pub mod test {

#[tokio::test]
async fn test_manual_prepared() {
crate::logger();
let mut server = test_server().await;

let mut prep = PreparedStatements::new();
let mut parse = Parse::named("test", "SELECT 1::bigint");
prep.insert_prepare(&mut parse);
assert_eq!(parse.name(), "__pgdog_1");
let name = "test";
let query = Bytes::from("SELECT 1::bigint".to_owned());
let prepare = prep.insert_prepare(name, query.clone(), &RewritePlan::default());
assert_eq!(prepare.name(), "__pgdog_1");

server
.send(
&vec![ProtocolMessage::from(Query::new(format!(
"PREPARE {} AS {}",
parse.name(),
parse.query()
)))]
&vec![ProtocolMessage::Query(Query::new(
"PREPARE __pgdog_1 AS SELECT 1::bigint",
))]
.into(),
)
.await
Expand Down Expand Up @@ -4362,6 +4367,126 @@ pub mod test {
);
}

#[tokio::test]
async fn test_prepare_from_client() {
let mut server = test_server().await;

// The last 2 will be simulated
// and we won't receive a "prepared statement already exists" error.
for _ in 0..3 {
server
.send(
&vec![ProtocolMessage::PrepareFromClient(Prepare::new(
"__stmt_1",
"PREPARE __pgdog_template_name AS SELECT $1",
))]
.into(),
)
.await
.unwrap();

for c in ['C', 'Z'] {
let msg = server.read().await.unwrap();
assert_eq!(msg.code(), c);
}
}

assert!(server.prepared_statements_mut().contains("__stmt_1"));
}

#[tokio::test]
async fn test_prepared_execute() {
let mut server = test_server().await;

for _ in 0..3 {
let req = vec![
ProtocolMessage::EnsurePrepared(Prepare::new(
"__stmt_1",
"PREPARE __pgdog_template_name (int) AS SELECT $1",
)),
ProtocolMessage::Query(Query::new("EXECUTE __stmt_1 (1)")),
];

server.send(&req.into()).await.unwrap();

for c in ['T', 'D', 'C', 'Z'] {
let msg = server.read().await.unwrap();
assert_eq!(msg.code(), c);
}
}
}

#[tokio::test]
async fn test_prepare_in_transaction() {
let mut server = test_server().await;

server.execute("BEGIN").await.unwrap();

for _ in 0..3 {
server
.send(
&vec![ProtocolMessage::PrepareFromClient(Prepare::new(
"__stmt_1",
"PREPARE __pgdog_template_name AS SELECT $1",
))]
.into(),
)
.await
.unwrap();

let cmd = server.read().await.unwrap();
assert_eq!(cmd.code(), 'C');
let rfq = server.read().await.unwrap();
assert!(rfq.in_transaction());
}

server.execute("ROLLBACK").await.unwrap();
}

#[tokio::test]
async fn test_prepare_in_transaction_error() {
let mut server = test_server().await;

server
.send(
&vec![ProtocolMessage::PrepareFromClient(Prepare::new(
"__stmt_1",
"PREPARE __pgdog_template_name AS SELECT $1",
))]
.into(),
)
.await
.unwrap();

let cmd = server.read().await.unwrap();
assert_eq!(cmd.code(), 'C');
let rfq = server.read().await.unwrap();
assert!(!rfq.in_transaction());

server.execute("BEGIN").await.unwrap();

let _ = server.execute("SELECT asd").await;

for _ in 0..3 {
server
.send(
&vec![ProtocolMessage::PrepareFromClient(Prepare::new(
"__stmt_1",
"PREPARE __pgdog_template_name AS SELECT $1",
))]
.into(),
)
.await
.unwrap();

let err = ErrorResponse::try_from(server.read().await.unwrap()).unwrap();
assert_eq!(err.code, "25P02");

let rfq = ReadyForQuery::try_from(server.read().await.unwrap()).unwrap();
assert!(rfq.is_transaction_aborted());
}
}

#[test]
fn test_effective_max_age_default_is_base() {
let server = Server::default();
Expand Down
2 changes: 1 addition & 1 deletion pgdog/src/frontend/client/query_engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -280,7 +280,7 @@ impl QueryEngine {
self.stats.state = state;

self.stats
.prepared_statements(context.prepared_statements.len_local());
.prepared_statements(context.prepared_statements.num_statements());
self.stats.memory_used(context.memory_stats);

self.comms.update_stats(self.stats);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ async fn test_close_parse_same_name_global_cache() {
assert_eq!(cached_query, "SELECT $1");

// Verify the client's local cache
assert_eq!(client.client().prepared_statements.len_local(), 1);
assert_eq!(client.client().prepared_statements.num_statements(), 1);
assert!(
client
.client()
Expand Down
1 change: 1 addition & 0 deletions pgdog/src/frontend/client/query_engine/test/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ mod schema_changed;
mod set;
mod set_schema_sharding;
mod sharded;
mod sharded_prepared;
mod spliced;
mod test_omnisharded;
mod transaction_state;
Expand Down
Loading
Loading