Coverage for app/backend/src/tests/test_bg_jobs.py: 99%

885 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-19 15:47 +0000

1from datetime import date, timedelta 

2from typing import Any 

3from unittest.mock import call, patch 

4 

5import pytest 

6import requests 

7from google.protobuf import empty_pb2 

8from google.protobuf.empty_pb2 import Empty 

9from prometheus_client import REGISTRY 

10from sqlalchemy import select, text 

11from sqlalchemy.sql import delete, func 

12 

13import couchers.jobs.worker 

14from couchers import experimentation 

15from couchers.config import config 

16from couchers.constants import ( 

17 HOST_REQUEST_MAX_REMINDERS, 

18 HOST_REQUEST_REMINDER_INTERVAL, 

19 MISSED_MESSAGES_DELAY, 

20 MISSED_MESSAGES_DELAY_WITH_PUSH, 

21) 

22from couchers.crypto import urlsafe_secure_token 

23from couchers.db import session_scope 

24from couchers.email.dev import print_dev_email 

25from couchers.email.queuing import queue_email 

26from couchers.jobs import handlers 

27from couchers.jobs.definitions import Job 

28from couchers.jobs.enqueue import queue_job 

29from couchers.jobs.handlers import ( 

30 add_users_to_email_list, 

31 enforce_community_membership, 

32 purge_account_deletion_tokens, 

33 purge_login_tokens, 

34 purge_password_reset_tokens, 

35 send_host_request_reminders, 

36 send_message_notifications, 

37 send_onboarding_emails, 

38 send_reference_reminders, 

39 send_request_notifications, 

40 update_badges, 

41 update_recommendation_scores, 

42) 

43from couchers.jobs.worker import _run_job_and_schedule, process_job, run_scheduler, service_jobs 

44from couchers.materialized_views import refresh_materialized_views 

45from couchers.metrics import create_prometheus_server 

46from couchers.models import ( 

47 AccountDeletionToken, 

48 BackgroundJob, 

49 BackgroundJobState, 

50 DeviceType, 

51 Email, 

52 HostRequest, 

53 HostRequestStatus, 

54 LoginToken, 

55 Message, 

56 MessageType, 

57 PasswordResetToken, 

58 PostalVerificationAttempt, 

59 PostalVerificationStatus, 

60 PushNotificationPlatform, 

61 PushNotificationSubscription, 

62 User, 

63 UserBadge, 

64 UserBlock, 

65 Volunteer, 

66) 

67from couchers.proto import conversations_pb2, messages_pb2, requests_pb2 

68from couchers.proto.internal import jobs_pb2 

69from couchers.utils import now, today 

70from tests.fixtures.db import generate_user, make_friends, make_user_block, make_volunteer 

71from tests.fixtures.misc import PushCollector, process_jobs 

72from tests.fixtures.sessions import conversations_session, requests_session 

73from tests.fixtures.timewarp import Timewarp 

74from tests.test_references import create_host_reference, create_host_request, create_host_request_by_date 

75from tests.test_requests import valid_request_text 

76 

77 

78def _add_mobile_push_subscription(user_id: int, *, disabled: bool = False) -> None: 

79 with session_scope() as session: 

80 sub = PushNotificationSubscription( 

81 user_id=user_id, 

82 platform=PushNotificationPlatform.expo, 

83 token=f"ExponentPushToken[{user_id}]", 

84 device_name="Test phone", 

85 device_type=DeviceType.ios, 

86 ) 

87 session.add(sub) 

88 if disabled: 

89 session.flush() 

90 sub.disabled_at = now() 

91 

92 

93def _count_queued_emails() -> int: 

94 with session_scope() as session: 

95 return session.execute( 

96 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

97 ).scalar_one() 

98 

99 

100def _check_job_counter(port, job, status, attempt, exception): 

101 metrics_string = requests.get(f"http://localhost:{port}").text 

102 string_to_check = f'attempt="{attempt}",exception="{exception}",job="{job}",status="{status}"' 

103 assert string_to_check in metrics_string 

104 

105 

106def _queued_seconds_sum(): 

107 total = 0.0 

108 for metric in REGISTRY.collect(): 

109 if metric.name == "couchers_background_jobs_queued_seconds": 

110 total += sum(sample.value for sample in metric.samples if sample.name.endswith("_sum")) 

111 return total 

112 

113 

114def test_email_job(db): 

115 with session_scope() as session: 

116 queue_email( 

117 session, 

118 jobs_pb2.SendEmailPayload( 

119 sender_name="sender_name", 

120 sender_email="sender_email", 

121 recipient="recipient", 

122 subject="subject", 

123 plain="plain", 

124 html="html", 

125 ), 

126 ) 

127 

128 def mock_print_dev_email(payload): 

129 assert payload.sender_name == "sender_name" 

130 assert payload.sender_email == "sender_email" 

131 assert payload.recipient == "recipient" 

132 assert payload.subject == "subject" 

133 assert payload.plain == "plain" 

134 assert payload.html == "html" 

135 return print_dev_email(payload) 

136 

137 with patch("couchers.jobs.handlers.print_dev_email", mock_print_dev_email): 

138 process_job() 

139 

140 with session_scope() as session: 

141 assert ( 

142 session.execute( 

143 select(func.count()) 

144 .select_from(BackgroundJob) 

145 .where(BackgroundJob.state == BackgroundJobState.completed) 

146 ).scalar_one() 

147 == 1 

148 ) 

149 assert ( 

150 session.execute( 

151 select(func.count()) 

152 .select_from(BackgroundJob) 

153 .where(BackgroundJob.state != BackgroundJobState.completed) 

154 ).scalar_one() 

155 == 0 

156 ) 

157 

158 

159def test_purge_login_tokens(db): 

160 user, api_token = generate_user() 

161 

162 with session_scope() as session: 

163 login_token = LoginToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now()) 

164 session.add(login_token) 

165 assert session.execute(select(func.count()).select_from(LoginToken)).scalar_one() == 1 

166 

167 queue_job(session, job=purge_login_tokens, payload=empty_pb2.Empty()) 

168 process_job() 

169 

170 with session_scope() as session: 

171 assert session.execute(select(func.count()).select_from(LoginToken)).scalar_one() == 0 

172 

173 with session_scope() as session: 

174 assert ( 

175 session.execute( 

176 select(func.count()) 

177 .select_from(BackgroundJob) 

178 .where(BackgroundJob.state == BackgroundJobState.completed) 

179 ).scalar_one() 

180 == 1 

181 ) 

182 assert ( 

183 session.execute( 

184 select(func.count()) 

185 .select_from(BackgroundJob) 

186 .where(BackgroundJob.state != BackgroundJobState.completed) 

187 ).scalar_one() 

188 == 0 

189 ) 

190 

191 

192def test_purge_password_reset_tokens(db): 

193 user, api_token = generate_user() 

194 

195 with session_scope() as session: 

196 password_reset_token = PasswordResetToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now()) 

197 session.add(password_reset_token) 

198 assert session.execute(select(func.count()).select_from(PasswordResetToken)).scalar_one() == 1 

199 

200 queue_job(session, job=purge_password_reset_tokens, payload=empty_pb2.Empty()) 

201 process_job() 

202 

203 with session_scope() as session: 

204 assert session.execute(select(func.count()).select_from(PasswordResetToken)).scalar_one() == 0 

205 

206 with session_scope() as session: 

207 assert ( 

208 session.execute( 

209 select(func.count()) 

210 .select_from(BackgroundJob) 

211 .where(BackgroundJob.state == BackgroundJobState.completed) 

212 ).scalar_one() 

213 == 1 

214 ) 

215 assert ( 

216 session.execute( 

217 select(func.count()) 

218 .select_from(BackgroundJob) 

219 .where(BackgroundJob.state != BackgroundJobState.completed) 

220 ).scalar_one() 

221 == 0 

222 ) 

223 

224 

225def test_purge_account_deletion_tokens(db): 

226 user, api_token = generate_user() 

227 user2, api_token2 = generate_user() 

228 user3, api_token3 = generate_user() 

229 

230 with session_scope() as session: 

231 """ 

232 3 cases: 

233 1) Token is valid 

234 2) Token expired but account retrievable 

235 3) Account is irretrievable (and expired) 

236 """ 

237 account_deletion_tokens = [ 

238 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now() - timedelta(hours=2)), 

239 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user2.id, expiry=now()), 

240 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user3.id, expiry=now() + timedelta(hours=5)), 

241 ] 

242 for token in account_deletion_tokens: 

243 session.add(token) 

244 assert session.execute(select(func.count()).select_from(AccountDeletionToken)).scalar_one() == 3 

245 

246 queue_job(session, job=purge_account_deletion_tokens, payload=empty_pb2.Empty()) 

247 process_job() 

248 

249 with session_scope() as session: 

250 assert session.execute(select(func.count()).select_from(AccountDeletionToken)).scalar_one() == 1 

251 

252 with session_scope() as session: 

253 assert ( 

254 session.execute( 

255 select(func.count()) 

256 .select_from(BackgroundJob) 

257 .where(BackgroundJob.state == BackgroundJobState.completed) 

258 ).scalar_one() 

259 == 1 

260 ) 

261 assert ( 

262 session.execute( 

263 select(func.count()) 

264 .select_from(BackgroundJob) 

265 .where(BackgroundJob.state != BackgroundJobState.completed) 

266 ).scalar_one() 

267 == 0 

268 ) 

269 

270 

271def test_enforce_community_memberships(db): 

272 with session_scope() as session: 

