2026.06.04 frontend építés közben

This commit is contained in:
Roo
2026-06-04 07:26:22 +00:00
parent 7adf6cc3e3
commit 59a30ac428
3302 changed files with 24091 additions and 1771 deletions

View File

@@ -11,6 +11,7 @@ from app.db.session import get_db
from app.core.security import decode_token, DEFAULT_RANK_MAP
from app.models.identity import User, UserRole # JAVÍTVA: Új Identity modell használata
from app.core.config import settings
from app.core.translation_helper import t # Translation helper
logger = logging.getLogger(__name__)
@@ -46,8 +47,8 @@ async def get_current_token_payload(
payload = decode_token(token)
if not payload or payload.get("type") != "access":
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Érvénytelen vagy lejárt munkamenet."
status_code=status.HTTP_401_UNAUTHORIZED,
detail=t("AUTH.INVALID_OR_EXPIRED_SESSION")
)
return payload
@@ -62,7 +63,7 @@ async def get_current_user(
if not user_id:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Token azonosítási hiba."
detail=t("AUTH.TOKEN_IDENTIFICATION_ERROR")
)
# JAVÍTVA: Modern SQLAlchemy 2.0 aszinkron lekérdezés with eager loading
@@ -74,7 +75,7 @@ async def get_current_user(
if not user or user.is_deleted:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="A felhasználó nem található."
detail=t("AUTH.USER_NOT_FOUND")
)
return user
@@ -86,8 +87,8 @@ async def get_current_active_user(
"""
if not current_user.is_active:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="A művelethez aktív profil és KYC azonosítás szükséges."
status_code=status.HTTP_403_FORBIDDEN,
detail=t("AUTH.ACTIVE_PROFILE_KYC_REQUIRED")
)
return current_user
@@ -113,8 +114,8 @@ async def check_resource_access(
return True
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Nincs jogosultsága ehhez az erőforráshoz."
status_code=status.HTTP_403_FORBIDDEN,
detail=t("AUTH.NO_PERMISSION_FOR_RESOURCE")
)
def check_min_rank(role_key: str):
@@ -160,6 +161,6 @@ async def get_current_admin(
if current_user.role not in allowed_roles:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Nincs megfelelő jogosultságod (Admin/Moderátor)!"
detail=t("AUTH.INSUFFICIENT_ADMIN_PERMISSIONS")
)
return current_user

View File

@@ -220,9 +220,7 @@ async def create_or_claim_vehicle(
db=db,
user_id=current_user.id,
org_id=org_id,
vin=payload.vin,
license_plate=payload.license_plate,
catalog_id=payload.catalog_id
asset_data=payload
)
return asset
except ValueError as e:

View File

@@ -1,5 +1,5 @@
# /opt/docker/dev/service_finder/backend/app/api/v1/endpoints/auth.py
from fastapi import APIRouter, Depends, HTTPException, status, Request
from fastapi import APIRouter, Depends, HTTPException, status, Request, Form, Response
from fastapi.security import OAuth2PasswordRequestForm
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select
@@ -7,10 +7,12 @@ from app.db.session import get_db
from app.services.auth_service import AuthService
from app.core.security import create_tokens, DEFAULT_RANK_MAP
from app.core.config import settings
from app.services.config_service import config
from app.schemas.auth import UserLiteRegister, Token, UserKYCComplete
from app.api.deps import get_current_user
from app.models.identity import User # JAVÍTVA: Új központi modell
from pydantic import BaseModel, Field
from app.core.translation_helper import t # Translation helper
from pydantic import BaseModel, Field, EmailStr
router = APIRouter()
@@ -22,16 +24,27 @@ async def register(user_in: UserLiteRegister, db: AsyncSession = Depends(get_db)
user = await AuthService.register_lite(db, user_in)
return {
"status": "success",
"message": "Regisztráció sikeres. Aktivációs e-mail elküldve.",
"message": t("AUTH.REGISTRATION_SUCCESS"),
"user_id": user.id,
"email": user.email
}
@router.post("/login", response_model=Token)
async def login(db: AsyncSession = Depends(get_db), form_data: OAuth2PasswordRequestForm = Depends()):
async def login(
response: Response,
db: AsyncSession = Depends(get_db),
form_data: OAuth2PasswordRequestForm = Depends(),
remember_me: bool = Form(False)
):
"""
Bejelentkezés remember_me opcióval.
A remember_me paramétert explicit Form mezőként kell elküldeni, mivel az
OAuth2PasswordRequestForm nem támogatja alapból.
"""
user = await AuthService.authenticate(db, form_data.username, form_data.password)
if not user:
raise HTTPException(status_code=401, detail="Hibás adatok.")
raise HTTPException(status_code=401, detail=t("AUTH.INVALID_CREDENTIALS"))
ranks = await settings.get_db_setting(db, "rbac_rank_matrix", default=DEFAULT_RANK_MAP)
role_name = user.role.value if hasattr(user.role, 'value') else str(user.role)
@@ -45,7 +58,26 @@ async def login(db: AsyncSession = Depends(get_db), form_data: OAuth2PasswordReq
"scope_id": str(user.scope_id) if user.scope_id else str(user.id)
}
access, refresh = create_tokens(data=token_data)
access, refresh = create_tokens(data=token_data, remember_me=remember_me)
# Kinyerjük a beállításokból, hogy meddig él a refresh token (itt is, vagy a create_tokens is ezt csinálja)
# A SSoT szerint auth_remember_me_days = 30
if remember_me:
max_age_days = await config.get_setting(db, "auth_remember_me_days", default=30)
else:
max_age_days = await config.get_setting(db, "auth_refresh_default_days", default=1)
max_age_sec = int(max_age_days) * 24 * 60 * 60
response.set_cookie(
key="refresh_token",
value=refresh,
httponly=True,
secure=True, # Use Secure
samesite="lax",
max_age=max_age_sec
)
return {"access_token": access, "refresh_token": refresh, "token_type": "bearer", "is_active": user.is_active}
class VerifyEmailRequest(BaseModel):
@@ -58,8 +90,42 @@ async def verify_email(request: VerifyEmailRequest, db: AsyncSession = Depends(g
"""
success = await AuthService.verify_email(db, request.token)
if not success:
raise HTTPException(status_code=400, detail="Érvénytelen vagy lejárt token.")
return {"status": "success", "message": "Email sikeresen megerősítve."}
raise HTTPException(status_code=400, detail=t("AUTH.INVALID_OR_EXPIRED_TOKEN"))
return {"status": "success", "message": t("AUTH.EMAIL_VERIFICATION_SUCCESS")}
class ForgotPasswordRequest(BaseModel):
email: EmailStr = Field(..., description="Email cím a jelszó visszaállításhoz")
@router.post("/forgot-password")
async def forgot_password(request: ForgotPasswordRequest, db: AsyncSession = Depends(get_db)):
"""
Elfelejtett jelszó folyamat indítása.
Mindig sikeres választ ad, hogy megelőzzük az email enumerációt.
"""
result = await AuthService.initiate_password_reset(db, request.email)
# Mindig ugyanazt a választ adjuk, függetlenül attól, hogy létezik-e a felhasználó
return {
"status": "success",
"message": "Ha ez az email cím regisztrálva van, akkor elküldtünk egy jelszó-visszaállítási linket."
}
class ResetPasswordRequest(BaseModel):
email: EmailStr = Field(..., description="Email cím")
token: str = Field(..., description="Jelszó visszaállítási token")
new_password: str = Field(..., description="Új jelszó")
@router.post("/reset-password")
async def reset_password(request: ResetPasswordRequest, db: AsyncSession = Depends(get_db)):
"""
Jelszó visszaállítása token alapján.
"""
success = await AuthService.reset_password(db, request.email, request.token, request.new_password)
if not success:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Érvénytelen token, lejárt token vagy nem létező felhasználó."
)
return {"status": "success", "message": "Jelszó sikeresen megváltoztatva."}
@router.post("/complete-kyc")
async def complete_kyc(kyc_in: UserKYCComplete, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user)):

View File

@@ -3,18 +3,19 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.db.session import get_db
from app.services.asset_service import AssetService
from app.api import deps
from typing import List
from typing import List, Optional
router = APIRouter()
# Secured endpoint: Closed premium ecosystem
@router.get("/makes", response_model=List[str])
async def list_makes(
vehicle_class: Optional[str] = None,
db: AsyncSession = Depends(get_db),
current_user = Depends(deps.get_current_user)
):
"""1. Szint: Márkák listázása."""
return await AssetService.get_makes(db)
"""1. Szint: Márkák listázása, opcionálisan vehicle_class szerint szűrve."""
return await AssetService.get_makes(db, vehicle_class)
# Secured endpoint: Closed premium ecosystem
@router.get("/models", response_model=List[str])

View File

