From a6e7a2e75ca28d5bf7486292b9fb31186042f5e5 Mon Sep 17 00:00:00 2001 From: Jiang Date: Fri, 12 Jun 2026 15:08:37 +0800 Subject: [PATCH] feat(admin): add project metadata config --- .env.example | 8 +- AUTHENTICATION_AND_USER_MANAGEMENT.md | 95 +++ app/api/v1/endpoints/admin_metadata.py | 804 ++++++++++++++++++ app/api/v1/router.py | 4 + app/auth/metadata_dependencies.py | 50 +- app/domain/schemas/admin_metadata.py | 156 ++++ app/infra/db/dynamic_manager.py | 4 +- .../repositories/metadata_repository.py | 357 +++++++- .../sql/004_metadata_auth_management.sql | 62 ++ .../005_metadata_project_configuration.sql | 78 ++ scripts/migrate_local_users_to_metadata.py | 63 ++ tests/api/test_admin_metadata_endpoints.py | 646 ++++++++++++++ tests/auth/test_metadata_dependencies.py | 83 ++ tests/unit/test_dynamic_manager.py | 12 + .../test_metadata_repository_dsn_decrypt.py | 48 ++ 15 files changed, 2460 insertions(+), 10 deletions(-) create mode 100644 AUTHENTICATION_AND_USER_MANAGEMENT.md create mode 100644 app/api/v1/endpoints/admin_metadata.py create mode 100644 app/domain/schemas/admin_metadata.py create mode 100644 resources/sql/004_metadata_auth_management.sql create mode 100644 resources/sql/005_metadata_project_configuration.sql create mode 100644 scripts/migrate_local_users_to_metadata.py create mode 100644 tests/api/test_admin_metadata_endpoints.py create mode 100644 tests/auth/test_metadata_dependencies.py create mode 100644 tests/unit/test_dynamic_manager.py diff --git a/.env.example b/.env.example index 9b3b90c..6822043 100644 --- a/.env.example +++ b/.env.example @@ -11,10 +11,12 @@ NETWORK_NAME="tjwater" # 生成方式: openssl rand -hex 32 SECRET_KEY=your-secret-key-here-change-in-production-use-openssl-rand-hex-32 -# 数据加密密钥 - 用于敏感数据加密 +# 数据加密密钥 - Fernet 格式,生产环境必须替换为独立密钥 # 生成方式: python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())" -ENCRYPTION_KEY= -DATABASE_ENCRYPTION_KEY="rJC2VqLg4KrlSq+DGJcYm869q4v5KB2dFAeuQTe0I50=" +# ENCRYPTION_KEY 用于 GeoServer 管理密码等通用敏感配置 +ENCRYPTION_KEY="replace-with-generated-fernet-key" +# DATABASE_ENCRYPTION_KEY 专用于 project_databases.dsn_encrypted +DATABASE_ENCRYPTION_KEY="replace-with-generated-fernet-key" # ============================================ # 数据库配置 (PostgreSQL) diff --git a/AUTHENTICATION_AND_USER_MANAGEMENT.md b/AUTHENTICATION_AND_USER_MANAGEMENT.md new file mode 100644 index 0000000..b4d01a3 --- /dev/null +++ b/AUTHENTICATION_AND_USER_MANAGEMENT.md @@ -0,0 +1,95 @@ +# TJWater Authentication and Metadata Management + +## Ownership + +Keycloak owns login identity, credentials, token issuance, and token expiry. +TJWater metadata stores only business snapshots and authorization data: + +- `users.keycloak_id` is the stable identity binding. +- `users.username`, `users.email`, and `users.last_login_at` are Keycloak claim caches. +- `users.role`, `users.is_active`, and `users.is_superuser` control TJWater system access. +- `user_project_membership.project_role` controls project access. + +The backend does not accept passwords, does not issue local JWTs, and does not +trust frontend-supplied user IDs. + +## Login Snapshot Refresh + +Every authenticated metadata-user resolution validates the Keycloak access token +and reads `sub`, `preferred_username` or `username`, and `email` claims. The +backend finds `users` by `keycloak_id = sub`, rejects inactive or missing users, +then refreshes `username`, `email`, and `last_login_at`. + +This keeps local display data current without changing the identity binding. +There is no Keycloak webhook requirement; second-level user or permission sync is +out of scope unless explicitly requested later. + +## Admin APIs + +All admin APIs require metadata admin access: `users.is_superuser = true` or +`users.role = 'admin'`. + +User and membership management: + +- `GET /api/v1/admin/me` +- `POST /api/v1/admin/users/sync` +- `POST /api/v1/admin/users/sync/batch` +- `GET /api/v1/admin/users` +- `GET /api/v1/admin/users/{user_id}` +- `PATCH /api/v1/admin/users/{user_id}` +- `GET /api/v1/admin/projects/{project_id}/members` +- `POST /api/v1/admin/projects/{project_id}/members` +- `PATCH /api/v1/admin/projects/{project_id}/members/{user_id}` +- `DELETE /api/v1/admin/projects/{project_id}/members/{user_id}` + +Project configuration: + +- `GET /api/v1/admin/projects` +- `POST /api/v1/admin/projects` +- `PATCH /api/v1/admin/projects/{project_id}` +- `GET /api/v1/admin/projects/{project_id}/databases` +- `PUT /api/v1/admin/projects/{project_id}/databases` +- `DELETE /api/v1/admin/projects/{project_id}/databases/{db_role}` +- `POST /api/v1/admin/projects/{project_id}/databases/{db_role}/health` +- `GET /api/v1/admin/projects/{project_id}/geoserver` +- `PUT /api/v1/admin/projects/{project_id}/geoserver` + +## Secret Handling + +Admins submit plaintext DSNs and GeoServer passwords only through HTTPS admin +APIs. Operators should not write encrypted columns manually. + +- `project_databases.dsn_encrypted` is encrypted with `DATABASE_ENCRYPTION_KEY`. +- `project_geoserver_configs.gs_admin_password_encrypted` is encrypted with + `ENCRYPTION_KEY`. +- Admin responses return only `has_dsn` or `has_password`. +- Audit logs record whether a secret was updated, but never store plaintext DSNs + or passwords. + +Generate both keys with: + +```bash +python -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())" +``` + +Keep keys stable for the lifetime of encrypted metadata. Rotating a key requires +decrypting with the old key and re-encrypting with the new key. + +## Metadata Schema Patches + +Apply metadata patches in order: + +1. `resources/sql/004_metadata_auth_management.sql` +2. `resources/sql/005_metadata_project_configuration.sql` + +`004` creates Keycloak-backed metadata users and project memberships. `005` +creates project, project database routing, and GeoServer configuration tables +with uniqueness, role/type, and pool-size constraints. + +## Frontend System Management + +`/system-admin` is shown only after `GET /api/v1/admin/me` confirms metadata +admin access. The page lets admins maintain metadata users, project members, +projects, project database routing for `biz_data` and `iot_data`, connection +health checks, and GeoServer config. This replaces direct SQL editing for normal +project onboarding. diff --git a/app/api/v1/endpoints/admin_metadata.py b/app/api/v1/endpoints/admin_metadata.py new file mode 100644 index 0000000..7648a25 --- /dev/null +++ b/app/api/v1/endpoints/admin_metadata.py @@ -0,0 +1,804 @@ +from typing import List +from uuid import UUID + +from fastapi import APIRouter, Depends, HTTPException, Path, Query, Response, status +from sqlalchemy import text +from sqlalchemy.engine.url import make_url +from sqlalchemy.exc import IntegrityError, SQLAlchemyError +from sqlalchemy.ext.asyncio import create_async_engine + +from app.auth.metadata_dependencies import ( + get_current_metadata_admin, + get_metadata_repository, +) +from app.core.audit import AuditAction, log_audit_event +from app.domain.schemas.admin_metadata import ( + AdminProjectCreateRequest, + AdminProjectResponse, + AdminProjectUpdateRequest, + MetadataUsersBatchSyncRequest, + MetadataUserResponse, + MetadataUserSyncRequest, + MetadataUserSyncResult, + MetadataUserUpdateRequest, + ProjectDatabaseHealthResponse, + ProjectDatabaseHealthRequest, + ProjectDatabaseResponse, + ProjectDatabaseUpsertRequest, + ProjectDbRole, + ProjectGeoServerConfigResponse, + ProjectGeoServerConfigUpsertRequest, + ProjectMemberCreateRequest, + ProjectMemberResponse, + ProjectMemberUpdateRequest, +) +from app.infra.db.metadb import models +from app.infra.db.metadb.repositories.metadata_repository import ( + MetadataRepository, + ProjectDbRouting, +) + +router = APIRouter() + + +def _project_response(project: models.Project) -> AdminProjectResponse: + return AdminProjectResponse( + project_id=project.id, + name=project.name, + code=project.code, + description=project.description, + gs_workspace=project.gs_workspace, + map_extent=project.map_extent, + status=project.status, + created_at=project.created_at, + updated_at=project.updated_at, + ) + + +def _project_database_response( + record: models.ProjectDatabase, +) -> ProjectDatabaseResponse: + return ProjectDatabaseResponse( + id=record.id, + project_id=record.project_id, + db_role=record.db_role, + db_type=record.db_type, + pool_min_size=record.pool_min_size, + pool_max_size=record.pool_max_size, + has_dsn=bool(record.dsn_encrypted), + ) + + +def _geoserver_config_response( + record: models.ProjectGeoServerConfig, +) -> ProjectGeoServerConfigResponse: + return ProjectGeoServerConfigResponse( + id=record.id, + project_id=record.project_id, + gs_base_url=record.gs_base_url, + gs_admin_user=record.gs_admin_user, + gs_datastore_name=record.gs_datastore_name, + default_extent=record.default_extent, + srid=record.srid, + configured=True, + has_password=bool(record.gs_admin_password_encrypted), + updated_at=record.updated_at, + ) + + +def _database_audit_payload(payload: ProjectDatabaseUpsertRequest) -> dict: + return { + "db_role": payload.db_role, + "db_type": _db_type_for_role(payload.db_role), + "pool_min_size": payload.pool_min_size, + "pool_max_size": payload.pool_max_size, + "dsn_updated": payload.dsn is not None, + } + + +def _geoserver_audit_payload( + payload: ProjectGeoServerConfigUpsertRequest, +) -> dict: + return { + "gs_base_url": payload.gs_base_url, + "gs_admin_user": payload.gs_admin_user, + "gs_datastore_name": payload.gs_datastore_name, + "default_extent": payload.default_extent, + "srid": payload.srid, + "password_updated": "gs_admin_password" in payload.model_fields_set, + } + + +def _to_async_sqlalchemy_url(dsn: str) -> str: + parsed = make_url(dsn) + if parsed.drivername in {"postgresql", "postgres"}: + parsed = parsed.set(drivername="postgresql+psycopg") + return parsed.render_as_string(hide_password=False) + + +def _db_type_for_role(db_role: str) -> str: + if db_role == "iot_data": + return "timescaledb" + return "postgresql" + + +def _status_for_config_value_error(exc: ValueError) -> int: + if "ENCRYPTION_KEY" in str(exc): + return status.HTTP_503_SERVICE_UNAVAILABLE + return status.HTTP_400_BAD_REQUEST + + +async def _check_database_connection(routing: ProjectDbRouting) -> None: + engine = create_async_engine( + _to_async_sqlalchemy_url(routing.dsn), + pool_size=1, + max_overflow=0, + pool_pre_ping=True, + ) + try: + async with engine.connect() as conn: + await conn.execute(text("SELECT 1")) + finally: + await engine.dispose() + + +def _database_health_error_detail(exc: Exception) -> str: + message = str(exc) + lower_message = message.lower() + if "password authentication failed" in lower_message: + return "连通性测试失败:用户名或密码错误,请检查 DSN 中的账号密码。" + if "connection refused" in lower_message: + return "连通性测试失败:目标主机或端口拒绝连接,请检查地址、端口和服务状态。" + if "timeout" in lower_message or "timed out" in lower_message: + return "连通性测试失败:连接超时,请检查网络、防火墙和数据库服务状态。" + if "could not translate host name" in lower_message or "name or service not known" in lower_message: + return "连通性测试失败:数据库主机名无法解析,请检查 DSN 中的主机地址。" + first_line = message.splitlines()[0] if message else exc.__class__.__name__ + return f"连通性测试失败:{first_line}" + + +async def _upsert_and_audit_metadata_user( + payload: MetadataUserSyncRequest, + *, + current_user, + metadata_repo: MetadataRepository, + response_status: int, +) -> MetadataUserResponse: + user = await metadata_repo.upsert_user_from_keycloak( + keycloak_id=payload.keycloak_id, + username=payload.username, + email=str(payload.email), + role=payload.role, + is_active=payload.is_active, + ) + await log_audit_event( + action=AuditAction.UPDATE, + user_id=current_user.id, + resource_type="metadata_user", + resource_id=str(user.id), + request_data=payload.model_dump(mode="json"), + response_status=response_status, + session=metadata_repo.session, + ) + return MetadataUserResponse.model_validate(user) + + +@router.get("/me", response_model=MetadataUserResponse) +async def get_metadata_admin_me( + current_user=Depends(get_current_metadata_admin), +) -> MetadataUserResponse: + return MetadataUserResponse.model_validate(current_user) + + +@router.post("/users/sync", response_model=MetadataUserResponse) +async def sync_metadata_user( + payload: MetadataUserSyncRequest, + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> MetadataUserResponse: + try: + return await _upsert_and_audit_metadata_user( + payload, + current_user=current_user, + metadata_repo=metadata_repo, + response_status=status.HTTP_200_OK, + ) + except IntegrityError as exc: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="User keycloak_id, username, or email conflicts with an existing user", + ) from exc + except SQLAlchemyError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc + + + +@router.post("/users/sync/batch", response_model=List[MetadataUserSyncResult]) +async def sync_metadata_users_batch( + payload: MetadataUsersBatchSyncRequest, + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> List[MetadataUserSyncResult]: + results: list[MetadataUserSyncResult] = [] + for item in payload.users: + try: + user = await _upsert_and_audit_metadata_user( + item, + current_user=current_user, + metadata_repo=metadata_repo, + response_status=status.HTTP_200_OK, + ) + except IntegrityError as exc: + results.append( + MetadataUserSyncResult( + keycloak_id=item.keycloak_id, + success=False, + error="User keycloak_id, username, or email conflicts with an existing user", + ) + ) + await metadata_repo.session.rollback() + except SQLAlchemyError as exc: + results.append( + MetadataUserSyncResult( + keycloak_id=item.keycloak_id, + success=False, + error=f"Metadata database error: {exc}", + ) + ) + await metadata_repo.session.rollback() + else: + results.append( + MetadataUserSyncResult( + keycloak_id=item.keycloak_id, + success=True, + user=user, + ) + ) + return results + + +@router.get("/users", response_model=List[MetadataUserResponse]) +async def list_metadata_users( + skip: int = Query(0, ge=0), + limit: int = Query(100, ge=1, le=1000), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> List[MetadataUserResponse]: + users = await metadata_repo.list_users(skip=skip, limit=limit) + return [MetadataUserResponse.model_validate(user) for user in users] + + +@router.get("/projects", response_model=List[AdminProjectResponse]) +async def list_admin_projects( + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> List[AdminProjectResponse]: + projects = await metadata_repo.list_project_records() + return [_project_response(project) for project in projects] + + +@router.post( + "/projects", + response_model=AdminProjectResponse, + status_code=status.HTTP_201_CREATED, +) +async def create_admin_project( + payload: AdminProjectCreateRequest, + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> AdminProjectResponse: + try: + project = await metadata_repo.create_project( + name=payload.name, + code=payload.code, + description=payload.description, + gs_workspace=payload.gs_workspace, + map_extent=payload.map_extent, + status=payload.status, + ) + except IntegrityError as exc: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Project code or GeoServer workspace conflicts with an existing project", + ) from exc + except SQLAlchemyError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc + + await log_audit_event( + action=AuditAction.CREATE, + user_id=current_user.id, + project_id=project.id, + resource_type="project", + resource_id=str(project.id), + request_data=payload.model_dump(mode="json"), + response_status=status.HTTP_201_CREATED, + session=metadata_repo.session, + ) + return _project_response(project) + + +@router.patch( + "/projects/{project_id}", + response_model=AdminProjectResponse, +) +async def update_admin_project( + payload: AdminProjectUpdateRequest, + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> AdminProjectResponse: + updates = payload.model_dump(mode="json", exclude_unset=True) + try: + project = await metadata_repo.update_project(project_id, updates=updates) + except IntegrityError as exc: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Project code or GeoServer workspace conflicts with an existing project", + ) from exc + except SQLAlchemyError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc + if project is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found") + + await log_audit_event( + action=AuditAction.UPDATE, + user_id=current_user.id, + project_id=project.id, + resource_type="project", + resource_id=str(project.id), + request_data=updates, + response_status=status.HTTP_200_OK, + session=metadata_repo.session, + ) + return _project_response(project) + + +@router.get( + "/projects/{project_id}/databases", + response_model=List[ProjectDatabaseResponse], +) +async def list_project_databases( + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> List[ProjectDatabaseResponse]: + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found") + records = await metadata_repo.list_project_databases(project_id) + return [_project_database_response(record) for record in records] + + +@router.put( + "/projects/{project_id}/databases", + response_model=ProjectDatabaseResponse, +) +async def upsert_project_database( + payload: ProjectDatabaseUpsertRequest, + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectDatabaseResponse: + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found") + try: + routing = ( + ProjectDbRouting( + project_id=project_id, + db_role=payload.db_role, + db_type=_db_type_for_role(payload.db_role), + dsn=payload.dsn, + pool_min_size=payload.pool_min_size, + pool_max_size=payload.pool_max_size, + ) + if payload.dsn + else await metadata_repo.get_project_db_routing(project_id, payload.db_role) + ) + if routing is None: + raise ValueError("dsn is required when creating project database config") + await _check_database_connection(routing) + except ValueError as exc: + raise HTTPException( + status_code=_status_for_config_value_error(exc), + detail=str(exc), + ) from exc + except Exception as exc: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail=_database_health_error_detail(exc), + ) from exc + + try: + record = await metadata_repo.upsert_project_database_config( + project_id, + db_role=payload.db_role, + db_type=_db_type_for_role(payload.db_role), + dsn=payload.dsn, + pool_min_size=payload.pool_min_size, + pool_max_size=payload.pool_max_size, + ) + except IntegrityError as exc: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Project database role conflicts with an existing config", + ) from exc + except SQLAlchemyError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc + + await log_audit_event( + action=AuditAction.CONFIG_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_database", + resource_id=payload.db_role, + request_data=_database_audit_payload(payload), + response_status=status.HTTP_200_OK, + session=metadata_repo.session, + ) + return _project_database_response(record) + + +@router.delete( + "/projects/{project_id}/databases/{db_role}", + status_code=status.HTTP_204_NO_CONTENT, +) +async def delete_project_database( + project_id: UUID = Path(...), + db_role: ProjectDbRole = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> None: + removed = await metadata_repo.delete_project_database_config(project_id, db_role) + if not removed: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail="Project database config not found", + ) + await log_audit_event( + action=AuditAction.CONFIG_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_database", + resource_id=db_role, + request_data={"deleted": True}, + response_status=status.HTTP_204_NO_CONTENT, + session=metadata_repo.session, + ) + + +@router.post( + "/projects/{project_id}/databases/{db_role}/health", + response_model=ProjectDatabaseHealthResponse, +) +async def check_project_database_health( + response: Response, + project_id: UUID = Path(...), + db_role: ProjectDbRole = Path(...), + payload: ProjectDatabaseHealthRequest | None = None, + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectDatabaseHealthResponse: + dsn_to_test = payload.dsn if payload and payload.dsn else None + if dsn_to_test: + routing = ProjectDbRouting( + project_id=project_id, + db_role=db_role, + db_type=_db_type_for_role(db_role), + dsn=dsn_to_test, + pool_min_size=1, + pool_max_size=1, + ) + else: + try: + routing = await metadata_repo.get_project_db_routing(project_id, db_role) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Project database routing DSN is invalid: {exc}", + ) from exc + if routing is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail="Project database config not found", + ) + + try: + await _check_database_connection(routing) + except Exception as exc: # health endpoint should return diagnostic status + response.status_code = status.HTTP_503_SERVICE_UNAVAILABLE + return ProjectDatabaseHealthResponse( + project_id=project_id, + db_role=db_role, + db_type=routing.db_type, + ok=False, + detail=_database_health_error_detail(exc), + ) + return ProjectDatabaseHealthResponse( + project_id=project_id, + db_role=db_role, + db_type=routing.db_type, + ok=True, + detail="连通性测试通过", + ) + + +@router.get( + "/projects/{project_id}/geoserver", + response_model=ProjectGeoServerConfigResponse, +) +async def get_project_geoserver_config( + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectGeoServerConfigResponse: + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found") + record = await metadata_repo.get_geoserver_config_record(project_id) + if record is None: + return ProjectGeoServerConfigResponse( + id=None, + project_id=project_id, + gs_base_url=None, + gs_admin_user=None, + gs_datastore_name="ds_postgis", + default_extent=None, + srid=4326, + configured=False, + has_password=False, + updated_at=None, + ) + return _geoserver_config_response(record) + + +@router.put( + "/projects/{project_id}/geoserver", + response_model=ProjectGeoServerConfigResponse, +) +async def upsert_project_geoserver_config( + payload: ProjectGeoServerConfigUpsertRequest, + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectGeoServerConfigResponse: + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Project not found") + try: + record = await metadata_repo.upsert_geoserver_config( + project_id, + gs_base_url=payload.gs_base_url, + gs_admin_user=payload.gs_admin_user, + gs_admin_password=payload.gs_admin_password, + password_update_requested="gs_admin_password" in payload.model_fields_set, + gs_datastore_name=payload.gs_datastore_name, + default_extent=payload.default_extent, + srid=payload.srid, + ) + except ValueError as exc: + raise HTTPException( + status_code=_status_for_config_value_error(exc), + detail=str(exc), + ) from exc + except SQLAlchemyError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc + + await log_audit_event( + action=AuditAction.CONFIG_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_geoserver", + resource_id=str(project_id), + request_data=_geoserver_audit_payload(payload), + response_status=status.HTTP_200_OK, + session=metadata_repo.session, + ) + return _geoserver_config_response(record) + + +@router.get("/users/{user_id}", response_model=MetadataUserResponse) +async def get_metadata_user( + user_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> MetadataUserResponse: + user = await metadata_repo.get_user_by_id(user_id) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found") + return MetadataUserResponse.model_validate(user) + + +@router.patch("/users/{user_id}", response_model=MetadataUserResponse) +async def update_metadata_user( + payload: MetadataUserUpdateRequest, + user_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> MetadataUserResponse: + updates = payload.model_dump(mode="json", exclude_unset=True) + if user_id == current_user.id: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Users cannot modify themselves", + ) + user = await metadata_repo.update_user_admin( + user_id, + updates=updates, + ) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found") + + await log_audit_event( + action=AuditAction.UPDATE, + user_id=current_user.id, + resource_type="metadata_user", + resource_id=str(user.id), + request_data=updates, + response_status=status.HTTP_200_OK, + session=metadata_repo.session, + ) + return MetadataUserResponse.model_validate(user) + + +@router.get( + "/projects/{project_id}/members", + response_model=List[ProjectMemberResponse], +) +async def list_project_members( + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> List[ProjectMemberResponse]: + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="Project not found" + ) + members = await metadata_repo.list_project_members(project_id) + return [ProjectMemberResponse(**member.__dict__) for member in members] + + +@router.post( + "/projects/{project_id}/members", + response_model=ProjectMemberResponse, + status_code=status.HTTP_201_CREATED, +) +async def add_project_member( + payload: ProjectMemberCreateRequest, + project_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectMemberResponse: + if payload.user_id == current_user.id: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Users cannot modify their own project membership", + ) + project = await metadata_repo.get_project_by_id(project_id) + if project is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="Project not found" + ) + user = await metadata_repo.get_user_by_id(payload.user_id) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found") + existing = await metadata_repo.get_project_membership(project_id, payload.user_id) + if existing is not None: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="User is already a project member", + ) + + membership = await metadata_repo.add_project_member( + project_id, payload.user_id, payload.project_role + ) + await log_audit_event( + action=AuditAction.PERMISSION_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_member", + resource_id=str(payload.user_id), + request_data=payload.model_dump(mode="json"), + response_status=status.HTTP_201_CREATED, + session=metadata_repo.session, + ) + return ProjectMemberResponse( + id=membership.id, + user_id=membership.user_id, + project_id=membership.project_id, + project_role=membership.project_role, + username=user.username, + email=user.email, + is_active=user.is_active, + ) + + +@router.patch( + "/projects/{project_id}/members/{user_id}", + response_model=ProjectMemberResponse, +) +async def update_project_member( + payload: ProjectMemberUpdateRequest, + project_id: UUID = Path(...), + user_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> ProjectMemberResponse: + if user_id == current_user.id: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Users cannot modify their own project membership", + ) + user = await metadata_repo.get_user_by_id(user_id) + if user is None: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="User not found") + membership = await metadata_repo.update_project_member_role( + project_id, user_id, payload.project_role + ) + if membership is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="Project member not found" + ) + await log_audit_event( + action=AuditAction.PERMISSION_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_member", + resource_id=str(user_id), + request_data=payload.model_dump(mode="json"), + response_status=status.HTTP_200_OK, + session=metadata_repo.session, + ) + return ProjectMemberResponse( + id=membership.id, + user_id=membership.user_id, + project_id=membership.project_id, + project_role=membership.project_role, + username=user.username, + email=user.email, + is_active=user.is_active, + ) + + +@router.delete("/projects/{project_id}/members/{user_id}", status_code=status.HTTP_204_NO_CONTENT) +async def remove_project_member( + project_id: UUID = Path(...), + user_id: UUID = Path(...), + current_user=Depends(get_current_metadata_admin), + metadata_repo: MetadataRepository = Depends(get_metadata_repository), +) -> None: + if user_id == current_user.id: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Users cannot modify their own project membership", + ) + removed = await metadata_repo.remove_project_member(project_id, user_id) + if not removed: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="Project member not found" + ) + await log_audit_event( + action=AuditAction.PERMISSION_CHANGE, + user_id=current_user.id, + project_id=project_id, + resource_type="project_member", + resource_id=str(user_id), + response_status=status.HTTP_204_NO_CONTENT, + session=metadata_repo.session, + ) diff --git a/app/api/v1/router.py b/app/api/v1/router.py index 3e4e7d3..99da712 100644 --- a/app/api/v1/router.py +++ b/app/api/v1/router.py @@ -1,5 +1,6 @@ from fastapi import APIRouter from app.api.v1.endpoints import ( + admin_metadata, agent_auth, project, simulation, @@ -54,6 +55,9 @@ api_router = APIRouter() # Core Services api_router.include_router(agent_auth.router, tags=["Agent Auth"]) +api_router.include_router( + admin_metadata.router, prefix="/admin", tags=["Metadata Admin"] +) api_router.include_router(audit.router, prefix="/audit", tags=["Audit Logs"]) # 新增 api_router.include_router(meta.router, tags=["Metadata"]) api_router.include_router(project.router, tags=["Project"]) diff --git a/app/auth/metadata_dependencies.py b/app/auth/metadata_dependencies.py index 8424429..021063e 100644 --- a/app/auth/metadata_dependencies.py +++ b/app/auth/metadata_dependencies.py @@ -6,8 +6,7 @@ from fastapi import Depends, HTTPException, status from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.ext.asyncio import AsyncSession -from app.auth.keycloak_dependencies import get_current_keycloak_sub -from app.core.config import settings +from app.auth.keycloak_dependencies import get_current_keycloak_payload from app.infra.db.metadb.database import get_metadata_session from app.infra.db.metadb.repositories.metadata_repository import MetadataRepository @@ -20,10 +19,40 @@ async def get_metadata_repository( return MetadataRepository(session) +def _keycloak_sub_from_payload(payload: dict) -> UUID: + sub = payload.get("sub") + if not sub: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="Missing subject claim", + headers={"WWW-Authenticate": "Bearer"}, + ) + + try: + return UUID(str(sub)) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="Invalid subject claim", + headers={"WWW-Authenticate": "Bearer"}, + ) from exc + + +def _username_from_payload(payload: dict) -> str | None: + username = payload.get("preferred_username") or payload.get("username") + return str(username) if username else None + + +def _email_from_payload(payload: dict) -> str | None: + email = payload.get("email") + return str(email) if email else None + + async def get_current_metadata_user( - keycloak_sub: UUID = Depends(get_current_keycloak_sub), + keycloak_payload: dict = Depends(get_current_keycloak_payload), metadata_repo: MetadataRepository = Depends(get_metadata_repository), ): + keycloak_sub = _keycloak_sub_from_payload(keycloak_payload) try: user = await metadata_repo.get_user_by_keycloak_id(keycloak_sub) except SQLAlchemyError as exc: @@ -39,6 +68,21 @@ async def get_current_metadata_user( raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail="Inactive user" ) + try: + user = await metadata_repo.refresh_user_keycloak_snapshot( + user, + username=_username_from_payload(keycloak_payload), + email=_email_from_payload(keycloak_payload), + ) + except SQLAlchemyError as exc: + logger.error( + "Metadata DB error while refreshing current user snapshot", + exc_info=True, + ) + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=f"Metadata database error: {exc}", + ) from exc return user diff --git a/app/domain/schemas/admin_metadata.py b/app/domain/schemas/admin_metadata.py new file mode 100644 index 0000000..0138348 --- /dev/null +++ b/app/domain/schemas/admin_metadata.py @@ -0,0 +1,156 @@ +from datetime import datetime +from typing import Literal +from uuid import UUID + +from pydantic import BaseModel, ConfigDict, Field, model_validator + + +BusinessRole = Literal["admin", "user", "operator", "viewer"] +ProjectRole = Literal["owner", "admin", "member", "viewer"] +ProjectStatus = Literal["active", "inactive", "archived"] +ProjectDbRole = Literal["biz_data", "iot_data"] + + +class MetadataUserSyncRequest(BaseModel): + keycloak_id: UUID + username: str = Field(..., min_length=1, max_length=50) + email: str = Field(..., min_length=1, max_length=100) + role: BusinessRole = "user" + is_active: bool = True + + +class MetadataUsersBatchSyncRequest(BaseModel): + users: list[MetadataUserSyncRequest] = Field(..., min_length=1, max_length=500) + + +class MetadataUserUpdateRequest(BaseModel): + role: BusinessRole | None = None + is_active: bool | None = None + + +class MetadataUserResponse(BaseModel): + id: UUID + keycloak_id: UUID + username: str + email: str + role: str + is_active: bool + is_superuser: bool + created_at: datetime + updated_at: datetime + last_login_at: datetime | None = None + + model_config = ConfigDict(from_attributes=True) + + +class MetadataUserSyncResult(BaseModel): + keycloak_id: UUID + user: MetadataUserResponse | None = None + success: bool + error: str | None = None + + +class ProjectMemberCreateRequest(BaseModel): + user_id: UUID + project_role: ProjectRole = "viewer" + + +class ProjectMemberUpdateRequest(BaseModel): + project_role: ProjectRole + + +class ProjectMemberResponse(BaseModel): + id: UUID + user_id: UUID + project_id: UUID + project_role: str + username: str + email: str + is_active: bool + + +class AdminProjectCreateRequest(BaseModel): + name: str = Field(..., min_length=1, max_length=100) + code: str = Field(..., min_length=1, max_length=50) + description: str | None = None + gs_workspace: str = Field(..., min_length=1, max_length=100) + map_extent: dict | None = None + status: ProjectStatus = "active" + + +class AdminProjectUpdateRequest(BaseModel): + name: str | None = Field(default=None, min_length=1, max_length=100) + code: str | None = Field(default=None, min_length=1, max_length=50) + description: str | None = None + gs_workspace: str | None = Field(default=None, min_length=1, max_length=100) + map_extent: dict | None = None + status: ProjectStatus | None = None + + +class AdminProjectResponse(BaseModel): + project_id: UUID + name: str + code: str + description: str | None = None + gs_workspace: str + map_extent: dict | None = None + status: str + created_at: datetime + updated_at: datetime + + +class ProjectDatabaseUpsertRequest(BaseModel): + db_role: ProjectDbRole + dsn: str | None = Field(default=None, min_length=1) + pool_min_size: int = Field(default=2, ge=1) + pool_max_size: int = Field(default=10, ge=1) + + @model_validator(mode="after") + def validate_pool_bounds(self): + if self.pool_max_size < self.pool_min_size: + raise ValueError("pool_max_size must be greater than or equal to pool_min_size") + return self + + +class ProjectDatabaseResponse(BaseModel): + id: UUID + project_id: UUID + db_role: str + db_type: str + pool_min_size: int + pool_max_size: int + has_dsn: bool + + +class ProjectDatabaseHealthRequest(BaseModel): + dsn: str | None = Field(default=None, min_length=1) + + +class ProjectDatabaseHealthResponse(BaseModel): + project_id: UUID + db_role: str + db_type: str + ok: bool + detail: str + + +class ProjectGeoServerConfigUpsertRequest(BaseModel): + gs_base_url: str | None = None + gs_admin_user: str | None = Field(default=None, max_length=50) + gs_admin_password: str | None = Field(default=None, min_length=1) + gs_datastore_name: str = Field(default="ds_postgis", min_length=1, max_length=100) + default_extent: dict | None = None + srid: int = Field(default=4326, ge=1) + + +class ProjectGeoServerConfigResponse(BaseModel): + id: UUID | None = None + project_id: UUID + gs_base_url: str | None = None + gs_admin_user: str | None = None + gs_datastore_name: str + default_extent: dict | None = None + srid: int + configured: bool = True + has_password: bool + updated_at: datetime | None = None diff --git a/app/infra/db/dynamic_manager.py b/app/infra/db/dynamic_manager.py index e444a6e..c78a2e5 100644 --- a/app/infra/db/dynamic_manager.py +++ b/app/infra/db/dynamic_manager.py @@ -54,9 +54,9 @@ class ProjectConnectionManager: def _normalize_pg_url(self, url: str) -> str: parsed = make_url(url) - if parsed.drivername == "postgresql": + if parsed.drivername in {"postgresql", "postgres"}: parsed = parsed.set(drivername="postgresql+psycopg") - return str(parsed) + return parsed.render_as_string(hide_password=False) async def get_pg_sessionmaker( self, diff --git a/app/infra/db/metadb/repositories/metadata_repository.py b/app/infra/db/metadb/repositories/metadata_repository.py index 9631d7b..2555fc1 100644 --- a/app/infra/db/metadb/repositories/metadata_repository.py +++ b/app/infra/db/metadb/repositories/metadata_repository.py @@ -1,9 +1,10 @@ from dataclasses import dataclass +from datetime import datetime, timezone from typing import Optional, List -from uuid import UUID +from uuid import UUID, uuid4 from cryptography.fernet import InvalidToken -from sqlalchemy import select +from sqlalchemy import delete, select from sqlalchemy.ext.asyncio import AsyncSession from app.core.encryption import ( @@ -78,6 +79,33 @@ class ProjectDetail: geoserver: Optional[ProjectGeoServerInfo] +@dataclass(frozen=True) +class ProjectMemberSummary: + id: UUID + user_id: UUID + project_id: UUID + project_role: str + username: str + email: str + is_active: bool + + +def _utcnow() -> datetime: + return datetime.now(timezone.utc) + + +def _encrypt_database_secret(value: str) -> str: + if not is_database_encryption_configured(): + raise ValueError("DATABASE_ENCRYPTION_KEY is not configured") + return get_database_encryptor().encrypt(value) + + +def _encrypt_general_secret(value: str) -> str: + if not is_encryption_configured(): + raise ValueError("ENCRYPTION_KEY is not configured") + return get_encryptor().encrypt(value) + + class MetadataRepository: """元数据访问层(system_hub)""" @@ -96,6 +124,86 @@ class MetadataRepository: ) return result.scalar_one_or_none() + async def get_user_by_id(self, user_id: UUID) -> Optional[models.User]: + result = await self.session.execute( + select(models.User).where(models.User.id == user_id) + ) + return result.scalar_one_or_none() + + async def list_users(self, skip: int = 0, limit: int = 100) -> List[models.User]: + result = await self.session.execute( + select(models.User) + .order_by(models.User.created_at.desc()) + .offset(skip) + .limit(limit) + ) + return list(result.scalars().all()) + + async def upsert_user_from_keycloak( + self, + *, + keycloak_id: UUID, + username: str, + email: str, + role: str, + is_active: bool, + ) -> models.User: + user = await self.get_user_by_keycloak_id(keycloak_id) + if user is None: + user = models.User( + id=uuid4(), + keycloak_id=keycloak_id, + username=username, + email=email, + role=role, + is_active=is_active, + is_superuser=False, + ) + self.session.add(user) + else: + user.username = username + user.email = email + user.role = role + user.is_active = is_active + await self.session.commit() + await self.session.refresh(user) + return user + + async def refresh_user_keycloak_snapshot( + self, + user: models.User, + *, + username: str | None, + email: str | None, + last_login_at: datetime | None = None, + ) -> models.User: + if username: + user.username = username + if email: + user.email = email + user.last_login_at = last_login_at or _utcnow() + user.updated_at = _utcnow() + await self.session.commit() + await self.session.refresh(user) + return user + + async def update_user_admin( + self, + user_id: UUID, + *, + updates: dict, + ) -> Optional[models.User]: + user = await self.get_user_by_id(user_id) + if user is None: + return None + if "role" in updates: + user.role = updates["role"] + if "is_active" in updates: + user.is_active = updates["is_active"] + await self.session.commit() + await self.session.refresh(user) + return user + async def get_project_by_id(self, project_id: UUID) -> Optional[models.Project]: result = await self.session.execute( select(models.Project).where(models.Project.id == project_id) @@ -108,6 +216,62 @@ class MetadataRepository: ) return result.scalar_one_or_none() + async def list_project_records(self) -> List[models.Project]: + result = await self.session.execute( + select(models.Project).order_by(models.Project.name) + ) + return list(result.scalars().all()) + + async def create_project( + self, + *, + name: str, + code: str, + description: str | None, + gs_workspace: str, + map_extent: dict | None, + status: str, + ) -> models.Project: + project = models.Project( + id=uuid4(), + name=name, + code=code, + description=description, + gs_workspace=gs_workspace, + map_extent=map_extent, + status=status, + created_at=_utcnow(), + updated_at=_utcnow(), + ) + self.session.add(project) + await self.session.commit() + await self.session.refresh(project) + return project + + async def update_project( + self, + project_id: UUID, + *, + updates: dict, + ) -> Optional[models.Project]: + project = await self.get_project_by_id(project_id) + if project is None: + return None + for field in ( + "name", + "code", + "description", + "gs_workspace", + "map_extent", + "status", + ): + if field in updates: + setattr(project, field, updates[field]) + project.updated_at = _utcnow() + await self.session.commit() + await self.session.refresh(project) + return project + async def get_project_detail_by_code(self, code: str) -> Optional[ProjectDetail]: project = await self.get_project_by_code(code) if not project: @@ -137,6 +301,142 @@ class MetadataRepository: ) return result.scalar_one_or_none() + async def list_project_members( + self, project_id: UUID + ) -> List[ProjectMemberSummary]: + stmt = ( + select(models.UserProjectMembership, models.User) + .join(models.User, models.User.id == models.UserProjectMembership.user_id) + .where(models.UserProjectMembership.project_id == project_id) + .order_by(models.User.username) + ) + result = await self.session.execute(stmt) + return [ + ProjectMemberSummary( + id=membership.id, + user_id=membership.user_id, + project_id=membership.project_id, + project_role=membership.project_role, + username=user.username, + email=user.email, + is_active=user.is_active, + ) + for membership, user in result.all() + ] + + async def get_project_membership( + self, project_id: UUID, user_id: UUID + ) -> Optional[models.UserProjectMembership]: + result = await self.session.execute( + select(models.UserProjectMembership).where( + models.UserProjectMembership.project_id == project_id, + models.UserProjectMembership.user_id == user_id, + ) + ) + return result.scalar_one_or_none() + + async def add_project_member( + self, project_id: UUID, user_id: UUID, project_role: str + ) -> models.UserProjectMembership: + membership = models.UserProjectMembership( + id=uuid4(), + user_id=user_id, + project_id=project_id, + project_role=project_role, + ) + self.session.add(membership) + await self.session.commit() + await self.session.refresh(membership) + return membership + + async def update_project_member_role( + self, project_id: UUID, user_id: UUID, project_role: str + ) -> Optional[models.UserProjectMembership]: + membership = await self.get_project_membership(project_id, user_id) + if membership is None: + return None + membership.project_role = project_role + await self.session.commit() + await self.session.refresh(membership) + return membership + + async def remove_project_member(self, project_id: UUID, user_id: UUID) -> bool: + result = await self.session.execute( + delete(models.UserProjectMembership).where( + models.UserProjectMembership.project_id == project_id, + models.UserProjectMembership.user_id == user_id, + ) + ) + await self.session.commit() + return bool(result.rowcount) + + async def list_project_databases( + self, project_id: UUID + ) -> List[models.ProjectDatabase]: + result = await self.session.execute( + select(models.ProjectDatabase) + .where(models.ProjectDatabase.project_id == project_id) + .order_by(models.ProjectDatabase.db_role) + ) + return list(result.scalars().all()) + + async def get_project_database_config( + self, project_id: UUID, db_role: str + ) -> Optional[models.ProjectDatabase]: + result = await self.session.execute( + select(models.ProjectDatabase).where( + models.ProjectDatabase.project_id == project_id, + models.ProjectDatabase.db_role == db_role, + ) + ) + return result.scalar_one_or_none() + + async def upsert_project_database_config( + self, + project_id: UUID, + *, + db_role: str, + db_type: str, + dsn: str | None, + pool_min_size: int, + pool_max_size: int, + ) -> models.ProjectDatabase: + record = await self.get_project_database_config(project_id, db_role) + if record is None: + if dsn is None: + raise ValueError("dsn is required when creating project database config") + record = models.ProjectDatabase( + id=uuid4(), + project_id=project_id, + db_role=db_role, + db_type=db_type, + dsn_encrypted=_encrypt_database_secret(dsn), + pool_min_size=pool_min_size, + pool_max_size=pool_max_size, + ) + self.session.add(record) + else: + record.db_type = db_type + if dsn is not None: + record.dsn_encrypted = _encrypt_database_secret(dsn) + record.pool_min_size = pool_min_size + record.pool_max_size = pool_max_size + await self.session.commit() + await self.session.refresh(record) + return record + + async def delete_project_database_config( + self, project_id: UUID, db_role: str + ) -> bool: + result = await self.session.execute( + delete(models.ProjectDatabase).where( + models.ProjectDatabase.project_id == project_id, + models.ProjectDatabase.db_role == db_role, + ) + ) + await self.session.commit() + return bool(result.rowcount) + async def get_project_db_routing( self, project_id: UUID, db_role: str ) -> Optional[ProjectDbRouting]: @@ -198,6 +498,59 @@ class MetadataRepository: srid=record.srid, ) + async def get_geoserver_config_record( + self, project_id: UUID + ) -> Optional[models.ProjectGeoServerConfig]: + result = await self.session.execute( + select(models.ProjectGeoServerConfig).where( + models.ProjectGeoServerConfig.project_id == project_id + ) + ) + return result.scalar_one_or_none() + + async def upsert_geoserver_config( + self, + project_id: UUID, + *, + gs_base_url: str | None, + gs_admin_user: str | None, + gs_admin_password: str | None, + password_update_requested: bool, + gs_datastore_name: str, + default_extent: dict | None, + srid: int, + ) -> models.ProjectGeoServerConfig: + record = await self.get_geoserver_config_record(project_id) + encrypted_password: str | None = None + if password_update_requested and gs_admin_password is not None: + encrypted_password = _encrypt_general_secret(gs_admin_password) + + if record is None: + record = models.ProjectGeoServerConfig( + id=uuid4(), + project_id=project_id, + gs_base_url=gs_base_url, + gs_admin_user=gs_admin_user, + gs_admin_password_encrypted=encrypted_password, + gs_datastore_name=gs_datastore_name, + default_extent=default_extent, + srid=srid, + updated_at=_utcnow(), + ) + self.session.add(record) + else: + record.gs_base_url = gs_base_url + record.gs_admin_user = gs_admin_user + if password_update_requested: + record.gs_admin_password_encrypted = encrypted_password + record.gs_datastore_name = gs_datastore_name + record.default_extent = default_extent + record.srid = srid + record.updated_at = _utcnow() + await self.session.commit() + await self.session.refresh(record) + return record + async def list_projects_for_user(self, user_id: UUID) -> List[ProjectSummary]: stmt = ( select(models.Project, models.UserProjectMembership.project_role) diff --git a/resources/sql/004_metadata_auth_management.sql b/resources/sql/004_metadata_auth_management.sql new file mode 100644 index 0000000..e3d682b --- /dev/null +++ b/resources/sql/004_metadata_auth_management.sql @@ -0,0 +1,62 @@ +-- Metadata auth management schema patch. +-- Keycloak owns login credentials; TJWater stores only business identity and access. + +CREATE EXTENSION IF NOT EXISTS pgcrypto; + +DO $$ +DECLARE + users_id_type text; +BEGIN + SELECT data_type INTO users_id_type + FROM information_schema.columns + WHERE table_schema = 'public' + AND table_name = 'users' + AND column_name = 'id'; + + IF users_id_type IS NULL THEN + CREATE TABLE users ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + keycloak_id UUID UNIQUE NOT NULL, + username VARCHAR(50) UNIQUE NOT NULL, + email VARCHAR(100) UNIQUE NOT NULL, + role VARCHAR(20) DEFAULT 'user' NOT NULL, + is_active BOOLEAN DEFAULT TRUE NOT NULL, + is_superuser BOOLEAN DEFAULT FALSE NOT NULL, + attributes JSONB, + created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP NOT NULL, + updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP NOT NULL, + last_login_at TIMESTAMP WITH TIME ZONE + ); + ELSIF users_id_type <> 'uuid' THEN + RAISE EXCEPTION + 'Existing public.users.id is %, not uuid. Export old local users, create Keycloak accounts, then migrate to metadata UUID users before applying this patch.', + users_id_type; + END IF; +END $$; + +ALTER TABLE users + ADD COLUMN IF NOT EXISTS keycloak_id UUID, + ADD COLUMN IF NOT EXISTS attributes JSONB, + ADD COLUMN IF NOT EXISTS last_login_at TIMESTAMP WITH TIME ZONE; + +ALTER TABLE users + ALTER COLUMN role SET DEFAULT 'user'; + +CREATE UNIQUE INDEX IF NOT EXISTS idx_users_keycloak_id ON users(keycloak_id); +CREATE INDEX IF NOT EXISTS idx_users_role ON users(role); +CREATE INDEX IF NOT EXISTS idx_users_is_active ON users(is_active); + +CREATE TABLE IF NOT EXISTS user_project_membership ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + user_id UUID NOT NULL, + project_id UUID NOT NULL, + project_role VARCHAR(20) DEFAULT 'viewer' NOT NULL, + CONSTRAINT user_project_membership_role_check + CHECK (project_role IN ('owner', 'admin', 'member', 'viewer')), + CONSTRAINT user_project_membership_unique UNIQUE (user_id, project_id) +); + +CREATE INDEX IF NOT EXISTS idx_user_project_membership_user_id + ON user_project_membership(user_id); +CREATE INDEX IF NOT EXISTS idx_user_project_membership_project_id + ON user_project_membership(project_id); diff --git a/resources/sql/005_metadata_project_configuration.sql b/resources/sql/005_metadata_project_configuration.sql new file mode 100644 index 0000000..1d6f910 --- /dev/null +++ b/resources/sql/005_metadata_project_configuration.sql @@ -0,0 +1,78 @@ +-- Metadata project configuration schema patch. +-- Admin APIs write these tables; operators should not hand-edit encrypted values. + +CREATE EXTENSION IF NOT EXISTS pgcrypto; + +CREATE OR REPLACE FUNCTION update_updated_at_column() +RETURNS TRIGGER AS $$ +BEGIN + NEW.updated_at = CURRENT_TIMESTAMP; + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +CREATE TABLE IF NOT EXISTS projects ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + name VARCHAR(100) NOT NULL, + code VARCHAR(50) UNIQUE NOT NULL, + description TEXT, + gs_workspace VARCHAR(100) UNIQUE NOT NULL, + map_extent JSONB, + status VARCHAR(20) DEFAULT 'active' NOT NULL, + created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP NOT NULL, + updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP NOT NULL, + CONSTRAINT projects_status_check CHECK (status IN ('active', 'inactive', 'archived')) +); + +CREATE INDEX IF NOT EXISTS idx_projects_status ON projects(status); +CREATE INDEX IF NOT EXISTS idx_projects_code ON projects(code); + +DROP TRIGGER IF EXISTS update_projects_updated_at ON projects; +CREATE TRIGGER update_projects_updated_at + BEFORE UPDATE ON projects + FOR EACH ROW + EXECUTE FUNCTION update_updated_at_column(); + +CREATE TABLE IF NOT EXISTS project_databases ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + project_id UUID NOT NULL REFERENCES projects(id) ON DELETE CASCADE, + db_role VARCHAR(20) NOT NULL, + db_type VARCHAR(20) NOT NULL, + dsn_encrypted TEXT NOT NULL, + pool_min_size INTEGER DEFAULT 2 NOT NULL, + pool_max_size INTEGER DEFAULT 10 NOT NULL, + CONSTRAINT project_databases_unique_role UNIQUE (project_id, db_role), + CONSTRAINT project_databases_role_check CHECK (db_role IN ('biz_data', 'iot_data')), + CONSTRAINT project_databases_type_check CHECK (db_type IN ('postgresql', 'timescaledb')), + CONSTRAINT project_databases_pool_check CHECK ( + pool_min_size >= 1 AND pool_max_size >= pool_min_size + ) +); + +CREATE INDEX IF NOT EXISTS idx_project_databases_project_id + ON project_databases(project_id); +CREATE INDEX IF NOT EXISTS idx_project_databases_role + ON project_databases(db_role); + +CREATE TABLE IF NOT EXISTS project_geoserver_configs ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + project_id UUID UNIQUE NOT NULL REFERENCES projects(id) ON DELETE CASCADE, + gs_base_url TEXT, + gs_admin_user VARCHAR(50), + gs_admin_password_encrypted TEXT, + gs_datastore_name VARCHAR(100) DEFAULT 'ds_postgis' NOT NULL, + default_extent JSONB, + srid INTEGER DEFAULT 4326 NOT NULL, + updated_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP NOT NULL, + CONSTRAINT project_geoserver_configs_srid_check CHECK (srid >= 1) +); + +CREATE INDEX IF NOT EXISTS idx_project_geoserver_configs_project_id + ON project_geoserver_configs(project_id); + +DROP TRIGGER IF EXISTS update_project_geoserver_configs_updated_at + ON project_geoserver_configs; +CREATE TRIGGER update_project_geoserver_configs_updated_at + BEFORE UPDATE ON project_geoserver_configs + FOR EACH ROW + EXECUTE FUNCTION update_updated_at_column(); diff --git a/scripts/migrate_local_users_to_metadata.py b/scripts/migrate_local_users_to_metadata.py new file mode 100644 index 0000000..55edbc3 --- /dev/null +++ b/scripts/migrate_local_users_to_metadata.py @@ -0,0 +1,63 @@ +#!/usr/bin/env python3 +"""Build metadata user sync payloads from an old-user to Keycloak mapping CSV.""" + +from __future__ import annotations + +import argparse +import csv +import json +from pathlib import Path + + +REQUIRED_COLUMNS = {"keycloak_id", "username", "email"} + + +def parse_bool(value: str | None) -> bool: + if value is None or value == "": + return True + return value.strip().lower() not in {"0", "false", "no", "n", "disabled"} + + +def build_payload(mapping_csv: Path) -> dict: + with mapping_csv.open(newline="", encoding="utf-8") as handle: + reader = csv.DictReader(handle) + missing = REQUIRED_COLUMNS.difference(reader.fieldnames or []) + if missing: + raise SystemExit(f"missing required CSV columns: {', '.join(sorted(missing))}") + + users = [] + for row in reader: + users.append( + { + "keycloak_id": row["keycloak_id"].strip(), + "username": row["username"].strip(), + "email": row["email"].strip(), + "role": (row.get("role") or "user").strip().lower(), + "is_active": parse_bool(row.get("is_active")), + } + ) + + return {"users": users} + + +def main() -> None: + parser = argparse.ArgumentParser( + description=( + "Convert old local user mappings into a JSON body for " + "POST /api/v1/admin/users/sync/batch. Passwords are never migrated." + ) + ) + parser.add_argument("mapping_csv", type=Path) + parser.add_argument("-o", "--output", type=Path) + args = parser.parse_args() + + payload = build_payload(args.mapping_csv) + content = json.dumps(payload, ensure_ascii=False, indent=2) + if args.output: + args.output.write_text(content + "\n", encoding="utf-8") + else: + print(content) + + +if __name__ == "__main__": + main() diff --git a/tests/api/test_admin_metadata_endpoints.py b/tests/api/test_admin_metadata_endpoints.py new file mode 100644 index 0000000..e76c4fa --- /dev/null +++ b/tests/api/test_admin_metadata_endpoints.py @@ -0,0 +1,646 @@ +from datetime import datetime, timezone +from types import SimpleNamespace +from unittest.mock import AsyncMock +from uuid import uuid4 + +import pytest +from fastapi import HTTPException +from fastapi import Response + +from app.auth.metadata_dependencies import get_current_metadata_admin +from app.api.v1.endpoints import admin_metadata +from app.domain.schemas.admin_metadata import ( + AdminProjectCreateRequest, + MetadataUsersBatchSyncRequest, + MetadataUserSyncRequest, + MetadataUserUpdateRequest, + ProjectDatabaseUpsertRequest, + ProjectGeoServerConfigUpsertRequest, + ProjectMemberCreateRequest, + ProjectMemberUpdateRequest, +) +from app.infra.db.metadb.repositories.metadata_repository import ProjectDbRouting + + +@pytest.fixture +def anyio_backend(): + return "asyncio" + + +def _user(**overrides): + data = { + "id": uuid4(), + "keycloak_id": uuid4(), + "username": "alice", + "email": "alice@example.com", + "role": "user", + "is_active": True, + "is_superuser": False, + "created_at": datetime(2026, 1, 1, tzinfo=timezone.utc), + "updated_at": datetime(2026, 1, 1, tzinfo=timezone.utc), + "last_login_at": None, + } + data.update(overrides) + return SimpleNamespace(**data) + + +def _project(**overrides): + data = {"id": uuid4(), "name": "Demo"} + data.update(overrides) + return SimpleNamespace(**data) + + +def _membership(**overrides): + data = { + "id": uuid4(), + "user_id": uuid4(), + "project_id": uuid4(), + "project_role": "viewer", + } + data.update(overrides) + return SimpleNamespace(**data) + + +def _database_config(**overrides): + data = { + "id": uuid4(), + "project_id": uuid4(), + "db_role": "biz_data", + "db_type": "postgresql", + "dsn_encrypted": "encrypted-dsn", + "pool_min_size": 1, + "pool_max_size": 5, + } + data.update(overrides) + return SimpleNamespace(**data) + + +def _geoserver_config(**overrides): + data = { + "id": uuid4(), + "project_id": uuid4(), + "gs_base_url": "http://geoserver", + "gs_admin_user": "admin", + "gs_admin_password_encrypted": "encrypted-password", + "gs_datastore_name": "ds_postgis", + "default_extent": {"bbox": [1, 2, 3, 4]}, + "srid": 4326, + "updated_at": datetime(2026, 1, 1, tzinfo=timezone.utc), + } + data.update(overrides) + return SimpleNamespace(**data) + + +def test_to_async_sqlalchemy_url_preserves_password(): + url = admin_metadata._to_async_sqlalchemy_url( + "postgresql://tjwater:secret@192.168.1.114:5433/tjwater" + ) + + assert url == "postgresql+psycopg://tjwater:secret@192.168.1.114:5433/tjwater" + assert "***" not in url + + +@pytest.mark.anyio +async def test_sync_metadata_user_upserts_without_password(monkeypatch): + keycloak_id = uuid4() + synced_user = _user(keycloak_id=keycloak_id, username="new-user") + admin = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=object(), + upsert_user_from_keycloak=AsyncMock(return_value=synced_user), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.sync_metadata_user( + MetadataUserSyncRequest( + keycloak_id=keycloak_id, + username="new-user", + email="new-user@example.com", + role="user", + is_active=True, + ), + current_user=admin, + metadata_repo=repo, + ) + + repo.upsert_user_from_keycloak.assert_awaited_once() + kwargs = repo.upsert_user_from_keycloak.await_args.kwargs + assert kwargs["keycloak_id"] == keycloak_id + assert "password" not in kwargs + assert response.username == "new-user" + admin_metadata.log_audit_event.assert_awaited_once() + + +@pytest.mark.anyio +async def test_batch_sync_metadata_users_returns_per_user_results(monkeypatch): + users = [_user(username="alice"), _user(username="bob")] + admin = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=SimpleNamespace(rollback=AsyncMock()), + upsert_user_from_keycloak=AsyncMock(side_effect=users), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.sync_metadata_users_batch( + MetadataUsersBatchSyncRequest( + users=[ + MetadataUserSyncRequest( + keycloak_id=users[0].keycloak_id, + username="alice", + email="alice@example.com", + role="user", + is_active=True, + ), + MetadataUserSyncRequest( + keycloak_id=users[1].keycloak_id, + username="bob", + email="bob@example.com", + role="viewer", + is_active=True, + ), + ] + ), + current_user=admin, + metadata_repo=repo, + ) + + assert [item.success for item in response] == [True, True] + assert [item.user.username for item in response] == ["alice", "bob"] + assert repo.upsert_user_from_keycloak.await_count == 2 + assert admin_metadata.log_audit_event.await_count == 2 + + +@pytest.mark.anyio +async def test_update_metadata_user_updates_role_and_active_status(monkeypatch): + user_id = uuid4() + updated = _user(id=user_id, role="operator", is_active=False) + repo = SimpleNamespace( + session=object(), + update_user_admin=AsyncMock(return_value=updated), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.update_metadata_user( + MetadataUserUpdateRequest( + role="operator", + is_active=False, + ), + user_id=user_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + repo.update_user_admin.assert_awaited_once_with( + user_id, + updates={"role": "operator", "is_active": False}, + ) + assert response.role == "operator" + admin_metadata.log_audit_event.assert_awaited_once() + + +@pytest.mark.anyio +async def test_update_metadata_user_rejects_self_update(monkeypatch): + current_user = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=object(), + update_user_admin=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.update_metadata_user( + MetadataUserUpdateRequest(role="viewer"), + user_id=current_user.id, + current_user=current_user, + metadata_repo=repo, + ) + + assert exc.value.status_code == 403 + repo.update_user_admin.assert_not_called() + admin_metadata.log_audit_event.assert_not_called() + + +@pytest.mark.anyio +async def test_create_project_audits_metadata_admin_change(monkeypatch): + project = SimpleNamespace( + id=uuid4(), + name="Demo Project", + code="demo", + description="desc", + gs_workspace="demo_ws", + map_extent={"bbox": [1, 2, 3, 4]}, + status="active", + created_at=datetime(2026, 1, 1, tzinfo=timezone.utc), + updated_at=datetime(2026, 1, 1, tzinfo=timezone.utc), + ) + repo = SimpleNamespace( + session=object(), + create_project=AsyncMock(return_value=project), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.create_admin_project( + AdminProjectCreateRequest( + name="Demo Project", + code="demo", + description="desc", + gs_workspace="demo_ws", + map_extent={"bbox": [1, 2, 3, 4]}, + status="active", + ), + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.project_id == project.id + repo.create_project.assert_awaited_once() + admin_metadata.log_audit_event.assert_awaited_once() + + +@pytest.mark.anyio +async def test_upsert_project_database_hides_dsn_and_audits_without_plaintext(monkeypatch): + project_id = uuid4() + record = _database_config(project_id=project_id) + repo = SimpleNamespace( + session=object(), + get_project_by_id=AsyncMock(return_value=_project(id=project_id)), + upsert_project_database_config=AsyncMock(return_value=record), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + monkeypatch.setattr(admin_metadata, "_check_database_connection", AsyncMock()) + + response = await admin_metadata.upsert_project_database( + ProjectDatabaseUpsertRequest( + db_role="biz_data", + dsn="postgresql://user:secret@localhost/db", + pool_min_size=1, + pool_max_size=5, + ), + project_id=project_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.has_dsn is True + assert "dsn" not in response.model_dump() + admin_metadata._check_database_connection.assert_awaited_once() + repo.upsert_project_database_config.assert_awaited_once() + assert repo.upsert_project_database_config.await_args.kwargs["db_type"] == "postgresql" + request_data = admin_metadata.log_audit_event.await_args.kwargs["request_data"] + assert request_data["dsn_updated"] is True + assert request_data["db_type"] == "postgresql" + assert "dsn" not in request_data + assert "postgresql://user:secret@localhost/db" not in str(request_data) + + +@pytest.mark.anyio +async def test_upsert_project_database_rejects_unhealthy_connection(monkeypatch): + project_id = uuid4() + repo = SimpleNamespace( + session=object(), + get_project_by_id=AsyncMock(return_value=_project(id=project_id)), + upsert_project_database_config=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + monkeypatch.setattr( + admin_metadata, + "_check_database_connection", + AsyncMock( + side_effect=Exception( + 'FATAL: password authentication failed for user "tjwater"' + ) + ), + ) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.upsert_project_database( + ProjectDatabaseUpsertRequest( + db_role="iot_data", + dsn="postgresql://tjwater:bad@192.168.1.114:5433/tjwater", + pool_min_size=1, + pool_max_size=5, + ), + project_id=project_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert exc.value.status_code == 400 + assert exc.value.detail == "连通性测试失败:用户名或密码错误,请检查 DSN 中的账号密码。" + repo.upsert_project_database_config.assert_not_called() + admin_metadata.log_audit_event.assert_not_called() + + +@pytest.mark.anyio +async def test_project_database_health_returns_ok(monkeypatch): + project_id = uuid4() + repo = SimpleNamespace( + get_project_db_routing=AsyncMock( + return_value=ProjectDbRouting( + project_id=project_id, + db_role="biz_data", + db_type="postgresql", + dsn="postgresql://user:secret@localhost/db", + pool_min_size=1, + pool_max_size=5, + ) + ) + ) + monkeypatch.setattr(admin_metadata, "_check_database_connection", AsyncMock()) + + response = await admin_metadata.check_project_database_health( + project_id=project_id, + db_role="biz_data", + response=Response(), + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.ok is True + assert response.detail == "连通性测试通过" + admin_metadata._check_database_connection.assert_awaited_once() + + +@pytest.mark.anyio +async def test_project_database_health_can_test_unsaved_plaintext_dsn(monkeypatch): + project_id = uuid4() + repo = SimpleNamespace( + get_project_db_routing=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "_check_database_connection", AsyncMock()) + + response = await admin_metadata.check_project_database_health( + project_id=project_id, + db_role="iot_data", + payload=admin_metadata.ProjectDatabaseHealthRequest( + dsn="postgresql://tjwater:secret@192.168.1.114:5433/tjwater" + ), + response=Response(), + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.ok is True + assert response.db_type == "timescaledb" + routing = admin_metadata._check_database_connection.await_args.args[0] + assert routing.dsn == "postgresql://tjwater:secret@192.168.1.114:5433/tjwater" + repo.get_project_db_routing.assert_not_called() + + +@pytest.mark.anyio +async def test_project_database_health_sanitizes_password_failures(monkeypatch): + project_id = uuid4() + repo = SimpleNamespace( + get_project_db_routing=AsyncMock( + return_value=ProjectDbRouting( + project_id=project_id, + db_role="iot_data", + db_type="timescaledb", + dsn="postgresql://tjwater:bad-password@192.168.1.114:5433/db", + pool_min_size=1, + pool_max_size=5, + ) + ) + ) + monkeypatch.setattr( + admin_metadata, + "_check_database_connection", + AsyncMock( + side_effect=Exception( + '(psycopg.OperationalError) connection failed: FATAL: ' + 'password authentication failed for user "tjwater"' + ) + ), + ) + + fastapi_response = Response() + response = await admin_metadata.check_project_database_health( + project_id=project_id, + db_role="iot_data", + response=fastapi_response, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert fastapi_response.status_code == 503 + assert response.ok is False + assert response.db_type == "timescaledb" + assert response.detail == "连通性测试失败:用户名或密码错误,请检查 DSN 中的账号密码。" + assert "psycopg" not in response.detail + + +@pytest.mark.anyio +async def test_upsert_geoserver_config_hides_password_and_audits_without_plaintext(monkeypatch): + project_id = uuid4() + record = _geoserver_config(project_id=project_id) + repo = SimpleNamespace( + session=object(), + get_project_by_id=AsyncMock(return_value=_project(id=project_id)), + upsert_geoserver_config=AsyncMock(return_value=record), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.upsert_project_geoserver_config( + ProjectGeoServerConfigUpsertRequest( + gs_base_url="http://geoserver", + gs_admin_user="admin", + gs_admin_password="secret-password", + gs_datastore_name="ds_postgis", + default_extent={"bbox": [1, 2, 3, 4]}, + srid=4326, + ), + project_id=project_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.has_password is True + assert "password" not in response.model_dump() + request_data = admin_metadata.log_audit_event.await_args.kwargs["request_data"] + assert request_data["password_updated"] is True + assert "secret-password" not in str(request_data) + + +@pytest.mark.anyio +async def test_get_geoserver_config_returns_empty_state_when_unconfigured(): + project_id = uuid4() + repo = SimpleNamespace( + get_project_by_id=AsyncMock(return_value=_project(id=project_id)), + get_geoserver_config_record=AsyncMock(return_value=None), + ) + + response = await admin_metadata.get_project_geoserver_config( + project_id=project_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.project_id == project_id + assert response.configured is False + assert response.has_password is False + assert response.gs_datastore_name == "ds_postgis" + + +@pytest.mark.anyio +async def test_metadata_admin_dependency_rejects_non_admin_user(): + with pytest.raises(HTTPException) as exc: + await get_current_metadata_admin(_user(role="user", is_superuser=False)) + + assert exc.value.status_code == 403 + assert exc.value.detail == "Admin access required" + + +@pytest.mark.anyio +async def test_add_project_member_rejects_duplicate(monkeypatch): + project_id = uuid4() + user_id = uuid4() + repo = SimpleNamespace( + session=object(), + get_project_by_id=AsyncMock(return_value=_project(id=project_id)), + get_user_by_id=AsyncMock(return_value=_user(id=user_id)), + get_project_membership=AsyncMock(return_value=_membership()), + add_project_member=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.add_project_member( + ProjectMemberCreateRequest(user_id=user_id, project_role="viewer"), + project_id=project_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert exc.value.status_code == 409 + repo.add_project_member.assert_not_called() + + +@pytest.mark.anyio +async def test_add_project_member_rejects_self_membership_change(monkeypatch): + project_id = uuid4() + current_user = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=object(), + get_project_by_id=AsyncMock(), + add_project_member=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.add_project_member( + ProjectMemberCreateRequest( + user_id=current_user.id, + project_role="viewer", + ), + project_id=project_id, + current_user=current_user, + metadata_repo=repo, + ) + + assert exc.value.status_code == 403 + repo.get_project_by_id.assert_not_called() + repo.add_project_member.assert_not_called() + admin_metadata.log_audit_event.assert_not_called() + + +@pytest.mark.anyio +async def test_update_project_member_role_audits_change(monkeypatch): + project_id = uuid4() + user_id = uuid4() + user = _user(id=user_id, username="bob", email="bob@example.com") + membership = _membership( + user_id=user_id, + project_id=project_id, + project_role="admin", + ) + repo = SimpleNamespace( + session=object(), + get_user_by_id=AsyncMock(return_value=user), + update_project_member_role=AsyncMock(return_value=membership), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.update_project_member( + ProjectMemberUpdateRequest(project_role="admin"), + project_id=project_id, + user_id=user_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response.project_role == "admin" + repo.update_project_member_role.assert_awaited_once_with( + project_id, user_id, "admin" + ) + admin_metadata.log_audit_event.assert_awaited_once() + + +@pytest.mark.anyio +async def test_update_project_member_rejects_self_membership_change(monkeypatch): + project_id = uuid4() + current_user = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=object(), + get_user_by_id=AsyncMock(), + update_project_member_role=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.update_project_member( + ProjectMemberUpdateRequest(project_role="admin"), + project_id=project_id, + user_id=current_user.id, + current_user=current_user, + metadata_repo=repo, + ) + + assert exc.value.status_code == 403 + repo.get_user_by_id.assert_not_called() + repo.update_project_member_role.assert_not_called() + admin_metadata.log_audit_event.assert_not_called() + + +@pytest.mark.anyio +async def test_remove_project_member_audits_change(monkeypatch): + project_id = uuid4() + user_id = uuid4() + repo = SimpleNamespace( + session=object(), + remove_project_member=AsyncMock(return_value=True), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + response = await admin_metadata.remove_project_member( + project_id=project_id, + user_id=user_id, + current_user=_user(role="admin", is_superuser=True), + metadata_repo=repo, + ) + + assert response is None + repo.remove_project_member.assert_awaited_once_with(project_id, user_id) + admin_metadata.log_audit_event.assert_awaited_once() + + +@pytest.mark.anyio +async def test_remove_project_member_rejects_self_membership_change(monkeypatch): + project_id = uuid4() + current_user = _user(role="admin", is_superuser=True) + repo = SimpleNamespace( + session=object(), + remove_project_member=AsyncMock(), + ) + monkeypatch.setattr(admin_metadata, "log_audit_event", AsyncMock()) + + with pytest.raises(HTTPException) as exc: + await admin_metadata.remove_project_member( + project_id=project_id, + user_id=current_user.id, + current_user=current_user, + metadata_repo=repo, + ) + + assert exc.value.status_code == 403 + repo.remove_project_member.assert_not_called() + admin_metadata.log_audit_event.assert_not_called() diff --git a/tests/auth/test_metadata_dependencies.py b/tests/auth/test_metadata_dependencies.py new file mode 100644 index 0000000..3abfbd8 --- /dev/null +++ b/tests/auth/test_metadata_dependencies.py @@ -0,0 +1,83 @@ +from datetime import datetime, timezone +from types import SimpleNamespace +from unittest.mock import AsyncMock +from uuid import uuid4 + +import pytest +from fastapi import HTTPException + +from app.auth import metadata_dependencies + + +@pytest.fixture +def anyio_backend(): + return "asyncio" + + +def _user(**overrides): + data = { + "id": uuid4(), + "keycloak_id": uuid4(), + "username": "old-name", + "email": "old@example.com", + "role": "user", + "is_active": True, + "is_superuser": False, + "created_at": datetime(2026, 1, 1, tzinfo=timezone.utc), + "updated_at": datetime(2026, 1, 1, tzinfo=timezone.utc), + "last_login_at": None, + } + data.update(overrides) + return SimpleNamespace(**data) + + +@pytest.mark.anyio +async def test_current_metadata_user_refreshes_keycloak_claim_snapshot(): + keycloak_id = uuid4() + user = _user(keycloak_id=keycloak_id) + refreshed = _user( + id=user.id, + keycloak_id=keycloak_id, + username="alice", + email="alice@example.com", + last_login_at=datetime(2026, 6, 12, tzinfo=timezone.utc), + ) + repo = SimpleNamespace( + get_user_by_keycloak_id=AsyncMock(return_value=user), + refresh_user_keycloak_snapshot=AsyncMock(return_value=refreshed), + ) + + response = await metadata_dependencies.get_current_metadata_user( + { + "sub": str(keycloak_id), + "preferred_username": "alice", + "email": "alice@example.com", + }, + metadata_repo=repo, + ) + + assert response.username == "alice" + repo.get_user_by_keycloak_id.assert_awaited_once_with(keycloak_id) + repo.refresh_user_keycloak_snapshot.assert_awaited_once_with( + user, + username="alice", + email="alice@example.com", + ) + + +@pytest.mark.anyio +async def test_current_metadata_user_rejects_invalid_keycloak_sub(): + repo = SimpleNamespace( + get_user_by_keycloak_id=AsyncMock(), + refresh_user_keycloak_snapshot=AsyncMock(), + ) + + with pytest.raises(HTTPException) as exc: + await metadata_dependencies.get_current_metadata_user( + {"sub": "not-a-uuid"}, + metadata_repo=repo, + ) + + assert exc.value.status_code == 401 + repo.get_user_by_keycloak_id.assert_not_called() + repo.refresh_user_keycloak_snapshot.assert_not_called() diff --git a/tests/unit/test_dynamic_manager.py b/tests/unit/test_dynamic_manager.py new file mode 100644 index 0000000..e8dd81d --- /dev/null +++ b/tests/unit/test_dynamic_manager.py @@ -0,0 +1,12 @@ +from app.infra.db.dynamic_manager import ProjectConnectionManager + + +def test_normalize_pg_url_preserves_password(): + manager = ProjectConnectionManager() + + url = manager._normalize_pg_url( + "postgresql://tjwater:secret@192.168.1.114:5433/tjwater" + ) + + assert url == "postgresql+psycopg://tjwater:secret@192.168.1.114:5433/tjwater" + assert "***" not in url diff --git a/tests/unit/test_metadata_repository_dsn_decrypt.py b/tests/unit/test_metadata_repository_dsn_decrypt.py index e1a74f8..0548d9e 100644 --- a/tests/unit/test_metadata_repository_dsn_decrypt.py +++ b/tests/unit/test_metadata_repository_dsn_decrypt.py @@ -23,6 +23,11 @@ class _DummyEncryptor: self._raise_invalid_token = raise_invalid_token self.encrypted_values = [] + def encrypt(self, value): + encrypted = f"encrypted::{value}" + self.encrypted_values.append(value) + return encrypted + def decrypt(self, _value): if self._raise_invalid_token: raise InvalidToken() @@ -117,3 +122,46 @@ def test_encrypted_dsn_decrypts_without_migration(monkeypatch): assert routing.dsn == "postgresql://u:p%40ss@host/db" session.commit.assert_not_awaited() + + +def test_upsert_project_database_config_encrypts_plaintext_dsn(monkeypatch): + project_id = uuid4() + session = SimpleNamespace( + execute=None, + add=None, + commit=None, + refresh=None, + ) + added = [] + session.execute = AsyncMock(return_value=_DummyResult(None)) + session.add = lambda item: added.append(item) + session.commit = AsyncMock() + session.refresh = AsyncMock() + encryptor = _DummyEncryptor() + repo = MetadataRepository(session) + + monkeypatch.setattr( + "app.infra.db.metadb.repositories.metadata_repository.is_database_encryption_configured", + lambda: True, + ) + monkeypatch.setattr( + "app.infra.db.metadb.repositories.metadata_repository.get_database_encryptor", + lambda: encryptor, + ) + + record = asyncio.run( + repo.upsert_project_database_config( + project_id, + db_role="biz_data", + db_type="postgresql", + dsn="postgresql://user:secret@localhost/db", + pool_min_size=1, + pool_max_size=5, + ) + ) + + assert encryptor.encrypted_values == ["postgresql://user:secret@localhost/db"] + assert record.dsn_encrypted == "encrypted::postgresql://user:secret@localhost/db" + assert added == [record] + session.commit.assert_awaited_once() + session.refresh.assert_awaited_once_with(record)