diff --git a/src/Simplex/Messaging/Server/MsgStore/Postgres.hs b/src/Simplex/Messaging/Server/MsgStore/Postgres.hs index a5b05b8e5..a79643f08 100644 --- a/src/Simplex/Messaging/Server/MsgStore/Postgres.hs +++ b/src/Simplex/Messaging/Server/MsgStore/Postgres.hs @@ -114,7 +114,9 @@ instance MsgStoreClass PostgresMsgStore where where st = dbStore $ queueStore_ ms oldMsg = now - ttl - batchSize = 10000 :: Int + -- expired messages read per page in expire_old_messages, and the page is one + -- transaction: queues in it stay row-locked against SEND and ACK until it commits. + batchSize = 100 :: Int toMessageStats (expiredMsgsCount, storedMsgsCount, storedQueues) = MessageStats {expiredMsgsCount, storedMsgsCount, storedQueues} @@ -363,8 +365,8 @@ deleteAllMessages ms = db [sql| UPDATE msg_queues - SET msg_queue_size = 0, msg_can_write = TRUE, msg_queue_expire = FALSE - WHERE msg_queue_size != 0 OR msg_can_write = FALSE OR msg_queue_expire = TRUE + SET msg_queue_size = 0, msg_can_write = TRUE + WHERE msg_queue_size != 0 OR msg_can_write = FALSE |] updateQueueCounts :: PostgresMsgStore -> IO () @@ -384,16 +386,15 @@ updateQueueCounts ms = db [sql| UPDATE msg_queues - SET msg_queue_size = 0, msg_can_write = TRUE, msg_queue_expire = FALSE - WHERE msg_queue_size != 0 OR msg_can_write = FALSE OR msg_queue_expire = TRUE + SET msg_queue_size = 0, msg_can_write = TRUE + WHERE msg_queue_size != 0 OR msg_can_write = FALSE |] void $ DB.execute_ db [sql| UPDATE msg_queues q SET msg_queue_size = s.size, - msg_can_write = s.quota_count = 0, - msg_queue_expire = s.size > s.quota_count + msg_can_write = s.quota_count = 0 FROM queue_stats s WHERE q.recipient_id = s.recipient_id |] diff --git a/src/Simplex/Messaging/Server/QueueStore/Postgres/Migrations.hs b/src/Simplex/Messaging/Server/QueueStore/Postgres/Migrations.hs index a42561430..165623258 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Postgres/Migrations.hs +++ b/src/Simplex/Messaging/Server/QueueStore/Postgres/Migrations.hs @@ -22,7 +22,8 @@ serverSchemaMigrations = ("20250903_store_messages", m20250903_store_messages, Just down_m20250903_store_messages), ("20250915_queue_ids_hash", m20250915_queue_ids_hash, Just down_m20250915_queue_ids_hash), ("20260916_prometheus_indexes", m20260916_prometheus_indexes, Just down_m20260916_prometheus_indexes), - ("20260917_msg_queues_hot", m20260917_msg_queues_hot, Just down_m20260917_msg_queues_hot) + ("20260917_msg_queues_hot", m20260917_msg_queues_hot, Just down_m20260917_msg_queues_hot), + ("20260918_expire_messages", m20260918_expire_messages, Just down_m20260918_expire_messages) ] -- | The list of migrations in ascending order by date @@ -734,3 +735,504 @@ ALTER TABLE msg_queues RESET (fillfactor, autovacuum_vacuum_scale_factor, autova ALTER TABLE messages RESET (autovacuum_vacuum_scale_factor, autovacuum_analyze_scale_factor, toast.autovacuum_vacuum_scale_factor); ALTER TABLE services RESET (fillfactor, autovacuum_vacuum_threshold, autovacuum_vacuum_scale_factor); |] + +m20260918_expire_messages :: Text +m20260918_expire_messages = + [r| +CREATE INDEX idx_messages_expire ON messages (msg_ts, recipient_id) WHERE NOT msg_quota; + +DROP INDEX idx_messages_recipient_id_msg_ts; + +DROP INDEX idx_msg_queues_expire; + +DROP PROCEDURE expire_old_messages(bigint, integer); + +CREATE PROCEDURE expire_old_messages(IN p_old_ts bigint, IN batch_size integer, OUT r_expired_msgs_count bigint, OUT r_stored_msgs_count bigint, OUT r_stored_queues bigint) + LANGUAGE plpgsql + AS $$ +DECLARE + rids BYTEA[]; + rid BYTEA; + last_ts BIGINT := -1; + last_rid BYTEA := '\x'; + next_ts BIGINT; + next_rid BYTEA; + del_count BIGINT; + total_deleted BIGINT := 0; +BEGIN + LOOP + -- The page is scanned in (msg_ts, recipient_id) order, so its last row is the next + -- keyset cursor. Advancing it past every row read, including queues left unexpired + -- because delete_expired_msgs skipped a locked row or raised, is what terminates the + -- loop; re-reading from the start instead would repeat those queues forever. + SELECT array_agg(DISTINCT recipient_id), + (array_agg(msg_ts ORDER BY msg_ts DESC, recipient_id DESC))[1], + (array_agg(recipient_id ORDER BY msg_ts DESC, recipient_id DESC))[1] + INTO rids, next_ts, next_rid + FROM ( + SELECT msg_ts, recipient_id + FROM messages + WHERE NOT msg_quota + AND msg_ts < p_old_ts + AND (msg_ts, recipient_id) > (last_ts, last_rid) + ORDER BY msg_ts ASC, recipient_id ASC + LIMIT batch_size + ) m; + + EXIT WHEN rids IS NULL; + + FOREACH rid IN ARRAY rids + LOOP + BEGIN + del_count := delete_expired_msgs(rid, p_old_ts); + total_deleted := total_deleted + del_count; + EXCEPTION WHEN OTHERS THEN + RAISE WARNING 'STORE, expire_old_messages, error expiring queue %: %', encode(rid, 'base64'), SQLERRM; + CONTINUE; + END; + END LOOP; + COMMIT; + + last_ts := next_ts; + last_rid := next_rid; + END LOOP; + + r_expired_msgs_count := total_deleted; + r_stored_msgs_count := (SELECT COUNT(1) FROM messages); + r_stored_queues := (SELECT COUNT(1) FROM msg_queues WHERE deleted_at IS NULL); +END; +$$; + +CREATE OR REPLACE FUNCTION delete_expired_msgs(p_recipient_id bytea, p_old_ts bigint) RETURNS bigint + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + keep_min_id BIGINT; + del_count BIGINT; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE SKIP LOCKED; + + IF NOT FOUND OR q_size = 0 THEN + RETURN 0; + END IF; + + SELECT MIN(message_id) INTO keep_min_id + FROM messages WHERE recipient_id = p_recipient_id AND msg_ts >= p_old_ts AND msg_quota = FALSE; + + IF keep_min_id IS NULL THEN + DELETE FROM messages WHERE recipient_id = p_recipient_id AND msg_quota = FALSE; + ELSE + DELETE FROM messages WHERE recipient_id = p_recipient_id AND message_id < keep_min_id AND msg_quota = FALSE; + END IF; + + GET DIAGNOSTICS del_count = ROW_COUNT; + IF del_count > 0 THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= del_count, + msg_queue_size = GREATEST(msg_queue_size - del_count, 0) + WHERE recipient_id = p_recipient_id; + END IF; + RETURN del_count; +END; +$$; + +CREATE OR REPLACE FUNCTION write_message(p_recipient_id bytea, p_msg_id bytea, p_msg_ts bigint, p_msg_quota boolean, p_msg_ntf_flag boolean, p_msg_body bytea, p_quota integer) RETURNS TABLE(quota_written boolean, was_empty boolean) + LANGUAGE plpgsql + AS $$ +DECLARE + q_can_write BOOLEAN; + q_size BIGINT; +BEGIN + SELECT msg_can_write, msg_queue_size INTO q_can_write, q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF q_can_write OR q_size = 0 THEN + quota_written := p_msg_quota OR q_size >= p_quota; + was_empty := q_size = 0; + + INSERT INTO messages(recipient_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body) + VALUES (p_recipient_id, p_msg_id, p_msg_ts, quota_written, p_msg_ntf_flag AND NOT quota_written, CASE WHEN quota_written THEN '' :: BYTEA ELSE p_msg_body END); + + UPDATE msg_queues + SET msg_can_write = NOT quota_written, + msg_queue_size = msg_queue_size + 1 + WHERE recipient_id = p_recipient_id; + + RETURN QUERY VALUES (quota_written, was_empty); + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION try_del_msg(p_recipient_id bytea, p_msg_id bytea) RETURNS TABLE(r_msg_id bytea, r_msg_ts bigint, r_msg_quota boolean, r_msg_ntf_flag boolean, r_msg_body bytea) + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + msg RECORD; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF NOT FOUND THEN + RETURN; + END IF; + + SELECT message_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF NOT FOUND THEN + IF q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + RETURN; + END IF; + + IF msg.msg_id = p_msg_id THEN + DELETE FROM messages WHERE message_id = msg.message_id; + IF FOUND THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= 1, + msg_queue_size = GREATEST(msg_queue_size - 1, 0) + WHERE recipient_id = p_recipient_id; + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION try_del_peek_msg(p_recipient_id bytea, p_msg_id bytea) RETURNS TABLE(r_msg_id bytea, r_msg_ts bigint, r_msg_quota boolean, r_msg_ntf_flag boolean, r_msg_body bytea) + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + msg RECORD; + msg_deleted BOOLEAN; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF NOT FOUND THEN + RETURN; + END IF; + + SELECT message_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF NOT FOUND THEN + IF q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + RETURN; + END IF; + + IF msg.msg_id = p_msg_id THEN + DELETE FROM messages WHERE message_id = msg.message_id; + + msg_deleted := FOUND; + IF msg_deleted THEN + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; + + SELECT msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF FOUND THEN + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + IF msg_deleted THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= 1, + msg_queue_size = GREATEST(msg_queue_size - 1, 0) + WHERE recipient_id = p_recipient_id; + END IF; + ELSIF msg_deleted OR q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + ELSE + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; +END; +$$; + +ALTER TABLE msg_queues DROP COLUMN msg_queue_expire; + |] + +down_m20260918_expire_messages :: Text +down_m20260918_expire_messages = + [r| +ALTER TABLE msg_queues ADD COLUMN msg_queue_expire boolean NOT NULL DEFAULT FALSE; + +-- ADD COLUMN already defaulted every row to FALSE, so the backfill only needs the rows +-- that become TRUE, which avoids rewriting the whole table. The flag is only read to pick +-- queues for delete_expired_msgs, and that returns 0 without deleting when msg_queue_size +-- is 0, so excluding empty queues here cannot change what expires. +UPDATE msg_queues q +SET msg_queue_expire = TRUE +WHERE msg_queue_size > 0 + AND EXISTS (SELECT 1 FROM messages m WHERE m.recipient_id = q.recipient_id AND NOT m.msg_quota); + +CREATE INDEX idx_msg_queues_expire ON msg_queues (recipient_id) WHERE deleted_at IS NULL AND msg_queue_expire; + +CREATE INDEX idx_messages_recipient_id_msg_ts ON messages (recipient_id, msg_ts); + +DROP INDEX idx_messages_expire; + +CREATE OR REPLACE FUNCTION delete_expired_msgs(p_recipient_id bytea, p_old_ts bigint) RETURNS bigint + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + keep_min_id BIGINT; + del_count BIGINT; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE SKIP LOCKED; + + IF NOT FOUND OR q_size = 0 THEN + RETURN 0; + END IF; + + SELECT MIN(message_id) INTO keep_min_id + FROM messages WHERE recipient_id = p_recipient_id AND msg_ts >= p_old_ts AND msg_quota = FALSE; + + IF keep_min_id IS NULL THEN + DELETE FROM messages WHERE recipient_id = p_recipient_id AND msg_quota = FALSE; + ELSE + DELETE FROM messages WHERE recipient_id = p_recipient_id AND message_id < keep_min_id AND msg_quota = FALSE; + END IF; + + GET DIAGNOSTICS del_count = ROW_COUNT; + IF del_count > 0 THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= del_count, + msg_queue_expire = msg_queue_size > del_count AND keep_min_id IS NOT NULL, + msg_queue_size = GREATEST(msg_queue_size - del_count, 0) + WHERE recipient_id = p_recipient_id; + END IF; + RETURN del_count; +END; +$$; + +CREATE OR REPLACE FUNCTION write_message(p_recipient_id bytea, p_msg_id bytea, p_msg_ts bigint, p_msg_quota boolean, p_msg_ntf_flag boolean, p_msg_body bytea, p_quota integer) RETURNS TABLE(quota_written boolean, was_empty boolean) + LANGUAGE plpgsql + AS $$ +DECLARE + q_can_write BOOLEAN; + q_size BIGINT; +BEGIN + SELECT msg_can_write, msg_queue_size INTO q_can_write, q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF q_can_write OR q_size = 0 THEN + quota_written := p_msg_quota OR q_size >= p_quota; + was_empty := q_size = 0; + + INSERT INTO messages(recipient_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body) + VALUES (p_recipient_id, p_msg_id, p_msg_ts, quota_written, p_msg_ntf_flag AND NOT quota_written, CASE WHEN quota_written THEN '' :: BYTEA ELSE p_msg_body END); + + UPDATE msg_queues + SET msg_can_write = NOT quota_written, + msg_queue_expire = TRUE, + msg_queue_size = msg_queue_size + 1 + WHERE recipient_id = p_recipient_id; + + RETURN QUERY VALUES (quota_written, was_empty); + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION try_del_msg(p_recipient_id bytea, p_msg_id bytea) RETURNS TABLE(r_msg_id bytea, r_msg_ts bigint, r_msg_quota boolean, r_msg_ntf_flag boolean, r_msg_body bytea) + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + msg RECORD; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF NOT FOUND THEN + RETURN; + END IF; + + SELECT message_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF NOT FOUND THEN + IF q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_expire = FALSE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + RETURN; + END IF; + + IF msg.msg_id = p_msg_id THEN + DELETE FROM messages WHERE message_id = msg.message_id; + IF FOUND THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= 1, + msg_queue_expire = msg_queue_size > 1, + msg_queue_size = GREATEST(msg_queue_size - 1, 0) + WHERE recipient_id = p_recipient_id; + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION try_del_peek_msg(p_recipient_id bytea, p_msg_id bytea) RETURNS TABLE(r_msg_id bytea, r_msg_ts bigint, r_msg_quota boolean, r_msg_ntf_flag boolean, r_msg_body bytea) + LANGUAGE plpgsql + AS $$ +DECLARE + q_size BIGINT; + msg RECORD; + msg_deleted BOOLEAN; +BEGIN + SELECT msg_queue_size INTO q_size + FROM msg_queues + WHERE recipient_id = p_recipient_id AND deleted_at IS NULL + FOR UPDATE; + + IF NOT FOUND THEN + RETURN; + END IF; + + SELECT message_id, msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF NOT FOUND THEN + IF q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_expire = FALSE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + RETURN; + END IF; + + IF msg.msg_id = p_msg_id THEN + DELETE FROM messages WHERE message_id = msg.message_id; + + msg_deleted := FOUND; + IF msg_deleted THEN + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; + + SELECT msg_id, msg_ts, msg_quota, msg_ntf_flag, msg_body + INTO msg + FROM messages + WHERE recipient_id = p_recipient_id + ORDER BY message_id ASC LIMIT 1; + + IF FOUND THEN + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + IF msg_deleted THEN + UPDATE msg_queues + SET msg_can_write = msg_can_write OR msg_queue_size <= 1, + msg_queue_expire = msg_queue_size > 1, + msg_queue_size = GREATEST(msg_queue_size - 1, 0) + WHERE recipient_id = p_recipient_id; + END IF; + ELSIF msg_deleted OR q_size != 0 THEN + UPDATE msg_queues + SET msg_can_write = TRUE, + msg_queue_expire = FALSE, + msg_queue_size = 0 + WHERE recipient_id = p_recipient_id; + END IF; + ELSE + RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); + END IF; +END; +$$; + +DROP PROCEDURE expire_old_messages(bigint, integer); + +CREATE PROCEDURE expire_old_messages(IN p_old_ts bigint, IN batch_size integer, OUT r_expired_msgs_count bigint, OUT r_stored_msgs_count bigint, OUT r_stored_queues bigint) + LANGUAGE plpgsql + AS $$ +DECLARE + rids BYTEA[]; + rid BYTEA; + last_rid BYTEA := '\x'; + del_count BIGINT; + total_deleted BIGINT := 0; +BEGIN + LOOP + SELECT array_agg(recipient_id) + INTO rids + FROM ( + SELECT recipient_id + FROM msg_queues + WHERE deleted_at IS NULL + AND msg_queue_expire = TRUE + AND recipient_id > last_rid + ORDER BY recipient_id ASC + LIMIT batch_size + ) qs; + + EXIT WHEN rids IS NULL OR cardinality(rids) = 0; + + FOREACH rid IN ARRAY rids + LOOP + BEGIN + del_count := delete_expired_msgs(rid, p_old_ts); + total_deleted := total_deleted + del_count; + EXCEPTION WHEN OTHERS THEN + RAISE WARNING 'STORE, expire_old_messages, error expiring queue %: %', encode(rid, 'base64'), SQLERRM; + CONTINUE; + END; + COMMIT; + END LOOP; + last_rid := rids[cardinality(rids)]; + END LOOP; + + r_expired_msgs_count := total_deleted; + r_stored_msgs_count := (SELECT COUNT(1) FROM messages); + r_stored_queues := (SELECT COUNT(1) FROM msg_queues WHERE deleted_at IS NULL); +END; +$$; + |] diff --git a/src/Simplex/Messaging/Server/QueueStore/Postgres/server_schema.sql b/src/Simplex/Messaging/Server/QueueStore/Postgres/server_schema.sql index 6c6cd27b0..c39b82990 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Postgres/server_schema.sql +++ b/src/Simplex/Messaging/Server/QueueStore/Postgres/server_schema.sql @@ -47,7 +47,6 @@ BEGIN IF del_count > 0 THEN UPDATE msg_queues SET msg_can_write = msg_can_write OR msg_queue_size <= del_count, - msg_queue_expire = msg_queue_size > del_count AND keep_min_id IS NOT NULL, msg_queue_size = GREATEST(msg_queue_size - del_count, 0) WHERE recipient_id = p_recipient_id; END IF; @@ -63,24 +62,33 @@ CREATE PROCEDURE smp_server.expire_old_messages(IN p_old_ts bigint, IN batch_siz DECLARE rids BYTEA[]; rid BYTEA; + last_ts BIGINT := -1; last_rid BYTEA := '\x'; + next_ts BIGINT; + next_rid BYTEA; del_count BIGINT; total_deleted BIGINT := 0; BEGIN LOOP - SELECT array_agg(recipient_id) - INTO rids + -- The page is scanned in (msg_ts, recipient_id) order, so its last row is the next + -- keyset cursor. Advancing it past every row read, including queues left unexpired + -- because delete_expired_msgs skipped a locked row or raised, is what terminates the + -- loop; re-reading from the start instead would repeat those queues forever. + SELECT array_agg(DISTINCT recipient_id), + (array_agg(msg_ts ORDER BY msg_ts DESC, recipient_id DESC))[1], + (array_agg(recipient_id ORDER BY msg_ts DESC, recipient_id DESC))[1] + INTO rids, next_ts, next_rid FROM ( - SELECT recipient_id - FROM msg_queues - WHERE deleted_at IS NULL - AND msg_queue_expire = TRUE - AND recipient_id > last_rid - ORDER BY recipient_id ASC + SELECT msg_ts, recipient_id + FROM messages + WHERE NOT msg_quota + AND msg_ts < p_old_ts + AND (msg_ts, recipient_id) > (last_ts, last_rid) + ORDER BY msg_ts ASC, recipient_id ASC LIMIT batch_size - ) qs; + ) m; - EXIT WHEN rids IS NULL OR cardinality(rids) = 0; + EXIT WHEN rids IS NULL; FOREACH rid IN ARRAY rids LOOP @@ -91,9 +99,11 @@ BEGIN RAISE WARNING 'STORE, expire_old_messages, error expiring queue %: %', encode(rid, 'base64'), SQLERRM; CONTINUE; END; - COMMIT; END LOOP; - last_rid := rids[cardinality(rids)]; + COMMIT; + + last_ts := next_ts; + last_rid := next_rid; END LOOP; r_expired_msgs_count := total_deleted; @@ -195,7 +205,6 @@ BEGIN IF q_size != 0 THEN UPDATE msg_queues SET msg_can_write = TRUE, - msg_queue_expire = FALSE, msg_queue_size = 0 WHERE recipient_id = p_recipient_id; END IF; @@ -207,7 +216,6 @@ BEGIN IF FOUND THEN UPDATE msg_queues SET msg_can_write = msg_can_write OR msg_queue_size <= 1, - msg_queue_expire = msg_queue_size > 1, msg_queue_size = GREATEST(msg_queue_size - 1, 0) WHERE recipient_id = p_recipient_id; RETURN QUERY VALUES (msg.msg_id, msg.msg_ts, msg.msg_quota, msg.msg_ntf_flag, msg.msg_body); @@ -245,7 +253,6 @@ BEGIN IF q_size != 0 THEN UPDATE msg_queues SET msg_can_write = TRUE, - msg_queue_expire = FALSE, msg_queue_size = 0 WHERE recipient_id = p_recipient_id; END IF; @@ -271,14 +278,12 @@ BEGIN IF msg_deleted THEN UPDATE msg_queues SET msg_can_write = msg_can_write OR msg_queue_size <= 1, - msg_queue_expire = msg_queue_size > 1, msg_queue_size = GREATEST(msg_queue_size - 1, 0) WHERE recipient_id = p_recipient_id; END IF; ELSIF msg_deleted OR q_size != 0 THEN UPDATE msg_queues SET msg_can_write = TRUE, - msg_queue_expire = FALSE, msg_queue_size = 0 WHERE recipient_id = p_recipient_id; END IF; @@ -348,7 +353,6 @@ BEGIN UPDATE msg_queues SET msg_can_write = NOT quota_written, - msg_queue_expire = TRUE, msg_queue_size = msg_queue_size + 1 WHERE recipient_id = p_recipient_id; @@ -440,7 +444,6 @@ CREATE TABLE smp_server.msg_queues ( rcv_service_id bytea, ntf_service_id bytea, msg_can_write boolean DEFAULT true NOT NULL, - msg_queue_expire boolean DEFAULT false NOT NULL, msg_queue_size bigint DEFAULT 0 NOT NULL ) WITH (fillfactor='80', autovacuum_vacuum_scale_factor='0.02', autovacuum_analyze_scale_factor='0.01', autovacuum_vacuum_cost_limit='1000'); @@ -485,6 +488,10 @@ ALTER TABLE ONLY smp_server.services +CREATE INDEX idx_messages_expire ON smp_server.messages USING btree (msg_ts, recipient_id) WHERE (NOT msg_quota); + + + CREATE INDEX idx_messages_recipient_id_message_id ON smp_server.messages USING btree (recipient_id, message_id); @@ -493,14 +500,6 @@ CREATE INDEX idx_messages_recipient_id_msg_quota ON smp_server.messages USING bt -CREATE INDEX idx_messages_recipient_id_msg_ts ON smp_server.messages USING btree (recipient_id, msg_ts); - - - -CREATE INDEX idx_msg_queues_expire ON smp_server.msg_queues USING btree (recipient_id) WHERE ((deleted_at IS NULL) AND msg_queue_expire); - - - CREATE UNIQUE INDEX idx_msg_queues_link_id ON smp_server.msg_queues USING btree (link_id); diff --git a/tests/CoreTests/MsgStoreTests.hs b/tests/CoreTests/MsgStoreTests.hs index 48fc7810a..05878cc14 100644 --- a/tests/CoreTests/MsgStoreTests.hs +++ b/tests/CoreTests/MsgStoreTests.hs @@ -79,6 +79,7 @@ msgStoreTests = do someMsgStoreTests it "should correctly update message counts and canWrite flag" testUpdateMessageCounts it "tryDelPeekMsg (ACK not from NSE) should reset message counts when queue is empty" testResetMessageCounts + it "should expire messages across commit batches" testExpireMessagesInBatches #endif describe "Journal message store: queue state backup expiration" $ do it "should remove old queue state backups" testRemoveQueueStateBackups @@ -327,35 +328,34 @@ testUpdateMessageCounts ms = do q <- ExceptT $ addQueue ms rId qr let write s = writeMsg ms q True =<< mkMessage s hasSize = checkQueueSize ms - q `hasSize` (0, True, False) + q `hasSize` (0, True) Just (Message {msgId = mId1}, True) <- write "message 1" - q `hasSize` (1, True, True) + q `hasSize` (1, True) Just (Message {msgId = mId2}, False) <- write "message 2" - q `hasSize` (2, True, True) + q `hasSize` (2, True) Just (Message {msgId = mId3}, False) <- write "message 3" - q `hasSize` (3, True, True) + q `hasSize` (3, True) Nothing <- write "message 4" - q `hasSize` (4, False, True) + q `hasSize` (4, False) Msg "message 1" <- tryPeekMsg ms q - q `hasSize` (4, False, True) + q `hasSize` (4, False) Msg "message 1" <- tryDelMsg ms q mId1 - q `hasSize` (3, False, True) + q `hasSize` (3, False) Msg "message 2" <- tryPeekMsg ms q (Msg "message 2", Msg "message 3") <- tryDelPeekMsg ms q mId2 - q `hasSize` (2, False, True) + q `hasSize` (2, False) (Msg "message 3", Just MessageQuota {msgId = mId4}) <- tryDelPeekMsg ms q mId3 - q `hasSize` (1, False, True) + q `hasSize` (1, False) (Just MessageQuota {}, Nothing) <- tryDelPeekMsg ms q mId4 - q `hasSize` (0, True, False) + q `hasSize` (0, True) -checkQueueSize :: PostgresMsgStore -> PostgresQueue -> (Int64, Bool, Bool) -> ExceptT ErrorType IO () -checkQueueSize ms q (size, canWrt, expire) = liftIO $ do - [(size', canWrt', expire')] <- +checkQueueSize :: PostgresMsgStore -> PostgresQueue -> (Int64, Bool) -> ExceptT ErrorType IO () +checkQueueSize ms q (size, canWrt) = liftIO $ do + [(size', canWrt')] <- withTransaction (dbStore $ queueStore ms) $ \db -> - DB.query db "SELECT msg_queue_size, msg_can_write, msg_queue_expire FROM msg_queues WHERE recipient_id = ?" (Only (recipientId q)) + DB.query db "SELECT msg_queue_size, msg_can_write FROM msg_queues WHERE recipient_id = ?" (Only (recipientId q)) size' `shouldBe` size canWrt' `shouldBe` canWrt - expire' `shouldBe` expire testResetMessageCounts :: PostgresMsgStore -> IO () testResetMessageCounts ms = do @@ -369,25 +369,63 @@ testResetMessageCounts ms = do Just (Message {msgId = mId2}, False) <- write "message 2" Just (Message {msgId = mId3}, False) <- write "message 3" Nothing <- write "message 4" - q `hasSize` (4, False, True) + q `hasSize` (4, False) liftIO $ setIncorrectSize q (10, True) Nothing <- write "message 5" - q `hasSize` (11, False, True) + q `hasSize` (11, False) (Msg "message 1", Msg "message 2") <- tryDelPeekMsg ms q mId1 - q `hasSize` (10, False, True) + q `hasSize` (10, False) (Msg "message 2", Msg "message 3") <- tryDelPeekMsg ms q mId2 - q `hasSize` (9, False, True) + q `hasSize` (9, False) (Msg "message 3", Just MessageQuota {msgId = mId4}) <- tryDelPeekMsg ms q mId3 - q `hasSize` (8, False, True) + q `hasSize` (8, False) (Just MessageQuota {}, Just MessageQuota {msgId = mId5}) <- tryDelPeekMsg ms q mId4 - q `hasSize` (7, False, True) + q `hasSize` (7, False) (Just MessageQuota {}, Nothing) <- tryDelPeekMsg ms q mId5 - q `hasSize` (0, True, False) -- reset + q `hasSize` (0, True) -- reset where setIncorrectSize :: PostgresQueue -> (Int64, Bool) -> IO () setIncorrectSize q (size, canWrt) = void $ withTransaction (dbStore $ queueStore ms) $ \db -> DB.execute db "UPDATE msg_queues SET msg_queue_size = ?, msg_can_write = ? WHERE recipient_id = ?" (size, canWrt, recipientId q) + +testExpireMessagesInBatches :: PostgresMsgStore -> IO () +testExpireMessagesInBatches ms = do + g <- C.newRandom + emptiedQs <- replicateM emptiedCount $ newQueue g + partialQs <- replicateM partialCount $ newQueue g + overQuotaQs <- replicateM overQuotaCount $ newQueue g + quotaMsgs <- runRight $ do + forM_ (emptiedQs <> partialQs) $ \q -> void $ write q "old 1" + forM_ emptiedQs $ \q -> void $ write q "old 2" + mapM fillPastQuota overQuotaQs + -- msg_ts has second granularity, so the recent messages need a new second to be kept + threadDelay 1100000 + boundary <- systemSeconds <$> getSystemTime + runRight_ $ forM_ partialQs $ \q -> void $ write q "recent" + + MessageStats {expiredMsgsCount, storedMsgsCount, storedQueues} <- expireOldMessages False ms boundary 0 + expiredMsgsCount `shouldBe` (emptiedCount * 2 + partialCount + sum quotaMsgs) + storedMsgsCount `shouldBe` (partialCount + overQuotaCount) -- recent messages and quota markers + storedQueues `shouldBe` (emptiedCount + partialCount + overQuotaCount) + runRight_ $ do + forM_ emptiedQs $ \q -> checkQueueSize ms q (0, True) + forM_ partialQs $ \q -> checkQueueSize ms q (1, True) + -- the quota marker is never expired, and the queue stays blocked until it is acked + forM_ overQuotaQs $ \q -> checkQueueSize ms q (1, False) + where + -- expire_old_messages pages through expired messages and commits per page, so these + -- counts put a page boundary inside each group of queues. + emptiedCount = 120 :: Int + partialCount = 40 :: Int + overQuotaCount = 10 :: Int + newQueue g = do + (rId, qr) <- testNewQueueRec g QMMessaging + runRight $ ExceptT $ addQueue ms rId qr + write q s = writeMsg ms q True =<< mkMessage s + fillPastQuota q = go 0 + where + go n = write q "fill" >>= maybe (pure n) (const $ go (n + 1)) #endif testQueueState :: JournalMsgStore s -> IO () diff --git a/tests/Test.hs b/tests/Test.hs index 5c135df4f..16375bcc6 100644 --- a/tests/Test.hs +++ b/tests/Test.hs @@ -109,7 +109,8 @@ main = do describe "SMP server schema dump" $ postgresSchemaDumpTest serverMigrations - [ "20250320_short_links" -- snd_secure moves to the bottom on down migration + [ "20250320_short_links", -- snd_secure moves to the bottom on down migration + "20260918_expire_messages" -- msg_queue_expire moves to the bottom on down migration ] -- skipComparisonForDownMigrations testStoreDBOpts "src/Simplex/Messaging/Server/QueueStore/Postgres/server_schema.sql"