Replaces AI-DnD's RPG relative-delta world state with the genre-neutral typed
narrative state of ADR 010: explicit, absolute, allowlisted events proposed by
the model, validated by the application, applied to one authoritative document,
and snapshotted per position so restore stays a row read.
This commit includes the corrective pass that followed the independent review
in planning/reports/M5-IMPLEMENTATION-REPORT.md. The invariant it exists to
hold is:
visible active transcript position == stored head == authoritative state
Narrator editing (D10, STORY-BRANCH-SEMANTICS §§14-15)
A narrator edit no longer rewrites a row. It returns to the state before the
turn, takes the reader's exact text as the accepted narration, re-derives the
state that text implies, and becomes a new active continuation — while the
original narration keeps its words, its live flag and its whole future as
retained history. At the tip the correction is another take; with story below
it, it forks. No new history machinery: this is the existing fork/take/head
path with the reader's text in place of a generated reply. The §14A refusal
is therefore gone for narrator turns, and remains only for player input.
Pre-M5 positions
Migration 88 backfills the empty narrative document onto every action written
before M5, and a missing snapshot now restores the empty document instead of
leaving the previous position's state standing. Restoring to an old Save
Point no longer leaves a later position's entities and facts on screen.
Narrator context
Replayed history carries prose only; the machine-readable block is no longer
reconstructed into past turns, where it contradicted the authoritative state
in the same prompt. A fact withdrawn by a manual correction is now named as
no longer true, with the reader's reason, rather than silently dropped.
Also
- state_changes joins the action-list bulk read, removing one query per row.
- Extraction takes only the application's own protocol payload: an ordinary
```json or ```python block in a story survives, and a mangled proposal
still does not reach the reader.
Planning: ADR 013 records the authoritative document shape; §§14-15/14A, D10,
C04 and BUILD-MILESTONES are updated to describe what exists. Debt is recorded
against M8 (scenario editor UX) and M9 (export of the audit trail).
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PWU4gTfLYY6Qq9U7aa9Qw2
400 lines
16 KiB
Python
400 lines
16 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 build_context, cursors
|
|
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
|
|
)
|
|
system_text, story_text, snapshot = build_context(
|
|
adventure, settings, memories, exclude_action_id=replacing_id
|
|
)
|
|
|
|
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,
|
|
)
|