Skip to content
Open
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
5 changes: 5 additions & 0 deletions .changeset/fix-ai-sdk-middleware-metrics.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"braintrust": minor
---

feat: Fix token metrics with AI SDK middleware
146 changes: 106 additions & 40 deletions e2e/helpers/mock-braintrust-server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,45 +90,111 @@ describe("production forwarding", () => {
},
);

it("reports a failed write without preventing later queued writes", async () => {
const received: number[] = [];
const upstream = createServer(async (req, res) => {
let body = "";
for await (const chunk of req) body += chunk;
const { sequence } = JSON.parse(body);
received.push(sequence);
res.statusCode = sequence === 1 ? 500 : 200;
res.end(sequence === 1 ? "initial write failed" : "{}");
});
await new Promise<void>((resolve) =>
upstream.listen(0, "127.0.0.1", resolve),
);
const url = `http://127.0.0.1:${(upstream.address() as AddressInfo).port}`;
const server = await startMockBraintrustServer({
prodForwarding: {
apiKey: "test-only-key",
apiUrl: url,
appUrl: url,
orgId: "org",
orgName: "org",
projectId: "project",
projectName: "tmp-luca-forwarding-test",
},
});
try {
for (const sequence of [1, 2]) {
const response = await fetch(`${server.url}/logs3`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ sequence, api_version: 2, rows: [] }),
});
await response.text();
it.each(["/logs3", "/otel/v1/traces"])(
"retries transient failures for %s before forwarding later writes",
async (path) => {
const received: Array<{ sequence: number; rows: unknown[] }> = [];
let attempts = 0;
const upstream = createServer(async (req, res) => {
let body = "";
for await (const chunk of req) body += chunk;
const payload = JSON.parse(body);
received.push(payload);
// Reproduce a gateway outage followed by successful ingestion.
res.statusCode = payload.sequence === 1 && ++attempts <= 2 ? 502 : 200;
res.end(res.statusCode === 502 ? "Bad Gateway" : "{}");
});
await new Promise<void>((resolve) =>
upstream.listen(0, "127.0.0.1", resolve),
);
const url = `http://127.0.0.1:${(upstream.address() as AddressInfo).port}`;
const server = await startMockBraintrustServer({
prodForwarding: {
apiKey: "test-only-key",
apiUrl: url,
appUrl: url,
orgId: "org",
orgName: "org",
projectId: "project",
projectName: "tmp-luca-forwarding-test",
},
});
try {
for (const sequence of [1, 2]) {
const response = await fetch(`${server.url}${path}`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
sequence,
api_version: 2,
rows: [{ id: "span", _is_merge: sequence === 2 }],
}),
});
expect(response.ok).toBe(true);
await response.text();
}
await expect(server.close()).resolves.toBeUndefined();
expect(received.map((payload) => payload.sequence)).toEqual([
1, 1, 1, 2,
]);
expect(received[1]).toEqual(received[0]);
expect(received[2]).toEqual(received[0]);
} finally {
upstream.closeAllConnections();
await new Promise<void>((resolve) => upstream.close(() => resolve()));
}
await expect(server.close()).rejects.toThrow("initial write failed");
expect(received).toEqual([1, 2]);
} finally {
upstream.closeAllConnections();
await new Promise<void>((resolve) => upstream.close(() => resolve()));
}
});
},
);

it.each([
{ status: 400, expectedSequence: [1, 2] },
{ status: 401, expectedSequence: [1, 2] },
{ status: 500, expectedSequence: [1, 1, 1, 2] },
{ status: 502, expectedSequence: [1, 1, 1, 2] },
{ status: 503, expectedSequence: [1, 1, 1, 2] },
{ status: 504, expectedSequence: [1, 1, 1, 2] },
])(
"reports HTTP $status failures without preventing later queued writes",
async ({ status, expectedSequence }) => {
const received: number[] = [];
const upstream = createServer(async (req, res) => {
let body = "";
for await (const chunk of req) body += chunk;
const { sequence } = JSON.parse(body);
received.push(sequence);
res.statusCode = sequence === 1 ? status : 200;
res.end(sequence === 1 ? "initial write failed" : "{}");
});
await new Promise<void>((resolve) =>
upstream.listen(0, "127.0.0.1", resolve),
);
const url = `http://127.0.0.1:${(upstream.address() as AddressInfo).port}`;
const server = await startMockBraintrustServer({
prodForwarding: {
apiKey: "test-only-key",
apiUrl: url,
appUrl: url,
orgId: "org",
orgName: "org",
projectId: "project",
projectName: "tmp-luca-forwarding-test",
},
});
try {
for (const sequence of [1, 2]) {
const response = await fetch(`${server.url}/logs3`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ sequence, api_version: 2, rows: [] }),
});
await response.text();
}
await expect(server.close()).rejects.toThrow("initial write failed");
expect(received).toEqual(expectedSequence);
} finally {
upstream.closeAllConnections();
await new Promise<void>((resolve) => upstream.close(() => resolve()));
}
},
);
});
22 changes: 20 additions & 2 deletions e2e/helpers/mock-braintrust-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import type {
ServerResponse,
} from "node:http";
import type { AddressInfo } from "node:net";
import { setTimeout } from "node:timers/promises";
import type { ProdForwarding } from "./prod-forwarding";

