feat: add personal user blocking (#429)

## Summary

- add idempotent block, unblock, and paginated block-list APIs across
HTTP, gRPC, service, and repository layers
- hide blocked creators from authenticated personal discovery while
leaving anonymous discovery unchanged
- preserve joined, favorited, and directly accessed rooms with a
creator_blocked marker
- filter blocked users from message history, search, context, pinned
messages, replay, unread counts, live events, and reconnect replay
- add the user_blocks migration, SQLx offline metadata, and
repository/API coverage

## Behavior

Blocking is private to the current account. It does not alter the
blocked account or public anonymous discovery. Detailed personal room
lists remain complete and expose creator_blocked so clients can label
the creator. Observer-specific realtime filtering advances event cursors
and fails closed when block-state lookup is unavailable.

## Verification

- make check
- make fmt-check
- cargo test --workspace: 5222 passed
- synctv-api-common library tests: 638 passed, 165 ignored
- PostgreSQL favorite-room creator_blocked integration test
- SQLx prepare across the workspace
- real two-account HTTP, realtime, and room-flow integration test

Client: https://github.com/synctv-org/synctv-app/pull/47
pull/430/head
zijiren 1 month ago committed by GitHub
parent 61a837a424
commit 52b004bab2
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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"
}

@ -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);

@ -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)?,

@ -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::<Result<Vec<_>, _>>()?,
is_public: Some(room.is_public),
creator_blocked: false,
})
}

@ -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<synctv_proto::client::Room, ApiError> {
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<DiscoveryRoomWindow, ApiError> {
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<synctv_proto::client::DiscoverRoomsResponse, ApiError> {
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<synctv_proto::client::DiscoverRoomsResponse, ApiError> {
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<synctv_core::models::Room> =
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::<HashSet<_>>())
.map_err(ApiError::from)
},
)?;
let presence_by_room: HashMap<synctv_core::models::RoomId, _> = 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::<HashSet<_>>();
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)?;

@ -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<synctv_proto::client::BlockUserResponse, ApiError> {
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<synctv_proto::client::UnblockUserResponse, ApiError> {
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<synctv_proto::client::ListBlockedUsersResponse, ApiError> {
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::<Vec<_>>();
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,

@ -1548,6 +1548,51 @@ impl MediaResourceHub {
event_cursor: Option<synctv_proto::client::EventCursor>,
) -> 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::<Vec<_>>();
let viewer_ids = observers
.iter()
.filter_map(|observer| observer.actor.user_id())
.collect::<HashSet<_>>()
.into_iter()
.collect::<Vec<_>>();
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<ResourceObserver>, 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::<HashSet<_>>(),
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(

@ -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<UserAvatarObjectResponse>;
@ -707,6 +709,71 @@ impl UserService for ClientServiceImpl {
Ok(Response::new(response))
}
async fn block_user(
&self,
request: Request<BlockUserRequest>,
) -> Result<Response<BlockUserResponse>, 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<UnblockUserRequest>,
) -> Result<Response<UnblockUserResponse>, 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<ListBlockedUsersRequest>,
) -> Result<Response<ListBlockedUsersResponse>, 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<CreateRoomRequest>,

@ -775,6 +775,14 @@ fn register_playlist_cover_object_routes() -> Router<AppState> {
fn register_extracted_user_routes() -> Router<AppState> {
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",

@ -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<AppState>,
Json(req): Json<BlockUserRequest>,
) -> AppResult<Json<BlockUserResponse>> {
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<AppState>,
Path(path): Path<BlockedUserPath>,
) -> AppResult<Json<UnblockUserResponse>> {
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<AppState>,
ProtoQuery(req): ProtoQuery<ListBlockedUsersRequest>,
) -> AppResult<Json<ListBlockedUsersResponse>> {
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(

@ -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,

@ -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,

@ -363,6 +363,8 @@ pub struct RoomListQuery {
pub is_public: Option<bool>,
/// Filter by creator
pub creator_id: Option<super::UserId>,
#[serde(default)]
pub excluded_creator_ids: Vec<super::UserId>,
pub category_id: Option<RoomCategoryId>,
#[serde(default)]
pub label_ids: Vec<RoomLabelId>,
@ -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,

@ -299,6 +299,12 @@ pub struct User {
pub deleted_at: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone)]
pub struct BlockedUser {
pub user: User,
pub blocked_at: DateTime<Utc>,
}
/// Administrative metadata for the account deletion and recovery lifecycle.
///
/// Authentication and normal user lookups only need [`User`]. Keeping this

@ -505,6 +505,7 @@ impl ChatRepository {
viewer_user_id: Option<&UserId>,
) -> Result<Vec<ChatPinnedMessage>> {
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<ChatReactionUsersCursor>,
limit: i32,
) -> Result<ChatReactionUsersPage> {
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<ChatMessage> {
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<ChatMessageWithAttachments>, Option<ChatHistoryCursor>)> {
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<bool> {
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<Vec<UserId>> {
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<Vec<UserId>> {
if viewer_user_ids.is_empty() {
return Ok(Vec::new());
}
let viewer_ids = viewer_user_ids
.iter()
.map(UserId::as_i64)
.collect::<Vec<_>>();
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<'_>,

@ -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::<Vec<_>>();
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);

@ -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"]

@ -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<DateTime<Utc>>,
}
#[derive(sqlx::FromRow)]
struct BlockedUserListRow {
id: UserId,
username: String,
signup_method: SignupMethod,
role: UserRole,
avatar_file_reference_id: Option<i64>,
status: UserStatus,
is_banned: bool,
banned_at: Option<DateTime<Utc>>,
banned_by: Option<UserId>,
banned_reason: Option<String>,
created_at: DateTime<Utc>,
updated_at: DateTime<Utc>,
version: i32,
deleted_at: Option<DateTime<Utc>>,
blocked_at: DateTime<Utc>,
}
impl From<BlockedUserListRow> 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<UserListRow> 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<DateTime<Utc>> {
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<bool> {
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<bool> {
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<Vec<UserId>> {
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<Vec<UserId>> {
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<BlockedUser>, 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::<Postgres>::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::<i64>()
.fetch_one(self.pools.primary())
.await?;
let mut list_builder = QueryBuilder::<Postgres>::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::<BlockedUserListRow>()
.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 {

@ -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"));
}

@ -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<ChatReactionUsersCursor>,
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<Vec<UserId>> {
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<Vec<UserId>> {
self.chat_repository
.blocking_viewer_ids(viewer_user_ids, blocked_user_id)
.await
}
async fn ensure_reply_target_visible(
&self,
room_id: &RoomId,

@ -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]

@ -157,6 +157,7 @@ impl std::fmt::Debug for UserService {
}
mod avatar;
mod blocking;
mod constructor;
mod deletion;
pub use deletion::{

@ -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<DateTime<Utc>> {
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<bool> {
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<bool> {
self.repository
.is_blocking(blocker_user_id, blocked_user_id)
.await
}
pub async fn blocked_user_ids(&self, blocker_user_id: &UserId) -> Result<Vec<UserId>> {
self.repository.blocked_user_ids(blocker_user_id).await
}
pub async fn blocked_user_ids_eventually_consistent(
&self,
blocker_user_id: &UserId,
) -> Result<Vec<UserId>> {
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<BlockedUser>, i64)> {
self.repository
.list_blocked_users(blocker_user_id, pagination, search)
.await
}
}

@ -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()

@ -906,6 +906,7 @@ fn build_main_protos(out_dir: &Path) -> Result<(), Box<dyn std::error::Error>> {
builder = add_query_params_attrs(
builder,
&[
".synctv.client.ListBlockedUsersRequest",
".synctv.client.ListMyRoomsRequest",
".synctv.client.ListNotificationsRequest",
".synctv.client.GetRoomMembersRequest",

@ -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 {

@ -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,

Loading…
Cancel
Save