mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-24 11:13:10 +08:00
feat(authz): add roles, role-assignments, teams, and me/permissions APIs
Completes the OSS admin surface for RBAC: full CRUD on authz_role,
authz_role_assignment, authz_team, and authz_team_member tables, plus a
per-user effective-permissions endpoint backing the frontend permission gate.
Plugin (enterprise Casbin) is invalidated on every write so the next enforce
sees the change.
New endpoints:
* GET/POST/PATCH/DELETE /api/v1/authz/roles -- custom-role CRUD with parent
cycle detection, system-role protection, and FK-aware delete (409 when
assignments still reference the role).
* GET/POST/DELETE /api/v1/authz/role-assignments -- list filterable by
user/role/domain; superuser-only for write, self-read allowed.
* GET/POST/PATCH/DELETE /api/v1/authz/teams + GET/POST/DELETE
/api/v1/authz/teams/{id}/members -- backs the share-with-team flow and
the team admin UI.
* POST /api/v1/authz/me/permissions -- returns per-resource allowed actions
for the current user (capped at 500 ids), backs the FE permission gate.
Base service additions (lfx.services.authorization.base):
* list_visible_resource_ids -- plugin SQL prefilter for list endpoints; OSS
returns None ("no prefilter, fall through to filter_visible_resources").
* get_effective_permissions -- default impl uses batch_enforce so plugins
inherit it for free; Casbin plugin can override for a tighter query.
This commit is contained in:
@ -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
|
||||
|
||||
@ -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",
|
||||
|
||||
97
src/backend/base/langflow/api/v1/authz_me.py
Normal file
97
src/backend/base/langflow/api/v1/authz_me.py
Normal file
@ -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,
|
||||
)
|
||||
130
src/backend/base/langflow/api/v1/authz_role_assignments.py
Normal file
130
src/backend/base/langflow/api/v1/authz_role_assignments.py
Normal file
@ -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)
|
||||
221
src/backend/base/langflow/api/v1/authz_roles.py
Normal file
221
src/backend/base/langflow/api/v1/authz_roles.py
Normal file
@ -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)
|
||||
249
src/backend/base/langflow/api/v1/authz_teams.py
Normal file
249
src/backend/base/langflow/api/v1/authz_teams.py
Normal file
@ -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)
|
||||
@ -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}
|
||||
50
src/backend/base/langflow/api/v1/schemas/authz_roles.py
Normal file
50
src/backend/base/langflow/api/v1/schemas/authz_roles.py
Normal file
@ -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 ``<resource_type>:<obj_pattern>:<action>`` — "
|
||||
"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}
|
||||
65
src/backend/base/langflow/api/v1/schemas/authz_teams.py
Normal file
65
src/backend/base/langflow/api/v1/schemas/authz_teams.py
Normal file
@ -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}
|
||||
@ -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."""
|
||||
|
||||
|
||||
Reference in New Issue
Block a user