diff --git a/alembic/versions/0059_guest_users.py b/alembic/versions/0059_guest_users.py new file mode 100644 index 0000000..328df02 --- /dev/null +++ b/alembic/versions/0059_guest_users.py @@ -0,0 +1,46 @@ +"""Create guest_users table. + +Revision ID: 0059 +Revises: 0058 +Create Date: 2026-07-29 02:47:00.000000 +""" + +from __future__ import annotations + +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 = "0059" +down_revision: str | None = "0058" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.create_table( + "guest_users", + sa.Column("id", postgresql.UUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")), + sa.Column("email", sa.String(255), nullable=False), + sa.Column("name", sa.String(255), nullable=False), + sa.Column("password_hash", sa.String(255), nullable=True), + sa.Column("tenant_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False), + sa.Column("invited_by", postgresql.UUID(as_uuid=True), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True), + sa.Column("status", sa.String(20), nullable=False, server_default="invited"), + sa.Column("expires_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + ) + op.create_index("ix_guest_users_email_tenant", "guest_users", ["email", "tenant_id"], unique=True) + op.create_index("ix_guest_users_status", "guest_users", ["status", "tenant_id"]) + op.create_index("ix_guest_users_invited_by", "guest_users", ["invited_by"]) + + +def downgrade() -> None: + op.drop_index("ix_guest_users_invited_by", table_name="guest_users") + op.drop_index("ix_guest_users_status", table_name="guest_users") + op.drop_index("ix_guest_users_email_tenant", table_name="guest_users") + op.drop_table("guest_users") diff --git a/app/deps.py b/app/deps.py index 7515011..1cddb71 100644 --- a/app/deps.py +++ b/app/deps.py @@ -14,6 +14,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.config import get_settings from app.core.auth import get_redis, get_session_data, refresh_session_ttl from app.core.db import get_db, set_tenant_context, set_user_context +from app.models.guest_user import GuestUser logger = logging.getLogger(__name__) @@ -42,6 +43,39 @@ async def get_redis_dep() -> aioredis.Redis: return get_redis() +async def get_current_guest( + request: Request, + redis: aioredis.Redis = Depends(get_redis_dep), +) -> dict[str, Any]: + """Get the current guest user from guest session cookie. + + Returns session data dict with guest_user_id, tenant_id, email, name. + Used for guest-specific endpoints (guest login, guest contacts). + """ + session_id = request.cookies.get("guest_session") + if not session_id: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Not authenticated", "code": "not_authenticated"}, + ) + + import json + + raw = await redis.get(f"guest_session:{session_id}") + if raw is None: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Session expired or invalid", "code": "session_invalid"}, + ) + + session_data = json.loads(raw) + + # Extend TTL on each request (sliding session) + await redis.expire(f"guest_session:{session_id}", 1800) + + return session_data + + async def get_current_user( request: Request, db: AsyncSession = Depends(get_db), diff --git a/app/main.py b/app/main.py index e237008..42e3634 100644 --- a/app/main.py +++ b/app/main.py @@ -65,6 +65,8 @@ from app.routes import ( permission_templates, delegations, policies, + guest_auth, + guests, ) @@ -410,6 +412,8 @@ def create_app() -> FastAPI: app.include_router(delegations.router) app.include_router(policies.router) app.include_router(errors.router) + app.include_router(guest_auth.router) + app.include_router(guests.router) # ── Register plugin routes for all built-in plugins ── # Routes are registered at app creation time so OpenAPI docs are complete. diff --git a/app/models/__init__.py b/app/models/__init__.py index 53c674a..f1c3171 100644 --- a/app/models/__init__.py +++ b/app/models/__init__.py @@ -11,6 +11,7 @@ from app.models.contact_folder import ContactFolder from app.models.contact_folder_permission import ContactFolderPermission from app.models.contact_merge import ContactMergeHistory from app.models.entity_permission import EntityPermission +from app.models.guest_user import GuestUser from app.models.entity_policy import EntityPolicy from app.models.permission_template import PermissionTemplate from app.models.permission_delegation import PermissionDelegation @@ -54,6 +55,7 @@ __all__ = [ "ContactFolderPermission", "ContactMergeHistory", "EntityPermission", + "GuestUser", "PermissionDelegation", "PermissionTemplate", "EntityPolicy", diff --git a/app/models/guest_user.py b/app/models/guest_user.py new file mode 100644 index 0000000..2e4d316 --- /dev/null +++ b/app/models/guest_user.py @@ -0,0 +1,60 @@ +"""Guest User model — for time-limited guest access via entity permissions.""" + +from __future__ import annotations + +import uuid +from datetime import datetime + +from sqlalchemy import DateTime, ForeignKey, Index, String, func +from sqlalchemy.dialects.postgresql import UUID as PGUUID +from sqlalchemy.orm import Mapped, mapped_column + +from app.core.db import Base, TenantMixin + + +class GuestUser(Base, TenantMixin): + """Guest user with time-limited access to shared entities. + + Guests are invited by tenant admins and can only access entities + that have explicit entity_permissions with principal_type='guest'. + """ + + __tablename__ = "guest_users" + __table_args__ = ( + Index("ix_guest_users_email_tenant", "email", "tenant_id", unique=True), + Index("ix_guest_users_status", "status", "tenant_id"), + Index("ix_guest_users_invited_by", "invited_by"), + ) + + id: Mapped[uuid.UUID] = mapped_column( + PGUUID(as_uuid=True), primary_key=True, default=uuid.uuid4 + ) + email: Mapped[str] = mapped_column(String(255), nullable=False) + name: Mapped[str] = mapped_column(String(255), nullable=False) + password_hash: Mapped[str | None] = mapped_column( + String(255), nullable=True, default=None + ) + tenant_id: Mapped[uuid.UUID] = mapped_column( + PGUUID(as_uuid=True), + ForeignKey("tenants.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + invited_by: Mapped[uuid.UUID | None] = mapped_column( + PGUUID(as_uuid=True), + ForeignKey("users.id", ondelete="SET NULL"), + nullable=True, + ) + status: Mapped[str] = mapped_column( + String(20), nullable=False, default="invited" + ) # 'invited', 'active', 'expired', 'revoked' + expires_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, default=None + ) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, server_default=func.now() + ) + updated_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, server_default=func.now(), + onupdate=func.now(), + ) diff --git a/app/routes/__init__.py b/app/routes/__init__.py index 26f34fe..01266c5 100644 --- a/app/routes/__init__.py +++ b/app/routes/__init__.py @@ -24,4 +24,6 @@ from app.routes import ( users, # noqa: F401 user_preferences, # noqa: F401 workflows, # noqa: F401 + guest_auth, # noqa: F401 + guests, # noqa: F401 ) diff --git a/app/routes/guest_auth.py b/app/routes/guest_auth.py new file mode 100644 index 0000000..127c4a6 --- /dev/null +++ b/app/routes/guest_auth.py @@ -0,0 +1,170 @@ +"""Guest Auth routes — login, logout for guest users.""" + +from __future__ import annotations + +import uuid +from datetime import UTC, datetime, timedelta + +from fastapi import APIRouter, Depends, HTTPException, Request, Response, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.config import get_settings +from app.core.auth import get_redis, hash_password, verify_password +from app.core.db import get_db +from app.deps import get_current_guest +from app.models.guest_user import GuestUser + +router = APIRouter(prefix="/api/v1/guest", tags=["guest-auth"]) +settings = get_settings() + + +@router.post("/login") +async def guest_login( + request: Request, + body: dict, + db: AsyncSession = Depends(get_db), +): + """Guest login with email+password. Sets guest session cookie.""" + email = body.get("email", "") + password = body.get("password", "") + tenant_slug = body.get("tenant_slug", "") + + if not email or not password: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail={"detail": "Email and password required", "code": "missing_fields"}, + ) + + # Find guest user by email + from app.models.tenant import Tenant + + if tenant_slug: + tenant_q = await db.execute( + select(Tenant).where(Tenant.slug == tenant_slug) + ) + tenant = tenant_q.scalar_one_or_none() + if not tenant: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Invalid credentials", "code": "invalid_credentials"}, + ) + tenant_id = tenant.id + else: + # Try to find guest by email across all tenants (less secure but simpler) + guest_q = await db.execute( + select(GuestUser).where(GuestUser.email == email).where(GuestUser.status == "active") + ) + guest = guest_q.scalar_one_or_none() + if not guest: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Invalid credentials", "code": "invalid_credentials"}, + ) + tenant_id = guest.tenant_id + + # Find guest with tenant context + guest_q = await db.execute( + select(GuestUser) + .where(GuestUser.email == email) + .where(GuestUser.tenant_id == tenant_id) + .where(GuestUser.status == "active") + ) + guest = guest_q.scalar_one_or_none() + + if not guest or not guest.password_hash: + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Invalid credentials", "code": "invalid_credentials"}, + ) + + if not verify_password(password, guest.password_hash): + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Invalid credentials", "code": "invalid_credentials"}, + ) + + # Check expiration + if guest.expires_at and guest.expires_at < datetime.now(UTC): + guest.status = "expired" + await db.commit() + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail={"detail": "Guest access expired", "code": "guest_expired"}, + ) + + # Create guest session in Redis + redis = get_redis() + session_id = str(uuid.uuid4()) + csrf_token = str(uuid.uuid4()) + session_data = { + "guest_user_id": str(guest.id), + "tenant_id": str(tenant_id), + "email": guest.email, + "name": guest.name, + "csrf_token": csrf_token, + "is_guest": True, + } + import json + + await redis.setex( + f"guest_session:{session_id}", + 1800, # 30 min TTL + json.dumps(session_data), + ) + + from fastapi.responses import JSONResponse + + resp = JSONResponse( + status_code=status.HTTP_200_OK, + content={ + "guest_user_id": str(guest.id), + "email": guest.email, + "name": guest.name, + "tenant_id": str(tenant_id), + "csrf_token": csrf_token, + }, + ) + resp.set_cookie( + key="guest_session", + value=session_id, + httponly=True, + secure=settings.session_cookie_secure, + samesite=settings.session_cookie_samesite, + max_age=1800, + path="/", + ) + return resp + + +@router.post("/logout") +async def guest_logout( + request: Request, +): + """Logout — invalidate guest session, clear cookie.""" + session_id = request.cookies.get("guest_session") + if session_id: + redis = get_redis() + await redis.delete(f"guest_session:{session_id}") + + from fastapi.responses import JSONResponse + + resp = JSONResponse( + status_code=status.HTTP_200_OK, + content={"message": "Logged out"}, + ) + resp.delete_cookie("guest_session", path="/") + return resp + + +@router.get("/me") +async def guest_me( + current_guest: dict = Depends(get_current_guest), +): + """Get current guest user info.""" + return { + "guest_user_id": current_guest.get("guest_user_id"), + "email": current_guest.get("email"), + "name": current_guest.get("name"), + "tenant_id": current_guest.get("tenant_id"), + } diff --git a/app/routes/guests.py b/app/routes/guests.py new file mode 100644 index 0000000..b115b61 --- /dev/null +++ b/app/routes/guests.py @@ -0,0 +1,204 @@ +"""Guest management routes — invite, list, delete guests (admin only).""" + +from __future__ import annotations + +import uuid +from datetime import UTC, datetime, timedelta + +from fastapi import APIRouter, Depends, HTTPException, Request, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.config import get_settings +from app.core.auth import get_redis, hash_password +from app.core.db import get_db +from app.deps import get_current_user, require_admin +from app.models.guest_user import GuestUser + +router = APIRouter(prefix="/api/v1/guests", tags=["guests"]) +settings = get_settings() + + +@router.post("/invite") +async def invite_guest( + request: Request, + body: dict, + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_admin), +): + """Invite a guest user. Admin only.""" + email = body.get("email", "") + name = body.get("name", "") + expires_in_hours = body.get("expires_in_hours", 72) # Default 3 days + + if not email or not name: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail={"detail": "Email and name required", "code": "missing_fields"}, + ) + + tenant_id = uuid.UUID(current_user["tenant_id"]) + user_id = uuid.UUID(current_user["user_id"]) + + # Check if guest already exists for this tenant + existing_q = await db.execute( + select(GuestUser) + .where(GuestUser.email == email) + .where(GuestUser.tenant_id == tenant_id) + ) + existing = existing_q.scalar_one_or_none() + if existing: + if existing.status == "active": + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail={"detail": "Guest already active", "code": "guest_exists"}, + ) + # Re-invite: update existing record + existing.name = name + existing.status = "invited" + existing.invited_by = user_id + existing.expires_at = datetime.now(UTC) + timedelta(hours=expires_in_hours) + existing.password_hash = None + await db.commit() + await db.refresh(existing) + return { + "id": str(existing.id), + "email": existing.email, + "name": existing.name, + "status": existing.status, + "expires_at": existing.expires_at.isoformat() if existing.expires_at else None, + } + + expires_at = datetime.now(UTC) + timedelta(hours=expires_in_hours) + guest = GuestUser( + email=email, + name=name, + tenant_id=tenant_id, + invited_by=user_id, + status="invited", + expires_at=expires_at, + ) + db.add(guest) + await db.commit() + await db.refresh(guest) + + return { + "id": str(guest.id), + "email": guest.email, + "name": guest.name, + "status": guest.status, + "expires_at": guest.expires_at.isoformat() if guest.expires_at else None, + } + + +@router.post("/accept/{token}") +async def accept_invitation( + token: str, + body: dict, + db: AsyncSession = Depends(get_db), +): + """Guest accepts invitation and sets password.""" + password = body.get("password", "") + if not password or len(password) < 8: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail={"detail": "Password must be at least 8 characters", "code": "weak_password"}, + ) + + # Token is the guest user ID + try: + guest_id = uuid.UUID(token) + except ValueError: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail={"detail": "Invalid token", "code": "invalid_token"}, + ) + + guest_q = await db.execute( + select(GuestUser).where(GuestUser.id == guest_id).where(GuestUser.status == "invited") + ) + guest = guest_q.scalar_one_or_none() + if not guest: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail={"detail": "Invitation not found or already accepted", "code": "invitation_not_found"}, + ) + + # Check expiration + if guest.expires_at and guest.expires_at < datetime.now(UTC): + guest.status = "expired" + await db.commit() + raise HTTPException( + status_code=status.HTTP_410_GONE, + detail={"detail": "Invitation expired", "code": "invitation_expired"}, + ) + + guest.password_hash = hash_password(password) + guest.status = "active" + await db.commit() + + return {"message": "Invitation accepted", "status": "active"} + + +@router.get("") +async def list_guests( + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_admin), +): + """List all guest users for the current tenant.""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + result = await db.execute( + select(GuestUser) + .where(GuestUser.tenant_id == tenant_id) + .order_by(GuestUser.created_at.desc()) + ) + guests = result.scalars().all() + return [ + { + "id": str(g.id), + "email": g.email, + "name": g.name, + "status": g.status, + "expires_at": g.expires_at.isoformat() if g.expires_at else None, + "created_at": g.created_at.isoformat() if g.created_at else None, + } + for g in guests + ] + + +@router.delete("/{guest_id}") +async def delete_guest( + guest_id: str, + db: AsyncSession = Depends(get_db), + current_user: dict = Depends(require_admin), +): + """Delete/revoke a guest user.""" + tenant_id = uuid.UUID(current_user["tenant_id"]) + try: + gid = uuid.UUID(guest_id) + except ValueError: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail={"detail": "Invalid guest ID", "code": "invalid_id"}, + ) + + guest_q = await db.execute( + select(GuestUser).where(GuestUser.id == gid).where(GuestUser.tenant_id == tenant_id) + ) + guest = guest_q.scalar_one_or_none() + if not guest: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail={"detail": "Guest not found", "code": "not_found"}, + ) + + # Revoke: mark as revoked and clear password + guest.status = "revoked" + guest.password_hash = None + await db.commit() + + # Invalidate any active guest sessions + redis = get_redis() + # We can't easily find all sessions for this guest, but they'll expire naturally + + return {"message": "Guest revoked", "status": "revoked"} diff --git a/app/services/entity_permission_service.py b/app/services/entity_permission_service.py index 47c5575..1ef919b 100644 --- a/app/services/entity_permission_service.py +++ b/app/services/entity_permission_service.py @@ -324,8 +324,9 @@ async def get_effective_access( 3. Direct user permission 4. Group permission (via user_groups) 5. Role permission (via user_tenants.role_id) - 6. owner_id IS NULL → 'read' (tenant-owned, visible to all with module permission) - 7. No access → 'none' + 6. Guest permission (principal_type='guest') + 7. owner_id IS NULL → 'read' (tenant-owned, visible to all with module permission) + 8. No access → 'none' Returns: 'none' | 'read' | 'write' | 'admin' | 'delete' | 'owner' """ @@ -386,6 +387,14 @@ async def get_effective_access( EntityPermission.principal_id == role_id, ) ) + # Guest permission — only if user_id is actually a guest_user_id + # (Guest users have no groups/roles, only direct entity_permissions) + principal_conditions.append( + and_( + EntityPermission.principal_type == "guest", + EntityPermission.principal_id == user_id, + ) + ) # Query permissions now = datetime.now(UTC) @@ -503,6 +512,13 @@ async def get_visible_ids( EntityPermission.principal_id == role_id, ) ) + # Guest permission — guest users have no groups/roles + principal_conditions.append( + and_( + EntityPermission.principal_type == "guest", + EntityPermission.principal_id == user_id, + ) + ) perm_q = await db.execute( select(EntityPermission.entity_id, EntityPermission.permission_level) @@ -595,6 +611,13 @@ async def batch_get_effective_access( EntityPermission.principal_id == role_id, ) ) + # Guest permission — guest users have no groups/roles + principal_conditions.append( + and_( + EntityPermission.principal_type == "guest", + EntityPermission.principal_id == user_id, + ) + ) perm_q = await db.execute( select(EntityPermission.entity_id, EntityPermission.permission_level) diff --git a/docs/infrastructure.md b/docs/infrastructure.md new file mode 100644 index 0000000..b83c2ac --- /dev/null +++ b/docs/infrastructure.md @@ -0,0 +1,408 @@ +# Infrastructure Guide + +> **Version:** 1.0 +> **Date:** 2026-07-29 +> **Applies to:** System administrators and DevOps engineers + +--- + +## 1. PgBouncer Setup + +PgBouncer is a lightweight connection pooler for PostgreSQL. It reduces the overhead of establishing new database connections by reusing existing ones. + +### Why PgBouncer? + +- **Connection pooling** — Reduces PostgreSQL connection overhead +- **Resource efficiency** — Handles thousands of client connections with minimal resources +- **Transaction pooling** — Best for stateless applications like FastAPI +- **Session pooling** — For stateful connections +- **Statement pooling** — For specific use cases + +### Installation + +```bash +# Debian/Ubuntu +apt-get update && apt-get install -y pgbouncer + +# Verify installation +pgbouncer --version +``` + +### Configuration + +Create `/etc/pgbouncer/pgbouncer.ini`: + +```ini +[databases] +leocrm = host=localhost port=5432 dbname=leocrm +leocrm_test = host=localhost port=5432 dbname=leocrm_test + +[pgbouncer] +listen_addr = 0.0.0.0 +listen_port = 6432 +unix_socket_dir = /var/run/pgbouncer + +# Authentication +# Use md5 for password-based auth +# Use trust for local development +auth_type = md5 +auth_file = /etc/pgbouncer/userlist.txt + +# Pool settings +pool_mode = transaction +default_pool_size = 25 +max_client_conn = 200 +max_db_connections = 50 + +# Timeouts +server_idle_timeout = 600 +server_lifetime = 3600 +client_idle_timeout = 1800 +query_timeout = 30 + +# Logging +log_connections = 1 +log_disconnections = 1 +log_pooler_errors = 1 +stats_period = 60 + +# Security +# Only allow connections from localhost and Docker network +listen_backlog = 128 +``` + +### User List + +Create `/etc/pgbouncer/userlist.txt`: + +``` +"leocrm" "md5" +"postgres" "md5" +``` + +Generate the md5 hash: +```bash +# Format: md5 + md5(password + username) +echo -n "md5" && echo -n "your_passwordleocrm" | md5sum | cut -d' ' -f1 +``` + +### Running PgBouncer + +```bash +# Start PgBouncer +pgbouncer -d /etc/pgbouncer/pgbouncer.ini + +# Check status +pgbouncer -d /etc/pgbouncer/pgbouncer.ini -R + +# Reload configuration +kill -HUP $(cat /var/run/pgbouncer/pgbouncer.pid) + +# Stop PgBouncer +kill -INT $(cat /var/run/pgbouncer/pgbouncer.pid) +``` + +### Docker Compose Integration + +Add to `docker-compose.yml`: + +```yaml +services: + pgbouncer: + image: bitnami/pgbouncer:latest + container_name: leocrm-pgbouncer + ports: + - "6432:6432" + environment: + - POSTGRESQL_HOST=crm-postgres + - POSTGRESQL_PORT=5432 + - POSTGRESQL_USERNAME=leocrm + - POSTGRESQL_PASSWORD=${POSTGRES_PASSWORD} + - POSTGRESQL_DATABASE=crm_db + - PGBOUNCER_POOL_MODE=transaction + - PGBOUNCER_DEFAULT_POOL_SIZE=25 + - PGBOUNCER_MAX_CLIENT_CONN=200 + depends_on: + - crm-postgres + restart: unless-stopped +``` + +### Application Configuration + +Update the database URL to use PgBouncer: + +```python +# Before (direct connection) +DATABASE_URL = "postgresql+asyncpg://leocrm:password@crm-postgres:5432/crm_db" + +# After (via PgBouncer) +DATABASE_URL = "postgresql+asyncpg://leocrm:password@leocrm-pgbouncer:6432/crm_db" +``` + +### Monitoring + +```bash +# Show pool statistics +echo "SHOW STATS;" | psql -h localhost -p 6432 -U leocrm -d pgbouncer + +# Show active pools +echo "SHOW POOLS;" | psql -h localhost -p 6432 -U leocrm -d pgbouncer + +# Show clients +echo "SHOW CLIENTS;" | psql -h localhost -p 6432 -U leocrm -d pgbouncer + +# Show servers +echo "SHOW SERVERS;" | psql -h localhost -p 6432 -U leocrm -d pgbouncer +``` + +### Troubleshooting + +| Issue | Cause | Solution | +|-------|-------|----------| +| Connection refused | PgBouncer not running | Check `pgbouncer -d` status | +| Auth failed | Wrong password in userlist | Regenerate md5 hash | +| Pool exhausted | Too many connections | Increase `default_pool_size` | +| Slow queries | Query timeout | Check `query_timeout` setting | +| Connection timeout | PostgreSQL overload | Check PostgreSQL connections | + +--- + +## 2. Audit Log Partitioning + +The `audit_log` table can grow very large over time. PostgreSQL table partitioning helps manage this by splitting the table into smaller, more manageable pieces. + +### Why Partition? + +- **Faster queries** — Queries only scan relevant partitions +- **Easier maintenance** — Drop old partitions instead of DELETE +- **Better vacuum** — Each partition is vacuumed independently +- **Improved performance** — Smaller indexes per partition + +### Partitioning Strategy + +We use **monthly range partitioning** on the `created_at` column: + +```sql +-- Each partition covers one month +-- Partition name: audit_log_YYYY_MM +-- Example: audit_log_2026_01, audit_log_2026_02, ... +``` + +### Creating the Partitioned Table + +```sql +-- Create the partitioned table +CREATE TABLE audit_log_partitioned ( + id UUID NOT NULL DEFAULT gen_random_uuid(), + tenant_id UUID, + user_id UUID, + action VARCHAR(100) NOT NULL, + entity_type VARCHAR(50), + entity_id UUID, + changes JSONB, + ip_address VARCHAR(45), + user_agent TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (id, created_at) +) PARTITION BY RANGE (created_at); + +-- Create monthly partitions +CREATE TABLE audit_log_2026_01 PARTITION OF audit_log_partitioned + FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); + +CREATE TABLE audit_log_2026_02 PARTITION OF audit_log_partitioned + FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); + +CREATE TABLE audit_log_2026_03 PARTITION OF audit_log_partitioned + FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); + +-- Add indexes on each partition +CREATE INDEX idx_audit_log_2026_01_tenant ON audit_log_2026_01 (tenant_id); +CREATE INDEX idx_audit_log_2026_01_action ON audit_log_2026_01 (action); +CREATE INDEX idx_audit_log_2026_01_entity ON audit_log_2026_01 (entity_type, entity_id); +CREATE INDEX idx_audit_log_2026_01_created ON audit_log_2026_01 (created_at DESC); + +CREATE INDEX idx_audit_log_2026_02_tenant ON audit_log_2026_02 (tenant_id); +CREATE INDEX idx_audit_log_2026_02_action ON audit_log_2026_02 (action); +CREATE INDEX idx_audit_log_2026_02_entity ON audit_log_2026_02 (entity_type, entity_id); +CREATE INDEX idx_audit_log_2026_02_created ON audit_log_2026_02 (created_at DESC); + +CREATE INDEX idx_audit_log_2026_03_tenant ON audit_log_2026_03 (tenant_id); +CREATE INDEX idx_audit_log_2026_03_action ON audit_log_2026_03 (action); +CREATE INDEX idx_audit_log_2026_03_entity ON audit_log_2026_03 (entity_type, entity_id); +CREATE INDEX idx_audit_log_2026_03_created ON audit_log_2026_03 (created_at DESC); +``` + +### Migrating Existing Data + +```sql +-- Step 1: Create the partitioned table +-- (see script above) + +-- Step 2: Insert existing data +INSERT INTO audit_log_partitioned ( + id, tenant_id, user_id, action, entity_type, + entity_id, changes, ip_address, user_agent, created_at +) +SELECT id, tenant_id, user_id, action, entity_type, + entity_id, changes, ip_address, user_agent, created_at +FROM audit_log; + +-- Step 3: Verify data integrity +SELECT COUNT(*) FROM audit_log_partitioned; +SELECT COUNT(*) FROM audit_log; + +-- Step 4: Rename tables +ALTER TABLE audit_log RENAME TO audit_log_old; +ALTER TABLE audit_log_partitioned RENAME TO audit_log; + +-- Step 5: Update sequences and indexes +-- (handled by the partitioned table definition) + +-- Step 6: Drop old table after verification +-- DROP TABLE audit_log_old; +``` + +### Automating Partition Creation + +Use the `setup_audit_partitioning.sql` script to automate partition management: + +```bash +# Run the setup script +psql -h localhost -U leocrm -d crm_db -f scripts/setup_audit_partitioning.sql +``` + +### Cron Job for Partition Maintenance + +Add to crontab to create partitions automatically: + +```bash +# Run on the 1st of each month at 2 AM +0 2 1 * * /usr/bin/psql -h localhost -U leocrm -d crm_db -c "SELECT create_monthly_audit_partition();" +``` + +### Querying Partitioned Data + +```sql +-- Query a specific month (fast, only scans one partition) +SELECT * FROM audit_log +WHERE created_at >= '2026-01-01' + AND created_at < '2026-02-01' + AND tenant_id = '...'; + +-- Query across months (scans multiple partitions) +SELECT * FROM audit_log +WHERE created_at >= '2026-01-01' + AND created_at < '2026-03-01' + AND action = 'permission_grant'; + +-- Check which partitions will be scanned +EXPLAIN SELECT * FROM audit_log +WHERE created_at >= '2026-01-01' + AND created_at < '2026-02-01'; +``` + +### Dropping Old Partitions + +```sql +-- Drop partitions older than retention period +DROP TABLE IF EXISTS audit_log_2025_01; +DROP TABLE IF EXISTS audit_log_2025_02; +-- ... + +-- Or use a function +SELECT drop_old_audit_partitions(12); -- Keep last 12 months +``` + +### Performance Considerations + +- **Index each partition** — Don't rely on parent table indexes +- **Use `created_at` in WHERE** — Always filter by date for partition pruning +- **Monitor partition count** — Too many partitions can slow planning +- **Archive old partitions** — Consider moving to cheaper storage +- **Vacuum partitions** — Each partition needs independent vacuum + +### Monitoring Partition Health + +```sql +-- Check partition sizes +SELECT + relname AS partition_name, + pg_size_pretty(pg_total_relation_size(relid)) AS total_size +FROM pg_catalog.pg_statio_user_tables +WHERE relname LIKE 'audit_log_%' +ORDER BY relname; + +-- Check row counts per partition +SELECT + relname AS partition_name, + n_live_tup AS row_count +FROM pg_catalog.pg_stat_user_tables +WHERE relname LIKE 'audit_log_%' +ORDER BY relname; + +-- List all partitions +SELECT + inhrelid::regclass AS partition_name +FROM pg_catalog.pg_inherits +WHERE inhparent = 'audit_log'::regclass +ORDER BY partition_name; +``` + +--- + +## 3. Backup and Recovery + +### Database Backup + +```bash +# Full backup +pg_dump -h localhost -U leocrm -d crm_db -F c -f /backups/crm_db_$(date +%Y%m%d).dump + +# Backup with compression +pg_dump -h localhost -U leocrm -d crm_db -F c -Z 9 -f /backups/crm_db_$(date +%Y%m%d).dump.gz + +# Backup specific schema only +pg_dump -h localhost -U leocrm -d crm_db -n public -F c -f /backups/crm_db_schema_$(date +%Y%m%d).dump +``` + +### Database Restore + +```bash +# Restore full backup +pg_restore -h localhost -U leocrm -d crm_db -c /backups/crm_db_20260701.dump + +# Restore with parallel workers (faster) +pg_restore -h localhost -U leocrm -d crm_db -j 4 -c /backups/crm_db_20260701.dump +``` + +### Automated Backup Script + +See `scripts/backup.py` for the automated backup solution. + +--- + +## 4. Monitoring and Alerts + +### Key Metrics + +| Metric | Target | Alert Threshold | +|--------|--------|----------------| +| Database connections | < 50 | > 80% of max | +| Query response time | < 100ms | > 500ms | +| Cache hit ratio | > 95% | < 90% | +| Partition size | < 10GB | > 50GB | +| PgBouncer pool usage | < 80% | > 90% | + +### Health Checks + +```bash +# Check PgBouncer status +echo "SHOW STATS;" | psql -h localhost -p 6432 -U leocrm -d pgbouncer | grep -E "total_|avg_" + +# Check partition health +psql -h localhost -U leocrm -d crm_db -c "SELECT count(*) FROM audit_log WHERE created_at < NOW() - INTERVAL '3 months';" + +# Check database size +psql -h localhost -U leocrm -d crm_db -c "SELECT pg_size_pretty(pg_database_size('crm_db'));" +``` diff --git a/docs/permissions.md b/docs/permissions.md new file mode 100644 index 0000000..a112487 --- /dev/null +++ b/docs/permissions.md @@ -0,0 +1,331 @@ +# LeoCRM Permission System + +> **Version:** 1.0 +> **Date:** 2026-07-29 +> **Applies to:** All developers and system administrators + +--- + +## 1. Architecture Overview + +The LeoCRM permission system is a multi-layered access control framework that combines: + +- **Row-Level Ownership** — Each entity can have an `owner_id` (user who owns it) +- **Entity Permissions (ACL)** — Explicit permission entries for users, groups, or roles +- **ABAC Policies** — Attribute-based policies for fine-grained access control +- **Role-Based Access Control (RBAC)** — Module-level permissions via user roles +- **System Admin Override** — System administrators see everything + +### Permission Resolution Order + +When checking access to an entity, the system resolves in this order (highest wins): + +1. **System Admin** → `delete` (full access to everything) +2. **Owner** → `owner` (from `owner_id` on the entity) +3. **Direct User Permission** → explicit ACL entry for the user +4. **Group Permission** → ACL entry for a group the user belongs to +5. **Role Permission** → ACL entry for the user's role +6. **Tenant-Owned** → `read` (if `owner_id IS NULL`, visible to all with module permission) +7. **No Access** → `none` + +### Permission Levels + +| Level | Value | Description | +|-------|-------|-------------| +| `none` | 0 | Explicit deny (overrides allow) | +| `read` | 1 | View the entity | +| `write` | 2 | Read + edit entity fields | +| `admin` | 3 | Write + delete + manage permissions | +| `delete` | 4 | Admin + transfer ownership | +| `owner` | 5 | Full control (automatic for owner) | + +--- + +## 2. Data Model + +### EntityPermission + +Stored in the `entity_permissions` table: + +| Column | Type | Description | +|--------|------|-------------| +| `id` | UUID | Primary key | +| `tenant_id` | UUID | Tenant scope | +| `entity_type` | String(50) | Entity type (e.g., 'contact', 'dms_file') | +| `entity_id` | UUID | The specific entity | +| `principal_type` | String(10) | 'user', 'group', 'role', 'guest' | +| `principal_id` | UUID | The user/group/role ID | +| `permission_level` | String(20) | 'none', 'read', 'write', 'admin', 'delete' | +| `expires_at` | DateTime | Optional expiration | +| `created_by` | UUID | Who created this permission | +| `created_at` | DateTime | Creation timestamp | +| `updated_at` | DateTime | Last update timestamp | + +### EntityPolicy (ABAC) + +Stored in the `entity_policies` table: + +| Column | Type | Description | +|--------|------|-------------| +| `id` | UUID | Primary key | +| `tenant_id` | UUID | Tenant scope | +| `name` | String(200) | Policy name | +| `entity_type` | String(50) | Target entity type | +| `principal_type` | String(10) | 'user', 'group', 'role' | +| `principal_id` | UUID | Target principal | +| `effect` | String(10) | 'allow' or 'deny' | +| `conditions` | JSONB | Attribute-based conditions | +| `priority` | Integer | Evaluation priority (higher = first) | +| `enabled` | Boolean | Whether the policy is active | +| `created_at` | DateTime | Creation timestamp | +| `updated_at` | DateTime | Last update timestamp | + +### OwnedMixin + +Adds `owner_id` to any model: + +- `NULL` → Tenant-owned (visible to all with module permission) +- `UUID` → Owned by that user +- Set automatically by service layer on creation +- Transfer requires owner, admin, or system_admin role + +--- + +## 3. Entity Permissions API + +### List Permissions + +``` +GET /api/v1/{entity_type}/{entity_id}/permissions +``` + +Returns all permission entries for an entity. + +### Create/Update Permission + +``` +POST /api/v1/{entity_type}/{entity_id}/permissions +``` + +```json +{ + "user_id": "uuid", + "access_level": "read", + "expires_at": "2026-12-31T23:59:59Z" +} +``` + +### Revoke Permission + +``` +DELETE /api/v1/{entity_type}/{entity_id}/permissions/{user_id} +``` + +### Check Access + +``` +GET /api/v1/{entity_type}/{entity_id}/access?user_id={uuid} +``` + +Returns the effective access level for a user. + +--- + +## 4. ABAC Policies + +### Policy Conditions Format + +```json +{ + "operator": "AND", + "rules": [ + {"field": "status", "op": "eq", "value": "active"}, + {"field": "amount", "op": "gte", "value": 1000}, + {"field": "tags", "op": "contains", "value": "vip"} + ] +} +``` + +### Supported Operators + +| Operator | Description | Example | +|----------|-------------|--------| +| `eq` | Equals | `{"field": "type", "op": "eq", "value": "company"}` | +| `neq` | Not equals | `{"field": "status", "op": "neq", "value": "archived"}` | +| `gt` | Greater than | `{"field": "amount", "op": "gt", "value": 100}` | +| `gte` | Greater or equal | `{"field": "amount", "op": "gte", "value": 50}` | +| `lt` | Less than | `{"field": "amount", "op": "lt", "value": 10000}` | +| `lte` | Less or equal | `{"field": "amount", "op": "lte", "value": 500}` | +| `in` | In list | `{"field": "status", "op": "in", "value": ["active", "pending"]}` | +| `not_in` | Not in list | `{"field": "status", "op": "not_in", "value": ["deleted"]}` | +| `contains` | String contains | `{"field": "name", "op": "contains", "value": "VIP"}` | +| `starts_with` | String starts with | `{"field": "name", "op": "starts_with", "value": "Confidential"}` | +| `is_null` | Is NULL | `{"field": "email", "op": "is_null"}` | +| `is_not_null` | Is not NULL | `{"field": "email", "op": "is_not_null"}` | + +### Policy Evaluation + +1. **Allow policies**: OR-joined (at least one must match for access) +2. **Deny policies**: NOT (none may match — deny takes precedence) +3. **Priority**: Higher priority policies evaluated first +4. **Enabled flag**: Disabled policies are skipped + +--- + +## 5. Service Layer + +### Entity Permission Service (`app/services/entity_permission_service.py`) + +| Function | Description | +|----------|-------------| +| `get_effective_access()` | Get access level for a user on a specific entity | +| `get_visible_ids()` | Get all visible entity IDs for a user | +| `batch_get_effective_access()` | Batch resolve access for multiple entities | +| `check_entity_access()` | Check if user has at least required level | +| `get_cached_visible_ids()` | Get visible IDs with Redis caching | +| `create_permission()` | Create or update a permission entry | +| `update_permission()` | Update an existing permission | +| `delete_permission()` | Delete a permission entry | +| `list_permissions()` | List all permissions for an entity | +| `list_all_permissions()` | List all permissions for a tenant | +| `cleanup_expired_permissions()` | Remove expired permission entries | +| `invalidate_all_user_entity_cache()` | Clear all cached permissions for a user | +| `get_permission_analytics()` | Get permission statistics | + +### Policy Service (`app/services/policy_service.py`) + +| Function | Description | +|----------|-------------| +| `create_policy()` | Create a new ABAC policy | +| `update_policy()` | Update an existing policy | +| `delete_policy()` | Delete a policy | +| `list_policies()` | List policies for a tenant | +| `build_sql_condition()` | Translate JSONB conditions to SQLAlchemy filters | +| `apply_policy_filter()` | Apply ABAC policies to a query | + +--- + +## 6. Caching + +The permission system uses Redis for caching visibility results: + +- **Cache key**: `ep_vis:{user_id}:{tenant_id}:{entity_type}` +- **Cache value**: JSON with `visible_ids` and `access_map` +- **TTL**: 5 minutes (300 seconds) +- **Invalidation**: Automatic on permission create/update/delete + +### Cache Flow + +1. Check Redis cache for user + entity type +2. Cache hit → return cached visible IDs +3. Cache miss → resolve from database, store in cache +4. Permission changes → invalidate affected user caches + +--- + +## 7. Performance Considerations + +- **Batch resolution** (`batch_get_effective_access`) is preferred over individual `get_effective_access` calls +- **Redis caching** reduces database load for repeated visibility checks +- **Bitmap optimization** for large entity sets (planned) +- **Indexes** on `entity_type + entity_id`, `principal_type + principal_id`, `tenant_id`, `expires_at` +- **Partitioning** recommended for `entity_permissions` table at scale + +--- + +## 8. Security Considerations + +- **Deny takes precedence** over allow in both ACL and ABAC +- **Expired permissions** are automatically excluded from resolution +- **System admin** bypasses all permission checks +- **Audit logging** for all permission changes +- **Notifications** sent to users when permissions are granted/revoked +- **Tenant isolation** enforced via `tenant_id` on all permission entries + +--- + +## 9. Examples + +### Grant Read Access to a User + +```python +from app.services import entity_permission_service as eps + +await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(user_id), + permission_level="read", + created_by=current_user.id, +) +``` + +### Check Access + +```python +access = await eps.get_effective_access( + db_session, tenant_id, user_id, "contact", contact_id +) +if access in ("read", "write", "admin", "delete", "owner"): + # User has access + pass +``` + +### Create ABAC Policy + +```python +from app.services import policy_service as ps + +policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="VIP Only", + entity_type="contact", + principal_type="user", + principal_id=str(user_id), + effect="allow", + conditions={ + "operator": "AND", + "rules": [ + {"field": "type", "op": "eq", "value": "company"}, + {"field": "name", "op": "contains", "value": "VIP"}, + ] + }, +) +``` + +### Batch Resolve Access + +```python +result = await eps.batch_get_effective_access( + db_session, tenant_id, user_id, "contact", entity_ids +) +for entity_id, level in result.items(): + print(f"Entity {entity_id}: {level}") +``` + +--- + +## 10. Troubleshooting + +### Common Issues + +| Issue | Cause | Solution | +|-------|-------|----------| +| User sees nothing | No permissions, not owner, not system admin | Grant explicit permission or set owner_id | +| Expired permission still works | Cache not invalidated | Wait for TTL or invalidate cache manually | +| ABAC policy not applied | Policy disabled or no conditions match | Check `enabled` flag and conditions | +| System admin can't see entity | Entity deleted or wrong tenant | Check `deleted_at` and `tenant_id` | +| Permission creation fails | Duplicate unique constraint | Use upsert (create_permission handles this) | + +### Debugging + +Enable debug logging: +```python +import logging +logging.getLogger("app.services.entity_permission_service").setLevel(logging.DEBUG) +logging.getLogger("app.services.policy_service").setLevel(logging.DEBUG) +``` diff --git a/docs/permissions_plugin_dev.md b/docs/permissions_plugin_dev.md new file mode 100644 index 0000000..9d7b510 --- /dev/null +++ b/docs/permissions_plugin_dev.md @@ -0,0 +1,509 @@ +# Permission System Plugin Development Guide + +> **Version:** 1.0 +> **Date:** 2026-07-29 +> **Applies to:** Plugin developers integrating with the LeoCRM permission system + +--- + +## 1. Overview + +Plugins can leverage the LeoCRM permission system to add row-level access control to their entities. This guide covers: + +- Adding `OwnedMixin` to plugin models +- Using `apply_visibility_filter()` for list queries +- Registering entity types for permission management +- Defining field-level permissions +- Integrating with ABAC policies + +--- + +## 2. Adding OwnedMixin to Plugin Models + +To enable ownership tracking for your plugin's entities, add `OwnedMixin` to your SQLAlchemy model: + +```python +from app.models.owned_mixin import OwnedMixin +from app.core.db import Base, TenantMixin + +class MyEntity(Base, TenantMixin, OwnedMixin): + __tablename__ = "my_entities" + + id: Mapped[uuid.UUID] = mapped_column( + PGUUID(as_uuid=True), primary_key=True, default=uuid.uuid4 + ) + name: Mapped[str] = mapped_column(String(200), nullable=False) + # ... other fields +``` + +### Owner ID Semantics + +- `NULL` → Tenant-owned (visible to all with module permission) +- `UUID` → Owned by that user (visible to owner + shared via entity_permissions) +- Set automatically by service layer on creation +- Transfer requires owner, admin, or system_admin role + +### Setting Owner on Creation + +```python +from app.models.owned_mixin import OwnedMixin + +async def create_entity(db: AsyncSession, data: dict, current_user: User): + entity = MyEntity( + tenant_id=current_user.tenant_id, + owner_id=current_user.id, # Set owner automatically + **data, + ) + db.add(entity) + await db.commit() + await db.refresh(entity) + return entity +``` + +--- + +## 3. Using apply_visibility_filter() + +The `apply_visibility_filter()` function filters a query to only return entities the user can see. This is the recommended way to implement list endpoints. + +### Basic Usage + +```python +from app.services import entity_permission_service as eps + +async def list_entities( + db: AsyncSession, + tenant_id: uuid.UUID, + user_id: uuid.UUID, +) -> list[MyEntity]: + # Get visible entity IDs + visible_ids, access_map = await eps.get_visible_ids( + db, tenant_id, user_id, "my_entity" + ) + + if not visible_ids: + return [] + + # Query only visible entities + result = await db.execute( + select(MyEntity) + .where(MyEntity.id.in_(visible_ids)) + .where(MyEntity.tenant_id == tenant_id) + ) + return result.scalars().all() +``` + +### With Caching + +```python +async def list_entities_cached( + db: AsyncSession, + redis: aioredis.Redis, + tenant_id: uuid.UUID, + user_id: uuid.UUID, +) -> list[MyEntity]: + # Use cached version for better performance + visible_ids, access_map = await eps.get_cached_visible_ids( + db, redis, tenant_id, user_id, "my_entity" + ) + + if not visible_ids: + return [] + + result = await db.execute( + select(MyEntity) + .where(MyEntity.id.in_(visible_ids)) + .where(MyEntity.tenant_id == tenant_id) + ) + return result.scalars().all() +``` + +### With Access Level in Response + +```python +async def list_entities_with_access( + db: AsyncSession, + tenant_id: uuid.UUID, + user_id: uuid.UUID, +) -> list[dict]: + visible_ids, access_map = await eps.get_visible_ids( + db, tenant_id, user_id, "my_entity" + ) + + if not visible_ids: + return [] + + result = await db.execute( + select(MyEntity) + .where(MyEntity.id.in_(visible_ids)) + .where(MyEntity.tenant_id == tenant_id) + ) + entities = result.scalars().all() + + return [ + { + **entity.to_dict(), + "access_level": access_map.get(entity.id, "none"), + } + for entity in entities + ] +``` + +--- + +## 4. Entity Registration + +Register your entity type so it appears in the permission management UI and API. + +### In Your Plugin Manifest + +```python +from app.plugins.manifest import PluginManifest + +class MyPluginManifest(PluginManifest): + # ... other fields + entity_types: list[str] = ["my_entity"] +``` + +### In Your Plugin Class + +```python +from app.plugins.base import BasePlugin + +class MyPlugin(BasePlugin): + @property + def entity_types(self) -> list[str]: + return ["my_entity"] +``` + +### Registering Routes for Permission Management + +```python +from fastapi import APIRouter, Depends + +router = APIRouter(prefix="/api/v1/my-entities") + +@router.get("/{entity_id}/permissions") +async def list_my_entity_permissions( + entity_id: uuid.UUID, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_current_user), +): + """List all permissions for a specific entity.""" + return await eps.list_permissions( + db, current_user.tenant_id, "my_entity", str(entity_id) + ) + +@router.post("/{entity_id}/permissions") +async def create_my_entity_permission( + entity_id: uuid.UUID, + body: PermissionCreate, + db: AsyncSession = Depends(get_db), + current_user: User = Depends(get_current_user), +): + """Grant permission to a user for this entity.""" + return await eps.create_permission( + db, + tenant_id=current_user.tenant_id, + entity_type="my_entity", + entity_id=str(entity_id), + principal_type="user", + principal_id=str(body.user_id), + permission_level=body.access_level, + expires_at=body.expires_at, + created_by=current_user.id, + ) +``` + +--- + +## 5. Field Definitions + +Define which fields of your entity are visible at each permission level. + +### Field Permission Schema + +```python +from pydantic import BaseModel +from typing import Any + +class EntityFieldDefinition(BaseModel): + """Define field visibility per permission level.""" + name: str + type: str # "string", "number", "boolean", "date", "reference" + required: bool = False + readable: dict[str, bool] = { + "read": True, + "write": True, + "admin": True, + "delete": True, + "owner": True, + } + writable: dict[str, bool] = { + "read": False, + "write": True, + "admin": True, + "delete": True, + "owner": True, + } +``` + +### Registering Field Definitions + +```python +from app.core.permission_registry import register_entity_fields + +FIELD_DEFINITIONS = [ + EntityFieldDefinition( + name="name", + type="string", + required=True, + readable={"read": True, "write": True, "admin": True, "delete": True, "owner": True}, + writable={"read": False, "write": True, "admin": True, "delete": True, "owner": True}, + ), + EntityFieldDefinition( + name="sensitive_data", + type="string", + required=False, + readable={"read": False, "write": False, "admin": True, "delete": True, "owner": True}, + writable={"read": False, "write": False, "admin": True, "delete": True, "owner": True}, + ), +] + +# Register during plugin activation +register_entity_fields("my_entity", FIELD_DEFINITIONS) +``` + +--- + +## 6. Integrating with ABAC Policies + +Your plugin can create and manage ABAC policies for its entities. + +### Creating Policies + +```python +from app.services import policy_service as ps + +# Create an allow policy +policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="Allow VIP my_entities", + entity_type="my_entity", + principal_type="user", + principal_id=str(user_id), + effect="allow", + conditions={ + "operator": "AND", + "rules": [ + {"field": "priority", "op": "gte", "value": 10}, + ] + }, +) +``` + +### Applying Policies to Queries + +```python +query = select(MyEntity).where(MyEntity.tenant_id == tenant_id) +query = await ps.apply_policy_filter( + db_session, query, "my_entity", user_id, tenant_id, MyEntity +) +result = await db_session.execute(query) +entities = result.scalars().all() +``` + +--- + +## 7. Best Practices + +### Do's + +- ✅ Always set `owner_id` on entity creation +- ✅ Use `get_visible_ids()` for list endpoints +- ✅ Use `batch_get_effective_access()` for bulk operations +- ✅ Register entity types for permission management +- ✅ Define field-level permissions for sensitive data +- ✅ Use Redis caching for frequently accessed permissions +- ✅ Handle permission expiration gracefully + +### Don'ts + +- ❌ Don't bypass permission checks for list endpoints +- ❌ Don't expose `owner_id` changes without authorization +- ❌ Don't create permissions without audit logging +- ❌ Don't forget to invalidate cache after permission changes +- ❌ Don't use `get_effective_access()` in loops — use batch instead + +--- + +## 8. Testing + +### Test Fixtures + +```python +import pytest +from app.services import entity_permission_service as eps +from tests.conftest import seed_tenant_and_users + +@pytest.mark.asyncio +async def test_my_entity_permissions(db_session): + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create entity with owner + entity = MyEntity( + tenant_id=tenant_id, + owner_id=user_id, + name="Test Entity", + ) + db_session.add(entity) + await db_session.commit() + + # Check access + access = await eps.get_effective_access( + db_session, tenant_id, user_id, "my_entity", entity.id + ) + assert access == "owner" +``` + +### Mocking Permissions + +```python +from unittest.mock import AsyncMock, patch + +async def test_list_with_mock_permissions(): + with patch( + "app.services.entity_permission_service.get_visible_ids", + new=AsyncMock(return_value=({uuid.UUID(int=1)}, {uuid.UUID(int=1): "read"})), + ): + # Your test code + pass +``` + +--- + +## 9. API Reference + +### Entity Permission Service + +| Function | Import | Description | +|----------|--------|-------------| +| `get_effective_access()` | `from app.services import entity_permission_service as eps` | Get access level for a user on a specific entity | +| `get_visible_ids()` | Same | Get all visible entity IDs for a user | +| `batch_get_effective_access()` | Same | Batch resolve access for multiple entities | +| `check_entity_access()` | Same | Check if user has at least required level | +| `get_cached_visible_ids()` | Same | Get visible IDs with Redis caching | +| `create_permission()` | Same | Create or update a permission entry | +| `update_permission()` | Same | Update an existing permission | +| `delete_permission()` | Same | Delete a permission entry | +| `list_permissions()` | Same | List all permissions for an entity | + +### Policy Service + +| Function | Import | Description | +|----------|--------|-------------| +| `create_policy()` | `from app.services import policy_service as ps` | Create a new ABAC policy | +| `update_policy()` | Same | Update an existing policy | +| `delete_policy()` | Same | Delete a policy | +| `list_policies()` | Same | List policies for a tenant | +| `build_sql_condition()` | Same | Translate JSONB conditions to SQLAlchemy filters | +| `apply_policy_filter()` | Same | Apply ABAC policies to a query | + +--- + +## 10. Example: Complete Plugin Integration + +```python +"""Example plugin with full permission system integration.""" + +from __future__ import annotations + +import uuid +from typing import Any + +from sqlalchemy import String +from sqlalchemy.dialects.postgresql import UUID as PGUUID +from sqlalchemy.ext.asyncio import AsyncSession +from sqlalchemy.orm import Mapped, mapped_column + +from app.core.db import Base, TenantMixin +from app.models.owned_mixin import OwnedMixin +from app.services import entity_permission_service as eps + + +class MyEntity(Base, TenantMixin, OwnedMixin): + """Example entity with permission support.""" + + __tablename__ = "my_entities" + + id: Mapped[uuid.UUID] = mapped_column( + PGUUID(as_uuid=True), primary_key=True, default=uuid.uuid4 + ) + name: Mapped[str] = mapped_column(String(200), nullable=False) + priority: Mapped[int] = mapped_column(nullable=False, default=0) + + +async def create_my_entity( + db: AsyncSession, + tenant_id: uuid.UUID, + user_id: uuid.UUID, + name: str, + priority: int = 0, +) -> MyEntity: + """Create a new entity with owner set.""" + entity = MyEntity( + tenant_id=tenant_id, + owner_id=user_id, + name=name, + priority=priority, + ) + db.add(entity) + await db.commit() + await db.refresh(entity) + return entity + + +async def list_visible_entities( + db: AsyncSession, + tenant_id: uuid.UUID, + user_id: uuid.UUID, +) -> list[MyEntity]: + """List entities visible to the user.""" + visible_ids, _ = await eps.get_visible_ids( + db, tenant_id, user_id, "my_entity" + ) + if not visible_ids: + return [] + + from sqlalchemy import select + result = await db.execute( + select(MyEntity) + .where(MyEntity.id.in_(visible_ids)) + .where(MyEntity.tenant_id == tenant_id) + ) + return result.scalars().all() + + +async def share_entity( + db: AsyncSession, + tenant_id: uuid.UUID, + entity_id: uuid.UUID, + target_user_id: uuid.UUID, + level: str, + created_by: uuid.UUID, +) -> dict: + """Share an entity with another user.""" + return await eps.create_permission( + db, + tenant_id=tenant_id, + entity_type="my_entity", + entity_id=str(entity_id), + principal_type="user", + principal_id=str(target_user_id), + permission_level=level, + created_by=created_by, + ) +``` diff --git a/frontend/src/pages/GuestContacts.tsx b/frontend/src/pages/GuestContacts.tsx new file mode 100644 index 0000000..93d747e --- /dev/null +++ b/frontend/src/pages/GuestContacts.tsx @@ -0,0 +1,224 @@ +import React, { useState, useEffect } from 'react'; +import { useNavigate } from 'react-router-dom'; +import { Users, LogOut, Search, Eye, Mail, Phone, Building2, User, Calendar } from 'lucide-react'; + +interface Contact { + id: string; + first_name: string; + last_name: string; + email?: string; + phone?: string; + company?: string; + position?: string; + created_at: string; +} + +export function GuestContactsPage() { + const navigate = useNavigate(); + const [contacts, setContacts] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const [search, setSearch] = useState(''); + const [guestInfo, setGuestInfo] = useState<{ name: string; email: string } | null>(null); + + useEffect(() => { + // Fetch guest info + fetch('/api/v1/guest/me', { credentials: 'include' }) + .then((res) => { + if (!res.ok) throw new Error('Not authenticated'); + return res.json(); + }) + .then((data) => setGuestInfo({ name: data.name, email: data.email })) + .catch(() => { + navigate('/guest/login'); + }); + }, [navigate]); + + useEffect(() => { + fetchContacts(); + }, []); + + const fetchContacts = async (query?: string) => { + setLoading(true); + setError(null); + try { + const url = query + ? `/api/v1/contacts?search=${encodeURIComponent(query)}&limit=50` + : '/api/v1/contacts?limit=50'; + const response = await fetch(url, { + credentials: 'include', + headers: { + 'X-CSRF-Token': '', + }, + }); + if (!response.ok) { + throw new Error('Failed to load contacts'); + } + const data = await response.json(); + setContacts(data.items || data.data || data || []); + } catch (err: any) { + setError(err.message || 'Failed to load contacts'); + } finally { + setLoading(false); + } + }; + + const handleSearch = (e: React.FormEvent) => { + e.preventDefault(); + fetchContacts(search); + }; + + const handleLogout = async () => { + try { + await fetch('/api/v1/guest/logout', { + method: 'POST', + credentials: 'include', + }); + } catch { + // ignore + } + navigate('/guest/login'); + }; + + return ( +
+ {/* Header */} +
+
+
+
+ +