273 queue_job(session, job=enforce_community_membership, payload=empty_pb2.Empty()) 

274 process_job() 

275 

276 with session_scope() as session: 

277 assert ( 

278 session.execute( 

279 select(func.count()) 

280 .select_from(BackgroundJob) 

281 .where(BackgroundJob.state == BackgroundJobState.completed) 

282 ).scalar_one() 

283 == 1 

284 ) 

285 assert ( 

286 session.execute( 

287 select(func.count()) 

288 .select_from(BackgroundJob) 

289 .where(BackgroundJob.state != BackgroundJobState.completed) 

290 ).scalar_one() 

291 == 0 

292 ) 

293 

294 

295def test_refresh_materialized_views(db): 

296 with session_scope() as session: 

297 queue_job(session, job=refresh_materialized_views, payload=empty_pb2.Empty()) 

298 

299 process_job() 

300 

301 with session_scope() as session: 

302 assert ( 

303 session.execute( 

304 select(func.count()) 

305 .select_from(BackgroundJob) 

306 .where(BackgroundJob.state == BackgroundJobState.completed) 

307 ).scalar_one() 

308 == 1 

309 ) 

310 assert ( 

311 session.execute( 

312 select(func.count()) 

313 .select_from(BackgroundJob) 

314 .where(BackgroundJob.state != BackgroundJobState.completed) 

315 ).scalar_one() 

316 == 0 

317 ) 

318 

319 

320def test_service_jobs(db): 

321 with session_scope() as session: 

322 queue_email( 

323 session, 

324 jobs_pb2.SendEmailPayload( 

325 sender_name="sender_name", 

326 sender_email="sender_email", 

327 recipient="recipient", 

328 subject="subject", 

329 plain="plain", 

330 html="html", 

331 ), 

332 ) 

333 

334 # we create this HitSleep exception here, and mock out the normal sleep(1) in the infinite loop to instead raise 

335 # this. that allows us to conveniently get out of the infinite loop and know we had no more jobs left 

336 class HitSleep(Exception): 

337 pass 

338 

339 # the mock `sleep` function that instead raises the aforementioned exception 

340 def raising_sleep(seconds): 

341 raise HitSleep() 

342 

343 with pytest.raises(HitSleep): 

344 with patch("couchers.jobs.worker.sleep", raising_sleep): 

345 service_jobs() 

346 

347 with session_scope() as session: 

348 assert ( 

349 session.execute( 

350 select(func.count()) 

351 .select_from(BackgroundJob) 

352 .where(BackgroundJob.state == BackgroundJobState.completed) 

353 ).scalar_one() 

354 == 1 

355 ) 

356 assert ( 

357 session.execute( 

358 select(func.count()) 

359 .select_from(BackgroundJob) 

360 .where(BackgroundJob.state != BackgroundJobState.completed) 

361 ).scalar_one() 

362 == 0 

363 ) 

364 

365 

366def test_scheduler(db, monkeypatch): 

367 def purge_login_tokens(payload: empty_pb2.Empty): 

368 return 

369 

370 def send_message_notifications(payload: empty_pb2.Empty): 

371 return 

372 

373 MOCK_JOBS = { 

374 "purge_login_tokens": Job(purge_login_tokens, timedelta(seconds=7)), 

375 "send_message_notifications": Job(send_message_notifications, timedelta(seconds=11)), 

376 } 

377 

378 current_time = 0 

379 end_time = 70 

380 

381 class EndOfTime(Exception): 

382 pass 

383 

384 def mock_monotonic(): 

385 return current_time 

386 

387 def mock_sleep(seconds): 

388 nonlocal current_time 

389 current_time += seconds 

390 if current_time > end_time: 

391 raise EndOfTime() 

392 

393 realized_schedule = [] 

394 

395 def mock_run_job_and_schedule(sched, job: Job[Any], frequency: timedelta) -> None: 

396 realized_schedule.append((current_time, job.name)) 

397 _run_job_and_schedule(sched, job, frequency) 

398 

399 monkeypatch.setattr(couchers.jobs.worker, "_run_job_and_schedule", mock_run_job_and_schedule) 

400 monkeypatch.setattr(couchers.jobs.worker, "JOBS", MOCK_JOBS) 

401 monkeypatch.setattr(couchers.jobs.worker, "monotonic", mock_monotonic) 

402 monkeypatch.setattr(couchers.jobs.worker, "sleep", mock_sleep) 

403 

404 with pytest.raises(EndOfTime): 

405 run_scheduler() 

406 

407 # Convert to job indices for comparison (to maintain test compatibility) 

408 job_order = ["purge_login_tokens", "send_message_notifications"] 

409 realized_schedule_indices = [(time, job_order.index(job_name)) for time, job_name in realized_schedule] 

410 

411 assert realized_schedule_indices == [ 

412 (0.0, 0), 

413 (0.0, 1), 

414 (7.0, 0), 

415 (11.0, 1), 

416 (14.0, 0), 

417 (21.0, 0), 

418 (22.0, 1), 

419 (28.0, 0), 

420 (33.0, 1), 

421 (35.0, 0), 

422 (42.0, 0), 

423 (44.0, 1), 

424 (49.0, 0), 

425 (55.0, 1), 

426 (56.0, 0), 

427 (63.0, 0), 

428 (66.0, 1), 

429 (70.0, 0), 

430 ] 

431 

432 with session_scope() as session: 

433 assert ( 

434 session.execute( 

435 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state == BackgroundJobState.pending) 

436 ).scalar_one() 

437 == 18 

438 ) 

439 assert ( 

440 session.execute( 

441 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state != BackgroundJobState.pending) 

442 ).scalar_one() 

443 == 0 

444 ) 

445 

446 

447def test_job_retry(db): 

448 called_count = 0 

449 

450 def mock_job(payload: empty_pb2.Empty) -> None: 

451 nonlocal called_count 

452 called_count += 1 

453 raise Exception() 

454 

455 with session_scope() as session: 

456 queue_job(session, job=mock_job, payload=empty_pb2.Empty()) 

457 

458 MOCK_JOBS: dict[str, Job[Any]] = { 

459 "mock_job": Job(mock_job), 

460 } 

461 # port 0 so parallel test runs don't fight over the port 

462 metrics_server = create_prometheus_server(port=0) 

463 

464 # if IN_TEST is true, then the bg worker will raise on exceptions 

465 config.IN_TEST = False 

466 

467 with patch("couchers.jobs.worker.JOBS", MOCK_JOBS): 

468 process_job() 

469 with session_scope() as session: 

470 assert ( 

471 session.execute( 

472 select(func.count()) 

473 .select_from(BackgroundJob) 

474 .where(BackgroundJob.state == BackgroundJobState.error) 

475 ).scalar_one() 

476 == 1 

477 ) 

478 assert ( 

479 session.execute( 

480 select(func.count()) 

481 .select_from(BackgroundJob) 

482 .where(BackgroundJob.state != BackgroundJobState.error) 

483 ).scalar_one() 

484 == 0 

485 ) 

486 

487 job = session.execute(select(BackgroundJob)).scalar_one() 

488 assert job.next_attempt_after > now() + timedelta(seconds=25) 

489 

490 job.next_attempt_after = func.now() 

491 process_job() 

492 with session_scope() as session: 

493 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now() 

494 process_job() 

495 with session_scope() as session: 

496 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now() 

497 process_job() 

498 with session_scope() as session: 

499 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now() 

500 process_job() 

501 

502 with session_scope() as session: 

503 assert ( 

504 session.execute( 

505 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state == BackgroundJobState.failed) 

506 ).scalar_one() 

507 == 1 

508 ) 

509 assert ( 

510 session.execute( 

511 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state != BackgroundJobState.failed) 

512 ).scalar_one() 

513 == 0 

514 ) 

515 

516 _check_job_counter(metrics_server.server_port, "mock_job", "error", "4", "Exception") 

517 _check_job_counter(metrics_server.server_port, "mock_job", "failed", "5", "Exception") 

518 metrics_server.shutdown() 

519 metrics_server.server_close() 

520 

521 

522def test_job_retry_backs_off_from_now_not_from_a_stale_next_attempt_after(db): 

523 def mock_job(payload: empty_pb2.Empty) -> None: 

524 raise Exception() 

525 

526 MOCK_JOBS: dict[str, Job[Any]] = {"mock_job": Job(mock_job)} 

527 

528 with session_scope() as session: 

529 queue_job(session, job=mock_job, payload=empty_pb2.Empty()) 

530 session.flush() 

531 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = now() - timedelta(hours=1) 

532 

533 config.IN_TEST = False 

534 

535 with patch("couchers.jobs.worker.JOBS", MOCK_JOBS): 

536 assert process_job() 

537 

538 with session_scope() as session: 

539 job = session.execute(select(BackgroundJob)).scalar_one() 

540 assert job.state == BackgroundJobState.error 

541 assert job.try_count == 1 

542 assert now() + timedelta(seconds=25) < job.next_attempt_after < now() + timedelta(seconds=35) 

543 

544 assert not process_job() 

545 

546 

547def test_job_queued_latency_excludes_retry_backoff(db): 

548 def mock_job(payload: empty_pb2.Empty) -> None: 

549 pass 

550 

551 MOCK_JOBS: dict[str, Job[Any]] = {"mock_job": Job(mock_job)} 

552 

553 with session_scope() as session: 

554 queue_job(session, job=mock_job, payload=empty_pb2.Empty()) 

555 session.flush() 

556 job = session.execute(select(BackgroundJob)).scalar_one() 

557 job.queued = now() - timedelta(hours=1) 

