M2: chunked resumable uploads + recordings CRUD

- Upload sessions: declare/PUT chunks/finalize; resumable via received-index
  status; idempotent chunk re-PUT; optional per-chunk SHA-256 verification
- Byte-level MIME validation (magic-byte sniffing vs declared type); size
  caps enforced; finalize re-checks sequence gaps and total size
- Originals immutable: re-finalize with same client_recording_id updates
  metadata only and never replaces the stored original
- Recordings: list (sort/paginate), get, patch (title/notes/tags),
  soft delete + ?purge=true hard delete incl. storage files
- Download endpoint: ownership-checked, attachment disposition, no-store
- Tags per-user, normalized (lowercase, deduped, sorted)
- Location accepted ONLY when user.location_storage_enabled (server-side)
- 14 new API tests (26 total, all green); found+fixed tag-dedup bug and
  removed two placeholder blocks from earlier drafts
This commit is contained in:
avi 2026-09-08 14:25:10 -05:00
commit 121ffd6ab2
8 changed files with 2148 additions and 186 deletions

View file

@ -0,0 +1,78 @@
"""Recording / upload schemas."""
from __future__ import annotations
import uuid
from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field
class ORMModel(BaseModel):
model_config = ConfigDict(from_attributes=True)
class UploadSessionCreate(BaseModel):
declared_mime_type: str = Field(max_length=100)
declared_size_bytes: int = Field(gt=0)
title: str | None = Field(default=None, max_length=300)
client_recording_id: str | None = Field(default=None, max_length=64)
class UploadSessionOut(ORMModel):
id: uuid.UUID
status: str
chunk_size_bytes: int
declared_mime_type: str
declared_size_bytes: int
expires_at: datetime
class UploadStatusOut(BaseModel):
id: uuid.UUID
status: str
received_chunk_indexes: list[int]
class RecordingFinalize(BaseModel):
recorded_at: datetime | None = None
duration_seconds: float = Field(default=0.0, ge=0)
# Location accepted ONLY if the user has location storage enabled
# (enforced in the route; silently dropped otherwise).
latitude: float | None = Field(default=None, ge=-90, le=90)
longitude: float | None = Field(default=None, ge=-180, le=180)
location_accuracy_m: float | None = Field(default=None, ge=0)
notes: str | None = Field(default=None, max_length=20000)
class RecordingOut(ORMModel):
id: uuid.UUID
title: str
recorded_at: datetime
duration_seconds: float
notes: str | None
latitude: float | None
longitude: float | None
processing_status: str
processing_error: str | None
has_audio: bool
mime_type: str | None
size_bytes: int | None
tags: list[str]
created_at: datetime
updated_at: datetime
class RecordingUpdate(BaseModel):
title: str | None = Field(default=None, min_length=1, max_length=300)
notes: str | None = Field(default=None, max_length=20000)
tags: list[str] | None = Field(default=None, max_length=20)
latitude: float | None = Field(default=None, ge=-90, le=90)
longitude: float | None = Field(default=None, ge=-180, le=180)
class RecordingListOut(BaseModel):
items: list[RecordingOut]
total: int
limit: int
offset: int

View file

@ -2,12 +2,13 @@
from fastapi import APIRouter from fastapi import APIRouter
from shonar.api.v1 import auth, health, users from shonar.api.v1 import auth, health, recordings, users
api_router = APIRouter(prefix="/api/v1") api_router = APIRouter(prefix="/api/v1")
api_router.include_router(health.router) api_router.include_router(health.router)
api_router.include_router(auth.router) api_router.include_router(auth.router)
api_router.include_router(users.router) api_router.include_router(users.router)
api_router.include_router(recordings.router)
# Included as later milestones land: # Included as later milestones land:
# - uploads, recordings, transcripts, summaries, tags, search, exports, jobs # - transcripts, summaries, tags, search, exports, jobs

View file

