Coverage for api/users/services.py: 95.58%

294 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-03 15:30 +0000

1import logging 

2 

3import redis 

4from sqlalchemy import case, delete, func, or_, select, update 

5from sqlalchemy.ext.asyncio import AsyncSession 

6 

7from core.config import settings 

8from core.permissions import Permission 

9from core.rbac import ( 

10 check_user_has_super_role, 

11 get_user_attributes, 

12 get_user_role_level, 

13 is_super_admin_role_name, 

14) 

15from core.security import clear_user_all_sessions, hash_password 

16from models.email_verification_tokens import EmailVerificationTokens 

17from models.login_logs import LoginLogs 

18from models.password_reset_tokens import PasswordResetTokens 

19from models.role_mapper import RoleMapper 

20from models.roles import Roles 

21from models.user_sessions import UserSessions 

22from models.users import Users 

23from utils.custom_exception import ( 

24 AuthorizationException, 

25 ConflictException, 

26 NotFoundException, 

27 ServerException, 

28) 

29 

30from .schema import ( 

31 UserCreate, 

32 UserDeleteBatchResponse, 

33 UserDeleteResult, 

34 UserPagination, 

35 UserResponse, 

36 UserUpdate, 

37) 

38 

39logger = logging.getLogger("users") 

40 

41 

42def _super_admin_user_ids_subquery(): 

43 """Subquery of user IDs assigned to the system super-admin role.""" 

44 return ( 

45 select(RoleMapper.user_id) 

46 .join(Roles, Roles.id == RoleMapper.role_id) 

47 .where(Roles.name == settings.DEFAULT_SUPER_ADMIN_ROLE) 

48 ) 

49 

50 

51async def get_all_users( 

52 db: AsyncSession, 

53 keyword: str | None = None, 

54 status: str | None = None, 

55 role: str | None = None, 

56 page: int = 1, 

57 per_page: int = 10, 

58 sort_by: str | None = None, 

59 desc: bool = False, 

60) -> UserPagination: 

61 """Get all users list""" 

62 try: 

63 query = select(Users) 

64 

65 if keyword: 

66 query = query.where( 

67 or_( 

68 Users.first_name.ilike(f"%{keyword}%"), 

69 Users.last_name.ilike(f"%{keyword}%"), 

70 Users.email.ilike(f"%{keyword}%"), 

71 ) 

72 ) 

73 

74 if status: 

75 status_list = [s.strip().lower() == "true" for s in status.split(",")] 

76 if len(status_list) == 1: 

77 query = query.where(Users.status == status_list[0]) 

78 else: 

79 query = query.where(Users.status.in_(status_list)) 

80 

81 has_role_join = False 

82 

83 if role: 

84 role_list = [r.strip() for r in role.split(",")] 

85 query = query.join(RoleMapper, Users.id == RoleMapper.user_id) 

86 query = query.join(Roles, RoleMapper.role_id == Roles.id) 

87 query = query.where(Roles.name.in_(role_list)) 

88 has_role_join = True 

89 

90 if not settings.SHOW_SUPER_ADMIN: 

91 query = query.where(Users.id.not_in(_super_admin_user_ids_subquery())) 

92 

93 if sort_by: 

94 if sort_by == "role": 

95 # For role sorting, use LEFT JOIN with distinct to avoid duplicates 

96 if not has_role_join: 

97 query = query.outerjoin(RoleMapper, Users.id == RoleMapper.user_id) 

98 query = query.outerjoin(Roles, RoleMapper.role_id == Roles.id) 

99 # Use distinct to avoid duplicate rows when a user has multiple roles 

100 query = query.distinct() 

101 null_order = case((Roles.name.is_(None), 1), else_=0) 

102 if desc: 

103 query = query.order_by(null_order.asc(), Roles.name.desc(), Users.id) 

104 else: 

105 query = query.order_by(null_order.asc(), Roles.name.asc(), Users.id) 

106 else: 

107 # For other fields, use direct column access 

108 sort_column = getattr(Users, sort_by, None) 

109 if sort_column: 

110 if desc: 

111 query = query.order_by(sort_column.desc()) 

112 else: 

113 query = query.order_by(sort_column.asc()) 

114 else: 

115 query = query.order_by(Users.created_at.desc()) 

116 else: 

117 query = query.order_by(Users.id.asc()) 

118 

119 count_query = select(func.count(Users.id)) 

120 if keyword: 

