Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,11 @@ GEMINI_MODEL=gemini-2.5-flash
# LITELLM_MODEL=gemini/gemini-2.5-flash
# ANTHROPIC_API_KEY=sk-ant-...

# ── Langfuse observability ────────────────────────────────────────────────────
# Get keys from http://localhost:3000 after running docker compose up
# ── Langfuse observability (v3) ───────────────────────────────────────────────
# Get keys from http://localhost:3010 after running docker compose up
LANGFUSE_PUBLIC_KEY=pk-lf-your-public-key
LANGFUSE_SECRET_KEY=sk-lf-your-secret-key
LANGFUSE_HOST=http://localhost:3000
LANGFUSE_HOST=http://localhost:3010

# ── Infrastructure ────────────────────────────────────────────────────────────
REDIS_URL=redis://localhost:6379/0
Expand Down
28 changes: 18 additions & 10 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,15 +17,18 @@
## Features

- **Multi-provider LLM support** via LiteLLM — swap between Gemini, DeepSeek, Claude, GPT-4o, Groq, Ollama, and 100+ more with a single env var change, zero code changes
- **Split planner model** — route the one-token routing decision to a cheap/fast model (`PLANNER_MODEL`) while keeping a capable model for synthesis; falls back to the same model if unset
- **LLM retry with jitter** — all LLM calls automatically retry up to 3× with exponential backoff and random jitter, handling transient rate-limit and network errors
- **FastAPI** production API with structured JSON logging and request-ID tracing
- **LangGraph** workflow: `planner → tool → synthesizer` state machine with persistent memory across turns
- **History truncation** — only the most recent 10 turns are sent to the LLM; full history is preserved in the Postgres checkpoint
- **DuckDuckGo search** — real web search with no API key required
- **Redis** session memory — conversation context survives API restarts
- **Postgres** run history — every agent run persisted for audit and replay
- **Langfuse** LLM tracing — full trace context, per-request handlers, generation-level spans
- **Prometheus + Grafana** — request rate, latency histograms (p50/p95/p99), agent run duration, error rate
- **Nuxt 3 frontend** — dark-mode chat UI + live status dashboard
- **Docker Compose** 8-container stack — API, frontend, Redis, Postgres, Langfuse, Prometheus, Grafana
- **Postgres** run history and LangGraph checkpoints — every agent run persisted for audit and replay; persistent multi-turn memory per session
- **Alembic** database migrations — schema changes are versioned and applied automatically at startup
- **Langfuse** LLM tracing — full trace context, node-level spans nested under a root trace per request
- **Prometheus + Grafana** — request rate, latency histograms (p50/p95/p99), agent run duration, tool execution counts
- **Nuxt 3 frontend** — dark-mode chat UI + live status dashboard with 120s request timeout and clear error messages
- **Docker Compose** 11-container stack — API, frontend, Redis, Postgres, Langfuse (web + worker + ClickHouse + MinIO), Prometheus, Grafana
- **uv** — fast, reproducible Python dependency management
- **Ruff + Pyright** — lint, format, and strict type checking
- **pytest** — unit, integration, and eval suites with automatic LLM mocking
Expand Down Expand Up @@ -243,8 +246,12 @@ agent-platform/
│ ├── tools/ # DuckDuckGo search tool
│ └── llm.py # Unified LLM factory (all providers)
├── infra/
│ └── docker/
│ └── Dockerfile # Multistage Python build
│ ├── docker/
│ │ └── Dockerfile # Multistage Python build
│ ├── migrations/ # Alembic migration scripts
│ │ ├── env.py # Migration env (excludes LangGraph tables)
│ │ └── versions/ # Versioned schema migrations
│ └── clickhouse/ # ClickHouse config (used by Langfuse)
├── monitoring/
│ ├── prometheus/ # prometheus.yml scrape config
│ └── grafana/ # Provisioned datasources + dashboards
Expand Down Expand Up @@ -291,6 +298,7 @@ just ci # full local CI gate
| `OLLAMA_BASE_URL` | `http://localhost:11434` | Ollama server URL |
| `OLLAMA_MODEL` | `llama3.2` | Ollama model name |
| `LITELLM_MODEL` | `gemini/gemini-2.5-flash` | Any LiteLLM model string |
| `PLANNER_MODEL` | *(same as LITELLM_MODEL)* | Cheaper model for the one-token routing step; falls back to `LITELLM_MODEL` if blank |

