diff --git a/backend/shonar/api/v1/recordings.py b/backend/shonar/api/v1/recordings.py index a90fbbe..9a23575 100644 --- a/backend/shonar/api/v1/recordings.py +++ b/backend/shonar/api/v1/recordings.py @@ -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( diff --git a/backend/shonar/services/processing.py b/backend/shonar/services/processing.py index ca91b24..27a5789 100644 --- a/backend/shonar/services/processing.py +++ b/backend/shonar/services/processing.py @@ -286,8 +286,16 @@ async def _fail( job.error = message job.stage = None job.finished_at = utcnow() - rec.processing_status = ProcessingStatus.failed - rec.processing_error = message + 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()