intermittent commit

This commit is contained in:
2026-07-14 13:22:09 +02:00
parent f2afcaa09c
commit ab07075a67
17 changed files with 1821 additions and 586 deletions

View File

@@ -41,15 +41,23 @@ from govoplan_access.backend.admin.service import (
update_system_role,
)
from govoplan_access.backend.api.v1.admin_common import (
_accounts_by_id,
_accounts_by_user_id,
_api_key_item,
_group_member_ids_by_group_id,
_group_summary,
_groups_by_user_id,
_http_admin_error,
_require_permission,
_require_system_role_assignment,
_resolve_tenant,
_roles_by_group_id,
_roles_by_user_id,
_role_summary,
_set_system_memberships,
_system_role_assignment_counts,
_system_account_item,
_tenant_role_assignment_counts,
_user_item,
)
from govoplan_access.backend.api.v1.admin_schemas import (
@@ -213,6 +221,45 @@ ACCESS_FUNCTION_DELEGATIONS_COLLECTION = "access.admin.function_delegations"
ACCESS_SYSTEM_ROLES_COLLECTION = "access.admin.system_roles"
ACCESS_SYSTEM_ACCOUNTS_COLLECTION = "access.admin.system_accounts"
ACCESS_API_KEYS_COLLECTION = "access.admin.api_keys"
ACCESS_FULL_CURSOR_PREFIX = "full:"
def _page_query(query, *, page: int, page_size: int):
total = query.order_by(None).count()
pages = max(1, (total + page_size - 1) // page_size)
items = query.offset((page - 1) * page_size).limit(page_size).all()
return items, {"total": total, "page": page, "page_size": page_size, "pages": pages}
def _encode_full_delta_cursor(scope: str, *, page: int, snapshot_sequence: int) -> str:
return f"{ACCESS_FULL_CURSOR_PREFIX}{scope}:{int(page)}:{int(snapshot_sequence)}"
def _decode_full_delta_cursor(value: str | None, *, scope: str) -> tuple[int, int] | None:
if not value or not value.startswith(ACCESS_FULL_CURSOR_PREFIX):
return None
parts = value.split(":", 3)
if len(parts) != 4 or parts[1] != scope:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid full snapshot cursor")
try:
page = int(parts[2])
snapshot_sequence = int(parts[3])
except ValueError as exc:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid full snapshot cursor") from exc
if page < 1 or snapshot_sequence < 0:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid full snapshot cursor")
return page, snapshot_sequence
def _full_delta_page(query, *, page: int, page_size: int, scope: str, snapshot_sequence: int):
items, pagination = _page_query(query, page=page, page_size=page_size)
has_more = page < pagination["pages"]
watermark = (
_encode_full_delta_cursor(scope, page=page + 1, snapshot_sequence=snapshot_sequence)
if has_more
else encode_sequence_watermark(snapshot_sequence)
)
return items, pagination, watermark, has_more
def _access_delta_watermark(session: Session, tenant_id: str | None, collections: tuple[str, ...]) -> str:
@@ -318,14 +365,35 @@ def _require_any_permission(principal: ApiPrincipal, *scopes: str) -> None:
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=f"Requires one of: {', '.join(scopes)}")
def _identity_item(session: Session, identity: Identity) -> IdentityAdminItem:
def _identity_account_links_by_identity(session: Session, identities: list[Identity]) -> dict[str, list[tuple[IdentityAccountLink, Account]]]:
identity_ids = [identity.id for identity in identities]
if not identity_ids:
return {}
grouped: dict[str, list[tuple[IdentityAccountLink, Account]]] = {}
rows = (
session.query(IdentityAccountLink, Account)
.join(Account, Account.id == IdentityAccountLink.account_id)
.filter(IdentityAccountLink.identity_id == identity.id)
.order_by(IdentityAccountLink.is_primary.desc(), Account.email.asc())
.filter(IdentityAccountLink.identity_id.in_(identity_ids))
.order_by(IdentityAccountLink.identity_id.asc(), IdentityAccountLink.is_primary.desc(), Account.email.asc())
.all()
)
for link, account in rows:
grouped.setdefault(link.identity_id, []).append((link, account))
return grouped
def _identity_item(session: Session, identity: Identity, *, account_links_by_identity: dict[str, list[tuple[IdentityAccountLink, Account]]] | None = None) -> IdentityAdminItem:
rows = (
account_links_by_identity.get(identity.id, [])
if account_links_by_identity is not None
else (
session.query(IdentityAccountLink, Account)
.join(Account, Account.id == IdentityAccountLink.account_id)
.filter(IdentityAccountLink.identity_id == identity.id)
.order_by(IdentityAccountLink.is_primary.desc(), Account.email.asc())
.all()
)
)
return IdentityAdminItem(
id=identity.id,
display_name=identity.display_name,
@@ -1073,12 +1141,16 @@ def approve_configuration_change_request_endpoint(
@router.get("/identities", response_model=IdentityListResponse)
def list_identities(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_any_scope("system:accounts:read", "access:account:read")),
):
del principal
identities = session.query(Identity).order_by(Identity.display_name.asc(), Identity.created_at.asc()).all()
return IdentityListResponse(identities=[_identity_item(session, item) for item in identities])
query = session.query(Identity).order_by(Identity.display_name.asc(), Identity.created_at.asc())
identities, pagination = _page_query(query, page=page, page_size=page_size)
account_links_by_identity = _identity_account_links_by_identity(session, identities)
return IdentityListResponse(identities=[_identity_item(session, item, account_links_by_identity=account_links_by_identity) for item in identities], **pagination)
@router.post("/identities", response_model=IdentityAdminItem, status_code=status.HTTP_201_CREATED)
@@ -1952,24 +2024,67 @@ def revoke_function_delegation(
return None
def _full_users_delta_response(session: Session, tenant: Tenant) -> UserListDeltaResponse:
users = session.query(User).filter(User.tenant_id == tenant.id).order_by(User.display_name.asc(), User.email.asc()).all()
def _user_items_for_response(session: Session, tenant: Tenant, users: list[User]) -> list[UserAdminItem]:
user_ids = [user.id for user in users]
accounts_by_id = _accounts_by_id(session, [user.account_id for user in users])
groups_by_user = _groups_by_user_id(session, tenant_id=tenant.id, user_ids=user_ids)
roles_by_user = _roles_by_user_id(session, tenant_id=tenant.id, user_ids=user_ids)
group_ids = sorted({group.id for groups in groups_by_user.values() for group in groups})
group_member_ids_by_group = _group_member_ids_by_group_id(session, tenant_id=tenant.id, group_ids=group_ids)
group_roles_by_group = _roles_by_group_id(session, tenant_id=tenant.id, group_ids=group_ids)
role_ids = sorted(
{
role.id
for roles in roles_by_user.values()
for role in roles
}
| {
role.id
for roles in group_roles_by_group.values()
for role in roles
}
)
tenant_role_counts = _tenant_role_assignment_counts(session, role_ids)
owner_ids = tenant_owner_user_ids(session, tenant.id)
idm_directory = _optional_idm_directory()
organization_directory = _optional_organization_directory()
return [
_user_item_for_response(
session,
user,
owner_ids=owner_ids,
idm_directory=idm_directory,
organization_directory=organization_directory,
accounts_by_id=accounts_by_id,
groups_by_user=groups_by_user,
roles_by_user=roles_by_user,
group_member_ids_by_group=group_member_ids_by_group,
group_roles_by_group=group_roles_by_group,
tenant_role_assignment_counts=tenant_role_counts,
)
for user in users
]
def _full_users_delta_response(session: Session, tenant: Tenant, *, cursor: tuple[int, int] | None = None, limit: int = 500) -> UserListDeltaResponse:
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=tenant.id, module_id=ACCESS_MODULE_ID, collections=(ACCESS_USERS_COLLECTION,))
page = cursor[0] if cursor is not None else 1
query = session.query(User).filter(User.tenant_id == tenant.id).order_by(User.display_name.asc(), User.email.asc(), User.id.asc())
users, pagination, watermark, has_more = _full_delta_page(query, page=page, page_size=limit, scope="users", snapshot_sequence=snapshot_sequence)
return UserListDeltaResponse(
users=[_user_item_for_response(session, user, owner_ids=owner_ids, idm_directory=idm_directory, organization_directory=organization_directory) for user in users],
users=_user_items_for_response(session, tenant, users),
deleted=[],
watermark=_access_delta_watermark(session, tenant.id, (ACCESS_USERS_COLLECTION,)),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _users_delta_response(session: Session, tenant: Tenant, *, since: str, limit: int) -> UserListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=tenant.id, collections=(ACCESS_USERS_COLLECTION,), since=since, limit=limit)
if entries is None:
return _full_users_delta_response(session, tenant)
return _full_users_delta_response(session, tenant, limit=limit)
changed_ids = _changed_ids(entries, "access_user")
visible = {
user.id: user
@@ -1978,16 +2093,13 @@ def _users_delta_response(session: Session, tenant: Tenant, *, since: str, limit
if changed_ids else []
)
}
owner_ids = tenant_owner_user_ids(session, tenant.id)
idm_directory = _optional_idm_directory()
organization_directory = _optional_organization_directory()
deleted = [
_delta_deleted_item(entry)
for entry in entries
if entry.resource_type == "access_user" and entry.resource_id not in visible
]
return UserListDeltaResponse(
users=[_user_item_for_response(session, user, owner_ids=owner_ids, idm_directory=idm_directory, organization_directory=organization_directory) for user in visible.values()],
users=_user_items_for_response(session, tenant, list(visible.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=tenant.id, collections=(ACCESS_USERS_COLLECTION,), entries=entries, has_more=has_more),
has_more=has_more,
@@ -2002,6 +2114,12 @@ def _user_item_for_response(
owner_ids: set[str] | None = None,
idm_directory: IdmDirectory | None = None,
organization_directory: OrganizationDirectory | None = None,
accounts_by_id: dict[str, Account] | None = None,
groups_by_user: dict[str, list[Group]] | None = None,
roles_by_user: dict[str, list[Role]] | None = None,
group_member_ids_by_group: dict[str, list[str]] | None = None,
group_roles_by_group: dict[str, list[Role]] | None = None,
tenant_role_assignment_counts: dict[str, tuple[int, int]] | None = None,
) -> UserAdminItem:
idm_assignments = _idm_assignments_for_user(idm_directory, user)
return _user_item(
@@ -2010,6 +2128,12 @@ def _user_item_for_response(
owner_ids=owner_ids,
idm_assignments=idm_assignments,
organization_directory=organization_directory,
accounts_by_id=accounts_by_id,
groups_by_user=groups_by_user,
roles_by_user=roles_by_user,
group_member_ids_by_group=group_member_ids_by_group,
group_roles_by_group=group_roles_by_group,
tenant_role_assignment_counts=tenant_role_assignment_counts,
)
@@ -2022,23 +2146,27 @@ def list_users_delta(
principal: ApiPrincipal = Depends(require_scope("admin:users:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
if since is None:
return _full_users_delta_response(session, tenant)
full_cursor = _decode_full_delta_cursor(since, scope="users")
if since is None or full_cursor is not None:
return _full_users_delta_response(session, tenant, cursor=full_cursor, limit=limit)
return _users_delta_response(session, tenant, since=since, limit=limit)
@router.get("/users", response_model=UserListResponse)
def list_users(
tenant_id: str | None = Query(default=None),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("admin:users:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
users = session.query(User).filter(User.tenant_id == tenant.id).order_by(User.display_name.asc(), User.email.asc()).all()
owner_ids = tenant_owner_user_ids(session, tenant.id)
idm_directory = _optional_idm_directory()
organization_directory = _optional_organization_directory()
return UserListResponse(users=[_user_item_for_response(session, user, owner_ids=owner_ids, idm_directory=idm_directory, organization_directory=organization_directory) for user in users])
query = session.query(User).filter(User.tenant_id == tenant.id).order_by(User.display_name.asc(), User.email.asc())
users, pagination = _page_query(query, page=page, page_size=page_size)
return UserListResponse(
users=_user_items_for_response(session, tenant, users),
**pagination,
)
@router.get("/users/{user_id}/access-explanation", response_model=UserAccessExplanationResponse)
@@ -2198,14 +2326,7 @@ def create_user(
)
@router.patch("/users/{user_id}", response_model=UserAdminItem)
def update_user(
user_id: str,
payload: UserUpdateRequest,
tenant_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
def _require_user_update_permissions(principal: ApiPrincipal, payload: UserUpdateRequest) -> None:
if "display_name" in payload.model_fields_set:
_require_permission(principal, "admin:users:update")
if payload.is_active is not None:
@@ -2214,10 +2335,16 @@ def update_user(
_require_permission(principal, "admin:groups:manage_members")
if payload.role_ids is not None:
_require_permission(principal, "admin:roles:assign")
tenant = _resolve_tenant(session, principal, tenant_id)
user = session.query(User).filter(User.id == user_id, User.tenant_id == tenant.id).one_or_none()
if user is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found")
def _apply_user_update(
session: Session,
*,
tenant: Tenant,
user: User,
payload: UserUpdateRequest,
principal: ApiPrincipal,
) -> tuple[set[str], set[str]]:
previous_group_ids = _current_user_group_ids(session, user) if payload.group_ids is not None else set()
previous_role_ids = _current_user_role_ids(session, user) if payload.role_ids is not None else set()
try:
@@ -2234,6 +2361,19 @@ def update_user(
assert_tenant_owner_exists(session, tenant.id)
except (AdminConflictError, AdminValidationError) as exc:
raise _http_admin_error(exc) from exc
return previous_group_ids, previous_role_ids
def _record_user_update_changes(
session: Session,
*,
tenant: Tenant,
user: User,
payload: UserUpdateRequest,
principal: ApiPrincipal,
previous_group_ids: set[str],
previous_role_ids: set[str],
) -> None:
_record_access_change(
session,
collection=ACCESS_USERS_COLLECTION,
@@ -2263,6 +2403,37 @@ def update_user(
tenant_id=tenant.id,
principal=principal,
)
@router.patch("/users/{user_id}", response_model=UserAdminItem)
def update_user(
user_id: str,
payload: UserUpdateRequest,
tenant_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
_require_user_update_permissions(principal, payload)
tenant = _resolve_tenant(session, principal, tenant_id)
user = session.query(User).filter(User.id == user_id, User.tenant_id == tenant.id).one_or_none()
if user is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found")
previous_group_ids, previous_role_ids = _apply_user_update(
session,
tenant=tenant,
user=user,
payload=payload,
principal=principal,
)
_record_user_update_changes(
session,
tenant=tenant,
user=user,
payload=payload,
principal=principal,
previous_group_ids=previous_group_ids,
previous_role_ids=previous_role_ids,
)
audit_event(
session,
tenant_id=tenant.id,
@@ -2281,21 +2452,43 @@ def update_user(
)
def _full_groups_delta_response(session: Session, tenant: Tenant) -> GroupListDeltaResponse:
groups = session.query(Group).filter(Group.tenant_id == tenant.id).order_by(Group.name.asc()).all()
def _group_summaries_for_response(session: Session, tenant: Tenant, groups: list[Group]) -> list[GroupSummary]:
group_ids = [group.id for group in groups]
member_ids_by_group = _group_member_ids_by_group_id(session, tenant_id=tenant.id, group_ids=group_ids)
roles_by_group = _roles_by_group_id(session, tenant_id=tenant.id, group_ids=group_ids)
role_ids = sorted({role.id for roles in roles_by_group.values() for role in roles})
tenant_role_counts = _tenant_role_assignment_counts(session, role_ids)
return [
_group_summary(
session,
group,
member_ids_by_group=member_ids_by_group,
roles_by_group=roles_by_group,
tenant_role_assignment_counts=tenant_role_counts,
)
for group in groups
]
def _full_groups_delta_response(session: Session, tenant: Tenant, *, cursor: tuple[int, int] | None = None, limit: int = 500) -> GroupListDeltaResponse:
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=tenant.id, module_id=ACCESS_MODULE_ID, collections=(ACCESS_GROUPS_COLLECTION,))
page = cursor[0] if cursor is not None else 1
query = session.query(Group).filter(Group.tenant_id == tenant.id).order_by(Group.name.asc(), Group.id.asc())
groups, pagination, watermark, has_more = _full_delta_page(query, page=page, page_size=limit, scope="groups", snapshot_sequence=snapshot_sequence)
return GroupListDeltaResponse(
groups=[_group_summary(session, group) for group in groups],
groups=_group_summaries_for_response(session, tenant, groups),
deleted=[],
watermark=_access_delta_watermark(session, tenant.id, (ACCESS_GROUPS_COLLECTION,)),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _groups_delta_response(session: Session, tenant: Tenant, *, since: str, limit: int) -> GroupListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=tenant.id, collections=(ACCESS_GROUPS_COLLECTION,), since=since, limit=limit)
if entries is None:
return _full_groups_delta_response(session, tenant)
return _full_groups_delta_response(session, tenant, limit=limit)
changed_ids = _changed_ids(entries, "access_group")
visible = {
group.id: group
@@ -2310,7 +2503,7 @@ def _groups_delta_response(session: Session, tenant: Tenant, *, since: str, limi
if entry.resource_type == "access_group" and entry.resource_id not in visible
]
return GroupListDeltaResponse(
groups=[_group_summary(session, group) for group in visible.values()],
groups=_group_summaries_for_response(session, tenant, list(visible.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=tenant.id, collections=(ACCESS_GROUPS_COLLECTION,), entries=entries, has_more=has_more),
has_more=has_more,
@@ -2327,20 +2520,27 @@ def list_groups_delta(
principal: ApiPrincipal = Depends(require_scope("admin:groups:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
if since is None:
return _full_groups_delta_response(session, tenant)
full_cursor = _decode_full_delta_cursor(since, scope="groups")
if since is None or full_cursor is not None:
return _full_groups_delta_response(session, tenant, cursor=full_cursor, limit=limit)
return _groups_delta_response(session, tenant, since=since, limit=limit)
@router.get("/groups", response_model=GroupListResponse)
def list_groups(
tenant_id: str | None = Query(default=None),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("admin:groups:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
groups = session.query(Group).filter(Group.tenant_id == tenant.id).order_by(Group.name.asc()).all()
return GroupListResponse(groups=[_group_summary(session, group) for group in groups])
query = session.query(Group).filter(Group.tenant_id == tenant.id).order_by(Group.name.asc())
groups, pagination = _page_query(query, page=page, page_size=page_size)
return GroupListResponse(
groups=_group_summaries_for_response(session, tenant, groups),
**pagination,
)
@router.post("/groups", response_model=GroupSummary, status_code=status.HTTP_201_CREATED)
@@ -2419,24 +2619,23 @@ def create_group(
return _group_summary(session, group)
@router.patch("/groups/{group_id}", response_model=GroupSummary)
def update_group(
group_id: str,
payload: GroupUpdateRequest,
tenant_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
def _require_group_update_permissions(principal: ApiPrincipal, payload: GroupUpdateRequest) -> None:
if any(field in payload.model_fields_set for field in ("name", "description", "is_active")):
_require_permission(principal, "admin:groups:write")
if payload.member_ids is not None:
_require_permission(principal, "admin:groups:manage_members")
if payload.role_ids is not None:
_require_permission(principal, "admin:roles:assign")
tenant = _resolve_tenant(session, principal, tenant_id)
group = session.query(Group).filter(Group.id == group_id, Group.tenant_id == tenant.id).one_or_none()
if group is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Group not found")
def _apply_group_update(
session: Session,
*,
tenant: Tenant,
group: Group,
payload: GroupUpdateRequest,
principal: ApiPrincipal,
) -> tuple[set[str], set[str]]:
previous_member_ids = _current_group_member_ids(session, group) if payload.member_ids is not None else set()
previous_role_ids = _current_group_role_ids(session, group) if payload.role_ids is not None else set()
try:
@@ -2456,6 +2655,19 @@ def update_group(
assert_tenant_owner_exists(session, tenant.id)
except (AdminConflictError, AdminValidationError) as exc:
raise _http_admin_error(exc) from exc
return previous_member_ids, previous_role_ids
def _record_group_update_changes(
session: Session,
*,
tenant: Tenant,
group: Group,
payload: GroupUpdateRequest,
principal: ApiPrincipal,
previous_member_ids: set[str],
previous_role_ids: set[str],
) -> None:
_record_access_change(
session,
collection=ACCESS_GROUPS_COLLECTION,
@@ -2485,6 +2697,37 @@ def update_group(
tenant_id=tenant.id,
principal=principal,
)
@router.patch("/groups/{group_id}", response_model=GroupSummary)
def update_group(
group_id: str,
payload: GroupUpdateRequest,
tenant_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
_require_group_update_permissions(principal, payload)
tenant = _resolve_tenant(session, principal, tenant_id)
group = session.query(Group).filter(Group.id == group_id, Group.tenant_id == tenant.id).one_or_none()
if group is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Group not found")
previous_member_ids, previous_role_ids = _apply_group_update(
session,
tenant=tenant,
group=group,
payload=payload,
principal=principal,
)
_record_group_update_changes(
session,
tenant=tenant,
group=group,
payload=payload,
principal=principal,
previous_member_ids=previous_member_ids,
previous_role_ids=previous_role_ids,
)
audit_event(
session,
tenant_id=tenant.id,
@@ -2498,23 +2741,32 @@ def update_group(
return _group_summary(session, group)
def _full_roles_delta_response(session: Session, tenant: Tenant) -> RoleListDeltaResponse:
def _tenant_role_summaries_for_response(session: Session, roles: list[Role]) -> list[RoleSummary]:
tenant_role_counts = _tenant_role_assignment_counts(session, [role.id for role in roles])
return [_role_summary(session, role, tenant_role_assignment_counts=tenant_role_counts) for role in roles]
def _full_roles_delta_response(session: Session, tenant: Tenant, *, cursor: tuple[int, int] | None = None, limit: int = 500) -> RoleListDeltaResponse:
ensure_default_roles(session, tenant)
session.commit()
roles = session.query(Role).filter(Role.tenant_id == tenant.id).order_by(Role.is_builtin.desc(), Role.name.asc()).all()
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=tenant.id, module_id=ACCESS_MODULE_ID, collections=(ACCESS_ROLES_COLLECTION,))
page = cursor[0] if cursor is not None else 1
query = session.query(Role).filter(Role.tenant_id == tenant.id).order_by(Role.is_builtin.desc(), Role.name.asc(), Role.id.asc())
roles, pagination, watermark, has_more = _full_delta_page(query, page=page, page_size=limit, scope="roles", snapshot_sequence=snapshot_sequence)
return RoleListDeltaResponse(
roles=[_role_summary(session, role) for role in roles],
roles=_tenant_role_summaries_for_response(session, roles),
deleted=[],
watermark=_access_delta_watermark(session, tenant.id, (ACCESS_ROLES_COLLECTION,)),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _roles_delta_response(session: Session, tenant: Tenant, *, since: str, limit: int) -> RoleListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=tenant.id, collections=(ACCESS_ROLES_COLLECTION,), since=since, limit=limit)
if entries is None:
return _full_roles_delta_response(session, tenant)
return _full_roles_delta_response(session, tenant, limit=limit)
changed_ids = _changed_ids(entries, "access_role")
visible = {
role.id: role
@@ -2529,7 +2781,7 @@ def _roles_delta_response(session: Session, tenant: Tenant, *, since: str, limit
if entry.resource_type == "access_role" and entry.resource_id not in visible
]
return RoleListDeltaResponse(
roles=[_role_summary(session, role) for role in visible.values()],
roles=_tenant_role_summaries_for_response(session, list(visible.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=tenant.id, collections=(ACCESS_ROLES_COLLECTION,), entries=entries, has_more=has_more),
has_more=has_more,
@@ -2546,22 +2798,29 @@ def list_roles_delta(
principal: ApiPrincipal = Depends(require_scope("admin:roles:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
if since is None:
return _full_roles_delta_response(session, tenant)
full_cursor = _decode_full_delta_cursor(since, scope="roles")
if since is None or full_cursor is not None:
return _full_roles_delta_response(session, tenant, cursor=full_cursor, limit=limit)
return _roles_delta_response(session, tenant, since=since, limit=limit)
@router.get("/roles", response_model=RoleListResponse)
def list_roles(
tenant_id: str | None = Query(default=None),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("admin:roles:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
ensure_default_roles(session, tenant)
session.commit()
roles = session.query(Role).filter(Role.tenant_id == tenant.id).order_by(Role.is_builtin.desc(), Role.name.asc()).all()
return RoleListResponse(roles=[_role_summary(session, role) for role in roles])
query = session.query(Role).filter(Role.tenant_id == tenant.id).order_by(Role.is_builtin.desc(), Role.name.asc())
roles, pagination = _page_query(query, page=page, page_size=page_size)
return RoleListResponse(
roles=_tenant_role_summaries_for_response(session, roles),
**pagination,
)
@router.post("/roles", response_model=RoleSummary, status_code=status.HTTP_201_CREATED)
@@ -2718,23 +2977,32 @@ def delete_role(
return None
def _full_system_roles_delta_response(session: Session) -> RoleListDeltaResponse:
def _system_role_summaries_for_response(session: Session, roles: list[Role]) -> list[RoleSummary]:
system_role_counts = _system_role_assignment_counts(session, [role.id for role in roles])
return [_role_summary(session, role, system_role_assignment_counts=system_role_counts) for role in roles]
def _full_system_roles_delta_response(session: Session, *, cursor: tuple[int, int] | None = None, limit: int = 500) -> RoleListDeltaResponse:
ensure_default_roles(session, None)
session.commit()
roles = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc()).all()
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=None, module_id=ACCESS_MODULE_ID, collections=(ACCESS_SYSTEM_ROLES_COLLECTION,))
page = cursor[0] if cursor is not None else 1
query = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc(), Role.id.asc())
roles, pagination, watermark, has_more = _full_delta_page(query, page=page, page_size=limit, scope="system-roles", snapshot_sequence=snapshot_sequence)
return RoleListDeltaResponse(
roles=[_role_summary(session, role) for role in roles],
roles=_system_role_summaries_for_response(session, roles),
deleted=[],
watermark=_access_delta_watermark(session, None, (ACCESS_SYSTEM_ROLES_COLLECTION,)),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _system_roles_delta_response(session: Session, *, since: str, limit: int) -> RoleListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=None, collections=(ACCESS_SYSTEM_ROLES_COLLECTION,), since=since, limit=limit)
if entries is None:
return _full_system_roles_delta_response(session)
return _full_system_roles_delta_response(session, limit=limit)
changed_ids = _changed_ids(entries, "access_system_role")
visible = {
role.id: role
@@ -2749,7 +3017,7 @@ def _system_roles_delta_response(session: Session, *, since: str, limit: int) ->
if entry.resource_type == "access_system_role" and entry.resource_id not in visible
]
return RoleListDeltaResponse(
roles=[_role_summary(session, role) for role in visible.values()],
roles=_system_role_summaries_for_response(session, list(visible.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=None, collections=(ACCESS_SYSTEM_ROLES_COLLECTION,), entries=entries, has_more=has_more),
has_more=has_more,
@@ -2765,20 +3033,27 @@ def list_system_roles_delta(
principal: ApiPrincipal = Depends(require_any_scope("system:roles:read", "system:access:read")),
):
del principal
if since is None:
return _full_system_roles_delta_response(session)
full_cursor = _decode_full_delta_cursor(since, scope="system-roles")
if since is None or full_cursor is not None:
return _full_system_roles_delta_response(session, cursor=full_cursor, limit=limit)
return _system_roles_delta_response(session, since=since, limit=limit)
@router.get("/system/roles", response_model=RoleListResponse)
def list_system_roles(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_any_scope("system:roles:read", "system:access:read")),
):
ensure_default_roles(session, None)
session.commit()
roles = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc()).all()
return RoleListResponse(roles=[_role_summary(session, role) for role in roles])
query = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc())
roles, pagination = _page_query(query, page=page, page_size=page_size)
return RoleListResponse(
roles=_system_role_summaries_for_response(session, roles),
**pagination,
)
@router.post("/system/roles", response_model=RoleSummary, status_code=status.HTTP_201_CREATED)
@@ -2901,25 +3176,29 @@ def delete_system_role_endpoint(
_SYSTEM_ACCOUNTS_DELTA_COLLECTIONS = (ACCESS_SYSTEM_ACCOUNTS_COLLECTION, ACCESS_SYSTEM_ROLES_COLLECTION)
def _full_system_accounts_delta_response(session: Session) -> SystemAccountListDeltaResponse:
def _full_system_accounts_delta_response(session: Session, *, cursor: tuple[int, int] | None = None, limit: int = 500) -> SystemAccountListDeltaResponse:
ensure_default_roles(session, None)
session.commit()
system_roles = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc()).all()
accounts = session.query(Account).order_by(Account.email.asc()).all()
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=None, module_id=ACCESS_MODULE_ID, collections=_SYSTEM_ACCOUNTS_DELTA_COLLECTIONS)
page = cursor[0] if cursor is not None else 1
query = session.query(Account).order_by(Account.email.asc(), Account.id.asc())
accounts, pagination, watermark, has_more = _full_delta_page(query, page=page, page_size=limit, scope="system-accounts", snapshot_sequence=snapshot_sequence)
return SystemAccountListDeltaResponse(
accounts=[_system_account_item(session, account) for account in accounts],
roles=[_role_summary(session, role) for role in system_roles],
roles=_system_role_summaries_for_response(session, system_roles),
deleted=[],
watermark=_access_delta_watermark(session, None, _SYSTEM_ACCOUNTS_DELTA_COLLECTIONS),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _system_accounts_delta_response(session: Session, *, since: str, limit: int) -> SystemAccountListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=None, collections=_SYSTEM_ACCOUNTS_DELTA_COLLECTIONS, since=since, limit=limit)
if entries is None:
return _full_system_accounts_delta_response(session)
return _full_system_accounts_delta_response(session, limit=limit)
account_ids = _changed_ids(entries, "access_system_account")
role_ids = _changed_ids(entries, "access_system_role")
visible_accounts = {
@@ -2946,7 +3225,7 @@ def _system_accounts_delta_response(session: Session, *, since: str, limit: int)
]
return SystemAccountListDeltaResponse(
accounts=[_system_account_item(session, account) for account in visible_accounts.values()],
roles=[_role_summary(session, role) for role in visible_roles.values()],
roles=_system_role_summaries_for_response(session, list(visible_roles.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=None, collections=_SYSTEM_ACCOUNTS_DELTA_COLLECTIONS, entries=entries, has_more=has_more),
has_more=has_more,
@@ -2962,42 +3241,58 @@ def list_system_accounts_delta(
principal: ApiPrincipal = Depends(require_any_scope("system:accounts:read", "system:access:read")),
):
del principal
if since is None:
return _full_system_accounts_delta_response(session)
full_cursor = _decode_full_delta_cursor(since, scope="system-accounts")
if since is None or full_cursor is not None:
return _full_system_accounts_delta_response(session, cursor=full_cursor, limit=limit)
return _system_accounts_delta_response(session, since=since, limit=limit)
@router.get("/system/accounts", response_model=SystemAccountListResponse)
def list_system_accounts(
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_any_scope("system:accounts:read", "system:access:read")),
):
ensure_default_roles(session, None)
session.commit()
system_roles = session.query(Role).filter(Role.tenant_id.is_(None)).order_by(Role.name.asc()).all()
accounts = session.query(Account).order_by(Account.email.asc()).all()
query = session.query(Account).order_by(Account.email.asc())
accounts, pagination = _page_query(query, page=page, page_size=page_size)
return SystemAccountListResponse(
accounts=[_system_account_item(session, account) for account in accounts],
roles=[_role_summary(session, role) for role in system_roles],
roles=_system_role_summaries_for_response(session, system_roles),
**pagination,
)
@router.patch("/system/accounts/{account_id}", response_model=SystemAccountItem)
def update_system_account(
account_id: str,
payload: SystemAccountUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
def _require_system_account_update_permissions(principal: ApiPrincipal, payload: SystemAccountUpdateRequest) -> None:
if "display_name" in payload.model_fields_set and not has_scope(principal, "system:accounts:update"):
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: system:accounts:update")
if payload.is_active is not None and not has_scope(principal, "system:accounts:suspend"):
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: system:accounts:suspend")
if payload.role_ids is not None:
_require_system_role_assignment(principal)
account = session.get(Account, account_id)
if account is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Account not found")
def _active_tenant_ids_for_account(session: Session, account: Account) -> list[str]:
return [
row[0]
for row in session.query(User.tenant_id)
.join(Tenant, Tenant.id == User.tenant_id)
.filter(User.account_id == account.id, User.is_active.is_(True), Tenant.is_active.is_(True))
.distinct()
.all()
]
def _apply_system_account_update(
session: Session,
*,
account: Account,
payload: SystemAccountUpdateRequest,
principal: ApiPrincipal,
) -> set[str]:
previous_role_ids = _current_system_role_ids(session, account) if payload.role_ids is not None else set()
if "display_name" in payload.model_fields_set:
account.display_name = payload.display_name.strip() if payload.display_name else None
@@ -3013,18 +3308,21 @@ def update_system_account(
session.flush()
assert_system_owner_exists(session)
if not account.is_active:
tenant_ids = [
row[0]
for row in session.query(User.tenant_id)
.join(Tenant, Tenant.id == User.tenant_id)
.filter(User.account_id == account.id, User.is_active.is_(True), Tenant.is_active.is_(True))
.distinct()
.all()
]
for tenant_id in tenant_ids:
for tenant_id in _active_tenant_ids_for_account(session, account):
assert_tenant_owner_exists(session, tenant_id)
except (AdminConflictError, AdminValidationError) as exc:
raise _http_admin_error(exc) from exc
return previous_role_ids
def _record_system_account_update_changes(
session: Session,
*,
account: Account,
payload: SystemAccountUpdateRequest,
principal: ApiPrincipal,
previous_role_ids: set[str],
) -> None:
_record_access_change(
session,
collection=ACCESS_SYSTEM_ACCOUNTS_COLLECTION,
@@ -3044,6 +3342,27 @@ def update_system_account(
tenant_id=None,
principal=principal,
)
@router.patch("/system/accounts/{account_id}", response_model=SystemAccountItem)
def update_system_account(
account_id: str,
payload: SystemAccountUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(get_api_principal),
):
_require_system_account_update_permissions(principal, payload)
account = session.get(Account, account_id)
if account is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Account not found")
previous_role_ids = _apply_system_account_update(session, account=account, payload=payload, principal=principal)
_record_system_account_update_changes(
session,
account=account,
payload=payload,
principal=principal,
previous_role_ids=previous_role_ids,
)
audit_from_principal(
session,
principal,
@@ -3266,24 +3585,38 @@ def update_system_account_memberships(
return _system_account_item(session, account)
def _full_api_keys_delta_response(session: Session, tenant: Tenant, *, include_revoked: bool) -> ApiKeyListDeltaResponse:
def _api_key_items_for_response(session: Session, keys: list[ApiKey]) -> list[ApiKeyAdminItem]:
accounts_by_user_id = _accounts_by_user_id(session, [item.user_id for item in keys])
return [_api_key_item(session, item, accounts_by_user_id=accounts_by_user_id) for item in keys]
def _full_api_keys_delta_response(session: Session, tenant: Tenant, *, include_revoked: bool, cursor: tuple[int, int] | None = None, limit: int = 500) -> ApiKeyListDeltaResponse:
query = session.query(ApiKey).filter(ApiKey.tenant_id == tenant.id)
if not include_revoked:
query = query.filter(ApiKey.revoked_at.is_(None))
keys = query.order_by(ApiKey.created_at.desc()).all()
snapshot_sequence = cursor[1] if cursor is not None else max_sequence_id(session, tenant_id=tenant.id, module_id=ACCESS_MODULE_ID, collections=(ACCESS_API_KEYS_COLLECTION,))
page = cursor[0] if cursor is not None else 1
keys, pagination, watermark, has_more = _full_delta_page(
query.order_by(ApiKey.created_at.desc(), ApiKey.id.asc()),
page=page,
page_size=limit,
scope="api-keys",
snapshot_sequence=snapshot_sequence,
)
return ApiKeyListDeltaResponse(
api_keys=[_api_key_item(session, item) for item in keys],
api_keys=_api_key_items_for_response(session, keys),
deleted=[],
watermark=_access_delta_watermark(session, tenant.id, (ACCESS_API_KEYS_COLLECTION,)),
has_more=False,
watermark=watermark,
has_more=has_more,
full=True,
**pagination,
)
def _api_keys_delta_response(session: Session, tenant: Tenant, *, include_revoked: bool, since: str, limit: int) -> ApiKeyListDeltaResponse:
entries, has_more = _access_delta_entries(session, tenant_id=tenant.id, collections=(ACCESS_API_KEYS_COLLECTION,), since=since, limit=limit)
if entries is None:
return _full_api_keys_delta_response(session, tenant, include_revoked=include_revoked)
return _full_api_keys_delta_response(session, tenant, include_revoked=include_revoked, limit=limit)
changed_ids = _changed_ids(entries, "access_api_key")
query = session.query(ApiKey).filter(ApiKey.tenant_id == tenant.id)
if not include_revoked:
@@ -3301,7 +3634,7 @@ def _api_keys_delta_response(session: Session, tenant: Tenant, *, include_revoke
if entry.resource_type == "access_api_key" and entry.resource_id not in visible
]
return ApiKeyListDeltaResponse(
api_keys=[_api_key_item(session, item) for item in visible.values()],
api_keys=_api_key_items_for_response(session, list(visible.values())),
deleted=deleted,
watermark=_access_delta_response_watermark(session, tenant_id=tenant.id, collections=(ACCESS_API_KEYS_COLLECTION,), entries=entries, has_more=has_more),
has_more=has_more,
@@ -3319,8 +3652,9 @@ def list_api_keys_delta(
principal: ApiPrincipal = Depends(require_scope("admin:api_keys:read")),
):
tenant = _resolve_tenant(session, principal, tenant_id)
if since is None:
return _full_api_keys_delta_response(session, tenant, include_revoked=include_revoked)
full_cursor = _decode_full_delta_cursor(since, scope="api-keys")
if since is None or full_cursor is not None:
return _full_api_keys_delta_response(session, tenant, include_revoked=include_revoked, cursor=full_cursor, limit=limit)
return _api_keys_delta_response(session, tenant, include_revoked=include_revoked, since=since, limit=limit)
@@ -3328,6 +3662,8 @@ def list_api_keys_delta(
def list_api_keys(
tenant_id: str | None = Query(default=None),
include_revoked: bool = Query(default=False),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("admin:api_keys:read")),
):
@@ -3335,8 +3671,8 @@ def list_api_keys(
query = session.query(ApiKey).filter(ApiKey.tenant_id == tenant.id)
if not include_revoked:
query = query.filter(ApiKey.revoked_at.is_(None))
keys = query.order_by(ApiKey.created_at.desc()).all()
return ApiKeyListResponse(api_keys=[_api_key_item(session, item) for item in keys])
keys, pagination = _page_query(query.order_by(ApiKey.created_at.desc()), page=page, page_size=page_size)
return ApiKeyListResponse(api_keys=_api_key_items_for_response(session, keys), **pagination)
@router.post("/api-keys", response_model=AdminApiKeyCreateResponse, status_code=status.HTTP_201_CREATED)