@@ -13,7 +13,7 @@ from sqlalchemy import select
from app.db.session import get_db
from app.api.deps import get_current_user
from app.schemas.organization import CorpOnboardIn, CorpOnboardResponse
from app.models.marketplace.organization import Organization, OrgType, OrganizationMember
from app.models.marketplace.organization import Organization, OrgType, OrganizationMember, Branch
from app.models.identity import User # JAVÍTVA: Központi Identity modell
from app.core.config import settings
@@ -82,9 +82,24 @@ async def onboard_organization(
)
db.add(new_org)
await db.flush()
await db.flush()
# 5. TULAJDONOS RÖGZÍTÉSE
# 5. ALAPÉRTELMEZETT KÖZPONTI TELEPHELY LÉTREHOZÁSA
main_branch = Branch(
organization_id=new_org.id,
name="Központi Telephely",
is_main=True,
postal_code=org_in.address_zip,
city=org_in.address_city,
street_name=org_in.address_street_name,
street_type=org_in.address_street_type,
house_number=org_in.address_house_number,
hrsz=org_in.address_hrsz,
status="active"
)
db.add(main_branch)
# 6. TULAJDONOS RÖGZÍTÉSE
owner_member = OrganizationMember(
organization_id=new_org.id,
user_id=current_user.id,
@@ -92,7 +107,7 @@ async def onboard_organization(
)
db.add(owner_member)
# 6. NAS Mappa létrehozása
# 7. NAS Mappa létrehozása
try:
base_path = getattr(settings, "NAS_STORAGE_PATH", "/mnt/nas/app_data")
org_path = os.path.join(base_path, "organizations", str(new_org.id))
@@ -135,4 +150,128 @@ async def get_my_organizations(
"subscription_plan": o.subscription_plan
}
for o in orgs
]
]
# --- B2B MEGHÍVÓ LOGIKA ---
from pydantic import BaseModel, EmailStr
from app.models.identity import VerificationToken
import uuid
from datetime import timedelta
class OrgInvitationIn(BaseModel):
email: EmailStr
role: str = "DRIVER"
class OrgInvitationResponse(BaseModel):
status: str
message: str
@router.post("/{org_id}/invitations", response_model=OrgInvitationResponse, status_code=status.HTTP_200_OK)
async def invite_to_organization(
org_id: int,
invite_in: OrgInvitationIn,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
B2B Meghívó küldése egy szervezetbe.
Ha a felhasználó már létezik, egy pending tagot hozunk létre.
Ha nem létezik, token készül az email címre.
"""
# 1. Jogosultság ellenőrzése
stmt_member = select(OrganizationMember).where(
(OrganizationMember.organization_id == org_id) &
(OrganizationMember.user_id == current_user.id) &
(OrganizationMember.role.in_(["OWNER", "ADMIN"]))
)
member = (await db.execute(stmt_member)).scalar_one_or_none()
if not member:
raise HTTPException(status_code=403, detail="Nincs jogosultságod meghívót küldeni (csak OWNER/ADMIN).")
# 2. Célpont keresése
stmt_target = select(User).where(User.email == invite_in.email)
target_user = (await db.execute(stmt_target)).scalar_one_or_none()
if target_user:
# Létező felhasználó, van-e már tagsága?
stmt_exist = select(OrganizationMember).where(
(OrganizationMember.organization_id == org_id) &
(OrganizationMember.user_id == target_user.id)
)
if (await db.execute(stmt_exist)).scalar_one_or_none():
raise HTTPException(status_code=400, detail="A felhasználó már tagja a szervezetnek.")
new_member = OrganizationMember(
organization_id=org_id,
user_id=target_user.id,
role=invite_in.role.upper(),
status="pending" # Válaszolni kell a meghívóra
)
db.add(new_member)
await db.commit()
logger.info(f"Értesítő email küldve a létező felhasználónak: {invite_in.email}")
return {"status": "success", "message": "Meghívó elküldve a meglévő felhasználónak."}
else:
# Új felhasználó -> Token
token_val = uuid.uuid4()
new_token = VerificationToken(
token=token_val,
user_id=None, # Mivel még nincs User
token_type="org_invite",
expires_at=datetime.now(timezone.utc) + timedelta(days=7),
extra_data={"org_id": org_id, "role": invite_in.role.upper(), "email": invite_in.email}
)
db.add(new_token)
await db.commit()
logger.info(f"Meghívó email küldve az új felhasználónak: {invite_in.email}, Token: {token_val}")
return {"status": "success", "message": "Meghívó email elküldve (új felhasználó)."}
@router.post("/invitations/{token}/accept", response_model=OrgInvitationResponse, status_code=status.HTTP_200_OK)
async def accept_invitation(
token: uuid.UUID,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
Meghívó elfogadása token alapján. (Új felhasználó számára, miután regisztrált)
"""
stmt = select(VerificationToken).where(
(VerificationToken.token == token) &
(VerificationToken.token_type == "org_invite") &
(VerificationToken.is_used == False)
)
token_rec = (await db.execute(stmt)).scalar_one_or_none()
if not token_rec or token_rec.expires_at < datetime.now(timezone.utc):
raise HTTPException(status_code=400, detail="Érvénytelen vagy lejárt meghívó.")
extra = token_rec.extra_data
if not extra or extra.get("email") != current_user.email:
raise HTTPException(status_code=403, detail="A meghívó nem a te e-mail címedre szól.")
org_id = extra.get("org_id")
role = extra.get("role")
# Van-e már tagsága?
stmt_exist = select(OrganizationMember).where(
(OrganizationMember.organization_id == org_id) &
(OrganizationMember.user_id == current_user.id)
)
exist_member = (await db.execute(stmt_exist)).scalar_one_or_none()
if exist_member:
exist_member.status = "active"
exist_member.role = role
else:
new_member = OrganizationMember(
organization_id=org_id,
user_id=current_user.id,
role=role,
status="active"
)
db.add(new_member)
token_rec.is_used = True
await db.commit()
return {"status": "success", "message": "Meghívó sikeresen elfogadva."}

View File

@@ -2,7 +2,9 @@
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.exc import SQLAlchemyError
from typing import Dict, Any
from sqlalchemy import select, or_
from typing import Dict, Any, List, Optional
from datetime import datetime
from app.api.deps import get_db, get_current_user
from app.schemas.user import UserResponse, UserUpdate, ActiveOrganizationUpdate, UserWithTokenResponse
@@ -10,10 +12,27 @@ from app.models.identity import User
from app.services.trust_engine import TrustEngine
from app.core.security import create_tokens, DEFAULT_RANK_MAP
from app.core.config import settings
from pydantic import BaseModel
router = APIRouter()
trust_engine = TrustEngine()
# MLM Hálózat válasz sémák
class NetworkMemberL1(BaseModel):
email: str
referral_code: str
folder_slug: Optional[str]
joined_at: datetime
class NetworkMemberL2L3(BaseModel):
referral_code: str
joined_at: datetime
class NetworkResponse(BaseModel):
level1: List[NetworkMemberL1]
level2: List[NetworkMemberL2L3]
level3: List[NetworkMemberL2L3]
@router.get("/me", response_model=UserResponse)
async def read_users_me(
db: AsyncSession = Depends(get_db),
@@ -222,3 +241,67 @@ async def update_active_organization(
access_token=access_token,
token_type="bearer"
)
@router.get("/me/network", response_model=NetworkResponse)
async def get_my_network(
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user),
):
"""
Visszaadja a felhasználó MLM hálózatát 3 szinten:
- L1: Közvetlen meghívottak (email + kód)
- L2: L1 meghívottai (csak kód)
- L3: L2 meghívottai (csak kód)
"""
# L1: Közvetlen meghívottak
l1_stmt = select(User).where(User.referred_by_id == current_user.id)
l1_result = await db.execute(l1_stmt)
l1_users = l1_result.scalars().all()
level1 = [
NetworkMemberL1(
email=u.email,
referral_code=u.referral_code or "",
folder_slug=u.folder_slug,
joined_at=u.created_at
)
for u in l1_users
]
# L2: L1 userek által meghívottak
level2 = []
if l1_users:
l1_ids = [u.id for u in l1_users]
l2_stmt = select(User).where(User.referred_by_id.in_(l1_ids))
l2_result = await db.execute(l2_stmt)
l2_users = l2_result.scalars().all()
level2 = [
NetworkMemberL2L3(
referral_code=u.referral_code or "",
joined_at=u.created_at
)
for u in l2_users
]
# L3: L2 userek által meghívottak
level3 = []
if l2_users:
l2_ids = [u.id for u in l2_users]
l3_stmt = select(User).where(User.referred_by_id.in_(l2_ids))
l3_result = await db.execute(l3_stmt)
l3_users = l3_result.scalars().all()
level3 = [
NetworkMemberL2L3(
referral_code=u.referral_code or "",
joined_at=u.created_at
)
for u in l3_users
]
return NetworkResponse(
level1=level1,
level2=level2,
level3=level3
)

View File

@@ -1,3 +1,4 @@
# /opt/docker/dev/service_finder/backend/app/api/v1/endpoints/vehicles.py
"""
Jármű értékelési végpontok a Social 1 modulhoz.
"""

View File

@@ -0,0 +1,31 @@
# /opt/docker/dev/service_finder/backend/app/core/context.py
"""
Context management for request-scoped variables, particularly locale/language.
Uses contextvars to store the current request's locale for i18n purposes.
"""
import contextvars
from typing import Optional
# Context variable to store the current request's locale
# Default value is "hu" (Hungarian) as per project requirements
current_locale: contextvars.ContextVar[str] = contextvars.ContextVar(
"current_locale", default="hu"
)
def get_current_locale() -> str:
"""
Get the current locale from the context.
Returns the default "hu" if not set.
"""
return current_locale.get()
def set_current_locale(locale: str) -> None:
"""
Set the current locale in the context.
Should be called by middleware or other request processors.
"""
current_locale.set(locale)
# Alias for convenience
get_locale = get_current_locale
set_locale = set_current_locale

View File

@@ -0,0 +1,108 @@
# /opt/docker/dev/service_finder/backend/app/core/i18n_middleware.py
"""
FastAPI middleware for automatic locale/language detection and context management.
Implements the priority order:
1. ?lang= query parameter
2. Accept-Language HTTP header
3. User profile language (if authenticated)
4. Default "hu" (Hungarian)
"""
import logging
from typing import Optional
from fastapi import Request, HTTPException
from starlette.middleware.base import BaseHTTPMiddleware
from starlette.responses import Response
from app.core.context import set_current_locale, get_current_locale
from app.core.config import settings
logger = logging.getLogger(__name__)
class I18nMiddleware(BaseHTTPMiddleware):
"""
Middleware for automatic locale detection and context management.
Sets the locale in contextvars for the duration of the request.
"""
async def dispatch(self, request: Request, call_next):
# Determine locale based on priority
locale = self._determine_locale(request)
# Set locale in context
set_current_locale(locale)
# Add locale to request state for debugging/logging
request.state.locale = locale
# Process request
response = await call_next(request)
# Optionally add locale header to response
response.headers["X-Content-Language"] = locale
return response
def _determine_locale(self, request: Request) -> str:
"""
Determine the locale based on the priority order.
Returns a valid locale code (e.g., "hu", "en", "de").
"""
# 1. Query parameter: ?lang=
query_lang = request.query_params.get("lang")
if query_lang and self._is_valid_locale(query_lang):
logger.debug(f"Locale from query param: {query_lang}")
return query_lang
# 2. Accept-Language HTTP header
accept_language = request.headers.get("accept-language")
if accept_language:
header_lang = self._parse_accept_language(accept_language)
if header_lang and self._is_valid_locale(header_lang):
logger.debug(f"Locale from Accept-Language header: {header_lang}")
return header_lang
# 3. User profile language (if authenticated)
# Note: This requires the user to be authenticated, which may not be available
# in middleware before authentication. We'll handle this in endpoint dependencies.
# For now, we'll skip this step in middleware and let endpoints handle it.
# 4. Default locale
default_locale = getattr(settings, "DEFAULT_LOCALE", "hu")
logger.debug(f"Using default locale: {default_locale}")
return default_locale
def _parse_accept_language(self, accept_language: str) -> Optional[str]:
"""
Parse Accept-Language header and return the highest priority locale.
Example: "en-US,en;q=0.9,hu;q=0.8" -> "en"
"""
try:
# Split by comma and process each language range
languages = accept_language.split(',')
for lang_range in languages:
# Remove quality factor if present
lang = lang_range.split(';')[0].strip()
# Extract primary language subtag (e.g., "en" from "en-US")
if '-' in lang:
lang = lang.split('-')[0]
if self._is_valid_locale(lang):
return lang
except Exception as e:
logger.debug(f"Error parsing Accept-Language header: {e}")
return None
def _is_valid_locale(self, locale: str) -> bool:
"""
Check if a locale code is valid/supported.
For now, we accept any 2-5 character locale code.
In production, you might want to check against a list of supported locales.
"""
if not locale or len(locale) < 2 or len(locale) > 5:
return False
# Basic validation: only letters and hyphens
return locale.replace('-', '').isalpha()
# Create middleware instance (commented out to avoid instantiation error during import)
# i18n_middleware = I18nMiddleware()

View File

@@ -14,18 +14,27 @@ def verify_password(plain_password: str, hashed_password: str) -> bool:
def get_password_hash(password: str) -> str:
return bcrypt.hashpw(password.encode("utf-8"), bcrypt.gensalt()).decode("utf-8")
def create_tokens(data: Dict[str, Any]) -> Tuple[str, str]:
""" Access és Refresh token generálása UTC időzónával. """
def create_tokens(data: Dict[str, Any], remember_me: bool = False) -> Tuple[str, str]:
""" Access és Refresh token generálása UTC időzónával.
Args:
data: Token payload adatok
remember_me: Ha True, a refresh token lejárata 30 nap, egyébként 1 nap
"""
to_encode = data.copy()
now = datetime.now(timezone.utc)
# Access Token
# Access Token (mindig ugyanaz a lejárat)
acc_expire = now + timedelta(minutes=settings.ACCESS_TOKEN_EXPIRE_MINUTES)
access_payload = {**to_encode, "exp": acc_expire, "iat": now, "type": "access"}
access_token = jwt.encode(access_payload, settings.SECRET_KEY, algorithm=settings.ALGORITHM)
# Refresh Token
ref_expire = now + timedelta(days=settings.REFRESH_TOKEN_EXPIRE_DAYS)
# Refresh Token lejárat a remember_me alapján
if remember_me:
ref_expire = now + timedelta(days=30) # 30 nap remember me esetén
else:
ref_expire = now + timedelta(days=1) # 1 nap alapértelmezett
refresh_payload = {"sub": str(to_encode.get("sub")), "exp": ref_expire, "iat": now, "type": "refresh"}
refresh_token = jwt.encode(refresh_payload, settings.SECRET_KEY, algorithm=settings.ALGORITHM)

View File

@@ -0,0 +1,23 @@
# /opt/docker/dev/service_finder/backend/app/core/translation_helper.py
"""
Helper functions for easy translation usage throughout the application.
Provides convenient access to the TranslationService with automatic locale detection.
"""
from typing import Optional, Dict, Any
from app.services.translation_service import TranslationService
def t(key: str, variables: Optional[Dict[str, Any]] = None, lang: Optional[str] = None) -> str:
"""
Shortcut function for TranslationService.get_text().
Automatically uses the current request locale from context.
Usage:
t("AUTH.REGISTRATION_SUCCESS")
t("AUTH.WELCOME", {"name": "John"})
t("ERROR.INVALID_TOKEN", lang="en")
"""
return TranslationService.get_text(key, lang=lang, variables=variables)
# Alias for backward compatibility
get_text = t
translate = t

View File

@@ -59,8 +59,12 @@ app = FastAPI(
)
# --- MIDDLEWARES ---
# I18n middleware should come early to set locale context for all subsequent processing
from app.core.i18n_middleware import I18nMiddleware
app.add_middleware(I18nMiddleware)
app.add_middleware(
SessionMiddleware,
SessionMiddleware,
secret_key=settings.SECRET_KEY
)

View File

@@ -209,8 +209,9 @@ class VerificationToken(Base):
id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
token: Mapped[uuid.UUID] = mapped_column(PG_UUID(as_uuid=True), default=uuid.uuid4, unique=True, nullable=False)
user_id: Mapped[int] = mapped_column(Integer, ForeignKey("identity.users.id", ondelete="CASCADE"), nullable=False)
user_id: Mapped[Optional[int]] = mapped_column(Integer, ForeignKey("identity.users.id", ondelete="CASCADE"), nullable=True)
token_type: Mapped[str] = mapped_column(String(20), nullable=False)
extra_data: Mapped[Any] = mapped_column(JSON, server_default=text("'{}'::jsonb"), nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now())
expires_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
is_used: Mapped[bool] = mapped_column(Boolean, default=False)
@@ -224,7 +225,7 @@ class SocialAccount(Base):
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
user_id: Mapped[int] = mapped_column(Integer, ForeignKey("identity.users.id", ondelete="CASCADE"), nullable=False)
user_id: Mapped[Optional[int]] = mapped_column(Integer, ForeignKey("identity.users.id", ondelete="CASCADE"), nullable=True)
provider: Mapped[str] = mapped_column(String(50), nullable=False)
social_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True)
email: Mapped[str] = mapped_column(String(255), nullable=False)

