mirror of
https://github.com/taylanbakircioglu/haproxy-openmanager.git
synced 2026-09-16 15:45:11 +00:00
bd6a31cb0d
Adds opt-in TOTP-based Multi-Factor Authentication that is fully
backwards compatible with existing logins. Operators choose to enable
MFA per account; nothing changes for users who do not opt in.
Highlights
==========
* RFC 6238 TOTP (6 digits, 30s period, SHA1) with ±30s skew tolerance,
compatible with Microsoft / Google Authenticator, Authy, Duo, 1Password.
* Per-step replay protection (`mfa_last_used_totp_step`) so a captured
code cannot be reused inside the same window.
* Fernet-encrypted TOTP secrets at rest, key resolution via
`MFA_ENCRYPTION_KEY` env (HKDF-derived from `SECRET_KEY` as fallback).
* 10 single-use, bcrypt-hashed backup codes per user, formatted
`XXXX-YYYY` from a confusion-free alphabet (no 0/O/1/I/L).
* Two-step login flow: `POST /api/auth/login` returns `mfa_required`
+ `mfa_token`, then `POST /api/auth/login/mfa-verify` accepts a TOTP
code OR a backup code. JWT is minted only after MFA succeeds.
* Self-service: users enable / disable MFA from their own row in the
Users page; admins reset (single user or bulk) but never enable on
behalf of someone else (matches AWS IAM / GitHub / Google Workspace).
* Bulk emergency reset CLI: `scripts/admin-mfa-reset-all.sh`.
Security hardening
==================
* Atomic transactions with `SELECT … FOR UPDATE` on `mfa_pending_logins`
and `users` rows so concurrent verify / enroll calls cannot race.
* `/api/mfa/enroll/start` refuses re-enrollment when MFA is already on
(prevents silent secret rotation via a stolen JWT).
* Pydantic `ValidationError` messages are sanitized before reaching the
audit log so request bodies (TOTP / backup codes in flight) never
appear in plaintext.
* Slowapi rate limits are per-USER, not per-IP, with a trusted-proxy
XFF strategy so a single ingress address cannot exhaust the bucket
for thousands of operators (`MFA_TRUSTED_PROXY_CIDRS`,
`MFA_RATE_LIMIT_*` env-overridable).
* Login query now scopes to `is_active = TRUE` so a soft-deleted row
with the same username can no longer occlude the active user
(also closes a small account-enumeration side channel).
Database
========
Additive migrations (idempotent `ADD COLUMN IF NOT EXISTS`,
`CREATE TABLE IF NOT EXISTS`):
- users: mfa_enabled, mfa_method, mfa_secret_encrypted,
mfa_enrolled_at, mfa_last_used_at, mfa_last_used_totp_step
- mfa_backup_codes (user_id ON DELETE CASCADE)
- mfa_pending_logins (user_id ON DELETE CASCADE, challenge_token,
attempts, expires_at)
- mfa_pending_enrollments (user_id ON DELETE CASCADE)
Frontend
========
* Login page becomes a 3-phase state machine
(credentials → MFA → submitting); legacy single-step login is
preserved for users who haven't enrolled.
* New MFAEnrollModal (3-step wizard: QR + secret → verify → backup
codes) using `qrcode.react`.
* Users page shows MFA column + per-row enable/disable/reset actions.
Admins viewing other users with MFA off see a non-actionable info
icon explaining that only the user themselves can enable MFA.
Deployment
==========
* `MFA_ENCRYPTION_KEY` is added to `k8s/manifests/03-secrets.yaml` as
a placeholder; `SECRET_KEY` is also placeholder-ized so both are
injected by the existing pipeline pattern (sed-replace + apply).
* No new build-time env vars are required for the frontend. The SPA
uses `window.location.host` for `/api/*` and is routed by the
existing nginx ingress configuration.
* `frontend/.dockerignore` ensures host `.env*` files cannot bleed
into the production bundle.
Tests
=====
* New unit suites:
- `test_mfa_service.py` (TOTP, encryption, backup codes)
- `test_mfa_backwards_compat.py` (regression — non-MFA flow unchanged)
- `test_mfa_rate_limits.py` (env override + dataclass immutability)
- `test_mfa_rate_limit_key.py` (JWT key, trusted-proxy XFF, fallbacks)
* All existing 1000+ unit tests continue to pass.
Documentation
=============
* README MFA section (overview, day-to-day operations, emergency
reset CLI, env variables, rate-limit tuning).
* `scripts/README.md` documents the bulk reset script.
Issue: #18
1303 lines
50 KiB
Python
1303 lines
50 KiB
Python
from fastapi import APIRouter, HTTPException, Header, Depends, Request
|
|
from typing import Optional, List
|
|
import logging
|
|
import json
|
|
|
|
from database.connection import get_database_connection, close_database_connection
|
|
from auth_middleware import get_current_user_from_token
|
|
from utils.activity_log import log_user_activity
|
|
from models.user import UserCreate, UserPasswordUpdate
|
|
|
|
router = APIRouter(prefix="/api", tags=["users", "servers"])
|
|
logger = logging.getLogger(__name__)
|
|
|
|
@router.get("/users")
|
|
async def get_users(authorization: str = Header(None)):
|
|
"""Get all users"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# R18c audit fix (round 5 #1 — KRITIK info leak): the only
|
|
# caller in the UI today is the admin User Management page;
|
|
# the endpoint exposes username, email, phone, full_name,
|
|
# is_admin, roles[], cluster_ids and timestamps for every
|
|
# active operator on the platform. Pre-fix any
|
|
# authenticated user (cluster reader, ssl reader, etc.)
|
|
# could enumerate the full operator roster, including
|
|
# admin emails for phishing and is_admin flags for target
|
|
# selection. Restrict to admins.
|
|
if not current_user.get("is_admin"):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Listing all users requires administrator privileges."
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Get users with their roles (only active users)
|
|
try:
|
|
users = await conn.fetch("""
|
|
SELECT u.id, u.username, u.email, u.full_name, u.phone, u.role, u.is_active,
|
|
u.is_admin, u.is_verified, u.created_at, u.updated_at, u.last_login_at,
|
|
COALESCE(u.mfa_enabled, FALSE) AS mfa_enabled
|
|
FROM users u
|
|
WHERE u.is_active = TRUE
|
|
ORDER BY u.username
|
|
""")
|
|
except Exception as schema_error:
|
|
logger.warning(f"Schema error in users query, using fallback: {schema_error}")
|
|
# Fallback query with minimal columns (pre-MFA-migration deploys)
|
|
users = await conn.fetch("""
|
|
SELECT id, username, email, is_active, is_admin, created_at
|
|
FROM users
|
|
WHERE is_active = TRUE
|
|
ORDER BY username
|
|
""")
|
|
|
|
# Get user roles for each user
|
|
users_list = []
|
|
for user in users:
|
|
user_dict = dict(user)
|
|
# Remove password_hash if it exists
|
|
user_dict.pop('password_hash', None)
|
|
|
|
# Format datetime fields with UTC timezone indicator
|
|
if user_dict.get('created_at'):
|
|
user_dict['created_at'] = user_dict['created_at'].isoformat().replace('+00:00', 'Z')
|
|
if user_dict.get('updated_at'):
|
|
user_dict['updated_at'] = user_dict['updated_at'].isoformat().replace('+00:00', 'Z')
|
|
if user_dict.get('last_login_at'):
|
|
user_dict['last_login_at'] = user_dict['last_login_at'].isoformat().replace('+00:00', 'Z')
|
|
|
|
# Get user roles
|
|
try:
|
|
user_roles = await conn.fetch("""
|
|
SELECT r.id, r.name, r.display_name, r.description
|
|
FROM user_roles ur
|
|
JOIN roles r ON ur.role_id = r.id
|
|
WHERE ur.user_id = $1 AND r.is_active = true
|
|
ORDER BY r.display_name
|
|
""", user['id'])
|
|
|
|
user_dict['roles'] = [dict(role) for role in user_roles]
|
|
except Exception as role_error:
|
|
logger.warning(f"Could not fetch roles for user {user['id']}: {role_error}")
|
|
user_dict['roles'] = []
|
|
|
|
users_list.append(user_dict)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"users": users_list,
|
|
"total": len(users_list)
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to get users: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to get users: {str(e)}")
|
|
|
|
@router.get("/roles")
|
|
async def get_roles(authorization: str = Header(None)):
|
|
"""Get all roles"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# R18c audit fix (round 5 #2 — KRITIK info leak): the role
|
|
# listing exposes the FULL `permissions` blob and
|
|
# `cluster_ids` for every role. Pre-fix any authenticated
|
|
# user could read the platform's RBAC layout — invaluable
|
|
# reconnaissance for an attacker planning a privilege
|
|
# escalation. Mutations on this endpoint are admin-only;
|
|
# the read path now matches.
|
|
if not current_user.get("is_admin"):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Listing roles requires administrator privileges."
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if roles table exists and get roles
|
|
try:
|
|
roles = await conn.fetch("""
|
|
SELECT id, name, display_name, description, permissions, cluster_ids,
|
|
is_active, is_system, created_at, updated_at
|
|
FROM roles
|
|
WHERE is_active = TRUE
|
|
ORDER BY name
|
|
""")
|
|
except Exception as schema_error:
|
|
logger.warning(f"Roles table not found or schema error: {schema_error}")
|
|
# Return default roles if table doesn't exist
|
|
roles = []
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Convert to list of dicts
|
|
roles_list = [dict(role) for role in roles]
|
|
|
|
return {
|
|
"roles": roles_list,
|
|
"total": len(roles_list)
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to get roles: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to get roles: {str(e)}")
|
|
|
|
@router.post("/users", summary="Create User", response_description="User created successfully")
|
|
async def create_user(user_data: UserCreate, authorization: str = Header(None)):
|
|
"""
|
|
# Create New User
|
|
|
|
Create a new user account with role assignments.
|
|
|
|
## Request Body
|
|
- **username**: Username (required, unique)
|
|
- **email**: Email address (required)
|
|
- **password**: Password (required, min 8 characters)
|
|
- **is_active**: Active status (default: true)
|
|
|
|
## Example Request
|
|
```bash
|
|
curl -X POST "{BASE_URL}/api/users/users" \\
|
|
-H "Authorization: Bearer eyJhbGciOiJIUz..." \\
|
|
-H "Content-Type: application/json" \\
|
|
-d '{
|
|
"username": "johndoe",
|
|
"email": "john@example.com",
|
|
"password": "securePass123",
|
|
"is_active": true
|
|
}'
|
|
```
|
|
|
|
## Example Response
|
|
```json
|
|
{
|
|
"message": "User 'johndoe' created successfully",
|
|
"user_id": 2
|
|
}
|
|
```
|
|
"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can create users
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can create new users"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if username already exists (only active users)
|
|
existing = await conn.fetchrow(
|
|
"SELECT id FROM users WHERE username = $1 AND is_active = TRUE",
|
|
user_data.username
|
|
)
|
|
if existing:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Username already exists")
|
|
|
|
# Check if email already exists (only active users)
|
|
if user_data.email:
|
|
existing_email = await conn.fetchrow(
|
|
"SELECT id FROM users WHERE email = $1 AND is_active = TRUE",
|
|
user_data.email
|
|
)
|
|
if existing_email:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Email already exists")
|
|
|
|
# Hash password
|
|
import bcrypt
|
|
salt = bcrypt.gensalt()
|
|
hashed_password = bcrypt.hashpw(user_data.password.encode('utf-8'), salt)
|
|
|
|
# Check if user should be admin based on roles
|
|
is_admin = False
|
|
role_ids = getattr(user_data, 'role_ids', None)
|
|
if role_ids:
|
|
# Check if any of the assigned roles is super_admin
|
|
admin_roles = await conn.fetch("""
|
|
SELECT id FROM roles
|
|
WHERE id = ANY($1) AND name = 'super_admin'
|
|
""", role_ids)
|
|
is_admin = len(admin_roles) > 0
|
|
|
|
# Create user
|
|
user_id = await conn.fetchval("""
|
|
INSERT INTO users (username, email, full_name, phone, password_hash,
|
|
role, is_active, is_admin, is_verified)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
|
RETURNING id
|
|
""", user_data.username, user_data.email, user_data.full_name, user_data.phone,
|
|
hashed_password.decode('utf-8'), 'user', # Default role (actual roles assigned via role_ids)
|
|
user_data.is_active, is_admin, # is_admin determined by roles
|
|
user_data.is_verified)
|
|
|
|
# Assign roles if provided
|
|
if role_ids:
|
|
for role_id in role_ids:
|
|
await conn.execute("""
|
|
INSERT INTO user_roles (user_id, role_id)
|
|
VALUES ($1, $2)
|
|
ON CONFLICT (user_id, role_id) DO NOTHING
|
|
""", user_id, role_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='create',
|
|
resource_type='user',
|
|
resource_id=str(user_id),
|
|
details={'username': user_data.username}
|
|
)
|
|
|
|
return {
|
|
"message": f"User '{user_data.username}' created successfully",
|
|
"user_id": user_id
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to create user: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to create user: {str(e)}")
|
|
|
|
@router.put("/users/{user_id}")
|
|
async def update_user(user_id: int, user_data: dict, authorization: str = Header(None)):
|
|
"""Update a user"""
|
|
try:
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can update other users
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can update user information"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user exists
|
|
existing_user = await conn.fetchrow("SELECT username, email FROM users WHERE id = $1", user_id)
|
|
if not existing_user:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="User not found")
|
|
|
|
# Build update query dynamically
|
|
update_fields = []
|
|
params = [user_id]
|
|
param_count = 1
|
|
|
|
if 'username' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"username = ${param_count}")
|
|
params.append(user_data['username'])
|
|
|
|
if 'email' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"email = ${param_count}")
|
|
params.append(user_data['email'])
|
|
|
|
if 'full_name' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"full_name = ${param_count}")
|
|
params.append(user_data['full_name'])
|
|
|
|
if 'phone' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"phone = ${param_count}")
|
|
params.append(user_data['phone'])
|
|
|
|
if 'role' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"role = ${param_count}")
|
|
params.append(user_data['role'])
|
|
|
|
if 'is_active' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"is_active = ${param_count}")
|
|
params.append(user_data['is_active'])
|
|
|
|
if 'is_admin' in user_data:
|
|
param_count += 1
|
|
update_fields.append(f"is_admin = ${param_count}")
|
|
params.append(user_data['is_admin'])
|
|
|
|
if 'password' in user_data and user_data['password']:
|
|
import bcrypt
|
|
salt = bcrypt.gensalt()
|
|
hashed_password = bcrypt.hashpw(user_data['password'].encode('utf-8'), salt)
|
|
param_count += 1
|
|
update_fields.append(f"password_hash = ${param_count}")
|
|
params.append(hashed_password.decode('utf-8'))
|
|
|
|
# Handle role assignment if provided and auto-update is_admin flag
|
|
if 'role_ids' in user_data:
|
|
role_ids = user_data['role_ids'] or []
|
|
|
|
# Check if user should be admin based on new roles
|
|
if role_ids:
|
|
admin_roles = await conn.fetch("""
|
|
SELECT id FROM roles
|
|
WHERE id = ANY($1) AND name = 'super_admin'
|
|
""", role_ids)
|
|
is_admin = len(admin_roles) > 0
|
|
else:
|
|
is_admin = False
|
|
|
|
# Add is_admin to update fields based on roles
|
|
param_count += 1
|
|
update_fields.append(f"is_admin = ${param_count}")
|
|
params.append(is_admin)
|
|
|
|
# Clear existing user roles
|
|
await conn.execute("DELETE FROM user_roles WHERE user_id = $1", user_id)
|
|
|
|
# Assign new roles
|
|
if role_ids:
|
|
for role_id in role_ids:
|
|
await conn.execute("""
|
|
INSERT INTO user_roles (user_id, role_id)
|
|
VALUES ($1, $2)
|
|
ON CONFLICT (user_id, role_id) DO NOTHING
|
|
""", user_id, role_id)
|
|
|
|
if not update_fields:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="No fields to update")
|
|
|
|
# Add updated_at
|
|
param_count += 1
|
|
update_fields.append(f"updated_at = CURRENT_TIMESTAMP")
|
|
|
|
# Execute update
|
|
query = f"UPDATE users SET {', '.join(update_fields)} WHERE id = $1"
|
|
await conn.execute(query, *params)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='update',
|
|
resource_type='user',
|
|
resource_id=str(user_id),
|
|
details={'updated_fields': list(user_data.keys())}
|
|
)
|
|
|
|
return {"message": f"User updated successfully"}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to update user: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to update user: {str(e)}")
|
|
|
|
@router.put("/users/{user_id}/password")
|
|
async def change_user_password(
|
|
user_id: int,
|
|
password_data: UserPasswordUpdate,
|
|
authorization: str = Header(None)
|
|
):
|
|
"""
|
|
Change user password
|
|
|
|
Users can change their own password by providing current password.
|
|
Admins can change any user's password without current password.
|
|
"""
|
|
try:
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if target user exists
|
|
target_user = await conn.fetchrow(
|
|
"SELECT id, username, password_hash, is_active FROM users WHERE id = $1",
|
|
user_id
|
|
)
|
|
if not target_user:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="User not found")
|
|
|
|
if not target_user['is_active']:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="User is inactive")
|
|
|
|
# Check permissions
|
|
is_self = current_user["id"] == user_id
|
|
is_admin = current_user.get("is_admin", False)
|
|
|
|
if not is_self and not is_admin:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="You can only change your own password"
|
|
)
|
|
|
|
# If changing own password, verify current password
|
|
if is_self:
|
|
import bcrypt
|
|
if not bcrypt.checkpw(
|
|
password_data.current_password.encode('utf-8'),
|
|
target_user['password_hash'].encode('utf-8')
|
|
):
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="Current password is incorrect"
|
|
)
|
|
|
|
# Hash new password
|
|
import bcrypt
|
|
salt = bcrypt.gensalt()
|
|
hashed_password = bcrypt.hashpw(
|
|
password_data.new_password.encode('utf-8'),
|
|
salt
|
|
)
|
|
|
|
# Update password
|
|
await conn.execute("""
|
|
UPDATE users
|
|
SET password_hash = $1, updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $2
|
|
""", hashed_password.decode('utf-8'), user_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='change_password',
|
|
resource_type='user',
|
|
resource_id=str(user_id),
|
|
details={
|
|
'target_username': target_user['username'],
|
|
'changed_by_self': is_self
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": "Password changed successfully",
|
|
"username": target_user['username']
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to change password: {e}")
|
|
raise HTTPException(
|
|
status_code=500,
|
|
detail=f"Failed to change password: {str(e)}"
|
|
)
|
|
|
|
@router.delete("/servers/{server_id}")
|
|
async def delete_server_global(server_id: int, request: Request, authorization: str = Header(None)):
|
|
"""Delete a server by ID - global endpoint for UI compatibility"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token, check_user_permission
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Risk-audit follow-up to Bulgu-#77: the legacy
|
|
# `delete_server` on `routers/backend.py` was upgraded to
|
|
# gate on `backends.update`, but BackendServers.js calls
|
|
# the compatibility alias `DELETE /api/servers/{id}` that
|
|
# routes here — so the FE delete path was still missing
|
|
# the per-action RBAC check. A read-only operator with
|
|
# `backends.read` and pool access could delete servers.
|
|
# Mirror the gate so both endpoints enforce the same
|
|
# contract.
|
|
has_permission = await check_user_permission(
|
|
current_user["id"], "backends", "update",
|
|
current_user=current_user,
|
|
)
|
|
if not has_permission:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="You don't have permission to delete servers",
|
|
)
|
|
|
|
# Get request body for cluster_id validation
|
|
request_body = await request.json() if hasattr(request, 'json') else {}
|
|
expected_cluster_id = request_body.get('cluster_id')
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Get server info before deletion
|
|
server = await conn.fetchrow("""
|
|
SELECT id, server_name, backend_name, server_address, server_port, cluster_id, is_active
|
|
FROM backend_servers WHERE id = $1
|
|
""", server_id)
|
|
|
|
if not server:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="Server not found")
|
|
|
|
cluster_id = server['cluster_id']
|
|
backend_name = server['backend_name']
|
|
server_name = server['server_name']
|
|
|
|
# Validate cluster ownership for multi-cluster security
|
|
if expected_cluster_id and cluster_id != expected_cluster_id:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail=f"Server belongs to cluster {cluster_id}, not cluster {expected_cluster_id}"
|
|
)
|
|
|
|
# Check if cluster exists
|
|
cluster_exists = await conn.fetchval("""
|
|
SELECT id FROM haproxy_clusters WHERE id = $1
|
|
""", cluster_id)
|
|
|
|
if not cluster_exists:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=404,
|
|
detail="Cluster not found"
|
|
)
|
|
|
|
# Check if user_pool_access table exists (for backward compatibility)
|
|
table_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.tables
|
|
WHERE table_name = 'user_pool_access'
|
|
)
|
|
""")
|
|
|
|
if table_exists:
|
|
# Check if expires_at column exists (for backward compatibility)
|
|
expires_at_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.columns
|
|
WHERE table_name = 'user_pool_access' AND column_name = 'expires_at'
|
|
)
|
|
""")
|
|
|
|
# Proper user-pool access validation
|
|
if expires_at_exists:
|
|
user_access = await conn.fetchrow("""
|
|
SELECT upa.access_level, hc.id, hc.name
|
|
FROM haproxy_clusters hc
|
|
JOIN haproxy_cluster_pools hcp ON hc.pool_id = hcp.id
|
|
JOIN user_pool_access upa ON hcp.id = upa.pool_id
|
|
WHERE upa.user_id = $1 AND hc.id = $2 AND upa.is_active = TRUE
|
|
AND (upa.expires_at IS NULL OR upa.expires_at > CURRENT_TIMESTAMP)
|
|
""", current_user['id'], cluster_id)
|
|
else:
|
|
# Fallback query without expires_at column
|
|
user_access = await conn.fetchrow("""
|
|
SELECT upa.access_level, hc.id, hc.name
|
|
FROM haproxy_clusters hc
|
|
JOIN haproxy_cluster_pools hcp ON hc.pool_id = hcp.id
|
|
JOIN user_pool_access upa ON hcp.id = upa.pool_id
|
|
WHERE upa.user_id = $1 AND hc.id = $2 AND upa.is_active = TRUE
|
|
""", current_user['id'], cluster_id)
|
|
|
|
if not user_access:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="You don't have access to this cluster"
|
|
)
|
|
|
|
logger.info(f"Global server delete: server_id={server_id}, name={server_name}, backend={backend_name}, cluster_id={cluster_id}")
|
|
|
|
# 1. Soft delete the server (mark as inactive and pending for deletion)
|
|
await conn.execute("""
|
|
UPDATE backend_servers
|
|
SET is_active = FALSE, last_config_status = 'PENDING', updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $1
|
|
""", server_id)
|
|
|
|
# 2. Don't mark backend as PENDING - let Agent Sync work based on server version
|
|
# This matches the behavior of server add operations
|
|
# await conn.execute("""
|
|
# UPDATE backends
|
|
# SET last_config_status = 'PENDING', updated_at = CURRENT_TIMESTAMP
|
|
# WHERE name = $1 AND cluster_id = $2
|
|
# """, backend_name, cluster_id)
|
|
|
|
# Create config version for server deletion (like frontend/SSL delete)
|
|
sync_results = []
|
|
if cluster_id:
|
|
try:
|
|
# Generate new HAProxy config without this server
|
|
from services.haproxy_config import generate_haproxy_config_for_cluster
|
|
import hashlib
|
|
import time
|
|
|
|
config_content = await generate_haproxy_config_for_cluster(cluster_id)
|
|
|
|
logger.info(f"SERVER DELETE DEBUG (user.py): Generated config content length: {len(config_content) if config_content else 0}")
|
|
if config_content:
|
|
logger.info(f"SERVER DELETE DEBUG (user.py): Config contains server line: {server_name in config_content}")
|
|
logger.info(f"SERVER DELETE DEBUG (user.py): Config contains 'DISABLED': {'DISABLED' in config_content}")
|
|
logger.info(f"SERVER DELETE DEBUG (user.py): Server marked for deletion - should be skipped in config")
|
|
else:
|
|
logger.error(f"SERVER DELETE DEBUG (user.py): Config content is empty or None!")
|
|
|
|
config_hash = hashlib.sha256(config_content.encode()).hexdigest()
|
|
version_name = f"server-{server_id}-delete-{int(time.time())}"
|
|
|
|
# Get system admin user ID for created_by
|
|
admin_user_id = await conn.fetchval("SELECT id FROM users WHERE username = 'admin' LIMIT 1") or 1
|
|
|
|
# Create new config version
|
|
config_version_id = await conn.fetchval("""
|
|
INSERT INTO config_versions
|
|
(cluster_id, version_name, config_content, checksum, created_by, is_active, status)
|
|
VALUES ($1, $2, $3, $4, $5, FALSE, 'PENDING')
|
|
RETURNING id
|
|
""", cluster_id, version_name, config_content, config_hash, admin_user_id)
|
|
|
|
logger.info(f"Created PENDING config version {version_name} for server delete")
|
|
sync_results = [{'node': 'pending', 'success': True, 'version': version_name, 'status': 'PENDING', 'message': 'Server deletion created. Click Apply to activate.'}]
|
|
|
|
except Exception as e:
|
|
logger.error(f"Cluster config update failed for server {server_name}: {e}")
|
|
sync_results = [{'node': 'cluster', 'success': False, 'error': str(e)}]
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log user activity
|
|
await log_user_activity(
|
|
user_id=current_user['id'],
|
|
action="delete_server",
|
|
resource_type="server",
|
|
resource_id=str(server_id),
|
|
details={
|
|
"server_name": server_name,
|
|
"backend_name": backend_name,
|
|
"cluster_id": cluster_id
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": f"Server '{server_name}' deleted successfully from backend '{backend_name}'",
|
|
"server": {
|
|
"id": server_id,
|
|
"name": server_name,
|
|
"backend": backend_name
|
|
},
|
|
"sync_results": sync_results
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to delete server {server_id}: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to delete server: {str(e)}")
|
|
|
|
@router.delete("/users/{user_id}")
|
|
async def delete_user(user_id: int, authorization: str = Header(None)):
|
|
"""Delete a user"""
|
|
try:
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can delete users
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can delete users"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user exists and get details
|
|
user_info = await conn.fetchrow("SELECT username, email FROM users WHERE id = $1", user_id)
|
|
if not user_info:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="User not found")
|
|
|
|
# Prevent self-deletion
|
|
if current_user["id"] == user_id:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Cannot delete your own account")
|
|
|
|
# Prevent deletion of admin user
|
|
if user_info["username"] == "admin":
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Cannot delete admin user")
|
|
|
|
# Soft delete - mark as inactive instead of hard delete
|
|
await conn.execute("""
|
|
UPDATE users
|
|
SET is_active = FALSE, updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $1
|
|
""", user_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='delete',
|
|
resource_type='user',
|
|
resource_id=str(user_id),
|
|
details={'username': user_info['username'], 'email': user_info['email']}
|
|
)
|
|
|
|
return {"message": f"User '{user_info['username']}' deleted successfully"}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to delete user: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to delete user: {str(e)}")
|
|
|
|
@router.get("/users/{user_id}/pool-access")
|
|
async def get_user_pool_access(user_id: int, authorization: str = Header(None)):
|
|
"""Get user's pool access permissions"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Only admin users or the user themselves can view pool access
|
|
if not current_user.get('is_admin') and current_user['id'] != user_id:
|
|
raise HTTPException(status_code=403, detail="Access denied")
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user_pool_access table exists
|
|
table_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.tables
|
|
WHERE table_name = 'user_pool_access'
|
|
)
|
|
""")
|
|
|
|
if not table_exists:
|
|
await close_database_connection(conn)
|
|
return {"pool_access": [], "message": "user_pool_access table not found"}
|
|
|
|
# Get user's pool access
|
|
pool_access = await conn.fetch("""
|
|
SELECT upa.id, upa.pool_id, upa.access_level, upa.granted_at, upa.expires_at, upa.is_active,
|
|
hcp.name as pool_name, hcp.description as pool_description,
|
|
u.username as granted_by_username
|
|
FROM user_pool_access upa
|
|
JOIN haproxy_cluster_pools hcp ON upa.pool_id = hcp.id
|
|
LEFT JOIN users u ON upa.granted_by = u.id
|
|
WHERE upa.user_id = $1
|
|
ORDER BY upa.granted_at DESC
|
|
""", user_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"pool_access": [dict(row) for row in pool_access]
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to get user pool access: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to get user pool access: {str(e)}")
|
|
|
|
@router.post("/users/{user_id}/pool-access")
|
|
async def grant_user_pool_access(
|
|
user_id: int,
|
|
access_data: dict,
|
|
authorization: str = Header(None)
|
|
):
|
|
"""Grant user access to pools"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Only admin users can grant pool access
|
|
if not current_user.get('is_admin'):
|
|
raise HTTPException(status_code=403, detail="Only admin users can grant pool access")
|
|
|
|
pool_ids = access_data.get('pool_ids', [])
|
|
access_level = access_data.get('access_level', 'read_write')
|
|
|
|
if not pool_ids:
|
|
raise HTTPException(status_code=400, detail="pool_ids is required")
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user exists
|
|
user_exists = await conn.fetchval("SELECT id FROM users WHERE id = $1", user_id)
|
|
if not user_exists:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="User not found")
|
|
|
|
# Check if user_pool_access table exists
|
|
table_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.tables
|
|
WHERE table_name = 'user_pool_access'
|
|
)
|
|
""")
|
|
|
|
if not table_exists:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=500, detail="user_pool_access table not found")
|
|
|
|
# Grant access to multiple pools
|
|
granted_pools = []
|
|
for pool_id in pool_ids:
|
|
# Check if pool exists
|
|
pool_exists = await conn.fetchval("SELECT name FROM haproxy_cluster_pools WHERE id = $1", pool_id)
|
|
if not pool_exists:
|
|
continue
|
|
|
|
# Grant access
|
|
await conn.execute("""
|
|
INSERT INTO user_pool_access (user_id, pool_id, access_level, granted_by)
|
|
VALUES ($1, $2, $3, $4)
|
|
ON CONFLICT (user_id, pool_id)
|
|
DO UPDATE SET
|
|
access_level = EXCLUDED.access_level,
|
|
granted_by = EXCLUDED.granted_by,
|
|
granted_at = CURRENT_TIMESTAMP,
|
|
is_active = TRUE
|
|
""", user_id, pool_id, access_level, current_user['id'])
|
|
|
|
granted_pools.append({"pool_id": pool_id, "pool_name": pool_exists})
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"message": f"Pool access granted successfully",
|
|
"user_id": user_id,
|
|
"granted_pools": granted_pools,
|
|
"access_level": access_level
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to grant user pool access: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to grant user pool access: {str(e)}")
|
|
|
|
@router.delete("/users/{user_id}/pool-access/{pool_id}")
|
|
async def revoke_user_pool_access(
|
|
user_id: int,
|
|
pool_id: int,
|
|
authorization: str = Header(None)
|
|
):
|
|
"""Revoke user access to a pool"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Only admin users can revoke pool access
|
|
if not current_user.get('is_admin'):
|
|
raise HTTPException(status_code=403, detail="Only admin users can revoke pool access")
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user_pool_access table exists
|
|
table_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.tables
|
|
WHERE table_name = 'user_pool_access'
|
|
)
|
|
""")
|
|
|
|
if not table_exists:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=500, detail="user_pool_access table not found")
|
|
|
|
# Revoke access (soft delete)
|
|
result = await conn.execute("""
|
|
UPDATE user_pool_access
|
|
SET is_active = FALSE
|
|
WHERE user_id = $1 AND pool_id = $2
|
|
""", user_id, pool_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"message": f"Pool access revoked successfully",
|
|
"user_id": user_id,
|
|
"pool_id": pool_id
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to revoke user pool access: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to revoke user pool access: {str(e)}")
|
|
|
|
|
|
@router.get("/user-activity")
|
|
async def get_user_activity(
|
|
limit: int = 50,
|
|
offset: int = 0,
|
|
user_id: Optional[int] = None,
|
|
authorization: str = Header(None)
|
|
):
|
|
"""Get user activity logs"""
|
|
try:
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# R18c audit fix (round 4 #21 — KRITIK info leak): pre-fix
|
|
# any authenticated user could:
|
|
# 1. Omit `user_id` and fetch the FULL activity log of
|
|
# every operator on the platform — including admin
|
|
# apply_changes, ACME orders, and (after R18b round 6)
|
|
# wizard `apply_error` / `acme_staging_error` blobs.
|
|
# 2. Pass an arbitrary `user_id` and read another
|
|
# operator's activity stream.
|
|
# The wizard's richer audit row makes this leak more
|
|
# consequential than before because the JSONB now carries
|
|
# operationally sensitive failure details. Restrict the
|
|
# endpoint to admins (full access) or to a user querying
|
|
# their own rows. Non-admin requests for someone else's
|
|
# activity → 403.
|
|
is_admin = bool(current_user.get("is_admin"))
|
|
own_id = current_user.get("id")
|
|
if not is_admin:
|
|
if user_id is None:
|
|
# Default to the caller's own rows for non-admins;
|
|
# the previous unfiltered listing is admin-only.
|
|
user_id = own_id
|
|
elif user_id != own_id:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only administrators can view another user's activity log."
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Build query with optional user filter
|
|
where_clause = ""
|
|
params = [limit, offset]
|
|
if user_id:
|
|
where_clause = "WHERE ua.user_id = $3"
|
|
params.append(user_id)
|
|
|
|
# Get activity logs
|
|
activities = await conn.fetch(f"""
|
|
SELECT ua.id, ua.user_id, ua.action, ua.resource_type, ua.resource_id,
|
|
ua.details, ua.created_at, ua.ip_address,
|
|
u.username, u.email
|
|
FROM user_activity_logs ua
|
|
LEFT JOIN users u ON ua.user_id = u.id
|
|
{where_clause}
|
|
ORDER BY ua.created_at DESC
|
|
LIMIT $1 OFFSET $2
|
|
""", *params)
|
|
|
|
# Get total count
|
|
count_query = f"SELECT COUNT(*) FROM user_activity_logs ua {where_clause}"
|
|
if user_id:
|
|
total = await conn.fetchval(count_query, user_id)
|
|
else:
|
|
total = await conn.fetchval(count_query)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"activities": [
|
|
{
|
|
"id": activity["id"],
|
|
"user_id": activity["user_id"],
|
|
"username": activity["username"],
|
|
"email": activity["email"],
|
|
"action": activity["action"],
|
|
"resource_type": activity["resource_type"],
|
|
"resource_id": activity["resource_id"],
|
|
"details": activity["details"],
|
|
"created_at": activity["created_at"].isoformat() if activity["created_at"] else None,
|
|
"ip_address": activity["ip_address"]
|
|
} for activity in activities
|
|
],
|
|
"total": total,
|
|
"limit": limit,
|
|
"offset": offset
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to get user activity: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to get user activity: {str(e)}")
|
|
|
|
# ==== ENHANCED ROLE MANAGEMENT ====
|
|
|
|
@router.post("/roles")
|
|
async def create_role(role_data: dict, authorization: str = Header(None)):
|
|
"""Create a new role with cluster-specific permissions"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can create roles
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can create roles"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if role name already exists
|
|
existing = await conn.fetchrow("SELECT id FROM roles WHERE name = $1", role_data['name'])
|
|
if existing:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Role name already exists")
|
|
|
|
# Create role with cluster_ids support
|
|
role_id = await conn.fetchval("""
|
|
INSERT INTO roles (name, display_name, description, permissions, cluster_ids, is_active)
|
|
VALUES ($1, $2, $3, $4, $5, $6)
|
|
RETURNING id
|
|
""",
|
|
role_data['name'],
|
|
role_data['display_name'],
|
|
role_data.get('description'),
|
|
json.dumps(role_data.get('permissions', [])) if isinstance(role_data.get('permissions'), list) else role_data.get('permissions', []),
|
|
json.dumps(role_data.get('cluster_ids')) if role_data.get('cluster_ids') else None,
|
|
role_data.get('is_active', True)
|
|
)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='create',
|
|
resource_type='role',
|
|
resource_id=str(role_id),
|
|
details={
|
|
'role_name': role_data['name'],
|
|
'display_name': role_data['display_name'],
|
|
'permissions_count': len(role_data.get('permissions', [])),
|
|
'cluster_specific': bool(role_data.get('cluster_ids'))
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": f"Role '{role_data['display_name']}' created successfully",
|
|
"id": role_id
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to create role: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to create role: {str(e)}")
|
|
|
|
@router.put("/roles/{role_id}")
|
|
async def update_role(role_id: int, role_data: dict, authorization: str = Header(None)):
|
|
"""Update an existing role"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can update roles
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can update roles"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if role exists
|
|
existing = await conn.fetchrow("SELECT * FROM roles WHERE id = $1", role_id)
|
|
if not existing:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="Role not found")
|
|
|
|
# Check if new name conflicts with other roles
|
|
if 'name' in role_data and role_data['name'] != existing['name']:
|
|
name_conflict = await conn.fetchrow("SELECT id FROM roles WHERE name = $1 AND id != $2", role_data['name'], role_id)
|
|
if name_conflict:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Role name already exists")
|
|
|
|
# Update role
|
|
await conn.execute("""
|
|
UPDATE roles SET
|
|
name = $1, display_name = $2, description = $3,
|
|
permissions = $4, cluster_ids = $5, is_active = $6,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $7
|
|
""",
|
|
role_data.get('name', existing['name']),
|
|
role_data.get('display_name', existing['display_name']),
|
|
role_data.get('description', existing['description']),
|
|
json.dumps(role_data.get('permissions', existing['permissions'])) if isinstance(role_data.get('permissions'), list) else role_data.get('permissions', existing['permissions']),
|
|
json.dumps(role_data.get('cluster_ids')) if 'cluster_ids' in role_data and role_data.get('cluster_ids') else (existing['cluster_ids'] if existing['cluster_ids'] else None),
|
|
role_data.get('is_active', existing['is_active']),
|
|
role_id
|
|
)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='update',
|
|
resource_type='role',
|
|
resource_id=str(role_id),
|
|
details={
|
|
'role_name': role_data.get('name', existing['name']),
|
|
'display_name': role_data.get('display_name', existing['display_name']),
|
|
'permissions_count': len(role_data.get('permissions', [])),
|
|
'cluster_specific': bool(role_data.get('cluster_ids'))
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": f"Role updated successfully"
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to update role: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to update role: {str(e)}")
|
|
|
|
@router.delete("/roles/{role_id}")
|
|
async def delete_role(role_id: int, authorization: str = Header(None)):
|
|
"""Delete a role (soft delete by setting is_active = false)"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can delete roles
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can delete roles"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if role exists and is not system role
|
|
role = await conn.fetchrow("SELECT * FROM roles WHERE id = $1", role_id)
|
|
if not role:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="Role not found")
|
|
|
|
if role.get('is_system'):
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=400, detail="Cannot delete system roles")
|
|
|
|
# Soft delete role
|
|
await conn.execute("""
|
|
UPDATE roles SET is_active = FALSE, updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $1
|
|
""", role_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='delete',
|
|
resource_type='role',
|
|
resource_id=str(role_id),
|
|
details={
|
|
'role_name': role['name'],
|
|
'display_name': role['display_name']
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": f"Role '{role['display_name']}' deleted successfully"
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to delete role: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to delete role: {str(e)}")
|
|
|
|
# ==== USER ROLE ASSIGNMENT ENDPOINTS ====
|
|
|
|
@router.post("/users/{user_id}/roles")
|
|
async def assign_user_roles(user_id: int, role_data: dict, authorization: str = Header(None)):
|
|
"""Assign roles to a user"""
|
|
try:
|
|
# Verify authentication
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# SECURITY: Only admin users can assign roles to users
|
|
if not current_user.get("is_admin", False):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Only admin users can assign roles to users"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if user exists
|
|
user = await conn.fetchrow("SELECT id, username FROM users WHERE id = $1", user_id)
|
|
if not user:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="User not found")
|
|
|
|
role_ids = role_data.get('role_ids', [])
|
|
|
|
# Clear existing user roles
|
|
await conn.execute("DELETE FROM user_roles WHERE user_id = $1", user_id)
|
|
|
|
# Assign new roles
|
|
if role_ids:
|
|
for role_id in role_ids:
|
|
await conn.execute("""
|
|
INSERT INTO user_roles (user_id, role_id)
|
|
VALUES ($1, $2)
|
|
ON CONFLICT (user_id, role_id) DO NOTHING
|
|
""", user_id, role_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log activity
|
|
await log_user_activity(
|
|
user_id=current_user["id"],
|
|
action='update',
|
|
resource_type='user_roles',
|
|
resource_id=str(user_id),
|
|
details={
|
|
'target_user': user['username'],
|
|
'assigned_roles': len(role_ids),
|
|
'role_ids': role_ids
|
|
}
|
|
)
|
|
|
|
return {
|
|
"message": f"Roles assigned to user '{user['username']}' successfully",
|
|
"assigned_roles": len(role_ids)
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to assign roles to user {user_id}: {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to assign roles: {str(e)}")
|
|
|
|
@router.post("/user-roles")
|
|
async def assign_roles_legacy(role_assignment: dict, authorization: str = Header(None)):
|
|
"""Legacy endpoint for role assignment (for compatibility)"""
|
|
try:
|
|
user_id = role_assignment.get('user_id')
|
|
role_ids = role_assignment.get('role_ids', [])
|
|
|
|
if not user_id:
|
|
raise HTTPException(status_code=400, detail="User ID is required")
|
|
|
|
return await assign_user_roles(user_id, {'role_ids': role_ids}, authorization)
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Failed to assign roles (legacy): {e}")
|
|
raise HTTPException(status_code=500, detail=f"Failed to assign roles: {str(e)}") |