chore: bump crates

pull/370/head
zijiren233 4 months ago
parent cd88d8535b
commit e4aad1d526
No known key found for this signature in database
GPG Key ID: 534E082AAA9B39DC

547
Cargo.lock generated

File diff suppressed because it is too large Load Diff

@ -96,9 +96,9 @@ tonic-prost-build = "0.14.6"
protoc-bin-vendored = "3.2.0"
prost = "0.14.3"
prost-types = "0.14.3"
prost-reflect = "0.16.3"
prost-reflect = "0.16.4"
prost-reflect-build = "0.16.0"
prost-protovalidate = "0.4.0"
prost-protovalidate = "0.4.2"
# HTTP Framework
axum = { version = "0.8.9", features = ["macros", "ws"] }
@ -108,7 +108,7 @@ tower = { version = "0.5.3", features = [
"limit",
"load-shed",
] }
tower-http = { version = "0.6.10", features = [
tower-http = { version = "0.6.11", features = [
"trace",
"on-early-drop",
"compression-full",
@ -116,13 +116,13 @@ tower-http = { version = "0.6.10", features = [
"timeout",
"cors",
] }
hyper = { version = "1.9.0", default-features = false, features = [
hyper = { version = "1.10.1", default-features = false, features = [
"http1",
"http2",
"server",
"client",
] }
reqwest = { version = "0.13.3", default-features = false, features = [
reqwest = { version = "0.13.4", default-features = false, features = [
"json",
"cookies",
"query",
@ -130,7 +130,7 @@ reqwest = { version = "0.13.3", default-features = false, features = [
] }
# Database
sqlx = { version = "0.8.6", default-features = false, features = [
sqlx = { version = "0.9.0", default-features = false, features = [
"runtime-tokio",
"postgres",
"migrate",
@ -141,7 +141,7 @@ sqlx = { version = "0.8.6", default-features = false, features = [
] }
# Redis
redis = { version = "1.2.1", default-features = false, features = [
redis = { version = "1.2.2", default-features = false, features = [
"tokio-comp",
"connection-manager",
"streams",
@ -166,7 +166,7 @@ subtle = "2.6.1"
# Serialization
serde = { version = "1.0.228", features = ["derive"] }
serde_json = "1.0.149"
serde_json = "1.0.150"
serde_yaml = "0.9.34"
serde_ignored = "0.1.14"
toml = "1.1.2"
@ -176,12 +176,12 @@ bincode = { version = "2.0.1", features = ["serde"] }
# Caching
moka = { version = "0.12.15", features = ["future", "sync"] }
dashmap = "6.1.0"
dashmap = "6.2.1"
lru = "0.18.0"
# IDs & Crypto
nanoid = "0.5.0"
uuid = { version = "1.23.1", features = ["v4", "v5", "serde"] }
uuid = { version = "1.23.2", features = ["v4", "v5", "serde"] }
sqids = "0.4.2"
rand = "0.10.1"
rand_08 = { package = "rand", version = "0.8.6" }
@ -197,7 +197,7 @@ iana-time-zone = "0.1.65"
synctv-xiu = { path = "synctv-xiu", default-features = false }
# Object Storage
opendal = { version = "0.56.0", default-features = false, features = [
opendal = { version = "0.57.0", default-features = false, features = [
"executors-tokio",
"services-s3",
] }
@ -232,15 +232,15 @@ urlencoding = "2.1.3"
regex = "1.12.3"
paste = "1.0.15"
byteorder = "1.5.0"
http = "1.4.0"
garde = { version = "0.22.1", features = ["derive"] }
http = "1.4.1"
garde = { version = "0.23.0", features = ["derive"] }
zxcvbn = { version = "3.1.1", default-features = false }
# Monitoring
prometheus = { version = "0.14.0", features = ["process"] }
# Time
chrono = { version = "0.4.44", features = ["serde"] }
chrono = { version = "0.4.45", features = ["serde"] }
# Kubernetes (in-cluster client for leader election & service discovery)
@ -271,7 +271,7 @@ handlebars = "6.4.1"
# Compression
flate2 = "1.1.9"
brotli = "8.0.2"
brotli = "8.0.3"
# Resilience
governor = "0.10.4"
@ -303,8 +303,8 @@ tokio-tungstenite = { version = "0.29.0", default-features = false, features = [
criterion = { version = "0.8.2", features = ["html_reports", "async_tokio"] }
# Allocators
tikv-jemallocator = "0.6.1"
mimalloc = { version = "0.1.50", features = ["v2"] }
tikv-jemallocator = "0.7.0"
mimalloc = { version = "0.1.52", features = ["v2"] }
# Internal workspace crates
synctv-proto = { path = "synctv-proto", default-features = false }

@ -9,6 +9,7 @@ use tokio::task::JoinHandle;
use tracing::Level;
use tracing::{error, info};
use crate::repository::query_builder::trusted_dynamic_sql;
use crate::resilience::timeout::DB_QUERY_TIMEOUT;
use crate::Config;
@ -119,10 +120,14 @@ async fn apply_session_settings(
statement_timeout_ms: u128,
client_min_messages: &'static str,
) -> std::result::Result<(), sqlx::Error> {
conn.execute(format!("SET statement_timeout = {statement_timeout_ms}").as_str())
.await?;
conn.execute(format!("SET client_min_messages = '{client_min_messages}'").as_str())
.await?;
conn.execute(trusted_dynamic_sql(format!(
"SET statement_timeout = {statement_timeout_ms}"
)))
.await?;
conn.execute(trusted_dynamic_sql(format!(
"SET client_min_messages = '{client_min_messages}'"
)))
.await?;
Ok(())
}

@ -85,7 +85,7 @@ impl AuditLogRepository {
/// The caller must have already pushed the prefix (e.g. `"SELECT ... WHERE "`)
/// so that the first condition can be appended directly.
fn push_filters<'q>(
builder: &mut QueryBuilder<'q, Postgres>,
builder: &mut QueryBuilder<Postgres>,
query: &'q AuditLogQuery,
effective_from: DateTime<Utc>,
) {

@ -139,7 +139,7 @@ impl MediaRepository {
}
fn push_media_list_order_by(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &MediaListQuery,
) {
use crate::models::{MediaListSortBy, SortDirection};
@ -196,7 +196,7 @@ impl MediaRepository {
}
fn push_media_scope_filters(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
room_id: &RoomId,
playlist_id: Option<&PlaylistId>,
query: &MediaListQuery,

@ -65,13 +65,14 @@ fn test_push_media_scope_filters_treats_empty_provider_instance_as_default() {
let built = builder.build();
assert!(built
.sql()
.as_str()
.contains("NULLIF(m.provider_instance_name, '') IS NULL"));
}
fn media_order_by_sql(query: &MediaListQuery) -> String {
let mut builder = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
MediaRepository::push_media_list_order_by(&mut builder, query);
builder.sql().to_string()
builder.sql().as_str().to_string()
}
#[test]

@ -91,7 +91,7 @@ impl NotificationRepository {
}
fn push_list_order_by(
qb: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
qb: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &NotificationListQuery,
) {
use crate::models::SortDirection;

@ -64,7 +64,7 @@ fn test_notification_list_query_with_filters() {
fn notification_order_by_sql(query: &NotificationListQuery) -> String {
let mut builder = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
NotificationRepository::push_list_order_by(&mut builder, query);
builder.sql().to_string()
builder.sql().as_str().to_string()
}
#[test]

@ -140,7 +140,7 @@ impl PlaylistRepository {
}
fn push_playlist_list_order_by(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &PlaylistListQuery,
) {
use crate::models::{PlaylistListSortBy, SortDirection};
@ -179,7 +179,7 @@ impl PlaylistRepository {
}
fn push_playlist_scope_filters(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
room_id: &RoomId,
parent_id: Option<&PlaylistId>,
query: &PlaylistListQuery,

@ -36,7 +36,7 @@ fn assert_position_eq(actual: f64, expected: f64) {
fn playlist_order_by_sql(query: &PlaylistListQuery) -> String {
let mut builder = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
PlaylistRepository::push_playlist_list_order_by(&mut builder, query);
builder.sql().to_string()
builder.sql().as_str().to_string()
}
#[test]
@ -143,6 +143,7 @@ fn test_push_playlist_scope_filters_treats_empty_provider_instance_as_default()
let built = builder.build();
assert!(built
.sql()
.as_str()
.contains("NULLIF(p.provider_instance_name, '') IS NULL"));
}

@ -126,7 +126,7 @@ impl ProviderInstanceRepository {
}
fn push_list_order_by(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &ProviderInstanceListQuery,
) {
use crate::models::SortDirection;
@ -161,7 +161,7 @@ impl ProviderInstanceRepository {
}
fn push_list_filters(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &ProviderInstanceListQuery,
) -> Result<()> {
builder.push(" WHERE TRUE");

@ -9,7 +9,7 @@ use serde_json::json;
fn order_by_sql(query: &ProviderInstanceListQuery) -> String {
let mut builder = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
ProviderInstanceRepository::push_list_order_by(&mut builder, query);
builder.sql().to_string()
builder.sql().as_str().to_string()
}
#[test]

@ -11,6 +11,10 @@
use crate::{Error, Result};
pub(crate) fn trusted_dynamic_sql(sql: String) -> sqlx::AssertSqlSafe<String> {
sqlx::AssertSqlSafe(sql)
}
/// A single condition in the WHERE clause.
enum Condition {
/// A static SQL fragment with no bound parameter (e.g. `r.deleted_at IS NULL`).

@ -472,7 +472,7 @@ impl RoomRepository {
Ok(deleted)
}
fn push_room_projection(builder: &mut QueryBuilder<'_, Postgres>) {
fn push_room_projection(builder: &mut QueryBuilder<Postgres>) {
builder.push(
r"
r.id,
@ -496,7 +496,7 @@ impl RoomRepository {
);
}
fn push_where_prefix(builder: &mut QueryBuilder<'_, Postgres>, has_condition: &mut bool) {
fn push_where_prefix(builder: &mut QueryBuilder<Postgres>, has_condition: &mut bool) {
if *has_condition {
builder.push(" AND ");
} else {
@ -506,7 +506,7 @@ impl RoomRepository {
}
fn push_room_list_filters<'q>(
builder: &mut QueryBuilder<'q, Postgres>,
builder: &mut QueryBuilder<Postgres>,
query: &'q RoomListQuery,
search_pattern: Option<&'q String>,
has_condition: &mut bool,

@ -6,6 +6,7 @@ use crate::{
PageParams, RoomId, RoomMember, RoomMemberListQuery, RoomMemberListSortBy,
RoomMemberWithUser, RoomRole, RoomStatus, UserId,
},
repository::query_builder::trusted_dynamic_sql,
Error, Result,
};
@ -282,7 +283,7 @@ impl RoomMemberRepository {
}
fn push_room_member_order_by(
builder: &mut sqlx::QueryBuilder<'_, sqlx::Postgres>,
builder: &mut sqlx::QueryBuilder<sqlx::Postgres>,
query: &RoomMemberListQuery,
) {
builder.push(" ORDER BY ");
@ -2402,7 +2403,7 @@ impl RoomMemberRepository {
);
let rows = Self::bind_my_room_filters(
sqlx::query_as::<_, MyRoomListRow>(&sql)
sqlx::query_as::<_, MyRoomListRow>(trusted_dynamic_sql(sql))
.bind(user_id)
.bind(limit)
.bind(offset),
@ -2461,7 +2462,7 @@ impl RoomMemberRepository {
);
let rows = Self::bind_my_room_filters(
sqlx::query_as::<_, MyRoomListRow>(&sql)
sqlx::query_as::<_, MyRoomListRow>(trusted_dynamic_sql(sql))
.bind(user_id)
.bind(limit)
.bind(offset),
@ -2685,7 +2686,7 @@ mod tests {
fn room_member_order_by_sql(query: &RoomMemberListQuery) -> String {
let mut builder = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
RoomMemberRepository::push_room_member_order_by(&mut builder, query);
builder.sql().to_string()
builder.sql().as_str().to_string()
}
#[test]

@ -783,7 +783,7 @@ impl UserRepository {
}
fn push_user_list_from_and_filters<'a>(
builder: &mut QueryBuilder<'a, Postgres>,
builder: &mut QueryBuilder<Postgres>,
query: &'a UserListQuery,
role_scope: UserListRoleScope,
search_pattern: Option<&'a str>,

@ -10,6 +10,7 @@ use tracing::{info, warn};
use super::LeaderCheck;
use crate::bootstrap::acquire_unbounded_ddl_connection;
use crate::repository::query_builder::trusted_dynamic_sql;
use crate::service::partitioning::{
add_months, current_database_date, quote_ident, size_centi_gib, size_centi_mib, start_of_month,
table_exists,
@ -171,42 +172,42 @@ impl AuditPartitionManager {
.await
.internal_with_err("Failed to acquire DDL connection for single partition creation")?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE TABLE IF NOT EXISTS {partition_ident} PARTITION OF audit_logs \
FOR VALUES FROM ('{start_date}') TO ('{end_date}')"
))
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to create audit partition")?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(actor_id, created_at DESC) WHERE actor_id IS NOT NULL",
quote_ident(&format!("{partition_name}_idx_actor_created"))
))
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to create audit partition index")?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(action, created_at DESC)",
quote_ident(&format!("{partition_name}_idx_action_created"))
))
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to create audit partition index")?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(target_type, target_id, created_at DESC) WHERE target_type IS NOT NULL",
quote_ident(&format!("{partition_name}_idx_target_created"))
))
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to create audit partition index")?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(ip_address) WHERE ip_address IS NOT NULL",
quote_ident(&format!("{partition_name}_idx_ip_address"))
))
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to create audit partition index")?;
@ -256,10 +257,13 @@ impl AuditPartitionManager {
.collect::<Vec<_>>();
for partition in &dropped_partitions {
sqlx::query(&format!("DROP TABLE IF EXISTS {}", quote_ident(partition)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to drop audit partition")?;
sqlx::query(trusted_dynamic_sql(format!(
"DROP TABLE IF EXISTS {}",
quote_ident(partition)
)))
.execute(&mut *conn)
.await
.internal_with_err("Failed to drop audit partition")?;
}
let dropped_count = len_to_i32(dropped_partitions.len(), "dropped audit partition count")?;

@ -11,6 +11,7 @@ use tracing::{error, info, warn};
use super::LeaderCheck;
use crate::bootstrap::acquire_unbounded_ddl_connection;
use crate::repository::query_builder::trusted_dynamic_sql;
use crate::service::global_settings::SettingsRegistry;
use crate::service::partitioning::{
current_database_date, quote_ident, size_centi_mib, table_exists,
@ -104,42 +105,42 @@ impl ChatPartitionManager {
let partition_name = format!("chat_messages_{}", start_date.format("%Y_%m_%d"));
let partition_ident = quote_ident(&partition_name);
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE TABLE IF NOT EXISTS {partition_ident} PARTITION OF chat_messages \
FOR VALUES FROM ('{start_date}') TO ('{end_date}')"
))
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to create chat partition: {e}")))?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(room_id, created_at DESC, id DESC)",
quote_ident(&format!("{partition_name}_idx_room_pagination"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to create chat partition index: {e}")))?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(user_id, created_at DESC)",
quote_ident(&format!("{partition_name}_idx_user_created"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to create chat partition index: {e}")))?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(created_at DESC)",
quote_ident(&format!("{partition_name}_idx_created_at"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to create chat partition index: {e}")))?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(room_id, status, created_at DESC, id DESC)",
quote_ident(&format!("{partition_name}_idx_status"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to create chat partition index: {e}")))?;
@ -188,10 +189,13 @@ impl ChatPartitionManager {
.collect::<Vec<_>>();
for partition in &partitions {
sqlx::query(&format!("DROP TABLE IF EXISTS {}", quote_ident(partition)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to drop old chat partition: {e}")))?;
sqlx::query(trusted_dynamic_sql(format!(
"DROP TABLE IF EXISTS {}",
quote_ident(partition)
)))
.execute(&mut *conn)
.await
.map_err(|e| Error::Internal(format!("Failed to drop old chat partition: {e}")))?;
}
let dropped_count = len_to_i64(partitions.len(), "dropped chat partition count")?;

@ -11,6 +11,7 @@ use tracing::{error, info};
use super::LeaderCheck;
use crate::bootstrap::acquire_unbounded_ddl_connection;
use crate::repository::query_builder::trusted_dynamic_sql;
use crate::service::partitioning::{
add_months, current_database_date, quote_ident, size_centi_mib, start_of_month, table_exists,
};
@ -93,40 +94,40 @@ impl NotificationPartitionManager {
let partition_name = format!("notifications_{}", start_date.format("%Y_%m"));
let partition_ident = quote_ident(&partition_name);
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE TABLE IF NOT EXISTS {partition_ident} PARTITION OF notifications \
FOR VALUES FROM ('{start_date}') TO ('{end_date}')"
))
)))
.execute(&mut *conn)
.await
.map_err(|e| {
Error::Internal(format!("Failed to create notification partition: {e}"))
})?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(user_id, is_read, created_at DESC)",
quote_ident(&format!("{partition_name}_idx_user_read_created"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| {
Error::Internal(format!("Failed to create notification partition index: {e}"))
})?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(user_id, created_at DESC) WHERE is_read = FALSE",
quote_ident(&format!("{partition_name}_idx_user_unread"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| {
Error::Internal(format!("Failed to create notification partition index: {e}"))
})?;
sqlx::query(&format!(
sqlx::query(trusted_dynamic_sql(format!(
"CREATE INDEX IF NOT EXISTS {} ON {partition_ident}(user_id, type, created_at DESC) WHERE is_read = FALSE",
quote_ident(&format!("{partition_name}_idx_user_type_created"))
))
)))
.execute(&mut *conn)
.await
.map_err(|e| {
@ -180,12 +181,15 @@ impl NotificationPartitionManager {
.collect::<Vec<_>>();
for partition in &partitions {
sqlx::query(&format!("DROP TABLE IF EXISTS {}", quote_ident(partition)))
.execute(&mut *conn)
.await
.map_err(|e| {
Error::Internal(format!("Failed to drop old notification partition: {e}"))
})?;
sqlx::query(trusted_dynamic_sql(format!(
"DROP TABLE IF EXISTS {}",
quote_ident(partition)
)))
.execute(&mut *conn)
.await
.map_err(|e| {
Error::Internal(format!("Failed to drop old notification partition: {e}"))
})?;
}
let dropped_count = len_to_i64(partitions.len(), "dropped notification partition count")?;

@ -26,7 +26,7 @@ use crate::{
cache::KeyBuilder,
models::{oauth2_client::OAuth2Provider, SignupMethod, User, UserId},
oauth2::Provider as OAuth2ProviderTrait,
repository::UserOAuthProviderRepository,
repository::{query_builder::trusted_dynamic_sql, UserOAuthProviderRepository},
service::{
user::PendingRegistrationConflict, OAuth2SignupPolicy, SettingsRegistry, UserService,
},
@ -1310,7 +1310,7 @@ impl OAuth2Service {
let mut new_user = None;
for (attempt, candidate) in candidates.iter().enumerate() {
let savepoint = format!("oauth2_user_create_{attempt}");
sqlx::query(&format!("SAVEPOINT {savepoint}"))
sqlx::query(trusted_dynamic_sql(format!("SAVEPOINT {savepoint}")))
.execute(&mut *tx)
.await
.internal_with_err("Failed to create OAuth2 user savepoint")?;
@ -1335,12 +1335,14 @@ impl OAuth2Service {
)
.await
{
sqlx::query(&format!("ROLLBACK TO SAVEPOINT {savepoint}"))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after email create error",
)?;
sqlx::query(trusted_dynamic_sql(format!(
"ROLLBACK TO SAVEPOINT {savepoint}"
)))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after email create error",
)?;
match error {
Error::AlreadyExists(_) => {
return Err(Error::AlreadyExists(
@ -1351,10 +1353,12 @@ impl OAuth2Service {
err => return Err(err),
}
}
sqlx::query(&format!("RELEASE SAVEPOINT {savepoint}"))
.execute(&mut *tx)
.await
.internal_with_err("Failed to release OAuth2 user savepoint")?;
sqlx::query(trusted_dynamic_sql(format!(
"RELEASE SAVEPOINT {savepoint}"
)))
.execute(&mut *tx)
.await
.internal_with_err("Failed to release OAuth2 user savepoint")?;
user_service
.cache_oauth2_username_best_effort(&created_user.id, candidate)
@ -1385,20 +1389,24 @@ impl OAuth2Service {
break;
}
Err(error) if UserService::is_username_conflict(&error) => {
sqlx::query(&format!("ROLLBACK TO SAVEPOINT {savepoint}"))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after username collision",
)?;
sqlx::query(trusted_dynamic_sql(format!(
"ROLLBACK TO SAVEPOINT {savepoint}"
)))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after username collision",
)?;
}
Err(err) => {
sqlx::query(&format!("ROLLBACK TO SAVEPOINT {savepoint}"))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after create error",
)?;
sqlx::query(trusted_dynamic_sql(format!(
"ROLLBACK TO SAVEPOINT {savepoint}"
)))
.execute(&mut *tx)
.await
.internal_with_err(
"Failed to roll back OAuth2 user savepoint after create error",
)?;
return Err(err);
}
}

@ -26,6 +26,10 @@ use crate::docker::{
ProcessLock, TEST_RUN_LABEL,
};
fn trusted_dynamic_sql(sql: String) -> sqlx::AssertSqlSafe<String> {
sqlx::AssertSqlSafe(sql)
}
/// Default `PostgreSQL` version for test containers
pub const POSTGRES_VERSION: &str = "18";
const DEFAULT_DOCKER_STARTUP_TIMEOUT_SECS: u64 = 300;
@ -182,7 +186,10 @@ impl SharedPostgresServer {
quote_identifier(database_name)
);
if let Err(err) = sqlx::query(&sql).execute(&self.admin_pool).await {
if let Err(err) = sqlx::query(trusted_dynamic_sql(sql))
.execute(&self.admin_pool)
.await
{
eprintln!("warning: failed to drop postgres test database {database_name}: {err}");
}
}
@ -218,13 +225,13 @@ impl TestContainer {
let database = quote_identifier(&self.database_name);
let drop_sql = format!("DROP DATABASE IF EXISTS {database} WITH (FORCE)");
sqlx::query(&drop_sql)
sqlx::query(trusted_dynamic_sql(drop_sql))
.execute(&self.shared.admin_pool)
.await
.expect("test database should be dropped before empty recreation");
let create_sql = format!("CREATE DATABASE {database}");
sqlx::query(&create_sql)
sqlx::query(trusted_dynamic_sql(create_sql))
.execute(&self.shared.admin_pool)
.await
.expect("test database should be recreated empty");
@ -673,19 +680,21 @@ async fn recreate_template_database(
"ALTER DATABASE {} WITH IS_TEMPLATE false",
quote_identifier(template_database)
);
let _ = sqlx::query(&untemplate_sql).execute(admin_pool).await;
let _ = sqlx::query(trusted_dynamic_sql(untemplate_sql))
.execute(admin_pool)
.await;
let drop_sql = format!(
"DROP DATABASE IF EXISTS {} WITH (FORCE)",
quote_identifier(template_database)
);
sqlx::query(&drop_sql)
sqlx::query(trusted_dynamic_sql(drop_sql))
.execute(admin_pool)
.await
.expect("template database cleanup should succeed");
let create_sql = format!("CREATE DATABASE {}", quote_identifier(template_database));
sqlx::query(&create_sql)
sqlx::query(trusted_dynamic_sql(create_sql))
.execute(admin_pool)
.await
.expect("template database creation should succeed");
@ -708,7 +717,7 @@ async fn recreate_template_database(
"ALTER DATABASE {} WITH ALLOW_CONNECTIONS false",
quote_identifier(template_database)
);
sqlx::query(&no_connections_sql)
sqlx::query(trusted_dynamic_sql(no_connections_sql))
.execute(admin_pool)
.await
.expect("template database should disallow new connections");
@ -716,7 +725,7 @@ async fn recreate_template_database(
let terminate_sql = format!(
"SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = '{template_database}' AND pid <> pg_backend_pid()"
);
sqlx::query(&terminate_sql)
sqlx::query(trusted_dynamic_sql(terminate_sql))
.execute(admin_pool)
.await
.expect("template database should terminate lingering sessions");
@ -725,7 +734,7 @@ async fn recreate_template_database(
"ALTER DATABASE {} WITH IS_TEMPLATE true",
quote_identifier(template_database)
);
sqlx::query(&mark_template_sql)
sqlx::query(trusted_dynamic_sql(mark_template_sql))
.execute(admin_pool)
.await
.expect("template database should be marked reusable");
@ -879,7 +888,7 @@ async fn provision_test_database(requested_db_name: &str, label: &str) -> TestCo
.await
.expect("direct postgres admin connection for template clone should succeed");
sqlx::query(&create_sql)
sqlx::query(trusted_dynamic_sql(create_sql))
.execute(&mut admin_connection)
.await
.expect("test database creation from template should succeed");

@ -9,6 +9,10 @@
use synctv_core::service::{auth::token_blacklist::PgTokenBlacklistStore, TokenBlacklistStore};
use synctv_core_testing::create_test_pool;
fn trusted_dynamic_sql(sql: String) -> sqlx::AssertSqlSafe<String> {
sqlx::AssertSqlSafe(sql)
}
#[tokio::test]
#[ignore = "Requires Docker"]
async fn test_pg_family_revocation_survives_cleanup_until_marker_expires() {
@ -62,6 +66,7 @@ async fn test_pg_family_revocation_is_atomic_when_timestamp_write_fails() {
let store = PgTokenBlacklistStore::new(pool.clone());
let key = format!("family:pg_atomicity_guard:{}", synctv_common::snanoid!(8));
let timestamp = chrono::Utc::now().timestamp();
let key_sql_literal = key.replace('\'', "''");
let trigger_fn_sql = r"
CREATE OR REPLACE FUNCTION fail_token_blacklist_family_insert()
@ -74,9 +79,12 @@ async fn test_pg_family_revocation_is_atomic_when_timestamp_write_fails() {
END;
$$ LANGUAGE plpgsql;
"
.replace("REPLACE_ME", &key);
.replace("REPLACE_ME", &key_sql_literal);
sqlx::query(&trigger_fn_sql).execute(&pool).await.unwrap();
sqlx::query(trusted_dynamic_sql(trigger_fn_sql))
.execute(&pool)
.await
.unwrap();
sqlx::query(
"DROP TRIGGER IF EXISTS trg_fail_token_blacklist_family_insert ON auth_token_blacklist",

@ -74,7 +74,7 @@ async fn run_migrate_with_connection(
.map_err(|e| anyhow::anyhow!("Failed to disable statement_timeout for migrations: {e}"))?;
sqlx::migrate!("../migrations")
.run_direct(&mut **conn)
.run_direct(None, &mut **conn, false)
.await
.map_err(|e| {
error!("Failed to run migrations: {}", e);

@ -54,6 +54,18 @@ const BOOTSTRAP_ROOT_USERNAME: &str = "e2e_root";
const TEST_CREDENTIAL_ENCRYPTION_KEY: &str =
"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef";
fn trusted_dynamic_sql(sql: String) -> sqlx::AssertSqlSafe<String> {
sqlx::AssertSqlSafe(sql)
}
fn quote_pg_ident(identifier: &str) -> String {
format!("\"{}\"", identifier.replace('"', "\"\""))
}
fn quote_pg_literal(value: &str) -> String {
format!("'{}'", value.replace('\'', "''"))
}
struct TestOpaqueCipherSuite;
impl CipherSuite for TestOpaqueCipherSuite {
@ -5999,23 +6011,25 @@ async fn full_stack_cli_db_status_fails_when_migration_metadata_is_unreadable()
let limited_role = format!("status_reader_{}", unique_test_suffix());
let limited_password = "StatusPwd12345!";
let limited_role_ident = quote_pg_ident(&limited_role);
let database_ident = quote_pg_ident(postgres.database_name());
let limited_password_literal = quote_pg_literal(limited_password);
let admin_pool = connect_test_pool_url(&database_url).await;
sqlx::query(&format!(
"CREATE ROLE \"{limited_role}\" LOGIN PASSWORD '{limited_password}'"
))
sqlx::query(trusted_dynamic_sql(format!(
"CREATE ROLE {limited_role_ident} LOGIN PASSWORD {limited_password_literal}"
)))
.execute(&admin_pool)
.await
.expect("test should create a limited role");
sqlx::query(&format!(
"GRANT CONNECT ON DATABASE \"{}\" TO \"{limited_role}\"",
postgres.database_name()
))
sqlx::query(trusted_dynamic_sql(format!(
"GRANT CONNECT ON DATABASE {database_ident} TO {limited_role_ident}"
)))
.execute(&admin_pool)
.await
.expect("test should grant database connect");
sqlx::query(&format!(
"GRANT USAGE ON SCHEMA public TO \"{limited_role}\""
))
sqlx::query(trusted_dynamic_sql(format!(
"GRANT USAGE ON SCHEMA public TO {limited_role_ident}"
)))
.execute(&admin_pool)
.await
.expect("test should grant schema usage without table reads");

Loading…
Cancel
Save