Skip to content
Open
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
8 changes: 5 additions & 3 deletions .github/workflows/integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -297,12 +297,14 @@ jobs:

# The control plane for vector indexes is not reachable over the wire while
# this backend declares no vector search capability, so its tests drive the
# storage layer directly against this job's PostgreSQL. They build their own
# throwaway databases; the connection string is the server, not a database.
# storage layer directly against this job's PostgreSQL. The PutItem create
# race tests hold a create open in an outside transaction, which also needs
# direct access. They build their own throwaway databases; the connection
# string is the server, not a database.
- name: Run PostgreSQL storage-level tests
env:
EXTENDDB_TEST_PG_CONNECTION_STRING: postgresql://postgres:devpass@127.0.0.1:5432
run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane
run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane --test put_create_race

# The daemonized server logs to syslog; dump it so server-side failures
# are diagnosable from the job log.
Expand Down
136 changes: 81 additions & 55 deletions crates/storage-postgres/src/data/put_item.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,25 +112,28 @@ impl PostgresEngine {
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;

let old: Option<(serde_json::Value,)> =
let mut old: Option<(serde_json::Value,)> =
bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?;

if let Some((ref old_json,)) = old {
let old_item: Item = json_to_item(old_json.clone())?;
match check_condition(condition, &old_item, maps) {
Ok(()) => {}
Err(StorageError::ConditionFailed(_)) => {
return Err(StorageError::ConditionFailed(Some(old_item)));
let mut attempt: u32 = 0;
loop {
if let Some((ref old_json,)) = old {
let old_item: Item = json_to_item(old_json.clone())?;
match check_condition(condition, &old_item, maps) {
Ok(()) => {}
Err(StorageError::ConditionFailed(_)) => {
return Err(StorageError::ConditionFailed(Some(old_item)));
}
Err(e) => return Err(e),
}
Err(e) => return Err(e),
// Row exists, condition passed: update in place.
let update_sql = format!(
"UPDATE {ddb_table} SET item_data = $3 WHERE pk = $1 AND {sk_col} = $2"
);
bind_sk_execute!(&update_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?;
break;
}
// Row exists, condition passed — update in place.
let update_sql = format!(
"UPDATE {ddb_table} SET item_data = $3 WHERE pk = $1 AND {sk_col} = $2"
);
bind_sk_execute!(&update_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?;
} else {
// No existing item — condition checks against empty item
// No existing item: the condition checks against an empty item
let empty = std::collections::BTreeMap::new();
match check_condition(condition, &empty, maps) {
Ok(()) => {}
Expand All @@ -139,21 +142,25 @@ impl PostgresEngine {
}
Err(e) => return Err(e),
}
// Condition passed against empty — atomic insert, fail if someone beat us.
let insert_sql = format!(
"INSERT INTO {ddb_table} (pk, {sk_col}, item_data) VALUES ($1, $2, $3) \
ON CONFLICT (pk, {sk_col}) DO NOTHING"
);
let result =
bind_sk_execute!(&insert_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?;
if result.rows_affected() == 0 {
// Another transaction inserted between our SELECT and INSERT.
// Fetch the winner to return with ConditionFailed.
let winner: Option<(serde_json::Value,)> =
bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?;
let winner_item = winner.map(|(v,)| json_to_item(v)).transpose()?;
return Err(StorageError::ConditionFailed(winner_item));
if result.rows_affected() == 1 {
break;
}
// Lost the create race. The winner has committed (the insert
// waited for it if it was in flight). The locking read returns
// its row, and this put overwrites it after checking its
// condition against it. If a delete committed since, the read
// returns no row and the insert is retried.
attempt += 1;
if attempt >= MAX_CREATE_RACE_ATTEMPTS {
return Err(create_race_exhausted(attempt));
}
old = bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?;
}

// Sync GSI/LSI update within transaction (D-4).
Expand Down Expand Up @@ -269,30 +276,37 @@ impl PostgresEngine {
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;

let old: Option<(serde_json::Value,)> = sqlx::query_as(&select_sql)
.bind(pk_text.as_str())
.fetch_optional(&mut *tx)
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;

if let Some((ref old_json,)) = old {
let old_item: Item = json_to_item(old_json.clone())?;
match check_condition(condition, &old_item, maps) {
Ok(()) => {}
Err(StorageError::ConditionFailed(_)) => {
return Err(StorageError::ConditionFailed(Some(old_item)));
}
Err(e) => return Err(e),
}
// Row exists, condition passed — update in place.
let update_sql = format!("UPDATE {ddb_table} SET item_data = $2 WHERE pk = $1");
sqlx::query(&update_sql)
let fetch_old = async |tx: &mut sqlx::PgConnection| {
sqlx::query_as::<_, (serde_json::Value,)>(&select_sql)
.bind(pk_text.as_str())
.bind(&item_json)
.execute(&mut *tx)
.fetch_optional(tx)
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;
} else {
.map_err(|e| StorageError::Internal(e.to_string()))
};
let mut old: Option<(serde_json::Value,)> = fetch_old(&mut tx).await?;

let mut attempt: u32 = 0;
loop {
if let Some((ref old_json,)) = old {
let old_item: Item = json_to_item(old_json.clone())?;
match check_condition(condition, &old_item, maps) {
Ok(()) => {}
Err(StorageError::ConditionFailed(_)) => {
return Err(StorageError::ConditionFailed(Some(old_item)));
}
Err(e) => return Err(e),
}
// Row exists, condition passed: update in place.
let update_sql =
format!("UPDATE {ddb_table} SET item_data = $2 WHERE pk = $1");
sqlx::query(&update_sql)
.bind(pk_text.as_str())
.bind(&item_json)
.execute(&mut *tx)
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;
break;
}
let empty = std::collections::BTreeMap::new();
match check_condition(condition, &empty, maps) {
Ok(()) => {}
Expand All @@ -301,7 +315,6 @@ impl PostgresEngine {
}
Err(e) => return Err(e),
}
// Condition passed against empty — atomic insert, fail if someone beat us.
let insert_sql = format!(
"INSERT INTO {ddb_table} (pk, item_data) VALUES ($1, $2) \
ON CONFLICT (pk) DO NOTHING"
Expand All @@ -312,16 +325,15 @@ impl PostgresEngine {
.execute(&mut *tx)
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;
if result.rows_affected() == 0 {
// Another transaction inserted between our SELECT and INSERT.
let winner: Option<(serde_json::Value,)> = sqlx::query_as(&select_sql)
.bind(pk_text.as_str())
.fetch_optional(&mut *tx)
.await
.map_err(|e| StorageError::Internal(e.to_string()))?;
let winner_item = winner.map(|(v,)| json_to_item(v)).transpose()?;
return Err(StorageError::ConditionFailed(winner_item));
if result.rows_affected() == 1 {
break;
}
// Lost the create race: overwrite the winner, as above.
attempt += 1;
if attempt >= MAX_CREATE_RACE_ATTEMPTS {
return Err(create_race_exhausted(attempt));
}
old = fetch_old(&mut tx).await?;
}

// Sync GSI/LSI update within transaction (D-4).
Expand Down Expand Up @@ -466,3 +478,17 @@ impl PostgresEngine {
json_opt.map(json_to_item).transpose()
}
}

/// Bound on insert retries when a put to a missing item keeps losing the
/// create race to a winner that is deleted again before the re-read. Same
/// bound as `UpdateItem`.
const MAX_CREATE_RACE_ATTEMPTS: u32 = 5;

/// The error after `MAX_CREATE_RACE_ATTEMPTS` lost create races. Like the one
/// `UpdateItem` returns, it is an internal error (HTTP 500).
fn create_race_exhausted(attempt: u32) -> StorageError {
StorageError::Internal(format!(
"PutItem could not create the item after {attempt} attempts: a concurrent writer \
repeatedly created and deleted the row"
))
}
Loading
Loading