diff --git a/src/backend/base/langflow/api/router.py b/src/backend/base/langflow/api/router.py index 3aec88f70b..e87f0874bb 100644 --- a/src/backend/base/langflow/api/router.py +++ b/src/backend/base/langflow/api/router.py @@ -5,7 +5,11 @@ from lfx.services.settings.feature_flags import FEATURE_FLAGS from langflow.api.v1 import ( api_key_router, authz_audit_router, + authz_me_router, + authz_role_assignments_router, + authz_roles_router, authz_shares_router, + authz_teams_router, chat_router, endpoints_router, extensions_router, @@ -81,6 +85,10 @@ router_v1.include_router(models_router) router_v1.include_router(model_options_router) router_v1.include_router(authz_shares_router) router_v1.include_router(authz_audit_router) +router_v1.include_router(authz_roles_router) +router_v1.include_router(authz_role_assignments_router) +router_v1.include_router(authz_teams_router) +router_v1.include_router(authz_me_router) # Extension reload is Mode A (local-dev / pip-installed) only. The route is diff --git a/src/backend/base/langflow/api/v1/__init__.py b/src/backend/base/langflow/api/v1/__init__.py index 4b289bcca9..a24d5ca9e7 100644 --- a/src/backend/base/langflow/api/v1/__init__.py +++ b/src/backend/base/langflow/api/v1/__init__.py @@ -1,6 +1,10 @@ from langflow.api.v1.api_key import router as api_key_router from langflow.api.v1.authz_audit import router as authz_audit_router +from langflow.api.v1.authz_me import router as authz_me_router +from langflow.api.v1.authz_role_assignments import router as authz_role_assignments_router +from langflow.api.v1.authz_roles import router as authz_roles_router from langflow.api.v1.authz_shares import router as authz_shares_router +from langflow.api.v1.authz_teams import router as authz_teams_router from langflow.api.v1.chat import router as chat_router from langflow.api.v1.endpoints import router as endpoints_router from langflow.api.v1.extensions import router as extensions_router @@ -30,7 +34,11 @@ from langflow.api.v1.voice_mode import router as voice_mode_router __all__ = [ "api_key_router", "authz_audit_router", + "authz_me_router", + "authz_role_assignments_router", + "authz_roles_router", "authz_shares_router", + "authz_teams_router", "chat_router", "endpoints_router", "extensions_router", diff --git a/src/backend/base/langflow/api/v1/authz_me.py b/src/backend/base/langflow/api/v1/authz_me.py new file mode 100644 index 0000000000..2dfe1357b0 --- /dev/null +++ b/src/backend/base/langflow/api/v1/authz_me.py @@ -0,0 +1,97 @@ +"""Per-user effective-permissions endpoint, used by the frontend permission gate. + +The UI calls this once per page load with the list of resource IDs it wants to +render and learns which actions to enable/disable per resource — without making +a 403-triggering request for each one. Backed by +:meth:`BaseAuthorizationService.get_effective_permissions`; OSS pass-through +returns every action for every ID (no policy applied) and the Casbin plugin +overrides it with a tighter implementation. +""" + +from __future__ import annotations + +from typing import Literal +from uuid import UUID + +from fastapi import APIRouter, HTTPException, status +from pydantic import BaseModel, Field + +from langflow.api.utils import CurrentActiveUser +from langflow.services.deps import get_authorization_service + +router = APIRouter(prefix="/authz/me", tags=["Authorization"]) + +# Match the resource slugs used by ensure_*_permission helpers. +ResourceTypeLiteral = Literal[ + "flow", + "deployment", + "project", + "knowledge_base", + "variable", + "file", + "component", +] + +# Default action vocabulary — matches Casbin's KNOWN_ACTIONS in EE roles.py. +_DEFAULT_ACTIONS: tuple[str, ...] = ("read", "write", "execute", "delete", "create") +_MAX_RESOURCE_IDS = 500 + + +class EffectivePermissionsRequest(BaseModel): + """Body for :func:`get_effective_permissions`.""" + + resource_type: ResourceTypeLiteral + resource_ids: list[UUID] = Field( + ..., + description="Resource IDs to evaluate. Capped at 500 per request to keep batch_enforce bounded.", + ) + actions: list[str] | None = Field( + default=None, + description="Actions to check. Defaults to read/write/execute/delete/create.", + ) + domain: str = Field( + default="*", + description="Casbin domain — typically ``project:{folder_id}`` or ``*``.", + ) + + +class EffectivePermissionsResponse(BaseModel): + """Response: ``{resource_id: [allowed_actions]}``.""" + + resource_type: ResourceTypeLiteral + permissions: dict[UUID, list[str]] + + +@router.post("/permissions", response_model=EffectivePermissionsResponse) +async def get_effective_permissions( + body: EffectivePermissionsRequest, + current_user: CurrentActiveUser, +) -> EffectivePermissionsResponse: + """Return per-resource allowed actions for the current user. + + Use this to render the UI permission gate (greyed-out buttons etc.) without + flooding the audit log with denied probes. Empty list for a resource_id + means the user cannot perform any of the requested actions on that resource. + """ + if not body.resource_ids: + return EffectivePermissionsResponse(resource_type=body.resource_type, permissions={}) + if len(body.resource_ids) > _MAX_RESOURCE_IDS: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail=f"resource_ids capped at {_MAX_RESOURCE_IDS}", + ) + + authz = get_authorization_service() + actions = tuple(body.actions) if body.actions else _DEFAULT_ACTIONS + permissions = await authz.get_effective_permissions( + user_id=current_user.id, + resource_type=body.resource_type, + resource_ids=body.resource_ids, + actions=actions, + domain=body.domain, + context={"is_superuser": getattr(current_user, "is_superuser", False)}, + ) + return EffectivePermissionsResponse( + resource_type=body.resource_type, + permissions=permissions, + ) diff --git a/src/backend/base/langflow/api/v1/authz_role_assignments.py b/src/backend/base/langflow/api/v1/authz_role_assignments.py new file mode 100644 index 0000000000..842d7be035 --- /dev/null +++ b/src/backend/base/langflow/api/v1/authz_role_assignments.py @@ -0,0 +1,130 @@ +"""CRUD API for authz_role_assignment rows. + +Assignments bind a user to a role within an optional domain. The actual policy +compilation (rule rows in ``casbin_rule``) is performed by the authorization +plugin — OSS keeps the assignment table and invalidates the plugin's cache on +write so the next ``enforce()`` picks up the change. +""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Annotated +from uuid import UUID + +from fastapi import APIRouter, HTTPException, Query, status +from lfx.log.logger import logger +from sqlalchemy.exc import IntegrityError +from sqlmodel import select + +from langflow.api.utils import CurrentActiveUser, DbSession +from langflow.api.v1.schemas.authz_role_assignments import ( + RoleAssignmentCreate, + RoleAssignmentRead, +) +from langflow.services.database.models.auth import AuthzRole, AuthzRoleAssignment +from langflow.services.database.models.user.model import User +from langflow.services.deps import get_authorization_service + +router = APIRouter(prefix="/authz/role-assignments", tags=["Authorization"]) + + +def _require_superuser(user) -> None: + if not getattr(user, "is_superuser", False): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Superuser required to administer role assignments.", + ) + + +@router.get("", response_model=list[RoleAssignmentRead]) +@router.get("/", response_model=list[RoleAssignmentRead]) +async def list_assignments( + session: DbSession, + current_user: CurrentActiveUser, + user_id: Annotated[UUID | None, Query(description="Filter by user")] = None, + role_id: Annotated[UUID | None, Query(description="Filter by role")] = None, + domain_type: Annotated[str | None, Query()] = None, + domain_id: Annotated[UUID | None, Query()] = None, +) -> list[RoleAssignmentRead]: + """List role assignments. Users may query their own assignments; superusers see all. + + A user querying with ``user_id != self.id`` and lacking superuser triggers 403. + """ + if user_id is None or user_id != current_user.id: + _require_superuser(current_user) + stmt = select(AuthzRoleAssignment) + if user_id is not None: + stmt = stmt.where(AuthzRoleAssignment.user_id == user_id) + if role_id is not None: + stmt = stmt.where(AuthzRoleAssignment.role_id == role_id) + if domain_type is not None: + stmt = stmt.where(AuthzRoleAssignment.domain_type == domain_type) + if domain_id is not None: + stmt = stmt.where(AuthzRoleAssignment.domain_id == domain_id) + rows = (await session.exec(stmt.order_by(AuthzRoleAssignment.assigned_at.desc()))).all() + return [RoleAssignmentRead.model_validate(row) for row in rows] + + +@router.post("", response_model=RoleAssignmentRead, status_code=status.HTTP_201_CREATED) +@router.post("/", response_model=RoleAssignmentRead, status_code=status.HTTP_201_CREATED) +async def create_assignment( + payload: RoleAssignmentCreate, + current_user: CurrentActiveUser, + session: DbSession, +) -> RoleAssignmentRead: + """Assign a role to a user. Superuser-only.""" + _require_superuser(current_user) + + user = await session.get(User, payload.user_id) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="user_id not found") + role = await session.get(AuthzRole, payload.role_id) + if role is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="role_id not found") + + assignment = AuthzRoleAssignment( + user_id=payload.user_id, + role_id=payload.role_id, + domain_type=payload.domain_type, + domain_id=payload.domain_id, + assigned_at=datetime.now(timezone.utc), + assigned_by=current_user.id, + ) + session.add(assignment) + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Assignment already exists for this user/role/domain", + ) from exc + await session.refresh(assignment) + await get_authorization_service().invalidate_user(payload.user_id) + logger.info( + "Assigned role=%s to user=%s (domain=%s/%s)", + role.name, + payload.user_id, + payload.domain_type, + payload.domain_id, + ) + return RoleAssignmentRead.model_validate(assignment) + + +@router.delete("/{assignment_id}", status_code=status.HTTP_204_NO_CONTENT) +async def delete_assignment( + assignment_id: UUID, + current_user: CurrentActiveUser, + session: DbSession, +) -> None: + """Revoke a role assignment. Superuser-only.""" + _require_superuser(current_user) + assignment = await session.get(AuthzRoleAssignment, assignment_id) + if assignment is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Assignment not found") + user_id = assignment.user_id + await session.delete(assignment) + await session.commit() + await get_authorization_service().invalidate_user(user_id) + logger.info("Revoked role assignment id=%s (user=%s)", assignment_id, user_id) diff --git a/src/backend/base/langflow/api/v1/authz_roles.py b/src/backend/base/langflow/api/v1/authz_roles.py new file mode 100644 index 0000000000..f682fd98bc --- /dev/null +++ b/src/backend/base/langflow/api/v1/authz_roles.py @@ -0,0 +1,221 @@ +"""CRUD API for authz_role rows (enforcement is delegated to authorization plugins).""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Annotated +from uuid import UUID + +from fastapi import APIRouter, HTTPException, Query, status +from lfx.log.logger import logger +from sqlalchemy.exc import IntegrityError +from sqlmodel import select + +from langflow.api.utils import CurrentActiveUser, DbSession +from langflow.api.v1.schemas.authz_roles import RoleCreate, RoleRead, RoleUpdate +from langflow.services.database.models.auth import AuthzRole, AuthzRoleAssignment +from langflow.services.deps import get_authorization_service + +router = APIRouter(prefix="/authz/roles", tags=["Authorization"]) + + +def _require_superuser(user) -> None: + """Superuser-only gate. Role admin is an operations action.""" + if not getattr(user, "is_superuser", False): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Superuser required to administer roles.", + ) + + +async def _detect_parent_cycle( + session: DbSession, + *, + role_id: UUID, + proposed_parent_id: UUID, +) -> bool: + """Walk the parent chain from ``proposed_parent_id``; True if ``role_id`` appears. + + Used to reject ``PATCH`` requests that would set a role as its own ancestor. + Walks at most ``len(all_roles)`` steps so a pre-existing cycle terminates. + """ + visited: set[UUID] = set() + cursor: UUID | None = proposed_parent_id + while cursor is not None and cursor not in visited: + if cursor == role_id: + return True + visited.add(cursor) + parent = await session.get(AuthzRole, cursor) + if parent is None: + return False + cursor = parent.parent_role_id + return False + + +@router.get("", response_model=list[RoleRead]) +@router.get("/", response_model=list[RoleRead]) +async def list_roles( + session: DbSession, + current_user: CurrentActiveUser, # noqa: ARG001 — any authenticated user can list + is_system: Annotated[bool | None, Query(description="Filter by is_system flag")] = None, + name: Annotated[str | None, Query(description="Substring match on role name")] = None, +) -> list[RoleRead]: + """List all roles. Open to authenticated users so the UI can populate dropdowns.""" + stmt = select(AuthzRole) + if is_system is not None: + stmt = stmt.where(AuthzRole.is_system == is_system) + if name: + stmt = stmt.where(AuthzRole.name.ilike(f"%{name}%")) + rows = (await session.exec(stmt.order_by(AuthzRole.name))).all() + return [RoleRead.model_validate(row) for row in rows] + + +@router.get("/{role_id}", response_model=RoleRead) +async def read_role( + role_id: UUID, + session: DbSession, + current_user: CurrentActiveUser, # noqa: ARG001 — any authenticated user can read +) -> RoleRead: + role = await session.get(AuthzRole, role_id) + if role is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Role not found") + return RoleRead.model_validate(role) + + +@router.post("", response_model=RoleRead, status_code=status.HTTP_201_CREATED) +@router.post("/", response_model=RoleRead, status_code=status.HTTP_201_CREATED) +async def create_role( + payload: RoleCreate, + current_user: CurrentActiveUser, + session: DbSession, +) -> RoleRead: + """Create a custom (non-system) role. Superuser-only.""" + _require_superuser(current_user) + + if payload.parent_role_id is not None: + parent = await session.get(AuthzRole, payload.parent_role_id) + if parent is None: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="parent_role_id does not reference an existing role", + ) + + role = AuthzRole( + name=payload.name, + description=payload.description, + is_system=False, + permissions=list(payload.permissions), + parent_role_id=payload.parent_role_id, + created_by=current_user.id, + ) + session.add(role) + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail=f"Role with name {payload.name!r} already exists", + ) from exc + await session.refresh(role) + await get_authorization_service().invalidate_all() + logger.info("Created role %s (id=%s)", role.name, role.id) + return RoleRead.model_validate(role) + + +@router.patch("/{role_id}", response_model=RoleRead) +async def update_role( + role_id: UUID, + payload: RoleUpdate, + current_user: CurrentActiveUser, + session: DbSession, +) -> RoleRead: + """Update fields on a custom role. System roles are read-only.""" + _require_superuser(current_user) + + role = await session.get(AuthzRole, role_id) + if role is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Role not found") + if role.is_system: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="System roles cannot be modified", + ) + + if payload.parent_role_id is not None: + # Validate parent exists and would not create a cycle. + if payload.parent_role_id == role.id: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="A role cannot be its own parent", + ) + parent = await session.get(AuthzRole, payload.parent_role_id) + if parent is None: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="parent_role_id does not reference an existing role", + ) + if await _detect_parent_cycle(session, role_id=role.id, proposed_parent_id=payload.parent_role_id): + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="Setting this parent would create a role hierarchy cycle", + ) + role.parent_role_id = payload.parent_role_id + + if payload.name is not None: + role.name = payload.name + if payload.description is not None: + role.description = payload.description + if payload.permissions is not None: + role.permissions = list(payload.permissions) + role.updated_at = datetime.now(timezone.utc) + + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Name conflict — another role already uses this name", + ) from exc + await session.refresh(role) + await get_authorization_service().invalidate_role(role.id) + logger.info("Updated role %s (id=%s)", role.name, role.id) + return RoleRead.model_validate(role) + + +@router.delete("/{role_id}", status_code=status.HTTP_204_NO_CONTENT) +async def delete_role( + role_id: UUID, + current_user: CurrentActiveUser, + session: DbSession, +) -> None: + """Delete a custom role. + + System roles cannot be deleted; roles with active assignments return 409 + (delete the assignments first). + """ + _require_superuser(current_user) + + role = await session.get(AuthzRole, role_id) + if role is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Role not found") + if role.is_system: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="System roles cannot be deleted", + ) + + assigned = ( + await session.exec(select(AuthzRoleAssignment).where(AuthzRoleAssignment.role_id == role_id).limit(1)) + ).first() + if assigned is not None: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Role still has active assignments — revoke them before deleting", + ) + + await session.delete(role) + await session.commit() + await get_authorization_service().invalidate_role(role_id) + logger.info("Deleted role id=%s", role_id) diff --git a/src/backend/base/langflow/api/v1/authz_teams.py b/src/backend/base/langflow/api/v1/authz_teams.py new file mode 100644 index 0000000000..bc7701b961 --- /dev/null +++ b/src/backend/base/langflow/api/v1/authz_teams.py @@ -0,0 +1,249 @@ +"""CRUD API for authz_team and authz_team_member rows. + +Teams group users for bulk role assignment and share targeting. The plugin +compiles team memberships to its own representation (e.g. enterprise Casbin +emits ``g, user:{id}, team:{id}, *`` rules during share PolicySync). +""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Annotated +from uuid import UUID + +from fastapi import APIRouter, HTTPException, Query, status +from lfx.log.logger import logger +from sqlalchemy.exc import IntegrityError +from sqlmodel import select + +from langflow.api.utils import CurrentActiveUser, DbSession +from langflow.api.v1.schemas.authz_teams import ( + TeamCreate, + TeamMemberCreate, + TeamMemberRead, + TeamRead, + TeamUpdate, +) +from langflow.services.database.models.auth import AuthzTeam, AuthzTeamMember +from langflow.services.database.models.user.model import User +from langflow.services.deps import get_authorization_service + +router = APIRouter(prefix="/authz/teams", tags=["Authorization"]) + + +def _require_superuser(user) -> None: + if not getattr(user, "is_superuser", False): + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Superuser required to administer teams.", + ) + + +# --- teams ---------------------------------------------------------------- # + + +@router.get("", response_model=list[TeamRead]) +@router.get("/", response_model=list[TeamRead]) +async def list_teams( + session: DbSession, + current_user: CurrentActiveUser, # noqa: ARG001 — any authenticated user can list + search: Annotated[str | None, Query(description="Substring match on team_name or adom_name")] = None, + is_active: Annotated[bool | None, Query()] = None, +) -> list[TeamRead]: + """List teams. Open to any authenticated user (for the share dialog's team picker).""" + stmt = select(AuthzTeam) + if search: + like = f"%{search}%" + stmt = stmt.where((AuthzTeam.team_name.ilike(like)) | (AuthzTeam.adom_name.ilike(like))) + if is_active is not None: + stmt = stmt.where(AuthzTeam.is_active == is_active) + rows = (await session.exec(stmt.order_by(AuthzTeam.team_name))).all() + return [TeamRead.model_validate(row) for row in rows] + + +@router.get("/{team_id}", response_model=TeamRead) +async def read_team( + team_id: UUID, + session: DbSession, + current_user: CurrentActiveUser, # noqa: ARG001 +) -> TeamRead: + team = await session.get(AuthzTeam, team_id) + if team is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Team not found") + return TeamRead.model_validate(team) + + +@router.post("", response_model=TeamRead, status_code=status.HTTP_201_CREATED) +@router.post("/", response_model=TeamRead, status_code=status.HTTP_201_CREATED) +async def create_team( + payload: TeamCreate, + current_user: CurrentActiveUser, + session: DbSession, +) -> TeamRead: + _require_superuser(current_user) + team = AuthzTeam( + team_name=payload.team_name, + adom_name=payload.adom_name, + description=payload.description, + is_active=payload.is_active, + ) + session.add(team) + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail=f"Team with adom_name {payload.adom_name!r} already exists", + ) from exc + await session.refresh(team) + logger.info("Created team %s (id=%s)", team.team_name, team.id) + return TeamRead.model_validate(team) + + +@router.patch("/{team_id}", response_model=TeamRead) +async def update_team( + team_id: UUID, + payload: TeamUpdate, + current_user: CurrentActiveUser, + session: DbSession, +) -> TeamRead: + _require_superuser(current_user) + team = await session.get(AuthzTeam, team_id) + if team is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Team not found") + + if payload.team_name is not None: + team.team_name = payload.team_name + if payload.adom_name is not None: + team.adom_name = payload.adom_name + if payload.description is not None: + team.description = payload.description + if payload.is_active is not None: + team.is_active = payload.is_active + team.updated_at = datetime.now(timezone.utc) + + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="adom_name conflict — another team already uses this slug", + ) from exc + await session.refresh(team) + # Team metadata change doesn't require policy reload — only memberships do. + logger.info("Updated team %s (id=%s)", team.team_name, team.id) + return TeamRead.model_validate(team) + + +@router.delete("/{team_id}", status_code=status.HTTP_204_NO_CONTENT) +async def delete_team( + team_id: UUID, + current_user: CurrentActiveUser, + session: DbSession, +) -> None: + _require_superuser(current_user) + team = await session.get(AuthzTeam, team_id) + if team is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Team not found") + # Cascade on team_members handles cleanup; share rows targeting this team + # are left in place (caller may want to migrate them before deleting). + await session.delete(team) + await session.commit() + await get_authorization_service().invalidate_all() + logger.info("Deleted team id=%s", team_id) + + +# --- team members --------------------------------------------------------- # + + +@router.get("/{team_id}/members", response_model=list[TeamMemberRead]) +async def list_members( + team_id: UUID, + session: DbSession, + current_user: CurrentActiveUser, # noqa: ARG001 +) -> list[TeamMemberRead]: + """List members of a team. Any authenticated user (so the UI can render team rosters).""" + team = await session.get(AuthzTeam, team_id) + if team is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Team not found") + rows = ( + await session.exec( + select(AuthzTeamMember) + .where(AuthzTeamMember.team_id == team_id) + .order_by( + AuthzTeamMember.created_at, + ) + ) + ).all() + return [TeamMemberRead.model_validate(row) for row in rows] + + +@router.post( + "/{team_id}/members", + response_model=TeamMemberRead, + status_code=status.HTTP_201_CREATED, +) +async def add_member( + team_id: UUID, + payload: TeamMemberCreate, + current_user: CurrentActiveUser, + session: DbSession, +) -> TeamMemberRead: + _require_superuser(current_user) + team = await session.get(AuthzTeam, team_id) + if team is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Team not found") + user = await session.get(User, payload.user_id) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="user_id not found") + + member = AuthzTeamMember( + team_id=team_id, + user_id=payload.user_id, + source=payload.source, + ) + session.add(member) + try: + await session.commit() + except IntegrityError as exc: + await session.rollback() + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="User is already a member of this team", + ) from exc + await session.refresh(member) + await get_authorization_service().invalidate_user(payload.user_id) + logger.info("Added user=%s to team=%s", payload.user_id, team_id) + return TeamMemberRead.model_validate(member) + + +@router.delete( + "/{team_id}/members/{user_id}", + status_code=status.HTTP_204_NO_CONTENT, +) +async def remove_member( + team_id: UUID, + user_id: UUID, + current_user: CurrentActiveUser, + session: DbSession, +) -> None: + _require_superuser(current_user) + member = ( + await session.exec( + select(AuthzTeamMember).where( + AuthzTeamMember.team_id == team_id, + AuthzTeamMember.user_id == user_id, + ) + ) + ).first() + if member is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail="Membership not found", + ) + await session.delete(member) + await session.commit() + await get_authorization_service().invalidate_user(user_id) + logger.info("Removed user=%s from team=%s", user_id, team_id) diff --git a/src/backend/base/langflow/api/v1/schemas/authz_role_assignments.py b/src/backend/base/langflow/api/v1/schemas/authz_role_assignments.py new file mode 100644 index 0000000000..91f5a6c7c3 --- /dev/null +++ b/src/backend/base/langflow/api/v1/schemas/authz_role_assignments.py @@ -0,0 +1,37 @@ +"""Pydantic schemas for /api/v1/authz/role-assignments.""" + +from __future__ import annotations + +from datetime import datetime +from uuid import UUID + +from pydantic import BaseModel, Field + + +class RoleAssignmentCreate(BaseModel): + """Payload for assigning a role to a user.""" + + user_id: UUID + role_id: UUID + domain_type: str = Field( + default="global", + description="``global``, ``org``, ``workspace`` — matches authz_role_assignment.domain_type", + ) + domain_id: UUID | None = Field( + default=None, + description="Required when ``domain_type`` scopes to a specific org/workspace/project", + ) + + +class RoleAssignmentRead(BaseModel): + """Serialized authz_role_assignment row returned by the API.""" + + id: UUID + user_id: UUID + role_id: UUID + domain_type: str + domain_id: UUID | None + assigned_at: datetime + assigned_by: UUID | None + + model_config = {"from_attributes": True} diff --git a/src/backend/base/langflow/api/v1/schemas/authz_roles.py b/src/backend/base/langflow/api/v1/schemas/authz_roles.py new file mode 100644 index 0000000000..67ee8d0fff --- /dev/null +++ b/src/backend/base/langflow/api/v1/schemas/authz_roles.py @@ -0,0 +1,50 @@ +"""Pydantic schemas for /api/v1/authz/roles.""" + +from __future__ import annotations + +from datetime import datetime +from uuid import UUID + +from pydantic import BaseModel, Field + + +class RoleCreate(BaseModel): + """Payload for creating an authz_role row.""" + + name: str = Field(..., min_length=1, max_length=255) + description: str | None = Field(default=None) + permissions: list[str] = Field( + default_factory=list, + description=( + "Permission strings in the form ``::`` — " + "for example ``flow:*:read``, ``deployment:*:execute``. Plugin (e.g. enterprise " + "Casbin) is responsible for compiling these to its policy format." + ), + ) + parent_role_id: UUID | None = Field(default=None) + + +class RoleUpdate(BaseModel): + """Payload for updating an authz_role row (PATCH semantics — only set fields apply).""" + + name: str | None = Field(default=None, min_length=1, max_length=255) + description: str | None = None + permissions: list[str] | None = None + parent_role_id: UUID | None = None + + +class RoleRead(BaseModel): + """Serialized authz_role row returned by the API.""" + + id: UUID + name: str + description: str | None + is_system: bool + permissions: list[str] + parent_role_id: UUID | None + workspace_id: UUID | None + created_at: datetime + updated_at: datetime + created_by: UUID | None + + model_config = {"from_attributes": True} diff --git a/src/backend/base/langflow/api/v1/schemas/authz_teams.py b/src/backend/base/langflow/api/v1/schemas/authz_teams.py new file mode 100644 index 0000000000..de33e13045 --- /dev/null +++ b/src/backend/base/langflow/api/v1/schemas/authz_teams.py @@ -0,0 +1,65 @@ +"""Pydantic schemas for /api/v1/authz/teams and team members.""" + +from __future__ import annotations + +from datetime import datetime +from typing import Literal +from uuid import UUID + +from pydantic import BaseModel, Field + + +class TeamCreate(BaseModel): + """Payload for creating an authz_team.""" + + team_name: str = Field(..., min_length=1, max_length=255) + adom_name: str = Field( + ..., + min_length=1, + max_length=255, + description="Administrative-domain slug, unique across all teams (often the SSO group name).", + ) + description: str | None = None + is_active: bool = True + + +class TeamUpdate(BaseModel): + """Payload for updating an authz_team (PATCH semantics).""" + + team_name: str | None = Field(default=None, min_length=1, max_length=255) + adom_name: str | None = Field(default=None, min_length=1, max_length=255) + description: str | None = None + is_active: bool | None = None + + +class TeamRead(BaseModel): + """Serialized authz_team row returned by the API.""" + + id: UUID + team_name: str + adom_name: str + description: str | None + is_active: bool + created_at: datetime + updated_at: datetime + + model_config = {"from_attributes": True} + + +class TeamMemberCreate(BaseModel): + """Payload for adding a user to a team.""" + + user_id: UUID + source: Literal["manual", "sso"] = "manual" + + +class TeamMemberRead(BaseModel): + """Serialized authz_team_member row.""" + + id: UUID + team_id: UUID + user_id: UUID + source: str + created_at: datetime + + model_config = {"from_attributes": True} diff --git a/src/lfx/src/lfx/services/authorization/base.py b/src/lfx/src/lfx/services/authorization/base.py index a6094359a3..e36333669c 100644 --- a/src/lfx/src/lfx/services/authorization/base.py +++ b/src/lfx/src/lfx/services/authorization/base.py @@ -86,6 +86,67 @@ class BaseAuthorizationService(Service, abc.ABC): ) return [action for action, allowed in zip(actions, results, strict=True) if allowed] + async def list_visible_resource_ids( + self, + *, + user_id: UUID, + resource_type: str, + domain: str = "*", + act: str = "read", + context: dict[str, Any] | None = None, + ) -> list[UUID] | None: + """Return resource IDs of `resource_type` the user can `act` on, or ``None``. + + Plugin override returns a concrete list — typically by querying its + policy store (e.g. SQL join on ``casbin_rule``) so list endpoints can + prefilter at the DB layer and avoid fetching invisible rows. + + Base returns ``None`` meaning "no prefilter available; caller should + fetch all candidates and apply :func:`filter_visible_resources` for + in-memory filtering via :meth:`batch_enforce`." OSS pass-through stays + on ``None`` so the existing list endpoints behave unchanged. + """ + _ = (user_id, resource_type, domain, act, context) + return None + + async def get_effective_permissions( + self, + *, + user_id: UUID, + resource_type: str, + resource_ids: Sequence[UUID], + actions: Sequence[str], + domain: str = "*", + context: dict[str, Any] | None = None, + ) -> dict[UUID, list[str]]: + """Return per-resource allowed actions for ``user_id``. + + Used by the frontend permission-gating layer to grey out buttons + without round-tripping to a 403. Default implementation issues a single + :meth:`batch_enforce` over the cartesian product of ``resource_ids`` x + ``actions``; plugins can override with a tighter query. + """ + if not resource_ids or not actions: + return {rid: [] for rid in resource_ids} + + requests: list[tuple[str, str]] = [ + (f"{resource_type}:{rid}", action) for rid in resource_ids for action in actions + ] + flat = await self.batch_enforce( + user_id=user_id, + domain=domain, + requests=requests, + context=context, + ) + + result: dict[UUID, list[str]] = {} + action_count = len(actions) + for index, rid in enumerate(resource_ids): + start = index * action_count + slice_ = flat[start : start + action_count] + result[rid] = [action for action, allowed in zip(actions, slice_, strict=True) if allowed] + return result + async def invalidate_user(self, user_id: UUID) -> None: """Drop cached policy for a single user. Plugin override; OSS no-op."""