import asyncio import uuid from sqlalchemy import delete from sqlalchemy.ext.asyncio import AsyncSession from app.models.adventure_chunk import AdventureChunk from app.rag.chunking import chunk_text from app.rag.embeddings import embed_texts EMBED_BATCH_SIZE = 32 async def ingest_adventure_text(session: AsyncSession, game_id: uuid.UUID, text: str) -> int: """Chunks and embeds a player-supplied adventure text, scoped to one game. Safe to call again for the same game (e.g. future re-ingest) — previous chunks are cleared first.""" await session.execute(delete(AdventureChunk).where(AdventureChunk.game_id == game_id)) chunks = chunk_text(text) if not chunks: await session.commit() return 0 for batch_start in range(0, len(chunks), EMBED_BATCH_SIZE): batch = chunks[batch_start : batch_start + EMBED_BATCH_SIZE] # embed_texts is a synchronous, CPU-bound sentence-transformers call — run it off the # event loop so a large adventure doesn't stall other concurrent requests/WS connections. embeddings = await asyncio.to_thread(embed_texts, batch) for i, (content, embedding) in enumerate(zip(batch, embeddings)): session.add( AdventureChunk( game_id=game_id, chunk_index=batch_start + i, content=content, embedding=embedding, ) ) await session.commit() return len(chunks)