S.H.O.N.A.R._Desktop_Companion/backend/shonar/services/uploads.py
avi 76c867fca4 Standalone Shonar Desktop: vendor portable sources + local engine; decouple from ~/Projects/Shonar
- 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
2026-09-14 17:14:54 -05:00

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