A campaign can import local .txt and .md files as Canon, Reference or Inspiration, and the class is load-bearing rather than a label: it decides the words a passage is framed with in the prompt, the weight it carries when passages are ranked, and which budget it competes in when the context is tight. This is a separate subsystem, which is the Phase 0B decision (IMPORTED-KNOWLEDGE-DESIGN.md §73). Story Cards do not carry classification, provenance, content identity, chunking, an index or a lifecycle, and they were not promoted into something that does. Nothing here reads or writes one. The subsystem, in backend/app/knowledge/: classes the three classes, their weights, and the prompt framing chunking deterministic, heading-aware, 60-800 tokens, no overlap fts SQLite FTS5 with porter stemming; scoped and bounded in SQL importer validate, hash, store, chunk, index — in one transaction embeddings local Ollama vectors through the shared provider retrieval query construction, hybrid merge, rerank inject the budgeted cut and the rendered prompt sections Relevance admission is a separate stage from ranking, and that separation is the milestone's most expensive lesson. An independent review found the first implementation deciding relevance with a floor expressed as a share of the best candidate — which the best clears by construction — so a passage was admitted on every turn regardless of the scene. A query about tide tables and container tonnage retrieved all five sources of a fantasy campaign, narrator-only hidden Canon among them. So the pipeline is now: candidate generation -> admission -> ranking -> class weighting -> budget Admission reads raw, candidate-set-independent signals: the cosine the model returned, and how many distinct meaningful query terms a passage contains. Ranking reads normalized ones, because bm25 has no fixed range and cosine's zero is not zero. Normalization decides order among things that matched; it can never decide whether anything matched. Authority is applied after admission, so a class orders what matched and never rescues what did not. Retrieval may therefore return nothing, and on a scene unrelated to the library it does. The other decisions that each replaced an obvious wrong one: - The class multiplies relevance rather than adding to it. An additive bonus satisfies "Canon outranks Reference" and makes "do not include irrelevant Canon" impossible, because a large enough constant wins on its own. - The semantic floor is measured, not guessed: 113 production-path pairs against nomic-embed-text put targeted matches at 0.55-0.85 and off-topic pairs at 0.36-0.56, and 0.58 sits between them. Because it is a property of that model and not of cosine similarity, it is keyed to the model rather than applied to whatever is configured: an embedding model with no measured calibration in this build does not borrow the number. Semantic admission is skipped, the campaign retrieves lexically, and the reason is stated in the knowledge status and in the turn's provenance. Degrading to lexical keeps the library usable; lending the threshold to an unmeasured model is how the admitted-everything defect would return. - One lexical term is not evidence. Two distinct meaningful terms, or one that is neither a standing campaign entity nor a negligible share of the query. The stop list grew from 42 words to 261, all function words — no subject matter, because a stop list that removes subject matter stops finding "The Silver Key". - Lexical retrieval is a production path, not a fallback. It finds the proper nouns and invented terms a setting bible is made of, and the library is fully usable with no embedding model configured. Safety is structural rather than filtered. Imported text reaches the prompt whole, inside a section that says what it is, under a rule stating the authority order in words and refusing every instruction inside it. No endpoint accepts a filesystem path, so H08 has no mechanism to escape from. Nothing renders imported content as HTML, so a script tag is five visible characters and a remote image is never fetched. Import, chunking, indexing, retrieval and a turn open no socket at all; only embeddings do, through the endpoint allowlist the memory bank already uses. Provenance is the rendered text, not a foreign key: deleting a source cannot turn a historical turn's evidence into dangling ids. Schema: knowledge_sources, knowledge_chunks, knowledge_embeddings, and an FTS5 virtual table attached to knowledge_chunks as a DDL hook so it is created and dropped with the table it indexes. Migration 92. A pre-M7 database opens unchanged and needs no sources to play. Bundle: the source content and the reader's judgements about it travel; the passages, index rows and vectors are rebuilt on import, so a restored campaign is searchable immediately without a reindex step. One runtime dependency: python-multipart, Starlette's multipart parser. It is what makes the upload surface possible, and the upload surface is why no pathname is ever accepted. The test doubles were the reason the defect shipped, so they were corrected too. The retrieval stub scored unrelated text at 0.06-0.20 where the real model scores it at 0.43-0.44, and its docstring said it had deliberately removed the constant component that "would put a similarity floor under every pair" — which is exactly the property real models have. The stub now has that floor, one test fails if it is ever removed, and another reproduces the superseded rule and asserts it is still fooled by the same fixture. Run against the pre-corrective implementation, the new suite fails 13 of 18. Tests: 939 passed, 14 skipped (836/7 at M6). 110 new across seven files, one of which mocks nothing between itself and Ollama and re-measures the similarity separation on every run. 43/43 checks in a real Firefox, reproduced. Docker build clean. Four other defects found by review or by the browser run were fixed here rather than carried: an unreachable relevance constant that appeared to enforce something and did not; acceptance tests using the wrong fixture files, so G07's trap was never exercised; a bidirectional override surviving into displayed filenames; and, from the implementation pass, the Insights panel showing M5's two state sections as raw keys and the source inspector refetching on every keystroke. M7 was independently reviewed, which returned PASS WITH CORRECTIVE WORK REQUIRED. Both blocking findings are closed, and closeout resolved the embedding-model calibration boundary the corrective pass had left as debt. planning/reports/M7-IMPLEMENTATION-REPORT.md carries the review, the corrective closeout and the closeout verification in sequence, none overwriting another. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017HdaXiFbscatQaLS7dJk6b
421 lines
18 KiB
Python
421 lines
18 KiB
Python
"""Playing a turn: the model call, the SSE stream, and the one-turn-at-a-time lock.
|
|
|
|
Everything a test needs to intercept lives here, and other modules reach it as
|
|
`turns.<name>` rather than importing it by value. That matters twice. The turn
|
|
lock guards one set only while one module owns it. And a test that replaces
|
|
`OpenAICompatibleProvider` or `generate_turn` patches this module, which every
|
|
caller reads through.
|
|
"""
|
|
import threading
|
|
|
|
from fastapi import Depends, HTTPException, Request
|
|
from fastapi.responses import StreamingResponse
|
|
from sqlalchemy.orm import Session
|
|
|
|
from ... import (
|
|
attempts, head, limits, memorybank, models, narrative, schemas, tree,
|
|
worldstate,
|
|
)
|
|
from ...context import ContextOverflow, build_context, cursors
|
|
from ...knowledge import retrieval as knowledge_retrieval
|
|
from ...database import get_db
|
|
from ...providers import OpenAICompatibleProvider, PromptParts, ProviderError
|
|
from ...sse import SSE_HEADERS, sse, turn_error
|
|
from ..settings import get_settings
|
|
|
|
from .deps import CurrentUser, current_adventure, router
|
|
from .nodes import _move_to_after, next_depth
|
|
from .paging import annotate_takes
|
|
|
|
|
|
def world_delta_of(snapshot: dict | None) -> dict | None:
|
|
"""Returns the bulk-read slice of a context snapshot, for `Action.world_delta`.
|
|
|
|
`context_snapshot` is deferred because it holds the whole assembled prompt.
|
|
The parts that every action needs get their own small column instead: the
|
|
world-change chips, the emit block replayed into history, and the refusal
|
|
note fed back to the model. Update this function wherever a snapshot is
|
|
written.
|
|
|
|
Carry all three report lists, not just `applied`. `Action.world_changes`
|
|
marks a chip from `clamped` and builds its refusal chips from `rejected`,
|
|
and `worldstate.refusals` reads both. Storing `applied` alone left every
|
|
consumer unable to tell a refused change from one that worked, which is the
|
|
distinction this column exists to carry. The two extra lists are subsets of
|
|
one turn's block, so they cost a few hundred bytes per action at most.
|
|
"""
|
|
ws = (snapshot or {}).get("world_state")
|
|
if not isinstance(ws, dict):
|
|
return None
|
|
report = ws.get("report") or {}
|
|
return {
|
|
"delta": ws.get("delta") or {},
|
|
"applied": report.get("applied") or [],
|
|
"clamped": report.get("clamped") or [],
|
|
"rejected": report.get("rejected") or [],
|
|
}
|
|
|
|
|
|
# One turn at a time per adventure. The set lives in memory, which is enough for
|
|
# a single-process local app. Sync endpoints run in a threadpool, so the
|
|
# check-and-add needs a lock. The check also has to run during the request
|
|
# rather than when the SSE generator first runs. Otherwise two rapid requests
|
|
# both pass the check and generate concurrently.
|
|
_active_turns: set[int] = set()
|
|
_active_turns_guard = threading.Lock()
|
|
|
|
|
|
def acquire_turn_lock(adventure_id: int):
|
|
"""Claims the adventure's turn slot atomically. `with_turn_lock` releases it."""
|
|
with _active_turns_guard:
|
|
if adventure_id in _active_turns:
|
|
raise HTTPException(409, "A turn is already generating for this adventure.")
|
|
_active_turns.add(adventure_id)
|
|
|
|
|
|
async def with_turn_lock(adventure_id: int, gen):
|
|
"""Wraps an SSE generator so that it releases the `acquire_turn_lock` lock."""
|
|
try:
|
|
async for event in gen:
|
|
yield event
|
|
finally:
|
|
_active_turns.discard(adventure_id)
|
|
|
|
|
|
def format_player_input(action_type: str, text: str) -> str:
|
|
"""Formats player input the way AI Dungeon does."""
|
|
text = text.strip()
|
|
if action_type == "say":
|
|
text = text.strip('"')
|
|
if text and text[-1] not in ".!?…":
|
|
text += "."
|
|
return f'> You say "{text}"'
|
|
if action_type == "do":
|
|
if text.lower().startswith("you "):
|
|
text = text[4:]
|
|
if text and text[-1] not in ".!?…":
|
|
text += "."
|
|
return f"> You {text}"
|
|
return text # The "story" type is appended as raw text.
|
|
|
|
|
|
def action_json(action: models.Action, db: Session | None = None) -> dict:
|
|
"""Serializes one action for the wire.
|
|
|
|
Passing `db` fills in the pager numbers, and a turn that was just played must
|
|
pass it. The attempt that turn created is often the second one at its
|
|
coordinate, so the message needs a pager that the client cannot infer from a
|
|
count of one. Without `db`, a retry showed no pager until the page reloaded.
|
|
The adventure GET has the same requirement and builds its window a third way.
|
|
"""
|
|
if db is not None:
|
|
annotate_takes(db, action.adventure_id, [action])
|
|
return schemas.ActionOut.model_validate(action).model_dump(mode="json")
|
|
|
|
|
|
async def generate_turn(
|
|
adventure: models.Adventure,
|
|
db: Session,
|
|
user: models.User,
|
|
retry_of: models.Action | None = None,
|
|
):
|
|
"""Streams the AI continuation as SSE, then stores the result.
|
|
|
|
If `retry_of` is set, the result is stored as a sibling of that AI action, at
|
|
the same turn and the same coordinate, and the discarded attempt stays where
|
|
it was written. Before calling, the caller must roll the adventure back to
|
|
the state before that turn. See `retry_action`. If this generator ends
|
|
without saving, the rollback is undone, so the state cannot diverge from the
|
|
text on screen.
|
|
"""
|
|
saved = False
|
|
try:
|
|
async for event in _generate_turn(adventure, db, user, retry_of):
|
|
if event is _SAVED:
|
|
saved = True
|
|
continue
|
|
yield event
|
|
finally:
|
|
if retry_of is not None and not saved:
|
|
# The turn failed with a provider error, an empty reply, or a
|
|
# disconnected client. No sibling was written, so the
|
|
# attempt on screen is still the live one. Restore the state it
|
|
# produced.
|
|
attempts.restore_state(adventure, retry_of)
|
|
db.commit()
|
|
|
|
|
|
# `_generate_turn` yields this sentinel once the action is committed. It tells
|
|
# the wrapper above to leave the rollback in place rather than reverse it.
|
|
_SAVED = object()
|
|
|
|
|
|
async def _generate_turn(
|
|
adventure: models.Adventure,
|
|
db: Session,
|
|
user: models.User,
|
|
retry_of: models.Action | None = None,
|
|
):
|
|
settings = get_settings(db, user)
|
|
# On a retry, the attempt being replaced is still the live node of its turn,
|
|
# because it stays live until a replacement exists. Filter it out of the
|
|
# context. Otherwise the model reads the attempt it is replacing as
|
|
# established story and writes a sequel to it.
|
|
replacing_id = retry_of.id if retry_of is not None else None
|
|
memories = await memorybank.retrieve_memories(
|
|
adventure, settings, update_stats=True, exclude_action_id=replacing_id
|
|
)
|
|
# M7: the imported library, retrieved for the position being read. Excluding
|
|
# the attempt being replaced matters here for the same reason it does for
|
|
# memories — the query is built from the recent story, and a discarded
|
|
# attempt must not steer which passages the replacement is given.
|
|
knowledge = await knowledge_retrieval.retrieve(
|
|
adventure, settings, exclude_action_id=replacing_id
|
|
)
|
|
try:
|
|
system_text, story_text, snapshot = build_context(
|
|
adventure,
|
|
settings,
|
|
memories,
|
|
exclude_action_id=replacing_id,
|
|
knowledge=knowledge,
|
|
)
|
|
except ContextOverflow as exc:
|
|
# M6: the protected context does not fit in the configured budget, so
|
|
# there is no prompt to send. This is a settings problem the reader can
|
|
# fix, and the message says how — reporting it as a failed turn keeps
|
|
# the story intact and tells them what to change, where building the
|
|
# prompt anyway would return a silently truncated reply.
|
|
yield turn_error(str(exc))
|
|
return
|
|
|
|
parts = PromptParts(system=system_text, story=story_text)
|
|
|
|
provider = OpenAICompatibleProvider(
|
|
settings.endpoint_url, settings.model, settings.api_mode,
|
|
settings.model_timeout_seconds,
|
|
)
|
|
chunks: list[str] = []
|
|
reasoning_chunks: list[str] = []
|
|
try:
|
|
async for kind, chunk in provider.generate(
|
|
parts, temperature=settings.temperature, max_tokens=settings.max_output_tokens
|
|
):
|
|
if kind == "reasoning":
|
|
reasoning_chunks.append(chunk)
|
|
yield sse({"type": "reasoning", "text": chunk})
|
|
else:
|
|
chunks.append(chunk)
|
|
yield sse({"type": "chunk", "text": chunk})
|
|
except ProviderError as exc:
|
|
yield turn_error(str(exc))
|
|
return
|
|
|
|
text = "".join(chunks).strip()
|
|
# The model's literal reply, kept for the Insights "Raw AI output" view. It
|
|
# still contains the world-state block, which the code below strips.
|
|
raw_output = text
|
|
if not text:
|
|
# The model streamed reasoning but no story text, so it spent its whole
|
|
# budget on reasoning. Report that rather than "empty response".
|
|
if reasoning_chunks:
|
|
detail = (
|
|
"The model used its entire token budget on reasoning and returned no "
|
|
'story text. Raise "Max output tokens" in Settings, set a "Reasoning '
|
|
'max tokens" cap, or switch to a non-reasoning model.'
|
|
)
|
|
else:
|
|
detail = "The AI returned an empty response."
|
|
yield turn_error(detail)
|
|
return
|
|
|
|
# M5: read the typed state proposal out of the reply, validate it, apply
|
|
# what survives, and strip the block from the displayed text.
|
|
#
|
|
# This replaced the Phase 12 relative-delta pipeline. The shape of the turn
|
|
# is unchanged — extract, referee, snapshot — because ADR 010 changed the
|
|
# protocol, not the lifecycle. What changed is that the referee now works on
|
|
# explicit typed events with absolute values, so an accepted proposal cannot
|
|
# mean something other than it says.
|
|
#
|
|
# A retry re-runs the same turn, so it is played at that turn's depth. This
|
|
# was `retry_of.index`, which held the same number until SP4. Depth stays
|
|
# correct once a branch has its own numbering.
|
|
ai_depth = retry_of.depth if retry_of is not None else next_depth(adventure)
|
|
|
|
text, parsed, raw_block = narrative.extract.split(text)
|
|
if not text.strip():
|
|
yield turn_error("The AI returned only a state update and no story text.")
|
|
return
|
|
review = narrative.validate.review(
|
|
parsed if parsed is not None else {"events": []},
|
|
narrative.store.current(adventure),
|
|
narrative.store.canon_of(adventure),
|
|
)
|
|
# Held until the action exists, because a proposal record names the node
|
|
# whose narration produced it and the node has no id yet. Everything lands
|
|
# in the single commit below (L01).
|
|
# The coordinate is read off the node after it is placed, not guessed here:
|
|
# `tree.place_action` assigns the branch, and a retry inherits the branch of
|
|
# the attempt it replaces.
|
|
pending_state = {
|
|
"review": review,
|
|
"parsed": parsed,
|
|
"raw_block": raw_block,
|
|
"unparseable": parsed is None and bool(raw_block),
|
|
}
|
|
snapshot["narrative_state"] = {
|
|
"accepted": review.accepted,
|
|
"rejected": [r.as_dict() for r in review.rejected],
|
|
"status": review.status,
|
|
}
|
|
|
|
snapshot["raw_output"] = raw_output
|
|
# The cost the endpoint reports for the call, including how much of the
|
|
# prompt came from cache rather than being billed in full. This is recorded
|
|
# per attempt, next to the prompt it priced.
|
|
snapshot["usage"] = provider.last_usage
|
|
|
|
reasoning = "".join(reasoning_chunks).strip() or None
|
|
ai_action = models.Action(
|
|
adventure_id=adventure.id,
|
|
depth=ai_depth,
|
|
type="ai",
|
|
text=text,
|
|
reasoning=reasoning,
|
|
context_snapshot=snapshot,
|
|
world_delta=world_delta_of(snapshot),
|
|
)
|
|
if retry_of is not None:
|
|
attempts.add_attempt(db, adventure, retry_of, ai_action)
|
|
db.add(ai_action)
|
|
# The text at this coordinate changed, so anything derived from it no
|
|
# longer describes the story. Withdraw the memory attached to the node
|
|
# and return that stretch to both passes. Before SP4 this code was
|
|
# unreachable, because the summarizer held the newest action back until
|
|
# a turn landed on top of it. `memorybank.SETTLE_SLACK` keeps a memory
|
|
# off the tip again, for cost rather than for correctness, so this is
|
|
# now the rare case: undo or delete can carry a summarized node back to
|
|
# the tip, and then a retry of it lands here. See `memorybank`.
|
|
memorybank.forget_node(db, adventure, retry_of)
|
|
cursors.rewind_all(adventure, retry_of.branch_id, ai_depth - 1)
|
|
# Flush so the new attempt has an id. The session does not autoflush,
|
|
# and attempts page in id order, so a read taken before this point puts
|
|
# the newest attempt nowhere.
|
|
db.flush()
|
|
else:
|
|
tree.place_action(db, adventure, ai_action)
|
|
db.add(ai_action)
|
|
db.flush()
|
|
# The state lands after the node exists and before the one commit, so the
|
|
# narration, the head, the accepted events, the provenance and the snapshot
|
|
# are one transaction. L01 forbids any window in which a turn looks accepted
|
|
# while its state is half-written, and the cheapest guarantee is to have a
|
|
# single commit rather than two that could get out of step.
|
|
new_state, _proposal = narrative.store.record(
|
|
db, adventure,
|
|
review=pending_state["review"],
|
|
raw_block=pending_state["raw_block"],
|
|
parsed=pending_state["parsed"],
|
|
action=ai_action,
|
|
branch_id=ai_action.branch_id,
|
|
depth=ai_action.depth,
|
|
model_name=settings.model or "",
|
|
source="accepted_story",
|
|
)
|
|
if pending_state["unparseable"]:
|
|
_proposal.status = "unparseable"
|
|
before_state = narrative.store.current(adventure)
|
|
narrative.store.set_current(adventure, new_state)
|
|
ai_action.state_changes = {
|
|
"accepted": pending_state["review"].accepted,
|
|
"rejected": [r.as_dict() for r in pending_state["review"].rejected],
|
|
"summary": narrative.apply.diff(before_state, new_state),
|
|
}
|
|
attempts.snapshot_outcome(adventure, ai_action)
|
|
adventure.updated_at = models.utcnow()
|
|
db.commit()
|
|
db.refresh(ai_action)
|
|
yield _SAVED
|
|
yield sse({"type": "done", "action": action_json(ai_action, db)})
|
|
# Phase 6: schedule summarization and embedding without waiting for them.
|
|
# The task opens its own database session.
|
|
memorybank.schedule_post_turn(adventure)
|
|
|
|
|
|
|
|
async def run_player_turn(
|
|
adventure: models.Adventure,
|
|
db: Session,
|
|
payload: schemas.ActionCreate,
|
|
user: models.User,
|
|
preformatted: bool = False,
|
|
):
|
|
"""Plays a player's turn: their action, then the reply to it.
|
|
|
|
`preformatted` means the text already carries the `> You ...` conventions and
|
|
is written as-is. That applies when the player retakes a turn they already
|
|
played (SP9). The editor is seeded with the stored text, which is already
|
|
formatted, and a plain edit puts that same text in the box and writes it back
|
|
verbatim. Formatting it a second time produces `> You > You ...`.
|
|
"""
|
|
# An empty do, say, or story action behaves as a continue.
|
|
if payload.type != "continue" and payload.text.strip():
|
|
formatted = (
|
|
payload.text.strip() if preformatted
|
|
else format_player_input(payload.type, payload.text)
|
|
)
|
|
player_action = models.Action(
|
|
adventure_id=adventure.id,
|
|
depth=next_depth(adventure),
|
|
type=payload.type,
|
|
text=formatted,
|
|
)
|
|
# The state this node leaves behind. The AI turn after it starts here,
|
|
# and a retry of that turn rolls back to here.
|
|
attempts.snapshot_outcome(adventure, player_action)
|
|
tree.place_action(db, adventure, player_action)
|
|
db.add(player_action)
|
|
db.commit()
|
|
db.refresh(player_action)
|
|
# The new action was added through its foreign key, so the loaded
|
|
# `adventure.actions` collection is stale. Without this expire,
|
|
# `build_context` for the AI action does not see the player action
|
|
# that was just saved.
|
|
db.expire(adventure, ["actions"])
|
|
yield sse({"type": "player", "action": action_json(player_action, db)})
|
|
|
|
async for event in generate_turn(adventure, db, user):
|
|
yield event
|
|
|
|
|
|
@router.post("/{adventure_id}/actions")
|
|
def create_action(
|
|
adventure_id: int,
|
|
payload: schemas.ActionCreate,
|
|
request: Request,
|
|
db: Session = Depends(get_db),
|
|
user: models.User = CurrentUser,
|
|
adventure: models.Adventure = Depends(current_adventure),
|
|
):
|
|
limits.check_row_cap("actions", db, user, adventure=adventure)
|
|
acquire_turn_lock(adventure_id)
|
|
try:
|
|
_move_to_after(db, adventure, payload.after_id)
|
|
# The first write below a moved-back head is where a divergence happens
|
|
# (M3). Undo alone does not fork — the user may be reading, or about to
|
|
# Redo — so this is the moment the story states which continuation it
|
|
# means. The displaced future keeps its rows on the branch being left.
|
|
# A head already at the tip, which is every ordinary turn, forks nothing.
|
|
if head.fork_if_behind_head(db, adventure):
|
|
db.commit()
|
|
db.refresh(adventure)
|
|
except BaseException:
|
|
_active_turns.discard(adventure_id)
|
|
raise
|
|
return StreamingResponse(
|
|
with_turn_lock(adventure_id, run_player_turn(adventure, db, payload, user)),
|
|
media_type="text/event-stream",
|
|
headers=SSE_HEADERS,
|
|
)
|