Instrumentation only: TRACE log lines for the tone-mismatch + progress investigation — reprocess (rid, requested tone, prior tone/status), summarize-start (rid, job, tone the worker actually read), progress (every distinct pct the provider emits), store_summary (rid, tone saved). rid minted per click joins all lines; zero behavior change.
This commit is contained in:
parent
e8bd5623d4
commit
768b75915a
2 changed files with 51 additions and 0 deletions
|
|
@ -7,6 +7,7 @@ return 404 (no existence leaks). Storage keys are never exposed.
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import contextlib
|
import contextlib
|
||||||
|
import logging
|
||||||
import uuid
|
import uuid
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
|
||||||
|
|
@ -42,6 +43,7 @@ from shonar.db.models import (
|
||||||
from shonar.services import uploads as up
|
from shonar.services import uploads as up
|
||||||
|
|
||||||
router = APIRouter(tags=["uploads", "recordings"])
|
router = APIRouter(tags=["uploads", "recordings"])
|
||||||
|
logger = logging.getLogger("shonar.recordings")
|
||||||
|
|
||||||
|
|
||||||
async def _recording_out(session, rec: Recording) -> RecordingOut: # noqa: ANN001
|
async def _recording_out(session, rec: Recording) -> RecordingOut: # noqa: ANN001
|
||||||
|
|
@ -400,6 +402,8 @@ async def reprocess_recording(
|
||||||
JobStatus.running,
|
JobStatus.running,
|
||||||
):
|
):
|
||||||
raise HTTPException(409, "That stage is already running.")
|
raise HTTPException(409, "That stage is already running.")
|
||||||
|
prior_tone = existing.tone if existing is not None else None
|
||||||
|
prior_status = existing.status.value if existing is not None else "new"
|
||||||
if existing is None:
|
if existing is None:
|
||||||
session.add(ProcessingJob(recording_id=rec.id, job_type=job_type,
|
session.add(ProcessingJob(recording_id=rec.id, job_type=job_type,
|
||||||
tone=tone))
|
tone=tone))
|
||||||
|
|
@ -414,6 +418,15 @@ async def reprocess_recording(
|
||||||
# Tone rides the job row (the queue carries only ids): a plain
|
# Tone rides the job row (the queue carries only ids): a plain
|
||||||
# re-summarize must clear a previous tone, not inherit it.
|
# re-summarize must clear a previous tone, not inherit it.
|
||||||
existing.tone = tone
|
existing.tone = tone
|
||||||
|
# TRACE (logging only): one line per click. rid joins this line to the
|
||||||
|
# worker's start/done lines; prior_tone/prior_status expose rapid double
|
||||||
|
# clicks overwriting a queued tone before pickup.
|
||||||
|
logger.info(
|
||||||
|
"TRACE reprocess rid=%s rec=%s job=%s tone=%r prior_tone=%r "
|
||||||
|
"prior_status=%s",
|
||||||
|
proc.trace_request(rec.id, job_type, tone),
|
||||||
|
rec.id, job_type, tone, prior_tone, prior_status,
|
||||||
|
)
|
||||||
if job_type is JT.transcribe:
|
if job_type is JT.transcribe:
|
||||||
if model is not None:
|
if model is not None:
|
||||||
# A re-transcribe may switch models; the saved per-recording
|
# A re-transcribe may switch models; the saved per-recording
|
||||||
|
|
|
||||||
|
|
@ -49,6 +49,27 @@ from shonar.storage import get_storage
|
||||||
|
|
||||||
logger = logging.getLogger("shonar.processing")
|
logger = logging.getLogger("shonar.processing")
|
||||||
|
|
||||||
|
# --- tone/request-integrity trace (pure logging; remove when the
|
||||||
|
# tone-mismatch investigation closes) ---------------------------------------
|
||||||
|
# POST /reprocess mints an rid per request and stashes it here keyed by
|
||||||
|
# (recording, job); the worker and store_summary read it back so all the
|
||||||
|
# lines for one click share the same rid= and can be grepped together.
|
||||||
|
# Single-process inline queue makes this safe; it carries no behavior.
|
||||||
|
_trace_rid: dict = {}
|
||||||
|
|
||||||
|
|
||||||
|
def trace_request(recording_id, job_type, tone):
|
||||||
|
"""Called from POST /reprocess. Returns the rid for the request log."""
|
||||||
|
rid = uuid.uuid4().hex[:8]
|
||||||
|
_trace_rid[(str(recording_id), str(job_type))] = rid
|
||||||
|
logger.info("TRACE reprocess rid=%s rec=%s job=%s tone=%r",
|
||||||
|
rid, recording_id, job_type, tone)
|
||||||
|
return rid
|
||||||
|
|
||||||
|
|
||||||
|
def trace_peek(recording_id, job_type):
|
||||||
|
return _trace_rid.get((str(recording_id), str(job_type)), "-")
|
||||||
|
|
||||||
|
|
||||||
def _accepts_on_progress(provider) -> bool:
|
def _accepts_on_progress(provider) -> bool:
|
||||||
"""True when the provider's transcribe() accepts on_progress=."""
|
"""True when the provider's transcribe() accepts on_progress=."""
|
||||||
|
|
@ -109,6 +130,10 @@ def _async_progress_writer(job_id):
|
||||||
if int(pct) == last:
|
if int(pct) == last:
|
||||||
return
|
return
|
||||||
last = int(pct)
|
last = int(pct)
|
||||||
|
# TRACE (logging only): every distinct progress value the provider
|
||||||
|
# emits, with wall time — proves whether/when the 99 is emitted and
|
||||||
|
# written vs. the UI only ever showing 0%.
|
||||||
|
logger.info("TRACE progress job=%s pct=%s", job_id, pct)
|
||||||
|
|
||||||
async def _write() -> None:
|
async def _write() -> None:
|
||||||
from sqlalchemy import update
|
from sqlalchemy import update
|
||||||
|
|
@ -481,6 +506,13 @@ async def run_summarize(ctx: dict, recording_id: str) -> None:
|
||||||
job.stage = "summarizing"
|
job.stage = "summarizing"
|
||||||
rec.processing_status = ProcessingStatus.processing
|
rec.processing_status = ProcessingStatus.processing
|
||||||
await session.commit() # visible before the long LLM call
|
await session.commit() # visible before the long LLM call
|
||||||
|
# TRACE (logging only): the tone the worker actually reads from the
|
||||||
|
# row NOW — differs from the reprocess line's tone if a later click
|
||||||
|
# overwrote it before pickup.
|
||||||
|
logger.info("TRACE summarize-start rid=%s job=%s rec=%s type=%s "
|
||||||
|
"tone=%r attempt=%s",
|
||||||
|
trace_peek(rec.id, job.job_type), job.id, rec.id,
|
||||||
|
job.job_type, job.tone, job.attempt)
|
||||||
fallback_note: str | None = None
|
fallback_note: str | None = None
|
||||||
result_provider = provider.name # overwritten only on the rescue path
|
result_provider = provider.name # overwritten only on the rescue path
|
||||||
try:
|
try:
|
||||||
|
|
@ -585,6 +617,12 @@ async def store_summary(
|
||||||
model: str,
|
model: str,
|
||||||
tone: str | None = None,
|
tone: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
# TRACE (logging only): the tone arriving at storage. Compare with the
|
||||||
|
# summarize-start line's tone for the same rid: a mismatch here means
|
||||||
|
# the value changed between worker pickup and the DB write.
|
||||||
|
logger.info("TRACE store_summary rid=%s rec=%s tone=%r provider=%s",
|
||||||
|
trace_peek(rec.id, "summarize"), rec.id, tone,
|
||||||
|
provider_name)
|
||||||
existing = (
|
existing = (
|
||||||
await session.scalars(
|
await session.scalars(
|
||||||
select(Summary)
|
select(Summary)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue