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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
31 changes: 15 additions & 16 deletions src/core/eval.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -84,10 +84,12 @@ import { Transform } from "node:stream";
import { setTimeout as sleep } from "node:timers/promises";
import {
AgentCoreCLIError,
CloudWatchQueryError,
ERROR_SOURCE,
FileWriteError,
InputValidationError,
NetworkingError,
ResourceNotFoundError,
} from "../errors";
import type {
BatchEvaluationDetail,
Expand DownExpand Up@@ -426,8 +428,6 @@ export class EvalClient implements CoreEvalClient {
const logGroupName = runtimeLogGroup(runtimeId, qualifier);
const serviceName = runtimeServiceName(runtimeName, qualifier);

// CloudWatch Insights takes epoch seconds. Discovery defaults to now-7d when
// no explicit window is given (matches the batch service's default).
const endMs = input.window ? +input.window.endTime : Date.now();
const startMs = input.window ? +input.window.startTime : endMs - SEVEN_DAYS_MS;
const startSec = Math.floor(startMs / 1000);
Expand All@@ -441,10 +441,10 @@ export class EvalClient implements CoreEvalClient {
const [runtimeRows, sharedRows] = await Promise.all([
runInsightsQuery(logs, [logGroupName], queryString, startSec, endSec).catch((error) => {
if (error instanceof ResourceNotFoundException) {
throw new InputValidationError(
throw new ResourceNotFoundError(
`No telemetry found for agent "${input.agent}": its runtime log group ${logGroupName} ` +
`does not exist. Ensure the agent has been invoked and emits traces.`,
{ meta: { agent: input.agent, logGroupName } },
{ cause: error, meta: { agent: input.agent, logGroupName } },
);
}
throw error;
Expand All@@ -454,7 +454,7 @@ export class EvalClient implements CoreEvalClient {
throw error;
}),
]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows]);
const traces = groupSpansBySession([...sharedRows, ...runtimeRows], this.logger);

// Warn when explicitly requested sessions never showed up in the logs (aged
// out, wrong id, or never emitted) so a caller isn't misled by a partial run.
Expand DownExpand Up@@ -1390,12 +1390,6 @@ function sanitizeQueryValue(value: string): string {
return value.replace(/'/g, "");
}

// buildSpanQuery is the single-phase Insights query: scope to one runtime by its
// OTel service.name, optionally narrow to specific sessions and/or one trace, and
// select the full span JSON (@message) plus the session id to group by. It does
// NOT over-filter on ispresent(kind) — that span-only predicate is what forced the
// old CLI's second query for log records; the looser scope returns everything for
// the session in one pass.
function buildSpanQuery(serviceName: string, sessionIds?: string[], traceId?: string): string {
let query = `fields @message, attributes.session.id as sessionId, traceId, spanId
| filter resource.attributes.service.name in ['${sanitizeQueryValue(serviceName)}']`;
Expand DownExpand Up@@ -1433,15 +1427,15 @@ async function runInsightsQuery(
const result = await logs.send(new GetQueryResultsCommand({ queryId }));
status = result.status ?? "Unknown";
if (status === "Failed" || status === "Cancelled" || status === "Timeout") {
throw new NetworkingError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId },
throw new CloudWatchQueryError(`CloudWatch Logs Insights query ${status.toLowerCase()}`, {
meta: { queryId, status },
});
}
if (status !== "Complete") await new Promise((resolve) => setTimeout(resolve, 1000));
}
if (status !== "Complete") {
throw new NetworkingError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId },
throw new CloudWatchQueryError("CloudWatch Logs Insights query did not finish in time", {
meta: { queryId, status },
});
}

Expand All@@ -1466,9 +1460,10 @@ async function runInsightsQuery(

// Group parsed @message docs by session, keeping only sessions with >=1 span
// (Evaluate rejects log-only sessions), and derive each session's trace/tool ids.
function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
function groupSpansBySession(rows: ResultField[][], logger: Logger): SessionTrace[] {
const docsBySession = new Map<string, SpanRecord[]>();
const sessionsWithSpans = new Set<string>();
let warnedAboutMalformedTelemetry = false;
for (const row of rows) {
const message = row.find((f) => f.field === "@message")?.value;
const sessionId = row.find((f) => f.field === "sessionId")?.value;
Expand All@@ -1479,6 +1474,10 @@ function groupSpansBySession(rows: ResultField[][]): SessionTrace[] {
try {
doc = JSON.parse(message) as SpanRecord;
} catch {
if (!warnedAboutMalformedTelemetry) {
logger.warn("skipping malformed telemetry records");
warnedAboutMalformedTelemetry = true;
}
continue;
}
const list = docsBySession.get(sessionId);
Expand Down
7 changes: 7 additions & 0 deletions src/errors/errors.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -183,6 +183,13 @@ export class NetworkingError extends AgentCoreCLIError {
}
}

/** A CloudWatch Logs Insights query reached a terminal failure state. */
export class CloudWatchQueryError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
super(message, { ...options, source: ERROR_SOURCE.SERVICE });
}
}

/** Service data was returned successfully, but did not match the expected contract. */
export class MalformedServiceResponseError extends AgentCoreCLIError {
constructor(message: string, options?: Omit<AgentCoreCLIErrorOptions, "source">) {
Expand Down
1 change: 1 addition & 0 deletions src/errors/index.tsx
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
export {
AgentCoreCLIError,
CloudWatchQueryError,
CommandInterruptedError,
DeserializationError,
EmbeddedAssetNotFoundError,
Expand Down
15 changes: 0 additions & 15 deletions src/handlers/eval/ondemand/ondemand.fixture.test.tsx
Original file line numberDiff line numberDiff line change
Expand Up@@ -13,21 +13,6 @@ import { createRootHandler } from "../../index";
const REGION = "us-west-2";
const FIXTURES = join(import.meta.dir, "__fixtures__");

// Record with: RECORD=1 bun test src/handlers/eval/ondemand/ondemand.fixture.test.tsx
//
// This exercises the real seam end to end: parsing → handler → CoreClient →
// getTracesForAgent (GetAgentRuntime + CloudWatch Logs Insights StartQuery /
// GetQueryResults, read from aws/spans and the runtime group) → evaluate (the
// Evaluate data-plane API) → rendered scores.
//
// Determinism: the window is PINNED (not --lookback-days) so the StartQuery input —
// which embeds startTime/endTime epoch seconds — hashes to the same fixture on
// record and replay. --session-ids bounds the fetch to the two sessions recorded
// against the live agent below.
//
// Re-recording needs the agent to still exist AND those sessions' spans to still be
// within CloudWatch retention (they age out). If they've aged out, invoke the agent
// to create fresh sessions, then repoint FIXTURE_SESSION_IDS + the window at them.
const FIXTURE_AGENT = "asdf_MyAgent-3s5axvBC6Q";
const FIXTURE_SESSION_IDS = [
"67ebf93b-65e3-4127-9e13-483b239f256a",
Expand Down
Loading
Loading