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
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-03 15:30 +0000
1import logging
3import redis
4from sqlalchemy import case, delete, func, or_, select, update
5from sqlalchemy.ext.asyncio import AsyncSession
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)
30from .schema import (
31 UserCreate,
32 UserDeleteBatchResponse,
33 UserDeleteResult,
34 UserPagination,
35 UserResponse,
36 UserUpdate,
37)
39logger = logging.getLogger("users")
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 )
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)
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 )
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))
81 has_role_join = False
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
90 if not settings.SHOW_SUPER_ADMIN:
91 query = query.where(Users.id.not_in(_super_admin_user_ids_subquery()))
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())
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))
140 if not settings.SHOW_SUPER_ADMIN:
141 count_query = count_query.where(Users.id.not_in(_super_admin_user_ids_subquery()))
143 total_result = await db.execute(count_query)
144 total = total_result.scalar()
146 offset = (page - 1) * per_page
147 query = query.offset(offset).limit(per_page)
149 result = await db.execute(query)
150 users = result.scalars().all()
152 if not users:
153 return UserPagination(users=[], total=0, page=page, per_page=per_page, total_pages=0)
155 user_roles = await _get_user_roles_map(db, [user.id for user in users])
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)
173 total_pages = (total + per_page - 1) // per_page
175 return UserPagination(
176 users=user_responses, total=total, page=page, per_page=per_page, total_pages=total_pages
177 )
179 except Exception as e:
180 raise ServerException(f"Failed to retrieve users: {str(e)}")
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")
200 if user_data.role:
201 await _assert_can_manage_user_role(db, actor_user_id, user_data.role)
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 )
212 db.add(user)
213 await db.commit()
214 await db.refresh(user)
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)
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 )
235 except ConflictException, NotFoundException, AuthorizationException:
236 raise
237 except Exception as e:
238 raise ServerException(f"Failed to create user: {str(e)}")
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")
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")
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")
269 if actor_user_id != user_id:
270 await _assert_can_manage_target_user(db, actor_user_id, user_id)
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 )
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
283 for field, value in update_data.items():
284 setattr(user, field, value)
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 )
300 await db.commit()
301 await db.refresh(user)
303 if role_update_requested:
304 await _update_user_role(db, user_id, user_data.role)
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
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 )
330 except ConflictException, NotFoundException, AuthorizationException:
331 raise
332 except Exception as e:
333 raise ServerException(f"Failed to update user: {str(e)}")
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
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)
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())
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
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
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
401 # Clear user sessions and tokens before deletion
402 await clear_user_all_sessions(db, redis_client, user_id)
404 # Delete related records first to avoid foreign key constraints
405 await _delete_user_related_records(db, user_id)
407 # Delete the user
408 await db.execute(delete(Users).where(Users.id == user_id))
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
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
430 # Commit all successful deletions
431 if success_count > 0:
432 await db.commit()
434 return UserDeleteBatchResponse(
435 results=results,
436 total_users=len(user_ids),
437 success_count=success_count,
438 failed_count=failed_count,
439 )
441 except Exception as e:
442 raise ServerException(f"Failed to delete users: {str(e)}")
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")
455 user.hash_password = await hash_password(new_password)
456 user.password_reset_required = True
457 await db.commit()
459 await clear_user_all_sessions(db, redis_client, user_id)
461 return True
463 except NotFoundException:
464 raise
465 except Exception as e:
466 raise ServerException(f"Failed to reset password: {str(e)}")
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 {}
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
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()
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
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
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")
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")
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")
544 if await check_user_has_super_role(actor_user_id, db):
545 return
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")
551 actor_level = await get_user_role_level(actor_user_id, db)
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 )
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")
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")
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
582 role_mapping = RoleMapper(user_id=user_id, role_id=role.id)
583 db.add(role_mapping)
584 await db.commit()
586 except NotFoundException:
587 raise
588 except Exception as e:
589 raise ServerException(f"Failed to assign role: {str(e)}")
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))
597 if role_name:
598 await _assign_user_role(db, user_id, role_name)
600 await db.commit()
602 except Exception as e:
603 raise ServerException(f"Failed to update user role: {str(e)}")
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))
612 # Delete user sessions
613 await db.execute(delete(UserSessions).where(UserSessions.user_id == user_id))
615 # Delete role mappings
616 await db.execute(delete(RoleMapper).where(RoleMapper.user_id == user_id))
618 # Delete password reset tokens
619 await db.execute(delete(PasswordResetTokens).where(PasswordResetTokens.user_id == user_id))
621 # Delete email verification tokens
622 await db.execute(
623 delete(EmailVerificationTokens).where(EmailVerificationTokens.user_id == user_id)
624 )
626 except Exception as e:
627 raise ServerException(f"Failed to delete user related records: {str(e)}")