Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 66 additions & 0 deletions backend/alembic/versions/4315ee2f32f5_media_service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
"""media service

Revision ID: 4315ee2f32f5
Revises: ce8c0c720681
Create Date: 2026-08-14 21:54:24.487092
"""

from collections.abc import Sequence

from alembic import op

revision: str = "4315ee2f32f5"
down_revision: str | None = "ce8c0c720681"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None

UPGRADE = """
CREATE TABLE media (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
owner_id UUID NOT NULL,
s3_key TEXT NOT NULL UNIQUE,
original_filename TEXT NOT NULL,
mime_type TEXT NOT NULL,
size_bytes BIGINT NOT NULL,
purpose TEXT NOT NULL,
visibility TEXT NOT NULL CHECK (visibility IN ('public', 'private')),
status TEXT NOT NULL,
metadata JSONB,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),

UNIQUE (id, visibility)
);

CREATE INDEX ix_media_owner ON media (owner_id, created_at DESC);

CREATE TABLE media_variants (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
media_id UUID NOT NULL,
variant_type TEXT NOT NULL,
s3_key TEXT NOT NULL UNIQUE,
mime_type TEXT NOT NULL,
size_bytes BIGINT NOT NULL,
visibility TEXT NOT NULL,
metadata JSONB,
status TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),

FOREIGN KEY (media_id, visibility) REFERENCES media (id, visibility) ON DELETE CASCADE,
UNIQUE (media_id, variant_type)
);
"""

DOWNGRADE = """
DROP TABLE IF EXISTS media_variants;
DROP TABLE IF EXISTS media;
"""


def upgrade() -> None:
op.execute(UPGRADE)


def downgrade() -> None:
op.execute(DOWNGRADE)
23 changes: 22 additions & 1 deletion backend/src/admin/api/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from fastapi import Depends, Request
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer

from admin.core.audit import AuditLog
from admin.core.cache import Cache as CacheProtocol
from admin.core.config import Settings, get_settings
from admin.core.errors import (
Expand All @@ -22,6 +23,7 @@
AccessRequestRepository,
AuditRepository,
InvitationRepository,
MediaRepository,
RoleRepository,
UserRepository,
)
Expand All @@ -30,6 +32,7 @@
from admin.services.access import AccessService
from admin.services.access_request import AccessRequestService
from admin.services.invitation import InvitationService
from admin.services.media import MediaService
from admin.services.member import MemberService

bearer_scheme = HTTPBearer(auto_error=True)
Expand Down Expand Up @@ -91,11 +94,24 @@ def get_audit_repository(connection: Connection) -> AuditRepository:
return AuditRepository(connection)


async def get_audit_log(connection: Connection) -> AsyncIterator[AuditLog]:
log = AuditLog()
yield log
if log.entries:
await AuditRepository(connection).record_many(log.entries)


def get_media_repository(connection: Connection) -> MediaRepository:
return MediaRepository(connection)


Users = Annotated[UserRepository, Depends(get_user_repository)]
Roles = Annotated[RoleRepository, Depends(get_role_repository)]
Invitations = Annotated[InvitationRepository, Depends(get_invitation_repository)]
AccessRequests = Annotated[AccessRequestRepository, Depends(get_access_request_repository)]
Audit = Annotated[AuditRepository, Depends(get_audit_repository)]
Audit = Annotated[AuditLog, Depends(get_audit_log)]
AuditEntries = Annotated[AuditRepository, Depends(get_audit_repository)]
Media = Annotated[MediaRepository, Depends(get_media_repository)]


