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
36 changes: 36 additions & 0 deletions ROADMAP.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
# Observability Roadmap

## Goal

Build an observability stack for agent execution that uses:

- pino + Loki for structured logs
- OpenTelemetry for spans and traces
- ClickHouse as the backend for trace/span storage and query
- Grafana for dashboards and correlation views

## Current Status

- Structured runtime span metadata is now emitted from the agent runtime.
- Runtime tests cover span metadata and tool span duration.
- The next step is to wire actual OpenTelemetry span emission and connect it to the tracing backend.

## Planned Phases

### Phase 1 — Runtime instrumentation
- Keep structured logs in pino/Loki.
- Emit root and child spans from the agent runtime.
- Preserve trace/span identifiers and duration metadata.

### Phase 2 — OpenTelemetry integration
- Add an OpenTelemetry tracer provider.
- Create spans for reasoning, tool calls, tool results, final output, and root execution.
- Propagate trace context across agent execution steps.

### Phase 3 — Backend wiring
- Configure an OpenTelemetry collector/exporter path for ClickHouse.
- Ensure trace data is ingested and queryable from Grafana.

### Phase 4 — Grafana experience
- Add dashboards for tool latency, error spans, and per-run execution views.
- Correlate logs and traces through shared trace identifiers.
9 changes: 9 additions & 0 deletions TODO.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# TODO

