Files
service-finder/backend/app/api/v1/endpoints/admin.py

737 lines
27 KiB
Python
Executable File

# /opt/docker/dev/service_finder/backend/app/api/v1/endpoints/admin.py
from fastapi import APIRouter, Depends, HTTPException, status, Body, BackgroundTasks, Query
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func, text, delete, or_
from sqlalchemy.orm import selectinload
from typing import List, Any, Dict, Optional
from datetime import datetime, timedelta, timezone
from app.api import deps
from app.models.identity import User, UserRole, Person # JAVÍTVA: Központi import
from app.models.identity.address import Address
from app.models.system import SystemParameter, ParameterScope
from app.services.system_service import system_service
# JAVÍTVA: Security audit modellek
from app.models import SecurityAuditLog, OperationalLog
# JAVÍTVA: Ezek a modellek a security.py-ból jönnek (ha ott vannak)
from app.models import PendingAction, ActionStatus
from app.services.security_service import security_service
from app.services.translation_service import TranslationService
from pydantic import BaseModel, Field
from typing import Optional as Opt
class ConfigUpdate(BaseModel):
key: str
value: Any
scope_level: ParameterScope = ParameterScope.GLOBAL
scope_id: Optional[str] = None
category: str = "general"
router = APIRouter()
@router.get("/health-monitor", tags=["Sentinel Monitoring"])
async def get_system_health(
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
stats = {}
# Adatbázis statisztikák (Nyers SQL marad, mert hatékony)
user_stats = await db.execute(text("SELECT subscription_plan, count(*) FROM identity.users GROUP BY subscription_plan"))
stats["user_distribution"] = {row[0]: row[1] for row in user_stats}
asset_count = await db.execute(text("SELECT count(*) FROM vehicle.assets"))
stats["total_assets"] = asset_count.scalar()
org_count = await db.execute(text("SELECT count(*) FROM fleet.organizations"))
stats["total_organizations"] = org_count.scalar()
# JAVÍTVA: Biztonsági státusz az új SecurityAuditLog alapján
day_ago = datetime.now() - timedelta(days=1)
crit_logs = await db.execute(
select(func.count(SecurityAuditLog.id))
.where(
SecurityAuditLog.is_critical == True,
SecurityAuditLog.created_at >= day_ago
)
)
stats["critical_alerts_24h"] = crit_logs.scalar() or 0
return stats
@router.get("/pending-actions", response_model=List[Any], tags=["Sentinel Security"])
async def list_pending_actions(
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
stmt = select(PendingAction).where(PendingAction.status == ActionStatus.pending)
result = await db.execute(stmt)
return result.scalars().all()
@router.post("/approve/{action_id}", tags=["Sentinel Security"])
async def approve_action(
action_id: int,
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
try:
await security_service.approve_action(db, admin.id, action_id)
return {"status": "success", "message": "Művelet végrehajtva."}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@router.get("/parameters", tags=["Dynamic Configuration"])
async def list_all_parameters(
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
result = await db.execute(select(SystemParameter))
return result.scalars().all()
@router.post("/parameters", tags=["Dynamic Configuration"])
async def set_parameter(
config: ConfigUpdate,
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
query = text("""
INSERT INTO system.system_parameters (key, value, scope_level, scope_id, category, last_modified_by)
VALUES (:key, :val, :sl, :sid, :cat, :user)
ON CONFLICT (key, scope_level, scope_id)
DO UPDATE SET
value = EXCLUDED.value,
category = EXCLUDED.category,
last_modified_by = EXCLUDED.last_modified_by,
updated_at = now()
""")
await db.execute(query, {
"key": config.key,
"val": config.value,
"sl": config.scope_level,
"sid": config.scope_id,
"cat": config.category,
"user": admin.email
})
await db.commit()
return {"status": "success", "message": f"'{config.key}' frissítve."}
@router.get("/parameters/scoped", tags=["Dynamic Configuration"])
async def get_scoped_parameter(
key: str,
user_id: Optional[str] = None,
region_id: Optional[str] = None,
country_code: Optional[str] = None,
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
"""
Hierarchikus paraméterlekérdezés a következő prioritással:
User > Region > Country > Global.
"""
value = await system_service.get_scoped_parameter(
db, key, user_id, region_id, country_code, default=None
)
if value is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"Paraméter '{key}' nem található a megadott scope-okban."
)
return {"key": key, "value": value}
@router.post("/translations/sync", tags=["System Utilities"])
async def sync_translations_to_json(
db: AsyncSession = Depends(deps.get_db),
admin: User = Depends(deps.get_current_admin)
):
await TranslationService.export_to_json(db)
return {"message": "JSON fájlok frissítve."}
@router.get("/ping", tags=["Admin Test"])
async def admin_ping(
current_user: User = Depends(deps.get_current_admin)
):
"""
Egyszerű ping végpont admin jogosultság ellenőrzéséhez.
"""
return {
"message": "Admin felület aktív",
"role": current_user.role.value if hasattr(current_user.role, "value") else current_user.role
}
@router.post("/users/{user_id}/ban", tags=["Admin Security"])
async def ban_user(
user_id: int,
reason: str = Body(..., embed=True),
current_admin: User = Depends(deps.get_current_admin),
db: AsyncSession = Depends(deps.get_db)
):
"""
Felhasználó tiltása (Ban Hammer).
- Megkeresi a usert (identity.users táblában).
- Ha nincs -> 404
- Ha a user.role == superadmin -> 403 (Saját magát/másik admint ne tiltson le).
- Állítja be a tiltást (is_active = False).
- Audit logba rögzíti a reason-t.
"""
from sqlalchemy import select
# 1. Keresd meg a usert
stmt = select(User).where(User.id == user_id)
result = await db.execute(stmt)
user = result.scalar_one_or_none()
if not user:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"User not found with ID: {user_id}"
)
# 2. Ellenőrizd, hogy nem superadmin-e
if user.role == UserRole.superadmin:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Cannot ban a superadmin user"
)
# 3. Tiltás beállítása
user.is_active = False
# Opcionálisan: banned_until mező kitöltése, ha létezik a modellben
# user.banned_until = datetime.now() + timedelta(days=30)
# 4. Audit log létrehozása
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action="ban_user",
target_user_id=user_id,
details=f"User banned. Reason: {reason}",
is_critical=True,
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"message": f"User {user_id} banned successfully.",
"reason": reason
}
@router.post("/marketplace/services/{staging_id}/approve", tags=["Marketplace Moderation"])
async def approve_staged_service(
staging_id: int,
current_admin: User = Depends(deps.get_current_admin),
db: AsyncSession = Depends(deps.get_db)
):
"""
Szerviz jóváhagyása a Piactéren (Kék Pipa).
- Megkeresi a marketplace.service_staging rekordot.
- Ha nincs -> 404
- Állítja a validation_level-t 100-ra, a status-t 'approved'-ra.
"""
from sqlalchemy import select
from app.models.staged_data import ServiceStaging
stmt = select(ServiceStaging).where(ServiceStaging.id == staging_id)
result = await db.execute(stmt)
staging = result.scalar_one_or_none()
if not staging:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"Service staging record not found with ID: {staging_id}"
)
# Jóváhagyás
staging.validation_level = 100
staging.status = "approved"
# Audit log
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action="approve_service",
target_staging_id=staging_id,
details=f"Service staging approved: {staging.service_name}",
is_critical=False,
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"message": f"Service staging {staging_id} approved.",
"service_name": staging.service_name
}
# ==================== EPIC 10: ADMIN FRONTEND API ENDPOINTS ====================
from app.workers.service.validation_pipeline import ValidationPipeline
from app.models.marketplace.service import ServiceProfile
from app.models.gamification.gamification import UserStats
class LocationUpdate(BaseModel):
latitude: float = Field(..., ge=-90, le=90)
longitude: float = Field(..., ge=-180, le=180)
class PenaltyRequest(BaseModel):
penalty_level: int = Field(..., ge=-10, le=-1, description="Negatív szint (-1 a legkisebb, -10 a legnagyobb büntetés)")
reason: str = Field(..., min_length=5, max_length=500)
@router.post("/services/{service_id}/trigger-ai", tags=["AI Pipeline"])
async def trigger_ai_pipeline(
service_id: int,
background_tasks: BackgroundTasks,
current_admin: User = Depends(deps.get_current_admin),
db: AsyncSession = Depends(deps.get_db)
):
"""
AI Pipeline manuális indítása egy adott szerviz profilra.
A végpont azonnal visszatér, és a validációt háttérfeladatként futtatja.
"""
# Ellenőrizzük, hogy létezik-e a szerviz profil
stmt = select(ServiceProfile).where(ServiceProfile.id == service_id)
result = await db.execute(stmt)
profile = result.scalar_one_or_none()
if not profile:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"Service profile not found with ID: {service_id}"
)
# Háttérfeladat hozzáadása
background_tasks.add_task(run_validation_pipeline, service_id)
# Audit log
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action="trigger_ai_pipeline",
target_service_id=service_id,
details=f"AI pipeline manually triggered for service {service_id}",
is_critical=False,
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"message": f"AI pipeline started for service {service_id}",
"service_name": profile.service_name,
"note": "Validation runs in background, check logs for results."
}
async def run_validation_pipeline(profile_id: int):
"""Háttérfeladat a ValidationPipeline futtatásához."""
try:
pipeline = ValidationPipeline()
success = await pipeline.run(profile_id)
logger = logging.getLogger("Service-AI-Pipeline")
if success:
logger.info(f"Pipeline successful for profile {profile_id}")
else:
logger.warning(f"Pipeline failed for profile {profile_id}")
except Exception as e:
logger.error(f"Pipeline error for profile {profile_id}: {e}")
@router.patch("/services/{service_id}/location", tags=["Service Management"])
async def update_service_location(
service_id: int,
location: LocationUpdate,
current_admin: User = Depends(deps.get_current_admin),
db: AsyncSession = Depends(deps.get_db)
):
"""
Szerviz térképes mozgatása (Koordináta frissítés).
A Nuxt Leaflet térkép drag-and-drop funkciójához használható.
"""
stmt = select(ServiceProfile).where(ServiceProfile.id == service_id)
result = await db.execute(stmt)
profile = result.scalar_one_or_none()
if not profile:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"Service profile not found with ID: {service_id}"
)
# Frissítjük a koordinátákat
profile.latitude = location.latitude
profile.longitude = location.longitude
profile.updated_at = datetime.now()
# Audit log
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action="update_service_location",
target_service_id=service_id,
details=f"Service location updated to lat={location.latitude}, lon={location.longitude}",
is_critical=False,
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"message": f"Service location updated for {service_id}",
"latitude": location.latitude,
"longitude": location.longitude
}
@router.patch("/users/{user_id}/penalty", tags=["Gamification Admin"])
async def apply_gamification_penalty(
user_id: int,
penalty: PenaltyRequest,
current_admin: User = Depends(deps.get_current_admin),
db: AsyncSession = Depends(deps.get_db)
):
"""
Gamification büntetés kiosztása egy felhasználónak.
Negatív szintek alkalmazása a frissen létrehozott Gamification rendszerben.
"""
# Ellenőrizzük, hogy létezik-e a felhasználó
user_stmt = select(User).where(User.id == user_id)
user_result = await db.execute(user_stmt)
user = user_result.scalar_one_or_none()
if not user:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=f"User not found with ID: {user_id}"
)
# Megkeressük a felhasználó gamification profilját (vagy létrehozzuk)
gamification_stmt = select(UserStats).where(UserStats.user_id == user_id)
gamification_result = await db.execute(gamification_stmt)
gamification = gamification_result.scalar_one_or_none()
if not gamification:
# Ha nincs profil, létrehozzuk alapértelmezett értékekkel
gamification = UserStats(
user_id=user_id,
level=0,
xp=0,
reputation_score=100,
created_at=datetime.now(),
updated_at=datetime.now()
)
db.add(gamification)
await db.flush()
# Alkalmazzuk a büntetést (negatív szint módosítása)
# A level mező lehet negatív is a büntetések miatt
new_level = gamification.level + penalty.penalty_level
gamification.level = new_level
gamification.updated_at = datetime.now()
# Audit log
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action="apply_gamification_penalty",
target_user_id=user_id,
details=f"Gamification penalty applied: level change {penalty.penalty_level}, reason: {penalty.reason}",
is_critical=False,
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"message": f"Gamification penalty applied to user {user_id}",
"user_id": user_id,
"penalty_level": penalty.penalty_level,
"new_level": new_level,
"reason": penalty.reason
}
# ==================== USER MANAGEMENT (EPIC 10: Admin Frontend) ====================
class BulkActionRequest(BaseModel):
"""Csoportos művelet kérés sémája."""
user_ids: List[int] = Field(..., min_length=1, description="Felhasználók ID listája")
action: str = Field(..., pattern="^(ban|unban|soft_delete|restore|hard_delete)$", description="Végrehajtandó művelet")
@router.get("/users", tags=["User Management"])
async def list_users(
skip: int = Query(0, ge=0, description="Hány rekordot hagyjunk ki (lapozás)"),
limit: int = Query(50, ge=1, le=200, description="Maximum visszaadott rekordok száma"),
search_term: Optional[str] = Query(None, description="Keresőszó a célzott kereséshez"),
search_category: str = Query('all', description="Keresési kategória: email, name, phone, id, address, all"),
role: Optional[str] = Query(None, description="Szerepkör szűrés (pl. admin, user)"),
is_active: Optional[bool] = Query(None, description="Aktív státusz szűrés"),
is_deleted: Optional[bool] = Query(None, description="Törölt státusz szűrés"),
db: AsyncSession = Depends(deps.get_db),
current_admin: User = Depends(deps.get_current_admin),
):
"""
Felhasználók listázása lapozással és célzott kereséssel.
Támogatott keresési kategóriák (search_category):
- **email**: Keresés e-mail címben (ILIKE)
- **name**: Keresés Person vezeték- VAGY keresztnévben (ILIKE)
- **phone**: Keresés Person telefonszámában (ILIKE)
- **id**: Pontos egyezés User ID alapján
- **address**: Keresés címmezőkben (postal_code, city, street, house_number)
NINCS automatikus debounce keresés! A keresés csak explicit paraméterrel indul.
"""
# Alap lekérdezés: User + Person (outerjoin) + Address (outerjoin)
# KRITIKUS: selectinload helyett explicit outerjoin + contains_eager kell,
# mert a selectinload KÜLÖN query-ben tölti be a kapcsolódó adatokat,
# így a WHERE feltételek a Person/Address oszlopokra NEM működnének!
from sqlalchemy.orm import contains_eager, joinedload
from app.models.identity.address import GeoPostalCode
query = (
select(User)
.outerjoin(Person, User.person_id == Person.id)
.outerjoin(Address, Person.address_id == Address.id)
.outerjoin(GeoPostalCode, Address.postal_code_id == GeoPostalCode.id)
.options(
contains_eager(User.person),
contains_eager(User.person, Person.address),
)
)
# Célzott keresés - csak akkor alkalmazzuk, ha van search_term
if search_term:
if search_category == 'all':
query = query.where(
or_(
User.email.ilike(f"%{search_term}%"),
Person.first_name.ilike(f"%{search_term}%"),
Person.last_name.ilike(f"%{search_term}%"),
Person.phone.ilike(f"%{search_term}%")
)
)
elif search_category == 'email':
query = query.where(User.email.ilike(f"%{search_term}%"))
elif search_category == 'name':
query = query.where(
or_(
Person.first_name.ilike(f"%{search_term}%"),
Person.last_name.ilike(f"%{search_term}%")
)
)
elif search_category == 'phone':
query = query.where(Person.phone.ilike(f"%{search_term}%"))
elif search_category == 'id':
try:
user_id = int(search_term)
query = query.where(User.id == user_id)
except ValueError:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Az 'id' keresési kategóriához numerikus érték szükséges."
)
elif search_category == 'address':
# FIGYELEM: Address.zip és Address.city PYTHON property-k (a GeoPostalCode kapcsolatból)
# Az SQL lekérdezésben a GeoPostalCode tábla oszlopait kell használni!
query = query.where(
or_(
GeoPostalCode.zip_code.ilike(f"%{search_term}%"),
GeoPostalCode.city.ilike(f"%{search_term}%"),
Address.street_name.ilike(f"%{search_term}%"),
Address.house_number.ilike(f"%{search_term}%")
)
)
else:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"Érvénytelen keresési kategória: {search_category}. Érvényes: email, name, phone, id, address, all"
)
if role:
try:
role_enum = UserRole(role)
query = query.where(User.role == role_enum)
except ValueError:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"Érvénytelen szerepkör: {role}. Érvényes értékek: {[r.value for r in UserRole]}"
)
if is_active is not None:
query = query.where(User.is_active == is_active)
if is_deleted is not None:
query = query.where(User.is_deleted == is_deleted)
count_query = select(func.count()).select_from(query.subquery())
total_result = await db.execute(count_query)
total = total_result.scalar() or 0
query = query.order_by(User.id.desc()).offset(skip).limit(limit)
result = await db.execute(query)
# Unique users (outerjoin may duplicate)
seen_ids = set()
unique_users = []
for row in result.unique().scalars().all():
if row.id not in seen_ids:
seen_ids.add(row.id)
unique_users.append(row)
user_list = []
for u in unique_users:
person = u.person
address = person.address if person else None
# Cím összefűzése
address_str = None
if address:
parts = []
if address.zip:
parts.append(str(address.zip))
if address.city:
parts.append(str(address.city))
if address.street_name:
parts.append(str(address.street_name))
if address.house_number:
parts.append(str(address.house_number))
address_str = ', '.join(parts) if parts else None
user_list.append({
"id": u.id,
"email": u.email,
"role": u.role.value if hasattr(u.role, "value") else str(u.role),
"is_active": u.is_active,
"is_deleted": u.is_deleted,
"subscription_plan": u.subscription_plan,
"region_code": u.region_code,
"preferred_language": u.preferred_language,
"created_at": u.created_at.isoformat() if u.created_at else None,
"deleted_at": u.deleted_at.isoformat() if u.deleted_at else None,
# Person adatok
"first_name": person.first_name if person else None,
"last_name": person.last_name if person else None,
"phone": person.phone if person else None,
# Címadatok
"postal_code": str(address.zip) if address and address.zip else None,
"city": address.city if address and address.city else None,
"street": address.street_name if address and address.street_name else None,
"house_number": address.house_number if address and address.house_number else None,
"address": address_str,
})
return {
"total": total,
"skip": skip,
"limit": limit,
"users": user_list,
}
@router.post("/users/bulk-action", tags=["User Management"])
async def bulk_user_action(
request: BulkActionRequest,
db: AsyncSession = Depends(deps.get_db),
current_admin: User = Depends(deps.get_current_admin),
):
"""
Csoportos műveletek felhasználókon.
Támogatott akciók:
- **ban**: Felhasználó tiltása (is_active=False). Rang >= 50 (Moderator) szükséges.
- **unban**: Tiltás feloldása (is_active=True). Rang >= 90 (Admin) szükséges.
- **soft_delete**: Lágy törlés (is_deleted=True, deleted_at=now). Rang >= 90 (Admin) szükséges.
- **restore**: Lágy törlés visszaállítása (is_deleted=False, deleted_at=None). Rang >= 90 (Admin) szükséges.
- **hard_delete**: Végleges törlés az adatbázisból. Rang >= 100 (Superadmin) szükséges!
"""
role_key = current_admin.role.value.upper() if hasattr(current_admin.role, "value") else str(current_admin.role).upper()
from app.core.security import DEFAULT_RANK_MAP
admin_rank = DEFAULT_RANK_MAP.get(role_key, 0)
if request.action == "hard_delete":
if admin_rank < 100:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Végleges törléshez Superadmin jogosultság szükséges (rank >= 100)."
)
elif request.action in ("soft_delete", "restore", "unban"):
if admin_rank < 90:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Ehhez a művelethez Admin jogosultság szükséges (rank >= 90)."
)
elif request.action in ("ban",):
if admin_rank < 50:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Tiltáshoz Moderator jogosultság szükséges (rank >= 50)."
)
stmt = select(User).where(User.id.in_(request.user_ids))
result = await db.execute(stmt)
users = result.scalars().all()
if not users:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Egy felhasználó sem található a megadott ID-kkel."
)
found_ids = {u.id for u in users}
missing_ids = [uid for uid in request.user_ids if uid not in found_ids]
now = datetime.now(timezone.utc)
affected_count = 0
for user in users:
if user.role == UserRole.superadmin and user.id != current_admin.id:
continue
if request.action == "ban":
user.is_active = False
elif request.action == "unban":
user.is_active = True
elif request.action == "soft_delete":
user.is_deleted = True
user.deleted_at = now
elif request.action == "restore":
user.is_deleted = False
user.deleted_at = None
elif request.action == "hard_delete":
await db.delete(user)
affected_count += 1
audit_log = SecurityAuditLog(
user_id=current_admin.id,
action=f"bulk_{request.action}",
details=f"Bulk action '{request.action}' on {affected_count} users. IDs: {request.user_ids}",
is_critical=(request.action in ("hard_delete", "ban")),
ip_address="admin_api"
)
db.add(audit_log)
await db.commit()
return {
"status": "success",
"action": request.action,
"affected_count": affected_count,
"total_requested": len(request.user_ids),
"missing_ids": missing_ids if missing_ids else None,
"message": f"{request.action} művelet végrehajtva {affected_count} felhasználón."
}