558 job.next_attempt_after = now() 

559 

560 before = _queued_seconds_sum() 

561 

562 with patch("couchers.jobs.worker.JOBS", MOCK_JOBS): 

563 assert process_job() 

564 

565 assert _queued_seconds_sum() - before < 60 

566 

567 

568def test_job_dequeue_steps_over_other_workers_jobs(db): 

569 handled = [] 

570 

571 def mock_job(payload: jobs_pb2.SendEmailPayload) -> None: 

572 handled.append(payload.subject) 

573 

574 MOCK_JOBS: dict[str, Job[Any]] = {"mock_job": Job(mock_job)} 

575 

576 with session_scope() as session: 

577 for subject in ["first", "second", "third"]: 

578 queue_job(session, job=mock_job, payload=jobs_pb2.SendEmailPayload(subject=subject)) 

579 session.flush() 

580 # queued in one transaction, so they'd otherwise all share a next_attempt_after and tie in the ordering 

581 for i, job in enumerate(session.execute(select(BackgroundJob).order_by(BackgroundJob.id)).scalars()): 

582 job.next_attempt_after = now() - timedelta(seconds=3 - i) 

583 

584 # another worker already finished "first" and committed 

585 with session_scope() as session: 

586 finished = ( 

587 session.execute(select(BackgroundJob).order_by(BackgroundJob.next_attempt_after).limit(1)).scalars().one() 

588 ) 

589 finished.state = BackgroundJobState.completed 

590 

591 with session_scope() as holder: 

592 # another worker is holding "second", mid-flight 

593 held = ( 

594 holder.execute( 

595 select(BackgroundJob) 

596 .where(BackgroundJob.ready_for_retry) 

597 .order_by(BackgroundJob.next_attempt_after) 

598 .limit(1) 

599 .with_for_update(skip_locked=True) 

600 ) 

601 .scalars() 

602 .one() 

603 ) 

604 assert held.payload == jobs_pb2.SendEmailPayload(subject="second").SerializeToString() 

605 

606 # the dequeue must run at READ COMMITTED: under a stricter isolation level it can't follow the update chain of 

607 # the row the other worker just completed, and aborts the whole transaction rather than stepping over it 

608 assert holder.execute(text("show transaction_isolation")).scalar_one() == "read committed" 

609 

610 with patch("couchers.jobs.worker.JOBS", MOCK_JOBS): 

611 assert process_job() 

612 

613 assert handled == ["third"] 

614 

615 

616def test_no_jobs_no_problem(db): 

617 with session_scope() as session: 

618 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0 

619 

620 assert not process_job() 

621 

622 with session_scope() as session: 

623 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0 

624 

625 

626def test_send_message_notifications_basic(db, moderator, timewarp: Timewarp): 

627 user1, token1 = generate_user() 

628 user2, token2 = generate_user() 

629 user3, token3 = generate_user() 

630 

631 make_friends(user1, user2) 

632 make_friends(user1, user3) 

633 make_friends(user2, user3) 

634 

635 send_message_notifications(empty_pb2.Empty()) 

636 process_jobs() 

637 

638 # should find no jobs, since there's no messages 

639 with session_scope() as session: 

640 assert ( 

641 session.execute( 

642 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

643 ).scalar_one() 

644 == 0 

645 ) 

646 

647 with conversations_session(token1) as c: 

648 group_chat_id1 = c.CreateGroupChat( 

649 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id, user3.id]) 

650 ).group_chat_id 

651 moderator.approve_group_chat(group_chat_id1) 

652 

653 with conversations_session(token1) as c: 

654 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 1")) 

655 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 2")) 

656 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 3")) 

657 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 4")) 

658 

659 with conversations_session(token3) as c: 

660 group_chat_id2 = c.CreateGroupChat( 

661 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id]) 

662 ).group_chat_id 

663 moderator.approve_group_chat(group_chat_id2) 

664 

665 with conversations_session(token3) as c: 

666 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id2, text="Test message 5")) 

667 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id2, text="Test message 6")) 

668 

669 send_message_notifications(empty_pb2.Empty()) 

670 process_jobs() 

671 

672 # no emails sent out 

673 with session_scope() as session: 

674 assert ( 

675 session.execute( 

676 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

677 ).scalar_one() 

678 == 0 

679 ) 

680 

681 timewarp.advance(MISSED_MESSAGES_DELAY) 

682 

683 # this should generate emails for both user2 and user3 

684 send_message_notifications(empty_pb2.Empty()) 

685 process_jobs() 

686 

687 with session_scope() as session: 

688 assert ( 

689 session.execute( 

690 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

691 ).scalar_one() 

692 == 2 

693 ) 

694 # delete them all 

695 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

696 

697 # shouldn't generate any more emails 

698 send_message_notifications(empty_pb2.Empty()) 

699 process_jobs() 

700 

701 with session_scope() as session: 

702 assert ( 

703 session.execute( 

704 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

705 ).scalar_one() 

706 == 0 

707 ) 

708 

709 

710def test_send_message_notifications_ignores_abandoned_subscription(db, moderator, timewarp: Timewarp): 

711 """ 

712 Rejoining a chat leaves the earlier subscription behind with its own last-seen state, and the job 

713 used to notify off that stale one even once the user had caught up on their newest subscription. 

714 """ 

715 user1, token1 = generate_user() 

716 user2, token2 = generate_user() 

717 user3, token3 = generate_user() 

718 

719 make_friends(user1, user2) 

720 make_friends(user1, user3) 

721 make_friends(user2, user3) 

722 

723 with conversations_session(token1) as c: 

724 group_chat_id = c.CreateGroupChat( 

725 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id, user3.id]) 

726 ).group_chat_id 

727 moderator.approve_group_chat(group_chat_id) 

728 

729 with conversations_session(token1) as c: 

730 # unread and un-notified for user2 when they get removed below 

731 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1")) 

732 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2")) 

733 c.RemoveGroupChatUser(conversations_pb2.RemoveGroupChatUserReq(group_chat_id=group_chat_id, user_id=user2.id)) 

734 c.InviteToGroupChat(conversations_pb2.InviteToGroupChatReq(group_chat_id=group_chat_id, user_id=user2.id)) 

735 

736 for token in [token2, token3]: 

737 with conversations_session(token) as c: 

738 latest_message_id = c.GetGroupChat( 

739 conversations_pb2.GetGroupChatReq(group_chat_id=group_chat_id) 

740 ).latest_message.message_id 

741 c.MarkLastSeenGroupChat( 

742 conversations_pb2.MarkLastSeenGroupChatReq( 

743 group_chat_id=group_chat_id, last_seen_message_id=latest_message_id 

744 ) 

745 ) 

746 

747 timewarp.advance(MISSED_MESSAGES_DELAY) 

748 

749 send_message_notifications(empty_pb2.Empty()) 

750 process_jobs() 

751 

752 with session_scope() as session: 

753 assert ( 

754 session.execute( 

755 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

756 ).scalar_one() 

757 == 0 

758 ) 

759 

760 

761def test_send_message_notifications_muted(db, moderator, timewarp: Timewarp): 

762 user1, token1 = generate_user() 

763 user2, token2 = generate_user() 

764 user3, token3 = generate_user() 

765 

766 make_friends(user1, user2) 

767 make_friends(user1, user3) 

768 make_friends(user2, user3) 

769 

770 send_message_notifications(empty_pb2.Empty()) 

771 process_jobs() 

772 

773 # should find no jobs, since there's no messages 

774 with session_scope() as session: 

775 assert ( 

776 session.execute( 

777 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

778 ).scalar_one() 

779 == 0 

780 ) 

781 

782 with conversations_session(token1) as c: 

783 group_chat_id = c.CreateGroupChat( 

784 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id, user3.id]) 

785 ).group_chat_id 

786 moderator.approve_group_chat(group_chat_id) 

787 

788 with conversations_session(token3) as c: 

789 # mute it for user 3 

790 c.MuteGroupChat(conversations_pb2.MuteGroupChatReq(group_chat_id=group_chat_id, forever=True)) 

791 

792 with conversations_session(token1) as c: 

793 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1")) 

794 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2")) 

795 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 3")) 

796 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 4")) 

797 

798 with conversations_session(token3) as c: 

799 group_chat_id = c.CreateGroupChat( 

800 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id]) 

801 ).group_chat_id 

802 moderator.approve_group_chat(group_chat_id) 

803 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 5")) 

804 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 6")) 

805 

806 send_message_notifications(empty_pb2.Empty()) 

807 process_jobs() 

808 

809 # no emails sent out 

810 with session_scope() as session: 

811 assert ( 

812 session.execute( 

813 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

814 ).scalar_one() 

815 == 0 

816 ) 

817 

818 timewarp.advance(MISSED_MESSAGES_DELAY) 

819 

820 # this should generate emails for both user2 and NOT user3 

821 send_message_notifications(empty_pb2.Empty()) 

822 process_jobs() 

823 

824 with session_scope() as session: 

825 assert ( 

826 session.execute( 

827 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

828 ).scalar_one() 

829 == 1 

830 ) 

831 # delete them all 

832 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

833 

834 # shouldn't generate any more emails 

835 send_message_notifications(empty_pb2.Empty()) 

836 process_jobs() 

837 

838 with session_scope() as session: 

839 assert ( 

840 session.execute( 

841 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

842 ).scalar_one() 

843 == 0 

844 ) 

845 

846 

847def _send_one_chat_message(token: str, recipient_id: int, moderator) -> None: 

848 with conversations_session(token) as c: 

849 group_chat_id = c.CreateGroupChat( 

850 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[recipient_id]) 