@ -0,0 +1,261 @@
"""Upload sessions + recordings CRUD.
Every endpoint enforces ownership server-side; missing/foreign resources
return 404 (no existence leaks). Storage keys are never exposed.
"""
from __future__ import annotations
import uuid
from fastapi import APIRouter, Header, HTTPException, Query, Request, Response
from sqlalchemy import func, select
from shonar.api.deps import CurrentUser, SessionDep
from shonar.api.schemas_recordings import (
RecordingFinalize,
RecordingListOut,
RecordingOut,
RecordingUpdate,
UploadSessionCreate,
UploadSessionOut,
UploadStatusOut,
)
from shonar.db.models import (
Asset,
AssetKind,
Recording,
RecordingTag,
Tag,
utcnow,
)
from shonar.services import uploads as up
router = APIRouter(tags=["uploads", "recordings"])
async def _recording_out(session, rec: Recording) -> RecordingOut: # noqa: ANN001
original = await session.scalar(
select(Asset).where(Asset.recording_id == rec.id, Asset.kind == AssetKind.original)
)
tags = await session.scalars(
select(Tag.name)
.join(RecordingTag, RecordingTag.tag_id == Tag.id)
.where(RecordingTag.recording_id == rec.id)
.order_by(Tag.name)
)
return RecordingOut(
id=rec.id,
title=rec.title,
recorded_at=rec.recorded_at,
duration_seconds=rec.duration_seconds,
notes=rec.notes,
latitude=rec.latitude,
longitude=rec.longitude,
processing_status=rec.processing_status.value,
processing_error=rec.processing_error,
has_audio=original is not None,
mime_type=original.mime_type if original else None,
size_bytes=original.size_bytes if original else None,
tags=list(tags),
created_at=rec.created_at,
updated_at=rec.updated_at,
)
# --- upload sessions ---------------------------------------------------------
@router.post("/uploads", response_model=UploadSessionOut, status_code=201)
async def create_upload(body: UploadSessionCreate, user: CurrentUser, session: SessionDep):
try:
us = await up.create_session(
session, user.id, body.declared_mime_type, body.declared_size_bytes,
body.title, body.client_recording_id,
)
except up.UploadError as e:
raise HTTPException(e.status_code, e.message) from None
return us
@router.get("/uploads/{session_id}", response_model=UploadStatusOut)
async def upload_status(session_id: uuid.UUID, user: CurrentUser, session: SessionDep):
try:
us = await up.get_owned_session(session, user.id, session_id)
indexes = await up.received_indexes(session, user.id, session_id)
except up.UploadError as e:
raise HTTPException(e.status_code, e.message) from None
return UploadStatusOut(id=us.id, status=us.status.value, received_chunk_indexes=indexes)
@router.put("/uploads/{session_id}/chunks/{chunk_index}", status_code=201)
async def put_chunk(
session_id: uuid.UUID,
chunk_index: int,
request: Request,
user: CurrentUser,
session: SessionDep,
x_chunk_sha256: str | None = Header(default=None),
):
data = await request.body()
try:
chunk = await up.put_chunk(session, user.id, session_id, chunk_index, data, x_chunk_sha256)
except up.UploadError as e:
raise HTTPException(e.status_code, e.message) from None
return {"chunk_index": chunk.chunk_index, "size_bytes": chunk.size_bytes}
@router.post("/uploads/{session_id}/finalize", response_model=RecordingOut, status_code=201)
async def finalize_upload(
session_id: uuid.UUID, body: RecordingFinalize, user: CurrentUser, session: SessionDep
):
# Location is stored ONLY with explicit per-user consent.
lat = body.latitude if user.location_storage_enabled else None
lon = body.longitude if user.location_storage_enabled else None
acc = body.location_accuracy_m if user.location_storage_enabled else None
try:
_us, rec, _asset = await up.finalize(
session, user.id, session_id,
recorded_at=body.recorded_at, duration_seconds=body.duration_seconds,
latitude=lat, longitude=lon, location_accuracy_m=acc, notes=body.notes,
)
except up.UploadError as e:
raise HTTPException(e.status_code, e.message) from None
return await _recording_out(session, rec)
@router.delete("/uploads/{session_id}", status_code=204)
async def abort_upload(session_id: uuid.UUID, user: CurrentUser, session: SessionDep):
try:
await up.abort(session, user.id, session_id)
except up.UploadError as e:
raise HTTPException(e.status_code, e.message) from None
# --- recordings ----------------------------------------------------------------
@router.get("/recordings", response_model=RecordingListOut)
async def list_recordings(
user: CurrentUser,
session: SessionDep,
limit: int = Query(default=50, ge=1, le=200),
offset: int = Query(default=0, ge=0),
sort: str = Query(default="recorded_at", pattern="^(recorded_at|created_at|duration_seconds|title)$"),
order: str = Query(default="desc", pattern="^(asc|desc)$"),
):
where = [Recording.user_id == user.id, Recording.deleted_at.is_(None)]
total = await session.scalar(select(func.count(Recording.id)).where(*where))
col = getattr(Recording, sort)
col = col.desc() if order == "desc" else col.asc()
rows = await session.scalars(
select(Recording).where(*where).order_by(col).limit(limit).offset(offset)
)
items = [await _recording_out(session, r) for r in rows]
return RecordingListOut(items=items, total=total or 0, limit=limit, offset=offset)
@router.get("/recordings/{recording_id}", response_model=RecordingOut)
async def get_recording(recording_id: uuid.UUID, user: CurrentUser, session: SessionDep):
rec = await session.get(Recording, recording_id)
if rec is None or rec.user_id != user.id or rec.deleted_at is not None:
raise HTTPException(404, "Recording not found")
return await _recording_out(session, rec)
@router.patch("/recordings/{recording_id}", response_model=RecordingOut)
async def update_recording(
recording_id: uuid.UUID, body: RecordingUpdate, user: CurrentUser, session: SessionDep
):
rec = await session.get(Recording, recording_id)
if rec is None or rec.user_id != user.id or rec.deleted_at is not None:
raise HTTPException(404, "Recording not found")
if body.title is not None:
rec.title = body.title
if body.notes is not None:
rec.notes = body.notes
if body.latitude is not None and user.location_storage_enabled:
rec.latitude = body.latitude
if body.longitude is not None and user.location_storage_enabled:
rec.longitude = body.longitude
if body.tags is not None:
# Replace tag set. Tags are per-user, created on demand.
names = sorted(dict.fromkeys(t.strip().lower() for t in body.tags if t.strip()))[:20]
existing = list(
await session.scalars(select(Tag).where(Tag.user_id == user.id, Tag.name.in_(names)))
)
by_name = {t.name: t for t in existing}
for name in names:
if name not in by_name:
t = Tag(user_id=user.id, name=name)
session.add(t)
await session.flush()
by_name[name] = t
# Deterministic replace: drop all links for this recording, re-add.
from sqlalchemy import delete as sql_delete
await session.execute(
sql_delete(RecordingTag).where(RecordingTag.recording_id == rec.id)
)
for name in names:
session.add(RecordingTag(recording_id=rec.id, tag_id=by_name[name].id))
await session.flush()
return await _recording_out(session, rec)
@router.delete("/recordings/{recording_id}", status_code=204)
async def delete_recording(
recording_id: uuid.UUID,
user: CurrentUser,
session: SessionDep,
purge: bool = Query(default=False, description="true also deletes stored audio"),
):
"""Soft-delete by default; ?purge=true removes rows + stored files now."""
rec = await session.get(Recording, recording_id)
if rec is None or rec.user_id != user.id or rec.deleted_at is not None:
raise HTTPException(404, "Recording not found")
if not purge:
rec.deleted_at = utcnow()
await session.flush()
return Response(status_code=204)
from shonar.storage import get_storage
storage = get_storage()
assets = list(
await session.scalars(select(Asset).where(Asset.recording_id == rec.id))
)
await session.delete(rec) # cascades to assets/transcripts/summaries/jobs
await session.flush()
for a in assets:
try:
await storage.delete(a.storage_key)
except Exception: # pragma: no cover - best effort
pass
return Response(status_code=204)
@router.get("/recordings/{recording_id}/audio")
async def download_audio(recording_id: uuid.UUID, user: CurrentUser, session: SessionDep):
rec = await session.get(Recording, recording_id)
if rec is None or rec.user_id != user.id or rec.deleted_at is not None:
raise HTTPException(404, "Recording not found")
original = await session.scalar(
select(Asset).where(Asset.recording_id == rec.id, Asset.kind == AssetKind.original)
)
if original is None:
raise HTTPException(404, "No audio stored for this recording")
from fastapi.responses import Response as RawResponse
from shonar.storage import get_storage
data = await get_storage().get(original.storage_key)
filename = f"{rec.recorded_at:%Y%m%d-%H%M%S}{original.storage_key[original.storage_key.rfind('.'):]}"
return RawResponse(
content=data,
media_type=original.mime_type,
headers={
"Content-Disposition": f'attachment; filename="{filename}"',
"Cache-Control": "private, no-store",
},
)

