Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
5858229
fix(langfuse): attribute generations to the wire model, not the agent…
erosika Jul 29, 2026
4526f58
fix(langfuse): send explicit cost total alongside per-type breakdown
erosika Jul 29, 2026
be5a0ea
test(langfuse): cover the explicit cost total on both cost paths
erosika Aug 10, 2026
d935b82
feat(langfuse): capture modes, api error coverage, session-finalize f…
erosika Aug 10, 2026
462e7a9
feat(langfuse): trace delegated subagents as spans under the parent turn
erosika Aug 10, 2026
89bfd56
feat(langfuse): record each MoA advisor as its own generation
erosika Aug 10, 2026
f09abcf
fix(langfuse): shutdown client on session finalize to avoid interpret…
bgodlin Aug 7, 2026
124a3bb
fix(langfuse): close root-observation CMs to prevent interpreter-tear…
bgodlin Aug 10, 2026
bf7c75b
fix(langfuse): finalize open root spans at exit so short-lived proces…
aldoeliacim Aug 8, 2026
c97dd64
fix(langfuse): scope client shutdown to process exit, not session rot…
erosika Aug 10, 2026
9620bdf
fix(langfuse): unwind root contexts and subagent spans in atexit fina…
erosika Aug 10, 2026
217071f
fix(langfuse): guard _get_langfuse() against concurrent-init TOCTOU
erosika Aug 10, 2026
4e767d3
fix(langfuse): surface reasoning_content in traces
erosika Aug 10, 2026
09e7745
fix(langfuse): include system prompt in generation input
erosika Aug 10, 2026
631f8f3
fix(langfuse): use update_trace so turn Input/Output columns fill
erosika Aug 10, 2026
63144ee
fix(langfuse): export canonical generation total
erosika Aug 10, 2026
756c78a
fix(langfuse): omit cost_details for subscription-included providers
erosika Aug 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
45 changes: 45 additions & 0 deletions agent/conversation_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,24 @@ def _apply_active_turn_redirect(agent: Any, messages: List[Dict[str, Any]], text
agent._stream_needs_break = True


def _moa_reference_metrics_for_hook(agent: Any) -> Any:
"""Per-advisor metrics for post_api_request, or None off the MoA path.

MoA runs N advisor models before its aggregator and returns only the
aggregator's response, so an observability plugin sees one generation for
the whole fan-out. The advisor spend is already computed per slot (see
``_RefAccounting``); this only carries it across the hook boundary.
"""
client = getattr(agent, "client", None)
getter = getattr(client, "last_reference_metrics", None)
if not callable(getter):
return None
try:
return getter()
except Exception:
return None


def _is_copilot_provider(agent: Any) -> bool:
"""Delegate to ``AIAgent._is_copilot_provider`` (single owner of the check).

Expand Down Expand Up @@ -356,6 +374,24 @@ def _print_nous_entitlement_guidance(agent, capability: str) -> bool:
return True


def _system_prompt_for_hooks(api_kwargs: Any, request_messages: Any) -> Any:
"""System prompt as actually sent to the provider, for observability hooks.

Providers move it out of ``messages``: Anthropic Messages uses a separate
``system`` kwarg (str or content-block list), the Responses/Codex API uses
top-level ``instructions``; Chat Completions keeps it as ``messages[0]``.
Returns None when the request carries no system prompt.
"""
system_prompt = api_kwargs.get("system")
if system_prompt is None:
system_prompt = api_kwargs.get("instructions")
if system_prompt is None and isinstance(request_messages, list) and request_messages:
first = request_messages[0]
if isinstance(first, dict) and first.get("role") == "system":
system_prompt = first.get("content")
return system_prompt


def _is_nous_inference_route(provider: str, base_url: str) -> bool:
provider = (provider or "").strip().lower()
if provider == "nous":
Expand Down Expand Up @@ -2270,6 +2306,13 @@ def run_conversation(
request_messages = api_kwargs.get("input")
if not isinstance(request_messages, list):
request_messages = api_messages
# Anthropic (``system``) and Responses/Codex
# (``instructions``) move the system prompt out of
# messages; pass it explicitly for observability
# plugins (Langfuse).
system_prompt_for_hooks = _system_prompt_for_hooks(
api_kwargs, request_messages
)
# Shallow-copy the outer list so plugins that retain the
# reference for async snapshotting don't observe later
# mutations of api_messages. The inner dicts are not
Expand Down Expand Up @@ -2305,6 +2348,7 @@ def run_conversation(
request_messages=list(request_messages)
if isinstance(request_messages, list)
else [],
system_prompt=system_prompt_for_hooks,
message_count=len(api_messages),
tool_count=len(agent.tools or []),
approx_input_tokens=approx_tokens,
Expand Down Expand Up @@ -5758,6 +5802,7 @@ def _perform_api_call(next_api_kwargs):
assistant_message=assistant_message,
assistant_content_chars=len(_assistant_text),
assistant_tool_call_count=len(_assistant_tool_calls),
moa_references=_moa_reference_metrics_for_hook(agent),
)
except Exception:
pass
Expand Down
33 changes: 33 additions & 0 deletions agent/moa_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -1562,6 +1562,11 @@ def __init__(self, preset_name: str, reference_callback: Any = None, agent: Any
# caller to stitch in the live session_id + resolved aggregator output
# and flush to the trace file (only when moa.save_traces is on).
self._pending_trace: Any = None
# Per-advisor metrics for observability hooks. Unlike _pending_trace
# this is NOT consumed — post_api_request fires on a different branch
# than consume_and_save_trace, so a consuming read would race it. Holds
# until the next fan-out replaces it.
self._last_reference_metrics: Any = None
# every_n fan-out cadence state. The iteration counter is scoped to a
# single USER TURN (not the facade lifetime): it counts create() calls
# since the last new user message and resets whenever the user-turn
Expand Down Expand Up @@ -1593,6 +1598,14 @@ def consume_reference_usage(self) -> tuple[Any, Any]:
self._pending_reference_cost = None
return usage, cost

def last_reference_metrics(self) -> Any:
"""Per-advisor metrics from the most recent fan-out, or None.

Read-only: a MoA turn's post_api_request hook must not disturb the
accounting that consume_reference_usage and consume_and_save_trace own.
"""
return self._last_reference_metrics

def _record_late_reference_accounting(self, label: str, accounting: Any) -> None:
"""Fold a late-completing interrupted reference's real spend in.

Expand Down Expand Up @@ -2141,6 +2154,18 @@ def _progress(done: int, total: int, label: str) -> None:
"aggregator_slot": aggregator,
"aggregator_temperature": aggregator_temperature,
}
# Derived from the same privacy-redacted _trace_refs, so an active
# privacy mode redacts the observability payload too.
try:
from agent.moa_trace import slot_metrics