851 ).group_chat_id 

852 moderator.approve_group_chat(group_chat_id) 

853 

854 with conversations_session(token) as c: 

855 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1")) 

856 

857 with session_scope() as session: 

858 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

859 

860 

861def test_send_message_notifications_delayed_for_push_capable_users( 

862 db, moderator, push_collector: PushCollector, timewarp: Timewarp 

863): 

864 user1, token1 = generate_user() 

865 user2, token2 = generate_user() 

866 make_friends(user1, user2) 

867 _add_mobile_push_subscription(user2.id) 

868 

869 _send_one_chat_message(token1, user2.id, moderator) 

870 

871 # user2 already got a push about this, so the usual 5 minute delay doesn't apply to them 

872 timewarp.advance(MISSED_MESSAGES_DELAY) 

873 send_message_notifications(empty_pb2.Empty()) 

874 process_jobs() 

875 assert _count_queued_emails() == 0 

876 

877 timewarp.advance(MISSED_MESSAGES_DELAY_WITH_PUSH) 

878 send_message_notifications(empty_pb2.Empty()) 

879 process_jobs() 

880 assert _count_queued_emails() == 1 

881 

882 

883def test_send_message_notifications_not_delayed_when_push_subscription_disabled( 

884 db, moderator, push_collector: PushCollector, timewarp: Timewarp 

885): 

886 user1, token1 = generate_user() 

887 user2, token2 = generate_user() 

888 make_friends(user1, user2) 

889 # e.g. the device unregistered, so we can't reach user2 by push and the email is all they'll get 

890 _add_mobile_push_subscription(user2.id, disabled=True) 

891 

892 _send_one_chat_message(token1, user2.id, moderator) 

893 

894 timewarp.advance(MISSED_MESSAGES_DELAY) 

895 send_message_notifications(empty_pb2.Empty()) 

896 process_jobs() 

897 assert _count_queued_emails() == 1 

898 

899 

900def test_send_request_notifications_host_request(db, moderator, timewarp: Timewarp): 

901 user1, token1 = generate_user() 

902 user2, token2 = generate_user() 

903 

904 today_plus_2 = (today() + timedelta(days=2)).isoformat() 

905 today_plus_3 = (today() + timedelta(days=3)).isoformat() 

906 

907 send_request_notifications(empty_pb2.Empty()) 

908 process_jobs() 

909 

910 # should find no jobs, since there's no messages 

911 with session_scope() as session: 

912 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0 

913 

914 with requests_session(token1) as requests: 

915 host_request_id = requests.CreateHostRequest( 

916 requests_pb2.CreateHostRequestReq( 

917 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

918 ) 

919 ).host_request_id 

920 moderator.approve_host_request(host_request_id) 

921 

922 with session_scope() as session: 

923 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

924 

925 timewarp.advance(MISSED_MESSAGES_DELAY) 

926 

927 # the only unseen message is the creation message, which the host was already 

928 # notified about via host_request__create — no missed_messages email 

929 send_request_notifications(empty_pb2.Empty()) 

930 process_jobs() 

931 assert _count_queued_emails() == 0 

932 

933 # test that responding to host request creates email 

934 with requests_session(token2) as requests: 

935 requests.RespondHostRequest( 

936 requests_pb2.RespondHostRequestReq( 

937 host_request_id=host_request_id, 

938 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED, 

939 text="Test request", 

940 ) 

941 ) 

942 

943 with session_scope() as session: 

944 # delete send_email BackgroundJob created by RespondHostRequest 

945 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

946 

947 timewarp.advance(MISSED_MESSAGES_DELAY) 

948 

949 # check send_request_notifications successfully creates background job 

950 send_request_notifications(empty_pb2.Empty()) 

951 process_jobs() 

952 assert _count_queued_emails() == 1 

953 

954 with session_scope() as session: 

955 # delete all BackgroundJobs 

956 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

957 

958 send_request_notifications(empty_pb2.Empty()) 

959 process_jobs() 

960 # should find no messages since guest has already been notified 

961 assert _count_queued_emails() == 0 

962 

963 

964def test_send_request_notifications_host_request_with_followup(db, moderator, timewarp: Timewarp): 

965 """ 

966 When the surfer sends a follow-up message after creating the host request, 

967 the host should get a missed_messages notification (even though the initial 

968 creation message alone would be skipped). 

969 """ 

970 user1, token1 = generate_user() 

971 user2, token2 = generate_user() 

972 

973 today_plus_2 = (today() + timedelta(days=2)).isoformat() 

974 today_plus_3 = (today() + timedelta(days=3)).isoformat() 

975 

976 with requests_session(token1) as requests: 

977 host_request_id = requests.CreateHostRequest( 

978 requests_pb2.CreateHostRequestReq( 

979 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

980 ) 

981 ).host_request_id 

982 moderator.approve_host_request(host_request_id) 

983 

984 # surfer sends a follow-up message 

985 with requests_session(token1) as requests: 

986 requests.SendHostRequestMessage( 

987 requests_pb2.SendHostRequestMessageReq(host_request_id=host_request_id, text="Following up on my request!") 

988 ) 

989 

990 with session_scope() as session: 

991 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

992 

993 timewarp.advance(MISSED_MESSAGES_DELAY) 

994 

995 # now there are two unseen text messages for the host, so missed_messages should fire 

996 send_request_notifications(empty_pb2.Empty()) 

997 process_jobs() 

998 assert _count_queued_emails() == 1 

999 

1000 

1001def test_send_request_notifications_two_requests_one_with_followup(db, moderator, timewarp: Timewarp): 

1002 """ 

1003 A host (user2) receives two requests: first from user1 (with a follow-up message), 

1004 then from user3 (creation only). Because request B is created after request A's 

1005 follow-up, it has a higher message ID. If the background job processes B first and 

1006 advances last_notified_request_message_id past A's messages, one might expect A's 

1007 notification to be lost — but it isn't, because the query results are already 

1008 materialized before the loop begins. 

1009 """ 

1010 user1, token1 = generate_user() 

1011 user2, token2 = generate_user() 

1012 user3, token3 = generate_user() 

1013 

1014 today_plus_2 = (today() + timedelta(days=2)).isoformat() 

1015 today_plus_3 = (today() + timedelta(days=3)).isoformat() 

1016 

1017 # request A: user1 -> user2, with a follow-up 

1018 with requests_session(token1) as requests: 

1019 host_request_a = requests.CreateHostRequest( 

1020 requests_pb2.CreateHostRequestReq( 

1021 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

1022 ) 

1023 ).host_request_id 

1024 moderator.approve_host_request(host_request_a) 

1025 

1026 with requests_session(token1) as requests: 

1027 requests.SendHostRequestMessage( 

1028 requests_pb2.SendHostRequestMessageReq(host_request_id=host_request_a, text="Sorry, meant Tuesday night!") 

1029 ) 

1030 

1031 # request B: user3 -> user2, creation only (higher message IDs than A's follow-up) 

1032 with requests_session(token3) as requests: 

1033 host_request_b = requests.CreateHostRequest( 

1034 requests_pb2.CreateHostRequestReq( 

1035 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

1036 ) 

1037 ).host_request_id 

1038 moderator.approve_host_request(host_request_b) 

1039 

1040 with session_scope() as session: 

1041 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1042 

1043 timewarp.advance(MISSED_MESSAGES_DELAY) 

1044 

1045 # should get exactly 1 missed_messages email: for request A (has follow-up), 

1046 # not request B (creation only, skipped) 

1047 send_request_notifications(empty_pb2.Empty()) 

1048 process_jobs() 

1049 assert _count_queued_emails() == 1 

1050 

1051 

1052def test_send_request_notifications_delayed_for_push_capable_users( 

1053 db, moderator, push_collector: PushCollector, timewarp: Timewarp 

1054): 

1055 user1, token1 = generate_user() 

1056 user2, token2 = generate_user() 

1057 _add_mobile_push_subscription(user1.id) 

1058 

1059 today_plus_2 = (today() + timedelta(days=2)).isoformat() 

1060 today_plus_3 = (today() + timedelta(days=3)).isoformat() 

1061 

1062 with requests_session(token1) as requests: 

1063 host_request_id = requests.CreateHostRequest( 

1064 requests_pb2.CreateHostRequestReq( 

1065 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

1066 ) 

1067 ).host_request_id 

1068 moderator.approve_host_request(host_request_id) 

1069 

1070 with requests_session(token2) as requests: 

1071 requests.RespondHostRequest( 

1072 requests_pb2.RespondHostRequestReq( 

1073 host_request_id=host_request_id, 

1074 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED, 

1075 text="Test request", 

1076 ) 

1077 ) 

1078 

1079 with session_scope() as session: 

1080 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1081 

1082 timewarp.advance(MISSED_MESSAGES_DELAY) 

1083 send_request_notifications(empty_pb2.Empty()) 

1084 process_jobs() 

1085 assert _count_queued_emails() == 0 

1086 

1087 timewarp.advance(MISSED_MESSAGES_DELAY_WITH_PUSH) 

1088 send_request_notifications(empty_pb2.Empty()) 

1089 process_jobs() 

1090 assert _count_queued_emails() == 1 

1091 

1092 

1093def test_send_message_notifications_seen(db, moderator, timewarp: Timewarp): 

1094 user1, token1 = generate_user() 

1095 user2, token2 = generate_user() 

1096 

1097 make_friends(user1, user2) 

1098 

1099 send_message_notifications(empty_pb2.Empty()) 

1100 

1101 # should find no jobs, since there's no messages 

1102 with session_scope() as session: 

1103 assert ( 

1104 session.execute( 

1105 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1106 ).scalar_one() 

1107 == 0 

1108 ) 

1109 

1110 with conversations_session(token1) as c: 

1111 group_chat_id = c.CreateGroupChat( 

1112 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id]) 

1113 ).group_chat_id 

1114 moderator.approve_group_chat(group_chat_id) 

1115 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1")) 

1116 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2")) 

1117 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 3")) 

1118 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 4")) 

1119 

1120 # user 2 now marks those messages as seen 

1121 with conversations_session(token2) as c: 

1122 m_id = c.GetGroupChat(conversations_pb2.GetGroupChatReq(group_chat_id=group_chat_id)).latest_message.message_id 

1123 c.MarkLastSeenGroupChat( 

1124 conversations_pb2.MarkLastSeenGroupChatReq(group_chat_id=group_chat_id, last_seen_message_id=m_id) 

1125 ) 

1126 

1127 send_message_notifications(empty_pb2.Empty()) 

1128 

1129 # no emails sent out 

1130 with session_scope() as session: 

1131 assert ( 

1132 session.execute( 

1133 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1134 ).scalar_one() 

1135 == 0 

1136 ) 

1137 

1138 timewarp.advance(timedelta(minutes=30)) 

1139 

1140 # still shouldn't generate emails as user2 has seen all messages 

1141 send_message_notifications(empty_pb2.Empty()) 

1142 

1143 with session_scope() as session: 

1144 assert ( 

1145 session.execute( 

1146 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1147 ).scalar_one() 

1148 == 0 

1149 ) 

1150 

1151 

1152def test_send_onboarding_emails(db): 

1153 # needs to get first onboarding email 

1154 user1, token1 = generate_user(onboarding_emails_sent=0, last_onboarding_email_sent=None, complete_profile=False) 

1155 

1156 send_onboarding_emails(empty_pb2.Empty()) 

1157 process_jobs() 

1158 

1159 with session_scope() as session: 

1160 assert ( 

1161 session.execute( 

1162 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1163 ).scalar_one() 

1164 == 1 

1165 ) 

1166 

1167 # needs to get second onboarding email, but not yet 

1168 user2, token2 = generate_user( 

1169 onboarding_emails_sent=1, last_onboarding_email_sent=now() - timedelta(days=6), complete_profile=False 

1170 ) 

1171 

1172 send_onboarding_emails(empty_pb2.Empty()) 

1173 process_jobs() 

1174 

1175 with session_scope() as session: 

1176 assert ( 

1177 session.execute( 

1178 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1179 ).scalar_one() 

1180 == 1 

1181 ) 

1182 

1183 # needs to get second onboarding email 

1184 user3, token3 = generate_user( 

1185 onboarding_emails_sent=1, last_onboarding_email_sent=now() - timedelta(days=8), complete_profile=False 

1186 ) 

1187 

1188 send_onboarding_emails(empty_pb2.Empty()) 

1189 process_jobs() 

1190 

1191 with session_scope() as session: 

1192 assert ( 

1193 session.execute( 

1194 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email") 

1195 ).scalar_one() 

1196 == 2 

1197 ) 

1198 

1199 

1200def test_send_reference_reminders(db): 

1201 # need to test: 

1202 # case 1: bidirectional (no emails) 

1203 # case 2: host left ref (surfer needs an email) 

1204 # case 3: surfer left ref (host needs an email) 

1205 # case 4: neither left ref (host & surfer need an email) 

1206 # case 5: neither left ref, but host blocked surfer, so neither should get an email 

1207 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email) 

1208 

1209 send_reference_reminders(empty_pb2.Empty()) 

1210 

1211 # case 1: bidirectional (no emails) 

1212 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1") 

1213 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2") 

1214 

1215 # case 2: host left ref (surfer needs an email) 

1216 # host 

1217 user3, token3 = generate_user(email="user3@couchers.org.invalid", name="User 3") 

1218 # surfer 

1219 user4, token4 = generate_user(email="user4@couchers.org.invalid", name="User 4") 

1220 

1221 # case 3: surfer left ref (host needs an email) 

1222 # host 

1223 user5, token5 = generate_user(email="user5@couchers.org.invalid", name="User 5") 

1224 # surfer 

1225 user6, token6 = generate_user(email="user6@couchers.org.invalid", name="User 6") 

1226 

1227 # case 4: neither left ref (host & surfer need an email) 

1228 # surfer 

1229 user7, token7 = generate_user(email="user7@couchers.org.invalid", name="User 7") 

1230 # host 

1231 user8, token8 = generate_user(email="user8@couchers.org.invalid", name="User 8") 

1232 

1233 # case 5: neither left ref, but host blocked surfer, so neither should get an email 

1234 # surfer 

1235 user9, token9 = generate_user(email="user9@couchers.org.invalid", name="User 9") 

1236 # host 

1237 user10, token10 = generate_user(email="user10@couchers.org.invalid", name="User 10") 

1238 

1239 make_user_block(user9, user10) 

1240 

1241 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email) 

1242 # host 

1243 user11, token11 = generate_user(email="user11@couchers.org.invalid", name="User 11") 

1244 # surfer 

1245 user12, token12 = generate_user(email="user12@couchers.org.invalid", name="User 12") 

1246 

1247 with session_scope() as session: 

1248 # note that create_host_reference creates a host request whose age is one day older than the timedelta here 

1249 

1250 # case 1: bidirectional (no emails) 

1251 ref1, hr1 = create_host_reference(session, user2.id, user1.id, timedelta(days=7), surfing=True) 

1252 create_host_reference(session, user1.id, user2.id, timedelta(days=7), host_request_id=hr1) 

1253 

1254 # case 2: host left ref (surfer needs an email) 

1255 ref2, hr2 = create_host_reference(session, user3.id, user4.id, timedelta(days=11), surfing=False) 

1256 

1257 # case 3: surfer left ref (host needs an email) 

1258 ref3, hr3 = create_host_reference(session, user6.id, user5.id, timedelta(days=9), surfing=True) 

1259 

1260 # case 4: neither left ref (host & surfer need an email) 

1261 hr4 = create_host_request(session, user7.id, user8.id, timedelta(days=4)) 

1262 

1263 # case 5: neither left ref, but host blocked surfer, so neither should get an email 

1264 hr5 = create_host_request(session, user9.id, user10.id, timedelta(days=7)) 

1265 

1266 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email) 

1267 hr6 = create_host_request(session, user12.id, user11.id, timedelta(days=6), surfer_reason_didnt_meetup="") 

1268 

1269 expected_emails = [ 

1270 ( 

1271 "user11@couchers.org.invalid", 

1272 "[TEST] You have 14 days to write a reference for User 12", 

1273 ("from when you hosted them", "/leave-reference/hosted/"), 

1274 ), 

1275 ( 

1276 "user4@couchers.org.invalid", 

1277 "[TEST] You have 3 days to write a reference for User 3", 

1278 ("from when you surfed with them", "/leave-reference/surfed/"), 

1279 ), 

1280 ( 

1281 "user5@couchers.org.invalid", 

1282 "[TEST] You have 7 days to write a reference for User 6", 

1283 ("from when you hosted them", "/leave-reference/hosted/"), 

1284 ), 

1285 ( 

1286 "user7@couchers.org.invalid", 

1287 "[TEST] You have 14 days to write a reference for User 8", 

1288 ("from when you surfed with them", "/leave-reference/surfed/"), 

1289 ), 

1290 ( 

1291 "user8@couchers.org.invalid", 

1292 "[TEST] You have 14 days to write a reference for User 7", 

1293 ("from when you hosted them", "/leave-reference/hosted/"), 

1294 ), 

1295 ] 

1296 

1297 send_reference_reminders(empty_pb2.Empty()) 

1298 

1299 while process_job(): 

1300 pass 

1301 

1302 with session_scope() as session: 

1303 emails = [ 

1304 (email.recipient, email.subject, email.plain, email.html) 

1305 for email in session.execute(select(Email).order_by(Email.recipient.asc())).scalars().all() 

1306 ] 

1307 

1308 actual_addresses_and_subjects = [email[:2] for email in emails] 

1309 expected_addresses_and_subjects = [email[:2] for email in expected_emails] 

1310 

1311 print(actual_addresses_and_subjects) 

1312 print(expected_addresses_and_subjects) 

1313 

1314 assert actual_addresses_and_subjects == expected_addresses_and_subjects 

1315 

1316 for (address, subject, plain, html), (_, _, search_strings) in zip(emails, expected_emails): 

1317 for find in search_strings: 

1318 assert find in plain, f"Expected to find string {find} in PLAIN email {subject} to {address}, didn't" 

1319 assert find in html, f"Expected to find string {find} in HTML email {subject} to {address}, didn't" 

1320 

1321 

1322def test_send_reference_reminders_public_trip_offer(db): 

1323 """An offer reverses initiator/recipient, so the reminders have to be picked by the stay role.""" 

1324 surfer, _ = generate_user(email="surfer@couchers.org.invalid", name="Surfer") 

1325 host, _ = generate_user(email="host@couchers.org.invalid", name="Host") 

1326 

1327 with session_scope() as session: 

1328 create_host_request(session, surfer.id, host.id, timedelta(days=4), is_offer=True) 

1329 

1330 send_reference_reminders(empty_pb2.Empty()) 

1331 

1332 while process_job(): 

1333 pass 

1334 

1335 with session_scope() as session: 

1336 emails = { 

1337 email.recipient: (email.subject, email.plain) for email in session.execute(select(Email)).scalars().all() 

1338 } 

