Compare commits
11 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e61bd9cb08 | |||
| c6e0130e59 | |||
| 67d6f3fe68 | |||
| dbc9690358 | |||
| 6d98728a2e | |||
| bfb2ffb6f6 | |||
| bd13b64959 | |||
| f2a57005e5 | |||
| 88fae33152 | |||
| c874883a84 | |||
| 6f22e86f54 |
@@ -174,13 +174,7 @@ Deferred to Phase 2: second bot, group node, scene configurations, witness filte
|
|||||||
|
|
||||||
### Phase 1.5 cleanup backlog
|
### Phase 1.5 cleanup backlog
|
||||||
|
|
||||||
Small follow-ups identified during Phase 1 reviews. Pick up at any time; none are blocking.
|
All items shipped — see Phase 2.5 status below.
|
||||||
|
|
||||||
- **`open_db` refactor.** `chat/web/bots.py:get_conn()` duplicates the context-manager body to add `check_same_thread=False`. Extend `open_db(path, *, check_same_thread=True)` and have `get_conn` call it directly — eliminates the duplicated PRAGMA setup and ensures any future PRAGMA tweak only happens in one place.
|
|
||||||
- **Regenerate broadcasts `turn_html` over SSE.** Currently a refresh is needed (see T29 limitation above). Mirror the broadcast logic from `chat/web/turns.py:post_turn` after the new `assistant_turn` lands.
|
|
||||||
- **`bot_reset` purges orphaned "you" activity rows** (see limitation above). Either delete `activity` rows by chat-membership or accept the noise indefinitely; the projection-layer fix is one extra `DELETE FROM activity WHERE entity_id='you' AND container_id IN (SELECT id FROM containers WHERE chat_id IN (...))` clause inside `_apply_bot_reset`.
|
|
||||||
- **Drawer edits for the deferred v1 fields**: edge_trust slider, edge_summary textarea, memory pov_summary textarea, knowledge_facts add/remove. The `manual_edit` projector already supports `edge_trust` / `edge_summary` / `memory_pov_summary` target_kinds — only the routes are missing. Knowledge_facts needs a new dispatch branch.
|
|
||||||
- **NICE trim order in prompt assembly** drops previous-scene first instead of last (T18 review). Greedy-cuts heuristic vs spec listing order; revisit if v1 play surfaces a real regression.
|
|
||||||
|
|
||||||
## Phase 2 status
|
## Phase 2 status
|
||||||
|
|
||||||
@@ -194,15 +188,28 @@ Phase 2 shipped end-to-end across **13 tasks** (T36–T48 wave). The multi-entit
|
|||||||
|
|
||||||
### Phase 2.5 / 3 backlog
|
### Phase 2.5 / 3 backlog
|
||||||
|
|
||||||
Carry-overs from Phase 2 reviews and implementer notes. None are blocking; pick up at any time.
|
All items shipped — see Phase 2.5 status below.
|
||||||
|
|
||||||
- **Interjection regenerate**: regenerate currently only acts on the addressee turn. Phase 2.5 should extend regenerate to cover the interjection turn too.
|
## Phase 2.5 status
|
||||||
- **Classifier-based addressee detection**: substring match is brittle (e.g., names that are common English words, or names appearing inside a quoted aside). A small classifier call could disambiguate.
|
|
||||||
- **LLM-merged group meta-summary**: current `group_node.summary` is a naive concat of host + guest per-POV summaries. Phase 2.5 should polish with an LLM-merged group view.
|
Phase 2.5 cleanup shipped end-to-end across 8 tasks (T68–T75). Two CLAUDE.md backlogs (Phase 1.5 cleanup, Phase 2.5/3) are now empty; deferred follow-ups discovered during execution are tracked in a new "Phase 2.6 / 3 backlog" section below.
|
||||||
- **First-meeting gate**: the drawer's "have they met?" textarea fires every time. Phase 2.5 should check whether the host→guest edge already exists and offer a "they already know each other" toggle to skip re-seeding.
|
|
||||||
- **Witness flag editing**: drawer doesn't allow editing memory witness flags (read-only). Phase 2.5+ may expose this.
|
- **`open_db` with check_same_thread parameter (T68)**: refactored `chat/db/connection.py` so `chat/web/bots.py:get_conn` no longer duplicates the PRAGMA setup. Default behavior preserved.
|
||||||
- **Significance for interjection memories**: the interjection's `memory_written` event doesn't enqueue a `SignificanceJob` (per the T44 implementer note). Phase 2.5 should wire this in so interjection memories are scored alongside primary turns.
|
- **`bot_reset` cross-chat cleanup (T69)**: now purges orphaned "you" activity rows. Note: this also fixed a latent FK constraint crash that was lurking in the projector — `activity.container_id` is FK-referenced and the prior code would have crashed on any reset of a bot whose chat had a non-NULL `container_id` "you" activity row. The bug was masked because no prior test seeded such a row.
|
||||||
- **Stale guest reference defensive degrade in `post_turn`**: T44 added a degrade-to-1:1 when `chat.guest_bot_id` points at a deleted bot. T47 fixes the root cause (resets clear the reference); the degrade can probably be removed but is harmless.
|
- **LLM-merged group meta-summary (T70)**: replaces Phase 2 T45's naive concat with a classifier merge call. Falls back to the naive concat on classifier failure.
|
||||||
- **Scene close on cancel**: scene close runs even when the primary turn is cancelled. Behavior may be intentional but could be argued either way; revisit if it surfaces a real UX regression.
|
- **`prompt.py` polish (T71)**: witness role parametric (`host` vs `guest` derived from chat membership); single `ACTIVITIES:` block with bullet-level trim; NICE trim order kept with documented rationale (greedy cheapest-impact-first beats spec-listing order in practice).
|
||||||
- **Dual `ACTIVITIES:` block**: T43's prompt assembly adds a second `ACTIVITIES:` block for guest activity. Cleaner would be a single block with three bullets and per-bullet trim.
|
- **Drawer polish (T72)**: deferred v1 edits (edge_trust slider, edge_summary textarea, memory pov_summary textarea, knowledge_facts add/remove) + first-meeting gate (Add-guest form disables prose textarea when host→guest edge already exists; "re-seed anyway" toggle re-enables) + witness flag inline-edit (per-memory checkboxes for [you, host, guest] flags). Two new `manual_edit` projector branches: `edge_knowledge_fact` and `memory_witness`.
|
||||||
- **Witness role hardcoded in prompt assembly**: `chat/services/prompt.py:436` hardcodes `witness_role="host"` regardless of which bot is speaking. Phase 2.5 should derive the role from chat membership (e.g. `"host" if speaker_bot_id == chat.host_bot_id else "guest"`) so guest-as-speaker prompts retrieve the right memory slice. Test contract pinned in `tests/test_witness_filter_multi.py`.
|
- **Regenerate polish (T73)**: regenerate now broadcasts `turn_html_replace` over SSE (NEW event distinct from `turn_html` to avoid breaking the existing append-semantic consumer); regenerate covers interjection turns (re-detects + re-streams or supersedes); defensive stale-guest degrade removed.
|
||||||
|
- **Turn-flow polish + addressee service (T74)**: classifier-based addressee detection (substring helper kept as no-guest fast path); SignificanceJob enqueued for interjection memories; scene-close-on-cancel pinned with comment + regression test (close detection is genuinely user-prose-only); defensive stale-guest degrade removed.
|
||||||
|
|
||||||
|
### Phase 2.6 / 3 backlog
|
||||||
|
|
||||||
|
New follow-ups discovered during Phase 2.5 execution. None are blocking; pick up at any time.
|
||||||
|
|
||||||
|
- **Frontend handler for `turn_html_replace` SSE event (from T73.1 review)**: regenerate's backend broadcast lands, but no live tab swaps the regenerated turn until a JS handler is wired. The existing `turn_html` event uses HTMX `sse-swap` to append; `turn_html_replace` ships JSON with `supersedes_id` for replacement semantics. Phase 2.6 should wire the JS to swap the prior turn's DOM node in place.
|
||||||
|
- **Cancel/stop hook for in-flight regenerate streams (from T73 review)**: `post_turn` registers stream tasks in `_in_flight_tasks` so the user can stop them. Regenerate doesn't. A user clicking "Stop" mid-regenerate has no cancel hook today.
|
||||||
|
- **DRY: regenerate vs post_turn (from T73 review)**: recent-dialogue assembly and prior-edges block are duplicated between `chat/services/regenerate.py` and `chat/web/turns.py`. Extract to shared helpers analogous to `_gather_state_update_inputs`.
|
||||||
|
- **Sibling-discovery query optimization (from T73 review)**: `regenerate.py`'s sibling-assistant-turn lookup scans all non-superseded `assistant_turn` rows globally. Adding a `chat_id` predicate via JSON extraction (or a denormalized column) bounds the cost to per-chat scale.
|
||||||
|
- **`_witness_role_for` defensive coding (from T71 review)**: helper returns `"guest"` when `host_bot_id is None`, which is wrong for Phase-1 chats. Defensive: `return "host" if host_bot_id is None or speaker_bot_id == host_bot_id else "guest"`. Not exercised by current tests; harden as a precaution.
|
||||||
|
- **Confidence type tightening (from T74 review)**: `chat/services/addressee.py::AddresseeDecision.confidence` could be typed as `Literal["high","medium","low"]` for stricter validation. Currently `str` with a comment.
|
||||||
|
- **Scene-close-on-cancel UX revisit**: T74.3 pinned the existing behavior (close fires even on cancel). If real play-testing surfaces a regression, revisit.
|
||||||
|
|||||||
@@ -0,0 +1,108 @@
|
|||||||
|
"""Addressee classifier service (T74.1).
|
||||||
|
|
||||||
|
Phase 2 (T44) detected the addressee — host vs. guest — with a simple
|
||||||
|
case-insensitive whole-word substring match against the bots' names.
|
||||||
|
That worked for the obvious case ("BotB, what do you think?") but lost
|
||||||
|
the long tail: pronouns, paraphrases, indirect address, narrative
|
||||||
|
focus on a particular party. T74.1 swaps the substring helper for a
|
||||||
|
classifier call that reads the prose holistically.
|
||||||
|
|
||||||
|
The substring helper in :mod:`chat.web.turns` is kept as a fast-path
|
||||||
|
for the no-guest case (only one bot present means there is nothing to
|
||||||
|
classify) and as a non-breaking fallback for the regenerate path. The
|
||||||
|
multi-entity branch in :func:`chat.web.turns.post_turn` calls
|
||||||
|
:func:`detect_addressee` from this module.
|
||||||
|
|
||||||
|
Failure mode: classifier flake or low-confidence response degrades to
|
||||||
|
the host (the default speaker per Phase 2's host-keeps-the-floor
|
||||||
|
bias). The decision carries ``confidence`` and ``reason`` so callers
|
||||||
|
that want to log degraded decisions can distinguish a real "host" call
|
||||||
|
from a fallback.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
from chat.llm.classify import classify
|
||||||
|
from chat.llm.client import LLMClient
|
||||||
|
|
||||||
|
|
||||||
|
class AddresseeDecision(BaseModel):
|
||||||
|
"""Which present bot the user is addressing.
|
||||||
|
|
||||||
|
``addressee_id`` is the chosen bot's id. ``confidence`` is one of
|
||||||
|
``"high"`` / ``"medium"`` / ``"low"`` — callers may treat ``"low"``
|
||||||
|
as a soft fallback to the host. ``reason`` is a short free-form
|
||||||
|
string. The classifier-failure fallback uses ``reason="fallback"``
|
||||||
|
so it's distinguishable from a real low-confidence call.
|
||||||
|
"""
|
||||||
|
|
||||||
|
addressee_id: str
|
||||||
|
confidence: str = "medium" # "high" | "medium" | "low"
|
||||||
|
reason: str = ""
|
||||||
|
|
||||||
|
|
||||||
|
_SYSTEM = (
|
||||||
|
"Given a user's turn prose and the names of present bots, decide "
|
||||||
|
"which bot the user is addressing. If the user is speaking to no "
|
||||||
|
"specific bot (descriptive narration, action without dialogue), "
|
||||||
|
"default to the host. Output strict JSON matching the schema. "
|
||||||
|
"The addressee_id MUST be one of the ids supplied in the user "
|
||||||
|
"message — do not invent ids."
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def detect_addressee(
|
||||||
|
client: LLMClient,
|
||||||
|
*,
|
||||||
|
classifier_model: str,
|
||||||
|
user_prose: str,
|
||||||
|
host_id: str,
|
||||||
|
host_name: str,
|
||||||
|
guest_id: str | None,
|
||||||
|
guest_name: str | None,
|
||||||
|
timeout_s: float = 30.0,
|
||||||
|
) -> AddresseeDecision:
|
||||||
|
"""Classify which present bot the user is addressing.
|
||||||
|
|
||||||
|
Defaults to host on classifier failure or when the classifier picks
|
||||||
|
an id that isn't one of the supplied ids. The caller is expected to
|
||||||
|
only invoke this in the multi-entity case (a guest is present);
|
||||||
|
when no guest is present the substring fast-path in
|
||||||
|
:mod:`chat.web.turns` is used instead and this function is not
|
||||||
|
called.
|
||||||
|
"""
|
||||||
|
fallback = AddresseeDecision(
|
||||||
|
addressee_id=host_id, confidence="low", reason="fallback"
|
||||||
|
)
|
||||||
|
user = (
|
||||||
|
f"Host: {host_name} (id={host_id})\n"
|
||||||
|
+ (
|
||||||
|
f"Guest: {guest_name} (id={guest_id})\n"
|
||||||
|
if guest_id is not None
|
||||||
|
else ""
|
||||||
|
)
|
||||||
|
+ f"\nUser prose:\n{user_prose}"
|
||||||
|
)
|
||||||
|
decision = await classify(
|
||||||
|
client,
|
||||||
|
model=classifier_model,
|
||||||
|
system=_SYSTEM,
|
||||||
|
user=user,
|
||||||
|
schema=AddresseeDecision,
|
||||||
|
default=fallback,
|
||||||
|
timeout_s=timeout_s,
|
||||||
|
)
|
||||||
|
# Defensive: if the classifier returned an id outside the supplied
|
||||||
|
# set, treat it as a fallback to the host. This catches pathological
|
||||||
|
# outputs that pass schema validation but pick a phantom id.
|
||||||
|
valid_ids = {host_id}
|
||||||
|
if guest_id is not None:
|
||||||
|
valid_ids.add(guest_id)
|
||||||
|
if decision.addressee_id not in valid_ids:
|
||||||
|
return fallback
|
||||||
|
return decision
|
||||||
|
|
||||||
|
|
||||||
|
__all__ = ["AddresseeDecision", "detect_addressee"]
|
||||||
+306
-6
@@ -26,6 +26,7 @@ Phase 1 simplifications (per the plan's "bound it" guidance):
|
|||||||
so affinity/trust/knowledge reflect the new output.
|
so affinity/trust/knowledge reflect the new output.
|
||||||
- The route does not broadcast a fresh ``turn_html`` SSE event; T34
|
- The route does not broadcast a fresh ``turn_html`` SSE event; T34
|
||||||
polishes UI swaps. The user refreshes the page to see the new turn.
|
polishes UI swaps. The user refreshes the page to see the new turn.
|
||||||
|
*(T73.1 closed this gap — see Phase 2.5 changes below.)*
|
||||||
|
|
||||||
Phase 2 changes (T44):
|
Phase 2 changes (T44):
|
||||||
|
|
||||||
@@ -42,6 +43,27 @@ Phase 2 changes (T44):
|
|||||||
is not invoked here. If the prior turn fired an interjection it
|
is not invoked here. If the prior turn fired an interjection it
|
||||||
remains attached to the original assistant_turn (which is superseded
|
remains attached to the original assistant_turn (which is superseded
|
||||||
alongside the regenerated turn) — Phase 2.5 will revisit.
|
alongside the regenerated turn) — Phase 2.5 will revisit.
|
||||||
|
|
||||||
|
Phase 2.5 changes:
|
||||||
|
|
||||||
|
- T73.1: After the new ``assistant_turn`` lands we publish a
|
||||||
|
``turn_html_replace`` SSE event carrying the rendered HTML for the
|
||||||
|
regenerated turn plus the original assistant_turn's event_id as
|
||||||
|
``supersedes_id`` so connected tabs can swap the prior DOM node
|
||||||
|
in-place. We use a NEW event name (rather than re-using ``turn_html``)
|
||||||
|
because the existing HTMX ``sse-swap="turn_html"`` consumer expects a
|
||||||
|
raw-HTML body and an *append* semantic; ``turn_html_replace`` is a
|
||||||
|
JSON payload (sse.py auto-serialises when extra keys accompany
|
||||||
|
``data``) so the front-end JS can read ``supersedes_id`` and replace
|
||||||
|
the right node.
|
||||||
|
- T73.2: Interjection regeneration. When the original assistant_turn
|
||||||
|
group included an interjection beat we redo BOTH the primary and the
|
||||||
|
interjection — re-running ``detect_interjection`` against the new
|
||||||
|
primary text. If the classifier returns False this time we supersede
|
||||||
|
the original interjection without appending a replacement.
|
||||||
|
- T73.3: The defensive degrade-to-1:1 for stale ``guest_bot_id``
|
||||||
|
references was removed — Phase 2 T47 fixed the root cause (resets
|
||||||
|
clear the reference) so the guard is dead code.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
@@ -51,6 +73,7 @@ from sqlite3 import Connection
|
|||||||
|
|
||||||
from chat.config import Settings
|
from chat.config import Settings
|
||||||
from chat.eventlog.log import append_and_apply, append_event
|
from chat.eventlog.log import append_and_apply, append_event
|
||||||
|
from chat.services.interjection import detect_interjection
|
||||||
from chat.services.memory_write import record_turn_memory_for_present
|
from chat.services.memory_write import record_turn_memory_for_present
|
||||||
from chat.services.multi_state_update import compute_state_updates_for_present
|
from chat.services.multi_state_update import compute_state_updates_for_present
|
||||||
from chat.services.prompt import assemble_narrative_prompt
|
from chat.services.prompt import assemble_narrative_prompt
|
||||||
@@ -58,6 +81,7 @@ from chat.state.edges import get_edge
|
|||||||
from chat.state.entities import get_bot, get_you
|
from chat.state.entities import get_bot, get_you
|
||||||
from chat.state.world import active_scene, get_chat
|
from chat.state.world import active_scene, get_chat
|
||||||
from chat.web.pubsub import publish
|
from chat.web.pubsub import publish
|
||||||
|
from chat.web.render import render_turn_html
|
||||||
|
|
||||||
|
|
||||||
async def regenerate_assistant_turn(
|
async def regenerate_assistant_turn(
|
||||||
@@ -90,13 +114,13 @@ async def regenerate_assistant_turn(
|
|||||||
|
|
||||||
# Phase 2: surface the guest (if any) so the prompt assembler and
|
# Phase 2: surface the guest (if any) so the prompt assembler and
|
||||||
# downstream multi-entity passes see the same shape post_turn does.
|
# downstream multi-entity passes see the same shape post_turn does.
|
||||||
|
# Phase 2 T47 made bot_reset cascade-clear ``chat.guest_bot_id`` when
|
||||||
|
# the referenced bot is purged (verified by tests/test_reset.py), so
|
||||||
|
# we trust the column here: it's either a valid bot id or NULL.
|
||||||
guest_bot_id = chat.get("guest_bot_id")
|
guest_bot_id = chat.get("guest_bot_id")
|
||||||
guest_bot: dict | None = None
|
guest_bot: dict | None = (
|
||||||
if guest_bot_id is not None:
|
get_bot(conn, guest_bot_id) if guest_bot_id is not None else None
|
||||||
guest_bot = get_bot(conn, guest_bot_id)
|
)
|
||||||
if guest_bot is None:
|
|
||||||
# Stale guest reference — degrade to single-bot regenerate.
|
|
||||||
guest_bot_id = None
|
|
||||||
|
|
||||||
# 1. Locate the original assistant_turn event.
|
# 1. Locate the original assistant_turn event.
|
||||||
row = conn.execute(
|
row = conn.execute(
|
||||||
@@ -108,6 +132,33 @@ async def regenerate_assistant_turn(
|
|||||||
raise ValueError("assistant_turn event not found")
|
raise ValueError("assistant_turn event not found")
|
||||||
original_assistant_payload = json.loads(row[0])
|
original_assistant_payload = json.loads(row[0])
|
||||||
original_user_turn_id = original_assistant_payload.get("user_turn_id")
|
original_user_turn_id = original_assistant_payload.get("user_turn_id")
|
||||||
|
|
||||||
|
# 1a. Look up any sibling interjection beat in the same turn group
|
||||||
|
# (T73.2). The original group is (primary + optional interjection),
|
||||||
|
# both pinned to the same ``user_turn_id``. The interjection has a
|
||||||
|
# populated ``interjection_of`` field in its payload — its speaker is
|
||||||
|
# the silent witness (the bot that wasn't the primary addressee).
|
||||||
|
# Filter on ``superseded_by IS NULL`` so prior regenerates of this
|
||||||
|
# group don't reappear as siblings.
|
||||||
|
original_interjection_event_id: int | None = None
|
||||||
|
original_interjection_payload: dict | None = None
|
||||||
|
if original_user_turn_id is not None:
|
||||||
|
sibling_cur = conn.execute(
|
||||||
|
"SELECT id, payload_json FROM event_log "
|
||||||
|
"WHERE kind = 'assistant_turn' "
|
||||||
|
" AND id != ? "
|
||||||
|
" AND superseded_by IS NULL",
|
||||||
|
(original_assistant_event_id,),
|
||||||
|
)
|
||||||
|
for sib_id, sib_payload_json in sibling_cur.fetchall():
|
||||||
|
sib_payload = json.loads(sib_payload_json)
|
||||||
|
if sib_payload.get("user_turn_id") != original_user_turn_id:
|
||||||
|
continue
|
||||||
|
if not sib_payload.get("interjection_of"):
|
||||||
|
continue
|
||||||
|
original_interjection_event_id = sib_id
|
||||||
|
original_interjection_payload = sib_payload
|
||||||
|
break
|
||||||
# Phase 2 v2 regenerates only the addressee turn — preserve whichever
|
# Phase 2 v2 regenerates only the addressee turn — preserve whichever
|
||||||
# bot the original turn was attributed to, falling back to the host
|
# bot the original turn was attributed to, falling back to the host
|
||||||
# for legacy rows that pre-date multi-entity support.
|
# for legacy rows that pre-date multi-entity support.
|
||||||
@@ -238,6 +289,27 @@ async def regenerate_assistant_turn(
|
|||||||
(new_assistant_event_id, original_assistant_event_id),
|
(new_assistant_event_id, original_assistant_event_id),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# 7a. Broadcast a turn_html_replace SSE event so connected tabs can
|
||||||
|
# swap the prior assistant_turn DOM node in-place (T73.1, Phase 1.5
|
||||||
|
# backlog #2). Uses a separate event name from post_turn's
|
||||||
|
# ``turn_html`` (which is append-only) because regenerate is a
|
||||||
|
# *replace* operation — see module docstring for the rationale.
|
||||||
|
speaker_name_for_render = (
|
||||||
|
speaker_bot.get("name", "bot") if speaker_bot is not None else "bot"
|
||||||
|
)
|
||||||
|
new_turn_html = render_turn_html(
|
||||||
|
speaker_name_for_render, new_text, role="bot"
|
||||||
|
)
|
||||||
|
await publish(
|
||||||
|
chat_id,
|
||||||
|
{
|
||||||
|
"event": "turn_html_replace",
|
||||||
|
"data": new_turn_html,
|
||||||
|
"turn_id": new_assistant_event_id,
|
||||||
|
"supersedes_id": original_assistant_event_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
# 8. Re-run downstream classifier passes (memory write + state update
|
# 8. Re-run downstream classifier passes (memory write + state update
|
||||||
# for every directed pair across present entities). Significance is
|
# for every directed pair across present entities). Significance is
|
||||||
# intentionally skipped on regenerate (the prior score remains
|
# intentionally skipped on regenerate (the prior score remains
|
||||||
@@ -317,6 +389,234 @@ async def regenerate_assistant_turn(
|
|||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# 9. Interjection regenerate branch (T73.2). When the original
|
||||||
|
# assistant_turn group included a follow-on interjection beat we need
|
||||||
|
# to revisit that beat against the regenerated primary. Three outcomes:
|
||||||
|
#
|
||||||
|
# - No original interjection: nothing to do; we already short-circuit
|
||||||
|
# above by leaving ``original_interjection_event_id`` as None.
|
||||||
|
# - Original interjection + classifier returns True: stream a fresh
|
||||||
|
# interjection from the silent witness, append it (with
|
||||||
|
# ``interjection_of`` linking to the new primary speaker), and
|
||||||
|
# supersede the original interjection's row. Also re-run memory
|
||||||
|
# + state-update so the second beat moves edges + writes memories.
|
||||||
|
# - Original interjection + classifier returns False: supersede the
|
||||||
|
# original interjection without appending a replacement. The
|
||||||
|
# regenerated group becomes "primary only" because the new primary
|
||||||
|
# no longer warrants a follow-on. No memory / state work needed
|
||||||
|
# for the absent beat.
|
||||||
|
#
|
||||||
|
# ``superseded_by`` on the original interjection's row points at the
|
||||||
|
# *new primary* in the no-replacement case (rather than NULL or a
|
||||||
|
# nonexistent id) so the row is consistently hidden by the standard
|
||||||
|
# ``superseded_by IS NULL`` timeline filter and the back-pointer
|
||||||
|
# leads somewhere meaningful for an "originally said …" affordance.
|
||||||
|
if original_interjection_event_id is not None and guest_bot is not None:
|
||||||
|
# Identify the silent witness from the original interjection's
|
||||||
|
# speaker_id (which is the bot that interjected last time). When
|
||||||
|
# we regenerate we keep the *same pair of present entities*, so
|
||||||
|
# the silent witness is whichever bot isn't the new primary
|
||||||
|
# speaker — derive it from present rather than reusing the prior
|
||||||
|
# speaker_id verbatim, in case the regenerated primary swapped
|
||||||
|
# who held the floor.
|
||||||
|
if speaker_bot_id == host_bot_id:
|
||||||
|
silent_witness = guest_bot
|
||||||
|
else:
|
||||||
|
silent_witness = host_bot
|
||||||
|
silent_witness_id = silent_witness.get("id")
|
||||||
|
|
||||||
|
edge_w_to_addr = get_edge(conn, silent_witness_id, speaker_bot_id) or {
|
||||||
|
"affinity": 50,
|
||||||
|
"trust": 50,
|
||||||
|
"summary": "",
|
||||||
|
}
|
||||||
|
edge_w_to_you = get_edge(conn, silent_witness_id, "you") or {
|
||||||
|
"affinity": 50,
|
||||||
|
"trust": 50,
|
||||||
|
"summary": "",
|
||||||
|
}
|
||||||
|
|
||||||
|
decision = await detect_interjection(
|
||||||
|
client,
|
||||||
|
classifier_model=settings.classifier_model,
|
||||||
|
addressee_name=speaker_bot.get("name", "bot"),
|
||||||
|
addressee_just_said=new_text,
|
||||||
|
silent_witness_name=silent_witness.get("name", "bot"),
|
||||||
|
silent_witness_persona=silent_witness.get("persona") or "",
|
||||||
|
silent_witness_edge_to_addressee=edge_w_to_addr,
|
||||||
|
silent_witness_edge_to_you=edge_w_to_you,
|
||||||
|
you_just_said=prose_for_prompt or "",
|
||||||
|
timeout_s=settings.classifier_timeout_s,
|
||||||
|
)
|
||||||
|
|
||||||
|
if decision.should_interject:
|
||||||
|
# Re-read recent so the just-appended primary is in the prompt.
|
||||||
|
interject_cur = conn.execute(
|
||||||
|
"SELECT id, kind, payload_json FROM event_log "
|
||||||
|
"WHERE kind IN ('user_turn', 'user_turn_edit', 'assistant_turn') "
|
||||||
|
" AND superseded_by IS NULL AND hidden = 0 "
|
||||||
|
"ORDER BY id DESC LIMIT 20",
|
||||||
|
)
|
||||||
|
interject_rows = list(reversed(interject_cur.fetchall()))
|
||||||
|
interject_recent: list[dict] = []
|
||||||
|
for _eid, kind, payload_json in interject_rows:
|
||||||
|
p = json.loads(payload_json)
|
||||||
|
if p.get("chat_id") != chat_id:
|
||||||
|
continue
|
||||||
|
if kind in ("user_turn", "user_turn_edit"):
|
||||||
|
interject_recent.append(
|
||||||
|
{"speaker": you_name, "text": p.get("prose", "")}
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
spk = p.get("speaker_id", "bot")
|
||||||
|
if spk == host_bot_id:
|
||||||
|
spk_name = host_bot.get("name", "bot")
|
||||||
|
elif spk == guest_bot.get("id"):
|
||||||
|
spk_name = guest_bot.get("name", "bot")
|
||||||
|
else:
|
||||||
|
spk_name = "bot"
|
||||||
|
interject_recent.append(
|
||||||
|
{"speaker": spk_name, "text": p.get("text", "")}
|
||||||
|
)
|
||||||
|
if interject_recent and interject_recent[-1].get("speaker") == you_name:
|
||||||
|
interject_recent = interject_recent[:-1]
|
||||||
|
|
||||||
|
interject_messages = assemble_narrative_prompt(
|
||||||
|
conn,
|
||||||
|
chat_id=chat_id,
|
||||||
|
speaker_bot_id=silent_witness_id,
|
||||||
|
addressee=speaker_bot_id,
|
||||||
|
user_turn_prose=prose_for_prompt or None,
|
||||||
|
recent_dialogue=interject_recent,
|
||||||
|
budget_soft=settings.narrative_budget_soft,
|
||||||
|
budget_hard=settings.narrative_budget_hard,
|
||||||
|
guest_id=guest_bot_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
interject_accumulated: list[str] = []
|
||||||
|
async for chunk in client.stream(
|
||||||
|
interject_messages,
|
||||||
|
model=settings.narrative_model,
|
||||||
|
max_tokens=settings.narrative_max_tokens,
|
||||||
|
temperature=settings.narrative_temperature,
|
||||||
|
):
|
||||||
|
interject_accumulated.append(chunk)
|
||||||
|
await publish(
|
||||||
|
chat_id,
|
||||||
|
{
|
||||||
|
"event": "token",
|
||||||
|
"text": chunk,
|
||||||
|
"speaker_id": silent_witness_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
interject_text = "".join(interject_accumulated)
|
||||||
|
|
||||||
|
new_interjection_event_id = append_event(
|
||||||
|
conn,
|
||||||
|
kind="assistant_turn",
|
||||||
|
payload={
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"speaker_id": silent_witness_id,
|
||||||
|
"text": interject_text,
|
||||||
|
"truncated": False,
|
||||||
|
"user_turn_id": (
|
||||||
|
new_user_event_id
|
||||||
|
if new_user_event_id is not None
|
||||||
|
else original_user_turn_id
|
||||||
|
),
|
||||||
|
"regenerated_from": original_interjection_event_id,
|
||||||
|
"interjection_of": speaker_bot_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
# Supersede the original interjection by the new one.
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE event_log SET superseded_by = ? WHERE id = ?",
|
||||||
|
(new_interjection_event_id, original_interjection_event_id),
|
||||||
|
)
|
||||||
|
|
||||||
|
# Broadcast a replace event so connected tabs swap the prior
|
||||||
|
# interjection node in-place (mirrors T73.1's primary swap).
|
||||||
|
interject_html = render_turn_html(
|
||||||
|
silent_witness.get("name", "bot"), interject_text, role="bot"
|
||||||
|
)
|
||||||
|
await publish(
|
||||||
|
chat_id,
|
||||||
|
{
|
||||||
|
"event": "turn_html_replace",
|
||||||
|
"data": interject_html,
|
||||||
|
"turn_id": new_interjection_event_id,
|
||||||
|
"supersedes_id": original_interjection_event_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
# Memory write for the new interjection beat (one event per
|
||||||
|
# present witness).
|
||||||
|
record_turn_memory_for_present(
|
||||||
|
conn,
|
||||||
|
chat_id=chat_id,
|
||||||
|
host_bot_id=host_bot_id,
|
||||||
|
guest_bot_id=guest_bot_id,
|
||||||
|
narrative_text=interject_text,
|
||||||
|
scene_id=scene["id"] if scene else None,
|
||||||
|
chat_clock_at=chat.get("time"),
|
||||||
|
)
|
||||||
|
|
||||||
|
# Re-run the multi-pair state-update with the post-interjection
|
||||||
|
# dialogue tail so deltas land on the post-primary baseline.
|
||||||
|
recent_post_interject = recent_for_update + [
|
||||||
|
{
|
||||||
|
"speaker": silent_witness.get("name", "bot"),
|
||||||
|
"text": interject_text,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
prior_edges_post: dict[tuple[str, str], dict] = {}
|
||||||
|
for src in present_ids:
|
||||||
|
for tgt in present_ids:
|
||||||
|
if src == tgt:
|
||||||
|
continue
|
||||||
|
edge = get_edge(conn, src, tgt) or {
|
||||||
|
"affinity": 50,
|
||||||
|
"trust": 50,
|
||||||
|
"summary": "",
|
||||||
|
}
|
||||||
|
prior_edges_post[(src, tgt)] = edge
|
||||||
|
|
||||||
|
state_updates_post = await compute_state_updates_for_present(
|
||||||
|
client,
|
||||||
|
classifier_model=settings.classifier_model,
|
||||||
|
present_ids=present_ids,
|
||||||
|
present_names=present_names,
|
||||||
|
personas=personas,
|
||||||
|
prior_edges=prior_edges_post,
|
||||||
|
recent_dialogue=recent_post_interject,
|
||||||
|
timeout_s=settings.classifier_timeout_s,
|
||||||
|
)
|
||||||
|
for src_id, tgt_id, update in state_updates_post:
|
||||||
|
append_and_apply(
|
||||||
|
conn,
|
||||||
|
kind="edge_update",
|
||||||
|
payload={
|
||||||
|
"source_id": src_id,
|
||||||
|
"target_id": tgt_id,
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"affinity_delta": update.affinity_delta,
|
||||||
|
"trust_delta": update.trust_delta,
|
||||||
|
"knowledge_facts": update.knowledge_facts,
|
||||||
|
"last_interaction_at": last_at,
|
||||||
|
"last_interaction_chat_id": chat_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
# Classifier said "no follow-on this time" — supersede the
|
||||||
|
# original interjection without a replacement. Point the
|
||||||
|
# back-pointer at the new primary so the row is consistently
|
||||||
|
# hidden by the standard timeline filter.
|
||||||
|
conn.execute(
|
||||||
|
"UPDATE event_log SET superseded_by = ? WHERE id = ?",
|
||||||
|
(new_assistant_event_id, original_interjection_event_id),
|
||||||
|
)
|
||||||
|
|
||||||
return new_text
|
return new_text
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+62
-7
@@ -55,6 +55,7 @@ from fastapi import APIRouter, Depends, Form, HTTPException, Request
|
|||||||
from fastapi.responses import HTMLResponse, RedirectResponse, Response
|
from fastapi.responses import HTMLResponse, RedirectResponse, Response
|
||||||
|
|
||||||
from chat.eventlog.log import append_and_apply, append_event
|
from chat.eventlog.log import append_and_apply, append_event
|
||||||
|
from chat.services.addressee import detect_addressee
|
||||||
from chat.services.background import SignificanceJob
|
from chat.services.background import SignificanceJob
|
||||||
from chat.services.interjection import detect_interjection
|
from chat.services.interjection import detect_interjection
|
||||||
from chat.services.memory_write import record_turn_memory_for_present
|
from chat.services.memory_write import record_turn_memory_for_present
|
||||||
@@ -235,11 +236,12 @@ async def post_turn(
|
|||||||
guest_bot = None
|
guest_bot = None
|
||||||
guest_bot_id = chat.get("guest_bot_id")
|
guest_bot_id = chat.get("guest_bot_id")
|
||||||
if guest_bot_id is not None:
|
if guest_bot_id is not None:
|
||||||
|
# T47's bot_reset cascade clears guest_bot_id from any chat that
|
||||||
|
# referenced the deleted bot, so by the time we read it here it's
|
||||||
|
# either None or a live bot id. The previous defensive
|
||||||
|
# degrade-to-1:1 block (T44) was rendered dead by T47 and removed
|
||||||
|
# in T74.4 — get_bot now returns a real row.
|
||||||
guest_bot = get_bot(conn, guest_bot_id)
|
guest_bot = get_bot(conn, guest_bot_id)
|
||||||
# If the chat references a deleted guest we degrade to single-bot
|
|
||||||
# rather than 404 — the chat is still usable as a 1:1.
|
|
||||||
if guest_bot is None:
|
|
||||||
guest_bot_id = None
|
|
||||||
|
|
||||||
settings = request.app.state.settings
|
settings = request.app.state.settings
|
||||||
|
|
||||||
@@ -262,8 +264,25 @@ async def post_turn(
|
|||||||
|
|
||||||
# 3. Determine the addressee. Done before assistant_turn_started so the
|
# 3. Determine the addressee. Done before assistant_turn_started so the
|
||||||
# placeholder reflects the bot the user is actually talking to (host
|
# placeholder reflects the bot the user is actually talking to (host
|
||||||
# in 1:1, host-or-guest in multi-entity).
|
# in 1:1, host-or-guest in multi-entity). T74.1 routes the multi-entity
|
||||||
addressee_id = _detect_addressee_id(prose, host_bot, guest_bot)
|
# case through the addressee classifier; the no-guest case still uses
|
||||||
|
# the substring fast-path because there is nothing to classify when
|
||||||
|
# only one bot is present (and a classifier round-trip there would
|
||||||
|
# just be throughput overhead).
|
||||||
|
if guest_bot is None:
|
||||||
|
addressee_id = _detect_addressee_id(prose, host_bot, guest_bot)
|
||||||
|
else:
|
||||||
|
decision = await detect_addressee(
|
||||||
|
client,
|
||||||
|
classifier_model=settings.classifier_model,
|
||||||
|
user_prose=prose,
|
||||||
|
host_id=host_bot["id"],
|
||||||
|
host_name=host_bot["name"],
|
||||||
|
guest_id=guest_bot["id"],
|
||||||
|
guest_name=guest_bot["name"],
|
||||||
|
timeout_s=settings.classifier_timeout_s,
|
||||||
|
)
|
||||||
|
addressee_id = decision.addressee_id
|
||||||
addressee_bot = (
|
addressee_bot = (
|
||||||
guest_bot if (guest_bot is not None and addressee_id == guest_bot["id"])
|
guest_bot if (guest_bot is not None and addressee_id == guest_bot["id"])
|
||||||
else host_bot
|
else host_bot
|
||||||
@@ -598,7 +617,7 @@ async def post_turn(
|
|||||||
|
|
||||||
# Memory write for the interjection beat — a second pair
|
# Memory write for the interjection beat — a second pair
|
||||||
# of memory_written events (host + guest POVs).
|
# of memory_written events (host + guest POVs).
|
||||||
record_turn_memory_for_present(
|
interject_memory_results = record_turn_memory_for_present(
|
||||||
conn,
|
conn,
|
||||||
chat_id=chat_id,
|
chat_id=chat_id,
|
||||||
host_bot_id=host_bot["id"],
|
host_bot_id=host_bot["id"],
|
||||||
@@ -608,6 +627,33 @@ async def post_turn(
|
|||||||
chat_clock_at=chat.get("time"),
|
chat_clock_at=chat.get("time"),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# T74.2: enqueue a significance pass for the interjection
|
||||||
|
# memory. Mirrors the primary-turn enqueue pattern above —
|
||||||
|
# we score on the host's memory id since the prose is
|
||||||
|
# identical across both POVs (per-POV rewrite happens at
|
||||||
|
# scene close in T45). Without this enqueue the
|
||||||
|
# interjection beat lands in memory but never gets scored,
|
||||||
|
# so it can never auto-pin even when it carries a pivotal
|
||||||
|
# moment.
|
||||||
|
interject_host_event = interject_memory_results.get(
|
||||||
|
host_bot["id"]
|
||||||
|
)
|
||||||
|
interject_host_memory_id = (
|
||||||
|
interject_host_event[1] if interject_host_event else None
|
||||||
|
)
|
||||||
|
if (
|
||||||
|
worker is not None
|
||||||
|
and interject_host_memory_id is not None
|
||||||
|
):
|
||||||
|
worker.enqueue(
|
||||||
|
SignificanceJob(
|
||||||
|
memory_id=interject_host_memory_id,
|
||||||
|
narrative_text=interjection_text,
|
||||||
|
prior_dialogue=recent_post_interject,
|
||||||
|
host_bot_id=host_bot["id"],
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
# 9. Scene-close detection (Plan §7.2, T26). Runs AFTER assistant_turn
|
# 9. Scene-close detection (Plan §7.2, T26). Runs AFTER assistant_turn
|
||||||
# and the optional interjection so the bots' responses are part of
|
# and the optional interjection so the bots' responses are part of
|
||||||
# the closing scene's final beat — closing before narrative would
|
# the closing scene's final beat — closing before narrative would
|
||||||
@@ -623,6 +669,15 @@ async def post_turn(
|
|||||||
# close in the same chat) — we have nothing to close. T13 (kickoff)
|
# close in the same chat) — we have nothing to close. T13 (kickoff)
|
||||||
# is the only scene-opener path in v1; Phase 2-3 will handle
|
# is the only scene-opener path in v1; Phase 2-3 will handle
|
||||||
# automatic re-opening with the next container.
|
# automatic re-opening with the next container.
|
||||||
|
#
|
||||||
|
# T74.3: this branch deliberately runs even when ``cancelled`` is
|
||||||
|
# True. Close detection consumes only the user's prose (which is
|
||||||
|
# fully appended to the event_log BEFORE streaming starts) and the
|
||||||
|
# current container name; it does NOT consume the bot's output.
|
||||||
|
# A user who types "we're done here, fade out" and then hits Stop
|
||||||
|
# mid-stream still meant to close the scene — the cancelled bot
|
||||||
|
# beat doesn't invalidate that intent. Pinned by
|
||||||
|
# test_cancelled_turn_still_closes_scene_when_user_prose_signals_close.
|
||||||
if scene is not None and prose.strip():
|
if scene is not None and prose.strip():
|
||||||
container = None
|
container = None
|
||||||
if scene.get("container_id") is not None:
|
if scene.get("container_id") is not None:
|
||||||
|
|||||||
@@ -0,0 +1,99 @@
|
|||||||
|
"""Addressee classifier service tests (T74.1).
|
||||||
|
|
||||||
|
Covers :func:`chat.services.addressee.detect_addressee`:
|
||||||
|
|
||||||
|
- Classifier picks the guest -> ``addressee_id == guest_id``.
|
||||||
|
- Classifier picks the host -> ``addressee_id == host_id``.
|
||||||
|
- Classifier flakes (3 bad-JSON responses, exhausting the built-in
|
||||||
|
retry budget in :func:`chat.llm.classify.classify`) -> fallback to
|
||||||
|
the host with ``reason="fallback"``.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from chat.llm.mock import MockLLMClient
|
||||||
|
from chat.services.addressee import AddresseeDecision, detect_addressee
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_classifier_picks_guest():
|
||||||
|
"""Classifier returns the guest id verbatim — caller propagates it."""
|
||||||
|
canned = [
|
||||||
|
json.dumps(
|
||||||
|
{
|
||||||
|
"addressee_id": "bot_b",
|
||||||
|
"confidence": "high",
|
||||||
|
"reason": "user named BotB",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
]
|
||||||
|
client = MockLLMClient(canned=canned)
|
||||||
|
|
||||||
|
result = await detect_addressee(
|
||||||
|
client,
|
||||||
|
classifier_model="test-model",
|
||||||
|
user_prose="BotB, what do you think?",
|
||||||
|
host_id="bot_a",
|
||||||
|
host_name="BotA",
|
||||||
|
guest_id="bot_b",
|
||||||
|
guest_name="BotB",
|
||||||
|
)
|
||||||
|
|
||||||
|
assert isinstance(result, AddresseeDecision)
|
||||||
|
assert result.addressee_id == "bot_b"
|
||||||
|
assert result.confidence == "high"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_classifier_picks_host():
|
||||||
|
"""Classifier returns the host id — caller propagates it."""
|
||||||
|
canned = [
|
||||||
|
json.dumps(
|
||||||
|
{
|
||||||
|
"addressee_id": "bot_a",
|
||||||
|
"confidence": "medium",
|
||||||
|
"reason": "narration aimed at host",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
]
|
||||||
|
client = MockLLMClient(canned=canned)
|
||||||
|
|
||||||
|
result = await detect_addressee(
|
||||||
|
client,
|
||||||
|
classifier_model="test-model",
|
||||||
|
user_prose="I lean back and stretch.",
|
||||||
|
host_id="bot_a",
|
||||||
|
host_name="BotA",
|
||||||
|
guest_id="bot_b",
|
||||||
|
guest_name="BotB",
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result.addressee_id == "bot_a"
|
||||||
|
assert result.confidence == "medium"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_classifier_failure_falls_back_to_host():
|
||||||
|
"""Three bad-JSON responses exhaust the retry budget and the
|
||||||
|
classifier-failure fallback returns ``host_id`` with
|
||||||
|
``reason="fallback"``."""
|
||||||
|
canned = ["not json", "still not json", "garbage"]
|
||||||
|
client = MockLLMClient(canned=canned)
|
||||||
|
|
||||||
|
result = await detect_addressee(
|
||||||
|
client,
|
||||||
|
classifier_model="test-model",
|
||||||
|
user_prose="anything",
|
||||||
|
host_id="bot_a",
|
||||||
|
host_name="BotA",
|
||||||
|
guest_id="bot_b",
|
||||||
|
guest_name="BotB",
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result.addressee_id == "bot_a"
|
||||||
|
assert result.reason == "fallback"
|
||||||
|
assert result.confidence == "low"
|
||||||
@@ -271,3 +271,394 @@ def test_regenerate_404_when_assistant_turn_missing(client, tmp_path):
|
|||||||
assert response.status_code == 404
|
assert response.status_code == 404
|
||||||
finally:
|
finally:
|
||||||
app.dependency_overrides.clear()
|
app.dependency_overrides.clear()
|
||||||
|
|
||||||
|
|
||||||
|
def _seed_with_interjection_group(db_path):
|
||||||
|
"""Seed a multi-entity scene with a (primary + interjection) group.
|
||||||
|
|
||||||
|
Returns ``(user_turn_id, primary_at_id, interjection_at_id)``.
|
||||||
|
|
||||||
|
The primary speaker is the host (bot_a); the silent witness who
|
||||||
|
interjected is the guest (bot_b). Mirrors the convention in
|
||||||
|
chat/web/turns.py — both assistant_turns share the same
|
||||||
|
``user_turn_id`` and the interjection's payload carries
|
||||||
|
``interjection_of=<primary speaker_id>``.
|
||||||
|
"""
|
||||||
|
with open_db(db_path) as conn:
|
||||||
|
for bot_id, name, persona in (
|
||||||
|
("bot_a", "BotA", "thoughtful"),
|
||||||
|
("bot_b", "BotB", "loud"),
|
||||||
|
):
|
||||||
|
append_event(
|
||||||
|
conn,
|
||||||
|
kind="bot_authored",
|
||||||
|
payload={
|
||||||
|
"id": bot_id,
|
||||||
|
"name": name,
|
||||||
|
"persona": persona,
|
||||||
|
"voice_samples": [],
|
||||||
|
"traits": [],
|
||||||
|
"backstory": "",
|
||||||
|
"initial_relationship_to_you": "",
|
||||||
|
"kickoff_prose": "",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
append_event(
|
||||||
|
conn,
|
||||||
|
kind="chat_created",
|
||||||
|
payload={
|
||||||
|
"id": "chat_multi",
|
||||||
|
"host_bot_id": "bot_a",
|
||||||
|
"guest_bot_id": "bot_b",
|
||||||
|
"initial_time": "2026-04-26T20:00:00+00:00",
|
||||||
|
"narrative_anchor": "Day 1",
|
||||||
|
"weather": "",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
for src, tgt in (
|
||||||
|
("bot_a", "you"),
|
||||||
|
("you", "bot_a"),
|
||||||
|
("bot_b", "you"),
|
||||||
|
("you", "bot_b"),
|
||||||
|
("bot_a", "bot_b"),
|
||||||
|
("bot_b", "bot_a"),
|
||||||
|
):
|
||||||
|
append_event(
|
||||||
|
conn,
|
||||||
|
kind="edge_update",
|
||||||
|
payload={
|
||||||
|
"source_id": src,
|
||||||
|
"target_id": tgt,
|
||||||
|
"chat_id": "chat_multi",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
for entity_id in ("you", "bot_a", "bot_b"):
|
||||||
|
append_event(
|
||||||
|
conn,
|
||||||
|
kind="activity_change",
|
||||||
|
payload={
|
||||||
|
"entity_id": entity_id,
|
||||||
|
"posture": "sitting",
|
||||||
|
"action": {"verb": "talking"},
|
||||||
|
"attention": "",
|
||||||
|
"holding": [],
|
||||||
|
"status": {},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
ut_id = append_event(
|
||||||
|
conn,
|
||||||
|
kind="user_turn",
|
||||||
|
payload={
|
||||||
|
"chat_id": "chat_multi",
|
||||||
|
"prose": "hello",
|
||||||
|
"segments": [],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
primary_id = append_event(
|
||||||
|
conn,
|
||||||
|
kind="assistant_turn",
|
||||||
|
payload={
|
||||||
|
"chat_id": "chat_multi",
|
||||||
|
"speaker_id": "bot_a",
|
||||||
|
"text": "Original primary.",
|
||||||
|
"truncated": False,
|
||||||
|
"user_turn_id": ut_id,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
interjection_id = append_event(
|
||||||
|
conn,
|
||||||
|
kind="assistant_turn",
|
||||||
|
payload={
|
||||||
|
"chat_id": "chat_multi",
|
||||||
|
"speaker_id": "bot_b",
|
||||||
|
"text": "Original interjection!",
|
||||||
|
"truncated": False,
|
||||||
|
"user_turn_id": ut_id,
|
||||||
|
"interjection_of": "bot_a",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
project(conn)
|
||||||
|
return ut_id, primary_id, interjection_id
|
||||||
|
|
||||||
|
|
||||||
|
def test_regenerate_broadcasts_turn_html_over_sse(
|
||||||
|
tmp_path, monkeypatch
|
||||||
|
):
|
||||||
|
"""T73.1: regenerate publishes a ``turn_html_replace`` SSE event so
|
||||||
|
connected tabs swap the prior turn's DOM node in place.
|
||||||
|
|
||||||
|
The event carries:
|
||||||
|
- ``data``: rendered HTML for the new turn
|
||||||
|
- ``turn_id``: event_id of the new assistant_turn
|
||||||
|
- ``supersedes_id``: event_id of the original assistant_turn
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from chat.config import Settings
|
||||||
|
from chat.db.migrate import apply_migrations
|
||||||
|
from chat.services import regenerate as regenerate_module
|
||||||
|
from chat.services.regenerate import regenerate_assistant_turn
|
||||||
|
|
||||||
|
db_path = tmp_path / "test.db"
|
||||||
|
cfg = tmp_path / "config.toml"
|
||||||
|
cfg.write_text('featherless_api_key = "test"\n')
|
||||||
|
monkeypatch.setenv("CHAT_CONFIG_PATH", str(cfg))
|
||||||
|
monkeypatch.setenv("CHAT_DB_PATH", str(db_path))
|
||||||
|
apply_migrations(db_path)
|
||||||
|
|
||||||
|
ut_id, at_id = _seed_with_one_turn(db_path)
|
||||||
|
|
||||||
|
published: list[tuple[str, dict]] = []
|
||||||
|
|
||||||
|
async def _capture(chat_id, event):
|
||||||
|
published.append((chat_id, event))
|
||||||
|
|
||||||
|
# Patch the imported reference inside the regenerate module so the
|
||||||
|
# service's call site goes through our spy.
|
||||||
|
monkeypatch.setattr(regenerate_module, "publish", _capture)
|
||||||
|
|
||||||
|
narrative_canned = "Refreshed reply."
|
||||||
|
state_canned = json.dumps(
|
||||||
|
{"affinity_delta": 0, "trust_delta": 0, "knowledge_facts": []}
|
||||||
|
)
|
||||||
|
canned = [narrative_canned, state_canned, state_canned]
|
||||||
|
mock_client = MockLLMClient(canned=list(canned))
|
||||||
|
|
||||||
|
settings = Settings(featherless_api_key="test")
|
||||||
|
|
||||||
|
with open_db(db_path) as conn:
|
||||||
|
new_text = asyncio.run(
|
||||||
|
regenerate_assistant_turn(
|
||||||
|
conn,
|
||||||
|
mock_client,
|
||||||
|
settings=settings,
|
||||||
|
chat_id="chat_bot_a",
|
||||||
|
original_assistant_event_id=at_id,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert new_text == narrative_canned
|
||||||
|
|
||||||
|
# Find the new assistant_turn event_id for cross-checking.
|
||||||
|
cur = conn.execute(
|
||||||
|
"SELECT id FROM event_log "
|
||||||
|
"WHERE kind = 'assistant_turn' AND id != ? "
|
||||||
|
"AND superseded_by IS NULL",
|
||||||
|
(at_id,),
|
||||||
|
).fetchone()
|
||||||
|
new_at_id = cur[0]
|
||||||
|
|
||||||
|
# Filter out per-token publishes; we want the replace broadcast.
|
||||||
|
replace_calls = [
|
||||||
|
ev for (_cid, ev) in published if ev.get("event") == "turn_html_replace"
|
||||||
|
]
|
||||||
|
assert len(replace_calls) == 1
|
||||||
|
payload = replace_calls[0]
|
||||||
|
assert payload["supersedes_id"] == at_id
|
||||||
|
assert payload["turn_id"] == new_at_id
|
||||||
|
# The HTML carries the new narrative text and the speaker name.
|
||||||
|
assert "Refreshed reply." in payload["data"]
|
||||||
|
assert "BotA" in payload["data"]
|
||||||
|
# Sanity: every publish targeted this chat.
|
||||||
|
for cid, _ev in published:
|
||||||
|
assert cid == "chat_bot_a"
|
||||||
|
|
||||||
|
|
||||||
|
def test_regenerate_with_interjection_redoes_both_turns(tmp_path, monkeypatch):
|
||||||
|
"""T73.2: when the original turn group included an interjection, both
|
||||||
|
the primary and the interjection are regenerated.
|
||||||
|
|
||||||
|
Setup: 3-entity scene (host BotA + guest BotB + you) with a prior
|
||||||
|
(primary by BotA + interjection by BotB) group. Mock the
|
||||||
|
interjection classifier to return ``should_interject=True`` so the
|
||||||
|
follow-on regenerates too.
|
||||||
|
|
||||||
|
Assert: 2 new assistant_turns exist for the same user_turn_id, the
|
||||||
|
second carrying ``interjection_of`` pointing at the new primary's
|
||||||
|
speaker_id. Both originals are superseded.
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from chat.config import Settings
|
||||||
|
from chat.db.migrate import apply_migrations
|
||||||
|
from chat.services import regenerate as regenerate_module
|
||||||
|
from chat.services.interjection import InterjectionDecision
|
||||||
|
from chat.services.regenerate import regenerate_assistant_turn
|
||||||
|
|
||||||
|
db_path = tmp_path / "test.db"
|
||||||
|
cfg = tmp_path / "config.toml"
|
||||||
|
cfg.write_text('featherless_api_key = "test"\n')
|
||||||
|
monkeypatch.setenv("CHAT_CONFIG_PATH", str(cfg))
|
||||||
|
monkeypatch.setenv("CHAT_DB_PATH", str(db_path))
|
||||||
|
apply_migrations(db_path)
|
||||||
|
|
||||||
|
ut_id, primary_id, interjection_id = _seed_with_interjection_group(db_path)
|
||||||
|
|
||||||
|
# Stub detect_interjection so the classifier "fires" with new prose.
|
||||||
|
async def _stub_should_interject(*_args, **_kwargs):
|
||||||
|
return InterjectionDecision(should_interject=True, reason="fired")
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
regenerate_module, "detect_interjection", _stub_should_interject
|
||||||
|
)
|
||||||
|
|
||||||
|
# Canned queue:
|
||||||
|
# 1. New primary narrative stream.
|
||||||
|
# 2-7. Six state-update classifier calls (one per directed pair
|
||||||
|
# across host/you/guest = 6 pairs) for the primary pass.
|
||||||
|
# 8. New interjection narrative stream.
|
||||||
|
# 9-14. Six state-update classifier calls for the post-interjection
|
||||||
|
# pass.
|
||||||
|
state_canned = json.dumps(
|
||||||
|
{"affinity_delta": 0, "trust_delta": 0, "knowledge_facts": []}
|
||||||
|
)
|
||||||
|
canned: list[str] = []
|
||||||
|
canned.append("New primary text.")
|
||||||
|
canned.extend([state_canned] * 6)
|
||||||
|
canned.append("New interjection text!")
|
||||||
|
canned.extend([state_canned] * 6)
|
||||||
|
mock_client = MockLLMClient(canned=list(canned))
|
||||||
|
|
||||||
|
settings = Settings(featherless_api_key="test")
|
||||||
|
|
||||||
|
with open_db(db_path) as conn:
|
||||||
|
new_text = asyncio.run(
|
||||||
|
regenerate_assistant_turn(
|
||||||
|
conn,
|
||||||
|
mock_client,
|
||||||
|
settings=settings,
|
||||||
|
chat_id="chat_multi",
|
||||||
|
original_assistant_event_id=primary_id,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert new_text == "New primary text."
|
||||||
|
|
||||||
|
# Both originals are superseded.
|
||||||
|
primary_super = conn.execute(
|
||||||
|
"SELECT superseded_by FROM event_log WHERE id = ?", (primary_id,)
|
||||||
|
).fetchone()[0]
|
||||||
|
interjection_super = conn.execute(
|
||||||
|
"SELECT superseded_by FROM event_log WHERE id = ?",
|
||||||
|
(interjection_id,),
|
||||||
|
).fetchone()[0]
|
||||||
|
assert primary_super is not None
|
||||||
|
assert interjection_super is not None
|
||||||
|
|
||||||
|
# Two NEW assistant_turn events exist (the regenerated primary
|
||||||
|
# and the regenerated interjection), both pinned to the same
|
||||||
|
# user_turn_id as the originals.
|
||||||
|
cur = conn.execute(
|
||||||
|
"SELECT id, payload_json FROM event_log "
|
||||||
|
"WHERE kind = 'assistant_turn' AND id NOT IN (?, ?) "
|
||||||
|
"ORDER BY id",
|
||||||
|
(primary_id, interjection_id),
|
||||||
|
).fetchall()
|
||||||
|
assert len(cur) == 2
|
||||||
|
new_primary_id, new_primary_payload_json = cur[0]
|
||||||
|
new_interjection_id, new_interjection_payload_json = cur[1]
|
||||||
|
new_primary_payload = json.loads(new_primary_payload_json)
|
||||||
|
new_interjection_payload = json.loads(new_interjection_payload_json)
|
||||||
|
|
||||||
|
assert new_primary_payload["text"] == "New primary text."
|
||||||
|
assert new_primary_payload["speaker_id"] == "bot_a"
|
||||||
|
assert new_primary_payload["user_turn_id"] == ut_id
|
||||||
|
assert new_primary_payload["regenerated_from"] == primary_id
|
||||||
|
assert "interjection_of" not in new_primary_payload
|
||||||
|
|
||||||
|
assert new_interjection_payload["text"] == "New interjection text!"
|
||||||
|
assert new_interjection_payload["speaker_id"] == "bot_b"
|
||||||
|
assert new_interjection_payload["user_turn_id"] == ut_id
|
||||||
|
assert new_interjection_payload["regenerated_from"] == interjection_id
|
||||||
|
# interjection_of links to the new primary's speaker (matches
|
||||||
|
# the existing convention in chat/web/turns.py).
|
||||||
|
assert new_interjection_payload["interjection_of"] == "bot_a"
|
||||||
|
|
||||||
|
# The originals' supersede pointers reach the new ones.
|
||||||
|
assert primary_super == new_primary_id
|
||||||
|
assert interjection_super == new_interjection_id
|
||||||
|
|
||||||
|
|
||||||
|
def test_regenerate_drops_interjection_when_classifier_returns_false(
|
||||||
|
tmp_path, monkeypatch
|
||||||
|
):
|
||||||
|
"""T73.2: when the original group included an interjection but the
|
||||||
|
classifier returns False this time, the new group is primary-only.
|
||||||
|
|
||||||
|
The original interjection is still superseded (we don't leave it
|
||||||
|
visible in the timeline alongside a regenerated primary it no longer
|
||||||
|
follows from), but no replacement assistant_turn is appended.
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from chat.config import Settings
|
||||||
|
from chat.db.migrate import apply_migrations
|
||||||
|
from chat.services import regenerate as regenerate_module
|
||||||
|
from chat.services.interjection import InterjectionDecision
|
||||||
|
from chat.services.regenerate import regenerate_assistant_turn
|
||||||
|
|
||||||
|
db_path = tmp_path / "test.db"
|
||||||
|
cfg = tmp_path / "config.toml"
|
||||||
|
cfg.write_text('featherless_api_key = "test"\n')
|
||||||
|
monkeypatch.setenv("CHAT_CONFIG_PATH", str(cfg))
|
||||||
|
monkeypatch.setenv("CHAT_DB_PATH", str(db_path))
|
||||||
|
apply_migrations(db_path)
|
||||||
|
|
||||||
|
ut_id, primary_id, interjection_id = _seed_with_interjection_group(db_path)
|
||||||
|
|
||||||
|
async def _stub_no_interject(*_args, **_kwargs):
|
||||||
|
return InterjectionDecision(
|
||||||
|
should_interject=False, reason="quiet"
|
||||||
|
)
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
regenerate_module, "detect_interjection", _stub_no_interject
|
||||||
|
)
|
||||||
|
|
||||||
|
# Canned queue: primary narrative + 6 state-update calls. No
|
||||||
|
# interjection stream because the classifier short-circuits.
|
||||||
|
state_canned = json.dumps(
|
||||||
|
{"affinity_delta": 0, "trust_delta": 0, "knowledge_facts": []}
|
||||||
|
)
|
||||||
|
canned: list[str] = ["New primary text."] + [state_canned] * 6
|
||||||
|
mock_client = MockLLMClient(canned=list(canned))
|
||||||
|
|
||||||
|
settings = Settings(featherless_api_key="test")
|
||||||
|
|
||||||
|
with open_db(db_path) as conn:
|
||||||
|
new_text = asyncio.run(
|
||||||
|
regenerate_assistant_turn(
|
||||||
|
conn,
|
||||||
|
mock_client,
|
||||||
|
settings=settings,
|
||||||
|
chat_id="chat_multi",
|
||||||
|
original_assistant_event_id=primary_id,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert new_text == "New primary text."
|
||||||
|
|
||||||
|
# Original primary superseded by the new primary.
|
||||||
|
primary_super = conn.execute(
|
||||||
|
"SELECT superseded_by FROM event_log WHERE id = ?", (primary_id,)
|
||||||
|
).fetchone()[0]
|
||||||
|
# Original interjection ALSO superseded — we don't leave a
|
||||||
|
# dangling beat attached to a regenerated primary that no longer
|
||||||
|
# warrants a follow-on. Back-pointer goes to the new primary.
|
||||||
|
interjection_super = conn.execute(
|
||||||
|
"SELECT superseded_by FROM event_log WHERE id = ?",
|
||||||
|
(interjection_id,),
|
||||||
|
).fetchone()[0]
|
||||||
|
assert primary_super is not None
|
||||||
|
assert interjection_super is not None
|
||||||
|
assert interjection_super == primary_super # both point at new primary
|
||||||
|
|
||||||
|
# Exactly ONE new assistant_turn — the primary; no replacement
|
||||||
|
# interjection.
|
||||||
|
cur = conn.execute(
|
||||||
|
"SELECT payload_json FROM event_log "
|
||||||
|
"WHERE kind = 'assistant_turn' AND id NOT IN (?, ?) "
|
||||||
|
"AND superseded_by IS NULL",
|
||||||
|
(primary_id, interjection_id),
|
||||||
|
).fetchall()
|
||||||
|
assert len(cur) == 1
|
||||||
|
new_primary_payload = json.loads(cur[0][0])
|
||||||
|
assert new_primary_payload["text"] == "New primary text."
|
||||||
|
assert "interjection_of" not in new_primary_payload
|
||||||
|
|||||||
+236
-26
@@ -405,14 +405,15 @@ def test_multi_bot_turn_no_interjection(app_state_setup, tmp_path):
|
|||||||
1 user_turn + 1 assistant_turn + 6 *post-turn* edge_updates + 2
|
1 user_turn + 1 assistant_turn + 6 *post-turn* edge_updates + 2
|
||||||
memory_written events. Single turn_html broadcast.
|
memory_written events. Single turn_html broadcast.
|
||||||
|
|
||||||
Canned queue (8 calls):
|
Canned queue (11 calls):
|
||||||
1. parse_turn
|
1. parse_turn
|
||||||
2. narrative stream (primary, addressee = host because the prose
|
2. detect_addressee (T74.1) -> host
|
||||||
|
3. narrative stream (primary, addressee = host because the prose
|
||||||
doesn't name the guest)
|
doesn't name the guest)
|
||||||
3-8. 6 state-update calls (one per directed pair across {you,
|
4-9. 6 state-update calls (one per directed pair across {you,
|
||||||
bot_a, bot_b})
|
bot_a, bot_b})
|
||||||
9. detect_interjection -> should_interject=False
|
10. detect_interjection -> should_interject=False
|
||||||
10. detect_scene_close -> should_close=False
|
11. detect_scene_close -> should_close=False
|
||||||
"""
|
"""
|
||||||
_seed_chat_with_guest(tmp_path / "test.db")
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
canned_parse = json.dumps(
|
canned_parse = json.dumps(
|
||||||
@@ -420,6 +421,9 @@ def test_multi_bot_turn_no_interjection(app_state_setup, tmp_path):
|
|||||||
)
|
)
|
||||||
canned = [
|
canned = [
|
||||||
canned_parse,
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{"addressee_id": "bot_a", "confidence": "medium", "reason": "host"}
|
||||||
|
),
|
||||||
"Greetings.",
|
"Greetings.",
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
@@ -474,14 +478,15 @@ def test_multi_bot_turn_with_interjection(app_state_setup, tmp_path):
|
|||||||
1 user_turn + 2 assistant_turns + (6 + 6) post-turn edge_updates +
|
1 user_turn + 2 assistant_turns + (6 + 6) post-turn edge_updates +
|
||||||
4 memory_written events.
|
4 memory_written events.
|
||||||
|
|
||||||
Canned queue (16 calls):
|
Canned queue (17 calls):
|
||||||
1. parse_turn
|
1. parse_turn
|
||||||
2. narrative stream (primary)
|
2. detect_addressee (T74.1) -> host
|
||||||
3-8. 6 state-update calls (post-primary)
|
3. narrative stream (primary)
|
||||||
9. detect_interjection -> should_interject=True
|
4-9. 6 state-update calls (post-primary)
|
||||||
10. narrative stream (interjection)
|
10. detect_interjection -> should_interject=True
|
||||||
11-16. 6 state-update calls (post-interjection)
|
11. narrative stream (interjection)
|
||||||
17. detect_scene_close -> should_close=False
|
12-17. 6 state-update calls (post-interjection)
|
||||||
|
18. detect_scene_close -> should_close=False
|
||||||
"""
|
"""
|
||||||
_seed_chat_with_guest(tmp_path / "test.db")
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
canned_parse = json.dumps(
|
canned_parse = json.dumps(
|
||||||
@@ -489,6 +494,9 @@ def test_multi_bot_turn_with_interjection(app_state_setup, tmp_path):
|
|||||||
)
|
)
|
||||||
canned = [
|
canned = [
|
||||||
canned_parse,
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{"addressee_id": "bot_a", "confidence": "medium", "reason": "host"}
|
||||||
|
),
|
||||||
"Primary beat.",
|
"Primary beat.",
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
@@ -555,14 +563,15 @@ def test_multi_bot_turn_scene_close_writes_per_pov_summaries(
|
|||||||
rewrites fire for both bots (memory.pov_summary changes for each).
|
rewrites fire for both bots (memory.pov_summary changes for each).
|
||||||
Interjection short-circuits at False so the queue stays compact.
|
Interjection short-circuits at False so the queue stays compact.
|
||||||
|
|
||||||
Canned queue (12 calls):
|
Canned queue (13 calls):
|
||||||
1. parse_turn
|
1. parse_turn
|
||||||
2. narrative stream (primary)
|
2. detect_addressee (T74.1) -> host
|
||||||
3-8. 6 state-update calls
|
3. narrative stream (primary)
|
||||||
9. detect_interjection -> False (no follow-on stream)
|
4-9. 6 state-update calls
|
||||||
10. detect_scene_close -> True
|
10. detect_interjection -> False (no follow-on stream)
|
||||||
11. apply_scene_close_summary host POV
|
11. detect_scene_close -> True
|
||||||
12. apply_scene_close_summary guest POV
|
12. apply_scene_close_summary host POV
|
||||||
|
13. apply_scene_close_summary guest POV
|
||||||
"""
|
"""
|
||||||
_seed_chat_with_guest(tmp_path / "test.db")
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
canned_parse = json.dumps(
|
canned_parse = json.dumps(
|
||||||
@@ -588,6 +597,9 @@ def test_multi_bot_turn_scene_close_writes_per_pov_summaries(
|
|||||||
)
|
)
|
||||||
canned = [
|
canned = [
|
||||||
canned_parse,
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{"addressee_id": "bot_a", "confidence": "medium", "reason": "host"}
|
||||||
|
),
|
||||||
"Goodnight.",
|
"Goodnight.",
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
@@ -639,12 +651,20 @@ def test_multi_bot_turn_scene_close_writes_per_pov_summaries(
|
|||||||
|
|
||||||
|
|
||||||
def test_addressee_detection_routes_to_named_bot(app_state_setup, tmp_path):
|
def test_addressee_detection_routes_to_named_bot(app_state_setup, tmp_path):
|
||||||
"""Prose that names the guest by name routes the primary turn to the
|
"""T74.1: the multi-entity addressee call goes through the classifier;
|
||||||
guest. Interjection (when fired) makes the host the silent witness
|
when the classifier returns the guest, the primary turn routes there.
|
||||||
and the second assistant_turn carries the host as speaker.
|
Interjection (when fired) makes the host the silent witness and the
|
||||||
|
second assistant_turn carries the host as speaker.
|
||||||
|
|
||||||
Canned queue: same shape as the with-interjection test (16 calls)
|
Canned queue (with classifier-led addressee = guest):
|
||||||
plus the trailing scene_close decision.
|
1. parse_turn
|
||||||
|
2. detect_addressee -> bot_b (the guest)
|
||||||
|
3. narrative stream (primary, addressee = guest)
|
||||||
|
4-9. 6 state-update calls
|
||||||
|
10. detect_interjection -> True
|
||||||
|
11. interjection narrative stream
|
||||||
|
12-17. 6 state-update calls (post-interjection)
|
||||||
|
18. detect_scene_close -> False
|
||||||
"""
|
"""
|
||||||
_seed_chat_with_guest(tmp_path / "test.db")
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
canned_parse = json.dumps(
|
canned_parse = json.dumps(
|
||||||
@@ -652,6 +672,13 @@ def test_addressee_detection_routes_to_named_bot(app_state_setup, tmp_path):
|
|||||||
)
|
)
|
||||||
canned = [
|
canned = [
|
||||||
canned_parse,
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{
|
||||||
|
"addressee_id": "bot_b",
|
||||||
|
"confidence": "high",
|
||||||
|
"reason": "user named BotB",
|
||||||
|
}
|
||||||
|
),
|
||||||
"BotB pondering.",
|
"BotB pondering.",
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
_zero_state(), _zero_state(), _zero_state(),
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
@@ -680,9 +707,192 @@ def test_addressee_detection_routes_to_named_bot(app_state_setup, tmp_path):
|
|||||||
primary_payload = json.loads(rows[0][0])
|
primary_payload = json.loads(rows[0][0])
|
||||||
interjection_payload = json.loads(rows[1][0])
|
interjection_payload = json.loads(rows[1][0])
|
||||||
|
|
||||||
# Primary speaker is the guest because the prose names BotB and not
|
# Primary speaker is the guest because the addressee classifier
|
||||||
# BotA (case-insensitive whole-word match).
|
# picked bot_b for the prose ("BotB, what do you think?").
|
||||||
assert primary_payload["speaker_id"] == "bot_b"
|
assert primary_payload["speaker_id"] == "bot_b"
|
||||||
# Interjection follow-on goes to the silent witness — the host.
|
# Interjection follow-on goes to the silent witness — the host.
|
||||||
assert interjection_payload["speaker_id"] == "bot_a"
|
assert interjection_payload["speaker_id"] == "bot_a"
|
||||||
assert interjection_payload["interjection_of"] == "bot_b"
|
assert interjection_payload["interjection_of"] == "bot_b"
|
||||||
|
|
||||||
|
|
||||||
|
def test_cancelled_turn_still_closes_scene_when_user_prose_signals_close(
|
||||||
|
app_state_setup, tmp_path
|
||||||
|
):
|
||||||
|
"""T74.3 regression: a cancelled primary stream still triggers scene
|
||||||
|
close when the user prose carries a hard close signal.
|
||||||
|
|
||||||
|
Rationale (also documented in turns.py near the close-detection
|
||||||
|
branch): close detection only consumes the user's prose, which is
|
||||||
|
fully appended to the event_log BEFORE streaming starts. The
|
||||||
|
cancelled bot beat doesn't invalidate the user's intent to close.
|
||||||
|
|
||||||
|
Implementation: install a MockLLMClient whose ``stream`` raises
|
||||||
|
CancelledError on the first iteration. The classifier calls (parse,
|
||||||
|
addressee, scene_close, per-POV summaries) are still served from
|
||||||
|
the canned queue. The post_turn route ultimately re-raises
|
||||||
|
CancelledError after recording the partial — TestClient surfaces
|
||||||
|
that as an exception, so we drive the request inside ``with
|
||||||
|
pytest.raises``. Despite the exception, the scene_closed event
|
||||||
|
must land in the event_log.
|
||||||
|
"""
|
||||||
|
from typing import AsyncIterator, Sequence
|
||||||
|
|
||||||
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
|
canned_parse = json.dumps(
|
||||||
|
{"segments": [{"kind": "narration", "text": "we are done here, fade out"}]}
|
||||||
|
)
|
||||||
|
pov_payload = json.dumps(
|
||||||
|
{
|
||||||
|
"summary": "BotA noticed the day winding down.",
|
||||||
|
"knowledge_facts": [],
|
||||||
|
"relationship_summary": "warmer",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
pov_payload_guest = json.dumps(
|
||||||
|
{
|
||||||
|
"summary": "BotB watched the scene close.",
|
||||||
|
"knowledge_facts": [],
|
||||||
|
"relationship_summary": "warmer",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
# Canned queue: parse + addressee + 6 state-updates +
|
||||||
|
# scene_close=True + 2 per-POV summaries. NO interjection slot
|
||||||
|
# because the cancel path short-circuits the interjection branch.
|
||||||
|
canned = [
|
||||||
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{"addressee_id": "bot_a", "confidence": "medium", "reason": "host"}
|
||||||
|
),
|
||||||
|
# NOTE: no narrative slot — the stream is hijacked below to
|
||||||
|
# raise CancelledError on first iteration; it never pulls a
|
||||||
|
# canned response.
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
json.dumps({"should_close": True, "reason": "fade out signaled"}),
|
||||||
|
pov_payload,
|
||||||
|
pov_payload_guest,
|
||||||
|
]
|
||||||
|
|
||||||
|
class _CancelOnStreamMock:
|
||||||
|
"""Mock LLM client that serves ``generate`` from a canned queue
|
||||||
|
and raises CancelledError on the FIRST iteration of ``stream``.
|
||||||
|
|
||||||
|
Mirrors :class:`chat.llm.mock.MockLLMClient` for ``generate`` but
|
||||||
|
diverges on ``stream`` to simulate a mid-stream cancel.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, canned: list[str]) -> None:
|
||||||
|
self._canned = list(canned)
|
||||||
|
|
||||||
|
async def generate(
|
||||||
|
self, messages: Sequence, *, model: str, **params
|
||||||
|
) -> str:
|
||||||
|
return self._canned.pop(0)
|
||||||
|
|
||||||
|
async def stream(
|
||||||
|
self, messages: Sequence, *, model: str, **params
|
||||||
|
) -> AsyncIterator[str]:
|
||||||
|
# Yield a CancelledError on first iteration to simulate the
|
||||||
|
# /turns/cancel route firing mid-stream.
|
||||||
|
raise asyncio.CancelledError
|
||||||
|
yield # pragma: no cover — keeps this an async generator.
|
||||||
|
|
||||||
|
from chat.web.kickoff import get_llm_client
|
||||||
|
|
||||||
|
mock = _CancelOnStreamMock(canned=list(canned))
|
||||||
|
app.dependency_overrides[get_llm_client] = lambda: mock
|
||||||
|
|
||||||
|
try:
|
||||||
|
# FastAPI/Starlette handles the re-raised CancelledError as an
|
||||||
|
# internal failure — TestClient surfaces it as a 500 response.
|
||||||
|
# We don't assert on the status here; the regression is whether
|
||||||
|
# the scene_closed event still landed in the event_log.
|
||||||
|
try:
|
||||||
|
app_state_setup.post(
|
||||||
|
"/chats/chat_bot_a/turns",
|
||||||
|
data={"prose": "we are done here, fade out"},
|
||||||
|
)
|
||||||
|
except BaseException:
|
||||||
|
# Some Starlette/asyncio versions propagate the
|
||||||
|
# CancelledError out of the test client; that's fine — the
|
||||||
|
# partial-record + scene-close still ran before the raise.
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
app.dependency_overrides.clear()
|
||||||
|
|
||||||
|
with open_db(tmp_path / "test.db") as conn:
|
||||||
|
scene_close_count = conn.execute(
|
||||||
|
"SELECT COUNT(*) FROM event_log WHERE kind = 'scene_closed'"
|
||||||
|
).fetchone()[0]
|
||||||
|
assistant_payload = conn.execute(
|
||||||
|
"SELECT payload_json FROM event_log "
|
||||||
|
"WHERE kind = 'assistant_turn' ORDER BY id"
|
||||||
|
).fetchall()
|
||||||
|
|
||||||
|
# Scene close lands despite the cancel.
|
||||||
|
assert scene_close_count == 1
|
||||||
|
# The cancelled assistant_turn was still recorded (truncated=True).
|
||||||
|
assert len(assistant_payload) == 1
|
||||||
|
assert json.loads(assistant_payload[0][0])["truncated"] is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_interjection_enqueues_significance_job(app_state_setup, tmp_path):
|
||||||
|
"""T74.2: when an interjection fires, the interjection memory is
|
||||||
|
enqueued for significance scoring just like the primary memory.
|
||||||
|
|
||||||
|
Capture enqueued ``SignificanceJob``s by replacing the background
|
||||||
|
worker's ``enqueue`` method with a list-append. Without T74.2, the
|
||||||
|
interjection memory would never be scored — only the primary's
|
||||||
|
enqueue would land. We therefore expect TWO jobs after a turn that
|
||||||
|
has both a primary and an interjection beat: one for the primary
|
||||||
|
memory, one for the interjection memory.
|
||||||
|
"""
|
||||||
|
_seed_chat_with_guest(tmp_path / "test.db")
|
||||||
|
canned_parse = json.dumps(
|
||||||
|
{"segments": [{"kind": "dialogue", "text": "tell me"}]}
|
||||||
|
)
|
||||||
|
canned = [
|
||||||
|
canned_parse,
|
||||||
|
json.dumps(
|
||||||
|
{"addressee_id": "bot_a", "confidence": "medium", "reason": "host"}
|
||||||
|
),
|
||||||
|
"Primary beat.",
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
json.dumps({"should_interject": True, "reason": "jealous"}),
|
||||||
|
"Interjection beat!",
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
_zero_state(), _zero_state(), _zero_state(),
|
||||||
|
json.dumps({"should_close": False, "reason": "no signal"}),
|
||||||
|
]
|
||||||
|
_override_llm(canned)
|
||||||
|
|
||||||
|
captured_jobs: list = []
|
||||||
|
worker = app.state.background_worker
|
||||||
|
# Re-enable enqueue capture even though the worker's loop is disabled
|
||||||
|
# — we want to count enqueues without the loop running classifier work.
|
||||||
|
worker.enabled = True
|
||||||
|
original_enqueue = worker.enqueue
|
||||||
|
worker.enqueue = captured_jobs.append # type: ignore[assignment]
|
||||||
|
|
||||||
|
try:
|
||||||
|
response = app_state_setup.post(
|
||||||
|
"/chats/chat_bot_a/turns", data={"prose": "tell me"}
|
||||||
|
)
|
||||||
|
assert response.status_code == 204
|
||||||
|
finally:
|
||||||
|
worker.enqueue = original_enqueue # type: ignore[assignment]
|
||||||
|
worker.enabled = False
|
||||||
|
app.dependency_overrides.clear()
|
||||||
|
|
||||||
|
# Expect 2 enqueues: 1 for the primary memory + 1 for the
|
||||||
|
# interjection memory.
|
||||||
|
assert len(captured_jobs) == 2
|
||||||
|
|
||||||
|
# Both jobs should reference distinct memory ids — the primary's
|
||||||
|
# host-POV memory and the interjection's host-POV memory.
|
||||||
|
memory_ids = [job.memory_id for job in captured_jobs]
|
||||||
|
assert len(set(memory_ids)) == 2
|
||||||
|
# The two narrative texts should be the two streamed beats.
|
||||||
|
narrative_texts = sorted(job.narrative_text for job in captured_jobs)
|
||||||
|
assert narrative_texts == ["Interjection beat!", "Primary beat."]
|
||||||
|
|||||||
Reference in New Issue
Block a user