Any API key supported by LiteLLM (`ANTHROPIC_API_KEY`, `DEEPSEEK_API_KEY`, etc.) is passed through automatically.

Expand Down Expand Up @@ -344,8 +352,8 @@ Full LLM traces with per-request handlers. Open http://localhost:3010. Create a
| Orchestration | LangGraph, LangChain |
| LLM | LiteLLM (100+ providers) · Gemini 2.5 Flash default |
| Search | DuckDuckGo (no API key) |
| Memory | Redis + LangGraph MemorySaver |
| History | Postgres + SQLAlchemy async + asyncpg |
| Memory | LangGraph AsyncPostgresSaver (falls back to MemorySaver) |
| History | Postgres + SQLAlchemy async + asyncpg + Alembic |
| Tracing | Langfuse |
| Metrics | Prometheus, Grafana |
| Logging | structlog (JSON in production) |
Expand Down
149 changes: 149 additions & 0 deletions alembic.ini
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
# A generic, single database configuration.

[alembic]
# path to migration scripts.
# this is typically a path given in POSIX (e.g. forward slashes)
# format, relative to the token %(here)s which refers to the location of this
# ini file
script_location = %(here)s/infra/migrations

# template used to generate migration file names; The default value is %%(rev)s_%%(slug)s
# Uncomment the line below if you want the files to be prepended with date and time
# see https://alembic.sqlalchemy.org/en/latest/tutorial.html#editing-the-ini-file
# for all available tokens
# file_template = %%(year)d_%%(month).2d_%%(day).2d_%%(hour).2d%%(minute).2d-%%(rev)s_%%(slug)s
# Or organize into date-based subdirectories (requires recursive_version_locations = true)
# file_template = %%(year)d/%%(month).2d/%%(day).2d_%%(hour).2d%%(minute).2d_%%(second).2d_%%(rev)s_%%(slug)s

# sys.path path, will be prepended to sys.path if present.
# defaults to the current working directory. for multiple paths, the path separator
# is defined by "path_separator" below.
prepend_sys_path = .


# timezone to use when rendering the date within the migration file
# as well as the filename.
# If specified, requires the tzdata library which can be installed by adding
# `alembic[tz]` to the pip requirements.
# string value is passed to ZoneInfo()
# leave blank for localtime
# timezone =

# max length of characters to apply to the "slug" field
# truncate_slug_length = 40

# set to 'true' to run the environment during
# the 'revision' command, regardless of autogenerate
# revision_environment = false

# set to 'true' to allow .pyc and .pyo files without
# a source .py file to be detected as revisions in the
# versions/ directory
# sourceless = false

# version location specification; This defaults
# to <script_location>/versions. When using multiple version
# directories, initial revisions must be specified with --version-path.
# The path separator used here should be the separator specified by "path_separator"
# below.
# version_locations = %(here)s/bar:%(here)s/bat:%(here)s/alembic/versions

# path_separator; This indicates what character is used to split lists of file
# paths, including version_locations and prepend_sys_path within configparser
# files such as alembic.ini.
# The default rendered in new alembic.ini files is "os", which uses os.pathsep
# to provide os-dependent path splitting.
#
# Note that in order to support legacy alembic.ini files, this default does NOT
# take place if path_separator is not present in alembic.ini. If this
# option is omitted entirely, fallback logic is as follows:
#
# 1. Parsing of the version_locations option falls back to using the legacy
# "version_path_separator" key, which if absent then falls back to the legacy
# behavior of splitting on spaces and/or commas.
# 2. Parsing of the prepend_sys_path option falls back to the legacy
# behavior of splitting on spaces, commas, or colons.
#
# Valid values for path_separator are:
#
# path_separator = :
# path_separator = ;
# path_separator = space
# path_separator = newline
#
# Use os.pathsep. Default configuration used for new projects.
path_separator = os

