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
71 changes: 68 additions & 3 deletions src/lib/providers/anthropic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,12 @@ import {
} from "../types/index.js";
import { logger } from "../utils/logger.js";
import { redactUrlCredentials } from "../utils/logSanitize.js";
import {
ANTHROPIC_MAX_CACHE_BREAKPOINTS,
applyAnthropicHistoryCacheBreakpoints,
countAnthropicCacheMarkers,
} from "../utils/anthropicCacheBreakpoints.js";
import type { VertexAnthropicMessage } from "../types/index.js";
import { calculateCost } from "../utils/pricing.js";
import {
createAnthropicConfig,
Expand Down Expand Up @@ -1446,9 +1452,27 @@ export class AnthropicProvider extends BaseProvider {
| { type: "enabled"; budget_tokens: number }
| undefined;

// Prompt-cache parity with the native Vertex+Claude path: upstream
// layers mark only the stable prefix (system via MessageBuilder,
// last tool via GenerationHandler) — the growing conversation
// history has no breakpoint, so on every turn it falls after the
// last marker and is re-billed as fresh input. Add rolling history
// breakpoints in whatever budget remains under Anthropic's
// four-marker ceiling; pre-existing markers are counted so the
// request can never exceed the cap.
const cacheMarkersUsed = countAnthropicCacheMarkers({
system,
tools,
messages: messages as unknown as VertexAnthropicMessage[],
});
const cachedMessages = applyAnthropicHistoryCacheBreakpoints(
messages as unknown as VertexAnthropicMessage[],
ANTHROPIC_MAX_CACHE_BREAKPOINTS - cacheMarkersUsed,
) as unknown as Anthropic.Messages.MessageParam[];
Comment on lines +1463 to +1471

const params: Anthropic.Messages.MessageCreateParamsNonStreaming = {
model: modelId,
messages,
messages: cachedMessages,
max_tokens: resolveClaudeMaxTokens(modelId, options.maxOutputTokens),
...(system ? { system } : {}),
...(options.temperature !== undefined &&
Expand Down Expand Up @@ -1741,10 +1765,22 @@ export class AnthropicProvider extends BaseProvider {
"gen_ai.usage.output_tokens",
usage.completionTokens || 0,
);
if (usage.cacheReadTokens) {
streamSpan.setAttribute(
"gen_ai.usage.cached_input_tokens",
usage.cacheReadTokens,
);
}
const cost = calculateCost(this.providerName, this.modelName, {
input: usage.promptTokens || 0,
output: usage.completionTokens || 0,
total: usage.totalTokens || 0,
...(usage.cacheReadTokens
? { cacheReadTokens: usage.cacheReadTokens }
: {}),
...(usage.cacheCreationTokens
? { cacheCreationTokens: usage.cacheCreationTokens }
: {}),
});
if (cost && cost > 0) {
streamSpan.setAttribute("neurolink.cost", cost);
Expand Down Expand Up @@ -1773,12 +1809,30 @@ export class AnthropicProvider extends BaseProvider {
const conversation = payload.messages.slice();
let totalInput = 0;
let totalOutput = 0;
let totalCacheRead = 0;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 The comment claims conversation itself is never mutated, but the loop later calls conversation.push(...) twice per tool step. The breakpoint helper is pure (it clones the array), but the statement about conversation is misleading. Consider rephrasing to clarify that the helper’s purity is what keeps re-counting stable, e.g.:

// Pure: applyAnthropicHistoryCacheBreakpoints clones `conversation`,
// so re-counting the original array at the top of each step stays stable.

let totalCacheWrite = 0;
let lastStop: string | null = null;

for (let step = 0; step < maxSteps; step++) {
// Prompt-cache parity with the native Vertex+Claude path — rolling
// history breakpoints, re-applied per step so the stable prefix
// stays byte-identical while the breakpoint follows the growing
// tail. Budget respects markers upstream layers already placed
// (system / last tool / message blocks) so the request never
// exceeds Anthropic's four-marker cap. Pure: `conversation` itself
// is never mutated, so re-counting per step stays stable.
const cacheMarkersUsed = countAnthropicCacheMarkers({
system: payload.system,
tools: anthropicTools,
messages: conversation as unknown as VertexAnthropicMessage[],
});
const cachedConversation = applyAnthropicHistoryCacheBreakpoints(
conversation as unknown as VertexAnthropicMessage[],
ANTHROPIC_MAX_CACHE_BREAKPOINTS - cacheMarkersUsed,
) as unknown as Anthropic.Messages.MessageParam[];
Comment on lines +1824 to +1832
const params: Anthropic.Messages.MessageCreateParamsStreaming = {
model: modelId,
messages: conversation,
messages: cachedConversation,
max_tokens: resolveClaudeMaxTokens(modelId, options.maxTokens),
stream: true,
...(payload.system ? { system: payload.system } : {}),
Expand Down Expand Up @@ -1816,6 +1870,12 @@ export class AnthropicProvider extends BaseProvider {
if (event.type === "message_start") {
totalInput += event.message.usage.input_tokens ?? 0;
totalOutput += event.message.usage.output_tokens ?? 0;
// Anthropic reports cache reads/writes SEPARATELY from
// input_tokens on the same message_start event — without these
// the streaming path silently drops all cache accounting.
totalCacheRead += event.message.usage.cache_read_input_tokens ?? 0;
totalCacheWrite +=
event.message.usage.cache_creation_input_tokens ?? 0;
} else if (event.type === "content_block_start") {
blockTypes.set(event.index, event.content_block.type);
if (event.content_block.type === "tool_use") {
Expand Down Expand Up @@ -2016,7 +2076,12 @@ export class AnthropicProvider extends BaseProvider {
resolveUsage({
promptTokens: totalInput,
completionTokens: totalOutput,
totalTokens: totalInput + totalOutput,
totalTokens:
totalInput + totalCacheRead + totalCacheWrite + totalOutput,
...(totalCacheRead > 0 ? { cacheReadTokens: totalCacheRead } : {}),
...(totalCacheWrite > 0
? { cacheCreationTokens: totalCacheWrite }
: {}),
});
resolveFinish(lastStop ?? "stop");
};
Expand Down
13 changes: 3 additions & 10 deletions src/lib/providers/openaiChatCompletionsClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import type {
OpenAICompatToolChoiceWire,
OpenAICompatV3CallToolChoice,
OpenAICompatV3CallTools,
DeferredUsage,
Tool,
} from "../types/index.js";
import { convertZodToJsonSchema } from "../utils/schemaConversion.js";
Expand Down Expand Up @@ -639,16 +640,8 @@ export const buildAPIError = async (
// collector resolves with the actual aggregated values after the multi-step
// loop ends, not the zeros they had at result-construction time.
export const createDeferredAnalytics = () => {
let resolveUsage: (u: {
promptTokens: number;
completionTokens: number;
totalTokens: number;
}) => void = () => {};
const usagePromise = new Promise<{
promptTokens: number;
completionTokens: number;
totalTokens: number;
}>((r) => {
let resolveUsage: (u: DeferredUsage) => void = () => {};

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 The new DeferredUsage type is a good consolidation. Just confirming: createDeferredAnalytics is also consumed by openaiChatCompletionsBase.ts (via this module). The base class currently resolves usage with only promptTokens/completionTokens/totalTokens. Since DeferredUsage now carries optional cache fields, consider whether openaiChatCompletionsBase.runStreamLoop should forward cache values when providers populate them. Not blocking for this PR, but worth a follow-up so OpenAI-compatible cache reads aren't dropped at the resolve boundary.

const usagePromise = new Promise<DeferredUsage>((r) => {
resolveUsage = r;
});
let resolveFinish: (reason: string) => void = () => {};
Expand Down
23 changes: 23 additions & 0 deletions src/lib/types/common.ts
Original file line number Diff line number Diff line change
Expand Up @@ -357,6 +357,16 @@ export type RawUsageObject = {
cacheCreationTokens?: number;
cacheReadTokens?: number;

// ai@6 normalized usage (generateText/streamText result.usage) — the SDK's
// asLanguageModelUsage() exposes cache reads flat as `cachedInputTokens`
// and the read/write split nested under `inputTokenDetails`. The native
// direct-Anthropic V3 model reports cache tokens ONLY through this shape.
cachedInputTokens?: number;
inputTokenDetails?: {
cacheReadTokens?: number;
cacheWriteTokens?: number;
};

// OpenAI/DeepSeek/NIM/OpenAI-compatible nested cache field (overlapping:
// cached_tokens is a SUBSET already included in prompt_tokens)
prompt_tokens_details?: { cached_tokens?: number };
Expand All @@ -371,6 +381,19 @@ export type RawUsageObject = {
usage?: RawUsageObject;
};

/**
* Aggregated usage resolved by a provider's deferred-analytics pair after a
* multi-step stream loop ends. The cache fields are optional — only providers
* with prompt caching (Anthropic) populate them.
*/
export type DeferredUsage = {
promptTokens: number;
completionTokens: number;
totalTokens: number;
cacheReadTokens?: number;
cacheCreationTokens?: number;
};

/**
* Options for token extraction from raw usage objects.
*/
Expand Down
69 changes: 68 additions & 1 deletion src/lib/utils/anthropicCacheBreakpoints.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,8 @@ import type {
const EPHEMERAL: VertexAnthropicCacheControl = { type: "ephemeral" };

/** Anthropic allows at most four `cache_control` breakpoints per request. */
const MAX_BREAKPOINTS = 4;
export const ANTHROPIC_MAX_CACHE_BREAKPOINTS = 4;
const MAX_BREAKPOINTS = ANTHROPIC_MAX_CACHE_BREAKPOINTS;

/**
* Annotate a native Vertex+Claude request with prompt-cache breakpoints.
Expand Down Expand Up @@ -82,6 +83,72 @@ export function applyVertexAnthropicCacheBreakpoints(
return { system, tools, messages };
}

/**
* Count the `cache_control` markers already present on a request. The direct
* Anthropic path (AI-SDK pipeline) arrives with markers the upstream layers
* placed — MessageBuilder tags the system prompt, GenerationHandler tags the
* last tool definition, and message content blocks may carry translated
* AI-SDK markers. Anthropic rejects requests with more than four markers, so
* any additional history breakpoints must fit in the remaining budget.
*/
export function countAnthropicCacheMarkers(input: {
system?: string | ReadonlyArray<{ cache_control?: unknown }> | undefined;
tools?: ReadonlyArray<{ cache_control?: unknown }> | undefined;
messages: ReadonlyArray<VertexAnthropicMessage>;
}): number {
let count = 0;
if (Array.isArray(input.system)) {
count += input.system.filter((b) => b.cache_control).length;
}
if (input.tools) {
count += input.tools.filter((t) => t.cache_control).length;
}
for (const message of input.messages) {
if (Array.isArray(message.content)) {
count += message.content.filter(
(b) => (b as { cache_control?: unknown }).cache_control,
).length;
}
}
return count;
}

/**
* Rolling history breakpoints for request paths whose stable-prefix markers
* are managed upstream (direct Anthropic: system via MessageBuilder, last
* tool via GenerationHandler). Marks the last content block of up to `budget`
* tail messages; a tail block that already carries a marker is left as-is
* without consuming budget (it already serves as that breakpoint). Pure —
* the input array is cloned, never mutated.
*/
export function applyAnthropicHistoryCacheBreakpoints(
input: ReadonlyArray<VertexAnthropicMessage>,
budget: number,
): VertexAnthropicMessage[] {
const messages = input.map((m) => ({ ...m }));
let remaining = Math.max(0, Math.min(budget, MAX_BREAKPOINTS));
for (let i = messages.length - 1; i >= 0 && remaining > 0; i--) {
if (lastContentBlockHasMarker(messages[i])) {
continue;
}
if (markLastContentBlock(messages, i)) {
remaining--;
}
}
return messages;
}

/** True when the last content block of a message already carries `cache_control`. */
function lastContentBlockHasMarker(message: VertexAnthropicMessage): boolean {
if (!Array.isArray(message.content) || message.content.length === 0) {
return false;
}
const last = message.content[message.content.length - 1] as {
cache_control?: unknown;
};
return !!last.cache_control;
}

/**
* Place a cache breakpoint on the last content block of `messages[i]`.
* Anthropic attaches `cache_control` to a content block, not the message
Expand Down
27 changes: 25 additions & 2 deletions src/lib/utils/tokenUtils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,9 @@ export function extractReasoningTokens(

/**
* Extract cache creation token count from various provider formats
* Supports: cacheCreationInputTokens, cacheCreationTokens
* Supports: cacheCreationInputTokens, cacheCreationTokens, and the ai@6
* normalized `inputTokenDetails.cacheWriteTokens` (the only shape the native
* direct-Anthropic path reports through generateText).
Comment on lines 94 to +98

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

printf '\n== tokenUtils.ts ==\n'
sed -n '1,240p' src/lib/utils/tokenUtils.ts

printf '\n== search extractInputTokens / calculateCost usages ==\n'
rg -n "extractInputTokens|calculateCost|extractCacheReadTokens|extractCacheCreationTokens|cachedInputTokens|inputTokenDetails|cacheReadTokens|cacheWriteTokens" src/lib -S

Repository: juspay/neurolink

Length of output: 26702


🏁 Script executed:

#!/bin/bash
set -euo pipefail

printf '\n== ai sdk anthopic source/docs search terms ==\n'
rg -n "cacheReadTokens|cacheWriteTokens|cachedInputTokens|inputTokens.*cache|inputTokenDetails|LanguageModelUsage" . -S

Repository: juspay/neurolink

Length of output: 21836


🌐 Web query:

AI SDK 6 Anthropic inputTokens inputTokenDetails cacheReadTokens cacheWriteTokens inclusive exclusive semantics

💡 Result:

In the Vercel AI SDK (version 6 and later), token usage reporting for Anthropic has been standardized to use the provider-agnostic usage object, moving away from provider-specific metadata fields [1]. Key components of the usage object related to Anthropic prompt caching are found within result.usage.inputTokenDetails [2][1][3][4]: - cacheReadTokens: Represents the number of input tokens served from the cache (cache hit) [1][5][3]. - cacheWriteTokens: Represents the number of input tokens written to the cache (cache creation/miss) [1][5][3]. - noCacheTokens: Represents the number of input tokens that were not cached [3]. These values are accessible directly from the usage object returned by generateText and streamText [1][4]. The AI SDK previously exposed cache creation tokens via providerMetadata.anthropic (as cacheCreationInputTokens), but this field was removed to unify the interface across all providers [1]. Regarding semantics: - The inputTokens total in the usage object generally reflects the total number of tokens processed [6]. - Cache read and cache write tokens are distinct; the first turn often involves cacheWriteTokens > 0 (and cacheReadTokens = 0), while subsequent turns with a byte-stable prefix will show cacheReadTokens > 0 [5]. - "Inclusive" or "exclusive" phrasing is not standard terminology for these fields in the SDK; rather, these fields are additive components that categorize the total inputTokens (i.e., inputTokens ≈ noCacheTokens + cacheReadTokens + cacheWriteTokens) [6][3]. For implementation, you set cache breakpoints using providerOptions in your message content [4]: { providerOptions: { anthropic: { cacheControl: { type: 'ephemeral' } } } } [4]

Citations:


🏁 Script executed:

#!/bin/bash
set -euo pipefail

sed -n '730,770p' src/lib/utils/pricing.ts
printf '\n== relevant tokenUtils block ==\n'
sed -n '236,320p' src/lib/utils/tokenUtils.ts
printf '\n== tests around cache breakpoints ==\n'
sed -n '470,525p' test/continuous-test-suite-cache-breakpoints.ts

Repository: juspay/neurolink

Length of output: 6499


🏁 Script executed:

#!/bin/bash
set -euo pipefail

printf '\n== src/lib/types/common.ts slice ==\n'
sed -n '330,390p' src/lib/types/common.ts

printf '\n== test/continuous-test-suite-cache-breakpoints.ts slice ==\n'
sed -n '490,535p' test/continuous-test-suite-cache-breakpoints.ts

printf '\n== GenerationHandler usage/cost slice ==\n'
sed -n '650,950p' src/lib/core/modules/GenerationHandler.ts

Repository: juspay/neurolink

Length of output: 14236


🏁 Script executed:

#!/bin/bash
set -euo pipefail

python3 - <<'PY'
from pathlib import Path
p = Path('test/continuous-test-suite-cache-breakpoints.ts')
text = p.read_text()
for i, line in enumerate(text.splitlines(), 1):
    if 525 <= i <= 575:
        print(f"{i}: {line}")
PY

Repository: juspay/neurolink

Length of output: 715


🏁 Script executed:

#!/bin/bash
set -euo pipefail

python3 - <<'PY'
from pathlib import Path
p = Path('test/continuous-test-suite-cache-breakpoints.ts')
for i, line in enumerate(p.read_text().splitlines(), 1):
    if 525 <= i <= 575:
        print(f"{i}: {line}")
PY

Repository: juspay/neurolink

Length of output: 715


Treat ai@6 cache tokens as overlapping with inputTokens
inputTokens is the inclusive total for the direct-Anthropic generateText shape, so cachedInputTokens, inputTokenDetails.cacheReadTokens, and inputTokenDetails.cacheWriteTokens need to be removed from input before pricing. As written, they’re handled like additive buckets and calculateCost will bill the cached/write portion twice.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/lib/utils/tokenUtils.ts` around lines 94 - 98, The token extraction and
pricing logic must treat ai@6 cache token fields as subsets of total input, not
additive tokens. Update the relevant functions in tokenUtils.ts, including the
cache extraction and calculateCost paths, to subtract cachedInputTokens,
inputTokenDetails.cacheReadTokens, and inputTokenDetails.cacheWriteTokens from
inputTokens before pricing, while preserving existing provider-specific
handling.

*/
export function extractCacheCreationTokens(
usage: RawUsageObject,
Expand All @@ -110,12 +112,21 @@ export function extractCacheCreationTokens(
) {
return usage.cacheCreationTokens;
}
if (
typeof usage.inputTokenDetails?.cacheWriteTokens === "number" &&
usage.inputTokenDetails.cacheWriteTokens > 0
) {
return usage.inputTokenDetails.cacheWriteTokens;
}
return undefined;
}

/**
* Extract cache read token count from various provider formats
* Supports: cacheReadInputTokens, cacheReadTokens
* Supports: cacheReadInputTokens, cacheReadTokens, and the ai@6 normalized
* shapes — flat `cachedInputTokens` and nested
* `inputTokenDetails.cacheReadTokens` (the only shapes the native
* direct-Anthropic path reports through generateText).
*/
export function extractCacheReadTokens(
usage: RawUsageObject,
Expand All @@ -129,6 +140,18 @@ export function extractCacheReadTokens(
if (typeof usage.cacheReadTokens === "number" && usage.cacheReadTokens > 0) {
return usage.cacheReadTokens;
}
if (
typeof usage.cachedInputTokens === "number" &&
usage.cachedInputTokens > 0
) {
return usage.cachedInputTokens;
}
if (
typeof usage.inputTokenDetails?.cacheReadTokens === "number" &&
usage.inputTokenDetails.cacheReadTokens > 0
) {
return usage.inputTokenDetails.cacheReadTokens;
}
return undefined;
}

Expand Down
Loading
Loading