Skip to content
Open
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
140 changes: 92 additions & 48 deletions nisshi-storage-sql/src/pg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,10 +72,10 @@ use opentelemetry::{KeyValue, metrics::Counter};
use rand::{prelude::*, rng};
use serde_json::Value;
use tokio_postgres::{
Config, Row, RowStream,
Config, Row,
binary_copy::BinaryCopyInWriter,
error::SqlState,
types::{BorrowToSql, ToSql, Type},
types::{ToSql, Type},
};
use tracing::{debug, error, instrument};
use url::Url;
Expand Down Expand Up @@ -804,28 +804,6 @@ impl Postgres {
.map_err(Into::into)
}

#[instrument(skip(self, tx, params))]
async fn tx_prepare_query_raw<P, I>(
&self,
tx: &Transaction<'_>,
sql: &str,
params: I,
) -> Result<RowStream, Error>
where
P: BorrowToSql,
I: IntoIterator<Item = P>,
I::IntoIter: ExactSizeIterator,
{
let sql = self.sql_lookup(sql)?;

let prepared = tx
.prepare_cached(sql)
.await
.inspect_err(|err| error!(?err))?;

tx.query_raw(&prepared, params).await.map_err(Into::into)
}

#[instrument(skip_all)]
async fn produce_in_tx(
&self,
Expand Down Expand Up @@ -1916,30 +1894,31 @@ impl Storage for Postgres {

debug!(?topic_uuid, cluster = self.cluster, ?topic);

_ = future::try_join_all(
(0..topic.num_partitions)
.map(|partition| {
let cluster = Box::new(self.cluster.clone()) as Box<dyn ToSql + Sync + Send>;
let name = Box::new(topic.name.clone()) as Box<dyn ToSql + Sync + Send>;
let partition = Box::new(partition) as Box<dyn ToSql + Sync + Send>;
[cluster, name, partition]
})
.map(|parameters| self.tx_prepare_query_raw(&tx, "topition_insert.sql", parameters))
.chain(
(0..topic.num_partitions)
.map(|partition| {
let cluster =
Box::new(self.cluster.clone()) as Box<dyn ToSql + Sync + Send>;
let name = Box::new(topic.name.clone()) as Box<dyn ToSql + Sync + Send>;
let partition = Box::new(partition) as Box<dyn ToSql + Sync + Send>;
[cluster, name, partition]
})
.map(|parameters| {
self.tx_prepare_query_raw(&tx, "watermark_insert.sql", parameters)
}),
),
)
.await?;
// The watermark insert selects the topition rows that the topition insert creates, so
// every topition insert must finish before the first watermark insert is sent. A
// statement that is not yet prepared on this connection takes a round trip first,
// which would otherwise let a watermark insert overtake its topition and match no
// rows, leaving a topic that rejects every produce.
for sql in ["topition_insert.sql", "watermark_insert.sql"] {
let inserted = future::try_join_all((0..topic.num_partitions).map(|partition| {
let tx = &tx;
let cluster = &self.cluster;
let name = &topic.name;

async move {
self.tx_prepare_execute(tx, sql, &[cluster, name, &partition])
.await
}
}))
.await?;

if inserted.iter().any(|rows| *rows != 1) {
return Err(Error::Message(format!(
"{sql} did not insert one row per partition of {} in {}",
topic.name, self.cluster
)));
}
}

if let Some(configs) = topic.configs {
for config in configs {
Expand Down Expand Up @@ -6565,6 +6544,71 @@ mod tests {
Ok(())
}

/// The watermark insert selects the topition row that the topition insert creates, so
/// `create_topic` must finish the second before sending the first. When it sent both
/// together, a connection holding only the watermark statement in its statement cache
/// sent the watermark insert first; it matched no rows without an error, and the topic
/// had no watermark, so every produce to it failed with `UnknownTopicOrPartition`. No
/// current caller leaves the cache in that state; this test sets it up directly.
#[tokio::test]
async fn create_topic_writes_the_watermark_when_only_its_statement_is_cached() -> Result<()> {
let cluster = alphanumeric_string(15);

let storage = Postgres::builder(CONNECTION)?
.cluster(cluster.as_str())
.node(rng().random_range(0..i32::MAX))
.build();

if let Err(err) = storage.connection().await {
eprintln!(
"skipping create_topic_writes_the_watermark_when_only_its_statement_is_cached: {err:?}"
);
return Ok(());
}

storage
.register_broker(BrokerRegistrationRequest {
broker_id: 111,
cluster_id: cluster.clone(),
incarnation_id: Uuid::now_v7(),
rack: None,
})
.await?;

{
let c = storage.connection().await?;
_ = c
.prepare_cached(storage.sql_lookup("watermark_insert.sql")?)
.await?;
}

assert_eq!(
1,
storage.pool.status().size,
"one connection holds the cache"
);

let topic_name = alphanumeric_string(15);

_ = storage
.create_topic(
CreatableTopic::default()
.name(topic_name.clone())
.num_partitions(1)
.replication_factor(0)
.assignments(Some([].into()))
.configs(Some([].into())),
false,
)
.await?;

let stage = storage.offset_stage(&Topition::new(topic_name, 0)).await?;

assert_eq!(0, stage.high_watermark);

Ok(())
}

/// Deleting one topic must leave every other topic in the cluster untouched.
/// `topic_delete_by.sql` narrowed the delete on `topic.cluster` instead of
/// `topic.id`, so the subquery only established that the named topic existed
Expand Down
Loading