Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e85c9a123d | ||
|
|
3516376d92 | ||
|
|
ffa984e271 | ||
|
|
7d13be6052 | ||
|
|
faff45db90 |
@@ -5,6 +5,40 @@ All notable changes to this project will be documented in this file.
|
||||
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
|
||||
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
|
||||
|
||||
## [1.5.0] - 2026-01-01
|
||||
|
||||
### Added
|
||||
|
||||
- **Groups Management** - Authentik group synchronization
|
||||
- `GET /auth/groups` - List all groups with search and pagination
|
||||
- `POST /auth/groups/sync-from-authentik` - Bulk sync groups from Authentik admin API
|
||||
- Database model and migration for groups table
|
||||
|
||||
## [1.4.6] - 2026-01-01
|
||||
|
||||
### Changed
|
||||
|
||||
- Code cleanup: move inline `re` import to top of auth/service.py
|
||||
|
||||
## [1.4.5] - 2026-01-01
|
||||
|
||||
### Fixed
|
||||
|
||||
- Separate httpx and SQLAlchemy async contexts in bulk sync (fixes greenlet error)
|
||||
|
||||
## [1.4.4] - 2026-01-01
|
||||
|
||||
### Fixed
|
||||
|
||||
- Use `uuid` field instead of `pk` for Authentik user sync (pk is integer, uuid is proper UUID)
|
||||
- Skip internal_service_account type users during bulk sync
|
||||
|
||||
## [1.4.3] - 2026-01-01
|
||||
|
||||
### Fixed
|
||||
|
||||
- Manually extract and send session cookies for Authentik flow auth (fixes cross-domain cookie handling)
|
||||
|
||||
## [1.4.2] - 2026-01-01
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
"""Create groups table
|
||||
|
||||
Revision ID: 002
|
||||
Revises: 001
|
||||
Create Date: 2026-01-01
|
||||
|
||||
Creates the groups table for syncing Authentik groups.
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
from sqlalchemy.dialects import postgresql
|
||||
|
||||
# revision identifiers, used by Alembic.
|
||||
revision: str = "002"
|
||||
down_revision: Union[str, None] = "001"
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
# Groups table
|
||||
op.create_table(
|
||||
"groups",
|
||||
sa.Column("id", postgresql.UUID(as_uuid=True), primary_key=True),
|
||||
sa.Column("authentik_id", postgresql.UUID(as_uuid=True), nullable=False, comment="Authentik group UUID"),
|
||||
sa.Column("name", sa.String(255), nullable=False),
|
||||
sa.Column("is_superuser", sa.Boolean(), nullable=False, server_default="false", comment="Whether members have superuser privileges"),
|
||||
sa.Column("parent_name", sa.String(255), nullable=True, comment="Parent group name for hierarchy"),
|
||||
sa.Column("member_count", sa.Integer(), nullable=False, server_default="0", comment="Number of users in this group"),
|
||||
sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False),
|
||||
sa.Column("synced_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False, comment="Last sync from Authentik"),
|
||||
)
|
||||
op.create_index("ix_groups_authentik_id", "groups", ["authentik_id"], unique=True)
|
||||
op.create_index("ix_groups_name", "groups", ["name"], unique=True)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_table("groups")
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "core-api"
|
||||
version = "1.4.2"
|
||||
version = "1.5.0"
|
||||
description = "Core Code API - Infrastructure management and tools API"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.12"
|
||||
|
||||
+64
-1
@@ -11,7 +11,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from src.controllers.base import BaseController
|
||||
from src.logging_config import get_logger
|
||||
from src.db import get_async_session
|
||||
from src.auth.schemas import AuthSyncRequest, AuthSyncResponse, UsersListResponse, BulkSyncResultSchema
|
||||
from src.auth.schemas import AuthSyncRequest, AuthSyncResponse, UsersListResponse, BulkSyncResultSchema, GroupsListResponse
|
||||
from src.auth.service import AuthService
|
||||
|
||||
logger = get_logger(__name__)
|
||||
@@ -152,6 +152,69 @@ class AuthController(BaseController):
|
||||
logger.error(f"Bulk sync failed: {e}")
|
||||
raise HTTPException(status_code=401, detail=str(e))
|
||||
|
||||
@router.get(
|
||||
"/groups",
|
||||
summary="List all groups",
|
||||
response_model=GroupsListResponse,
|
||||
responses={
|
||||
200: {"description": "List of groups"},
|
||||
},
|
||||
)
|
||||
async def list_groups(
|
||||
search: Optional[str] = Query(None, description="Search by group name"),
|
||||
offset: int = Query(0, ge=0, description="Number of records to skip"),
|
||||
limit: int = Query(50, ge=1, le=100, description="Maximum records to return"),
|
||||
session: AsyncSession = Depends(get_async_session),
|
||||
) -> GroupsListResponse:
|
||||
"""
|
||||
List all groups synced from Authentik
|
||||
|
||||
Returns paginated list of groups with their details.
|
||||
Supports search filtering by name.
|
||||
"""
|
||||
service = AuthService(session)
|
||||
items, total = await service.list_groups(
|
||||
search=search,
|
||||
offset=offset,
|
||||
limit=limit,
|
||||
)
|
||||
return GroupsListResponse(items=items, total=total)
|
||||
|
||||
@router.post(
|
||||
"/groups/sync-from-authentik",
|
||||
summary="Bulk sync groups from Authentik",
|
||||
response_model=BulkSyncResultSchema,
|
||||
responses={
|
||||
200: {"description": "Sync completed"},
|
||||
401: {"description": "Authentik API credentials invalid"},
|
||||
503: {"description": "Authentik service unavailable"},
|
||||
},
|
||||
)
|
||||
async def sync_groups_from_authentik(
|
||||
session: AsyncSession = Depends(get_async_session),
|
||||
) -> BulkSyncResultSchema:
|
||||
"""
|
||||
Fetch all groups from Authentik and sync to local database
|
||||
|
||||
This endpoint uses the Authentik admin API to fetch all groups
|
||||
and create/update them in the local database. Requires
|
||||
AUTHENTIK_USERNAME and AUTHENTIK_PASSWORD to be configured.
|
||||
|
||||
Use this to populate groups or to re-sync after changes in Authentik.
|
||||
"""
|
||||
service = AuthService(session)
|
||||
|
||||
try:
|
||||
result = await service.bulk_sync_groups_from_authentik()
|
||||
logger.info(
|
||||
f"Groups bulk sync completed: {result.created} created, "
|
||||
f"{result.updated} updated, {result.failed} failed"
|
||||
)
|
||||
return result
|
||||
except ValueError as e:
|
||||
logger.error(f"Groups bulk sync failed: {e}")
|
||||
raise HTTPException(status_code=401, detail=str(e))
|
||||
|
||||
@router.get(
|
||||
"/me",
|
||||
summary="Get current user profile",
|
||||
|
||||
@@ -108,3 +108,22 @@ class BulkSyncResultSchema(BaseSchema):
|
||||
failed: int = Field(..., description="Number of users that failed to sync")
|
||||
total_in_authentik: int = Field(..., description="Total users in Authentik")
|
||||
errors: list[str] = Field(default_factory=list, description="Error messages for failed syncs")
|
||||
|
||||
|
||||
class GroupListItemSchema(BaseSchema):
|
||||
"""Group item for list display"""
|
||||
|
||||
id: uuid.UUID = Field(..., description="Internal group ID")
|
||||
authentik_id: uuid.UUID = Field(..., description="Authentik group ID")
|
||||
name: str = Field(..., description="Group name")
|
||||
is_superuser: bool = Field(default=False, description="Whether group has superuser privileges")
|
||||
parent_name: Optional[str] = Field(None, description="Parent group name")
|
||||
member_count: int = Field(default=0, description="Number of users in this group")
|
||||
synced_at: datetime = Field(..., description="Last sync timestamp")
|
||||
|
||||
|
||||
class GroupsListResponse(BaseSchema):
|
||||
"""Response from GET /auth/groups"""
|
||||
|
||||
items: list[GroupListItemSchema] = Field(..., description="List of groups")
|
||||
total: int = Field(..., description="Total count of groups")
|
||||
|
||||
+265
-87
@@ -3,6 +3,7 @@ Authentication Service
|
||||
|
||||
Business logic for user synchronization from Authentik.
|
||||
"""
|
||||
import re
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
@@ -14,8 +15,8 @@ from sqlalchemy.orm import selectinload
|
||||
|
||||
from src.config import get_settings
|
||||
from src.logging_config import get_logger
|
||||
from src.db.models import User, Role, UserPreferences
|
||||
from src.auth.schemas import TokenInfoSchema, UserSchema, RoleSchema, UserPreferencesSchema, UserListItemSchema, BulkSyncResultSchema
|
||||
from src.db.models import User, Role, UserPreferences, Group
|
||||
from src.auth.schemas import TokenInfoSchema, UserSchema, RoleSchema, UserPreferencesSchema, UserListItemSchema, BulkSyncResultSchema, GroupListItemSchema
|
||||
|
||||
logger = get_logger(__name__)
|
||||
settings = get_settings()
|
||||
@@ -249,24 +250,29 @@ class AuthService:
|
||||
|
||||
return items, total
|
||||
|
||||
def _get_csrf_token(self, client: httpx.AsyncClient) -> str:
|
||||
"""Extract CSRF token from cookies"""
|
||||
for cookie in client.cookies.jar:
|
||||
if cookie.name == "authentik_csrf":
|
||||
return cookie.value
|
||||
def _extract_cookie(self, headers: httpx.Headers, cookie_name: str) -> str:
|
||||
"""Extract a specific cookie value from Set-Cookie headers"""
|
||||
for header in headers.get_list('set-cookie'):
|
||||
if header.startswith(f'{cookie_name}='):
|
||||
match = re.match(rf'{cookie_name}=([^;]+)', header)
|
||||
if match:
|
||||
return match.group(1)
|
||||
return ""
|
||||
|
||||
async def _authentik_session_login(self, client: httpx.AsyncClient) -> None:
|
||||
async def _authentik_session_login(self, client: httpx.AsyncClient) -> str:
|
||||
"""
|
||||
Authenticate with Authentik using the flow API to establish a session
|
||||
|
||||
Authentik's flow API requires:
|
||||
1. Cookie persistence between requests
|
||||
1. Cookie persistence between requests (manually handled due to domain restrictions)
|
||||
2. X-authentik-CSRF header set to the authentik_csrf cookie value
|
||||
3. Multi-stage flow handling (identification -> password -> done)
|
||||
|
||||
Args:
|
||||
client: httpx client with cookie persistence
|
||||
client: httpx client
|
||||
|
||||
Returns:
|
||||
Session cookie value for subsequent API calls
|
||||
|
||||
Raises:
|
||||
ValueError: If authentication fails
|
||||
@@ -278,55 +284,66 @@ class AuthService:
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
|
||||
logger.debug(f"Flow initial response: component={data.get('component')}, type={data.get('type')}")
|
||||
# Extract cookies manually from Set-Cookie headers (bypasses domain restrictions)
|
||||
session_cookie = self._extract_cookie(resp.headers, "authentik_session")
|
||||
csrf_cookie = self._extract_cookie(resp.headers, "authentik_csrf")
|
||||
|
||||
# Get CSRF token for subsequent requests
|
||||
csrf_token = self._get_csrf_token(client)
|
||||
logger.debug(f"CSRF token obtained: {bool(csrf_token)}")
|
||||
logger.debug(f"Flow initial: component={data.get('component')}, session={bool(session_cookie)}, csrf={bool(csrf_cookie)}")
|
||||
|
||||
# Build headers with CSRF token
|
||||
headers = {
|
||||
"Accept": "application/json",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
if csrf_token:
|
||||
headers["X-authentik-CSRF"] = csrf_token
|
||||
# Build headers with manual cookie and CSRF token
|
||||
def build_headers():
|
||||
hdrs = {
|
||||
"Accept": "application/json",
|
||||
"Content-Type": "application/json",
|
||||
"Cookie": f"authentik_session={session_cookie}",
|
||||
}
|
||||
if csrf_cookie:
|
||||
hdrs["Cookie"] += f"; authentik_csrf={csrf_cookie}"
|
||||
hdrs["X-authentik-CSRF"] = csrf_cookie
|
||||
return hdrs
|
||||
|
||||
# Step 2: Handle identification stage - submit username
|
||||
if data.get("component") == "ak-stage-identification":
|
||||
resp = await client.post(
|
||||
flow_url,
|
||||
json={"uid_field": settings.authentik_username},
|
||||
headers=headers,
|
||||
headers=build_headers(),
|
||||
)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
logger.debug(f"After username: component={data.get('component')}, type={data.get('type')}")
|
||||
|
||||
# Update CSRF token (might change between stages)
|
||||
csrf_token = self._get_csrf_token(client)
|
||||
if csrf_token:
|
||||
headers["X-authentik-CSRF"] = csrf_token
|
||||
# Update session cookie if new one received
|
||||
new_session = self._extract_cookie(resp.headers, "authentik_session")
|
||||
if new_session:
|
||||
session_cookie = new_session
|
||||
|
||||
logger.debug(f"After username: component={data.get('component')}")
|
||||
|
||||
# Step 3: Handle password stage if required
|
||||
if data.get("component") == "ak-stage-password":
|
||||
resp = await client.post(
|
||||
flow_url,
|
||||
json={"password": settings.authentik_password},
|
||||
headers=headers,
|
||||
headers=build_headers(),
|
||||
)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
logger.debug(f"After password: component={data.get('component')}, type={data.get('type')}")
|
||||
|
||||
# Update session cookie if new one received
|
||||
new_session = self._extract_cookie(resp.headers, "authentik_session")
|
||||
if new_session:
|
||||
session_cookie = new_session
|
||||
|
||||
logger.debug(f"After password: component={data.get('component')}")
|
||||
|
||||
# Check for access denied
|
||||
if data.get("component") == "ak-stage-access-denied":
|
||||
raise ValueError("Authentik authentication failed: access denied")
|
||||
|
||||
# Check for redirect (successful auth)
|
||||
if data.get("type") == "redirect" or data.get("to"):
|
||||
if data.get("component") == "xak-flow-redirect" or data.get("to"):
|
||||
logger.info("Successfully authenticated with Authentik via flow")
|
||||
return
|
||||
return session_cookie
|
||||
|
||||
# If we're still in identification stage, the username might be wrong
|
||||
if data.get("component") == "ak-stage-identification":
|
||||
@@ -334,6 +351,7 @@ class AuthService:
|
||||
raise ValueError(f"Authentication stuck at identification stage: {response_errors}")
|
||||
|
||||
logger.info(f"Authentik flow completed with component: {data.get('component')}")
|
||||
return session_cookie
|
||||
|
||||
async def bulk_sync_from_authentik(self) -> BulkSyncResultSchema:
|
||||
"""
|
||||
@@ -351,16 +369,21 @@ class AuthService:
|
||||
errors = []
|
||||
total_in_authentik = 0
|
||||
|
||||
# Step 1: Fetch all user data from Authentik API
|
||||
authentik_users = []
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30.0, follow_redirects=True) as client:
|
||||
# Authenticate with Authentik to get session
|
||||
await self._authentik_session_login(client)
|
||||
# Authenticate with Authentik to get session cookie
|
||||
session_cookie = await self._authentik_session_login(client)
|
||||
|
||||
# Fetch users from Authentik admin API using session
|
||||
# Fetch users from Authentik admin API using session cookie
|
||||
response = await client.get(
|
||||
f"{settings.authentik_url}/api/v3/core/users/",
|
||||
params={"page_size": 500},
|
||||
headers={"Accept": "application/json"},
|
||||
headers={
|
||||
"Accept": "application/json",
|
||||
"Cookie": f"authentik_session={session_cookie}",
|
||||
},
|
||||
)
|
||||
|
||||
if response.status_code == 401:
|
||||
@@ -372,72 +395,227 @@ class AuthService:
|
||||
authentik_users = data.get("results", [])
|
||||
total_in_authentik = data.get("pagination", {}).get("count", len(authentik_users))
|
||||
|
||||
for auth_user in authentik_users:
|
||||
try:
|
||||
# Skip service accounts and inactive users
|
||||
if auth_user.get("type") == "service_account":
|
||||
continue
|
||||
if not auth_user.get("is_active", True):
|
||||
continue
|
||||
except httpx.HTTPStatusError as e:
|
||||
raise ValueError(f"Authentik API error: {e.response.status_code}")
|
||||
except httpx.RequestError as e:
|
||||
raise ValueError(f"Failed to connect to Authentik: {str(e)}")
|
||||
|
||||
# Extract user data from Authentik
|
||||
authentik_id = uuid.UUID(auth_user["pk"])
|
||||
email = auth_user.get("email") or f"{auth_user['username']}@local"
|
||||
name = auth_user.get("name") or auth_user.get("username", "Unknown")
|
||||
avatar_url = auth_user.get("avatar")
|
||||
# Step 2: Sync users to database (outside of httpx context to avoid greenlet issues)
|
||||
for auth_user in authentik_users:
|
||||
try:
|
||||
# Skip service accounts and inactive users
|
||||
if auth_user.get("type") in ("service_account", "internal_service_account"):
|
||||
continue
|
||||
if not auth_user.get("is_active", True):
|
||||
continue
|
||||
|
||||
# Get user's groups for role mapping
|
||||
groups = []
|
||||
groups_summary = auth_user.get("groups_obj", [])
|
||||
for group in groups_summary:
|
||||
groups.append(group.get("name", ""))
|
||||
# Extract user data from Authentik
|
||||
authentik_id = uuid.UUID(auth_user["uuid"])
|
||||
email = auth_user.get("email") or f"{auth_user['username']}@local"
|
||||
name = auth_user.get("name") or auth_user.get("username", "Unknown")
|
||||
avatar_url = auth_user.get("avatar")
|
||||
|
||||
# Check if user exists
|
||||
stmt = select(User).where(User.authentik_id == authentik_id)
|
||||
result = await self.session.execute(stmt)
|
||||
user = result.scalar_one_or_none()
|
||||
# Get user's groups for role mapping
|
||||
groups = []
|
||||
groups_summary = auth_user.get("groups_obj", [])
|
||||
for group in groups_summary:
|
||||
groups.append(group.get("name", ""))
|
||||
|
||||
if user is None:
|
||||
# Create new user
|
||||
user = User(
|
||||
authentik_id=authentik_id,
|
||||
email=email,
|
||||
name=name,
|
||||
avatar_url=avatar_url,
|
||||
)
|
||||
self.session.add(user)
|
||||
await self.session.flush()
|
||||
# Check if user exists
|
||||
stmt = select(User).where(User.authentik_id == authentik_id)
|
||||
result = await self.session.execute(stmt)
|
||||
user = result.scalar_one_or_none()
|
||||
|
||||
# Create default preferences
|
||||
preferences = UserPreferences(user_id=user.id)
|
||||
self.session.add(preferences)
|
||||
created += 1
|
||||
logger.info(f"Created user from Authentik: {email}")
|
||||
else:
|
||||
# Update existing user
|
||||
user.email = email
|
||||
user.name = name
|
||||
user.avatar_url = avatar_url
|
||||
updated += 1
|
||||
logger.info(f"Updated user from Authentik: {email}")
|
||||
if user is None:
|
||||
# Create new user
|
||||
user = User(
|
||||
authentik_id=authentik_id,
|
||||
email=email,
|
||||
name=name,
|
||||
avatar_url=avatar_url,
|
||||
)
|
||||
self.session.add(user)
|
||||
await self.session.flush()
|
||||
|
||||
# Sync roles from groups
|
||||
await self.sync_roles(user, groups)
|
||||
# Create default preferences
|
||||
preferences = UserPreferences(user_id=user.id)
|
||||
self.session.add(preferences)
|
||||
created += 1
|
||||
logger.info(f"Created user from Authentik: {email}")
|
||||
else:
|
||||
# Update existing user
|
||||
user.email = email
|
||||
user.name = name
|
||||
user.avatar_url = avatar_url
|
||||
updated += 1
|
||||
logger.info(f"Updated user from Authentik: {email}")
|
||||
|
||||
except Exception as e:
|
||||
failed += 1
|
||||
error_msg = f"Failed to sync user {auth_user.get('username', 'unknown')}: {str(e)}"
|
||||
errors.append(error_msg)
|
||||
logger.warning(error_msg)
|
||||
# Sync roles from groups
|
||||
await self.sync_roles(user, groups)
|
||||
|
||||
# Commit all changes
|
||||
await self.session.commit()
|
||||
except Exception as e:
|
||||
failed += 1
|
||||
error_msg = f"Failed to sync user {auth_user.get('username', 'unknown')}: {str(e)}"
|
||||
errors.append(error_msg)
|
||||
logger.warning(error_msg)
|
||||
|
||||
# Commit all changes
|
||||
await self.session.commit()
|
||||
|
||||
return BulkSyncResultSchema(
|
||||
created=created,
|
||||
updated=updated,
|
||||
failed=failed,
|
||||
total_in_authentik=total_in_authentik,
|
||||
errors=errors,
|
||||
)
|
||||
|
||||
async def list_groups(
|
||||
self,
|
||||
search: Optional[str] = None,
|
||||
offset: int = 0,
|
||||
limit: int = 50,
|
||||
) -> tuple[list[GroupListItemSchema], int]:
|
||||
"""
|
||||
List all groups with optional search and pagination
|
||||
|
||||
Args:
|
||||
search: Optional search query (matches name)
|
||||
offset: Number of records to skip
|
||||
limit: Maximum number of records to return
|
||||
|
||||
Returns:
|
||||
Tuple of (list of group schemas, total count)
|
||||
"""
|
||||
from sqlalchemy import func
|
||||
|
||||
# Base query
|
||||
base_query = select(Group)
|
||||
|
||||
# Apply search filter if provided
|
||||
if search:
|
||||
search_filter = f"%{search}%"
|
||||
base_query = base_query.where(Group.name.ilike(search_filter))
|
||||
|
||||
# Get total count
|
||||
count_query = select(func.count()).select_from(base_query.subquery())
|
||||
total_result = await self.session.execute(count_query)
|
||||
total = total_result.scalar() or 0
|
||||
|
||||
# Apply pagination and ordering
|
||||
query = base_query.order_by(Group.name).offset(offset).limit(limit)
|
||||
result = await self.session.execute(query)
|
||||
groups = list(result.scalars().all())
|
||||
|
||||
# Convert to schemas
|
||||
items = [
|
||||
GroupListItemSchema(
|
||||
id=group.id,
|
||||
authentik_id=group.authentik_id,
|
||||
name=group.name,
|
||||
is_superuser=group.is_superuser,
|
||||
parent_name=group.parent_name,
|
||||
member_count=group.member_count,
|
||||
synced_at=group.synced_at,
|
||||
)
|
||||
for group in groups
|
||||
]
|
||||
|
||||
return items, total
|
||||
|
||||
async def bulk_sync_groups_from_authentik(self) -> BulkSyncResultSchema:
|
||||
"""
|
||||
Fetch all groups from Authentik admin API and sync to local database
|
||||
|
||||
Returns:
|
||||
BulkSyncResultSchema with counts of created/updated/failed groups
|
||||
"""
|
||||
if not settings.authentik_username or not settings.authentik_password:
|
||||
raise ValueError("AUTHENTIK_USERNAME and AUTHENTIK_PASSWORD must be configured")
|
||||
|
||||
created = 0
|
||||
updated = 0
|
||||
failed = 0
|
||||
errors = []
|
||||
total_in_authentik = 0
|
||||
|
||||
# Step 1: Fetch all group data from Authentik API
|
||||
authentik_groups = []
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30.0, follow_redirects=True) as client:
|
||||
# Authenticate with Authentik to get session cookie
|
||||
session_cookie = await self._authentik_session_login(client)
|
||||
|
||||
# Fetch groups from Authentik admin API using session cookie
|
||||
response = await client.get(
|
||||
f"{settings.authentik_url}/api/v3/core/groups/",
|
||||
params={"page_size": 500},
|
||||
headers={
|
||||
"Accept": "application/json",
|
||||
"Cookie": f"authentik_session={session_cookie}",
|
||||
},
|
||||
)
|
||||
|
||||
if response.status_code == 401:
|
||||
raise ValueError("Authentik API token is invalid or expired")
|
||||
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
|
||||
authentik_groups = data.get("results", [])
|
||||
total_in_authentik = data.get("pagination", {}).get("count", len(authentik_groups))
|
||||
|
||||
except httpx.HTTPStatusError as e:
|
||||
raise ValueError(f"Authentik API error: {e.response.status_code}")
|
||||
except httpx.RequestError as e:
|
||||
raise ValueError(f"Failed to connect to Authentik: {str(e)}")
|
||||
|
||||
# Step 2: Sync groups to database (outside of httpx context to avoid greenlet issues)
|
||||
for auth_group in authentik_groups:
|
||||
try:
|
||||
# Extract group data from Authentik
|
||||
authentik_id = uuid.UUID(auth_group["pk"])
|
||||
name = auth_group.get("name", "Unknown")
|
||||
is_superuser = auth_group.get("is_superuser", False)
|
||||
parent_name = auth_group.get("parent_name")
|
||||
# users field contains list of user PKs
|
||||
member_count = len(auth_group.get("users", []))
|
||||
|
||||
# Check if group exists
|
||||
stmt = select(Group).where(Group.authentik_id == authentik_id)
|
||||
result = await self.session.execute(stmt)
|
||||
group = result.scalar_one_or_none()
|
||||
|
||||
if group is None:
|
||||
# Create new group
|
||||
group = Group(
|
||||
authentik_id=authentik_id,
|
||||
name=name,
|
||||
is_superuser=is_superuser,
|
||||
parent_name=parent_name,
|
||||
member_count=member_count,
|
||||
)
|
||||
self.session.add(group)
|
||||
created += 1
|
||||
logger.info(f"Created group from Authentik: {name}")
|
||||
else:
|
||||
# Update existing group
|
||||
group.name = name
|
||||
group.is_superuser = is_superuser
|
||||
group.parent_name = parent_name
|
||||
group.member_count = member_count
|
||||
updated += 1
|
||||
logger.info(f"Updated group from Authentik: {name}")
|
||||
|
||||
except Exception as e:
|
||||
failed += 1
|
||||
error_msg = f"Failed to sync group {auth_group.get('name', 'unknown')}: {str(e)}"
|
||||
errors.append(error_msg)
|
||||
logger.warning(error_msg)
|
||||
|
||||
# Commit all changes
|
||||
await self.session.commit()
|
||||
|
||||
return BulkSyncResultSchema(
|
||||
created=created,
|
||||
updated=updated,
|
||||
|
||||
@@ -7,6 +7,7 @@ from src.db.models.user import User
|
||||
from src.db.models.role import Role, UserRole
|
||||
from src.db.models.user_preferences import UserPreferences
|
||||
from src.db.models.api_key import ApiKey
|
||||
from src.db.models.group import Group
|
||||
|
||||
__all__ = [
|
||||
"User",
|
||||
@@ -14,4 +15,5 @@ __all__ = [
|
||||
"UserRole",
|
||||
"UserPreferences",
|
||||
"ApiKey",
|
||||
"Group",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
"""
|
||||
Group Model
|
||||
|
||||
Represents groups synced from Authentik.
|
||||
Groups are used for access control and user organization.
|
||||
"""
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
from typing import TYPE_CHECKING, List
|
||||
|
||||
from sqlalchemy import String, Boolean, DateTime, func, Table, Column, ForeignKey
|
||||
from sqlalchemy.dialects.postgresql import UUID
|
||||
from sqlalchemy.orm import Mapped, mapped_column, relationship
|
||||
|
||||
from src.db.database import Base
|
||||
|
||||
|
||||
# Association table for User-Group many-to-many relationship
|
||||
user_groups = Table(
|
||||
"user_groups",
|
||||
Base.metadata,
|
||||
Column("user_id", UUID(as_uuid=True), ForeignKey("users.id", ondelete="CASCADE"), primary_key=True),
|
||||
Column("group_id", UUID(as_uuid=True), ForeignKey("groups.id", ondelete="CASCADE"), primary_key=True),
|
||||
)
|
||||
|
||||
|
||||
class Group(Base):
|
||||
"""
|
||||
Group model synced from Authentik
|
||||
|
||||
Groups are fetched from Authentik admin API and cached locally.
|
||||
They represent organizational units for access control.
|
||||
"""
|
||||
|
||||
__tablename__ = "groups"
|
||||
|
||||
id: Mapped[uuid.UUID] = mapped_column(
|
||||
UUID(as_uuid=True),
|
||||
primary_key=True,
|
||||
default=uuid.uuid4,
|
||||
)
|
||||
authentik_id: Mapped[uuid.UUID] = mapped_column(
|
||||
UUID(as_uuid=True),
|
||||
unique=True,
|
||||
nullable=False,
|
||||
index=True,
|
||||
comment="Authentik group UUID",
|
||||
)
|
||||
name: Mapped[str] = mapped_column(
|
||||
String(255),
|
||||
unique=True,
|
||||
nullable=False,
|
||||
index=True,
|
||||
)
|
||||
is_superuser: Mapped[bool] = mapped_column(
|
||||
Boolean,
|
||||
default=False,
|
||||
nullable=False,
|
||||
comment="Whether members of this group have superuser privileges",
|
||||
)
|
||||
parent_name: Mapped[str | None] = mapped_column(
|
||||
String(255),
|
||||
nullable=True,
|
||||
comment="Parent group name for hierarchy",
|
||||
)
|
||||
member_count: Mapped[int] = mapped_column(
|
||||
default=0,
|
||||
nullable=False,
|
||||
comment="Number of users in this group",
|
||||
)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
server_default=func.now(),
|
||||
nullable=False,
|
||||
)
|
||||
synced_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
server_default=func.now(),
|
||||
onupdate=func.now(),
|
||||
nullable=False,
|
||||
comment="Last sync from Authentik",
|
||||
)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"<Group {self.name}>"
|
||||
Reference in New Issue
Block a user