1339 

1340 assert set(emails) == {"surfer@couchers.org.invalid", "host@couchers.org.invalid"} 

1341 surfer_subject, surfer_plain = emails["surfer@couchers.org.invalid"] 

1342 assert surfer_subject == "[TEST] You have 14 days to write a reference for Host" 

1343 assert "from when you surfed with them" in surfer_plain 

1344 host_subject, host_plain = emails["host@couchers.org.invalid"] 

1345 assert host_subject == "[TEST] You have 14 days to write a reference for Surfer" 

1346 assert "from when you hosted them" in host_plain 

1347 

1348 

1349def test_send_host_request_reminders(db, moderator): 

1350 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1") 

1351 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2") 

1352 user3, token3 = generate_user(email="user3@couchers.org.invalid", name="User 3") 

1353 user4, token4 = generate_user(email="user4@couchers.org.invalid", name="User 4") 

1354 user5, token5 = generate_user(email="user5@couchers.org.invalid", name="User 5") 

1355 user6, token6 = generate_user(email="user6@couchers.org.invalid", name="User 6") 

1356 user7, token7 = generate_user(email="user7@couchers.org.invalid", name="User 7") 

1357 user8, token8 = generate_user(email="user8@couchers.org.invalid", name="User 8") 

1358 user9, token9 = generate_user(email="user9@couchers.org.invalid", name="User 9") 

1359 user10, token10 = generate_user(email="user10@couchers.org.invalid", name="User 10") 

1360 user11, token11 = generate_user(email="user11@couchers.org.invalid", name="User 11") 

1361 user12, token12 = generate_user(email="user12@couchers.org.invalid", name="User 12") 

1362 user13, token13 = generate_user(email="user13@couchers.org.invalid", name="User 13") 

1363 user14, token14 = generate_user(email="user14@couchers.org.invalid", name="User 14") 

1364 

1365 with session_scope() as session: 

1366 # case 1: pending, future, interval elapsed => notify 

1367 hr1 = create_host_request_by_date( 

1368 session=session, 

1369 surfer_user_id=user1.id, 

1370 host_user_id=user2.id, 

1371 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1), 

1372 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2), 

1373 status=HostRequestStatus.pending, 

1374 host_sent_request_reminders=0, 

1375 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1376 ) 

1377 

1378 # case 2: max reminders reached => do not notify 

1379 hr2 = create_host_request_by_date( 

1380 session=session, 

1381 surfer_user_id=user3.id, 

1382 host_user_id=user4.id, 

1383 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1), 

1384 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2), 

1385 status=HostRequestStatus.pending, 

1386 host_sent_request_reminders=HOST_REQUEST_MAX_REMINDERS, 

1387 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1388 ) 

1389 

1390 # case 3: interval not yet elapsed => do not notify 

1391 hr3 = create_host_request_by_date( 

1392 session=session, 

1393 surfer_user_id=user5.id, 

1394 host_user_id=user6.id, 

1395 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1), 

1396 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2), 

1397 status=HostRequestStatus.pending, 

1398 host_sent_request_reminders=0, 

1399 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL + timedelta(hours=1), 

1400 ) 

1401 

1402 # case 4: start date is today => do not notify 

1403 hr4 = create_host_request_by_date( 

1404 session=session, 

1405 surfer_user_id=user7.id, 

1406 host_user_id=user8.id, 

1407 from_date=today(), 

1408 to_date=today() + timedelta(days=2), 

1409 status=HostRequestStatus.pending, 

1410 host_sent_request_reminders=0, 

1411 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1412 ) 

1413 

1414 # case 5: from_date in the past => do not notify 

1415 hr5 = create_host_request_by_date( 

1416 session=session, 

1417 surfer_user_id=user9.id, 

1418 host_user_id=user10.id, 

1419 from_date=today() - timedelta(days=1), 

1420 to_date=today() + timedelta(days=1), 

1421 status=HostRequestStatus.pending, 

1422 host_sent_request_reminders=0, 

1423 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1424 ) 

1425 

1426 # case 6: non-pending status => do not notify 

1427 hr6 = create_host_request_by_date( 

1428 session=session, 

1429 surfer_user_id=user11.id, 

1430 host_user_id=user12.id, 

1431 from_date=today() + timedelta(days=3), 

1432 to_date=today() + timedelta(days=4), 

1433 status=HostRequestStatus.accepted, 

1434 host_sent_request_reminders=0, 

1435 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1436 ) 

1437 

1438 # case 7: host already sent a message => do not notify 

1439 hr7 = create_host_request_by_date( 

1440 session=session, 

1441 surfer_user_id=user13.id, 

1442 host_user_id=user14.id, 

1443 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1), 

1444 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2), 

1445 status=HostRequestStatus.pending, 

1446 host_sent_request_reminders=0, 

1447 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1448 ) 

1449 

1450 msg = Message( 

1451 conversation_id=hr7, 

1452 author_id=user14.id, 

1453 text="Looking forward to hosting you!", 

1454 message_type=MessageType.text, 

1455 ) 

1456 msg.time = now() 

1457 session.add(msg) 

1458 

1459 # Approve host requests so they're visible for notifications 

1460 moderator.approve_host_request(hr1) 

1461 moderator.approve_host_request(hr2) 

1462 moderator.approve_host_request(hr3) 

1463 moderator.approve_host_request(hr4) 

1464 moderator.approve_host_request(hr5) 

1465 moderator.approve_host_request(hr6) 

1466 moderator.approve_host_request(hr7) 

1467 

1468 send_host_request_reminders(empty_pb2.Empty()) 

1469 

1470 while process_job(): 

1471 pass 

1472 

1473 with session_scope() as session: 

1474 emails = [ 

1475 (email.recipient, email.subject, email.plain, email.html) 

1476 for email in session.execute(select(Email).order_by(Email.recipient.asc())).scalars().all() 

1477 ] 

1478 

1479 expected_emails = [ 

1480 ( 

1481 "user2@couchers.org.invalid", 

1482 "[TEST] You have a pending host request from User 1", 

1483 ("User 1", "is waiting for your response"), 

1484 ) 

1485 ] 

1486 

1487 actual_addresses_and_subjects = [email[:2] for email in emails] 

1488 expected_addresses_and_subjects = [email[:2] for email in expected_emails] 

1489 

1490 print(actual_addresses_and_subjects) 

1491 print(expected_addresses_and_subjects) 

1492 

1493 assert actual_addresses_and_subjects == expected_addresses_and_subjects 

1494 

1495 for (address, subject, plain, html), (_, _, search_strings) in zip(emails, expected_emails): 

1496 for find in search_strings: 

1497 assert find in plain, f"Expected to find string {find} in PLAIN email {subject} to {address}, didn't" 

1498 assert find in html, f"Expected to find string {find} in HTML email {subject} to {address}, didn't" 

1499 

1500 

1501def test_add_users_to_email_list(db, feature_flags): 

1502 feature_flags.set("listmonk_enabled", True) 

1503 config.LISTMONK_BASE_URL = "https://example.com" 

1504 config.LISTMONK_API_USERNAME = "test_user" 

1505 config.LISTMONK_API_KEY = "dummy_api_key" 

1506 config.LISTMONK_LIST_ID = 6 

1507 

1508 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls: 

1509 mock_session_cls.return_value.post.return_value.status_code = 200 

1510 add_users_to_email_list(empty_pb2.Empty()) 

1511 mock_session_cls.return_value.post.assert_not_called() 

1512 

1513 generate_user(in_sync_with_newsletter=False, email="testing1@couchers.invalid", name="Tester1", id=15) 

1514 generate_user(in_sync_with_newsletter=True, email="testing2@couchers.invalid", name="Tester2") 

1515 generate_user(in_sync_with_newsletter=False, email="testing3@couchers.invalid", name="Tester3 von test", id=17) 

1516 generate_user( 

1517 in_sync_with_newsletter=False, email="testing4@couchers.invalid", name="Tester4", opt_out_of_newsletter=True 

1518 ) 

1519 

1520 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls: 

1521 mock_sess = mock_session_cls.return_value 

1522 mock_sess.post.return_value.status_code = 200 

1523 add_users_to_email_list(empty_pb2.Empty()) 

1524 mock_sess.post.assert_has_calls( 

1525 [ 

1526 call( 

1527 "https://example.com/api/subscribers", 

1528 json={ 

1529 "email": "testing1@couchers.invalid", 

1530 "name": "Tester1", 

1531 "lists": [6], 

1532 "preconfirm_subscriptions": True, 

1533 "attribs": {"couchers_user_id": 15}, 

1534 "status": "enabled", 

1535 }, 

1536 timeout=10, 

1537 ), 

1538 call( 

1539 "https://example.com/api/subscribers", 

1540 json={ 

1541 "email": "testing3@couchers.invalid", 

1542 "name": "Tester3 von test", 

1543 "lists": [6], 

1544 "preconfirm_subscriptions": True, 

1545 "attribs": {"couchers_user_id": 17}, 

1546 "status": "enabled", 

1547 }, 

1548 timeout=10, 

1549 ), 

1550 ], 

1551 any_order=True, 

1552 ) 

1553 

1554 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls: 

1555 mock_session_cls.return_value.post.return_value.status_code = 200 

1556 add_users_to_email_list(empty_pb2.Empty()) 

1557 mock_session_cls.return_value.post.assert_not_called() 

1558 

1559 

1560def test_update_recommendation_scores(db): 

1561 update_recommendation_scores(empty_pb2.Empty()) 

