import json import logging import uuid from collections.abc import Awaitable, Callable from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.llm.client import get_dm_system_prompt, get_llm_client from app.llm.context import build_context from app.llm.json_utils import fix_double_escaped_unicode from app.llm.tools import character_sheet, dice, end_game, link_existing_character, monster, world_state from app.models.character import Character from app.models.game import Game, GameParticipant from app.models.message import Message from app.models.user import User from app.rag.adventure_retrieval import build_adventure_rag_block from app.rag.retrieval import build_rag_block logger = logging.getLogger("app.llm.orchestrator") MAX_TOOL_ROUNDS = 5 MAX_TOKENS = 4096 TOOLS = [ {"type": "function", "function": dice.TOOL_SCHEMA}, {"type": "function", "function": character_sheet.TOOL_SCHEMA}, {"type": "function", "function": end_game.TOOL_SCHEMA}, {"type": "function", "function": monster.TOOL_SCHEMA}, {"type": "function", "function": world_state.TOOL_SCHEMA}, {"type": "function", "function": link_existing_character.TOOL_SCHEMA}, ] async def _execute_tool_call( session: AsyncSession, game_id: uuid.UUID, tool_name: str, tool_input: dict ) -> dict: """Runs one tool call. Errors are returned as a payload (not raised), so the DM sees them as a tool result and can recover instead of the whole turn crashing.""" try: if tool_name == "roll_dice": return dice.roll(tool_input["notation"]) if tool_name == "upsert_character_sheet": return await character_sheet.upsert(session, game_id, tool_input) if tool_name == "end_game": return await end_game.end(session, game_id, tool_input) if tool_name == "update_monster_hp": return await monster.update(session, game_id, tool_input) if tool_name == "update_world_state": return await world_state.update(session, game_id, tool_input) if tool_name == "link_existing_character": return await link_existing_character.link(session, game_id, tool_input) return {"error": f"Unknown tool {tool_name!r}"} except Exception as exc: # noqa: BLE001 logger.warning("Tool call %s failed: %s", tool_name, exc) await session.rollback() return {"error": str(exc)} async def _request_player_roll( session: AsyncSession, participant: GameParticipant, notation: str, reason: str ) -> dict: """Handles a roll_dice call for a player who rolls their own physical dice: validates the notation without touching the RNG, records it as the participant's pending roll (picked up and cleared at the top of their next turn), and returns a tool result that tells the DM to ask for it and wait instead of rolling itself.""" try: dice.parse_notation(notation) except ValueError as exc: return {"error": str(exc)} participant.pending_roll = {"notation": notation, "reason": reason} await session.commit() return { "awaiting_player_roll": True, "notation": notation, "reason": reason, "note": ( f"Der Spieler würfelt selbst mit physischen Würfeln. Nenne ihm klar und knapp, was er " f"würfeln soll ({notation}" + (f", {reason}" if reason else "") + "), und warte auf " "seine Antwort mit dem Ergebnis in einer künftigen Nachricht — würfle nicht selbst und " "erfinde kein Ergebnis." ), } def _format_character_sheet(character: Character) -> str: parts = [ f"- {character.name} ({character.race or '?'} {character.char_class or '?'}, Stufe {character.level})" ] if character.stats: parts.append(f" Werte: {character.stats}") hp = None if character.current_hp is not None or character.max_hp is not None: hp = f"{character.current_hp if character.current_hp is not None else '?'}/{character.max_hp if character.max_hp is not None else '?'}" if hp: parts.append(f" TP: {hp}") if character.combat_stats: parts.append(f" Kampfwerte: {character.combat_stats}") if character.abilities: ability_names = ", ".join(a.get("name", "") for a in character.abilities if isinstance(a, dict)) parts.append(f" Fähigkeiten: {ability_names}") if character.equipment: parts.append(f" Ausrüstung: {', '.join(character.equipment)}") if character.description: parts.append(f" Hintergrund: {character.description}") return "\n".join(parts) async def _check_hp_game_over(session: AsyncSession, game_id: uuid.UUID, result: dict) -> str | None: """After an upsert_character_sheet call, auto-ends the game if the character it touched dropped to 0 HP or below. Returns the end reason if the game was just ended, else None.""" current_hp = result.get("current_hp") if current_hp is None or current_hp > 0: return None game = await session.get(Game, game_id) if game is None or game.status == "ended": return None reason = f"{result.get('name', 'Ein Charakter')} ist bei {current_hp} Trefferpunkten zusammengebrochen." game.status = "ended" game.ended_reason = reason await session.commit() return reason async def run_dm_turn( session: AsyncSession, game_id: uuid.UUID, latest_player_message: str | None = None, user_id: uuid.UUID | None = None, on_roll: Callable[[str], Awaitable[None]] | None = None, on_awaiting_roll: Callable[[str, str], Awaitable[None]] | None = None, on_game_ended: Callable[[str], Awaitable[None]] | None = None, ) -> Message: client = get_llm_client() system_prompt = get_dm_system_prompt() participant: GameParticipant | None = None if user_id is not None: participant = ( await session.execute( select(GameParticipant).where( GameParticipant.game_id == game_id, GameParticipant.user_id == user_id ) ) ).scalar_one_or_none() if participant is not None and participant.pending_roll: pending = participant.pending_roll system_prompt += ( f"\n\nOffene Würfelanfrage: Du hast {pending.get('notation')} " f"({pending.get('reason') or 'ohne Angabe'}) angefragt, der Spieler würfelt selbst. " "Seine gerade eingegangene Nachricht enthält vermutlich das Ergebnis dieses Wurfs — " "nimm die genannte Zahl direkt als Wurfergebnis, würfle NICHT selbst und frag nicht " "erneut danach, außer die Nachricht beantwortet die Anfrage erkennbar nicht." ) # Consumed either way — if the message didn't actually answer it, the DM asking again # via roll_dice creates a fresh pending_roll for next time. participant.pending_roll = None await session.commit() game = await session.get(Game, game_id) if game is not None: system_prompt = ( f"{system_prompt}\n\n" "Spiel-Rahmendaten (vom Ersteller beim Anlegen bereits festgelegt — NICHT erneut " "abfragen, in Session 0 höchstens knapp bestätigen):\n" f"Name: {game.name}\nBeschreibung: {game.description}" ) participant_names = ( await session.execute( select(User.name) .join(GameParticipant, GameParticipant.user_id == User.id) .where(GameParticipant.game_id == game_id) ) ).scalars().all() if participant_names: system_prompt += ( "\n\nBereits beigetretene Spieler (über Login/Einladungslink — NICHT nach Anzahl " "oder Namen fragen): " + ", ".join(participant_names) + ". Weitere können jederzeit " "über den Einladungslink dazukommen; du erkennst sie automatisch am Namen, sobald sie schreiben." ) linked_characters = ( await session.execute( select(Character) .join(GameParticipant, GameParticipant.character_id == Character.id) .where(GameParticipant.game_id == game_id) ) ).scalars().all() if linked_characters: sheets = "\n".join(_format_character_sheet(c) for c in linked_characters) system_prompt += ( "\n\nBekannte Charakterblätter (bereits vollständig vorhanden — NICHT erneut nach " "Werten/Ausrüstung/Hintergrund fragen, nur bei tatsächlichen Änderungen über " f"upsert_character_sheet aktualisieren):\n{sheets}" ) await world_state.refresh_if_due(session, game_id, client, settings.dm_model) world_summary = await world_state.get_summary(session, game_id) if world_summary: system_prompt = ( f"{system_prompt}\n\nBisheriger Weltzustand (automatisch/von dir über update_world_state " f"aktuell gehalten — der Chatverlauf unten zeigt nur ein aktuelles Fenster, das hier " f"trägt alles Ältere):\n{world_summary}" ) if latest_player_message: rag_block = await build_rag_block(session, latest_player_message) if rag_block: system_prompt = f"{system_prompt}\n\n{rag_block}" adventure_block = await build_adventure_rag_block(session, game_id, latest_player_message) if adventure_block: system_prompt = f"{system_prompt}\n\n{adventure_block}" messages = [{"role": "system", "content": system_prompt}] messages.extend(await build_context(session, game_id)) final_text = "" for _ in range(MAX_TOOL_ROUNDS): response = await client.chat.completions.create( model=settings.dm_model, max_tokens=MAX_TOKENS, messages=messages, tools=TOOLS, ) choice = response.choices[0] message = choice.message final_text = message.content or "" if not message.tool_calls: break messages.append(message.model_dump(exclude_unset=True)) for tool_call in message.tool_calls: tool_input = fix_double_escaped_unicode(json.loads(tool_call.function.arguments)) if tool_call.function.name == "roll_dice" and participant is not None and participant.self_rolls: notation = tool_input.get("notation", "") reason = tool_input.get("reason", "") result = await _request_player_roll(session, participant, notation, reason) if "error" not in result and on_awaiting_roll is not None: await on_awaiting_roll(notation, reason) else: if tool_call.function.name == "roll_dice" and on_roll is not None: await on_roll(tool_input.get("notation", "")) result = await _execute_tool_call(session, game_id, tool_call.function.name, tool_input) end_reason = None if tool_call.function.name == "end_game" and "reason" in result: end_reason = result["reason"] elif tool_call.function.name == "upsert_character_sheet": end_reason = await _check_hp_game_over(session, game_id, result) if end_reason: result = {**result, "game_ended": True, "end_reason": end_reason} if end_reason and on_game_ended is not None: await on_game_ended(end_reason) messages.append( { "role": "tool", "tool_call_id": tool_call.id, "content": json.dumps(result, ensure_ascii=False), } ) else: if not final_text: final_text = "(Der Dungeon Master braucht einen Moment länger als erwartet — bitte versuche es erneut.)" dm_message = Message(game_id=game_id, sender_type="dm", content=final_text) session.add(dm_message) await session.commit() await session.refresh(dm_message) return dm_message