# set to 'true' to search source files recursively
# in each "version_locations" directory
# new in Alembic version 1.10
# recursive_version_locations = false

# the output encoding used when revision files
# are written from script.py.mako
# output_encoding = utf-8

# database URL. This is consumed by the user-maintained env.py script only.
# other means of configuring database URLs may be customized within the env.py
# file.
sqlalchemy.url = driver://user:pass@localhost/dbname


[post_write_hooks]
# post_write_hooks defines scripts or Python functions that are run
# on newly generated revision scripts. See the documentation for further
# detail and examples

# format using "black" - use the console_scripts runner, against the "black" entrypoint
# hooks = black
# black.type = console_scripts
# black.entrypoint = black
# black.options = -l 79 REVISION_SCRIPT_FILENAME

# lint with attempts to fix using "ruff" - use the module runner, against the "ruff" module
# hooks = ruff
# ruff.type = module
# ruff.module = ruff
# ruff.options = check --fix REVISION_SCRIPT_FILENAME

# Alternatively, use the exec runner to execute a binary found on your PATH
# hooks = ruff
# ruff.type = exec
# ruff.executable = ruff
# ruff.options = check --fix REVISION_SCRIPT_FILENAME

# Logging configuration. This is also consumed by the user-maintained
# env.py script only.
[loggers]
keys = root,sqlalchemy,alembic

[handlers]
keys = console

[formatters]
keys = generic

[logger_root]
level = WARNING
handlers = console
qualname =

[logger_sqlalchemy]
level = WARNING
handlers =
qualname = sqlalchemy.engine

[logger_alembic]
level = INFO
handlers =
qualname = alembic

[handler_console]
class = StreamHandler
args = (sys.stderr,)
level = NOTSET
formatter = generic

[formatter_generic]
format = %(levelname)-5.5s [%(name)s] %(message)s
datefmt = %H:%M:%S
66 changes: 43 additions & 23 deletions apps/api/agent.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import contextlib
import time
import uuid
from typing import cast
from typing import Any, cast

