Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
0cbf22c
Add AI usage tracking via HolmesUsageEvents (account/cluster/user/fea…
alonelish Apr 29, 2026
c934138
Surface request_id in /api/chat response metadata for feedback
alonelish Apr 30, 2026
6c80253
Add is_internal flag to HolmesUsageEvents to filter internal calls
alonelish Apr 30, 2026
4567656
Address CodeRabbit feedback on AI usage tracking PR
alonelish Apr 30, 2026
06ef21b
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish Apr 30, 2026
81b7bb2
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 3, 2026
0f8bd38
Address second-round CodeRabbit feedback
alonelish May 3, 2026
f14450b
Auto-detect Slack-driven /api/chat calls and tag as request_type='sla…
alonelish May 4, 2026
6c7474c
Also auto-set request_source='slack' when Slack prefix is detected
alonelish May 4, 2026
b7969f9
Make user_id optional on POST /api/feedback
alonelish May 4, 2026
4826e65
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 5, 2026
df209d5
Wire usage recorder into ConversationWorker so worker chats are tracked
alonelish May 5, 2026
ec2236e
Fall back to Conversations row for user_id / request_source on worker…
alonelish May 5, 2026
63ab3e9
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 5, 2026
e776bd0
Use UTC for feedback_at and let worker request_type auto-detect
alonelish May 5, 2026
0101f1f
Dedupe _resolve_provider via shared usage_recorder.resolve_provider
alonelish May 5, 2026
0f188f0
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 6, 2026
321f973
Drop redundant 'Grafana' prefix from Loki / Tempo catalog names
alonelish May 6, 2026
a17fc1b
Drop /api/feedback endpoint — FE will call record_feedback() RPC dire…
alonelish May 6, 2026
ab4958b
Document each UsageRecorderState field inline
alonelish May 6, 2026
f942de1
Replace string-literal status values with RequestStatus enum
alonelish May 6, 2026
32fc36d
record_usage_event takes UsageRecorderState directly; drop to_kwargs()
alonelish May 6, 2026
cbdfc54
record_from_llm_result: filter stats via Pydantic, not a hardcoded set
alonelish May 6, 2026
e22fa33
Move _fire and _capture_terminal into UsageRecorderState
alonelish May 6, 2026
be0ac0c
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 7, 2026
ede8721
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 7, 2026
1921789
Capture partial token costs from TOKEN_COUNT events
alonelish May 7, 2026
91b5638
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 7, 2026
9b94de2
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 7, 2026
40d040a
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 10, 2026
2ef9cab
Merge remote-tracking branch 'upstream/master' into claude/confident-…
alonelish May 10, 2026
e629689
Replace per-request Thread with bounded ThreadPoolExecutor for recorder
alonelish May 10, 2026
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
4 changes: 2 additions & 2 deletions datasource-catalog.json
Original file line number Diff line number Diff line change
Expand Up @@ -222,7 +222,7 @@
},
{
"id": "grafana/loki",
"name": "Grafana Loki",
"name": "Loki",
"description": "Search and query Grafana Loki logs",
"categories": ["logs"],
"iconUrl": "https://raw.githubusercontent.com/gilbarbara/logos/de2c1f96ff6e74ea7ea979b43202e8d4b863c655/logos/grafana.svg",
Expand Down Expand Up @@ -294,7 +294,7 @@
},
{
"id": "grafana/tempo",
"name": "Grafana Tempo",
"name": "Tempo",
"description": "Search and query Grafana Tempo traces",
"categories": ["traces"],
"iconUrl": "https://raw.githubusercontent.com/gilbarbara/logos/de2c1f96ff6e74ea7ea979b43202e8d4b863c655/logos/grafana.svg",
Expand Down
49 changes: 46 additions & 3 deletions experimental/ag-ui/server-agui.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,11 @@
from starlette.responses import PlainTextResponse

from holmes.utils.stream import StreamMessage, StreamEvents
from holmes.core.usage_recorder import (
UsageRecorderState,
resolve_provider,
stream_with_usage_recording,
)
from holmes.common.env_vars import (
HOLMES_HOST,
HOLMES_PORT,
Expand Down Expand Up @@ -139,9 +144,47 @@ async def event_generator(message_history):
run_id=input_data.run_id,
)
)
hgpt_chat_stream_response: StreamMessage = ai.call_stream(
msgs=message_history,
enable_tool_approval=chat_request.enable_tool_approval or False,
# Build the AI usage recorder state. AG-UI clients pass FE flags
# via input_data.context if they want to populate them; otherwise
# request_source / source_ref / user_id stay NULL.
ctx = input_data.context or {}
agui_user_id = None
try:
agui_user_id = getattr(input_data, "user_id", None) or (
ctx.get("user_id") if isinstance(ctx, dict) else None
)
except (AttributeError, TypeError) as e:
# Telemetry-only fallback: if input_data shape changed and
# neither attribute nor mapping access works, log at debug
# so we notice during dev rather than silently dropping
# user attribution. Don't break the response path.
logging.debug(
"agui_chat: failed to extract user_id (input_data=%r ctx=%r): %s",
type(input_data).__name__,
type(ctx).__name__,
e,
)
ai_model = getattr(ai.llm, "model", None) or chat_request.model or "unknown"
ai_provider = resolve_provider(ai_model)
recorder_state = UsageRecorderState(
dal=dal,
request_type="agui_chat",
request_source=ctx.get("request_source") if isinstance(ctx, dict) else None,
source_ref=ctx.get("source_ref") if isinstance(ctx, dict) else None,
conversation_id=getattr(input_data, "thread_id", None),
conversation_source=None, # AG-UI doesn't write Conversations or ChatHistory
user_id=agui_user_id,
is_streaming=True,
model=ai_model,
provider=ai_provider,
is_robusta_model=getattr(ai.llm, "is_robusta_model", False),
)
hgpt_chat_stream_response: StreamMessage = stream_with_usage_recording(
ai.call_stream(
msgs=message_history,
enable_tool_approval=chat_request.enable_tool_approval or False,
),
recorder_state,
)
for chunk in hgpt_chat_stream_response:
if hasattr(chunk, "event"):
Expand Down
16 changes: 16 additions & 0 deletions holmes/checks/checks.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,11 @@
from holmes.config import Config
from holmes.core.issue import Issue, IssueStatus
from holmes.core.tool_calling_llm import LLMResult, ToolCallingLLM
from holmes.core.usage_recorder import (
UsageRecorderState,
record_error,
record_from_llm_result,
)
from holmes.plugins.destinations.pagerduty.plugin import PagerDutyDestination
from holmes.plugins.destinations.slack.plugin import SlackDestination

Expand Down Expand Up @@ -89,6 +94,7 @@ def execute_check(
ai: ToolCallingLLM,
verbose: bool = False,
console: Optional[Console] = None,
recorder_state: Optional[UsageRecorderState] = None,
) -> CheckResult:
"""
Execute a single health check.
Expand All @@ -101,6 +107,10 @@ def execute_check(
ai: The LLM instance to use for evaluation
verbose: Whether to print verbose output
console: Optional console for output (only used if verbose=True)
recorder_state: Optional UsageRecorderState. When supplied (e.g. by the
/api/checks/execute endpoint) a usage event is recorded for this
LLM call. The CLI runner doesn't pass one and is therefore not
tracked, by design.

Returns:
CheckResult with status, message, and metadata
Expand All @@ -114,6 +124,10 @@ def execute_check(
start_time = time.time()
try:
response = _execute_ai_check(check, ai)
if recorder_state is not None:
# Fire the usage recorder. response IS-A LLMResult (RequestStats
# subclass) so cost/token fields come straight off it.
record_from_llm_result(recorder_state, response)
check_response = _parse_check_response(response)
if verbose and console:
status_str = "PASS" if check_response.passed else "FAIL"
Expand Down Expand Up @@ -142,6 +156,8 @@ def execute_check(
)

except Exception as e:
if recorder_state is not None:
record_error(recorder_state, e)
duration = time.time() - start_time
result = CheckResult(
check_name=check.name,
Expand Down
21 changes: 21 additions & 0 deletions holmes/checks/checks_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from holmes.core.issue import Issue, IssueStatus
from holmes.core.tool_calling_llm import LLMResult, ToolCallingLLM
from holmes.core.tools import PrerequisiteCacheMode, ToolsetTag
from holmes.core.usage_recorder import UsageRecorderState, resolve_provider
from holmes.plugins.destinations.slack.plugin import SlackDestination

checks_app = FastAPI()
Expand Down Expand Up @@ -110,12 +111,32 @@ def execute_health_check(
destinations=destination_names,
)

# Build the recorder state so the operator-driven check shows up in
# HolmesUsageEvents with request_type='health_check' and source_ref
# set to the check name (per-check cost reporting key).
ai_model = getattr(ai.llm, "model", None) or request.model or "unknown"
recorder_state = UsageRecorderState(
dal=_CONFIG.dal,
request_type="health_check",
request_source="operator",
source_ref=request.name or "api-check",
conversation_id=None,
conversation_source=None,
user_id=None,
is_streaming=False,
model=ai_model,
provider=resolve_provider(ai_model),
is_robusta_model=getattr(ai.llm, "is_robusta_model", False),
meta={"check_mode": request.mode.value, "timeout": request.timeout},
)

# Execute the check using the shared function
result: CheckResult = execute_check(
check=check,
ai=ai,
verbose=False,
console=None,
recorder_state=recorder_state,
)

# Track notification statuses
Expand Down
6 changes: 6 additions & 0 deletions holmes/core/conversations_worker/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,12 @@ class ConversationTask(BaseModel):
request_sequence: int
metadata: Dict[str, Any] = Field(default_factory=dict)
title: Optional[str] = None
# The Conversations row's user_id column (the human who started the
# chat). Used as a fallback for HolmesUsageEvents.user_id when the FE
# didn't include user_id in the user_message event's data — common
# because the runner-side Conversations row already has the value, so
# the FE has no reason to duplicate it into every per-turn event.
user_id: Optional[str] = None

# Hydrated post-construction from events; not part of the validated row schema.
_user_message_data: Dict[str, Any] = PrivateAttr(default_factory=dict)
Expand Down
60 changes: 58 additions & 2 deletions holmes/core/conversations_worker/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,10 @@
inject_frontend_tools,
)
from holmes.core.tracing import TracingFactory
from holmes.core.usage_recorder import (
build_chat_recorder_state,
stream_with_usage_recording,
)
from holmes.utils.holmes_status import update_holmes_status_in_db
from holmes.utils.stream import StreamEvents

Expand Down Expand Up @@ -483,6 +487,11 @@ def _build_task_from_conversation_row(
request_sequence=int(conv.get("request_sequence", 1)),
metadata=conv.get("metadata") or {},
title=conv.get("title"),
# Conversations.user_id (set by the FE when it created the row)
# — surfaced on the task so per-turn ChatRequest construction
# can use it as a fallback when the user_message event's data
# doesn't carry user_id explicitly.
user_id=conv.get("user_id"),
)
except Exception:
logging.exception(
Expand Down Expand Up @@ -617,6 +626,20 @@ def _process_conversation(self, task: ConversationTask) -> None:
if data.get("tool_decisions"):
enable_tool_approval = True

# AI usage tracking (HolmesUsageEvents) — resolve user_id and
# request_source with row-level fallbacks. The FE writes both onto
# the Conversations row when it creates the chat (user_id as a
# column, request_source under metadata) but doesn't necessarily
# repeat them in every user_message event's data. Without this
# fallback, follow-up turns produce HolmesUsageEvents rows with
# NULL user_id / request_source even though the values are known.
# Per-event data still wins so the FE can override per-turn (e.g.
# an alert-investigation chat that pivots to a freeform question).
resolved_user_id = data.get("user_id") or task.user_id
resolved_request_source = data.get("request_source") or (
task.metadata.get("request_source") if task.metadata else None
)

chat_request = ChatRequest(
ask=ask,
images=data.get("images"),
Expand All @@ -630,7 +653,27 @@ def _process_conversation(self, task: ConversationTask) -> None:
frontend_tool_results=data.get("frontend_tool_results"), # type: ignore[arg-type]
response_format=data.get("response_format"),
behavior_controls=data.get("behavior_controls"),
user_id=data.get("user_id"),
# source_ref / meta / is_internal still come from the per-event
# blob only — they're per-turn signals (which alert this
# follow-up question was about, etc.), not Conversation-level
# state. user_id / request_source fall back to the Conversations
# row when the FE didn't repeat them in the event.
user_id=resolved_user_id,
# request_type: pass through whatever the FE sent (None if absent)
# rather than hard-coding 'user_chat' here. The recorder helper
# (build_chat_recorder_state) handles the default and runs Slack
# auto-detection — hard-coding 'user_chat' would defeat the
# auto-detection because the helper bails out if request_type is
# already truthy. Today only /api/chat hits the Slack-prefix
# path, but the runner could route Slack through Conversations
# at any time without a code change here.
request_type=data.get("request_type"),
request_source=resolved_request_source,
source_ref=data.get("source_ref"),
conversation_id=task.conversation_id,
conversation_source="conversations",
meta=data.get("meta"),
is_internal=data.get("is_internal"),
)

self._run_chat_and_publish(
Expand Down Expand Up @@ -787,7 +830,19 @@ def _run_chat_and_publish(
request_context = {"user_id": chat_request.user_id}

try:
stream = request_ai.call_stream(
# Wrap the raw stream with the usage recorder BEFORE the
# publisher consumes it, so the recorder sees Holmes' native
# StreamMessage events (TOOL_RESULT / ANSWER_END / etc.) and
# can fire one HolmesUsageEvents row per worker-driven turn.
# Mirrors the wiring in server.py::chat() for the streaming
# path; without this the worker bypasses the recorder entirely.
recorder_state = build_chat_recorder_state(
chat_request,
request_ai,
dal=self.dal,
is_streaming=True,
)
raw_stream = request_ai.call_stream(
msgs=messages,
enable_tool_approval=chat_request.enable_tool_approval or False,
tool_decisions=chat_request.tool_decisions,
Expand All @@ -796,6 +851,7 @@ def _run_chat_and_publish(
request_context=request_context,
trace_span=trace_span,
)
stream = stream_with_usage_recording(raw_stream, recorder_state)

terminal = publisher.consume(stream)
if terminal is None:
Expand Down
63 changes: 63 additions & 0 deletions holmes/core/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,69 @@ class ChatRequestBaseModel(BaseModel):
)
user_id: Optional[str] = None # User ID from relay session token validation

# ── AI usage tracking fields (HolmesUsageEvents). All optional / additive;
# old clients that don't supply them keep working unchanged. ──
request_type: Optional[str] = Field(
default=None,
description=(
"Backend-set classification: 'user_chat' (default for /api/chat), "
"'scheduled_prompt' (set by ScheduledPromptsExecutor), 'agui_chat' "
"(set by AG-UI handler), 'health_check' (set by /api/checks/execute)."
),
)
request_source: Optional[str] = Field(
default=None,
description=(
"FE-supplied UI flow label, free-form. Examples: 'freeform', "
"'followup_logs', 'alert_investigation', 'resource_chat'."
),
)
source_ref: Optional[str] = Field(
default=None,
description=(
"FE-supplied opaque pointer to the entity the chat is about "
"(e.g. an issue id when request_source='alert_investigation'). "
"Meaning is implied by request_source."
),
)
conversation_id: Optional[str] = Field(
default=None,
description=(
"Stable id grouping multi-turn chats. Soft reference (NOT a FK): "
"matches Conversations.conversation_id when worker handles the chat, "
"or the FE-owned ChatHistory id for direct /api/chat traffic. NULL for "
"single-turn / non-UI flows."
),
)
conversation_source: Optional[str] = Field(
default=None,
description=(
"Discriminator telling dashboards which table conversation_id targets: "
"'conversations' (worker path) or 'chat_history' (direct /api/chat). "
"Worker sets it explicitly; chat() defaults to 'chat_history' when "
"conversation_id is non-NULL and not already set."
),
)
meta: Optional[Dict[str, Any]] = Field(
default=None,
description=(
"Forward-compatibility metadata bag. FE-supplied opaque dict; the "
"server shallow-merges with backend-derived keys (backend wins on "
"collision). Keep small; promote stable keys to real columns over time. "
"Do NOT put PII / large strings (prompts, completions, tool outputs) here."
),
)
is_internal: Optional[bool] = Field(
default=None,
description=(
"Marks server-internal calls (title generation, classification, "
"summarization, etc.) so dashboards can filter them out of user-facing "
"metrics. FE sets True for those. When unset, the server defaults it "
"to True if request_source starts with 'internal_' (backwards compat "
"with the prefix convention) — otherwise False."
),
)

# In our setup with litellm, the first message in conversation_history
# should follow the structure [{"role": "system", "content": ...}],
# where the "role" field is expected to be "system".
Expand Down
4 changes: 4 additions & 0 deletions holmes/core/scheduled_prompts/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,10 @@ def _execute_prompt(
additional_system_prompt=additional_system_prompt,
trace_span=heartbeat_span,
behavior_controls=behavior_controls,
# AI usage tracking — these runs are server-driven, not user-driven.
request_type="scheduled_prompt",
request_source="scheduler",
source_ref=sp.id,
)

empty_request = Request(scope={"type": "http", "headers": []})
Expand Down
Loading
Loading