diff --git a/backend/shonar/services/ai/__init__.py b/backend/shonar/services/ai/__init__.py index cba72b7..700254c 100644 --- a/backend/shonar/services/ai/__init__.py +++ b/backend/shonar/services/ai/__init__.py @@ -59,6 +59,7 @@ class TranscriptionProvider(Protocol): mime: str, *, language_hint: str | None = None, + on_progress=None, # optional Callable[[int], None], 0..99 ) -> TranscriptResult: ... diff --git a/backend/shonar/services/ai/faster_whisper.py b/backend/shonar/services/ai/faster_whisper.py index 7ed0ca6..3c2b9b7 100644 --- a/backend/shonar/services/ai/faster_whisper.py +++ b/backend/shonar/services/ai/faster_whisper.py @@ -61,11 +61,17 @@ class FasterWhisperProvider: mime: str, *, language_hint: str | None = None, + on_progress=None, # Callable[[int], None] | None — 0..99 percent ) -> TranscriptResult: # faster-whisper is blocking CPU work: keep it off the event loop. - return await asyncio.to_thread(self._run, audio, language_hint) + return await asyncio.to_thread(self._run, audio, language_hint, on_progress) - def _run(self, audio: bytes, language_hint: str | None) -> TranscriptResult: + def _run( + self, + audio: bytes, + language_hint: str | None, + on_progress=None, + ) -> TranscriptResult: path: Path | None = None try: with tempfile.NamedTemporaryFile(suffix=".m4a", delete=False) as f: @@ -78,10 +84,18 @@ class FasterWhisperProvider: beam_size=5, language=language_hint, ) - segments = [ - Segment(start=s.start, end=s.end, text=s.text.strip()) - for s in segments_iter - ] + duration = float(getattr(info, "duration", 0.0) or 0.0) + segments = [] + last_pct = -1 + for s in segments_iter: + segments.append(Segment(start=s.start, end=s.end, text=s.text.strip())) + if on_progress is not None and duration > 0: + # Throttle: report only on whole-percent gains. Capped + # at 99 — the caller commits 100 when the row finishes. + pct = min(99, int(s.end / duration * 100)) + if pct > last_pct: + last_pct = pct + on_progress(pct) text = " ".join(s.text for s in segments).strip() return TranscriptResult( text=text, diff --git a/backend/shonar/services/ai/whisper_http.py b/backend/shonar/services/ai/whisper_http.py index a5e791a..b25e81f 100644 --- a/backend/shonar/services/ai/whisper_http.py +++ b/backend/shonar/services/ai/whisper_http.py @@ -53,6 +53,7 @@ class WhisperHttpProvider: mime: str, *, language_hint: str | None = None, + on_progress=None, # accepted for protocol parity; not reported ) -> TranscriptResult: headers = ( {"Authorization": f"Bearer {self.api_key}"} if self.api_key else {} diff --git a/backend/shonar/services/processing.py b/backend/shonar/services/processing.py index 8090cac..ca91b24 100644 --- a/backend/shonar/services/processing.py +++ b/backend/shonar/services/processing.py @@ -16,6 +16,7 @@ Entry points: from __future__ import annotations +import asyncio import logging import uuid from datetime import timedelta @@ -46,6 +47,49 @@ from shonar.storage import get_storage logger = logging.getLogger("shonar.processing") + +def _accepts_on_progress(provider) -> bool: + """True when the provider's transcribe() accepts on_progress=.""" + import inspect + + try: + return "on_progress" in inspect.signature(provider.transcribe).parameters + except (TypeError, ValueError): # builtins / exotic callables + return False + + +def _thread_progress_reporter(job_id): + """Callback safe to invoke from a worker thread (faster-whisper runs + via asyncio.to_thread): schedules a tiny session update on the loop. + + Progress writes are best-effort display state — failures are logged + and swallowed, never allowed to disturb the transcription itself.""" + loop = asyncio.get_running_loop() + + def report(pct: int) -> None: + async def _write() -> None: + from sqlalchemy import update + + from shonar.db.session import session_factory + + try: + async with session_factory()() as s: + await s.execute( + update(ProcessingJob) + .where(ProcessingJob.id == job_id) + .values(progress=int(pct)) + ) + await s.commit() + except Exception: # pragma: no cover - display state only + logger.debug("progress write failed for job %s", job_id, exc_info=True) + + try: + asyncio.run_coroutine_threadsafe(_write(), loop) + except RuntimeError: # loop already gone (shutdown race) + logger.debug("progress dropped for job %s (loop gone)", job_id) + + return report + MAX_TRIES = 3 # A `running` job younger than this is treated as live work, not a crash @@ -230,8 +274,13 @@ async def _fail( (returning normally), otherwise they raise for arq retry.""" transient = exc is None or isinstance(exc, ProviderTransientError) if transient and _job_try(ctx) < MAX_TRIES: + # Running state was committed before the long phase; put the row + # back to queued for the retry and persist that (a raise no longer + # rolls the pre-phase commit back). + job.status = JobStatus.queued job.attempt = _job_try(ctx) - await session.flush() + job.stage = None + await session.commit() raise ProviderTransientError(message) job.status = JobStatus.failed job.error = message @@ -293,12 +342,17 @@ async def run_transcribe(ctx: dict, recording_id: str) -> None: job.progress = None rec.processing_status = ProcessingStatus.processing rec.processing_error = None - await session.flush() + # Commit before the long CPU phase: an open transaction is invisible + # to other readers (and on SQLite it locks out the progress writer). + await session.commit() try: audio = await get_storage().get(original.storage_key) job.stage = "transcribing" - await session.flush() - result = await provider.transcribe(audio, original.mime_type) + await session.commit() + kwargs = {} + if _accepts_on_progress(provider): + kwargs["on_progress"] = _thread_progress_reporter(job.id) + result = await provider.transcribe(audio, original.mime_type, **kwargs) except AIError as e: await _fail(session, rec, job, str(e), ctx, e) await session.commit() @@ -371,8 +425,9 @@ async def run_summarize(ctx: dict, recording_id: str) -> None: job.status = JobStatus.running job.attempt = _job_try(ctx) job.started_at = utcnow() + job.stage = "summarizing" rec.processing_status = ProcessingStatus.processing - await session.flush() + await session.commit() # visible before the long LLM call try: result = await provider.summarize(text, title=rec.title) except AIError as e: @@ -381,6 +436,8 @@ async def run_summarize(ctx: dict, recording_id: str) -> None: return await store_summary(session, rec, result.to_dict(), provider.name, result.model) job.status = JobStatus.succeeded + job.stage = None + job.progress = 100 job.finished_at = utcnow() rec.processing_status = ProcessingStatus.completed rec.processing_error = None diff --git a/desktop/app/src/main/kotlin/com/shonar/desktop/DesktopState.kt b/desktop/app/src/main/kotlin/com/shonar/desktop/DesktopState.kt index a6f3aff..f0a35cc 100644 --- a/desktop/app/src/main/kotlin/com/shonar/desktop/DesktopState.kt +++ b/desktop/app/src/main/kotlin/com/shonar/desktop/DesktopState.kt @@ -48,6 +48,8 @@ class DesktopState(private val appDir: File = defaultAppDir()) { val hasReport: Boolean, val status: FileStatus = FileStatus.NEW, val statusNote: String? = null, + /** 0..1 while transcribing (backend-reported); null otherwise. */ + val progress: Float? = null, ) data class DetailUi( @@ -303,7 +305,18 @@ class DesktopState(private val appDir: File = defaultAppDir()) { "SHONAR_TRANSCRIPTION_PROVIDER" to "faster_whisper", "SHONAR_TRANSCRIPTION_MODEL" to "base", ) - if (ollamaUp()) { + // Summarizer preference: LAN inference server (H200) when + // its key is present in the environment, then local Ollama. + // LAN-only: the key never leaves the network boundary and + // is read from the environment, never stored by this app. + val lanKey = System.getenv(ENV_LAN_LLM_KEY)?.takeIf { it.isNotBlank() } + ?: readLanKeyFromHermesEnv() + if (lanKey.isNotBlank()) { + env["SHONAR_LLM_PROVIDER"] = "openai_compat" + env["SHONAR_LLM_BASE_URL"] = LAN_LLM_BASE_URL + env["SHONAR_LLM_MODEL"] = LAN_LLM_MODEL + env["SHONAR_LLM_API_KEY"] = lanKey + } else if (ollamaUp()) { env["SHONAR_LLM_PROVIDER"] = "ollama" env["SHONAR_LLM_BASE_URL"] = "http://127.0.0.1:11434" env["SHONAR_LLM_MODEL"] = ollamaModel() @@ -352,6 +365,17 @@ class DesktopState(private val appDir: File = defaultAppDir()) { apiProc = null } + /** The H200 key lives in ~/.hermes/.env (0600, same machine); desktop + * launches from a plain session don't carry it as an env var, so read + * it from there. Never logged, never written anywhere else. */ + private fun readLanKeyFromHermesEnv(): String = runCatching { + File(System.getProperty("user.home"), ".hermes/.env") + .readLines() + .firstOrNull { it.trimStart().startsWith("$ENV_LAN_LLM_KEY=") } + ?.substringAfter('=')?.trim() + .orEmpty() + }.getOrDefault("") + private fun ollamaUp(): Boolean = runCatching { java.net.URL("http://127.0.0.1:11434/api/tags").readText().contains("\"models\"") }.getOrDefault(false) @@ -483,14 +507,28 @@ class DesktopState(private val appDir: File = defaultAppDir()) { } } + /** file name -> live activity from job polling. */ + private val _liveProgress = MutableStateFlow>(emptyMap()) + val liveProgress = _liveProgress.asStateFlow() + + private fun setLive(name: String, live: LiveProgress?) { + _liveProgress.value = _liveProgress.value.toMutableMap().apply { + if (live == null) remove(name) else put(name, live) + } + } + private fun rescanStatuses() { + val live = _liveProgress.value _entries.value = _entries.value.map { e -> val done = reportFile(e.file).exists() e.copy( hasReport = done, + statusNote = live[e.file.name]?.label, + progress = live[e.file.name]?.fraction, status = when { done -> FileStatus.DONE e.file.name in failedFiles -> FileStatus.FAILED + live.containsKey(e.file.name) -> FileStatus.RUNNING e.file.name in inFlight -> FileStatus.QUEUED else -> FileStatus.NEW }, @@ -524,14 +562,17 @@ class DesktopState(private val appDir: File = defaultAppDir()) { val ref = provider.upload(draft) {} var jobs: List = emptyList() var polls = 0 - while (polls < 600) { // up to ~30 min per file at 3s polls + while (polls < 600) { // up to ~30 min per file at 2s polls polls++ - delay(3000) + delay(2000) jobs = runCatching { parseJobs(provider.fetchJobs(ref.key)) }.getOrNull().orEmpty() + setLive(f.name, jobProgress(jobs)) + rescanStatuses() val t = jobs.firstOrNull { it.jobType == "transcribe" }?.status val s = jobs.firstOrNull { it.jobType == "summarize" }?.status if ((t == null || t in TERMINAL) && (s == null || s in TERMINAL)) break } + setLive(f.name, null) val transcribeFailed = jobs.any { it.jobType == "transcribe" && it.status == "failed" } @@ -542,6 +583,7 @@ class DesktopState(private val appDir: File = defaultAppDir()) { saveReport(DetailUi(file = f, transcript = transcript, summary = summary)) true } catch (e: Exception) { + setLive(f.name, null) // Surface pump failures for diagnosis instead of swallowing them. runCatching { File(appDir, "pump-error.log").appendText( @@ -561,6 +603,19 @@ class DesktopState(private val appDir: File = defaultAppDir()) { return } val hits = mutableListOf() + // Name matches come first: they find untranscribed recordings too + // (searching "tycos" should surface "…Tycos Space.m4a" immediately, + // whether or not a report exists yet). + _entries.value.forEach { e -> + if (e.file.nameWithoutExtension.lowercase().contains(needle)) { + hits += SearchHit( + report = e.file.nameWithoutExtension, + line = 0, // 0 = name match (no transcript line) + snippet = "recording name match", + audioPath = e.file.absolutePath, + ) + } + } dir.listFiles() ?.filter { it.isFile && it.name.endsWith(".transcript.md") } ?.sortedBy { it.name.lowercase() } @@ -661,10 +716,11 @@ class DesktopState(private val appDir: File = defaultAppDir()) { _detail.value = d.copy() // Poll jobs until the transcribe step reaches a terminal state. while (true) { - delay(3000) + delay(2000) val jobs = runCatching { parseJobs(provider.fetchJobs(ref.key)) }.getOrNull() .orEmpty() d.jobs = jobs + jobLabel(jobs)?.let { d.busy = it } _detail.value = d.copy() val t = jobs.firstOrNull { it.jobType == "transcribe" }?.status if (t == null || t in TERMINAL) break @@ -698,6 +754,26 @@ class DesktopState(private val appDir: File = defaultAppDir()) { } companion object { + /** Live pipeline activity for a recording: human label plus the + * 0..1 fraction when the backend reports one. Null = idle. */ + fun jobProgress(jobs: List): LiveProgress? { + val t = jobs.firstOrNull { it.jobType == "transcribe" } + val s = jobs.firstOrNull { it.jobType == "summarize" } + return when { + t != null && t.status == "running" -> when { + t.stage == "loading-model" -> LiveProgress("loading model…", null) + t.progress != null -> + LiveProgress("transcribing… ${t.progress}%", t.progress / 100f) + else -> LiveProgress("transcribing…", null) + } + s != null && s.status == "running" -> LiveProgress("summarizing…", null) + else -> null + } + } + + /** Label-only convenience for the Detail screen's busy text. */ + fun jobLabel(jobs: List): String? = jobProgress(jobs)?.label + const val KEY_URL = "server.url" const val KEY_FOLDER = "library.folder" const val KEY_PASSWORD = "local.password" @@ -705,6 +781,10 @@ class DesktopState(private val appDir: File = defaultAppDir()) { const val KEY_AUTO = "library.auto_transcribe" /** Preferred local summarization model (Ollama). */ const val DEFAULT_LLM_MODEL = "qwen3:4b" + /** LAN H200 inference server (private 10.x network, never internet). */ + const val LAN_LLM_BASE_URL = "http://10.50.200.200:8100" + const val LAN_LLM_MODEL = "qwen3.8-flash-next" + const val ENV_LAN_LLM_KEY = "HERMES_CUSTOM_10_50_200_200_8100_API_KEY" // Must satisfy the engine's email validation (period in domain). const val LOCAL_EMAIL = "desktop@app.shonar" @@ -724,4 +804,13 @@ class DesktopState(private val appDir: File = defaultAppDir()) { } } -data class SearchHit(val report: String, val line: Int, val snippet: String) +data class SearchHit( + val report: String, + val line: Int, + val snippet: String, + /** Set for name matches so the hit can open the recording. */ + val audioPath: String? = null, +) + +/** What the pipeline is doing for a library file right now. */ +data class LiveProgress(val label: String, val fraction: Float?) diff --git a/desktop/app/src/main/kotlin/com/shonar/desktop/Screens.kt b/desktop/app/src/main/kotlin/com/shonar/desktop/Screens.kt index da1c285..84d82cc 100644 --- a/desktop/app/src/main/kotlin/com/shonar/desktop/Screens.kt +++ b/desktop/app/src/main/kotlin/com/shonar/desktop/Screens.kt @@ -133,12 +133,21 @@ fun LibraryScreen(state: DesktopState) { } else { LazyColumn(Modifier.fillMaxWidth().weight(1f), verticalArrangement = Arrangement.spacedBy(4.dp)) { - items(hits) { h -> - Card(Modifier.fillMaxWidth()) { + items(hits, key = { "${it.audioPath ?: "text"}:${it.report}:${it.line}" }) { h -> + val openable = h.audioPath?.let { File(it).exists() } == true + Card( + Modifier.fillMaxWidth().then( + if (openable) Modifier.clickable { + state.openDetail(File(h.audioPath)) + } else Modifier + ), + ) { Column(Modifier.padding(10.dp)) { Text(h.report, style = MaterialTheme.typography.titleSmall) - Text("line ${h.line}: …${h.snippet}…", - style = MaterialTheme.typography.bodySmall) + Text( + if (h.line == 0) h.snippet else "line ${h.line}: …${h.snippet}…", + style = MaterialTheme.typography.bodySmall, + ) } } } @@ -166,7 +175,10 @@ fun LibraryScreen(state: DesktopState) { when (e.status) { FileStatus.DONE -> "transcript saved" FileStatus.QUEUED -> "queued for transcription" - FileStatus.RUNNING -> "transcribing…" + // statusNote carries the live + // "transcribing… 42%" label. + FileStatus.RUNNING -> + e.statusNote ?: "transcribing…" FileStatus.FAILED -> "failed" + (e.statusNote?.let { " ($it)" } ?: "") FileStatus.NEW -> "not transcribed" @@ -179,6 +191,13 @@ fun LibraryScreen(state: DesktopState) { else -> MaterialTheme.colorScheme.onSurfaceVariant }, ) + if (e.progress != null) { + Spacer(Modifier.height(6.dp)) + LinearProgressIndicator( + progress = { e.progress }, + modifier = Modifier.fillMaxWidth().height(4.dp), + ) + } } if (e.status == FileStatus.NEW || e.status == FileStatus.FAILED) { TextButton({ state.pumpFile(e.file) }) { Text("Transcribe") } diff --git a/desktop/app/src/test/kotlin/com/shonar/desktop/JobProgressTest.kt b/desktop/app/src/test/kotlin/com/shonar/desktop/JobProgressTest.kt new file mode 100644 index 0000000..d600e18 --- /dev/null +++ b/desktop/app/src/test/kotlin/com/shonar/desktop/JobProgressTest.kt @@ -0,0 +1,48 @@ +package com.shonar.desktop + +import com.shonar.recording.JobInfo +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Test + +class JobProgressTest { + + private fun job(type: String, status: String, stage: String? = null, progress: Int? = null) = + JobInfo(jobType = type, status = status, stage = stage, progress = progress) + + @Test fun `idle when nothing running`() { + assertNull(DesktopState.jobProgress(emptyList())) + assertNull(DesktopState.jobProgress(listOf(job("transcribe", "succeeded")))) + } + + @Test fun `loading model has no percentage`() { + val p = DesktopState.jobProgress( + listOf(job("transcribe", "running", stage = "loading-model")), + ) + assertEquals("loading model…", p?.label) + assertNull(p?.fraction) + } + + @Test fun `transcribing shows backend percentage`() { + val p = DesktopState.jobProgress( + listOf(job("transcribe", "running", stage = "transcribing", progress = 42)), + ) + assertEquals("transcribing… 42%", p?.label) + assertEquals(0.42f, p!!.fraction!!, 0.0001f) + } + + @Test fun `transcribe failure falls back to indeterminate`() { + val p = DesktopState.jobProgress( + listOf(job("transcribe", "running", stage = "transcribing")), + ) + assertEquals("transcribing…", p?.label) + assertNull(p?.fraction) + } + + @Test fun `summarizing stage`() { + val p = DesktopState.jobProgress( + listOf(job("transcribe", "succeeded"), job("summarize", "running")), + ) + assertEquals("summarizing…", p?.label) + } +}