import structlog
from fastapi import APIRouter, BackgroundTasks, Request
Expand Down Expand Up @@ -66,29 +67,48 @@ async def run_agent(

invoke_config: RunnableConfig = {"configurable": {"thread_id": thread_id}}

langfuse_handler = None
if settings.langfuse_public_key and settings.langfuse_secret_key:
try:
from langfuse.callback import CallbackHandler # type: ignore[import-untyped]

lf_kwargs: dict[str, object] = {
"public_key": settings.langfuse_public_key,
"secret_key": settings.langfuse_secret_key,
"host": settings.langfuse_host,
"user_id": body.user_id,
"session_id": thread_id,
"trace_name": _LANGFUSE_TRACE_NAME,
"tags": [settings.environment, settings.llm_provider],
}
langfuse_handler = CallbackHandler(**lf_kwargs) # type: ignore[arg-type]
invoke_config["callbacks"] = [langfuse_handler]
except Exception as exc:
log.warning("langfuse.callback.unavailable", error=str(exc))
# ── Run agent ────────────────────────────────────────────────────────────
# Use the single Langfuse client initialised at startup (app.state.langfuse).
# CallbackHandler is created INSIDE start_as_current_observation so that the
# active OTel span is already set when the handler builds its child spans —
# this is what makes node/LLM observations appear nested under the root trace.
lf_client: Any = getattr(request.app.state, "langfuse", None)

start = time.perf_counter()
result = await request.app.state.agent.ainvoke(cast(AgentState, state), config=invoke_config)
duration = time.perf_counter() - start
result: dict[str, Any] = {}

if lf_client is not None:
with lf_client.start_as_current_observation(
name=_LANGFUSE_TRACE_NAME,
as_type="span",
input={"user_message": body.input},
) as root_obs:
with contextlib.suppress(Exception):
lf_client.update_current_trace(
name=_LANGFUSE_TRACE_NAME,
user_id=body.user_id,
session_id=thread_id,
tags=[settings.environment, settings.llm_provider],
metadata={"thread_id": thread_id, "llm_provider": settings.llm_provider},
)
# Create CallbackHandler here, after the parent span is active,
# so LangGraph node spans are captured as children of root_obs.
with contextlib.suppress(Exception):
from langfuse.langchain import CallbackHandler # type: ignore[import-untyped]

invoke_config["callbacks"] = [CallbackHandler()]

result = await request.app.state.agent.ainvoke(
cast(AgentState, state), config=invoke_config
)
with contextlib.suppress(Exception):
root_obs.update(output=result.get("final_output"))
else:
result = await request.app.state.agent.ainvoke(
cast(AgentState, state), config=invoke_config
)

duration = time.perf_counter() - start
agent_run_duration_seconds.observe(duration)
duration_ms = round(duration * 1000, 2)

Expand All @@ -101,8 +121,8 @@ async def run_agent(
thread_id=thread_id,
)

if langfuse_handler is not None:
background_tasks.add_task(langfuse_handler.langfuse.flush)
if lf_client is not None:
background_tasks.add_task(lf_client.flush)

background_tasks.add_task(
_persist_run,
Expand Down
9 changes: 8 additions & 1 deletion apps/api/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@


class Settings(BaseSettings):
model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8")
model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore")

app_name: str = "Agent Platform"
environment: str = "development"
Expand Down Expand Up @@ -31,6 +31,13 @@ class Settings(BaseSettings):
# gemini/gemini-2.5-flash, anthropic/claude-opus-4-7, openai/gpt-4o
# groq/llama-3.1-8b-instant, bedrock/anthropic.claude-3-5-sonnet-20241022-v2:0
litellm_model: str = "gemini/gemini-2.5-flash"
# Cheap fast model used only for planner routing (outputs one token).
# Defaults to same as litellm_model so it works out-of-the-box without config.
planner_model: str = ""

@property
def resolved_planner_model(self) -> str:
return self.planner_model or self.litellm_model

# ── Langfuse observability ────────────────────────────────────────────────
langfuse_public_key: str | None = None
Expand Down
32 changes: 29 additions & 3 deletions apps/api/main.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import contextlib
from contextlib import asynccontextmanager

import structlog
Expand Down Expand Up @@ -26,11 +27,31 @@ async def lifespan(app: FastAPI):
log.info("startup", app=settings.app_name, env=settings.environment, llm=settings.llm_provider)

try:
from apps.api.db.session import create_tables
import asyncio

await create_tables()
from alembic import command as alembic_command
from alembic.config import Config as AlembicConfig

alembic_cfg = AlembicConfig("alembic.ini")
await asyncio.to_thread(alembic_command.upgrade, alembic_cfg, "head")
log.info("db.migrations.applied")
except Exception as exc:
log.warning("db.init.skipped", reason=str(exc))
log.warning("db.migrations.skipped", reason=str(exc))

# Initialize Langfuse once at startup so the OTel tracer provider is set up
# exactly once. Re-creating Langfuse() per request resets the OTel context
# and breaks parent-span propagation for child observations.
app.state.langfuse = None
if settings.langfuse_public_key and settings.langfuse_secret_key:
with contextlib.suppress(Exception):
from langfuse import Langfuse # type: ignore[import-untyped]

app.state.langfuse = Langfuse(
public_key=settings.langfuse_public_key,
secret_key=settings.langfuse_secret_key,
host=settings.langfuse_host,
)
log.info("tracing.backend", backend="langfuse", host=settings.langfuse_host)

# Build agent graph with Postgres checkpointer for persistent memory.
# Uses the same agentdb already running; falls back to MemorySaver if unavailable.
Expand All @@ -54,6 +75,11 @@ async def lifespan(app: FastAPI):

yield

lf = getattr(app.state, "langfuse", None)
if lf is not None:
with contextlib.suppress(Exception):
lf.flush()

ctx = getattr(app.state, "checkpointer_ctx", None)
if ctx is not None:
try: # noqa: SIM105
Expand Down
Loading
Loading