"""
title: Openlayer Tracing
author: Openlayer
author_url: https://www.openlayer.com
version: 1.0.0
required_open_webui_version: 0.11.0
requirements: openlayer>=0.36.0
license: MIT
description: Sends each Open WebUI chat turn to Openlayer as a trace, including tool calls, knowledge sources, token usage, user, session, and (optionally) the images, audio, and files in the conversation.
"""
import asyncio
import base64
import json
import logging
import mimetypes
import os
import re
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field
log = logging.getLogger("openlayer-filter")
_START_KEY = "openlayer_start_time"
_INLET_MESSAGES_KEY = "openlayer_inlet_messages"
_TURN_FILES_KEY = "openlayer_turn_files"
_DATA_URL = re.compile(r"^data:([^;,]+)[^,]*;base64,(.+)$", re.DOTALL)
_FILE_ID = re.compile(r"/api/v1/files/([0-9a-fA-F-]{36})")
def _text(content: Any) -> str:
"""Flatten OpenAI-style message content (str or list of parts) to text."""
if isinstance(content, str):
return content
if isinstance(content, list):
parts = []
for part in content:
if not isinstance(part, dict):
continue
if part.get("type") in ("text", "input_text", "output_text"):
parts.append(part.get("text", ""))
elif part.get("type") in ("image_url", "input_image"):
parts.append("[image]")
elif part.get("type") == "input_audio":
parts.append("[audio]")
elif part.get("type") == "file":
parts.append("[file]")
return " ".join(p for p in parts if p)
return "" if content is None else str(content)
def _file_refs(entries: Any) -> List[tuple]:
"""(file id, name) pairs from Open WebUI file entries; skips collections and web pages."""
refs = []
for entry in entries or []:
if not isinstance(entry, dict) or entry.get("type") not in (None, "file", "image"):
continue
info = entry.get("file") if isinstance(entry.get("file"), dict) else {}
file_id = info.get("id") or entry.get("id")
if file_id:
refs.append((file_id, info.get("filename") or entry.get("name") or "file"))
return refs
def _usage(usage: Dict[str, Any]) -> Dict[str, Optional[int]]:
"""Read Open WebUI's normalized usage.
Since v0.11, input_tokens/output_tokens/total_tokens are cumulative for the
whole turn (every model call in a tool loop), while prompt_tokens and
completion_tokens only describe the last call. Prefer the cumulative fields.
"""
usage = usage or {}
prompt = usage.get("input_tokens", usage.get("prompt_tokens", usage.get("prompt_eval_count")))
completion = usage.get("output_tokens", usage.get("completion_tokens", usage.get("eval_count")))
total = usage.get("total_tokens")
if total is None and prompt is not None and completion is not None:
total = int(prompt) + int(completion)
return {"prompt": prompt, "completion": completion, "total": total}
def _stamp(step, start: float, end: Optional[float] = None) -> None:
"""Pin a step's timing (create_step only fills in times that are unset)."""
step.start_time = start
step.end_time = end if end is not None else start
step.latency = (step.end_time - start) * 1000
class _Media:
"""Turns Open WebUI media into Openlayer attachments, once per unique file."""
def __init__(self, max_bytes: int, user_id: Optional[str]):
from openlayer.lib.tracing import Attachment
from openlayer.lib.tracing.content import AudioContent, FileContent, ImageContent
self._attachment = Attachment
self._kinds = {"image": ImageContent, "audio": AudioContent}
self._file = FileContent
self.max_bytes = max_bytes
self.user_id = user_id
self.by_checksum: Dict[str, Any] = {}
self.skipped: List[str] = []
def item(self, data: bytes, name: str, media_type: str):
"""Wrap bytes in the content item the platform renders for this media type."""
if len(data) > self.max_bytes:
self.skipped.append(f"{name} ({len(data)} bytes)")
return None
# from_bytes keeps the bytes out of the published trace if the upload fails.
attachment = self._attachment.from_bytes(data, name=name, media_type=media_type)
attachment = self.by_checksum.setdefault(attachment.checksum_md5, attachment)
content_class = self._kinds.get(media_type.split("/")[0], self._file)
return content_class(attachment)
def from_data_url(self, url: str, kind: str):
match = _DATA_URL.match(url or "")
if not match:
return None
media_type = match.group(1)
try:
data = base64.b64decode(match.group(2))
except (ValueError, TypeError):
return None
extension = (mimetypes.guess_extension(media_type) or ".bin").lstrip(".")
return self.item(data, f"{kind}.{extension}", media_type)
async def from_file_id(self, file_id: str, fallback_name: str = "file"):
"""Read a file from Open WebUI's storage (local disk, S3, GCS, or Azure)."""
from open_webui.models.files import Files
from open_webui.storage.provider import Storage
record = await Files.get_file_by_id(file_id)
# Only the chat user's own files: a prompt or tool result could mention any file id.
if not record or not record.path or record.user_id != self.user_id:
return None
meta = record.meta or {}
size = meta.get("size")
name = record.filename or meta.get("name") or fallback_name
if isinstance(size, int) and size > self.max_bytes:
self.skipped.append(f"{name} ({size} bytes)")
return None
media_type = meta.get("content_type") or mimetypes.guess_type(name)[0] or "application/octet-stream"
local_path = await asyncio.to_thread(Storage.get_file, record.path)
data = await asyncio.to_thread(Path(local_path).read_bytes)
return self.item(data, name, media_type)
class Filter:
class Valves(BaseModel):
priority: int = Field(default=0, description="Run order among filters (lower runs first).")
openlayer_api_key: str = Field(
default_factory=lambda: os.getenv("OPENLAYER_API_KEY", ""),
description="Openlayer API key. Keep on Default to use the OPENLAYER_API_KEY environment variable.",
)
inference_pipeline_id: str = Field(
default_factory=lambda: os.getenv("OPENLAYER_INFERENCE_PIPELINE_ID", ""),
description="ID of the Openlayer data source that receives traces.",
)
base_url: str = Field(
default_factory=lambda: os.getenv("OPENLAYER_BASE_URL", ""),
description="Only for self-hosted Openlayer, e.g. https://openlayer.example.com/v1.",
)
capture_tool_calls: bool = Field(default=True, description="Add a step for every tool call and its result.")
capture_sources: bool = Field(default=True, description="Add a retrieval step with knowledge and web search sources.")
upload_attachments: bool = Field(
default=False,
description="Upload the images, audio, and files in each chat turn to your Openlayer workspace storage so they show up in the trace.",
)
max_attachment_mb: int = Field(default=20, description="Skip attachments larger than this.")
routes: str = Field(
default="",
description=(
"Optional JSON map from a model ID or a run type (subagent, automation, api, ui) to a data source ID, "
'e.g. {"subagent": "<id>", "support-agent": "<id>"}. Turns that match nothing go to the default data source.'
),
)
debug: bool = Field(default=False, description="Log what is sent to Openlayer.")
def __init__(self):
self.valves = self.Valves()
self._configured_for = None
def _ready(self) -> bool:
v = self.valves
if not v.openlayer_api_key or not v.inference_pipeline_id:
return False
key = (v.openlayer_api_key, v.inference_pipeline_id, v.base_url, v.upload_attachments)
if key != self._configured_for:
from openlayer.lib import init
kwargs = {
"api_key": v.openlayer_api_key,
"inference_pipeline_id": v.inference_pipeline_id,
"attachment_upload_enabled": v.upload_attachments,
# Open WebUI makes its own provider calls; don't patch SDK clients here.
"auto_instrument": False,
}
if v.base_url:
kwargs["base_url"] = v.base_url
init(**kwargs)
self._configured_for = key
return True
async def inlet(
self, body: dict, __metadata__: Optional[dict] = None, __model__: Optional[dict] = None
) -> dict:
# __metadata__ is the same dict object in outlet, so values stashed here travel with the turn.
if __metadata__ is not None:
__metadata__.setdefault(_START_KEY, time.time())
# Outlet only gets text content, so keep the request's media parts
# (Open WebUI has already inlined images as data URLs at this point).
__metadata__[_INLET_MESSAGES_KEY] = [dict(m) for m in body.get("messages") or []]
if self.valves.upload_attachments:
# Open WebUI resends every file in the chat on each turn; keep only this turn's.
user_message = __metadata__.get("user_message")
if isinstance(user_message, dict): # UI chats
turn_files = _file_refs(user_message.get("files"))
else: # API callers: the request's own files, minus the model's knowledge
knowledge = ((__model__ or {}).get("info") or {}).get("meta", {}).get("knowledge") or []
knowledge_ids = {k.get("id") for k in knowledge if isinstance(k, dict)}
turn_files = [r for r in _file_refs(body.get("files")) if r[0] not in knowledge_ids]
__metadata__[_TURN_FILES_KEY] = turn_files
return body
async def outlet(
self,
body: dict,
__user__: Optional[dict] = None,
__metadata__: Optional[dict] = None,
__model__: Optional[dict] = None,
) -> dict:
try:
if self._ready():
await self._trace(body, __user__ or {}, __metadata__ or {}, __model__ or {})
except Exception:
# Never break the chat because tracing failed.
log.exception("Openlayer: failed to trace chat turn")
finally:
if __metadata__ is not None:
__metadata__.pop(_INLET_MESSAGES_KEY, None)
__metadata__.pop(_TURN_FILES_KEY, None)
return body
def _routes(self) -> Dict[str, str]:
if not self.valves.routes.strip():
return {}
try:
routes = json.loads(self.valves.routes)
except ValueError:
log.error("Openlayer: the routes valve is not valid JSON; using the default data source")
return {}
return {str(k): str(v) for k, v in routes.items() if v} if isinstance(routes, dict) else {}
@staticmethod
async def _run_context(metadata: dict, chat_id: str) -> Dict[str, Any]:
"""What started this turn: a sub-agent, an automation, an API call, or the UI."""
if metadata.get("automation_id"):
return {"run_type": "automation", "automation_id": metadata["automation_id"]}
if metadata.get("internal") and chat_id:
try:
from open_webui.models.chats import Chats
chat = await Chats.get_chat_by_id(chat_id)
chat_meta = (chat.meta if chat else None) or {}
except Exception:
chat_meta = {}
if chat_meta.get("type") == "subagent":
return {
"run_type": "subagent",
"parent_chat_id": chat_meta.get("parent_chat_id"),
"parent_message_id": chat_meta.get("parent_message_id"),
"delegation_id": chat_meta.get("delegation_id"),
}
return {"run_type": "ui" if metadata.get("session_id") else "api"}
async def _prompt(self, messages: List[dict], metadata: dict, media: Optional[_Media]) -> List[dict]:
"""Build the prompt, with media as content items when uploads are on."""
from openlayer.lib.tracing.content import TextContent
prompt = []
for message in messages:
content = message.get("content")
if media is None or not isinstance(content, list):
prompt.append({"role": message.get("role"), "content": _text(content)})
continue
items = []
for part in content:
kind = part.get("type") if isinstance(part, dict) else None
item = None
if kind in ("text", "input_text"):
item = TextContent(text=part.get("text", ""))
elif kind in ("image_url", "input_image"):
image = part.get("image_url")
url = image.get("url") if isinstance(image, dict) else image
item = media.from_data_url(url or "", "image")
elif kind == "input_audio":
audio = part.get("input_audio") or {}
fmt = audio.get("format") or "wav"
item = media.from_data_url(f"data:audio/{fmt};base64,{audio.get('data', '')}", "audio")
elif kind == "file":
item = media.from_data_url((part.get("file") or {}).get("file_data", ""), "file")
items.append(item or TextContent(text=_text([part])))
prompt.append({"role": message.get("role"), "content": items})
if media is not None:
# Files the user attached in this turn (documents, audio, images) go on the latest user message.
files = []
for file_id, name in metadata.get(_TURN_FILES_KEY) or []:
item = await media.from_file_id(file_id, name)
if item and all(f.attachment is not item.attachment for f in files):
files.append(item)
last_user = next((p for p in reversed(prompt) if p["role"] == "user"), None)
if last_user and files:
if isinstance(last_user["content"], str):
last_user["content"] = [TextContent(text=last_user["content"])]
present = {id(getattr(i, "attachment", None)) for i in last_user["content"]}
last_user["content"].extend(f for f in files if id(f.attachment) not in present)
return prompt
async def _tool_output(self, result: dict, media: Optional[_Media]):
"""Tool results as text, plus any files the tool produced (e.g. generated images)."""
from openlayer.lib.tracing.content import TextContent
text = _text(result.get("output"))
if media is None:
return text
# Built-in tools (e.g. generate_image) reference their files inside the result text
# as /api/v1/files/<id>/content; some tools also return a separate files list.
urls = [e.get("url", "") if isinstance(e, dict) else str(e) for e in result.get("files") or []]
file_ids = list(dict.fromkeys(_FILE_ID.findall(" ".join(urls) + " " + text)))
items = []
for file_id in file_ids:
item = await media.from_file_id(file_id, "generated")
if item:
items.append(item)
for url in urls:
if url.startswith("data:"):
item = media.from_data_url(url, "image")
if item:
items.append(item)
return [TextContent(text=text), *items] if items else text
async def _trace(self, body: dict, user: dict, metadata: dict, model: dict) -> None:
from openlayer.lib import update_current_trace, update_trace_user_session
from openlayer.lib.tracing import tracer
from openlayer.lib.tracing.enums import StepType
messages: List[dict] = body.get("messages") or []
if not messages or messages[-1].get("role") != "assistant":
return
assistant = messages[-1]
end = time.time()
start = metadata.get(_START_KEY) or end
media = (
_Media(self.valves.max_attachment_mb * 1024 * 1024, user.get("id"))
if self.valves.upload_attachments
else None
)
prompt = await self._prompt(metadata.get(_INLET_MESSAGES_KEY) or messages[:-1], metadata, media)
question = next((_text_of(p["content"]) for p in reversed(prompt) if p["role"] == "user"), "")
answer = _text(assistant.get("content"))
chat_id = body.get("chat_id") or metadata.get("chat_id") or ""
message_id = body.get("id") or metadata.get("message_id")
model_id = body.get("model") or model.get("id") or "unknown"
# Custom Workspace models ("agents") wrap a base model; price the call by the base model.
base_model_id = (model.get("info") or {}).get("base_model_id") or model_id
provider = model.get("owned_by") or "unknown"
run = await self._run_context(metadata, chat_id)
routes = self._routes()
destination = (
(routes.get(run["run_type"]) if run["run_type"] in ("subagent", "automation") else None)
or routes.get(model_id)
or routes.get(run["run_type"])
or self.valves.inference_pipeline_id
)
usage = _usage(assistant.get("usage") or {})
output_items = assistant.get("output") or []
sources = assistant.get("sources") or metadata.get("sources") or []
reasoning = [
_text(item.get("content") or item.get("summary"))
for item in output_items
if item.get("type") == "reasoning"
]
results = {
item.get("call_id"): item for item in output_items if item.get("type") == "function_call_output"
}
tool_calls = []
if self.valves.capture_tool_calls:
for item in output_items:
if item.get("type") != "function_call":
continue
try:
arguments = json.loads(item.get("arguments") or "{}")
except (TypeError, ValueError):
arguments = {"arguments": item.get("arguments")}
result = results.get(item.get("call_id"), {})
tool_calls.append((item, arguments, result, await self._tool_output(result, media)))
root_metadata = {
"chat_id": chat_id,
"message_id": message_id,
"model": model_id,
"model_name": model.get("name"),
"features": {k: v for k, v in (metadata.get("features") or {}).items() if v},
"tool_ids": metadata.get("tool_ids") or [],
"files": [
(f.get("file") or f).get("filename") or (f.get("file") or f).get("name")
for f in metadata.get("files") or []
if isinstance(f, dict)
],
**{k: v for k, v in run.items() if v},
"user_name": user.get("name"),
"user_role": user.get("role"),
}
if media is not None and media.skipped:
root_metadata["attachments_skipped"] = media.skipped
with tracer.create_step(
name="Open WebUI chat turn",
step_type=StepType.USER_CALL,
inputs={"prompt": prompt},
metadata=root_metadata,
# Always explicit: the SDK default is process-wide and shared with other functions.
inference_pipeline_id=destination,
) as root:
_stamp(root, start, end)
update_trace_user_session(
user_id=user.get("email") or user.get("id") or "anonymous",
# Sub-agent runs join the session of the chat that delegated them.
session_id=run.get("parent_chat_id") or chat_id or None,
)
if message_id:
# Lets the feedback Event function update this exact row later.
update_current_trace(inferenceId=message_id)
if self.valves.capture_sources and sources:
with tracer.create_step(
name="Knowledge retrieval",
step_type=StepType.RETRIEVER,
inputs={"query": question},
) as step:
_stamp(step, start)
step.log(output=[self._source(s) for s in sources])
with tracer.create_step(
name="LLM chat completion",
step_type=StepType.CHAT_COMPLETION,
inputs={"prompt": prompt},
metadata={"reasoning": "\n\n".join(reasoning)} if reasoning else None,
) as llm:
_stamp(llm, start + 0.001, end)
llm.model = base_model_id
llm.provider = provider
llm.prompt_tokens = usage["prompt"]
llm.completion_tokens = usage["completion"]
llm.tokens = usage["total"]
llm.log(output=answer)
for index, (item, arguments, result, output) in enumerate(tool_calls):
with tracer.create_step(
name=item.get("name") or "tool",
step_type=StepType.TOOL,
inputs=arguments,
metadata={"call_id": item.get("call_id"), "status": result.get("status")},
) as step:
# Outlet has no per-call timings; keep tool steps in call order.
_stamp(step, start + 0.002 + 0.001 * index)
step.log(output=output)
root.log(output=answer)
if self.valves.debug:
count = len(media.by_checksum) if media is not None else 0
log.info(
"Openlayer: traced %s message %s to %s (%s tokens, %d attachments)",
run["run_type"], message_id, destination, usage["total"], count,
)
@staticmethod
def _source(source: dict) -> dict:
info = source.get("source") or {}
return {
"name": info.get("name") or info.get("id"),
"type": info.get("type"),
"documents": [str(d)[:500] for d in (source.get("document") or [])][:5],
}
def _text_of(content: Any) -> str:
"""Text of a prompt message whose content may hold Openlayer content items."""
if isinstance(content, list):
return " ".join(getattr(item, "text", "") for item in content if getattr(item, "text", None))
return _text(content)