Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 9 additions & 5 deletions src/backend/base/langflow/api/v1/authz_me.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from langflow.api.utils import CurrentActiveUser, DbSessionReadOnly
from langflow.services.auth.context import current_auth_context_for_authz
from langflow.services.authorization.access_ceiling import filter_actions_by_external_access_ceiling
from langflow.services.authorization.actions import FlowAction
from langflow.services.authorization.guards import should_apply_owner_override
from langflow.services.database.models.deployment.model import Deployment
from langflow.services.database.models.file.model import File as UserFile
Expand All @@ -44,11 +45,11 @@
]

# Default action vocabulary aligned with the authorization plugin's known actions.
_DEFAULT_ACTIONS: tuple[str, ...] = ("read", "write", "execute", "delete", "create")
_DEFAULT_ACTIONS: tuple[str, ...] = ("read", "write", "execute", "delete", "create", "deploy")
_MAX_RESOURCE_IDS = 500
# Cap actions per request to bound the batch_enforce cartesian product
# (resource_ids x actions). Headroom over the 6 known actions
# (read/write/execute/delete/create/manage) covers future additions without
# (resource_ids x actions). Headroom over the 7 known actions
# (read/write/execute/delete/create/deploy/manage) covers future additions without
# letting a client request `["read"] * 100000` to flood the enforcer.
_MAX_ACTIONS = 10