- [x] Add structured span metadata in the agent runtime
- [x] Add regression tests for runtime span metadata and tool span duration
- [x] Add Grafana/Loki query examples for span-based observability
- [x] Add OpenTelemetry span emission for agent runtime events
- [ ] Propagate trace context across agent runtime steps
- [ ] Configure OpenTelemetry exporter path for ClickHouse
- [ ] Add Grafana dashboards that correlate logs and traces
118 changes: 64 additions & 54 deletions apps/cli/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import {
import { logger } from '@sup/infra/logger';
import type { AgentStep } from '@sup/types/agent-types';
import { ProtocolResolver } from '@sup/lib/protocol-resolver';
import { AgentRuntime } from '@sup/lib';
import { adapters } from '@sup/tools';
import { JsonFileProvider } from '@sup/infra/adapters/JsonFileProvider';

Expand Down Expand Up @@ -83,17 +84,36 @@ export async function startChat() {
if (userInput.trim().toLowerCase() === 'exit') break;

try {
const runtime = new AgentRuntime({
name: AGENT,
sessionId: session.id,
input: userInput,
});

let generator = agent(userInput, session, {
resolver: ProtocolResolver,
tools: adapters,
});

for (const adapterFn of activeAdapters) {
generator = adapterFn(generator);
}
for (const adapterFn of activeAdapters) {
generator = adapterFn(generator);
}

await renderStream(generator);
} catch (error) {
const renderState = { accumulated: '', firstToken: true };
await runtime.run(() => generator, {
onStep: (step) => {
if (process.env.LOG_STEPS === 'true') {
console.log(step);
}
renderStep(step, renderState);
},
onSpan: (span) => {
if (process.env.LOG_STEPS === 'true') {
console.log('[span]', span);
}
},
});
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
logger.error('Error in agent execution:', message);
console.error('Full error:', error);
Expand All @@ -107,58 +127,48 @@ export async function startChat() {
}
}

async function renderStream(
generator: AsyncGenerator<AgentStep, void, unknown>
): Promise<void> {
let accumulated = '';
let firstToken = true;

for await (const step of generator) {
if (process.env.LOG_STEPS === "true") {
console.log(step);
}
switch (step.type) {
case 'thinking':
if (STREAMING) process.stdout.write(`\n${step.message}\n`);
break;

case 'text_delta':
if (STREAMING) {
if (firstToken) {
process.stdout.write('\nAgent: ');
firstToken = false;
}
if (step.delta) process.stdout.write(step.delta);
} else {
if (step.delta) accumulated += step.delta;
function renderStep(step: AgentStep, state: { accumulated: string; firstToken: boolean }) {
switch (step.type) {
case 'thinking':
if (STREAMING) process.stdout.write(`\n${step.message}\n`);
break;

case 'text_delta':
if (STREAMING) {
if (state.firstToken) {
process.stdout.write('\nAgent: ');
state.firstToken = false;
}
break;
if (step.delta) process.stdout.write(step.delta);
} else {
if (step.delta) state.accumulated += step.delta;
}
break;

case 'tool_call':
if (STREAMING) {
if (!firstToken) process.stdout.write('\n');
process.stdout.write(`[${step.toolId}] `);
firstToken = true;
}
break;

case 'tool_result':
if (STREAMING) process.stdout.write(`✓\n`);
break;

case 'final':
if (STREAMING) {
if (firstToken) {
process.stdout.write(`\nAgent: ${step.text}`);
}
process.stdout.write('\n\n');
} else {
process.stdout.write(`\nAgent: ${accumulated || step.text}\n\n`);
case 'tool_call':
if (STREAMING) {
if (!state.firstToken) process.stdout.write('\n');
process.stdout.write(`[${step.toolId}] `);
state.firstToken = true;
}
break;

case 'tool_result':
if (STREAMING) process.stdout.write(`✓\n`);
break;

case 'final':
if (STREAMING) {
if (state.firstToken) {
process.stdout.write(`\nAgent: ${step.text}`);
}
accumulated = '';
firstToken = true;
break;
}
process.stdout.write('\n\n');
} else {
process.stdout.write(`\nAgent: ${state.accumulated || step.text}\n\n`);
}
state.accumulated = '';
state.firstToken = true;
break;
}
}

Expand Down
48 changes: 32 additions & 16 deletions apps/tui/index.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
} from "@openfeature/server-sdk";

import { ProtocolResolver } from "@sup/lib/protocol-resolver";
import { AgentRuntime } from "@sup/lib";
import { adapters } from "@sup/tools";
import { JsonFileProvider } from "@sup/infra/adapters/JsonFileProvider";

Expand Down Expand Up @@ -78,27 +79,42 @@ function App() {
setMessages((m) => [...m, { role: "user", text }]);

try {
const runtime = new AgentRuntime({
name: AGENT,
sessionId: session.id,
input: text,
});

const generator = currentAgent(text, session, {
resolver: ProtocolResolver,
tools: adapters,
});

for await (const step of generator) {
if (step.type === "text_delta" && step.delta) {
setStreaming((s) => s + step.delta);
} else if (step.type === "tool_call") {
setActiveTool(step.toolId);
} else if (step.type === "tool_result") {
setActiveTool("");
} else if (step.type === "final") {
setMessages((m) => [
...m,
{ role: "agent", text: step.text || streaming() },
]);
setStreaming("");
setActiveTool("");
}
}
await runtime.run(() => generator, {
onStep: (step) => {
if (step.type === "text_delta" && step.delta) {
setStreaming((s) => s + step.delta);
} else if (step.type === "tool_call") {
setActiveTool(step.toolId);
} else if (step.type === "tool_result") {
setActiveTool("");
} else if (step.type === "final") {
setMessages((m) => [
...m,
{ role: "agent", text: step.text || streaming() },
]);
setStreaming("");
setActiveTool("");
}
},
onSpan: (span) => {
if (span.name === "reasoning") {
setActiveTool(`reasoning: ${span.message}`);
} else if (span.name === "tool") {
setActiveTool(span.message);
}
},
});
} catch (err) {
setMessages((m) => [...m, { role: "agent", text: `Error: ${err}` }]);
}
Expand Down
42 changes: 29 additions & 13 deletions apps/tui/tui.tsx
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { render } from "@opentui/solid";
import { createSignal, createEffect, For } from "solid-js";
import { useKeyboard } from "@opentui/solid";
import { AgentRuntime } from "@sup/lib";

type Message = { role: "user" | "agent"; text: string };

Expand All @@ -16,21 +17,36 @@ function App({ agent, session, resolver, adapters }) {
setMessages((m) => [...m, { role: "user", text }]);
setInput("");

const runtime = new AgentRuntime({
name: "tui-agent",
sessionId: session.id,
input: text,
});

const generator = agent(text, session, { resolver, tools: adapters });

for await (const step of generator) {
if (step.type === "text_delta" && step.delta) {
setStreaming((s) => s + step.delta);
} else if (step.type === "tool_call") {
setActiveTool(step.toolId);
} else if (step.type === "tool_result") {
setActiveTool("");
} else if (step.type === "final") {
setMessages((m) => [...m, { role: "agent", text: step.text || streaming() }]);
setStreaming("");
setActiveTool("");
}
}
await runtime.run(() => generator, {
onStep: (step) => {
if (step.type === "text_delta" && step.delta) {
setStreaming((s) => s + step.delta);
} else if (step.type === "tool_call") {
setActiveTool(step.toolId);
} else if (step.type === "tool_result") {
setActiveTool("");
} else if (step.type === "final") {
setMessages((m) => [...m, { role: "agent", text: step.text || streaming() }]);
setStreaming("");
setActiveTool("");
}
},
onSpan: (span) => {
if (span.name === "reasoning") {
setActiveTool(`reasoning: ${span.message}`);
} else if (span.name === "tool") {
setActiveTool(span.message);
}
},
});
}

return (
Expand Down
78 changes: 78 additions & 0 deletions config/grafana/agent-runtime-tracing-queries.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
# Grafana / Loki Queries for Agent Runtime Spans

These queries assume the runtime is emitting structured JSON logs through the existing Loki transport and that your Loki datasource is using the `app="sup"` label.

## 1. Tool latency percentiles

### p50

```logql
{app="sup"}
| json
| eventType="span_end"
| name="tool"
| label_format toolId="{{if .toolId}}{{.toolId}}{{else}}unknown{{end}}"
| unwrap durationMs [5m]
```

### p95

```logql
quantile_over_time(0.95,
{app="sup"}
| json
| eventType="span_end"
| name="tool"
| label_format toolId="{{if .toolId}}{{.toolId}}{{else}}unknown{{end}}"
| unwrap durationMs [5m]
) by (toolId)
```

### p99

```logql
quantile_over_time(0.99,
{app="sup"}
| json
| eventType="span_end"
| name="tool"
| label_format toolId="{{if .toolId}}{{.toolId}}{{else}}unknown{{end}}"
| unwrap durationMs [5m]
) by (toolId)
```

## 2. Trace waterfall for a single run

Use a Grafana variable named `trace_id` and set it to the runtime `traceId` value.

```logql
{app="sup"}
| json
| traceId="$trace_id"
| sort_by_timestamp asc
| line_format "{{.name}} [{{.durationMs}}ms] -> {{.message}}"
```

## 3. Failed runs and error spans

```logql
{app="sup"}
| json
| eventType="span_end"
| status="error"
| line_format "{{.name}} failed in {{.durationMs}}ms :: {{.message}}"
```

## 4. Grafana derived field for trace drill-down

In Grafana, add a derived field for the Loki datasource:

- Name: `traceId`
- Regex: `"traceId":"([^"]+)"`
- URL/Link: use the Explore page or your trace viewer with a query like:

```text
{app="sup"} | json | traceId="${__value.raw}"
```

This makes each span log entry clickable and lets you jump directly into the execution trace for that run.
Loading
Loading