Skip to content

Commit 776b9e6

Browse files
authored
fix: share policy refresh and owner permissions (langflow-ai#13830)
Fix share policy refresh and owner permissions
1 parent ffd119c commit 776b9e6

8 files changed

Lines changed: 225 additions & 23 deletions

File tree

src/backend/base/langflow/api/v1/authz_me.py

Lines changed: 74 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,10 +15,19 @@
1515

1616
from fastapi import APIRouter, HTTPException, status
1717
from pydantic import BaseModel, Field, field_validator
18+
from sqlmodel import col, select
19+
from sqlmodel.ext.asyncio.session import AsyncSession
1820

19-
from langflow.api.utils import CurrentActiveUser
21+
from langflow.api.utils import CurrentActiveUser, DbSessionReadOnly
2022
from langflow.services.auth.context import current_auth_context_for_authz
2123
from langflow.services.authorization.access_ceiling import filter_actions_by_external_access_ceiling
24+
from langflow.services.authorization.guards import should_apply_owner_override
25+
from langflow.services.database.models.deployment.model import Deployment
26+
from langflow.services.database.models.file.model import File as UserFile
27+
from langflow.services.database.models.flow.model import Flow
28+
from langflow.services.database.models.folder.model import Folder
29+
from langflow.services.database.models.knowledge_base.model import KnowledgeBaseRecord
30+
from langflow.services.database.models.variable.model import Variable
2231
from langflow.services.deps import get_authorization_service
2332

2433
router = APIRouter(prefix="/authz/me", tags=["Authorization"])
@@ -43,6 +52,15 @@
4352
# letting a client request `["read"] * 100000` to flood the enforcer.
4453
_MAX_ACTIONS = 10
4554

55+
_RESOURCE_OWNER_LOOKUPS: dict[str, tuple[type, str]] = {
56+
"flow": (Flow, "user_id"),
57+
"deployment": (Deployment, "user_id"),
58+
"project": (Folder, "user_id"),
59+
"knowledge_base": (KnowledgeBaseRecord, "user_id"),
60+
"variable": (Variable, "user_id"),
61+
"file": (UserFile, "user_id"),
62+
}
63+
4664

4765
class EffectivePermissionsRequest(BaseModel):
4866
"""Body for :func:`get_effective_permissions`."""
@@ -102,10 +120,57 @@ class EffectivePermissionsResponse(BaseModel):
102120
permissions: dict[UUID, list[str]]
103121

104122

123+
async def _owned_resource_ids(
124+
*,
125+
session: AsyncSession,
126+
resource_type: str,
127+
resource_ids: list[UUID],
128+
user_id: UUID,
129+
) -> set[UUID]:
130+
"""Return requested resource IDs owned by ``user_id``."""
131+
lookup = _RESOURCE_OWNER_LOOKUPS.get(resource_type)
132+
if lookup is None or not resource_ids:
133+
return set()
134+
model, owner_attr = lookup
135+
stmt = select(model.id).where(
136+
col(model.id).in_(resource_ids),
137+
getattr(model, owner_attr) == user_id,
138+
)
139+
return set((await session.exec(stmt)).all())
140+
141+
142+
async def _apply_owner_permissions(
143+
*,
144+
session: AsyncSession,
145+
permissions: dict[UUID, list[str]],
146+
resource_type: str,
147+
resource_ids: list[UUID],
148+
actions: tuple[str, ...],
149+
user_id: UUID,
150+
) -> dict[UUID, list[str]]:
151+
"""Mirror route guard owner override for the UI permission endpoint."""
152+
normalized = {resource_id: list(permissions.get(resource_id, [])) for resource_id in resource_ids}
153+
if not await should_apply_owner_override():
154+
return normalized
155+
156+
owned_ids = await _owned_resource_ids(
157+
session=session,
158+
resource_type=resource_type,
159+
resource_ids=resource_ids,
160+
user_id=user_id,
161+
)
162+
for resource_id in owned_ids:
163+
allowed = dict.fromkeys(normalized.get(resource_id, []))
164+
allowed.update(dict.fromkeys(actions))
165+
normalized[resource_id] = list(allowed)
166+
return normalized
167+
168+
105169
@router.post("/permissions", response_model=EffectivePermissionsResponse)
106170
async def get_effective_permissions(
107171
body: EffectivePermissionsRequest,
108172
current_user: CurrentActiveUser,
173+
session: DbSessionReadOnly,
109174
) -> EffectivePermissionsResponse:
110175
"""Return per-resource allowed actions for the current user.
111176
@@ -134,6 +199,14 @@ async def get_effective_permissions(
134199
"is_superuser": current_user.is_superuser,
135200
},
136201
)
202+
permissions = await _apply_owner_permissions(
203+
session=session,
204+
permissions=permissions,
205+
resource_type=body.resource_type,
206+
resource_ids=body.resource_ids,
207+
actions=actions,
208+
user_id=current_user.id,
209+
)
137210
permissions = {
138211
resource_id: filter_actions_by_external_access_ceiling(allowed_actions)
139212
for resource_id, allowed_actions in permissions.items()

src/backend/base/langflow/api/v1/authz_shares.py

Lines changed: 40 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -8,11 +8,13 @@
88

99
from fastapi import APIRouter, HTTPException, Query, status
1010
from lfx.log.logger import logger
11+
from lfx.services.authorization.base import BaseAuthorizationService
1112
from sqlmodel import select
1213

1314
from langflow.api.utils import CurrentActiveUser, DbSession
1415
from langflow.api.v1.schemas.authz_shares import ShareCreate, ShareRead, ShareUpdate
1516
from langflow.services.authorization import ShareAction, ensure_share_permission
17+
from langflow.services.authorization.invalidation import safe_invalidate_all, safe_invalidate_user
1618
from langflow.services.authorization.utils import audit_decision
1719
from langflow.services.database.models.auth import AuthzShare, AuthzTeamMember, SharePermissionLevel, ShareScope
1820
from langflow.services.database.models.deployment.model import Deployment
@@ -104,13 +106,32 @@ async def _user_can_see_share(
104106
)
105107

106108

107-
async def _invalidate_for_share(scope: str, target_id: UUID | None) -> None:
109+
async def _invalidate_for_share(scope: str, target_id: UUID | None, *, op: str = "share:write") -> None:
108110
"""Invalidate cached policy after a share write (user scope vs invalidate_all)."""
109111
authz = get_authorization_service()
110112
if scope == ShareScope.USER.value and target_id is not None:
111-
await authz.invalidate_user(target_id)
113+
await safe_invalidate_user(authz, target_id, op=op)
112114
else:
113-
await authz.invalidate_all()
115+
await safe_invalidate_all(authz, op=op)
116+
117+
118+
def _uses_base_sync_shares(authz: BaseAuthorizationService) -> bool:
119+
"""Return True when the service only has the OSS no-op sync_shares hook."""
120+
return getattr(type(authz), "sync_shares", None) is BaseAuthorizationService.sync_shares
121+
122+
123+
async def _refresh_policy_for_share(scope: str, target_id: UUID | None, *, op: str) -> None:
124+
"""Refresh share-derived policy after the share DB transaction is durable."""
125+
authz = get_authorization_service()
126+
sync_shares = getattr(authz, "sync_shares", None)
127+
if sync_shares is not None and not _uses_base_sync_shares(authz):
128+
try:
129+
await sync_shares()
130+
except Exception as exc: # noqa: BLE001 - plugin hooks are best-effort post-commit work
131+
logger.warning("sync_shares failed after %s; falling back to targeted invalidation: %s", op, exc)
132+
else:
133+
return
134+
await _invalidate_for_share(scope, target_id, op=op)
114135

115136

116137
async def _ensure_can_administer_share(
@@ -189,23 +210,26 @@ async def create_share(
189210
detail="Share could not be created: it may already exist or conflict with an existing share.",
190211
) from exc
191212
await session.refresh(row)
213+
response = ShareRead.model_validate(row, from_attributes=True)
214+
await session.commit()
192215

193-
# Invalidate policy cache for the share audience.
194-
await _invalidate_for_share(payload.scope, payload.target_id)
216+
# Refresh policy after commit so plugins using a separate DB connection see
217+
# the durable authz_share row instead of the pre-commit transaction state.
218+
await _refresh_policy_for_share(payload.scope, payload.target_id, op="share:create")
195219

196220
await audit_decision(
197221
user_id=current_user.id,
198222
action="share:create",
199223
obj=f"{payload.resource_type}:{payload.resource_id}",
200224
result="allow",
201225
details={
202-
"share_id": str(row.id),
226+
"share_id": str(response.id),
203227
"scope": payload.scope,
204228
"target_id": str(payload.target_id) if payload.target_id else None,
205229
"permission_level": payload.permission_level,
206230
},
207231
)
208-
return ShareRead.model_validate(row, from_attributes=True)
232+
return response
209233

210234

211235
_LIST_SHARES_MAX_LIMIT = 200
@@ -381,20 +405,22 @@ async def update_share(
381405
detail="Share could not be updated: it may conflict with an existing share.",
382406
) from exc
383407
await session.refresh(row)
408+
response = ShareRead.model_validate(row, from_attributes=True)
409+
await session.commit()
384410

385-
await _invalidate_for_share(row.scope, row.target_id)
411+
await _refresh_policy_for_share(response.scope, response.target_id, op="share:update")
386412

387413
await audit_decision(
388414
user_id=current_user.id,
389415
action="share:update",
390-
obj=f"{row.resource_type}:{row.resource_id}",
416+
obj=f"{response.resource_type}:{response.resource_id}",
391417
result="allow",
392418
details={
393-
"share_id": str(row.id),
394-
"permission_level": row.permission_level,
419+
"share_id": str(response.id),
420+
"permission_level": response.permission_level,
395421
},
396422
)
397-
return ShareRead.model_validate(row, from_attributes=True)
423+
return response
398424

399425

400426
@router.delete("/{share_id}", status_code=status.HTTP_204_NO_CONTENT)
@@ -429,8 +455,9 @@ async def delete_share(
429455
scope = row.scope
430456
await session.delete(row)
431457
await session.flush()
458+
await session.commit()
432459

433-
await _invalidate_for_share(scope, target_id)
460+
await _refresh_policy_for_share(scope, target_id, op="share:delete")
434461

435462
await audit_decision(
436463
user_id=current_user.id,

src/backend/base/langflow/services/authorization/__init__.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
ensure_project_permission,
2525
ensure_share_permission,
2626
ensure_variable_permission,
27+
should_apply_owner_override,
2728
)
2829
from langflow.services.authorization.listing import (
2930
filter_visible_resources,
@@ -57,5 +58,6 @@
5758
"requires_flow_permission",
5859
"requires_resource_permission",
5960
"restrict_to_owned_or_visible",
61+
"should_apply_owner_override",
6062
"visible_id_prefilter",
6163
]

src/backend/base/langflow/services/authorization/guards.py

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,11 @@ async def _api_key_scopes_require_plugin_enforcement() -> bool:
9696
return False
9797

9898

99+
async def should_apply_owner_override() -> bool:
100+
"""Return True when Langflow should apply the built-in owner override."""
101+
return not await _api_key_scopes_require_plugin_enforcement()
102+
103+
99104
def _coerce_action(
100105
act: DeploymentAction
101106
| FlowAction
@@ -224,11 +229,7 @@ async def _ensure_resource_permission(
224229
detail="External credentials do not allow this action",
225230
)
226231

227-
if (
228-
owner_id is not None
229-
and getattr(user, "id", None) == owner_id
230-
and not await _api_key_scopes_require_plugin_enforcement()
231-
):
232+
if owner_id is not None and getattr(user, "id", None) == owner_id and await should_apply_owner_override():
232233
await _audit.audit_decision(
233234
user_id=user.id,
234235
action=f"{resource_type}:{act_str}",

src/backend/base/langflow/services/authorization/utils.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
ensure_project_permission,
5252
ensure_share_permission,
5353
ensure_variable_permission,
54+
should_apply_owner_override,
5455
)
5556
from langflow.services.authorization.listing import (
5657
_default_resource_id_getter,
@@ -104,4 +105,5 @@ def permission_denied_to_http(exc):
104105
"get_authorization_service",
105106
"get_settings_service",
106107
"permission_denied_to_http",
108+
"should_apply_owner_override",
107109
]

src/backend/tests/unit/api/v1/test_authz_admin_routes.py

Lines changed: 47 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -995,12 +995,55 @@ async def test_me_permissions_returns_per_resource_actions(stub_authz):
995995
resource_type="flow",
996996
resource_ids=resource_ids,
997997
)
998-
result = await authz_me.get_effective_permissions(body=body, current_user=user)
998+
result = await authz_me.get_effective_permissions(body=body, current_user=user, session=_FakeAsyncSession())
999999
assert result.resource_type == "flow"
10001000
assert set(result.permissions[resource_ids[0]]) == {"read", "execute"}
10011001
assert "delete" in result.permissions[resource_ids[1]]
10021002

10031003

1004+
@pytest.mark.asyncio
1005+
async def test_me_permissions_adds_owner_override_actions(stub_authz):
1006+
"""Owners get the same built-in owner override in the UI gate as route guards."""
1007+
from langflow.api.v1 import authz_me
1008+
from langflow.api.v1.authz_me import EffectivePermissionsRequest
1009+
1010+
authz = stub_authz()
1011+
owned_flow_id = uuid4()
1012+
other_flow_id = uuid4()
1013+
authz.effective_perms_payload = {
1014+
owned_flow_id: [],
1015+
other_flow_id: [],
1016+
}
1017+
user = _make_user()
1018+
session = _FakeAsyncSession(exec_results=[[owned_flow_id]])
1019+
1020+
body = EffectivePermissionsRequest(
1021+
resource_type="flow",
1022+
resource_ids=[owned_flow_id, other_flow_id],
1023+
actions=["read", "write", "execute", "delete"],
1024+
)
1025+
result = await authz_me.get_effective_permissions(body=body, current_user=user, session=session)
1026+
1027+
assert set(result.permissions[owned_flow_id]) == {"read", "write", "execute", "delete"}
1028+
assert result.permissions[other_flow_id] == []
1029+
1030+
1031+
@pytest.mark.asyncio
1032+
async def test_me_permissions_returns_empty_entries_for_plugin_omissions(stub_authz):
1033+
"""A plugin omission is normalized to [] so the response shape stays stable."""
1034+
from langflow.api.v1 import authz_me
1035+
from langflow.api.v1.authz_me import EffectivePermissionsRequest
1036+
1037+
stub_authz()
1038+
flow_id = uuid4()
1039+
user = _make_user()
1040+
body = EffectivePermissionsRequest(resource_type="flow", resource_ids=[flow_id])
1041+
1042+
result = await authz_me.get_effective_permissions(body=body, current_user=user, session=_FakeAsyncSession())
1043+
1044+
assert result.permissions == {flow_id: []}
1045+
1046+
10041047
@pytest.mark.asyncio
10051048
async def test_me_permissions_caps_resource_ids_at_500(stub_authz):
10061049
from langflow.api.v1 import authz_me
@@ -1014,7 +1057,7 @@ async def test_me_permissions_caps_resource_ids_at_500(stub_authz):
10141057
resource_ids=[uuid4() for _ in range(501)],
10151058
)
10161059
with pytest.raises(HTTPException) as excinfo:
1017-
await authz_me.get_effective_permissions(body=body, current_user=user)
1060+
await authz_me.get_effective_permissions(body=body, current_user=user, session=_FakeAsyncSession())
10181061
assert excinfo.value.status_code == 400
10191062
assert "capped at 500" in excinfo.value.detail
10201063

@@ -1028,7 +1071,7 @@ async def test_me_permissions_empty_request_returns_empty(stub_authz):
10281071
user = _make_user()
10291072

10301073
body = EffectivePermissionsRequest(resource_type="flow", resource_ids=[])
1031-
result = await authz_me.get_effective_permissions(body=body, current_user=user)
1074+
result = await authz_me.get_effective_permissions(body=body, current_user=user, session=_FakeAsyncSession())
10321075
assert result.permissions == {}
10331076

10341077

@@ -1116,7 +1159,7 @@ async def _capture(**kwargs):
11161159
resource_ids=[rid],
11171160
actions=["READ", "read", " Write "],
11181161
)
1119-
await authz_me.get_effective_permissions(body=body, current_user=user)
1162+
await authz_me.get_effective_permissions(body=body, current_user=user, session=_FakeAsyncSession())
11201163
# Normalization happened at the model layer; handler sees the bounded set.
11211164
assert tuple(captured["actions"]) == ("read", "write")
11221165

0 commit comments

Comments
 (0)