fix(chat): de-duplicate the double turn execution at its source
The UI POSTs the SSE stream and, when nothing streams to the browser (iOS can't read a fetch-stream body → the fetch throws in ~1s), falls back to the blocking endpoint. But the server-side stream runs to completion regardless, so BOTH turns executed — double-persisting the message and (once logging became guaranteed) double-logging the hand. Make a turn idempotent instead of chasing why the client bails: the first request for a (session, message) owns it; a concurrent duplicate waits on the owner's Event and reuses its reply rather than running a second full turn. respond and respond_stream both claim/await; a finally always releases waiters. Short window so a genuine later resend still runs fresh. Verified with a threaded race: two simultaneous calls, body runs once, both get the same reply. Also fixes the duplicate user-message persistence (the same double-execution) that was polluting reconstructed history. 6 dedup tests + concurrency check; suite 232. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+154
-85
@@ -10,11 +10,55 @@ deliberate) and hands back a ready message list + the active mode. Then:
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
import time
|
||||
|
||||
from lyra import config, llm, logbus, memory, mind, modes, poker_prompts, summary
|
||||
from lyra import tools as toolkit
|
||||
from lyra.llm import Backend
|
||||
|
||||
MAX_TOOL_ROUNDS = 5 # cap tool-call iterations per turn
|
||||
|
||||
# --- turn de-duplication --------------------------------------------------
|
||||
# The web UI hits TWO endpoints for one message: it POSTs the SSE stream, and if
|
||||
# nothing streams to the browser (iOS can't read a fetch stream body → the fetch
|
||||
# throws in ~1s) it falls back to the blocking endpoint. But the server-side stream
|
||||
# runs to completion regardless, so BOTH turns execute — double-persisting the
|
||||
# message and (once logging became guaranteed) double-logging the hand. This guard
|
||||
# makes a turn idempotent: the first request for a given (session, message) owns it;
|
||||
# a duplicate that arrives while it's in flight waits for and reuses that result
|
||||
# instead of running a second full turn. Window is short so a genuine re-send later
|
||||
# still runs fresh.
|
||||
_TURN_TTL = 20.0
|
||||
_turn_lock = threading.Lock()
|
||||
_turns: dict[tuple, dict] = {} # (session_id, msg) -> {event, reply, ts}
|
||||
|
||||
|
||||
def _claim_turn(session_id: str, user_msg: str):
|
||||
"""(is_owner, rec). Owner executes the turn then calls _finish_turn; a non-owner
|
||||
(the near-simultaneous duplicate) waits on rec['event'] and reuses rec['reply']."""
|
||||
key = (session_id, (user_msg or "").strip())
|
||||
now = time.monotonic()
|
||||
with _turn_lock:
|
||||
for k in [k for k, r in _turns.items() if now - r["ts"] > _TURN_TTL]:
|
||||
del _turns[k]
|
||||
rec = _turns.get(key)
|
||||
if rec is not None:
|
||||
return False, rec
|
||||
rec = {"event": threading.Event(), "reply": None, "ts": now}
|
||||
_turns[key] = rec
|
||||
return True, rec
|
||||
|
||||
|
||||
def _finish_turn(rec: dict, reply: str) -> None:
|
||||
rec["reply"] = reply
|
||||
rec["ts"] = time.monotonic()
|
||||
rec["event"].set()
|
||||
|
||||
|
||||
def _await_duplicate(rec: dict) -> str:
|
||||
rec["event"].wait(timeout=_TURN_TTL)
|
||||
return rec["reply"] or _TANGLED
|
||||
# Which backends get function-calling tools is config-driven (cfg.tool_backends,
|
||||
# env TOOL_BACKENDS, default "cloud"). The MI50's llama.cpp server only does tools
|
||||
# when launched with --jinja + a tool-capable model, else it 500s on the tools
|
||||
@@ -133,26 +177,37 @@ def respond(session_id: str, user_msg: str, backend: Backend = "cloud",
|
||||
logbus.log("info", "chat request", session=session_id, backend=backend,
|
||||
model=model, embed=cfg.embed_backend)
|
||||
|
||||
turn = mind.assemble(session_id, user_msg, backend, model)
|
||||
messages = turn.messages
|
||||
tool_specs = toolkit.specs(turn.mode.tools) if backend in cfg.tool_backends else None
|
||||
ctx = {"session_id": session_id, "backend": backend}
|
||||
# A concurrent duplicate (the UI's stream + blocking fallback for one message)
|
||||
# reuses the owner's result instead of running a second full turn.
|
||||
is_owner, rec = _claim_turn(session_id, user_msg)
|
||||
if not is_owner:
|
||||
logbus.log("info", "duplicate turn deduped", session=session_id, path="respond")
|
||||
return _await_duplicate(rec)
|
||||
|
||||
# Persist the user turn before the tool loop so its timestamp precedes any
|
||||
# tool events fired mid-turn (keeps the transcript export in true order).
|
||||
memory.remember(session_id, "user", user_msg)
|
||||
reply, tools_run = _mind_loop(messages, backend, model, tool_specs, ctx, session_id)
|
||||
_ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run, backend, model, ctx, session_id)
|
||||
mouth = _mouth_target(cfg, backend, model)
|
||||
if mouth and reply:
|
||||
reply = _voice_pass(messages, reply, *mouth)
|
||||
if not reply:
|
||||
reply = _TANGLED
|
||||
logbus.log("info", "reply", session=session_id, chars=len(reply), voiced=bool(mouth))
|
||||
reply = _TANGLED
|
||||
try:
|
||||
turn = mind.assemble(session_id, user_msg, backend, model)
|
||||
messages = turn.messages
|
||||
tool_specs = toolkit.specs(turn.mode.tools) if backend in cfg.tool_backends else None
|
||||
ctx = {"session_id": session_id, "backend": backend}
|
||||
|
||||
memory.remember(session_id, "assistant", reply)
|
||||
summary.maybe_summarize_async(session_id) # compact once enough new turns pile up
|
||||
return reply
|
||||
# Persist the user turn before the tool loop so its timestamp precedes any
|
||||
# tool events fired mid-turn (keeps the transcript export in true order).
|
||||
memory.remember(session_id, "user", user_msg)
|
||||
reply, tools_run = _mind_loop(messages, backend, model, tool_specs, ctx, session_id)
|
||||
_ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run, backend, model, ctx, session_id)
|
||||
mouth = _mouth_target(cfg, backend, model)
|
||||
if mouth and reply:
|
||||
reply = _voice_pass(messages, reply, *mouth)
|
||||
if not reply:
|
||||
reply = _TANGLED
|
||||
logbus.log("info", "reply", session=session_id, chars=len(reply), voiced=bool(mouth))
|
||||
|
||||
memory.remember(session_id, "assistant", reply)
|
||||
summary.maybe_summarize_async(session_id) # compact once enough new turns pile up
|
||||
return reply
|
||||
finally:
|
||||
_finish_turn(rec, reply)
|
||||
|
||||
|
||||
def respond_stream(session_id: str, user_msg: str, backend: Backend = "cloud",
|
||||
@@ -164,73 +219,87 @@ def respond_stream(session_id: str, user_msg: str, backend: Backend = "cloud",
|
||||
logbus.log("info", "chat request (stream)", session=session_id, backend=backend,
|
||||
model=model, embed=cfg.embed_backend)
|
||||
|
||||
turn = mind.assemble(session_id, user_msg, backend, model)
|
||||
messages = turn.messages
|
||||
tool_specs = toolkit.specs(turn.mode.tools) if backend in cfg.tool_backends else None
|
||||
ctx = {"session_id": session_id, "backend": backend}
|
||||
mouth = _mouth_target(cfg, backend, model)
|
||||
# A concurrent duplicate (this stream + the UI's blocking fallback for one
|
||||
# message) reuses the owner's result instead of running a second full turn.
|
||||
is_owner, rec = _claim_turn(session_id, user_msg)
|
||||
if not is_owner:
|
||||
logbus.log("info", "duplicate turn deduped", session=session_id, path="stream")
|
||||
reply = _await_duplicate(rec)
|
||||
yield ("delta", reply)
|
||||
yield ("done", reply)
|
||||
return
|
||||
|
||||
# Persist the user turn up front (see respond): keeps tool events, which fire
|
||||
# mid-turn, chronologically after the user message in the exported transcript.
|
||||
memory.remember(session_id, "user", user_msg)
|
||||
reply = _TANGLED
|
||||
try:
|
||||
turn = mind.assemble(session_id, user_msg, backend, model)
|
||||
messages = turn.messages
|
||||
tool_specs = toolkit.specs(turn.mode.tools) if backend in cfg.tool_backends else None
|
||||
ctx = {"session_id": session_id, "backend": backend}
|
||||
mouth = _mouth_target(cfg, backend, model)
|
||||
|
||||
if mouth is None:
|
||||
# No separate voice: stream the mind directly (the original path, unchanged).
|
||||
parts: list[str] = []
|
||||
tools_run: list[str] = []
|
||||
for _ in range(MAX_TOOL_ROUNDS):
|
||||
assistant_msg = None
|
||||
tool_calls = None
|
||||
for ev, payload in llm.chat_call_stream(
|
||||
messages, backend=backend, model=model, tools=tool_specs
|
||||
):
|
||||
if ev == "delta":
|
||||
parts.append(payload)
|
||||
yield ("delta", payload)
|
||||
elif ev == "message":
|
||||
assistant_msg = payload
|
||||
elif ev == "tool_calls":
|
||||
tool_calls = payload
|
||||
if not tool_calls:
|
||||
break
|
||||
messages.append(assistant_msg)
|
||||
for tc in tool_calls:
|
||||
result = toolkit.dispatch(tc["name"], tc["arguments"], ctx)
|
||||
memory.add_tool_event(session_id, tc["name"], tc["arguments"], result)
|
||||
logbus.log("info", "tool call", session=session_id, tool=tc["name"], result=result[:80])
|
||||
messages.append({"role": "tool", "tool_call_id": tc["id"], "content": result})
|
||||
_maybe_switch_mode(session_id, tc["name"])
|
||||
tools_run.append(tc["name"])
|
||||
yield ("tool", tc["name"])
|
||||
for name in _ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run,
|
||||
backend, model, ctx, session_id):
|
||||
yield ("tool", name)
|
||||
reply = "".join(parts)
|
||||
if not reply:
|
||||
reply = _TANGLED
|
||||
yield ("delta", reply)
|
||||
else:
|
||||
# Mind decides + runs tools (non-streamed); mouth re-voices, streamed.
|
||||
draft, tools_run = _mind_loop(messages, backend, model, tool_specs, ctx, session_id)
|
||||
tools_run += _ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run,
|
||||
backend, model, ctx, session_id)
|
||||
for name in tools_run:
|
||||
yield ("tool", name)
|
||||
parts = []
|
||||
try:
|
||||
for ev, payload in llm.chat_call_stream(
|
||||
mind.voice_messages(messages, draft), backend=mouth[0], model=mouth[1], tools=None
|
||||
):
|
||||
if ev == "delta":
|
||||
parts.append(payload)
|
||||
yield ("delta", payload)
|
||||
except Exception as exc:
|
||||
logbus.log("error", "voice stream failed", error=str(exc)[:160])
|
||||
reply = "".join(parts).strip() or draft or _TANGLED
|
||||
if not parts:
|
||||
yield ("delta", reply)
|
||||
# Persist the user turn up front (see respond): keeps tool events, which fire
|
||||
# mid-turn, chronologically after the user message in the exported transcript.
|
||||
memory.remember(session_id, "user", user_msg)
|
||||
|
||||
logbus.log("info", "reply", session=session_id, chars=len(reply), voiced=bool(mouth))
|
||||
memory.remember(session_id, "assistant", reply)
|
||||
summary.maybe_summarize_async(session_id)
|
||||
yield ("done", reply)
|
||||
if mouth is None:
|
||||
# No separate voice: stream the mind directly (the original path, unchanged).
|
||||
parts: list[str] = []
|
||||
tools_run: list[str] = []
|
||||
for _ in range(MAX_TOOL_ROUNDS):
|
||||
assistant_msg = None
|
||||
tool_calls = None
|
||||
for ev, payload in llm.chat_call_stream(
|
||||
messages, backend=backend, model=model, tools=tool_specs
|
||||
):
|
||||
if ev == "delta":
|
||||
parts.append(payload)
|
||||
yield ("delta", payload)
|
||||
elif ev == "message":
|
||||
assistant_msg = payload
|
||||
elif ev == "tool_calls":
|
||||
tool_calls = payload
|
||||
if not tool_calls:
|
||||
break
|
||||
messages.append(assistant_msg)
|
||||
for tc in tool_calls:
|
||||
result = toolkit.dispatch(tc["name"], tc["arguments"], ctx)
|
||||
memory.add_tool_event(session_id, tc["name"], tc["arguments"], result)
|
||||
logbus.log("info", "tool call", session=session_id, tool=tc["name"], result=result[:80])
|
||||
messages.append({"role": "tool", "tool_call_id": tc["id"], "content": result})
|
||||
_maybe_switch_mode(session_id, tc["name"])
|
||||
tools_run.append(tc["name"])
|
||||
yield ("tool", tc["name"])
|
||||
for name in _ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run,
|
||||
backend, model, ctx, session_id):
|
||||
yield ("tool", name)
|
||||
reply = "".join(parts)
|
||||
if not reply:
|
||||
reply = _TANGLED
|
||||
yield ("delta", reply)
|
||||
else:
|
||||
# Mind decides + runs tools (non-streamed); mouth re-voices, streamed.
|
||||
draft, tools_run = _mind_loop(messages, backend, model, tool_specs, ctx, session_id)
|
||||
tools_run += _ensure_hand_logged(messages, user_msg, turn.msg_type, tools_run,
|
||||
backend, model, ctx, session_id)
|
||||
for name in tools_run:
|
||||
yield ("tool", name)
|
||||
parts = []
|
||||
try:
|
||||
for ev, payload in llm.chat_call_stream(
|
||||
mind.voice_messages(messages, draft), backend=mouth[0], model=mouth[1], tools=None
|
||||
):
|
||||
if ev == "delta":
|
||||
parts.append(payload)
|
||||
yield ("delta", payload)
|
||||
except Exception as exc:
|
||||
logbus.log("error", "voice stream failed", error=str(exc)[:160])
|
||||
reply = "".join(parts).strip() or draft or _TANGLED
|
||||
if not parts:
|
||||
yield ("delta", reply)
|
||||
|
||||
logbus.log("info", "reply", session=session_id, chars=len(reply), voiced=bool(mouth))
|
||||
memory.remember(session_id, "assistant", reply)
|
||||
summary.maybe_summarize_async(session_id)
|
||||
yield ("done", reply)
|
||||
finally:
|
||||
_finish_turn(rec, reply)
|
||||
|
||||
Reference in New Issue
Block a user