diff --git a/backend/open_webui/main.py b/backend/open_webui/main.py
index 747859db88..2bbede4141 100644
--- a/backend/open_webui/main.py
+++ b/backend/open_webui/main.py
@@ -68,6 +68,7 @@ from open_webui.socket.main import (
get_models_in_use,
)
from open_webui.routers import (
+ analytics,
audio,
images,
ollama,
@@ -1459,6 +1460,9 @@ app.include_router(functions.router, prefix="/api/v1/functions", tags=["function
app.include_router(
evaluations.router, prefix="/api/v1/evaluations", tags=["evaluations"]
)
+app.include_router(
+ analytics.router, prefix="/api/v1/analytics", tags=["analytics"]
+)
app.include_router(utils.router, prefix="/api/v1/utils", tags=["utils"])
# SCIM 2.0 API for identity management
diff --git a/backend/open_webui/migrations/versions/8452d01d26d7_add_chat_message_table.py b/backend/open_webui/migrations/versions/8452d01d26d7_add_chat_message_table.py
new file mode 100644
index 0000000000..5a139db3e2
--- /dev/null
+++ b/backend/open_webui/migrations/versions/8452d01d26d7_add_chat_message_table.py
@@ -0,0 +1,173 @@
+"""Add chat_message table
+
+Revision ID: 8452d01d26d7
+Revises: 374d2f66af06
+Create Date: 2026-02-01 04:00:00.000000
+
+"""
+
+import time
+import json
+import logging
+from typing import Sequence, Union
+
+from alembic import op
+import sqlalchemy as sa
+
+log = logging.getLogger(__name__)
+
+revision: str = "8452d01d26d7"
+down_revision: Union[str, None] = "374d2f66af06"
+branch_labels: Union[str, Sequence[str], None] = None
+depends_on: Union[str, Sequence[str], None] = None
+
+
+def upgrade() -> None:
+ # Step 1: Create table
+ op.create_table(
+ "chat_message",
+ sa.Column("id", sa.Text(), primary_key=True),
+ sa.Column("chat_id", sa.Text(), nullable=False, index=True),
+ sa.Column("user_id", sa.Text(), index=True),
+ sa.Column("role", sa.Text(), nullable=False),
+ sa.Column("parent_id", sa.Text(), nullable=True),
+ sa.Column("content", sa.JSON(), nullable=True),
+ sa.Column("output", sa.JSON(), nullable=True),
+ sa.Column("model_id", sa.Text(), nullable=True, index=True),
+ sa.Column("files", sa.JSON(), nullable=True),
+ sa.Column("sources", sa.JSON(), nullable=True),
+ sa.Column("embeds", sa.JSON(), nullable=True),
+ sa.Column("done", sa.Boolean(), default=True),
+ sa.Column("status_history", sa.JSON(), nullable=True),
+ sa.Column("error", sa.JSON(), nullable=True),
+ sa.Column("usage", sa.JSON(), nullable=True),
+ sa.Column("created_at", sa.BigInteger(), index=True),
+ sa.Column("updated_at", sa.BigInteger()),
+ sa.ForeignKeyConstraint(["chat_id"], ["chat.id"], ondelete="CASCADE"),
+ )
+
+ # Create composite indexes
+ op.create_index(
+ "chat_message_chat_parent_idx", "chat_message", ["chat_id", "parent_id"]
+ )
+ op.create_index(
+ "chat_message_model_created_idx", "chat_message", ["model_id", "created_at"]
+ )
+ op.create_index(
+ "chat_message_user_created_idx", "chat_message", ["user_id", "created_at"]
+ )
+
+ # Step 2: Backfill from existing chats
+ conn = op.get_bind()
+
+ chat_table = sa.table(
+ "chat",
+ sa.column("id", sa.Text()),
+ sa.column("user_id", sa.Text()),
+ sa.column("chat", sa.JSON()),
+ )
+
+ chat_message_table = sa.table(
+ "chat_message",
+ sa.column("id", sa.Text()),
+ sa.column("chat_id", sa.Text()),
+ sa.column("user_id", sa.Text()),
+ sa.column("role", sa.Text()),
+ sa.column("parent_id", sa.Text()),
+ sa.column("content", sa.JSON()),
+ sa.column("output", sa.JSON()),
+ sa.column("model_id", sa.Text()),
+ sa.column("files", sa.JSON()),
+ sa.column("sources", sa.JSON()),
+ sa.column("embeds", sa.JSON()),
+ sa.column("done", sa.Boolean()),
+ sa.column("status_history", sa.JSON()),
+ sa.column("error", sa.JSON()),
+ sa.column("usage", sa.JSON()),
+ sa.column("created_at", sa.BigInteger()),
+ sa.column("updated_at", sa.BigInteger()),
+ )
+
+ # Fetch all chats
+ chats = conn.execute(
+ sa.select(chat_table.c.id, chat_table.c.user_id, chat_table.c.chat)
+ ).fetchall()
+
+ now = int(time.time())
+ messages_inserted = 0
+ messages_failed = 0
+
+ for chat_row in chats:
+ chat_id = chat_row[0]
+ user_id = chat_row[1]
+ chat_data = chat_row[2]
+
+ if not chat_data:
+ continue
+
+ # Handle both string and dict chat data
+ if isinstance(chat_data, str):
+ try:
+ chat_data = json.loads(chat_data)
+ except Exception:
+ continue
+
+ history = chat_data.get("history", {})
+ messages = history.get("messages", {})
+
+ for message_id, message in messages.items():
+ if not isinstance(message, dict):
+ continue
+
+ role = message.get("role")
+ if not role:
+ continue
+
+ timestamp = message.get("timestamp", now)
+
+ # Normalize timestamp: convert ms to seconds, validate range
+ if timestamp > 10_000_000_000:
+ timestamp = timestamp // 1000
+ # Must be after 2020 and not too far in the future
+ if timestamp < 1577836800 or timestamp > now + 86400:
+ timestamp = now
+
+ # Use savepoint to allow individual insert failures without aborting transaction
+ savepoint = conn.begin_nested()
+ try:
+ conn.execute(
+ sa.insert(chat_message_table).values(
+ id=f"{chat_id}-{message_id}",
+ chat_id=chat_id,
+ user_id=user_id,
+ role=role,
+ parent_id=message.get("parentId"),
+ content=message.get("content"),
+ output=message.get("output"),
+ model_id=message.get("model"),
+ files=message.get("files"),
+ sources=message.get("sources"),
+ embeds=message.get("embeds"),
+ done=message.get("done", True),
+ status_history=message.get("statusHistory"),
+ error=message.get("error"),
+ created_at=timestamp,
+ updated_at=timestamp,
+ )
+ )
+ savepoint.commit()
+ messages_inserted += 1
+ except Exception as e:
+ savepoint.rollback()
+ messages_failed += 1
+ log.warning(f"Failed to insert message {message_id}: {e}")
+ continue
+
+ log.info(f"Backfilled {messages_inserted} messages into chat_message table ({messages_failed} failed)")
+
+
+def downgrade() -> None:
+ op.drop_index("chat_message_user_created_idx", table_name="chat_message")
+ op.drop_index("chat_message_model_created_idx", table_name="chat_message")
+ op.drop_index("chat_message_chat_parent_idx", table_name="chat_message")
+ op.drop_table("chat_message")
diff --git a/backend/open_webui/models/chat_messages.py b/backend/open_webui/models/chat_messages.py
new file mode 100644
index 0000000000..9254baf5d2
--- /dev/null
+++ b/backend/open_webui/models/chat_messages.py
@@ -0,0 +1,545 @@
+import json
+import time
+import uuid
+from typing import Any, Optional
+
+from sqlalchemy.orm import Session
+from open_webui.internal.db import Base, get_db_context
+
+from pydantic import BaseModel, ConfigDict
+from sqlalchemy import (
+ BigInteger,
+ Boolean,
+ Column,
+ ForeignKey,
+ Text,
+ JSON,
+ Index,
+)
+
+####################
+# Helpers
+####################
+
+
+def _normalize_timestamp(timestamp: int) -> float:
+ """Normalize and validate timestamp. Returns current time if invalid."""
+ now = time.time()
+
+ # Convert milliseconds to seconds if needed
+ if timestamp > 10_000_000_000:
+ timestamp = timestamp / 1000
+
+ # Validate: must be after 2020 and not in the future (with 1 day tolerance)
+ min_valid = 1577836800 # 2020-01-01 00:00:00 UTC
+ max_valid = now + 86400 # 1 day in the future (clock skew tolerance)
+
+ if timestamp < min_valid or timestamp > max_valid:
+ return now
+
+ return timestamp
+
+
+####################
+# ChatMessage DB Schema
+####################
+
+
+class ChatMessage(Base):
+ __tablename__ = "chat_message"
+
+ # Identity
+ id = Column(Text, primary_key=True)
+ chat_id = Column(
+ Text, ForeignKey("chat.id", ondelete="CASCADE"), nullable=False, index=True
+ )
+ user_id = Column(Text, index=True)
+
+ # Structure
+ role = Column(Text, nullable=False) # user, assistant, system
+ parent_id = Column(Text, nullable=True)
+
+ # Content
+ content = Column(JSON, nullable=True) # Can be str or list of blocks
+ output = Column(JSON, nullable=True)
+
+ # Model (for assistant messages)
+ model_id = Column(Text, nullable=True, index=True)
+
+ # Attachments
+ files = Column(JSON, nullable=True)
+ sources = Column(JSON, nullable=True)
+ embeds = Column(JSON, nullable=True)
+
+ # Status
+ done = Column(Boolean, default=True)
+ status_history = Column(JSON, nullable=True)
+ error = Column(JSON, nullable=True)
+
+ # Usage (tokens, timing, etc.)
+ usage = Column(JSON, nullable=True)
+
+ # Timestamps
+ created_at = Column(BigInteger, index=True)
+ updated_at = Column(BigInteger)
+
+ __table_args__ = (
+ Index("chat_message_chat_parent_idx", "chat_id", "parent_id"),
+ Index("chat_message_model_created_idx", "model_id", "created_at"),
+ Index("chat_message_user_created_idx", "user_id", "created_at"),
+ )
+
+
+####################
+# Pydantic Models
+####################
+
+
+class ChatMessageModel(BaseModel):
+ model_config = ConfigDict(from_attributes=True)
+
+ id: str
+ chat_id: str
+ user_id: str
+ role: str
+ parent_id: Optional[str] = None
+ content: Optional[Any] = None # str or list of blocks
+ output: Optional[list] = None
+ model_id: Optional[str] = None
+ files: Optional[list] = None
+ sources: Optional[list] = None
+ embeds: Optional[list] = None
+ done: bool = True
+ status_history: Optional[list] = None
+ error: Optional[dict] = None
+ usage: Optional[dict] = None
+ created_at: int
+ updated_at: int
+
+
+####################
+# Table Operations
+####################
+
+
+class ChatMessageTable:
+ def upsert_message(
+ self,
+ message_id: str,
+ chat_id: str,
+ user_id: str,
+ data: dict,
+ db: Optional[Session] = None,
+ ) -> Optional[ChatMessageModel]:
+ """Insert or update a chat message."""
+ with get_db_context(db) as db:
+ now = int(time.time())
+ timestamp = data.get("timestamp", now)
+
+ # Use composite ID: {chat_id}-{message_id}
+ composite_id = f"{chat_id}-{message_id}"
+
+ existing = db.get(ChatMessage, composite_id)
+ if existing:
+ # Update existing
+ if "role" in data:
+ existing.role = data["role"]
+ if "parent_id" in data:
+ existing.parent_id = data.get("parent_id") or data.get("parentId")
+ if "content" in data:
+ existing.content = data.get("content")
+ if "output" in data:
+ existing.output = data.get("output")
+ if "model_id" in data or "model" in data:
+ existing.model_id = data.get("model_id") or data.get("model")
+ if "files" in data:
+ existing.files = data.get("files")
+ if "sources" in data:
+ existing.sources = data.get("sources")
+ if "embeds" in data:
+ existing.embeds = data.get("embeds")
+ if "done" in data:
+ existing.done = data.get("done", True)
+ if "status_history" in data or "statusHistory" in data:
+ existing.status_history = data.get("status_history") or data.get(
+ "statusHistory"
+ )
+ if "error" in data:
+ existing.error = data.get("error")
+ # Extract usage - check direct field first, then info.usage
+ usage = data.get("usage")
+ if not usage:
+ info = data.get("info", {})
+ usage = info.get("usage") if info else None
+ if usage:
+ existing.usage = usage
+ existing.updated_at = now
+ db.commit()
+ db.refresh(existing)
+ return ChatMessageModel.model_validate(existing)
+ else:
+ # Insert new
+ # Extract usage - check direct field first, then info.usage
+ usage = data.get("usage")
+ if not usage:
+ info = data.get("info", {})
+ usage = info.get("usage") if info else None
+ message = ChatMessage(
+ id=composite_id,
+ chat_id=chat_id,
+ user_id=user_id,
+ role=data.get("role", "user"),
+ parent_id=data.get("parent_id") or data.get("parentId"),
+ content=data.get("content"),
+ output=data.get("output"),
+ model_id=data.get("model_id") or data.get("model"),
+ files=data.get("files"),
+ sources=data.get("sources"),
+ embeds=data.get("embeds"),
+ done=data.get("done", True),
+ status_history=data.get("status_history")
+ or data.get("statusHistory"),
+ error=data.get("error"),
+ usage=usage,
+ created_at=timestamp,
+ updated_at=now,
+ )
+ db.add(message)
+ db.commit()
+ db.refresh(message)
+ return ChatMessageModel.model_validate(message)
+
+ def get_message_by_id(
+ self, id: str, db: Optional[Session] = None
+ ) -> Optional[ChatMessageModel]:
+ with get_db_context(db) as db:
+ message = db.get(ChatMessage, id)
+ return ChatMessageModel.model_validate(message) if message else None
+
+ def get_messages_by_chat_id(
+ self, chat_id: str, db: Optional[Session] = None
+ ) -> list[ChatMessageModel]:
+ with get_db_context(db) as db:
+ messages = (
+ db.query(ChatMessage)
+ .filter_by(chat_id=chat_id)
+ .order_by(ChatMessage.created_at.asc())
+ .all()
+ )
+ return [ChatMessageModel.model_validate(message) for message in messages]
+
+ def get_messages_by_user_id(
+ self,
+ user_id: str,
+ skip: int = 0,
+ limit: int = 50,
+ db: Optional[Session] = None,
+ ) -> list[ChatMessageModel]:
+ with get_db_context(db) as db:
+ messages = (
+ db.query(ChatMessage)
+ .filter_by(user_id=user_id)
+ .order_by(ChatMessage.created_at.desc())
+ .offset(skip)
+ .limit(limit)
+ .all()
+ )
+ return [ChatMessageModel.model_validate(message) for message in messages]
+
+ def get_messages_by_model_id(
+ self,
+ model_id: str,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ skip: int = 0,
+ limit: int = 100,
+ db: Optional[Session] = None,
+ ) -> list[ChatMessageModel]:
+ with get_db_context(db) as db:
+ query = db.query(ChatMessage).filter_by(model_id=model_id)
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+ messages = (
+ query.order_by(ChatMessage.created_at.desc())
+ .offset(skip)
+ .limit(limit)
+ .all()
+ )
+ return [ChatMessageModel.model_validate(message) for message in messages]
+
+ def delete_messages_by_chat_id(
+ self, chat_id: str, db: Optional[Session] = None
+ ) -> bool:
+ with get_db_context(db) as db:
+ db.query(ChatMessage).filter_by(chat_id=chat_id).delete()
+ db.commit()
+ return True
+
+ # Analytics methods
+ def get_message_count_by_model(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, int]:
+ with get_db_context(db) as db:
+ from sqlalchemy import func
+
+ query = db.query(
+ ChatMessage.model_id, func.count(ChatMessage.id).label("count")
+ ).filter(ChatMessage.role == "assistant", ChatMessage.model_id.isnot(None))
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.group_by(ChatMessage.model_id).all()
+ return {row.model_id: row.count for row in results}
+
+ def get_token_usage_by_model(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, dict]:
+ """Aggregate token usage by model using database-level aggregation."""
+ with get_db_context(db) as db:
+ from sqlalchemy import func, cast, Integer
+
+ dialect = db.bind.dialect.name
+
+ if dialect == "sqlite":
+ input_tokens = cast(
+ func.json_extract(ChatMessage.usage, "$.input_tokens"), Integer
+ )
+ output_tokens = cast(
+ func.json_extract(ChatMessage.usage, "$.output_tokens"), Integer
+ )
+ elif dialect == "postgresql":
+ # Use json_extract_path_text for PostgreSQL JSON columns
+ input_tokens = cast(
+ func.json_extract_path_text(ChatMessage.usage, "input_tokens"), Integer
+ )
+ output_tokens = cast(
+ func.json_extract_path_text(ChatMessage.usage, "output_tokens"), Integer
+ )
+ else:
+ raise NotImplementedError(f"Unsupported dialect: {dialect}")
+
+ query = db.query(
+ ChatMessage.model_id,
+ func.coalesce(func.sum(input_tokens), 0).label("input_tokens"),
+ func.coalesce(func.sum(output_tokens), 0).label("output_tokens"),
+ func.count(ChatMessage.id).label("message_count"),
+ ).filter(
+ ChatMessage.role == "assistant",
+ ChatMessage.model_id.isnot(None),
+ ChatMessage.usage.isnot(None),
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.group_by(ChatMessage.model_id).all()
+
+ return {
+ row.model_id: {
+ "input_tokens": row.input_tokens,
+ "output_tokens": row.output_tokens,
+ "total_tokens": row.input_tokens + row.output_tokens,
+ "message_count": row.message_count,
+ }
+ for row in results
+ }
+
+ def get_token_usage_by_user(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, dict]:
+ """Aggregate token usage by user using database-level aggregation."""
+ with get_db_context(db) as db:
+ from sqlalchemy import func, cast, Integer
+
+ dialect = db.bind.dialect.name
+
+ if dialect == "sqlite":
+ input_tokens = cast(
+ func.json_extract(ChatMessage.usage, "$.input_tokens"), Integer
+ )
+ output_tokens = cast(
+ func.json_extract(ChatMessage.usage, "$.output_tokens"), Integer
+ )
+ elif dialect == "postgresql":
+ # Use json_extract_path_text for PostgreSQL JSON columns
+ input_tokens = cast(
+ func.json_extract_path_text(ChatMessage.usage, "input_tokens"), Integer
+ )
+ output_tokens = cast(
+ func.json_extract_path_text(ChatMessage.usage, "output_tokens"), Integer
+ )
+ else:
+ raise NotImplementedError(f"Unsupported dialect: {dialect}")
+
+ query = db.query(
+ ChatMessage.user_id,
+ func.coalesce(func.sum(input_tokens), 0).label("input_tokens"),
+ func.coalesce(func.sum(output_tokens), 0).label("output_tokens"),
+ func.count(ChatMessage.id).label("message_count"),
+ ).filter(
+ ChatMessage.role == "assistant",
+ ChatMessage.user_id.isnot(None),
+ ChatMessage.usage.isnot(None),
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.group_by(ChatMessage.user_id).all()
+
+ return {
+ row.user_id: {
+ "input_tokens": row.input_tokens,
+ "output_tokens": row.output_tokens,
+ "total_tokens": row.input_tokens + row.output_tokens,
+ "message_count": row.message_count,
+ }
+ for row in results
+ }
+
+ def get_message_count_by_user(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, int]:
+ with get_db_context(db) as db:
+ from sqlalchemy import func
+
+ query = db.query(
+ ChatMessage.user_id, func.count(ChatMessage.id).label("count")
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.group_by(ChatMessage.user_id).all()
+ return {row.user_id: row.count for row in results}
+
+ def get_message_count_by_chat(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, int]:
+ with get_db_context(db) as db:
+ from sqlalchemy import func
+
+ query = db.query(
+ ChatMessage.chat_id, func.count(ChatMessage.id).label("count")
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.group_by(ChatMessage.chat_id).all()
+ return {row.chat_id: row.count for row in results}
+
+ def get_daily_message_counts_by_model(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, dict[str, int]]:
+ """Get message counts grouped by day and model."""
+ with get_db_context(db) as db:
+ from datetime import datetime, timedelta
+
+ query = db.query(ChatMessage.created_at, ChatMessage.model_id).filter(
+ ChatMessage.role == "assistant",
+ ChatMessage.model_id.isnot(None)
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.all()
+
+ # Group by date -> model -> count
+ daily_counts: dict[str, dict[str, int]] = {}
+ for timestamp, model_id in results:
+ date_str = datetime.fromtimestamp(_normalize_timestamp(timestamp)).strftime("%Y-%m-%d")
+ if date_str not in daily_counts:
+ daily_counts[date_str] = {}
+ daily_counts[date_str][model_id] = daily_counts[date_str].get(model_id, 0) + 1
+
+ # Fill in missing days
+ if start_date and end_date:
+ current = datetime.fromtimestamp(_normalize_timestamp(start_date))
+ end_dt = datetime.fromtimestamp(_normalize_timestamp(end_date))
+ while current <= end_dt:
+ date_str = current.strftime("%Y-%m-%d")
+ if date_str not in daily_counts:
+ daily_counts[date_str] = {}
+ current += timedelta(days=1)
+
+ return daily_counts
+
+ def get_hourly_message_counts_by_model(
+ self,
+ start_date: Optional[int] = None,
+ end_date: Optional[int] = None,
+ db: Optional[Session] = None,
+ ) -> dict[str, dict[str, int]]:
+ """Get message counts grouped by hour and model."""
+ with get_db_context(db) as db:
+ from datetime import datetime, timedelta
+
+ query = db.query(ChatMessage.created_at, ChatMessage.model_id).filter(
+ ChatMessage.role == "assistant",
+ ChatMessage.model_id.isnot(None)
+ )
+
+ if start_date:
+ query = query.filter(ChatMessage.created_at >= start_date)
+ if end_date:
+ query = query.filter(ChatMessage.created_at <= end_date)
+
+ results = query.all()
+
+ # Group by hour -> model -> count
+ hourly_counts: dict[str, dict[str, int]] = {}
+ for timestamp, model_id in results:
+ hour_str = datetime.fromtimestamp(_normalize_timestamp(timestamp)).strftime("%Y-%m-%d %H:00")
+ if hour_str not in hourly_counts:
+ hourly_counts[hour_str] = {}
+ hourly_counts[hour_str][model_id] = hourly_counts[hour_str].get(model_id, 0) + 1
+
+ # Fill in missing hours
+ if start_date and end_date:
+ current = datetime.fromtimestamp(_normalize_timestamp(start_date)).replace(minute=0, second=0, microsecond=0)
+ end_dt = datetime.fromtimestamp(_normalize_timestamp(end_date))
+ while current <= end_dt:
+ hour_str = current.strftime("%Y-%m-%d %H:00")
+ if hour_str not in hourly_counts:
+ hourly_counts[hour_str] = {}
+ current += timedelta(hours=1)
+
+ return hourly_counts
+
+
+ChatMessages = ChatMessageTable()
diff --git a/backend/open_webui/models/chats.py b/backend/open_webui/models/chats.py
index eb0763048b..51a714cea7 100644
--- a/backend/open_webui/models/chats.py
+++ b/backend/open_webui/models/chats.py
@@ -8,6 +8,7 @@ from sqlalchemy.orm import Session
from open_webui.internal.db import Base, JSONField, get_db, get_db_context
from open_webui.models.tags import TagModel, Tag, Tags
from open_webui.models.folders import Folders
+from open_webui.models.chat_messages import ChatMessages
from open_webui.utils.misc import sanitize_data_for_db, sanitize_text_for_db
from pydantic import BaseModel, ConfigDict
@@ -314,6 +315,22 @@ class ChatTable:
db.add(chat_item)
db.commit()
db.refresh(chat_item)
+
+ # Dual-write initial messages to chat_message table
+ try:
+ history = form_data.chat.get("history", {})
+ messages = history.get("messages", {})
+ for message_id, message in messages.items():
+ if isinstance(message, dict) and message.get("role"):
+ ChatMessages.upsert_message(
+ message_id=message_id,
+ chat_id=id,
+ user_id=user_id,
+ data=message,
+ )
+ except Exception as e:
+ log.warning(f"Failed to write initial messages to chat_message table: {e}")
+
return ChatModel.model_validate(chat_item) if chat_item else None
def _chat_import_form_to_chat_model(
@@ -356,6 +373,23 @@ class ChatTable:
db.add_all(chats)
db.commit()
+
+ # Dual-write messages to chat_message table
+ try:
+ for form_data, chat_obj in zip(chat_import_forms, chats):
+ history = form_data.chat.get("history", {})
+ messages = history.get("messages", {})
+ for message_id, message in messages.items():
+ if isinstance(message, dict) and message.get("role"):
+ ChatMessages.upsert_message(
+ message_id=message_id,
+ chat_id=chat_obj.id,
+ user_id=user_id,
+ data=message,
+ )
+ except Exception as e:
+ log.warning(f"Failed to write imported messages to chat_message table: {e}")
+
return [ChatModel.model_validate(chat) for chat in chats]
def update_chat_by_id(
@@ -458,6 +492,18 @@ class ChatTable:
history["currentId"] = message_id
chat["history"] = history
+
+ # Dual-write to chat_message table
+ try:
+ ChatMessages.upsert_message(
+ message_id=message_id,
+ chat_id=id,
+ user_id=self.get_chat_by_id(id).user_id,
+ data=history["messages"][message_id],
+ )
+ except Exception as e:
+ log.warning(f"Failed to write to chat_message table: {e}")
+
return self.update_chat_by_id(id, chat)
def add_message_status_to_chat_by_id_and_message_id(
diff --git a/backend/open_webui/models/users.py b/backend/open_webui/models/users.py
index deeda29a85..43c5a94f81 100644
--- a/backend/open_webui/models/users.py
+++ b/backend/open_webui/models/users.py
@@ -243,6 +243,7 @@ class UsersTable:
email: str,
profile_image_url: str = "/user.png",
role: str = "pending",
+ username: Optional[str] = None,
oauth: Optional[dict] = None,
db: Optional[Session] = None,
) -> Optional[UserModel]:
@@ -257,6 +258,7 @@ class UsersTable:
"last_active_at": int(time.time()),
"created_at": int(time.time()),
"updated_at": int(time.time()),
+ "username": username,
"oauth": oauth,
}
)
diff --git a/backend/open_webui/routers/analytics.py b/backend/open_webui/routers/analytics.py
new file mode 100644
index 0000000000..d8f6928ffc
--- /dev/null
+++ b/backend/open_webui/routers/analytics.py
@@ -0,0 +1,247 @@
+from typing import Optional
+import logging
+from fastapi import APIRouter, Depends, Query
+from pydantic import BaseModel
+
+from open_webui.models.chat_messages import ChatMessages, ChatMessageModel
+from open_webui.utils.auth import get_admin_user
+from open_webui.internal.db import get_session
+from sqlalchemy.orm import Session
+
+log = logging.getLogger(__name__)
+
+
+router = APIRouter()
+
+
+####################
+# Response Models
+####################
+
+
+class ModelAnalyticsEntry(BaseModel):
+ model_id: str
+ count: int
+
+
+class ModelAnalyticsResponse(BaseModel):
+ models: list[ModelAnalyticsEntry]
+
+
+class UserAnalyticsEntry(BaseModel):
+ user_id: str
+ name: Optional[str] = None
+ email: Optional[str] = None
+ count: int
+ input_tokens: int = 0
+ output_tokens: int = 0
+ total_tokens: int = 0
+
+
+class UserAnalyticsResponse(BaseModel):
+ users: list[UserAnalyticsEntry]
+
+
+####################
+# Endpoints
+####################
+
+
+@router.get("/models", response_model=ModelAnalyticsResponse)
+async def get_model_analytics(
+ start_date: Optional[int] = Query(None, description="Start timestamp (epoch)"),
+ end_date: Optional[int] = Query(None, description="End timestamp (epoch)"),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Get message counts per model."""
+ counts = ChatMessages.get_message_count_by_model(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ models = [
+ ModelAnalyticsEntry(model_id=model_id, count=count)
+ for model_id, count in sorted(counts.items(), key=lambda x: -x[1])
+ ]
+ return ModelAnalyticsResponse(models=models)
+
+
+@router.get("/users", response_model=UserAnalyticsResponse)
+async def get_user_analytics(
+ start_date: Optional[int] = Query(None, description="Start timestamp (epoch)"),
+ end_date: Optional[int] = Query(None, description="End timestamp (epoch)"),
+ limit: int = Query(50, description="Max users to return"),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Get message counts and token usage per user with user info."""
+ from open_webui.models.users import Users
+
+ counts = ChatMessages.get_message_count_by_user(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ token_usage = ChatMessages.get_token_usage_by_user(
+ start_date=start_date, end_date=end_date, db=db
+ )
+
+ # Get user info for top users
+ top_user_ids = [uid for uid, _ in sorted(counts.items(), key=lambda x: -x[1])[:limit]]
+ user_info = {u.id: u for u in Users.get_users_by_user_ids(top_user_ids, db=db)}
+
+ users = []
+ for user_id in top_user_ids:
+ u = user_info.get(user_id)
+ tokens = token_usage.get(user_id, {})
+ users.append(UserAnalyticsEntry(
+ user_id=user_id,
+ name=u.name if u else None,
+ email=u.email if u else None,
+ count=counts[user_id],
+ input_tokens=tokens.get("input_tokens", 0),
+ output_tokens=tokens.get("output_tokens", 0),
+ total_tokens=tokens.get("total_tokens", 0),
+ ))
+
+ return UserAnalyticsResponse(users=users)
+
+
+@router.get("/messages", response_model=list[ChatMessageModel])
+async def get_messages(
+ model_id: Optional[str] = Query(None, description="Filter by model ID"),
+ user_id: Optional[str] = Query(None, description="Filter by user ID"),
+ chat_id: Optional[str] = Query(None, description="Filter by chat ID"),
+ start_date: Optional[int] = Query(None, description="Start timestamp (epoch)"),
+ end_date: Optional[int] = Query(None, description="End timestamp (epoch)"),
+ skip: int = Query(0),
+ limit: int = Query(50, le=100),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Query messages with filters."""
+ if chat_id:
+ return ChatMessages.get_messages_by_chat_id(chat_id=chat_id, db=db)
+ elif model_id:
+ return ChatMessages.get_messages_by_model_id(
+ model_id=model_id,
+ start_date=start_date,
+ end_date=end_date,
+ skip=skip,
+ limit=limit,
+ db=db,
+ )
+ elif user_id:
+ return ChatMessages.get_messages_by_user_id(
+ user_id=user_id, skip=skip, limit=limit, db=db
+ )
+ else:
+ # Return empty if no filter specified
+ return []
+
+
+class SummaryResponse(BaseModel):
+ total_messages: int
+ total_chats: int
+ total_models: int
+ total_users: int
+
+
+@router.get("/summary", response_model=SummaryResponse)
+async def get_summary(
+ start_date: Optional[int] = Query(None, description="Start timestamp (epoch)"),
+ end_date: Optional[int] = Query(None, description="End timestamp (epoch)"),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Get summary statistics for the dashboard."""
+ model_counts = ChatMessages.get_message_count_by_model(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ user_counts = ChatMessages.get_message_count_by_user(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ chat_counts = ChatMessages.get_message_count_by_chat(
+ start_date=start_date, end_date=end_date, db=db
+ )
+
+ return SummaryResponse(
+ total_messages=sum(model_counts.values()),
+ total_chats=len(chat_counts),
+ total_models=len(model_counts),
+ total_users=len(user_counts),
+ )
+
+
+class DailyStatsEntry(BaseModel):
+ date: str
+ models: dict[str, int]
+
+
+class DailyStatsResponse(BaseModel):
+ data: list[DailyStatsEntry]
+
+
+@router.get("/daily", response_model=DailyStatsResponse)
+async def get_daily_stats(
+ start_date: Optional[int] = Query(None, description="Start timestamp (epoch)"),
+ end_date: Optional[int] = Query(None, description="End timestamp (epoch)"),
+ granularity: str = Query("daily", description="Granularity: 'hourly' or 'daily'"),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Get message counts grouped by model for time-series chart."""
+ if granularity == "hourly":
+ counts = ChatMessages.get_hourly_message_counts_by_model(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ else:
+ counts = ChatMessages.get_daily_message_counts_by_model(
+ start_date=start_date, end_date=end_date, db=db
+ )
+ return DailyStatsResponse(
+ data=[
+ DailyStatsEntry(date=date, models=models)
+ for date, models in sorted(counts.items())
+ ]
+ )
+
+
+class TokenUsageEntry(BaseModel):
+ model_id: str
+ input_tokens: int
+ output_tokens: int
+ total_tokens: int
+ message_count: int
+
+
+class TokenUsageResponse(BaseModel):
+ models: list[TokenUsageEntry]
+ total_input_tokens: int
+ total_output_tokens: int
+ total_tokens: int
+
+
+@router.get("/tokens", response_model=TokenUsageResponse)
+async def get_token_usage(
+ start_date: Optional[int] = Query(None),
+ end_date: Optional[int] = Query(None),
+ user=Depends(get_admin_user),
+ db: Session = Depends(get_session),
+):
+ """Get token usage aggregated by model."""
+ usage = ChatMessages.get_token_usage_by_model(
+ start_date=start_date, end_date=end_date, db=db
+ )
+
+ models = [
+ TokenUsageEntry(model_id=model_id, **data)
+ for model_id, data in sorted(usage.items(), key=lambda x: -x[1]["total_tokens"])
+ ]
+
+ total_input = sum(m.input_tokens for m in models)
+ total_output = sum(m.output_tokens for m in models)
+
+ return TokenUsageResponse(
+ models=models,
+ total_input_tokens=total_input,
+ total_output_tokens=total_output,
+ total_tokens=total_input + total_output,
+ )
diff --git a/backend/open_webui/routers/openai.py b/backend/open_webui/routers/openai.py
index 44575e57f2..b1d31afb8f 100644
--- a/backend/open_webui/routers/openai.py
+++ b/backend/open_webui/routers/openai.py
@@ -794,6 +794,79 @@ def convert_to_azure_payload(url, payload: dict, api_version: str):
return url, payload
+def convert_to_responses_payload(payload: dict) -> dict:
+ """
+ Convert Chat Completions payload to Responses API format.
+
+ Chat Completions: { messages: [{role, content}], ... }
+ Responses API: { input: [{type: "message", role, content: [...]}], instructions: "system" }
+ """
+ messages = payload.pop("messages", [])
+
+ system_content = ""
+ input_items = []
+
+ for msg in messages:
+ role = msg.get("role", "user")
+ content = msg.get("content", "")
+
+ # Check for stored output items (from previous Responses API turn)
+ stored_output = msg.get("output")
+ if stored_output and isinstance(stored_output, list):
+ input_items.extend(stored_output)
+ continue
+
+ if role == "system":
+ if isinstance(content, str):
+ system_content = content
+ elif isinstance(content, list):
+ system_content = "\n".join(p.get("text", "") for p in content if p.get("type") == "text")
+ continue
+
+ # Convert content format
+ text_type = "output_text" if role == "assistant" else "input_text"
+
+ if isinstance(content, str):
+ content_parts = [{"type": text_type, "text": content}]
+ elif isinstance(content, list):
+ content_parts = []
+ for part in content:
+ if part.get("type") == "text":
+ content_parts.append({"type": text_type, "text": part.get("text", "")})
+ elif part.get("type") == "image_url":
+ url_data = part.get("image_url", {})
+ url = url_data.get("url", "") if isinstance(url_data, dict) else url_data
+ content_parts.append({"type": "input_image", "image_url": url})
+ else:
+ content_parts = [{"type": text_type, "text": str(content)}]
+
+ input_items.append({
+ "type": "message",
+ "role": role,
+ "content": content_parts
+ })
+
+ responses_payload = {**payload, "input": input_items}
+
+ if system_content:
+ responses_payload["instructions"] = system_content
+
+ if "max_tokens" in responses_payload:
+ responses_payload["max_output_tokens"] = responses_payload.pop("max_tokens")
+
+ return responses_payload
+
+
+
+def convert_responses_result(response: dict) -> dict:
+ """
+ Convert non-streaming Responses API result.
+ Just add done flag - pass through raw response, frontend handles output.
+ """
+ response["done"] = True
+ return response
+
+
@router.post("/chat/completions")
async def generate_chat_completion(
request: Request,
@@ -915,6 +988,8 @@ async def generate_chat_completion(
request, url, key, api_config, metadata, user=user
)
+ is_responses = api_config.get("api_type") == "responses"
+
if api_config.get("azure", False):
api_version = api_config.get("api_version", "2023-03-15-preview")
request_url, payload = convert_to_azure_payload(url, payload, api_version)
@@ -925,9 +1000,18 @@ async def generate_chat_completion(
headers["api-key"] = key
headers["api-version"] = api_version
- request_url = f"{request_url}/chat/completions?api-version={api_version}"
+
+ if is_responses:
+ payload = convert_to_responses_payload(payload)
+ request_url = f"{request_url}/responses?api-version={api_version}"
+ else:
+ request_url = f"{request_url}/chat/completions?api-version={api_version}"
else:
- request_url = f"{url}/chat/completions"
+ if is_responses:
+ payload = convert_to_responses_payload(payload)
+ request_url = f"{url}/responses"
+ else:
+ request_url = f"{url}/chat/completions"
payload = json.dumps(payload)
@@ -974,6 +1058,10 @@ async def generate_chat_completion(
else:
return PlainTextResponse(status_code=r.status, content=response)
+ # Convert Responses API result to simple format
+ if is_responses and isinstance(response, dict):
+ response = convert_responses_result(response)
+
return response
except Exception as e:
log.exception(e)
diff --git a/backend/open_webui/utils/middleware.py b/backend/open_webui/utils/middleware.py
index 81e07df94e..0f40ce5941 100644
--- a/backend/open_webui/utils/middleware.py
+++ b/backend/open_webui/utils/middleware.py
@@ -107,6 +107,7 @@ from open_webui.utils.filter import (
)
from open_webui.utils.code_interpreter import execute_code_jupyter
from open_webui.utils.payload import apply_system_prompt_to_body
+from open_webui.utils.response import normalize_usage
from open_webui.utils.mcp.client import MCPClient
@@ -293,6 +294,476 @@ def get_citation_source_from_tool_result(
]
+def split_content_and_whitespace(content):
+ content_stripped = content.rstrip()
+ original_whitespace = (
+ content[len(content_stripped) :] if len(content) > len(content_stripped) else ""
+ )
+ return content_stripped, original_whitespace
+
+
+def is_opening_code_block(content):
+ backtick_segments = content.split("```")
+ # Even number of segments means the last backticks are opening a new block
+ return len(backtick_segments) > 1 and len(backtick_segments) % 2 == 0
+
+
+def serialize_output(output: list) -> str:
+ """
+ Convert OR-aligned output items to HTML for display.
+ For LLM consumption, use convert_output_to_messages() instead.
+ """
+ content = ""
+
+ # First pass: collect function_call_output items by call_id for lookup
+ tool_outputs = {}
+ for item in output:
+ if item.get("type") == "function_call_output":
+ tool_outputs[item.get("call_id")] = item
+
+ # Second pass: render items in order
+ for idx, item in enumerate(output):
+ item_type = item.get("type", "")
+
+ if item_type == "message":
+ for content_part in item.get("content", []):
+ if "text" in content_part:
+ text = content_part.get("text", "").strip()
+ if text:
+ content = f"{content}{text}\n"
+
+ elif item_type == "function_call":
+ # Render tool call inline with its result (if available)
+ if content and not content.endswith("\n"):
+ content += "\n"
+
+ call_id = item.get("call_id", "")
+ name = item.get("name", "")
+ arguments = item.get("arguments", "")
+
+ result_item = tool_outputs.get(call_id)
+ if result_item:
+ result_text = ""
+ for out in result_item.get("output", []):
+ if "text" in out:
+ result_text += out.get("text", "")
+ files = result_item.get("files")
+ embeds = result_item.get("embeds", "")
+
+ content += f'Tool Executed
\nExecuting...
\nThought for {duration or 0} seconds
\n{display}\nThinking…
\n{display}\nTool Executed
\nExecuting...
\nThought for {duration or 0} seconds
\n{display}\nThinking…
\n{display}\nAnalyzed
\n```{lang}\n{code}\n```\nAnalyzing...
\n```{lang}\n{code}\n```\n