121 count_query = count_query.where( 

122 or_( 

123 Users.first_name.ilike(f"%{keyword}%"), 

124 Users.last_name.ilike(f"%{keyword}%"), 

125 Users.email.ilike(f"%{keyword}%"), 

126 ) 

127 ) 

128 if status: 

129 status_list = [s.strip().lower() == "true" for s in status.split(",")] 

130 if len(status_list) == 1: 

131 count_query = count_query.where(Users.status == status_list[0]) 

132 else: 

133 count_query = count_query.where(Users.status.in_(status_list)) 

134 if role: 

135 role_list = [r.strip() for r in role.split(",")] 

136 count_query = count_query.join(RoleMapper, Users.id == RoleMapper.user_id) 

137 count_query = count_query.join(Roles, RoleMapper.role_id == Roles.id) 

138 count_query = count_query.where(Roles.name.in_(role_list)) 

139 

140 if not settings.SHOW_SUPER_ADMIN: 

141 count_query = count_query.where(Users.id.not_in(_super_admin_user_ids_subquery())) 

142 

143 total_result = await db.execute(count_query) 

144 total = total_result.scalar() 

145 

146 offset = (page - 1) * per_page 

147 query = query.offset(offset).limit(per_page) 

148 

149 result = await db.execute(query) 

150 users = result.scalars().all() 

151 

152 if not users: 

153 return UserPagination(users=[], total=0, page=page, per_page=per_page, total_pages=0) 

154 

155 user_roles = await _get_user_roles_map(db, [user.id for user in users]) 

156 

157 user_responses = [] 

158 for user in users: 

159 role_name, role_level = user_roles.get(user.id, (None, None)) 

160 user_response = UserResponse( 

161 id=user.id, 

162 email=user.email, 

163 first_name=user.first_name, 

164 last_name=user.last_name, 

165 phone=user.phone, 

166 status=user.status, 

167 created_at=user.created_at, 

168 role=role_name, 

169 role_level=role_level, 

170 ) 

171 user_responses.append(user_response) 

172 

173 total_pages = (total + per_page - 1) // per_page 

174 

175 return UserPagination( 

176 users=user_responses, total=total, page=page, per_page=per_page, total_pages=total_pages 

177 ) 

178 

179 except Exception as e: 

180 raise ServerException(f"Failed to retrieve users: {str(e)}") 

181 

182 

183async def create_user( 

184 db: AsyncSession, 

185 user_data: UserCreate, 

186 actor_user_id: str, 

187) -> UserResponse: 

188 """Create a new user""" 

189 try: 

190 # Check if the email already exists 

191 result = await db.execute( 

192 select(Users).where( 

193 or_(Users.email == user_data.email, Users.pending_email == user_data.email) 

194 ) 

195 ) 

196 existing_user = result.scalar_one_or_none() 

197 if existing_user: 

198 raise ConflictException("Email already exists") 

199 

200 if user_data.role: 

201 await _assert_can_manage_user_role(db, actor_user_id, user_data.role) 

202 

203 user = Users( 

204 first_name=user_data.first_name, 

205 last_name=user_data.last_name, 

206 email=user_data.email, 

207 phone=user_data.phone, 

208 hash_password=await hash_password(user_data.password), 

209 status=user_data.status, 

210 ) 

211 

212 db.add(user) 

213 await db.commit() 

214 await db.refresh(user) 

215 

216 user_role = None 

217 user_role_level = None 

218 if user_data.role: 

219 await _assign_user_role(db, user.id, user_data.role) 

220 user_role = user_data.role 

221 user_role_level = await _get_role_level_by_name(db, user_data.role) 

222 

223 return UserResponse( 

224 id=user.id, 

225 email=user.email, 

226 first_name=user.first_name, 

227 last_name=user.last_name, 

228 phone=user.phone, 

229 status=user.status, 

230 created_at=user.created_at, 

231 role=user_role, 

232 role_level=user_role_level, 

233 ) 

234 

235 except ConflictException, NotFoundException, AuthorizationException: 

236 raise 

237 except Exception as e: 

238 raise ServerException(f"Failed to create user: {str(e)}") 

239 

240 

241async def update_user( 

242 db: AsyncSession, 

243 user_id: str, 

244 user_data: UserUpdate, 

245 actor_user_id: str, 

246) -> UserResponse: 

247 """Update user information""" 

248 try: 

249 result = await db.execute(select(Users).where(Users.id == user_id)) 

250 user = result.scalar_one_or_none() 

251 if not user: 

252 raise NotFoundException("User not found") 

253 