+ Shared Contacts +

+
+
+ {guestInfo && ( + + {guestInfo.name} ({guestInfo.email}) + + )} + +
+
+
+
+ + {/* Search */} +
+
+
+
+ +
+ setSearch(e.target.value)} + placeholder="Search shared contacts..." + className="block w-full pl-10 pr-3 py-2 border border-gray-300 dark:border-gray-600 rounded-md shadow-sm placeholder-gray-400 dark:placeholder-gray-500 text-gray-900 dark:text-white dark:bg-gray-800 focus:outline-none focus:ring-primary-500 focus:border-primary-500 sm:text-sm" + /> +
+
+ + {/* Error */} + {error && ( +
+

{error}

+
+ )} + + {/* Loading */} + {loading && ( +
+ + + + +
+ )} + + {/* Contacts Grid */} + {!loading && !error && ( + <> +

+ {contacts.length} contact{contacts.length !== 1 ? 's' : ''} shared with you +

+ {contacts.length === 0 ? ( +
+ +

No contacts

+

+ No contacts have been shared with you yet. +

+
+ ) : ( +
+ {contacts.map((contact) => ( +
+
+
+
+ +
+
+

+ {contact.first_name} {contact.last_name} +

+ {contact.company && ( +

+ + {contact.company} +

+ )} +
+
+ +
+ +
+ {contact.email && ( +

+ + {contact.email} +

+ )} + {contact.phone && ( +

+ + {contact.phone} +

+ )} + {contact.position && ( +

+ + {contact.position} +

+ )} +
+ +
+

+ + Created: {new Date(contact.created_at).toLocaleDateString()} +

+
+
+ ))} +
+ )} + + )} +
+
+ ); +} diff --git a/frontend/src/pages/GuestLogin.tsx b/frontend/src/pages/GuestLogin.tsx new file mode 100644 index 0000000..df15e79 --- /dev/null +++ b/frontend/src/pages/GuestLogin.tsx @@ -0,0 +1,173 @@ +import React, { useState } from 'react'; +import { useNavigate, useSearchParams } from 'react-router-dom'; +import { Eye, EyeOff, LogIn, Mail, Lock, Building2 } from 'lucide-react'; + +interface GuestLoginResponse { + guest_user_id: string; + email: string; + name: string; + tenant_id: string; + csrf_token: string; +} + +export function GuestLoginPage() { + const navigate = useNavigate(); + const [searchParams] = useSearchParams(); + const [email, setEmail] = useState(searchParams.get('email') || ''); + const [password, setPassword] = useState(''); + const [tenantSlug, setTenantSlug] = useState(''); + const [showPassword, setShowPassword] = useState(false); + const [error, setError] = useState(null); + const [loading, setLoading] = useState(false); + + const handleSubmit = async (e: React.FormEvent) => { + e.preventDefault(); + setError(null); + setLoading(true); + + try { + const response = await fetch('/api/v1/guest/login', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ email, password, tenant_slug: tenantSlug }), + credentials: 'include', + }); + + if (!response.ok) { + const data = await response.json(); + throw new Error(data.detail?.detail || 'Login failed'); + } + + const data: GuestLoginResponse = await response.json(); + // Redirect to guest contacts page + navigate('/guest/contacts'); + } catch (err: any) { + setError(err.message || 'An error occurred'); + } finally { + setLoading(false); + } + }; + + return ( +
+
+
+

+ Guest Login +

+

+ Sign in with your guest credentials to access shared contacts +

+
+ +
+ {error && ( +
+

{error}

+
+ )} + +
+
+ +
+
+ +
+ setEmail(e.target.value)} + className="appearance-none rounded-none relative block w-full px-3 py-2 pl-10 border border-gray-300 dark:border-gray-600 placeholder-gray-500 dark:placeholder-gray-400 text-gray-900 dark:text-white dark:bg-gray-800 rounded-t-md focus:outline-none focus:ring-primary-500 focus:border-primary-500 focus:z-10 sm:text-sm" + placeholder="Email address" + /> +
+
+ +
+ +
+
+ +
+ setPassword(e.target.value)} + className="appearance-none rounded-none relative block w-full px-3 py-2 pl-10 pr-10 border border-gray-300 dark:border-gray-600 placeholder-gray-500 dark:placeholder-gray-400 text-gray-900 dark:text-white dark:bg-gray-800 focus:outline-none focus:ring-primary-500 focus:border-primary-500 focus:z-10 sm:text-sm" + placeholder="Password" + /> + +
+
+ +
+ +
+
+ +
+ setTenantSlug(e.target.value)} + className="appearance-none rounded-none relative block w-full px-3 py-2 pl-10 border border-gray-300 dark:border-gray-600 placeholder-gray-500 dark:placeholder-gray-400 text-gray-900 dark:text-white dark:bg-gray-800 rounded-b-md focus:outline-none focus:ring-primary-500 focus:border-primary-500 focus:z-10 sm:text-sm" + placeholder="Tenant slug (optional)" + /> +
+
+
+ +
+ +
+
+
+
+ ); +} diff --git a/frontend/src/routes/index.tsx b/frontend/src/routes/index.tsx index 816aafd..267ab02 100644 --- a/frontend/src/routes/index.tsx +++ b/frontend/src/routes/index.tsx @@ -4,6 +4,8 @@ import { AppShell } from '@/components/layout/AppShell'; import { ProtectedRoute } from './ProtectedRoute'; import { ProtectedRoute as PermissionRoute } from '@/components/common/ProtectedRoute'; import { LoginPage } from '@/pages/Login'; +import { GuestLoginPage } from '@/pages/GuestLogin'; +import { GuestContactsPage } from '@/pages/GuestContacts'; import { PasswordResetRequestPage } from '@/pages/PasswordResetRequest'; import { PasswordResetConfirmPage } from '@/pages/PasswordResetConfirm'; import { Loader2 } from 'lucide-react'; @@ -110,6 +112,14 @@ const router = createBrowserRouter([ path: '/password-reset/confirm', element: , }, + { + path: '/guest/login', + element: , + }, + { + path: '/guest/contacts', + element: , + }, { element: ( diff --git a/scripts/setup_audit_partitioning.sql b/scripts/setup_audit_partitioning.sql new file mode 100644 index 0000000..e482ab0 --- /dev/null +++ b/scripts/setup_audit_partitioning.sql @@ -0,0 +1,298 @@ +-- +-- Audit Log Partitioning Setup for LeoCRM +-- ========================================= +-- +-- This script creates a partitioned audit_log table with monthly partitions. +-- It includes functions for automatic partition creation and maintenance. +-- +-- Usage: +-- psql -h localhost -U leocrm -d crm_db -f scripts/setup_audit_partitioning.sql +-- +-- The script is idempotent — safe to run multiple times. +-- + +BEGIN; + +-- ═══════════════════════════════════════════════════════════════ +-- 1. Create the partitioned audit_log table +-- ═══════════════════════════════════════════════════════════════ + +-- Check if the partitioned table already exists +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_class WHERE relname = 'audit_log' AND relkind = 'p' + ) THEN + -- Check if there's an existing non-partitioned table + IF EXISTS (SELECT 1 FROM pg_class WHERE relname = 'audit_log' AND relkind = 'r') THEN + -- Rename existing table + ALTER TABLE audit_log RENAME TO audit_log_old; + RAISE NOTICE 'Renamed existing audit_log to audit_log_old'; + END IF; + + -- Create the partitioned table + CREATE TABLE audit_log ( + id UUID NOT NULL DEFAULT gen_random_uuid(), + tenant_id UUID, + user_id UUID, + action VARCHAR(100) NOT NULL, + entity_type VARCHAR(50), + entity_id UUID, + changes JSONB, + ip_address VARCHAR(45), + user_agent TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (id, created_at) + ) PARTITION BY RANGE (created_at); + + RAISE NOTICE 'Created partitioned audit_log table'; + END IF; +END +$$; + +-- ═══════════════════════════════════════════════════════════════ +-- 2. Function: Create a single monthly partition +-- ═══════════════════════════════════════════════════════════════ + +CREATE OR REPLACE FUNCTION create_monthly_audit_partition( + partition_date DATE DEFAULT CURRENT_DATE +) +RETURNS void +LANGUAGE plpgsql +AS $$ +DECLARE + year_str TEXT; + month_str TEXT; + partition_name TEXT; + start_date TEXT; + end_date TEXT; +BEGIN + year_str := TO_CHAR(partition_date, 'YYYY'); + month_str := TO_CHAR(partition_date, 'MM'); + partition_name := 'audit_log_' || year_str || '_' || month_str; + start_date := year_str || '-' || month_str || '-01'; + end_date := TO_CHAR(partition_date + INTERVAL '1 month', 'YYYY-MM-DD'); + + -- Check if partition already exists + IF NOT EXISTS ( + SELECT 1 FROM pg_class WHERE relname = partition_name + ) THEN + EXECUTE format(' + CREATE TABLE %I PARTITION OF audit_log + FOR VALUES FROM (%L) TO (%L)', + partition_name, start_date, end_date + ); + + -- Create indexes on the new partition + EXECUTE format(' + CREATE INDEX %I ON %I (tenant_id)', + 'idx_' || partition_name || '_tenant', partition_name + ); + EXECUTE format(' + CREATE INDEX %I ON %I (action)', + 'idx_' || partition_name || '_action', partition_name + ); + EXECUTE format(' + CREATE INDEX %I ON %I (entity_type, entity_id)', + 'idx_' || partition_name || '_entity', partition_name + ); + EXECUTE format(' + CREATE INDEX %I ON %I (created_at DESC)', + 'idx_' || partition_name || '_created', partition_name + ); + + RAISE NOTICE 'Created partition: %', partition_name; + ELSE + RAISE NOTICE 'Partition already exists: %', partition_name; + END IF; +END; +$$; + +-- ═══════════════════════════════════════════════════════════════ +-- 3. Function: Create partitions for the next N months +-- ═══════════════════════════════════════════════════════════════ + +CREATE OR REPLACE FUNCTION create_future_audit_partitions( + months_ahead INT DEFAULT 3 +) +RETURNS void +LANGUAGE plpgsql +AS $$ +DECLARE + i INT; + partition_date DATE; +BEGIN + -- Create current month partition first + partition_date := DATE_TRUNC('month', CURRENT_DATE)::DATE; + PERFORM create_monthly_audit_partition(partition_date); + + -- Create future month partitions + FOR i IN 1..months_ahead LOOP + partition_date := (DATE_TRUNC('month', CURRENT_DATE) + (i || ' months')::INTERVAL)::DATE; + PERFORM create_monthly_audit_partition(partition_date); + END LOOP; + + RAISE NOTICE 'Created % future partitions (current + % months)', months_ahead + 1, months_ahead; +END; +$$; + +-- ═══════════════════════════════════════════════════════════════ +-- 4. Function: Drop old partitions beyond retention period +-- ═══════════════════════════════════════════════════════════════ + +CREATE OR REPLACE FUNCTION drop_old_audit_partitions( + retention_months INT DEFAULT 12 +) +RETURNS INT +LANGUAGE plpgsql +AS $$ +DECLARE + partition_record RECORD; + drop_count INT := 0; + cutoff_date DATE; +BEGIN + cutoff_date := DATE_TRUNC('month', CURRENT_DATE - (retention_months || ' months')::INTERVAL)::DATE; + + FOR partition_record IN + SELECT + inhrelid::regclass AS partition_name, + pg_get_expr(relpartbound, relid) AS partition_bound + FROM pg_catalog.pg_inherits + JOIN pg_class ON pg_class.oid = inhrelid + WHERE inhparent = 'audit_log'::regclass + LOOP + -- Extract the upper bound date from the partition bound expression + -- Format: FOR VALUES FROM ('2026-01-01') TO ('2026-02-01') + IF partition_record.partition_bound ~ 'TO \(''([0-9]{4}-[0-9]{2}-[0-9]{2})' THEN + DECLARE + upper_date DATE; + BEGIN + upper_date := SUBSTRING( + partition_record.partition_bound + FROM 'TO \(''([0-9]{4}-[0-9]{2}-[0-9]{2})' + )::DATE; + + IF upper_date <= cutoff_date THEN + EXECUTE format('DROP TABLE IF EXISTS %I', partition_record.partition_name); + drop_count := drop_count + 1; + RAISE NOTICE 'Dropped old partition: %', partition_record.partition_name; + END IF; + END; + END IF; + END LOOP; + + RETURN drop_count; +END; +$$; + +-- ═══════════════════════════════════════════════════════════════ +-- 5. Create initial partitions (current + next 3 months) +-- ═══════════════════════════════════════════════════════════════ + +SELECT create_future_audit_partitions(3); + +-- ═══════════════════════════════════════════════════════════════ +-- 6. Migrate existing data (if any) +-- ═══════════════════════════════════════════════════════════════ + +DO $$ +DECLARE + old_count BIGINT; + new_count BIGINT; +BEGIN + -- Check if old table exists and has data + IF EXISTS (SELECT 1 FROM pg_class WHERE relname = 'audit_log_old') THEN + EXECUTE 'SELECT COUNT(*) FROM audit_log_old' INTO old_count; + + IF old_count > 0 THEN + RAISE NOTICE 'Migrating % rows from audit_log_old...', old_count; + + -- Insert data into partitioned table + EXECUTE ' + INSERT INTO audit_log ( + id, tenant_id, user_id, action, entity_type, + entity_id, changes, ip_address, user_agent, created_at + ) + SELECT + id, tenant_id, user_id, action, entity_type, + entity_id, changes, ip_address, user_agent, created_at + FROM audit_log_old + '; + + EXECUTE 'SELECT COUNT(*) FROM audit_log' INTO new_count; + RAISE NOTICE 'Migration complete: % rows migrated', new_count; + + -- Verify data integrity + IF old_count = new_count THEN + RAISE NOTICE 'Data integrity verified: % rows match', old_count; + ELSE + RAISE WARNING 'Data mismatch: old=% rows, new=% rows', old_count, new_count; + END IF; + ELSE + RAISE NOTICE 'No data to migrate in audit_log_old'; + END IF; + END IF; +END; +$$; + +-- ═══════════════════════════════════════════════════════════════ +-- 7. Verify setup +-- ═══════════════════════════════════════════════════════════════ + +DO $$ +DECLARE + partition_count INT; + table_type TEXT; +BEGIN + -- Check table type + SELECT relkind INTO table_type FROM pg_class WHERE relname = 'audit_log'; + + IF table_type = 'p' THEN + RAISE NOTICE 'audit_log is correctly set up as a partitioned table'; + ELSE + RAISE WARNING 'audit_log is NOT a partitioned table (relkind=%)', table_type; + END IF; + + -- Count partitions + SELECT COUNT(*) INTO partition_count + FROM pg_catalog.pg_inherits + WHERE inhparent = 'audit_log'::regclass; + + RAISE NOTICE 'Number of partitions: %', partition_count; +END; +$$; + +COMMIT; + +-- ═══════════════════════════════════════════════════════════════ +-- Usage Examples +-- ═══════════════════════════════════════════════════════════════ +-- +-- Create partitions for the next 6 months: +-- SELECT create_future_audit_partitions(6); +-- +-- Drop partitions older than 12 months: +-- SELECT drop_old_audit_partitions(12); +-- +-- Create a specific month partition: +-- SELECT create_monthly_audit_partition('2026-07-01'::DATE); +-- +-- Check partition sizes: +-- SELECT +-- relname AS partition_name, +-- pg_size_pretty(pg_total_relation_size(relid)) AS total_size +-- FROM pg_catalog.pg_statio_user_tables +-- WHERE relname LIKE 'audit_log_%' +-- ORDER BY relname; +-- +-- Check row counts per partition: +-- SELECT +-- relname AS partition_name, +-- n_live_tup AS row_count +-- FROM pg_catalog.pg_stat_user_tables +-- WHERE relname LIKE 'audit_log_%' +-- ORDER BY relname; +-- +-- Cron job (run on 1st of each month at 2 AM): +-- 0 2 1 * * /usr/bin/psql -h localhost -U leocrm -d crm_db -c "SELECT create_future_audit_partitions(3);" +-- diff --git a/scripts/setup_pgbouncer.sh b/scripts/setup_pgbouncer.sh new file mode 100755 index 0000000..2e30710 --- /dev/null +++ b/scripts/setup_pgbouncer.sh @@ -0,0 +1,279 @@ +#!/bin/bash +# +# PgBouncer Setup Script for LeoCRM +# =================================== +# +# This script installs and configures PgBouncer for PostgreSQL connection pooling. +# It creates the configuration files and starts the service. +# +# Usage: +# ./scripts/setup_pgbouncer.sh [--docker] [--password PASSWORD] +# +# Options: +# --docker Configure for Docker Compose environment +# --password PASS Set PostgreSQL password (default: from .env or prompt) +# + +set -euo pipefail + +# ─── Color Output ─── +RED='\033[0;31m' +GREEN='\033[0;32m' +YELLOW='\033[1;33m' +NC='\033[0m' # No Color + +log_info() { echo -e "${GREEN}[INFO]${NC} $1"; } +log_warn() { echo -e "${YELLOW}[WARN]${NC} $1"; } +log_error() { echo -e "${RED}[ERROR]${NC} $1"; } + +# ─── Default Values ─── +PGBOUNCER_VERSION="1.23.1" +PGBOUNCER_PORT="6432" +POOL_MODE="transaction" +DEFAULT_POOL_SIZE="25" +MAX_CLIENT_CONN="200" + +# ─── Parse Arguments ─── +DOCKER_MODE=false +POSTGRES_PASSWORD="" + +while [[ $# -gt 0 ]]; do + case $1 in + --docker) + DOCKER_MODE=true + shift + ;; + --password) + POSTGRES_PASSWORD="$2" + shift 2 + ;; + --help) + echo "Usage: $0 [--docker] [--password PASSWORD]" + echo "" + echo "Options:" + echo " --docker Configure for Docker Compose environment" + echo " --password PASS Set PostgreSQL password (default: from .env or prompt)" + exit 0 + ;; + *) + log_error "Unknown option: $1" + exit 1 + ;; + esac +done + +# ─── Detect Environment ─── +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PROJECT_DIR="$(cd "$SCRIPT_DIR/.." && pwd)" + +if [ -f "$PROJECT_DIR/.env" ]; then + source "$PROJECT_DIR/.env" +fi + +if [ -z "$POSTGRES_PASSWORD" ]; then + POSTGRES_PASSWORD="${POSTGRES_PASSWORD:-}" +fi + +if [ -z "$POSTGRES_PASSWORD" ]; then + read -s -p "Enter PostgreSQL password: " POSTGRES_PASSWORD + echo "" +fi + +# ─── Docker Mode ─── +if [ "$DOCKER_MODE" = true ]; then + log_info "Configuring PgBouncer for Docker Compose..." + + # Check if docker-compose.yml exists + if [ ! -f "$PROJECT_DIR/docker-compose.yml" ]; then + log_error "docker-compose.yml not found in $PROJECT_DIR" + exit 1 + fi + + # Check if PgBouncer service already exists + if grep -q "pgbouncer" "$PROJECT_DIR/docker-compose.yml" 2>/dev/null; then + log_warn "PgBouncer service already exists in docker-compose.yml" + else + log_info "Adding PgBouncer service to docker-compose.yml..." + + # Add PgBouncer service before the last line of services + cat >> "$PROJECT_DIR/docker-compose.yml" << 'EOF' + + pgbouncer: + image: bitnami/pgbouncer:latest + container_name: leocrm-pgbouncer + ports: + - "6432:6432" + environment: + - POSTGRESQL_HOST=crm-postgres + - POSTGRESQL_PORT=5432 + - POSTGRESQL_USERNAME=leocrm + - POSTGRESQL_PASSWORD=${POSTGRES_PASSWORD} + - POSTGRESQL_DATABASE=crm_db + - PGBOUNCER_POOL_MODE=transaction + - PGBOUNCER_DEFAULT_POOL_SIZE=25 + - PGBOUNCER_MAX_CLIENT_CONN=200 + depends_on: + - crm-postgres + restart: unless-stopped +EOF + log_info "PgBouncer service added to docker-compose.yml" + fi + + log_info "Docker PgBouncer configuration complete!" + log_info "Run 'docker-compose up -d pgbouncer' to start PgBouncer" + exit 0 +fi + +# ─── Native Installation ─── +log_info "Installing PgBouncer v${PGBOUNCER_VERSION}..." + +# Check if PgBouncer is already installed +if command -v pgbouncer &>/dev/null; then + log_info "PgBouncer is already installed: $(pgbouncer --version)" +else + # Install PgBouncer + if command -v apt-get &>/dev/null; then + apt-get update + apt-get install -y pgbouncer + elif command -v yum &>/dev/null; then + yum install -y pgbouncer + else + log_error "Unsupported package manager. Install PgBouncer manually." + exit 1 + fi + log_info "PgBouncer installed successfully" +fi + +# ─── Create Configuration ─── +log_info "Creating PgBouncer configuration..." + +PGBOUNCER_CONF_DIR="/etc/pgbouncer" +mkdir -p "$PGBOUNCER_CONF_DIR" + +# Generate md5 password hash +PG_MD5_HASH=$(echo -n "md5$(echo -n "${POSTGRES_PASSWORD}leocrm" | md5sum | cut -d' ' -f1)") + +# Create pgbouncer.ini +cat > "${PGBOUNCER_CONF_DIR}/pgbouncer.ini" << INI +[databases] +leocrm = host=localhost port=5432 dbname=crm_db +leocrm_test = host=localhost port=5432 dbname=leocrm_test + +[pgbouncer] +listen_addr = 0.0.0.0 +listen_port = ${PGBOUNCER_PORT} +unix_socket_dir = /var/run/pgbouncer + +auth_type = md5 +auth_file = ${PGBOUNCER_CONF_DIR}/userlist.txt + +pool_mode = ${POOL_MODE} +default_pool_size = ${DEFAULT_POOL_SIZE} +max_client_conn = ${MAX_CLIENT_CONN} +max_db_connections = 50 + +server_idle_timeout = 600 +server_lifetime = 3600 +client_idle_timeout = 1800 +query_timeout = 30 + +log_connections = 1 +log_disconnections = 1 +log_pooler_errors = 1 +stats_period = 60 + +listen_backlog = 128 +INI + +log_info "Created ${PGBOUNCER_CONF_DIR}/pgbouncer.ini" + +# Create userlist.txt +cat > "${PGBOUNCER_CONF_DIR}/userlist.txt" << USERLIST +"leocrm" "${PG_MD5_HASH}" +"postgres" "${PG_MD5_HASH}" +USERLIST + +log_info "Created ${PGBOUNCER_CONF_DIR}/userlist.txt" + +# Set proper permissions +chmod 640 "${PGBOUNCER_CONF_DIR}/pgbouncer.ini" +chmod 640 "${PGBOUNCER_CONF_DIR}/userlist.txt" +chown -R pgbouncer:pgbouncer "$PGBOUNCER_CONF_DIR" 2>/dev/null || true + +# ─── Create Systemd Service ─── +log_info "Creating systemd service..." + +cat > /etc/systemd/system/pgbouncer.service << 'SYSTEMD' +[Unit] +Description=PgBouncer PostgreSQL Connection Pooler +After=network.target + +[Service] +Type=forking +User=pgbouncer +ExecStart=/usr/sbin/pgbouncer -d /etc/pgbouncer/pgbouncer.ini +ExecReload=/bin/kill -HUP $MAINPID +ExecStop=/bin/kill -INT $MAINPID +PIDFile=/var/run/pgbouncer/pgbouncer.pid +Restart=always +RestartSec=5 + +[Install] +WantedBy=multi-user.target +SYSTEMD + +log_info "Created systemd service" + +# ─── Start PgBouncer ─── +log_info "Starting PgBouncer..." + +systemctl daemon-reload +systemctl enable pgbouncer +systemctl start pgbouncer + +# Wait for PgBouncer to start +sleep 2 + +# Check status +if systemctl is-active --quiet pgbouncer; then + log_info "PgBouncer is running on port ${PGBOUNCER_PORT}" +else + log_error "PgBouncer failed to start. Check logs: journalctl -u pgbouncer" + exit 1 +fi + +# ─── Verify Connection ─── +log_info "Verifying PgBouncer connection..." + +if command -v psql &>/dev/null; then + if PGPASSWORD="$POSTGRES_PASSWORD" psql -h localhost -p "$PGBOUNCER_PORT" -U leocrm -d leocrm -c "SELECT 1 AS pgbouncer_test;" &>/dev/null; then + log_info "PgBouncer connection verified successfully!" + else + log_warn "Could not verify PgBouncer connection. Check PostgreSQL credentials." + fi +fi + +# ─── Summary ─── +echo "" +echo "═══════════════════════════════════════════════" +echo " PgBouncer Setup Complete" +echo "═══════════════════════════════════════════════" +echo "" +echo " Configuration:" +echo " Config file: ${PGBOUNCER_CONF_DIR}/pgbouncer.ini" +echo " User list: ${PGBOUNCER_CONF_DIR}/userlist.txt" +echo " Listen port: ${PGBOUNCER_PORT}" +echo " Pool mode: ${POOL_MODE}" +echo " Pool size: ${DEFAULT_POOL_SIZE}" +echo "" +echo " Commands:" +echo " Status: systemctl status pgbouncer" +echo " Restart: systemctl restart pgbouncer" +echo " Reload: systemctl reload pgbouncer" +echo " Stop: systemctl stop pgbouncer" +echo "" +echo " Monitoring:" +echo " Pools: echo 'SHOW POOLS;' | psql -h localhost -p ${PGBOUNCER_PORT} -U leocrm -d pgbouncer" +echo " Stats: echo 'SHOW STATS;' | psql -h localhost -p ${PGBOUNCER_PORT} -U leocrm -d pgbouncer" +echo "" +echo "═══════════════════════════════════════════════" diff --git a/tests/test_abac.py b/tests/test_abac.py new file mode 100644 index 0000000..45e5060 --- /dev/null +++ b/tests/test_abac.py @@ -0,0 +1,424 @@ +"""Tests for ABAC policy service — attribute-based access control policies. + +Covers: +- Policy allow (conditions match → access granted) +- Policy deny (conditions match → access blocked) +- Multi-condition AND (all conditions must match) +- Multi-condition OR (any condition must match) +- Policy CRUD operations +- build_sql_condition translation +- apply_policy_filter integration +""" + +from __future__ import annotations + +import uuid + +import pytest +import pytest_asyncio +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.contact import Contact +from app.models.entity_policy import EntityPolicy +from app.services import policy_service as ps +from tests.conftest import seed_tenant_and_users + + +@pytest.mark.asyncio +class TestABACPolicyService: + """Tests for ABAC policy service — attribute-based access control.""" + + async def test_policy_allow(self, db_session: AsyncSession): + """Allow policy with matching conditions grants access.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create an allow policy for contacts where type == 'company' + policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="Allow company contacts", + entity_type="contact", + principal_type="user", + principal_id=str(user_id), + effect="allow", + conditions={ + "operator": "AND", + "rules": [ + {"field": "type", "op": "eq", "value": "company"} + ] + }, + ) + + assert policy["effect"] == "allow" + assert policy["name"] == "Allow company contacts" + assert policy["entity_type"] == "contact" + assert policy["conditions"]["rules"][0]["field"] == "type" + + # Verify the policy was persisted + db_policy = await db_session.get(EntityPolicy, uuid.UUID(policy["id"])) + assert db_policy is not None + assert db_policy.effect == "allow" + + async def test_policy_deny(self, db_session: AsyncSession): + """Deny policy blocks access when conditions match.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create a deny policy for contacts where name contains 'Confidential' + policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="Deny confidential contacts", + entity_type="contact", + principal_type="user", + principal_id=str(user_id), + effect="deny", + conditions={ + "operator": "AND", + "rules": [ + {"field": "name", "op": "starts_with", "value": "Confidential"} + ] + }, + ) + + assert policy["effect"] == "deny" + assert policy["enabled"] is True + + async def test_multi_condition_and(self, db_session: AsyncSession): + """AND conditions: all rules must match for the policy to apply.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create policy with AND conditions + policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="VIP company contacts", + entity_type="contact", + principal_type="user", + principal_id=str(user_id), + effect="allow", + conditions={ + "operator": "AND", + "rules": [ + {"field": "type", "op": "eq", "value": "company"}, + {"field": "name", "op": "contains", "value": "VIP"}, + ] + }, + ) + + assert policy["conditions"]["operator"] == "AND" + assert len(policy["conditions"]["rules"]) == 2 + + # Test build_sql_condition translation + condition = ps.build_sql_condition(policy["conditions"], Contact) + assert condition is not None + + async def test_multi_condition_or(self, db_session: AsyncSession): + """OR conditions: any rule matching triggers the policy.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create policy with OR conditions + policy = await ps.create_policy( + db_session, + tenant_id=tenant_id, + name="High value or VIP contacts", + entity_type="contact", + principal_type="user", + principal_id=str(user_id), + effect="allow", + conditions={ + "operator": "OR", + "rules": [ + {"field": "name", "op": "contains", "value": "VIP"}, + {"field": "name", "op": "contains", "value": "Premium"}, + ] + }, + ) + + assert policy["conditions"]["operator"] == "OR" + assert len(policy["conditions"]["rules"]) == 2 + + # Test build_sql_condition translation + condition = ps.build_sql_condition(policy["conditions"], Contact) + assert condition is not None + + async def test_build_sql_condition_eq(self, db_session: AsyncSession): + """build_sql_condition translates 'eq' operator correctly.""" + conditions = { + "operator": "AND", + "rules": [ + {"field": "type", "op": "eq", "value": "person"} + ] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is not None + + async def test_build_sql_condition_gt(self, db_session: AsyncSession): + """build_sql_condition translates 'gt' operator correctly.""" + conditions = { + "operator": "AND", + "rules": [ + {"field": "id", "op": "gt", "value": 100} + ] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is not None + + async def test_build_sql_condition_contains(self, db_session: AsyncSession): + """build_sql_condition translates 'contains' operator correctly.""" + conditions = { + "operator": "AND", + "rules": [ + {"field": "name", "op": "contains", "value": "test"} + ] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is not None + + async def test_build_sql_condition_nested(self, db_session: AsyncSession): + """build_sql_condition handles nested conditions blocks.""" + conditions = { + "operator": "AND", + "rules": [ + { + "operator": "OR", + "rules": [ + {"field": "type", "op": "eq", "value": "company"}, + {"field": "type", "op": "eq", "value": "person"}, + ] + }, + {"field": "name", "op": "is_not_null", "value": None} + ] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is not None + + async def test_build_sql_condition_empty_rules(self, db_session: AsyncSession): + """build_sql_condition returns None for empty rules.""" + conditions = { + "operator": "AND", + "rules": [] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is None + + async def test_build_sql_condition_unsupported_op(self, db_session: AsyncSession): + """build_sql_condition skips unsupported operators gracefully.""" + conditions = { + "operator": "AND", + "rules": [ + {"field": "name", "op": "unsupported_op", "value": "test"} + ] + } + condition = ps.build_sql_condition(conditions, Contact) + assert condition is None + + async def test_list_policies(self, db_session: AsyncSession): + """list_policies returns all policies for a tenant.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create two policies + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Policy 1", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Policy 2", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="deny", + ) + + policies = await ps.list_policies(db_session, tenant_id) + assert len(policies) >= 2 + names = [p["name"] for p in policies] + assert "Policy 1" in names + assert "Policy 2" in names + + async def test_list_policies_filtered_by_entity_type(self, db_session: AsyncSession): + """list_policies filters by entity_type when provided.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Contact Policy", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + await ps.create_policy( + db_session, tenant_id=tenant_id, name="File Policy", + entity_type="dms_file", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + + contact_policies = await ps.list_policies(db_session, tenant_id, entity_type="contact") + assert len(contact_policies) >= 1 + assert all(p["entity_type"] == "contact" for p in contact_policies) + + async def test_update_policy(self, db_session: AsyncSession): + """update_policy modifies an existing policy.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + policy = await ps.create_policy( + db_session, tenant_id=tenant_id, name="Original", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + + updated = await ps.update_policy( + db_session, tenant_id, policy["id"], + name="Updated", effect="deny", + ) + + assert updated["name"] == "Updated" + assert updated["effect"] == "deny" + + async def test_delete_policy(self, db_session: AsyncSession): + """delete_policy removes a policy.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + policy = await ps.create_policy( + db_session, tenant_id=tenant_id, name="To Delete", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + + await ps.delete_policy(db_session, tenant_id, policy["id"]) + + # Verify it's gone + policies = await ps.list_policies(db_session, tenant_id) + ids = [p["id"] for p in policies] + assert policy["id"] not in ids + + async def test_apply_policy_filter_no_policies(self, db_session: AsyncSession): + """apply_policy_filter returns query unchanged when no policies match.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + query = select(Contact).where(Contact.tenant_id == tenant_id) + result = await ps.apply_policy_filter( + db_session, query, "contact", user_id, tenant_id, Contact + ) + + # Query should be unchanged (no ABAC restriction) + rows = await db_session.execute(result) + contacts = rows.scalars().all() + assert len(contacts) >= 1 + + async def test_apply_policy_filter_with_allow(self, db_session: AsyncSession): + """apply_policy_filter with allow policy filters correctly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create allow policy for company-type contacts + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Only companies", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + conditions={ + "operator": "AND", + "rules": [ + {"field": "type", "op": "eq", "value": "company"} + ] + }, + ) + + query = select(Contact).where(Contact.tenant_id == tenant_id) + result = await ps.apply_policy_filter( + db_session, query, "contact", user_id, tenant_id, Contact + ) + + rows = await db_session.execute(result) + contacts = rows.scalars().all() + assert len(contacts) >= 1 + for c in contacts: + assert c.type == "company", f"Expected 'company', got '{c.type}'" + + async def test_apply_policy_filter_with_deny(self, db_session: AsyncSession): + """apply_policy_filter with deny policy blocks matching entities.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create deny policy for contacts named 'Company Alpha' + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Block Company Alpha", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="deny", + conditions={ + "operator": "AND", + "rules": [ + {"field": "name", "op": "eq", "value": "Company Alpha"} + ] + }, + ) + + query = select(Contact).where(Contact.tenant_id == tenant_id) + result = await ps.apply_policy_filter( + db_session, query, "contact", user_id, tenant_id, Contact + ) + + rows = await db_session.execute(result) + contacts = rows.scalars().all() + # Company Alpha should be filtered out + names = [c.name for c in contacts] + assert "Company Alpha" not in names, "Deny policy should have blocked Company Alpha" + + async def test_policy_priority_ordering(self, db_session: AsyncSession): + """Policies are ordered by priority (higher first).""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + await ps.create_policy( + db_session, tenant_id=tenant_id, name="Low Priority", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", priority=1, + ) + await ps.create_policy( + db_session, tenant_id=tenant_id, name="High Priority", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", priority=100, + ) + + policies = await ps.list_policies(db_session, tenant_id, entity_type="contact") + # High priority should come first + high_idx = next(i for i, p in enumerate(policies) if p["name"] == "High Priority") + low_idx = next(i for i, p in enumerate(policies) if p["name"] == "Low Priority") + assert high_idx < low_idx, "High priority policy should come first" + + async def test_policy_enabled_flag(self, db_session: AsyncSession): + """Disabled policies are not applied.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + policy = await ps.create_policy( + db_session, tenant_id=tenant_id, name="Disabled Policy", + entity_type="contact", principal_type="user", + principal_id=str(user_id), effect="allow", + ) + + # Disable the policy + updated = await ps.update_policy( + db_session, tenant_id, policy["id"], enabled=False + ) + assert updated["enabled"] is False +} diff --git a/tests/test_entity_permissions.py b/tests/test_entity_permissions.py new file mode 100644 index 0000000..eacda88 --- /dev/null +++ b/tests/test_entity_permissions.py @@ -0,0 +1,406 @@ +"""Tests for entity permission service — ACL resolution, ownership, sharing, expiration. + +Covers: +- Owner sees own contacts +- Non-owner doesn't see others' contacts +- Shared user sees shared contacts +- Permission expiration +- System admin sees all +- Batch resolution +""" + +from __future__ import annotations + +import uuid +from datetime import UTC, datetime, timedelta + +import pytest +import pytest_asyncio +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.contact import Contact +from app.models.entity_permission import EntityPermission +from app.models.user import User +from app.services import entity_permission_service as eps +from tests.conftest import seed_tenant_and_users + + +@pytest.mark.asyncio +class TestEntityPermissions: + """Tests for entity permission service — ACL resolution, ownership, sharing.""" + + async def test_owner_sees_own_contacts(self, db_session: AsyncSession): + """Owner has 'owner' access level on their own contacts.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = user_id + await db_session.commit() + + access = await eps.get_effective_access( + db_session, tenant_id, user_id, "contact", contact_id + ) + assert access == "owner", f"Expected 'owner', got '{access}'" + + async def test_non_owner_doesnt_see_others_contacts(self, db_session: AsyncSession): + """Non-owner without explicit permission gets 'none' access.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + other_user_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + access = await eps.get_effective_access( + db_session, tenant_id, other_user_id, "contact", contact_id + ) + assert access == "none", f"Expected 'none', got '{access}'" + + async def test_shared_user_sees_shared_contacts(self, db_session: AsyncSession): + """User with explicit 'read' permission sees the contact.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + shared_user_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read permission to viewer + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(shared_user_id), + permission_level="read", + created_by=owner_id, + ) + + access = await eps.get_effective_access( + db_session, tenant_id, shared_user_id, "contact", contact_id + ) + assert access == "read", f"Expected 'read', got '{access}'" + + async def test_permission_expiration(self, db_session: AsyncSession): + """Expired permission returns 'none' access.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + shared_user_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read permission that expired 1 hour ago + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(shared_user_id), + permission_level="read", + expires_at=datetime.now(UTC) - timedelta(hours=1), + created_by=owner_id, + ) + + access = await eps.get_effective_access( + db_session, tenant_id, shared_user_id, "contact", contact_id + ) + assert access == "none", f"Expected 'none' for expired permission, got '{access}'" + + async def test_system_admin_sees_all(self, db_session: AsyncSession): + """System admin gets 'delete' access on any entity.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + contact_id = seed["company_a"].id + + # Create a system admin user + from app.core.auth import hash_password + sys_admin = User( + email="sysadmin@test.com", + name="System Admin", + password_hash=hash_password("TestPass123!"), + is_active=True, + is_system_admin=True, + preferences={}, + ) + db_session.add(sys_admin) + await db_session.commit() + + access = await eps.get_effective_access( + db_session, tenant_id, sys_admin.id, "contact", contact_id + ) + assert access == "delete", f"Expected 'delete' for system admin, got '{access}'" + + async def test_batch_resolution(self, db_session: AsyncSession): + """Batch resolution returns correct access levels for multiple entities.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Create a second contact owned by viewer + contact2 = Contact( + tenant_id=tenant_id, + type="company", + name="Company Viewer", + displayname="Company Viewer", + owner_id=viewer_id, + created_by=viewer_id, + updated_by=viewer_id, + ) + db_session.add(contact2) + await db_session.commit() + + # Set owner on first contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read permission to viewer on first contact + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="read", + created_by=owner_id, + ) + + # Batch resolve for viewer + result = await eps.batch_get_effective_access( + db_session, tenant_id, viewer_id, "contact", [contact_id, contact2.id] + ) + + assert result[contact_id] == "read", f"Expected 'read' for shared contact, got '{result[contact_id]}'" + assert result[contact2.id] == "owner", f"Expected 'owner' for own contact, got '{result[contact2.id]}'" + + async def test_tenant_owned_contact_visible_to_all(self, db_session: AsyncSession): + """Tenant-owned contact (owner_id IS NULL) is visible as 'read' to all tenant users.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Ensure owner_id is NULL + contact = await db_session.get(Contact, contact_id) + contact.owner_id = None + await db_session.commit() + + access = await eps.get_effective_access( + db_session, tenant_id, viewer_id, "contact", contact_id + ) + assert access == "read", f"Expected 'read' for tenant-owned contact, got '{access}'" + + async def test_group_permission_propagation(self, db_session: AsyncSession): + """User inherits permissions from group membership.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Create a group and add viewer to it + from app.models.group import Group, UserGroup + group = Group(tenant_id=tenant_id, name="Viewers") + db_session.add(group) + await db_session.flush() + + ug = UserGroup(tenant_id=tenant_id, user_id=viewer_id, group_id=group.id) + db_session.add(ug) + await db_session.commit() + + # Grant permission to the group + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="group", + principal_id=str(group.id), + permission_level="write", + created_by=owner_id, + ) + + access = await eps.get_effective_access( + db_session, tenant_id, viewer_id, "contact", contact_id + ) + assert access == "write", f"Expected 'write' via group, got '{access}'" + + async def test_highest_permission_wins(self, db_session: AsyncSession): + """When multiple permissions exist, the highest level wins.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read + write permissions (write should win) + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="read", + created_by=owner_id, + ) + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="write", + created_by=owner_id, + ) + + access = await eps.get_effective_access( + db_session, tenant_id, viewer_id, "contact", contact_id + ) + assert access == "write", f"Expected 'write' (highest wins), got '{access}'" + + async def test_visible_ids_returns_owned_and_shared(self, db_session: AsyncSession): + """get_visible_ids returns owned + shared + tenant-owned entities.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Create a second contact owned by viewer + contact2 = Contact( + tenant_id=tenant_id, + type="company", + name="Viewer Owned", + displayname="Viewer Owned", + owner_id=viewer_id, + created_by=viewer_id, + updated_by=viewer_id, + ) + db_session.add(contact2) + await db_session.commit() + + # Set owner on first contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read to viewer on first contact + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="read", + created_by=owner_id, + ) + + visible, access_map = await eps.get_visible_ids( + db_session, tenant_id, viewer_id, "contact" + ) + + assert contact_id in visible, "Shared contact should be visible" + assert contact2.id in visible, "Owned contact should be visible" + assert access_map[contact_id] == "read" + assert access_map[contact2.id] == "owner" + + async def test_check_entity_access_enforces_required_level(self, db_session: AsyncSession): + """check_entity_access returns True/False based on required level.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id on the contact + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Grant read to viewer + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="read", + created_by=owner_id, + ) + + # Should have read access + assert await eps.check_entity_access( + db_session, tenant_id, viewer_id, "contact", contact_id, "read" + ) is True + + # Should NOT have write access + assert await eps.check_entity_access( + db_session, tenant_id, viewer_id, "contact", contact_id, "write" + ) is False + + async def test_cleanup_expired_permissions(self, db_session: AsyncSession): + """cleanup_expired_permissions removes expired entries.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Create expired permission + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level="read", + expires_at=datetime.now(UTC) - timedelta(hours=1), + created_by=owner_id, + ) + + count = await eps.cleanup_expired_permissions(db_session) + assert count >= 1, "Expected at least 1 expired permission cleaned up" + + # Verify it's gone + access = await eps.get_effective_access( + db_session, tenant_id, viewer_id, "contact", contact_id + ) + assert access == "none", "Permission should be gone after cleanup" +} diff --git a/tests/test_permission_performance.py b/tests/test_permission_performance.py new file mode 100644 index 0000000..a7aa6f4 --- /dev/null +++ b/tests/test_permission_performance.py @@ -0,0 +1,289 @@ +"""Performance tests for permission system — visibility filter, caching. + +Covers: +- Visibility filter performance with 1000 mock contacts +- Cache hit vs miss performance comparison +""" + +from __future__ import annotations + +import time +import uuid +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest +import pytest_asyncio +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.contact import Contact +from app.services import entity_permission_service as eps +from tests.conftest import seed_tenant_and_users + + +@pytest.mark.asyncio +class TestPermissionPerformance: + """Performance tests for permission system.""" + + async def test_visibility_filter_performance(self, db_session: AsyncSession): + """get_visible_ids with 1000 contacts completes in reasonable time.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create 1000 contacts owned by the user + contacts = [] + for i in range(1000): + contacts.append(Contact( + tenant_id=tenant_id, + type="person", + first_name=f"Perf{i}", + last_name=f"Test{i}", + owner_id=user_id, + created_by=user_id, + updated_by=user_id, + )) + db_session.add_all(contacts) + await db_session.commit() + + # Measure time for get_visible_ids + start = time.perf_counter() + visible, access_map = await eps.get_visible_ids( + db_session, tenant_id, user_id, "contact" + ) + elapsed = time.perf_counter() - start + + # Should complete in under 2 seconds for 1000 contacts + assert elapsed < 2.0, f"get_visible_ids took {elapsed:.3f}s (expected <2.0s)" + assert len(visible) >= 1000, f"Expected >=1000 visible, got {len(visible)}" + + # All should be at 'owner' level + for eid, level in access_map.items(): + if eid in [c.id for c in contacts]: + assert level == "owner", f"Expected 'owner', got '{level}'" + + async def test_cache_hit_vs_miss(self, db_session: AsyncSession, redis_client): + """Cache hit is significantly faster than cache miss.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create 100 contacts owned by the user + contacts = [] + for i in range(100): + contacts.append(Contact( + tenant_id=tenant_id, + type="person", + first_name=f"Cache{i}", + last_name=f"Test{i}", + owner_id=user_id, + created_by=user_id, + updated_by=user_id, + )) + db_session.add_all(contacts) + await db_session.commit() + + # First call — cache miss (populates cache) + miss_start = time.perf_counter() + visible_miss, access_map_miss = await eps.get_cached_visible_ids( + db_session, redis_client, tenant_id, user_id, "contact" + ) + miss_elapsed = time.perf_counter() - miss_start + + assert len(visible_miss) >= 100, f"Expected >=100 visible, got {len(visible_miss)}" + + # Second call — cache hit + hit_start = time.perf_counter() + visible_hit, access_map_hit = await eps.get_cached_visible_ids( + db_session, redis_client, tenant_id, user_id, "contact" + ) + hit_elapsed = time.perf_counter() - hit_start + + # Cache hit should be faster (at least 2x speedup) + assert hit_elapsed < miss_elapsed, f"Cache hit {hit_elapsed:.4f}s should be faster than miss {miss_elapsed:.4f}s" + assert len(visible_hit) == len(visible_miss), "Cache hit should return same results" + + async def test_batch_resolution_performance(self, db_session: AsyncSession): + """batch_get_effective_access with 100 entities completes quickly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Create 100 contacts + contacts = [] + for i in range(100): + contacts.append(Contact( + tenant_id=tenant_id, + type="person", + first_name=f"Batch{i}", + last_name=f"Test{i}", + owner_id=user_id, + created_by=user_id, + updated_by=user_id, + )) + db_session.add_all(contacts) + await db_session.commit() + + entity_ids = [c.id for c in contacts] + + start = time.perf_counter() + result = await eps.batch_get_effective_access( + db_session, tenant_id, user_id, "contact", entity_ids + ) + elapsed = time.perf_counter() - start + + assert elapsed < 1.0, f"batch_get_effective_access took {elapsed:.3f}s (expected <1.0s)" + assert len(result) == 100, f"Expected 100 results, got {len(result)}" + + # All should be 'owner' + for eid, level in result.items(): + assert level == "owner", f"Expected 'owner', got '{level}'" + + async def test_cache_invalidation_performance(self, db_session: AsyncSession, redis_client): + """Cache invalidation for all entity types completes quickly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + + # Populate cache for multiple entity types + for entity_type in ["contact", "dms_file", "mailbox", "calendar_event", "task"]: + await eps.get_cached_visible_ids( + db_session, redis_client, tenant_id, user_id, entity_type + ) + + # Measure invalidation time + start = time.perf_counter() + await eps.invalidate_all_user_entity_cache(redis_client, tenant_id, user_id) + elapsed = time.perf_counter() - start + + assert elapsed < 1.0, f"Cache invalidation took {elapsed:.3f}s (expected <1.0s)" + + # Verify cache is cleared + for entity_type in ["contact", "dms_file", "mailbox", "calendar_event", "task"]: + cache_key = f"ep_vis:{user_id}:{tenant_id}:{entity_type}" + cached = await redis_client.get(cache_key) + assert cached is None, f"Cache for {entity_type} should be cleared" + + async def test_get_effective_access_performance(self, db_session: AsyncSession): + """get_effective_access for a single entity completes quickly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + contact_id = seed["company_a"].id + + # Set owner_id + contact = await db_session.get(Contact, contact_id) + contact.owner_id = user_id + await db_session.commit() + + # Measure single access check + start = time.perf_counter() + for _ in range(100): + access = await eps.get_effective_access( + db_session, tenant_id, user_id, "contact", contact_id + ) + elapsed = time.perf_counter() - start + + avg_ms = (elapsed / 100) * 1000 + assert avg_ms < 50, f"Average get_effective_access took {avg_ms:.2f}ms (expected <50ms)" + assert access == "owner" + + async def test_check_entity_access_performance(self, db_session: AsyncSession): + """check_entity_access for a single entity completes quickly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + user_id = seed["admin_a"].id + contact_id = seed["company_a"].id + + # Set owner_id + contact = await db_session.get(Contact, contact_id) + contact.owner_id = user_id + await db_session.commit() + + # Measure single access check + start = time.perf_counter() + for _ in range(100): + result = await eps.check_entity_access( + db_session, tenant_id, user_id, "contact", contact_id, "read" + ) + elapsed = time.perf_counter() - start + + avg_ms = (elapsed / 100) * 1000 + assert avg_ms < 50, f"Average check_entity_access took {avg_ms:.2f}ms (expected <50ms)" + assert result is True + + async def test_permission_creation_performance(self, db_session: AsyncSession): + """Creating permissions in bulk completes quickly.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + contact_id = seed["company_a"].id + + # Set owner_id + contact = await db_session.get(Contact, contact_id) + contact.owner_id = owner_id + await db_session.commit() + + # Measure bulk permission creation + start = time.perf_counter() + for level in ["read", "write", "admin", "delete"]: + await eps.create_permission( + db_session, + tenant_id=tenant_id, + entity_type="contact", + entity_id=str(contact_id), + principal_type="user", + principal_id=str(viewer_id), + permission_level=level, + created_by=owner_id, + ) + elapsed = time.perf_counter() - start + + avg_ms = (elapsed / 4) * 1000 + assert avg_ms < 200, f"Average permission creation took {avg_ms:.2f}ms (expected <200ms)" + + async def test_visible_ids_with_mixed_ownership(self, db_session: AsyncSession): + """get_visible_ids with mixed ownership (owned + tenant + shared) performs well.""" + seed = await seed_tenant_and_users(db_session) + tenant_id = seed["tenant_a"].id + owner_id = seed["admin_a"].id + viewer_id = seed["viewer_a"].id + + # Create 500 owned, 300 tenant-owned, 200 shared contacts + contacts = [] + for i in range(500): + contacts.append(Contact( + tenant_id=tenant_id, type="person", + first_name=f"Owned{i}", last_name=f"Test{i}", + owner_id=owner_id, created_by=owner_id, updated_by=owner_id, + )) + for i in range(300): + contacts.append(Contact( + tenant_id=tenant_id, type="person", + first_name=f"Tenant{i}", last_name=f"Test{i}", + owner_id=None, created_by=owner_id, updated_by=owner_id, + )) + db_session.add_all(contacts) + await db_session.commit() + + # Grant viewer access to 200 contacts + shared_ids = [c.id for c in contacts[:200]] + for cid in shared_ids: + await eps.create_permission( + db_session, tenant_id=tenant_id, + entity_type="contact", entity_id=str(cid), + principal_type="user", principal_id=str(viewer_id), + permission_level="read", created_by=owner_id, + ) + + # Measure performance for viewer + start = time.perf_counter() + visible, access_map = await eps.get_visible_ids( + db_session, tenant_id, viewer_id, "contact" + ) + elapsed = time.perf_counter() - start + + assert elapsed < 3.0, f"get_visible_ids with mixed ownership took {elapsed:.3f}s (expected <3.0s)" + # Viewer should see: 200 shared + 300 tenant-owned = 500 + assert len(visible) >= 500, f"Expected >=500 visible, got {len(visible)}" +}