diff --git a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts index 726e77b39..d3f97b80a 100644 --- a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts +++ b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts @@ -693,6 +693,19 @@ WHERE idx.indrelid = to_regclass($1) } } +async function postgresKvRowKeyIsCCollated(db: PostgresQueryable, kvTable: string): Promise { + const result = await db.query<{c_collated: boolean}>( + `SELECT col.collname = 'C' AND col.collnamespace = 'pg_catalog'::regnamespace AS c_collated +FROM pg_attribute att +JOIN pg_collation col ON col.oid = att.attcollation +WHERE att.attrelid = to_regclass($1) + AND att.attname = 'row_key' + AND NOT att.attisdropped`, + [kvTable], + ); + return result.rows[0]?.c_collated === true; +} + export async function ensurePostgresKvSchema(client: IPostgresClient): Promise { const kvTable = client.kvTable(); const table = quoteIdentifier(kvTable); @@ -703,8 +716,8 @@ export async function ensurePostgresKvSchema(client: IPostgresClient): Promise { ) d`); const schemaDiff = await raw.query<{n: string}>(` SELECT count(*) AS n FROM ( - (SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes WHERE tablename = '${OLD}' + (SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes + WHERE tablename = '${OLD}' AND indexname <> '${OLD}_row_key_c_idx' EXCEPT SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}') UNION ALL (SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}' EXCEPT SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes WHERE tablename = '${OLD}') ) d`); + const collations = await raw.query<{tablename: string; attname: string; collname: string}>(` + SELECT cls.relname AS tablename, att.attname, col.collname + FROM pg_attribute att + JOIN pg_class cls ON cls.oid = att.attrelid + JOIN pg_collation col ON col.oid = att.attcollation + WHERE att.attrelid IN ('${OLD}'::regclass, '${NEW}'::regclass) + AND att.attname IN ('partition_key', 'row_key') + ORDER BY cls.relname, att.attname`); + const cIndexes = await raw.query<{tablename: string}>( + `SELECT tablename FROM pg_indexes WHERE indexname IN ('${OLD}_row_key_c_idx', '${NEW}_row_key_c_idx')`, + ); console.log(`key diffs=${diff.rows[0]!.n} index-definition diffs=${schemaDiff.rows[0]!.n}`); + console.log(collations.rows.map((r) => `${r.tablename}.${r.attname}=${r.collname}`).join(' ')); expect(Number(diff.rows[0]!.n)).toBe(0); expect(Number(schemaDiff.rows[0]!.n)).toBe(0); + expect(collations.rows.filter((r) => r.tablename === OLD).map((r) => r.collname)).toEqual(['default', 'default']); + expect(collations.rows.filter((r) => r.tablename === NEW).map((r) => r.collname)).toEqual(['C', 'C']); + expect(cIndexes.rows.map((r) => r.tablename)).toEqual([OLD]); + }); + + it('never recollates an existing table and keeps its C index', async () => { + const OLD = `${KV}_keep_old`; + await raw.query(`DROP TABLE IF EXISTS ${OLD}`); + const oldClient = new TableClient(raw, OLD); + await legacyEnsurePostgresKvSchema(oldClient); + await ensurePostgresKvSchema(oldClient); + await ensurePostgresKvSchema(oldClient); + const collations = await raw.query<{attname: string; collname: string}>(` + SELECT att.attname, col.collname + FROM pg_attribute att + JOIN pg_collation col ON col.oid = att.attcollation + WHERE att.attrelid = '${OLD}'::regclass AND att.attname IN ('partition_key', 'row_key') + ORDER BY att.attname`); + const indexes = await raw.query<{indexname: string}>( + `SELECT indexname FROM pg_indexes WHERE tablename = '${OLD}' ORDER BY indexname`, + ); + console.log(`after boot on a legacy table: ${indexes.rows.map((r) => r.indexname).join(', ')}`); + expect(collations.rows.map((r) => r.collname)).toEqual(['default', 'default']); + expect(indexes.rows.map((r) => r.indexname)).toContain(`${OLD}_row_key_c_idx`); }); it('survives three simultaneous boots (two api replicas and a worker)', async () => { diff --git a/fluxer_svc/src/postgres.rs b/fluxer_svc/src/postgres.rs index 48bf9ca08..decc0908d 100644 --- a/fluxer_svc/src/postgres.rs +++ b/fluxer_svc/src/postgres.rs @@ -3,7 +3,7 @@ use crate::config::ServiceConfig; use anyhow::Context; use chrono::{DateTime, Utc}; -use deadpool_postgres::{Client, Manager, Pool, Runtime}; +use deadpool_postgres::{Client, Manager, Pool, Runtime, Transaction}; use rustls::RootCertStore; use serde_json::{Map, Number, Value}; use std::io::Cursor; @@ -129,6 +129,25 @@ fn build_disabled_tls_connector() -> MakeRustlsConnect { MakeRustlsConnect::new(tls_config) } +async fn row_key_is_c_collated( + transaction: &Transaction<'_>, + kv_table: &str, +) -> anyhow::Result { + let row = transaction + .query_opt( + r#"SELECT col.collname = 'C' AND col.collnamespace = 'pg_catalog'::regnamespace AS c_collated +FROM pg_attribute att +JOIN pg_collation col ON col.oid = att.attcollation +WHERE att.attrelid = to_regclass($1) + AND att.attname = 'row_key' + AND NOT att.attisdropped"#, + &[&kv_table], + ) + .await + .context("failed to inspect Postgres KV row_key collation")?; + Ok(row.and_then(|row| row.get::<_, Option>("c_collated")) == Some(true)) +} + pub async fn ensure_kv_schema(pool: &Pool, kv_table: &str) -> anyhow::Result<()> { let table = quote_identifier(kv_table)?; let old_partition_index = quote_identifier(&format!("{kv_table}_partition_idx"))?; @@ -166,15 +185,29 @@ pub async fn ensure_kv_schema(pool: &Pool, kv_table: &str) -> anyhow::Result<()> r#" CREATE TABLE IF NOT EXISTS {table} ( table_name text NOT NULL, - partition_key text NOT NULL, - row_key text NOT NULL, + partition_key text COLLATE "C" NOT NULL, + row_key text COLLATE "C" NOT NULL, row_data jsonb NOT NULL, expires_at timestamptz, updated_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (table_name, row_key) ); CREATE INDEX IF NOT EXISTS {partition_row_index} ON {table} (table_name, partition_key, row_key); -CREATE INDEX IF NOT EXISTS {row_key_c_index} ON {table} (table_name, row_key COLLATE "C"); +"# + )) + .await + .context("failed to ensure Postgres KV schema")?; + if !row_key_is_c_collated(&transaction, kv_table).await? { + transaction + .batch_execute(&format!( + r#"CREATE INDEX IF NOT EXISTS {row_key_c_index} ON {table} (table_name, row_key COLLATE "C");"# + )) + .await + .context("failed to ensure Postgres KV schema")?; + } + transaction + .batch_execute(&format!( + r#" CREATE INDEX IF NOT EXISTS {expires_index} ON {table} (expires_at) WHERE expires_at IS NOT NULL; CREATE INDEX IF NOT EXISTS {messages_message_index} ON {table} (partition_key, ((CASE WHEN row_data -> 'message_id' ->> 'value' ~ '^-?[0-9]+$' THEN (row_data -> 'message_id' ->> 'value')::bigint END))) WHERE table_name = 'messages'; CREATE INDEX IF NOT EXISTS {message_reactions_message_index} ON {table} (partition_key, ((CASE WHEN row_data -> 'message_id' ->> 'value' ~ '^-?[0-9]+$' THEN (row_data -> 'message_id' ->> 'value')::bigint END))) WHERE table_name = 'message_reactions';