export type JsonValue =
Expand Down Expand Up @@ -470,14 +471,31 @@ export async function startMockBraintrustServer(
}
headers.set("authorization", `Bearer ${prodForwarding.apiKey}`);

const response = await fetch(url, {
const requestInit = {
body:
prodRequest.method === "GET" || prodRequest.method === "HEAD"
? undefined
: prodRequest.rawBody,
headers,
method: prodRequest.method,
});
};
// Retry ingestion inside the forwarding queue so later merges cannot
// overtake the initial upsert. Registration requests remain single-attempt.
const maxAttempts =
prodRequest.method === "POST" &&
["/logs3", "/otel/v1/traces"].includes(prodRequest.path)
? 3
: 1;
let response = await fetch(url, requestInit);
for (
let attempt = 1;
attempt < maxAttempts && [500, 502, 503, 504].includes(response.status);
attempt++
) {
await response.arrayBuffer().catch(() => {});
await setTimeout(500 * 2 ** (attempt - 1));
response = await fetch(url, requestInit);
}

if (!response.ok) {
const responseText = await response.text().catch(() => "");
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
Expand Up @@ -2893,8 +2893,8 @@ span_tree:
│ "completion_tokens": 7,
│ "prompt_cache_creation_tokens": 0,
│ "prompt_cached_tokens": 4516,
│ "prompt_tokens": 3,
│ "tokens": 10
│ "prompt_tokens": 4519,
│ "tokens": 4526
│ }
├── ai-sdk-anthropic-cache-operation
│ metadata: {
Expand Down Expand Up @@ -3134,8 +3134,8 @@ span_tree:
│ │ "completion_tokens": 7,
│ │ "prompt_cache_creation_tokens": 0,
│ │ "prompt_cached_tokens": 4516,
│ │ "prompt_tokens": 3,
│ │ "tokens": 10
│ │ "prompt_tokens": 4519,
│ │ "tokens": 4526
│ │ }
│ └── generateText [function]
│ input: {
Expand Down Expand Up @@ -3370,8 +3370,8 @@ span_tree:
│ "completion_tokens": 7,
│ "prompt_cache_creation_tokens": 0,
│ "prompt_cached_tokens": 4516,
│ "prompt_tokens": 3,
│ "tokens": 10
│ "prompt_tokens": 4519,
│ "tokens": 4526
│ }
├── ai-sdk-deny-output-override-operation
│ metadata: {
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
Expand Up @@ -2975,8 +2975,8 @@ span_tree:
│ "completion_tokens": 7,
│ "prompt_cache_creation_tokens": 0,
│ "prompt_cached_tokens": 4516,
│ "prompt_tokens": 3,
│ "tokens": 10
│ "prompt_tokens": 4519,
│ "tokens": 4526
│ }
├── ai-sdk-anthropic-cache-operation
│ metadata: {
Expand Down Expand Up @@ -3216,8 +3216,8 @@ span_tree:
│ │ "completion_tokens": 7,
│ │ "prompt_cache_creation_tokens": 0,
│ │ "prompt_cached_tokens": 4516,
│ │ "prompt_tokens": 3,
│ │ "tokens": 10
│ │ "prompt_tokens": 4519,
│ │ "tokens": 4526
│ │ }
│ └── generateText [function]
│ input: {
Expand Down Expand Up @@ -3452,8 +3452,8 @@ span_tree:
│ "completion_tokens": 7,
│ "prompt_cache_creation_tokens": 0,
│ "prompt_cached_tokens": 4516,
│ "prompt_tokens": 3,
│ "tokens": 10
│ "prompt_tokens": 4519,
│ "tokens": 4526
│ }
├── ai-sdk-deny-output-override-operation
│ metadata: {
Expand Down
Loading
Loading