# /opt/docker/dev/service_finder/backend/app/api/v1/endpoints/organizations.py import os import re import uuid import hashlib import random import logging from typing import List from datetime import datetime, timedelta, timezone from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy.ext.asyncio import AsyncSession 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, OrganizationUpdate, OrganizationResponse from app.models.marketplace.organization import Organization, OrgType, OrganizationMember, Branch, OrgUserRole from app.models.identity import User, OneTimePassword # JAVÍTVA: Központi Identity modell from app.core.config import settings from app.services.security_service import security_service from app.models import LogSeverity router = APIRouter() logger = logging.getLogger(__name__) @router.post("/onboard", response_model=CorpOnboardResponse, status_code=status.HTTP_201_CREATED) async def onboard_organization( org_in: CorpOnboardIn, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """ Új szervezet (cég/szerviz) rögzítése. Automatikusan generál slug-ot és létrehozza a NAS mappa-struktúrát. """ # 1. Magyar adószám validáció (XXXXXXXX-Y-ZZ) if org_in.country_code == "HU": if not re.match(r"^\d{8}-\d-\d{2}$", org_in.tax_number): raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail="Érvénytelen magyar adószám formátum!" ) # 2. Duplikáció ellenőrzés – adószám alapján stmt_exist = select(Organization).where(Organization.tax_number == org_in.tax_number) result_exist = await db.execute(stmt_exist) if result_exist.scalar_one_or_none(): raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="Ez a cég (adószám) már regisztrálva van a rendszerben." ) # 3. KÖTELEZŐ MEZŐ: folder_slug generálása # Mivel az adatbázisban NOT NULL, itt muszáj létrehozni temp_slug = hashlib.md5(f"{org_in.tax_number}-{uuid.uuid4()}".encode()).hexdigest()[:12] # 4. Mentés new_org = Organization( full_name=org_in.full_name, name=org_in.name, display_name=org_in.display_name, tax_number=org_in.tax_number, reg_number=org_in.reg_number, folder_slug=temp_slug, # JAVÍTVA: Kötelező mező beillesztve address_zip=org_in.address_zip, address_city=org_in.address_city, address_street_name=org_in.address_street_name, address_street_type=org_in.address_street_type, address_house_number=org_in.address_house_number, address_hrsz=org_in.address_hrsz, country_code=org_in.country_code, org_type=OrgType.business, status="pending_verification", # --- EXPLICIT IDŐBÉLYEGEK A DB HIBA ELKERÜLÉSÉRE --- first_registered_at=datetime.now(timezone.utc), current_lifecycle_started_at=datetime.now(timezone.utc), created_at=datetime.now(timezone.utc), subscription_plan="FREE", base_asset_limit=1, purchased_extra_slots=0, notification_settings={}, external_integration_config={}, is_ownership_transferable=True ) db.add(new_org) await db.flush() # 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, role="OWNER" # JAVÍTVA: Enum kompatibilis nagybetűs forma ) db.add(owner_member) # 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)) os.makedirs(os.path.join(org_path, "documents"), exist_ok=True) logger.info(f"NAS mappa kész: {org_path}") except Exception as e: logger.error(f"NAS hiba: {e}") await db.commit() await db.refresh(new_org) return {"organization_id": new_org.id, "status": new_org.status} @router.get("/lookup-tax/{tax_number}") async def lookup_tax_number( tax_number: str, db: AsyncSession = Depends(get_db), ): """ Adószám alapú cégadat lekérdezés a hivatalos NAV Online Számla API v3 rendszerén keresztül. A queryTaxpayer végpontot használja valós idejű cégadatok lekérésére. Duplikáció ellenőrzés: Mielőtt a NAV-hoz fordulnánk, ellenőrizzük, hogy az adószám első 8 számjegye (törzsszám) alapján létezik-e már aktív cég az adatbázisban. Hibakezelés: - 409: Az adószám már regisztrálva van a rendszerben - 503: NAV rendszer nem elérhető (hálózati hiba, DNS, timeout) - 404: Adószám nem található a NAV adatbázisában """ from app.services.nav_service import NavService, NavServiceUnavailableError # Duplikáció ellenőrzés: adószám első 8 számjegye (törzsszám) alapján tax_core = tax_number.strip().replace("-", "").replace(" ", "")[:8] stmt_existing = select(Organization).where( Organization.tax_number.like(f"{tax_core}%"), Organization.is_deleted == False, ) result_existing = await db.execute(stmt_existing) existing_org = result_existing.scalar_one_or_none() if existing_org: logger.warning(f"Duplikáció észlelve: adószám {tax_core} már létezik (org_id={existing_org.id})") raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="Ez az adószám már regisztrálva van a rendszerben." ) try: result = await NavService.query_taxpayer(tax_number) except NavServiceUnavailableError as e: logger.error(f"NAV szolgáltatás hiba: {e}") raise HTTPException( status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail="A NAV rendszere jelenleg nem elérhető." ) if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Cég nem található a NAV adatbázisában" ) return result @router.get("/my", response_model=List[dict]) async def get_my_organizations( db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """ A bejelentkezett felhasználóhoz tartozó összes szervezet listázása. """ stmt = ( select(Organization) .join(OrganizationMember) .where(OrganizationMember.user_id == current_user.id) ) result = await db.execute(stmt) orgs = result.scalars().all() # Return full organization details return [ { "organization_id": o.id, "status": o.status, "name": o.name, "full_name": o.full_name, "display_name": o.display_name, "tax_number": o.tax_number, "country_code": o.country_code, "is_active": o.is_active, "is_deleted": o.is_deleted, "subscription_plan": o.subscription_plan, "org_type": o.org_type.value if hasattr(o.org_type, 'value') else str(o.org_type), "visual_settings": o.visual_settings if hasattr(o, 'visual_settings') else None } for o in orgs ] @router.patch("/{org_id}/visual-settings", response_model=OrganizationResponse) async def update_organization_visual_settings( org_id: int, update_data: OrganizationUpdate, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """ Szervezet vizuális beállításainak részleges frissítése (JSONB merge). A frontend küldhet csak egyetlen kulcsot is (pl. {"theme": "dark"}), és az nem írja felül null-lal a többi mezőt. A meglévő visual_settings objektum mélyen merge-elődik a beérkező adatokkal. Csak OWNER vagy ADMIN jogosultságú felhasználó módosíthatja. """ # 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=status.HTTP_403_FORBIDDEN, detail="Nincs jogosultságod a szervezet beállításainak módosításához (csak OWNER/ADMIN)." ) # 2. Szervezet lekérése stmt_org = select(Organization).where( Organization.id == org_id, Organization.is_deleted == False ) result = await db.execute(stmt_org) org = result.scalar_one_or_none() if not org: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Szervezet nem található." ) # 3. JSONB merge: csak a kapott kulcsokat frissítjük, a többit megtartjuk update_dict = update_data.dict(exclude_unset=True) if "visual_settings" in update_dict and update_dict["visual_settings"] is not None: incoming_vs = update_dict["visual_settings"] current_vs = dict(org.visual_settings) if org.visual_settings else {} # Mély merge: a meglévő kulcsok megmaradnak, csak a kapottak frissülnek current_vs.update(incoming_vs) org.visual_settings = current_vs del update_dict["visual_settings"] # 4. Egyéb mezők frissítése (display_name, language, default_currency) for field, value in update_dict.items(): if hasattr(org, field) and value is not None: setattr(org, field, value) try: await db.commit() await db.refresh(org) except Exception as e: await db.rollback() logger.error(f"Hiba a szervezet frissítésekor (org_id={org_id}): {e}") raise HTTPException( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Adatbázis hiba: {str(e)}" ) return OrganizationResponse.model_validate(org) @router.patch("/{org_id}", response_model=OrganizationResponse) async def update_organization( org_id: int, update_data: OrganizationUpdate, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """ Szervezet adatainak részleges frissítése (általános PATCH). Támogatja a cégadatok (name, full_name, tax_number, display_name), cím adatok (address_zip, address_city, address_street_name, stb.), valamint a visual_settings, language, default_currency mezők frissítését. Csak OWNER vagy ADMIN jogosultságú felhasználó módosíthatja. """ # 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=status.HTTP_403_FORBIDDEN, detail="Nincs jogosultságod a szervezet adatainak módosításához (csak OWNER/ADMIN)." ) # 2. Szervezet lekérése stmt_org = select(Organization).where( Organization.id == org_id, Organization.is_deleted == False ) result = await db.execute(stmt_org) org = result.scalar_one_or_none() if not org: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Szervezet nem található." ) # 3. JSONB merge visual_settings esetén update_dict = update_data.model_dump(exclude_unset=True) if "visual_settings" in update_dict and update_dict["visual_settings"] is not None: incoming_vs = update_dict["visual_settings"] current_vs = dict(org.visual_settings) if org.visual_settings else {} current_vs.update(incoming_vs) org.visual_settings = current_vs del update_dict["visual_settings"] # 4. Többi mező frissítése for field, value in update_dict.items(): if hasattr(org, field) and value is not None: setattr(org, field, value) try: await db.commit() await db.refresh(org) except Exception as e: await db.rollback() logger.error(f"Hiba a szervezet frissítésekor (org_id={org_id}): {e}") raise HTTPException( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Adatbázis hiba: {str(e)}" ) return OrganizationResponse.model_validate(org) # --- 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."} # ── CÉG CSATLAKOZÁS ÉS ÁRVA CÉG ÁTVÉTEL ── from app.models.identity import OneTimePassword import random class JoinRequestIn(BaseModel): """Csatlakozási kérelem egy céghez, amelynek van aktív adminja.""" message: Optional[str] = Field(None, description="Üzenet az adminnak") class ClaimRequestIn(BaseModel): """Árva cég átvételi kérelem - email ellenőrzés.""" email: EmailStr = Field(..., description="A cég eredeti email címe") class ClaimVerifyIn(BaseModel): """Árva cég átvételi kérelem - OTP kód ellenőrzés.""" email: EmailStr = Field(..., description="A cég eredeti email címe") code: str = Field(..., min_length=6, max_length=6, description="6-jegyű megerősítő kód") @router.post("/{org_id}/join-request") async def request_join_organization( org_id: int, request: JoinRequestIn, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user), ): """ Csatlakozási kérelem egy szervezethez, amelynek van aktív adminja. A kérelem 7 napig él, amíg az admin jóvá nem hagyja. """ # Ellenőrizzük, hogy létezik-e a szervezet stmt = select(Organization).where( Organization.id == org_id, Organization.is_deleted == False ) result = await db.execute(stmt) org = result.scalar_one_or_none() if not org: raise HTTPException(status_code=404, detail="Szervezet nem található.") # Ellenőrizzük, hogy van-e aktív admin stmt = select(OrganizationMember).where( OrganizationMember.organization_id == org_id, OrganizationMember.role.in_([OrgUserRole.OWNER, OrgUserRole.ADMIN]) ).join(User, OrganizationMember.user_id == User.id).where(User.is_active == True, User.is_deleted == False) result = await db.execute(stmt) active_admins = result.scalars().all() if not active_admins: raise HTTPException( status_code=400, detail="A szervezetnek nincs aktív adminja. Használja a /claim/request végpontot." ) # Ellenőrizzük, hogy a felhasználó már nem tag-e stmt = select(OrganizationMember).where( OrganizationMember.organization_id == org_id, OrganizationMember.user_id == current_user.id ) result = await db.execute(stmt) existing = result.scalar_one_or_none() if existing: raise HTTPException(status_code=400, detail="Már tagja vagy ennek a szervezetnek.") # Létrehozzuk a pending tagot new_member = OrganizationMember( organization_id=org_id, user_id=current_user.id, role=OrgUserRole.DRIVER, status="pending" ) db.add(new_member) # Naplózás await security_service.log_event( db, user_id=current_user.id, action="ORG_JOIN_REQUEST", severity=LogSeverity.info, target_type="Organization", target_id=str(org_id), new_data={"message": request.message or ""} ) await db.commit() return { "status": "pending", "message": "Csatlakozási kérelmed rögzítettük. Az admin jóváhagyására vársz." } @router.post("/{org_id}/claim/request") async def claim_orphaned_organization_request( org_id: int, request: ClaimRequestIn, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user), ): """ Árva cég átvételi kérelem indítása. Ellenőrzi, hogy a megadott email cím egyezik-e a cég eredeti email címével, majd 6-jegyű OTP kódot küld ki. """ # Ellenőrizzük a szervezetet stmt = select(Organization).where( Organization.id == org_id, Organization.is_deleted == False ) result = await db.execute(stmt) org = result.scalar_one_or_none() if not org: raise HTTPException(status_code=404, detail="Szervezet nem található.") # Ellenőrizzük, hogy tényleg árva-e (nincs aktív admin) stmt = select(OrganizationMember).where( OrganizationMember.organization_id == org_id, OrganizationMember.role.in_([OrgUserRole.OWNER, OrgUserRole.ADMIN]) ).join(User, OrganizationMember.user_id == User.id).where(User.is_active == True, User.is_deleted == False) result = await db.execute(stmt) active_admins = result.scalars().all() if active_admins: raise HTTPException( status_code=400, detail="A szervezetnek van aktív adminja. Használja a /join-request végpontot." ) # Ellenőrizzük az email címet - keressük a tulajdonos usert # A tulajdonos email-jét kell ellenőrizni owner_stmt = select(User).where( User.id == org.owner_id, User.email == request.email ) owner_result = await db.execute(owner_stmt) owner_user = owner_result.scalar_one_or_none() if not owner_user: # Lehet, hogy a tulajdonos törölve van, és az email át lett írva # Keressük a deleted_ prefix-es email-ben owner_stmt = select(User).where( User.id == org.owner_id, User.email.like(f"%{request.email}") ) owner_result = await db.execute(owner_stmt) owner_user = owner_result.scalar_one_or_none() if not owner_user: raise HTTPException( status_code=400, detail="Az email cím nem egyezik a cég tulajdonosának email címével." ) # Generáljunk 6-jegyű OTP kódot code = f"{random.randint(0, 999999):06d}" expires_at = datetime.now(timezone.utc) + timedelta(minutes=15) otp = OneTimePassword( email=request.email, code=code, otp_type="company_claim", extra_data={"org_id": org_id, "user_id": current_user.id}, expires_at=expires_at ) db.add(otp) await db.commit() # Szimulált email küldés logger.info(f"=== CÉG ÁTVÉTELI OTP ===") logger.info(f"Email: {request.email}") logger.info(f"Kód: {code}") logger.info(f"Lejárat: {expires_at}") logger.info(f"=== A kódot email-ben kell elküldeni a cég tulajdonosának ===") return { "status": "otp_sent", "message": "Egy 6-jegyű megerősítő kódot küldtünk a cég tulajdonosának email címére.", "expires_at": expires_at.isoformat() } @router.post("/{org_id}/claim/verify") async def claim_orphaned_organization_verify( org_id: int, request: ClaimVerifyIn, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user), ): """ Árva cég átvételének véglegesítése OTP kód alapján. Sikeres ellenőrzés után a próbálkozó lesz az új admin. """ now = datetime.now(timezone.utc) # Ellenőrizzük a szervezetet stmt = select(Organization).where( Organization.id == org_id, Organization.is_deleted == False ) result = await db.execute(stmt) org = result.scalar_one_or_none() if not org: raise HTTPException(status_code=404, detail="Szervezet nem található.") # Ellenőrizzük az OTP kódot otp_stmt = select(OneTimePassword).where( OneTimePassword.email == request.email, OneTimePassword.code == request.code, OneTimePassword.otp_type == "company_claim", OneTimePassword.is_used == False, OneTimePassword.expires_at > now ).order_by(OneTimePassword.created_at.desc()).limit(1) result = await db.execute(otp_stmt) otp = result.scalar_one_or_none() if not otp: raise HTTPException( status_code=400, detail="Érvénytelen vagy lejárt kód." ) # Ellenőrizzük, hogy az OTP ehhez a szervezethez tartozik-e extra = otp.extra_data or {} if extra.get("org_id") != org_id: raise HTTPException(status_code=400, detail="A kód nem ehhez a szervezethez tartozik.") # Mark OTP as used otp.is_used = True # Ellenőrizzük, hogy a felhasználó már nem tag-e member_stmt = select(OrganizationMember).where( OrganizationMember.organization_id == org_id, OrganizationMember.user_id == current_user.id ) result = await db.execute(member_stmt) existing_member = result.scalar_one_or_none() if existing_member: # Frissítjük a szerepkört existing_member.role = OrgUserRole.ADMIN else: # Új tag létrehozása ADMIN szerepkörrel new_member = OrganizationMember( organization_id=org_id, user_id=current_user.id, role=OrgUserRole.ADMIN, status="active" ) db.add(new_member) # Átírjuk a tulajdonost (owner_id) org.owner_id = current_user.id # Naplózás await security_service.log_event( db, user_id=current_user.id, action="ORG_CLAIMED", severity=LogSeverity.warning, target_type="Organization", target_id=str(org_id), new_data={"previous_owner": str(extra.get("user_id")), "claim_type": "orphaned"} ) await db.commit() return { "status": "success", "message": "Sikeresen átvetted a szervezet irányítását. Mostantól te vagy az admin." }