diff --git a/.env.example b/.env.example index 67fcba7..dfb1a60 100644 --- a/.env.example +++ b/.env.example @@ -96,7 +96,8 @@ BACKUP_STORAGE_PATH=/var/lib/tranquil/backups # Discord notifications (via webhook) # DISCORD_WEBHOOK_URL=https://discord.com/api/webhooks/... # Telegram notifications (via bot) -# TELEGRAM_BOT_TOKEN=your-bot-token +# TELEGRAM_BOT_TOKEN=bot-token +# TELEGRAM_WEBHOOK_SECRET=random-secret # Signal notifications (via signal-cli) # SIGNAL_CLI_PATH=/usr/local/bin/signal-cli # SIGNAL_SENDER_NUMBER=+1234567890 diff --git a/.sqlx/query-25309f4a08845a49557d694ad9b5b9a137be4dcce28e9293551c8c3fd40fdd86.json b/.sqlx/query-25309f4a08845a49557d694ad9b5b9a137be4dcce28e9293551c8c3fd40fdd86.json index b0716fa..0c47103 100644 --- a/.sqlx/query-25309f4a08845a49557d694ad9b5b9a137be4dcce28e9293551c8c3fd40fdd86.json +++ b/.sqlx/query-25309f4a08845a49557d694ad9b5b9a137be4dcce28e9293551c8c3fd40fdd86.json @@ -44,7 +44,8 @@ "channel_verification", "passkey_recovery", "legacy_login_alert", - "migration_verification" + "migration_verification", + "channel_verified" ] } } diff --git a/.sqlx/query-4445cc86cdf04894b340e67661b79a3c411917144a011f50849b737130b24dbe.json b/.sqlx/query-4445cc86cdf04894b340e67661b79a3c411917144a011f50849b737130b24dbe.json index 7004d69..ac57c00 100644 --- a/.sqlx/query-4445cc86cdf04894b340e67661b79a3c411917144a011f50849b737130b24dbe.json +++ b/.sqlx/query-4445cc86cdf04894b340e67661b79a3c411917144a011f50849b737130b24dbe.json @@ -32,7 +32,8 @@ "channel_verification", "passkey_recovery", "legacy_login_alert", - "migration_verification" + "migration_verification", + "channel_verified" ] } } diff --git a/.sqlx/query-d8fd97c8be3211b2509669dd859245b14e15f81a42d7e0c4c428b65f466af5ee.json b/.sqlx/query-63c2a9079c147be6d04bf02c63ea7a3d0b0db3f35438380c0fe7e4c60b420c67.json similarity index 63% rename from .sqlx/query-d8fd97c8be3211b2509669dd859245b14e15f81a42d7e0c4c428b65f466af5ee.json rename to .sqlx/query-63c2a9079c147be6d04bf02c63ea7a3d0b0db3f35438380c0fe7e4c60b420c67.json index cef549d..035bbb8 100644 --- a/.sqlx/query-d8fd97c8be3211b2509669dd859245b14e15f81a42d7e0c4c428b65f466af5ee.json +++ b/.sqlx/query-63c2a9079c147be6d04bf02c63ea7a3d0b0db3f35438380c0fe7e4c60b420c67.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT email, handle, preferred_comms_channel as \"preferred_channel!: CommsChannel\", preferred_locale\n FROM users WHERE id = $1", + "query": "SELECT email, handle, preferred_comms_channel as \"preferred_channel!: CommsChannel\", preferred_locale, telegram_chat_id, discord_id, signal_number\n FROM users WHERE id = $1", "describe": { "columns": [ { @@ -34,6 +34,21 @@ "ordinal": 3, "name": "preferred_locale", "type_info": "Varchar" + }, + { + "ordinal": 4, + "name": "telegram_chat_id", + "type_info": "Int8" + }, + { + "ordinal": 5, + "name": "discord_id", + "type_info": "Text" + }, + { + "ordinal": 6, + "name": "signal_number", + "type_info": "Text" } ], "parameters": { @@ -45,8 +60,11 @@ true, false, false, + true, + true, + true, true ] }, - "hash": "d8fd97c8be3211b2509669dd859245b14e15f81a42d7e0c4c428b65f466af5ee" + "hash": "63c2a9079c147be6d04bf02c63ea7a3d0b0db3f35438380c0fe7e4c60b420c67" } diff --git a/.sqlx/query-8047fda41bd94f819213decb8b3e0aba49a8dbdb10217eefd77e3567f8c9694a.json b/.sqlx/query-8047fda41bd94f819213decb8b3e0aba49a8dbdb10217eefd77e3567f8c9694a.json index 8cc9194..70ae98c 100644 --- a/.sqlx/query-8047fda41bd94f819213decb8b3e0aba49a8dbdb10217eefd77e3567f8c9694a.json +++ b/.sqlx/query-8047fda41bd94f819213decb8b3e0aba49a8dbdb10217eefd77e3567f8c9694a.json @@ -49,7 +49,8 @@ "channel_verification", "passkey_recovery", "legacy_login_alert", - "migration_verification" + "migration_verification", + "channel_verified" ] } } diff --git a/.sqlx/query-9d17e25776c67f96022010840c1d04bdd542b5bcc511b1778bab0159a4566e9c.json b/.sqlx/query-9d17e25776c67f96022010840c1d04bdd542b5bcc511b1778bab0159a4566e9c.json new file mode 100644 index 0000000..21f7511 --- /dev/null +++ b/.sqlx/query-9d17e25776c67f96022010840c1d04bdd542b5bcc511b1778bab0159a4566e9c.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE users SET telegram_chat_id = $2, telegram_verified = TRUE, updated_at = NOW()\n WHERE id = (\n SELECT id FROM users\n WHERE LOWER(telegram_username) = LOWER($1) AND telegram_username IS NOT NULL AND deactivated_at IS NULL\n LIMIT 1\n ) RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Text", + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "9d17e25776c67f96022010840c1d04bdd542b5bcc511b1778bab0159a4566e9c" +} diff --git a/.sqlx/query-a71a76724a3a7406e30a998c03d52554efaf649bb10b35b6e5d64ac59a479023.json b/.sqlx/query-a71a76724a3a7406e30a998c03d52554efaf649bb10b35b6e5d64ac59a479023.json new file mode 100644 index 0000000..075e698 --- /dev/null +++ b/.sqlx/query-a71a76724a3a7406e30a998c03d52554efaf649bb10b35b6e5d64ac59a479023.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT telegram_chat_id FROM users WHERE id = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "telegram_chat_id", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [ + true + ] + }, + "hash": "a71a76724a3a7406e30a998c03d52554efaf649bb10b35b6e5d64ac59a479023" +} diff --git a/.sqlx/query-247470d26a90617e7dc9b5b3a2146ee3f54448e3c24943f7005e3a8e28820d43.json b/.sqlx/query-c48c9af71d8dab70ea2f11df30d2cc4976732a47430bd471f5aa47afc16e5740.json similarity index 82% rename from .sqlx/query-247470d26a90617e7dc9b5b3a2146ee3f54448e3c24943f7005e3a8e28820d43.json rename to .sqlx/query-c48c9af71d8dab70ea2f11df30d2cc4976732a47430bd471f5aa47afc16e5740.json index ac92ad5..9d34348 100644 --- a/.sqlx/query-247470d26a90617e7dc9b5b3a2146ee3f54448e3c24943f7005e3a8e28820d43.json +++ b/.sqlx/query-c48c9af71d8dab70ea2f11df30d2cc4976732a47430bd471f5aa47afc16e5740.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT\n email,\n preferred_comms_channel as \"preferred_channel!: CommsChannel\",\n discord_id,\n discord_verified,\n telegram_username,\n telegram_verified,\n signal_number,\n signal_verified\n FROM users WHERE did = $1", + "query": "SELECT\n email,\n preferred_comms_channel as \"preferred_channel!: CommsChannel\",\n discord_id,\n discord_verified,\n telegram_username,\n telegram_verified,\n telegram_chat_id,\n signal_number,\n signal_verified\n FROM users WHERE did = $1", "describe": { "columns": [ { @@ -47,11 +47,16 @@ }, { "ordinal": 6, + "name": "telegram_chat_id", + "type_info": "Int8" + }, + { + "ordinal": 7, "name": "signal_number", "type_info": "Text" }, { - "ordinal": 7, + "ordinal": 8, "name": "signal_verified", "type_info": "Bool" } @@ -69,8 +74,9 @@ true, false, true, + true, false ] }, - "hash": "247470d26a90617e7dc9b5b3a2146ee3f54448e3c24943f7005e3a8e28820d43" + "hash": "c48c9af71d8dab70ea2f11df30d2cc4976732a47430bd471f5aa47afc16e5740" } diff --git a/.sqlx/query-d54660032ca184c20d9cf4563d9a34934ff717796d4f564f1384c1c31fd76594.json b/.sqlx/query-d54660032ca184c20d9cf4563d9a34934ff717796d4f564f1384c1c31fd76594.json index c22e22b..c257c14 100644 --- a/.sqlx/query-d54660032ca184c20d9cf4563d9a34934ff717796d4f564f1384c1c31fd76594.json +++ b/.sqlx/query-d54660032ca184c20d9cf4563d9a34934ff717796d4f564f1384c1c31fd76594.json @@ -41,7 +41,8 @@ "channel_verification", "passkey_recovery", "legacy_login_alert", - "migration_verification" + "migration_verification", + "channel_verified" ] } } diff --git a/.sqlx/query-d7f32a31b4edeebbbf54a1878dfa05f2fcb3c57fe063be4e6a78fe5e74fb9dc3.json b/.sqlx/query-d7f32a31b4edeebbbf54a1878dfa05f2fcb3c57fe063be4e6a78fe5e74fb9dc3.json new file mode 100644 index 0000000..46d16e7 --- /dev/null +++ b/.sqlx/query-d7f32a31b4edeebbbf54a1878dfa05f2fcb3c57fe063be4e6a78fe5e74fb9dc3.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE users SET\n telegram_username = $1,\n telegram_verified = CASE WHEN LOWER(telegram_username) = LOWER($1) THEN telegram_verified ELSE FALSE END,\n telegram_chat_id = CASE WHEN LOWER(telegram_username) = LOWER($1) THEN telegram_chat_id ELSE NULL END,\n updated_at = NOW()\n WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "d7f32a31b4edeebbbf54a1878dfa05f2fcb3c57fe063be4e6a78fe5e74fb9dc3" +} diff --git a/.sqlx/query-c7ebbeca2ba26ef7b5a7c00441f51f68eb47f1421fa0a937eaa5a79aec75001d.json b/.sqlx/query-e49cbb17eb279cd12874fc3e7f5cbdbf0072c42127e225b79c0fe90de78670bd.json similarity index 58% rename from .sqlx/query-c7ebbeca2ba26ef7b5a7c00441f51f68eb47f1421fa0a937eaa5a79aec75001d.json rename to .sqlx/query-e49cbb17eb279cd12874fc3e7f5cbdbf0072c42127e225b79c0fe90de78670bd.json index e6ec9a2..fc1606f 100644 --- a/.sqlx/query-c7ebbeca2ba26ef7b5a7c00441f51f68eb47f1421fa0a937eaa5a79aec75001d.json +++ b/.sqlx/query-e49cbb17eb279cd12874fc3e7f5cbdbf0072c42127e225b79c0fe90de78670bd.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE users SET telegram_username = NULL, telegram_verified = FALSE, updated_at = NOW() WHERE id = $1", + "query": "UPDATE users SET telegram_username = NULL, telegram_verified = FALSE, telegram_chat_id = NULL, updated_at = NOW() WHERE id = $1", "describe": { "columns": [], "parameters": { @@ -10,5 +10,5 @@ }, "nullable": [] }, - "hash": "c7ebbeca2ba26ef7b5a7c00441f51f68eb47f1421fa0a937eaa5a79aec75001d" + "hash": "e49cbb17eb279cd12874fc3e7f5cbdbf0072c42127e225b79c0fe90de78670bd" } diff --git a/.sqlx/query-fc1ef8c0979206bf95cccee7c39614f7c882860c39cfe94e411cf6e1d55b63f4.json b/.sqlx/query-fc1ef8c0979206bf95cccee7c39614f7c882860c39cfe94e411cf6e1d55b63f4.json new file mode 100644 index 0000000..dcb8a19 --- /dev/null +++ b/.sqlx/query-fc1ef8c0979206bf95cccee7c39614f7c882860c39cfe94e411cf6e1d55b63f4.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE users SET telegram_chat_id = $2, telegram_verified = TRUE, updated_at = NOW() WHERE LOWER(telegram_username) = LOWER($1) AND telegram_username IS NOT NULL AND handle = $3 RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Text", + "Int8", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "fc1ef8c0979206bf95cccee7c39614f7c882860c39cfe94e411cf6e1d55b63f4" +} diff --git a/crates/tranquil-comms/src/locale.rs b/crates/tranquil-comms/src/locale.rs index c0bf56a..d34dfa9 100644 --- a/crates/tranquil-comms/src/locale.rs +++ b/crates/tranquil-comms/src/locale.rs @@ -31,6 +31,8 @@ pub struct NotificationStrings { pub legacy_login_body: &'static str, pub migration_verification_subject: &'static str, pub migration_verification_body: &'static str, + pub channel_verified_subject: &'static str, + pub channel_verified_body: &'static str, } pub fn get_strings(locale: &str) -> &'static NotificationStrings { @@ -66,6 +68,8 @@ static STRINGS_EN: NotificationStrings = NotificationStrings { legacy_login_body: "Hello @{handle},\n\nA login to your account was detected using a legacy app (like Bluesky) that doesn't support TOTP verification.\n\nDetails:\n- Time: {timestamp}\n- IP Address: {ip}\n\nYour TOTP protection was bypassed for this login. The session has limited permissions for sensitive operations.\n\nIf this wasn't you, please:\n1. Change your password immediately\n2. Review your active sessions\n3. Consider disabling legacy app logins in your security settings\n\nStay safe,\n{hostname}", migration_verification_subject: "Verify your email - {hostname}", migration_verification_body: "Welcome to {hostname}!\n\nYour account has been migrated successfully. To complete the setup, please verify your email address.\n\nYour verification code is:\n{code}\n\nCopy the code above and enter it at:\n{verify_page}\n\nThis code will expire in 48 hours.\n\nOr if you like to live dangerously:\n{verify_link}\n\nIf you did not migrate your account, please ignore this email.", + channel_verified_subject: "Channel verified - {hostname}", + channel_verified_body: "Hello {handle},\n\n{channel} has been verified as a notification channel for your account on {hostname}.", }; static STRINGS_ZH: NotificationStrings = NotificationStrings { @@ -90,6 +94,8 @@ static STRINGS_ZH: NotificationStrings = NotificationStrings { legacy_login_body: "您好 @{handle},\n\n检测到使用不支持 TOTP 验证的传统应用(如 Bluesky)登录您的账户。\n\n详细信息:\n- 时间:{timestamp}\n- IP 地址:{ip}\n\n此次登录绕过了 TOTP 保护。该会话对敏感操作的权限有限。\n\n如果这不是您的操作,请:\n1. 立即更改密码\n2. 检查您的活跃会话\n3. 考虑在安全设置中禁用传统应用登录\n\n请注意安全,\n{hostname}", migration_verification_subject: "验证您的邮箱 - {hostname}", migration_verification_body: "欢迎来到 {hostname}!\n\n您的账户已成功迁移。要完成设置,请验证您的邮箱地址。\n\n您的验证码是:\n{code}\n\n复制上述验证码并在此输入:\n{verify_page}\n\n此验证码将在 48 小时后过期。\n\n或者直接点击链接:\n{verify_link}\n\n如果您没有迁移账户,请忽略此邮件。", + channel_verified_subject: "通知渠道已验证 - {hostname}", + channel_verified_body: "您好 {handle},\n\n{channel} 已被验证为您在 {hostname} 上的通知渠道。", }; static STRINGS_JA: NotificationStrings = NotificationStrings { @@ -114,6 +120,8 @@ static STRINGS_JA: NotificationStrings = NotificationStrings { legacy_login_body: "@{handle} 様\n\nTOTP 認証に対応していないレガシーアプリ(Bluesky など)からのログインが検出されました。\n\n詳細:\n- 時刻:{timestamp}\n- IP アドレス:{ip}\n\nこのログインでは TOTP 保護がバイパスされました。このセッションは機密操作に対する権限が制限されています。\n\n心当たりがない場合は:\n1. 直ちにパスワードを変更してください\n2. アクティブなセッションを確認してください\n3. セキュリティ設定でレガシーアプリのログインを無効にすることを検討してください\n\nご注意ください。\n{hostname}", migration_verification_subject: "メールアドレスの認証 - {hostname}", migration_verification_body: "{hostname} へようこそ!\n\nアカウントの移行が完了しました。設定を完了するには、メールアドレスを認証してください。\n\n認証コードは:\n{code}\n\n上記のコードをコピーして、こちらで入力してください:\n{verify_page}\n\nこのコードは48時間後に期限切れとなります。\n\n自己責任でワンクリック認証:\n{verify_link}\n\nアカウントを移行していない場合は、このメールを無視してください。", + channel_verified_subject: "通知チャンネル認証完了 - {hostname}", + channel_verified_body: "{handle} 様\n\n{channel} が {hostname} の通知チャンネルとして認証されました。", }; static STRINGS_KO: NotificationStrings = NotificationStrings { @@ -138,6 +146,8 @@ static STRINGS_KO: NotificationStrings = NotificationStrings { legacy_login_body: "안녕하세요 @{handle}님,\n\nTOTP 인증을 지원하지 않는 레거시 앱(예: Bluesky)을 사용한 로그인이 감지되었습니다.\n\n세부 정보:\n- 시간: {timestamp}\n- IP 주소: {ip}\n\n이 로그인에서 TOTP 보호가 우회되었습니다. 이 세션은 민감한 작업에 대한 권한이 제한됩니다.\n\n본인이 아닌 경우:\n1. 즉시 비밀번호를 변경하세요\n2. 활성 세션을 검토하세요\n3. 보안 설정에서 레거시 앱 로그인 비활성화를 고려하세요\n\n{hostname} 드림", migration_verification_subject: "이메일 인증 - {hostname}", migration_verification_body: "{hostname}에 오신 것을 환영합니다!\n\n계정 마이그레이션이 완료되었습니다. 설정을 완료하려면 이메일 주소를 인증하세요.\n\n인증 코드는:\n{code}\n\n위 코드를 복사하여 여기에 입력하세요:\n{verify_page}\n\n이 코드는 48시간 후에 만료됩니다.\n\n위험을 감수하고 원클릭 인증:\n{verify_link}\n\n계정을 마이그레이션하지 않았다면 이 이메일을 무시하세요.", + channel_verified_subject: "알림 채널 인증 완료 - {hostname}", + channel_verified_body: "안녕하세요 {handle}님,\n\n{channel}이(가) {hostname}의 알림 채널로 인증되었습니다.", }; static STRINGS_SV: NotificationStrings = NotificationStrings { @@ -162,6 +172,8 @@ static STRINGS_SV: NotificationStrings = NotificationStrings { legacy_login_body: "Hej @{handle},\n\nEn inloggning till ditt konto upptäcktes med en äldre app (som Bluesky) som inte stöder TOTP-verifiering.\n\nDetaljer:\n- Tid: {timestamp}\n- IP-adress: {ip}\n\nDitt TOTP-skydd kringgicks för denna inloggning. Sessionen har begränsade behörigheter för känsliga operationer.\n\nOm detta inte var du:\n1. Ändra ditt lösenord omedelbart\n2. Granska dina aktiva sessioner\n3. Överväg att inaktivera äldre appinloggningar i dina säkerhetsinställningar\n\nVar försiktig,\n{hostname}", migration_verification_subject: "Verifiera din e-post - {hostname}", migration_verification_body: "Välkommen till {hostname}!\n\nDitt konto har migrerats framgångsrikt. För att slutföra installationen, verifiera din e-postadress.\n\nDin verifieringskod är:\n{code}\n\nKopiera koden ovan och ange den på:\n{verify_page}\n\nDenna kod upphör om 48 timmar.\n\nEller om du gillar att leva farligt:\n{verify_link}\n\nOm du inte migrerade ditt konto kan du ignorera detta meddelande.", + channel_verified_subject: "Aviseringskanal verifierad - {hostname}", + channel_verified_body: "Hej {handle},\n\n{channel} har verifierats som aviseringskanal för ditt konto på {hostname}.", }; static STRINGS_FI: NotificationStrings = NotificationStrings { @@ -186,6 +198,8 @@ static STRINGS_FI: NotificationStrings = NotificationStrings { legacy_login_body: "Hei @{handle},\n\nTilillesi havaittiin kirjautuminen vanhalla sovelluksella (kuten Bluesky), joka ei tue TOTP-vahvistusta.\n\nTiedot:\n- Aika: {timestamp}\n- IP-osoite: {ip}\n\nTOTP-suojauksesi ohitettiin tässä kirjautumisessa. Istunnolla on rajoitetut oikeudet arkaluontoisiin toimintoihin.\n\nJos tämä et ollut sinä:\n1. Vaihda salasanasi välittömästi\n2. Tarkista aktiiviset istuntosi\n3. Harkitse vanhojen sovellusten kirjautumisen poistamista käytöstä turvallisuusasetuksissa\n\nOle varovainen,\n{hostname}", migration_verification_subject: "Vahvista sähköpostisi - {hostname}", migration_verification_body: "Tervetuloa palveluun {hostname}!\n\nTilisi on siirretty onnistuneesti. Viimeistele asennus vahvistamalla sähköpostiosoitteesi.\n\nVahvistuskoodisi on:\n{code}\n\nKopioi koodi yllä ja syötä se osoitteessa:\n{verify_page}\n\nTämä koodi vanhenee 48 tunnissa.\n\nTai jos pidät vaarallisesta elämästä:\n{verify_link}\n\nJos et siirtänyt tiliäsi, voit jättää tämän viestin huomiotta.", + channel_verified_subject: "Ilmoituskanava vahvistettu - {hostname}", + channel_verified_body: "Hei {handle},\n\n{channel} on vahvistettu ilmoituskanavaksi tilillesi palvelussa {hostname}.", }; pub fn format_message(template: &str, vars: &[(&str, &str)]) -> String { diff --git a/crates/tranquil-comms/src/sender.rs b/crates/tranquil-comms/src/sender.rs index d71f8f0..bd0dab8 100644 --- a/crates/tranquil-comms/src/sender.rs +++ b/crates/tranquil-comms/src/sender.rs @@ -67,6 +67,12 @@ pub fn mime_encode_header(value: &str) -> String { } } +pub fn escape_html(text: &str) -> String { + text.replace('&', "&") + .replace('<', "<") + .replace('>', ">") +} + pub fn is_valid_phone_number(number: &str) -> bool { if number.len() < 2 || number.len() > 20 { return false; @@ -244,6 +250,60 @@ impl TelegramSender { let bot_token = std::env::var("TELEGRAM_BOT_TOKEN").ok()?; Some(Self::new(bot_token)) } + + pub async fn set_webhook( + &self, + webhook_url: &str, + secret_token: Option<&str>, + ) -> Result<(), SendError> { + let url = format!("https://api.telegram.org/bot{}/setWebhook", self.bot_token); + let mut payload = json!({ "url": webhook_url }); + if let Some(secret) = secret_token { + payload["secret_token"] = json!(secret); + } + let response = self + .http_client + .post(&url) + .json(&payload) + .send() + .await + .map_err(|e| SendError::ExternalService(format!("setWebhook request failed: {}", e)))?; + if !response.status().is_success() { + let body = response.text().await.unwrap_or_default(); + return Err(SendError::ExternalService(format!( + "setWebhook returned error: {}", + body + ))); + } + Ok(()) + } + + pub async fn resolve_bot_username(&self) -> Result { + let url = format!("https://api.telegram.org/bot{}/getMe", self.bot_token); + let response = self.http_client.get(&url).send().await.map_err(|e| { + SendError::ExternalService(format!("Telegram getMe request failed: {}", e)) + })?; + + if !response.status().is_success() { + let body = response.text().await.unwrap_or_default(); + return Err(SendError::ExternalService(format!( + "Telegram getMe returned error: {}", + body + ))); + } + + let data: serde_json::Value = response.json().await.map_err(|e| { + SendError::ExternalService(format!("Failed to parse getMe response: {}", e)) + })?; + + data.get("result") + .and_then(|r| r.get("username")) + .and_then(|u| u.as_str()) + .map(|s| s.to_string()) + .ok_or_else(|| { + SendError::ExternalService("getMe response missing username".to_string()) + }) + } } #[async_trait] @@ -254,13 +314,14 @@ impl CommsSender for TelegramSender { async fn send(&self, notification: &QueuedComms) -> Result<(), SendError> { let chat_id = ¬ification.recipient; - let subject = notification.subject.as_deref().unwrap_or("Notification"); - let text = format!("*{}*\n\n{}", subject, notification.body); + let subject = escape_html(notification.subject.as_deref().unwrap_or("Notification")); + let body = escape_html(¬ification.body); + let text = format!("{}\n\n{}", subject, body); let url = format!("https://api.telegram.org/bot{}/sendMessage", self.bot_token); let payload = json!({ "chat_id": chat_id, "text": text, - "parse_mode": "Markdown" + "parse_mode": "HTML" }); let mut last_error = None; for attempt in 0..MAX_RETRIES { diff --git a/crates/tranquil-db-traits/src/infra.rs b/crates/tranquil-db-traits/src/infra.rs index 82a11b1..86008e5 100644 --- a/crates/tranquil-db-traits/src/infra.rs +++ b/crates/tranquil-db-traits/src/infra.rs @@ -78,6 +78,7 @@ pub enum CommsType { LegacyLoginAlert, MigrationVerification, ChannelVerification, + ChannelVerified, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, sqlx::Type)] diff --git a/crates/tranquil-db-traits/src/user.rs b/crates/tranquil-db-traits/src/user.rs index c7e7f28..c34c150 100644 --- a/crates/tranquil-db-traits/src/user.rs +++ b/crates/tranquil-db-traits/src/user.rs @@ -196,6 +196,12 @@ pub trait UserRepository: Send + Sync { identifier: &str, ) -> Result, DbError>; + async fn check_channel_verified_by_did( + &self, + did: &Did, + channel: CommsChannel, + ) -> Result, DbError>; + async fn admin_update_email(&self, did: &Did, email: &str) -> Result; async fn admin_update_handle(&self, did: &Did, handle: &Handle) -> Result; @@ -222,6 +228,21 @@ pub trait UserRepository: Send + Sync { async fn clear_signal(&self, user_id: Uuid) -> Result<(), DbError>; + async fn set_unverified_telegram( + &self, + user_id: Uuid, + telegram_username: &str, + ) -> Result<(), DbError>; + + async fn store_telegram_chat_id( + &self, + telegram_username: &str, + chat_id: i64, + handle: Option<&str>, + ) -> Result, DbError>; + + async fn get_telegram_chat_id(&self, user_id: Uuid) -> Result, DbError>; + async fn get_verification_info( &self, did: &Did, @@ -575,6 +596,9 @@ pub struct UserCommsPrefs { pub handle: Handle, pub preferred_channel: CommsChannel, pub preferred_locale: Option, + pub telegram_chat_id: Option, + pub discord_id: Option, + pub signal_number: Option, } #[derive(Debug, Clone)] @@ -635,6 +659,7 @@ pub struct NotificationPrefs { pub discord_verified: bool, pub telegram_username: Option, pub telegram_verified: bool, + pub telegram_chat_id: Option, pub signal_number: Option, pub signal_verified: bool, } diff --git a/crates/tranquil-db/src/postgres/user.rs b/crates/tranquil-db/src/postgres/user.rs index ca23487..644ad15 100644 --- a/crates/tranquil-db/src/postgres/user.rs +++ b/crates/tranquil-db/src/postgres/user.rs @@ -311,7 +311,7 @@ impl UserRepository for PostgresUserRepository { async fn get_comms_prefs(&self, user_id: Uuid) -> Result, DbError> { let row = sqlx::query!( - r#"SELECT email, handle, preferred_comms_channel as "preferred_channel!: CommsChannel", preferred_locale + r#"SELECT email, handle, preferred_comms_channel as "preferred_channel!: CommsChannel", preferred_locale, telegram_chat_id, discord_id, signal_number FROM users WHERE id = $1"#, user_id ) @@ -323,6 +323,9 @@ impl UserRepository for PostgresUserRepository { handle: Handle::from(r.handle), preferred_channel: r.preferred_channel, preferred_locale: r.preferred_locale, + telegram_chat_id: r.telegram_chat_id, + discord_id: r.discord_id, + signal_number: r.signal_number, })) } @@ -561,6 +564,33 @@ impl UserRepository for PostgresUserRepository { Ok(row) } + async fn check_channel_verified_by_did( + &self, + did: &Did, + channel: CommsChannel, + ) -> Result, DbError> { + let row = sqlx::query!( + r#"SELECT + email_verified, + discord_verified, + telegram_verified, + signal_verified + FROM users + WHERE did = $1"#, + did.as_str() + ) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(row.map(|r| match channel { + CommsChannel::Email => r.email_verified, + CommsChannel::Discord => r.discord_verified, + CommsChannel::Telegram => r.telegram_verified, + CommsChannel::Signal => r.signal_verified, + })) + } + async fn admin_update_email(&self, did: &Did, email: &str) -> Result { let result = sqlx::query!( "UPDATE users SET email = $1 WHERE did = $2", @@ -609,6 +639,7 @@ impl UserRepository for PostgresUserRepository { discord_verified, telegram_username, telegram_verified, + telegram_chat_id, signal_number, signal_verified FROM users WHERE did = $1"#, @@ -624,6 +655,7 @@ impl UserRepository for PostgresUserRepository { discord_verified: r.discord_verified, telegram_username: r.telegram_username, telegram_verified: r.telegram_verified, + telegram_chat_id: r.telegram_chat_id, signal_number: r.signal_number, signal_verified: r.signal_verified, })) @@ -676,7 +708,7 @@ impl UserRepository for PostgresUserRepository { async fn clear_telegram(&self, user_id: Uuid) -> Result<(), DbError> { sqlx::query!( - "UPDATE users SET telegram_username = NULL, telegram_verified = FALSE, updated_at = NOW() WHERE id = $1", + "UPDATE users SET telegram_username = NULL, telegram_verified = FALSE, telegram_chat_id = NULL, updated_at = NOW() WHERE id = $1", user_id ) .execute(&self.pool) @@ -3136,4 +3168,66 @@ impl UserRepository for PostgresUserRepository { passkeys_deleted: deleted.rows_affected(), }) } + + async fn set_unverified_telegram( + &self, + user_id: Uuid, + telegram_username: &str, + ) -> Result<(), DbError> { + sqlx::query!( + r#"UPDATE users SET + telegram_username = $1, + telegram_verified = CASE WHEN LOWER(telegram_username) = LOWER($1) THEN telegram_verified ELSE FALSE END, + telegram_chat_id = CASE WHEN LOWER(telegram_username) = LOWER($1) THEN telegram_chat_id ELSE NULL END, + updated_at = NOW() + WHERE id = $2"#, + telegram_username, + user_id + ) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } + + async fn store_telegram_chat_id( + &self, + telegram_username: &str, + chat_id: i64, + handle: Option<&str>, + ) -> Result, DbError> { + let result = match handle { + Some(h) => sqlx::query_scalar!( + "UPDATE users SET telegram_chat_id = $2, telegram_verified = TRUE, updated_at = NOW() WHERE LOWER(telegram_username) = LOWER($1) AND telegram_username IS NOT NULL AND handle = $3 RETURNING id", + telegram_username, + chat_id, + h + ) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?, + None => sqlx::query_scalar!( + r#"UPDATE users SET telegram_chat_id = $2, telegram_verified = TRUE, updated_at = NOW() + WHERE id = ( + SELECT id FROM users + WHERE LOWER(telegram_username) = LOWER($1) AND telegram_username IS NOT NULL AND deactivated_at IS NULL + LIMIT 1 + ) RETURNING id"#, + telegram_username, + chat_id + ) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?, + }; + Ok(result) + } + + async fn get_telegram_chat_id(&self, user_id: Uuid) -> Result, DbError> { + let row = sqlx::query_scalar!("SELECT telegram_chat_id FROM users WHERE id = $1", user_id) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(row.flatten()) + } } diff --git a/crates/tranquil-pds/src/api/identity/account.rs b/crates/tranquil-pds/src/api/identity/account.rs index 4ff1a0a..73b4054 100644 --- a/crates/tranquil-pds/src/api/identity/account.rs +++ b/crates/tranquil-pds/src/api/identity/account.rs @@ -208,7 +208,15 @@ pub async fn create_account( _ => return ApiError::MissingDiscordId.into_response(), }, "telegram" => match &input.telegram_username { - Some(username) if !username.trim().is_empty() => username.trim().to_string(), + Some(username) if !username.trim().is_empty() => { + let clean = username.trim().trim_start_matches('@'); + if !crate::api::validation::is_valid_telegram_username(clean) { + return ApiError::InvalidRequest( + "Invalid Telegram username. Must be 5-32 characters, alphanumeric or underscore".into(), + ).into_response(); + } + clean.to_string() + } _ => return ApiError::MissingTelegramUsername.into_response(), }, "signal" => match &input.signal_number { @@ -634,7 +642,7 @@ pub async fn create_account( telegram_username: input .telegram_username .as_deref() - .map(|s| s.trim()) + .map(|s| s.trim().trim_start_matches('@')) .filter(|s| !s.is_empty()) .map(String::from), signal_number: input diff --git a/crates/tranquil-pds/src/api/mod.rs b/crates/tranquil-pds/src/api/mod.rs index 51bbfc0..493f1fc 100644 --- a/crates/tranquil-pds/src/api/mod.rs +++ b/crates/tranquil-pds/src/api/mod.rs @@ -12,6 +12,7 @@ pub mod proxy_client; pub mod repo; pub mod responses; pub mod server; +pub mod telegram_webhook; pub mod temp; pub mod validation; pub mod verification; diff --git a/crates/tranquil-pds/src/api/notification_prefs.rs b/crates/tranquil-pds/src/api/notification_prefs.rs index e8dda02..f47dd03 100644 --- a/crates/tranquil-pds/src/api/notification_prefs.rs +++ b/crates/tranquil-pds/src/api/notification_prefs.rs @@ -165,15 +165,37 @@ pub async fn request_channel_verification( "signal" => tranquil_db_traits::CommsChannel::Signal, _ => return Err("Invalid channel".to_string()), }; + let hostname = pds_hostname(); + let encoded_token = urlencoding::encode(&formatted_token); + let encoded_identifier = urlencoding::encode(identifier); + let verify_link = format!( + "https://{}/app/verify?token={}&identifier={}", + hostname, encoded_token, encoded_identifier + ); + let body = format!( + "Your verification code is: {}\n\nOr verify directly:\n{}", + formatted_token, verify_link + ); + let recipient = match comms_channel { + tranquil_db_traits::CommsChannel::Telegram => state + .user_repo + .get_telegram_chat_id(user_id) + .await + .ok() + .flatten() + .map(|id| id.to_string()) + .unwrap_or_else(|| identifier.to_string()), + _ => identifier.to_string(), + }; state .infra_repo .enqueue_comms( Some(user_id), comms_channel, tranquil_db_traits::CommsType::ChannelVerification, - identifier, + &recipient, Some("Verify your channel"), - &format!("Your verification code is: {}", formatted_token), + &body, Some(json!({"code": formatted_token})), ) .await @@ -199,6 +221,28 @@ pub async fn update_notification_prefs( let handle = user_row.handle; let current_email = user_row.email; + let current_prefs = state + .user_repo + .get_notification_prefs(&auth.did) + .await + .map_err(|e| ApiError::InternalError(Some(format!("Database error: {}", e))))? + .ok_or(ApiError::AccountNotFound)?; + + let effective_channel = input + .preferred_channel + .as_deref() + .map(|ch| match ch { + "email" => Ok(CommsChannel::Email), + "discord" => Ok(CommsChannel::Discord), + "telegram" => Ok(CommsChannel::Telegram), + "signal" => Ok(CommsChannel::Signal), + _ => Err(ApiError::InvalidRequest( + "Invalid channel. Must be one of: email, discord, telegram, signal".into(), + )), + }) + .transpose()? + .unwrap_or(current_prefs.preferred_channel); + let mut verification_required: Vec = Vec::new(); if let Some(ref channel_str) = input.preferred_channel { @@ -249,6 +293,11 @@ pub async fn update_notification_prefs( if let Some(ref discord_id) = input.discord_id { if discord_id.is_empty() { + if effective_channel == CommsChannel::Discord { + return Err(ApiError::InvalidRequest( + "Cannot remove Discord while it is the preferred notification channel".into(), + )); + } state .user_repo .clear_discord(user_id) @@ -267,30 +316,40 @@ pub async fn update_notification_prefs( if let Some(ref telegram) = input.telegram_username { let telegram_clean = telegram.trim_start_matches('@'); if telegram_clean.is_empty() { + if effective_channel == CommsChannel::Telegram { + return Err(ApiError::InvalidRequest( + "Cannot remove Telegram while it is the preferred notification channel".into(), + )); + } state .user_repo .clear_telegram(user_id) .await .map_err(|e| ApiError::InternalError(Some(format!("Database error: {}", e))))?; info!(did = %auth.did, "Cleared Telegram username"); + } else if !crate::api::validation::is_valid_telegram_username(telegram_clean) { + return Err(ApiError::InvalidRequest( + "Invalid Telegram username. Must be 5-32 characters, alphanumeric or underscore" + .into(), + )); } else { - request_channel_verification( - &state, - user_id, - &auth.did, - "telegram", - telegram_clean, - None, - ) - .await - .map_err(|e| ApiError::InternalError(Some(e)))?; + state + .user_repo + .set_unverified_telegram(user_id, telegram_clean) + .await + .map_err(|e| ApiError::InternalError(Some(format!("Database error: {}", e))))?; verification_required.push("telegram".to_string()); - info!(did = %auth.did, "Requested Telegram verification"); + info!(did = %auth.did, telegram_username = %telegram_clean, "Stored unverified Telegram username"); } } if let Some(ref signal) = input.signal_number { if signal.is_empty() { + if effective_channel == CommsChannel::Signal { + return Err(ApiError::InvalidRequest( + "Cannot remove Signal while it is the preferred notification channel".into(), + )); + } state .user_repo .clear_signal(user_id) diff --git a/crates/tranquil-pds/src/api/server/email.rs b/crates/tranquil-pds/src/api/server/email.rs index 3f51602..986f002 100644 --- a/crates/tranquil-pds/src/api/server/email.rs +++ b/crates/tranquil-pds/src/api/server/email.rs @@ -419,6 +419,45 @@ pub async fn check_email_verified( } } +#[derive(Deserialize)] +pub struct CheckChannelVerifiedInput { + pub did: String, + pub channel: String, +} + +pub async fn check_channel_verified( + State(state): State, + _rate_limit: RateLimited, + Json(input): Json, +) -> Response { + let channel = match input.channel.to_lowercase().as_str() { + "email" => CommsChannel::Email, + "discord" => CommsChannel::Discord, + "telegram" => CommsChannel::Telegram, + "signal" => CommsChannel::Signal, + _ => { + return ApiError::InvalidRequest("invalid channel".into()).into_response(); + } + }; + + let did = match crate::Did::new(input.did) { + Ok(d) => d, + Err(_) => return ApiError::InvalidRequest("invalid did".into()).into_response(), + }; + match state + .user_repo + .check_channel_verified_by_did(&did, channel) + .await + { + Ok(Some(verified)) => VerifiedResponse::response(verified).into_response(), + Ok(None) => ApiError::AccountNotFound.into_response(), + Err(e) => { + error!("DB error checking channel verified: {:?}", e); + ApiError::InternalError(None).into_response() + } + } +} + #[derive(Deserialize)] pub struct AuthorizeEmailUpdateQuery { pub token: String, diff --git a/crates/tranquil-pds/src/api/server/meta.rs b/crates/tranquil-pds/src/api/server/meta.rs index f48a8ce..26384f9 100644 --- a/crates/tranquil-pds/src/api/server/meta.rs +++ b/crates/tranquil-pds/src/api/server/meta.rs @@ -1,5 +1,5 @@ use crate::state::AppState; -use crate::util::pds_hostname; +use crate::util::{pds_hostname, telegram_bot_username}; use axum::{Json, extract::State, http::StatusCode, response::IntoResponse}; use serde_json::json; @@ -52,7 +52,7 @@ pub async fn describe_server() -> impl IntoResponse { if let Some(email) = contact_email { contact.insert("email".to_string(), json!(email)); } - Json(json!({ + let mut response = json!({ "availableUserDomains": domains, "inviteCodeRequired": invite_code_required, "did": format!("did:web:{}", pds_hostname), @@ -61,7 +61,11 @@ pub async fn describe_server() -> impl IntoResponse { "version": env!("CARGO_PKG_VERSION"), "availableCommsChannels": get_available_comms_channels(), "selfHostedDidWebEnabled": is_self_hosted_did_web_enabled() - })) + }); + if let Some(bot_username) = telegram_bot_username() { + response["telegramBotUsername"] = json!(bot_username); + } + Json(response) } pub async fn health(State(state): State) -> impl IntoResponse { match state.infra_repo.health_check().await { diff --git a/crates/tranquil-pds/src/api/server/mod.rs b/crates/tranquil-pds/src/api/server/mod.rs index 3b4499a..2d051aa 100644 --- a/crates/tranquil-pds/src/api/server/mod.rs +++ b/crates/tranquil-pds/src/api/server/mod.rs @@ -23,7 +23,7 @@ pub use account_status::{ }; pub use app_password::{create_app_password, list_app_passwords, revoke_app_password}; pub use email::{ - authorize_email_update, check_comms_channel_in_use, check_email_in_use, + authorize_email_update, check_channel_verified, check_comms_channel_in_use, check_email_in_use, check_email_update_status, check_email_verified, confirm_email, request_email_update, update_email, }; diff --git a/crates/tranquil-pds/src/api/server/passkey_account.rs b/crates/tranquil-pds/src/api/server/passkey_account.rs index 38312d8..461d120 100644 --- a/crates/tranquil-pds/src/api/server/passkey_account.rs +++ b/crates/tranquil-pds/src/api/server/passkey_account.rs @@ -172,7 +172,15 @@ pub async fn create_passkey_account( _ => return ApiError::MissingDiscordId.into_response(), }, "telegram" => match &input.telegram_username { - Some(username) if !username.trim().is_empty() => username.trim().to_string(), + Some(username) if !username.trim().is_empty() => { + let clean = username.trim().trim_start_matches('@'); + if !crate::api::validation::is_valid_telegram_username(clean) { + return ApiError::InvalidRequest( + "Invalid Telegram username. Must be 5-32 characters, alphanumeric or underscore".into(), + ).into_response(); + } + clean.to_string() + } _ => return ApiError::MissingTelegramUsername.into_response(), }, "signal" => match &input.signal_number { @@ -410,7 +418,7 @@ pub async fn create_passkey_account( telegram_username: input .telegram_username .as_deref() - .map(|s| s.trim()) + .map(|s| s.trim().trim_start_matches('@')) .filter(|s| !s.is_empty()) .map(String::from), signal_number: input diff --git a/crates/tranquil-pds/src/api/server/verify_token.rs b/crates/tranquil-pds/src/api/server/verify_token.rs index 724b5fd..57ed01d 100644 --- a/crates/tranquil-pds/src/api/server/verify_token.rs +++ b/crates/tranquil-pds/src/api/server/verify_token.rs @@ -1,5 +1,7 @@ use crate::api::error::{ApiError, DbResultExt}; +use crate::comms::comms_repo; use crate::types::Did; +use crate::util::pds_hostname; use axum::{Json, extract::State}; use serde::{Deserialize, Serialize}; use tracing::{info, warn}; @@ -161,6 +163,20 @@ async fn handle_channel_update( info!(did = %did, channel = %channel, "Channel verified successfully"); + let recipient = resolve_verified_recipient(state, user_id, channel, identifier).await; + if let Err(e) = comms_repo::enqueue_channel_verified( + state.user_repo.as_ref(), + state.infra_repo.as_ref(), + user_id, + channel, + &recipient, + pds_hostname(), + ) + .await + { + warn!(error = %e, "Failed to enqueue channel verified notification"); + } + Ok(Json(VerifyTokenOutput { success: true, did: did.to_string().into(), @@ -169,11 +185,30 @@ async fn handle_channel_update( })) } +async fn resolve_verified_recipient( + state: &AppState, + user_id: uuid::Uuid, + channel: &str, + identifier: &str, +) -> String { + match channel { + "telegram" => state + .user_repo + .get_telegram_chat_id(user_id) + .await + .ok() + .flatten() + .map(|id| id.to_string()) + .unwrap_or_else(|| identifier.to_string()), + _ => identifier.to_string(), + } +} + async fn handle_signup_verification( state: &AppState, did: &str, channel: &str, - _identifier: &str, + identifier: &str, ) -> Result, ApiError> { let did_typed: Did = did .parse() @@ -232,6 +267,20 @@ async fn handle_signup_verification( info!(did = %did, channel = %channel, "Signup verified successfully"); + let recipient = resolve_verified_recipient(state, user.id, channel, identifier).await; + if let Err(e) = comms_repo::enqueue_channel_verified( + state.user_repo.as_ref(), + state.infra_repo.as_ref(), + user.id, + channel, + &recipient, + pds_hostname(), + ) + .await + { + warn!(error = %e, "Failed to enqueue channel verified notification"); + } + Ok(Json(VerifyTokenOutput { success: true, did: did.to_string().into(), diff --git a/crates/tranquil-pds/src/api/telegram_webhook.rs b/crates/tranquil-pds/src/api/telegram_webhook.rs new file mode 100644 index 0000000..dc84b4e --- /dev/null +++ b/crates/tranquil-pds/src/api/telegram_webhook.rs @@ -0,0 +1,172 @@ +use axum::{ + extract::State, + http::{HeaderMap, StatusCode}, + response::IntoResponse, +}; +use serde::Deserialize; +use tracing::{debug, info, warn}; + +use crate::comms::comms_repo; +use crate::state::AppState; +use crate::util::pds_hostname; + +#[derive(Deserialize)] +struct TelegramUpdate { + message: Option, +} + +#[derive(Deserialize)] +struct TelegramMessage { + text: Option, + from: Option, +} + +#[derive(Deserialize)] +struct TelegramUser { + id: i64, + username: Option, +} + +pub async fn handle_telegram_webhook( + State(state): State, + headers: HeaderMap, + body: String, +) -> impl IntoResponse { + let expected_secret = match std::env::var("TELEGRAM_WEBHOOK_SECRET") { + Ok(s) => s, + Err(_) => { + warn!("Telegram webhook called but TELEGRAM_WEBHOOK_SECRET is not configured"); + return StatusCode::FORBIDDEN; + } + }; + let provided = headers + .get("x-telegram-bot-api-secret-token") + .and_then(|v| v.to_str().ok()) + .unwrap_or_default(); + if provided != expected_secret { + warn!("Telegram webhook received with invalid secret token"); + return StatusCode::UNAUTHORIZED; + } + + let update: TelegramUpdate = match serde_json::from_str(&body) { + Ok(u) => u, + Err(_) => return StatusCode::OK, + }; + + if let Some(message) = update.message { + let is_start = message + .text + .as_deref() + .is_some_and(|t| t.starts_with("/start")); + + if is_start + && let Some(from) = message.from + && let Some(username) = from.username + { + let handle = parse_start_handle(message.text.as_deref()); + + debug!( + telegram_username = %username, + chat_id = from.id, + handle = ?handle, + "Received /start from Telegram user" + ); + match state + .user_repo + .store_telegram_chat_id(&username, from.id, handle.as_deref()) + .await + { + Ok(Some(user_id)) => { + info!( + telegram_username = %username, + chat_id = from.id, + "Verified Telegram user and stored chat_id" + ); + if let Err(e) = comms_repo::enqueue_channel_verified( + state.user_repo.as_ref(), + state.infra_repo.as_ref(), + user_id, + "telegram", + &from.id.to_string(), + pds_hostname(), + ) + .await + { + warn!(error = %e, "Failed to enqueue channel verified notification"); + } + } + Ok(None) => { + debug!( + telegram_username = %username, + "No matching user found for Telegram username" + ); + } + Err(e) => { + warn!( + telegram_username = %username, + error = %e, + "Failed to store Telegram chat_id" + ); + } + } + } + } + + StatusCode::OK +} + +fn parse_start_handle(text: Option<&str>) -> Option { + text.and_then(|t| t.strip_prefix("/start ")) + .map(|payload| payload.trim()) + .filter(|p| !p.is_empty()) + .map(|payload| payload.replace('_', ".")) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn deep_link_underscores_decoded_to_dots() { + assert_eq!( + parse_start_handle(Some("/start lewis_buttercup_wizardry_systems")), + Some("lewis.buttercup.wizardry.systems".to_string()), + ); + } + + #[test] + fn manual_handle_with_dots_passes_through() { + assert_eq!( + parse_start_handle(Some("/start lewis.buttercup.wizardry.systems")), + Some("lewis.buttercup.wizardry.systems".to_string()), + ); + } + + #[test] + fn bare_start_returns_none() { + assert_eq!(parse_start_handle(Some("/start")), None); + } + + #[test] + fn start_with_trailing_space_returns_none() { + assert_eq!(parse_start_handle(Some("/start ")), None); + } + + #[test] + fn none_text_returns_none() { + assert_eq!(parse_start_handle(None), None); + } + + #[test] + fn non_start_command_returns_none() { + assert_eq!(parse_start_handle(Some("/help")), None); + } + + #[test] + fn payload_with_extra_whitespace_trimmed() { + assert_eq!( + parse_start_handle(Some("/start alice_example_com ")), + Some("alice.example.com".to_string()), + ); + } +} diff --git a/crates/tranquil-pds/src/api/validation.rs b/crates/tranquil-pds/src/api/validation.rs index d1b121c..5978934 100644 --- a/crates/tranquil-pds/src/api/validation.rs +++ b/crates/tranquil-pds/src/api/validation.rs @@ -349,6 +349,14 @@ pub fn is_valid_email(email: &str) -> bool { }) } +pub fn is_valid_telegram_username(username: &str) -> bool { + let clean = username.strip_prefix('@').unwrap_or(username); + (5..=32).contains(&clean.len()) + && clean + .chars() + .all(|c| c.is_ascii_alphanumeric() || c == '_') +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/tranquil-pds/src/comms/service.rs b/crates/tranquil-pds/src/comms/service.rs index e5d23d9..4e76308 100644 --- a/crates/tranquil-pds/src/comms/service.rs +++ b/crates/tranquil-pds/src/comms/service.rs @@ -10,7 +10,7 @@ use tranquil_comms::{ CommsChannel, CommsSender, CommsStatus, CommsType, NewComms, SendError, format_message, get_strings, }; -use tranquil_db_traits::{InfraRepository, QueuedComms, UserRepository}; +use tranquil_db_traits::{InfraRepository, QueuedComms, UserCommsPrefs, UserRepository}; use uuid::Uuid; pub struct CommsService { @@ -75,6 +75,7 @@ impl CommsService { tranquil_db_traits::CommsType::MigrationVerification } CommsType::ChannelVerification => tranquil_db_traits::CommsType::ChannelVerification, + CommsType::ChannelVerified => tranquil_db_traits::CommsType::ChannelVerified, }; let id = self .infra_repo @@ -170,6 +171,7 @@ impl CommsService { tranquil_db_traits::CommsType::ChannelVerification => { CommsType::ChannelVerification } + tranquil_db_traits::CommsType::ChannelVerified => CommsType::ChannelVerified, }, status: match item.status { tranquil_db_traits::CommsStatus::Pending => CommsStatus::Pending, @@ -247,6 +249,49 @@ pub fn channel_display_name(channel: CommsChannel) -> &'static str { } } +struct ResolvedRecipient { + channel: tranquil_db_traits::CommsChannel, + recipient: String, +} + +fn resolve_recipient( + prefs: &UserCommsPrefs, + channel: tranquil_db_traits::CommsChannel, +) -> ResolvedRecipient { + let email_fallback = || ResolvedRecipient { + channel: tranquil_db_traits::CommsChannel::Email, + recipient: prefs.email.clone().unwrap_or_default(), + }; + match channel { + tranquil_db_traits::CommsChannel::Email => email_fallback(), + tranquil_db_traits::CommsChannel::Telegram => prefs + .telegram_chat_id + .map(|id| ResolvedRecipient { + channel, + recipient: id.to_string(), + }) + .unwrap_or_else(email_fallback), + tranquil_db_traits::CommsChannel::Discord => prefs + .discord_id + .as_ref() + .filter(|id| !id.is_empty()) + .map(|id| ResolvedRecipient { + channel, + recipient: id.clone(), + }) + .unwrap_or_else(email_fallback), + tranquil_db_traits::CommsChannel::Signal => prefs + .signal_number + .as_ref() + .filter(|n| !n.is_empty()) + .map(|n| ResolvedRecipient { + channel, + recipient: n.clone(), + }) + .unwrap_or_else(email_fallback), + } +} + fn channel_from_str(s: &str) -> tranquil_db_traits::CommsChannel { match s { "discord" => tranquil_db_traits::CommsChannel::Discord, @@ -276,13 +321,13 @@ pub mod repo { &[("hostname", hostname), ("handle", &prefs.handle)], ); let subject = format_message(strings.welcome_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::Welcome, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -307,13 +352,13 @@ pub mod repo { &[("handle", &prefs.handle), ("code", code)], ); let subject = format_message(strings.password_reset_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::PasswordReset, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -471,13 +516,13 @@ pub mod repo { &[("handle", &prefs.handle), ("code", code)], ); let subject = format_message(strings.account_deletion_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::AccountDeletion, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -502,13 +547,13 @@ pub mod repo { &[("handle", &prefs.handle), ("token", token)], ); let subject = format_message(strings.plc_operation_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::PlcOperation, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -533,13 +578,13 @@ pub mod repo { &[("handle", &prefs.handle), ("url", recovery_url)], ); let subject = format_message(strings.passkey_recovery_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::PasskeyRecovery, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -663,13 +708,13 @@ pub mod repo { &[("handle", &prefs.handle), ("code", code)], ); let subject = format_message(strings.two_factor_code_subject, &[("hostname", hostname)]); - let channel = prefs.preferred_channel; + let resolved = resolve_recipient(&prefs, prefs.preferred_channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::TwoFactorCode, - &prefs.email.unwrap_or_default(), + &resolved.recipient, Some(&subject), &body, None, @@ -703,12 +748,56 @@ pub mod repo { ], ); let subject = format_message(strings.legacy_login_subject, &[("hostname", hostname)]); + let resolved = resolve_recipient(&prefs, channel); infra_repo .enqueue_comms( Some(user_id), - channel, + resolved.channel, CommsType::LegacyLoginAlert, - &prefs.email.unwrap_or_default(), + &resolved.recipient, + Some(&subject), + &body, + None, + ) + .await + } + + pub async fn enqueue_channel_verified( + user_repo: &dyn UserRepository, + infra_repo: &dyn InfraRepository, + user_id: Uuid, + channel_name: &str, + recipient: &str, + hostname: &str, + ) -> Result { + let prefs = user_repo + .get_comms_prefs(user_id) + .await? + .ok_or(DbError::NotFound)?; + let strings = get_strings(prefs.preferred_locale.as_deref().unwrap_or("en")); + let display_name = match channel_name { + "email" => "Email", + "discord" => "Discord", + "telegram" => "Telegram", + "signal" => "Signal", + other => other, + }; + let body = format_message( + strings.channel_verified_body, + &[ + ("handle", &prefs.handle), + ("channel", display_name), + ("hostname", hostname), + ], + ); + let subject = format_message(strings.channel_verified_subject, &[("hostname", hostname)]); + let comms_channel = channel_from_str(channel_name); + infra_repo + .enqueue_comms( + Some(user_id), + comms_channel, + CommsType::ChannelVerified, + recipient, Some(&subject), &body, None, diff --git a/crates/tranquil-pds/src/lib.rs b/crates/tranquil-pds/src/lib.rs index fcb9f1e..f4b354e 100644 --- a/crates/tranquil-pds/src/lib.rs +++ b/crates/tranquil-pds/src/lib.rs @@ -281,6 +281,10 @@ pub fn app(state: AppState) -> Router { "/_checkEmailVerified", post(api::server::check_email_verified), ) + .route( + "/_checkChannelVerified", + post(api::server::check_channel_verified), + ) .route( "/com.atproto.server.confirmEmail", post(api::server::confirm_email), @@ -639,6 +643,10 @@ pub fn app(state: AppState) -> Router { .route("/robots.txt", get(api::server::robots_txt)) .route("/logo", get(api::server::get_logo)) .route("/u/{handle}/did.json", get(api::identity::user_did_doc)) + .route( + "/webhook/telegram", + post(api::telegram_webhook::handle_telegram_webhook), + ) .layer(DefaultBodyLimit::max(util::get_max_blob_size())) .layer(middleware::from_fn(metrics::metrics_middleware)) .layer( diff --git a/crates/tranquil-pds/src/main.rs b/crates/tranquil-pds/src/main.rs index 20b28a2..5a97615 100644 --- a/crates/tranquil-pds/src/main.rs +++ b/crates/tranquil-pds/src/main.rs @@ -78,7 +78,34 @@ async fn run() -> Result<(), Box> { } if let Some(telegram_sender) = TelegramSender::from_env() { + let secret_token = match std::env::var("TELEGRAM_WEBHOOK_SECRET") { + Ok(s) => s, + Err(_) => { + return Err( + "TELEGRAM_BOT_TOKEN is set but TELEGRAM_WEBHOOK_SECRET is missing. Both are required for secure Telegram integration.".into() + ); + } + }; info!("Telegram comms enabled"); + match telegram_sender.resolve_bot_username().await { + Ok(username) => { + info!(bot_username = %username, "Resolved Telegram bot username"); + tranquil_pds::util::set_telegram_bot_username(username); + let hostname = + std::env::var("PDS_HOSTNAME").unwrap_or_else(|_| "localhost".to_string()); + let webhook_url = format!("https://{}/webhook/telegram", hostname); + match telegram_sender + .set_webhook(&webhook_url, Some(&secret_token)) + .await + { + Ok(()) => info!(url = %webhook_url, "Telegram webhook registered"), + Err(e) => warn!("Failed to register Telegram webhook: {}", e), + } + } + Err(e) => { + warn!("Failed to resolve Telegram bot username: {}", e); + } + } comms_service = comms_service.register_sender(telegram_sender); } diff --git a/crates/tranquil-pds/src/sso/endpoints.rs b/crates/tranquil-pds/src/sso/endpoints.rs index 0fcb7ed..0f636b0 100644 --- a/crates/tranquil-pds/src/sso/endpoints.rs +++ b/crates/tranquil-pds/src/sso/endpoints.rs @@ -119,8 +119,8 @@ pub async fn sso_initiate( let auth_header = headers .get(axum::http::header::AUTHORIZATION) .and_then(|v| v.to_str().ok()); - let extracted = extract_auth_token_from_header(auth_header) - .ok_or(ApiError::SsoNotAuthenticated)?; + let extracted = + extract_auth_token_from_header(auth_header).ok_or(ApiError::SsoNotAuthenticated)?; let auth_user = validate_bearer_token_cached( state.user_repo.as_ref(), state.cache.as_ref(), @@ -899,7 +899,15 @@ pub async fn complete_registration( _ => return Err(ApiError::MissingDiscordId), }, "telegram" => match &input.telegram_username { - Some(username) if !username.trim().is_empty() => username.trim().to_string(), + Some(username) if !username.trim().is_empty() => { + let clean = username.trim().trim_start_matches('@'); + if !crate::api::validation::is_valid_telegram_username(clean) { + return Err(ApiError::InvalidRequest( + "Invalid Telegram username. Must be 5-32 characters, alphanumeric or underscore".into(), + )); + } + clean.to_string() + } _ => return Err(ApiError::MissingTelegramUsername), }, "signal" => match &input.signal_number { @@ -1104,7 +1112,7 @@ pub async fn complete_registration( telegram_username: input .telegram_username .clone() - .map(|s| s.trim().to_string()) + .map(|s| s.trim().trim_start_matches('@').to_string()) .filter(|s| !s.is_empty()), signal_number: input .signal_number diff --git a/crates/tranquil-pds/src/util.rs b/crates/tranquil-pds/src/util.rs index 058bfa0..463609a 100644 --- a/crates/tranquil-pds/src/util.rs +++ b/crates/tranquil-pds/src/util.rs @@ -14,6 +14,7 @@ const DEFAULT_MAX_BLOB_SIZE: usize = 10 * 1024 * 1024 * 1024; static MAX_BLOB_SIZE: OnceLock = OnceLock::new(); static PDS_HOSTNAME: OnceLock = OnceLock::new(); static PDS_HOSTNAME_WITHOUT_PORT: OnceLock = OnceLock::new(); +static TELEGRAM_BOT_USERNAME: OnceLock = OnceLock::new(); pub fn get_max_blob_size() -> usize { *MAX_BLOB_SIZE.get_or_init(|| { @@ -104,6 +105,14 @@ pub fn pds_hostname_without_port() -> &'static str { }) } +pub fn set_telegram_bot_username(username: String) { + TELEGRAM_BOT_USERNAME.set(username).ok(); +} + +pub fn telegram_bot_username() -> Option<&'static str> { + TELEGRAM_BOT_USERNAME.get().map(|s| s.as_str()) +} + pub fn pds_public_url() -> String { format!("https://{}", pds_hostname()) } diff --git a/crates/tranquil-pds/tests/helpers/mod.rs b/crates/tranquil-pds/tests/helpers/mod.rs index 8ae6cfc..522663a 100644 --- a/crates/tranquil-pds/tests/helpers/mod.rs +++ b/crates/tranquil-pds/tests/helpers/mod.rs @@ -4,6 +4,87 @@ use serde_json::{Value, json}; pub use crate::common::*; +#[allow(dead_code)] +pub async fn paginate_records( + client: &reqwest::Client, + base: &str, + jwt: &str, + did: &str, + collection: &str, + limit: usize, +) -> Vec { + paginate_records_inner(client, base, jwt, did, collection, limit, None, Vec::new()).await +} + +#[allow(clippy::too_many_arguments)] +async fn paginate_records_inner( + client: &reqwest::Client, + base: &str, + jwt: &str, + did: &str, + collection: &str, + limit: usize, + cursor: Option, + mut acc: Vec, +) -> Vec { + let limit_str = limit.to_string(); + let mut query: Vec<(&str, &str)> = vec![ + ("repo", did), + ("collection", collection), + ("limit", &limit_str), + ]; + if let Some(ref c) = cursor { + query.push(("cursor", c.as_str())); + } + + let res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(jwt) + .query(&query) + .send() + .await; + + let Ok(response) = res else { return acc }; + let Ok(body) = response.json::().await else { + return acc; + }; + + let Some(records) = body["records"].as_array() else { + return acc; + }; + acc.extend(records.iter().cloned()); + + match body["cursor"].as_str() { + Some(next) => { + Box::pin(paginate_records_inner( + client, + base, + jwt, + did, + collection, + limit, + Some(next.to_string()), + acc, + )) + .await + } + None => acc, + } +} + +#[allow(dead_code)] +pub async fn count_records( + client: &reqwest::Client, + base: &str, + jwt: &str, + did: &str, + collection: &str, +) -> usize { + paginate_records(client, base, jwt, did, collection, 100) + .await + .len() +} + fn unique_id() -> String { uuid::Uuid::new_v4().simple().to_string()[..12].to_string() } @@ -246,3 +327,150 @@ pub async fn set_account_deactivated(did: &str, deactivated: bool) { .await .expect("Failed to update deactivated_at"); } + +#[allow(dead_code)] +pub fn make_cid(data: &[u8]) -> cid::Cid { + use sha2::{Digest, Sha256}; + let hash = Sha256::digest(data); + let multihash = multihash::Multihash::wrap(0x12, &hash).unwrap(); + cid::Cid::new_v1(0x71, multihash) +} + +#[allow(dead_code)] +pub fn write_varint(buf: &mut Vec, value: u64) { + buf.extend(encode_varint_bytes(value)); +} + +fn encode_varint_bytes(value: u64) -> Vec { + match value < 0x80 { + true => vec![value as u8], + false => { + let mut rest = encode_varint_bytes(value >> 7); + rest.insert(0, ((value & 0x7F) as u8) | 0x80); + rest + } + } +} + +#[allow(dead_code)] +pub fn encode_car_block(cid: &cid::Cid, data: &[u8]) -> Vec { + let cid_bytes = cid.to_bytes(); + let mut result = Vec::new(); + write_varint(&mut result, (cid_bytes.len() + data.len()) as u64); + result.extend_from_slice(&cid_bytes); + result.extend_from_slice(data); + result +} + +#[allow(dead_code)] +pub fn create_test_record() -> (Vec, cid::Cid) { + use ipld_core::ipld::Ipld; + use std::collections::BTreeMap; + let record = Ipld::Map(BTreeMap::from([ + ( + "$type".to_string(), + Ipld::String("app.bsky.feed.post".to_string()), + ), + ( + "text".to_string(), + Ipld::String("Test post for verification".to_string()), + ), + ( + "createdAt".to_string(), + Ipld::String("2024-01-01T00:00:00Z".to_string()), + ), + ])); + let bytes = serde_ipld_dagcbor::to_vec(&record).unwrap(); + let cid = make_cid(&bytes); + (bytes, cid) +} + +#[allow(dead_code)] +pub fn create_mst_node(entries: Vec<(String, cid::Cid)>) -> (Vec, cid::Cid) { + use ipld_core::ipld::Ipld; + use std::collections::BTreeMap; + let ipld_entries: Vec = entries + .into_iter() + .map(|(key, value_cid)| { + Ipld::Map(BTreeMap::from([ + ("k".to_string(), Ipld::Bytes(key.into_bytes())), + ("v".to_string(), Ipld::Link(value_cid)), + ("p".to_string(), Ipld::Integer(0)), + ])) + }) + .collect(); + let node = Ipld::Map(BTreeMap::from([( + "e".to_string(), + Ipld::List(ipld_entries), + )])); + let bytes = serde_ipld_dagcbor::to_vec(&node).unwrap(); + let cid = make_cid(&bytes); + (bytes, cid) +} + +#[allow(dead_code)] +pub fn create_car_signed_commit( + did: &str, + data_cid: &cid::Cid, + signing_key: &k256::ecdsa::SigningKey, +) -> (Vec, cid::Cid) { + use jacquard_common::types::{integer::LimitedU32, string::Tid}; + use jacquard_repo::commit::Commit; + let rev = Tid::now(LimitedU32::MIN); + let did = jacquard_common::types::string::Did::new(did).expect("valid DID"); + let unsigned = Commit::new_unsigned(did, *data_cid, rev, None); + let signed = unsigned.sign(signing_key).expect("signing failed"); + let signed_bytes = signed.to_cbor().expect("serialization failed"); + let cid = make_cid(&signed_bytes); + (signed_bytes, cid) +} + +#[allow(dead_code)] +pub fn build_car_with_signature( + did: &str, + signing_key: &k256::ecdsa::SigningKey, +) -> (Vec, cid::Cid) { + let (record_bytes, record_cid) = create_test_record(); + let (mst_bytes, mst_cid) = + create_mst_node(vec![("app.bsky.feed.post/test123".to_string(), record_cid)]); + let (commit_bytes, commit_cid) = create_car_signed_commit(did, &mst_cid, signing_key); + let header = iroh_car::CarHeader::new_v1(vec![commit_cid]); + let header_bytes = header.encode().unwrap(); + let mut car = Vec::new(); + write_varint(&mut car, header_bytes.len() as u64); + car.extend_from_slice(&header_bytes); + car.extend(encode_car_block(&commit_cid, &commit_bytes)); + car.extend(encode_car_block(&mst_cid, &mst_bytes)); + car.extend(encode_car_block(&record_cid, &record_bytes)); + (car, commit_cid) +} + +#[allow(dead_code)] +pub fn get_multikey_from_signing_key(signing_key: &k256::ecdsa::SigningKey) -> String { + let public_key = signing_key.verifying_key(); + let compressed = public_key.to_sec1_bytes(); + let buf: Vec = encode_varint_bytes(0xE7) + .into_iter() + .chain(compressed.iter().copied()) + .collect(); + multibase::encode(multibase::Base::Base58Btc, buf) +} + +#[allow(dead_code)] +pub async fn get_user_signing_key(did: &str) -> Option> { + let db_url = get_db_connection_string().await; + let pool = sqlx::PgPool::connect(&db_url).await.ok()?; + let row = sqlx::query!( + r#" + SELECT k.key_bytes, k.encryption_version + FROM user_keys k + JOIN users u ON k.user_id = u.id + WHERE u.did = $1 + "#, + did + ) + .fetch_optional(&pool) + .await + .ok()??; + tranquil_pds::config::decrypt_key(&row.key_bytes, row.encryption_version).ok() +} diff --git a/crates/tranquil-pds/tests/import_with_verification.rs b/crates/tranquil-pds/tests/import_with_verification.rs index 18dab9b..97d4aca 100644 --- a/crates/tranquil-pds/tests/import_with_verification.rs +++ b/crates/tranquil-pds/tests/import_with_verification.rs @@ -1,66 +1,13 @@ mod common; -use cid::Cid; +mod helpers; use common::*; -use ipld_core::ipld::Ipld; -use jacquard_common::types::{integer::LimitedU32, string::Tid}; -use jacquard_repo::commit::Commit; +use helpers::*; use k256::ecdsa::SigningKey; use reqwest::StatusCode; use serde_json::json; -use sha2::{Digest, Sha256}; -use sqlx::PgPool; -use std::collections::BTreeMap; use wiremock::matchers::{method, path}; use wiremock::{Mock, MockServer, ResponseTemplate}; -fn make_cid(data: &[u8]) -> Cid { - let mut hasher = Sha256::new(); - hasher.update(data); - let hash = hasher.finalize(); - let multihash = multihash::Multihash::wrap(0x12, &hash).unwrap(); - Cid::new_v1(0x71, multihash) -} - -fn write_varint(buf: &mut Vec, mut value: u64) { - loop { - let mut byte = (value & 0x7F) as u8; - value >>= 7; - if value != 0 { - byte |= 0x80; - } - buf.push(byte); - if value == 0 { - break; - } - } -} - -fn encode_car_block(cid: &Cid, data: &[u8]) -> Vec { - let cid_bytes = cid.to_bytes(); - let mut result = Vec::new(); - write_varint(&mut result, (cid_bytes.len() + data.len()) as u64); - result.extend_from_slice(&cid_bytes); - result.extend_from_slice(data); - result -} - -fn get_multikey_from_signing_key(signing_key: &SigningKey) -> String { - let public_key = signing_key.verifying_key(); - let compressed = public_key.to_sec1_bytes(); - fn encode_uvarint(mut x: u64) -> Vec { - let mut out = Vec::new(); - while x >= 0x80 { - out.push(((x as u8) & 0x7F) | 0x80); - x >>= 7; - } - out.push(x as u8); - out - } - let mut buf = encode_uvarint(0xE7); - buf.extend_from_slice(&compressed); - multibase::encode(multibase::Base::Base58Btc, buf) -} - fn create_did_document( did: &str, handle: &str, @@ -89,70 +36,6 @@ fn create_did_document( }) } -fn create_signed_commit(did: &str, data_cid: &Cid, signing_key: &SigningKey) -> (Vec, Cid) { - let rev = Tid::now(LimitedU32::MIN); - let did = jacquard_common::types::string::Did::new(did).expect("valid DID"); - let unsigned = Commit::new_unsigned(did, *data_cid, rev, None); - let signed = unsigned.sign(signing_key).expect("signing failed"); - let signed_bytes = signed.to_cbor().expect("serialization failed"); - let cid = make_cid(&signed_bytes); - (signed_bytes, cid) -} - -fn create_mst_node(entries: Vec<(String, Cid)>) -> (Vec, Cid) { - let ipld_entries: Vec = entries - .into_iter() - .map(|(key, value_cid)| { - Ipld::Map(BTreeMap::from([ - ("k".to_string(), Ipld::Bytes(key.into_bytes())), - ("v".to_string(), Ipld::Link(value_cid)), - ("p".to_string(), Ipld::Integer(0)), - ])) - }) - .collect(); - let node = Ipld::Map(BTreeMap::from([( - "e".to_string(), - Ipld::List(ipld_entries), - )])); - let bytes = serde_ipld_dagcbor::to_vec(&node).unwrap(); - let cid = make_cid(&bytes); - (bytes, cid) -} - -fn create_record() -> (Vec, Cid) { - let record = Ipld::Map(BTreeMap::from([ - ( - "$type".to_string(), - Ipld::String("app.bsky.feed.post".to_string()), - ), - ( - "text".to_string(), - Ipld::String("Test post for verification".to_string()), - ), - ( - "createdAt".to_string(), - Ipld::String("2024-01-01T00:00:00Z".to_string()), - ), - ])); - let bytes = serde_ipld_dagcbor::to_vec(&record).unwrap(); - let cid = make_cid(&bytes); - (bytes, cid) -} -fn build_car_with_signature(did: &str, signing_key: &SigningKey) -> (Vec, Cid) { - let (record_bytes, record_cid) = create_record(); - let (mst_bytes, mst_cid) = - create_mst_node(vec![("app.bsky.feed.post/test123".to_string(), record_cid)]); - let (commit_bytes, commit_cid) = create_signed_commit(did, &mst_cid, signing_key); - let header = iroh_car::CarHeader::new_v1(vec![commit_cid]); - let header_bytes = header.encode().unwrap(); - let mut car = Vec::new(); - write_varint(&mut car, header_bytes.len() as u64); - car.extend_from_slice(&header_bytes); - car.extend(encode_car_block(&commit_cid, &commit_bytes)); - car.extend(encode_car_block(&mst_cid, &mst_bytes)); - car.extend(encode_car_block(&record_cid, &record_bytes)); - (car, commit_cid) -} async fn setup_mock_plc_directory(did: &str, did_doc: serde_json::Value) -> MockServer { let mock_server = MockServer::start().await; let did_encoded = urlencoding::encode(did); @@ -164,23 +47,7 @@ async fn setup_mock_plc_directory(did: &str, did_doc: serde_json::Value) -> Mock .await; mock_server } -async fn get_user_signing_key(did: &str) -> Option> { - let db_url = get_db_connection_string().await; - let pool = PgPool::connect(&db_url).await.ok()?; - let row = sqlx::query!( - r#" - SELECT k.key_bytes, k.encryption_version - FROM user_keys k - JOIN users u ON k.user_id = u.id - WHERE u.did = $1 - "#, - did - ) - .fetch_optional(&pool) - .await - .ok()??; - tranquil_pds::config::decrypt_key(&row.key_bytes, row.encryption_version).ok() -} + #[tokio::test] #[ignore = "requires exclusive env var access; run with: cargo test test_import_with_valid_signature_and_mock_plc -- --ignored --test-threads=1"] async fn test_import_with_valid_signature_and_mock_plc() { diff --git a/crates/tranquil-pds/tests/whole_story.rs b/crates/tranquil-pds/tests/whole_story.rs new file mode 100644 index 0000000..5c0449c --- /dev/null +++ b/crates/tranquil-pds/tests/whole_story.rs @@ -0,0 +1,2013 @@ +mod common; +mod helpers; + +use chrono::Utc; +use common::*; +use futures::{StreamExt, future::join_all}; +use helpers::*; +use k256::ecdsa::SigningKey; +use reqwest::{StatusCode, header}; +use serde_json::{Value, json}; + +#[tokio::test] +async fn test_complete_user_journey_signup_to_deletion() { + let client = client(); + let base = base_url().await; + let uid = uuid::Uuid::new_v4().simple().to_string(); + let handle = format!("journey{}", &uid[..8]); + let email = format!("journey{}@test.com", &uid[..8]); + let password = "JourneyPass123!"; + + let create_res = client + .post(format!("{}/xrpc/com.atproto.server.createAccount", base)) + .json(&json!({ + "handle": handle, + "email": email, + "password": password + })) + .send() + .await + .expect("Account creation failed"); + assert_eq!(create_res.status(), StatusCode::OK); + let account: Value = create_res.json().await.unwrap(); + let did = account["did"].as_str().unwrap().to_string(); + + let jwt = verify_new_account(&client, &did).await; + + let blob_data = b"This is my avatar image data for the complete journey test"; + let upload_res = client + .post(format!("{}/xrpc/com.atproto.repo.uploadBlob", base)) + .header(header::CONTENT_TYPE, "image/png") + .bearer_auth(&jwt) + .body(blob_data.to_vec()) + .send() + .await + .expect("Blob upload failed"); + assert_eq!(upload_res.status(), StatusCode::OK); + let upload_body: Value = upload_res.json().await.unwrap(); + let avatar_blob = upload_body["blob"].clone(); + + let profile_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.actor.profile", + "rkey": "self", + "record": { + "$type": "app.bsky.actor.profile", + "displayName": "Journey Test User", + "description": "Testing the complete user journey", + "avatar": avatar_blob + } + })) + .send() + .await + .expect("Profile creation failed"); + assert_eq!(profile_res.status(), StatusCode::OK); + + let post1_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "My first post on this journey!", + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("First post failed"); + assert_eq!(post1_res.status(), StatusCode::OK); + let post1_body: Value = post1_res.json().await.unwrap(); + let post1_uri = post1_body["uri"].as_str().unwrap(); + let post1_cid = post1_body["cid"].as_str().unwrap(); + + let post2_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Second post in my journey", + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Second post failed"); + assert_eq!(post2_res.status(), StatusCode::OK); + + let list_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .expect("List records failed"); + assert_eq!(list_res.status(), StatusCode::OK); + let list_body: Value = list_res.json().await.unwrap(); + assert_eq!(list_body["records"].as_array().unwrap().len(), 2); + + let post1_rkey = post1_uri.split('/').next_back().unwrap(); + let edit_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": post1_rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "My first post on this journey! (edited)", + "createdAt": Utc::now().to_rfc3339() + }, + "swapRecord": post1_cid + })) + .send() + .await + .expect("Edit post failed"); + assert_eq!(edit_res.status(), StatusCode::OK); + + let backup_res = client + .post(format!("{}/xrpc/_backup.createBackup", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("Backup creation failed"); + assert_eq!(backup_res.status(), StatusCode::OK); + let backup_body: Value = backup_res.json().await.unwrap(); + let backup_id = backup_body["id"].as_str().unwrap(); + + let download_res = client + .get(format!("{}/xrpc/_backup.getBackup?id={}", base, backup_id)) + .bearer_auth(&jwt) + .send() + .await + .expect("Backup download failed"); + assert_eq!(download_res.status(), StatusCode::OK); + let backup_bytes = download_res.bytes().await.unwrap(); + assert!(backup_bytes.len() > 100, "Backup should have content"); + + let delete_res = client + .post(format!("{}/xrpc/com.atproto.server.deleteSession", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("Logout failed"); + assert!(delete_res.status() == StatusCode::OK || delete_res.status() == StatusCode::NO_CONTENT); + + let login_res = client + .post(format!("{}/xrpc/com.atproto.server.createSession", base)) + .json(&json!({ + "identifier": handle, + "password": password + })) + .send() + .await + .expect("Re-login failed"); + assert_eq!(login_res.status(), StatusCode::OK); + let login_body: Value = login_res.json().await.unwrap(); + let new_jwt = login_body["accessJwt"].as_str().unwrap(); + + let verify_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(new_jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .expect("Verify after re-login failed"); + assert_eq!(verify_res.status(), StatusCode::OK); + let verify_body: Value = verify_res.json().await.unwrap(); + assert_eq!(verify_body["records"].as_array().unwrap().len(), 2); + + let request_delete_res = client + .post(format!( + "{}/xrpc/com.atproto.server.requestAccountDelete", + base + )) + .bearer_auth(new_jwt) + .send() + .await + .expect("Request delete failed"); + assert_eq!(request_delete_res.status(), StatusCode::OK); + + let pool = get_test_db_pool().await; + let row = sqlx::query!( + "SELECT token FROM account_deletion_requests WHERE did = $1", + did + ) + .fetch_one(pool) + .await + .expect("Failed to get deletion token"); + + let final_delete_res = client + .post(format!("{}/xrpc/com.atproto.server.deleteAccount", base)) + .json(&json!({ + "did": did, + "password": password, + "token": row.token + })) + .send() + .await + .expect("Final delete failed"); + assert_eq!(final_delete_res.status(), StatusCode::OK); + + let user_gone = sqlx::query!("SELECT id FROM users WHERE did = $1", did) + .fetch_optional(pool) + .await + .expect("Failed to check user"); + assert!(user_gone.is_none(), "User should be deleted"); +} + +#[tokio::test] +async fn test_multi_user_social_graph_lifecycle() { + let client = client(); + let base = base_url().await; + + let (alice_did, alice_jwt) = setup_new_user("alice-social").await; + let (bob_did, bob_jwt) = setup_new_user("bob-social").await; + let (carol_did, carol_jwt) = setup_new_user("carol-social").await; + + let (alice_post_uri, alice_post_cid) = + create_post(&client, &alice_did, &alice_jwt, "Hello from Alice!").await; + + let (bob_post_uri, bob_post_cid) = + create_post(&client, &bob_did, &bob_jwt, "Hello from Bob!").await; + + let (_bob_follows_alice_uri, _) = create_follow(&client, &bob_did, &bob_jwt, &alice_did).await; + let (_carol_follows_alice_uri, _) = + create_follow(&client, &carol_did, &carol_jwt, &alice_did).await; + let (_carol_follows_bob_uri, _) = + create_follow(&client, &carol_did, &carol_jwt, &bob_did).await; + + let (_bob_likes_alice_uri, _) = create_like( + &client, + &bob_did, + &bob_jwt, + &alice_post_uri, + &alice_post_cid, + ) + .await; + let (_carol_likes_alice_uri, _) = create_like( + &client, + &carol_did, + &carol_jwt, + &alice_post_uri, + &alice_post_cid, + ) + .await; + + let (_bob_reposts_alice_uri, _) = create_repost( + &client, + &bob_did, + &bob_jwt, + &alice_post_uri, + &alice_post_cid, + ) + .await; + + let reply_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&carol_jwt) + .json(&json!({ + "repo": carol_did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Great post Alice!", + "createdAt": Utc::now().to_rfc3339(), + "reply": { + "root": { "uri": alice_post_uri, "cid": alice_post_cid }, + "parent": { "uri": alice_post_uri, "cid": alice_post_cid } + } + } + })) + .send() + .await + .expect("Reply failed"); + assert_eq!(reply_res.status(), StatusCode::OK); + let reply_body: Value = reply_res.json().await.unwrap(); + let carol_reply_uri = reply_body["uri"].as_str().unwrap(); + let carol_reply_cid = reply_body["cid"].as_str().unwrap(); + + let alice_reply_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&alice_jwt) + .json(&json!({ + "repo": alice_did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Thanks Carol!", + "createdAt": Utc::now().to_rfc3339(), + "reply": { + "root": { "uri": alice_post_uri, "cid": alice_post_cid }, + "parent": { "uri": carol_reply_uri, "cid": carol_reply_cid } + } + } + })) + .send() + .await + .expect("Alice reply failed"); + assert_eq!(alice_reply_res.status(), StatusCode::OK); + + let alice_follows_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", alice_did.as_str()), + ("collection", "app.bsky.graph.follow"), + ]) + .send() + .await + .unwrap(); + let alice_follows: Value = alice_follows_res.json().await.unwrap(); + assert_eq!( + alice_follows["records"].as_array().unwrap().len(), + 0, + "Alice follows nobody" + ); + + let bob_follows_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", bob_did.as_str()), + ("collection", "app.bsky.graph.follow"), + ]) + .send() + .await + .unwrap(); + let bob_follows: Value = bob_follows_res.json().await.unwrap(); + assert_eq!( + bob_follows["records"].as_array().unwrap().len(), + 1, + "Bob follows 1 person" + ); + + let carol_follows_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", carol_did.as_str()), + ("collection", "app.bsky.graph.follow"), + ]) + .send() + .await + .unwrap(); + let carol_follows: Value = carol_follows_res.json().await.unwrap(); + assert_eq!( + carol_follows["records"].as_array().unwrap().len(), + 2, + "Carol follows 2 people" + ); + + let bob_likes_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", bob_did.as_str()), + ("collection", "app.bsky.feed.like"), + ]) + .send() + .await + .unwrap(); + let bob_likes: Value = bob_likes_res.json().await.unwrap(); + assert_eq!(bob_likes["records"].as_array().unwrap().len(), 1); + + let alice_posts_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", alice_did.as_str()), + ("collection", "app.bsky.feed.post"), + ]) + .send() + .await + .unwrap(); + let alice_posts: Value = alice_posts_res.json().await.unwrap(); + assert_eq!( + alice_posts["records"].as_array().unwrap().len(), + 2, + "Alice has 2 posts (original + reply)" + ); + + let bob_likes_rkey = bob_likes["records"][0]["uri"] + .as_str() + .unwrap() + .split('/') + .next_back() + .unwrap(); + let unlike_res = client + .post(format!("{}/xrpc/com.atproto.repo.deleteRecord", base)) + .bearer_auth(&bob_jwt) + .json(&json!({ + "repo": bob_did, + "collection": "app.bsky.feed.like", + "rkey": bob_likes_rkey + })) + .send() + .await + .expect("Unlike failed"); + assert_eq!(unlike_res.status(), StatusCode::OK); + + let bob_likes_after = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", bob_did.as_str()), + ("collection", "app.bsky.feed.like"), + ]) + .send() + .await + .unwrap(); + let bob_likes_after_body: Value = bob_likes_after.json().await.unwrap(); + assert_eq!(bob_likes_after_body["records"].as_array().unwrap().len(), 0); + + let relike_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&bob_jwt) + .json(&json!({ + "repo": bob_did, + "collection": "app.bsky.feed.like", + "record": { + "$type": "app.bsky.feed.like", + "subject": { "uri": alice_post_uri, "cid": alice_post_cid }, + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Relike failed"); + assert_eq!(relike_res.status(), StatusCode::OK); + + let (_alice_likes_bob_uri, _) = create_like( + &client, + &alice_did, + &alice_jwt, + &bob_post_uri, + &bob_post_cid, + ) + .await; + + let mutual_follow_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&alice_jwt) + .json(&json!({ + "repo": alice_did, + "collection": "app.bsky.graph.follow", + "record": { + "$type": "app.bsky.graph.follow", + "subject": bob_did, + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Mutual follow failed"); + assert_eq!(mutual_follow_res.status(), StatusCode::OK); + + let alice_final_follows = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", alice_did.as_str()), + ("collection", "app.bsky.graph.follow"), + ]) + .send() + .await + .unwrap(); + let alice_final: Value = alice_final_follows.json().await.unwrap(); + assert_eq!(alice_final["records"].as_array().unwrap().len(), 1); +} + +#[tokio::test] +async fn test_blob_lifecycle_upload_use_remove() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("blob-lifecycle").await; + + let blob1_data = b"First blob for testing lifecycle"; + let upload1_res = client + .post(format!("{}/xrpc/com.atproto.repo.uploadBlob", base)) + .header(header::CONTENT_TYPE, "text/plain") + .bearer_auth(&jwt) + .body(blob1_data.to_vec()) + .send() + .await + .expect("Upload 1 failed"); + assert_eq!(upload1_res.status(), StatusCode::OK); + let upload1_body: Value = upload1_res.json().await.unwrap(); + let blob1 = upload1_body["blob"].clone(); + let blob1_cid = blob1["ref"]["$link"].as_str().unwrap(); + + let blob2_data = b"Second blob for testing lifecycle"; + let upload2_res = client + .post(format!("{}/xrpc/com.atproto.repo.uploadBlob", base)) + .header(header::CONTENT_TYPE, "text/plain") + .bearer_auth(&jwt) + .body(blob2_data.to_vec()) + .send() + .await + .expect("Upload 2 failed"); + assert_eq!(upload2_res.status(), StatusCode::OK); + let upload2_body: Value = upload2_res.json().await.unwrap(); + let blob2 = upload2_body["blob"].clone(); + let _blob2_cid = blob2["ref"]["$link"].as_str().unwrap(); + + let post_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Post with images", + "createdAt": Utc::now().to_rfc3339(), + "embed": { + "$type": "app.bsky.embed.images", + "images": [ + { "alt": "First image", "image": blob1 }, + { "alt": "Second image", "image": blob2 } + ] + } + } + })) + .send() + .await + .expect("Post with blobs failed"); + assert_eq!(post_res.status(), StatusCode::OK); + let post_body: Value = post_res.json().await.unwrap(); + let post_uri = post_body["uri"].as_str().unwrap(); + let post_cid = post_body["cid"].as_str().unwrap(); + let post_rkey = post_uri.split('/').next_back().unwrap(); + + let get_blob1 = client + .get(format!("{}/xrpc/com.atproto.sync.getBlob", base)) + .query(&[("did", did.as_str()), ("cid", blob1_cid)]) + .send() + .await + .expect("Get blob 1 failed"); + assert_eq!(get_blob1.status(), StatusCode::OK); + let blob1_content = get_blob1.bytes().await.unwrap(); + assert_eq!(blob1_content.as_ref(), blob1_data); + + let get_record = client + .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", post_rkey), + ]) + .send() + .await + .expect("Get record failed"); + let record_body: Value = get_record.json().await.unwrap(); + let images = record_body["value"]["embed"]["images"].as_array().unwrap(); + assert_eq!(images.len(), 2); + + let edit_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": post_rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "Post with single image now", + "createdAt": Utc::now().to_rfc3339(), + "embed": { + "$type": "app.bsky.embed.images", + "images": [ + { "alt": "Only first image", "image": blob1 } + ] + } + }, + "swapRecord": post_cid + })) + .send() + .await + .expect("Edit failed"); + assert_eq!(edit_res.status(), StatusCode::OK); + + let get_edited = client + .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", post_rkey), + ]) + .send() + .await + .unwrap(); + let edited_body: Value = get_edited.json().await.unwrap(); + let edited_images = edited_body["value"]["embed"]["images"].as_array().unwrap(); + assert_eq!(edited_images.len(), 1); + + let profile_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.actor.profile", + "rkey": "self", + "record": { + "$type": "app.bsky.actor.profile", + "displayName": "Blob Test User", + "avatar": blob1 + } + })) + .send() + .await + .expect("Profile failed"); + assert_eq!(profile_res.status(), StatusCode::OK); + + let list_blobs = client + .get(format!("{}/xrpc/com.atproto.sync.listBlobs", base)) + .query(&[("did", did.as_str())]) + .send() + .await + .expect("List blobs failed"); + assert_eq!(list_blobs.status(), StatusCode::OK); + let blobs_body: Value = list_blobs.json().await.unwrap(); + let cids = blobs_body["cids"].as_array().unwrap(); + assert!( + cids.iter().any(|c| c.as_str() == Some(blob1_cid)), + "blob1 should still exist (referenced by profile and post)" + ); +} + +#[tokio::test] +async fn test_session_and_record_interaction() { + let client = client(); + let base = base_url().await; + let uid = uuid::Uuid::new_v4().simple().to_string(); + let handle = format!("sess{}", &uid[..8]); + let email = format!("sess{}@test.com", &uid[..8]); + let password = "SessionTest123!"; + + let create_res = client + .post(format!("{}/xrpc/com.atproto.server.createAccount", base)) + .json(&json!({ + "handle": handle, + "email": email, + "password": password + })) + .send() + .await + .expect("Account creation failed"); + assert_eq!(create_res.status(), StatusCode::OK); + let account: Value = create_res.json().await.unwrap(); + let did = account["did"].as_str().unwrap().to_string(); + let jwt = verify_new_account(&client, &did).await; + + let (post1_uri, _) = create_post(&client, &did, &jwt, "Post from session 1").await; + + let session2_res = client + .post(format!("{}/xrpc/com.atproto.server.createSession", base)) + .json(&json!({ + "identifier": handle, + "password": password + })) + .send() + .await + .expect("Session 2 creation failed"); + assert_eq!(session2_res.status(), StatusCode::OK); + let session2: Value = session2_res.json().await.unwrap(); + let jwt2 = session2["accessJwt"].as_str().unwrap(); + let refresh2 = session2["refreshJwt"].as_str().unwrap(); + + let (post2_uri, _) = create_post(&client, &did, jwt2, "Post from session 2").await; + + let list_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let list_body: Value = list_res.json().await.unwrap(); + assert_eq!( + list_body["records"].as_array().unwrap().len(), + 2, + "Both posts visible from session 1" + ); + + let refresh_res = client + .post(format!("{}/xrpc/com.atproto.server.refreshSession", base)) + .bearer_auth(refresh2) + .send() + .await + .expect("Refresh failed"); + assert_eq!(refresh_res.status(), StatusCode::OK); + let refresh_body: Value = refresh_res.json().await.unwrap(); + let new_jwt2 = refresh_body["accessJwt"].as_str().unwrap(); + + let verify_posts = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(new_jwt2) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let verify_body: Value = verify_posts.json().await.unwrap(); + assert_eq!( + verify_body["records"].as_array().unwrap().len(), + 2, + "Posts still visible after token refresh" + ); + + let post1_rkey = post1_uri.split('/').next_back().unwrap(); + let delete_res = client + .post(format!("{}/xrpc/com.atproto.repo.deleteRecord", base)) + .bearer_auth(new_jwt2) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": post1_rkey + })) + .send() + .await + .expect("Delete from session 2 failed"); + assert_eq!(delete_res.status(), StatusCode::OK); + + let final_list = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let final_body: Value = final_list.json().await.unwrap(); + let remaining_posts = final_body["records"].as_array().unwrap(); + assert_eq!(remaining_posts.len(), 1); + assert!( + remaining_posts[0]["uri"] + .as_str() + .unwrap() + .contains(post2_uri.split('/').next_back().unwrap()) + ); +} + +#[tokio::test] +async fn test_app_password_record_lifecycle() { + let client = client(); + let base = base_url().await; + let uid = uuid::Uuid::new_v4().simple().to_string(); + let handle = format!("apprec{}", &uid[..8]); + let email = format!("apprec{}@test.com", &uid[..8]); + let password = "AppRecTest123!"; + + let create_res = client + .post(format!("{}/xrpc/com.atproto.server.createAccount", base)) + .json(&json!({ + "handle": handle, + "email": email, + "password": password + })) + .send() + .await + .expect("Account creation failed"); + assert_eq!(create_res.status(), StatusCode::OK); + let account: Value = create_res.json().await.unwrap(); + let did = account["did"].as_str().unwrap().to_string(); + let main_jwt = verify_new_account(&client, &did).await; + + let create_app_pass = client + .post(format!( + "{}/xrpc/com.atproto.server.createAppPassword", + base + )) + .bearer_auth(&main_jwt) + .json(&json!({ "name": "Test App" })) + .send() + .await + .expect("App password creation failed"); + assert_eq!(create_app_pass.status(), StatusCode::OK); + let app_pass_body: Value = create_app_pass.json().await.unwrap(); + let app_password = app_pass_body["password"].as_str().unwrap(); + + let app_login = client + .post(format!("{}/xrpc/com.atproto.server.createSession", base)) + .json(&json!({ + "identifier": handle, + "password": app_password + })) + .send() + .await + .expect("App password login failed"); + assert_eq!(app_login.status(), StatusCode::OK); + let app_session: Value = app_login.json().await.unwrap(); + let app_jwt = app_session["accessJwt"].as_str().unwrap(); + + create_post(&client, &did, app_jwt, "Post from app password session").await; + + let verify_main = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&main_jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let verify_body: Value = verify_main.json().await.unwrap(); + assert_eq!( + verify_body["records"].as_array().unwrap().len(), + 1, + "Post visible from main session" + ); + + let revoke_res = client + .post(format!( + "{}/xrpc/com.atproto.server.revokeAppPassword", + base + )) + .bearer_auth(&main_jwt) + .json(&json!({ "name": "Test App" })) + .send() + .await + .expect("Revoke failed"); + assert_eq!(revoke_res.status(), StatusCode::OK); + + let post_after_revoke = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(app_jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Should fail", + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Post attempt failed"); + assert!( + post_after_revoke.status() == StatusCode::UNAUTHORIZED + || post_after_revoke.status() == StatusCode::BAD_REQUEST, + "Revoked app password should not create posts" + ); + + let final_list = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&main_jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let final_body: Value = final_list.json().await.unwrap(); + assert_eq!( + final_body["records"].as_array().unwrap().len(), + 1, + "Only the valid post should exist" + ); +} + +#[tokio::test] +async fn test_handle_change_with_existing_content() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("handlechange").await; + + let (post_uri, post_cid) = create_post(&client, &did, &jwt, "Post before handle change").await; + + let profile_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.actor.profile", + "rkey": "self", + "record": { + "$type": "app.bsky.actor.profile", + "displayName": "Original Handle User" + } + })) + .send() + .await + .expect("Profile creation failed"); + assert_eq!(profile_res.status(), StatusCode::OK); + + let (other_did, other_jwt) = setup_new_user("other-user").await; + let (_like_uri, _) = create_like(&client, &other_did, &other_jwt, &post_uri, &post_cid).await; + let (_follow_uri, _) = create_follow(&client, &other_did, &other_jwt, &did).await; + + let new_handle = format!("newh{}", &uuid::Uuid::new_v4().simple().to_string()[..8]); + let update_res = client + .post(format!("{}/xrpc/com.atproto.identity.updateHandle", base)) + .bearer_auth(&jwt) + .json(&json!({ "handle": new_handle })) + .send() + .await + .expect("Handle update failed"); + assert_eq!(update_res.status(), StatusCode::OK); + + let session_res = client + .get(format!("{}/xrpc/com.atproto.server.getSession", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("Get session failed"); + let session_body: Value = session_res.json().await.unwrap(); + assert!( + session_body["handle"] + .as_str() + .unwrap() + .starts_with(&new_handle), + "Handle should be updated" + ); + + let posts_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let posts_body: Value = posts_res.json().await.unwrap(); + assert_eq!( + posts_body["records"].as_array().unwrap().len(), + 1, + "Post should still exist after handle change" + ); + + let profile_check = client + .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.actor.profile"), + ("rkey", "self"), + ]) + .send() + .await + .unwrap(); + let profile_body: Value = profile_check.json().await.unwrap(); + assert_eq!( + profile_body["value"]["displayName"], "Original Handle User", + "Profile should be intact" + ); + + let other_follows = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[ + ("repo", other_did.as_str()), + ("collection", "app.bsky.graph.follow"), + ]) + .send() + .await + .unwrap(); + let follows_body: Value = other_follows.json().await.unwrap(); + let follows = follows_body["records"].as_array().unwrap(); + assert_eq!(follows.len(), 1); + assert_eq!( + follows[0]["value"]["subject"], did, + "Follow should still point to the DID" + ); + + let resolve_res = client + .get(format!("{}/xrpc/com.atproto.identity.resolveHandle", base)) + .query(&[("handle", session_body["handle"].as_str().unwrap())]) + .send() + .await + .expect("Resolve failed"); + assert_eq!(resolve_res.status(), StatusCode::OK); + let resolve_body: Value = resolve_res.json().await.unwrap(); + assert_eq!( + resolve_body["did"], did, + "New handle should resolve to same DID" + ); +} + +#[tokio::test] +async fn test_deactivation_preserves_data() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("deactivate-data").await; + + create_post(&client, &did, &jwt, "Post 1 before deactivation").await; + create_post(&client, &did, &jwt, "Post 2 before deactivation").await; + + let profile_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.actor.profile", + "rkey": "self", + "record": { + "$type": "app.bsky.actor.profile", + "displayName": "Deactivation Test User" + } + })) + .send() + .await + .unwrap(); + assert_eq!(profile_res.status(), StatusCode::OK); + + let deactivate_res = client + .post(format!( + "{}/xrpc/com.atproto.server.deactivateAccount", + base + )) + .bearer_auth(&jwt) + .json(&json!({})) + .send() + .await + .expect("Deactivation failed"); + assert_eq!(deactivate_res.status(), StatusCode::OK); + + let posts_while_deactivated = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + assert_eq!(posts_while_deactivated.status(), StatusCode::OK); + let posts_body: Value = posts_while_deactivated.json().await.unwrap(); + assert_eq!( + posts_body["records"].as_array().unwrap().len(), + 2, + "Posts should still be readable while deactivated" + ); + + let activate_res = client + .post(format!("{}/xrpc/com.atproto.server.activateAccount", base)) + .bearer_auth(&jwt) + .json(&json!({})) + .send() + .await + .expect("Activation failed"); + assert_eq!(activate_res.status(), StatusCode::OK); + + create_post(&client, &did, &jwt, "Post 3 after reactivation").await; + + let final_posts = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let final_body: Value = final_posts.json().await.unwrap(); + assert_eq!( + final_body["records"].as_array().unwrap().len(), + 3, + "All three posts should exist" + ); +} + +#[tokio::test] +async fn test_password_change_session_behavior() { + let client = client(); + let base = base_url().await; + let uid = uuid::Uuid::new_v4().simple().to_string(); + let handle = format!("pwch{}", &uid[..8]); + let email = format!("pwch{}@test.com", &uid[..8]); + let old_password = "OldPassword123!"; + let new_password = "NewPassword456!"; + + let create_res = client + .post(format!("{}/xrpc/com.atproto.server.createAccount", base)) + .json(&json!({ + "handle": handle, + "email": email, + "password": old_password + })) + .send() + .await + .expect("Account creation failed"); + assert_eq!(create_res.status(), StatusCode::OK); + let account: Value = create_res.json().await.unwrap(); + let did = account["did"].as_str().unwrap().to_string(); + let jwt1 = verify_new_account(&client, &did).await; + + create_post(&client, &did, &jwt1, "Post before password change").await; + + let change_pw_res = client + .post(format!("{}/xrpc/_account.changePassword", base)) + .bearer_auth(&jwt1) + .json(&json!({ + "currentPassword": old_password, + "newPassword": new_password + })) + .send() + .await + .expect("Password change failed"); + assert_eq!(change_pw_res.status(), StatusCode::OK); + + let old_pw_login = client + .post(format!("{}/xrpc/com.atproto.server.createSession", base)) + .json(&json!({ + "identifier": handle, + "password": old_password + })) + .send() + .await + .unwrap(); + assert!( + old_pw_login.status() == StatusCode::UNAUTHORIZED + || old_pw_login.status() == StatusCode::BAD_REQUEST, + "Old password should not work" + ); + + let new_pw_login = client + .post(format!("{}/xrpc/com.atproto.server.createSession", base)) + .json(&json!({ + "identifier": handle, + "password": new_password + })) + .send() + .await + .expect("New password login failed"); + assert_eq!(new_pw_login.status(), StatusCode::OK); + let new_session: Value = new_pw_login.json().await.unwrap(); + let new_jwt = new_session["accessJwt"].as_str().unwrap(); + + let posts_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(new_jwt) + .query(&[("repo", did.as_str()), ("collection", "app.bsky.feed.post")]) + .send() + .await + .unwrap(); + let posts_body: Value = posts_res.json().await.unwrap(); + assert_eq!( + posts_body["records"].as_array().unwrap().len(), + 1, + "Post should still exist after password change" + ); +} + +#[tokio::test] +async fn test_backup_restore_workflow() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("backup-restore").await; + + futures::future::join_all((0..3).map(|i| { + let client = client.clone(); + let did = did.clone(); + let jwt = jwt.clone(); + async move { + create_post(&client, &did, &jwt, &format!("Post {} for backup test", i)).await; + } + })) + .await; + + let blob_data = b"Blob data for backup test"; + let upload_res = client + .post(format!("{}/xrpc/com.atproto.repo.uploadBlob", base)) + .header(header::CONTENT_TYPE, "text/plain") + .bearer_auth(&jwt) + .body(blob_data.to_vec()) + .send() + .await + .expect("Blob upload failed"); + assert_eq!(upload_res.status(), StatusCode::OK); + let upload_body: Value = upload_res.json().await.unwrap(); + let blob = upload_body["blob"].clone(); + + let profile_res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.actor.profile", + "rkey": "self", + "record": { + "$type": "app.bsky.actor.profile", + "displayName": "Backup Test User", + "avatar": blob + } + })) + .send() + .await + .unwrap(); + assert_eq!(profile_res.status(), StatusCode::OK); + + let backup1_res = client + .post(format!("{}/xrpc/_backup.createBackup", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("Backup 1 failed"); + assert_eq!(backup1_res.status(), StatusCode::OK); + let backup1: Value = backup1_res.json().await.unwrap(); + let backup1_id = backup1["id"].as_str().unwrap(); + let backup1_rev = backup1["repoRev"].as_str().unwrap(); + + create_post(&client, &did, &jwt, "Post 4 after first backup").await; + create_post(&client, &did, &jwt, "Post 5 after first backup").await; + + let backup2_res = client + .post(format!("{}/xrpc/_backup.createBackup", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("Backup 2 failed"); + assert_eq!(backup2_res.status(), StatusCode::OK); + let backup2: Value = backup2_res.json().await.unwrap(); + let backup2_id = backup2["id"].as_str().unwrap(); + let backup2_rev = backup2["repoRev"].as_str().unwrap(); + + assert_ne!( + backup1_rev, backup2_rev, + "Backups should have different revs" + ); + + let list_res = client + .get(format!("{}/xrpc/_backup.listBackups", base)) + .bearer_auth(&jwt) + .send() + .await + .expect("List backups failed"); + let list_body: Value = list_res.json().await.unwrap(); + let backups = list_body["backups"].as_array().unwrap(); + assert_eq!(backups.len(), 2, "Should have 2 backups"); + + let download1 = client + .get(format!("{}/xrpc/_backup.getBackup?id={}", base, backup1_id)) + .bearer_auth(&jwt) + .send() + .await + .expect("Download backup 1 failed"); + assert_eq!(download1.status(), StatusCode::OK); + let backup1_bytes = download1.bytes().await.unwrap(); + + let download2 = client + .get(format!("{}/xrpc/_backup.getBackup?id={}", base, backup2_id)) + .bearer_auth(&jwt) + .send() + .await + .expect("Download backup 2 failed"); + assert_eq!(download2.status(), StatusCode::OK); + let backup2_bytes = download2.bytes().await.unwrap(); + + assert!( + backup2_bytes.len() > backup1_bytes.len(), + "Second backup should be larger (more posts)" + ); + + let delete_old = client + .post(format!( + "{}/xrpc/_backup.deleteBackup?id={}", + base, backup1_id + )) + .bearer_auth(&jwt) + .send() + .await + .expect("Delete backup failed"); + assert_eq!(delete_old.status(), StatusCode::OK); + + let final_list = client + .get(format!("{}/xrpc/_backup.listBackups", base)) + .bearer_auth(&jwt) + .send() + .await + .unwrap(); + let final_body: Value = final_list.json().await.unwrap(); + assert_eq!(final_body["backups"].as_array().unwrap().len(), 1); +} + +#[tokio::test] +async fn test_scale_100_posts_with_pagination() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("scale-posts").await; + + let post_count = 1000; + let post_futures: Vec<_> = (0..post_count) + .map(|i| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + async move { + let res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": format!("Scale test post number {}", i), + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Post creation failed"); + let status = res.status(); + let body: Value = res.json().await.unwrap_or_default(); + assert_eq!( + status, + StatusCode::OK, + "Failed to create post {}: {:?}", + i, + body + ); + } + }) + .collect(); + + join_all(post_futures).await; + + let count_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("limit", "1"), + ]) + .send() + .await + .unwrap(); + assert_eq!(count_res.status(), StatusCode::OK); + + let all_uris: Vec = + paginate_records(&client, base, &jwt, &did, "app.bsky.feed.post", 25) + .await + .iter() + .filter_map(|r| r["uri"].as_str().map(String::from)) + .collect(); + + assert_eq!( + all_uris.len(), + post_count, + "Should have paginated through all {} posts", + post_count + ); + + let unique_uris: std::collections::HashSet<_> = all_uris.iter().collect(); + assert_eq!( + unique_uris.len(), + post_count, + "All posts should have unique URIs" + ); + + let delete_futures: Vec<_> = all_uris + .iter() + .take(500) + .map(|uri| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + let rkey = uri.split('/').next_back().unwrap().to_string(); + async move { + let res = client + .post(format!("{}/xrpc/com.atproto.repo.deleteRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey + })) + .send() + .await + .expect("Delete failed"); + assert_eq!(res.status(), StatusCode::OK); + } + }) + .collect(); + + join_all(delete_futures).await; + + let final_count = count_records(&client, base, &jwt, &did, "app.bsky.feed.post").await; + assert_eq!( + final_count, 500, + "Should have 500 posts remaining after deleting 500" + ); +} + +#[tokio::test] +async fn test_scale_many_users_social_graph() { + let client = client(); + let base = base_url().await; + + let user_count = 50; + let user_futures: Vec<_> = (0..user_count) + .map(|i| async move { setup_new_user(&format!("graph{}", i)).await }) + .collect(); + + let users: Vec<(String, String)> = join_all(user_futures).await; + + let follow_futures: Vec<_> = users + .iter() + .enumerate() + .flat_map(|(i, (follower_did, follower_jwt))| { + let client = client.clone(); + let base = base.to_string(); + users.iter().enumerate().filter(move |(j, _)| *j != i).map({ + let client = client.clone(); + let base = base.clone(); + let follower_did = follower_did.clone(); + let follower_jwt = follower_jwt.clone(); + move |(_, (followee_did, _))| { + let client = client.clone(); + let base = base.clone(); + let follower_did = follower_did.clone(); + let follower_jwt = follower_jwt.clone(); + let followee_did = followee_did.clone(); + async move { + let rkey = format!( + "follow_{}", + &uuid::Uuid::new_v4().simple().to_string()[..12] + ); + let res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&follower_jwt) + .json(&json!({ + "repo": follower_did, + "collection": "app.bsky.graph.follow", + "rkey": rkey, + "record": { + "$type": "app.bsky.graph.follow", + "subject": followee_did, + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Follow failed"); + let status = res.status(); + let body: Value = res.json().await.unwrap_or_default(); + assert_eq!(status, StatusCode::OK, "Follow failed: {:?}", body); + } + } + }) + }) + .collect(); + + join_all(follow_futures).await; + + let expected_follows_per_user = user_count - 1; + let verify_futures: Vec<_> = users + .iter() + .map(|(did, jwt)| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + async move { + let res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.graph.follow"), + ("limit", "100"), + ]) + .send() + .await + .unwrap(); + let body: Value = res.json().await.unwrap(); + body["records"].as_array().unwrap().len() + } + }) + .collect(); + + let follow_counts: Vec = join_all(verify_futures).await; + let total_follows: usize = follow_counts.iter().sum(); + + assert_eq!( + total_follows, + user_count * expected_follows_per_user, + "Each of {} users should follow {} others = {} total follows", + user_count, + expected_follows_per_user, + user_count * expected_follows_per_user + ); + + let (poster_did, poster_jwt) = &users[0]; + let (post_uri, post_cid) = create_post(&client, poster_did, poster_jwt, "Popular post").await; + + let like_futures: Vec<_> = users + .iter() + .skip(1) + .map(|(liker_did, liker_jwt)| { + let client = client.clone(); + let base = base.to_string(); + let liker_did = liker_did.clone(); + let liker_jwt = liker_jwt.clone(); + let post_uri = post_uri.clone(); + let post_cid = post_cid.clone(); + async move { + let rkey = format!("like_{}", &uuid::Uuid::new_v4().simple().to_string()[..12]); + let res = client + .post(format!("{}/xrpc/com.atproto.repo.putRecord", base)) + .bearer_auth(&liker_jwt) + .json(&json!({ + "repo": liker_did, + "collection": "app.bsky.feed.like", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.like", + "subject": { "uri": post_uri, "cid": post_cid }, + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Like failed"); + assert_eq!(res.status(), StatusCode::OK); + } + }) + .collect(); + + join_all(like_futures).await; +} + +#[tokio::test] +async fn test_scale_many_blobs_in_repo() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("scale-blobs").await; + + let blob_count = 300; + let blob_futures: Vec<_> = (0..blob_count) + .map(|i| { + let client = client.clone(); + let base = base.to_string(); + let jwt = jwt.clone(); + async move { + let blob_data = format!("Blob data number {} with some padding to make it realistic size for testing purposes", i); + let res = client + .post(format!("{}/xrpc/com.atproto.repo.uploadBlob", base)) + .header(header::CONTENT_TYPE, "text/plain") + .bearer_auth(&jwt) + .body(blob_data) + .send() + .await + .expect("Blob upload failed"); + assert_eq!(res.status(), StatusCode::OK, "Failed to upload blob {}", i); + let body: Value = res.json().await.unwrap(); + body["blob"].clone() + } + }) + .collect(); + + let blobs: Vec = join_all(blob_futures).await; + + let post_futures: Vec<_> = blobs + .chunks(3) + .enumerate() + .map(|(i, blob_chunk)| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + let images: Vec = blob_chunk + .iter() + .enumerate() + .map(|(j, blob)| { + json!({ + "alt": format!("Image {} in post {}", j, i), + "image": blob + }) + }) + .collect(); + async move { + let res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": format!("Post {} with {} images", i, images.len()), + "createdAt": Utc::now().to_rfc3339(), + "embed": { + "$type": "app.bsky.embed.images", + "images": images + } + } + })) + .send() + .await + .expect("Post with blobs failed"); + let status = res.status(); + let body: Value = res.json().await.unwrap_or_default(); + assert_eq!(status, StatusCode::OK, "Post with blobs failed: {:?}", body); + } + }) + .collect(); + + join_all(post_futures).await; + + let list_blobs_res = client + .get(format!("{}/xrpc/com.atproto.sync.listBlobs", base)) + .query(&[("did", did.as_str())]) + .send() + .await + .expect("List blobs failed"); + assert_eq!(list_blobs_res.status(), StatusCode::OK); + let blobs_body: Value = list_blobs_res.json().await.unwrap(); + let cids = blobs_body["cids"].as_array().unwrap(); + assert_eq!( + cids.len(), + blob_count, + "Should have {} blobs in repo", + blob_count + ); + + let verify_futures: Vec<_> = cids + .iter() + .take(10) + .map(|cid| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let cid_str = cid.as_str().unwrap().to_string(); + async move { + let res = client + .get(format!("{}/xrpc/com.atproto.sync.getBlob", base)) + .query(&[("did", did.as_str()), ("cid", cid_str.as_str())]) + .send() + .await + .expect("Get blob failed"); + assert_eq!( + res.status(), + StatusCode::OK, + "Failed to get blob {}", + cid_str + ); + let bytes = res.bytes().await.unwrap(); + assert!(bytes.len() > 50, "Blob should have content"); + } + }) + .collect(); + + join_all(verify_futures).await; +} + +#[tokio::test] +async fn test_scale_batch_operations() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("scale-batch").await; + + let batch_size = 200; + let writes: Vec = (0..batch_size) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": format!("batch_{:03}", i), + "value": { + "$type": "app.bsky.feed.post", + "text": format!("Batch created post {}", i), + "createdAt": Utc::now().to_rfc3339() + } + }) + }) + .collect(); + + let apply_res = client + .post(format!("{}/xrpc/com.atproto.repo.applyWrites", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "writes": writes + })) + .send() + .await + .expect("Batch create failed"); + assert_eq!( + apply_res.status(), + StatusCode::OK, + "Batch create of {} posts should succeed", + batch_size + ); + + let batch_count = count_records(&client, base, &jwt, &did, "app.bsky.feed.post").await; + assert_eq!( + batch_count, batch_size, + "Should have {} posts after batch create", + batch_size + ); + + let update_writes: Vec = (0..batch_size) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": format!("batch_{:03}", i), + "value": { + "$type": "app.bsky.feed.post", + "text": format!("UPDATED batch post {}", i), + "createdAt": Utc::now().to_rfc3339() + } + }) + }) + .collect(); + + let update_res = client + .post(format!("{}/xrpc/com.atproto.repo.applyWrites", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "writes": update_writes + })) + .send() + .await + .expect("Batch update failed"); + assert_eq!( + update_res.status(), + StatusCode::OK, + "Batch update of {} posts should succeed", + batch_size + ); + + let verify_res = client + .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", "batch_000"), + ]) + .send() + .await + .unwrap(); + let verify_body: Value = verify_res.json().await.unwrap(); + assert!( + verify_body["value"]["text"] + .as_str() + .unwrap() + .starts_with("UPDATED"), + "Post should be updated" + ); + + let delete_writes: Vec = (0..batch_size) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": format!("batch_{:03}", i) + }) + }) + .collect(); + + let delete_res = client + .post(format!("{}/xrpc/com.atproto.repo.applyWrites", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "writes": delete_writes + })) + .send() + .await + .expect("Batch delete failed"); + assert_eq!( + delete_res.status(), + StatusCode::OK, + "Batch delete of {} posts should succeed", + batch_size + ); + + let final_res = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("limit", "100"), + ]) + .send() + .await + .unwrap(); + let final_body: Value = final_res.json().await.unwrap(); + assert_eq!( + final_body["records"].as_array().unwrap().len(), + 0, + "Should have 0 posts after batch delete" + ); +} + +#[tokio::test] +async fn test_scale_reply_thread_depth() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("deep-thread").await; + + let root_res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": "Root post of deep thread", + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Root post failed"); + assert_eq!(root_res.status(), StatusCode::OK); + let root_body: Value = root_res.json().await.unwrap(); + let root_uri = root_body["uri"].as_str().unwrap().to_string(); + let root_cid = root_body["cid"].as_str().unwrap().to_string(); + + let thread_depth = 500; + + struct ReplyState { + parent_uri: String, + parent_cid: String, + depth: usize, + } + + let initial_state = ReplyState { + parent_uri: root_uri.clone(), + parent_cid: root_cid.clone(), + depth: 1, + }; + + let (parent_uri, reply_count) = futures::stream::unfold(initial_state, |state| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + let root_uri = root_uri.clone(); + let root_cid = root_cid.clone(); + async move { + if state.depth > thread_depth { + return None; + } + + let res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": format!("Reply at depth {}", state.depth), + "createdAt": Utc::now().to_rfc3339(), + "reply": { + "root": { "uri": root_uri, "cid": root_cid }, + "parent": { "uri": &state.parent_uri, "cid": &state.parent_cid } + } + } + })) + .send() + .await + .expect("Reply failed"); + assert_eq!( + res.status(), + StatusCode::OK, + "Reply at depth {} failed", + state.depth + ); + let body: Value = res.json().await.unwrap(); + let new_uri = body["uri"].as_str().unwrap().to_string(); + let new_cid = body["cid"].as_str().unwrap().to_string(); + + let next_state = ReplyState { + parent_uri: new_uri.clone(), + parent_cid: new_cid, + depth: state.depth + 1, + }; + + Some((new_uri, next_state)) + } + }) + .fold((String::new(), 0usize), |(_, count), uri| async move { + (uri, count + 1) + }) + .await; + + assert_eq!( + reply_count, thread_depth, + "Should have created {} replies", + thread_depth + ); + + let thread_count = count_records(&client, base, &jwt, &did, "app.bsky.feed.post").await; + assert_eq!( + thread_count, + thread_depth + 1, + "Should have root + {} replies", + thread_depth + ); + + let deepest_rkey = parent_uri.split('/').next_back().unwrap(); + let deep_res = client + .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", deepest_rkey), + ]) + .send() + .await + .unwrap(); + assert_eq!(deep_res.status(), StatusCode::OK); + let deep_body: Value = deep_res.json().await.unwrap(); + assert_eq!( + deep_body["value"]["reply"]["root"]["uri"], root_uri, + "Deepest reply should reference root" + ); +} + +#[tokio::test] +async fn test_concurrent_import_and_writes() { + let client = client(); + let base = base_url().await; + let (did, jwt) = setup_new_user("import-conc").await; + + let key_bytes = get_user_signing_key(&did) + .await + .expect("Failed to get user signing key"); + let signing_key = SigningKey::from_slice(&key_bytes).expect("Failed to create signing key"); + + let (car_bytes, _root_cid) = build_car_with_signature(&did, &signing_key); + + let write_count = 10; + + let import_future = { + let client = client.clone(); + let base = base.to_string(); + let jwt = jwt.clone(); + let car_bytes = car_bytes.clone(); + async move { + let res = client + .post(format!("{}/xrpc/com.atproto.repo.importRepo", base)) + .bearer_auth(&jwt) + .header("Content-Type", "application/vnd.ipld.car") + .body(car_bytes) + .send() + .await + .expect("Import request failed"); + let status = res.status(); + let body: Value = res.json().await.unwrap_or_default(); + assert_eq!(status, StatusCode::OK, "Import should succeed: {:?}", body); + } + }; + + let write_futures: Vec<_> = (0..write_count) + .map(|i| { + let client = client.clone(); + let base = base.to_string(); + let did = did.clone(); + let jwt = jwt.clone(); + async move { + let res = client + .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "record": { + "$type": "app.bsky.feed.post", + "text": format!("Concurrent post {}", i), + "createdAt": Utc::now().to_rfc3339() + } + })) + .send() + .await + .expect("Write request failed"); + let status = res.status(); + let body: Value = res.json().await.unwrap_or_default(); + assert_eq!( + status, + StatusCode::OK, + "Write {} should succeed: {:?}", + i, + body + ); + } + }) + .collect(); + + tokio::join!(import_future, join_all(write_futures)); + + let final_posts = client + .get(format!("{}/xrpc/com.atproto.repo.listRecords", base)) + .bearer_auth(&jwt) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("limit", "100"), + ]) + .send() + .await + .unwrap(); + let final_body: Value = final_posts.json().await.unwrap(); + let record_count = final_body["records"].as_array().unwrap().len(); + + let min_expected = write_count; + assert!( + record_count >= min_expected, + "Expected at least {} records (from writes), got {} (import may also contribute records)", + min_expected, + record_count + ); +} diff --git a/deploy/nginx/nginx-quadlet.conf b/deploy/nginx/nginx-quadlet.conf index 5b65cb7..10629c0 100644 --- a/deploy/nginx/nginx-quadlet.conf +++ b/deploy/nginx/nginx-quadlet.conf @@ -92,6 +92,15 @@ http { proxy_set_header X-Forwarded-Proto $scheme; } + location /webhook/ { + proxy_pass http://127.0.0.1:3000; + proxy_http_version 1.1; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + location = /metrics { proxy_pass http://127.0.0.1:3000; proxy_http_version 1.1; diff --git a/docs/install-debian.md b/docs/install-debian.md index 70ecf5e..21ac261 100644 --- a/docs/install-debian.md +++ b/docs/install-debian.md @@ -203,6 +203,15 @@ server { proxy_set_header X-Forwarded-Proto $scheme; } + location /webhook/ { + proxy_pass http://127.0.0.1:3000; + proxy_http_version 1.1; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + location = /metrics { proxy_pass http://127.0.0.1:3000; proxy_http_version 1.1; diff --git a/frontend/src/components/dashboard/CommsContent.svelte b/frontend/src/components/dashboard/CommsContent.svelte index cfc723c..b13546f 100644 --- a/frontend/src/components/dashboard/CommsContent.svelte +++ b/frontend/src/components/dashboard/CommsContent.svelte @@ -17,6 +17,7 @@ let saving = $state(false) let preferredChannel = $state('email') let availableCommsChannels = $state(['email']) + let telegramBotUsername = $state(undefined) let email = $state('') let discordId = $state('') let discordVerified = $state(false) @@ -24,6 +25,9 @@ let telegramVerified = $state(false) let signalNumber = $state('') let signalVerified = $state(false) + let savedDiscordId = $state('') + let savedTelegramUsername = $state('') + let savedSignalNumber = $state('') let verifyingChannel = $state(null) let verificationCode = $state('') let historyLoading = $state(true) @@ -59,7 +63,11 @@ telegramVerified = prefs.telegramVerified signalNumber = prefs.signalNumber ?? '' signalVerified = prefs.signalVerified + savedDiscordId = discordId + savedTelegramUsername = telegramUsername + savedSignalNumber = signalNumber availableCommsChannels = serverInfo.availableCommsChannels ?? ['email'] + telegramBotUsername = serverInfo.telegramBotUsername } catch (e) { toast.error(e instanceof ApiError ? e.message : $_('comms.failedToLoad')) } finally { @@ -71,15 +79,24 @@ e.preventDefault() saving = true try { - await api.updateNotificationPrefs(session.accessJwt, { + const result = await api.updateNotificationPrefs(session.accessJwt, { preferredChannel, - discordId: discordId || undefined, - telegramUsername: telegramUsername || undefined, - signalNumber: signalNumber || undefined, + discordId: discordId !== savedDiscordId ? discordId : undefined, + telegramUsername: telegramUsername !== savedTelegramUsername ? telegramUsername : undefined, + signalNumber: signalNumber !== savedSignalNumber ? signalNumber : undefined, }) await refreshSession() toast.success($_('comms.preferencesSaved')) - await loadPrefs() + savedDiscordId = discordId + savedTelegramUsername = telegramUsername + savedSignalNumber = signalNumber + const channelToVerify = result.verificationRequired?.find( + (ch: string) => ch === 'discord' || ch === 'telegram' || ch === 'signal' + ) + if (channelToVerify) { + verifyingChannel = channelToVerify + verificationCode = '' + } } catch (e) { toast.error(e instanceof ApiError ? e.message : $_('comms.failedToSave')) } finally { @@ -218,7 +235,9 @@
- {$_('comms.primary')} + + {preferredChannel === 'email' ? $_('comms.primary') : $_('comms.verified')} +
@@ -229,7 +248,7 @@ {#if discordId} - {discordVerified ? $_('comms.verified') : $_('comms.notVerified')} + {preferredChannel === 'discord' && discordVerified ? $_('comms.primary') : discordVerified ? $_('comms.verified') : $_('comms.notVerified')} {/if} @@ -242,7 +261,7 @@ placeholder={$_('register.discordIdPlaceholder')} disabled={saving} /> - {#if discordId && !discordVerified} + {#if discordId && discordId === savedDiscordId && !discordVerified} {/if} @@ -251,7 +270,7 @@ {/if} {#if verifyingChannel === 'discord'}
- +
@@ -265,7 +284,7 @@ {#if telegramUsername} - {telegramVerified ? $_('comms.verified') : $_('comms.notVerified')} + {preferredChannel === 'telegram' && telegramVerified ? $_('comms.primary') : telegramVerified ? $_('comms.verified') : $_('comms.notVerified')} {/if} @@ -278,18 +297,15 @@ placeholder={$_('register.telegramUsernamePlaceholder')} disabled={saving} /> - {#if telegramUsername && !telegramVerified} - - {/if} {#if telegramInUse}

{$_('comms.telegramInUseWarning')}

{/if} - {#if verifyingChannel === 'telegram'} -
- - - + {#if telegramUsername && telegramUsername === savedTelegramUsername && !telegramVerified && telegramBotUsername} + {@const encodedHandle = session.handle.replaceAll('.', '_')} +
+ {$_('comms.telegramOpenLink')} + {$_('comms.telegramStartBot', { values: { botUsername: telegramBotUsername, handle: session.handle } })}
{/if}
@@ -301,7 +317,7 @@ {#if signalNumber} - {signalVerified ? $_('comms.verified') : $_('comms.notVerified')} + {preferredChannel === 'signal' && signalVerified ? $_('comms.primary') : signalVerified ? $_('comms.verified') : $_('comms.notVerified')} {/if} @@ -314,7 +330,7 @@ placeholder={$_('register.signalNumberPlaceholder')} disabled={saving} /> - {#if signalNumber && !signalVerified} + {#if signalNumber && signalNumber === savedSignalNumber && !signalVerified} {/if} @@ -323,7 +339,7 @@ {/if} {#if verifyingChannel === 'signal'}
- +
@@ -505,6 +521,23 @@ color: var(--warning-text); } + .telegram-verify-prompt { + display: flex; + flex-direction: column; + gap: var(--space-2); + padding: var(--space-3) var(--space-4); + background: var(--accent-bg, var(--bg-card)); + border: 1px solid var(--accent, var(--border-color)); + border-radius: var(--radius-md); + font-size: var(--text-sm); + color: var(--text-primary); + } + + .manual-hint { + font-size: var(--text-xs); + color: var(--text-secondary); + } + .verify-btn { padding: var(--space-2) var(--space-3); font-size: var(--text-sm); @@ -517,7 +550,8 @@ } .verify-form input { - width: 120px; + flex: 1; + min-width: 0; } .verify-form button { diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index 1ebc2e0..ebac18d 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -80,6 +80,7 @@ import type { TotpStatus, UpdateLegacyLoginResponse, UpdateLocaleResponse, + UpdateNotificationPrefsResponse, UploadBlobResponse, VerificationChannel, VerifyMigrationEmailResponse, @@ -479,6 +480,16 @@ export const api = { }); }, + checkChannelVerified( + did: string, + channel: string, + ): Promise<{ verified: boolean }> { + return xrpc("_checkChannelVerified", { + method: "POST", + body: { did, channel }, + }); + }, + checkEmailInUse(email: string): Promise<{ inUse: boolean }> { return xrpc("_account.checkEmailInUse", { method: "POST", @@ -648,7 +659,7 @@ export const api = { discordId?: string; telegramUsername?: string; signalNumber?: string; - }): Promise { + }): Promise { return xrpc("_account.updateNotificationPrefs", { method: "POST", token, @@ -1847,7 +1858,7 @@ export const typedApi = { telegramUsername?: string; signalNumber?: string; }, - ): Promise> { + ): Promise> { return xrpcResult("_account.updateNotificationPrefs", { method: "POST", token, diff --git a/frontend/src/lib/registration/VerificationStep.svelte b/frontend/src/lib/registration/VerificationStep.svelte index dd76658..8aba59b 100644 --- a/frontend/src/lib/registration/VerificationStep.svelte +++ b/frontend/src/lib/registration/VerificationStep.svelte @@ -1,5 +1,5 @@
-

- We've sent a verification code to your {channelLabel(flow.info.verificationChannel)}. - Enter it below to continue. -

+ {#if isTelegram && telegramBotUsername} + {@const handle = flow.account?.handle ?? `${flow.info.handle.trim()}.${flow.state.pdsHostname}`} + {@const encodedHandle = handle.replaceAll('.', '_')} +

+ Open Telegram to verify, + or send /start {handle} to @{telegramBotUsername} manually. +

+

Waiting for verification...

+ {:else} +

+ We've sent a verification code to your {channelLabel(flow.info.verificationChannel)}. + Enter it below to continue. +

- {#if resendMessage} -
{resendMessage}
+ {#if resendMessage} +
{resendMessage}
+ {/if} + +
+
+ + + Copy the entire code from your message, including dashes. +
+ + + + +
{/if} - -
-
- - - Copy the entire code from your message, including dashes. -
- - - - -
diff --git a/migrations/20260203_telegram_chat_id.sql b/migrations/20260203_telegram_chat_id.sql new file mode 100644 index 0000000..704e8b6 --- /dev/null +++ b/migrations/20260203_telegram_chat_id.sql @@ -0,0 +1 @@ +ALTER TABLE users ADD COLUMN telegram_chat_id BIGINT; diff --git a/migrations/20260204_channel_verified_comms_type.sql b/migrations/20260204_channel_verified_comms_type.sql new file mode 100644 index 0000000..067de80 --- /dev/null +++ b/migrations/20260204_channel_verified_comms_type.sql @@ -0,0 +1 @@ +ALTER TYPE comms_type ADD VALUE IF NOT EXISTS 'channel_verified'; diff --git a/nginx.frontend.conf b/nginx.frontend.conf index 65f0fc0..6b2467e 100644 --- a/nginx.frontend.conf +++ b/nginx.frontend.conf @@ -122,6 +122,15 @@ http { proxy_set_header X-Forwarded-Proto $scheme; } + location /webhook/ { + proxy_pass http://backend; + proxy_http_version 1.1; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + location = /metrics { proxy_pass http://backend; proxy_http_version 1.1; diff --git a/scripts/run-tests.sh b/scripts/run-tests.sh index c91e64c..8834398 100755 --- a/scripts/run-tests.sh +++ b/scripts/run-tests.sh @@ -18,4 +18,5 @@ sqlx migrate run --source "$PROJECT_DIR/migrations" echo "" echo "Running tests..." echo "" +ulimit -n 65536 cargo nextest run "$@"