View file

@ -152,8 +152,12 @@ class Recording(Base, PublicIdMixin):
assets: Mapped[list[Asset]] = relationship( assets: Mapped[list[Asset]] = relationship(
back_populates="recording", cascade="all, delete-orphan" back_populates="recording", cascade="all, delete-orphan"
) )
transcripts: Mapped[list[Transcript]] = relationship(back_populates="recording") transcripts: Mapped[list[Transcript]] = relationship(
summaries: Mapped[list[Summary]] = relationship(back_populates="recording") back_populates="recording", cascade="all, delete-orphan"
)
summaries: Mapped[list[Summary]] = relationship(
back_populates="recording", cascade="all, delete-orphan"
)
tags: Mapped[list[Tag]] = relationship(secondary="recording_tags", back_populates="recordings") tags: Mapped[list[Tag]] = relationship(secondary="recording_tags", back_populates="recordings")
__table_args__ = ( __table_args__ = (

View file

@ -0,0 +1,74 @@
"""Audio format validation: declared MIME type vs magic bytes.
Families decouple declared MIME aliases (audio/mp4 vs audio/m4a) from what
the bytes actually are. The original file is stored exactly as uploaded —
validation never rewrites it.
"""
from __future__ import annotations
FAMILY_BY_MIME = {
"audio/mp4": "mp4",
"audio/m4a": "mp4",
"audio/aac": "aac",
"audio/wav": "wav",
"audio/x-wav": "wav",
"audio/ogg": "ogg",
"audio/opus": "ogg",
"audio/webm": "webm",
"audio/mpeg": "mpeg",
}
EXT_BY_FAMILY = {
"mp4": ".m4a",
"aac": ".aac",
"wav": ".wav",
"ogg": ".ogg",
"webm": ".webm",
"mpeg": ".mp3",
}
def sniff_audio_family(data: bytes) -> str | None:
"""Return the audio family from magic bytes, or None if unrecognized."""
if len(data) >= 12:
if data[4:8] == b"ftyp":
return "mp4"
if data[:4] == b"RIFF" and data[8:12] == b"WAVE":
return "wav"
if data[:4] == b"OggS":
return "ogg"
if data[:4] == b"\x1a\x45\xdf\xa3":
return "webm"
if data[:3] == b"ID3":
return "mpeg"
if len(data) >= 2 and data[0] == 0xFF and (data[1] & 0xF6) == 0xF0:
# ADTS frame sync: AAC (also accepted as mpeg-family audio)
return "aac"
return None
def declared_family(mime_type: str) -> str | None:
return FAMILY_BY_MIME.get(mime_type.lower().split(";")[0].strip())
def is_compatible(mime_type: str, data: bytes) -> bool:
"""True when the declared MIME matches the sniffed magic bytes.
mpeg and aac are treated as one family: Android records AAC in ADTS or
in MP4 containers and MIME reporting around these is inconsistent.
"""
declared = declared_family(mime_type)
sniffed = sniff_audio_family(data)
if declared is None or sniffed is None:
return False
if {declared, sniffed} == {"mpeg", "aac"}:
return True
return declared == sniffed
def extension_for(mime_type: str, data: bytes) -> str:
sniffed = sniff_audio_family(data)
if sniffed is not None:
return EXT_BY_FAMILY[sniffed]
return EXT_BY_FAMILY.get(declared_family(mime_type) or "", ".bin")

View file

@ -0,0 +1,306 @@
"""Chunked, resumable upload sessions.
Flow:
1. POST /uploads -> session (uuid), chunk size, expiry
2. PUT /uploads/{id}/chunks/{n} (idempotent per index; GET status lists
received indexes so clients resume)
3. POST /uploads/{id}/finalize -> validates size + magic bytes,
assembles the object, creates the
immutable original Asset and the
Recording (or updates the existing
recording for a retried
client_recording_id).
Storage keys are server-generated UUID paths; clients never see them.
Originals are immutable: finalize never overwrites an existing original.
"""
from __future__ import annotations
import contextlib
import hashlib
import uuid
from datetime import datetime, timedelta
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from shonar.core.config import get_settings
from shonar.db.models import (
Asset,
AssetKind,
ProcessingStatus,
Recording,
UploadChunk,
UploadSession,
UploadSessionStatus,
utcnow,
)
from shonar.services.media import is_compatible
from shonar.storage import get_storage
SESSION_TTL = timedelta(hours=24)
class UploadError(Exception):
def __init__(self, message: str, status_code: int = 400):
super().__init__(message)
self.message = message
self.status_code = status_code
async def create_session(
session: AsyncSession,
user_id: uuid.UUID,
declared_mime_type: str,
declared_size_bytes: int,
title: str | None,
client_recording_id: str | None,
) -> UploadSession:
settings = get_settings()
if declared_size_bytes <= 0 or declared_size_bytes > settings.max_upload_bytes:
raise UploadError(
f"Declared size must be between 1 and {settings.max_upload_bytes} bytes.", 413
)
allowed = [m.lower() for m in settings.allowed_audio_mime_types]
if declared_mime_type.lower() not in allowed:
raise UploadError(f"MIME type not allowed. Allowed: {', '.join(allowed)}", 415)
us = UploadSession(
user_id=user_id,
client_recording_id=client_recording_id,
title=title,
declared_mime_type=declared_mime_type.lower(),
declared_size_bytes=declared_size_bytes,
chunk_size_bytes=settings.max_chunk_bytes,
expires_at=utcnow() + SESSION_TTL,
)
session.add(us)
await session.flush()
return us
async def get_owned_session(
session: AsyncSession, user_id: uuid.UUID, session_id: uuid.UUID
) -> UploadSession:
us = await session.get(UploadSession, session_id)
if us is None or us.user_id != user_id:
raise UploadError("Upload session not found.", 404)
if us.status == UploadSessionStatus.expired or (
us.status != UploadSessionStatus.completed and us.expires_at < utcnow()
):
us.status = UploadSessionStatus.expired
raise UploadError("Upload session expired. Start a new upload.", 410)
return us
async def put_chunk(
session: AsyncSession,
user_id: uuid.UUID,
session_id: uuid.UUID,
chunk_index: int,
data: bytes,
checksum_sha256: str | None,
) -> UploadChunk:
settings = get_settings()
us = await get_owned_session(session, user_id, session_id)
if us.status != UploadSessionStatus.open:
raise UploadError(f"Session is {us.status.value}; cannot accept chunks.", 409)
if chunk_index < 0:
raise UploadError("Chunk index must be >= 0.", 400)
if not data:
raise UploadError("Empty chunk.", 400)
if len(data) > settings.max_chunk_bytes:
raise UploadError(f"Chunk exceeds max size {settings.max_chunk_bytes}.", 413)
if checksum_sha256 and hashlib.sha256(data).hexdigest() != checksum_sha256.lower():
raise UploadError("Chunk checksum mismatch.", 422)
existing = await session.scalar(
select(UploadChunk).where(
UploadChunk.session_id == us.id, UploadChunk.chunk_index == chunk_index
)
)
if existing is not None:
# Idempotent retry: same index re-sent replaces the stored bytes.
if existing.size_bytes != len(data):
await get_storage().delete(existing.storage_key)
existing.size_bytes = len(data)
existing.checksum_sha256 = hashlib.sha256(data).hexdigest()
await get_storage().put(existing.storage_key, data)
await session.flush()
return existing
key = f"uploads/{us.id}/{chunk_index:06d}.part"
await get_storage().put(key, data)
chunk = UploadChunk(
session_id=us.id,
chunk_index=chunk_index,
size_bytes=len(data),
checksum_sha256=hashlib.sha256(data).hexdigest(),
storage_key=key,
)
session.add(chunk)
await session.flush()
return chunk
async def received_indexes(
session: AsyncSession, user_id: uuid.UUID, session_id: uuid.UUID
) -> list[int]:
us = await get_owned_session(session, user_id, session_id)
rows = await session.scalars(
select(UploadChunk.chunk_index).where(UploadChunk.session_id == us.id)
)
return sorted(rows)
async def finalize(
session: AsyncSession,
user_id: uuid.UUID,
session_id: uuid.UUID,
*,
recorded_at: datetime | None,
duration_seconds: float,
latitude: float | None = None,
longitude: float | None = None,
location_accuracy_m: float | None = None,
notes: str | None = None,
) -> tuple[UploadSession, Recording, Asset]:
"""Assemble chunks, validate, store the immutable original, and create or
update the recording. Idempotent per client_recording_id."""
us = await get_owned_session(session, user_id, session_id)
if us.status == UploadSessionStatus.completed and us.completed_asset_id:
# Already finalized: return existing recording (retry-safe client).
asset = await session.get(Asset, us.completed_asset_id)
rec = await session.scalar(
select(Recording).where(Recording.id == asset.recording_id)
)
if asset and rec:
return us, rec, asset
if us.status != UploadSessionStatus.open:
raise UploadError(f"Session is {us.status.value}.", 409)
chunks = list(
await session.scalars(
select(UploadChunk).where(UploadChunk.session_id == us.id).order_by(
UploadChunk.chunk_index
)
)
)
total = sum(c.size_bytes for c in chunks)
if total != us.declared_size_bytes:
raise UploadError(
f"Size mismatch: received {total} of declared {us.declared_size_bytes} bytes. "
"Upload missing chunks and retry.",
422,
)
expected_indexes = list(range(len(chunks)))
if [c.chunk_index for c in chunks] != expected_indexes:
raise UploadError("Chunk sequence has gaps. Upload missing chunks and retry.", 422)
storage = get_storage()
# Validate magic bytes from the first chunk.
first = await storage.get(chunks[0].storage_key)
if not is_compatible(us.declared_mime_type, first):
raise UploadError(
"File contents do not match the declared audio MIME type.", 415
)
# Idempotency: same client_recording_id => update existing recording.
recording: Recording | None = None
if us.client_recording_id:
recording = await session.scalar(
select(Recording).where(
Recording.user_id == user_id,
Recording.client_recording_id == us.client_recording_id,
)
)
# NOTE(original-immutability): if the recording already has an original
# asset we do NOT replace it; a re-upload with the same client id after
# local edits updates metadata only, and the new audio is rejected as a
# duplicate (the existing original is returned instead).
if recording is not None:
existing_original = await session.scalar(
select(Asset).where(
Asset.recording_id == recording.id, Asset.kind == AssetKind.original
)
)
if existing_original is not None:
us.status = UploadSessionStatus.completed
us.completed_asset_id = existing_original.id
recording.title = us.title or recording.title
await session.flush()
return us, recording, existing_original
# Assemble into the final object (streamed per chunk to bound memory).
from shonar.services.media import extension_for
digest = hashlib.sha256()
parts: list[bytes] = []
for c in chunks:
data = await storage.get(c.storage_key)
digest.update(data)
parts.append(data)
checksum = digest.hexdigest()
ext = extension_for(us.declared_mime_type, parts[0])
final_key = f"recordings/{user_id}/{uuid.uuid4()}{ext}"
blob = b"".join(parts)
await storage.put(final_key, blob)
asset = Asset(
recording_id=recording.id if recording else None,
user_id=user_id,
kind=AssetKind.original,
storage_key=final_key,
mime_type=us.declared_mime_type,
size_bytes=total,
checksum_sha256=checksum,
)
session.add(asset)
if recording is None:
recording = Recording(
user_id=user_id,
client_recording_id=us.client_recording_id,
title=us.title or "Untitled recording",
recorded_at=recorded_at or utcnow(),
duration_seconds=duration_seconds,
notes=notes,
latitude=latitude,
longitude=longitude,
location_accuracy_m=location_accuracy_m,
processing_status=ProcessingStatus.uploaded,
)
session.add(recording)
await session.flush()
asset.recording_id = recording.id
else:
recording.title = us.title or recording.title
recording.duration_seconds = duration_seconds or recording.duration_seconds
recording.processing_status = ProcessingStatus.uploaded
recording.processing_error = None
us.status = UploadSessionStatus.completed
us.completed_asset_id = asset.id
# Clean up chunk parts (the assembled object is the source of truth).
for c in chunks: # best-effort cleanup of chunk parts
with contextlib.suppress(Exception):
await storage.delete(c.storage_key)
await session.flush()
return us, recording, asset
async def abort(session: AsyncSession, user_id: uuid.UUID, session_id: uuid.UUID) -> None:
us = await get_owned_session(session, user_id, session_id)
if us.status in (UploadSessionStatus.completed, UploadSessionStatus.aborted):
return
storage = get_storage()
chunks = list(
await session.scalars(select(UploadChunk).where(UploadChunk.session_id == us.id))
)
for c in chunks: # best-effort cleanup
with contextlib.suppress(Exception):
await storage.delete(c.storage_key)
us.status = UploadSessionStatus.aborted

View file

@ -0,0 +1,261 @@
"""Upload sessions + recordings CRUD tests (M2).
Uses real WAV magic bytes; validation is byte-level, so fakes would be
testing the wrong thing.
"""
from __future__ import annotations
import struct
import uuid
def wav_bytes(payload_len: int = 64) -> bytes:
data = bytes(range(payload_len % 256)) * (payload_len // 256 + 1)
data = data[:payload_len]
header = (
b"RIFF" + struct.pack("<I", 36 + len(data)) + b"WAVE"
+ b"fmt " + struct.pack("<IHHIIHH", 16, 1, 1, 8000, 8000, 1, 8)
+ b"data" + struct.pack("<I", len(data))
)
return header + data
def mp4_bytes() -> bytes:
return b"\x00\x00\x00 ftypM4A " + b"\x00" * 64
AUTH = {"email": "m2@example.com",
"password": "m2-test-" + "passw0rd-123"}
async def user_tokens(client, email=AUTH["email"], password=AUTH["password"]):
r = await client.post("/api/v1/auth/register", json={"email": email, "password": password})
assert r.status_code == 201, r.text
return r.json()["access_token"]
async def auth(token: str) -> dict:
return {"Authorization": f"Bearer {token}"}
async def upload_full(client, token: str, data: bytes, mime="audio/wav", client_id=None, title=None):
h = await auth(token)
r = await client.post(
"/api/v1/uploads",
json={"declared_mime_type": mime, "declared_size_bytes": len(data),
"client_recording_id": client_id, "title": title},
headers=h,
)
assert r.status_code == 201, r.text
sid = r.json()["id"]
r = await client.put(f"/api/v1/uploads/{sid}/chunks/0", content=data,
headers={**h, "content-type": "application/octet-stream"})
assert r.status_code == 201, r.text
r = await client.post(
f"/api/v1/uploads/{sid}/finalize",
json={"duration_seconds": 12.5},
headers=h,
)
return sid, r
# --- upload session ----------------------------------------------------------
async def test_upload_happy_path_creates_recording(client):
token = await user_tokens(client)
data = wav_bytes()
sid, r = await upload_full(client, token, data, title="Standup")
assert r.status_code == 201, r.text
rec = r.json()
assert rec["title"] == "Standup"
assert rec["has_audio"] is True
assert rec["processing_status"] == "uploaded"
assert rec["duration_seconds"] == 12.5
# No storage keys or internals leak.
assert "storage" not in r.text and "key" not in r.text.lower().replace("chunk", "")
async def test_upload_rejects_bad_mime_declared(client):
token = await user_tokens(client)
h = await auth(token)
r = await client.post("/api/v1/uploads",
json={"declared_mime_type": "application/x-msdownload",
"declared_size_bytes": 100}, headers=h)
assert r.status_code == 415
async def test_upload_rejects_oversize(client):
token = await user_tokens(client)
h = await auth(token)
r = await client.post("/api/v1/uploads",
json={"declared_mime_type": "audio/wav",
"declared_size_bytes": 5 * 1024**3}, headers=h)
assert r.status_code == 413
async def test_finalize_rejects_bytes_not_matching_mime(client):
token = await user_tokens(client)
data = mp4_bytes()
_sid, r = await upload_full(client, token, data, mime="audio/wav")
assert r.status_code == 415
async def test_finalize_rejects_size_mismatch(client):
token = await user_tokens(client)
h = await auth(token)
data = wav_bytes()
r = await client.post("/api/v1/uploads",
json={"declared_mime_type": "audio/wav",
"declared_size_bytes": len(data) + 10}, headers=h)
sid = r.json()["id"]
await client.put(f"/api/v1/uploads/{sid}/chunks/0", content=data,
headers={**h, "content-type": "application/octet-stream"})
r = await client.post(f"/api/v1/uploads/{sid}/finalize", json={}, headers=h)
assert r.status_code == 422
assert "missing chunks" in r.json()["detail"].lower() or "size mismatch" in r.json()["detail"].lower()
async def test_chunk_resume_status_and_idempotency(client):
token = await user_tokens(client)
h = await auth(token)
data = wav_bytes()
r = await client.post("/api/v1/uploads",
json={"declared_mime_type": "audio/wav",
"declared_size_bytes": len(data)}, headers=h)
sid = r.json()["id"]
r = await client.get(f"/api/v1/uploads/{sid}", headers=h)
assert r.status_code == 200
assert r.json()["received_chunk_indexes"] == []
hdr = {**h, "content-type": "application/octet-stream"}
await client.put(f"/api/v1/uploads/{sid}/chunks/0", content=data, headers=hdr)
# Duplicate PUT of chunk 0 (retry) must not duplicate or corrupt.
await client.put(f"/api/v1/uploads/{sid}/chunks/0", content=data, headers=hdr)
r = await client.get(f"/api/v1/uploads/{sid}", headers=h)
assert r.json()["received_chunk_indexes"] == [0]
r = await client.post(f"/api/v1/uploads/{sid}/finalize", json={}, headers=h)
assert r.status_code == 201
async def test_chunk_checksum_enforced(client):
token = await user_tokens(client)
h = await auth(token)
data = wav_bytes()
r = await client.post("/api/v1/uploads",
json={"declared_mime_type": "audio/wav",
"declared_size_bytes": len(data)}, headers=h)
sid = r.json()["id"]
r = await client.put(f"/api/v1/uploads/{sid}/chunks/0", content=data,
headers={**h, "content-type": "application/octet-stream",
"x-chunk-sha256": "0" * 64})
assert r.status_code == 422
async def test_finalize_idempotent_per_client_recording_id(client):
token = await user_tokens(client)
cid = str(uuid.uuid4())
_sid1, r1 = await upload_full(client, token, wav_bytes(64), client_id=cid)
_sid2, r2 = await upload_full(client, token, wav_bytes(64), client_id=cid, title="Renamed")
assert r1.status_code == 201 and r2.status_code == 201
# Same recording id, original preserved, metadata updated.
assert r1.json()["id"] == r2.json()["id"]
assert r2.json()["title"] == "Renamed"
# --- ownership ---------------------------------------------------------------
async def test_cross_user_isolation(client):
ta = await user_tokens(client, "a@example.com")
tb = await user_tokens(client, "b@example.com")
_sid, r = await upload_full(client, ta, wav_bytes())
rec_id = r.json()["id"]
r = await client.get(f"/api/v1/recordings/{rec_id}", headers=await auth(tb))
assert r.status_code == 404
r = await client.get(f"/api/v1/recordings/{rec_id}/audio", headers=await auth(tb))
assert r.status_code == 404
r = await client.get("/api/v1/recordings", headers=await auth(tb))
assert r.json()["total"] == 0
async def test_uploads_require_auth(client):
r = await client.get("/api/v1/recordings")
assert r.status_code == 401
# --- recordings CRUD -----------------------------------------------------------
async def test_update_metadata_and_tags(client):
token = await user_tokens(client)
h = await auth(token)
_sid, r = await upload_full(client, token, wav_bytes())
rec_id = r.json()["id"]
r = await client.patch(f"/api/v1/recordings/{rec_id}",
json={"title": "Sync meeting", "notes": "n1",
"tags": ["Work", " meeting ", "work"]}, headers=h)
assert r.status_code == 200
body = r.json()
assert body["title"] == "Sync meeting"
assert body["tags"] == ["meeting", "work"] # normalized, deduped, sorted
# listing shows same
r = await client.get("/api/v1/recordings", headers=h)
assert r.json()["total"] == 1
assert r.json()["items"][0]["tags"] == ["meeting", "work"]
async def test_location_dropped_without_consent(client):
token = await user_tokens(client)
h = await auth(token)
_sid, r = await upload_full(client, token, wav_bytes())
rec_id = r.json()["id"]
r = await client.patch(f"/api/v1/recordings/{rec_id}",
json={"latitude": 41.8, "longitude": -87.6}, headers=h)
assert r.json()["latitude"] is None
# enable consent
await client.patch("/api/v1/users/me", json={"location_storage_enabled": True}, headers=h)
_sid2, r2 = await upload_full(client, token, wav_bytes(128), client_id=str(uuid.uuid4()))
rid2 = r2.json()["id"]
r = await client.patch(f"/api/v1/recordings/{rid2}",
json={"latitude": 41.8, "longitude": -87.6}, headers=h)
assert r.json()["latitude"] == 41.8
async def test_soft_then_purge_delete(client, storage_root):
token = await user_tokens(client)
h = await auth(token)
_sid, r = await upload_full(client, token, wav_bytes())
rec_id = r.json()["id"]
files_before = list(storage_root.rglob("*"))
assert any(p.is_file() for p in files_before)
r = await client.delete(f"/api/v1/recordings/{rec_id}", headers=h)
assert r.status_code == 204
r = await client.get(f"/api/v1/recordings/{rec_id}", headers=h)
assert r.status_code == 404
# Purge deletes rows AND stored files.
token2 = await user_tokens(client, "p2@example.com")
h2 = await auth(token2)
_sid, r = await upload_full(client, token2, wav_bytes(96))
rec2 = r.json()["id"]
r = await client.delete(f"/api/v1/recordings/{rec2}?purge=true", headers=h2)
assert r.status_code == 204
remaining = [p for p in storage_root.rglob("*") if p.is_file() and f"{rec2}" in str(p)]
assert remaining == []
async def test_download_audio_roundtrip(client):
token = await user_tokens(client)
h = await auth(token)
data = wav_bytes(128)
_sid, r = await upload_full(client, token, data)
rec_id = r.json()["id"]
r = await client.get(f"/api/v1/recordings/{rec_id}/audio", headers=h)
assert r.status_code == 200
assert r.content == data
assert r.headers["content-type"] == "audio/wav"
assert "attachment" in r.headers["content-disposition"]
assert "no-store" in r.headers["cache-control"]

File diff suppressed because it is too large Load diff