Expand All @@ -74,7 +75,7 @@ class EffectivePermissionsRequest(BaseModel):
default=None,
description=(
"Actions to check. Each entry is lowercased and de-duplicated; the list is "
f"capped at {_MAX_ACTIONS}. Defaults to read/write/execute/delete/create."
f"capped at {_MAX_ACTIONS}. Defaults to read/write/execute/delete/create/deploy."
),
)
domain: str = Field(
Expand Down Expand Up @@ -159,9 +160,12 @@ async def _apply_owner_permissions(
resource_ids=resource_ids,
user_id=user_id,
)
owner_override_actions = (
tuple(action for action in actions if action != FlowAction.DEPLOY.value) if resource_type == "flow" else actions
)
for resource_id in owned_ids:
allowed = dict.fromkeys(normalized.get(resource_id, []))
allowed.update(dict.fromkeys(actions))
allowed.update(dict.fromkeys(owner_override_actions))
normalized[resource_id] = list(allowed)
return normalized

Expand Down
46 changes: 27 additions & 19 deletions src/backend/base/langflow/api/v1/deployments.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,6 @@
resolve_added_snapshot_bindings_for_update,
resolve_deployment_adapter,
resolve_flow_version_patch_for_update,
resolve_project_id_for_deployment_create,
resolve_snapshot_map_for_create,
rollback_provider_create,
rollback_provider_update,
Expand Down Expand Up @@ -74,6 +73,10 @@
filter_visible_resources,
visible_id_prefilter,
)
from langflow.services.authorization.deployment import (
authorize_flow_versions_for_deployment,
resolve_project_id_for_deployment_create,
)
from langflow.services.authorization.fetch import deny_to_404
from langflow.services.authorization.utils import _resolve_authz_domain
from langflow.services.database.models.deployment.crud import (
Expand Down Expand Up @@ -510,21 +513,25 @@ async def create_deployment(
db=session,
)
should_create_provider_resource = existing_resource_key is None
project_id = await resolve_project_id_for_deployment_create(payload=payload, user_id=current_user.id, db=session)
await ensure_deployment_permission(current_user, DeploymentAction.CREATE, project_id=project_id)
project_id = await resolve_project_id_for_deployment_create(
session=session,
current_user=current_user,
requested_project_id=payload.project_id,
)
flow_version_ids = deployment_mapper.util_create_flow_version_ids(payload)
await validate_project_scoped_flow_version_ids(
authorized_flow_version_ids = await authorize_flow_versions_for_deployment(
flow_version_ids=flow_version_ids,
user_id=current_user.id,
current_user=current_user,
project_id=project_id,
db=session,
session=session,
)
if should_create_provider_resource:
adapter_payload = await deployment_mapper.resolve_deployment_create(
user_id=current_user.id,
project_id=project_id,
db=session,
payload=payload,
authorized_flow_version_ids=authorized_flow_version_ids,
)
with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(provider_id):
provider_create_result = await deployment_adapter.create(
Expand Down Expand Up @@ -1496,26 +1503,27 @@ async def update_deployment(
deployment_row_id = deployment_row.id
deployment_resource_key = deployment_row.resource_key
deployment_provider_account_id = deployment_row.deployment_provider_account_id
# Owner-namespaced operations use ``deployment_row.user_id`` — the
# deployment, its flow versions, and its attachments all live in the
# owner's scope. ``current_user`` is still the actor for authorization
# and audit, but the data plane operates in the owner's namespace.
# Deployment/provider operations and attachment rows use the deployment
# owner's namespace. Referenced flow versions may belong to another user;
# ``current_user`` remains the actor for their authorization and audit.
owner_id = deployment_row.user_id
adapter_payload = await deployment_mapper.resolve_deployment_update(
user_id=owner_id,
deployment_db_id=deployment_row_id,
db=session,
payload=payload,
)
referenced_flow_version_ids = deployment_mapper.util_update_flow_version_ids(payload)
added_flow_version_ids, remove_flow_version_ids = resolve_flow_version_patch_for_update(
deployment_mapper=deployment_mapper,
payload=payload,
)
await validate_project_scoped_flow_version_ids(
flow_version_ids=list(dict.fromkeys([*added_flow_version_ids, *remove_flow_version_ids])),
user_id=owner_id,
authorized_flow_version_ids = await authorize_flow_versions_for_deployment(
flow_version_ids=referenced_flow_version_ids,
current_user=current_user,
project_id=deployment_row.project_id,
session=session,
)
adapter_payload = await deployment_mapper.resolve_deployment_update(
user_id=owner_id,
deployment_db_id=deployment_row_id,
db=session,
payload=payload,
authorized_flow_version_ids=authorized_flow_version_ids,
)
with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(deployment_provider_account_id):
update_result: DeploymentUpdateResult = await deployment_adapter.update(
Expand Down
16 changes: 14 additions & 2 deletions src/backend/base/langflow/api/v1/mappers/deployments/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -184,8 +184,9 @@ async def resolve_deployment_create(
project_id: UUID,
db: AsyncSession,
payload: DeploymentCreateRequest,
authorized_flow_version_ids: frozenset[UUID] = frozenset(),
) -> AdapterDeploymentCreate:
_ = (user_id, project_id, db, payload)
_ = (user_id, project_id, db, payload, authorized_flow_version_ids)
msg = "This deployment provider is not configured for creating deployments."
raise NotImplementedError(msg)

Expand All @@ -196,8 +197,9 @@ async def resolve_deployment_update(
deployment_db_id: UUID,
db: AsyncSession,
payload: DeploymentUpdateRequest,
authorized_flow_version_ids: frozenset[UUID] = frozenset(),
) -> AdapterDeploymentUpdate:
_ = (user_id, deployment_db_id, db, payload)
_ = (user_id, deployment_db_id, db, payload, authorized_flow_version_ids)
msg = "This deployment provider is not configured for updating deployments."
raise NotImplementedError(msg)

Expand Down Expand Up @@ -642,6 +644,16 @@ def util_flow_version_patch(self, payload: DeploymentUpdateRequest) -> FlowVersi
_ = payload
return FlowVersionPatch()

def util_update_flow_version_ids(self, payload: DeploymentUpdateRequest) -> list[UUID]:
"""Resolve every flow-version id referenced by an update payload.

The default contract derives references from the attachment patch.
Providers with flow operations that do not change attachments must
override this method so those references are authorized as well.
"""
patch = self.util_flow_version_patch(payload)
return list(dict.fromkeys([*patch.add_flow_version_ids, *patch.remove_flow_version_ids]))

def extract_snapshot_bindings(
self,
provider_view: DeploymentListResult,
Expand Down
71 changes: 26 additions & 45 deletions src/backend/base/langflow/api/v1/mappers/deployments/helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,7 @@
fetch_provider_resource_keys,
sync_attachment_snapshot_ids,
)
from langflow.api.v1.schemas.deployments import (
DeploymentCreateRequest,
DeploymentUpdateRequest,
)
from langflow.initial_setup.setup import get_or_create_default_folder
from langflow.api.v1.schemas.deployments import DeploymentUpdateRequest
from langflow.services.adapters.deployment.context import deployment_provider_scope
from langflow.services.database.models.deployment.crud import (
DeploymentMetadataUpdate,
Expand Down Expand Up @@ -76,7 +72,6 @@
list_deployment_attachments_with_versions,
update_deployment_attachment_provider_snapshot_id,
)
from langflow.services.database.models.folder.model import Folder
from langflow.services.database.utils import require_non_empty

if TYPE_CHECKING:
Expand Down Expand Up @@ -116,16 +111,31 @@ def _build_indexed_flow_version_ids_cte(*, flow_version_ids: list[UUID]):
)


def _require_authorized_flow_version_ids(
*,
flow_version_ids: list[UUID],
authorized_flow_version_ids: frozenset[UUID],
) -> None:
"""Reject artifact resolution that lacks authorization evidence."""
if not set(flow_version_ids).issubset(authorized_flow_version_ids):
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Flow version not found.")


async def build_flow_artifacts_from_flow_versions(
*,
db,
user_id: UUID,
deployment_user_id: UUID,
deployment_db_id: UUID,
flow_version_ids: list[UUID],
authorized_flow_version_ids: frozenset[UUID],
) -> list[tuple[UUID, int, UUID, BaseFlowArtifact]]:
"""Resolve deployment-scoped flow version ids into artifacts preserving input order."""
if not flow_version_ids:
return []
_require_authorized_flow_version_ids(
flow_version_ids=flow_version_ids,
authorized_flow_version_ids=authorized_flow_version_ids,
)

indexed_flow_version_ids_cte = _build_indexed_flow_version_ids_cte(flow_version_ids=flow_version_ids)

Expand All @@ -144,23 +154,17 @@ async def build_flow_artifacts_from_flow_versions(
.select_from(indexed_flow_version_ids_cte)
.join(
FlowVersion,
and_(
FlowVersion.id == indexed_flow_version_ids_cte.c.flow_version_id,
FlowVersion.user_id == user_id,
),
FlowVersion.id == indexed_flow_version_ids_cte.c.flow_version_id,
)
.join(
Flow,
and_(
Flow.id == FlowVersion.flow_id,
Flow.user_id == user_id,
),
Flow.id == FlowVersion.flow_id,
)
.join(
Deployment,
and_(
Deployment.id == deployment_db_id,
Deployment.user_id == user_id,
Deployment.user_id == deployment_user_id,
Deployment.project_id == Flow.folder_id,
),
)
Expand Down Expand Up @@ -199,14 +203,18 @@ async def build_flow_artifacts_from_flow_versions(
async def build_project_scoped_flow_artifacts_from_flow_versions(
*,
db,
user_id: UUID,
project_id: UUID,
reference_ids: Sequence[UUID | str],
authorized_flow_version_ids: frozenset[UUID],
) -> list[tuple[UUID, BaseFlowArtifact]]:
"""Resolve project-scoped flow version references preserving input order."""
flow_version_ids = parse_flow_version_reference_ids(reference_ids)
if not flow_version_ids:
return []
_require_authorized_flow_version_ids(
flow_version_ids=flow_version_ids,
authorized_flow_version_ids=authorized_flow_version_ids,
)
indexed_flow_version_ids_cte = _build_indexed_flow_version_ids_cte(flow_version_ids=flow_version_ids)

statement = (
Expand All @@ -222,16 +230,12 @@ async def build_project_scoped_flow_artifacts_from_flow_versions(
.select_from(indexed_flow_version_ids_cte)
.join(
FlowVersion,
and_(
FlowVersion.user_id == user_id,
FlowVersion.id == indexed_flow_version_ids_cte.c.flow_version_id,
),
FlowVersion.id == indexed_flow_version_ids_cte.c.flow_version_id,
)
.join(
Flow,
and_(
Flow.id == FlowVersion.flow_id,
Flow.user_id == user_id,
Flow.folder_id == project_id,
),
)
Expand Down Expand Up @@ -519,29 +523,6 @@ async def resolve_adapter_mapper_from_deployment(
)


async def resolve_project_id_for_deployment_create(
*,
payload: DeploymentCreateRequest,
user_id: UUID,
db: DbSession,
) -> UUID:
if payload.project_id is not None:
project = (
await db.exec(
select(Folder).where(
Folder.user_id == user_id,
Folder.id == payload.project_id,
)
)
).first()
if project is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found.")
return project.id

default_folder = await get_or_create_default_folder(db, user_id)
return default_folder.id


def resolve_snapshot_map_for_create(
*,
deployment_mapper: BaseDeploymentMapper,
Expand Down
Loading
Loading