Coverage for app/backend/src/couchers/servicers/requests.py: 92%
339 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-22 16:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-22 16:01 +0000
1import logging
2from datetime import timedelta
4import grpc
5from google.protobuf import empty_pb2
6from sqlalchemy import exists, select
7from sqlalchemy.orm import Session, aliased
8from sqlalchemy.sql import and_, func, or_
10from couchers.constants import HOST_REQUEST_MIN_LENGTH_UTF16
11from couchers.context import CouchersContext, make_notification_user_context
12from couchers.db import can_moderate_node
13from couchers.event_log import log_event
14from couchers.helpers.completed_profile import has_completed_profile
15from couchers.materialized_views import UserResponseRate
16from couchers.metrics import (
17 account_age_on_host_request_create_histogram,
18 host_request_first_response_histogram,
19 host_request_responses_counter,
20 host_requests_sent_counter,
21 sent_messages_counter,
22)
23from couchers.models import (
24 Conversation,
25 HostRequest,
26 HostRequestFeedback,
27 HostRequestQuality,
28 HostRequestStatus,
29 Message,
30 MessageType,
31 ModerationObjectType,
32 RateLimitAction,
33 User,
34)
35from couchers.models.notifications import NotificationTopicAction
36from couchers.models.public_trips import PublicTrip, PublicTripStatus
37from couchers.moderation.utils import create_moderation
38from couchers.notifications.notify import mark_notifications_seen, notify
39from couchers.proto import conversations_pb2, notification_data_pb2, requests_pb2, requests_pb2_grpc
40from couchers.rate_limits.check import process_rate_limits_and_check_abort
41from couchers.rate_limits.definitions import RATE_LIMIT_HOURS
42from couchers.servicers.api import response_rate_to_pb, user_model_to_pb
43from couchers.sql import to_bool, users_visible, where_moderated_content_visible, where_users_column_visible
44from couchers.utils import (
45 Timestamp_from_datetime,
46 date_to_api,
47 get_coordinates,
48 now,
49 parse_date,
50 today_in_timezone,
51)
53logger = logging.getLogger(__name__)
55DEFAULT_PAGINATION_LENGTH = 10
56MAX_PAGE_SIZE = 50
59hostrequeststatus2api = {
60 HostRequestStatus.pending: conversations_pb2.HOST_REQUEST_STATUS_PENDING,
61 HostRequestStatus.accepted: conversations_pb2.HOST_REQUEST_STATUS_ACCEPTED,
62 HostRequestStatus.rejected: conversations_pb2.HOST_REQUEST_STATUS_REJECTED,
63 HostRequestStatus.confirmed: conversations_pb2.HOST_REQUEST_STATUS_CONFIRMED,
64 HostRequestStatus.cancelled: conversations_pb2.HOST_REQUEST_STATUS_CANCELLED,
65}
67api2hostrequeststatus = {
68 conversations_pb2.HOST_REQUEST_STATUS_PENDING: HostRequestStatus.pending,
69 conversations_pb2.HOST_REQUEST_STATUS_ACCEPTED: HostRequestStatus.accepted,
70 conversations_pb2.HOST_REQUEST_STATUS_REJECTED: HostRequestStatus.rejected,
71 conversations_pb2.HOST_REQUEST_STATUS_CONFIRMED: HostRequestStatus.confirmed,
72 conversations_pb2.HOST_REQUEST_STATUS_CANCELLED: HostRequestStatus.cancelled,
73}
75hostrequestquality2sql = {
76 requests_pb2.HOST_REQUEST_QUALITY_UNSPECIFIED: HostRequestQuality.high_quality,
77 requests_pb2.HOST_REQUEST_QUALITY_LOW: HostRequestQuality.okay_quality,
78 requests_pb2.HOST_REQUEST_QUALITY_OKAY: HostRequestQuality.low_quality,
79}
82def message_to_pb(message: Message) -> conversations_pb2.Message:
83 """
84 Turns the given message to a protocol buffer
85 """
86 if message.is_normal_message:
87 return conversations_pb2.Message(
88 message_id=message.id,
89 author_user_id=message.author_id,
90 time=Timestamp_from_datetime(message.time),
91 text=conversations_pb2.MessageContentText(text=message.text),
92 )
93 else:
94 return conversations_pb2.Message(
95 message_id=message.id,
96 author_user_id=message.author_id,
97 time=Timestamp_from_datetime(message.time),
98 chat_created=(
99 conversations_pb2.MessageContentChatCreated()
100 if message.message_type == MessageType.chat_created
101 else None
102 ),
103 host_request_status_changed=(
104 conversations_pb2.MessageContentHostRequestStatusChanged(
105 status=hostrequeststatus2api[message.host_request_status_target] # type: ignore[index]
106 )
107 if message.message_type == MessageType.host_request_status_changed
108 else None
109 ),
110 )
113def host_request_to_pb(
114 host_request: HostRequest, session: Session, context: CouchersContext
115) -> requests_pb2.HostRequest:
116 initial_message = session.execute(
117 select(Message)
118 .where(Message.conversation_id == host_request.conversation_id)
119 .order_by(Message.id.asc())
120 .limit(1)
121 ).scalar_one()
123 latest_message = session.execute(
124 select(Message)
125 .where(Message.conversation_id == host_request.conversation_id)
126 .order_by(Message.id.desc())
127 .limit(1)
128 ).scalar_one()
130 lat, lng = get_coordinates(host_request.hosting_location)
132 need_feedback = False
133 if context.user_id == host_request.recipient_user_id and host_request.status == HostRequestStatus.rejected:
134 need_feedback = not session.execute(
135 select(
136 exists().where(
137 HostRequestFeedback.from_user_id == context.user_id,
138 HostRequestFeedback.host_request_id == host_request.conversation_id,
139 )
140 )
141 ).scalar_one()
143 return requests_pb2.HostRequest(
144 host_request_id=host_request.conversation_id,
145 surfer_user_id=host_request.initiator_user_id,
146 host_user_id=host_request.recipient_user_id,
147 status=hostrequeststatus2api[host_request.status],
148 created=Timestamp_from_datetime(initial_message.time),
149 from_date=date_to_api(host_request.from_date),
150 to_date=date_to_api(host_request.to_date),
151 last_seen_message_id=(
152 host_request.initiator_last_seen_message_id
153 if context.user_id == host_request.initiator_user_id
154 else host_request.recipient_last_seen_message_id
155 ),
156 latest_message=message_to_pb(latest_message),
157 hosting_city=host_request.hosting_city,
158 hosting_lat=lat,
159 hosting_lng=lng,
160 hosting_radius=host_request.hosting_radius,
161 need_host_request_feedback=need_feedback,
162 is_archived=(
163 host_request.is_recipient_archived
164 if context.user_id == host_request.recipient_user_id
165 else host_request.is_initiator_archived
166 ),
167 public_trip_id=host_request.public_trip_id,
168 )
171def _possibly_observe_first_response_time(
172 session: Session, host_request: HostRequest, user_id: int, response_type: str
173) -> None:
174 # if this is the first response then there's nothing by this user yet
175 assert host_request.recipient_user_id == user_id
177 number_messages_by_host = session.execute(
178 select(func.count())
179 .where(Message.conversation_id == host_request.conversation_id)
180 .where(Message.author_id == user_id)
181 ).scalar_one_or_none()
183 if number_messages_by_host == 0:
184 host_gender = session.execute(select(User.gender).where(User.id == host_request.recipient_user_id)).scalar_one()
185 surfer_gender = session.execute(
186 select(User.gender).where(User.id == host_request.initiator_user_id)
187 ).scalar_one()
188 host_request_first_response_histogram.labels(host_gender, surfer_gender, response_type).observe(
189 (now() - host_request.conversation.created).total_seconds()
190 )
193def _is_host_request_long_enough(text: str) -> bool:
194 # Python's len(str) does not match Javascript's string.length.
195 # e.g. len("é") == 2 but "é".length == 1.
196 # To match the frontend's validation, measure the string in utf16 code units.
197 text_length_utf16 = len(text.encode("utf-16-le")) // 2 # utf-16-le does not include a prefix BOM code unit.
198 return text_length_utf16 >= HOST_REQUEST_MIN_LENGTH_UTF16
201class Requests(requests_pb2_grpc.RequestsServicer):
202 def CreateHostRequest(
203 self, request: requests_pb2.CreateHostRequestReq, context: CouchersContext, session: Session
204 ) -> requests_pb2.CreateHostRequestRes:
205 user = session.execute(select(User).where(User.id == context.user_id)).scalar_one()
206 if not has_completed_profile(session, user):
207 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "incomplete_profile_send_request")
209 if request.host_user_id == context.user_id:
210 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "cant_request_self")
212 # just to check recipient exists and is visible
213 recipient = session.execute(
214 select(User).where(users_visible(context, User)).where(User.id == request.host_user_id)
215 ).scalar_one_or_none()
216 if not recipient:
217 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "user_not_found")
219 from_date = parse_date(request.from_date)
220 to_date = parse_date(request.to_date)
222 if not from_date or not to_date:
223 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_date")
225 today = today_in_timezone(recipient.timezone)
227 # request starts from the past
228 if from_date < today:
229 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_from_before_today")
231 # from_date is not >= to_date
232 if from_date >= to_date:
233 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_from_after_to")
235 # No need to check today > to_date
237 if from_date - today > timedelta(days=365):
238 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_from_after_one_year")
240 if to_date - from_date > timedelta(days=365):
241 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_to_after_one_year")
243 # Check minimum length
244 if not _is_host_request_long_enough(request.text):
245 context.abort_with_error_code(
246 grpc.StatusCode.INVALID_ARGUMENT,
247 "host_request_too_short2",
248 substitutions={"count": HOST_REQUEST_MIN_LENGTH_UTF16},
249 )
251 # Check if user has been sending host requests excessively
252 if process_rate_limits_and_check_abort(
253 session=session, user_id=context.user_id, action=RateLimitAction.host_request
254 ):
255 context.abort_with_error_code(
256 grpc.StatusCode.RESOURCE_EXHAUSTED,
257 "host_request_rate_limit2",
258 substitutions={"count": RATE_LIMIT_HOURS},
259 )
261 # If this is an offer in response to a public trip, validate it
262 public_trip_id = request.public_trip_id if request.HasField("public_trip_id") else None
263 if public_trip_id is not None:
264 public_trip = session.execute(
265 select(PublicTrip).where(PublicTrip.id == public_trip_id)
266 ).scalar_one_or_none()
267 if not public_trip:
268 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "public_trip_not_found")
269 # The trip's traveler must be the recipient of this host request (role reversal)
270 if public_trip.user_id != recipient.id:
271 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "public_trip_user_mismatch")
272 # Trip must still be active
273 if public_trip.status != PublicTripStatus.searching_for_host:
274 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "public_trip_not_active")
275 # Offered dates must fall within the trip's window (host can shorten, not extend)
276 if from_date < public_trip.from_date or to_date > public_trip.to_date:
277 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "public_trip_dates_out_of_range")
278 # Enforce same_gender_only restriction (community moderators bypass)
279 if (
280 public_trip.same_gender_only
281 and not can_moderate_node(session, context.user_id, public_trip.node_id)
282 and user.gender != recipient.gender
283 ):
284 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "public_trip_same_gender_only")
285 # Prevent duplicate offers on the same trip
286 existing_offer = session.execute(
287 select(HostRequest)
288 .where(HostRequest.public_trip_id == public_trip_id)
289 .where(HostRequest.initiator_user_id == context.user_id)
290 ).scalar_one_or_none()
291 if existing_offer:
292 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "duplicate_host_request_for_trip")
294 conversation = Conversation()
295 session.add(conversation)
296 session.flush()
298 session.add(
299 Message(
300 conversation_id=conversation.id,
301 author_id=context.user_id,
302 message_type=MessageType.chat_created,
303 )
304 )
306 message = Message(
307 conversation_id=conversation.id,
308 author_id=context.user_id,
309 text=request.text,
310 message_type=MessageType.text,
311 )
312 session.add(message)
313 session.flush()
315 # Create moderation state for UMS (starts as SHADOWED)
316 moderation_state = create_moderation(
317 session=session,
318 object_type=ModerationObjectType.host_request,
319 object_id=conversation.id,
320 creator_user_id=context.user_id,
321 )
323 host_request = HostRequest(
324 conversation_id=conversation.id,
325 initiator_user_id=context.user_id,
326 recipient_user_id=recipient.id,
327 moderation_state_id=moderation_state.id,
328 from_date=from_date,
329 to_date=to_date,
330 status=HostRequestStatus.pending,
331 initiator_last_seen_message_id=message.id,
332 # TODO: tz
333 # timezone=recipient.timezone,
334 hosting_city=recipient.city,
335 hosting_location=recipient.geom,
336 hosting_radius=recipient.geom_radius,
337 public_trip_id=public_trip_id,
338 )
339 session.add(host_request)
340 session.flush()
342 recipient_context = make_notification_user_context(user_id=host_request.recipient_user_id)
343 notify(
344 session,
345 user_id=host_request.recipient_user_id,
346 topic_action=NotificationTopicAction.host_request__create,
347 key=str(host_request.conversation_id),
348 data=notification_data_pb2.HostRequestCreate(
349 host_request=host_request_to_pb(host_request, session, recipient_context),
350 surfer=user_model_to_pb(host_request.initiator, session, recipient_context),
351 text=request.text,
352 ),
353 moderation_state_id=moderation_state.id,
354 )
356 host_requests_sent_counter.labels(user.gender, recipient.gender).inc()
357 sent_messages_counter.labels(user.gender, "host request send").inc()
358 account_age_on_host_request_create_histogram.labels(user.gender, recipient.gender).observe(
359 (now() - user.joined).total_seconds()
360 )
361 log_event(
362 context,
363 session,
364 "host_request.created",
365 {
366 "host_request_id": host_request.conversation_id,
367 "host_id": recipient.id,
368 "surfer_gender": user.gender,
369 "host_gender": recipient.gender,
370 "city": recipient.city,
371 "from_date": str(from_date),
372 "to_date": str(to_date),
373 "nights": (to_date - from_date).days,
374 },
375 )
377 return requests_pb2.CreateHostRequestRes(host_request_id=host_request.conversation_id)
379 def GetHostRequest(
380 self, request: requests_pb2.GetHostRequestReq, context: CouchersContext, session: Session
381 ) -> requests_pb2.HostRequest:
382 host_request = session.execute(
383 where_moderated_content_visible(
384 where_users_column_visible(
385 where_users_column_visible(
386 select(HostRequest),
387 context,
388 HostRequest.initiator_user_id,
389 ),
390 context,
391 HostRequest.recipient_user_id,
392 ),
393 context,
394 HostRequest,
395 is_list_operation=False,
396 )
397 .where(HostRequest.conversation_id == request.host_request_id)
398 .where(
399 or_(HostRequest.initiator_user_id == context.user_id, HostRequest.recipient_user_id == context.user_id)
400 )
401 ).scalar_one_or_none()
403 if not host_request:
404 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
406 return host_request_to_pb(host_request, session, context)
408 # TODO(#7722): remove after FE migrates to ListMessageThreads
409 def ListHostRequests(
410 self, request: requests_pb2.ListHostRequestsReq, context: CouchersContext, session: Session
411 ) -> requests_pb2.ListHostRequestsRes:
412 if request.only_sent and request.only_received: 412 ↛ 413line 412 didn't jump to line 413 because the condition on line 412 was never true
413 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "host_request_sent_or_received")
415 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH
416 pagination = min(pagination, MAX_PAGE_SIZE)
418 # By outer joining messages on itself where the second id is bigger, only the highest IDs will have
419 # none as message_2.id. So just filter for these to get the highest messages only.
420 # See https://stackoverflow.com/a/27802817/6115336
421 message_2 = aliased(Message)
422 statement = where_moderated_content_visible(
423 where_users_column_visible(
424 where_users_column_visible(
425 select(Message, HostRequest, Conversation)
426 .outerjoin(
427 message_2, and_(Message.conversation_id == message_2.conversation_id, Message.id < message_2.id)
428 )
429 .join(HostRequest, HostRequest.conversation_id == Message.conversation_id)
430 .join(Conversation, Conversation.id == HostRequest.conversation_id),
431 context,
432 HostRequest.initiator_user_id,
433 ),
434 context,
435 HostRequest.recipient_user_id,
436 ),
437 context,
438 HostRequest,
439 is_list_operation=True,
440 ).where(message_2.id == None)
442 sort_by_from_date = request.sort_by == requests_pb2.HOST_REQUEST_SORT_BY_FROM_DATE
444 if sort_by_from_date:
445 if request.page_token:
446 token_date_str, token_conv_id_str = request.page_token.split(":")
447 token_date = parse_date(token_date_str)
448 token_conv_id = int(token_conv_id_str)
449 statement = statement.where(
450 or_(
451 HostRequest.from_date > token_date,
452 and_(
453 HostRequest.from_date == token_date,
454 HostRequest.conversation_id > token_conv_id,
455 ),
456 )
457 )
458 else:
459 if request.page_token:
460 statement = statement.where(Message.id < int(request.page_token))
462 if request.only_sent:
463 statement = statement.where(HostRequest.initiator_user_id == context.user_id)
464 elif request.only_received:
465 statement = statement.where(HostRequest.recipient_user_id == context.user_id)
466 elif request.HasField("only_archived"):
467 statement = statement.where(
468 or_(
469 and_(
470 HostRequest.initiator_user_id == context.user_id,
471 HostRequest.is_initiator_archived == request.only_archived,
472 ),
473 and_(
474 HostRequest.recipient_user_id == context.user_id,
475 HostRequest.is_recipient_archived == request.only_archived,
476 ),
477 )
478 )
479 else:
480 statement = statement.where(
481 or_(HostRequest.recipient_user_id == context.user_id, HostRequest.initiator_user_id == context.user_id)
482 )
484 # TODO: I considered having the latest control message be the single source of truth for
485 # the HostRequest.status, but decided against it because of this filter.
486 # Another possibility is to filter in the python instead of SQL, but that's slower
487 if request.only_active:
488 statement = statement.where(
489 or_(
490 HostRequest.status == HostRequestStatus.pending,
491 HostRequest.status == HostRequestStatus.accepted,
492 HostRequest.status == HostRequestStatus.confirmed,
493 )
494 )
495 statement = statement.where(HostRequest.end_time >= func.now())
497 if request.status_in:
498 statement = statement.where(HostRequest.status.in_([api2hostrequeststatus[s] for s in request.status_in]))
500 if sort_by_from_date:
501 statement = statement.order_by(HostRequest.from_date.asc(), HostRequest.conversation_id.asc())
502 else:
503 statement = statement.order_by(Message.id.desc())
504 statement = statement.limit(pagination + 1)
505 results = session.execute(statement).all()
507 host_requests = []
508 for result in results[:pagination]:
509 lat, lng = get_coordinates(result.HostRequest.hosting_location)
510 host_requests.append(
511 requests_pb2.HostRequest(
512 host_request_id=result.HostRequest.conversation_id,
513 surfer_user_id=result.HostRequest.initiator_user_id,
514 host_user_id=result.HostRequest.recipient_user_id,
515 status=hostrequeststatus2api[result.HostRequest.status],
516 created=Timestamp_from_datetime(result.Conversation.created),
517 from_date=date_to_api(result.HostRequest.from_date),
518 to_date=date_to_api(result.HostRequest.to_date),
519 last_seen_message_id=(
520 result.HostRequest.initiator_last_seen_message_id
521 if context.user_id == result.HostRequest.initiator_user_id
522 else result.HostRequest.recipient_last_seen_message_id
523 ),
524 latest_message=message_to_pb(result.Message),
525 hosting_city=result.HostRequest.hosting_city,
526 hosting_lat=lat,
527 hosting_lng=lng,
528 hosting_radius=result.HostRequest.hosting_radius,
529 )
530 )
532 no_more = len(results) <= pagination
534 if len(results) > pagination:
535 if sort_by_from_date:
536 last = results[pagination - 1]
537 next_page_token = f"{date_to_api(last.HostRequest.from_date)}:{last.HostRequest.conversation_id}"
538 else:
539 next_page_token = str(min(g.Message.id for g in results[:pagination]))
540 else:
541 next_page_token = None
543 return requests_pb2.ListHostRequestsRes(
544 next_page_token=next_page_token, no_more=no_more, host_requests=host_requests
545 )
547 def RespondHostRequest(
548 self, request: requests_pb2.RespondHostRequestReq, context: CouchersContext, session: Session
549 ) -> empty_pb2.Empty:
550 def count_host_response(other_user_id: int, response_type: str) -> None:
551 user_gender = session.execute(select(User.gender).where(User.id == context.user_id)).scalar_one()
552 other_gender = session.execute(select(User.gender).where(User.id == other_user_id)).scalar_one()
553 host_request_responses_counter.labels(user_gender, other_gender, response_type).inc()
554 sent_messages_counter.labels(user_gender, "host request response").inc()
556 host_request = session.execute(
557 where_moderated_content_visible(
558 where_users_column_visible(
559 where_users_column_visible(
560 select(HostRequest),
561 context,
562 HostRequest.initiator_user_id,
563 ),
564 context,
565 HostRequest.recipient_user_id,
566 ),
567 context,
568 HostRequest,
569 is_list_operation=False,
570 ).where(HostRequest.conversation_id == request.host_request_id)
571 ).scalar_one_or_none()
573 if not host_request:
574 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
576 if host_request.initiator_user_id != context.user_id and host_request.recipient_user_id != context.user_id:
577 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
579 if request.status == conversations_pb2.HOST_REQUEST_STATUS_PENDING:
580 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
582 if host_request.end_time < now(): 582 ↛ 583line 582 didn't jump to line 583 because the condition on line 582 was never true
583 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "host_request_in_past")
585 control_message = Message(
586 message_type=MessageType.host_request_status_changed,
587 conversation_id=host_request.conversation_id,
588 author_id=context.user_id,
589 )
591 if request.status == conversations_pb2.HOST_REQUEST_STATUS_ACCEPTED:
592 # only host can accept
593 if context.user_id != host_request.recipient_user_id:
594 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "not_the_host")
595 # can't accept a cancelled or confirmed request (only reject), or already accepted
596 if ( 596 ↛ 601line 596 didn't jump to line 601 because the condition on line 596 was never true
597 host_request.status == HostRequestStatus.cancelled
598 or host_request.status == HostRequestStatus.confirmed
599 or host_request.status == HostRequestStatus.accepted
600 ):
601 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
602 _possibly_observe_first_response_time(session, host_request, context.user_id, "accepted")
603 control_message.host_request_status_target = HostRequestStatus.accepted
604 host_request.status = HostRequestStatus.accepted
605 session.flush()
607 recipient_context = make_notification_user_context(user_id=host_request.initiator_user_id)
608 notify(
609 session,
610 user_id=host_request.initiator_user_id,
611 topic_action=NotificationTopicAction.host_request__accept,
612 key=str(host_request.conversation_id),
613 data=notification_data_pb2.HostRequestAccept(
614 host_request=host_request_to_pb(host_request, session, recipient_context),
615 host=user_model_to_pb(host_request.recipient, session, recipient_context),
616 ),
617 moderation_state_id=host_request.moderation_state_id,
618 )
620 count_host_response(host_request.initiator_user_id, "accepted")
621 log_event(
622 context,
623 session,
624 "host_request.accepted",
625 {
626 "host_request_id": host_request.conversation_id,
627 "surfer_id": host_request.initiator_user_id,
628 "host_id": host_request.recipient_user_id,
629 "surfer_gender": host_request.initiator.gender,
630 "host_gender": host_request.recipient.gender,
631 "from_date": str(host_request.from_date),
632 "to_date": str(host_request.to_date),
633 "host_city": host_request.hosting_city,
634 },
635 )
637 if request.status == conversations_pb2.HOST_REQUEST_STATUS_REJECTED:
638 # only host can reject
639 if context.user_id != host_request.recipient_user_id: 639 ↛ 640line 639 didn't jump to line 640 because the condition on line 639 was never true
640 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
641 # can't reject a cancelled or already rejected request
642 if host_request.status == HostRequestStatus.cancelled or host_request.status == HostRequestStatus.rejected: 642 ↛ 643line 642 didn't jump to line 643 because the condition on line 642 was never true
643 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
644 _possibly_observe_first_response_time(session, host_request, context.user_id, "rejected")
645 control_message.host_request_status_target = HostRequestStatus.rejected
646 host_request.status = HostRequestStatus.rejected
647 session.flush()
649 recipient_context = make_notification_user_context(user_id=host_request.initiator_user_id)
650 notify(
651 session,
652 user_id=host_request.initiator_user_id,
653 topic_action=NotificationTopicAction.host_request__reject,
654 key=str(host_request.conversation_id),
655 data=notification_data_pb2.HostRequestReject(
656 host_request=host_request_to_pb(host_request, session, recipient_context),
657 host=user_model_to_pb(host_request.recipient, session, recipient_context),
658 ),
659 moderation_state_id=host_request.moderation_state_id,
660 )
662 count_host_response(host_request.initiator_user_id, "rejected")
664 log_event(
665 context,
666 session,
667 "host_request.rejected",
668 {
669 "host_request_id": host_request.conversation_id,
670 "surfer_id": host_request.initiator_user_id,
671 "host_id": host_request.recipient_user_id,
672 "surfer_gender": host_request.initiator.gender,
673 "host_gender": host_request.recipient.gender,
674 "from_date": str(host_request.from_date),
675 "to_date": str(host_request.to_date),
676 "host_city": host_request.hosting_city,
677 },
678 )
680 if request.status == conversations_pb2.HOST_REQUEST_STATUS_CONFIRMED:
681 # only surfer can confirm
682 if context.user_id != host_request.initiator_user_id:
683 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
684 # can only confirm an accepted request
685 if host_request.status != HostRequestStatus.accepted:
686 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
687 control_message.host_request_status_target = HostRequestStatus.confirmed
688 host_request.status = HostRequestStatus.confirmed
689 session.flush()
691 recipient_context = make_notification_user_context(user_id=host_request.recipient_user_id)
692 notify(
693 session,
694 user_id=host_request.recipient_user_id,
695 topic_action=NotificationTopicAction.host_request__confirm,
696 key=str(host_request.conversation_id),
697 data=notification_data_pb2.HostRequestConfirm(
698 host_request=host_request_to_pb(host_request, session, recipient_context),
699 surfer=user_model_to_pb(host_request.initiator, session, recipient_context),
700 ),
701 moderation_state_id=host_request.moderation_state_id,
702 )
704 count_host_response(host_request.recipient_user_id, "confirmed")
705 log_event(
706 context,
707 session,
708 "host_request.confirmed",
709 {
710 "host_request_id": host_request.conversation_id,
711 "surfer_id": host_request.initiator_user_id,
712 "host_id": host_request.recipient_user_id,
713 "surfer_gender": host_request.initiator.gender,
714 "host_gender": host_request.recipient.gender,
715 "from_date": str(host_request.from_date),
716 "to_date": str(host_request.to_date),
717 "host_city": host_request.hosting_city,
718 },
719 )
721 if request.status == conversations_pb2.HOST_REQUEST_STATUS_CANCELLED:
722 # only surfer can cancel
723 if context.user_id != host_request.initiator_user_id:
724 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
725 # can't' cancel an already cancelled or rejected request
726 if host_request.status == HostRequestStatus.rejected or host_request.status == HostRequestStatus.cancelled: 726 ↛ 727line 726 didn't jump to line 727 because the condition on line 726 was never true
727 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status")
728 control_message.host_request_status_target = HostRequestStatus.cancelled
729 host_request.status = HostRequestStatus.cancelled
730 session.flush()
732 recipient_context = make_notification_user_context(user_id=host_request.recipient_user_id)
733 notify(
734 session,
735 user_id=host_request.recipient_user_id,
736 topic_action=NotificationTopicAction.host_request__cancel,
737 key=str(host_request.conversation_id),
738 data=notification_data_pb2.HostRequestCancel(
739 host_request=host_request_to_pb(host_request, session, recipient_context),
740 surfer=user_model_to_pb(host_request.initiator, session, recipient_context),
741 ),
742 moderation_state_id=host_request.moderation_state_id,
743 )
745 count_host_response(host_request.recipient_user_id, "cancelled")
746 log_event(
747 context,
748 session,
749 "host_request.cancelled",
750 {
751 "host_request_id": host_request.conversation_id,
752 "surfer_id": host_request.initiator_user_id,
753 "host_id": host_request.recipient_user_id,
754 "surfer_gender": host_request.initiator.gender,
755 "host_gender": host_request.recipient.gender,
756 "from_date": str(host_request.from_date),
757 "to_date": str(host_request.to_date),
758 "host_city": host_request.hosting_city,
759 },
760 )
762 session.add(control_message)
764 if request.text:
765 latest_message = Message(
766 conversation_id=host_request.conversation_id,
767 text=request.text,
768 author_id=context.user_id,
769 message_type=MessageType.text,
770 )
772 session.add(latest_message)
773 else:
774 latest_message = control_message
776 session.flush()
778 if host_request.initiator_user_id == context.user_id:
779 host_request.initiator_last_seen_message_id = latest_message.id
780 else:
781 host_request.recipient_last_seen_message_id = latest_message.id
782 session.commit()
784 return empty_pb2.Empty()
786 def GetHostRequestMessages(
787 self, request: requests_pb2.GetHostRequestMessagesReq, context: CouchersContext, session: Session
788 ) -> requests_pb2.GetHostRequestMessagesRes:
789 host_request = session.execute(
790 where_moderated_content_visible(select(HostRequest), context, HostRequest, is_list_operation=False).where(
791 HostRequest.conversation_id == request.host_request_id
792 )
793 ).scalar_one_or_none()
795 if not host_request: 795 ↛ 796line 795 didn't jump to line 796 because the condition on line 795 was never true
796 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
798 if host_request.initiator_user_id != context.user_id and host_request.recipient_user_id != context.user_id: 798 ↛ 799line 798 didn't jump to line 799 because the condition on line 798 was never true
799 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
801 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH
802 pagination = min(pagination, MAX_PAGE_SIZE)
804 messages = (
805 session.execute(
806 select(Message)
807 .where(Message.conversation_id == host_request.conversation_id)
808 .where(or_(Message.id < request.last_message_id, to_bool(request.last_message_id == 0)))
809 .order_by(Message.id.desc())
810 .limit(pagination + 1)
811 )
812 .scalars()
813 .all()
814 )
816 no_more = len(messages) <= pagination
818 last_message_id = min(m.id if m else 1 for m in messages[:pagination]) if len(messages) > 0 else 0
820 return requests_pb2.GetHostRequestMessagesRes(
821 last_message_id=last_message_id,
822 no_more=no_more,
823 messages=[message_to_pb(message) for message in messages[:pagination]],
824 )
826 def SendHostRequestMessage(
827 self, request: requests_pb2.SendHostRequestMessageReq, context: CouchersContext, session: Session
828 ) -> empty_pb2.Empty:
829 if request.text == "":
830 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_message")
831 host_request = session.execute(
832 where_moderated_content_visible(select(HostRequest), context, HostRequest, is_list_operation=False).where(
833 HostRequest.conversation_id == request.host_request_id
834 )
835 ).scalar_one_or_none()
837 if not host_request:
838 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
840 if host_request.initiator_user_id != context.user_id and host_request.recipient_user_id != context.user_id:
841 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
843 if host_request.recipient_user_id == context.user_id:
844 _possibly_observe_first_response_time(session, host_request, context.user_id, "message")
846 message = Message(
847 conversation_id=host_request.conversation_id,
848 author_id=context.user_id,
849 message_type=MessageType.text,
850 text=request.text,
851 )
853 session.add(message)
854 session.flush()
856 if host_request.initiator_user_id == context.user_id:
857 host_request.initiator_last_seen_message_id = message.id
859 recipient_context = make_notification_user_context(user_id=host_request.recipient_user_id)
860 notify(
861 session,
862 user_id=host_request.recipient_user_id,
863 topic_action=NotificationTopicAction.host_request__message,
864 key=str(host_request.conversation_id),
865 data=notification_data_pb2.HostRequestMessage(
866 host_request=host_request_to_pb(host_request, session, recipient_context),
867 user=user_model_to_pb(host_request.initiator, session, recipient_context),
868 text=request.text,
869 am_host=True,
870 ),
871 moderation_state_id=host_request.moderation_state_id,
872 )
874 else:
875 host_request.recipient_last_seen_message_id = message.id
877 recipient_context = make_notification_user_context(user_id=host_request.initiator_user_id)
878 notify(
879 session,
880 user_id=host_request.initiator_user_id,
881 topic_action=NotificationTopicAction.host_request__message,
882 key=str(host_request.conversation_id),
883 data=notification_data_pb2.HostRequestMessage(
884 host_request=host_request_to_pb(host_request, session, recipient_context),
885 user=user_model_to_pb(host_request.recipient, session, recipient_context),
886 text=request.text,
887 am_host=False,
888 ),
889 moderation_state_id=host_request.moderation_state_id,
890 )
892 session.commit()
894 user_gender = session.execute(select(User.gender).where(User.id == context.user_id)).scalar_one()
895 sent_messages_counter.labels(user_gender, "host request").inc()
896 log_event(
897 context,
898 session,
899 "host_request.message_sent",
900 {
901 "host_request_id": host_request.conversation_id,
902 "surfer_id": host_request.initiator_user_id,
903 "host_id": host_request.recipient_user_id,
904 "role": "host" if context.user_id == host_request.recipient_user_id else "surfer",
905 "host_city": host_request.hosting_city,
906 },
907 )
909 return empty_pb2.Empty()
911 def GetHostRequestUpdates(
912 self, request: requests_pb2.GetHostRequestUpdatesReq, context: CouchersContext, session: Session
913 ) -> requests_pb2.GetHostRequestUpdatesRes:
914 if request.only_sent and request.only_received: 914 ↛ 915line 914 didn't jump to line 915 because the condition on line 914 was never true
915 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "host_request_sent_or_received")
917 if request.newest_message_id == 0:
918 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_message")
920 if not session.execute(select(Message).where(Message.id == request.newest_message_id)).scalar_one_or_none(): 920 ↛ 921line 920 didn't jump to line 921 because the condition on line 920 was never true
921 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_message")
923 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH
924 pagination = min(pagination, MAX_PAGE_SIZE)
926 statement = where_moderated_content_visible(
927 select(
928 Message,
929 HostRequest.status.label("host_request_status"),
930 HostRequest.conversation_id.label("host_request_id"),
931 )
932 .join(HostRequest, HostRequest.conversation_id == Message.conversation_id)
933 .where(Message.id > request.newest_message_id),
934 context,
935 HostRequest,
936 is_list_operation=False,
937 )
939 if request.only_sent: 939 ↛ 940line 939 didn't jump to line 940 because the condition on line 939 was never true
940 statement = statement.where(HostRequest.initiator_user_id == context.user_id)
941 elif request.only_received: 941 ↛ 942line 941 didn't jump to line 942 because the condition on line 941 was never true
942 statement = statement.where(HostRequest.recipient_user_id == context.user_id)
943 else:
944 statement = statement.where(
945 or_(HostRequest.recipient_user_id == context.user_id, HostRequest.initiator_user_id == context.user_id)
946 )
948 statement = statement.order_by(Message.id.asc()).limit(pagination + 1)
949 res = session.execute(statement).all()
951 no_more = len(res) <= pagination
953 last_message_id = min(m.Message.id if m else 1 for m in res[:pagination]) if len(res) > 0 else 0 # TODO
955 return requests_pb2.GetHostRequestUpdatesRes(
956 no_more=no_more,
957 updates=[
958 requests_pb2.HostRequestUpdate(
959 host_request_id=result.host_request_id,
960 status=hostrequeststatus2api[result.host_request_status],
961 message=message_to_pb(result.Message),
962 )
963 for result in res[:pagination]
964 ],
965 )
967 def MarkLastSeenHostRequest(
968 self, request: requests_pb2.MarkLastSeenHostRequestReq, context: CouchersContext, session: Session
969 ) -> empty_pb2.Empty:
970 host_request = session.execute(
971 where_moderated_content_visible(select(HostRequest), context, HostRequest, is_list_operation=False).where(
972 HostRequest.conversation_id == request.host_request_id
973 )
974 ).scalar_one_or_none()
976 if not host_request: 976 ↛ 977line 976 didn't jump to line 977 because the condition on line 976 was never true
977 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
979 if host_request.initiator_user_id != context.user_id and host_request.recipient_user_id != context.user_id: 979 ↛ 980line 979 didn't jump to line 980 because the condition on line 979 was never true
980 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
982 if host_request.initiator_user_id == context.user_id: 982 ↛ 983line 982 didn't jump to line 983 because the condition on line 982 was never true
983 if not host_request.initiator_last_seen_message_id <= request.last_seen_message_id:
984 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "cant_unsee_messages")
985 host_request.initiator_last_seen_message_id = request.last_seen_message_id
986 else:
987 if not host_request.recipient_last_seen_message_id <= request.last_seen_message_id:
988 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "cant_unsee_messages")
989 host_request.recipient_last_seen_message_id = request.last_seen_message_id
991 mark_notifications_seen(
992 session,
993 user_id=context.user_id,
994 key=str(host_request.conversation_id),
995 topic_actions=[
996 NotificationTopicAction.host_request__create,
997 NotificationTopicAction.host_request__accept,
998 NotificationTopicAction.host_request__reject,
999 NotificationTopicAction.host_request__confirm,
1000 NotificationTopicAction.host_request__cancel,
1001 NotificationTopicAction.host_request__message,
1002 NotificationTopicAction.host_request__missed_messages,
1003 NotificationTopicAction.host_request__reminder,
1004 ],
1005 )
1007 session.commit()
1008 return empty_pb2.Empty()
1010 def SetHostRequestArchiveStatus(
1011 self, request: requests_pb2.SetHostRequestArchiveStatusReq, context: CouchersContext, session: Session
1012 ) -> requests_pb2.SetHostRequestArchiveStatusRes:
1013 host_request = session.execute(
1014 where_moderated_content_visible(select(HostRequest), context, HostRequest, is_list_operation=False)
1015 .where(HostRequest.conversation_id == request.host_request_id)
1016 .where(
1017 or_(HostRequest.initiator_user_id == context.user_id, HostRequest.recipient_user_id == context.user_id)
1018 )
1019 ).scalar_one_or_none()
1021 if not host_request: 1021 ↛ 1022line 1021 didn't jump to line 1022 because the condition on line 1021 was never true
1022 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
1024 if context.user_id == host_request.initiator_user_id: 1024 ↛ 1027line 1024 didn't jump to line 1027 because the condition on line 1024 was always true
1025 host_request.is_initiator_archived = request.is_archived
1026 else:
1027 host_request.is_recipient_archived = request.is_archived
1029 return requests_pb2.SetHostRequestArchiveStatusRes(
1030 host_request_id=host_request.conversation_id,
1031 is_archived=request.is_archived,
1032 )
1034 def GetResponseRate(
1035 self, request: requests_pb2.GetResponseRateReq, context: CouchersContext, session: Session
1036 ) -> requests_pb2.GetResponseRateRes:
1037 user_res = session.execute(
1038 select(User.id, UserResponseRate)
1039 .outerjoin(UserResponseRate, UserResponseRate.user_id == User.id)
1040 .where(users_visible(context, User))
1041 .where(User.id == request.user_id)
1042 ).one_or_none()
1044 # if user doesn't exist, return None
1045 if not user_res:
1046 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "user_not_found")
1048 user, response_rates = user_res
1049 return requests_pb2.GetResponseRateRes(**response_rate_to_pb(response_rates)) # type: ignore[arg-type]
1051 def SendHostRequestFeedback(
1052 self, request: requests_pb2.SendHostRequestFeedbackReq, context: CouchersContext, session: Session
1053 ) -> empty_pb2.Empty:
1054 host_request = session.execute(
1055 where_moderated_content_visible(select(HostRequest), context, HostRequest, is_list_operation=False)
1056 .where(HostRequest.conversation_id == request.host_request_id)
1057 .where(HostRequest.recipient_user_id == context.user_id)
1058 ).scalar_one_or_none()
1060 if not host_request:
1061 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found")
1063 feedback = session.execute(
1064 select(HostRequestFeedback)
1065 .where(HostRequestFeedback.host_request_id == host_request.conversation_id)
1066 .where(HostRequestFeedback.from_user_id == context.user_id)
1067 ).scalar_one_or_none()
1069 if feedback:
1070 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "already_left_host_request_feedback")
1072 session.add(
1073 HostRequestFeedback(
1074 host_request_id=host_request.conversation_id,
1075 from_user_id=host_request.recipient_user_id,
1076 to_user_id=host_request.initiator_user_id,
1077 request_quality=hostrequestquality2sql.get(request.host_request_quality),
1078 decline_reason=request.decline_reason,
1079 )
1080 )
1081 quality = hostrequestquality2sql.get(request.host_request_quality)
1082 log_event(
1083 context,
1084 session,
1085 "host_request.feedback_submitted",
1086 {
1087 "host_request_id": host_request.conversation_id,
1088 "surfer_id": host_request.initiator_user_id,
1089 "host_id": host_request.recipient_user_id,
1090 "request_quality": quality.name if quality else None,
1091 "has_decline_reason": bool(request.decline_reason),
1092 "host_city": host_request.hosting_city,
1093 },
1094 )
1096 return empty_pb2.Empty()