254 # Check if the email is already used by another user 

255 if user_data.email and user_data.email != user.email: 

256 result = await db.execute( 

257 select(Users).where( 

258 or_(Users.email == user_data.email, Users.pending_email == user_data.email), 

259 Users.id != user_id, 

260 ) 

261 ) 

262 if result.scalar_one_or_none(): 

263 raise ConflictException("Email already exists") 

264 

265 role_update_requested = "role" in user_data.model_dump(exclude_unset=True) 

266 if role_update_requested and actor_user_id == user_id: 

267 raise AuthorizationException("Cannot change your own role") 

268 

269 if actor_user_id != user_id: 

270 await _assert_can_manage_target_user(db, actor_user_id, user_id) 

271 

272 if role_update_requested: 

273 await _assert_can_manage_user_role( 

274 db, 

275 actor_user_id, 

276 user_data.role, 

277 target_user_id=user_id, 

278 ) 

279 

280 update_data = user_data.model_dump(exclude_unset=True, exclude={"role"}) 

281 email_changed = "email" in update_data and update_data["email"] != user.email 

282 

283 for field, value in update_data.items(): 

284 setattr(user, field, value) 

285 

286 # Admin-set email is treated as confirmed: drop pending change and unused tokens 

287 if email_changed: 

288 user.pending_email = None 

289 user.email_verified = True 

290 await db.execute( 

291 update(EmailVerificationTokens) 

292 .where( 

293 EmailVerificationTokens.user_id == user.id, 

294 EmailVerificationTokens.token_type == "email_change", 

295 EmailVerificationTokens.is_used.is_(False), 

296 ) 

297 .values(is_used=True) 

298 ) 

299 

300 await db.commit() 

301 await db.refresh(user) 

302 

303 if role_update_requested: 

304 await _update_user_role(db, user_id, user_data.role) 

305 

306 role_query = ( 

307 select(Roles.name, Roles.level) 

308 .join(RoleMapper, Roles.id == RoleMapper.role_id) 

309 .where(RoleMapper.user_id == user.id) 

310 .order_by(Roles.level.desc(), Roles.name.asc()) 

311 .limit(1) 

312 ) 

313 role_result = await db.execute(role_query) 

314 role_row = role_result.one_or_none() 

315 user_role = role_row[0] if role_row else None 

316 user_role_level = role_row[1] if role_row else None 

317 

318 return UserResponse( 

319 id=user.id, 

320 email=user.email, 

321 first_name=user.first_name, 

322 last_name=user.last_name, 

323 phone=user.phone, 

324 status=user.status, 

325 created_at=user.created_at, 

326 role=user_role, 

327 role_level=user_role_level, 

328 ) 

329 

330 except ConflictException, NotFoundException, AuthorizationException: 

331 raise 

332 except Exception as e: 

333 raise ServerException(f"Failed to update user: {str(e)}") 

334 

335 

336async def delete_users( 

337 db: AsyncSession, redis_client: redis.Redis, user_ids: list[str], token: dict | None = None 

338) -> UserDeleteBatchResponse: 

339 """Delete multiple users with detailed batch processing results""" 

340 try: 

341 results = [] 

342 success_count = 0 

343 failed_count = 0 

344 

345 # Get current user ID from token 

346 current_user_id = token.get("sub") if token else None 

347 actor_is_super = False 

348 actor_level = 0 

349 if current_user_id: 

350 actor_is_super = await check_user_has_super_role(current_user_id, db) 

351 if not actor_is_super: 

352 actor_level = await get_user_role_level(current_user_id, db) 

353 

354 # Check which users exist 

355 result = await db.execute(select(Users.id).where(Users.id.in_(user_ids))) 

356 existing_ids = set(result.scalars().all()) 

357 

358 # Process each user ID 

359 for user_id in user_ids: 

360 try: 

361 # Skip if trying to delete own account 

362 if current_user_id and user_id == current_user_id: 

363 results.append( 

364 UserDeleteResult( 

365 user_id=user_id, 

366 status="failed", 

367 message="Cannot delete your own account", 

368 ) 

369 ) 

370 failed_count += 1 

371 continue 

372 

373 if user_id in existing_ids: 

374 if await check_user_has_super_role(user_id, db): 

375 results.append( 

376 UserDeleteResult( 

377 user_id=user_id, 

378 status="failed", 

379 message="Cannot delete a system super-admin user", 

380 ) 

381 ) 

382 failed_count += 1 

383 continue 

384 

385 if not actor_is_super: 

386 target_level = await get_user_role_level(user_id, db) 

387 if target_level > actor_level: 

388 results.append( 

389 UserDeleteResult( 

390 user_id=user_id, 

391 status="failed", 

392 message=( 

393 "Cannot delete a user with a higher role level " 

394 "than your own" 

395 ), 

396 ) 

397 ) 

398 failed_count += 1 

399 continue 

400 

401 # Clear user sessions and tokens before deletion 

402 await clear_user_all_sessions(db, redis_client, user_id) 

403 

404 # Delete related records first to avoid foreign key constraints 

405 await _delete_user_related_records(db, user_id) 

406 

407 # Delete the user 

408 await db.execute(delete(Users).where(Users.id == user_id)) 

409 

410 results.append( 

411 UserDeleteResult( 

412 user_id=user_id, status="success", message="User deleted successfully" 

413 ) 

414 ) 

415 success_count += 1 

416 else: 

417 results.append( 

418 UserDeleteResult(user_id=user_id, status="failed", message="User not found") 

419 ) 

420 failed_count += 1 

421 

422 except Exception as e: 

423 results.append( 

424 UserDeleteResult( 

425 user_id=user_id, status="failed", message=f"Failed to delete user: {str(e)}" 

426 ) 

427 ) 

428 failed_count += 1 

429 

430 # Commit all successful deletions 

431 if success_count > 0: 

432 await db.commit() 

433 

434 return UserDeleteBatchResponse( 

435 results=results, 

436 total_users=len(user_ids), 

437 success_count=success_count, 

438 failed_count=failed_count, 

439 ) 

440 

441 except Exception as e: 

442 raise ServerException(f"Failed to delete users: {str(e)}") 

443 

444 

445async def reset_user_password( 

446 db: AsyncSession, redis_client: redis.Redis, user_id: str, new_password: str 

447) -> bool: 

448 """Reset user password and logout all devices""" 

449 try: 

450 result = await db.execute(select(Users).where(Users.id == user_id)) 

451 user = result.scalar_one_or_none() 

452 if not user: 

453 raise NotFoundException("User not found") 

454 

455 user.hash_password = await hash_password(new_password) 

456 user.password_reset_required = True 

457 await db.commit() 

458 

459 await clear_user_all_sessions(db, redis_client, user_id) 

460 

461 return True 

462 

463 except NotFoundException: 

464 raise 

465 except Exception as e: 

466 raise ServerException(f"Failed to reset password: {str(e)}") 

467 

468 

469async def _get_user_roles_map( 

470 db: AsyncSession, user_ids: list[str] 

471) -> dict[str, tuple[str | None, int | None]]: 

472 """Batch-load primary role name/level per user (highest level wins).""" 

473 if not user_ids: 

474 return {} 

475 

476 roles_query = ( 

477 select(RoleMapper.user_id, Roles.name, Roles.level) 

478 .join(Roles, Roles.id == RoleMapper.role_id) 

479 .where(RoleMapper.user_id.in_(user_ids)) 

480 .order_by(Roles.level.desc(), Roles.name.asc()) 

481 ) 

482 role_result = await db.execute(roles_query) 

483 user_roles: dict[str, tuple[str | None, int | None]] = {} 

484 for user_id, role_name, role_level in role_result.all(): 

485 if user_id not in user_roles: 

486 user_roles[user_id] = (role_name, role_level) 

487 return user_roles 

488 

489 

490async def _get_user_role_name(db: AsyncSession, user_id: str) -> str | None: 

491 result = await db.execute( 

492 select(Roles.name) 

493 .join(RoleMapper, Roles.id == RoleMapper.role_id) 

494 .where(RoleMapper.user_id == user_id) 

495 .order_by(Roles.level.desc(), Roles.name.asc()) 

496 .limit(1) 

497 ) 

498 return result.scalar_one_or_none() 

499 

500 

501async def _get_role_level_by_name(db: AsyncSession, role_name: str) -> int | None: 

502 result = await db.execute(select(Roles.level).where(Roles.name == role_name).limit(1)) 

503 level = result.scalar_one_or_none() 

504 return int(level) if level is not None else None 

505 

506 

507async def _assert_can_manage_target_user( 

508 db: AsyncSession, 

509 actor_user_id: str, 

510 target_user_id: str, 

511) -> None: 

512 """Block update/delete when the target user's role level is higher than the actor's.""" 

513 if actor_user_id == target_user_id: 

514 return 

515 if await check_user_has_super_role(actor_user_id, db): 

516 return 

517 

518 actor_level = await get_user_role_level(actor_user_id, db) 

519 target_level = await get_user_role_level(target_user_id, db) 

520 if target_level > actor_level: 

521 raise AuthorizationException("Cannot manage a user with a higher role level than your own") 

522 

523 

524async def _assert_can_manage_user_role( 

525 db: AsyncSession, 

526 actor_user_id: str, 

527 role_name: str | None, 

528 *, 

529 target_user_id: str | None = None, 

530) -> None: 

531 """ 

532 Require manage-roles (or super-admin) to change roles. 

533 The system super-admin role cannot be assigned or removed via API. 

534 Assigned / target role levels cannot exceed the actor's level. 

535 """ 

536 if is_super_admin_role_name(role_name): 

537 raise AuthorizationException("Cannot assign the system super-admin role") 

538 

539 if target_user_id: 

540 current_role = await _get_user_role_name(db, target_user_id) 

541 if is_super_admin_role_name(current_role): 

542 raise AuthorizationException("Cannot change the role of a system super-admin user") 

543 

544 if await check_user_has_super_role(actor_user_id, db): 

545 return 

546 

547 attributes = await get_user_attributes(actor_user_id, db) 

548 if not attributes.get(Permission.MANAGE_ROLES.value, False): 

549 raise AuthorizationException("Permission denied to assign roles") 

550 

551 actor_level = await get_user_role_level(actor_user_id, db) 

552 

553 if target_user_id: 

554 target_level = await get_user_role_level(target_user_id, db) 

555 if target_level > actor_level: 

556 raise AuthorizationException( 

557 "Cannot manage a user with a higher role level than your own" 

558 ) 

559 

560 if role_name: 

561 new_role_level = await _get_role_level_by_name(db, role_name) 

562 if new_role_level is None: 

563 raise NotFoundException(f"Role '{role_name}' not found") 

564 if new_role_level > actor_level: 

565 raise AuthorizationException("Cannot assign a role with a higher level than your own") 

566 

567 

568async def _assign_user_role(db: AsyncSession, user_id: str, role_name: str) -> None: 

569 """Assign a role to a user""" 

570 try: 

571 role_result = await db.execute(select(Roles).where(Roles.name == role_name)) 

572 role = role_result.scalar_one_or_none() 

573 if not role: 

574 raise NotFoundException(f"Role '{role_name}' not found") 

575 

576 existing_mapping = await db.execute( 

577 select(RoleMapper).where(RoleMapper.user_id == user_id, RoleMapper.role_id == role.id) 

578 ) 

579 if existing_mapping.scalar_one_or_none(): 

580 return 

581 

582 role_mapping = RoleMapper(user_id=user_id, role_id=role.id) 

583 db.add(role_mapping) 

584 await db.commit() 

585 

586 except NotFoundException: 

587 raise 

588 except Exception as e: 

589 raise ServerException(f"Failed to assign role: {str(e)}") 

590 

591 

592async def _update_user_role(db: AsyncSession, user_id: str, role_name: str | None) -> None: 

593 """Update user role (remove existing and assign new one)""" 

594 try: 

595 await db.execute(delete(RoleMapper).where(RoleMapper.user_id == user_id)) 

596 

597 if role_name: 

598 await _assign_user_role(db, user_id, role_name) 

599 

600 await db.commit() 

601 

602 except Exception as e: 

603 raise ServerException(f"Failed to update user role: {str(e)}") 

604 

605 

606async def _delete_user_related_records(db: AsyncSession, user_id: str) -> None: 

607 """Delete all records related to a user to avoid foreign key constraints""" 

608 try: 

609 # Delete login logs 

610 await db.execute(delete(LoginLogs).where(LoginLogs.user_id == user_id)) 

611 

612 # Delete user sessions 

613 await db.execute(delete(UserSessions).where(UserSessions.user_id == user_id)) 

614 

615 # Delete role mappings 

616 await db.execute(delete(RoleMapper).where(RoleMapper.user_id == user_id)) 

617 

618 # Delete password reset tokens 

619 await db.execute(delete(PasswordResetTokens).where(PasswordResetTokens.user_id == user_id)) 

620 

621 # Delete email verification tokens 

622 await db.execute( 

623 delete(EmailVerificationTokens).where(EmailVerificationTokens.user_id == user_id) 

624 ) 

625 

626 except Exception as e: 

627 raise ServerException(f"Failed to delete user related records: {str(e)}")