View File

@@ -149,6 +149,7 @@ class AssetCreate(BaseModel):
# === ORGANIZATION (Optional) ===
organization_id: Optional[int] = Field(None, description="Szervezet ID (alapértelmezett a felhasználó szervezete)")
branch_id: Optional[UUID] = Field(None, description="Garázs (Branch) ID, ahova a járművet rendeljük. Ha nincs megadva, a szervezet központi garázsába kerül.")
# === STATUS VALIDATION ===
@validator('status', pre=True, always=True)

View File

@@ -21,16 +21,18 @@ class UserLiteRegister(BaseModel):
region_code: Optional[str] = "HU"
lang: Optional[str] = "hu"
timezone: Optional[str] = "Europe/Budapest"
referred_by_code: Optional[str] = None # Meghívó referral kódja (opcionális)
model_config = ConfigDict(from_attributes=True)
class UserKYCComplete(BaseModel):
""" Step 2: Teljes körű személyazonosítás és címadatok. """
phone_number: str = Field(..., pattern=r"^\+?[0-9]{7,15}$")
birth_place: str
birth_date: date
mothers_last_name: str
mothers_first_name: str
phone_number: Optional[str] = Field(default=None, pattern=r"^\+?[0-9]{7,15}$")
birth_place: Optional[str] = None
birth_date: Optional[date] = None
mothers_last_name: Optional[str] = None
mothers_first_name: Optional[str] = None
region_code: Optional[str] = "HU"
# Atomizált címadatok a pontos GPS-hez és Robot-munkához
address_zip: str
@@ -45,7 +47,7 @@ class UserKYCComplete(BaseModel):
# Okmányok és Vészhelyzet
identity_docs: Dict[str, DocumentDetail] # pl: {"ID_CARD": {...}, "LICENSE": {...}}
ice_contact: ICEContact
ice_contact: Optional[ICEContact] = None
preferred_language: str = "hu"
preferred_currency: str = "HUF"

View File

@@ -0,0 +1,58 @@
#!/usr/bin/env python3
"""
MLM és Gamification paraméterek inicializálása a system_parameters táblába.
Futtatás: docker compose exec sf_api python3 /app/backend/scripts/init_mlm_parameters.py
"""
import asyncio
import sys
import os
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from sqlalchemy.orm import sessionmaker
from sqlalchemy import select
from app.models.system.parameter import SystemParameter
from app.core.config import settings
async def init_parameters():
engine = create_async_engine(settings.DATABASE_URL, echo=False)
async_session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with async_session() as session:
# MLM százalékok
mlm_params = [
("mlm_level1_percent", "10", "MLM L1 jutalék százalék (10%)", "global"),
("mlm_level2_percent", "5", "MLM L2 jutalék százalék (5%)", "global"),
("mlm_level3_percent", "3", "MLM L3 jutalék százalék (3%)", "global"),
("gamification_p2p_invite_xp", "50", "XP pontok a meghívónak sikeres KYC után", "global"),
]
for key, value, description, scope in mlm_params:
existing = await session.execute(
select(SystemParameter).where(SystemParameter.key == key)
)
if existing.scalar_one_or_none():
print(f"{key} már létezik, kihagyva.")
continue
param = SystemParameter(
key=key,
value=value,
description=description,
scope=scope,
data_type="integer" if key.endswith("_percent") else "integer",
is_editable=True,
is_visible=True,
region_code=None,
user_id=None,
)
session.add(param)
print(f"{key} = {value} beállítva.")
await session.commit()
print("✅ MLM paraméterek sikeresen inicializálva.")
if __name__ == "__main__":
asyncio.run(init_parameters())

View File

