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
2 changes: 2 additions & 0 deletions changelogs/fragments/12846.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
feat:
- Add a Sessions tab to Agent Traces that groups traces by `gen_ai.conversation.id` ([#12846](https://github.com/opensearch-project/OpenSearch-Dashboards/pull/12846))
5 changes: 5 additions & 0 deletions src/plugins/agent_traces/common/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,11 @@ export const AGENT_TRACES_DEFAULT_LANGUAGE = 'PPL';
export const AGENT_TRACES_TRACES_TAB_ID = 'traces';
export const AGENT_TRACES_SPANS_TAB_ID = 'spans';
export const AGENT_TRACES_VISUALIZATION_TAB_ID = 'visualization';
export const AGENT_TRACES_SESSIONS_TAB_ID = 'sessions';

/** Span attribute that groups traces into a session (OTel GenAI semantic conventions).
* Data Prepper normalizes OpenInference/ADOT `session.id` to this key at ingest. */
export const AGENT_TRACES_SESSION_ID_FIELD = 'attributes.gen_ai.conversation.id';

export enum AgentTracesFlavor {
Traces = 'traces',
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

import { fetchSessions, sessionFilterFor } from './fetch_sessions';

const SESSION = 'attributes.gen_ai.conversation.id';
const ppl = (schema: string[], rows: unknown[][]) => ({
schema: schema.map((name) => ({ name })),
datarows: rows,
});

describe('fetchSessions', () => {
it('reports ignored commands and whether the query filters', () => {
expect(sessionFilterFor('source = t | where a = 1 | head 5')).toEqual({
whereQuery: 'source = t | where a = 1',
ignoredCommands: ['head'],
hasFilter: true,
});
expect(sessionFilterFor('source = t').hasFilter).toBe(false);
});

it('selects sessions under the time range and totals them over the whole session', async () => {
const executeQuery = jest.fn(async (_dataset: { timeFieldName?: string }, query: string) => {
if (query.includes('as total_sessions')) return ppl(['total_sessions'], [[7]]);
if (query.includes('as last_seen')) return ppl([SESSION, 'last_seen'], [['s1', 'x']]);
if (query.includes('as total_traces'))
return ppl(
[SESSION, 'total_traces', 'start_time', 'end_time'],
[['s1', 2, '2026-09-29 10:00:00', '2026-09-29 10:01:00']]
);
if (query.includes('dedup traceId')) return ppl(['traceId', SESSION], [['t1', 's1']]);
return ppl([], []);
});
const dataset = { id: 'd', title: 'spans', type: 'INDEX_PATTERN', timeFieldName: 'endTime' };

const result = await fetchSessions(
{ executeQuery },
dataset as never,
'source = spans',
(t) => t
);

expect(result.totalSessions).toBe(7);
expect(result.sessions.map((s) => [s.sessionId, s.totalTraces])).toEqual([['s1', 2]]);
const byQuery = (part: string) =>
executeQuery.mock.calls.find(([, q]) => q.includes(part))?.[0];
// The matching query keeps the time range; per-session totals drop it.
expect(byQuery('as last_seen')?.timeFieldName).toBe('endTime');
expect(byQuery('as total_traces')?.timeFieldName).toBeUndefined();
});

it('sizes the trace map by the sessions trace counts and flags partial results', async () => {
const executeQuery = jest.fn(async (_dataset: unknown, query: string) => {
if (query.includes('as total_sessions')) return ppl(['total_sessions'], [[1]]);
if (query.includes('as last_seen')) return ppl([SESSION, 'last_seen'], [['s1', 'x']]);
if (query.includes('as total_traces'))
return ppl(
[SESSION, 'total_traces', 'start_time', 'end_time'],
[['s1', 30, '2026-09-29 10:00:00', '2026-09-29 10:01:00']]
);
if (query.includes('as spans by traceId')) return ppl(['traceId', SESSION], [['t1', 's1']]);
return ppl([], []);
});
const dataset = { id: 'd', title: 'spans', type: 'INDEX_PATTERN' };

const result = await fetchSessions(
{ executeQuery },
dataset as never,
'source = spans',
(t) => t
);
const mapQuery = executeQuery.mock.calls.find(([, q]) => q.includes('as spans by traceId'));
expect(mapQuery?.[1]).toMatch(/\| head 31$/); // 30 traces + one per session of slack
expect(result.partial).toBe(false);

const capped = await fetchSessions(
{ executeQuery },
dataset as never,
'source = spans',
(t) => t,
{ maxTraces: 10 }
);
expect(capped.partial).toBe(true);
});
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

import { AGENT_TRACES_SESSION_ID_FIELD } from '../../../../common';
import { Dataset } from '../../../../../data/common';
import { extractSpanFilterQuery, splitPplCommands } from '../traces/table_shared';
import {
PPLResponse,
transformPPLDataToTraceHits,
} from '../traces/trace_details/traces/ppl_to_trace_hits';
import { hitsToAgentSpans, spanToRow } from '../traces/hooks/tree_utils';
import { Bucket } from '../../../components/fields_selector/types';
import {
SessionRow,
assembleSessionRows,
buildRootSpansQuery,
buildMatchingSessionIdsQuery,
buildMatchingSessionCountQuery,
withoutTimeRange,
buildSessionStatsQuery,
buildTraceSessionMapQuery,
buildSessionFacetQuery,
parseFacetBuckets,
SESSION_FACET_FIELDS,
SESSION_TRACES_LIMIT,
getSourceCommand,
parseSessionStats,
pplResponseToRecords,
} from './session_utils';

/** Anything that can run a PPL query against a dataset (PPLService in the app and embeddable). */
export interface PPLQueryRunner {
executeQuery(dataset: Dataset, pplQuery: string): Promise<unknown>;
}

export interface FetchSessionsResult {
sessions: SessionRow[];
/** Sessions matching the query and time range; `sessions` holds at most the page limit. */
totalSessions: number | null;
/** Commands in the query the Sessions view does not apply (stats, head, ...). */
ignoredCommands: string[];
/** Whether the query filters spans (so an empty list means "no match"). */
hasFilter: boolean;
/**
* The listed sessions hold more traces than `maxTraces`: per-session details (first and
* last message, tokens, trace list) come from a subset of their traces.
*/
partial: boolean;
}

export interface FetchSessionsOptions {
/** Max traces looked up across the listed sessions. */
maxTraces?: number;
}

/** The span-level filter a Sessions query applies, plus what it leaves out. */
export const sessionFilterFor = (baseQueryString: string) => {
const { filterQuery, ignoredCommands } = extractSpanFilterQuery(baseQueryString);
return {
whereQuery: filterQuery,
ignoredCommands,
hasFilter: splitPplCommands(filterQuery).length > 1,
};
};

/** Session-level facet values for the fields panel (distinct sessions per value). */
export const fetchSessionFacets = async (
ppl: PPLQueryRunner,
dataset: Dataset,
whereQuery: string
): Promise<Record<string, Bucket[]>> => {
const entries = await Promise.all(
SESSION_FACET_FIELDS.map(async (field) => {
try {
const response = await ppl.executeQuery(dataset, buildSessionFacetQuery(whereQuery, field));
return [field, parseFacetBuckets(pplResponseToRecords(response), field)] as const;
} catch {
return [field, [] as Bucket[]] as const; // e.g. the field is not mapped in this index
}
})
);
return Object.fromEntries(entries);
};

/**
* Fetch the sessions list for a query. Shared by the Sessions tab and the dashboard panel.
*
* 1. Session ids matching the query and time range (the filter selects sessions), plus a count.
* 2. Stats for those sessions (trace count, start/end), over the whole session.
* 3. Trace map: which traces belong to each session (any span may carry the id).
* 4. Root spans of those traces: first/last message, tokens, user id.
*/
export const fetchSessions = async (
ppl: PPLQueryRunner,
dataset: Dataset,
baseQueryString: string,
formatTs: (ts: string) => string,
{ maxTraces = SESSION_TRACES_LIMIT }: FetchSessionsOptions = {}
): Promise<FetchSessionsResult> => {
const { whereQuery, ignoredCommands, hasFilter } = sessionFilterFor(baseQueryString);
const source = getSourceCommand(whereQuery);

const [idsResponse, countResponse] = await Promise.all([
ppl.executeQuery(dataset, buildMatchingSessionIdsQuery(whereQuery)),
ppl.executeQuery(dataset, buildMatchingSessionCountQuery(whereQuery)).catch(() => null),
]);
const total = countResponse
? Number(pplResponseToRecords(countResponse)[0]?.total_sessions)
: NaN;
const wholeSession = withoutTimeRange(dataset);
const sessionIds = pplResponseToRecords(idsResponse)
.map((r) => r[AGENT_TRACES_SESSION_ID_FIELD])
.filter((id): id is string => typeof id === 'string' && id !== '');

let stats: ReturnType<typeof parseSessionStats> = [];
if (sessionIds.length > 0) {
const statsResponse = await ppl.executeQuery(
wholeSession,
buildSessionStatsQuery(source, sessionIds)
);
stats = parseSessionStats(pplResponseToRecords(statsResponse));
}

let sessions: SessionRow[] = [];
// Size the trace map from the sessions' own trace counts (plus one row of slack per
// session for traces that carry more than one session id), capped at maxTraces.
const tracesInSessions = stats.reduce((sum, s) => sum + s.totalTraces, 0);
const partial = tracesInSessions > maxTraces;
if (stats.length > 0) {
const mapResponse = await ppl.executeQuery(
wholeSession,
buildTraceSessionMapQuery(
source,
stats.map((s) => s.sessionId),
Math.min(tracesInSessions + stats.length, maxTraces)
)
);
const traceToSession = new Map<string, string>();
for (const rec of pplResponseToRecords(mapResponse)) {
const traceId = rec.traceId;
const sessionId = rec[AGENT_TRACES_SESSION_ID_FIELD];
if (traceId && sessionId) traceToSession.set(String(traceId), String(sessionId));
}

const traceIds = [...traceToSession.keys()];
let rootRows: Array<ReturnType<typeof spanToRow>> = [];
if (traceIds.length > 0) {
const rootsResponse = await ppl.executeQuery(
wholeSession,
buildRootSpansQuery(source, traceIds)
);
// The root-span query returns span rows in the PPL response shape.
const rootHits = transformPPLDataToTraceHits(rootsResponse as PPLResponse);
rootRows = hitsToAgentSpans(rootHits).map((span, i) => spanToRow(span, i, formatTs));
}
sessions = assembleSessionRows(stats, traceToSession, rootRows);
}

return {
sessions,
totalSessions: Number.isFinite(total) ? Math.max(total, sessions.length) : null,
ignoredCommands,
hasFilter,
partial,
};
};
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

import { useEffect, useState } from 'react';
import { usePPLQueryDeps } from '../../traces/hooks/use_ppl_query_deps';
import { transformPPLDataToTraceHits } from '../../traces/trace_details/traces/ppl_to_trace_hits';
import { TraceRow, buildFullSpanTree, hitsToAgentSpans } from '../../traces/hooks/tree_utils';
import { buildSessionSpansQuery, getSourceCommand, rawStart } from '../session_utils';

export interface SessionTrace {
traceId: string;
/** Root row of the trace (carries input/output, latency, tokens). */
root: TraceRow;
/** Full span tree for the trace, as the trace flyout expects. */
tree: TraceRow[];
/** All spans of the trace flattened, earliest first. */
spans: TraceRow[];
}

export interface UseSessionDetailResult {
traces: SessionTrace[];
loading: boolean;
error: string | null;
}

const flatten = (rows: TraceRow[], out: TraceRow[] = []): TraceRow[] => {
for (const row of rows) {
out.push(row);
if (row.children?.length) flatten(row.children as TraceRow[], out);
}
return out;
};

/** Load every span of a session's traces and build one tree per trace, earliest trace first. */
export const useSessionDetail = (
traceIds: string[] | null,
formatTs: (ts: string) => string
): UseSessionDetailResult => {
const { pplService, datasetParam, baseQueryString } = usePPLQueryDeps();
const [traces, setTraces] = useState<SessionTrace[]>([]);
const [loading, setLoading] = useState(false);
const [error, setError] = useState<string | null>(null);
const key = traceIds ? traceIds.join(',') : '';

useEffect(() => {
if (!traceIds || traceIds.length === 0 || !pplService || !datasetParam || !baseQueryString) {
setTraces([]);
return;
}
let cancelled = false;
setLoading(true);
setError(null);

(async () => {
try {
// No time filter: a session's early turns can fall outside the picked range.
const datasetWithoutTime = {
id: datasetParam.id,
title: datasetParam.title,
type: datasetParam.type,
...(datasetParam.dataSource && { dataSource: datasetParam.dataSource }),
};
const source = getSourceCommand(baseQueryString);
const response = await pplService.executeQuery(
datasetWithoutTime as typeof datasetParam,
buildSessionSpansQuery(source, traceIds)
);
const spans = hitsToAgentSpans(transformPPLDataToTraceHits(response));

const byTrace = new Map<string, typeof spans>();
for (const span of spans) {
const list = byTrace.get(span.traceId) ?? [];
list.push(span);
byTrace.set(span.traceId, list);
}

const result: SessionTrace[] = [];
for (const [traceId, traceSpans] of byTrace) {
const tree = buildFullSpanTree(traceSpans, formatTs) as TraceRow[];
const root = tree.find((r) => !r.parentSpanId) ?? tree[0];
if (!root) continue;
result.push({ traceId, root, tree, spans: flatten(tree) });
}
result.sort((a, b) => rawStart(a.root).localeCompare(rawStart(b.root)));

if (!cancelled) setTraces(result);
} catch (err) {
if (!cancelled) setError((err as Error).message || 'Failed to load session');
} finally {
if (!cancelled) setLoading(false);
}
})();

return () => {
cancelled = true;
};
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [key, pplService, datasetParam, baseQueryString, formatTs]);

return { traces, loading, error };
};
Loading
Loading