api: /reprocess endpoint; summarize failure keeps usable transcript

- POST /recordings/{id}/reprocess?job=transcribe|summarize re-runs one
  stage on demand (409 if running or nothing to summarize)
- _fail(): summarize failure on a recording with a usable transcript
  marks completed with 'Summary failed: ...' instead of failing the
  whole recording
This commit is contained in:
avi 2026-09-13 19:20:33 -05:00
commit 5b4a111960
2 changed files with 69 additions and 2 deletions

View file

@ -331,6 +331,65 @@ async def list_jobs(recording_id: uuid.UUID, user: CurrentUser, session: Session
return list(rows)
@router.post("/recordings/{recording_id}/reprocess", response_model=list[ProcessingJobOut])
async def reprocess_recording(
recording_id: uuid.UUID,
user: CurrentUser,
session: SessionDep,
job: str = Query(default="summarize", pattern="^(transcribe|summarize)$"),
):
"""Force one pipeline stage to run again (Summarize / Re-transcribe).
Unlike the enqueue-on-finalize path, this ignores prior success: a
summary the user wants regenerated (better model, new prompt) is a
deliberate request. Running jobs are left alone (409 instead of a
duplicate).
"""
from shonar.db.models import JobStatus, JobType as JT, ProcessingStatus
from shonar.services import processing as proc
rec = await _owned_recording(session, user, recording_id)
job_type = JT(job)
if job_type is JT.summarize and await proc.latest_transcript_text(
session, rec.id
) is None:
raise HTTPException(409, "Transcribe first — there is nothing to summarize.")
existing = await session.scalar(
select(ProcessingJob)
.where(
ProcessingJob.recording_id == rec.id,
ProcessingJob.job_type == job_type,
)
.order_by(ProcessingJob.id.desc())
)
if existing is not None and existing.status in (
JobStatus.queued,
JobStatus.running,
):
raise HTTPException(409, "That stage is already running.")
if existing is None:
session.add(ProcessingJob(recording_id=rec.id, job_type=job_type))
else:
existing.status = JobStatus.queued
existing.attempt = 0
existing.error = None
existing.stage = None
existing.progress = None
existing.started_at = None
existing.finished_at = None
if job_type is JT.transcribe:
rec.processing_status = ProcessingStatus.processing
rec.processing_error = None
await session.flush()
await proc.transport_enqueue(job_type, rec.id)
rows = await session.scalars(
select(ProcessingJob)
.where(ProcessingJob.recording_id == rec.id)
.order_by(ProcessingJob.id)
)
return list(rows)
# --- M8: user edits (new version, edited_by_user=True; pipeline won't clobber)
def _transcript_out(row: Transcript) -> TranscriptOut:
return TranscriptOut(

View file

@ -286,6 +286,14 @@ async def _fail(
job.error = message
job.stage = None
job.finished_at = utcnow()
if job.job_type == JobType.summarize and await latest_transcript_text(
session, rec.id
):
# The transcript is usable; a summary timeout must not mark the
# whole recording failed (the summary can be re-run separately).
rec.processing_status = ProcessingStatus.completed
rec.processing_error = f"Summary failed: {message}"
else:
rec.processing_status = ProcessingStatus.failed
rec.processing_error = message
await session.flush()