M9: search endpoint + list filters, exports (audio/txt/md/zip), retention sweep
This commit is contained in:
parent
4944fcd162
commit
6824d9fbf1
20 changed files with 1513 additions and 12 deletions
|
|
@ -0,0 +1,67 @@
|
||||||
|
package com.shonar.recording
|
||||||
|
|
||||||
|
import java.io.File
|
||||||
|
import java.text.SimpleDateFormat
|
||||||
|
import java.util.Date
|
||||||
|
import java.util.Locale
|
||||||
|
|
||||||
|
/**
|
||||||
|
* All filename policy in one place. Storage choice: app-specific
|
||||||
|
* `filesDir` (default) or the user-picked folder (local-only provider).
|
||||||
|
* MediaStore is intentionally NOT used: recordings must survive uninstall
|
||||||
|
* control, stay off the shared music library, and never need
|
||||||
|
* READ_MEDIA_AUDIO for our own files. Import scan in
|
||||||
|
* [RecordingRepository.importExistingFiles] adopts externally added files.
|
||||||
|
*/
|
||||||
|
object RecordingFiles {
|
||||||
|
const val EXT = "m4a"
|
||||||
|
const val MIME = "audio/mp4"
|
||||||
|
private const val MAX_STEM = 120
|
||||||
|
|
||||||
|
fun defaultDisplayName(nowMs: Long = System.currentTimeMillis()): String {
|
||||||
|
val s = SimpleDateFormat("yyyy-MM-dd HH-mm", Locale.US).format(Date(nowMs))
|
||||||
|
return "Shonar Recording - $s"
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Null = blank / "." / ".." / nothing usable after cleaning. */
|
||||||
|
fun sanitizeStem(raw: String): String? {
|
||||||
|
var stem = raw.trim()
|
||||||
|
.replace(Regex("[/\\\\]"), "_")
|
||||||
|
.replace(Regex("\\p{Cntrl}"), "")
|
||||||
|
.replace(Regex("\\s+"), " ").trim()
|
||||||
|
// Windows-hostile + MediaStore-hostile chars.
|
||||||
|
stem = stem.replace(Regex("[<>:\"|?*]"), "_").trim()
|
||||||
|
// No leading dots (hidden files) / trailing dots-spaces (Windows trim).
|
||||||
|
stem = stem.trim('.', ' ', '\t').trim()
|
||||||
|
if (stem.isEmpty() || stem == "." || stem == "..") return null
|
||||||
|
if (stem.length > MAX_STEM) stem = stem.take(MAX_STEM).trimEnd()
|
||||||
|
return stem.ifEmpty { null }
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Never overwrite: "Name.m4a", "Name (2).m4a", ... */
|
||||||
|
fun uniqueFinalFile(dir: File, stem: String, ext: String = EXT): File {
|
||||||
|
var target = File(dir, "$stem.$ext")
|
||||||
|
var n = 1
|
||||||
|
while (target.exists()) {
|
||||||
|
n++
|
||||||
|
target = File(dir, "$stem ($n).$ext")
|
||||||
|
}
|
||||||
|
return target
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Sweep crashed sessions: files in .in-progress older than 24h or orphaned. */
|
||||||
|
fun sweepOrphanedTempFiles(libraryRoot: File, validId: String?): Int {
|
||||||
|
val tmp = File(libraryRoot, ".in-progress")
|
||||||
|
if (!tmp.isDirectory) return 0
|
||||||
|
var removed = 0
|
||||||
|
val cutoff = System.currentTimeMillis() - 24 * 3600 * 1000L
|
||||||
|
tmp.listFiles()?.forEach { f ->
|
||||||
|
if (!f.isFile) return@forEach
|
||||||
|
val orphan = validId == null || f.nameWithoutExtension != validId
|
||||||
|
if (orphan && f.lastModified() < cutoff) {
|
||||||
|
if (f.delete()) removed++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return removed
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,112 @@
|
||||||
|
package com.shonar.recording
|
||||||
|
|
||||||
|
import android.app.Notification
|
||||||
|
import android.app.NotificationChannel
|
||||||
|
import android.app.NotificationManager
|
||||||
|
import android.app.PendingIntent
|
||||||
|
import android.content.Context
|
||||||
|
import android.content.Intent
|
||||||
|
import android.os.Build
|
||||||
|
import androidx.core.app.NotificationCompat
|
||||||
|
import com.shonar.MainActivity
|
||||||
|
import com.shonar.R
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Branding: all colors live in Theme.kt (Teal #4FD1C5 / DeepNavy #0B1220).
|
||||||
|
* Notification uses the app icon + [NotificationCompat] accent defaults so
|
||||||
|
* it follows the system theme. Replace R.drawable.* with Shonar's final
|
||||||
|
* logo assets — only this file + widget layout reference them.
|
||||||
|
*
|
||||||
|
* Icons: system drawables (android.R.drawable.*) are used for actions so
|
||||||
|
* no new assets are required. Swap with branded vectors later:
|
||||||
|
* ic_cancel, ic_pause, ic_play, ic_delete, ic_save_check.
|
||||||
|
*/
|
||||||
|
object RecordingNotificationHelper {
|
||||||
|
const val CHANNEL_ID = "shonar_recording"
|
||||||
|
const val NOTIFICATION_ID = 1001
|
||||||
|
|
||||||
|
fun ensureChannel(ctx: Context) {
|
||||||
|
if (Build.VERSION.SDK_INT < Build.VERSION_CODES.O) return
|
||||||
|
val mgr = ctx.getSystemService(NotificationManager::class.java)
|
||||||
|
mgr.createNotificationChannel(
|
||||||
|
NotificationChannel(
|
||||||
|
CHANNEL_ID,
|
||||||
|
ctx.getString(R.string.notif_channel_recording),
|
||||||
|
NotificationManager.IMPORTANCE_LOW, // no sound/vibration while recording
|
||||||
|
).apply { description = ctx.getString(R.string.notif_channel_recording_desc) },
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun actionIntent(ctx: Context, action: String, req: Int): PendingIntent {
|
||||||
|
val i = Intent(ctx, RecordingService::class.java).setAction(action)
|
||||||
|
return PendingIntent.getService(
|
||||||
|
ctx, req, i,
|
||||||
|
PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
fun build(ctx: Context, phase: RecordingSnapshot.Phase, elapsedMs: Long): Notification {
|
||||||
|
val openApp = PendingIntent.getActivity(
|
||||||
|
ctx, 10,
|
||||||
|
Intent(ctx, MainActivity::class.java)
|
||||||
|
.setAction(Intent.ACTION_MAIN)
|
||||||
|
.addCategory(android.content.Intent.CATEGORY_LAUNCHER),
|
||||||
|
PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE,
|
||||||
|
)
|
||||||
|
val paused = phase == RecordingSnapshot.Phase.PAUSED
|
||||||
|
val title = if (paused) {
|
||||||
|
ctx.getString(R.string.notif_paused)
|
||||||
|
} else {
|
||||||
|
ctx.getString(R.string.notif_recording)
|
||||||
|
}
|
||||||
|
val builder = NotificationCompat.Builder(ctx, CHANNEL_ID)
|
||||||
|
.setSmallIcon(R.drawable.ic_stat_mic)
|
||||||
|
.setContentTitle(title)
|
||||||
|
.setContentText(formatElapsed(elapsedMs))
|
||||||
|
.setSubText("SHONAR")
|
||||||
|
.setOngoing(true)
|
||||||
|
.setOnlyAlertOnce(true)
|
||||||
|
.setCategory(NotificationCompat.CATEGORY_STATUS)
|
||||||
|
.setVisibility(NotificationCompat.VISIBILITY_PUBLIC)
|
||||||
|
.setContentIntent(openApp)
|
||||||
|
.setShowWhen(false)
|
||||||
|
.setUsesChronometer(false)
|
||||||
|
// Chronometer-style elapsed is rendered as text so Glance/widget
|
||||||
|
// and notification stay in sync without re-posting every second.
|
||||||
|
.addAction(
|
||||||
|
android.R.drawable.ic_menu_close_clear_cancel, "Cancel",
|
||||||
|
actionIntent(ctx, RecordingService.ACTION_CANCEL, 21),
|
||||||
|
)
|
||||||
|
.addAction(
|
||||||
|
if (paused) android.R.drawable.ic_media_play else android.R.drawable.ic_media_pause,
|
||||||
|
if (paused) "Resume" else "Pause",
|
||||||
|
actionIntent(
|
||||||
|
ctx,
|
||||||
|
if (paused) RecordingService.ACTION_RESUME else RecordingService.ACTION_PAUSE,
|
||||||
|
22,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.addAction(
|
||||||
|
android.R.drawable.ic_menu_delete, "Delete",
|
||||||
|
actionIntent(ctx, RecordingService.ACTION_DELETE, 23),
|
||||||
|
)
|
||||||
|
.addAction(
|
||||||
|
android.R.drawable.checkbox_on_background, "Save",
|
||||||
|
actionIntent(ctx, RecordingService.ACTION_SAVE, 24),
|
||||||
|
)
|
||||||
|
// MediaStyle keeps Save (checkmark) prominent on all OEM skins.
|
||||||
|
builder.setStyle(
|
||||||
|
androidx.media.app.NotificationCompat.MediaStyle()
|
||||||
|
.setShowActionsInCompactView(1, 3),
|
||||||
|
)
|
||||||
|
return builder.build()
|
||||||
|
}
|
||||||
|
|
||||||
|
fun formatElapsed(ms: Long): String {
|
||||||
|
val s = (ms / 1000).coerceAtLeast(0)
|
||||||
|
val h = s / 3600
|
||||||
|
val m = (s % 3600) / 60
|
||||||
|
val sec = s % 60
|
||||||
|
return if (h > 0) "%d:%02d:%02d".format(h, m, sec) else "%02d:%02d".format(m, sec)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,41 @@
|
||||||
|
package com.shonar.recording
|
||||||
|
|
||||||
|
import android.Manifest
|
||||||
|
import android.content.Context
|
||||||
|
import android.content.pm.PackageManager
|
||||||
|
import android.os.Build
|
||||||
|
import androidx.core.content.ContextCompat
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Version matrix (12/13/14+):
|
||||||
|
* - RECORD_AUDIO: runtime on all versions, required before MediaRecorder.
|
||||||
|
* - POST_NOTIFICATIONS (API 33+): runtime; without it the foreground
|
||||||
|
* service + widget still work, but the notification is suppressed by
|
||||||
|
* the system. We keep recording + show in-app/widget state.
|
||||||
|
* - FOREGROUND_SERVICE (API 28+ manifest) + FOREGROUND_SERVICE_MICROPHONE
|
||||||
|
* (API 30+ manifest, enforced 34+): no runtime grant; declare in
|
||||||
|
* manifest + pass microphone FGS type at startForeground().
|
||||||
|
* - Background start (API 31+): a widget tap (system-bound PendingIntent)
|
||||||
|
* is an exempted FGS start; an in-app tap uses startForegroundService.
|
||||||
|
* - Microphone-in-use indicator (API 29+ green dot) is system-owned.
|
||||||
|
*/
|
||||||
|
object RecordingPermissionHelper {
|
||||||
|
fun hasRecordAudio(c: Context): Boolean =
|
||||||
|
ContextCompat.checkSelfPermission(c, Manifest.permission.RECORD_AUDIO) ==
|
||||||
|
PackageManager.PERMISSION_GRANTED
|
||||||
|
|
||||||
|
fun hasPostNotifications(c: Context): Boolean {
|
||||||
|
if (Build.VERSION.SDK_INT < 33) return true
|
||||||
|
return ContextCompat.checkSelfPermission(c, Manifest.permission.POST_NOTIFICATIONS) ==
|
||||||
|
PackageManager.PERMISSION_GRANTED
|
||||||
|
}
|
||||||
|
|
||||||
|
fun requiredMissing(c: Context): List<String> {
|
||||||
|
val out = mutableListOf<String>()
|
||||||
|
if (!hasRecordAudio(c)) out += Manifest.permission.RECORD_AUDIO
|
||||||
|
if (Build.VERSION.SDK_INT >= 33 && !hasPostNotifications(c)) {
|
||||||
|
out += Manifest.permission.POST_NOTIFICATIONS
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,42 @@
|
||||||
|
package com.shonar.recording
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Single source of truth for the in-progress recording session.
|
||||||
|
* Persists across process death via [RecordingStateStore]; mirrored to
|
||||||
|
* the widget + notification.
|
||||||
|
*/
|
||||||
|
enum class RecordingPhase { IDLE, RECORDING, PAUSED }
|
||||||
|
|
||||||
|
sealed interface RecordingCommand {
|
||||||
|
data object Start : RecordingCommand
|
||||||
|
data object Pause : RecordingCommand
|
||||||
|
data object Resume : RecordingCommand
|
||||||
|
/** Stop, finalize temp file -> Room, then open rename. */
|
||||||
|
data object Save : RecordingCommand
|
||||||
|
/** Stop + delete temp file, no Room row. */
|
||||||
|
data object Cancel : RecordingCommand
|
||||||
|
/** Delete = Cancel alias kept for notification UX (per spec). */
|
||||||
|
data object Delete : RecordingCommand
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Snapshot observed by UI / widget. Kept compatible with existing code. */
|
||||||
|
data class RecordingSnapshot(
|
||||||
|
val phase: Phase = Phase.IDLE,
|
||||||
|
val recordingId: String? = null,
|
||||||
|
val elapsedMs: Long = 0L,
|
||||||
|
) {
|
||||||
|
// Legacy alias so existing HomeScreen code keeps compiling.
|
||||||
|
enum class Phase { IDLE, RECORDING, PAUSED }
|
||||||
|
|
||||||
|
fun toPhase(): RecordingPhase = when (phase) {
|
||||||
|
Phase.IDLE -> RecordingPhase.IDLE
|
||||||
|
Phase.RECORDING -> RecordingPhase.RECORDING
|
||||||
|
Phase.PAUSED -> RecordingPhase.PAUSED
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun RecordingPhase.toSnapshotPhase(): RecordingSnapshot.Phase = when (this) {
|
||||||
|
RecordingPhase.IDLE -> RecordingSnapshot.Phase.IDLE
|
||||||
|
RecordingPhase.RECORDING -> RecordingSnapshot.Phase.RECORDING
|
||||||
|
RecordingPhase.PAUSED -> RecordingSnapshot.Phase.PAUSED
|
||||||
|
}
|
||||||
|
|
@ -0,0 +1,68 @@
|
||||||
|
package com.shonar.recording
|
||||||
|
|
||||||
|
import android.content.Context
|
||||||
|
import androidx.datastore.preferences.core.edit
|
||||||
|
import androidx.datastore.preferences.core.longPreferencesKey
|
||||||
|
import androidx.datastore.preferences.core.stringPreferencesKey
|
||||||
|
import androidx.datastore.preferences.preferencesDataStore
|
||||||
|
import kotlinx.coroutines.flow.first
|
||||||
|
import kotlinx.coroutines.flow.map
|
||||||
|
|
||||||
|
private val Context.sessionStore by preferencesDataStore("recording_session")
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Crash-safe session journal. The service writes before/after every
|
||||||
|
* transition; on reboot / process death we either resume UI state or
|
||||||
|
* sweep an orphaned temp file. Small, synchronous, no Room needed here
|
||||||
|
* because this must work even when the DB is closed.
|
||||||
|
*/
|
||||||
|
class RecordingStateStore(private val context: Context) {
|
||||||
|
private val ID = stringPreferencesKey("session_id")
|
||||||
|
private val PHASE = stringPreferencesKey("session_phase")
|
||||||
|
private val ACCUM = longPreferencesKey("session_accum_ms")
|
||||||
|
private val STARTED_AT = longPreferencesKey("session_started_at_ms")
|
||||||
|
private val CREATED_AT = longPreferencesKey("session_created_at_ms")
|
||||||
|
|
||||||
|
data class Persisted(
|
||||||
|
val id: String,
|
||||||
|
val phase: RecordingPhase,
|
||||||
|
val accumulatedMs: Long,
|
||||||
|
val resumedAtElapsedMs: Long,
|
||||||
|
val createdAtMs: Long,
|
||||||
|
)
|
||||||
|
|
||||||
|
suspend fun save(
|
||||||
|
id: String,
|
||||||
|
phase: RecordingPhase,
|
||||||
|
accumulatedMs: Long,
|
||||||
|
resumedAtElapsedMs: Long,
|
||||||
|
createdAtMs: Long,
|
||||||
|
) {
|
||||||
|
context.sessionStore.edit { p ->
|
||||||
|
p[ID] = id
|
||||||
|
p[PHASE] = phase.name
|
||||||
|
p[ACCUM] = accumulatedMs
|
||||||
|
p[STARTED_AT] = resumedAtElapsedMs
|
||||||
|
p[CREATED_AT] = createdAtMs
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
suspend fun load(): Persisted? {
|
||||||
|
val p = context.sessionStore.data.map { it }.first()
|
||||||
|
val id = p[ID] ?: return null
|
||||||
|
val phase = runCatching { RecordingPhase.valueOf(p[PHASE] ?: "IDLE") }
|
||||||
|
.getOrDefault(RecordingPhase.IDLE)
|
||||||
|
if (phase == RecordingPhase.IDLE) return null
|
||||||
|
return Persisted(
|
||||||
|
id = id,
|
||||||
|
phase = phase,
|
||||||
|
accumulatedMs = p[ACCUM] ?: 0L,
|
||||||
|
resumedAtElapsedMs = p[STARTED_AT] ?: 0L,
|
||||||
|
createdAtMs = p[CREATED_AT] ?: System.currentTimeMillis(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
suspend fun clear() {
|
||||||
|
context.sessionStore.edit { it.clear() }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -55,6 +55,10 @@ exclude = ["migrations"]
|
||||||
[tool.ruff.lint]
|
[tool.ruff.lint]
|
||||||
select = ["E", "F", "I", "UP", "B", "SIM"]
|
select = ["E", "F", "I", "UP", "B", "SIM"]
|
||||||
|
|
||||||
|
[tool.ruff.lint.per-file-ignores]
|
||||||
|
# FastAPI idiom: Query(...)/Depends() as parameter defaults.
|
||||||
|
"shonar/api/**" = ["B008"]
|
||||||
|
|
||||||
[tool.mypy]
|
[tool.mypy]
|
||||||
python_version = "3.11"
|
python_version = "3.11"
|
||||||
ignore_missing_imports = true
|
ignore_missing_imports = true
|
||||||
|
|
|
||||||
23
backend/shonar/api/schemas_search.py
Normal file
23
backend/shonar/api/schemas_search.py
Normal file
|
|
@ -0,0 +1,23 @@
|
||||||
|
"""Search + export schemas (M9)."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
|
||||||
|
class SearchHitOut(BaseModel):
|
||||||
|
id: uuid.UUID
|
||||||
|
title: str
|
||||||
|
field: str # title | tag | notes | summary | transcript
|
||||||
|
snippet: str
|
||||||
|
|
||||||
|
|
||||||
|
class SearchOut(BaseModel):
|
||||||
|
query: str
|
||||||
|
scope: str
|
||||||
|
items: list[SearchHitOut]
|
||||||
|
total: int
|
||||||
|
limit: int
|
||||||
|
offset: int
|
||||||
|
|
@ -2,7 +2,15 @@
|
||||||
|
|
||||||
from fastapi import APIRouter
|
from fastapi import APIRouter
|
||||||
|
|
||||||
from shonar.api.v1 import auth, health, models, provider_info, recordings, users
|
from shonar.api.v1 import (
|
||||||
|
auth,
|
||||||
|
health,
|
||||||
|
models,
|
||||||
|
provider_info,
|
||||||
|
recordings,
|
||||||
|
search_exports,
|
||||||
|
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)
|
||||||
|
|
@ -11,6 +19,7 @@ api_router.include_router(users.router)
|
||||||
api_router.include_router(recordings.router)
|
api_router.include_router(recordings.router)
|
||||||
api_router.include_router(models.router)
|
api_router.include_router(models.router)
|
||||||
api_router.include_router(provider_info.router)
|
api_router.include_router(provider_info.router)
|
||||||
|
api_router.include_router(search_exports.router)
|
||||||
|
|
||||||
# Included as later milestones land:
|
# Included as later milestones land:
|
||||||
# - tags, search, exports (M9)
|
# - tags (M9 follow-on)
|
||||||
|
|
|
||||||
|
|
@ -8,6 +8,7 @@ from __future__ import annotations
|
||||||
|
|
||||||
import contextlib
|
import contextlib
|
||||||
import uuid
|
import uuid
|
||||||
|
from datetime import datetime
|
||||||
|
|
||||||
from fastapi import APIRouter, Header, HTTPException, Query, Request, Response
|
from fastapi import APIRouter, Header, HTTPException, Query, Request, Response
|
||||||
from sqlalchemy import func, select
|
from sqlalchemy import func, select
|
||||||
|
|
@ -157,8 +158,35 @@ async def list_recordings(
|
||||||
pattern="^(recorded_at|created_at|duration_seconds|title)$",
|
pattern="^(recorded_at|created_at|duration_seconds|title)$",
|
||||||
),
|
),
|
||||||
order: str = Query(default="desc", pattern="^(asc|desc)$"),
|
order: str = Query(default="desc", pattern="^(asc|desc)$"),
|
||||||
|
# M9 filters (AND-combined).
|
||||||
|
tag: str | None = Query(default=None, max_length=80),
|
||||||
|
status: str | None = Query(default=None, max_length=20),
|
||||||
|
from_date: datetime | None = Query(default=None, description="recorded_at >= (UTC)"),
|
||||||
|
to_date: datetime | None = Query(default=None, description="recorded_at <= (UTC)"),
|
||||||
):
|
):
|
||||||
where = [Recording.user_id == user.id, Recording.deleted_at.is_(None)]
|
where = [Recording.user_id == user.id, Recording.deleted_at.is_(None)]
|
||||||
|
if tag:
|
||||||
|
from shonar.db.models import RecordingTag as RT
|
||||||
|
from shonar.db.models import Tag as T
|
||||||
|
|
||||||
|
where.append(
|
||||||
|
Recording.id.in_(
|
||||||
|
select(RT.recording_id)
|
||||||
|
.join(T, T.id == RT.tag_id)
|
||||||
|
.where(T.user_id == user.id, T.name == tag.strip().lower())
|
||||||
|
)
|
||||||
|
)
|
||||||
|
if status:
|
||||||
|
from shonar.db.models import ProcessingStatus
|
||||||
|
|
||||||
|
try:
|
||||||
|
where.append(Recording.processing_status == ProcessingStatus(status))
|
||||||
|
except ValueError:
|
||||||
|
raise HTTPException(422, f"Unknown status: {status}") from None
|
||||||
|
if from_date is not None:
|
||||||
|
where.append(Recording.recorded_at >= from_date)
|
||||||
|
if to_date is not None:
|
||||||
|
where.append(Recording.recorded_at <= to_date)
|
||||||
total = await session.scalar(select(func.count(Recording.id)).where(*where))
|
total = await session.scalar(select(func.count(Recording.id)).where(*where))
|
||||||
col = getattr(Recording, sort)
|
col = getattr(Recording, sort)
|
||||||
col = col.desc() if order == "desc" else col.asc()
|
col = col.desc() if order == "desc" else col.asc()
|
||||||
|
|
|
||||||
81
backend/shonar/api/v1/search_exports.py
Normal file
81
backend/shonar/api/v1/search_exports.py
Normal file
|
|
@ -0,0 +1,81 @@
|
||||||
|
"""Search + exports endpoints (M9).
|
||||||
|
|
||||||
|
Ownership is enforced everywhere: a search only ever sees the caller's own
|
||||||
|
non-deleted recordings, and an export 404s on anything the caller doesn't
|
||||||
|
own (no existence leak).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
|
||||||
|
from fastapi import APIRouter, HTTPException, Query
|
||||||
|
from fastapi.responses import Response
|
||||||
|
|
||||||
|
from shonar.api.deps import CurrentUser, SessionDep
|
||||||
|
from shonar.api.schemas_search import SearchHitOut, SearchOut
|
||||||
|
from shonar.db.models import Recording
|
||||||
|
from shonar.services import exports as ex
|
||||||
|
from shonar.services import search as sr
|
||||||
|
|
||||||
|
router = APIRouter(tags=["search", "exports"])
|
||||||
|
|
||||||
|
|
||||||
|
# --- search -------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/search", response_model=SearchOut)
|
||||||
|
async def search(
|
||||||
|
user: CurrentUser,
|
||||||
|
session: SessionDep,
|
||||||
|
q: str = Query(..., min_length=1, max_length=200, description="Search text"),
|
||||||
|
scope: str = Query(
|
||||||
|
default="all",
|
||||||
|
pattern="^(all|title|notes|transcript|summary|tag)$",
|
||||||
|
description="Restrict the search to one field (default: all)",
|
||||||
|
),
|
||||||
|
limit: int = Query(default=20, ge=1, le=100),
|
||||||
|
offset: int = Query(default=0, ge=0),
|
||||||
|
):
|
||||||
|
hits, total = await sr.search_recordings(
|
||||||
|
session, user.id, q, scope=scope, limit=limit, offset=offset
|
||||||
|
)
|
||||||
|
items = [
|
||||||
|
SearchHitOut(
|
||||||
|
id=h.recording.id,
|
||||||
|
title=h.recording.title,
|
||||||
|
field=h.field,
|
||||||
|
snippet=h.snippet,
|
||||||
|
)
|
||||||
|
for h in hits
|
||||||
|
]
|
||||||
|
return SearchOut(query=q, scope=scope, items=items, total=total, limit=limit, offset=offset)
|
||||||
|
|
||||||
|
|
||||||
|
# --- exports ------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/recordings/{recording_id}/export")
|
||||||
|
async def export_recording(
|
||||||
|
recording_id: uuid.UUID,
|
||||||
|
user: CurrentUser,
|
||||||
|
session: SessionDep,
|
||||||
|
fmt: str = Query(default="zip", pattern="^(audio|txt|md|zip)$"),
|
||||||
|
):
|
||||||
|
"""Return one export artifact inline (audio / txt / md / zip bundle)."""
|
||||||
|
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")
|
||||||
|
try:
|
||||||
|
result = await ex.build_export(session, rec, fmt)
|
||||||
|
except ex.ExportError as e:
|
||||||
|
raise HTTPException(e.status_code, e.message) from None
|
||||||
|
await ex.record_export(session, user.id, rec, fmt, result)
|
||||||
|
return Response(
|
||||||
|
content=result.data,
|
||||||
|
media_type=result.mime_type,
|
||||||
|
headers={
|
||||||
|
"Content-Disposition": f'attachment; filename="{result.filename}"',
|
||||||
|
"Cache-Control": "private, no-store",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
@ -93,6 +93,12 @@ class Settings(BaseSettings):
|
||||||
# --- Search ---------------------------------------------------------
|
# --- Search ---------------------------------------------------------
|
||||||
search_backend: str = "postgres_fts" # postgres_fts (meilisearch: TODO)
|
search_backend: str = "postgres_fts" # postgres_fts (meilisearch: TODO)
|
||||||
|
|
||||||
|
# --- Retention --------------------------------------------------------
|
||||||
|
# Grace window before the sweep hard-deletes soft-deleted recordings
|
||||||
|
# and accounts (rows + stored files). Cancellations are valid until
|
||||||
|
# the sweep fires.
|
||||||
|
retention_grace_days: int = 30
|
||||||
|
|
||||||
# --- Misc -------------------------------------------------------------
|
# --- Misc -------------------------------------------------------------
|
||||||
rate_limit_auth: str = "10/minute"
|
rate_limit_auth: str = "10/minute"
|
||||||
rate_limit_default: str = "120/minute"
|
rate_limit_default: str = "120/minute"
|
||||||
|
|
|
||||||
|
|
@ -162,6 +162,6 @@ async def delete_account(session: AsyncSession, user: User, password: str) -> No
|
||||||
.where(RefreshToken.user_id == user.id, RefreshToken.revoked_at.is_(None))
|
.where(RefreshToken.user_id == user.id, RefreshToken.revoked_at.is_(None))
|
||||||
.values(revoked_at=now)
|
.values(revoked_at=now)
|
||||||
)
|
)
|
||||||
# NOTE: hard deletion of rows/files is performed by a retention sweep so
|
# NOTE: hard deletion of rows/files is performed by the retention sweep
|
||||||
# an accidental deletion can be cancelled within the grace window (see
|
# (services/retention.py) so an accidental deletion can be cancelled
|
||||||
# docs/security.md). TODO: scheduled purge job (30-day grace).
|
# within the grace window (see docs/security.md).
|
||||||
|
|
|
||||||
251
backend/shonar/services/exports.py
Normal file
251
backend/shonar/services/exports.py
Normal file
|
|
@ -0,0 +1,251 @@
|
||||||
|
"""Exports (M9): audio, transcript txt, notes markdown, bundle zip.
|
||||||
|
|
||||||
|
Synchronous generation — every artifact is small enough (text, or one
|
||||||
|
audio file) that a background job adds failure modes, not speed. Each
|
||||||
|
successful export records an ExportJob row and stores the produced bytes
|
||||||
|
as an Asset(kind=export) so the audit trail exists; the response is the
|
||||||
|
file itself (no separate download-asset round trip).
|
||||||
|
|
||||||
|
Formats:
|
||||||
|
audio — the original upload, byte-identical, original mime/extension
|
||||||
|
txt — current transcript text
|
||||||
|
md — notes.md: title, metadata, notes, summary sections, transcript
|
||||||
|
zip — bundle: original audio + transcript.txt + notes.md
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import io
|
||||||
|
import uuid
|
||||||
|
import zipfile
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
from sqlalchemy import select
|
||||||
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
|
from shonar.db.models import (
|
||||||
|
Asset,
|
||||||
|
AssetKind,
|
||||||
|
ExportJob,
|
||||||
|
JobStatus,
|
||||||
|
Recording,
|
||||||
|
Summary,
|
||||||
|
Transcript,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class ExportError(Exception):
|
||||||
|
def __init__(self, status_code: int, message: str):
|
||||||
|
super().__init__(message)
|
||||||
|
self.status_code = status_code
|
||||||
|
self.message = message
|
||||||
|
|
||||||
|
|
||||||
|
EXPORT_FORMATS = ("audio", "txt", "md", "zip")
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class ExportResult:
|
||||||
|
filename: str
|
||||||
|
mime_type: str
|
||||||
|
data: bytes
|
||||||
|
|
||||||
|
|
||||||
|
def _safe_filename(rec: Recording) -> str:
|
||||||
|
"""Slug the title; fall back to the recorded_at stamp. Never leaks ids."""
|
||||||
|
base = "".join(
|
||||||
|
c if (c.isalnum() or c in "-_ ") else " " for c in (rec.title or "")
|
||||||
|
).strip()
|
||||||
|
if not base:
|
||||||
|
base = f"recording-{rec.recorded_at:%Y%m%d-%H%M%S}"
|
||||||
|
return base[:120]
|
||||||
|
|
||||||
|
|
||||||
|
async def _current_transcript(session: AsyncSession, rec_id: uuid.UUID) -> Transcript | None:
|
||||||
|
return await session.scalar(
|
||||||
|
select(Transcript)
|
||||||
|
.where(Transcript.recording_id == rec_id, Transcript.superseded_at.is_(None))
|
||||||
|
.order_by(Transcript.version.desc())
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def _current_summary(session: AsyncSession, rec_id: uuid.UUID) -> Summary | None:
|
||||||
|
return await session.scalar(
|
||||||
|
select(Summary)
|
||||||
|
.where(Summary.recording_id == rec_id, Summary.superseded_at.is_(None))
|
||||||
|
.order_by(Summary.version.desc())
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def _original_asset(session: AsyncSession, rec_id: uuid.UUID) -> Asset | None:
|
||||||
|
return await session.scalar(
|
||||||
|
select(Asset).where(Asset.recording_id == rec_id, Asset.kind == AssetKind.original)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _summary_md(summary: Summary | None) -> str:
|
||||||
|
if summary is None:
|
||||||
|
return ""
|
||||||
|
c = summary.content or {}
|
||||||
|
out = ["## Summary\n"]
|
||||||
|
if c.get("short"):
|
||||||
|
out.append(f"{c['short']}\n")
|
||||||
|
if c.get("detailed"):
|
||||||
|
out.append(f"### Detailed\n\n{c['detailed']}\n")
|
||||||
|
for key, header in (
|
||||||
|
("key_points", "Key Points"),
|
||||||
|
("decisions", "Decisions"),
|
||||||
|
("action_items", "Action Items"),
|
||||||
|
("questions", "Questions"),
|
||||||
|
):
|
||||||
|
items = c.get(key)
|
||||||
|
if isinstance(items, list) and items:
|
||||||
|
out.append(f"### {header}\n")
|
||||||
|
out.extend(f"- {x}" for x in items)
|
||||||
|
out.append("")
|
||||||
|
return "\n".join(out)
|
||||||
|
|
||||||
|
|
||||||
|
def _transcript_md(t: Transcript | None) -> str:
|
||||||
|
if t is None:
|
||||||
|
return ""
|
||||||
|
lines = ["## Transcript\n"]
|
||||||
|
segs = [s for s in (t.segments or []) if isinstance(s, dict)]
|
||||||
|
if segs:
|
||||||
|
for s in segs:
|
||||||
|
start = float(s.get("start", 0.0))
|
||||||
|
stamp = f"{int(start // 60):02d}:{start % 60:04.1f}"
|
||||||
|
speaker = f"**{s['speaker']}**: " if s.get("speaker") else ""
|
||||||
|
lines.append(f"- `[{stamp}]` {speaker}{s.get('text', '').strip()}")
|
||||||
|
else:
|
||||||
|
lines.append(t.text or "")
|
||||||
|
return "\n".join(lines) + "\n"
|
||||||
|
|
||||||
|
|
||||||
|
def _notes_md(
|
||||||
|
rec: Recording, t: Transcript | None, summary: Summary | None,
|
||||||
|
tag_names: list[str] | None = None,
|
||||||
|
) -> str:
|
||||||
|
parts = [
|
||||||
|
f"# {rec.title}\n",
|
||||||
|
f"- Recorded: {rec.recorded_at:%Y-%m-%d %H:%M} UTC",
|
||||||
|
f"- Duration: {rec.duration_seconds:.1f}s",
|
||||||
|
]
|
||||||
|
if tag_names:
|
||||||
|
parts.append("- Tags: " + ", ".join(f"`{x}`" for x in tag_names))
|
||||||
|
parts.append("")
|
||||||
|
if rec.notes:
|
||||||
|
parts.append(f"## Notes\n\n{rec.notes}\n")
|
||||||
|
s = _summary_md(summary)
|
||||||
|
if s:
|
||||||
|
parts.append(s)
|
||||||
|
tr = _transcript_md(t)
|
||||||
|
if tr:
|
||||||
|
parts.append(tr)
|
||||||
|
return "\n".join(parts)
|
||||||
|
|
||||||
|
|
||||||
|
def _zip(files: list[tuple[str, bytes]]) -> bytes:
|
||||||
|
buf = io.BytesIO()
|
||||||
|
with zipfile.ZipFile(buf, "w", zipfile.ZIP_DEFLATED) as z:
|
||||||
|
for name, data in files:
|
||||||
|
z.writestr(name, data)
|
||||||
|
return buf.getvalue()
|
||||||
|
|
||||||
|
|
||||||
|
async def build_export(
|
||||||
|
session: AsyncSession, rec: Recording, export_format: str
|
||||||
|
) -> ExportResult:
|
||||||
|
"""Build one export artifact for an owned, non-deleted recording."""
|
||||||
|
from shonar.storage import get_storage
|
||||||
|
|
||||||
|
if export_format not in EXPORT_FORMATS:
|
||||||
|
raise ExportError(422, f"Unknown export format. Use one of: {', '.join(EXPORT_FORMATS)}")
|
||||||
|
|
||||||
|
original = await _original_asset(session, rec.id)
|
||||||
|
t = await _current_transcript(session, rec.id)
|
||||||
|
summary = await _current_summary(session, rec.id)
|
||||||
|
from shonar.services.search import tag_names
|
||||||
|
|
||||||
|
tags = await tag_names(session, rec.id)
|
||||||
|
base = _safe_filename(rec)
|
||||||
|
|
||||||
|
if export_format == "audio":
|
||||||
|
if original is None:
|
||||||
|
raise ExportError(404, "No audio stored for this recording")
|
||||||
|
data = await get_storage().get(original.storage_key)
|
||||||
|
ext = original.storage_key[original.storage_key.rfind(".") :]
|
||||||
|
return ExportResult(filename=f"{base}{ext}", mime_type=original.mime_type, data=data)
|
||||||
|
|
||||||
|
if export_format == "txt":
|
||||||
|
if t is None:
|
||||||
|
raise ExportError(404, "No transcript yet — transcribe first")
|
||||||
|
return ExportResult(
|
||||||
|
filename=f"{base}.txt", mime_type="text/plain; charset=utf-8",
|
||||||
|
data=(t.text or "").encode("utf-8"),
|
||||||
|
)
|
||||||
|
|
||||||
|
if export_format == "md":
|
||||||
|
if t is None and summary is None and not rec.notes:
|
||||||
|
raise ExportError(
|
||||||
|
404, "Nothing to export — this recording has no notes, transcript, or summary"
|
||||||
|
)
|
||||||
|
return ExportResult(
|
||||||
|
filename=f"{base}.md", mime_type="text/markdown; charset=utf-8",
|
||||||
|
data=_notes_md(rec, t, summary, tags).encode("utf-8"),
|
||||||
|
)
|
||||||
|
|
||||||
|
# zip bundle: whatever exists, always at least the audio when present.
|
||||||
|
if original is None and t is None and summary is None and not rec.notes:
|
||||||
|
raise ExportError(404, "Nothing to export for this recording")
|
||||||
|
files: list[tuple[str, bytes]] = []
|
||||||
|
if original is not None:
|
||||||
|
audio = await get_storage().get(original.storage_key)
|
||||||
|
ext = original.storage_key[original.storage_key.rfind(".") :]
|
||||||
|
files.append((f"{base}{ext}", audio))
|
||||||
|
if t is not None:
|
||||||
|
files.append(("transcript.txt", (t.text or "").encode("utf-8")))
|
||||||
|
files.append(("notes.md", _notes_md(rec, t, summary, tags).encode("utf-8")))
|
||||||
|
return ExportResult(
|
||||||
|
filename=f"{base}.zip", mime_type="application/zip", data=_zip(files)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def record_export(
|
||||||
|
session: AsyncSession, user_id: uuid.UUID, rec: Recording,
|
||||||
|
export_format: str, result: ExportResult,
|
||||||
|
) -> None:
|
||||||
|
"""Persist the audit trail: ExportJob(succeeded) + Asset(kind=export).
|
||||||
|
|
||||||
|
Best-effort storage of the artifact bytes; a storage failure never
|
||||||
|
fails the download the user already received.
|
||||||
|
"""
|
||||||
|
from shonar.storage import get_storage
|
||||||
|
|
||||||
|
job = ExportJob(
|
||||||
|
user_id=user_id,
|
||||||
|
recording_id=rec.id,
|
||||||
|
export_type=export_format,
|
||||||
|
status=JobStatus.succeeded,
|
||||||
|
)
|
||||||
|
session.add(job)
|
||||||
|
try:
|
||||||
|
key = f"exports/{user_id}/{rec.id}/{export_format}-{uuid.uuid4().hex}"
|
||||||
|
await get_storage().put(key, result.data)
|
||||||
|
import hashlib
|
||||||
|
|
||||||
|
asset = Asset(
|
||||||
|
recording_id=rec.id,
|
||||||
|
user_id=user_id,
|
||||||
|
kind=AssetKind.export,
|
||||||
|
storage_key=key,
|
||||||
|
mime_type=result.mime_type,
|
||||||
|
size_bytes=len(result.data),
|
||||||
|
checksum_sha256=hashlib.sha256(result.data).hexdigest(),
|
||||||
|
)
|
||||||
|
session.add(asset)
|
||||||
|
await session.flush()
|
||||||
|
job.asset_id = asset.id
|
||||||
|
except Exception: # noqa: BLE001 — audit copy is best-effort
|
||||||
|
pass
|
||||||
|
await session.flush()
|
||||||
|
|
@ -23,9 +23,13 @@ from shonar.services import processing
|
||||||
logger = logging.getLogger("shonar.inline_queue")
|
logger = logging.getLogger("shonar.inline_queue")
|
||||||
|
|
||||||
RETRY_DELAY_SECONDS = 5.0
|
RETRY_DELAY_SECONDS = 5.0
|
||||||
|
# Desktop engine: hard-delete expired soft-deletes once the app has been up
|
||||||
|
# for a day, then daily. Delayed first run keeps startup snappy.
|
||||||
|
RETENTION_INTERVAL_SECONDS = 24 * 3600.0
|
||||||
|
|
||||||
_queue: asyncio.Queue[tuple[str, str, int]] | None = None
|
_queue: asyncio.Queue[tuple[str, str, int]] | None = None
|
||||||
_consumer: asyncio.Task | None = None
|
_consumer: asyncio.Task | None = None
|
||||||
|
_retention: asyncio.Task | None = None
|
||||||
|
|
||||||
|
|
||||||
def _get_queue() -> asyncio.Queue[tuple[str, str, int]]:
|
def _get_queue() -> asyncio.Queue[tuple[str, str, int]]:
|
||||||
|
|
@ -37,10 +41,12 @@ def _get_queue() -> asyncio.Queue[tuple[str, str, int]]:
|
||||||
|
|
||||||
async def start() -> None:
|
async def start() -> None:
|
||||||
"""Start the consumer and re-run anything the DB says is pending."""
|
"""Start the consumer and re-run anything the DB says is pending."""
|
||||||
global _consumer
|
global _consumer, _retention
|
||||||
_get_queue()
|
_get_queue()
|
||||||
if _consumer is None or _consumer.done():
|
if _consumer is None or _consumer.done():
|
||||||
_consumer = asyncio.create_task(_consume(), name="shonar-inline-queue")
|
_consumer = asyncio.create_task(_consume(), name="shonar-inline-queue")
|
||||||
|
if _retention is None or _retention.done():
|
||||||
|
_retention = asyncio.create_task(_retention_loop(), name="shonar-retention")
|
||||||
# Crash recovery: queued rows (and orphaned running rows requeued by
|
# Crash recovery: queued rows (and orphaned running rows requeued by
|
||||||
# the sweep) go back on the in-process queue.
|
# the sweep) go back on the in-process queue.
|
||||||
count = await processing.sweep_stale()
|
count = await processing.sweep_stale()
|
||||||
|
|
@ -48,13 +54,40 @@ async def start() -> None:
|
||||||
logger.info("inline queue startup sweep requeued %d jobs", count)
|
logger.info("inline queue startup sweep requeued %d jobs", count)
|
||||||
|
|
||||||
|
|
||||||
|
async def _retention_loop() -> None:
|
||||||
|
"""Daily hard-delete sweep for the desktop engine (no arq cron here).
|
||||||
|
|
||||||
|
First pass a few minutes after start (the app may only run for hours
|
||||||
|
at a time, so a full-day initial sleep could starve the sweep), then
|
||||||
|
once a day while running.
|
||||||
|
"""
|
||||||
|
from shonar.services import retention
|
||||||
|
|
||||||
|
await asyncio.sleep(120.0)
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
purged = await retention.sweep_deleted()
|
||||||
|
if purged["recordings"] or purged["users"]:
|
||||||
|
logger.info("inline retention sweep: %s", purged)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception: # noqa: BLE001 — the loop must survive any failure
|
||||||
|
logger.exception("inline retention sweep failed")
|
||||||
|
await asyncio.sleep(RETENTION_INTERVAL_SECONDS)
|
||||||
|
|
||||||
|
|
||||||
async def stop() -> None:
|
async def stop() -> None:
|
||||||
global _consumer
|
global _consumer, _retention
|
||||||
if _consumer is not None:
|
if _consumer is not None:
|
||||||
_consumer.cancel()
|
_consumer.cancel()
|
||||||
with contextlib.suppress(BaseException): # noqa: BLE001 — shutdown is best-effort
|
with contextlib.suppress(BaseException): # noqa: BLE001 — shutdown is best-effort
|
||||||
await _consumer
|
await _consumer
|
||||||
_consumer = None
|
_consumer = None
|
||||||
|
if _retention is not None:
|
||||||
|
_retention.cancel()
|
||||||
|
with contextlib.suppress(BaseException): # noqa: BLE001
|
||||||
|
await _retention
|
||||||
|
_retention = None
|
||||||
|
|
||||||
|
|
||||||
async def enqueue(job_type: JobType, recording_id: str) -> None:
|
async def enqueue(job_type: JobType, recording_id: str) -> None:
|
||||||
|
|
|
||||||
101
backend/shonar/services/retention.py
Normal file
101
backend/shonar/services/retention.py
Normal file
|
|
@ -0,0 +1,101 @@
|
||||||
|
"""Retention sweep (M9): hard-delete what passed its grace window.
|
||||||
|
|
||||||
|
Two policies, both driven by ``deleted_at``:
|
||||||
|
|
||||||
|
* **Recordings** soft-deleted more than ``retention_grace_days`` ago have
|
||||||
|
their rows removed (cascade cleans transcripts/summaries/jobs/tags) and
|
||||||
|
every stored asset file deleted best-effort.
|
||||||
|
* **Accounts** deleted more than ``retention_grace_days`` ago are hard
|
||||||
|
deleted (user cascade takes their recordings/assets/devices/tokens);
|
||||||
|
their storage files are collected the same way.
|
||||||
|
|
||||||
|
Running inside the grace window is a no-op, so an accidental delete stays
|
||||||
|
cancellable until the sweep actually fires. The sweep is idempotent and
|
||||||
|
safe to run on any cadence (worker cron + inline-queue timer).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import contextlib
|
||||||
|
import logging
|
||||||
|
from datetime import timedelta
|
||||||
|
|
||||||
|
from sqlalchemy import delete as sql_delete
|
||||||
|
from sqlalchemy import select
|
||||||
|
|
||||||
|
from shonar.core.config import get_settings
|
||||||
|
from shonar.db.models import Asset, Recording, User, utcnow
|
||||||
|
from shonar.db.session import session_factory
|
||||||
|
from shonar.storage import get_storage
|
||||||
|
|
||||||
|
logger = logging.getLogger("shonar.retention")
|
||||||
|
|
||||||
|
|
||||||
|
async def sweep_deleted(limit: int = 200) -> dict[str, int]:
|
||||||
|
"""Hard-purge expired recordings and accounts. Returns counts."""
|
||||||
|
grace = timedelta(days=get_settings().retention_grace_days)
|
||||||
|
cutoff = utcnow() - grace
|
||||||
|
purged = {"recordings": 0, "users": 0, "files": 0}
|
||||||
|
storage = get_storage()
|
||||||
|
|
||||||
|
async with session_factory()() as session:
|
||||||
|
# --- recordings (skip rows whose account is also expiring: the
|
||||||
|
# user cascade below collects their files in one pass) ---
|
||||||
|
expiring_users = select(User.id).where(
|
||||||
|
User.deleted_at.is_not(None), User.deleted_at < cutoff
|
||||||
|
)
|
||||||
|
recs = list(
|
||||||
|
await session.scalars(
|
||||||
|
select(Recording)
|
||||||
|
.where(
|
||||||
|
Recording.deleted_at.is_not(None),
|
||||||
|
Recording.deleted_at < cutoff,
|
||||||
|
Recording.user_id.notin_(expiring_users),
|
||||||
|
)
|
||||||
|
.limit(limit)
|
||||||
|
)
|
||||||
|
)
|
||||||
|
for rec in recs:
|
||||||
|
assets = list(
|
||||||
|
await session.scalars(select(Asset).where(Asset.recording_id == rec.id))
|
||||||
|
)
|
||||||
|
await session.delete(rec)
|
||||||
|
await session.flush()
|
||||||
|
for a in assets:
|
||||||
|
with contextlib.suppress(Exception): # best effort; DB row is gone
|
||||||
|
await storage.delete(a.storage_key)
|
||||||
|
purged["files"] += 1
|
||||||
|
purged["recordings"] += 1
|
||||||
|
|
||||||
|
# --- accounts ---
|
||||||
|
users = list(
|
||||||
|
await session.scalars(
|
||||||
|
select(User)
|
||||||
|
.where(User.deleted_at.is_not(None), User.deleted_at < cutoff)
|
||||||
|
.limit(limit)
|
||||||
|
)
|
||||||
|
)
|
||||||
|
for user in users:
|
||||||
|
assets = list(
|
||||||
|
await session.scalars(select(Asset).where(Asset.user_id == user.id))
|
||||||
|
)
|
||||||
|
keys = [a.storage_key for a in assets]
|
||||||
|
# DB-level delete: the ORM would null the NOT NULL FKs of the
|
||||||
|
# user's recordings before the ON DELETE CASCADE could fire.
|
||||||
|
await session.execute(sql_delete(User).where(User.id == user.id))
|
||||||
|
await session.flush()
|
||||||
|
for key in keys:
|
||||||
|
with contextlib.suppress(Exception):
|
||||||
|
await storage.delete(key)
|
||||||
|
purged["files"] += 1
|
||||||
|
purged["users"] += 1
|
||||||
|
|
||||||
|
await session.commit()
|
||||||
|
|
||||||
|
if purged["recordings"] or purged["users"]:
|
||||||
|
logger.info(
|
||||||
|
"retention sweep purged %d recordings, %d accounts, %d files (grace %dd)",
|
||||||
|
purged["recordings"], purged["users"], purged["files"],
|
||||||
|
get_settings().retention_grace_days,
|
||||||
|
)
|
||||||
|
return purged
|
||||||
292
backend/shonar/services/search.py
Normal file
292
backend/shonar/services/search.py
Normal file
|
|
@ -0,0 +1,292 @@
|
||||||
|
"""Full-text search across recordings (M9).
|
||||||
|
|
||||||
|
Postgres uses the tsvector columns from migration ``fts0000000001``
|
||||||
|
(title/notes, transcript text, summary content, tag names). SQLite (the
|
||||||
|
desktop bundled-lite engine) falls back to a substring scan; libraries
|
||||||
|
there are single-user and small, and the endpoint contract is identical.
|
||||||
|
|
||||||
|
This module is the whole SearchBackend seam — a Meilisearch/OpenSearch
|
||||||
|
implementation would replace it, not the callers.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
from sqlalchemy import select, text
|
||||||
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
|
from shonar.db.models import Recording, RecordingTag, Summary, Tag, Transcript
|
||||||
|
|
||||||
|
# ``scope`` values accepted by the endpoint.
|
||||||
|
SCOPES = ("all", "title", "notes", "transcript", "summary", "tag")
|
||||||
|
|
||||||
|
# Hit ordering: a title hit outranks a transcript hit.
|
||||||
|
_FIELD_RANK = {"title": 0, "tag": 1, "notes": 2, "summary": 3, "transcript": 4}
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class SearchHit:
|
||||||
|
recording: Recording
|
||||||
|
field: str # where the best match landed
|
||||||
|
snippet: str
|
||||||
|
|
||||||
|
|
||||||
|
def _snippet_around(value: str, at: int, width: int = 160) -> str:
|
||||||
|
"""~``width`` chars centred on ``at``, word-bounded, with ellipses."""
|
||||||
|
half = width // 2
|
||||||
|
start = max(0, at - half)
|
||||||
|
end = min(len(value), at + half)
|
||||||
|
if start > 0:
|
||||||
|
start = value.rfind(" ", 0, start) + 1 or start
|
||||||
|
if end < len(value):
|
||||||
|
nxt = value.find(" ", end)
|
||||||
|
end = nxt if nxt != -1 else end
|
||||||
|
prefix = "…" if start > 0 else ""
|
||||||
|
suffix = "…" if end < len(value) else ""
|
||||||
|
return f"{prefix}{value[start:end].strip()}{suffix}"
|
||||||
|
|
||||||
|
|
||||||
|
async def tag_names(session: AsyncSession, recording_id: uuid.UUID) -> list[str]:
|
||||||
|
return list(
|
||||||
|
await session.scalars(
|
||||||
|
select(Tag.name)
|
||||||
|
.join(RecordingTag, RecordingTag.tag_id == Tag.id)
|
||||||
|
.where(RecordingTag.recording_id == recording_id)
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def _current_transcript_text(session: AsyncSession, recording_id: uuid.UUID) -> str:
|
||||||
|
row = await session.scalar(
|
||||||
|
select(Transcript.text)
|
||||||
|
.where(Transcript.recording_id == recording_id, Transcript.superseded_at.is_(None))
|
||||||
|
.order_by(Transcript.version.desc())
|
||||||
|
)
|
||||||
|
return row or ""
|
||||||
|
|
||||||
|
|
||||||
|
def _summary_flat(content: dict | None) -> str:
|
||||||
|
if not content:
|
||||||
|
return ""
|
||||||
|
parts = [str(content.get("short", "")), str(content.get("detailed", ""))]
|
||||||
|
for key in ("key_points", "decisions", "action_items", "questions"):
|
||||||
|
v = content.get(key)
|
||||||
|
if isinstance(v, list):
|
||||||
|
parts.extend(str(x) for x in v)
|
||||||
|
return " ".join(p for p in parts if p)
|
||||||
|
|
||||||
|
|
||||||
|
# --- SQLite fallback ----------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def _search_sqlite(
|
||||||
|
session: AsyncSession, user_id: uuid.UUID, q: str, scope: str,
|
||||||
|
limit: int, offset: int,
|
||||||
|
) -> tuple[list[SearchHit], int]:
|
||||||
|
needles = [t for t in q.lower().split() if t]
|
||||||
|
if not needles:
|
||||||
|
return [], 0
|
||||||
|
recs = list(
|
||||||
|
await session.scalars(
|
||||||
|
select(Recording)
|
||||||
|
.where(Recording.user_id == user_id, Recording.deleted_at.is_(None))
|
||||||
|
.order_by(Recording.recorded_at.desc())
|
||||||
|
)
|
||||||
|
)
|
||||||
|
hits: list[SearchHit] = []
|
||||||
|
for rec in recs:
|
||||||
|
fields: list[tuple[str, str]] = []
|
||||||
|
if scope in ("all", "title"):
|
||||||
|
fields.append(("title", rec.title or ""))
|
||||||
|
if scope in ("all", "notes"):
|
||||||
|
fields.append(("notes", rec.notes or ""))
|
||||||
|
if scope in ("all", "tag"):
|
||||||
|
fields.append(("tag", " ".join(await tag_names(session, rec.id))))
|
||||||
|
if scope in ("all", "transcript"):
|
||||||
|
fields.append(("transcript", await _current_transcript_text(session, rec.id)))
|
||||||
|
if scope in ("all", "summary"):
|
||||||
|
content = await session.scalar(
|
||||||
|
select(Summary.content)
|
||||||
|
.where(Summary.recording_id == rec.id, Summary.superseded_at.is_(None))
|
||||||
|
.order_by(Summary.version.desc())
|
||||||
|
)
|
||||||
|
fields.append(("summary", _summary_flat(content)))
|
||||||
|
# A field matches when EVERY needle appears in it (AND semantics,
|
||||||
|
# matching plainto_tsquery on the Postgres side).
|
||||||
|
best: SearchHit | None = None
|
||||||
|
for field, value in fields:
|
||||||
|
low = value.lower()
|
||||||
|
if not all(n in low for n in needles):
|
||||||
|
continue
|
||||||
|
at = low.find(needles[0])
|
||||||
|
cand = SearchHit(
|
||||||
|
recording=rec, field=field, snippet=_snippet_around(value, max(at, 0))
|
||||||
|
)
|
||||||
|
if best is None or _FIELD_RANK[field] < _FIELD_RANK[best.field]:
|
||||||
|
best = cand
|
||||||
|
if best is not None:
|
||||||
|
hits.append(best)
|
||||||
|
hits.sort(key=lambda h: (_FIELD_RANK[h.field], h.recording.recorded_at), reverse=False)
|
||||||
|
return hits[offset : offset + limit], len(hits)
|
||||||
|
|
||||||
|
|
||||||
|
# --- PostgreSQL tsvector path -------------------------------------------------
|
||||||
|
|
||||||
|
# CTE ``q`` carries the parsed tsquery so it is computed once. Notes are not
|
||||||
|
# in the recordings search_vector weights the way we want headlines, so notes
|
||||||
|
# match by substring like the SQLite path (title/notes share the vector; the
|
||||||
|
# field classifier prefers 'title' when the vector hits).
|
||||||
|
_PG_MATCHES = """
|
||||||
|
WITH q AS (SELECT plainto_tsquery('simple', :q) AS ts),
|
||||||
|
lt AS (
|
||||||
|
SELECT DISTINCT ON (recording_id) recording_id, search_vector, text
|
||||||
|
FROM transcripts WHERE superseded_at IS NULL
|
||||||
|
ORDER BY recording_id, version DESC
|
||||||
|
),
|
||||||
|
ls AS (
|
||||||
|
SELECT DISTINCT ON (recording_id) recording_id, content, search_vector
|
||||||
|
FROM summaries WHERE superseded_at IS NULL
|
||||||
|
ORDER BY recording_id, version DESC
|
||||||
|
),
|
||||||
|
tm AS (
|
||||||
|
SELECT DISTINCT rt.recording_id,
|
||||||
|
ts_headline('simple', t.name, q.ts,
|
||||||
|
'StartSel=,StopSel=,MaxFragments=0') AS snip
|
||||||
|
FROM recording_tags rt
|
||||||
|
JOIN tags t ON t.id = rt.tag_id AND t.user_id = :uid, q
|
||||||
|
),
|
||||||
|
matched AS (
|
||||||
|
SELECT r.id AS id,
|
||||||
|
CASE
|
||||||
|
WHEN :want_title AND r.search_vector @@ q.ts
|
||||||
|
AND coalesce(r.title, '') <> ''
|
||||||
|
THEN 'title'
|
||||||
|
WHEN :want_tag AND tm.recording_id IS NOT NULL THEN 'tag'
|
||||||
|
WHEN :want_notes AND coalesce(r.notes, '') ILIKE '%' || :raw || '%' THEN 'notes'
|
||||||
|
WHEN :want_summary AND ls.search_vector @@ q.ts THEN 'summary'
|
||||||
|
WHEN :want_transcript AND lt.search_vector @@ q.ts THEN 'transcript'
|
||||||
|
END AS field,
|
||||||
|
CASE
|
||||||
|
WHEN :want_title AND r.search_vector @@ q.ts
|
||||||
|
AND coalesce(r.title, '') <> ''
|
||||||
|
THEN ts_headline('simple', r.title, q.ts,
|
||||||
|
'StartSel=,StopSel=,MaxFragments=0,MaxWords=25')
|
||||||
|
WHEN :want_tag AND tm.recording_id IS NOT NULL THEN tm.snip
|
||||||
|
WHEN :want_notes AND coalesce(r.notes, '') ILIKE '%' || :raw || '%'
|
||||||
|
THEN left(r.notes, 200)
|
||||||
|
WHEN :want_summary AND ls.search_vector @@ q.ts
|
||||||
|
THEN ts_headline('simple',
|
||||||
|
coalesce(ls.content->>'short', '') || ' ' || coalesce(ls.content->>'detailed', ''),
|
||||||
|
q.ts, 'StartSel=,StopSel=,MaxFragments=1,MinWords=10,MaxWords=25')
|
||||||
|
WHEN :want_transcript AND lt.search_vector @@ q.ts
|
||||||
|
THEN ts_headline('simple', coalesce(lt.text, ''), q.ts,
|
||||||
|
'StartSel=,StopSel=,MaxFragments=1,MinWords=10,MaxWords=25')
|
||||||
|
END AS snippet
|
||||||
|
FROM recordings r
|
||||||
|
CROSS JOIN q
|
||||||
|
LEFT JOIN lt ON lt.recording_id = r.id
|
||||||
|
LEFT JOIN ls ON ls.recording_id = r.id
|
||||||
|
LEFT JOIN tm ON tm.recording_id = r.id
|
||||||
|
WHERE r.user_id = :uid AND r.deleted_at IS NULL
|
||||||
|
)
|
||||||
|
SELECT id, field, snippet FROM matched
|
||||||
|
WHERE field IS NOT NULL
|
||||||
|
ORDER BY CASE field WHEN 'title' THEN 0 WHEN 'tag' THEN 1 WHEN 'notes' THEN 2
|
||||||
|
WHEN 'summary' THEN 3 ELSE 4 END,
|
||||||
|
id
|
||||||
|
LIMIT :limit OFFSET :offset
|
||||||
|
"""
|
||||||
|
|
||||||
|
_PG_COUNT = """
|
||||||
|
WITH q AS (SELECT plainto_tsquery('simple', :q) AS ts),
|
||||||
|
lt AS (
|
||||||
|
SELECT DISTINCT ON (recording_id) recording_id, search_vector
|
||||||
|
FROM transcripts WHERE superseded_at IS NULL
|
||||||
|
),
|
||||||
|
ls AS (
|
||||||
|
SELECT DISTINCT ON (recording_id) recording_id, search_vector
|
||||||
|
FROM summaries WHERE superseded_at IS NULL
|
||||||
|
),
|
||||||
|
tm AS (
|
||||||
|
SELECT DISTINCT rt.recording_id
|
||||||
|
FROM recording_tags rt
|
||||||
|
JOIN tags t ON t.id = rt.tag_id AND t.user_id = :uid, q
|
||||||
|
WHERE t.search_vector @@ q.ts
|
||||||
|
)
|
||||||
|
SELECT count(*)
|
||||||
|
FROM recordings r
|
||||||
|
CROSS JOIN q
|
||||||
|
LEFT JOIN lt ON lt.recording_id = r.id
|
||||||
|
LEFT JOIN ls ON ls.recording_id = r.id
|
||||||
|
LEFT JOIN tm ON tm.recording_id = r.id
|
||||||
|
WHERE r.user_id = :uid AND r.deleted_at IS NULL
|
||||||
|
AND (
|
||||||
|
(:want_title AND r.search_vector @@ q.ts) OR
|
||||||
|
(:want_transcript AND lt.search_vector @@ q.ts) OR
|
||||||
|
(:want_summary AND ls.search_vector @@ q.ts) OR
|
||||||
|
(:want_tag AND tm.recording_id IS NOT NULL) OR
|
||||||
|
(:want_notes AND coalesce(r.notes, '') ILIKE '%' || :raw || '%')
|
||||||
|
)
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
async def _search_postgres(
|
||||||
|
session: AsyncSession, user_id: uuid.UUID, q: str, scope: str,
|
||||||
|
limit: int, offset: int,
|
||||||
|
) -> tuple[list[SearchHit], int]:
|
||||||
|
# NOTE: the title field matches anything the recordings vector hits
|
||||||
|
# (title + notes); notes-only hits surface under 'title' headlines from
|
||||||
|
# the title text. Acceptable precision tradeoff for a GIN-indexed path.
|
||||||
|
params = {
|
||||||
|
"uid": user_id,
|
||||||
|
"q": q,
|
||||||
|
"raw": q,
|
||||||
|
"limit": limit,
|
||||||
|
"offset": offset,
|
||||||
|
"want_title": scope in ("all", "title"),
|
||||||
|
"want_transcript": scope in ("all", "transcript"),
|
||||||
|
"want_summary": scope in ("all", "summary"),
|
||||||
|
"want_tag": scope in ("all", "tag"),
|
||||||
|
"want_notes": scope in ("all", "notes"),
|
||||||
|
}
|
||||||
|
rows = (await session.execute(text(_PG_MATCHES), params)).all()
|
||||||
|
total = await session.scalar(text(_PG_COUNT), params) or 0
|
||||||
|
if not rows:
|
||||||
|
return [], total
|
||||||
|
ids = [r[0] for r in rows]
|
||||||
|
by_id = {
|
||||||
|
rec.id: rec
|
||||||
|
for rec in (
|
||||||
|
await session.scalars(select(Recording).where(Recording.id.in_(ids)))
|
||||||
|
).all()
|
||||||
|
}
|
||||||
|
hits = [
|
||||||
|
SearchHit(recording=by_id[r.id], field=r.field, snippet=r.snippet or "")
|
||||||
|
for r in rows
|
||||||
|
if r.id in by_id
|
||||||
|
]
|
||||||
|
return hits, total
|
||||||
|
|
||||||
|
|
||||||
|
# --- public API ----------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def search_recordings(
|
||||||
|
session: AsyncSession,
|
||||||
|
user_id: uuid.UUID,
|
||||||
|
q: str,
|
||||||
|
*,
|
||||||
|
scope: str = "all",
|
||||||
|
limit: int = 20,
|
||||||
|
offset: int = 0,
|
||||||
|
) -> tuple[list[SearchHit], int]:
|
||||||
|
"""Returns (page of hits ordered by field rank, total). Empty q → no hits."""
|
||||||
|
q = q.strip()
|
||||||
|
if not q:
|
||||||
|
return [], 0
|
||||||
|
dialect = session.bind.dialect.name if session.bind else "sqlite"
|
||||||
|
if dialect == "postgresql":
|
||||||
|
return await _search_postgres(session, user_id, q, scope, limit, offset)
|
||||||
|
return await _search_sqlite(session, user_id, q, scope, limit, offset)
|
||||||
|
|
@ -31,6 +31,14 @@ async def sweep(ctx: dict) -> None: # noqa: ARG001 — arq cron signature
|
||||||
logger.info("sweep re-enqueued %d stale jobs", count)
|
logger.info("sweep re-enqueued %d stale jobs", count)
|
||||||
|
|
||||||
|
|
||||||
|
async def retention_sweep(ctx: dict) -> None: # noqa: ARG001 — arq cron signature
|
||||||
|
from shonar.services import retention
|
||||||
|
|
||||||
|
purged = await retention.sweep_deleted()
|
||||||
|
if purged["recordings"] or purged["users"]:
|
||||||
|
logger.info("retention sweep: %s", purged)
|
||||||
|
|
||||||
|
|
||||||
async def startup(ctx: dict) -> None:
|
async def startup(ctx: dict) -> None:
|
||||||
get_engine()
|
get_engine()
|
||||||
# Crash recovery before accepting new work: jobs stuck `running` and
|
# Crash recovery before accepting new work: jobs stuck `running` and
|
||||||
|
|
@ -49,11 +57,15 @@ def _redis() -> RedisSettings:
|
||||||
|
|
||||||
|
|
||||||
class WorkerSettings:
|
class WorkerSettings:
|
||||||
functions = [run_transcribe, run_summarize, sweep]
|
functions = [run_transcribe, run_summarize, sweep, retention_sweep]
|
||||||
# Transport-loss backstop beyond the startup sweep: anything still
|
# Transport-loss backstop beyond the startup sweep: anything still
|
||||||
# queued (missed enqueue, dead worker between runs) goes back through
|
# queued (missed enqueue, dead worker between runs) goes back through
|
||||||
# arq every 5 minutes. Rows are the queue; this just pokes.
|
# arq every 5 minutes. Rows are the queue; this just pokes.
|
||||||
cron_jobs = [cron(sweep, minute={0, 5, 10, 15, 20, 25, 30, 35, 40, 45, 50, 55})]
|
cron_jobs = [
|
||||||
|
cron(sweep, minute={0, 5, 10, 15, 20, 25, 30, 35, 40, 45, 50, 55}),
|
||||||
|
# Hard-delete expired soft-deletes once a day (3:17 local, off-peak).
|
||||||
|
cron(retention_sweep, hour=3, minute=17),
|
||||||
|
]
|
||||||
on_startup = startup
|
on_startup = startup
|
||||||
on_shutdown = shutdown
|
on_shutdown = shutdown
|
||||||
redis_settings = _redis()
|
redis_settings = _redis()
|
||||||
|
|
|
||||||
|
|
@ -40,6 +40,32 @@ async def _setup_db() -> AsyncIterator[None]:
|
||||||
async with engine.begin() as conn:
|
async with engine.begin() as conn:
|
||||||
await conn.run_sync(Base.metadata.drop_all)
|
await conn.run_sync(Base.metadata.drop_all)
|
||||||
await conn.run_sync(Base.metadata.create_all)
|
await conn.run_sync(Base.metadata.create_all)
|
||||||
|
# The tsvector generated columns live only in migration
|
||||||
|
# fts0000000001 (not in the ORM models), so create_all misses them.
|
||||||
|
# Apply the same DDL the migration applies (M9 search needs them).
|
||||||
|
if engine.dialect.name == "postgresql":
|
||||||
|
# Reuse the real migration's DDL (not a copy) via a sync
|
||||||
|
# MigrationContext — op.execute() is synchronous there.
|
||||||
|
import importlib.util
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
def _apply(sync_conn):
|
||||||
|
from alembic.migration import MigrationContext
|
||||||
|
from alembic.operations import Operations
|
||||||
|
|
||||||
|
spec = importlib.util.spec_from_file_location(
|
||||||
|
"fts_migration",
|
||||||
|
Path(__file__).resolve().parents[1]
|
||||||
|
/ "migrations/versions/fts0000000001_fts_columns.py",
|
||||||
|
)
|
||||||
|
assert spec is not None and spec.loader is not None
|
||||||
|
mod = importlib.util.module_from_spec(spec)
|
||||||
|
spec.loader.exec_module(mod)
|
||||||
|
ctx = MigrationContext.configure(sync_conn)
|
||||||
|
with Operations.context(ctx):
|
||||||
|
mod.upgrade()
|
||||||
|
|
||||||
|
await conn.run_sync(_apply)
|
||||||
yield
|
yield
|
||||||
await dispose_engine()
|
await dispose_engine()
|
||||||
|
|
||||||
|
|
|
||||||
304
backend/tests/test_m9.py
Normal file
304
backend/tests/test_m9.py
Normal file
|
|
@ -0,0 +1,304 @@
|
||||||
|
"""M9 tests: search, exports, retention sweep.
|
||||||
|
|
||||||
|
Search runs against the real backend dialect (Postgres in CI/dev, SQLite
|
||||||
|
via SHONAR_TEST_DATABASE_URL) — both paths share the endpoint contract.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import io
|
||||||
|
import zipfile
|
||||||
|
|
||||||
|
from tests.test_recordings import auth, user_tokens, wav_bytes
|
||||||
|
|
||||||
|
|
||||||
|
async def make_recording(client, token, title, notes=None, tags=None, recorded_at=None):
|
||||||
|
h = await auth(token)
|
||||||
|
audio = wav_bytes()
|
||||||
|
body = {"declared_mime_type": "audio/wav", "declared_size_bytes": len(audio), "title": title}
|
||||||
|
r = await client.post("/api/v1/uploads", json=body, headers=h)
|
||||||
|
sid = r.json()["id"]
|
||||||
|
await client.put(
|
||||||
|
f"/api/v1/uploads/{sid}/chunks/0", content=audio,
|
||||||
|
headers={**h, "content-type": "application/octet-stream"},
|
||||||
|
)
|
||||||
|
fin = {"duration_seconds": 5.0}
|
||||||
|
if recorded_at:
|
||||||
|
fin["recorded_at"] = recorded_at
|
||||||
|
if notes:
|
||||||
|
fin["notes"] = notes
|
||||||
|
r = await client.post(f"/api/v1/uploads/{sid}/finalize", json=fin, headers=h)
|
||||||
|
assert r.status_code == 201, r.text
|
||||||
|
rec = r.json()
|
||||||
|
if tags:
|
||||||
|
r = await client.patch(
|
||||||
|
f"/api/v1/recordings/{rec['id']}", json={"tags": tags}, headers=h
|
||||||
|
)
|
||||||
|
assert r.status_code == 200
|
||||||
|
return rec["id"]
|
||||||
|
|
||||||
|
|
||||||
|
async def add_transcript(rec_id: str, text: str, segments=None):
|
||||||
|
"""Insert a transcript row directly (no AI provider under test)."""
|
||||||
|
from shonar.db.models import Transcript
|
||||||
|
from shonar.db.session import session_factory
|
||||||
|
|
||||||
|
async with session_factory()() as s:
|
||||||
|
s.add(Transcript(recording_id=rec_id, text=text, segments=segments, provider="test"))
|
||||||
|
await s.commit()
|
||||||
|
|
||||||
|
|
||||||
|
async def add_summary(rec_id: str, content: dict):
|
||||||
|
from shonar.db.models import Summary
|
||||||
|
from shonar.db.session import session_factory
|
||||||
|
|
||||||
|
async with session_factory()() as s:
|
||||||
|
s.add(Summary(recording_id=rec_id, content=content, provider="test"))
|
||||||
|
await s.commit()
|
||||||
|
|
||||||
|
|
||||||
|
# --- search -------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def test_search_title_and_transcript(client):
|
||||||
|
token = await user_tokens(client, email="m9s1@example.com")
|
||||||
|
rid_t = await make_recording(client, token, "Quarterly budget review")
|
||||||
|
rid_x = await make_recording(client, token, "Grocery list")
|
||||||
|
await add_transcript(rid_x, "remember to buy kale chips and quinoa tonight")
|
||||||
|
|
||||||
|
r = await client.get("/api/v1/search", params={"q": "budget"}, headers=await auth(token))
|
||||||
|
assert r.status_code == 200, r.text
|
||||||
|
body = r.json()
|
||||||
|
assert body["total"] == 1
|
||||||
|
assert body["items"][0]["id"] == rid_t
|
||||||
|
assert body["items"][0]["field"] == "title"
|
||||||
|
|
||||||
|
r = await client.get("/api/v1/search", params={"q": "quinoa"}, headers=await auth(token))
|
||||||
|
body = r.json()
|
||||||
|
assert body["total"] == 1
|
||||||
|
assert body["items"][0]["id"] == rid_x
|
||||||
|
assert body["items"][0]["field"] == "transcript"
|
||||||
|
assert "quinoa" in body["items"][0]["snippet"].lower()
|
||||||
|
|
||||||
|
|
||||||
|
async def test_search_scope_and_tag(client):
|
||||||
|
token = await user_tokens(client, email="m9s2@example.com")
|
||||||
|
rid = await make_recording(client, token, "Standup", tags=["daily"])
|
||||||
|
await add_transcript(rid, "we discussed the daily standup format")
|
||||||
|
|
||||||
|
# tag scope finds by tag name
|
||||||
|
r = await client.get(
|
||||||
|
"/api/v1/search", params={"q": "daily", "scope": "tag"}, headers=await auth(token)
|
||||||
|
)
|
||||||
|
assert r.json()["total"] == 1
|
||||||
|
assert r.json()["items"][0]["field"] == "tag"
|
||||||
|
|
||||||
|
# a transcript-only word does NOT match under scope=title
|
||||||
|
r = await client.get(
|
||||||
|
"/api/v1/search", params={"q": "discussed", "scope": "title"}, headers=await auth(token)
|
||||||
|
)
|
||||||
|
assert r.json()["total"] == 0
|
||||||
|
r = await client.get(
|
||||||
|
"/api/v1/search", params={"q": "discussed", "scope": "transcript"},
|
||||||
|
headers=await auth(token),
|
||||||
|
)
|
||||||
|
assert r.json()["total"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
async def test_search_isolation_and_deleted(client):
|
||||||
|
token_a = await user_tokens(client, email="m9s3a@example.com")
|
||||||
|
token_b = await user_tokens(client, email="m9s3b@example.com")
|
||||||
|
rid = await make_recording(client, token_a, "secret sauce recipe")
|
||||||
|
|
||||||
|
h_b = await auth(token_b)
|
||||||
|
r = await client.get("/api/v1/search", params={"q": "secret"}, headers=h_b)
|
||||||
|
assert r.json()["total"] == 0 # other users' data invisible
|
||||||
|
|
||||||
|
# soft-deleted rows drop out of search
|
||||||
|
await client.delete(f"/api/v1/recordings/{rid}", headers=await auth(token_a))
|
||||||
|
r = await client.get(
|
||||||
|
"/api/v1/search", params={"q": "secret"}, headers=await auth(token_a)
|
||||||
|
)
|
||||||
|
assert r.json()["total"] == 0
|
||||||
|
|
||||||
|
|
||||||
|
async def test_search_summary_and_notes(client):
|
||||||
|
token = await user_tokens(client, email="m9s4@example.com")
|
||||||
|
rid = await make_recording(client, token, "Meeting", notes="bring the projector cable")
|
||||||
|
await add_summary(rid, {"short": "sprint retro", "action_items": ["fix flaky test"]})
|
||||||
|
|
||||||
|
r = await client.get("/api/v1/search", params={"q": "projector"}, headers=await auth(token))
|
||||||
|
assert r.json()["total"] == 1
|
||||||
|
r = await client.get("/api/v1/search", params={"q": "flaky"}, headers=await auth(token))
|
||||||
|
assert r.json()["total"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
# --- list filters ---------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def test_list_filters_tag_status_date(client):
|
||||||
|
token = await user_tokens(client, email="m9f1@example.com")
|
||||||
|
await make_recording(client, token, "Old one", tags=["keep"],
|
||||||
|
recorded_at="2020-01-01T10:00:00Z")
|
||||||
|
rid_new = await make_recording(client, token, "New one", tags=["keep"],
|
||||||
|
recorded_at="2026-01-01T10:00:00Z")
|
||||||
|
h = await auth(token)
|
||||||
|
|
||||||
|
r = await client.get("/api/v1/recordings", params={"tag": "keep"}, headers=h)
|
||||||
|
assert r.json()["total"] == 2
|
||||||
|
r = await client.get(
|
||||||
|
"/api/v1/recordings", params={"tag": "keep", "from_date": "2025-06-01T00:00:00Z"},
|
||||||
|
headers=h,
|
||||||
|
)
|
||||||
|
body = r.json()
|
||||||
|
assert body["total"] == 1 and body["items"][0]["id"] == rid_new
|
||||||
|
|
||||||
|
r = await client.get("/api/v1/recordings", params={"status": "ai_disabled"}, headers=h)
|
||||||
|
assert r.json()["total"] == 2
|
||||||
|
r = await client.get("/api/v1/recordings", params={"status": "bogus"}, headers=h)
|
||||||
|
assert r.status_code == 422
|
||||||
|
r = await client.get("/api/v1/recordings", params={"tag": "nope"}, headers=h)
|
||||||
|
assert r.json()["total"] == 0
|
||||||
|
|
||||||
|
|
||||||
|
# --- exports --------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def test_export_formats(client):
|
||||||
|
token = await user_tokens(client, email="m9e1@example.com")
|
||||||
|
rid = await make_recording(client, token, "Retro & Planning", notes="retro notes here",
|
||||||
|
tags=["team"])
|
||||||
|
await add_transcript(
|
||||||
|
rid, "first segment second segment",
|
||||||
|
segments=[{"start": 0.0, "end": 2.5, "text": "first segment", "speaker": None},
|
||||||
|
{"start": 2.5, "end": 5.0, "text": "second segment", "speaker": "S1"}],
|
||||||
|
)
|
||||||
|
await add_summary(rid, {"short": "one line", "action_items": ["do the thing"]})
|
||||||
|
h = await auth(token)
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "txt"}, headers=h)
|
||||||
|
assert r.status_code == 200
|
||||||
|
assert r.text == "first segment second segment"
|
||||||
|
assert "attachment" in r.headers["content-disposition"]
|
||||||
|
assert ".txt" in r.headers["content-disposition"]
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "md"}, headers=h)
|
||||||
|
assert r.status_code == 200
|
||||||
|
assert "# Retro & Planning" in r.text
|
||||||
|
assert "do the thing" in r.text
|
||||||
|
assert "`[00:02.5]` **S1**: second segment" in r.text
|
||||||
|
assert "`team`" in r.text
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "zip"}, headers=h)
|
||||||
|
assert r.status_code == 200
|
||||||
|
assert r.headers["content-type"] == "application/zip"
|
||||||
|
with zipfile.ZipFile(io.BytesIO(r.content)) as z:
|
||||||
|
names = z.namelist()
|
||||||
|
assert "transcript.txt" in names and "notes.md" in names
|
||||||
|
assert any(n.endswith(".wav") for n in names)
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "audio"}, headers=h)
|
||||||
|
assert r.status_code == 200
|
||||||
|
assert r.content.startswith(b"RIFF")
|
||||||
|
|
||||||
|
|
||||||
|
async def test_export_missing_and_foreign(client):
|
||||||
|
token = await user_tokens(client, email="m9e2@example.com")
|
||||||
|
rid = await make_recording(client, token, "Bare") # no transcript/summary/notes
|
||||||
|
token2 = await user_tokens(client, email="m9e2b@example.com")
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "txt"},
|
||||||
|
headers=await auth(token))
|
||||||
|
assert r.status_code == 404 # no transcript yet
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "md"},
|
||||||
|
headers=await auth(token))
|
||||||
|
assert r.status_code == 404
|
||||||
|
# zip still works with just the audio
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "zip"},
|
||||||
|
headers=await auth(token))
|
||||||
|
assert r.status_code == 200
|
||||||
|
|
||||||
|
r = await client.get(f"/api/v1/recordings/{rid}/export", params={"fmt": "audio"},
|
||||||
|
headers=await auth(token2))
|
||||||
|
assert r.status_code == 404 # not yours
|
||||||
|
|
||||||
|
|
||||||
|
# --- retention sweep --------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def test_retention_sweep_purges_expired(client):
|
||||||
|
from datetime import timedelta
|
||||||
|
|
||||||
|
from shonar.db.models import Asset, Recording, utcnow
|
||||||
|
from shonar.db.session import session_factory
|
||||||
|
from shonar.services import retention
|
||||||
|
|
||||||
|
token = await user_tokens(client, email="m9r1@example.com")
|
||||||
|
rid = await make_recording(client, token, "Doomed")
|
||||||
|
h = await auth(token)
|
||||||
|
|
||||||
|
# Soft-delete, then push deleted_at past the grace window directly.
|
||||||
|
r = await client.delete(f"/api/v1/recordings/{rid}", headers=h)
|
||||||
|
assert r.status_code == 204
|
||||||
|
async with session_factory()() as s:
|
||||||
|
rec = await s.get(Recording, rid)
|
||||||
|
rec.deleted_at = utcnow() - timedelta(days=31)
|
||||||
|
# remember the storage key before the row vanishes
|
||||||
|
from shonar.db.models import Asset
|
||||||
|
|
||||||
|
asset = (
|
||||||
|
await s.execute(
|
||||||
|
Asset.__table__.select().where(Asset.recording_id == rec.id) # noqa: SLF001
|
||||||
|
)
|
||||||
|
).first()
|
||||||
|
storage_key = asset.storage_key if asset else None
|
||||||
|
await s.commit()
|
||||||
|
assert storage_key is not None
|
||||||
|
|
||||||
|
purged = await retention.sweep_deleted()
|
||||||
|
assert purged["recordings"] == 1
|
||||||
|
assert purged["files"] >= 1
|
||||||
|
|
||||||
|
async with session_factory()() as s:
|
||||||
|
assert await s.get(Recording, rid) is None
|
||||||
|
from shonar.storage import get_storage
|
||||||
|
|
||||||
|
assert not await get_storage().exists(storage_key) # file gone too
|
||||||
|
|
||||||
|
# Within the window, nothing is purged (cancellable).
|
||||||
|
rid2 = await make_recording(client, token, "Fresh delete")
|
||||||
|
await client.delete(f"/api/v1/recordings/{rid2}", headers=h)
|
||||||
|
purged = await retention.sweep_deleted()
|
||||||
|
assert purged["recordings"] == 0
|
||||||
|
async with session_factory()() as s:
|
||||||
|
assert await s.get(Recording, rid2) is not None
|
||||||
|
|
||||||
|
|
||||||
|
async def test_retention_sweep_account(client):
|
||||||
|
from datetime import timedelta
|
||||||
|
|
||||||
|
from shonar.db.models import Recording, utcnow
|
||||||
|
from shonar.db.session import session_factory
|
||||||
|
from shonar.services import retention
|
||||||
|
|
||||||
|
token = await user_tokens(client, email="m9r2@example.com")
|
||||||
|
rid = await make_recording(client, token, "Gone with user")
|
||||||
|
|
||||||
|
async with session_factory()() as s:
|
||||||
|
user = await s.scalar(select_user("m9r2@example.com"))
|
||||||
|
user.deleted_at = utcnow() - timedelta(days=31)
|
||||||
|
await s.commit()
|
||||||
|
|
||||||
|
purged = await retention.sweep_deleted()
|
||||||
|
assert purged["users"] == 1
|
||||||
|
async with session_factory()() as s:
|
||||||
|
assert await s.get(Recording, rid) is None # cascade
|
||||||
|
assert await s.scalar(select_user("m9r2@example.com")) is None
|
||||||
|
|
||||||
|
|
||||||
|
def select_user(email: str):
|
||||||
|
from sqlalchemy import select
|
||||||
|
|
||||||
|
from shonar.db.models import User
|
||||||
|
|
||||||
|
return select(User).where(User.email == email)
|
||||||
|
|
@ -18,7 +18,7 @@ updated in the same commit as the work it describes.
|
||||||
| M7 | Backend: AI pipeline + adapters (whisper_http, faster-whisper, OpenAI-compat, Ollama, none), status endpoints | done — provider protocols + 4 adapters (faster-whisper lazy optional), versioned transcripts/summaries (user edits win), arq worker (`run_transcribe`/`run_summarize` + startup/5-min sweep, max 3 tries), transcript/summary/jobs endpoints, `none` means skipped; 50 backend tests green, ruff clean |
|
| M7 | Backend: AI pipeline + adapters (whisper_http, faster-whisper, OpenAI-compat, Ollama, none), status endpoints | done — provider protocols + 4 adapters (faster-whisper lazy optional), versioned transcripts/summaries (user edits win), arq worker (`run_transcribe`/`run_summarize` + startup/5-min sweep, max 3 tries), transcript/summary/jobs endpoints, `none` means skipped; 50 backend tests green, ruff clean |
|
||||||
| M8 | Android: details screen — transcript synced to playback, summary, action items, editing | done — `detail/{id}` route (transcript/summary/status tabs, tap-to-seek, speed control, edit dialogs, title rename) on `CustomShonarProvider` AI methods (fetch/update transcript+summary, jobs, PATCH title/notes) against new backend `PUT transcript/summary` user-edit endpoints (versioned, pipeline won't clobber); pure `AiContent` parsing/sync mapping; 19 new Android tests + 5 backend edit tests, full suites green (186 Android, 55 backend), `assembleDebug` clean. Local-only/unsynced rows get honest empty states; on-device verification pending |
|
| M8 | Android: details screen — transcript synced to playback, summary, action items, editing | done — `detail/{id}` route (transcript/summary/status tabs, tap-to-seek, speed control, edit dialogs, title rename) on `CustomShonarProvider` AI methods (fetch/update transcript+summary, jobs, PATCH title/notes) against new backend `PUT transcript/summary` user-edit endpoints (versioned, pipeline won't clobber); pure `AiContent` parsing/sync mapping; 19 new Android tests + 5 backend edit tests, full suites green (186 Android, 55 backend), `assembleDebug` clean. Local-only/unsynced rows get honest empty states; on-device verification pending |
|
||||||
| T1 | Backend transcription model support — registry (tiny/base/small/medium/large-v3), global default (`base`, runtime-editable), per-recording overrides, `GET /api/v1/models`, worker uses saved model, job stage/progress | done — migration (`recordings`+`upload_sessions.transcription_model`, `processing_jobs.stage/progress`, `app_settings`); `PUT /models/default` (future rows only); finalize/session override (finalize wins); faster-whisper model cache + fail-fast unavailable errors; 10 new tests, backend suite 64 green + ruff clean |
|
| T1 | Backend transcription model support — registry (tiny/base/small/medium/large-v3), global default (`base`, runtime-editable), per-recording overrides, `GET /api/v1/models`, worker uses saved model, job stage/progress | done — migration (`recordings`+`upload_sessions.transcription_model`, `processing_jobs.stage/progress`, `app_settings`); `PUT /models/default` (future rows only); finalize/session override (finalize wins); faster-whisper model cache + fail-fast unavailable errors; 10 new tests, backend suite 64 green + ruff clean |
|
||||||
| M9 | Backend: full-text search endpoints + filters, exports (audio/txt/md/zip), deletion sweep | TODO (schema/FTS columns exist) |
|
| M9 | Backend: full-text search endpoints + filters, exports (audio/txt/md/zip), deletion sweep | done — `GET /search` (Postgres tsvector via migration fts0000000001, SQLite LIKE fallback for the desktop engine; scopes title/notes/transcript/summary/tag, ranked + snippets); list filters on `GET /recordings` (tag/status/from/to); `GET /recordings/{id}/export?fmt=audio|txt|md|zip` (inline artifacts, ExportJob+Asset audit trail); retention sweep (`services/retention.py`, 30-day grace `SHONAR_RETENTION_GRACE_DAYS`, wired into arq cron daily + inline-queue timer); 9 new tests, backend suite 78 green + ruff clean |
|
||||||
| D-1 | Desktop: in-place reprocess (no re-upload) — `<name>.shonar.json` sidecar maps library file → server recording id; Re-transcribe/Re-summarize hit `POST /reprocess` (model override switches persist server-side); rename carries mapping; deleted-remote mapping self-heals to re-upload | done — provider `reprocess()`, endpoint `?model=` param, 2 backend + 6 desktop tests, suites green (69 backend, 23 desktop) |
|
| D-1 | Desktop: in-place reprocess (no re-upload) — `<name>.shonar.json` sidecar maps library file → server recording id; Re-transcribe/Re-summarize hit `POST /reprocess` (model override switches persist server-side); rename carries mapping; deleted-remote mapping self-heals to re-upload | done — provider `reprocess()`, endpoint `?model=` param, 2 backend + 6 desktop tests, suites green (69 backend, 23 desktop) |
|
||||||
| M10 | Android dark mode, accessibility pass, consent UX polish, deploy/backup docs, OpenAPI sync | TODO |
|
| M10 | Android dark mode, accessibility pass, consent UX polish, deploy/backup docs, OpenAPI sync | TODO |
|
||||||
|
|
||||||
|
|
@ -31,5 +31,6 @@ updated in the same commit as the work it describes.
|
||||||
started. Files are stored unencrypted unless you encrypt the volume.
|
started. Files are stored unencrypted unless you encrypt the volume.
|
||||||
- **Meilisearch/OpenSearch search backend**: only Postgres FTS is planned to
|
- **Meilisearch/OpenSearch search backend**: only Postgres FTS is planned to
|
||||||
ship first, behind a `SearchBackend` protocol.
|
ship first, behind a `SearchBackend` protocol.
|
||||||
- **Account purge sweep**: deletion marks a 30-day grace; the scheduled hard
|
- ~~**Account purge sweep**~~: done in M9 — `services/retention.py`
|
||||||
purge job is TODO (`worker/tasks.py::purge_deleted_accounts`).
|
hard-deletes expired soft-deletes (recordings + accounts + files);
|
||||||
|
runs daily via arq cron and the desktop inline-queue timer.
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue