Coverage for app/backend/src/tests/test_bg_jobs.py: 99%
838 statements
« prev ^ index » next coverage.py v7.15.3, created at 2026-08-04 22:32 +0000
« prev ^ index » next coverage.py v7.15.3, created at 2026-08-04 22:32 +0000
1from datetime import date, datetime, timedelta
2from typing import Any
3from unittest.mock import call, patch
5import pytest
6import requests
7from google.protobuf import empty_pb2
8from google.protobuf.empty_pb2 import Empty
9from sqlalchemy import select, text
10from sqlalchemy.sql import delete, func
12import couchers.jobs.worker
13from couchers import experimentation
14from couchers.config import config
15from couchers.constants import HOST_REQUEST_MAX_REMINDERS, HOST_REQUEST_REMINDER_INTERVAL
16from couchers.crypto import urlsafe_secure_token
17from couchers.db import session_scope
18from couchers.email.dev import print_dev_email
19from couchers.email.queuing import queue_email
20from couchers.jobs import handlers
21from couchers.jobs.definitions import Job
22from couchers.jobs.enqueue import queue_job
23from couchers.jobs.handlers import (
24 add_users_to_email_list,
25 enforce_community_membership,
26 purge_account_deletion_tokens,
27 purge_login_tokens,
28 purge_password_reset_tokens,
29 send_host_request_reminders,
30 send_message_notifications,
31 send_onboarding_emails,
32 send_reference_reminders,
33 send_request_notifications,
34 update_badges,
35 update_recommendation_scores,
36)
37from couchers.jobs.worker import _run_job_and_schedule, process_job, run_scheduler, service_jobs
38from couchers.materialized_views import refresh_materialized_views
39from couchers.metrics import create_prometheus_server
40from couchers.models import (
41 AccountDeletionToken,
42 BackgroundJob,
43 BackgroundJobState,
44 DeviceType,
45 Email,
46 HostRequest,
47 HostRequestStatus,
48 LoginToken,
49 Message,
50 MessageType,
51 PasswordResetToken,
52 PushNotificationPlatform,
53 PushNotificationSubscription,
54 User,
55 UserBadge,
56 UserBlock,
57 Volunteer,
58)
59from couchers.proto import conversations_pb2, messages_pb2, requests_pb2
60from couchers.proto.internal import jobs_pb2
61from couchers.utils import now, today
62from tests.fixtures.db import generate_user, make_friends, make_user_block, make_volunteer
63from tests.fixtures.misc import PushCollector, now_5_min_in_future, process_jobs
64from tests.fixtures.sessions import conversations_session, requests_session
65from tests.test_references import create_host_reference, create_host_request, create_host_request_by_date
66from tests.test_requests import valid_request_text
69def now_1_day_in_future() -> datetime:
70 return now() + timedelta(hours=25)
73def _add_mobile_push_subscription(user_id: int, *, disabled: bool = False) -> None:
74 with session_scope() as session:
75 sub = PushNotificationSubscription(
76 user_id=user_id,
77 platform=PushNotificationPlatform.expo,
78 token=f"ExponentPushToken[{user_id}]",
79 device_name="Test phone",
80 device_type=DeviceType.ios,
81 )
82 session.add(sub)
83 if disabled:
84 session.flush()
85 sub.disabled_at = now()
88def _count_queued_emails() -> int:
89 with session_scope() as session:
90 return session.execute(
91 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
92 ).scalar_one()
95@pytest.fixture(autouse=True)
96def _(testconfig):
97 pass
100def _check_job_counter(job, status, attempt, exception):
101 metrics_string = requests.get("http://localhost:8000").text
102 string_to_check = f'attempt="{attempt}",exception="{exception}",job="{job}",status="{status}"'
103 assert string_to_check in metrics_string
106def test_email_job(db):
107 with session_scope() as session:
108 queue_email(
109 session,
110 jobs_pb2.SendEmailPayload(
111 sender_name="sender_name",
112 sender_email="sender_email",
113 recipient="recipient",
114 subject="subject",
115 plain="plain",
116 html="html",
117 ),
118 )
120 def mock_print_dev_email(payload):
121 assert payload.sender_name == "sender_name"
122 assert payload.sender_email == "sender_email"
123 assert payload.recipient == "recipient"
124 assert payload.subject == "subject"
125 assert payload.plain == "plain"
126 assert payload.html == "html"
127 return print_dev_email(payload)
129 with patch("couchers.jobs.handlers.print_dev_email", mock_print_dev_email):
130 process_job()
132 with session_scope() as session:
133 assert (
134 session.execute(
135 select(func.count())
136 .select_from(BackgroundJob)
137 .where(BackgroundJob.state == BackgroundJobState.completed)
138 ).scalar_one()
139 == 1
140 )
141 assert (
142 session.execute(
143 select(func.count())
144 .select_from(BackgroundJob)
145 .where(BackgroundJob.state != BackgroundJobState.completed)
146 ).scalar_one()
147 == 0
148 )
151def test_purge_login_tokens(db):
152 user, api_token = generate_user()
154 with session_scope() as session:
155 login_token = LoginToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now())
156 session.add(login_token)
157 assert session.execute(select(func.count()).select_from(LoginToken)).scalar_one() == 1
159 queue_job(session, job=purge_login_tokens, payload=empty_pb2.Empty())
160 process_job()
162 with session_scope() as session:
163 assert session.execute(select(func.count()).select_from(LoginToken)).scalar_one() == 0
165 with session_scope() as session:
166 assert (
167 session.execute(
168 select(func.count())
169 .select_from(BackgroundJob)
170 .where(BackgroundJob.state == BackgroundJobState.completed)
171 ).scalar_one()
172 == 1
173 )
174 assert (
175 session.execute(
176 select(func.count())
177 .select_from(BackgroundJob)
178 .where(BackgroundJob.state != BackgroundJobState.completed)
179 ).scalar_one()
180 == 0
181 )
184def test_purge_password_reset_tokens(db):
185 user, api_token = generate_user()
187 with session_scope() as session:
188 password_reset_token = PasswordResetToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now())
189 session.add(password_reset_token)
190 assert session.execute(select(func.count()).select_from(PasswordResetToken)).scalar_one() == 1
192 queue_job(session, job=purge_password_reset_tokens, payload=empty_pb2.Empty())
193 process_job()
195 with session_scope() as session:
196 assert session.execute(select(func.count()).select_from(PasswordResetToken)).scalar_one() == 0
198 with session_scope() as session:
199 assert (
200 session.execute(
201 select(func.count())
202 .select_from(BackgroundJob)
203 .where(BackgroundJob.state == BackgroundJobState.completed)
204 ).scalar_one()
205 == 1
206 )
207 assert (
208 session.execute(
209 select(func.count())
210 .select_from(BackgroundJob)
211 .where(BackgroundJob.state != BackgroundJobState.completed)
212 ).scalar_one()
213 == 0
214 )
217def test_purge_account_deletion_tokens(db):
218 user, api_token = generate_user()
219 user2, api_token2 = generate_user()
220 user3, api_token3 = generate_user()
222 with session_scope() as session:
223 """
224 3 cases:
225 1) Token is valid
226 2) Token expired but account retrievable
227 3) Account is irretrievable (and expired)
228 """
229 account_deletion_tokens = [
230 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user.id, expiry=now() - timedelta(hours=2)),
231 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user2.id, expiry=now()),
232 AccountDeletionToken(token=urlsafe_secure_token(), user_id=user3.id, expiry=now() + timedelta(hours=5)),
233 ]
234 for token in account_deletion_tokens:
235 session.add(token)
236 assert session.execute(select(func.count()).select_from(AccountDeletionToken)).scalar_one() == 3
238 queue_job(session, job=purge_account_deletion_tokens, payload=empty_pb2.Empty())
239 process_job()
241 with session_scope() as session:
242 assert session.execute(select(func.count()).select_from(AccountDeletionToken)).scalar_one() == 1
244 with session_scope() as session:
245 assert (
246 session.execute(
247 select(func.count())
248 .select_from(BackgroundJob)
249 .where(BackgroundJob.state == BackgroundJobState.completed)
250 ).scalar_one()
251 == 1
252 )
253 assert (
254 session.execute(
255 select(func.count())
256 .select_from(BackgroundJob)
257 .where(BackgroundJob.state != BackgroundJobState.completed)
258 ).scalar_one()
259 == 0
260 )
263def test_enforce_community_memberships(db):
264 with session_scope() as session:
265 queue_job(session, job=enforce_community_membership, payload=empty_pb2.Empty())
266 process_job()
268 with session_scope() as session:
269 assert (
270 session.execute(
271 select(func.count())
272 .select_from(BackgroundJob)
273 .where(BackgroundJob.state == BackgroundJobState.completed)
274 ).scalar_one()
275 == 1
276 )
277 assert (
278 session.execute(
279 select(func.count())
280 .select_from(BackgroundJob)
281 .where(BackgroundJob.state != BackgroundJobState.completed)
282 ).scalar_one()
283 == 0
284 )
287def test_refresh_materialized_views(db):
288 with session_scope() as session:
289 queue_job(session, job=refresh_materialized_views, payload=empty_pb2.Empty())
291 process_job()
293 with session_scope() as session:
294 assert (
295 session.execute(
296 select(func.count())
297 .select_from(BackgroundJob)
298 .where(BackgroundJob.state == BackgroundJobState.completed)
299 ).scalar_one()
300 == 1
301 )
302 assert (
303 session.execute(
304 select(func.count())
305 .select_from(BackgroundJob)
306 .where(BackgroundJob.state != BackgroundJobState.completed)
307 ).scalar_one()
308 == 0
309 )
312def test_service_jobs(db):
313 with session_scope() as session:
314 queue_email(
315 session,
316 jobs_pb2.SendEmailPayload(
317 sender_name="sender_name",
318 sender_email="sender_email",
319 recipient="recipient",
320 subject="subject",
321 plain="plain",
322 html="html",
323 ),
324 )
326 # we create this HitSleep exception here, and mock out the normal sleep(1) in the infinite loop to instead raise
327 # this. that allows us to conveniently get out of the infinite loop and know we had no more jobs left
328 class HitSleep(Exception):
329 pass
331 # the mock `sleep` function that instead raises the aforementioned exception
332 def raising_sleep(seconds):
333 raise HitSleep()
335 with pytest.raises(HitSleep):
336 with patch("couchers.jobs.worker.sleep", raising_sleep):
337 service_jobs()
339 with session_scope() as session:
340 assert (
341 session.execute(
342 select(func.count())
343 .select_from(BackgroundJob)
344 .where(BackgroundJob.state == BackgroundJobState.completed)
345 ).scalar_one()
346 == 1
347 )
348 assert (
349 session.execute(
350 select(func.count())
351 .select_from(BackgroundJob)
352 .where(BackgroundJob.state != BackgroundJobState.completed)
353 ).scalar_one()
354 == 0
355 )
358def test_scheduler(db, monkeypatch):
359 def purge_login_tokens(payload: empty_pb2.Empty):
360 return
362 def send_message_notifications(payload: empty_pb2.Empty):
363 return
365 MOCK_JOBS = {
366 "purge_login_tokens": Job(purge_login_tokens, timedelta(seconds=7)),
367 "send_message_notifications": Job(send_message_notifications, timedelta(seconds=11)),
368 }
370 current_time = 0
371 end_time = 70
373 class EndOfTime(Exception):
374 pass
376 def mock_monotonic():
377 return current_time
379 def mock_sleep(seconds):
380 nonlocal current_time
381 current_time += seconds
382 if current_time > end_time:
383 raise EndOfTime()
385 realized_schedule = []
387 def mock_run_job_and_schedule(sched, job: Job[Any], frequency: timedelta) -> None:
388 realized_schedule.append((current_time, job.name))
389 _run_job_and_schedule(sched, job, frequency)
391 monkeypatch.setattr(couchers.jobs.worker, "_run_job_and_schedule", mock_run_job_and_schedule)
392 monkeypatch.setattr(couchers.jobs.worker, "JOBS", MOCK_JOBS)
393 monkeypatch.setattr(couchers.jobs.worker, "monotonic", mock_monotonic)
394 monkeypatch.setattr(couchers.jobs.worker, "sleep", mock_sleep)
396 with pytest.raises(EndOfTime):
397 run_scheduler()
399 # Convert to job indices for comparison (to maintain test compatibility)
400 job_order = ["purge_login_tokens", "send_message_notifications"]
401 realized_schedule_indices = [(time, job_order.index(job_name)) for time, job_name in realized_schedule]
403 assert realized_schedule_indices == [
404 (0.0, 0),
405 (0.0, 1),
406 (7.0, 0),
407 (11.0, 1),
408 (14.0, 0),
409 (21.0, 0),
410 (22.0, 1),
411 (28.0, 0),
412 (33.0, 1),
413 (35.0, 0),
414 (42.0, 0),
415 (44.0, 1),
416 (49.0, 0),
417 (55.0, 1),
418 (56.0, 0),
419 (63.0, 0),
420 (66.0, 1),
421 (70.0, 0),
422 ]
424 with session_scope() as session:
425 assert (
426 session.execute(
427 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state == BackgroundJobState.pending)
428 ).scalar_one()
429 == 18
430 )
431 assert (
432 session.execute(
433 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state != BackgroundJobState.pending)
434 ).scalar_one()
435 == 0
436 )
439def test_job_retry(db):
440 called_count = 0
442 def mock_job(payload: empty_pb2.Empty) -> None:
443 nonlocal called_count
444 called_count += 1
445 raise Exception()
447 with session_scope() as session:
448 queue_job(session, job=mock_job, payload=empty_pb2.Empty())
450 MOCK_JOBS: dict[str, Job[Any]] = {
451 "mock_job": Job(mock_job),
452 }
453 create_prometheus_server(port=8000)
455 # if IN_TEST is true, then the bg worker will raise on exceptions
456 new_config = config.copy()
457 new_config.IN_TEST = False
459 with patch("couchers.jobs.worker.config", new_config), patch("couchers.jobs.worker.JOBS", MOCK_JOBS):
460 process_job()
461 with session_scope() as session:
462 assert (
463 session.execute(
464 select(func.count())
465 .select_from(BackgroundJob)
466 .where(BackgroundJob.state == BackgroundJobState.error)
467 ).scalar_one()
468 == 1
469 )
470 assert (
471 session.execute(
472 select(func.count())
473 .select_from(BackgroundJob)
474 .where(BackgroundJob.state != BackgroundJobState.error)
475 ).scalar_one()
476 == 0
477 )
479 job = session.execute(select(BackgroundJob)).scalar_one()
480 assert job.next_attempt_after > now() + timedelta(seconds=25)
482 job.next_attempt_after = func.now()
483 process_job()
484 with session_scope() as session:
485 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now()
486 process_job()
487 with session_scope() as session:
488 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now()
489 process_job()
490 with session_scope() as session:
491 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = func.now()
492 process_job()
494 with session_scope() as session:
495 assert (
496 session.execute(
497 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state == BackgroundJobState.failed)
498 ).scalar_one()
499 == 1
500 )
501 assert (
502 session.execute(
503 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.state != BackgroundJobState.failed)
504 ).scalar_one()
505 == 0
506 )
508 _check_job_counter("mock_job", "error", "4", "Exception")
509 _check_job_counter("mock_job", "failed", "5", "Exception")
512def test_job_retry_backs_off_from_now_not_from_a_stale_next_attempt_after(db):
513 def mock_job(payload: empty_pb2.Empty) -> None:
514 raise Exception()
516 MOCK_JOBS: dict[str, Job[Any]] = {"mock_job": Job(mock_job)}
518 with session_scope() as session:
519 queue_job(session, job=mock_job, payload=empty_pb2.Empty())
520 session.flush()
521 session.execute(select(BackgroundJob)).scalar_one().next_attempt_after = now() - timedelta(hours=1)
523 new_config = config.copy()
524 new_config.IN_TEST = False
526 with patch("couchers.jobs.worker.config", new_config), patch("couchers.jobs.worker.JOBS", MOCK_JOBS):
527 assert process_job()
529 with session_scope() as session:
530 job = session.execute(select(BackgroundJob)).scalar_one()
531 assert job.state == BackgroundJobState.error
532 assert job.try_count == 1
533 assert now() + timedelta(seconds=25) < job.next_attempt_after < now() + timedelta(seconds=35)
535 assert not process_job()
538def test_job_dequeue_steps_over_other_workers_jobs(db):
539 handled = []
541 def mock_job(payload: jobs_pb2.SendEmailPayload) -> None:
542 handled.append(payload.subject)
544 MOCK_JOBS: dict[str, Job[Any]] = {"mock_job": Job(mock_job)}
546 with session_scope() as session:
547 for subject in ["first", "second", "third"]:
548 queue_job(session, job=mock_job, payload=jobs_pb2.SendEmailPayload(subject=subject))
549 session.flush()
550 # queued in one transaction, so they'd otherwise all share a next_attempt_after and tie in the ordering
551 for i, job in enumerate(session.execute(select(BackgroundJob).order_by(BackgroundJob.id)).scalars()):
552 job.next_attempt_after = now() - timedelta(seconds=3 - i)
554 # another worker already finished "first" and committed
555 with session_scope() as session:
556 finished = (
557 session.execute(select(BackgroundJob).order_by(BackgroundJob.next_attempt_after).limit(1)).scalars().one()
558 )
559 finished.state = BackgroundJobState.completed
561 with session_scope() as holder:
562 # another worker is holding "second", mid-flight
563 held = (
564 holder.execute(
565 select(BackgroundJob)
566 .where(BackgroundJob.ready_for_retry)
567 .order_by(BackgroundJob.next_attempt_after)
568 .limit(1)
569 .with_for_update(skip_locked=True)
570 )
571 .scalars()
572 .one()
573 )
574 assert held.payload == jobs_pb2.SendEmailPayload(subject="second").SerializeToString()
576 # the dequeue must run at READ COMMITTED: under a stricter isolation level it can't follow the update chain of
577 # the row the other worker just completed, and aborts the whole transaction rather than stepping over it
578 assert holder.execute(text("show transaction_isolation")).scalar_one() == "read committed"
580 with patch("couchers.jobs.worker.JOBS", MOCK_JOBS):
581 assert process_job()
583 assert handled == ["third"]
586def test_no_jobs_no_problem(db):
587 with session_scope() as session:
588 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0
590 assert not process_job()
592 with session_scope() as session:
593 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0
596def test_send_message_notifications_basic(db, moderator):
597 user1, token1 = generate_user()
598 user2, token2 = generate_user()
599 user3, token3 = generate_user()
601 make_friends(user1, user2)
602 make_friends(user1, user3)
603 make_friends(user2, user3)
605 send_message_notifications(empty_pb2.Empty())
606 process_jobs()
608 # should find no jobs, since there's no messages
609 with session_scope() as session:
610 assert (
611 session.execute(
612 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
613 ).scalar_one()
614 == 0
615 )
617 with conversations_session(token1) as c:
618 group_chat_id1 = c.CreateGroupChat(
619 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id, user3.id])
620 ).group_chat_id
621 moderator.approve_group_chat(group_chat_id1)
623 with conversations_session(token1) as c:
624 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 1"))
625 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 2"))
626 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 3"))
627 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id1, text="Test message 4"))
629 with conversations_session(token3) as c:
630 group_chat_id2 = c.CreateGroupChat(
631 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id])
632 ).group_chat_id
633 moderator.approve_group_chat(group_chat_id2)
635 with conversations_session(token3) as c:
636 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id2, text="Test message 5"))
637 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id2, text="Test message 6"))
639 send_message_notifications(empty_pb2.Empty())
640 process_jobs()
642 # no emails sent out
643 with session_scope() as session:
644 assert (
645 session.execute(
646 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
647 ).scalar_one()
648 == 0
649 )
651 # this should generate emails for both user2 and user3
652 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
653 send_message_notifications(empty_pb2.Empty())
654 process_jobs()
656 with session_scope() as session:
657 assert (
658 session.execute(
659 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
660 ).scalar_one()
661 == 2
662 )
663 # delete them all
664 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
666 # shouldn't generate any more emails
667 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
668 send_message_notifications(empty_pb2.Empty())
669 process_jobs()
671 with session_scope() as session:
672 assert (
673 session.execute(
674 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
675 ).scalar_one()
676 == 0
677 )
680def test_send_message_notifications_muted(db, moderator):
681 user1, token1 = generate_user()
682 user2, token2 = generate_user()
683 user3, token3 = generate_user()
685 make_friends(user1, user2)
686 make_friends(user1, user3)
687 make_friends(user2, user3)
689 send_message_notifications(empty_pb2.Empty())
690 process_jobs()
692 # should find no jobs, since there's no messages
693 with session_scope() as session:
694 assert (
695 session.execute(
696 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
697 ).scalar_one()
698 == 0
699 )
701 with conversations_session(token1) as c:
702 group_chat_id = c.CreateGroupChat(
703 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id, user3.id])
704 ).group_chat_id
705 moderator.approve_group_chat(group_chat_id)
707 with conversations_session(token3) as c:
708 # mute it for user 3
709 c.MuteGroupChat(conversations_pb2.MuteGroupChatReq(group_chat_id=group_chat_id, forever=True))
711 with conversations_session(token1) as c:
712 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1"))
713 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2"))
714 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 3"))
715 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 4"))
717 with conversations_session(token3) as c:
718 group_chat_id = c.CreateGroupChat(
719 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id])
720 ).group_chat_id
721 moderator.approve_group_chat(group_chat_id)
722 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 5"))
723 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 6"))
725 send_message_notifications(empty_pb2.Empty())
726 process_jobs()
728 # no emails sent out
729 with session_scope() as session:
730 assert (
731 session.execute(
732 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
733 ).scalar_one()
734 == 0
735 )
737 # this should generate emails for both user2 and NOT user3
738 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
739 send_message_notifications(empty_pb2.Empty())
740 process_jobs()
742 with session_scope() as session:
743 assert (
744 session.execute(
745 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
746 ).scalar_one()
747 == 1
748 )
749 # delete them all
750 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
752 # shouldn't generate any more emails
753 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
754 send_message_notifications(empty_pb2.Empty())
755 process_jobs()
757 with session_scope() as session:
758 assert (
759 session.execute(
760 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
761 ).scalar_one()
762 == 0
763 )
766def _send_one_chat_message(token: str, recipient_id: int, moderator) -> None:
767 with conversations_session(token) as c:
768 group_chat_id = c.CreateGroupChat(
769 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[recipient_id])
770 ).group_chat_id
771 moderator.approve_group_chat(group_chat_id)
773 with conversations_session(token) as c:
774 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1"))
776 with session_scope() as session:
777 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
780def test_send_message_notifications_delayed_for_push_capable_users(db, moderator, push_collector: PushCollector):
781 user1, token1 = generate_user()
782 user2, token2 = generate_user()
783 make_friends(user1, user2)
784 _add_mobile_push_subscription(user2.id)
786 _send_one_chat_message(token1, user2.id, moderator)
788 # user2 already got a push about this, so the usual 5 minute delay doesn't apply to them
789 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
790 send_message_notifications(empty_pb2.Empty())
791 process_jobs()
792 assert _count_queued_emails() == 0
794 with patch("couchers.jobs.handlers.now", now_1_day_in_future):
795 send_message_notifications(empty_pb2.Empty())
796 process_jobs()
797 assert _count_queued_emails() == 1
800def test_send_message_notifications_not_delayed_when_push_subscription_disabled(
801 db, moderator, push_collector: PushCollector
802):
803 user1, token1 = generate_user()
804 user2, token2 = generate_user()
805 make_friends(user1, user2)
806 # e.g. the device unregistered, so we can't reach user2 by push and the email is all they'll get
807 _add_mobile_push_subscription(user2.id, disabled=True)
809 _send_one_chat_message(token1, user2.id, moderator)
811 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
812 send_message_notifications(empty_pb2.Empty())
813 process_jobs()
814 assert _count_queued_emails() == 1
817def test_send_request_notifications_host_request(db, moderator):
818 user1, token1 = generate_user()
819 user2, token2 = generate_user()
821 today_plus_2 = (today() + timedelta(days=2)).isoformat()
822 today_plus_3 = (today() + timedelta(days=3)).isoformat()
824 send_request_notifications(empty_pb2.Empty())
825 process_jobs()
827 # should find no jobs, since there's no messages
828 with session_scope() as session:
829 assert session.execute(select(func.count()).select_from(BackgroundJob)).scalar_one() == 0
831 with requests_session(token1) as requests:
832 host_request_id = requests.CreateHostRequest(
833 requests_pb2.CreateHostRequestReq(
834 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
835 )
836 ).host_request_id
837 moderator.approve_host_request(host_request_id)
839 with session_scope() as session:
840 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
842 # the only unseen message is the creation message, which the host was already
843 # notified about via host_request__create — no missed_messages email
844 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
845 send_request_notifications(empty_pb2.Empty())
846 process_jobs()
847 assert (
848 session.execute(
849 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
850 ).scalar_one()
851 == 0
852 )
854 # test that responding to host request creates email
855 with requests_session(token2) as requests:
856 requests.RespondHostRequest(
857 requests_pb2.RespondHostRequestReq(
858 host_request_id=host_request_id,
859 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED,
860 text="Test request",
861 )
862 )
864 with session_scope() as session:
865 # delete send_email BackgroundJob created by RespondHostRequest
866 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
868 # check send_request_notifications successfully creates background job
869 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
870 send_request_notifications(empty_pb2.Empty())
871 process_jobs()
872 assert (
873 session.execute(
874 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
875 ).scalar_one()
876 == 1
877 )
879 # delete all BackgroundJobs
880 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
882 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
883 send_request_notifications(empty_pb2.Empty())
884 process_jobs()
885 # should find no messages since guest has already been notified
886 assert (
887 session.execute(
888 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
889 ).scalar_one()
890 == 0
891 )
894def test_send_request_notifications_host_request_with_followup(db, moderator):
895 """
896 When the surfer sends a follow-up message after creating the host request,
897 the host should get a missed_messages notification (even though the initial
898 creation message alone would be skipped).
899 """
900 user1, token1 = generate_user()
901 user2, token2 = generate_user()
903 today_plus_2 = (today() + timedelta(days=2)).isoformat()
904 today_plus_3 = (today() + timedelta(days=3)).isoformat()
906 with requests_session(token1) as requests:
907 host_request_id = requests.CreateHostRequest(
908 requests_pb2.CreateHostRequestReq(
909 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
910 )
911 ).host_request_id
912 moderator.approve_host_request(host_request_id)
914 # surfer sends a follow-up message
915 with requests_session(token1) as requests:
916 requests.SendHostRequestMessage(
917 requests_pb2.SendHostRequestMessageReq(host_request_id=host_request_id, text="Following up on my request!")
918 )
920 with session_scope() as session:
921 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
923 # now there are two unseen text messages for the host, so missed_messages should fire
924 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
925 send_request_notifications(empty_pb2.Empty())
926 process_jobs()
927 assert (
928 session.execute(
929 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
930 ).scalar_one()
931 == 1
932 )
935def test_send_request_notifications_two_requests_one_with_followup(db, moderator):
936 """
937 A host (user2) receives two requests: first from user1 (with a follow-up message),
938 then from user3 (creation only). Because request B is created after request A's
939 follow-up, it has a higher message ID. If the background job processes B first and
940 advances last_notified_request_message_id past A's messages, one might expect A's
941 notification to be lost — but it isn't, because the query results are already
942 materialized before the loop begins.
943 """
944 user1, token1 = generate_user()
945 user2, token2 = generate_user()
946 user3, token3 = generate_user()
948 today_plus_2 = (today() + timedelta(days=2)).isoformat()
949 today_plus_3 = (today() + timedelta(days=3)).isoformat()
951 # request A: user1 -> user2, with a follow-up
952 with requests_session(token1) as requests:
953 host_request_a = requests.CreateHostRequest(
954 requests_pb2.CreateHostRequestReq(
955 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
956 )
957 ).host_request_id
958 moderator.approve_host_request(host_request_a)
960 with requests_session(token1) as requests:
961 requests.SendHostRequestMessage(
962 requests_pb2.SendHostRequestMessageReq(host_request_id=host_request_a, text="Sorry, meant Tuesday night!")
963 )
965 # request B: user3 -> user2, creation only (higher message IDs than A's follow-up)
966 with requests_session(token3) as requests:
967 host_request_b = requests.CreateHostRequest(
968 requests_pb2.CreateHostRequestReq(
969 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
970 )
971 ).host_request_id
972 moderator.approve_host_request(host_request_b)
974 with session_scope() as session:
975 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
977 # should get exactly 1 missed_messages email: for request A (has follow-up),
978 # not request B (creation only, skipped)
979 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
980 send_request_notifications(empty_pb2.Empty())
981 process_jobs()
982 assert (
983 session.execute(
984 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
985 ).scalar_one()
986 == 1
987 )
990def test_send_request_notifications_delayed_for_push_capable_users(db, moderator, push_collector: PushCollector):
991 user1, token1 = generate_user()
992 user2, token2 = generate_user()
993 _add_mobile_push_subscription(user1.id)
995 today_plus_2 = (today() + timedelta(days=2)).isoformat()
996 today_plus_3 = (today() + timedelta(days=3)).isoformat()
998 with requests_session(token1) as requests:
999 host_request_id = requests.CreateHostRequest(
1000 requests_pb2.CreateHostRequestReq(
1001 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
1002 )
1003 ).host_request_id
1004 moderator.approve_host_request(host_request_id)
1006 with requests_session(token2) as requests:
1007 requests.RespondHostRequest(
1008 requests_pb2.RespondHostRequestReq(
1009 host_request_id=host_request_id,
1010 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED,
1011 text="Test request",
1012 )
1013 )
1015 with session_scope() as session:
1016 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1018 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
1019 send_request_notifications(empty_pb2.Empty())
1020 process_jobs()
1021 assert _count_queued_emails() == 0
1023 with patch("couchers.jobs.handlers.now", now_1_day_in_future):
1024 send_request_notifications(empty_pb2.Empty())
1025 process_jobs()
1026 assert _count_queued_emails() == 1
1029def test_send_message_notifications_seen(db, moderator):
1030 user1, token1 = generate_user()
1031 user2, token2 = generate_user()
1033 make_friends(user1, user2)
1035 send_message_notifications(empty_pb2.Empty())
1037 # should find no jobs, since there's no messages
1038 with session_scope() as session:
1039 assert (
1040 session.execute(
1041 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1042 ).scalar_one()
1043 == 0
1044 )
1046 with conversations_session(token1) as c:
1047 group_chat_id = c.CreateGroupChat(
1048 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id])
1049 ).group_chat_id
1050 moderator.approve_group_chat(group_chat_id)
1051 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1"))
1052 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2"))
1053 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 3"))
1054 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 4"))
1056 # user 2 now marks those messages as seen
1057 with conversations_session(token2) as c:
1058 m_id = c.GetGroupChat(conversations_pb2.GetGroupChatReq(group_chat_id=group_chat_id)).latest_message.message_id
1059 c.MarkLastSeenGroupChat(
1060 conversations_pb2.MarkLastSeenGroupChatReq(group_chat_id=group_chat_id, last_seen_message_id=m_id)
1061 )
1063 send_message_notifications(empty_pb2.Empty())
1065 # no emails sent out
1066 with session_scope() as session:
1067 assert (
1068 session.execute(
1069 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1070 ).scalar_one()
1071 == 0
1072 )
1074 def now_30_min_in_future():
1075 return now() + timedelta(minutes=30)
1077 # still shouldn't generate emails as user2 has seen all messages
1078 with patch("couchers.jobs.handlers.now", now_30_min_in_future):
1079 send_message_notifications(empty_pb2.Empty())
1081 with session_scope() as session:
1082 assert (
1083 session.execute(
1084 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1085 ).scalar_one()
1086 == 0
1087 )
1090def test_send_onboarding_emails(db):
1091 # needs to get first onboarding email
1092 user1, token1 = generate_user(onboarding_emails_sent=0, last_onboarding_email_sent=None, complete_profile=False)
1094 send_onboarding_emails(empty_pb2.Empty())
1095 process_jobs()
1097 with session_scope() as session:
1098 assert (
1099 session.execute(
1100 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1101 ).scalar_one()
1102 == 1
1103 )
1105 # needs to get second onboarding email, but not yet
1106 user2, token2 = generate_user(
1107 onboarding_emails_sent=1, last_onboarding_email_sent=now() - timedelta(days=6), complete_profile=False
1108 )
1110 send_onboarding_emails(empty_pb2.Empty())
1111 process_jobs()
1113 with session_scope() as session:
1114 assert (
1115 session.execute(
1116 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1117 ).scalar_one()
1118 == 1
1119 )
1121 # needs to get second onboarding email
1122 user3, token3 = generate_user(
1123 onboarding_emails_sent=1, last_onboarding_email_sent=now() - timedelta(days=8), complete_profile=False
1124 )
1126 send_onboarding_emails(empty_pb2.Empty())
1127 process_jobs()
1129 with session_scope() as session:
1130 assert (
1131 session.execute(
1132 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1133 ).scalar_one()
1134 == 2
1135 )
1138def test_send_reference_reminders(db):
1139 # need to test:
1140 # case 1: bidirectional (no emails)
1141 # case 2: host left ref (surfer needs an email)
1142 # case 3: surfer left ref (host needs an email)
1143 # case 4: neither left ref (host & surfer need an email)
1144 # case 5: neither left ref, but host blocked surfer, so neither should get an email
1145 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email)
1147 send_reference_reminders(empty_pb2.Empty())
1149 # case 1: bidirectional (no emails)
1150 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1")
1151 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2")
1153 # case 2: host left ref (surfer needs an email)
1154 # host
1155 user3, token3 = generate_user(email="user3@couchers.org.invalid", name="User 3")
1156 # surfer
1157 user4, token4 = generate_user(email="user4@couchers.org.invalid", name="User 4")
1159 # case 3: surfer left ref (host needs an email)
1160 # host
1161 user5, token5 = generate_user(email="user5@couchers.org.invalid", name="User 5")
1162 # surfer
1163 user6, token6 = generate_user(email="user6@couchers.org.invalid", name="User 6")
1165 # case 4: neither left ref (host & surfer need an email)
1166 # surfer
1167 user7, token7 = generate_user(email="user7@couchers.org.invalid", name="User 7")
1168 # host
1169 user8, token8 = generate_user(email="user8@couchers.org.invalid", name="User 8")
1171 # case 5: neither left ref, but host blocked surfer, so neither should get an email
1172 # surfer
1173 user9, token9 = generate_user(email="user9@couchers.org.invalid", name="User 9")
1174 # host
1175 user10, token10 = generate_user(email="user10@couchers.org.invalid", name="User 10")
1177 make_user_block(user9, user10)
1179 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email)
1180 # host
1181 user11, token11 = generate_user(email="user11@couchers.org.invalid", name="User 11")
1182 # surfer
1183 user12, token12 = generate_user(email="user12@couchers.org.invalid", name="User 12")
1185 with session_scope() as session:
1186 # note that create_host_reference creates a host request whose age is one day older than the timedelta here
1188 # case 1: bidirectional (no emails)
1189 ref1, hr1 = create_host_reference(session, user2.id, user1.id, timedelta(days=7), surfing=True)
1190 create_host_reference(session, user1.id, user2.id, timedelta(days=7), host_request_id=hr1)
1192 # case 2: host left ref (surfer needs an email)
1193 ref2, hr2 = create_host_reference(session, user3.id, user4.id, timedelta(days=11), surfing=False)
1195 # case 3: surfer left ref (host needs an email)
1196 ref3, hr3 = create_host_reference(session, user6.id, user5.id, timedelta(days=9), surfing=True)
1198 # case 4: neither left ref (host & surfer need an email)
1199 hr4 = create_host_request(session, user7.id, user8.id, timedelta(days=4))
1201 # case 5: neither left ref, but host blocked surfer, so neither should get an email
1202 hr5 = create_host_request(session, user9.id, user10.id, timedelta(days=7))
1204 # case 6: neither left ref, surfer indicated they didn't meet up, (host still needs an email)
1205 hr6 = create_host_request(session, user12.id, user11.id, timedelta(days=6), surfer_reason_didnt_meetup="")
1207 expected_emails = [
1208 (
1209 "user11@couchers.org.invalid",
1210 "[TEST] You have 14 days to write a reference for User 12",
1211 ("from when you hosted them", "/leave-reference/hosted/"),
1212 ),
1213 (
1214 "user4@couchers.org.invalid",
1215 "[TEST] You have 3 days to write a reference for User 3",
1216 ("from when you surfed with them", "/leave-reference/surfed/"),
1217 ),
1218 (
1219 "user5@couchers.org.invalid",
1220 "[TEST] You have 7 days to write a reference for User 6",
1221 ("from when you hosted them", "/leave-reference/hosted/"),
1222 ),
1223 (
1224 "user7@couchers.org.invalid",
1225 "[TEST] You have 14 days to write a reference for User 8",
1226 ("from when you surfed with them", "/leave-reference/surfed/"),
1227 ),
1228 (
1229 "user8@couchers.org.invalid",
1230 "[TEST] You have 14 days to write a reference for User 7",
1231 ("from when you hosted them", "/leave-reference/hosted/"),
1232 ),
1233 ]
1235 send_reference_reminders(empty_pb2.Empty())
1237 while process_job():
1238 pass
1240 with session_scope() as session:
1241 emails = [
1242 (email.recipient, email.subject, email.plain, email.html)
1243 for email in session.execute(select(Email).order_by(Email.recipient.asc())).scalars().all()
1244 ]
1246 actual_addresses_and_subjects = [email[:2] for email in emails]
1247 expected_addresses_and_subjects = [email[:2] for email in expected_emails]
1249 print(actual_addresses_and_subjects)
1250 print(expected_addresses_and_subjects)
1252 assert actual_addresses_and_subjects == expected_addresses_and_subjects
1254 for (address, subject, plain, html), (_, _, search_strings) in zip(emails, expected_emails):
1255 for find in search_strings:
1256 assert find in plain, f"Expected to find string {find} in PLAIN email {subject} to {address}, didn't"
1257 assert find in html, f"Expected to find string {find} in HTML email {subject} to {address}, didn't"
1260def test_send_host_request_reminders(db, moderator):
1261 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1")
1262 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2")
1263 user3, token3 = generate_user(email="user3@couchers.org.invalid", name="User 3")
1264 user4, token4 = generate_user(email="user4@couchers.org.invalid", name="User 4")
1265 user5, token5 = generate_user(email="user5@couchers.org.invalid", name="User 5")
1266 user6, token6 = generate_user(email="user6@couchers.org.invalid", name="User 6")
1267 user7, token7 = generate_user(email="user7@couchers.org.invalid", name="User 7")
1268 user8, token8 = generate_user(email="user8@couchers.org.invalid", name="User 8")
1269 user9, token9 = generate_user(email="user9@couchers.org.invalid", name="User 9")
1270 user10, token10 = generate_user(email="user10@couchers.org.invalid", name="User 10")
1271 user11, token11 = generate_user(email="user11@couchers.org.invalid", name="User 11")
1272 user12, token12 = generate_user(email="user12@couchers.org.invalid", name="User 12")
1273 user13, token13 = generate_user(email="user13@couchers.org.invalid", name="User 13")
1274 user14, token14 = generate_user(email="user14@couchers.org.invalid", name="User 14")
1276 with session_scope() as session:
1277 # case 1: pending, future, interval elapsed => notify
1278 hr1 = create_host_request_by_date(
1279 session=session,
1280 surfer_user_id=user1.id,
1281 host_user_id=user2.id,
1282 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1),
1283 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2),
1284 status=HostRequestStatus.pending,
1285 host_sent_request_reminders=0,
1286 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1287 )
1289 # case 2: max reminders reached => do not notify
1290 hr2 = create_host_request_by_date(
1291 session=session,
1292 surfer_user_id=user3.id,
1293 host_user_id=user4.id,
1294 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1),
1295 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2),
1296 status=HostRequestStatus.pending,
1297 host_sent_request_reminders=HOST_REQUEST_MAX_REMINDERS,
1298 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1299 )
1301 # case 3: interval not yet elapsed => do not notify
1302 hr3 = create_host_request_by_date(
1303 session=session,
1304 surfer_user_id=user5.id,
1305 host_user_id=user6.id,
1306 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1),
1307 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2),
1308 status=HostRequestStatus.pending,
1309 host_sent_request_reminders=0,
1310 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL + timedelta(hours=1),
1311 )
1313 # case 4: start date is today => do not notify
1314 hr4 = create_host_request_by_date(
1315 session=session,
1316 surfer_user_id=user7.id,
1317 host_user_id=user8.id,
1318 from_date=today(),
1319 to_date=today() + timedelta(days=2),
1320 status=HostRequestStatus.pending,
1321 host_sent_request_reminders=0,
1322 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1323 )
1325 # case 5: from_date in the past => do not notify
1326 hr5 = create_host_request_by_date(
1327 session=session,
1328 surfer_user_id=user9.id,
1329 host_user_id=user10.id,
1330 from_date=today() - timedelta(days=1),
1331 to_date=today() + timedelta(days=1),
1332 status=HostRequestStatus.pending,
1333 host_sent_request_reminders=0,
1334 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1335 )
1337 # case 6: non-pending status => do not notify
1338 hr6 = create_host_request_by_date(
1339 session=session,
1340 surfer_user_id=user11.id,
1341 host_user_id=user12.id,
1342 from_date=today() + timedelta(days=3),
1343 to_date=today() + timedelta(days=4),
1344 status=HostRequestStatus.accepted,
1345 host_sent_request_reminders=0,
1346 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1347 )
1349 # case 7: host already sent a message => do not notify
1350 hr7 = create_host_request_by_date(
1351 session=session,
1352 surfer_user_id=user13.id,
1353 host_user_id=user14.id,
1354 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1),
1355 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2),
1356 status=HostRequestStatus.pending,
1357 host_sent_request_reminders=0,
1358 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1359 )
1361 msg = Message(
1362 conversation_id=hr7,
1363 author_id=user14.id,
1364 text="Looking forward to hosting you!",
1365 message_type=MessageType.text,
1366 )
1367 msg.time = now()
1368 session.add(msg)
1370 # Approve host requests so they're visible for notifications
1371 moderator.approve_host_request(hr1)
1372 moderator.approve_host_request(hr2)
1373 moderator.approve_host_request(hr3)
1374 moderator.approve_host_request(hr4)
1375 moderator.approve_host_request(hr5)
1376 moderator.approve_host_request(hr6)
1377 moderator.approve_host_request(hr7)
1379 send_host_request_reminders(empty_pb2.Empty())
1381 while process_job():
1382 pass
1384 with session_scope() as session:
1385 emails = [
1386 (email.recipient, email.subject, email.plain, email.html)
1387 for email in session.execute(select(Email).order_by(Email.recipient.asc())).scalars().all()
1388 ]
1390 expected_emails = [
1391 (
1392 "user2@couchers.org.invalid",
1393 "[TEST] You have a pending host request from User 1",
1394 ("User 1", "is waiting for your response"),
1395 )
1396 ]
1398 actual_addresses_and_subjects = [email[:2] for email in emails]
1399 expected_addresses_and_subjects = [email[:2] for email in expected_emails]
1401 print(actual_addresses_and_subjects)
1402 print(expected_addresses_and_subjects)
1404 assert actual_addresses_and_subjects == expected_addresses_and_subjects
1406 for (address, subject, plain, html), (_, _, search_strings) in zip(emails, expected_emails):
1407 for find in search_strings:
1408 assert find in plain, f"Expected to find string {find} in PLAIN email {subject} to {address}, didn't"
1409 assert find in html, f"Expected to find string {find} in HTML email {subject} to {address}, didn't"
1412def test_add_users_to_email_list(db, feature_flags):
1413 feature_flags.set("listmonk_enabled", True)
1414 new_config = config.copy()
1415 new_config.LISTMONK_BASE_URL = "https://example.com"
1416 new_config.LISTMONK_API_USERNAME = "test_user"
1417 new_config.LISTMONK_API_KEY = "dummy_api_key"
1418 new_config.LISTMONK_LIST_ID = 6
1420 with patch("couchers.jobs.handlers.config", new_config):
1421 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls:
1422 mock_session_cls.return_value.post.return_value.status_code = 200
1423 add_users_to_email_list(empty_pb2.Empty())
1424 mock_session_cls.return_value.post.assert_not_called()
1426 generate_user(in_sync_with_newsletter=False, email="testing1@couchers.invalid", name="Tester1", id=15)
1427 generate_user(in_sync_with_newsletter=True, email="testing2@couchers.invalid", name="Tester2")
1428 generate_user(in_sync_with_newsletter=False, email="testing3@couchers.invalid", name="Tester3 von test", id=17)
1429 generate_user(
1430 in_sync_with_newsletter=False, email="testing4@couchers.invalid", name="Tester4", opt_out_of_newsletter=True
1431 )
1433 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls:
1434 mock_sess = mock_session_cls.return_value
1435 mock_sess.post.return_value.status_code = 200
1436 add_users_to_email_list(empty_pb2.Empty())
1437 mock_sess.post.assert_has_calls(
1438 [
1439 call(
1440 "https://example.com/api/subscribers",
1441 json={
1442 "email": "testing1@couchers.invalid",
1443 "name": "Tester1",
1444 "lists": [6],
1445 "preconfirm_subscriptions": True,
1446 "attribs": {"couchers_user_id": 15},
1447 "status": "enabled",
1448 },
1449 timeout=10,
1450 ),
1451 call(
1452 "https://example.com/api/subscribers",
1453 json={
1454 "email": "testing3@couchers.invalid",
1455 "name": "Tester3 von test",
1456 "lists": [6],
1457 "preconfirm_subscriptions": True,
1458 "attribs": {"couchers_user_id": 17},
1459 "status": "enabled",
1460 },
1461 timeout=10,
1462 ),
1463 ],
1464 any_order=True,
1465 )
1467 with patch("couchers.jobs.handlers.requests.Session") as mock_session_cls:
1468 mock_session_cls.return_value.post.return_value.status_code = 200
1469 add_users_to_email_list(empty_pb2.Empty())
1470 mock_session_cls.return_value.post.assert_not_called()
1473def test_update_recommendation_scores(db):
1474 update_recommendation_scores(empty_pb2.Empty())
1477def test_update_badges(db, push_collector: PushCollector):
1478 user1, _ = generate_user(last_donated=None)
1479 user2, _ = generate_user(last_donated=None)
1480 user3, _ = generate_user(last_donated=None)
1481 user4, _ = generate_user(phone="+15555555555", phone_verification_verified=func.now(), last_donated=None)
1482 user5, _ = generate_user(phone="+15555555556", phone_verification_verified=func.now(), last_donated=None)
1483 user6, _ = generate_user(last_donated=None)
1485 with session_scope() as session:
1486 session.add(UserBadge(user_id=user5.id, badge_id="board_member"))
1488 update_badges(empty_pb2.Empty())
1489 process_jobs()
1491 with session_scope() as session:
1492 badge_tuples = session.execute(
1493 select(UserBadge.user_id, UserBadge.badge_id).order_by(UserBadge.user_id.asc(), UserBadge.id.asc())
1494 ).all()
1496 expected = [
1497 (user1.id, "founder"),
1498 (user1.id, "board_member"),
1499 (user2.id, "founder"),
1500 (user2.id, "board_member"),
1501 (user4.id, "phone_verified"),
1502 (user5.id, "phone_verified"),
1503 ]
1505 assert badge_tuples == expected # type: ignore[comparison-overlap]
1507 print(push_collector.by_user)
1509 push = push_collector.pop_for_user(user1.id, last=False)
1510 assert push.content.title == "New profile badge: Founder"
1511 assert push.content.body == "The Founder badge was added to your profile."
1513 push = push_collector.pop_for_user(user1.id, last=True)
1514 assert push.content.title == "New profile badge: Board Member"
1515 assert push.content.body == "The Board Member badge was added to your profile."
1517 push = push_collector.pop_for_user(user2.id, last=False)
1518 assert push.content.title == "New profile badge: Founder"
1519 assert push.content.body == "The Founder badge was added to your profile."
1521 push = push_collector.pop_for_user(user2.id, last=True)
1522 assert push.content.title == "New profile badge: Board Member"
1523 assert push.content.body == "The Board Member badge was added to your profile."
1525 push = push_collector.pop_for_user(user4.id, last=True)
1526 assert push.content.title == "New profile badge: Verified Phone"
1527 assert push.content.body == "The Verified Phone badge was added to your profile."
1529 push = push_collector.pop_for_user(user5.id, last=False)
1530 assert push.content.title == "Profile badge removed"
1531 assert push.content.body == "The Board Member badge was removed from your profile."
1533 push = push_collector.pop_for_user(user5.id, last=True)
1534 assert push.content.title == "New profile badge: Verified Phone"
1535 assert push.content.body == "The Verified Phone badge was added to your profile."
1538def test_update_badges_awards_moderator_to_superuser(db):
1539 """The show_moderator_badge flag defaults on, so superusers are awarded the moderator badge."""
1540 superuser, _ = generate_user(is_superuser=True, last_donated=None)
1542 update_badges(empty_pb2.Empty())
1544 with session_scope() as session:
1545 assert (
1546 session.execute(
1547 select(func.count())
1548 .select_from(UserBadge)
1549 .where(UserBadge.user_id == superuser.id, UserBadge.badge_id == "moderator")
1550 ).scalar()
1551 == 1
1552 )
1555def test_update_badges_skips_moderator_when_flag_off(db, monkeypatch):
1556 """With show_moderator_badge forced off, superusers are not awarded the moderator badge."""
1557 # force show_moderator_badge off for everyone (force rule with no coverage applies globally)
1558 monkeypatch.setattr(experimentation, "_initialized", True)
1559 monkeypatch.setattr(
1560 experimentation,
1561 "_state",
1562 {"features": {"show_moderator_badge": {"defaultValue": True, "rules": [{"force": False}]}}, "savedGroups": {}},
1563 )
1564 monkeypatch.setitem(config, "FEATURE_FLAGS_FILE_OVERRIDE_PATH", "")
1566 superuser, _ = generate_user(is_superuser=True, last_donated=None)
1568 update_badges(empty_pb2.Empty())
1570 with session_scope() as session:
1571 assert (
1572 session.execute(
1573 select(func.count())
1574 .select_from(UserBadge)
1575 .where(UserBadge.user_id == superuser.id, UserBadge.badge_id == "moderator")
1576 ).scalar()
1577 == 0
1578 )
1581def test_send_request_notifications_blocked_users_no_notification(db, moderator):
1582 """
1583 Regression test: send_request_notifications should not send notifications
1584 when the host and surfer are not visible to each other (e.g., one blocked the other).
1585 """
1586 user1, token1 = generate_user()
1587 user2, token2 = generate_user()
1589 today_plus_2 = (today() + timedelta(days=2)).isoformat()
1590 today_plus_3 = (today() + timedelta(days=3)).isoformat()
1592 # Create a host request
1593 with requests_session(token1) as requests:
1594 host_request_id = requests.CreateHostRequest(
1595 requests_pb2.CreateHostRequestReq(
1596 host_user_id=user2.id, from_date=today_plus_2, to_date=today_plus_3, text=valid_request_text()
1597 )
1598 ).host_request_id
1599 moderator.approve_host_request(host_request_id)
1601 with session_scope() as session:
1602 # delete send_email BackgroundJob created by CreateHostRequest
1603 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1605 # Now user2 (host) blocks user1 (surfer)
1606 make_user_block(user2, user1)
1608 with session_scope() as session:
1609 # check send_request_notifications does NOT create background job because users are blocked
1610 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
1611 send_request_notifications(empty_pb2.Empty())
1612 process_jobs()
1614 # Should be 0 emails because the host blocked the surfer
1615 assert (
1616 session.execute(
1617 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1618 ).scalar_one()
1619 == 0
1620 ), "No notification email should be sent when host has blocked surfer"
1622 # Also test the reverse direction: surfer sends message to host, host should not get notification
1623 # First unblock
1624 with session_scope() as session:
1625 session.execute(delete(UserBlock).execution_options(synchronize_session=False))
1626 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1628 # Host responds
1629 with requests_session(token2) as requests:
1630 requests.RespondHostRequest(
1631 requests_pb2.RespondHostRequestReq(
1632 host_request_id=host_request_id,
1633 status=messages_pb2.HOST_REQUEST_STATUS_ACCEPTED,
1634 text="Accepting your request",
1635 )
1636 )
1638 with session_scope() as session:
1639 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1641 # Now user1 (surfer) blocks user2 (host)
1642 make_user_block(user1, user2)
1644 with session_scope() as session:
1645 # check send_request_notifications does NOT create background job
1646 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
1647 send_request_notifications(empty_pb2.Empty())
1648 process_jobs()
1650 # Should be 0 emails because the surfer blocked the host
1651 assert (
1652 session.execute(
1653 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1654 ).scalar_one()
1655 == 0
1656 ), "No notification email should be sent when surfer has blocked host"
1659def test_send_host_request_reminders_blocked_users_no_notification(db, moderator):
1660 """
1661 send_host_request_reminders should not send notifications when the host and surfer are not visible to each other
1662 (e.g., one blocked the other).
1663 """
1664 user1, token1 = generate_user(email="user1@couchers.org.invalid", name="User 1")
1665 user2, token2 = generate_user(email="user2@couchers.org.invalid", name="User 2")
1667 with session_scope() as session:
1668 # Create a pending host request where the host has not replied
1669 hr = create_host_request_by_date(
1670 session=session,
1671 surfer_user_id=user1.id,
1672 host_user_id=user2.id,
1673 from_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=1),
1674 to_date=today() + HOST_REQUEST_REMINDER_INTERVAL + timedelta(days=2),
1675 status=HostRequestStatus.pending,
1676 host_sent_request_reminders=0,
1677 last_sent_request_reminder_time=now() - HOST_REQUEST_REMINDER_INTERVAL,
1678 )
1680 # Approve the host request so it's visible for notifications
1681 moderator.approve_host_request(hr)
1683 # Verify that without blocking, a reminder would be sent
1684 send_host_request_reminders(empty_pb2.Empty())
1686 while process_job():
1687 pass
1689 with session_scope() as session:
1690 emails = session.execute(select(Email)).scalars().all()
1691 assert len(emails) == 1, "Expected 1 reminder email before blocking"
1693 # Clean up emails and background jobs
1694 session.execute(delete(Email).execution_options(synchronize_session=False))
1695 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1697 # Reset the reminder counter so we can test again
1698 host_request = session.execute(select(HostRequest).where(HostRequest.conversation_id == hr)).scalar_one()
1699 host_request.recipient_sent_request_reminders = 0
1700 host_request.last_sent_request_reminder_time = now() - HOST_REQUEST_REMINDER_INTERVAL
1702 # Now have the host block the surfer
1703 make_user_block(user2, user1)
1705 send_host_request_reminders(empty_pb2.Empty())
1707 while process_job(): 1707 ↛ 1708line 1707 didn't jump to line 1708 because the condition on line 1707 was never true
1708 pass
1710 with session_scope() as session:
1711 emails = session.execute(select(Email)).scalars().all()
1712 assert len(emails) == 0, "No reminder email should be sent when host has blocked surfer"
1715def test_send_message_notifications_blocked_users_no_notification(db, moderator):
1716 """
1717 Regression test: send_message_notifications should not send notifications
1718 for messages from users who are blocked by the recipient.
1719 """
1720 user1, token1 = generate_user()
1721 user2, token2 = generate_user()
1723 make_friends(user1, user2)
1725 # Create a group chat and send messages
1726 with conversations_session(token1) as c:
1727 group_chat_id = c.CreateGroupChat(
1728 conversations_pb2.CreateGroupChatReq(recipient_user_ids=[user2.id])
1729 ).group_chat_id
1731 # Approve the group chat so it's visible for notifications
1732 moderator.approve_group_chat(group_chat_id)
1734 with conversations_session(token1) as c:
1735 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 1"))
1736 c.SendMessage(conversations_pb2.SendMessageReq(group_chat_id=group_chat_id, text="Test message 2"))
1738 # Verify that without blocking, a notification would be sent
1739 with session_scope() as session:
1740 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1742 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
1743 send_message_notifications(empty_pb2.Empty())
1744 process_jobs()
1746 with session_scope() as session:
1747 email_job_count = session.execute(
1748 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1749 ).scalar_one()
1750 assert email_job_count == 1, "Expected 1 notification email before blocking"
1752 # Clean up
1753 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1755 # Reset the notification state so user2 will receive notifications for old messages again
1756 with session_scope() as session:
1757 u2 = session.execute(select(User).where(User.id == user2.id)).scalar_one()
1758 u2.last_notified_message_id = 0
1760 # Now have user2 block user1
1761 make_user_block(user2, user1)
1763 # The existing messages from user1 should now NOT trigger notifications
1764 # since user2 has blocked user1
1765 with session_scope() as session:
1766 session.execute(delete(BackgroundJob).execution_options(synchronize_session=False))
1768 with patch("couchers.jobs.handlers.now", now_5_min_in_future):
1769 send_message_notifications(empty_pb2.Empty())
1770 process_jobs()
1772 with session_scope() as session:
1773 email_job_count = session.execute(
1774 select(func.count()).select_from(BackgroundJob).where(BackgroundJob.job_type == "send_email")
1775 ).scalar_one()
1776 assert email_job_count == 0, "No notification email should be sent when recipient has blocked sender"
1779def test_update_badges_volunteers(db, push_collector: PushCollector):
1780 """Test that volunteer and past_volunteer badges are automatically granted based on Volunteer model."""
1781 # Create 6 users - users 1 and 2 get founder/board_member badges from static_badges
1782 user1, _ = generate_user(last_donated=None)
1783 user2, _ = generate_user(last_donated=None)
1784 user3, _ = generate_user(last_donated=None)
1785 user4, _ = generate_user(last_donated=None)
1786 user5, _ = generate_user(last_donated=None)
1787 user6, _ = generate_user(last_donated=None)
1789 with session_scope() as session:
1790 # user3: active volunteer (stopped_volunteering is null)
1791 session.add(
1792 make_volunteer(
1793 user_id=user3.id,
1794 role="Developer",
1795 started_volunteering=date(2020, 1, 1),
1796 stopped_volunteering=None,
1797 )
1798 )
1800 # user4: past volunteer (stopped_volunteering is set)
1801 session.add(
1802 make_volunteer(
1803 user_id=user4.id,
1804 role="Designer",
1805 started_volunteering=date(2020, 1, 1),
1806 stopped_volunteering=date(2023, 6, 1),
1807 )
1808 )
1810 # user5: has old volunteer badge that should be removed (not a volunteer anymore)
1811 session.add(UserBadge(user_id=user5.id, badge_id="volunteer"))
1813 # user6: has old past_volunteer badge that should be removed
1814 session.add(UserBadge(user_id=user6.id, badge_id="past_volunteer"))
1816 update_badges(empty_pb2.Empty())
1817 process_jobs()
1819 with session_scope() as session:
1820 # Check user3 has volunteer badge
1821 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all()
1822 assert "volunteer" in user3_badges
1823 assert "past_volunteer" not in user3_badges
1825 # Check user4 has past_volunteer badge
1826 user4_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user4.id)).scalars().all()
1827 assert "past_volunteer" in user4_badges
1828 assert "volunteer" not in user4_badges
1830 # Check user5 lost the volunteer badge (not in Volunteer table)
1831 user5_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user5.id)).scalars().all()
1832 assert "volunteer" not in user5_badges
1834 # Check user6 lost the past_volunteer badge (not in Volunteer table)
1835 user6_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user6.id)).scalars().all()
1836 assert "past_volunteer" not in user6_badges
1838 # Check notifications for volunteer badge users
1839 push = push_collector.pop_for_user(user3.id, last=True)
1840 assert push.content.title == "New profile badge: Active Volunteer"
1841 assert push.content.body == "The Active Volunteer badge was added to your profile."
1843 push = push_collector.pop_for_user(user4.id, last=True)
1844 assert push.content.title == "New profile badge: Past Volunteer"
1845 assert push.content.body == "The Past Volunteer badge was added to your profile."
1847 push = push_collector.pop_for_user(user5.id, last=True)
1848 assert push.content.title == "Profile badge removed"
1849 assert push.content.body == "The Active Volunteer badge was removed from your profile."
1851 push = push_collector.pop_for_user(user6.id, last=True)
1852 assert push.content.title == "Profile badge removed"
1853 assert push.content.body == "The Past Volunteer badge was removed from your profile."
1856def test_update_badges_volunteer_status_change(db, push_collector: PushCollector):
1857 """Test that badge is updated when volunteer status changes from active to past."""
1858 # Create users - users 1 and 2 get founder/board_member badges from static_badges
1859 user1, _ = generate_user(last_donated=None)
1860 user2, _ = generate_user(last_donated=None)
1861 user3, _ = generate_user(last_donated=None)
1863 with session_scope() as session:
1864 # user3: start as active volunteer
1865 session.add(
1866 make_volunteer(
1867 user_id=user3.id,
1868 role="Developer",
1869 started_volunteering=date(2020, 1, 1),
1870 stopped_volunteering=None,
1871 show_on_team_page=True,
1872 )
1873 )
1875 update_badges(empty_pb2.Empty())
1876 process_jobs()
1878 with session_scope() as session:
1879 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all()
1880 assert "volunteer" in user3_badges
1881 assert "past_volunteer" not in user3_badges
1883 push = push_collector.pop_for_user(user3.id, last=True)
1884 assert push.content.title == "New profile badge: Active Volunteer"
1885 assert push.content.body == "The Active Volunteer badge was added to your profile."
1887 # Now change the volunteer to past volunteer
1888 with session_scope() as session:
1889 volunteer = session.execute(select(Volunteer).where(Volunteer.user_id == user3.id)).scalar_one()
1890 volunteer.stopped_volunteering = date(2023, 12, 1)
1892 update_badges(empty_pb2.Empty())
1893 process_jobs()
1895 with session_scope() as session:
1896 user3_badges = session.execute(select(UserBadge.badge_id).where(UserBadge.user_id == user3.id)).scalars().all()
1897 assert "volunteer" not in user3_badges
1898 assert "past_volunteer" in user3_badges
1900 # Check both badges were updated
1901 push = push_collector.pop_for_user(user3.id, last=False)
1902 assert push.content.title == "Profile badge removed"
1903 assert push.content.body == "The Active Volunteer badge was removed from your profile."
1905 push = push_collector.pop_for_user(user3.id, last=True)
1906 assert push.content.title == "New profile badge: Past Volunteer"
1907 assert push.content.body == "The Past Volunteer badge was added to your profile."
1910def test_send_message_notifications_empty_unseen_simple(monkeypatch):
1911 class DummyUser:
1912 id = 1
1913 is_visible = True
1914 last_notified_message_id = 0
1916 class FirstResult:
1917 def scalars(self):
1918 return self
1920 def unique(self):
1921 return [DummyUser()]
1923 class SecondResult:
1924 def all(self):
1925 return []
1927 class DummySession:
1928 def __init__(self):
1929 self.calls = 0
1931 def execute(self, *a, **k):
1932 self.calls += 1
1933 return FirstResult() if self.calls == 1 else SecondResult()
1935 def commit(self):
1936 pass
1938 def flush(self):
1939 pass
1941 def fake_session_scope():
1942 class Ctx:
1943 def __enter__(self):
1944 return DummySession()
1946 def __exit__(self, exc_type, exc, tb):
1947 pass
1949 return Ctx()
1951 monkeypatch.setattr(handlers, "session_scope", fake_session_scope)
1953 handlers.send_message_notifications(Empty())