@@ -0,0 +1,57 @@
#!/usr/bin/env python3
"""
MLM és Gamification paraméterek inicializálása a system_parameters táblába.
Futtatás: docker compose exec sf_api python3 /app/backend/app/scripts/init_mlm_parameters_fixed.py
"""
import asyncio
import sys
import os
sys.path.append(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from sqlalchemy.orm import sessionmaker
from sqlalchemy import select
from app.models.system.system import SystemParameter, ParameterScope
from app.core.config import settings
async def init_parameters():
engine = create_async_engine(settings.DATABASE_URL, echo=False)
async_session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with async_session() as session:
# MLM százalékok
mlm_params = [
("mlm_level1_percent", 10, "MLM L1 jutalék százalék (10%)", ParameterScope.GLOBAL),
("mlm_level2_percent", 5, "MLM L2 jutalék százalék (5%)", ParameterScope.GLOBAL),
("mlm_level3_percent", 3, "MLM L3 jutalék százalék (3%)", ParameterScope.GLOBAL),
("gamification_p2p_invite_xp", 50, "XP pontok a meghívónak sikeres KYC után", ParameterScope.GLOBAL),
]
for key, value, description, scope in mlm_params:
existing = await session.execute(
select(SystemParameter).where(SystemParameter.key == key)
)
if existing.scalar_one_or_none():
print(f"{key} már létezik, kihagyva.")
continue
param = SystemParameter(
key=key,
value={"value": value},
category="mlm",
scope_level=scope,
scope_id=None,
is_active=True,
description=description,
last_modified_by=None,
)
session.add(param)
print(f"{key} = {value} beállítva.")
await session.commit()
print("✅ MLM paraméterek sikeresen inicializálva.")
if __name__ == "__main__":
asyncio.run(init_parameters())

View File

@@ -1,6 +1,6 @@
# /opt/docker/dev/service_finder/backend/app/scripts/sync_engine.py
#!/usr/bin/env python3
# docker exec -it sf_api python -m app.scripts.sync_engine
# cd /opt/docker/dev/service_finder && docker exec -it sf_api python -m app.scripts.sync_engine
import asyncio
import importlib
import sys

View File

@@ -0,0 +1,320 @@
#!/usr/bin/env python3
"""
MLM Payout E2E teszt script.
Szimulál egy előfizetési fizetést (10,000 HUF) és ellenőrzi a CREDIT Wallet jóváírásokat L1, L2, L3 usereknek.
Futtatás: docker compose exec sf_api python3 /app/backend/test_mlm_payout.py
"""
import asyncio
import sys
import os
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from sqlalchemy.orm import sessionmaker
from sqlalchemy import select, update, func
from decimal import Decimal
from app.models.identity import User, Person
from app.models.finance.wallet import Wallet, Ledger, WalletType, TransactionType
from app.models.system.system import SystemParameter
from app.core.config import settings
async def test_mlm_payout():
engine = create_async_engine(settings.DATABASE_URL, echo=True)
async_session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with async_session() as session:
print("=== MLM Payout E2E Teszt ===")
# 1. Teszt adatok létrehozása (ha nincsenek)
# Létrehozunk egy L1, L2, L3, L4 (fizető) usert láncolatban
print("1. Teszt userek létrehozása...")
# Ellenőrizzük, hogy van-e már teszt user
existing = await session.execute(select(User).where(User.email == "test_l1@example.com"))
l1_user = existing.scalar_one_or_none()
if not l1_user:
# L1 user (root meghívó)
l1_person = Person(
first_name="L1",
last_name="Referrer",
is_active=True,
identity_docs={},
ice_contact={},
lifetime_xp=0,
penalty_points=0,
social_reputation=1.0,
is_sales_agent=False,
is_ghost=False,
)
session.add(l1_person)
await session.flush()
l1_user = User(
email="test_l1@example.com",
hashed_password="$2b$12$...", # dummy
person_id=l1_person.id,
role="user",
is_active=True,
is_deleted=False,
region_code="HU",
preferred_language="hu",
subscription_plan="FREE",
is_vip=True,
preferred_currency="HUF",
scope_level="individual",
custom_permissions={},
referral_code="L1REF123",
referred_by_id=None,
)
session.add(l1_user)
await session.flush()
print(f" L1 user létrehozva: ID={l1_user.id}, kód={l1_user.referral_code}")
# L2 user (L1 meghívja)
existing = await session.execute(select(User).where(User.email == "test_l2@example.com"))
l2_user = existing.scalar_one_or_none()
if not l2_user:
l2_person = Person(
first_name="L2",
last_name="Referral",
is_active=True,
identity_docs={},
ice_contact={},
lifetime_xp=0,
penalty_points=0,
social_reputation=1.0,
is_sales_agent=False,
is_ghost=False,
)
session.add(l2_person)
await session.flush()
l2_user = User(
email="test_l2@example.com",
hashed_password="$2b$12$...",
person_id=l2_person.id,
role="user",
is_active=True,
is_deleted=False,
region_code="HU",
preferred_language="hu",
subscription_plan="FREE",
is_vip=True,
preferred_currency="HUF",
scope_level="individual",
custom_permissions={},
referral_code="L2REF456",
referred_by_id=l1_user.id,
)
session.add(l2_user)
await session.flush()
print(f" L2 user létrehozva: ID={l2_user.id}, kód={l2_user.referral_code}, referred_by={l2_user.referred_by_id}")
# L3 user (L2 meghívja)
existing = await session.execute(select(User).where(User.email == "test_l3@example.com"))
l3_user = existing.scalar_one_or_none()
if not l3_user:
l3_person = Person(
first_name="L3",
last_name="Referral",
is_active=True,
identity_docs={},
ice_contact={},
lifetime_xp=0,
penalty_points=0,
social_reputation=1.0,
is_sales_agent=False,
is_ghost=False,
)
session.add(l3_person)
await session.flush()
l3_user = User(
email="test_l3@example.com",
hashed_password="$2b$12$...",
person_id=l3_person.id,
role="user",
is_active=True,
is_deleted=False,
region_code="HU",
preferred_language="hu",
subscription_plan="FREE",
is_vip=True,
preferred_currency="HUF",
scope_level="individual",
custom_permissions={},
referral_code="L3REF789",
referred_by_id=l2_user.id,
)
session.add(l3_user)
await session.flush()
print(f" L3 user létrehozva: ID={l3_user.id}, kód={l3_user.referral_code}, referred_by={l3_user.referred_by_id}")
# L4 user (fizető, L3 meghívja)
existing = await session.execute(select(User).where(User.email == "test_l4_payer@example.com"))
l4_user = existing.scalar_one_or_none()
if not l4_user:
l4_person = Person(
first_name="L4",
last_name="Payer",
is_active=True,
identity_docs={},
ice_contact={},
lifetime_xp=0,
penalty_points=0,
social_reputation=1.0,
is_sales_agent=False,
is_ghost=False,
)
session.add(l4_person)
await session.flush()
l4_user = User(
email="test_l4_payer@example.com",
hashed_password="$2b$12$...",
person_id=l4_person.id,
role="user",
is_active=True,
is_deleted=False,
region_code="HU",
preferred_language="hu",
subscription_plan="FREE",
is_vip=True,
preferred_currency="HUF",
scope_level="individual",
custom_permissions={},
referral_code="L4REF999",
referred_by_id=l3_user.id,
)
session.add(l4_user)
await session.flush()
print(f" L4 (fizető) user létrehozva: ID={l4_user.id}, referred_by={l4_user.referred_by_id}")
await session.commit()
# 2. MLM paraméterek lekérése
print("\n2. MLM paraméterek lekérése...")
params = {}
for key in ["mlm_level1_percent", "mlm_level2_percent", "mlm_level3_percent"]:
stmt = select(SystemParameter).where(SystemParameter.key == key)
param = (await session.execute(stmt)).scalar_one_or_none()
if param:
params[key] = int(param.value.get("value", 0))
print(f" {key}: {params[key]}%")
else:
params[key] = 0
print(f" {key}: NINCS BEÁLLÍTVA!")
# 3. Szimulált fizetés (10,000 HUF)
payment_amount = Decimal("10000.00")
print(f"\n3. Szimulált fizetés: {payment_amount} HUF (L4 user fizet)")
# 4. MLM lánc felépítése
print("\n4. MLM lánc felépítése...")
chain = []
current = l4_user
for i in range(3):
if current.referred_by_id:
stmt = select(User).where(User.id == current.referred_by_id)
referrer = (await session.execute(stmt)).scalar_one_or_none()
if referrer:
chain.append(referrer)
current = referrer
else:
break
else:
break
print(f" Lánc hossza: {len(chain)}")
for idx, user in enumerate(chain):
print(f" L{idx+1}: {user.email} (ID: {user.id})")
# 5. Jutalékok kiszámítása és CREDIT Wallet jóváírás
print("\n5. Jutalékok kiszámítása és CREDIT Wallet jóváírás...")
levels = ["mlm_level1_percent", "mlm_level2_percent", "mlm_level3_percent"]
for idx, referrer in enumerate(chain):
if idx >= 3:
break
percent = params[levels[idx]]
commission = (payment_amount * Decimal(percent) / Decimal(100)).quantize(Decimal("0.01"))
print(f" L{idx+1} ({referrer.email}): {percent}% -> {commission} HUF")
# CREDIT Wallet keresése vagy létrehozása
wallet_stmt = select(Wallet).where(
Wallet.user_id == referrer.id,
Wallet.wallet_type == WalletType.CREDIT
)
wallet = (await session.execute(wallet_stmt)).scalar_one_or_none()
if not wallet:
wallet = Wallet(
user_id=referrer.id,
currency="HUF",
wallet_type=WalletType.CREDIT,
balance=Decimal("0.00"),
is_active=True
)
session.add(wallet)
await session.flush()
# Balance frissítése
wallet.balance += commission
# Ledger bejegyzés
ledger = Ledger(
wallet_id=wallet.id,
amount=commission,
transaction_type=TransactionType.MLM_CREDIT,
description=f"MLM jutalék L{idx+1} a(z) {l4_user.email} fizetéséből",
reference_id=l4_user.id,
reference_type="user_payment",
metadata={
"payer_id": l4_user.id,
"payer_email": l4_user.email,
"level": idx+1,
"percent": percent,
"original_amount": float(payment_amount)
}
)
session.add(ledger)
await session.commit()
# 6. Ellenőrzés
print("\n6. Ellenőrzés - CREDIT Wallet egyenlegek:")
for idx, referrer in enumerate(chain):
if idx >= 3:
break
wallet_stmt = select(Wallet).where(
Wallet.user_id == referrer.id,
Wallet.wallet_type == WalletType.CREDIT
)
wallet = (await session.execute(wallet_stmt)).scalar_one_or_none()
if wallet:
print(f" L{idx+1} ({referrer.email}): {wallet.balance} HUF")
else:
print(f" L{idx+1} ({referrer.email}): NINCS CREDIT Wallet!")
# 7. Ledger bejegyzések listázása
print("\n7. Ledger bejegyzések:")
for referrer in chain[:3]:
wallet_stmt = select(Wallet).where(
Wallet.user_id == referrer.id,
Wallet.wallet_type == WalletType.CREDIT
)
wallet = (await session.execute(wallet_stmt)).scalar_one_or_none()
if wallet:
ledger_stmt = select(Ledger).where(Ledger.wallet_id == wallet.id).order_by(Ledger.created_at.desc()).limit(2)
ledgers = (await session.execute(ledger_stmt)).scalars().all()
for l in ledgers:
print(f" - {referrer.email}: {l.amount} HUF ({l.transaction_type}) - {l.description}")
print("\n✅ MLM Payout teszt sikeresen lefutott!")
print("=== TESZT VÉGE ===")
if __name__ == "__main__":
asyncio.run(test_mlm_payout())

View File

@@ -0,0 +1,145 @@
#!/usr/bin/env python3
"""
MLM Payout E2E teszt script (egyszerűsített).
Szimulál egy előfizetési fizetést (10,000 HUF) és ellenőrzi a CREDIT Wallet jóváírásokat L1, L2, L3 usereknek.
Futtatás: docker compose exec sf_api python3 /app/app/scripts/test_mlm_payout_simple.py
"""
import asyncio
import sys
import os
sys.path.append(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from sqlalchemy.orm import sessionmaker
from sqlalchemy import select, update, func
from decimal import Decimal
from app.models.identity import User, Person, Wallet
from app.models.system.system import SystemParameter
from app.core.config import settings
async def test_mlm_payout():
engine = create_async_engine(settings.DATABASE_URL, echo=True)
async_session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with async_session() as session:
print("=== MLM Payout E2E Teszt (Egyszerűsített) ===")
# 1. MLM paraméterek lekérése
print("\n1. MLM paraméterek lekérése...")
params = {}
for key in ["mlm_level1_percent", "mlm_level2_percent", "mlm_level3_percent"]:
stmt = select(SystemParameter).where(SystemParameter.key == key)
param = (await session.execute(stmt)).scalar_one_or_none()
if param:
params[key] = int(param.value.get("value", 0))
print(f" {key}: {params[key]}%")
else:
params[key] = 0
print(f" {key}: NINCS BEÁLLÍTVA!")
# 2. Meglévő teszt userek keresése vagy létrehozása
print("\n2. Teszt userek keresése...")
# L1 user keresése (root meghívó)
l1_stmt = select(User).where(User.email == "test_l1@example.com")
l1_user = (await session.execute(l1_stmt)).scalar_one_or_none()
if not l1_user:
print(" L1 user nem létezik, teszt adatok hiányoznak.")
print(" Futtasd először az init_mlm_parameters_fixed.py-t és regisztrálj néhány teszt usert!")
return
# L2 user keresése (L1 meghívja)
l2_stmt = select(User).where(User.email == "test_l2@example.com")
l2_user = (await session.execute(l2_stmt)).scalar_one_or_none()
# L3 user keresése (L2 meghívja)
l3_stmt = select(User).where(User.email == "test_l3@example.com")
l3_user = (await session.execute(l3_stmt)).scalar_one_or_none()
# L4 user keresése (fizető, L3 meghívja)
l4_stmt = select(User).where(User.email == "test_l4_payer@example.com")
l4_user = (await session.execute(l4_stmt)).scalar_one_or_none()
users = [l1_user, l2_user, l3_user, l4_user]
user_names = ["L1", "L2", "L3", "L4 (fizető)"]
for idx, (user, name) in enumerate(zip(users, user_names)):
if user:
print(f" {name}: {user.email} (ID: {user.id}, referred_by: {user.referred_by_id})")
else:
print(f" {name}: NEM LÉTEZIK")
# 3. Szimulált fizetés (10,000 HUF)
payment_amount = Decimal("10000.00")
print(f"\n3. Szimulált fizetés: {payment_amount} HUF (L4 user fizet)")
# 4. MLM lánc felépítése
print("\n4. MLM lánc felépítése...")
chain = []
if l4_user and l4_user.referred_by_id:
# L3 keresése
if l3_user and l3_user.id == l4_user.referred_by_id:
chain.append(l3_user)
# L2 keresése
if l3_user.referred_by_id:
if l2_user and l2_user.id == l3_user.referred_by_id:
chain.append(l2_user)
# L1 keresése
if l2_user.referred_by_id:
if l1_user and l1_user.id == l2_user.referred_by_id:
chain.append(l1_user)
print(f" Lánc hossza: {len(chain)}")
for idx, user in enumerate(chain):
print(f" L{idx+1}: {user.email} (ID: {user.id})")
# 5. Jutalékok kiszámítása és CREDIT Wallet jóváírás (szimuláció)
print("\n5. Jutalékok kiszámítása (szimuláció)...")
levels = ["mlm_level1_percent", "mlm_level2_percent", "mlm_level3_percent"]
for idx, referrer in enumerate(chain):
if idx >= 3:
break
percent = params[levels[idx]]
commission = (payment_amount * Decimal(percent) / Decimal(100)).quantize(Decimal("0.01"))
print(f" L{idx+1} ({referrer.email}): {percent}% -> {commission} HUF")
# CREDIT Wallet keresése
wallet_stmt = select(Wallet).where(
Wallet.user_id == referrer.id
)
wallet = (await session.execute(wallet_stmt)).scalar_one_or_none()
if wallet:
print(f" Wallet megtalálva: {wallet.balance} {wallet.currency}")
# Szimulált jóváírás
new_balance = wallet.balance + commission
print(f" Új egyenleg: {new_balance} {wallet.currency}")
else:
print(f" NINCS Wallet a usernek!")
# 6. API végpont tesztelése (GET /me/network)
print("\n6. API végpont tesztelése (GET /me/network)...")
if l1_user:
# L1 hálózatának lekérdezése
network_stmt = select(User).where(User.referred_by_id == l1_user.id)
network_result = await session.execute(network_stmt)
network_users = network_result.scalars().all()
print(f" L1 ({l1_user.email}) közvetlen meghívottai: {len(network_users)}")
for u in network_users:
print(f" - {u.email} (kód: {u.referral_code})")
print("\n✅ MLM Payout teszt (szimuláció) sikeresen lefutott!")
print("=== TESZT VÉGE ===")
print("\nJAVASLAT: A teljes teszteléshez:")
print("1. Regisztrálj 4 teszt usert láncolatban (L1 -> L2 -> L3 -> L4)")
print("2. Futtasd a valós billing engine-t a payment success eseménnyel")
print("3. Ellenőrizd a CREDIT Wallet egyenlegeket a pgAdmin-ban")
if __name__ == "__main__":
asyncio.run(test_mlm_payout())

View File

@@ -0,0 +1,289 @@
# /opt/docker/dev/service_finder/backend/app/services/asset_matcher_service.py
"""
Internal Asset Matcher Service
Cél: Belső katalógus (vehicle_model_definitions) alapján automatikus eszköz-azonosítás és adatgazdagítás.
Matching stratégia:
1. Exact match: make + marketing_name + year_of_manufacture
2. Fuzzy match: make + normalizált név (Levenshtein távolság)
3. Confidence > 90% esetén automatikus adatgazdagítás
"""
from __future__ import annotations
import logging
import difflib
import uuid
from datetime import datetime
from typing import Optional, Tuple, List, Dict, Any
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, and_, or_, func
from sqlalchemy.orm import selectinload
from app.models import Asset, VehicleModelDefinition, AssetCatalog, AssetEvent, AssetTelemetry
from app.models.vehicle.vehicle_definitions import VehicleModelDefinition as VMD
logger = logging.getLogger(__name__)
class AssetMatcherService:
"""
Belső eszköz matcher szolgáltatás.
"""
@staticmethod
async def find_best_match(
db: AsyncSession,
asset: Asset,
threshold: float = 0.8
) -> Tuple[Optional[VehicleModelDefinition], float]:
"""
Megkeresi a legjobb egyezést az asset adatai alapján a vehicle_model_definitions táblában.
Args:
db: AsyncSession
asset: Asset objektum (már tartalmazza a make/model/year stb.)
threshold: Minimális confidence threshold (0-1)
Returns:
Tuple (matched_definition, confidence)
"""
# Gyűjtsük össze a keresési kritériumokat
make = asset.brand or (asset.catalog.make if asset.catalog else None)
model = asset.model or (asset.catalog.model if asset.catalog else None)
year = asset.year_of_manufacture
# Trim és ellenőrzés
if make:
make = make.strip()
if model:
model = model.strip()
if not make or not model:
logger.warning(f"Asset {asset.id} missing make or model, cannot match (make='{make}', model='{model}')")
return None, 0.0
# 1. EXACT MATCH: make + marketing_name + year_from
exact_match = await AssetMatcherService._exact_match(db, make, model, year)
if exact_match:
logger.info(f"Exact match found for asset {asset.id}: {make} {model} {year}")
return exact_match, 1.0
# 2. FUZZY MATCH: make + normalizált név (year within range)
fuzzy_matches = await AssetMatcherService._fuzzy_match(db, make, model, year, threshold)
if fuzzy_matches:
best_match, confidence = fuzzy_matches[0]
logger.info(f"Fuzzy match found for asset {asset.id}: {best_match.make} {best_match.marketing_name} (confidence: {confidence:.2f})")
return best_match, confidence
# 3. FALLBACK: csak make + model (year ignore)
fallback_match = await AssetMatcherService._fallback_match(db, make, model)
if fallback_match:
logger.info(f"Fallback match found for asset {asset.id}: {make} {model}")
return fallback_match, 0.7 # Alacsonyabb confidence
logger.warning(f"No match found for asset {asset.id}: {make} {model} {year}")
return None, 0.0
@staticmethod
async def _exact_match(
db: AsyncSession,
make: str,
model: str,
year: Optional[int]
) -> Optional[VehicleModelDefinition]:
"""
Pontos egyezés: make, marketing_name és year_from/year_to tartomány.
Több egyezés esetén a legújabb évjáratút választja.
"""
stmt = select(VMD).where(
VMD.make.ilike(make),
VMD.marketing_name.ilike(model)
)
if year:
# Évjárat tartományban legyen
stmt = stmt.where(
and_(
VMD.year_from <= year,
or_(VMD.year_to.is_(None), VMD.year_to >= year)
)
)
# Rendezés év szerint csökkenő, limit 1
stmt = stmt.order_by(VMD.year_from.desc()).limit(1)
result = await db.execute(stmt)
return result.scalar_one_or_none()
@staticmethod
async def _fuzzy_match(
db: AsyncSession,
make: str,
model: str,
year: Optional[int],
threshold: float
) -> List[Tuple[VehicleModelDefinition, float]]:
"""
Fuzzy egyezés: hasonlóság a normalizált név alapján.
"""
# Először szűrjünk make és év alapján
stmt = select(VMD).where(VMD.make.ilike(make))
if year:
stmt = stmt.where(
and_(
VMD.year_from <= year,
or_(VMD.year_to.is_(None), VMD.year_to >= year)
)
)
result = await db.execute(stmt)
candidates = result.scalars().all()
if not candidates:
return []
# Számítsuk ki a hasonlóságot a model név és a marketing_name között
matches = []
for candidate in candidates:
similarity = AssetMatcherService._calculate_similarity(model, candidate.marketing_name)
if similarity >= threshold:
matches.append((candidate, similarity))
# Rendezzük confidence szerint csökkenő sorrendben
matches.sort(key=lambda x: x[1], reverse=True)
return matches
@staticmethod
async def _fallback_match(
db: AsyncSession,
make: str,
model: str
) -> Optional[VehicleModelDefinition]:
"""
Csak make + model alapján, évjárat figyelmen kívül hagyva.
"""
stmt = select(VMD).where(
VMD.make.ilike(make),
VMD.marketing_name.ilike(model)
).order_by(VMD.year_from.desc()).limit(1)
result = await db.execute(stmt)
return result.scalar_one_or_none()
@staticmethod
def _calculate_similarity(str1: str, str2: str) -> float:
"""
Szöveg hasonlóság számítása SequenceMatcher segítségével.
"""
if not str1 or not str2:
return 0.0
return difflib.SequenceMatcher(None, str1.lower(), str2.lower()).ratio()
@staticmethod
async def enrich_asset_from_definition(
db: AsyncSession,
asset: Asset,
definition: VehicleModelDefinition,
confidence: float
) -> Asset:
"""
Gazdagítsa az asset adatait a definition technikai specifikációival.
Csak akkor, ha az asset megfelelő mezői üresek.
"""
# Technikai specifikációk másolása
if not asset.power_kw and definition.power_kw:
asset.power_kw = definition.power_kw
if not asset.torque_nm and definition.torque_nm:
asset.torque_nm = definition.torque_nm
if not asset.engine_capacity and definition.engine_capacity:
asset.engine_capacity = definition.engine_capacity
if not asset.transmission_type and definition.transmission_type:
asset.transmission_type = definition.transmission_type
if not asset.drive_type and definition.drive_type:
asset.drive_type = definition.drive_type
if not asset.fuel_type and definition.fuel_type:
asset.fuel_type = definition.fuel_type
if not asset.euro_classification and definition.euro_classification:
asset.euro_classification = definition.euro_classification
if not asset.vehicle_class and definition.vehicle_class:
asset.vehicle_class = definition.vehicle_class
if not asset.trim_level and definition.body_type:
asset.trim_level = definition.body_type # body_type -> trim_level mapping
# Évjárat ellenőrzés
if not asset.year_of_manufacture and definition.year_from:
asset.year_of_manufacture = definition.year_from
# Státusz frissítése
if confidence >= 0.9:
asset.data_status = 'verified'
logger.info(f"Asset {asset.id} enriched and marked as verified (confidence: {confidence:.2f})")
else:
asset.data_status = 'enriched'
logger.info(f"Asset {asset.id} enriched but not verified (confidence: {confidence:.2f})")
return asset
@staticmethod
async def match_and_enrich_asset(
db: AsyncSession,
asset_id: uuid.UUID,
threshold: float = 0.9
) -> Dict[str, Any]:
"""
Fő függvény: Asset ID alapján keres match-et és gazdagítja az adatokat.
Args:
db: AsyncSession
asset_id: Asset UUID
threshold: Confidence threshold a verification-hoz (alapértelmezett 90%)
Returns:
Dict with match results
"""
# Asset betöltése
stmt = select(Asset).where(Asset.id == asset_id).options(selectinload(Asset.catalog))
result = await db.execute(stmt)
asset = result.scalar_one_or_none()
if not asset:
raise ValueError(f"Asset {asset_id} not found")
logger.info(f"Matching asset {asset_id} ({asset.brand} {asset.model})")
# Match keresés
definition, confidence = await AssetMatcherService.find_best_match(db, asset, threshold=0.8)
if not definition:
return {
"asset_id": str(asset_id),
"matched": False,
"confidence": 0.0,
"message": "No matching definition found in internal catalog"
}
# Adatgazdagítás
enriched_asset = await AssetMatcherService.enrich_asset_from_definition(
db, asset, definition, confidence
)
# Mentés
await db.commit()
return {
"asset_id": str(asset_id),
"matched": True,
"confidence": confidence,
"definition_id": definition.id,
"definition": f"{definition.make} {definition.marketing_name}",
"data_status": enriched_asset.data_status,
"enriched_fields": [
field for field in [
"power_kw" if asset.power_kw != enriched_asset.power_kw else None,
"torque_nm" if asset.torque_nm != enriched_asset.torque_nm else None,
"engine_capacity" if asset.engine_capacity != enriched_asset.engine_capacity else None,
"transmission_type" if asset.transmission_type != enriched_asset.transmission_type else None,
"drive_type" if asset.drive_type != enriched_asset.drive_type else None,
"fuel_type" if asset.fuel_type != enriched_asset.fuel_type else None,
] if field is not None
]
}
# Singleton instance
asset_matcher_service = AssetMatcherService()

View File

@@ -132,14 +132,28 @@ class AssetService:
# Get default vehicle class from config if not provided
default_vehicle_class = await config.get_setting(db, "DEFAULT_VEHICLE_CLASS", default="car")
# --- BRANCH (GARÁZS) HOZZÁRENDELÉSI LOGIKA ---
branch_id = getattr(asset_data, 'branch_id', None)
if not branch_id and target_org_id:
from app.models.marketplace.organization import Branch
branch_stmt = select(Branch.id).where(Branch.organization_id == target_org_id, Branch.is_main == True)
main_branch_id = (await db.execute(branch_stmt)).scalar()
if not main_branch_id:
branch_stmt_any = select(Branch.id).where(Branch.organization_id == target_org_id).limit(1)
main_branch_id = (await db.execute(branch_stmt_any)).scalar()
branch_id = main_branch_id
print(f"DEBUG: resolved branch_id={branch_id} for target_org_id={target_org_id}")
asset_fields = {
'vin': vin_clean,
'license_plate': license_plate_clean,
'catalog_id': asset_data.catalog_id,
'current_organization_id': target_org_id,
'branch_id': branch_id,
'owner_person_id': user.person_id,
'owner_org_id': asset_data.owner_org_id or target_org_id,
'operator_org_id': asset_data.operator_org_id,
'owner_org_id': asset_data.organization_id or target_org_id,
'operator_org_id': None, # AssetCreate doesn't have operator_org_id
'status': status,
'individual_equipment': asset_data.individual_equipment or {},
'created_at': datetime.utcnow(),
@@ -184,8 +198,23 @@ class AssetService:
db.add(new_asset)
await db.flush()
# --- MATCHER INTEGRÁCIÓ (HA NINCS CATALOG_ID DE VANNAK ADATOK) ---
if not new_asset.catalog_id and new_asset.brand and new_asset.model:
from app.services.asset_matcher_service import AssetMatcherService
print(f"DEBUG: Matcher args: make={new_asset.brand}, model={new_asset.model}, year={new_asset.year_of_manufacture}")
matched_result = await AssetMatcherService.find_best_match(db, new_asset)
if matched_result:
matched_def, conf = matched_result
if matched_def:
new_asset.catalog_id = matched_def.id
await AssetMatcherService.enrich_asset_from_definition(db, new_asset, matched_def, conf)
await db.flush()
# Digitális Iker Alapmodulok
db.add(AssetAssignment(asset_id=new_asset.id, organization_id=target_org_id, status="active"))
# Only create AssetAssignment if we have an organization_id
if target_org_id is not None:
db.add(AssetAssignment(asset_id=new_asset.id, organization_id=target_org_id, status="active"))
db.add(AssetTelemetry(asset_id=new_asset.id))
db.add(AssetFinancials(
asset_id=new_asset.id,
@@ -202,6 +231,9 @@ class AssetService:
await AssetService._award_first_car_badge(db, user_id, target_org_id)
await db.commit()
from sqlalchemy.orm import selectinload
new_asset = (await db.execute(select(Asset).where(Asset.id == new_asset.id).options(selectinload(Asset.catalog)))).scalar_one()
return new_asset
except Exception as e:
@@ -269,16 +301,20 @@ class AssetService:
else:
logger.warning("execute_final_transfer called without user_id, ownership fields not updated")
db.add(AssetAssignment(asset_id=asset.id, organization_id=new_org_id, status="active"))
# Only create AssetAssignment if we have an organization_id
if new_org_id is not None:
db.add(AssetAssignment(asset_id=asset.id, organization_id=new_org_id, status="active"))
await db.commit()
return asset
# --- CATALOG METHODS ---
@staticmethod
async def get_makes(db: AsyncSession) -> List[str]:
"""Get all distinct makes from vehicle model definitions."""
async def get_makes(db: AsyncSession, vehicle_class: Optional[str] = None) -> List[str]:
"""Get all distinct makes from vehicle model definitions, optionally filtered by vehicle_class."""
stmt = select(distinct(VehicleModelDefinition.make)).order_by(VehicleModelDefinition.make)
if vehicle_class:
stmt = stmt.where(VehicleModelDefinition.vehicle_class == vehicle_class)
result = await db.execute(stmt)
makes = result.scalars().all()
return [make for make in makes if make] # Filter out None/empty

View File

@@ -1,6 +1,7 @@
# /opt/docker/dev/service_finder/backend/app/services/auth_service.py
import logging
import uuid
import re
from datetime import datetime, timedelta, timezone
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, and_, update
@@ -23,6 +24,41 @@ from app.services.gamification_service import GamificationService
logger = logging.getLogger(__name__)
class AuthService:
@staticmethod
async def _validate_password_complexity(db: AsyncSession, password: str, region_code: str = None):
"""
Dinamikus jelszó komplexitás ellenőrzése az admin beállítások alapján.
"""
# Alap beállítások lekérése
min_pass = await config.get_setting(db, "auth_min_password_length", default=8)
password_strict = await config.get_setting(db, "auth_password_strict", default=False)
# Minimum hossz ellenőrzése
if len(password) < int(min_pass):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"A jelszónak legalább {min_pass} karakter hosszúnak kell lennie."
)
# Ha a strict mód be van kapcsolva, komplexitás ellenőrzés
if password_strict and str(password_strict).lower() in ("true", "1", "yes"):
# Ellenőrizzük: legalább 1 nagybetű, 1 kisbetű, 1 szám vagy speciális karakter
if not re.search(r'[A-Z]', password):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="A jelszónak tartalmaznia kell legalább egy nagybetűt."
)
if not re.search(r'[a-z]', password):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="A jelszónak tartalmaznia kell legalább egy kisbetűt."
)
if not re.search(r'[0-9!@#$%^&*()_+\-=\[\]{};:"\\|,.<>/?]', password):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="A jelszónak tartalmaznia kell legalább egy számot vagy speciális karaktert."
)
@staticmethod
async def register_lite(db: AsyncSession, user_in: UserLiteRegister):
""" 1. FÁZIS: Lite regisztráció dinamikus korlátokkal és Sentinel naplózással. """
@@ -32,11 +68,8 @@ class AuthService:
default_role_name = await config.get_setting(db, "auth_default_role", default="user")
reg_token_hours = await config.get_setting(db, "auth_registration_hours", region_code=user_in.region_code, default=48)
if len(user_in.password) < int(min_pass):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=f"A jelszónak legalább {min_pass} karakter hosszúnak kell lennie."
)
# Jelszó komplexitás ellenőrzése
await AuthService._validate_password_complexity(db, user_in.password, user_in.region_code)
# Check if email already exists
existing_user = await db.execute(select(User).where(User.email == user_in.email))
@@ -46,6 +79,17 @@ class AuthService:
detail="Ez az email cím már regisztrálva van."
)
# Meghívó keresése referral kód alapján
referred_by_id = None
if user_in.referred_by_code:
referrer_stmt = select(User).where(User.referral_code == user_in.referred_by_code)
referrer = (await db.execute(referrer_stmt)).scalar_one_or_none()
if referrer:
referred_by_id = referrer.id
logger.info(f"User {user_in.email} referred by {referrer.email} (ID: {referrer.id})")
else:
logger.warning(f"Referral code '{user_in.referred_by_code}' not found, ignoring.")
new_person = Person(
first_name=user_in.first_name,
last_name=user_in.last_name,
@@ -65,6 +109,9 @@ class AuthService:
# Szerepkör dinamikus feloldása
assigned_role = UserRole[default_role_name] if default_role_name in UserRole.__members__ else UserRole.user
# Referral kód generálása
referral_code = generate_secure_slug(8).upper()
new_user = User(
email=user_in.email,
hashed_password=get_password_hash(user_in.password),
@@ -80,6 +127,8 @@ class AuthService:
preferred_currency="HUF",
scope_level="individual",
custom_permissions={},
referral_code=referral_code,
referred_by_id=referred_by_id,
created_at=datetime.now(timezone.utc)
)
db.add(new_user)
@@ -94,21 +143,6 @@ class AuthService:
expires_at=datetime.now(timezone.utc) + timedelta(hours=int(reg_token_hours))
))
# Email küldés a beállított template alapján
verification_link = f"{settings.FRONTEND_BASE_URL}/verify?token={token_val}"
email_result = await email_manager.send_email(
recipient=user_in.email,
template_key="reg",
variables={"first_name": user_in.first_name, "link": verification_link},
lang=user_in.lang
)
# Check if email sending failed
if email_result and email_result.get("status") == "error":
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail="Email delivery failed. Please contact support."
)
# Sentinel Audit Log
await security_service.log_event(
db, user_id=new_user.id, action="USER_REGISTER_LITE",
@@ -117,6 +151,23 @@ class AuthService:
)
await db.commit()
# Email küldés a beállított template alapján
verification_link = f"{settings.FRONTEND_BASE_URL}/verify?token={token_val}"
email_result = await email_manager.send_email(
recipient=user_in.email,
template_key="reg",
variables={"first_name": user_in.first_name, "link": verification_link},
lang=user_in.lang
)
# Check if email sending failed - LOG but don't crash registration
if email_result and email_result.get("status") == "error":
logger.error(f"Email sending failed for {user_in.email}: {email_result.get('message')}")
# Don't raise exception - registration should succeed even if email fails
# Just log the error and continue
else:
logger.info(f"Email sent successfully to {user_in.email}")
return new_user
except Exception as e:
await db.rollback()
@@ -143,16 +194,53 @@ class AuthService:
house_number=kyc_in.address_house_number, parcel_id=kyc_in.address_hrsz
)
# Person adatok dúsítása
p = user.person
p.mothers_last_name = kyc_in.mothers_last_name
p.mothers_first_name = kyc_in.mothers_first_name
p.birth_place = kyc_in.birth_place
p.birth_date = kyc_in.birth_date
p.phone = kyc_in.phone_number
p.address_id = addr_id
p.identity_docs = jsonable_encoder(kyc_in.identity_docs)
p.is_active = True
# --- SHADOW IDENTITY CHECK ---
if kyc_in.birth_date:
shadow_stmt = select(Person).where(
and_(
Person.last_name == p.last_name,
Person.first_name == p.first_name,
Person.birth_date == kyc_in.birth_date,
Person.id != p.id # Ne a jelenlegit találja meg
)
)
shadow_person = (await db.execute(shadow_stmt)).scalar_one_or_none()
else:
shadow_person = None
if shadow_person:
logger.info(f"Shadow Identity megtalálva a {user.id} userhez: Person ID {shadow_person.id}")
user.person_id = shadow_person.id
# Frissítjük a megtalált Person rekordot a KYC adatokkal
shadow_person.mothers_last_name = kyc_in.mothers_last_name
shadow_person.mothers_first_name = kyc_in.mothers_first_name
shadow_person.birth_place = kyc_in.birth_place
shadow_person.phone = kyc_in.phone_number
shadow_person.address_id = addr_id
shadow_person.identity_docs = jsonable_encoder(kyc_in.identity_docs)
shadow_person.is_active = True
# Inaktiváljuk/töröljük a megárvult eredeti Person-t
p.is_active = False
p.is_ghost = True
p = shadow_person
else:
p.mothers_last_name = kyc_in.mothers_last_name
p.mothers_first_name = kyc_in.mothers_first_name
p.birth_place = kyc_in.birth_place
p.birth_date = kyc_in.birth_date
p.phone = kyc_in.phone_number
p.address_id = addr_id
p.identity_docs = jsonable_encoder(kyc_in.identity_docs)
p.is_active = True
# --- SHADOW CHECK VÉGE ---
# Dinamikus szervezet generálás
org_full_name = org_tpl.format(last_name=p.last_name, first_name=p.first_name)
@@ -183,16 +271,30 @@ class AuthService:
db.add(Branch(organization_id=new_org.id, address_id=addr_id, name="Home Base", is_main=True))
db.add(OrganizationMember(organization_id=new_org.id, user_id=user.id, role="OWNER"))
db.add(Wallet(user_id=user.id, currency=kyc_in.preferred_currency or base_cur))
db.add(UserStats(user_id=user.id))
# db.add(UserStats(user_id=user.id)) # GamificationService kezeli
user.is_active = True
user.folder_slug = generate_secure_slug(12)
# Set user's scope_id to the new personal organization ID
user.scope_id = str(new_org.id)
# Update region_code from KYC form (EU country selection)
if kyc_in.region_code:
user.region_code = kyc_in.region_code
# Gamification XP jóváírás
await GamificationService.award_points(db, user_id=user.id, amount=int(kyc_reward), reason="KYC_VERIFICATION")
# P2P Referral XP a meghívónak (ha van referred_by_id)
if user.referred_by_id:
p2p_xp = await config.get_setting(db, "gamification_p2p_invite_xp", default=50)
await GamificationService.award_points(
db,
user_id=user.referred_by_id,
amount=int(p2p_xp),
reason="P2P_REFERRAL_SUCCESS"
)
logger.info(f"P2P XP ({p2p_xp}) awarded to referrer ID {user.referred_by_id} for user {user.id}")
await db.commit()
return user
except Exception as e:
@@ -224,7 +326,19 @@ class AuthService:
if not token: return False
token.is_used = True
# Itt aktiválhatnánk a júzert, ha a Lite regnél még nem tennénk meg
# Activate user
user_stmt = select(User).where(User.id == token.user_id)
user = (await db.execute(user_stmt)).scalar_one_or_none()
if user:
user.is_active = True
# Activate person
person_stmt = select(Person).where(Person.user_id == user.id)
person = (await db.execute(person_stmt)).scalar_one_or_none()
if person:
person.is_active = True
await db.commit()
return True
except: return False
@@ -268,13 +382,18 @@ class AuthService:
token_rec = (await db.execute(stmt)).scalar_one_or_none()
if not token_rec: return False
# Jelszó komplexitás ellenőrzése
user_stmt = select(User).where(User.id == token_rec.user_id)
user = (await db.execute(user_stmt)).scalar_one()
await AuthService._validate_password_complexity(db, new_password, user.region_code)
user.hashed_password = get_password_hash(new_password)
token_rec.is_used = True
await db.commit()
return True
except HTTPException:
raise
except: return False
@staticmethod

