Files
service-finder/backend/app/api/v1/endpoints/organizations.py
2026-06-24 11:29:45 +00:00

1078 lines
38 KiB
Python
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# /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, Optional
from datetime import datetime, timedelta, timezone
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, or_
from sqlalchemy.orm import selectinload
from pydantic import BaseModel, Field, ConfigDict, EmailStr
from app.db.session import get_db
from app.api.deps import get_current_user, RequireSystemCapability
from app.schemas.organization import CorpOnboardIn, CorpOnboardResponse, OrganizationUpdate, OrganizationResponse
from app.schemas.subscription import SubscriptionTierResponse
from app.services.billing_engine import upgrade_org_subscription
from app.core.capabilities import Capability
from app.models.marketplace.organization import Organization, OrgType, OrganizationMember, Branch, OrgUserRole, OrgRole
from app.models.identity import User, OneTimePassword
from app.core.config import settings
from app.services.security_service import security_service
from app.models import LogSeverity
from app.services.translation_service import TranslationService
router = APIRouter()
logger = logging.getLogger(__name__)
def _(key: str, **kwargs) -> str:
"""Rövidítés a TranslationService.get_text() hívásához."""
return TranslationService.get_text(key, variables=kwargs if kwargs else None)
@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=_("ORGANIZATION.ERROR.INVALID_TAX_FORMAT")
)
# 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=_("ORGANIZATION.ERROR.TAX_ALREADY_REGISTERED")
)
# 3. KÖTELEZŐ MEZŐ: folder_slug generálása
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,
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",
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=_("ORGANIZATION.BRANCH.DEFAULT_NAME"),
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"
)
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.
"""
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=_("ORGANIZATION.ERROR.TAX_ALREADY_REGISTERED")
)
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=_("ORGANIZATION.ERROR.NAV_UNAVAILABLE")
)
if result is None:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=_("ORGANIZATION.ERROR.TAX_NOT_FOUND")
)
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ó szervezetek listázása.
P0: Minden szervezethez visszaadja a valós max_vehicles és max_garages limiteket
a subscription_tier JSONB rules alapján.
"""
from app.models.core_logic import SubscriptionTier, OrganizationSubscription
subq = (
select(Organization.id)
.outerjoin(OrganizationMember, OrganizationMember.organization_id == Organization.id)
.where(
or_(
Organization.owner_id == current_user.id,
OrganizationMember.user_id == current_user.id
)
)
.where(Organization.org_type.notin_([OrgType.service_provider, OrgType.service]))
.where(Organization.is_deleted == False)
.distinct()
.subquery()
)
stmt = (
select(Organization)
.where(Organization.id.in_(select(subq.c.id)))
.options(selectinload(Organization.members))
)
result = await db.execute(stmt)
orgs = result.scalars().all()
# ── P0: Pre-fetch all org subscription tiers for limit resolution ──
org_ids = [o.id for o in orgs]
org_tier_map: dict[int, tuple[int, int]] = {}
if org_ids:
org_subs_stmt = (
select(OrganizationSubscription.org_id, SubscriptionTier)
.join(SubscriptionTier, OrganizationSubscription.tier_id == SubscriptionTier.id)
.where(
OrganizationSubscription.org_id.in_(org_ids),
OrganizationSubscription.is_active == True
)
.distinct(OrganizationSubscription.org_id)
.order_by(OrganizationSubscription.org_id, OrganizationSubscription.valid_from.desc())
)
org_subs_result = await db.execute(org_subs_stmt)
for row in org_subs_result:
o_id = row.org_id
tier = row[1]
if tier and tier.rules:
allowances = tier.rules.get("allowances", {})
org_tier_map[o_id] = (
int(allowances.get("max_vehicles", 1)),
int(allowances.get("max_garages", 1)),
)
return [
{
"organization_id": o.id,
"owner_id": o.owner_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,
"subscription_expires_at": o.subscription_expires_at.isoformat() if o.subscription_expires_at else None,
"subscription_tier_id": o.subscription_tier_id,
# P0: Real limits from subscription tier
"max_vehicles": org_tier_map.get(o.id, (o.base_asset_limit or 1, 1))[0],
"max_garages": org_tier_map.get(o.id, (1, 1))[1],
"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,
"user_role": _get_user_role(o, current_user.id),
# ── Cím adatok (Address fields) ──
"address_zip": o.address_zip,
"address_city": o.address_city,
"address_street_name": o.address_street_name,
"address_street_type": o.address_street_type,
"address_house_number": o.address_house_number,
"address_hrsz": o.address_hrsz,
}
for o in orgs
]
def _get_user_role(org: Organization, user_id: int) -> Optional[str]:
"""
Segédfüggvény: visszaadja a felhasználó szerepkörét a szervezetben.
"""
if org.members:
for member in org.members:
if member.user_id == user_id:
return str(member.role)
if org.owner_id == user_id:
return "OWNER"
return None
# ── TEAM MANAGEMENT API VÉGPONTOK ──
class MemberResponse(BaseModel):
"""Tag adatainak visszaadása."""
id: int
organization_id: int
user_id: Optional[int] = None
person_id: Optional[int] = None
invited_email: Optional[str] = None
role: str
status: str
is_permanent: bool = False
is_verified: bool = False
expires_at: Optional[datetime] = None
created_at: Optional[datetime] = None
updated_at: Optional[datetime] = None
user_email: Optional[str] = None
user_display_name: Optional[str] = None
model_config = ConfigDict(from_attributes=True)
class MemberRoleUpdate(BaseModel):
"""Tag szerepkör módosítása."""
role: str = Field(..., description="Új szerepkör (OWNER, ADMIN, MANAGER, MEMBER, AGENT)")
class InvitationCreate(BaseModel):
"""Meghívó létrehozása."""
email: str = Field(..., description="Meghívott e-mail címe")
role: str = Field("MEMBER", description="Szerepkör a szervezetben")
message: Optional[str] = Field(None, description="Üzenet a meghívotthoz")
async def _get_org_settings(org: Organization) -> dict:
"""Visszaadja a szervezet dinamikus beállításait, alapértelmezett értékekkel."""
defaults = {"invite_expiry_days": 7, "max_members": 50, "allow_public_join": False}
if hasattr(org, 'settings') and org.settings:
merged = dict(defaults)
merged.update(org.settings)
return merged
return defaults
async def _check_org_admin_access(org_id: int, user_id: int, db: AsyncSession) -> Organization:
"""
Segédfüggvény: ellenőrzi, hogy a felhasználó OWNER vagy ADMIN a szervezetben.
Visszaadja a szervezet objektumot, vagy 403/404 hibát dob.
"""
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=404, detail=_("ORGANIZATION.ERROR.NOT_FOUND"))
# Tulajdonos automatikusan jogosult
if org.owner_id == user_id:
return org
stmt_member = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.user_id == user_id,
OrganizationMember.role.in_(["OWNER", "ADMIN"]),
OrganizationMember.status == "active"
)
member = (await db.execute(stmt_member)).scalar_one_or_none()
if not member:
raise HTTPException(status_code=403, detail=_("ORGANIZATION.ERROR.ACCESS_DENIED"))
return org
@router.get("/{org_id}/members", response_model=List[MemberResponse])
async def list_organization_members(
org_id: int,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
Szervezet tagjainak listázása.
Csak OWNER vagy ADMIN jogosultságú felhasználó érheti el.
"""
await _check_org_admin_access(org_id, current_user.id, db)
stmt = (
select(OrganizationMember, User)
.outerjoin(User, OrganizationMember.user_id == User.id)
.where(OrganizationMember.organization_id == org_id)
.order_by(OrganizationMember.created_at.desc())
)
result = await db.execute(stmt)
rows = result.all()
members = []
for member, user in rows:
member_dict = {
"id": member.id,
"organization_id": member.organization_id,
"user_id": member.user_id,
"person_id": member.person_id,
"invited_email": member.invited_email,
"role": str(member.role),
"status": member.status,
"is_permanent": member.is_permanent,
"is_verified": member.is_verified,
"expires_at": member.expires_at,
"created_at": member.created_at,
"updated_at": member.updated_at,
"user_email": user.email if user else None,
"user_display_name": user.email if user else None,
}
members.append(member_dict)
return members
@router.post("/{org_id}/invitations", response_model=dict, status_code=status.HTTP_201_CREATED)
async def create_invitation(
org_id: int,
invite_in: InvitationCreate,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
Meghívó küldése egy szervezetbe.
Az invite_expiry_days a szervezet dinamikus settings mezőjéből olvasódik ki.
Csak OWNER vagy ADMIN jogosultságú felhasználó hívhatja.
"""
org = await _check_org_admin_access(org_id, current_user.id, db)
# Dinamikus beállítások lekérése
org_settings = await _get_org_settings(org)
invite_expiry_days = org_settings.get("invite_expiry_days", 7)
# Ellenőrizzük, hogy a szerepkör érvényes-e
valid_roles = [r.value for r in OrgUserRole]
role_upper = invite_in.role.upper()
if role_upper not in valid_roles:
raise HTTPException(
status_code=400,
detail=_("ORGANIZATION.ERROR.INVALID_ROLE", roles=", ".join(valid_roles))
)
# Ellenőrizzük, hogy az email már létezik-e a szervezetben
stmt_existing = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
or_(
OrganizationMember.invited_email == invite_in.email,
OrganizationMember.user_id == select(User.id).where(User.email == invite_in.email).scalar_subquery()
)
)
existing = (await db.execute(stmt_existing)).scalar_one_or_none()
if existing:
raise HTTPException(status_code=400, detail=_("ORGANIZATION.ERROR.EMAIL_ALREADY_INVITED"))
# Keresés a felhasználók között
stmt_user = select(User).where(User.email == invite_in.email)
target_user = (await db.execute(stmt_user)).scalar_one_or_none()
expires_at = datetime.now(timezone.utc) + timedelta(days=invite_expiry_days)
if target_user:
new_member = OrganizationMember(
organization_id=org_id,
user_id=target_user.id,
role=role_upper,
status="pending",
expires_at=expires_at,
invited_email=invite_in.email,
)
db.add(new_member)
await db.commit()
logger.info(f"Meghívó elküldve létező felhasználónak: {invite_in.email} (org_id={org_id})")
return {
"status": "success",
"message": _("ORGANIZATION.INVITATION.SENT_EXISTING_USER"),
"member_id": new_member.id,
"expires_at": expires_at.isoformat()
}
else:
new_member = OrganizationMember(
organization_id=org_id,
invited_email=invite_in.email,
role=role_upper,
status="pending",
expires_at=expires_at,
)
db.add(new_member)
await db.commit()
logger.info(f"Meghívó elküldve új (nem regisztrált) felhasználónak: {invite_in.email} (org_id={org_id})")
return {
"status": "success",
"message": _("ORGANIZATION.INVITATION.SENT_NEW_USER"),
"member_id": new_member.id,
"expires_at": expires_at.isoformat()
}
@router.patch("/{org_id}/members/{member_id}", response_model=MemberResponse)
async def update_member_role(
org_id: int,
member_id: int,
role_update: MemberRoleUpdate,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
Tag szerepkörének módosítása a szervezetben.
Csak OWNER vagy ADMIN jogosultságú felhasználó módosíthatja.
"""
org = await _check_org_admin_access(org_id, current_user.id, db)
# Ellenőrizzük, hogy a szerepkör érvényes-e
valid_roles = [r.value for r in OrgUserRole]
role_upper = role_update.role.upper()
if role_upper not in valid_roles:
raise HTTPException(
status_code=400,
detail=_("ORGANIZATION.ERROR.INVALID_ROLE", roles=", ".join(valid_roles))
)
# Tag lekérése
stmt_member = select(OrganizationMember).where(
OrganizationMember.id == member_id,
OrganizationMember.organization_id == org_id
)
member = (await db.execute(stmt_member)).scalar_one_or_none()
if not member:
raise HTTPException(status_code=404, detail=_("ORGANIZATION.ERROR.MEMBER_NOT_FOUND"))
# Nem lehet az utolsó OWNER-t átállítani
if member.role == "OWNER" and role_upper != "OWNER":
stmt_other_owners = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.role == "OWNER",
OrganizationMember.id != member_id,
OrganizationMember.status == "active"
)
other_owners = (await db.execute(stmt_other_owners)).scalars().all()
if not other_owners:
raise HTTPException(
status_code=400,
detail=_("ORGANIZATION.ERROR.LAST_OWNER_ROLE")
)
# Ha valaki ADMIN-t akar csinálni, az csak OWNER lehet
if role_upper == "OWNER" and current_user.id != org.owner_id:
stmt_caller = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.user_id == current_user.id,
OrganizationMember.role == "OWNER",
OrganizationMember.status == "active"
)
caller = (await db.execute(stmt_caller)).scalar_one_or_none()
if not caller:
raise HTTPException(status_code=403, detail=_("ORGANIZATION.ERROR.OWNER_ONLY_TRANSFER"))
member.role = role_upper
await db.commit()
await db.refresh(member)
user = None
if member.user_id:
stmt_user = select(User).where(User.id == member.user_id)
user = (await db.execute(stmt_user)).scalar_one_or_none()
return {
"id": member.id,
"organization_id": member.organization_id,
"user_id": member.user_id,
"person_id": member.person_id,
"invited_email": member.invited_email,
"role": str(member.role),
"status": member.status,
"is_permanent": member.is_permanent,
"is_verified": member.is_verified,
"expires_at": member.expires_at,
"created_at": member.created_at,
"updated_at": member.updated_at,
"user_email": user.email if user else None,
"user_display_name": user.display_name if user else None,
}
@router.delete("/{org_id}/members/{member_id}", status_code=status.HTTP_200_OK)
async def remove_member(
org_id: int,
member_id: int,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""
Tag eltávolítása / meghívó visszavonása a szervezetből.
Csak OWNER vagy ADMIN jogosultságú felhasználó végezheti.
"""
org = await _check_org_admin_access(org_id, current_user.id, db)
stmt_member = select(OrganizationMember).where(
OrganizationMember.id == member_id,
OrganizationMember.organization_id == org_id
)
member = (await db.execute(stmt_member)).scalar_one_or_none()
if not member:
raise HTTPException(status_code=404, detail=_("ORGANIZATION.ERROR.MEMBER_NOT_FOUND"))
# Nem lehet az utolsó OWNER-t eltávolítani
if member.role == "OWNER":
stmt_other_owners = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.role == "OWNER",
OrganizationMember.id != member_id,
OrganizationMember.status == "active"
)
other_owners = (await db.execute(stmt_other_owners)).scalars().all()
if not other_owners:
raise HTTPException(
status_code=400,
detail=_("ORGANIZATION.ERROR.LAST_OWNER_REMOVE")
)
await security_service.log_event(
db, user_id=current_user.id, action="ORG_MEMBER_REMOVED",
severity=LogSeverity.info, target_type="OrganizationMember", target_id=str(member_id),
old_data={"role": str(member.role), "status": member.status, "user_id": member.user_id}
)
if member.status == "pending":
await db.delete(member)
action = "invitation_revoked"
message = _("ORGANIZATION.MEMBER.INVITATION_REVOKED")
else:
member.status = "removed"
action = "member_removed"
message = _("ORGANIZATION.MEMBER.REMOVED")
await db.commit()
return {
"status": "success",
"action": action,
"message": message
}
@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).
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=_("ORGANIZATION.ERROR.ACCESS_DENIED")
)
# 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=_("ORGANIZATION.ERROR.NOT_FOUND")
)
# 3. JSONB merge: csak a kapott kulcsokat frissítjük
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 {}
current_vs.update(incoming_vs)
org.visual_settings = current_vs
del update_dict["visual_settings"]
# 4. Egyéb mezők 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=_("ORGANIZATION.ERROR.DATABASE_ERROR", error=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).
Csak OWNER vagy ADMIN jogosultságú felhasználó módosíthatja.
"""
# 1. Jogosultság ellenőrzése
# Először lekérjük a szervezetet, hogy ellenőrizhessük az owner_id-t is
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=_("ORGANIZATION.ERROR.NOT_FOUND")
)
# Jogosultság: OWNER/ADMIN a OrganizationMember táblában VAGY owner_id a Organization-ben
is_owner_by_field = (org.owner_id == current_user.id)
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 and not is_owner_by_field:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail=_("ORGANIZATION.ERROR.ACCESS_DENIED")
)
# 2. 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=_("ORGANIZATION.ERROR.DATABASE_ERROR", error=str(e))
)
return OrganizationResponse.model_validate(org)
# ── P0 FEATURE: SUBSCRIPTION & PACKAGE ASSIGNMENT BRIDGE ──
class SubscriptionAssignIn(BaseModel):
"""PUT body a szervezet előfizetési csomagjának beállításához."""
tier_id: int = Field(..., description="SubscriptionTier ID a system.subscription_tiers táblából")
@router.put("/{org_id}/subscription", response_model=dict)
async def assign_organization_subscription(
org_id: int,
body: SubscriptionAssignIn,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user),
):
"""
Szervezet előfizetési csomagjának beállítása (Subscription & Package Assignment Bridge).
P0 Feature: This endpoint assigns a subscription_tier to an organization.
It requires org-level ADMIN or OWNER access (checked via _check_org_admin_access).
System-level SUPERADMIN capability is also accepted via the org admin check.
The endpoint:
1. Validates the tier exists in system.subscription_tiers
2. Sets subscription_tier_id on the Organization record
3. Updates subscription_plan and base_asset_limit from tier rules
4. Creates/updates a finance.org_subscriptions audit record
"""
# Allow org ADMIN/OWNER to manage their own subscription
await _check_org_admin_access(org_id, current_user.id, db)
try:
result = await upgrade_org_subscription(
db=db,
org_id=org_id,
tier_id=body.tier_id,
actor_user_id=current_user.id
)
await db.commit()
await security_service.log_event(
db, user_id=current_user.id, action="ORG_SUBSCRIPTION_ASSIGNED",
severity=LogSeverity.info, target_type="Organization", target_id=str(org_id),
new_data={"tier_id": body.tier_id, "tier_name": result.get("tier_name")}
)
return result
except ValueError as e:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=str(e)
)
except Exception as e:
await db.rollback()
logger.error(f"Subscription assignment failed: org_id={org_id}, tier_id={body.tier_id}: {e}")
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=_("ORGANIZATION.ERROR.DATABASE_ERROR", error=str(e))
)
# ── CÉG CSATLAKOZÁS ÉS ÁRVA CÉG ÁTVÉTEL ──
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.
"""
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=_("ORGANIZATION.ERROR.NOT_FOUND"))
# Ellenőrizzük, hogy van-e aktív admin
stmt = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.role.in_(["OWNER", "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=_("ORGANIZATION.ERROR.NO_ACTIVE_ADMIN")
)
# 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=_("ORGANIZATION.ERROR.ALREADY_MEMBER"))
# Létrehozzuk a pending tagot
new_member = OrganizationMember(
organization_id=org_id,
user_id=current_user.id,
role="MEMBER",
status="pending"
)
db.add(new_member)
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": _("ORGANIZATION.JOIN_REQUEST.SENT")
}
@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.
"""
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=_("ORGANIZATION.ERROR.NOT_FOUND"))
# Ellenőrizzük, hogy tényleg árva-e
stmt = select(OrganizationMember).where(
OrganizationMember.organization_id == org_id,
OrganizationMember.role.in_(["OWNER", "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=_("ORGANIZATION.ERROR.HAS_ACTIVE_ADMIN")
)
# Ellenőrizzük az email címet
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:
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=_("ORGANIZATION.ERROR.EMAIL_MISMATCH")
)
# 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()
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": _("ORGANIZATION.CLAIM.OTP_SENT"),
"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=_("ORGANIZATION.ERROR.NOT_FOUND"))
# 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=_("ORGANIZATION.ERROR.INVALID_CODE")
)
# 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=_("ORGANIZATION.ERROR.CODE_ORG_MISMATCH"))
# 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 = "ADMIN"
else:
# Új tag létrehozása ADMIN szerepkörrel
new_member = OrganizationMember(
organization_id=org_id,
user_id=current_user.id,
role="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": _("ORGANIZATION.CLAIM.SUCCESS")
}