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

1import logging 

2from datetime import timedelta 

3 

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_ 

9 

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) 

52 

53logger = logging.getLogger(__name__) 

54 

55DEFAULT_PAGINATION_LENGTH = 10 

56MAX_PAGE_SIZE = 50 

57 

58 

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} 

66 

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} 

74 

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} 

80 

81 

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 ) 

111 

112 

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

122 

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

129 

130 lat, lng = get_coordinates(host_request.hosting_location) 

131 

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

142 

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 ) 

169 

170 

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 

176 

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

182 

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 ) 

191 

192 

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 

199 

200 

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") 

208 

209 if request.host_user_id == context.user_id: 

210 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "cant_request_self") 

211 

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") 

218 

219 from_date = parse_date(request.from_date) 

220 to_date = parse_date(request.to_date) 

221 

222 if not from_date or not to_date: 

223 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_date") 

224 

225 today = today_in_timezone(recipient.timezone) 

226 

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") 

230 

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") 

234 

235 # No need to check today > to_date 

236 

237 if from_date - today > timedelta(days=365): 

238 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_from_after_one_year") 

239 

240 if to_date - from_date > timedelta(days=365): 

241 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "date_to_after_one_year") 

242 

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 ) 

250 

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 ) 

260 

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") 

293 

294 conversation = Conversation() 

295 session.add(conversation) 

296 session.flush() 

297 

298 session.add( 

299 Message( 

300 conversation_id=conversation.id, 

301 author_id=context.user_id, 

302 message_type=MessageType.chat_created, 

303 ) 

304 ) 

305 

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

314 

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 ) 

322 

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

341 

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 ) 

355 

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 ) 

376 

377 return requests_pb2.CreateHostRequestRes(host_request_id=host_request.conversation_id) 

378 

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

402 

403 if not host_request: 

404 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found") 

405 

406 return host_request_to_pb(host_request, session, context) 

407 

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") 

414 

415 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH 

416 pagination = min(pagination, MAX_PAGE_SIZE) 

417 

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) 

441 

442 sort_by_from_date = request.sort_by == requests_pb2.HOST_REQUEST_SORT_BY_FROM_DATE 

443 

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

461 

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 ) 

483 

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

496 

497 if request.status_in: 

498 statement = statement.where(HostRequest.status.in_([api2hostrequeststatus[s] for s in request.status_in])) 

499 

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

506 

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 ) 

531 

532 no_more = len(results) <= pagination 

533 

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 

542 

543 return requests_pb2.ListHostRequestsRes( 

544 next_page_token=next_page_token, no_more=no_more, host_requests=host_requests 

545 ) 

546 

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

555 

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

572 

573 if not host_request: 

574 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found") 

575 

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") 

578 

579 if request.status == conversations_pb2.HOST_REQUEST_STATUS_PENDING: 

580 context.abort_with_error_code(grpc.StatusCode.PERMISSION_DENIED, "invalid_host_request_status") 

581 

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") 

584 

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 ) 

590 

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

606 

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 ) 

619 

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 ) 

636 

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

648 

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 ) 

661 

662 count_host_response(host_request.initiator_user_id, "rejected") 

663 

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 ) 

679 

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

690 

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 ) 

703 

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 ) 

720 

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

731 

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 ) 

744 

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 ) 

761 

762 session.add(control_message) 

763 

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 ) 

771 

772 session.add(latest_message) 

773 else: 

774 latest_message = control_message 

775 

776 session.flush() 

777 

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

783 

784 return empty_pb2.Empty() 

785 

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

794 

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") 

797 

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") 

800 

801 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH 

802 pagination = min(pagination, MAX_PAGE_SIZE) 

803 

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 ) 

815 

816 no_more = len(messages) <= pagination 

817 

818 last_message_id = min(m.id if m else 1 for m in messages[:pagination]) if len(messages) > 0 else 0 

819 

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 ) 

825 

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

836 

837 if not host_request: 

838 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found") 

839 

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") 

842 

843 if host_request.recipient_user_id == context.user_id: 

844 _possibly_observe_first_response_time(session, host_request, context.user_id, "message") 

845 

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 ) 

852 

853 session.add(message) 

854 session.flush() 

855 

856 if host_request.initiator_user_id == context.user_id: 

857 host_request.initiator_last_seen_message_id = message.id 

858 

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 ) 

873 

874 else: 

875 host_request.recipient_last_seen_message_id = message.id 

876 

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 ) 

891 

892 session.commit() 

893 

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 ) 

908 

909 return empty_pb2.Empty() 

910 

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") 

916 

917 if request.newest_message_id == 0: 

918 context.abort_with_error_code(grpc.StatusCode.INVALID_ARGUMENT, "invalid_message") 

919 

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") 

922 

923 pagination = request.number if request.number > 0 else DEFAULT_PAGINATION_LENGTH 

924 pagination = min(pagination, MAX_PAGE_SIZE) 

925 

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 ) 

938 

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 ) 

947 

948 statement = statement.order_by(Message.id.asc()).limit(pagination + 1) 

949 res = session.execute(statement).all() 

950 

951 no_more = len(res) <= pagination 

952 

953 last_message_id = min(m.Message.id if m else 1 for m in res[:pagination]) if len(res) > 0 else 0 # TODO 

954 

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 ) 

966 

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

975 

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") 

978 

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") 

981 

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 

990 

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 ) 

1006 

1007 session.commit() 

1008 return empty_pb2.Empty() 

1009 

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

1020 

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") 

1023 

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 

1028 

1029 return requests_pb2.SetHostRequestArchiveStatusRes( 

1030 host_request_id=host_request.conversation_id, 

1031 is_archived=request.is_archived, 

1032 ) 

1033 

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

1043 

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") 

1047 

1048 user, response_rates = user_res 

1049 return requests_pb2.GetResponseRateRes(**response_rate_to_pb(response_rates)) # type: ignore[arg-type] 

1050 

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

1059 

1060 if not host_request: 

1061 context.abort_with_error_code(grpc.StatusCode.NOT_FOUND, "host_request_not_found") 

1062 

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

1068 

1069 if feedback: 

1070 context.abort_with_error_code(grpc.StatusCode.FAILED_PRECONDITION, "already_left_host_request_feedback") 

1071 

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 ) 

1095 

1096 return empty_pb2.Empty()