Progress reporting end-to-end + name search + LAN summarizer

Backend:
- faster-whisper reports 0-99 percent (segment end / audio duration,
  throttled to whole percents); processing.py persists it to the job
  row from the worker thread via a best-effort scheduled writer.
- Commit job state (running/stage) before the long CPU/LLM phases so
  readers see it (open tx was invisible + SQLite-locked); transient
  failures reset the row to queued explicitly before raising for retry.
- summarize jobs carry stage=summarizing; both jobs clear it at 100.

Desktop:
- Library rows show live 'transcribing… 42%' with a determinate bar;
  Detail screen busy text updates per poll; poll interval 2s.
- Search matches recording names too (name hits first, openable) —
  finds untranscribed files (the 'tycos' miss).
- Engine prefers the LAN H200 openai_compat summarizer when the key
  resolves (env or ~/.hermes/.env), falls back to local Ollama, then
  none. Key never stored by the app.

Verified: backend 67 pytest green on SQLite and Postgres, ruff clean;
desktop 13 tests green (5 new jobProgress cases); live probe through
the app's engine showed queued→running 33%→100 and summarize running;
H200 summarize call 5.9s vs ~30s local qwen3:4b.
This commit is contained in:
avi 2026-09-13 02:04:08 -05:00
commit 86e18c6b52
7 changed files with 250 additions and 21 deletions

View file

@ -59,6 +59,7 @@ class TranscriptionProvider(Protocol):
mime: str,
*,
language_hint: str | None = None,
on_progress=None, # optional Callable[[int], None], 0..99
) -> TranscriptResult: ...

View file

@ -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,

View file

@ -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 {}

View file

@ -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

View file

@ -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<Map<String, LiveProgress>>(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<JobInfo> = 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<SearchHit>()
// 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<JobInfo>): 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<JobInfo>): 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?)

View file

@ -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") }

View file

@ -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)
}
}