1562 

1563 

1564def test_update_badges(db, push_collector: PushCollector): 

1565 user1, _ = generate_user(last_donated=None) 

1566 user2, _ = generate_user(last_donated=None) 

1567 user3, _ = generate_user(last_donated=None) 

1568 user4, _ = generate_user(phone="+15555555555", phone_verification_verified=func.now(), last_donated=None) 

1569 user5, _ = generate_user(phone="+15555555556", phone_verification_verified=func.now(), last_donated=None) 

1570 user6, _ = generate_user(last_donated=None) 

1571 

1572 with session_scope() as session: 

1573 session.add(UserBadge(user_id=user5.id, badge_id="board_member")) 

1574 session.add( 

1575 PostalVerificationAttempt( 

1576 user_id=user6.id, 

1577 address_line_1="123 Main St", 

1578 city="Test City", 

1579 country_code="US", 

1580 status=PostalVerificationStatus.succeeded, 

1581 verification_code="ABC123", 

1582 postcard_sent_at=func.now(), 

1583 verified_at=func.now(), 

1584 ) 

1585 ) 

1586 

1587 update_badges(empty_pb2.Empty()) 

1588 process_jobs() 

1589 

1590 with session_scope() as session: 

1591 badge_tuples = session.execute( 

1592 select(UserBadge.user_id, UserBadge.badge_id).order_by(UserBadge.user_id.asc(), UserBadge.id.asc()) 

1593 ).all() 

1594 

1595 expected = [ 

1596 (user1.id, "founder"), 

1597 (user1.id, "board_member"), 

1598 (user2.id, "founder"), 

1599 (user2.id, "board_member"), 

1600 (user4.id, "phone_verified"), 

1601 (user5.id, "phone_verified"), 

1602 (user6.id, "postal_verified"), 

1603 ] 

1604 

1605 assert badge_tuples == expected # type: ignore[comparison-overlap] 

1606 

1607 print(push_collector.by_user) 

1608 

1609 push = push_collector.pop_for_user(user1.id, last=False) 

1610 assert push.content.title == "New profile badge: Founder" 

1611 assert push.content.body == "The Founder badge was added to your profile." 

1612 

1613 push = push_collector.pop_for_user(user1.id, last=True) 

1614 assert push.content.title == "New profile badge: Board Member" 

1615 assert push.content.body == "The Board Member badge was added to your profile." 

1616 

1617 push = push_collector.pop_for_user(user2.id, last=False) 

1618 assert push.content.title == "New profile badge: Founder" 

1619 assert push.content.body == "The Founder badge was added to your profile." 

1620 

1621 push = push_collector.pop_for_user(user2.id, last=True) 

1622 assert push.content.title == "New profile badge: Board Member" 

1623 assert push.content.body == "The Board Member badge was added to your profile." 

1624 

1625 push = push_collector.pop_for_user(user4.id, last=True) 

1626 assert push.content.title == "New profile badge: Verified Phone" 

1627 assert push.content.body == "The Verified Phone badge was added to your profile." 

1628 

1629 push = push_collector.pop_for_user(user5.id, last=False) 

1630 assert push.content.title == "Profile badge removed" 

1631 assert push.content.body == "The Board Member badge was removed from your profile." 

1632 

1633 push = push_collector.pop_for_user(user5.id, last=True) 

1634 assert push.content.title == "New profile badge: Verified Phone" 

1635 assert push.content.body == "The Verified Phone badge was added to your profile." 

1636 

1637 

1638def test_update_badges_awards_moderator_to_superuser(db): 

1639 """The show_moderator_badge flag defaults on, so superusers are awarded the moderator badge.""" 

1640 superuser, _ = generate_user(is_superuser=True, last_donated=None) 

1641 

1642 update_badges(empty_pb2.Empty()) 

1643 

1644 with session_scope() as session: 

1645 assert ( 

1646 session.execute( 

1647 select(func.count()) 

1648 .select_from(UserBadge) 

1649 .where(UserBadge.user_id == superuser.id, UserBadge.badge_id == "moderator") 

1650 ).scalar() 

1651 == 1 

1652 ) 

1653 

1654 

1655def test_update_badges_skips_moderator_when_flag_off(db, monkeypatch): 

1656 """With show_moderator_badge forced off, superusers are not awarded the moderator badge.""" 

1657 # force show_moderator_badge off for everyone (force rule with no coverage applies globally) 

1658 monkeypatch.setattr(experimentation, "_initialized", True) 

1659 monkeypatch.setattr( 

1660 experimentation, 

1661 "_state", 

1662 {"features": {"show_moderator_badge": {"defaultValue": True, "rules": [{"force": False}]}}, "savedGroups": {}}, 

1663 ) 

1664 monkeypatch.setitem(config, "FEATURE_FLAGS_FILE_OVERRIDE_PATH", "") 

1665 

1666 superuser, _ = generate_user(is_superuser=True, last_donated=None) 

1667 

1668 update_badges(empty_pb2.Empty()) 

1669 

1670 with session_scope() as session: 

1671 assert ( 

1672 session.execute( 

1673 select(func.count()) 

1674 .select_from(UserBadge) 

1675 .where(UserBadge.user_id == superuser.id, UserBadge.badge_id == "moderator") 

1676 ).scalar() 

1677 == 0 

1678 ) 

1679 

1680 

1681def test_send_request_notifications_blocked_users_no_notification(db, moderator, timewarp: Timewarp): 

1682 """ 

1683 Regression test: send_request_notifications should not send notifications 

1684 when the host and surfer are not visible to each other (e.g., one blocked the other). 

1685 """ 

1686 user1, token1 = generate_user() 

1687 user2, token2 = generate_user() 

1688 

1689 today_plus_2 = (today() + timedelta(days=2)).isoformat() 

1690 today_plus_3 = (today() + timedelta(days=3)).isoformat() 

1691 

1692 # Create a host request 

1693 with requests_session(token1) as requests: 

1694 host_request_id = requests.CreateHostRequest( 

1695 requests_pb2.CreateHostRequestReq( 

1696 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text() 

1697 ) 

1698 ).host_request_id 

1699 moderator.approve_host_request(host_request_id) 

1700 

1701 with session_scope() as session: 

1702 # delete send_email BackgroundJob created by CreateHostRequest 

1703 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1704 

1705 # Now user2 (host) blocks user1 (surfer) 

1706 make_user_block(user2, user1) 

1707 

1708 timewarp.advance(MISSED_MESSAGES_DELAY) 

1709 

1710 # check send_request_notifications does NOT create background job because users are blocked 

1711 send_request_notifications(empty_pb2.Empty()) 

1712 process_jobs() 

1713 

1714 # Should be 0 emails because the host blocked the surfer 

1715 assert _count_queued_emails() == 0, "No notification email should be sent when host has blocked surfer" 

1716 

1717 # Also test the reverse direction: surfer sends message to host, host should not get notification 

1718 # First unblock 

1719 with session_scope() as session: 

1720 session.execute(delete(UserBlock).execution_options(synchronize_session=False)) 

1721 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1722 

1723 # Host responds 

1724 with requests_session(token2) as requests: 

1725 requests.RespondHostRequest( 

1726 requests_pb2.RespondHostRequestReq( 

1727 host_request_id=host_request_id, 

1728 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED, 

1729 text="Accepting your request", 

1730 ) 

1731 ) 

1732 

1733 with session_scope() as session: 

1734 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1735 

1736 # Now user1 (surfer) blocks user2 (host) 

1737 make_user_block(user1, user2) 

1738 

1739 timewarp.advance(MISSED_MESSAGES_DELAY) 

1740 

1741 # check send_request_notifications does NOT create background job 

1742 send_request_notifications(empty_pb2.Empty()) 

1743 process_jobs() 

1744 

1745 # Should be 0 emails because the surfer blocked the host 

1746 assert _count_queued_emails() == 0, "No notification email should be sent when surfer has blocked host" 

1747 

1748 

1749def test_send_host_request_reminders_blocked_users_no_notification(db, moderator): 

1750 """ 

1751 send_host_request_reminders should not send notifications when the host and surfer are not visible to each other 

1752 (e.g., one blocked the other). 

1753 """ 

1754 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1") 

1755 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2") 

1756 

1757 with session_scope() as session: 

1758 # Create a pending host request where the host has not replied 

1759 hr = create_host_request_by_date( 

1760 session=session, 

1761 surfer_user_id=user1.id, 

1762 host_user_id=user2.id, 

1763 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1), 

1764 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2), 

1765 status=HostRequestStatus.pending, 

1766 host_sent_request_reminders=0, 

1767 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL, 

1768 ) 

1769 

1770 # Approve the host request so it's visible for notifications 

1771 moderator.approve_host_request(hr) 

1772 

1773 # Verify that without blocking, a reminder would be sent 

1774 send_host_request_reminders(empty_pb2.Empty()) 

1775 

1776 while process_job(): 

1777 pass 

1778 

1779 with session_scope() as session: 

1780 emails = session.execute(select(Email)).scalars().all() 

1781 assert len(emails) == 1, "Expected 1 reminder email before blocking" 

1782 

1783 # Clean up emails and background jobs 

1784 session.execute(delete(Email).execution_options(synchronize_session=False)) 

1785 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1786 

1787 # Reset the reminder counter so we can test again 

1788 host_request = session.execute(select(HostRequest).where(HostRequest.conversation_id == hr)).scalar_one() 

1789 host_request.recipient_sent_request_reminders = 0 

1790 host_request.last_sent_request_reminder_time = now() - HOST_REQUEST_REMINDER_INTERVAL 

1791 

1792 # Now have the host block the surfer 

1793 make_user_block(user2, user1) 

1794 

1795 send_host_request_reminders(empty_pb2.Empty()) 

1796 

1797 while process_job(): 1797 ↛ 1798line 1797 didn't jump to line 1798 because the condition on line 1797 was never true

1798 pass 

1799 

1800 with session_scope() as session: 

1801 emails = session.execute(select(Email)).scalars().all() 

1802 assert len(emails) == 0, "No reminder email should be sent when host has blocked surfer" 

1803 

1804 

1805def test_send_message_notifications_blocked_users_no_notification(db, moderator, timewarp: Timewarp): 

1806 """ 

1807 Regression test: send_message_notifications should not send notifications 

1808 for messages from users who are blocked by the recipient. 

1809 """ 

1810 user1, token1 = generate_user() 

1811 user2, token2 = generate_user() 

1812 

1813 make_friends(user1, user2) 

1814 

1815 # Create a group chat and send messages 

1816 with conversations_session(token1) as c: 

1817 group_chat_id = c.CreateGroupChat( 

1818 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id]) 

1819 ).group_chat_id 

1820 

1821 # Approve the group chat so it's visible for notifications 

1822 moderator.approve_group_chat(group_chat_id) 

1823 

1824 with conversations_session(token1) as c: 

1825 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1")) 

1826 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2")) 

1827 

1828 # Verify that without blocking, a notification would be sent 

1829 with session_scope() as session: 

1830 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1831 

1832 timewarp.advance(MISSED_MESSAGES_DELAY) 

1833 

1834 send_message_notifications(empty_pb2.Empty()) 

1835 process_jobs() 

1836 

1837 assert _count_queued_emails() == 1, "Expected 1 notification email before blocking" 

1838 

1839 with session_scope() as session: 

1840 # Clean up 

1841 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1842 

1843 # Reset the notification state so user2 will receive notifications for old messages again 

1844 with session_scope() as session: 

1845 u2 = session.execute(select(User).where(User.id == user2.id)).scalar_one() 

1846 u2.last_notified_message_id = 0 

1847 

1848 # Now have user2 block user1 

1849 make_user_block(user2, user1) 

1850 

1851 # The existing messages from user1 should now NOT trigger notifications 

1852 # since user2 has blocked user1 

1853 with session_scope() as session: 

1854 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False)) 

1855 

1856 send_message_notifications(empty_pb2.Empty()) 

1857 process_jobs() 

1858 

1859 assert _count_queued_emails() == 0, "No notification email should be sent when recipient has blocked sender" 

1860 

1861 

1862def test_update_badges_volunteers(db, push_collector: PushCollector): 

1863 """Test that volunteer and past_volunteer badges are automatically granted based on Volunteer model.""" 

1864 # Create 6 users - users 1 and 2 get founder/board_member badges from static_badges 

1865 user1, _ = generate_user(last_donated=None) 

1866 user2, _ = generate_user(last_donated=None) 

1867 user3, _ = generate_user(last_donated=None) 

1868 user4, _ = generate_user(last_donated=None) 

1869 user5, _ = generate_user(last_donated=None) 

1870 user6, _ = generate_user(last_donated=None) 

1871 

1872 with session_scope() as session: 

1873 # user3: active volunteer (stopped_volunteering is null) 

1874 session.add( 

1875 make_volunteer( 

1876 user_id=user3.id, 

1877 role="Developer", 

1878 started_volunteering=date(2020, 1, 1), 

1879 stopped_volunteering=None, 

1880 ) 

1881 ) 

1882 

1883 # user4: past volunteer (stopped_volunteering is set) 

1884 session.add( 

1885 make_volunteer( 

1886 user_id=user4.id, 

1887 role="Designer", 

1888 started_volunteering=date(2020, 1, 1), 

1889 stopped_volunteering=date(2023, 6, 1), 

1890 ) 

1891 ) 

1892 

1893 # user5: has old volunteer badge that should be removed (not a volunteer anymore) 

1894 session.add(UserBadge(user_id=user5.id, badge_id="volunteer")) 

1895 

1896 # user6: has old past_volunteer badge that should be removed 

1897 session.add(UserBadge(user_id=user6.id, badge_id="past_volunteer")) 

1898 

1899 update_badges(empty_pb2.Empty()) 

1900 process_jobs() 

1901 

1902 with session_scope() as session: 

1903 # Check user3 has volunteer badge 

1904 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all() 

1905 assert "volunteer" in user3_badges 

1906 assert "past_volunteer" not in user3_badges 

1907 

1908 # Check user4 has past_volunteer badge 

1909 user4_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user4.id)).scalars().all() 

1910 assert "past_volunteer" in user4_badges 

1911 assert "volunteer" not in user4_badges 

1912 

1913 # Check user5 lost the volunteer badge (not in Volunteer table) 

1914 user5_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user5.id)).scalars().all() 

1915 assert "volunteer" not in user5_badges 

1916 

1917 # Check user6 lost the past_volunteer badge (not in Volunteer table) 

1918 user6_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user6.id)).scalars().all() 

1919 assert "past_volunteer" not in user6_badges 

1920 

1921 # Check notifications for volunteer badge users 

1922 push = push_collector.pop_for_user(user3.id, last=True) 

1923 assert push.content.title == "New profile badge: Active Volunteer" 

1924 assert push.content.body == "The Active Volunteer badge was added to your profile." 

1925 

1926 push = push_collector.pop_for_user(user4.id, last=True) 

1927 assert push.content.title == "New profile badge: Past Volunteer" 

1928 assert push.content.body == "The Past Volunteer badge was added to your profile." 

1929 

1930 push = push_collector.pop_for_user(user5.id, last=True) 

1931 assert push.content.title == "Profile badge removed" 

1932 assert push.content.body == "The Active Volunteer badge was removed from your profile." 

1933 

1934 push = push_collector.pop_for_user(user6.id, last=True) 

1935 assert push.content.title == "Profile badge removed" 

1936 assert push.content.body == "The Past Volunteer badge was removed from your profile." 

1937 

1938 

1939def test_update_badges_volunteer_status_change(db, push_collector: PushCollector): 

1940 """Test that badge is updated when volunteer status changes from active to past.""" 

1941 # Create users - users 1 and 2 get founder/board_member badges from static_badges 

1942 user1, _ = generate_user(last_donated=None) 

1943 user2, _ = generate_user(last_donated=None) 

1944 user3, _ = generate_user(last_donated=None) 

1945 

1946 with session_scope() as session: 

1947 # user3: start as active volunteer 

1948 session.add( 

1949 make_volunteer( 

1950 user_id=user3.id, 

1951 role="Developer", 

1952 started_volunteering=date(2020, 1, 1), 

1953 stopped_volunteering=None, 

1954 show_on_team_page=True, 

1955 ) 

1956 ) 

1957 

1958 update_badges(empty_pb2.Empty()) 

1959 process_jobs() 

1960 

1961 with session_scope() as session: 

1962 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all() 

1963 assert "volunteer" in user3_badges 

1964 assert "past_volunteer" not in user3_badges 

1965 

1966 push = push_collector.pop_for_user(user3.id, last=True) 

1967 assert push.content.title == "New profile badge: Active Volunteer" 

1968 assert push.content.body == "The Active Volunteer badge was added to your profile." 

1969 

1970 # Now change the volunteer to past volunteer 

1971 with session_scope() as session: 

1972 volunteer = session.execute(select(Volunteer).where(Volunteer.user_id == user3.id)).scalar_one() 

1973 volunteer.stopped_volunteering = date(2023, 12, 1) 

1974 

1975 update_badges(empty_pb2.Empty()) 

1976 process_jobs() 

1977 

1978 with session_scope() as session: 

1979 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all() 

1980 assert "volunteer" not in user3_badges 

1981 assert "past_volunteer" in user3_badges 

1982 

1983 # Check both badges were updated 

1984 push = push_collector.pop_for_user(user3.id, last=False) 

1985 assert push.content.title == "Profile badge removed" 

1986 assert push.content.body == "The Active Volunteer badge was removed from your profile." 

1987 

1988 push = push_collector.pop_for_user(user3.id, last=True) 

1989 assert push.content.title == "New profile badge: Past Volunteer" 

1990 assert push.content.body == "The Past Volunteer badge was added to your profile." 

1991 

1992 

1993def test_send_message_notifications_empty_unseen_simple(monkeypatch): 

1994 class DummyUser: 

1995 id = 1 

1996 is_visible = True 

1997 last_notified_message_id = 0 

1998 

1999 class FirstResult: 

2000 def scalars(self): 

2001 return self 

2002 

2003 def unique(self): 

2004 return [DummyUser()] 

2005 

2006 class SecondResult: 

2007 def all(self): 

2008 return [] 

2009 

2010 class DummySession: 

2011 def __init__(self): 

2012 self.calls = 0 

2013 

2014 def execute(self, *a, **k): 

2015 self.calls += 1 

2016 return FirstResult() if self.calls == 1 else SecondResult() 

2017 

2018 def commit(self): 

2019 pass 

2020 

2021 def flush(self): 

2022 pass 

2023 

2024 def fake_session_scope(): 

2025 class Ctx: 

2026 def __enter__(self): 

2027 return DummySession() 

2028 

2029 def __exit__(self, exc_type, exc, tb): 

2030 pass 

2031 

2032 return Ctx() 

2033 

2034 monkeypatch.setattr(handlers, "session_scope", fake_session_scope) 

2035 

2036 handlers.send_message_notifications(Empty())