- shared/ = portable Android-origin sources vendored from deferred/desktop-server (app/build.gradle.kts srcDir repointed; PlaybackController.kt excluded as Android-only) - backend/ = bundled-lite engine (SQLite + inline queue); .venv symlinked from the old checkout, PYTHONPATH pins THIS backend's code over any editable install - repoRoot() resolves this project dir (env SHONAR_REPO still wins); desktop-dev.sh watches shared/ + backend/ - Verified: :app:compileKotlin + :app:test green (23 tests); engine boots on :8010, self-migrates, /healthz ok
342 lines
13 KiB
Python
342 lines
13 KiB
Python
"""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,
|
|
transcription_model: str | None = 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,
|
|
transcription_model=(transcription_model.strip() if transcription_model else None),
|
|
)
|
|
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,
|
|
transcription_model: 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."""
|
|
from shonar.services.ai import ProviderConfigError
|
|
from shonar.services.ai.model_registry import (
|
|
effective_model,
|
|
get_global_default_model,
|
|
validate_model_name,
|
|
)
|
|
|
|
us = await get_owned_session(session, user_id, session_id)
|
|
# Resolve + validate the transcription model before touching audio.
|
|
# Finalize body wins over the session's upload-screen choice; an
|
|
# explicit override updates history, the default never rewrites it.
|
|
override_raw = transcription_model or us.transcription_model
|
|
override: str | None = None
|
|
if override_raw:
|
|
try:
|
|
override = validate_model_name(override_raw)
|
|
except ProviderConfigError as e:
|
|
raise UploadError(str(e), 422) from None
|
|
model = effective_model(override, await get_global_default_model(session))
|
|
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
|
|
if override is not None:
|
|
# Explicit re-choice replaces the saved model; the default
|
|
# never rewrites history.
|
|
recording.transcription_model = override
|
|
await session.flush()
|
|
from shonar.services import processing as _processing
|
|
|
|
# Same audio, maybe new metadata — and a failed pipeline deserves
|
|
# another attempt. Idempotent: completed work is never redone.
|
|
await _processing.enqueue_for_recording(session, recording)
|
|
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,
|
|
transcription_model=model,
|
|
)
|
|
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()
|
|
from shonar.services import processing as _processing
|
|
|
|
# New audio on disk: queue whatever AI stages apply (none configured =
|
|
# ai_disabled, never an error).
|
|
await _processing.enqueue_for_recording(session, recording)
|
|
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
|