Adopt engine jobs at startup and on 'already running'; live work outranks a saved report
This commit is contained in:
parent
7fa99e8291
commit
ba54cb853b
6 changed files with 132 additions and 9 deletions
|
|
@ -192,6 +192,7 @@ class DesktopState(private val appDir: File = defaultAppDir()) {
|
||||||
_connected.value = true
|
_connected.value = true
|
||||||
refreshModels()
|
refreshModels()
|
||||||
if (_screen.value == Screen.ENGINE) _screen.value = Screen.LIBRARY
|
if (_screen.value == Screen.ENGINE) _screen.value = Screen.LIBRARY
|
||||||
|
adoptEngineJobs()
|
||||||
if (_autoTranscribe.value) pumpNewFiles()
|
if (_autoTranscribe.value) pumpNewFiles()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -599,9 +600,12 @@ class DesktopState(private val appDir: File = defaultAppDir()) {
|
||||||
statusNote = live[e.file.name]?.label,
|
statusNote = live[e.file.name]?.label,
|
||||||
progress = live[e.file.name]?.fraction,
|
progress = live[e.file.name]?.fraction,
|
||||||
status = when {
|
status = when {
|
||||||
|
// Live engine work outranks a saved report: a
|
||||||
|
// re-summarize on a finished file must still show
|
||||||
|
// its bar (the report on disk is simply stale).
|
||||||
|
live.containsKey(e.file.name) -> FileStatus.RUNNING
|
||||||
done -> FileStatus.DONE
|
done -> FileStatus.DONE
|
||||||
e.file.name in failedFiles -> FileStatus.FAILED
|
e.file.name in failedFiles -> FileStatus.FAILED
|
||||||
live.containsKey(e.file.name) -> FileStatus.RUNNING
|
|
||||||
e.file.name in inFlight -> FileStatus.QUEUED
|
e.file.name in inFlight -> FileStatus.QUEUED
|
||||||
else -> FileStatus.NEW
|
else -> FileStatus.NEW
|
||||||
},
|
},
|
||||||
|
|
@ -668,6 +672,70 @@ class DesktopState(private val appDir: File = defaultAppDir()) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Startup adoption: the engine's inline queue survives independently
|
||||||
|
* of the app — restart only the app mid-transcribe (or mid-re-summarize)
|
||||||
|
* and the job keeps running while the UI shows nothing (and auto-queue
|
||||||
|
* would re-upload a duplicate). Scan mapped files (mapping sidecar =
|
||||||
|
* ever uploaded, so a job may exist), fetch their engine jobs, and
|
||||||
|
* adopt the ones with active work: live row progress now, report when
|
||||||
|
* the jobs settle. No re-upload — the mapping already points at the
|
||||||
|
* engine's recording. Files with an existing report are scanned too:
|
||||||
|
* a running re-summarize must not be orphaned by a report on disk.
|
||||||
|
*
|
||||||
|
* Files the pump owns (in flight this session) are skipped: the pump
|
||||||
|
* is their owner and its upload must not be duplicated.
|
||||||
|
*/
|
||||||
|
private fun adoptEngineJobs() {
|
||||||
|
val dir = _folder.value ?: return
|
||||||
|
scope.launch {
|
||||||
|
val files = dir.listFiles()?.filter {
|
||||||
|
LibraryQueue.isAudioCandidate(it, AUDIO_EXTS.keys) &&
|
||||||
|
loadMapping(it) != null &&
|
||||||
|
it.name !in inFlight && !isRunning(it.name)
|
||||||
|
}.orEmpty()
|
||||||
|
for (f in files) {
|
||||||
|
val remoteId = loadMapping(f)?.recordingId ?: continue
|
||||||
|
val jobs = runCatching { parseJobs(provider.fetchJobs(remoteId)) }
|
||||||
|
.getOrNull() ?: continue
|
||||||
|
if (!LibraryQueue.hasActiveWork(jobs)) continue
|
||||||
|
setLive(f.name, jobProgress(jobs))
|
||||||
|
rescanStatuses()
|
||||||
|
adoptOne(f, remoteId)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Watch one adopted job to a terminal state, then write the report.
|
||||||
|
* Same budget as the pump (~30 min at 2s polls); a job that outlives
|
||||||
|
* it just stays shown as live until a later rescan settles it. If the
|
||||||
|
* auto-queue pump takes the file over mid-watch it owns completion —
|
||||||
|
* bail out so the report is written exactly once. */
|
||||||
|
private suspend fun adoptOne(f: File, remoteId: String) {
|
||||||
|
var jobs: List<JobInfo> = emptyList()
|
||||||
|
var polls = 0
|
||||||
|
while (polls < 600 && _connected.value && f.name !in inFlight) {
|
||||||
|
polls++
|
||||||
|
delay(2000)
|
||||||
|
jobs = runCatching { parseJobs(provider.fetchJobs(remoteId)) }
|
||||||
|
.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
|
||||||
|
}
|
||||||
|
if (f.name in inFlight) return // the pump finished it; its report stands
|
||||||
|
setLive(f.name, null)
|
||||||
|
val transcript = runCatching { parseTranscript(provider.fetchTranscript(remoteId)) }
|
||||||
|
.getOrNull()
|
||||||
|
val summary = runCatching { parseSummary(provider.fetchSummary(remoteId)) }.getOrNull()
|
||||||
|
if (transcript != null || summary != null) {
|
||||||
|
saveReport(DetailUi(file = f, transcript = transcript, summary = summary))
|
||||||
|
}
|
||||||
|
rescanStatuses()
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Rename a library file (extension fixed) and carry its report along.
|
* Rename a library file (extension fixed) and carry its report along.
|
||||||
* Refuses while the file is queued/transcribing so the pump can't
|
* Refuses while the file is queued/transcribing so the pump can't
|
||||||
|
|
@ -1173,6 +1241,13 @@ class DesktopState(private val appDir: File = defaultAppDir()) {
|
||||||
d0.remoteId = null
|
d0.remoteId = null
|
||||||
d0.busy = null
|
d0.busy = null
|
||||||
d0.error = "That server recording no longer exists — press Transcribe to upload again."
|
d0.error = "That server recording no longer exists — press Transcribe to upload again."
|
||||||
|
} else if (it is ProviderError.Transient &&
|
||||||
|
it.message?.contains("already running") == true) {
|
||||||
|
// The engine already works this stage (earlier click,
|
||||||
|
// auto-queue, or a job that survived an app restart):
|
||||||
|
// that IS the work the user asked for — attach to it
|
||||||
|
// and show live progress instead of an error.
|
||||||
|
d0.error = null
|
||||||
} else {
|
} else {
|
||||||
d0.busy = null
|
d0.busy = null
|
||||||
d0.error = it.message ?: "Reprocess failed."
|
d0.error = it.message ?: "Reprocess failed."
|
||||||
|
|
@ -1254,7 +1329,12 @@ class DesktopState(private val appDir: File = defaultAppDir()) {
|
||||||
LiveProgress("transcribing… ${t.progress}%", t.progress / 100f)
|
LiveProgress("transcribing… ${t.progress}%", t.progress / 100f)
|
||||||
else -> LiveProgress("transcribing…", null)
|
else -> LiveProgress("transcribing…", null)
|
||||||
}
|
}
|
||||||
s != null && s.status == "running" -> LiveProgress("summarizing…", null)
|
s != null && s.status == "running" ->
|
||||||
|
if (s.progress != null)
|
||||||
|
LiveProgress("summarizing… ${s.progress}%", s.progress / 100f)
|
||||||
|
// No tokens yet: the LLM is cold or the server is busy
|
||||||
|
// queueing. Saying so beats a bar that looks frozen.
|
||||||
|
else LiveProgress("summarizing… waiting for model", null)
|
||||||
else -> null
|
else -> null
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -132,6 +132,12 @@ object LibraryQueue {
|
||||||
): Pair<List<File>, List<File>> =
|
): Pair<List<File>, List<File>> =
|
||||||
selected.partition { it.name !in runningNames }
|
selected.partition { it.name !in runningNames }
|
||||||
|
|
||||||
|
/** Does the engine still have work for this recording (queued or
|
||||||
|
* running)? Startup adoption uses it to decide whether a mapped file
|
||||||
|
* is being processed without the app knowing. */
|
||||||
|
fun hasActiveWork(jobs: List<com.shonar.recording.JobInfo>): Boolean =
|
||||||
|
jobs.any { it.status == "queued" || it.status == "running" }
|
||||||
|
|
||||||
/** Audio candidates: regular files, known extension, not hidden. */
|
/** Audio candidates: regular files, known extension, not hidden. */
|
||||||
fun isAudioCandidate(f: File, audioExts: Set<String>): Boolean =
|
fun isAudioCandidate(f: File, audioExts: Set<String>): Boolean =
|
||||||
f.isFile && !f.name.startsWith(".") &&
|
f.isFile && !f.name.startsWith(".") &&
|
||||||
|
|
|
||||||
|
|
@ -694,13 +694,27 @@ fun DetailScreen(state: DesktopState) {
|
||||||
}
|
}
|
||||||
val sJob = detail.jobs.firstOrNull { it.jobType == "summarize" }
|
val sJob = detail.jobs.firstOrNull { it.jobType == "summarize" }
|
||||||
sJob?.takeIf { it.status == "running" || it.status == "queued" }?.let {
|
sJob?.takeIf { it.status == "running" || it.status == "queued" }?.let {
|
||||||
val label = if (it.status == "running") "Summarizing…" else "Summary queued"
|
val label = when {
|
||||||
|
it.status == "running" && it.progress != null ->
|
||||||
|
"Summarizing… ${it.progress}%"
|
||||||
|
it.status == "running" -> "Summarizing… waiting for model"
|
||||||
|
else -> "Summary queued"
|
||||||
|
}
|
||||||
Text(label, style = MaterialTheme.typography.bodyMedium,
|
Text(label, style = MaterialTheme.typography.bodyMedium,
|
||||||
color = MaterialTheme.colorScheme.primary)
|
color = MaterialTheme.colorScheme.primary)
|
||||||
Spacer(Modifier.height(4.dp))
|
if (it.status == "running") {
|
||||||
LinearProgressIndicator(
|
Spacer(Modifier.height(4.dp))
|
||||||
modifier = Modifier.fillMaxWidth().height(4.dp),
|
if (it.progress != null) {
|
||||||
)
|
LinearProgressIndicator(
|
||||||
|
progress = { it.progress / 100f },
|
||||||
|
modifier = Modifier.fillMaxWidth().height(4.dp),
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
LinearProgressIndicator(
|
||||||
|
modifier = Modifier.fillMaxWidth().height(4.dp),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
detail.error?.let {
|
detail.error?.let {
|
||||||
Text(it, color = MaterialTheme.colorScheme.error)
|
Text(it, color = MaterialTheme.colorScheme.error)
|
||||||
|
|
@ -819,7 +833,7 @@ fun DetailScreen(state: DesktopState) {
|
||||||
append(when {
|
append(when {
|
||||||
sStatus == "queued" -> "Queued…"
|
sStatus == "queued" -> "Queued…"
|
||||||
sProg != null && sProg > 0 -> "Summarizing… $sProg%"
|
sProg != null && sProg > 0 -> "Summarizing… $sProg%"
|
||||||
else -> "Summarizing…"
|
else -> "Summarizing… waiting for model"
|
||||||
})
|
})
|
||||||
sumTone?.let { append(" · $it") }
|
sumTone?.let { append(" · $it") }
|
||||||
if (sAttempt > 0)
|
if (sAttempt > 0)
|
||||||
|
|
|
||||||
|
|
@ -43,6 +43,16 @@ class JobProgressTest {
|
||||||
val p = DesktopState.jobProgress(
|
val p = DesktopState.jobProgress(
|
||||||
listOf(job("transcribe", "succeeded"), job("summarize", "running")),
|
listOf(job("transcribe", "succeeded"), job("summarize", "running")),
|
||||||
)
|
)
|
||||||
assertEquals("summarizing…", p?.label)
|
assertEquals("summarizing… waiting for model", p?.label)
|
||||||
|
assertNull(p?.fraction)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test fun `summarizing shows backend percentage`() {
|
||||||
|
val p = DesktopState.jobProgress(
|
||||||
|
listOf(job("transcribe", "succeeded"),
|
||||||
|
job("summarize", "running", progress = 37)),
|
||||||
|
)
|
||||||
|
assertEquals("summarizing… 37%", p?.label)
|
||||||
|
assertEquals(0.37f, p!!.fraction!!, 0.0001f)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -183,4 +183,17 @@ class TranscribeBatchTest {
|
||||||
assertEquals(b, LibraryQueue.batchRemove(b, "z.m4a"))
|
assertEquals(b, LibraryQueue.batchRemove(b, "z.m4a"))
|
||||||
assertEquals(null, LibraryQueue.batchRemove(null, "a.m4a"))
|
assertEquals(null, LibraryQueue.batchRemove(null, "a.m4a"))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test fun activeWorkDetectionForJobAdoption() {
|
||||||
|
fun job(type: String, status: String) =
|
||||||
|
com.shonar.recording.JobInfo(jobType = type, status = status)
|
||||||
|
assertTrue(LibraryQueue.hasActiveWork(listOf(job("transcribe", "queued"))))
|
||||||
|
assertTrue(LibraryQueue.hasActiveWork(listOf(job("transcribe", "running"))))
|
||||||
|
assertTrue(LibraryQueue.hasActiveWork(
|
||||||
|
listOf(job("transcribe", "succeeded"), job("summarize", "queued"))))
|
||||||
|
// Fully settled or empty: nothing to adopt.
|
||||||
|
assertTrue(!LibraryQueue.hasActiveWork(emptyList()))
|
||||||
|
assertTrue(!LibraryQueue.hasActiveWork(
|
||||||
|
listOf(job("transcribe", "succeeded"), job("summarize", "failed"))))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
0
backend/shonar/services/ai/.hermes-tmp.94QlWo
Normal file
0
backend/shonar/services/ai/.hermes-tmp.94QlWo
Normal file
Loading…
Add table
Add a link
Reference in a new issue