self._last_reference_metrics = [
slot_metrics(acct, label, output=text)
for label, text, acct in _trace_refs
]
except Exception as exc: # pragma: no cover - never break a turn
logger.debug("MoA reference metrics render failed: %s", exc)
self._last_reference_metrics = None

# Surface each reference model's answer to the display BEFORE the
# aggregator acts — once per turn (only on the iteration that
Expand Down Expand Up @@ -2285,6 +2310,14 @@ def consume_and_save_trace(
session_id, aggregator_output_fallback=aggregator_output_fallback
)

def last_reference_metrics(self) -> Any:
"""Per-advisor metrics from the most recent fan-out, or None.

Read-only, unlike the two consume_* methods above: the observability
hook fires on a different branch than the accounting they own.
"""
return getattr(self.chat.completions, "_last_reference_metrics", None)


def build_moa_facade(agent, preset_name: Any = None) -> MoAClient:
"""Build the MoA facade client for ``agent``, wiring the reference relay.
Expand Down
15 changes: 15 additions & 0 deletions agent/moa_trace.py
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,21 @@ def _slot_trace(acct: Any, label: str) -> dict[str, Any]:
}


def slot_metrics(acct: Any, label: str, output: Any = None) -> dict[str, Any]:
"""Render one reference's accounting for observability hooks.

Same fields as ``_slot_trace`` minus ``input_messages``, which is the bulk
of a trace record and would cross the plugin-hook boundary for every
advisor on every turn. ``output`` comes from the caller because the
privacy-redacted advisor text lives alongside the accounting, not on it.
"""
trace = _slot_trace(acct, label)
trace.pop("input_messages", None)
if output is not None:
trace["output"] = output
return trace


def save_moa_turn(
*,
session_id: Optional[str],
Expand Down
4 changes: 4 additions & 0 deletions hermes_cli/hooks.py
Original file line number Diff line number Diff line change
Expand Up @@ -211,6 +211,10 @@ def _cmd_list(_args) -> None:
"usage": {"input_tokens": 2048, "output_tokens": 512},
"assistant_content_chars": 1200,
"assistant_tool_call_count": 0,
# Per-advisor metrics on a MoA turn, None otherwise. MoA returns only
# the aggregator's response, so without this an observer cannot see the
# fan-out or price it at each advisor's own model.
"moa_references": None,
},
"subagent_stop": {
"parent_session_id": "parent-sess",
Expand Down
30 changes: 30 additions & 0 deletions plugins/observability/langfuse/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,16 +36,46 @@ hermes plugins list # observability/langfuse should show "enable
hermes chat -q "hello" # then check Langfuse for a "Hermes turn" trace
```

Generation observations include the Hermes system prompt when the provider
uses a separate `system` param (Anthropic Messages API). Open an **LLM call**
child span to inspect `role: system` (truncated via `HERMES_LANGFUSE_MAX_CHARS`).

## Optional tuning

```bash
HERMES_LANGFUSE_ENV=production # environment tag
HERMES_LANGFUSE_RELEASE=v1.0.0 # release tag
HERMES_LANGFUSE_SAMPLE_RATE=0.5 # sample 50% of traces
HERMES_LANGFUSE_MAX_CHARS=12000 # max chars per field (default: 12000)
HERMES_LANGFUSE_CAPTURE=sanitized # content capture mode (see below)
HERMES_LANGFUSE_DEBUG=true # verbose plugin logging
```

## Capture modes

`HERMES_LANGFUSE_CAPTURE` controls how much *content* (prompts, responses,
tool arguments/results) is exported. Structural metadata — IDs, roles, tool
names, token usage, cost, timing — is always captured in every mode.

| mode | behavior |
|------|----------|
| `metadata` | No content. Each content field is replaced by a shape/size stub (`{"omitted": true, "type": "text", "chars": N}`). |
| `sanitized` | **(default)** Content is exported after secret-pattern redaction (API keys, tokens, JWTs, private keys, `password=`-style assignments) and truncation. Redaction runs *before* truncation. |
| `full` | Raw content, truncated only. Explicit opt-in — traces will contain whatever passed through the conversation, including injected memory and file contents. |

The active mode is recorded on every trace as `metadata.capture_mode`.

Note: `sanitized` is pattern-based defense in depth, not a DLP guarantee.
For personal sessions or shared Langfuse projects, prefer `metadata`.

## Error + shutdown coverage

- Failed model requests (`api_request_error` hook) close their generation
with `level=ERROR`, status code, retry counters, and a capture-mode-scrubbed
error message. Non-retryable failures also finish the turn trace.
- Session end/finalize closes any still-open traces for that session and
flushes queued events, so interrupted or tool-only turns don't dangle.

## Disable

```bash
Expand Down
Loading
Loading