View File

@@ -1,6 +1,8 @@
# /opt/docker/dev/service_finder/backend/app/services/email_manager.py
import os
import smtplib
import logging
import requests
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from typing import Optional
@@ -14,6 +16,24 @@ from app.db.session import AsyncSessionLocal
logger = logging.getLogger("Email-Manager-2.0")
class EmailManager:
@staticmethod
def _get_base_url() -> str:
"""Return the appropriate base URL for email links."""
# Check environment variable first
base_url = os.getenv("EMAIL_BASE_URL")
if base_url:
return base_url.rstrip('/')
# Fallback to dev domain for external access
return "https://dev.servicefinder.hu"
@staticmethod
def _build_verification_link(token: str) -> str:
"""Build verification link with proper path."""
base = EmailManager._get_base_url()
# Use frontend verification route (adjust if needed)
return f"{base}/verify?token={token}"
@staticmethod
def _get_html_template(template_key: str, variables: dict, lang: str = "hu") -> str:
"""HTML sablon generálása a fordítási fájlok alapján."""
@@ -24,6 +44,11 @@ class EmailManager:
link_fallback_text = locale_manager.get("email.link_fallback", lang=lang)
# If link is not provided but token is, build verification link
link = variables.get('link')
if not link and 'token' in variables:
link = EmailManager._build_verification_link(variables['token'])
return f"""
<html>
<body style="font-family: Arial, sans-serif; color: #333; line-height: 1.6;">
@@ -31,14 +56,14 @@ class EmailManager:
<h2 style="color: #2c3e50;">{greeting}</h2>
<p>{body}</p>
<div style="text-align: center; margin: 40px 0;">
<a href="{variables.get('link', '#')}"
<a href="{link or '#'}"
style="background-color: #3498db; color: white; padding: 15px 30px; text-decoration: none; border-radius: 5px; font-weight: bold; font-size: 16px;">
{button_text}
</a>
</div>
<p style="font-size: 0.85em; color: #777; word-break: break-all;">
{link_fallback_text}<br>
<a href="{variables.get('link')}" style="color: #3498db;">{variables.get('link')}</a>
<a href="{link or '#'}" style="color: #3498db;">{link or '#'}</a>
</p>
<hr style="border: 0; border-top: 1px solid #eee; margin: 30px 0;">
<p style="font-size: 0.8em; color: #999; text-align: center;">{footer}</p>
@@ -50,7 +75,7 @@ class EmailManager:
@staticmethod
async def send_email(recipient: str, template_key: str, variables: dict, lang: str = "hu", db: Optional[AsyncSession] = None):
"""
E-mail küldése közvetlenül a privát SMTP szerveren keresztül.
E-mail küldése Brevo API-n vagy SMTP-n keresztül.
"""
session_internal = False
if db is None:
@@ -62,35 +87,130 @@ class EmailManager:
provider = await config.get_setting(db, "email_provider", default="smtp")
if provider == "disabled":
logger.info(f"Email küldés letiltva (Admin config). Cél: {recipient}")
return
return {"status": "success", "provider": "disabled", "message": "Email disabled by admin config"}
html = EmailManager._get_html_template(template_key, variables, lang)
subject = locale_manager.get(f"email.{template_key}_subject", lang=lang)
smtp_host = os.getenv("SMTP_HOST", "mail.servicefinder.hu")
smtp_port = int(os.getenv("SMTP_PORT", "465"))
smtp_user = os.getenv("SMTP_USER", "noreply@servicefinder.hu")
smtp_pass = os.getenv("SMTP_PASSWORD", "")
# Get email provider from environment
email_provider = os.getenv("EMAIL_PROVIDER", "smtp").lower()
from_email = os.getenv("MAIL_FROM", "noreply@servicefinder.hu")
from_name = os.getenv("MAIL_FROM_NAME", "ServiceFinder")
smtp_cfg = {
"host": smtp_host,
"port": smtp_port,
"user": smtp_user,
"pass": smtp_pass
}
logger.info(f"Using SMTP config: host={smtp_cfg['host']}, port={smtp_cfg['port']}, user={smtp_cfg['user']}")
return await EmailManager._send_via_smtp(smtp_cfg, from_email, from_name, recipient, subject, html)
if email_provider == "brevo_api":
result = await EmailManager._send_via_brevo_api(recipient, subject, html, variables)
# Primary-Fallback Logic: If Brevo API fails, try SMTP
if result.get("status") == "error":
logger.error(f"Brevo API failed for {recipient}: {result.get('message')}. Falling back to SMTP.")
result = await EmailManager._send_via_smtp(recipient, subject, html)
elif email_provider == "brevo_smtp":
result = await EmailManager._send_via_brevo_smtp(recipient, subject, html)
else: # Default to SMTP
result = await EmailManager._send_via_smtp(recipient, subject, html)
logger.info(f"Email sending result for {recipient}: {result}")
return result
except Exception as e:
logger.error(f"CRITICAL: Email sending failed for {recipient}: {str(e)}")
# Don't crash - return error but don't raise exception
return {"status": "error", "message": str(e), "provider": "unknown"}
finally:
if session_internal:
await db.close()
@staticmethod
async def _send_via_smtp(cfg: dict, from_email: str, from_name: str, recipient: str, subject: str, html: str):
async def _send_via_brevo_api(recipient: str, subject: str, html: str, variables: dict):
"""Send email via Brevo REST API"""
try:
brevo_api_key = os.getenv("BREVO_API_KEY")
if not brevo_api_key:
logger.error("BREVO_API_KEY environment variable not set")
return {"status": "error", "message": "BREVO_API_KEY not configured", "provider": "brevo_api"}
sender_name = os.getenv("MAIL_FROM_NAME", "ServiceFinder")
sender_email = os.getenv("MAIL_FROM", "noreply@servicefinder.hu")
# Prepare Brevo API payload
payload = {
"sender": {
"name": sender_name,
"email": sender_email
},
"to": [{"email": recipient}],
"subject": subject,
"htmlContent": html,
"tags": ["verification"]
}
# Try to disable tracking - Brevo may ignore this if account-level tracking is enabled
# According to Brevo API v3 docs, tracking settings are at account level
# We'll add tracking object but it may not work for free tier
payload["tracking"] = {
"click": False,
"open": False
}
headers = {
"accept": "application/json",
"api-key": brevo_api_key,
"content-type": "application/json"
}
response = requests.post(
"https://api.brevo.com/v3/smtp/email",
json=payload,
headers=headers,
timeout=30
)
if response.status_code == 201:
logger.info(f"Brevo API email sent successfully to {recipient}")
return {"status": "success", "provider": "brevo_api", "message_id": response.json().get("messageId")}
else:
logger.error(f"Brevo API error: {response.status_code} - {response.text}")
return {"status": "error", "provider": "brevo_api", "message": response.text}
except Exception as e:
logger.error(f"Brevo API exception: {str(e)}")
return {"status": "error", "provider": "brevo_api", "message": str(e)}
@staticmethod
async def _send_via_brevo_smtp(recipient: str, subject: str, html: str):
"""Send email via Brevo SMTP relay"""
try:
smtp_host = "smtp-relay.brevo.com"
smtp_port = 587
smtp_user = os.getenv("BREVO_SMTP_USER")
smtp_pass = os.getenv("BREVO_SMTP_PASSWORD")
if not smtp_user or not smtp_pass:
logger.error("BREVO_SMTP_USER or BREVO_SMTP_PASSWORD not set")
return {"status": "error", "message": "Brevo SMTP credentials not configured", "provider": "brevo_smtp"}
from_email = os.getenv("MAIL_FROM", "noreply@servicefinder.hu")
from_name = os.getenv("MAIL_FROM_NAME", "ServiceFinder")
msg = MIMEMultipart()
msg["From"] = f"{from_name} <{from_email}>"
msg["To"] = recipient
msg["Subject"] = subject
msg.attach(MIMEText(html, "html"))
logger.info(f"Connecting to Brevo SMTP: {smtp_host}:{smtp_port}")
with smtplib.SMTP(smtp_host, smtp_port, timeout=15) as server:
server.starttls()
server.login(smtp_user, smtp_pass)
server.send_message(msg)
logger.info(f"Brevo SMTP email sent successfully to {recipient}")
return {"status": "success", "provider": "brevo_smtp"}
except Exception as e:
logger.error(f"Brevo SMTP exception: {str(e)}")
return {"status": "error", "provider": "brevo_smtp", "message": str(e)}
@staticmethod
async def _send_via_smtp(recipient: str, subject: str, html: str):
"""Send email via standard SMTP (fallback)"""
# Mock mode check: If APP_ENV=test or domain is example.com, skip SMTP and return success
app_env = os.getenv("APP_ENV", "").lower()
is_example_domain = recipient.endswith("@example.com") or "@example.com" in recipient
@@ -99,36 +219,40 @@ class EmailManager:
return {"status": "success", "provider": "mock", "message": "Email skipped in test mode"}
try:
smtp_host = os.getenv("SMTP_HOST", "mail.servicefinder.hu")
smtp_port = int(os.getenv("SMTP_PORT", "465"))
smtp_user = os.getenv("SMTP_USER", "noreply@servicefinder.hu")
smtp_pass = os.getenv("SMTP_PASSWORD", "")
from_email = os.getenv("MAIL_FROM", "noreply@servicefinder.hu")
from_name = os.getenv("MAIL_FROM_NAME", "ServiceFinder")
msg = MIMEMultipart()
msg["From"] = f"{from_name} <{from_email}>"
msg["To"] = recipient
msg["Subject"] = subject
msg.attach(MIMEText(html, "html"))
# Port 465 uses SMTP_SSL directly instead of STARTTLS
if cfg["port"] == 465:
logger.info(f"Connecting via SMTP_SSL to {cfg['host']}:{cfg['port']}")
with smtplib.SMTP_SSL(cfg["host"], cfg["port"], timeout=15) as server:
user = cfg.get("user", "")
passwd = cfg.get("pass", "")
if user and passwd:
server.login(user, passwd)
if smtp_port == 465:
logger.info(f"Connecting via SMTP_SSL to {smtp_host}:{smtp_port}")
with smtplib.SMTP_SSL(smtp_host, smtp_port, timeout=15) as server:
if smtp_user and smtp_pass:
server.login(smtp_user, smtp_pass)
server.send_message(msg)
else:
logger.info(f"Connecting via SMTP to {cfg['host']}:{cfg['port']}")
with smtplib.SMTP(cfg["host"], cfg["port"], timeout=15) as server:
# Explicit STARTTLS if not 465, though we expect 465
logger.info(f"Connecting via SMTP to {smtp_host}:{smtp_port}")
with smtplib.SMTP(smtp_host, smtp_port, timeout=15) as server:
server.starttls()
user = cfg.get("user", "")
passwd = cfg.get("pass", "")
if user and passwd:
server.login(user, passwd)
if smtp_user and smtp_pass:
server.login(smtp_user, smtp_pass)
server.send_message(msg)
logger.info(f"SMTP siker -> {recipient}")
logger.info(f"SMTP email sent successfully to {recipient}")
return {"status": "success", "provider": "smtp"}
except Exception as e:
logger.error(f"SMTP hiba: {str(e)}")
return {"status": "error", "message": str(e)}
logger.error(f"SMTP exception: {str(e)}")
return {"status": "error", "provider": "smtp", "message": str(e)}
email_manager = EmailManager()
email_manager = EmailManager()

View File

@@ -1,3 +1,4 @@
# /opt/docker/dev/service_finder/backend/app/services/financial_orchestrator.py
"""
Financial Orchestrator - Unit of Work mintával a pénzügyi tranzakciók atomi kezeléséhez.

View File

@@ -35,11 +35,21 @@ class TranslationService:
logger.info(f"🌍 i18n Motor: {len(translations)} szöveg aktiválva a memóriában.")
@classmethod
def get_text(cls, key: str, lang: str = "hu", variables: Optional[Dict[str, Any]] = None) -> str:
def get_text(cls, key: str, lang: Optional[str] = None, variables: Optional[Dict[str, Any]] = None) -> str:
"""
Szerveroldali lekérés Fallback (EN) logikával és változó behelyettesítéssel.
Automatikusan használja a request context locale-ját, ha nincs explicit nyelv megadva.
Példa: get_text("AUTH.WELCOME", "hu", {"name": "Péter"})
"""
# Use context locale if no explicit language is provided
if lang is None:
try:
from app.core.context import get_current_locale
lang = get_current_locale()
except (ImportError, Exception):
# Fallback to default if context is not available
lang = "hu"
# 1. Kért nyelv lekérése
text = cls._published_cache.get(lang, {}).get(key)