diff --git a/Makefile b/Makefile index edc584f..03bc623 100644 --- a/Makefile +++ b/Makefile @@ -15,8 +15,9 @@ test: ## Run unit tests only (fast, no IO) test-all: ## Run all tests (unit + integration) uv run pytest -v -lint: ## Run linter (ruff check) +lint: ## Run linter (ruff check + mypy) uv run ruff check src/ tests/ + uv run mypy src/ format: ## Format code (ruff format + fix) uv run ruff format src/ tests/ @@ -37,7 +38,7 @@ run-api: ## Start the FastAPI server uv run uvicorn researchos.infrastructure.api.main:app --reload --port 8000 run-bot: ## Start the Telegram bot - uv run python -m researchos.infrastructure.bot.main + uv run python scripts/run_telegram_bot.py # ── Docker ── docker-up: ## Start local stack (api + chroma) diff --git a/docs/essays/2026-W33-grafo-vs-pipeline.md b/docs/essays/2026-W33-grafo-vs-pipeline.md index e083c79..fef0bd4 100644 --- a/docs/essays/2026-W33-grafo-vs-pipeline.md +++ b/docs/essays/2026-W33-grafo-vs-pipeline.md @@ -2,10 +2,10 @@ Cuando construímos soluciones dentro de ML es bastante común que como científico de datos tendamos a crear procesos síncronos, lineas y continuos de manera que cada paso está completamente determinado y el proceso fluye desde el punto A hasta el punto B en una serie de sucesos predecibles y replicables. -Ahora bien, en la transición a AI engineer me he encontrado con la necesidad de pensar en sistemas que se retroalimenten, que tengan la capacidad de volver sobre sí mismos, "repensarse" y luego seleccionar otros caminos si es el caso, además de tener que manejar sistemas con muchas bifurcaciones y caminos posibles. Así pues, surge como herramienta langGraph y en particular la filosofía de grafos: Aplicado al desarrollo que tenemos, un grafo permite estructuralmente volver hacia atrás si es necesario y además bifurcarse, de manera que el sistema se retro-alimenta y se mueve en flujos que no podrían ser posibles con una cadena de funciones ya que como lo mencioné arriba, éstas son secuenciales. +Ahora bien, en la transición a AI engineer me he encontrado con la necesidad de pensar en sistemas que se retroalimenten, que tengan la capacidad de volver sobre sí mismos, "repensarse" y luego seleccionar otros caminos si es el caso, además de tener que manejar sistemas con muchas bifurcaciones y caminos posibles. Así pues, surge como herramienta langGraph y en particular la filosofía de grafos: Aplicado al desarrollo que tenemos, un grafo permite estructuralmente volver hacia atrás si es necesario y además bifurcarse, de manera que el sistema se retro-alimenta y se mueve en flujos que aunque podrían ciclarse y bifurcarse con python puro tienen un costo de desarrollo manual además del control y observación del flujo de la información, cosa que LangGraph ya hace por sí mismo. -Este hecho lleva a cuestionar por qué es necesario en este punto integrar LangGraph cuando mi pipeline actualmnete solo tiene 2 pasos "recuperar->generate", pues la realidad es que ahora mismo un grafo no aporta mucho valor ya que el sistema es secuencial, sin embargo, previendo que en próximas tareas quiero que haya retroalimentación y que además espero bifurcar decisiones entonces introducir un grafo en este momento me ayuda a que la construcción de la V2 y futuras versiones esté sobre esta filosofía. Sí podría continuar como lo tengo ahora pero más adelante tendría que hacer mayor refactorización para alcanzar las funcionalidades que espero que el sistema tenga. +Este hecho lleva a cuestionar por qué es necesario en este punto integrar LangGraph cuando mi pipeline actualmnete solo tiene 2 pasos "recuperar->generate", pues la realidad es que ahora mismo un grafo no aporta mucho valor ya que el sistema es secuencial, sin embargo, previendo que en próximas tareas quiero que haya retroalimentación y que además espero bifurcar decisiones entonces introducir un grafo en este momento me ayuda a que la construcción de la V2 y futuras versiones esté sobre esta filosofía. Sí podría continuar como lo tengo ahora pero más adelante tendría que hacer mayor refactorización para alcanzar las funcionalidades que espero que el sistema tenga. Aún más, al conocer el flujo de la información de antemano, es más sencillo ir entendiendo la estructura y modo funcional de la librería sobre este grafo de 2 nodos, centrándome en el cómo orquestar con la librería más allá de la lógica en sí misma. Una ventaja adicional es que LangGraph al ser una librería de bajo nivel tiene implementaciones que permiten conocer el estado explícito del sistema, si mantuviera el pipeline entonces sí podría inspeccionar el estado pero tendría que implementar manualmente estrategias para conocer esos estados, lo cual en temas de observabilidad da una ventaja importante con langGraph. -Dicho esto, no es que no pueda continuar con un pipeline secuencia, sino que hacerlo ahora mismo es una decisión estratégica de cara a la evolución de funcionalidades que espero alcanzar con mi sistema, por lo que esta sobre-ingeniería tendrá sus beneficios en el futuro. +Dicho esto, no es que no pueda continuar con un pipeline secuencia, sino que hacerlo ahora mismo es una decisión estratégica de cara a la evolución de funcionalidades que espero alcanzar con mi sistema, por lo que esta inversión en infraestructura tendrá sus beneficios en el futuro. diff --git a/docs/learnings.md b/docs/learnings.md index 4c01c09..2ac6b7f 100644 --- a/docs/learnings.md +++ b/docs/learnings.md @@ -408,3 +408,97 @@ Regla simple para ResearchOS: cambiaron en la migración 0.x → 1.x) pero no verifiqué la API real todavía. --- + +**Fecha:** 18/08/2026 + +### ¿Qué aprendí? + +- **Un wrapper es una función de orden superior: recibe una función y + devuelve otra función.** Lo practiqué primero en + `taller_retorno_researchos.ipynb` con `with_exclamation(fn)`, que define + `wrapped(text)` cerrando sobre `fn` y devuelve `wrapped` sin ejecutarlo. + Es el mismo molde que `with_logging(answer_fn: AnswerFn) -> AnswerFn` en + `scripts/run_telegram_bot.py`: recibe un `AnswerFn`, devuelve otro. +- **Lo que importa no es que envuelva, sino cuándo corre cada parte.** En el + segundo experimento del notebook, `print("CONSTRUYENDO")` corre una sola + vez -- cuando se llama `with_exclamation(greet)` -- y `print("EJECUTANDO")` + corre una vez por cada llamada a la función envuelta. Esa separación + construcción-una-vez / ejecución-por-llamada es exactamente por qué + `with_logging` puede abrir el `logging.FileHandler` y hacer + `query_logger.addHandler(...)` en el cuerpo de la función envolvente: ese + código corre una sola vez, al armar `with_logging(answer)` en el + composition root, no en cada pregunta que llega por Telegram. +- **El wrapper no necesita saber nada de Telegram.** `with_logging` solo + conoce la firma `AnswerFn` (`Callable[[str], Awaitable[str]]`) — le es + indiferente si la función que envuelve viene de `answer_query` con + vector-only o con hybrid+rerank. Mismo patrón de composition root que ya + veníamos usando (`retrieve_hybrid_rerank`), aplicado ahora para agregar un + efecto secundario (logging) en vez de una estrategia de recuperación. + +**LangGraph: estado, nodo y edge, con evidencia de mis propios experimentos** + +- **Estado.** Es la estructura de datos explícita (`TypedDict`, `dataclass` o + modelo Pydantic) que se pasa por todo el grafo. Cada nodo recibe el + estado (o un subconjunto, si se declaran `InputState`/`OutputState`/ + `PrivateState` separados) y devuelve un **update parcial** — no el estado + completo. Probé esto con `InputState`/`OverallState`/`PrivateState`/ + `OutputState` en un grafo de 3 nodos donde cada uno lee de un canal y + escribe en otro, y funcionó exactamente así: `node_1` escribe en + `OverallState`, `node_2` lee de ahí y escribe en `PrivateState`, `node_3` + lee de `PrivateState` y arma el `OutputState` final. +- **Nodo.** Una función que recibe estado y devuelve una actualización de + estado. Nada más — no decide a dónde ir después (eso es trabajo del edge), + salvo que sea un nodo tipo `Command` (que mi tutor y yo ya distinguimos + esta semana para T19: si el nodo solo calcula, va con conditional edge). +- **Edge.** Conexión entre nodos. Puede ser fija (`add_edge`, siempre va de + A a B) o condicional (`add_conditional_edges`, una función del estado + decide el destino, como `should_continue` devolviendo `"tool_node"` o + `END` según si el último mensaje tiene `tool_calls`). +- **El ciclo es lo que hace posible el multi-tool-call, y lo comprobé + rompiéndolo.** Con el edge `tool_node → llm_call` presente, le pedí al + agente derivar una expresión, multiplicar y dividir el resultado — hizo + las tres cosas en secuencia porque después de cada tool call volvía a + `llm_call` a decidir el siguiente paso. Al comentar ese edge, el agente + ejecutó `multiply` una sola vez y terminó ahí — sin el edge de retorno, + `tool_node` no tiene a dónde ir y el grafo termina implícitamente. Es la + demostración empírica de por qué T22 (reescritura de query + reintento) + necesita un ciclo real, no solo un conditional edge de ida. +- **El orden en que declaro los edges no afecta el grafo compilado.** + Definí primero el `add_edge("tool_node", "llm_call")` y después el + conditional edge, y también al revés — mismo comportamiento en ambos + casos. El grafo se arma con la suma de todas las llamadas a + `add_edge`/`add_conditional_edges` antes de `compile()`, no importa en + qué secuencia se llamaron. +- **`TypedDict` no valida en runtime; Pydantic sí, en cada actualización de + estado.** Me había quedado la duda de qué significa "validación + recursiva" con Pydantic como estado — sí es lo que sospechaba: cada vez + que un nodo devuelve un update, si el estado es un modelo Pydantic, se + valida contra los tipos declarados en ese momento, no solo al construir + el estado inicial. Con `TypedDict` esa verificación no existe en + ejecución — es solo información para el type checker. +- **Un reducer decide cómo se combina el valor viejo con el nuevo, no lo + reemplaza por default.** Sin anotación, una clave se sobrescribe. Con + `Annotated[list[str], operator.add]` (o un reducer propio como + `append_strings(left, right)`), el update se acumula en vez de pisar el + valor anterior — así es como `messages` en `MessagesState` va creciendo + turno a turno en vez de perder el historial. +- **Ya tengo un primer borrador del estado para T18**, como `dataclass`: + `query: str`, `documents: list[Document]`, `answer: str`, + `messages: Annotated[list, add_messages]`, `rewritten_query: str` — con + `rewritten_query` ya pensando en T22 antes de empezar T18. + +### Errores interesantes + +- En el notebook, re-ejecutar la celda de `greet = with_exclamation(greet)` + varias veces sin reiniciar el kernel apiló wrappers uno sobre otro (el + output mostró tres `"EJECUTANDO"` y `"!!!"` en vez de uno) — cada + ejecución envolvía el `greet` ya envuelto de la ejecución anterior, no el + original. No es un bug de la función, es un recordatorio de que el + estado de un notebook persiste entre celdas y `x = f(x)` no es idempotente + si se re-corre la celda. + +### ¿Qué no entendí bien? + +- El mecanismo de flujo de la información ya que el patrón me muestra que la función que envuelve recibe los mismo argumentos de la función que quiero envolver, pero aún así, no asimilo muy bien cómo fluye la información ya que estoy acostrumbrado a un patrón más lineal (spaguetti) + +--- diff --git a/docs/work_log.md b/docs/work_log.md index f96ae1e..9da4a99 100644 --- a/docs/work_log.md +++ b/docs/work_log.md @@ -287,3 +287,44 @@ de preguntas en el ritual normal --- + +## 2026-08-18 + +### Trabajo desarrollado +- `scripts/run_telegram_bot.py`: agregado `with_logging(answer_fn) -> AnswerFn`, + un wrapper que registra cada query entrante en `data/raw/queries.jsonl` + (timestamp UTC + texto) — insumo real para T24 (eval V1 vs V2 sin data + leakage). Corregida de paso la duplicación de `AnswerFn`: ahora se importa + de `telegram_bot.py` en vez de redefinirse +- `mypy` conectado a `make lint` — ya estaba como dependencia dev y con + config básica desde antes, pero nunca se ejecutaba. Corrida completa: 38 + errores encontrados; arreglados los de configuración/ruido (`rank_bm25` y + `fitz` sin stubs de tipos) y dos `var-annotated` en `retrieval_service.py`; + quedan 34 errores reales (gaps de manejo de `None` en 6 archivos) sin + tocar, a la espera de decidir alcance +- `make run-bot` corregido: apuntaba a un módulo inexistente + (`researchos.infrastructure.bot.main`); ahora corre el script real + (`scripts/run_telegram_bot.py`) +- Verificado manualmente que el bot arranca sin errores con el wrapper de + logging activo: `data/raw/queries.jsonl` se crea al construir + `with_logging`, embedder/Chroma/BM25 se construyen sin fallas. Falta + confirmar con un mensaje real desde Telegram +- `notebooks/201-jmmz-langraph-study.ipynb`: práctica de `StateGraph` — + nodos, edges fijos y condicionales, ciclos (verificado quitando y + reordenando edges), reducers, y esquemas de estado separados + (`InputState`/`OutputState`/`PrivateState`) +- `docs/learnings.md`: entrada de hoy documenta el patrón wrapper/decorador + (con ejemplo propio del taller de retorno) y los conceptos de estado, + nodo y edge de LangGraph con evidencia de los experimentos del notebook 201 + +### Próximos pasos +- Decidir qué hacer con los 34 errores de mypy restantes (`arxiv.py`, + `telegram_bot.py`, `anthropic_llm.py`, `chroma.py`, `embedder.py`, + `ingestion_service.py`) — arreglar ahora, un subconjunto, o registrar + como deuda en `ROADMAP.md` +- Confirmar el logging de queries con un mensaje real por Telegram +- Arrancar T18 (LangGraph fundamentals) con el borrador de estado ya + escrito en el notebook 201 (`query`, `documents`, `answer`, `messages`, + `rewritten_query`) + +--- diff --git a/notebooks/201-jmmz-langraph-study.ipynb b/notebooks/201-jmmz-langraph-study.ipynb index 8a50702..835110c 100644 --- a/notebooks/201-jmmz-langraph-study.ipynb +++ b/notebooks/201-jmmz-langraph-study.ipynb @@ -703,6 +703,263 @@ "source": [ "Conclusión: en la lógica definida no hubo cambio si definía primero el edge y luego el conditional_edge o visceversa" ] + }, + { + "cell_type": "markdown", + "id": "38dee84d", + "metadata": {}, + "source": [ + "# Cómo de declara el estado " + ] + }, + { + "cell_type": "markdown", + "id": "0360ca43", + "metadata": {}, + "source": [ + "funciones reducer: Especifican cómo aplicar actualizaciones al \n", + "\n", + "El estado puede ser un TypeDict o un modelo Pydantic\n", + "\n", + "TODOS los nodos emitirán actualizaciones del estado\n", + "\n", + "- TypeDict es la manera princiapl de definicir el estado. \n", + "- Si se requieren valores por default entonces usar dataclass\n", + "- Pydantic solo si se requiere validación de datos recursiva (qué significa esto? tal vez que cada vez que el estado se actualiza entonces pydantic verifica los tipos de datos)\n", + "\n", + "Se pueden agregar estados privados a nodos específicos y definir esquemas específicos para engrada y salidas del grafo e incluso un estado interno " + ] + }, + { + "cell_type": "code", + "execution_count": 1, + "id": "6e2c5cbf", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "{'graph_output': 'My name is Lance'}" + ] + }, + "execution_count": 1, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "from typing import TypedDict\n", + "\n", + "from langgraph.graph import END, START, StateGraph\n", + "\n", + "\n", + "class InputState(TypedDict):\n", + " user_input: str\n", + "\n", + "\n", + "class OutputState(TypedDict):\n", + " graph_output: str\n", + "\n", + "\n", + "class OverallState(TypedDict):\n", + " foo: str\n", + " user_input: str\n", + " graph_output: str\n", + "\n", + "\n", + "class PrivateState(TypedDict):\n", + " bar: str\n", + "\n", + "\n", + "def node_1(state: InputState) -> OverallState:\n", + " # Write to OverallState\n", + " return {\"foo\": state[\"user_input\"] + \" name\"}\n", + "\n", + "\n", + "def node_2(state: OverallState) -> PrivateState:\n", + " # Read from OverallState, write to PrivateState\n", + " return {\"bar\": state[\"foo\"] + \" is\"}\n", + "\n", + "\n", + "def node_3(state: PrivateState) -> OutputState:\n", + " # Read from PrivateState, write to OutputState\n", + " return {\"graph_output\": state[\"bar\"] + \" Lance\"}\n", + "\n", + "\n", + "builder = StateGraph(OverallState, input_schema=InputState, output_schema=OutputState)\n", + "builder.add_node(\"node_1\", node_1)\n", + "builder.add_node(\"node_2\", node_2)\n", + "builder.add_node(\"node_3\", node_3)\n", + "builder.add_edge(START, \"node_1\")\n", + "builder.add_edge(\"node_1\", \"node_2\")\n", + "builder.add_edge(\"node_2\", \"node_3\")\n", + "builder.add_edge(\"node_3\", END)\n", + "\n", + "graph = builder.compile()\n", + "graph.invoke({\"user_input\": \"My\"})\n", + "# {'graph_output': 'My name is Lance'}" + ] + }, + { + "cell_type": "markdown", + "id": "c5a33504", + "metadata": {}, + "source": [ + "# Qué devuelve un nodo" + ] + }, + { + "cell_type": "markdown", + "id": "b6293a06", + "metadata": {}, + "source": [ + " un nodo puede escribir en cualquier canal de estados en el estado del grafo. El estado del grafo es la unión de los canales de estados definidos en la inicialización, que incluye OverallState y los filtros InputState y OutputState.\n", + "\n", + " Un nodo devuelve un estado actualizado. Cualquier canal de estado en el grafo puede ser escrito desde un nodo" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "cf3ec66c", + "metadata": {}, + "outputs": [], + "source": [ + "StateGraph()" + ] + }, + { + "cell_type": "markdown", + "id": "b44f1d72", + "metadata": {}, + "source": [ + "# Qué hace un reducer y cuándo se necesita" + ] + }, + { + "cell_type": "markdown", + "id": "8b0c9900", + "metadata": {}, + "source": [ + "Una función reductora establece cómo se actualiza el estado en un nodo. Cada clave en el estado tiene su propia función reductora. Se necestia cuando quiero customizar el comportamiento de actualización del estado.\n", + "\n", + "SI NO SE HACE ESPECÍFICA se asume que todas las actualizaciones de la clave deben sobreescribirla.\n", + "\n", + "Los reductores ason funciones con 2 funciones posiciones left y right, siendo left el estado anterior y right el resultado de la modificación del nodo. Así pues, por default se reemplaza, o uno podría hacer reductores personalizados que modifiquen el estado en una manera particular, por ejemplo sumándolo: \n", + "\n", + "```python\n", + "from typing import Annotated\n", + "\n", + "from typing_extensions import TypedDict\n", + "\n", + "\n", + "def append_strings(left: list[str], right: list[str]) -> list[str]:\n", + " \"\"\"Combine the existing state value (left) with a node update (right).\"\"\"\n", + " return left + right\n", + "\n", + "\n", + "class State(TypedDict):\n", + " tags: Annotated[list[str], append_strings]\n", + "\n", + "```\n", + "\n", + "Suppose the state is {\"tags\": [\"draft\"]} and a node returns {\"tags\": [\"review\"]}. LangGraph calls:\n", + "\n", + "\n", + "```python\n", + "append_strings(left=[\"draft\"], right=[\"review\"]) # returns [\"draft\", \"review\"]\n", + "```\n", + "\n", + "Para espeficiar un reductor personalizado se debe declarar en el estado:\n", + "\n", + "```python\n", + "from operator import add\n", + "from typing import Annotated\n", + "\n", + "from typing_extensions import TypedDict\n", + "\n", + "\n", + "class State(TypedDict):\n", + " foo: int\n", + " bar: Annotated[list[str], add]\n", + "```" + ] + }, + { + "cell_type": "markdown", + "id": "70755cac", + "metadata": {}, + "source": [ + "# Qué campos voy a declarar en mi estado" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "49337ccf", + "metadata": {}, + "outputs": [], + "source": [ + "from researchos.domain import Document\n", + "from dataclasses import dataclass, field\n", + "from langgraph.graph import add_messages\n", + "\n", + "@dataclass\n", + "class State:\n", + " query: str\n", + " documents: list[Document] = field(default_factory=list)\n", + " answer: str = \"\"\n", + " messages: Annotated[list, add_messages] = field(default_factory=list)\n", + " rewritten_query: str = \"\"" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9ca4ca80", + "metadata": {}, + "outputs": [], + "source": [] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "0dc6f992", + "metadata": {}, + "outputs": [], + "source": [] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "f0439666", + "metadata": {}, + "outputs": [], + "source": [] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "4acbf941", + "metadata": {}, + "outputs": [], + "source": [] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9f4d269f", + "metadata": {}, + "outputs": [], + "source": [] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "2acd339e", + "metadata": {}, + "outputs": [], + "source": [] } ], "metadata": { diff --git a/notebooks/taller_retorno_researchos.ipynb b/notebooks/taller_retorno_researchos.ipynb index cca92f0..512e8e4 100644 --- a/notebooks/taller_retorno_researchos.ipynb +++ b/notebooks/taller_retorno_researchos.ipynb @@ -1033,11 +1033,105 @@ "\n", "Con eso ajustamos el arranque de V2 y decidimos si conviene afianzar algún tema antes de tocar código nuevo.\n" ] + }, + { + "cell_type": "code", + "execution_count": 4, + "id": "2014f700", + "metadata": {}, + "outputs": [], + "source": [ + "def with_exclamation(fn):\n", + " def wrapped(text):\n", + " return fn(text) + \"!fjdsañf\"\n", + " return wrapped\n", + "\n", + "def greet(name):\n", + " return f\"Hola {name}\"\n", + "\n", + "new_greet = with_exclamation(greet)\n" + ] + }, + { + "cell_type": "code", + "execution_count": 5, + "id": "c58ec0b8", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + ".wrapped(text)>" + ] + }, + "execution_count": 5, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "new_greet" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "db3c6398", + "metadata": {}, + "outputs": [], + "source": [ + "print(new_greet(\"Mario\")) # Hola Mario!" + ] + }, + { + "cell_type": "code", + "execution_count": 8, + "id": "7e07b9ae", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "CONSTRUYENDO\n", + "EJECUTANDO\n", + "EJECUTANDO\n", + "EJECUTANDO\n", + "Hola Mario!!!\n", + "EJECUTANDO\n", + "EJECUTANDO\n", + "EJECUTANDO\n", + "Hola Ana!!!\n" + ] + } + ], + "source": [ + "def with_exclamation(fn):\n", + " print(\"CONSTRUYENDO\") # ← momento 1\n", + "\n", + " def wrapped(text):\n", + " print(\"EJECUTANDO\") # ← momento 2\n", + " return fn(text) + \"!\"\n", + "\n", + " return wrapped\n", + "\n", + "greet = with_exclamation(greet) # imprime \"CONSTRUYENDO\"\n", + "print(greet(\"Mario\")) # imprime \"EJECUTANDO\"\n", + "print(greet(\"Ana\") ) # imprime \"EJECUTANDO\"" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "73f7d4ff", + "metadata": {}, + "outputs": [], + "source": [] } ], "metadata": { "kernelspec": { - "display_name": "researchos (3.11.8.final.0)", + "display_name": "researchos (3.11.8)", "language": "python", "name": "python3" }, diff --git a/pyproject.toml b/pyproject.toml index 7af8a57..7649c99 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -77,3 +77,7 @@ python_version = "3.11" warn_return_any = true warn_unused_configs = true disallow_untyped_defs = false + +[[tool.mypy.overrides]] +module = ["rank_bm25", "fitz"] +ignore_missing_imports = true diff --git a/scripts/run_telegram_bot.py b/scripts/run_telegram_bot.py index 59a4c05..e72d7b9 100644 --- a/scripts/run_telegram_bot.py +++ b/scripts/run_telegram_bot.py @@ -4,24 +4,29 @@ __import__("pysqlite3") sys.modules["sqlite3"] = sys.modules.pop("pysqlite3") +import json +import logging +from datetime import UTC, datetime + import chromadb from researchos.application.services.rag_service import answer_query from researchos.application.services.retrieval_service import hybrid_rerank_search, hybrid_search from researchos.config import settings from researchos.domain.models import Document -from researchos.infrastructure.bot.telegram_bot import TelegramBot +from researchos.infrastructure.bot.telegram_bot import AnswerFn, TelegramBot from researchos.infrastructure.llm.anthropic_llm import AnthropicLLM from researchos.infrastructure.retrieval.bm25 import BM25Retriever from researchos.infrastructure.retrieval.chroma import ChromaVectorStore from researchos.infrastructure.retrieval.embedder import LocalEmbedder -from researchos.paths import CHROMA_DIR +from researchos.paths import CHROMA_DIR, DATA_DIR + +query_logger = logging.getLogger("researchos.queries") # ── Build retrievers ── COLLECTION_NAME = "papers" K = 5 - embedder = LocalEmbedder() chroma = ChromaVectorStore(embedder=embedder, collection_name=COLLECTION_NAME) llm = AnthropicLLM() @@ -47,5 +52,34 @@ async def answer(query: str) -> str: return await answer_query(query, llm, retrieve=retrieve_hybrid_rerank) -bot = TelegramBot(token=settings.telegram_bot_token, answer_fn=answer) +def with_logging(answer_fn: AnswerFn) -> AnswerFn: + """Wrap an AnswerFn so every incoming query is appended to a JSONL file. + + Setup runs once at construction; the inner function runs per query. + Used to collect real user queries for the V2 evaluation dataset (T24), + avoiding the data leakage of writing eval questions against a known corpus. + """ + log_path = DATA_DIR / "raw" / "queries.jsonl" + log_path.parent.mkdir(parents=True, exist_ok=True) + + file_handler = logging.FileHandler(log_path, encoding="utf-8") + file_handler.setFormatter(logging.Formatter("%(message)s")) + + query_logger.addHandler(file_handler) + query_logger.setLevel(logging.INFO) + query_logger.propagate = False + + async def logged_answer(query: str) -> str: + query_logger.info( + json.dumps( + {"ts": datetime.now(UTC).isoformat(), "query": query}, + ensure_ascii=False, + ) + ) + return await answer_fn(query) + + return logged_answer + + +bot = TelegramBot(token=settings.telegram_bot_token, answer_fn=with_logging(answer)) bot.run() diff --git a/src/researchos/application/services/retrieval_service.py b/src/researchos/application/services/retrieval_service.py index 54d1354..dabb476 100644 --- a/src/researchos/application/services/retrieval_service.py +++ b/src/researchos/application/services/retrieval_service.py @@ -51,8 +51,8 @@ async def hybrid_search( results = await asyncio.gather(*[retriever.search(query, n) for retriever in retrievers]) - rrf_scores = {} - docs_by_id = {} + rrf_scores: dict[str, float] = {} + docs_by_id: dict[str, Document] = {} for ranking in results: for rank, doc in enumerate(ranking, start=1): # 1-indexed