diff --git a/.sqlx/query-9f1b41e87c6a136f880ff89356befe073e9c994a370a63b5a44402f596847283.json b/.sqlx/query-09f137012e80545f7a8ff1e0351a98707061855241689ae26f08a89262d88762.json similarity index 92% rename from .sqlx/query-9f1b41e87c6a136f880ff89356befe073e9c994a370a63b5a44402f596847283.json rename to .sqlx/query-09f137012e80545f7a8ff1e0351a98707061855241689ae26f08a89262d88762.json index cab5c10d..2d30d9e3 100644 --- a/.sqlx/query-9f1b41e87c6a136f880ff89356befe073e9c994a370a63b5a44402f596847283.json +++ b/.sqlx/query-09f137012e80545f7a8ff1e0351a98707061855241689ae26f08a89262d88762.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $7\n AND ($2 OR status <> $3)\n AND (created_at, id) > ($4, $5)\n ORDER BY created_at ASC, id ASC\n LIMIT $6\n ", + "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $7\n AND ($2 OR status <> $3)\n AND (created_at, id) > ($4, $5)\n AND ($8::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $8 AND ub.blocked_user_id = user_id\n ))\n ORDER BY created_at ASC, id ASC\n LIMIT $6\n ", "describe": { "columns": [ { @@ -188,7 +188,8 @@ "Timestamptz", "Int8", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -210,5 +211,5 @@ false ] }, - "hash": "9f1b41e87c6a136f880ff89356befe073e9c994a370a63b5a44402f596847283" + "hash": "09f137012e80545f7a8ff1e0351a98707061855241689ae26f08a89262d88762" } diff --git a/.sqlx/query-0e75cb57e7d05f2dcee261da1b49ea171eefe562d61bfa69971af46a0ce5b732.json b/.sqlx/query-0e75cb57e7d05f2dcee261da1b49ea171eefe562d61bfa69971af46a0ce5b732.json new file mode 100644 index 00000000..695b07bd --- /dev/null +++ b/.sqlx/query-0e75cb57e7d05f2dcee261da1b49ea171eefe562d61bfa69971af46a0ce5b732.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM user_blocks WHERE blocker_user_id = $1 AND blocked_user_id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [] + }, + "hash": "0e75cb57e7d05f2dcee261da1b49ea171eefe562d61bfa69971af46a0ce5b732" +} diff --git a/.sqlx/query-54313833a6cee855e8de7b0c6dad45c278b6cf13481bcdf608016a46e08079f6.json b/.sqlx/query-1da0b11f8c3874d241457aea05bd30d4d17e5b7995f702cd839a3fd3b70eea04.json similarity index 63% rename from .sqlx/query-54313833a6cee855e8de7b0c6dad45c278b6cf13481bcdf608016a46e08079f6.json rename to .sqlx/query-1da0b11f8c3874d241457aea05bd30d4d17e5b7995f702cd839a3fd3b70eea04.json index d7d5fbcb..cfa30880 100644 --- a/.sqlx/query-54313833a6cee855e8de7b0c6dad45c278b6cf13481bcdf608016a46e08079f6.json +++ b/.sqlx/query-1da0b11f8c3874d241457aea05bd30d4d17e5b7995f702cd839a3fd3b70eea04.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $4\n AND status <> $2\n AND (user_id IS NULL OR user_id <> $3)\n AND created_at >= NOW() - INTERVAL '90 days'\n ", + "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $4\n AND status <> $2\n AND (user_id IS NULL OR user_id <> $3)\n AND NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $3 AND ub.blocked_user_id = user_id\n )\n AND created_at >= NOW() - INTERVAL '90 days'\n ", "describe": { "columns": [ { @@ -22,5 +22,5 @@ null ] }, - "hash": "54313833a6cee855e8de7b0c6dad45c278b6cf13481bcdf608016a46e08079f6" + "hash": "1da0b11f8c3874d241457aea05bd30d4d17e5b7995f702cd839a3fd3b70eea04" } diff --git a/.sqlx/query-faac538ae5f2f341021d114379930307d8291e161bc8609717e6c57dc72275a6.json b/.sqlx/query-23e80a3bfca3636d2364dcc63d6b31e4816c363fc32e78cdc61b6c2508cc621f.json similarity index 86% rename from .sqlx/query-faac538ae5f2f341021d114379930307d8291e161bc8609717e6c57dc72275a6.json rename to .sqlx/query-23e80a3bfca3636d2364dcc63d6b31e4816c363fc32e78cdc61b6c2508cc621f.json index 051ff21b..5563d41f 100644 --- a/.sqlx/query-faac538ae5f2f341021d114379930307d8291e161bc8609717e6c57dc72275a6.json +++ b/.sqlx/query-23e80a3bfca3636d2364dcc63d6b31e4816c363fc32e78cdc61b6c2508cc621f.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT p.room_id AS \"room_id!: RoomId\",\n p.message_id AS \"message_id!\",\n p.message_created_at AS \"message_created_at!\",\n p.pinned_by AS \"pinned_by?: UserId\",\n u.username AS pinned_by_username,\n p.note,\n p.pinned_at AS \"pinned_at!\"\n FROM chat_message_pins p\n LEFT JOIN users u ON u.id = p.pinned_by\n JOIN chat_messages m\n ON m.room_id = p.room_id\n AND m.id = p.message_id\n AND m.created_at = p.message_created_at\n WHERE p.room_id = $1\n AND m.status <> $2\n AND m.deletion_source IS DISTINCT FROM $4\n ORDER BY p.pinned_at DESC, p.message_id DESC\n LIMIT $3\n ", + "query": "\n SELECT p.room_id AS \"room_id!: RoomId\",\n p.message_id AS \"message_id!\",\n p.message_created_at AS \"message_created_at!\",\n p.pinned_by AS \"pinned_by?: UserId\",\n u.username AS pinned_by_username,\n p.note,\n p.pinned_at AS \"pinned_at!\"\n FROM chat_message_pins p\n LEFT JOIN users u ON u.id = p.pinned_by\n JOIN chat_messages m\n ON m.room_id = p.room_id\n AND m.id = p.message_id\n AND m.created_at = p.message_created_at\n WHERE p.room_id = $1\n AND m.status <> $2\n AND m.deletion_source IS DISTINCT FROM $4\n AND ($5::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $5 AND ub.blocked_user_id = m.user_id\n ))\n ORDER BY p.pinned_at DESC, p.message_id DESC\n LIMIT $3\n ", "describe": { "columns": [ { @@ -86,7 +86,8 @@ "Int8", "Int2", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -99,5 +100,5 @@ false ] }, - "hash": "faac538ae5f2f341021d114379930307d8291e161bc8609717e6c57dc72275a6" + "hash": "23e80a3bfca3636d2364dcc63d6b31e4816c363fc32e78cdc61b6c2508cc621f" } diff --git a/.sqlx/query-432cb4e024309fb2a20867254851910b82ee90135b1b2f4b02383d5ccf28ea64.json b/.sqlx/query-432cb4e024309fb2a20867254851910b82ee90135b1b2f4b02383d5ccf28ea64.json new file mode 100644 index 00000000..fec4fafa --- /dev/null +++ b/.sqlx/query-432cb4e024309fb2a20867254851910b82ee90135b1b2f4b02383d5ccf28ea64.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT blocker_user_id AS \"blocker_user_id!: UserId\"\n FROM user_blocks\n WHERE blocked_user_id = $1 AND blocker_user_id = ANY($2::bigint[])\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "blocker_user_id!: UserId", + "type_info": "Int8", + "origin": { + "Table": { + "table": "user_blocks", + "name": "blocker_user_id" + } + } + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8Array" + ] + }, + "nullable": [ + false + ] + }, + "hash": "432cb4e024309fb2a20867254851910b82ee90135b1b2f4b02383d5ccf28ea64" +} diff --git a/.sqlx/query-70d766384a03de012c33b406e9f9b7715a01132bc138cb5ad37e510c30ed543b.json b/.sqlx/query-4a9fb3744b1434c5a5f8235f7885892f18ed30fc1390748b21eb24faf54f84d2.json similarity index 92% rename from .sqlx/query-70d766384a03de012c33b406e9f9b7715a01132bc138cb5ad37e510c30ed543b.json rename to .sqlx/query-4a9fb3744b1434c5a5f8235f7885892f18ed30fc1390748b21eb24faf54f84d2.json index 395ae36d..f1f38118 100644 --- a/.sqlx/query-70d766384a03de012c33b406e9f9b7715a01132bc138cb5ad37e510c30ed543b.json +++ b/.sqlx/query-4a9fb3744b1434c5a5f8235f7885892f18ed30fc1390748b21eb24faf54f84d2.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH search_terms AS (\n SELECT websearch_to_tsquery('simple', $2) AS tsquery\n )\n SELECT m.id AS \"id!\",\n m.room_id AS \"room_id!: RoomId\",\n m.user_id AS \"user_id?: UserId\",\n m.client_message_id,\n m.content AS \"content!\",\n m.message_type AS \"message_type!: ChatMessageType\",\n m.status AS \"status!: ChatMessageStatus\",\n m.version AS \"version!\",\n m.reply_to_message_id,\n m.reply_to_message_created_at,\n m.metadata AS \"metadata?: ChatMetadata\",\n m.edited_at,\n m.deleted_at,\n m.deleted_by AS \"deleted_by?: UserId\",\n m.delete_reason,\n m.created_at AS \"created_at!\"\n FROM chat_messages m\n CROSS JOIN search_terms st\n WHERE m.room_id = $1\n AND m.deletion_source IS DISTINCT FROM $8\n AND (m.content_search @@ st.tsquery OR m.content ILIKE $3 ESCAPE '\\')\n AND ($4 OR m.status <> $5)\n AND ($6::bigint IS NULL OR m.user_id = $6)\n ORDER BY m.created_at DESC, m.id DESC\n LIMIT $7\n ", + "query": "\n WITH search_terms AS (\n SELECT websearch_to_tsquery('simple', $2) AS tsquery\n )\n SELECT m.id AS \"id!\",\n m.room_id AS \"room_id!: RoomId\",\n m.user_id AS \"user_id?: UserId\",\n m.client_message_id,\n m.content AS \"content!\",\n m.message_type AS \"message_type!: ChatMessageType\",\n m.status AS \"status!: ChatMessageStatus\",\n m.version AS \"version!\",\n m.reply_to_message_id,\n m.reply_to_message_created_at,\n m.metadata AS \"metadata?: ChatMetadata\",\n m.edited_at,\n m.deleted_at,\n m.deleted_by AS \"deleted_by?: UserId\",\n m.delete_reason,\n m.created_at AS \"created_at!\"\n FROM chat_messages m\n CROSS JOIN search_terms st\n WHERE m.room_id = $1\n AND m.deletion_source IS DISTINCT FROM $8\n AND (m.content_search @@ st.tsquery OR m.content ILIKE $3 ESCAPE '\\')\n AND ($4 OR m.status <> $5)\n AND ($6::bigint IS NULL OR m.user_id = $6)\n AND ($9::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $9 AND ub.blocked_user_id = m.user_id\n ))\n ORDER BY m.created_at DESC, m.id DESC\n LIMIT $7\n ", "describe": { "columns": [ { @@ -189,7 +189,8 @@ "Int2", "Int8", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -211,5 +212,5 @@ false ] }, - "hash": "70d766384a03de012c33b406e9f9b7715a01132bc138cb5ad37e510c30ed543b" + "hash": "4a9fb3744b1434c5a5f8235f7885892f18ed30fc1390748b21eb24faf54f84d2" } diff --git a/.sqlx/query-5d28adbefb0cf1072239211e30ba61f1fcea221319825a68d6721629a585db13.json b/.sqlx/query-5d28adbefb0cf1072239211e30ba61f1fcea221319825a68d6721629a585db13.json new file mode 100644 index 00000000..73319de1 --- /dev/null +++ b/.sqlx/query-5d28adbefb0cf1072239211e30ba61f1fcea221319825a68d6721629a585db13.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT EXISTS (\n SELECT 1 FROM user_blocks\n WHERE blocker_user_id = $1 AND blocked_user_id = $2\n ) AS \"is_blocked!\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "is_blocked!", + "type_info": "Bool", + "origin": "Expression" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [ + null + ] + }, + "hash": "5d28adbefb0cf1072239211e30ba61f1fcea221319825a68d6721629a585db13" +} diff --git a/.sqlx/query-a5209a475ad100c1d7348d9a6a9d9e5b9ad22fe88c407c16390c3c50325ac59a.json b/.sqlx/query-8a8b021567a21b203a9b1bd33e54f1cdb89c313b94909ade8f69f168011c05e7.json similarity index 75% rename from .sqlx/query-a5209a475ad100c1d7348d9a6a9d9e5b9ad22fe88c407c16390c3c50325ac59a.json rename to .sqlx/query-8a8b021567a21b203a9b1bd33e54f1cdb89c313b94909ade8f69f168011c05e7.json index d64e6218..6de07305 100644 --- a/.sqlx/query-a5209a475ad100c1d7348d9a6a9d9e5b9ad22fe88c407c16390c3c50325ac59a.json +++ b/.sqlx/query-8a8b021567a21b203a9b1bd33e54f1cdb89c313b94909ade8f69f168011c05e7.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_message_events e\n JOIN chat_messages m\n ON m.room_id = e.room_id\n AND m.id = e.message_id\n AND m.created_at = e.message_created_at\n WHERE e.room_id = $1\n AND e.sequence > $2\n AND e.event_type = $3\n AND m.status <> $4\n AND m.deletion_source IS DISTINCT FROM $6\n AND (m.user_id IS NULL OR m.user_id <> $5)\n ", + "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_message_events e\n JOIN chat_messages m\n ON m.room_id = e.room_id\n AND m.id = e.message_id\n AND m.created_at = e.message_created_at\n WHERE e.room_id = $1\n AND e.sequence > $2\n AND e.event_type = $3\n AND m.status <> $4\n AND m.deletion_source IS DISTINCT FROM $6\n AND (m.user_id IS NULL OR m.user_id <> $5)\n AND NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $5 AND ub.blocked_user_id = m.user_id\n )\n ", "describe": { "columns": [ { @@ -24,5 +24,5 @@ null ] }, - "hash": "a5209a475ad100c1d7348d9a6a9d9e5b9ad22fe88c407c16390c3c50325ac59a" + "hash": "8a8b021567a21b203a9b1bd33e54f1cdb89c313b94909ade8f69f168011c05e7" } diff --git a/.sqlx/query-ac367a4f3f53aa56fb3283be08d4c8ddd3773180bcf60d7552e8552911d9fc71.json b/.sqlx/query-ac367a4f3f53aa56fb3283be08d4c8ddd3773180bcf60d7552e8552911d9fc71.json new file mode 100644 index 00000000..2bd65cd8 --- /dev/null +++ b/.sqlx/query-ac367a4f3f53aa56fb3283be08d4c8ddd3773180bcf60d7552e8552911d9fc71.json @@ -0,0 +1,28 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT blocked_user_id AS \"blocked_user_id!: UserId\"\n FROM user_blocks\n WHERE blocker_user_id = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "blocked_user_id!: UserId", + "type_info": "Int8", + "origin": { + "Table": { + "table": "user_blocks", + "name": "blocked_user_id" + } + } + } + ], + "parameters": { + "Left": [ + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "ac367a4f3f53aa56fb3283be08d4c8ddd3773180bcf60d7552e8552911d9fc71" +} diff --git a/.sqlx/query-4ff5fa022f22d915a9e5b0a8a0779b5cba054363719b751f2942bfa7b36a9690.json b/.sqlx/query-c0aa17f25a96b767446606d62a5fc95bf2c997ff4743123f2539b50000a0b897.json similarity index 92% rename from .sqlx/query-4ff5fa022f22d915a9e5b0a8a0779b5cba054363719b751f2942bfa7b36a9690.json rename to .sqlx/query-c0aa17f25a96b767446606d62a5fc95bf2c997ff4743123f2539b50000a0b897.json index 6c4bda6d..4922792e 100644 --- a/.sqlx/query-4ff5fa022f22d915a9e5b0a8a0779b5cba054363719b751f2942bfa7b36a9690.json +++ b/.sqlx/query-c0aa17f25a96b767446606d62a5fc95bf2c997ff4743123f2539b50000a0b897.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH search_terms AS (\n SELECT websearch_to_tsquery('simple', $2) AS tsquery\n )\n SELECT m.id AS \"id!\",\n m.room_id AS \"room_id!: RoomId\",\n m.user_id AS \"user_id?: UserId\",\n m.client_message_id,\n m.content AS \"content!\",\n m.message_type AS \"message_type!: ChatMessageType\",\n m.status AS \"status!: ChatMessageStatus\",\n m.version AS \"version!\",\n m.reply_to_message_id,\n m.reply_to_message_created_at,\n m.metadata AS \"metadata?: ChatMetadata\",\n m.edited_at,\n m.deleted_at,\n m.deleted_by AS \"deleted_by?: UserId\",\n m.delete_reason,\n m.created_at AS \"created_at!\"\n FROM chat_messages m\n CROSS JOIN search_terms st\n WHERE m.room_id = $1\n AND m.deletion_source IS DISTINCT FROM $10\n AND (m.content_search @@ st.tsquery OR m.content ILIKE $3 ESCAPE '\\')\n AND ($4 OR m.status <> $5)\n AND ($6::bigint IS NULL OR m.user_id = $6)\n AND (m.created_at, m.id) < ($7, $8)\n ORDER BY m.created_at DESC, m.id DESC\n LIMIT $9\n ", + "query": "\n WITH search_terms AS (\n SELECT websearch_to_tsquery('simple', $2) AS tsquery\n )\n SELECT m.id AS \"id!\",\n m.room_id AS \"room_id!: RoomId\",\n m.user_id AS \"user_id?: UserId\",\n m.client_message_id,\n m.content AS \"content!\",\n m.message_type AS \"message_type!: ChatMessageType\",\n m.status AS \"status!: ChatMessageStatus\",\n m.version AS \"version!\",\n m.reply_to_message_id,\n m.reply_to_message_created_at,\n m.metadata AS \"metadata?: ChatMetadata\",\n m.edited_at,\n m.deleted_at,\n m.deleted_by AS \"deleted_by?: UserId\",\n m.delete_reason,\n m.created_at AS \"created_at!\"\n FROM chat_messages m\n CROSS JOIN search_terms st\n WHERE m.room_id = $1\n AND m.deletion_source IS DISTINCT FROM $10\n AND (m.content_search @@ st.tsquery OR m.content ILIKE $3 ESCAPE '\\')\n AND ($4 OR m.status <> $5)\n AND ($6::bigint IS NULL OR m.user_id = $6)\n AND (m.created_at, m.id) < ($7, $8)\n AND ($11::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $11 AND ub.blocked_user_id = m.user_id\n ))\n ORDER BY m.created_at DESC, m.id DESC\n LIMIT $9\n ", "describe": { "columns": [ { @@ -191,7 +191,8 @@ "Timestamptz", "Int8", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -213,5 +214,5 @@ false ] }, - "hash": "4ff5fa022f22d915a9e5b0a8a0779b5cba054363719b751f2942bfa7b36a9690" + "hash": "c0aa17f25a96b767446606d62a5fc95bf2c997ff4743123f2539b50000a0b897" } diff --git a/.sqlx/query-1f0885abb44aa0d2886957a27c29a0d3aa9dd070bd5724062a03e56ceb314bdc.json b/.sqlx/query-d8d3abfc9af6f5923299099d52af3b50983ffb1568559d8fc73ae662fe3aae55.json similarity index 91% rename from .sqlx/query-1f0885abb44aa0d2886957a27c29a0d3aa9dd070bd5724062a03e56ceb314bdc.json rename to .sqlx/query-d8d3abfc9af6f5923299099d52af3b50983ffb1568559d8fc73ae662fe3aae55.json index 4588bb92..ee544fb3 100644 --- a/.sqlx/query-1f0885abb44aa0d2886957a27c29a0d3aa9dd070bd5724062a03e56ceb314bdc.json +++ b/.sqlx/query-d8d3abfc9af6f5923299099d52af3b50983ffb1568559d8fc73ae662fe3aae55.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $6\n AND ($2 OR status <> $3)\n AND created_at >= NOW() - INTERVAL '90 days'\n AND message_type = ANY($5::smallint[])\n ORDER BY created_at DESC, id DESC\n LIMIT $4\n ", + "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $6\n AND ($2 OR status <> $3)\n AND created_at >= NOW() - INTERVAL '90 days'\n AND message_type = ANY($5::smallint[])\n AND ($7::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $7 AND ub.blocked_user_id = user_id\n ))\n ORDER BY created_at DESC, id DESC\n LIMIT $4\n ", "describe": { "columns": [ { @@ -187,7 +187,8 @@ "Int2", "Int8", "Int2Array", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -209,5 +210,5 @@ false ] }, - "hash": "1f0885abb44aa0d2886957a27c29a0d3aa9dd070bd5724062a03e56ceb314bdc" + "hash": "d8d3abfc9af6f5923299099d52af3b50983ffb1568559d8fc73ae662fe3aae55" } diff --git a/.sqlx/query-81a1e2c3230e42a42679934f5ec9a481d65e5f04e3310196af2975fc4566515f.json b/.sqlx/query-e0cba9a4ceaa4e76ae1631787821d6146765431d82465f9910e79e110a634e58.json similarity index 92% rename from .sqlx/query-81a1e2c3230e42a42679934f5ec9a481d65e5f04e3310196af2975fc4566515f.json rename to .sqlx/query-e0cba9a4ceaa4e76ae1631787821d6146765431d82465f9910e79e110a634e58.json index d0daaa1a..324c6a5a 100644 --- a/.sqlx/query-81a1e2c3230e42a42679934f5ec9a481d65e5f04e3310196af2975fc4566515f.json +++ b/.sqlx/query-e0cba9a4ceaa4e76ae1631787821d6146765431d82465f9910e79e110a634e58.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $7\n AND ($2 OR status <> $3)\n AND (created_at, id) < ($4, $5)\n ORDER BY created_at DESC, id DESC\n LIMIT $6\n ", + "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $7\n AND ($2 OR status <> $3)\n AND (created_at, id) < ($4, $5)\n AND ($8::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $8 AND ub.blocked_user_id = user_id\n ))\n ORDER BY created_at DESC, id DESC\n LIMIT $6\n ", "describe": { "columns": [ { @@ -188,7 +188,8 @@ "Timestamptz", "Int8", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -210,5 +211,5 @@ false ] }, - "hash": "81a1e2c3230e42a42679934f5ec9a481d65e5f04e3310196af2975fc4566515f" + "hash": "e0cba9a4ceaa4e76ae1631787821d6146765431d82465f9910e79e110a634e58" } diff --git a/.sqlx/query-cc1770204bc4418886bfffea97a54c3d7e2de4ed3aa0a06941aad7fde4e97689.json b/.sqlx/query-e6fcd030a26be30f0c10685a546d7875d62da5bdeace8ea475ef66e5a0ef4126.json similarity index 72% rename from .sqlx/query-cc1770204bc4418886bfffea97a54c3d7e2de4ed3aa0a06941aad7fde4e97689.json rename to .sqlx/query-e6fcd030a26be30f0c10685a546d7875d62da5bdeace8ea475ef66e5a0ef4126.json index 54562330..940994e6 100644 --- a/.sqlx/query-cc1770204bc4418886bfffea97a54c3d7e2de4ed3aa0a06941aad7fde4e97689.json +++ b/.sqlx/query-e6fcd030a26be30f0c10685a546d7875d62da5bdeace8ea475ef66e5a0ef4126.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $6\n AND status <> $2\n AND (user_id IS NULL OR user_id <> $3)\n AND (created_at, id) > ($4, $5)\n ", + "query": "\n SELECT COUNT(*) AS \"count!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $6\n AND status <> $2\n AND (user_id IS NULL OR user_id <> $3)\n AND (created_at, id) > ($4, $5)\n AND NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $3 AND ub.blocked_user_id = user_id\n )\n ", "describe": { "columns": [ { @@ -24,5 +24,5 @@ null ] }, - "hash": "cc1770204bc4418886bfffea97a54c3d7e2de4ed3aa0a06941aad7fde4e97689" + "hash": "e6fcd030a26be30f0c10685a546d7875d62da5bdeace8ea475ef66e5a0ef4126" } diff --git a/.sqlx/query-f8ebc75fcd508e0378a2ed629ff7a90ce40bc53dd7d45622f31c00682f225a2f.json b/.sqlx/query-f8ebc75fcd508e0378a2ed629ff7a90ce40bc53dd7d45622f31c00682f225a2f.json new file mode 100644 index 00000000..900cae9c --- /dev/null +++ b/.sqlx/query-f8ebc75fcd508e0378a2ed629ff7a90ce40bc53dd7d45622f31c00682f225a2f.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT EXISTS (\n SELECT 1\n FROM user_blocks\n WHERE blocker_user_id = $1 AND blocked_user_id = $2\n ) AS \"is_blocking!\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "is_blocking!", + "type_info": "Bool", + "origin": "Expression" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [ + null + ] + }, + "hash": "f8ebc75fcd508e0378a2ed629ff7a90ce40bc53dd7d45622f31c00682f225a2f" +} diff --git a/.sqlx/query-c37c8b961170e850d399cb2acd8b1443dd6131c58ea1bf063427cea5d32ac713.json b/.sqlx/query-fc90493f312b6e89150d75be2d68b78e7e643546abf0e28a724e3023f8b216d7.json similarity index 92% rename from .sqlx/query-c37c8b961170e850d399cb2acd8b1443dd6131c58ea1bf063427cea5d32ac713.json rename to .sqlx/query-fc90493f312b6e89150d75be2d68b78e7e643546abf0e28a724e3023f8b216d7.json index 0dd79942..1b8f0a6a 100644 --- a/.sqlx/query-c37c8b961170e850d399cb2acd8b1443dd6131c58ea1bf063427cea5d32ac713.json +++ b/.sqlx/query-fc90493f312b6e89150d75be2d68b78e7e643546abf0e28a724e3023f8b216d7.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $8\n AND ($2 OR status <> $3)\n AND (created_at, id) < ($4, $5)\n AND message_type = ANY($7::smallint[])\n ORDER BY created_at DESC, id DESC\n LIMIT $6\n ", + "query": "\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $8\n AND ($2 OR status <> $3)\n AND (created_at, id) < ($4, $5)\n AND message_type = ANY($7::smallint[])\n AND ($9::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $9 AND ub.blocked_user_id = user_id\n ))\n ORDER BY created_at DESC, id DESC\n LIMIT $6\n ", "describe": { "columns": [ { @@ -189,7 +189,8 @@ "Int8", "Int8", "Int2Array", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -211,5 +212,5 @@ false ] }, - "hash": "c37c8b961170e850d399cb2acd8b1443dd6131c58ea1bf063427cea5d32ac713" + "hash": "fc90493f312b6e89150d75be2d68b78e7e643546abf0e28a724e3023f8b216d7" } diff --git a/.sqlx/query-ff2c3fb9de287a18385f07764131b26a93791a0db392c19c8fce842c11c6dc0f.json b/.sqlx/query-ff2c3fb9de287a18385f07764131b26a93791a0db392c19c8fce842c11c6dc0f.json new file mode 100644 index 00000000..5904fc58 --- /dev/null +++ b/.sqlx/query-ff2c3fb9de287a18385f07764131b26a93791a0db392c19c8fce842c11c6dc0f.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO user_blocks (blocker_user_id, blocked_user_id)\n VALUES ($1, $2)\n ON CONFLICT (blocker_user_id, blocked_user_id)\n DO UPDATE SET blocker_user_id = EXCLUDED.blocker_user_id\n RETURNING created_at AS \"created_at!\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "created_at!", + "type_info": "Timestamptz", + "origin": { + "Table": { + "table": "user_blocks", + "name": "created_at" + } + } + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "ff2c3fb9de287a18385f07764131b26a93791a0db392c19c8fce842c11c6dc0f" +} diff --git a/.sqlx/query-8d4a4914427075127768d878a628470d8ae24168f9ac340ca3be9a30cadecc11.json b/.sqlx/query-ff5dbc825d2ff6b934b6346af2ef599df70a6d46ce6b434cafb4c26e171248b3.json similarity index 81% rename from .sqlx/query-8d4a4914427075127768d878a628470d8ae24168f9ac340ca3be9a30cadecc11.json rename to .sqlx/query-ff5dbc825d2ff6b934b6346af2ef599df70a6d46ce6b434cafb4c26e171248b3.json index 9975ff7f..889d286e 100644 --- a/.sqlx/query-8d4a4914427075127768d878a628470d8ae24168f9ac340ca3be9a30cadecc11.json +++ b/.sqlx/query-ff5dbc825d2ff6b934b6346af2ef599df70a6d46ce6b434cafb4c26e171248b3.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH candidates AS (\n SELECT id,\n room_id,\n user_id,\n client_message_id,\n content,\n message_type,\n status,\n version,\n reply_to_message_id,\n reply_to_message_created_at,\n metadata,\n edited_at,\n deleted_at,\n deleted_by,\n delete_reason,\n created_at,\n CASE\n WHEN jsonb_typeof(metadata #> '{playback,positionSeconds}') = 'number'\n THEN (metadata #>> '{playback,positionSeconds}')::double precision\n ELSE NULL\n END AS playback_position\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $11\n AND ($2 OR status <> $3)\n AND ($4::text IS NULL OR metadata #>> '{playback,mediaId}' = $4)\n AND ($5::text IS NULL OR metadata #>> '{playback,playlistId}' = $5)\n AND ($6::text IS NULL OR metadata #>> '{playback,targetHash}' = $6)\n AND message_type = ANY($7::smallint[])\n )\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM candidates\n WHERE playback_position BETWEEN $8 AND $9\n ORDER BY playback_position ASC, created_at ASC, id ASC\n LIMIT $10\n ", + "query": "\n WITH candidates AS (\n SELECT id,\n room_id,\n user_id,\n client_message_id,\n content,\n message_type,\n status,\n version,\n reply_to_message_id,\n reply_to_message_created_at,\n metadata,\n edited_at,\n deleted_at,\n deleted_by,\n delete_reason,\n created_at,\n CASE\n WHEN jsonb_typeof(metadata #> '{playback,positionSeconds}') = 'number'\n THEN (metadata #>> '{playback,positionSeconds}')::double precision\n ELSE NULL\n END AS playback_position\n FROM chat_messages\n WHERE room_id = $1\n AND deletion_source IS DISTINCT FROM $11\n AND ($2 OR status <> $3)\n AND ($4::text IS NULL OR metadata #>> '{playback,mediaId}' = $4)\n AND ($5::text IS NULL OR metadata #>> '{playback,playlistId}' = $5)\n AND ($6::text IS NULL OR metadata #>> '{playback,targetHash}' = $6)\n AND message_type = ANY($7::smallint[])\n AND ($12::bigint IS NULL OR NOT EXISTS (\n SELECT 1 FROM user_blocks ub\n WHERE ub.blocker_user_id = $12 AND ub.blocked_user_id = user_id\n ))\n )\n SELECT id AS \"id!\",\n room_id AS \"room_id!: RoomId\",\n user_id AS \"user_id?: UserId\",\n client_message_id,\n content AS \"content!\",\n message_type AS \"message_type!: ChatMessageType\",\n status AS \"status!: ChatMessageStatus\",\n version AS \"version!\",\n reply_to_message_id,\n reply_to_message_created_at,\n metadata AS \"metadata?: ChatMetadata\",\n edited_at,\n deleted_at,\n deleted_by AS \"deleted_by?: UserId\",\n delete_reason,\n created_at AS \"created_at!\"\n FROM candidates\n WHERE playback_position BETWEEN $8 AND $9\n ORDER BY playback_position ASC, created_at ASC, id ASC\n LIMIT $10\n ", "describe": { "columns": [ { @@ -192,7 +192,8 @@ "Float8", "Float8", "Int8", - "Int2" + "Int2", + "Int8" ] }, "nullable": [ @@ -214,5 +215,5 @@ false ] }, - "hash": "8d4a4914427075127768d878a628470d8ae24168f9ac340ca3be9a30cadecc11" + "hash": "ff5dbc825d2ff6b934b6346af2ef599df70a6d46ce6b434cafb4c26e171248b3" } diff --git a/migrations/20260821001_create_user_blocks.sql b/migrations/20260821001_create_user_blocks.sql new file mode 100644 index 00000000..8f08b353 --- /dev/null +++ b/migrations/20260821001_create_user_blocks.sql @@ -0,0 +1,13 @@ +CREATE TABLE IF NOT EXISTS user_blocks ( + blocker_user_id BIGINT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + blocked_user_id BIGINT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (blocker_user_id, blocked_user_id), + CONSTRAINT user_blocks_no_self_block CHECK (blocker_user_id <> blocked_user_id) +); + +CREATE INDEX IF NOT EXISTS idx_user_blocks_blocker_created + ON user_blocks(blocker_user_id, created_at DESC, blocked_user_id DESC); + +CREATE INDEX IF NOT EXISTS idx_user_blocks_blocked + ON user_blocks(blocked_user_id, blocker_user_id); diff --git a/synctv-api-common/src/impls/admin/rooms.rs b/synctv-api-common/src/impls/admin/rooms.rs index 723c762b..56117894 100644 --- a/synctv-api-common/src/impls/admin/rooms.rs +++ b/synctv-api-common/src/impls/admin/rooms.rs @@ -90,6 +90,7 @@ impl AdminApiImpl { crate::impls::proto_validated_user_id(creator_id, &self.public_id_codec) }) .transpose()?, + excluded_creator_ids: Vec::new(), category_id: parse_optional_room_category_id(&req.category_id, &self.public_id_codec)?, label_ids: parse_room_label_ids(&req.label_ids, &self.public_id_codec)?, sort_by: proto_admin_room_list_sort_by(req.sort_by)?, diff --git a/synctv-api-common/src/impls/client/convert.rs b/synctv-api-common/src/impls/client/convert.rs index f71fcdf3..82e5ace9 100644 --- a/synctv-api-common/src/impls/client/convert.rs +++ b/synctv-api-common/src/impls/client/convert.rs @@ -2581,6 +2581,7 @@ pub fn try_room_to_proto_with_availability_and_presence( .map(|label| room_label_to_proto(label, public_id_codec)) .collect::, _>>()?, is_public: Some(room.is_public), + creator_blocked: false, }) } diff --git a/synctv-api-common/src/impls/client/room.rs b/synctv-api-common/src/impls/client/room.rs index 8734434e..3f51375f 100644 --- a/synctv-api-common/src/impls/client/room.rs +++ b/synctv-api-common/src/impls/client/room.rs @@ -2,7 +2,7 @@ use crate::impls::ApiError; use futures::{stream, StreamExt as _, TryStreamExt as _}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use synctv_core::models::{ ChatMentionInput, ChatMessageEvent, ChatMessageType, ChatMessageWithAttachments, ChatPinEvent, ChatPlaybackMessagesQuery, CreateChatAttachmentUploadSession, MarkChatRead, PageParams, @@ -222,6 +222,7 @@ impl ClientApiImpl { async fn room_to_proto_for_favorite_view( &self, room: &synctv_core::models::Room, + creator_blocked: bool, ) -> Result { let (settings, presence, availability, member_count) = tokio::join!( self.room_service.get_room_settings(&room.id), @@ -233,15 +234,18 @@ impl ClientApiImpl { let presence = presence.map_err(ApiError::from)?; let availability = availability.map_err(ApiError::from)?; let member_count = member_count?; - self.room_to_proto_with_availability_presence_and_loaded_cover( - room, - Some(&settings), - member_count, - availability, - Some(&presence), - None, - ) - .await + let mut proto_room = self + .room_to_proto_with_availability_presence_and_loaded_cover( + room, + Some(&settings), + member_count, + availability, + Some(&presence), + None, + ) + .await?; + proto_room.creator_blocked = creator_blocked; + Ok(proto_room) } async fn load_room_playback_state_proto( @@ -576,11 +580,21 @@ impl ClientApiImpl { } else if room.status.is_closed() { return Err(ApiError::NotFound("Room not found".to_string())); } - self.rooms_to_discovery_items(vec![room], viewer_id, DiscoveryReadConsistency::Primary) + let creator_blocked = self + .user_service + .is_blocking(viewer_id, &room.created_by) + .await + .map_err(ApiError::from)?; + let mut item = self + .rooms_to_discovery_items(vec![room], viewer_id, DiscoveryReadConsistency::Primary) .await? .into_iter() .next() - .ok_or_else(|| ApiError::NotFound("Room not found".to_string())) + .ok_or_else(|| ApiError::NotFound("Room not found".to_string()))?; + if let Some(room) = item.room.as_mut() { + room.creator_blocked = creator_blocked; + } + Ok(item) } pub async fn get_public_room_discovery( @@ -613,9 +627,17 @@ impl ClientApiImpl { async fn discover_room_window( &self, + viewer_id: Option<&UserId>, req: synctv_proto::client::DiscoverRoomsRequest, ) -> Result { - let query = build_room_discovery_query(req, &self.public_id_codec)?; + let mut query = build_room_discovery_query(req, &self.public_id_codec)?; + if let Some(viewer_id) = viewer_id { + query.excluded_creator_ids = self + .user_service + .blocked_user_ids_eventually_consistent(viewer_id) + .await + .map_err(ApiError::from)?; + } let hot_stats = self .presence_service .hot_room_stats() @@ -701,7 +723,7 @@ impl ClientApiImpl { viewer_id: &UserId, req: synctv_proto::client::DiscoverRoomsRequest, ) -> Result { - let window = self.discover_room_window(req).await?; + let window = self.discover_room_window(Some(viewer_id), req).await?; let (featured_rooms, rooms) = tokio::try_join!( self.rooms_to_discovery_items( window.featured_rooms, @@ -725,7 +747,7 @@ impl ClientApiImpl { &self, req: synctv_proto::client::DiscoverRoomsRequest, ) -> Result { - let window = self.discover_room_window(req).await?; + let window = self.discover_room_window(None, req).await?; let guest_enabled = self.public_guest_access_enabled()?; let (featured_rooms, rooms) = tokio::try_join!( self.rooms_to_public_discovery_items( @@ -764,32 +786,45 @@ impl ClientApiImpl { rooms.iter().map(|(room, _, _)| room.id).collect(); let room_models: Vec = rooms.iter().map(|(room, _, _)| room.clone()).collect(); - let (room_settings_map, presence_stats, availability_map, creator_views, favorite_room_ids) = - tokio::try_join!( - async { - self.room_service - .get_room_settings_batch_eventually_consistent(&room_ids) - .await - .map_err(ApiError::from) - }, - async { - self.presence_service - .room_stats_batch(&room_ids) - .await - .map_err(ApiError::from) - }, - async { - self.room_service - .room_availability_batch_eventually_consistent(&room_models) - .await - .map_err(ApiError::from) - }, - self.load_room_creator_public_views( - &room_models, - DiscoveryReadConsistency::EventuallyConsistent, - ), - self.favorite_room_ids_for_viewer_eventually_consistent(Some(uid), &room_ids), - )?; + let ( + room_settings_map, + presence_stats, + availability_map, + creator_views, + favorite_room_ids, + blocked_creator_ids, + ) = tokio::try_join!( + async { + self.room_service + .get_room_settings_batch_eventually_consistent(&room_ids) + .await + .map_err(ApiError::from) + }, + async { + self.presence_service + .room_stats_batch(&room_ids) + .await + .map_err(ApiError::from) + }, + async { + self.room_service + .room_availability_batch_eventually_consistent(&room_models) + .await + .map_err(ApiError::from) + }, + self.load_room_creator_public_views( + &room_models, + DiscoveryReadConsistency::EventuallyConsistent, + ), + self.favorite_room_ids_for_viewer_eventually_consistent(Some(uid), &room_ids), + async { + self.user_service + .blocked_user_ids_eventually_consistent(&uid) + .await + .map(|ids| ids.into_iter().collect::>()) + .map_err(ApiError::from) + }, + )?; let presence_by_room: HashMap = presence_stats .iter() .map(|stats| (stats.room_id, stats)) @@ -809,6 +844,7 @@ impl ClientApiImpl { let presence_by_room = &presence_by_room; let availability_map = &availability_map; let favorite_room_ids = &favorite_room_ids; + let blocked_creator_ids = &blocked_creator_ids; async move { let settings = required_room_settings(room_settings_map, &room.id)?; let availability = required_room_availability(availability_map, &room.id)?; @@ -822,7 +858,7 @@ impl ClientApiImpl { } else { synctv_proto::client::MyRoomRelation::Participating as i32 }; - let proto_room = self + let mut proto_room = self .room_to_proto_with_availability_presence_and_loaded_cover( &room, Some(settings), @@ -832,6 +868,7 @@ impl ClientApiImpl { Some(creator), ) .await?; + proto_room.creator_blocked = blocked_creator_ids.contains(&room.created_by); Ok::<_, ApiError>(synctv_proto::client::MyRoom { room: Some(proto_room), permissions, @@ -864,7 +901,14 @@ impl ClientApiImpl { .favorite_room(&uid, &room_id) .await .map_err(ApiError::from)?; - let proto_room = self.room_to_proto_for_favorite_view(&room).await?; + let creator_blocked = self + .user_service + .is_blocking(user_id, &room.created_by) + .await + .map_err(ApiError::from)?; + let proto_room = self + .room_to_proto_for_favorite_view(&room, creator_blocked) + .await?; Ok(synctv_proto::client::FavoriteRoomResponse { room: Some(proto_room), }) @@ -883,7 +927,14 @@ impl ClientApiImpl { .unfavorite_room(&uid, &room_id) .await .map_err(ApiError::from)?; - let proto_room = self.room_to_proto_for_favorite_view(&room).await?; + let creator_blocked = self + .user_service + .is_blocking(user_id, &room.created_by) + .await + .map_err(ApiError::from)?; + let proto_room = self + .room_to_proto_for_favorite_view(&room, creator_blocked) + .await?; Ok(synctv_proto::client::UnfavoriteRoomResponse { room: Some(proto_room), }) @@ -909,11 +960,23 @@ impl ClientApiImpl { .list_favorite_rooms(&uid, pagination, search.as_deref()) .await .map_err(ApiError::from)?; + let blocked_creator_ids = self + .user_service + .blocked_user_ids(user_id) + .await + .map_err(ApiError::from)? + .into_iter() + .collect::>(); let rooms_ref = &rooms; + let blocked_creator_ids = &blocked_creator_ids; let response_rooms = stream::iter(0..rooms.len()) .map(|index| async move { let room = &rooms_ref[index]; - self.room_to_proto_for_favorite_view(room).await + self.room_to_proto_for_favorite_view( + room, + blocked_creator_ids.contains(&room.created_by), + ) + .await }) .buffered(16) .try_collect() @@ -1061,7 +1124,17 @@ impl ClientApiImpl { None => Ok(false), } }; - let (playback_state, settings, presence, member_count, favorited) = tokio::try_join!( + let creator_blocked = async { + match actor.user_id() { + Some(user_id) => self + .user_service + .is_blocking(&user_id, &room.created_by) + .await + .map_err(ApiError::from), + None => Ok(false), + } + }; + let (playback_state, settings, presence, member_count, favorited, creator_blocked) = tokio::try_join!( self.load_room_playback_state_proto(&rid), async { self.room_service @@ -1077,6 +1150,7 @@ impl ClientApiImpl { }, self.load_room_member_count(&rid), favorite, + creator_blocked, )?; let availability = self .room_service @@ -1084,7 +1158,7 @@ impl ClientApiImpl { .await .map_err(ApiError::from)?; - let proto_room = self + let mut proto_room = self .room_to_proto_with_availability_presence_and_loaded_cover( &room, Some(&settings), @@ -1094,6 +1168,7 @@ impl ClientApiImpl { None, ) .await?; + proto_room.creator_blocked = creator_blocked; Ok(synctv_proto::client::GetRoomResponse { room: Some(proto_room), playback_state: Some(playback_state), @@ -3120,6 +3195,13 @@ mod tests { ) .await, )?; + api_ok( + api.user_service + .block_user(&member.id, &owner.id) + .await + .map(|_| ()) + .map_err(crate::impls::ApiError::from), + )?; let favorites_before_leave = api_ok( api.list_favorite_rooms( @@ -3134,6 +3216,21 @@ mod tests { )?; assert_eq!(favorites_before_leave.total, 1); assert_eq!(favorites_before_leave.rooms.len(), 1); + assert!(favorites_before_leave.rooms[0].creator_blocked); + + let discovery = api_ok( + api.get_room_discovery( + &member.id, + synctv_proto::client::GetRoomDiscoveryRequest { + room_id: room_id.clone(), + }, + ) + .await, + )?; + assert!(discovery + .room + .as_ref() + .is_some_and(|room| room.creator_blocked)); api_ok(api.leave_room(&member.id, &room_id).await)?; diff --git a/synctv-api-common/src/impls/client/user.rs b/synctv-api-common/src/impls/client/user.rs index 7daab5f8..ed5b0eca 100644 --- a/synctv-api-common/src/impls/client/user.rs +++ b/synctv-api-common/src/impls/client/user.rs @@ -418,6 +418,98 @@ impl ClientApiImpl { self.user_to_proto_with_avatar(&user).await } + pub async fn block_user( + &self, + user_id: &UserId, + req: synctv_proto::client::BlockUserRequest, + ) -> Result { + crate::impls::validate_proto_request(&req)?; + let blocked_user_id = + crate::impls::proto_validated_user_id(req.user_id, &self.public_id_codec)?; + let blocked_at = self + .user_service + .block_user(user_id, &blocked_user_id) + .await + .map_err(ApiError::from)?; + let blocked_user = self + .user_service + .get_user(&blocked_user_id) + .await + .map_err(ApiError::from)?; + let public_user = self + .user_public_view_with_loaded_avatar(&blocked_user) + .await?; + + Ok(synctv_proto::client::BlockUserResponse { + blocked_user: Some(synctv_proto::client::BlockedUser { + user: Some(public_user), + blocked_at: blocked_at.timestamp(), + }), + }) + } + + pub async fn unblock_user( + &self, + user_id: &UserId, + req: synctv_proto::client::UnblockUserRequest, + ) -> Result { + crate::impls::validate_proto_request(&req)?; + let blocked_user_id = + crate::impls::proto_validated_user_id(req.user_id, &self.public_id_codec)?; + self.user_service + .unblock_user(user_id, &blocked_user_id) + .await + .map_err(ApiError::from)?; + Ok(synctv_proto::client::UnblockUserResponse { success: true }) + } + + pub async fn list_blocked_users( + &self, + user_id: &UserId, + req: synctv_proto::client::ListBlockedUsersRequest, + ) -> Result { + crate::impls::validate_proto_request(&req)?; + let page = if req.page > 0 { + req.page.cast_unsigned() + } else { + 1 + }; + let page_size = if req.page_size > 0 { + req.page_size.cast_unsigned().min(100) + } else { + 50 + }; + let pagination = PageParams::new(Some(page), Some(page_size)); + let search = (!req.search.is_empty()).then_some(req.search); + let (blocked_users, total) = self + .user_service + .list_blocked_users(user_id, pagination, search.as_deref()) + .await + .map_err(ApiError::from)?; + let users = blocked_users + .iter() + .map(|blocked| blocked.user.clone()) + .collect::>(); + let public_users = self + .batch_user_public_views_with_loaded_avatars(&users) + .await?; + let users = blocked_users + .into_iter() + .zip(public_users) + .map(|(blocked, user)| synctv_proto::client::BlockedUser { + user: Some(user), + blocked_at: blocked.blocked_at.timestamp(), + }) + .collect(); + + Ok(synctv_proto::client::ListBlockedUsersResponse { + users, + total: i32::try_from(total).map_err(|_| { + ApiError::Internal("blocked user total exceeds i32::MAX".to_string()) + })?, + }) + } + pub async fn get_user_preferences( &self, user_id: &UserId, diff --git a/synctv-api-common/src/impls/messaging/resource_observer.rs b/synctv-api-common/src/impls/messaging/resource_observer.rs index 68a685b6..83e532d2 100644 --- a/synctv-api-common/src/impls/messaging/resource_observer.rs +++ b/synctv-api-common/src/impls/messaging/resource_observer.rs @@ -1548,6 +1548,51 @@ impl MediaResourceHub { event_cursor: Option, ) -> ResourceRefreshOutcome { let subscriptions = self.snapshot_subscriptions().await; + let chat_author_id = match event { + RealtimeEvent::ChatMessageEvent { event, .. } => event.message.message.user_id, + RealtimeEvent::ChatPinEvent { event, .. } => event.message.message.user_id, + _ => None, + }; + let mut blocked_chat_viewers = HashSet::new(); + if let Some(chat_author_id) = chat_author_id { + let observers = subscriptions + .iter() + .filter_map(|(_, observer, _, _)| observer.upgrade()) + .collect::>(); + let viewer_ids = observers + .iter() + .filter_map(|observer| observer.actor.user_id()) + .collect::>() + .into_iter() + .collect::>(); + blocked_chat_viewers = if let Some(chat_service) = observers + .iter() + .find_map(|observer| observer.chat_service.as_ref()) + { + match chat_service + .blocking_viewer_ids(&viewer_ids, &chat_author_id) + .await + { + Ok(ids) => ids.into_iter().collect(), + Err(error) => { + tracing::warn!( + room_id = %self.room_id, + user_id = %chat_author_id, + error = %error, + "Failed to resolve chat block visibility; suppressing the event" + ); + viewer_ids.into_iter().collect() + } + } + } else { + tracing::warn!( + room_id = %self.room_id, + user_id = %chat_author_id, + "Chat service unavailable while resolving block visibility; suppressing the event" + ); + viewer_ids.into_iter().collect() + }; + } let mut refresh_plan = HashMap::< ResourceSubscriberKey, (Weak, ResourceObservation, bool, u64), @@ -1709,6 +1754,28 @@ impl MediaResourceHub { ) { continue; } + if observer + .actor + .user_id() + .is_some_and(|user_id| blocked_chat_viewers.contains(&user_id)) + { + if matches!( + self.send_and_commit_subscription_update( + key, + *revision, + &observer, + updated_observation.clone(), + None, + ) + .await, + SubscriptionRefreshCommit::Committed + ) { + observer + .replace_local_observation(updated_observation) + .await; + } + continue; + } let event_payload = match observer.chat_message_event_to_proto(chat_event).await { Ok(event) => event, @@ -1765,6 +1832,28 @@ impl MediaResourceHub { ) { continue; } + if observer + .actor + .user_id() + .is_some_and(|user_id| blocked_chat_viewers.contains(&user_id)) + { + if matches!( + self.send_and_commit_subscription_update( + key, + *revision, + &observer, + updated_observation.clone(), + None, + ) + .await, + SubscriptionRefreshCommit::Committed + ) { + observer + .replace_local_observation(updated_observation) + .await; + } + continue; + } let event_payload = match observer.chat_pin_event_to_proto(event).await { Ok(event) => event, Err(error) => { @@ -2534,6 +2623,15 @@ impl ResourceObserver { let Some(chat_service) = self.chat_service.as_ref() else { return Ok(()); }; + let blocked_user_ids = match self.actor.user_id() { + Some(user_id) => chat_service + .blocked_user_ids(&user_id) + .await + .map_err(|error| error.to_string())? + .into_iter() + .collect::>(), + None => HashSet::new(), + }; if !observe_id.is_empty() { if !chat_service .is_event_sequence_retained_for_room(&self.room_id, after_event_sequence) @@ -2572,6 +2670,16 @@ impl ResourceObserver { if !Self::apply_event_cursor_to_observation(&mut observation, &cursor) { continue; } + if event + .message + .message + .user_id + .is_some_and(|user_id| blocked_user_ids.contains(&user_id)) + { + self.replace_local_observation(observation.clone()).await; + self.room_hub.register_observation(self, observation).await; + continue; + } self.send_server_message(ServerMessage { message: Some( synctv_proto::client::server_message::Message::ResourceEvent( diff --git a/synctv-api-grpc/src/grpc/client_service/user.rs b/synctv-api-grpc/src/grpc/client_service/user.rs index 5facc836..5b64b391 100644 --- a/synctv-api-grpc/src/grpc/client_service/user.rs +++ b/synctv-api-grpc/src/grpc/client_service/user.rs @@ -4,18 +4,19 @@ use tonic::{Request, Response, Status}; use super::{map_api_error, ClientServiceImpl}; use synctv_api_common::impls::EndpointRateLimitCategory; use synctv_proto::client::{ - user_service_server::UserService, CloseAccountRequest, CloseAccountResponse, - CompleteUserAvatarUploadSessionRequest, CompleteUserAvatarUploadSessionResponse, - ConfirmEmailBindRequest, CreateRoomRequest, CreateUserAvatarUploadSessionRequest, - CreateUserAvatarUploadSessionResponse, DeletePasskeyRequest, DeletePasskeyResponse, - DeleteTotpRequest, DeleteTotpResponse, DiscoverRoomsRequest, DiscoverRoomsResponse, - FavoriteRoomRequest, FavoriteRoomResponse, FinishOpaquePasswordUpdateRequest, - FinishPasskeyBindRequest, FinishRoomPasswordLoginRequest, + user_service_server::UserService, BlockUserRequest, BlockUserResponse, CloseAccountRequest, + CloseAccountResponse, CompleteUserAvatarUploadSessionRequest, + CompleteUserAvatarUploadSessionResponse, ConfirmEmailBindRequest, CreateRoomRequest, + CreateUserAvatarUploadSessionRequest, CreateUserAvatarUploadSessionResponse, + DeletePasskeyRequest, DeletePasskeyResponse, DeleteTotpRequest, DeleteTotpResponse, + DiscoverRoomsRequest, DiscoverRoomsResponse, FavoriteRoomRequest, FavoriteRoomResponse, + FinishOpaquePasswordUpdateRequest, FinishPasskeyBindRequest, FinishRoomPasswordLoginRequest, FinishSensitiveOperationVerificationRequest, FinishTotpSetupRequest, GetProfileRequest, GetRoomDiscoveryRequest, GetRoomRequest, GetRoomResponse, GetUserAvatarObjectRequest, GetUserPreferencesRequest, GetUserPreferencesResponse, JoinRoomRequest, JoinRoomResponse, - ListFavoriteRoomsRequest, ListFavoriteRoomsResponse, ListMyRoomsRequest, ListMyRoomsResponse, - ListPasskeysRequest, ListPasskeysResponse, LogoutRequest, LogoutResponse, PasskeyCredential, + ListBlockedUsersRequest, ListBlockedUsersResponse, ListFavoriteRoomsRequest, + ListFavoriteRoomsResponse, ListMyRoomsRequest, ListMyRoomsResponse, ListPasskeysRequest, + ListPasskeysResponse, LogoutRequest, LogoutResponse, PasskeyCredential, RegenerateTotpRecoveryCodesRequest, RequestSensitiveOperationEmailCodeRequest, RequestSensitiveOperationEmailCodeResponse, Room, RoomDiscoveryItem, SensitiveOperationVerificationOutcome, SetTwoFactorEnabledRequest, SetUsernameRequest, @@ -24,9 +25,10 @@ use synctv_proto::client::{ StartRoomPasswordLoginRequest, StartRoomPasswordLoginResponse, StartSensitiveOperationPasskeyRequest, StartSensitiveOperationPasskeyResponse, StartSensitiveOperationVerificationRequest, StartTotpSetupRequest, StartTotpSetupResponse, - TotpRecoveryCodesResponse, UnbindEmailRequest, UnfavoriteRoomRequest, UnfavoriteRoomResponse, - UpdateUserAvatarRequest, UpdateUserPreferencesRequest, UpdateUserPreferencesResponse, - UploadUserAvatarObjectRequest, UploadUserAvatarObjectResponse, User, UserAvatarObjectResponse, + TotpRecoveryCodesResponse, UnbindEmailRequest, UnblockUserRequest, UnblockUserResponse, + UnfavoriteRoomRequest, UnfavoriteRoomResponse, UpdateUserAvatarRequest, + UpdateUserPreferencesRequest, UpdateUserPreferencesResponse, UploadUserAvatarObjectRequest, + UploadUserAvatarObjectResponse, User, UserAvatarObjectResponse, }; type UserAvatarObjectStream = super::GrpcStatusStream; @@ -707,6 +709,71 @@ impl UserService for ClientServiceImpl { Ok(Response::new(response)) } + async fn block_user( + &self, + request: Request, + ) -> Result, Status> { + let metadata = self.request_metadata(&request)?; + let req = request.into_inner(); + let executor = self.client_api.clone(); + let client_api = self.client_api.clone(); + let response = executor + .execute_user_endpoint( + &metadata, + EndpointRateLimitCategory::Write, + move |authenticated| async move { + client_api.block_user(&authenticated.user_id(), req).await + }, + ) + .await + .map_err(map_api_error)?; + Ok(Response::new(response)) + } + + async fn unblock_user( + &self, + request: Request, + ) -> Result, Status> { + let metadata = self.request_metadata(&request)?; + let req = request.into_inner(); + let executor = self.client_api.clone(); + let client_api = self.client_api.clone(); + let response = executor + .execute_user_endpoint( + &metadata, + EndpointRateLimitCategory::Write, + move |authenticated| async move { + client_api.unblock_user(&authenticated.user_id(), req).await + }, + ) + .await + .map_err(map_api_error)?; + Ok(Response::new(response)) + } + + async fn list_blocked_users( + &self, + request: Request, + ) -> Result, Status> { + let metadata = self.request_metadata(&request)?; + let req = request.into_inner(); + let executor = self.client_api.clone(); + let client_api = self.client_api.clone(); + let response = executor + .execute_user_endpoint( + &metadata, + EndpointRateLimitCategory::Read, + move |authenticated| async move { + client_api + .list_blocked_users(&authenticated.user_id(), req) + .await + }, + ) + .await + .map_err(map_api_error)?; + Ok(Response::new(response)) + } + async fn create_room( &self, request: Request, diff --git a/synctv-api-http/src/http/mod.rs b/synctv-api-http/src/http/mod.rs index 9aaab3da..5dea48ea 100644 --- a/synctv-api-http/src/http/mod.rs +++ b/synctv-api-http/src/http/mod.rs @@ -775,6 +775,14 @@ fn register_playlist_cover_object_routes() -> Router { fn register_extracted_user_routes() -> Router { Router::new() .route("/api/user", get(user::get_me)) + .route( + "/api/user/blocks", + get(user::list_blocked_users).post(user::block_user), + ) + .route( + "/api/user/blocks/{userId}", + axum::routing::delete(user::unblock_user), + ) .route("/api/user/rooms/discover", get(user::discover_rooms)) .route( "/api/user/rooms/{roomId}/discovery", diff --git a/synctv-api-http/src/http/user.rs b/synctv-api-http/src/http/user.rs index dd65761a..83fe9a93 100644 --- a/synctv-api-http/src/http/user.rs +++ b/synctv-api-http/src/http/user.rs @@ -12,17 +12,19 @@ use super::{middleware::RequestMetadata, validation::ProtoQuery, AppResult, AppS use synctv_api_common::impls::EndpointRateLimitCategory; use synctv_proto::client::User; use synctv_proto::client::{ - CloseAccountRequest, CloseAccountResponse, DeletePasskeyRequest, DeletePasskeyResponse, - DeleteTotpRequest, DeleteTotpResponse, DiscoverRoomsRequest, DiscoverRoomsResponse, - FavoriteRoomRequest, FavoriteRoomResponse, FinishPasskeyBindRequest, - FinishSensitiveOperationVerificationRequest, FinishTotpSetupRequest, GetRoomDiscoveryRequest, + BlockUserRequest, BlockUserResponse, CloseAccountRequest, CloseAccountResponse, + DeletePasskeyRequest, DeletePasskeyResponse, DeleteTotpRequest, DeleteTotpResponse, + DiscoverRoomsRequest, DiscoverRoomsResponse, FavoriteRoomRequest, FavoriteRoomResponse, + FinishPasskeyBindRequest, FinishSensitiveOperationVerificationRequest, FinishTotpSetupRequest, + GetRoomDiscoveryRequest, ListBlockedUsersRequest, ListBlockedUsersResponse, ListFavoriteRoomsRequest, ListFavoriteRoomsResponse, ListMyRoomsResponse, ListPasskeysResponse, PasskeyCredential, RequestSensitiveOperationEmailCodeRequest, RequestSensitiveOperationEmailCodeResponse, RoomDiscoveryItem, RoomPathRequest, SensitiveOperationVerificationOutcome, StartPasskeyBindRequest, StartPasskeyBindResponse, StartSensitiveOperationPasskeyRequest, StartSensitiveOperationPasskeyResponse, StartSensitiveOperationVerificationRequest, StartTotpSetupRequest, StartTotpSetupResponse, - TotpRecoveryCodesResponse, UnfavoriteRoomRequest, UnfavoriteRoomResponse, + TotpRecoveryCodesResponse, UnblockUserRequest, UnblockUserResponse, UnfavoriteRoomRequest, + UnfavoriteRoomResponse, }; use synctv_proto::client::{ CompleteUserAvatarUploadSessionRequest, CompleteUserAvatarUploadSessionResponse, @@ -94,6 +96,12 @@ pub struct PasskeyCredentialPath { pub credential_id: String, } +#[derive(Debug, serde::Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct BlockedUserPath { + pub user_id: String, +} + /// Get current user info #[cfg_attr( feature = "openapi", @@ -131,6 +139,102 @@ pub async fn get_me( Ok(Json(response)) } +#[cfg_attr( + feature = "openapi", + utoipa::path( + post, + path = "/api/user/blocks", + tag = "User", + request_body = BlockUserRequest, + responses((status = 200, description = "User blocked", body = BlockUserResponse)), + security(("bearer_auth" = [])) + ) +)] +pub async fn block_user( + request_meta: RequestMetadata, + State(state): State, + Json(req): Json, +) -> AppResult> { + let executor = state.shared_api_runtime.client_api.clone(); + let client_api = state.shared_api_runtime.client_api.clone(); + let response = executor + .execute_user_endpoint( + &request_meta.0, + EndpointRateLimitCategory::Write, + |auth| async move { client_api.block_user(&auth.user_id(), req).await }, + ) + .await + .map_err(super::error::map_api_error)?; + Ok(Json(response)) +} + +#[cfg_attr( + feature = "openapi", + utoipa::path( + delete, + path = "/api/user/blocks/{userId}", + tag = "User", + params(("userId" = String, Path, description = "Public user ID")), + responses((status = 200, description = "User unblocked", body = UnblockUserResponse)), + security(("bearer_auth" = [])) + ) +)] +pub async fn unblock_user( + request_meta: RequestMetadata, + State(state): State, + Path(path): Path, +) -> AppResult> { + let executor = state.shared_api_runtime.client_api.clone(); + let client_api = state.shared_api_runtime.client_api.clone(); + let response = executor + .execute_user_endpoint( + &request_meta.0, + EndpointRateLimitCategory::Write, + |auth| async move { + client_api + .unblock_user( + &auth.user_id(), + UnblockUserRequest { + user_id: path.user_id, + }, + ) + .await + }, + ) + .await + .map_err(super::error::map_api_error)?; + Ok(Json(response)) +} + +#[cfg_attr( + feature = "openapi", + utoipa::path( + get, + path = "/api/user/blocks", + tag = "User", + params(ListBlockedUsersRequest), + responses((status = 200, description = "Blocked users", body = ListBlockedUsersResponse)), + security(("bearer_auth" = [])) + ) +)] +pub async fn list_blocked_users( + request_meta: RequestMetadata, + State(state): State, + ProtoQuery(req): ProtoQuery, +) -> AppResult> { + let executor = state.shared_api_runtime.client_api.clone(); + let client_api = state.shared_api_runtime.client_api.clone(); + let response = executor + .execute_user_endpoint( + &request_meta.0, + EndpointRateLimitCategory::Read, + |auth| async move { client_api.list_blocked_users(&auth.user_id(), req).await }, + ) + .await + .map_err(super::error::map_api_error)?; + Ok(Json(response)) +} + #[cfg_attr( feature = "openapi", utoipa::path( diff --git a/synctv-api-http/src/openapi.rs b/synctv-api-http/src/openapi.rs index e785b7cc..3f1df6c4 100644 --- a/synctv-api-http/src/openapi.rs +++ b/synctv-api-http/src/openapi.rs @@ -67,6 +67,9 @@ pub struct GoogleRpcStatusSchema { ticket::create_ticket, webrtc::get_ice_servers, user::get_me, + user::block_user, + user::unblock_user, + user::list_blocked_users, user::update_user, user::get_user_preferences, user::update_user_preferences, diff --git a/synctv-core/src/models/mod.rs b/synctv-core/src/models/mod.rs index c306bcd0..d6ce5673 100644 --- a/synctv-core/src/models/mod.rs +++ b/synctv-core/src/models/mod.rs @@ -275,7 +275,8 @@ pub use source_config::{ YoutubePlaylistSourceConfig, }; pub use user::{ - SignupMethod, User, UserLifecycleMetadata, UserListQuery, UserListSortBy, UserRole, UserStatus, + BlockedUser, SignupMethod, User, UserLifecycleMetadata, UserListQuery, UserListSortBy, + UserRole, UserStatus, }; pub use user_preferences::{ UserAuthFactors, UserNotificationPreferences, UserPreferences, UserPreferencesUpdate, diff --git a/synctv-core/src/models/room.rs b/synctv-core/src/models/room.rs index 32596b00..1d468753 100644 --- a/synctv-core/src/models/room.rs +++ b/synctv-core/src/models/room.rs @@ -363,6 +363,8 @@ pub struct RoomListQuery { pub is_public: Option, /// Filter by creator pub creator_id: Option, + #[serde(default)] + pub excluded_creator_ids: Vec, pub category_id: Option, #[serde(default)] pub label_ids: Vec, @@ -381,6 +383,7 @@ impl Default for RoomListQuery { is_banned: Some(false), is_public: None, creator_id: None, + excluded_creator_ids: Vec::new(), category_id: None, label_ids: Vec::new(), sort_by: RoomListSortBy::CreatedAt, diff --git a/synctv-core/src/models/user.rs b/synctv-core/src/models/user.rs index f0453adc..22253ae7 100644 --- a/synctv-core/src/models/user.rs +++ b/synctv-core/src/models/user.rs @@ -299,6 +299,12 @@ pub struct User { pub deleted_at: Option>, } +#[derive(Debug, Clone)] +pub struct BlockedUser { + pub user: User, + pub blocked_at: DateTime, +} + /// Administrative metadata for the account deletion and recovery lifecycle. /// /// Authentication and normal user lookups only need [`User`]. Keeping this diff --git a/synctv-core/src/repository/chat.rs b/synctv-core/src/repository/chat.rs index 8ba82d89..c95c2f70 100644 --- a/synctv-core/src/repository/chat.rs +++ b/synctv-core/src/repository/chat.rs @@ -505,6 +505,7 @@ impl ChatRepository { viewer_user_id: Option<&UserId>, ) -> Result> { let limit = limit.clamp(1, 100); + let viewer_user_id_value = viewer_user_id.map(UserId::as_i64); let pool = self.eventually_consistent_pool(); let rows = sqlx::query_as!( ChatMessagePin, @@ -525,6 +526,10 @@ impl ChatRepository { WHERE p.room_id = $1 AND m.status <> $2 AND m.deletion_source IS DISTINCT FROM $4 + AND ($5::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $5 AND ub.blocked_user_id = m.user_id + )) ORDER BY p.pinned_at DESC, p.message_id DESC LIMIT $3 "#, @@ -532,6 +537,7 @@ impl ChatRepository { i16::from(ChatMessageStatus::Deleted), i64::from(limit), DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await?; @@ -961,18 +967,15 @@ impl ChatRepository { &self, room_id: &RoomId, message_id: i64, + viewer_user_id: &UserId, reaction_key: &str, cursor: Option, limit: i32, ) -> Result { let limit = limit.clamp(1, 100); let message = self - .get_by_room_and_id(room_id, message_id) - .await? - .ok_or_else(|| Error::NotFound("Message not found".to_string()))?; - if message.status == ChatMessageStatus::Deleted { - return Err(Error::Conflict("Message has been deleted".to_string())); - } + .reaction_message_for_viewer(room_id, message_id, viewer_user_id) + .await?; let pool = self.eventually_consistent_pool(); let total = sqlx::query_scalar!( @@ -1045,6 +1048,41 @@ impl ChatRepository { }) } + pub async fn ensure_reaction_message_visible_to_viewer( + &self, + room_id: &RoomId, + message_id: i64, + viewer_user_id: &UserId, + ) -> Result<()> { + self.reaction_message_for_viewer(room_id, message_id, viewer_user_id) + .await + .map(drop) + } + + async fn reaction_message_for_viewer( + &self, + room_id: &RoomId, + message_id: i64, + viewer_user_id: &UserId, + ) -> Result { + let message = self + .get_by_room_and_id(room_id, message_id) + .await? + .ok_or_else(|| Error::NotFound("Message not found".to_string()))?; + if let Some(message_user_id) = message.user_id.as_ref() { + if self + .is_user_blocked_by(viewer_user_id, message_user_id) + .await? + { + return Err(Error::NotFound("Message not found".to_string())); + } + } + if message.status == ChatMessageStatus::Deleted { + return Err(Error::Conflict("Message has been deleted".to_string())); + } + Ok(message) + } + pub async fn get_event_by_id( &self, room_id: &RoomId, @@ -1593,6 +1631,10 @@ impl ChatRepository { AND status <> $2 AND (user_id IS NULL OR user_id <> $3) AND (created_at, id) > ($4, $5) + AND NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $3 AND ub.blocked_user_id = user_id + ) "#, room_id.as_i64(), i16::from(ChatMessageStatus::Deleted), @@ -1904,6 +1946,10 @@ impl ChatRepository { AND m.status <> $4 AND m.deletion_source IS DISTINCT FROM $6 AND (m.user_id IS NULL OR m.user_id <> $5) + AND NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $5 AND ub.blocked_user_id = m.user_id + ) "#, room_id.as_i64(), sequence, @@ -1928,6 +1974,10 @@ impl ChatRepository { AND deletion_source IS DISTINCT FROM $4 AND status <> $2 AND (user_id IS NULL OR user_id <> $3) + AND NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $3 AND ub.blocked_user_id = user_id + ) AND created_at >= NOW() - INTERVAL '90 days' "#, room_id.as_i64(), @@ -2001,6 +2051,7 @@ impl ChatRepository { ) -> Result<(Vec, Option)> { let limit = request.limit.clamp(1, 100); let included_message_type_codes = request.selection.message_type_codes(); + let viewer_user_id = request.viewer_user_id.map(UserId::as_i64); let messages = if let Some(cursor) = request.cursor { sqlx::query_as!( ChatMessageRow, @@ -2027,6 +2078,10 @@ impl ChatRepository { AND ($2 OR status <> $3) AND (created_at, id) < ($4, $5) AND message_type = ANY($7::smallint[]) + AND ($9::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $9 AND ub.blocked_user_id = user_id + )) ORDER BY created_at DESC, id DESC LIMIT $6 "#, @@ -2038,6 +2093,7 @@ impl ChatRepository { i64::from(limit), &included_message_type_codes, DeletionSource::Account as DeletionSource, + viewer_user_id, ) .fetch_all(pool) .await? @@ -2067,6 +2123,10 @@ impl ChatRepository { AND ($2 OR status <> $3) AND created_at >= NOW() - INTERVAL '90 days' AND message_type = ANY($5::smallint[]) + AND ($7::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $7 AND ub.blocked_user_id = user_id + )) ORDER BY created_at DESC, id DESC LIMIT $4 "#, @@ -2076,6 +2136,7 @@ impl ChatRepository { i64::from(limit), &included_message_type_codes, DeletionSource::Account as DeletionSource, + viewer_user_id, ) .fetch_all(pool) .await? @@ -2146,6 +2207,7 @@ impl ChatRepository { let limit = query.limit.clamp(1, 100); let fetch_limit = i64::from(limit) + 1; let user_id = query.user_id.map(|id| id.as_i64()); + let viewer_user_id_value = viewer_user_id.map(UserId::as_i64); let pool = self.pool(); let messages = if let Some(cursor) = query.cursor { sqlx::query_as!( @@ -2178,6 +2240,10 @@ impl ChatRepository { AND ($4 OR m.status <> $5) AND ($6::bigint IS NULL OR m.user_id = $6) AND (m.created_at, m.id) < ($7, $8) + AND ($11::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $11 AND ub.blocked_user_id = m.user_id + )) ORDER BY m.created_at DESC, m.id DESC LIMIT $9 "#, @@ -2191,6 +2257,7 @@ impl ChatRepository { cursor.id, fetch_limit, DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await? @@ -2224,6 +2291,10 @@ impl ChatRepository { AND (m.content_search @@ st.tsquery OR m.content ILIKE $3 ESCAPE '\') AND ($4 OR m.status <> $5) AND ($6::bigint IS NULL OR m.user_id = $6) + AND ($9::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $9 AND ub.blocked_user_id = m.user_id + )) ORDER BY m.created_at DESC, m.id DESC LIMIT $7 "#, @@ -2235,6 +2306,7 @@ impl ChatRepository { user_id, fetch_limit, DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await? @@ -2282,6 +2354,7 @@ impl ChatRepository { .map(|target| crate::models::try_hash_playback_target(Some(target))) .transpose()?; let included_message_type_codes = query.selection.message_type_codes(); + let viewer_user_id_value = viewer_user_id.map(UserId::as_i64); let pool = self.eventually_consistent_pool(); let rows = sqlx::query_as!( ChatMessageRow, @@ -2316,6 +2389,10 @@ impl ChatRepository { AND ($5::text IS NULL OR metadata #>> '{playback,playlistId}' = $5) AND ($6::text IS NULL OR metadata #>> '{playback,targetHash}' = $6) AND message_type = ANY($7::smallint[]) + AND ($12::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $12 AND ub.blocked_user_id = user_id + )) ) SELECT id AS "id!", room_id AS "room_id!: RoomId", @@ -2349,6 +2426,7 @@ impl ChatRepository { end_seconds, i64::from(limit), DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await?; @@ -2370,6 +2448,16 @@ impl ChatRepository { let Some(anchor) = self.get_by_room_and_id(room_id, message_id).await? else { return Ok(None); }; + if let (Some(viewer_user_id), Some(message_user_id)) = + (viewer_user_id, anchor.user_id.as_ref()) + { + if self + .is_user_blocked_by(viewer_user_id, message_user_id) + .await? + { + return Ok(None); + } + } if anchor.status == ChatMessageStatus::Deleted && !include_deleted { return Ok(None); } @@ -2377,6 +2465,7 @@ impl ChatRepository { let before_limit = before_limit.clamp(0, 50); let after_limit = after_limit.clamp(0, 50); let pool = self.eventually_consistent_pool(); + let viewer_user_id_value = viewer_user_id.map(UserId::as_i64); let mut before = sqlx::query_as!( ChatMessageRow, r#" @@ -2401,6 +2490,10 @@ impl ChatRepository { AND deletion_source IS DISTINCT FROM $7 AND ($2 OR status <> $3) AND (created_at, id) < ($4, $5) + AND ($8::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $8 AND ub.blocked_user_id = user_id + )) ORDER BY created_at DESC, id DESC LIMIT $6 "#, @@ -2411,6 +2504,7 @@ impl ChatRepository { anchor.id, i64::from(before_limit), DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await?; @@ -2440,6 +2534,10 @@ impl ChatRepository { AND deletion_source IS DISTINCT FROM $7 AND ($2 OR status <> $3) AND (created_at, id) > ($4, $5) + AND ($8::bigint IS NULL OR NOT EXISTS ( + SELECT 1 FROM user_blocks ub + WHERE ub.blocker_user_id = $8 AND ub.blocked_user_id = user_id + )) ORDER BY created_at ASC, id ASC LIMIT $6 "#, @@ -2450,6 +2548,7 @@ impl ChatRepository { anchor.id, i64::from(after_limit), DeletionSource::Account as DeletionSource, + viewer_user_id_value, ) .fetch_all(pool) .await?; @@ -2601,6 +2700,16 @@ impl ChatRepository { else { return Ok(None); }; + if let (Some(viewer_user_id), Some(message_user_id)) = + (viewer_user_id, message.user_id.as_ref()) + { + if self + .is_user_blocked_by(viewer_user_id, message_user_id) + .await? + { + return Ok(None); + } + } let attachments = if message.status == ChatMessageStatus::Deleted { Vec::new() } else { @@ -2634,6 +2743,64 @@ impl ChatRepository { })) } + pub async fn is_user_blocked_by( + &self, + viewer_user_id: &UserId, + message_user_id: &UserId, + ) -> Result { + let is_blocked = sqlx::query_scalar!( + r#" + SELECT EXISTS ( + SELECT 1 FROM user_blocks + WHERE blocker_user_id = $1 AND blocked_user_id = $2 + ) AS "is_blocked!" + "#, + viewer_user_id as &UserId, + message_user_id as &UserId, + ) + .fetch_one(self.pool()) + .await?; + Ok(is_blocked) + } + + pub async fn blocked_user_ids(&self, viewer_user_id: &UserId) -> Result> { + sqlx::query_scalar!( + r#"SELECT blocked_user_id AS "blocked_user_id!: UserId" + FROM user_blocks + WHERE blocker_user_id = $1"#, + viewer_user_id as &UserId, + ) + .fetch_all(self.pool()) + .await + .map_err(Into::into) + } + + pub async fn blocking_viewer_ids( + &self, + viewer_user_ids: &[UserId], + blocked_user_id: &UserId, + ) -> Result> { + if viewer_user_ids.is_empty() { + return Ok(Vec::new()); + } + let viewer_ids = viewer_user_ids + .iter() + .map(UserId::as_i64) + .collect::>(); + sqlx::query_scalar!( + r#" + SELECT blocker_user_id AS "blocker_user_id!: UserId" + FROM user_blocks + WHERE blocked_user_id = $1 AND blocker_user_id = ANY($2::bigint[]) + "#, + blocked_user_id as &UserId, + &viewer_ids, + ) + .fetch_all(self.pool()) + .await + .map_err(Into::into) + } + pub async fn edit_with_event( &self, request: EditChatMessageEventRequest<'_>, diff --git a/synctv-core/src/repository/room.rs b/synctv-core/src/repository/room.rs index baaa5ba1..b4b6aafa 100644 --- a/synctv-core/src/repository/room.rs +++ b/synctv-core/src/repository/room.rs @@ -760,6 +760,19 @@ impl RoomRepository { builder.push("r.created_by = ").push_bind(creator_id); } + if !query.excluded_creator_ids.is_empty() { + let creator_ids = query + .excluded_creator_ids + .iter() + .map(UserId::as_i64) + .collect::>(); + Self::push_where_prefix(builder, has_condition); + builder + .push("r.created_by <> ALL(") + .push_bind(creator_ids) + .push("::bigint[])"); + } + if let Some(category_id) = query.category_id { Self::push_where_prefix(builder, has_condition); builder.push("r.category_id = ").push_bind(category_id); diff --git a/synctv-core/src/repository/room_tests.rs b/synctv-core/src/repository/room_tests.rs index a0744ce7..ed9a49b8 100644 --- a/synctv-core/src/repository/room_tests.rs +++ b/synctv-core/src/repository/room_tests.rs @@ -13,6 +13,7 @@ fn test_room_list_order_clause_supports_name_ascending() { sort_direction: crate::models::SortDirection::Asc, pagination: PageParams::default(), creator_id: None, + excluded_creator_ids: Vec::new(), category_id: None, label_ids: Vec::new(), }; @@ -31,6 +32,7 @@ fn test_room_list_order_clause_supports_last_activity_nulls_last() { sort_direction: crate::models::SortDirection::Desc, pagination: PageParams::default(), creator_id: None, + excluded_creator_ids: Vec::new(), category_id: None, label_ids: Vec::new(), }; @@ -659,6 +661,57 @@ async fn test_list_rooms_with_filters() { assert!(rooms.iter().all(|r| r.name.contains("Active"))); } +#[tokio::test] +#[ignore = "Requires Docker"] +async fn test_list_rooms_excludes_blocked_creators_before_pagination() { + use crate::repository::user::UserRepository; + use crate::test_helpers::{RoomFixture, UserFixture}; + + let (_postgres, pool) = create_test_pool().await; + let user_repo = UserRepository::new(pool.clone()); + let room_repo = RoomRepository::new(pool.clone()); + let visible_owner = user_repo + .create(&UserFixture::new().with_username("visible_owner").build()) + .await + .checked("visible owner should be created"); + let blocked_owner = user_repo + .create(&UserFixture::new().with_username("blocked_owner").build()) + .await + .checked("blocked owner should be created"); + + for (name, owner_id) in [ + ("Visible discovery room", visible_owner.id), + ("Blocked discovery room", blocked_owner.id), + ] { + room_repo + .create( + &RoomFixture::new() + .with_name(name) + .with_owner(owner_id) + .build(), + ) + .await + .checked("room should be created"); + } + + let query = RoomListQuery { + pagination: PageParams::new(Some(1), Some(1)), + search: Some("discovery room".to_string()), + excluded_creator_ids: vec![blocked_owner.id], + sort_by: crate::models::RoomListSortBy::Name, + sort_direction: crate::models::SortDirection::Asc, + ..Default::default() + }; + let (rooms, total) = room_repo + .list(&query) + .await + .checked("rooms should be filtered"); + + assert_eq!(total, 1); + assert_eq!(rooms.len(), 1); + assert_eq!(rooms[0].created_by, visible_owner.id); +} + /// Integration test: room member_count counts current member rows. #[tokio::test] #[ignore = "Requires Docker"] diff --git a/synctv-core/src/repository/user.rs b/synctv-core/src/repository/user.rs index 1f531fbb..d99c7f3b 100644 --- a/synctv-core/src/repository/user.rs +++ b/synctv-core/src/repository/user.rs @@ -5,8 +5,8 @@ use super::query_builder::ilike_contains_pattern; use crate::repository::pools::RepoPools; use crate::{ models::{ - DeletionSource, SignupMethod, User, UserId, UserLifecycleMetadata, UserListQuery, - UserListSortBy, UserRole, UserStatus, + BlockedUser, DeletionSource, PageParams, SignupMethod, User, UserId, UserLifecycleMetadata, + UserListQuery, UserListSortBy, UserRole, UserStatus, }, Error, Result, }; @@ -38,6 +38,49 @@ struct UserListRow { deleted_at: Option>, } +#[derive(sqlx::FromRow)] +struct BlockedUserListRow { + id: UserId, + username: String, + signup_method: SignupMethod, + role: UserRole, + avatar_file_reference_id: Option, + status: UserStatus, + is_banned: bool, + banned_at: Option>, + banned_by: Option, + banned_reason: Option, + created_at: DateTime, + updated_at: DateTime, + version: i32, + deleted_at: Option>, + blocked_at: DateTime, +} + +impl From for BlockedUser { + fn from(row: BlockedUserListRow) -> Self { + Self { + user: User { + id: row.id, + username: row.username, + role: row.role, + avatar_file_reference_id: row.avatar_file_reference_id, + status: row.status, + is_banned: row.is_banned, + banned_at: row.banned_at, + banned_by: row.banned_by, + banned_reason: row.banned_reason, + signup_method: row.signup_method, + created_at: row.created_at, + updated_at: row.updated_at, + version: row.version, + deleted_at: row.deleted_at, + }, + blocked_at: row.blocked_at, + } + } +} + impl From for User { fn from(row: UserListRow) -> Self { Self { @@ -80,6 +123,166 @@ impl UserRepository { } } + pub async fn block_user( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result> { + let blocked_at = sqlx::query_scalar!( + r#" + INSERT INTO user_blocks (blocker_user_id, blocked_user_id) + VALUES ($1, $2) + ON CONFLICT (blocker_user_id, blocked_user_id) + DO UPDATE SET blocker_user_id = EXCLUDED.blocker_user_id + RETURNING created_at AS "created_at!" + "#, + blocker_user_id as &UserId, + blocked_user_id as &UserId, + ) + .fetch_one(self.pools.primary()) + .await?; + + Ok(blocked_at) + } + + pub async fn unblock_user( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result { + let result = sqlx::query!( + "DELETE FROM user_blocks WHERE blocker_user_id = $1 AND blocked_user_id = $2", + blocker_user_id as &UserId, + blocked_user_id as &UserId, + ) + .execute(self.pools.primary()) + .await?; + + Ok(result.rows_affected() > 0) + } + + pub async fn is_blocking( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result { + let is_blocking = sqlx::query_scalar!( + r#" + SELECT EXISTS ( + SELECT 1 + FROM user_blocks + WHERE blocker_user_id = $1 AND blocked_user_id = $2 + ) AS "is_blocking!" + "#, + blocker_user_id as &UserId, + blocked_user_id as &UserId, + ) + .fetch_one(self.pools.primary()) + .await?; + Ok(is_blocking) + } + + pub async fn blocked_user_ids(&self, blocker_user_id: &UserId) -> Result> { + sqlx::query_scalar!( + r#"SELECT blocked_user_id AS "blocked_user_id!: UserId" + FROM user_blocks + WHERE blocker_user_id = $1"#, + blocker_user_id as &UserId, + ) + .fetch_all(self.pools.primary()) + .await + .map_err(Into::into) + } + + pub async fn blocked_user_ids_eventually_consistent( + &self, + blocker_user_id: &UserId, + ) -> Result> { + sqlx::query_scalar!( + r#"SELECT blocked_user_id AS "blocked_user_id!: UserId" + FROM user_blocks + WHERE blocker_user_id = $1"#, + blocker_user_id as &UserId, + ) + .fetch_all(self.pools.read()) + .await + .map_err(Into::into) + } + + pub async fn list_blocked_users( + &self, + blocker_user_id: &UserId, + pagination: PageParams, + search: Option<&str>, + ) -> Result<(Vec, i64)> { + let limit = pagination.limit_i64()?; + let offset = pagination.offset_i64()?; + let search_pattern = search.and_then(ilike_contains_pattern); + + let mut count_builder = QueryBuilder::::new( + "SELECT COUNT(*) FROM user_blocks ub \ + JOIN user_account_profiles p ON p.id = ub.blocked_user_id \ + WHERE ub.blocker_user_id = ", + ); + count_builder.push_bind(blocker_user_id); + count_builder.push(" AND p.deleted_at IS NULL"); + if let Some(pattern) = search_pattern.as_ref() { + count_builder + .push(" AND p.username ILIKE ") + .push_bind(pattern) + .push(" ESCAPE '\\'"); + } + let total = count_builder + .build_query_scalar::() + .fetch_one(self.pools.primary()) + .await?; + + let mut list_builder = QueryBuilder::::new( + r" + SELECT p.id, + p.username, + p.signup_method, + p.role, + p.avatar_file_reference_id, + p.status, + p.is_banned, + p.banned_at, + p.banned_by, + p.banned_reason, + p.created_at, + p.updated_at, + p.version, + p.deleted_at, + ub.created_at AS blocked_at + FROM user_blocks ub + JOIN user_account_profiles p ON p.id = ub.blocked_user_id + WHERE ub.blocker_user_id = + ", + ); + list_builder.push_bind(blocker_user_id); + list_builder.push(" AND p.deleted_at IS NULL"); + if let Some(pattern) = search_pattern.as_ref() { + list_builder + .push(" AND p.username ILIKE ") + .push_bind(pattern) + .push(" ESCAPE '\\'"); + } + list_builder + .push(" ORDER BY ub.created_at DESC, ub.blocked_user_id DESC LIMIT ") + .push_bind(limit) + .push(" OFFSET ") + .push_bind(offset); + let users = list_builder + .build_query_as::() + .fetch_all(self.pools.primary()) + .await? + .into_iter() + .map(Into::into) + .collect(); + + Ok((users, total)) + } + /// Get the database pool #[must_use] pub const fn pool(&self) -> &PgPool { diff --git a/synctv-core/src/repository/user_tests.rs b/synctv-core/src/repository/user_tests.rs index ab26b5d1..5127542e 100644 --- a/synctv-core/src/repository/user_tests.rs +++ b/synctv-core/src/repository/user_tests.rs @@ -222,3 +222,79 @@ async fn test_soft_delete_user() { .checked("operation should succeed") .is_none()); } + +#[tokio::test] +#[ignore = "Requires Docker"] +async fn test_user_blocking_is_personal_idempotent_and_searchable() { + let (_postgres, pool) = create_test_pool().await; + let repo = UserRepository::new(pool.clone()); + let blocker = repo + .create(&User::new("blocker".into(), SignupMethod::Email)) + .await + .checked("blocker should be created"); + let blocked_alpha = repo + .create(&User::new("blocked_alpha".into(), SignupMethod::Email)) + .await + .checked("first blocked user should be created"); + let blocked_beta = repo + .create(&User::new("blocked_beta".into(), SignupMethod::Email)) + .await + .checked("second blocked user should be created"); + + let first_blocked_at = repo + .block_user(&blocker.id, &blocked_alpha.id) + .await + .checked("user should be blocked"); + let repeated_blocked_at = repo + .block_user(&blocker.id, &blocked_alpha.id) + .await + .checked("blocking should be idempotent"); + repo.block_user(&blocker.id, &blocked_beta.id) + .await + .checked("second user should be blocked"); + + assert_eq!(repeated_blocked_at, first_blocked_at); + assert!(repo + .is_blocking(&blocker.id, &blocked_alpha.id) + .await + .checked("blocking relationship should be readable")); + assert!(!repo + .is_blocking(&blocked_alpha.id, &blocker.id) + .await + .checked("reverse blocking relationship should be readable")); + + let (first_page, total) = repo + .list_blocked_users( + &blocker.id, + PageParams::new(Some(1), Some(1)), + Some("blocked_"), + ) + .await + .checked("blocked users should be listed"); + assert_eq!(total, 2); + assert_eq!(first_page.len(), 1); + + let (search_result, search_total) = repo + .list_blocked_users( + &blocker.id, + PageParams::new(Some(1), Some(10)), + Some("alpha"), + ) + .await + .checked("blocked users should be searchable"); + assert_eq!(search_total, 1); + assert_eq!(search_result[0].user.id, blocked_alpha.id); + + assert!(repo + .unblock_user(&blocker.id, &blocked_alpha.id) + .await + .checked("user should be unblocked")); + assert!(!repo + .unblock_user(&blocker.id, &blocked_alpha.id) + .await + .checked("unblocking should be idempotent")); + assert!(!repo + .is_blocking(&blocker.id, &blocked_alpha.id) + .await + .checked("removed relationship should remain readable")); +} diff --git a/synctv-core/src/service/chat.rs b/synctv-core/src/service/chat.rs index 9ed51e5b..5e79fb06 100644 --- a/synctv-core/src/service/chat.rs +++ b/synctv-core/src/service/chat.rs @@ -55,6 +55,7 @@ const CHAT_REACTION_DETAIL_CACHE_CAPACITY: u64 = 1024; struct ChatReactionDetailCacheKey { room_id: RoomId, message_id: i64, + viewer_user_id: UserId, reaction_key: String, cursor: Option, limit: i32, @@ -1373,17 +1374,28 @@ impl ChatService { let key = ChatReactionDetailCacheKey { room_id: *room_id, message_id, + viewer_user_id: *viewer_user_id, reaction_key: reaction_key.to_string(), cursor, limit, }; if let Some(page) = self.reaction_detail_cache.get(&key).await { + self.chat_repository + .ensure_reaction_message_visible_to_viewer(room_id, message_id, viewer_user_id) + .await?; return Ok(page); } let page = self .chat_repository - .list_reaction_users(room_id, message_id, reaction_key, cursor, limit) + .list_reaction_users( + room_id, + message_id, + viewer_user_id, + reaction_key, + cursor, + limit, + ) .await?; self.reaction_detail_cache.insert(key, page.clone()).await; Ok(page) @@ -1525,6 +1537,20 @@ impl ChatService { Ok(username) } + pub async fn blocked_user_ids(&self, viewer_user_id: &UserId) -> Result> { + self.chat_repository.blocked_user_ids(viewer_user_id).await + } + + pub async fn blocking_viewer_ids( + &self, + viewer_user_ids: &[UserId], + blocked_user_id: &UserId, + ) -> Result> { + self.chat_repository + .blocking_viewer_ids(viewer_user_ids, blocked_user_id) + .await + } + async fn ensure_reply_target_visible( &self, room_id: &RoomId, diff --git a/synctv-core/src/service/chat_tests.rs b/synctv-core/src/service/chat_tests.rs index 8b378b54..92c2feb3 100644 --- a/synctv-core/src/service/chat_tests.rs +++ b/synctv-core/src/service/chat_tests.rs @@ -3310,6 +3310,26 @@ async fn chat_reactions_update_history_and_emit_reaction_events() { assert_eq!(next.total, 2); assert_eq!(next.users.len(), 1); assert_ne!(page.users[0].user_id, next.users[0].user_id); + + ok( + service.user_service.block_user(&member.id, &owner.id).await, + "member should block message owner", + ); + let blocked_message_error = service + .list_reaction_users( + &room.id, + message.message.message.id, + &member.id, + "like", + None, + 10, + ) + .await + .expect_err("blocked message reaction users should be hidden"); + assert!(matches!( + blocked_message_error, + crate::Error::NotFound(ref message) if message == "Message not found" + )); } #[tokio::test] diff --git a/synctv-core/src/service/user.rs b/synctv-core/src/service/user.rs index 03d5b359..15120900 100644 --- a/synctv-core/src/service/user.rs +++ b/synctv-core/src/service/user.rs @@ -157,6 +157,7 @@ impl std::fmt::Debug for UserService { } mod avatar; +mod blocking; mod constructor; mod deletion; pub use deletion::{ diff --git a/synctv-core/src/service/user/blocking.rs b/synctv-core/src/service/user/blocking.rs new file mode 100644 index 00000000..eb0bdda4 --- /dev/null +++ b/synctv-core/src/service/user/blocking.rs @@ -0,0 +1,70 @@ +use chrono::{DateTime, Utc}; + +use crate::{ + models::{BlockedUser, PageParams, UserId}, + Error, Result, +}; + +use super::UserService; + +impl UserService { + pub async fn block_user( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result> { + if blocker_user_id == blocked_user_id { + return Err(Error::InvalidInput("Cannot block yourself".to_string())); + } + if self.repository.get_by_id(blocked_user_id).await?.is_none() { + return Err(Error::NotFound("User not found".to_string())); + } + self.repository + .block_user(blocker_user_id, blocked_user_id) + .await + } + + pub async fn unblock_user( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result { + self.repository + .unblock_user(blocker_user_id, blocked_user_id) + .await + } + + pub async fn is_blocking( + &self, + blocker_user_id: &UserId, + blocked_user_id: &UserId, + ) -> Result { + self.repository + .is_blocking(blocker_user_id, blocked_user_id) + .await + } + + pub async fn blocked_user_ids(&self, blocker_user_id: &UserId) -> Result> { + self.repository.blocked_user_ids(blocker_user_id).await + } + + pub async fn blocked_user_ids_eventually_consistent( + &self, + blocker_user_id: &UserId, + ) -> Result> { + self.repository + .blocked_user_ids_eventually_consistent(blocker_user_id) + .await + } + + pub async fn list_blocked_users( + &self, + blocker_user_id: &UserId, + pagination: PageParams, + search: Option<&str>, + ) -> Result<(Vec, i64)> { + self.repository + .list_blocked_users(blocker_user_id, pagination, search) + .await + } +} diff --git a/synctv-management/src/mapping/response.rs b/synctv-management/src/mapping/response.rs index e26bf45f..8a89e185 100644 --- a/synctv-management/src/mapping/response.rs +++ b/synctv-management/src/mapping/response.rs @@ -1190,6 +1190,7 @@ pub(crate) fn created_room_to_client_proto( cover: None, presence: None, creator: Some(user_public_view_to_client_proto(creator, public_id_codec)?), + creator_blocked: false, category: room .category .as_ref() diff --git a/synctv-proto-build/src/lib.rs b/synctv-proto-build/src/lib.rs index 4674063d..8440bb4d 100644 --- a/synctv-proto-build/src/lib.rs +++ b/synctv-proto-build/src/lib.rs @@ -906,6 +906,7 @@ fn build_main_protos(out_dir: &Path) -> Result<(), Box> { builder = add_query_params_attrs( builder, &[ + ".synctv.client.ListBlockedUsersRequest", ".synctv.client.ListMyRoomsRequest", ".synctv.client.ListNotificationsRequest", ".synctv.client.GetRoomMembersRequest", diff --git a/synctv-proto/proto/client.proto b/synctv-proto/proto/client.proto index adb2f550..40ef2fb0 100644 --- a/synctv-proto/proto/client.proto +++ b/synctv-proto/proto/client.proto @@ -79,6 +79,9 @@ service UserService { rpc UpdateUserPreferences(UpdateUserPreferencesRequest) returns (UpdateUserPreferencesResponse); rpc SetTwoFactorEnabled(SetTwoFactorEnabledRequest) returns (GetUserPreferencesResponse); rpc CloseAccount(CloseAccountRequest) returns (CloseAccountResponse); + rpc BlockUser(BlockUserRequest) returns (BlockUserResponse); + rpc UnblockUser(UnblockUserRequest) returns (UnblockUserResponse); + rpc ListBlockedUsers(ListBlockedUsersRequest) returns (ListBlockedUsersResponse); // User-initiated room lifecycle operations outside room-scoped context rpc CreateRoom(CreateRoomRequest) returns (Room); @@ -291,6 +294,56 @@ message UserPublicView { FileObjectAccess avatar_access = 7; } +message BlockedUser { + UserPublicView user = 1; + int64 blocked_at = 2; +} + +message BlockUserRequest { + string user_id = 1 [(buf.validate.field).string = { + min_len: 1 + max_len: 64 + pattern: "^usr_[A-Za-z0-9]+$" + }]; +} + +message BlockUserResponse { + BlockedUser blocked_user = 1; +} + +message UnblockUserRequest { + string user_id = 1 [(buf.validate.field).string = { + min_len: 1 + max_len: 64 + pattern: "^usr_[A-Za-z0-9]+$" + }]; +} + +message UnblockUserResponse { + bool success = 1; +} + +message ListBlockedUsersRequest { + option (buf.validate.message).cel = { + id: "list_blocked_users.page" + message: "page must be 0 (use default) or at least 1" + expression: "this.page == 0 || this.page >= 1" + }; + option (buf.validate.message).cel = { + id: "list_blocked_users.page_size" + message: "page_size must be 0 (use default) or between 1 and 100" + expression: "this.page_size == 0 || (this.page_size >= 1 && this.page_size <= 100)" + }; + int32 page = 1; + int32 page_size = 2; + string search = 3 [(buf.validate.field).string = {max_len: 100}]; +} + +message ListBlockedUsersResponse { + repeated BlockedUser users = 1; + int32 total = 2; +} + enum PlayMode { PLAY_MODE_UNSPECIFIED = 0; PLAY_MODE_SEQUENTIAL = 1; @@ -569,6 +622,7 @@ message Room { RoomCategory category = 16; repeated RoomLabel labels = 17; optional bool is_public = 18; + bool creator_blocked = 19; } message RoomCategory { diff --git a/synctv/src/cli/tests.rs b/synctv/src/cli/tests.rs index 114583f7..5f058500 100644 --- a/synctv/src/cli/tests.rs +++ b/synctv/src/cli/tests.rs @@ -5783,6 +5783,7 @@ fn render_human_output_uses_room_and_member_enums_by_context() { category: None, labels: Vec::new(), is_public: Some(true), + creator_blocked: false, }), playback_state: None, requires_approval: false,