From 61f2124b1cb347671c697d213a4094aca43e1522 Mon Sep 17 00:00:00 2001 From: avi Date: Fri, 18 Sep 2026 15:00:04 -0500 Subject: [PATCH] AI progress: honest milestone bar from closed JSON contract units MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- backend/shonar/api/schemas_recordings.py | 3 + backend/shonar/services/ai/_llm.py | 159 ++++++++++++++++++-- backend/shonar/services/ai/ollama.py | 19 +-- backend/shonar/services/ai/openai_compat.py | 33 ++-- backend/tests/test_ai_adapters.py | 22 ++- shared/com/shonar/recording/AiContent.kt | 3 + 6 files changed, 197 insertions(+), 42 deletions(-) diff --git a/backend/shonar/api/schemas_recordings.py b/backend/shonar/api/schemas_recordings.py index 71327cf..f056f9b 100644 --- a/backend/shonar/api/schemas_recordings.py +++ b/backend/shonar/api/schemas_recordings.py @@ -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. diff --git a/backend/shonar/services/ai/_llm.py b/backend/shonar/services/ai/_llm.py index 7cb04b6..0757a82 100644 --- a/backend/shonar/services/ai/_llm.py +++ b/backend/shonar/services/ai/_llm.py @@ -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 diff --git a/backend/shonar/services/ai/ollama.py b/backend/shonar/services/ai/ollama.py index 7edf1ea..973c963 100644 --- a/backend/shonar/services/ai/ollama.py +++ b/backend/shonar/services/ai/ollama.py @@ -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) diff --git a/backend/shonar/services/ai/openai_compat.py b/backend/shonar/services/ai/openai_compat.py index 49cb1bb..c436490 100644 --- a/backend/shonar/services/ai/openai_compat.py +++ b/backend/shonar/services/ai/openai_compat.py @@ -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.") diff --git a/backend/tests/test_ai_adapters.py b/backend/tests/test_ai_adapters.py index ff295c7..463c2f0 100644 --- a/backend/tests/test_ai_adapters.py +++ b/backend/tests/test_ai_adapters.py @@ -155,10 +155,12 @@ async def test_openai_compat_streams_with_progress(): assert all(0 <= p <= 99 for p in seen_pcts) # never claims 100 early -async def test_openai_compat_thinking_moves_progress_before_content(): - """qwen3-on-llama.cpp streams reasoning_content for most of the run; - progress must climb during that phase, not sit at "waiting for model" - and then jump straight to done.""" +async def test_openai_compat_thinking_earns_only_first_token(): + """qwen3-on-llama.cpp streams reasoning_content for most of the run. + Thinking is real work but NOT contract output: it fires the 5% + first-token milestone once ("the model is alive") and then holds. + Progress only climbs when JSON contract units actually close — a + continuous percent over the thinking span would be a fake estimate.""" import json as _json def sse(request: httpx.Request) -> httpx.Response: @@ -177,10 +179,14 @@ async def test_openai_compat_thinking_moves_progress_before_content(): seen_pcts = [] res = await p.summarize("t", on_progress=seen_pcts.append) assert res.short == "Done." # reasoning never enters the reply - # The ticker must report real mid-run values while only reasoning - # streamed — otherwise the UI shows no percentage for the whole job. - assert len(seen_pcts) > 10 - assert seen_pcts[-1] > 50 + # Thinking earned exactly one tick: the first-token milestone. No + # fake crawl across the reasoning span. + assert seen_pcts[0] == 5 + assert seen_pcts.count(5) == 1 + assert max(seen_pcts) <= 99 # never claims 100 before storage + # Contract units closing does move the bar to the root milestone. + assert seen_pcts[-1] == 95 + assert seen_pcts == sorted(seen_pcts) async def test_openai_compat_partial_json_gets_defaults(): diff --git a/shared/com/shonar/recording/AiContent.kt b/shared/com/shonar/recording/AiContent.kt index 8041d83..e25d654 100644 --- a/shared/com/shonar/recording/AiContent.kt +++ b/shared/com/shonar/recording/AiContent.kt @@ -57,6 +57,8 @@ data class JobInfo( /** Display-only phase ("loading-model", "transcribing") + 0-100 progress. */ val stage: String? = null, val progress: Int? = null, + /** Voice a running summarize job writes in ("dry, witty"); labels the bar. */ + val tone: String? = null, /** ISO-8601 UTC when the job started running (null if unknown). */ val startedAt: String? = null, /** ISO-8601 UTC when the job was queued (null if unknown). The wait @@ -137,6 +139,7 @@ fun parseJobs(raw: String): List = runCatching { error = o.optString("error", null).takeUnless { it.isNullOrBlank() }, stage = o.optString("stage", null).takeUnless { it.isNullOrBlank() }, progress = if (o.has("progress") && !o.isNull("progress")) o.optInt("progress") else null, + tone = o.optString("tone", null).takeUnless { it.isNullOrBlank() }, startedAt = o.optString("started_at", null).takeUnless { it.isNullOrBlank() }, createdAt = o.optString("created_at", null).takeUnless { it.isNullOrBlank() }, )