def get_access_service(
Expand Down Expand Up @@ -138,10 +154,15 @@ def get_access_request_service(
)


def get_media_service(media: Media, storage: Storage, audit: Audit) -> MediaService:
return MediaService(media=media, storage=storage, audit=audit)


AccessServiceDep = Annotated[AccessService, Depends(get_access_service)]
MemberServiceDep = Annotated[MemberService, Depends(get_member_service)]
InvitationServiceDep = Annotated[InvitationService, Depends(get_invitation_service)]
AccessRequestServiceDep = Annotated[AccessRequestService, Depends(get_access_request_service)]
MediaServiceDep = Annotated[MediaService, Depends(get_media_service)]


async def get_identity(
Expand Down
4 changes: 4 additions & 0 deletions backend/src/admin/api/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,10 @@

from admin.api.v1 import (
access_requests,
audit,
health,
invitations,
media,
members,
roles,
session,
Expand All @@ -15,6 +17,8 @@
api_router.include_router(invitations.router)
api_router.include_router(access_requests.router)
api_router.include_router(roles.router)
api_router.include_router(media.router)
api_router.include_router(audit.router)

root_router = APIRouter()
root_router.include_router(health.router)
29 changes: 29 additions & 0 deletions backend/src/admin/api/v1/audit.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import uuid
from typing import Annotated

from fastapi import APIRouter, Depends

from admin.api.dependencies import AuditEntries, AuthContext, require
from admin.domain.permissions import Permission
from admin.schemas.audit import AuditLogRead
from admin.schemas.base import CursorPage, CursorParams

router = APIRouter(prefix="/audit", tags=["audit"])


@router.get("", response_model=CursorPage[AuditLogRead])
async def list_audit_entries(
entries: AuditEntries,
params: Annotated[CursorParams, Depends()],
_: Annotated[AuthContext, Depends(require(Permission.AUDIT_READ))],
actor_id: uuid.UUID | None = None,
resource_type: str | None = None,
resource_id: str | None = None,
) -> CursorPage[AuditLogRead]:
rows = await entries.list_entries(
params,
actor_id=actor_id,
resource_type=resource_type,
resource_id=resource_id,
)
return CursorPage[AuditLogRead].build(rows, params, lambda item: [item.created_at, item.id])
77 changes: 77 additions & 0 deletions backend/src/admin/api/v1/media.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
import uuid
from typing import Annotated

from fastapi import APIRouter, Query, status

from admin.api.dependencies import CurrentUser, MediaServiceDep
from admin.schemas.media import (
MediaPresetRead,
MediaRead,
MediaUploadRequest,
MediaUploadTicket,
)

router = APIRouter(prefix="/media", tags=["media"])


@router.get("/presets", response_model=list[MediaPresetRead])
async def list_presets(service: MediaServiceDep, _: CurrentUser) -> list[MediaPresetRead]:
return service.presets()


@router.post(
"/upload-tickets",
response_model=list[MediaUploadTicket],
status_code=status.HTTP_201_CREATED,
)
async def create_upload_tickets(
payload: MediaUploadRequest, service: MediaServiceDep, context: CurrentUser
) -> list[MediaUploadTicket]:
return await service.create_upload_tickets(actor=context.user, payload=payload)


@router.get("", response_model=list[MediaRead])
async def list_media(
service: MediaServiceDep,
context: CurrentUser,
ids: Annotated[list[uuid.UUID], Query()],
) -> list[MediaRead]:
return await service.get_files(
actor=context.user, actor_permissions=context.permissions, media_ids=ids
)


@router.delete("", status_code=status.HTTP_204_NO_CONTENT)
async def delete_media(
service: MediaServiceDep,
context: CurrentUser,
ids: Annotated[list[uuid.UUID], Query()],
) -> None:
await service.delete_files(
actor=context.user, actor_permissions=context.permissions, media_ids=ids
)


@router.post("/{media_id}/complete", response_model=MediaRead)
async def complete_upload(
media_id: uuid.UUID, service: MediaServiceDep, context: CurrentUser
) -> MediaRead:
return await service.complete_upload(actor=context.user, media_id=media_id)


@router.get("/{media_id}", response_model=MediaRead)
async def get_media(
media_id: uuid.UUID, service: MediaServiceDep, context: CurrentUser
) -> MediaRead:
return await service.get_file(
actor=context.user, actor_permissions=context.permissions, media_id=media_id
)


@router.delete("/{media_id}", status_code=status.HTTP_204_NO_CONTENT)
async def delete_single_media(
media_id: uuid.UUID, service: MediaServiceDep, context: CurrentUser
) -> None:
await service.delete_file(
actor=context.user, actor_permissions=context.permissions, media_id=media_id
)
11 changes: 11 additions & 0 deletions backend/src/admin/core/audit.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
from dataclasses import dataclass, field

from admin.schemas.audit import AuditEntry


@dataclass(slots=True)
class AuditLog:
entries: list[AuditEntry] = field(default_factory=list)

def add(self, *entries: AuditEntry) -> None:
self.entries.extend(entries)
44 changes: 38 additions & 6 deletions backend/src/admin/core/storage.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
import asyncio
import uuid
from dataclasses import dataclass
from datetime import UTC, datetime
from datetime import UTC, datetime, timedelta
from enum import StrEnum
from itertools import batched
from typing import Any

import boto3
Expand All @@ -24,10 +25,10 @@ class MediaVisibility(StrEnum):
PRESIGN_CACHE_RATIO = 0.8
PUBLIC_PREFIX = "public"
PRIVATE_PREFIX = "private"
S3_DELETE_BATCH_SIZE = 1000

IMMUTABLE_CACHE_CONTROL = "public, max-age=31536000, immutable"
PRIVATE_CACHE_CONTROL = "private, max-age=0, no-store"
STRING_ADAPTER = TypeAdapter(str)


def cache_control_for(visibility: MediaVisibility) -> str:
Expand All @@ -50,6 +51,15 @@ class ObjectMetadata:
content_type: str


@dataclass(frozen=True, slots=True)
class MediaUrl:
url: str
expires_at: datetime | None = None


MEDIA_URL_ADAPTER = TypeAdapter(MediaUrl)


def build_object_key(*, visibility: MediaVisibility, filename: str) -> str:
prefix = PUBLIC_PREFIX if visibility is MediaVisibility.PUBLIC else PRIVATE_PREFIX
stamp = datetime.now(UTC)
Expand All @@ -70,6 +80,10 @@ def __init__(self, config: StorageConfig, cache: Cache) -> None:
config=Config(signature_version="s3v4", s3={"addressing_style": "path"}),
)

@property
def max_upload_bytes(self) -> int:
return self._config.max_upload_bytes

def create_upload_ticket(
self,
*,
Expand Down Expand Up @@ -103,22 +117,28 @@ def create_upload_ticket(
def public_url(self, key: str) -> str:
return self._config.public_url_for(key)

async def signed_download_url(self, key: str) -> str:
async def url_for(self, *, key: str, visibility: MediaVisibility) -> MediaUrl:
if visibility is MediaVisibility.PUBLIC:
return MediaUrl(url=self.public_url(key))
return await self.signed_download_url(key)

async def signed_download_url(self, key: str) -> MediaUrl:
ttl = self._config.download_url_ttl_seconds

async def generate() -> str:
return await asyncio.to_thread(
async def generate() -> MediaUrl:
url = await asyncio.to_thread(
self._client.generate_presigned_url,
"get_object",
Params={"Bucket": self._config.bucket_name, "Key": key},
ExpiresIn=ttl,
)
return MediaUrl(url=url, expires_at=datetime.now(UTC) + timedelta(seconds=ttl))

return await self._cache.fetch(
CacheNamespace.MEDIA,
f"presign:get:{key}",
generate,
adapter=STRING_ADAPTER,
adapter=MEDIA_URL_ADAPTER,
ttl=ttl * PRESIGN_CACHE_RATIO,
)

Expand All @@ -140,3 +160,15 @@ async def delete(self, key: str) -> None:
self._client.delete_object, Bucket=self._config.bucket_name, Key=key
)
await self._cache.bump(CacheNamespace.MEDIA)

async def delete_many(self, keys: list[str]) -> None:
if not keys:
return

for batch in batched(keys, S3_DELETE_BATCH_SIZE):
await asyncio.to_thread(
self._client.delete_objects,
Bucket=self._config.bucket_name,
Delete={"Objects": [{"Key": key} for key in batch], "Quiet": True},
)
await self._cache.bump(CacheNamespace.MEDIA)
6 changes: 6 additions & 0 deletions backend/src/admin/domain/enums.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@ class AccessRequestStatus(StrEnum):
DENIED = "denied"


class MediaPurpose(StrEnum):
AVATAR = "avatar"


class AuditAction(StrEnum):
USER_PROVISIONED = "user.provisioned"
USER_SUSPENDED = "user.suspended"
Expand All @@ -40,3 +44,5 @@ class AuditAction(StrEnum):
ACCESS_REQUEST_CREATED = "access_request.created"
ACCESS_REQUEST_APPROVED = "access_request.approved"
ACCESS_REQUEST_DENIED = "access_request.denied"
MEDIA_UPLOADED = "media.uploaded"
MEDIA_DELETED = "media.deleted"
Loading
Loading