AI progress: honest milestone bar from closed JSON contract units
Replace the char-count estimate (EST_OUTPUT_CHARS) with ContractProgress:
points are earned only when real output units complete — first token (5),
each closed contract key / array item (5-90), root closed (95), stored (100).
Reasoning/thinking streams no longer fake progress: the first thinking
delta fires the first-token milestone ('the model is alive') and holds.
Surface the job's tone through ProcessingJobOut/JobInfo so the bar can
label itself 'Summarizing dry wit…'. Test rewritten to pin the new
honest-thinking semantics.
This commit is contained in:
parent
6ec811393a
commit
61f2124b1c
6 changed files with 197 additions and 42 deletions
|
|
@ -127,6 +127,9 @@ class ProcessingJobOut(ORMModel):
|
|||
# "transcribing") plus 0-100 progress when known.
|
||||
stage: str | None
|
||||
progress: int | None
|
||||
# The voice a running summarize job is writing in, so the UI can say
|
||||
# "Summarizing dry wit…" instead of a generic label.
|
||||
tone: str | None = None
|
||||
started_at: datetime | None
|
||||
finished_at: datetime | None
|
||||
# Queue wait is visible time too: the UI counts it for queued jobs.
|
||||
|
|
|
|||
|
|
@ -38,26 +38,161 @@ def build_user_message(transcript: str, title: str | None,
|
|||
return msg
|
||||
|
||||
|
||||
# Rough typical length of a summary JSON reply. The stream cannot know the
|
||||
# model's final size, so progress is an estimate that crawls toward 99 and
|
||||
# the job sets 100 on success — honest "almost there", never a fake jump.
|
||||
EST_OUTPUT_CHARS = 1200
|
||||
# Progress from REAL completed work units only — never a timer, never an
|
||||
# elapsed-time estimate, never a char-count guess against a made-up output
|
||||
# size. The summary contract is a fixed JSON shape (six keys), and JSON
|
||||
# structure is parseable from the stream prefix, so every point the bar
|
||||
# advances corresponds to a unit of output the model has actually
|
||||
# finished: the first token, each closed contract key, each closed array
|
||||
# item, the closing brace. The un-finishable part — the thinking phase
|
||||
# and the interior of an open value — honestly earns zero points: no
|
||||
# provider (llama.cpp, vLLM, Ollama) exposes a token target while
|
||||
# streaming, so any "continuous" percent over that span would be fake.
|
||||
# Milestones: 0 queued · 5 first token · 5→90 six contract keys
|
||||
# (arrays advance per closed item) · 95 root closed/parsed · 100 stored.
|
||||
_CONTRACT_FIRST_TOKEN = 5
|
||||
_CONTRACT_ROOT = 95
|
||||
# Fixed per-key budget (5..90 band). detailed is contractually the
|
||||
# longest field; arrays share a band with per-item steps inside.
|
||||
_KEY_POINTS_BUDGET = {
|
||||
"short": 12,
|
||||
"detailed": 20,
|
||||
"key_points": 13,
|
||||
"decisions": 13,
|
||||
"action_items": 13,
|
||||
"questions": 13,
|
||||
}
|
||||
_ARR_OPEN, _ARR_ITEM, _ARR_ITEM_MAX, _ARR_CLOSE = 3, 2, 7, 4
|
||||
|
||||
|
||||
class ContractProgress:
|
||||
"""Streaming-JSON milestone tracker. feed(text_prefix) accepts the
|
||||
accumulated reply content (repeated prefixes are fine — scanning
|
||||
resumes where it left off) and returns a new pct, or None when
|
||||
nothing new completed. Monotonic by construction."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._pos = 0 # scan offset into the prefix
|
||||
self._depth = 0 # {} and [] nesting
|
||||
self._in_str = False
|
||||
self._esc = False
|
||||
self._str_buf: list[str] = []
|
||||
self._cur_key: str | None = None # key whose value we're in
|
||||
self._at_key_slot = True # next depth-1 string is a key
|
||||
self._items_open = 0 # items closed in current array
|
||||
self._key_score = 0 # points earned for current key
|
||||
self._earned = 0 # points above the first-token step
|
||||
self._done_keys: set[str] = set()
|
||||
self._reported = 0
|
||||
self._root_closed = False
|
||||
|
||||
def _emit(self, pct: int) -> int | None:
|
||||
pct = min(99, max(0, pct))
|
||||
if pct > self._reported:
|
||||
self._reported = pct
|
||||
return pct
|
||||
return None
|
||||
|
||||
def _award(self, pts: int, cap: int) -> None:
|
||||
"""Add points to the current key, capped at that key's budget."""
|
||||
grant = min(pts, cap - self._key_score)
|
||||
if grant > 0:
|
||||
self._key_score += grant
|
||||
self._earned += grant
|
||||
|
||||
def feed(self, prefix: str) -> int | None:
|
||||
new = self._emit(_CONTRACT_FIRST_TOKEN) if self._reported == 0 else None
|
||||
result: int | None = new
|
||||
text = prefix
|
||||
while self._pos < len(text):
|
||||
c = text[self._pos]
|
||||
self._pos += 1
|
||||
if self._in_str:
|
||||
if self._esc:
|
||||
self._esc = False
|
||||
elif c == "\\":
|
||||
self._esc = True
|
||||
elif c == '"':
|
||||
self._in_str = False
|
||||
s = "".join(self._str_buf)
|
||||
self._str_buf = []
|
||||
if self._depth == 1:
|
||||
# Depth-1 strings alternate key, value, key, ...
|
||||
# Slot tracking (not lookahead) so a key closing
|
||||
# one chunk before its ':' is not misread.
|
||||
if self._at_key_slot:
|
||||
self._cur_key = s
|
||||
self._key_score = 0
|
||||
self._at_key_slot = False
|
||||
elif self._cur_key in ("short", "detailed"):
|
||||
self._award(_KEY_POINTS_BUDGET[self._cur_key],
|
||||
_KEY_POINTS_BUDGET[self._cur_key])
|
||||
self._done_keys.add(self._cur_key)
|
||||
self._cur_key = None
|
||||
result = self._emit(
|
||||
_CONTRACT_FIRST_TOKEN + self._earned) or result
|
||||
self._at_key_slot = True
|
||||
elif self._depth == 2:
|
||||
# a closed array item (arrays hold strings)
|
||||
if self._cur_key in _KEY_POINTS_BUDGET:
|
||||
self._items_open += 1
|
||||
remaining = _ARR_ITEM_MAX \
|
||||
- _ARR_ITEM * (self._items_open - 1)
|
||||
self._award(min(_ARR_ITEM, max(0, remaining)),
|
||||
_KEY_POINTS_BUDGET[self._cur_key])
|
||||
result = self._emit(
|
||||
_CONTRACT_FIRST_TOKEN + self._earned) or result
|
||||
else:
|
||||
self._str_buf.append(c)
|
||||
continue
|
||||
if c == '"':
|
||||
self._in_str = True
|
||||
self._str_buf = []
|
||||
elif c in "{[":
|
||||
self._depth += 1
|
||||
if c == "[" and self._cur_key:
|
||||
self._items_open = 0
|
||||
self._award(_ARR_OPEN, _KEY_POINTS_BUDGET[self._cur_key])
|
||||
result = self._emit(
|
||||
_CONTRACT_FIRST_TOKEN + self._earned) or result
|
||||
elif c in "}]":
|
||||
self._depth -= 1
|
||||
if self._depth == 1:
|
||||
if c == "]":
|
||||
if self._cur_key in _KEY_POINTS_BUDGET:
|
||||
self._award(_ARR_CLOSE,
|
||||
_KEY_POINTS_BUDGET[self._cur_key])
|
||||
self._done_keys.add(self._cur_key)
|
||||
self._cur_key = None
|
||||
result = self._emit(
|
||||
_CONTRACT_FIRST_TOKEN + self._earned) or result
|
||||
self._at_key_slot = True
|
||||
elif self._depth == 0:
|
||||
self._root_closed = True
|
||||
result = self._emit(_CONTRACT_ROOT) or result
|
||||
elif c == "," and self._depth == 1:
|
||||
self._at_key_slot = True
|
||||
return result
|
||||
|
||||
|
||||
def make_progress_ticker(on_progress):
|
||||
"""Wrap an optional on_progress callback into a feed(n_chars) sink.
|
||||
"""Wrap an optional on_progress callback into a feed(prefix) sink.
|
||||
|
||||
Reports a clamped 0..99 estimate from streamed output size. Silent
|
||||
no-op when the caller has no callback; failures never disturb the
|
||||
summary itself."""
|
||||
Receives the accumulated reply *content* (not reasoning) as it
|
||||
streams and reports 5..95 strictly from completed contract units —
|
||||
see ContractProgress for what each point means. Silent no-op without
|
||||
a callback; failures never disturb the summary itself."""
|
||||
if on_progress is None:
|
||||
return lambda n_chars: None
|
||||
return lambda prefix: None
|
||||
tracker = ContractProgress()
|
||||
|
||||
def feed(n_chars: int) -> None:
|
||||
def feed(prefix: str) -> None:
|
||||
try: # noqa: SIM105 — swallow deliberately: progress is display state
|
||||
on_progress(min(99, n_chars * 100 // EST_OUTPUT_CHARS))
|
||||
pct = tracker.feed(prefix or "")
|
||||
except Exception: # display state only
|
||||
pass
|
||||
return
|
||||
if pct is not None:
|
||||
on_progress(pct)
|
||||
|
||||
return feed
|
||||
|
||||
|
|
|
|||
|
|
@ -93,7 +93,8 @@ async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
|||
|
||||
Streaming answers are NDJSON (one JSON object per line, ``done`` on the
|
||||
last); a server honoring stream=false returns one JSON body — both work.
|
||||
The running character count feeds [ticker] for UI progress."""
|
||||
The accumulated content *prefix* feeds [ticker], which derives progress
|
||||
from completed JSON contract units only (see ContractProgress)."""
|
||||
import json as _json
|
||||
|
||||
ctype = resp.headers.get("content-type", "")
|
||||
|
|
@ -103,11 +104,10 @@ async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
|||
content = _json.loads(body)["message"]["content"]
|
||||
except (ValueError, KeyError, TypeError) as e:
|
||||
raise ProviderTransientError("Ollama sent an unreadable reply.") from e
|
||||
ticker(len(content or ""))
|
||||
ticker(content or "")
|
||||
return content or ""
|
||||
|
||||
parts: list[str] = []
|
||||
total = 0
|
||||
async for line in resp.aiter_lines():
|
||||
line = line.strip()
|
||||
if not line:
|
||||
|
|
@ -117,14 +117,15 @@ async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
|||
except ValueError:
|
||||
continue
|
||||
piece = (obj.get("message") or {}).get("content") or ""
|
||||
# Like the openai_compat adapter: a thinking channel is real
|
||||
# work and must move the progress ticker even when the server
|
||||
# ignored think:false.
|
||||
# Like the openai_compat adapter: a thinking channel means the
|
||||
# model started (the ticker's first-token milestone) but is not
|
||||
# contract output, so it moves nothing beyond that.
|
||||
thinking = (obj.get("message") or {}).get("thinking") or ""
|
||||
if piece or thinking:
|
||||
if piece:
|
||||
parts.append(piece)
|
||||
total += len(piece) + len(thinking)
|
||||
ticker(total)
|
||||
ticker("".join(parts))
|
||||
elif thinking:
|
||||
ticker("")
|
||||
if obj.get("done"):
|
||||
break
|
||||
content = "".join(parts)
|
||||
|
|
|
|||
|
|
@ -21,9 +21,12 @@ from shonar.services.ai._llm import (
|
|||
async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
||||
"""Reassemble the assistant reply from a (possibly streamed) response.
|
||||
|
||||
Feeds the running character count to [ticker] as deltas arrive so the
|
||||
UI can show progress. Servers that ignored "stream": true answer with
|
||||
a plain JSON body — that path is handled too."""
|
||||
Feeds the accumulated content *prefix* to [ticker] as deltas arrive;
|
||||
the ticker derives progress from completed JSON contract units only
|
||||
(see ContractProgress). Reasoning deltas are real work but not
|
||||
contract output — they stream before the JSON and intentionally do
|
||||
not move the bar. Servers that ignored "stream": true answer with a
|
||||
plain JSON body — that path is handled too."""
|
||||
import json as _json
|
||||
|
||||
ctype = resp.headers.get("content-type", "")
|
||||
|
|
@ -33,11 +36,10 @@ async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
|||
content = _json.loads(body)["choices"][0]["message"]["content"]
|
||||
except (ValueError, KeyError, IndexError, TypeError) as e:
|
||||
raise ProviderTransientError("LLM sent an unreadable reply.") from e
|
||||
ticker(len(content or ""))
|
||||
ticker(content or "")
|
||||
return content or ""
|
||||
|
||||
parts: list[str] = []
|
||||
total = 0
|
||||
async for line in resp.aiter_lines():
|
||||
if not line.startswith("data:"):
|
||||
continue
|
||||
|
|
@ -49,18 +51,23 @@ async def _collect_reply(resp: httpx.Response, ticker) -> str:
|
|||
delta = chunk["choices"][0].get("delta") or {}
|
||||
piece = delta.get("content") or ""
|
||||
# Thinking models (qwen3 on llama.cpp) stream a long
|
||||
# reasoning_content channel BEFORE any content: counting only
|
||||
# content left the progress bar frozen at "waiting for model"
|
||||
# for the entire run, then jumping straight to done. Reasoning
|
||||
# is real work — count it toward progress (never into the
|
||||
# reply itself).
|
||||
# reasoning_content channel BEFORE any content. It is real
|
||||
# work but not contract output: counting it made the bar
|
||||
# rocket to 99 while the JSON had not even started. The
|
||||
# honest signal is the JSON itself completing.
|
||||
thinking = delta.get("reasoning_content") or ""
|
||||
except (ValueError, KeyError, IndexError, TypeError):
|
||||
continue # keep-alives / usage chunks / odd frames: not content
|
||||
if piece or thinking:
|
||||
if piece:
|
||||
parts.append(piece)
|
||||
total += len(piece) + len(thinking)
|
||||
ticker(total)
|
||||
ticker("".join(parts))
|
||||
elif thinking:
|
||||
# First thinking delta proves the model started working:
|
||||
# the ticker's 5% "first token" milestone fires once (an
|
||||
# empty prefix is a no-op after that). The thinking phase
|
||||
# itself earns no further points — it is not contract
|
||||
# output — but "the model is alive" is real information.
|
||||
ticker("")
|
||||
content = "".join(parts)
|
||||
if not content:
|
||||
raise ProviderTransientError("LLM sent an unreadable reply.")
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue