feat: add API inference logging for debugging

- Add ApiInferenceLogger singleton class for logging API requests/responses
- Add inferenceLogger property to BaseProvider
- Update all providers to use the logger
- Configure logger in extension.ts with output channel sink
- Add payload sanitization (truncate strings, cap arrays, redact secrets)
- Enable via ROO_CODE_API_LOGGING=true environment variable

Closes ROO-271
This commit is contained in:
daniel-lxs 2025-12-22 17:29:08 -05:00 committed by Hannes Rudolph
parent c193f59819
commit 87546c9fce
27 changed files with 2014 additions and 634 deletions

View file

@ -0,0 +1,359 @@
/**
* Lightweight logger for API inference requests/responses.
*
* This logger is designed to capture raw inference inputs/outputs across providers
* for debugging purposes. It emits structured objects to a configurable sink.
*
* For streaming requests, only the final assembled response is logged (not individual chunks).
*
* Enable via environment variable: process.env.ROO_CODE_API_LOGGING === "true"
*/
export interface ApiInferenceLoggerConfig {
enabled: boolean
sink: (...args: unknown[]) => void
}
export interface ApiInferenceContext {
provider: string
operation: string
model?: string
taskId?: string
requestId?: string
}
export interface ApiInferenceHandle {
success: (responsePayload: unknown) => void
error: (errorPayload: unknown) => void
}
/**
* Configuration for payload size limiting to avoid freezing the Output Channel.
*/
const PAYLOAD_LIMITS = {
/** Maximum string length before truncation */
MAX_STRING_LENGTH: 10_000,
/** Maximum array entries to log */
MAX_ARRAY_LENGTH: 200,
/** Maximum object keys to log */
MAX_OBJECT_KEYS: 200,
}
/**
* Regex pattern for detecting base64 image data URLs.
*/
const BASE64_IMAGE_PATTERN = /^data:image\/[^;]+;base64,/
/**
* Secret field patterns to redact in logged payloads.
* Case-insensitive matching is applied.
* Note: Patterns are designed to avoid false positives (e.g., "inputTokens" should not be redacted).
*/
const SECRET_PATTERNS = [
"authorization",
"apikey",
"api_key",
"x-api-key",
"access_token",
"accesstoken",
"bearer",
"secret",
"password",
"credential",
]
/**
* Patterns that indicate a field is NOT a secret (allowlist).
* These are checked before secret patterns to prevent false positives.
*/
const NON_SECRET_PATTERNS = ["inputtokens", "outputtokens", "cachetokens", "reasoningtokens", "totaltokens"]
/**
* Check if a key name looks like a secret field.
*/
function isSecretKey(key: string): boolean {
const lowerKey = key.toLowerCase()
// Check allowlist first to avoid false positives
if (NON_SECRET_PATTERNS.some((pattern) => lowerKey.includes(pattern))) {
return false
}
return SECRET_PATTERNS.some((pattern) => lowerKey.includes(pattern))
}
/**
* Truncate a string if it exceeds the maximum length.
* Also replaces base64 image data with a placeholder.
*/
function sanitizeString(str: string): string {
// Check for base64 image data URLs first
if (BASE64_IMAGE_PATTERN.test(str)) {
return `[ImageData len=${str.length}]`
}
// Truncate long strings
if (str.length > PAYLOAD_LIMITS.MAX_STRING_LENGTH) {
return `[Truncated len=${str.length}]`
}
return str
}
/**
* Recursively sanitize and redact secrets from an object.
* Applies size limiting to prevent Output Channel from freezing:
* - Strings longer than MAX_STRING_LENGTH are truncated
* - Arrays longer than MAX_ARRAY_LENGTH are capped
* - Objects with more than MAX_OBJECT_KEYS are capped
* - Base64 image data URLs are replaced with placeholders
* - Secret fields are redacted
* Returns a sanitized copy of the object.
*/
function sanitizePayload(obj: unknown, visited = new WeakSet<object>()): unknown {
if (obj === null || obj === undefined) {
return obj
}
// Handle strings
if (typeof obj === "string") {
return sanitizeString(obj)
}
// Handle other primitives
if (typeof obj !== "object") {
return obj
}
// Prevent infinite recursion on circular references
if (visited.has(obj as object)) {
return "[Circular Reference]"
}
visited.add(obj as object)
// Handle arrays with length limiting
if (Array.isArray(obj)) {
const maxLen = PAYLOAD_LIMITS.MAX_ARRAY_LENGTH
if (obj.length > maxLen) {
const truncated = obj.slice(0, maxLen).map((item) => sanitizePayload(item, visited))
truncated.push(`[...${obj.length - maxLen} more items]`)
return truncated
}
return obj.map((item) => sanitizePayload(item, visited))
}
// Handle objects with key limiting
const entries = Object.entries(obj as Record<string, unknown>)
const maxKeys = PAYLOAD_LIMITS.MAX_OBJECT_KEYS
const result: Record<string, unknown> = {}
let keyCount = 0
for (const [key, value] of entries) {
if (keyCount >= maxKeys) {
result["[...]"] = `${entries.length - maxKeys} more keys omitted`
break
}
if (isSecretKey(key)) {
result[key] = "[REDACTED]"
} else if (typeof value === "string") {
result[key] = sanitizeString(value)
} else if (typeof value === "object" && value !== null) {
result[key] = sanitizePayload(value, visited)
} else {
result[key] = value
}
keyCount++
}
return result
}
/**
* Generate a unique request ID.
*/
function generateRequestId(): string {
return `req_${Date.now()}_${Math.random().toString(36).slice(2, 11)}`
}
/**
* Singleton logger class for API inference logging.
*/
class ApiInferenceLoggerSingleton {
private enabled = false
private sink: ((...args: unknown[]) => void) | null = null
/**
* Configure the logger with enabled state and output sink.
* Should be called once during extension activation.
*/
configure(config: ApiInferenceLoggerConfig): void {
this.enabled = config.enabled
this.sink = config.enabled ? config.sink : null
}
/**
* Check if logging is currently enabled.
*/
isEnabled(): boolean {
return this.enabled && this.sink !== null
}
/**
* Start logging an API inference request.
* Returns a handle to log the response or error.
*
* @param context - Context information about the request
* @param requestPayload - The request payload to log
* @returns A handle with success() and error() methods
*/
start(context: ApiInferenceContext, requestPayload: unknown): ApiInferenceHandle {
const requestId = context.requestId ?? generateRequestId()
const startTime = Date.now()
const startTimestamp = new Date().toISOString()
// Log the request
if (this.isEnabled()) {
this.logRequest({
...context,
requestId,
timestamp: startTimestamp,
payload: requestPayload,
})
}
return {
success: (responsePayload: unknown) => {
if (this.isEnabled()) {
const endTime = Date.now()
this.logResponse({
...context,
requestId,
timestamp: new Date().toISOString(),
durationMs: endTime - startTime,
payload: responsePayload,
})
}
},
error: (errorPayload: unknown) => {
if (this.isEnabled()) {
const endTime = Date.now()
this.logError({
...context,
requestId,
timestamp: new Date().toISOString(),
durationMs: endTime - startTime,
error: errorPayload,
})
}
},
}
}
/**
* Log a request with stable tag format.
*/
private logRequest(data: {
provider: string
operation: string
model?: string
taskId?: string
requestId: string
timestamp: string
payload: unknown
}): void {
if (!this.sink) return
try {
this.sink("[API][request]", {
provider: data.provider,
operation: data.operation,
model: data.model,
taskId: data.taskId,
requestId: data.requestId,
timestamp: data.timestamp,
payload: sanitizePayload(data.payload),
})
} catch {
// Silently ignore logging errors to avoid breaking the application
}
}
/**
* Log a successful response with stable tag format.
*/
private logResponse(data: {
provider: string
operation: string
model?: string
taskId?: string
requestId: string
timestamp: string
durationMs: number
payload: unknown
}): void {
if (!this.sink) return
try {
this.sink("[API][response]", {
provider: data.provider,
operation: data.operation,
model: data.model,
taskId: data.taskId,
requestId: data.requestId,
timestamp: data.timestamp,
durationMs: data.durationMs,
payload: sanitizePayload(data.payload),
})
} catch {
// Silently ignore logging errors to avoid breaking the application
}
}
/**
* Log an error response with stable tag format.
*/
private logError(data: {
provider: string
operation: string
model?: string
taskId?: string
requestId: string
timestamp: string
durationMs: number
error: unknown
}): void {
if (!this.sink) return
try {
// Handle Error objects specially
let errorData: unknown
if (data.error instanceof Error) {
errorData = {
name: data.error.name,
message: data.error.message,
stack: data.error.stack,
}
} else {
errorData = sanitizePayload(data.error)
}
this.sink("[API][error]", {
provider: data.provider,
operation: data.operation,
model: data.model,
taskId: data.taskId,
requestId: data.requestId,
timestamp: data.timestamp,
durationMs: data.durationMs,
error: errorData,
})
} catch {
// Silently ignore logging errors to avoid breaking the application
}
}
}
/**
* Singleton instance of the API inference logger.
*/
export const ApiInferenceLogger = new ApiInferenceLoggerSingleton()

View file

@ -0,0 +1,519 @@
import { ApiInferenceLogger } from "../ApiInferenceLogger"
describe("ApiInferenceLogger", () => {
let mockSink: ReturnType<typeof vi.fn>
beforeEach(() => {
mockSink = vi.fn()
// Reset the logger to disabled state before each test
ApiInferenceLogger.configure({ enabled: false, sink: () => {} })
})
describe("configure", () => {
it("should enable logging when configured with enabled=true", () => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
expect(ApiInferenceLogger.isEnabled()).toBe(true)
})
it("should disable logging when configured with enabled=false", () => {
ApiInferenceLogger.configure({ enabled: false, sink: mockSink })
expect(ApiInferenceLogger.isEnabled()).toBe(false)
})
it("should not log when disabled", () => {
ApiInferenceLogger.configure({ enabled: false, sink: mockSink })
const handle = ApiInferenceLogger.start(
{ provider: "test", operation: "createMessage" },
{ test: "payload" },
)
handle.success({ response: "data" })
expect(mockSink).not.toHaveBeenCalled()
})
})
describe("start", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should emit a request log with [API][request] tag", () => {
ApiInferenceLogger.start({ provider: "OpenAI", operation: "createMessage" }, { model: "gpt-4" })
expect(mockSink).toHaveBeenCalledTimes(1)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
provider: "OpenAI",
operation: "createMessage",
payload: expect.objectContaining({ model: "gpt-4" }),
}),
)
})
it("should include context fields in the log", () => {
ApiInferenceLogger.start(
{
provider: "Anthropic",
operation: "createMessage",
model: "claude-3",
taskId: "task-123",
requestId: "req-456",
},
{ test: "data" },
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
provider: "Anthropic",
operation: "createMessage",
model: "claude-3",
taskId: "task-123",
requestId: "req-456",
}),
)
})
it("should generate a requestId if not provided", () => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
requestId: expect.stringMatching(/^req_\d+_[a-z0-9]+$/),
}),
)
})
})
describe("success", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should emit a response log with [API][response] tag", () => {
const handle = ApiInferenceLogger.start({ provider: "OpenAI", operation: "createMessage" }, {})
mockSink.mockClear()
handle.success({ text: "Hello world", usage: { inputTokens: 10, outputTokens: 20 } })
expect(mockSink).toHaveBeenCalledTimes(1)
expect(mockSink).toHaveBeenCalledWith(
"[API][response]",
expect.objectContaining({
provider: "OpenAI",
operation: "createMessage",
payload: expect.objectContaining({
text: "Hello world",
usage: { inputTokens: 10, outputTokens: 20 },
}),
}),
)
})
it("should include duration in the response log", () => {
const handle = ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
mockSink.mockClear()
// Small delay to ensure measurable duration
handle.success({ response: "data" })
expect(mockSink).toHaveBeenCalledWith(
"[API][response]",
expect.objectContaining({
durationMs: expect.any(Number),
}),
)
})
})
describe("error", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should emit an error log with [API][error] tag", () => {
const handle = ApiInferenceLogger.start({ provider: "OpenAI", operation: "createMessage" }, {})
mockSink.mockClear()
handle.error(new Error("API request failed"))
expect(mockSink).toHaveBeenCalledTimes(1)
expect(mockSink).toHaveBeenCalledWith(
"[API][error]",
expect.objectContaining({
provider: "OpenAI",
operation: "createMessage",
error: expect.objectContaining({
name: "Error",
message: "API request failed",
}),
}),
)
})
it("should include duration in the error log", () => {
const handle = ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
mockSink.mockClear()
handle.error({ code: 500, message: "Internal error" })
expect(mockSink).toHaveBeenCalledWith(
"[API][error]",
expect.objectContaining({
durationMs: expect.any(Number),
}),
)
})
it("should handle non-Error objects", () => {
const handle = ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
mockSink.mockClear()
handle.error({ status: 401, message: "Unauthorized" })
expect(mockSink).toHaveBeenCalledWith(
"[API][error]",
expect.objectContaining({
error: expect.objectContaining({
status: 401,
message: "Unauthorized",
}),
}),
)
})
})
describe("secret redaction", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should redact Authorization header", () => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{ headers: { Authorization: "Bearer sk-secret-key-12345" } },
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
headers: { Authorization: "[REDACTED]" },
}),
}),
)
})
it("should redact apiKey field", () => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { apiKey: "sk-secret-12345" })
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
apiKey: "[REDACTED]",
}),
}),
)
})
it("should redact nested secret fields", () => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{
config: {
auth: {
access_token: "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...",
api_key: "secret-api-key",
},
},
},
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
config: {
auth: {
access_token: "[REDACTED]",
api_key: "[REDACTED]",
},
},
}),
}),
)
})
it("should redact secret fields in arrays", () => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{
items: [{ apiKey: "secret1" }, { apiKey: "secret2" }],
},
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
items: [{ apiKey: "[REDACTED]" }, { apiKey: "[REDACTED]" }],
}),
}),
)
})
it("should not redact non-secret fields", () => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{ model: "gpt-4", messages: [{ role: "user", content: "Hello" }] },
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
model: "gpt-4",
messages: [{ role: "user", content: "Hello" }],
}),
}),
)
})
})
describe("payload size limiting", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should truncate strings longer than 10,000 characters", () => {
const longString = "x".repeat(15_000)
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { content: longString })
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
content: "[Truncated len=15000]",
}),
}),
)
})
it("should not truncate strings within the limit", () => {
const normalString = "x".repeat(5_000)
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { content: normalString })
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
content: normalString,
}),
}),
)
})
it("should replace base64 image data with placeholder", () => {
const imageData =
"data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mNk+M9QDwADhgGAWjR9awAAAABJRU5ErkJggg=="
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { image: imageData })
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
image: expect.stringMatching(/^\[ImageData len=\d+\]$/),
}),
}),
)
})
it("should replace base64 image data for various image types", () => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{
png: "data:image/png;base64,abc123",
jpeg: "data:image/jpeg;base64,abc123",
gif: "data:image/gif;base64,abc123",
webp: "data:image/webp;base64,abc123",
},
)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
png: expect.stringMatching(/^\[ImageData len=\d+\]$/),
jpeg: expect.stringMatching(/^\[ImageData len=\d+\]$/),
gif: expect.stringMatching(/^\[ImageData len=\d+\]$/),
webp: expect.stringMatching(/^\[ImageData len=\d+\]$/),
}),
}),
)
})
it("should cap arrays longer than 200 entries", () => {
const longArray = Array.from({ length: 250 }, (_, i) => ({ id: i }))
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { items: longArray })
const call = mockSink.mock.calls[0]
const payload = call[1].payload as { items: any[] }
expect(payload.items.length).toBe(201)
expect(payload.items[200]).toBe("[...50 more items]")
})
it("should not cap arrays within the limit", () => {
const normalArray = Array.from({ length: 50 }, (_, i) => ({ id: i }))
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { items: normalArray })
const call = mockSink.mock.calls[0]
const payload = call[1].payload as { items: any[] }
expect(payload.items.length).toBe(50)
})
it("should cap objects with more than 200 keys", () => {
const bigObject: Record<string, number> = {}
for (let i = 0; i < 250; i++) {
bigObject[`key${i}`] = i
}
ApiInferenceLogger.start({ provider: "test", operation: "test" }, bigObject)
const call = mockSink.mock.calls[0]
const payload = call[1].payload as Record<string, unknown>
const keys = Object.keys(payload)
expect(keys.length).toBe(201)
expect(payload["[...]"]).toBe("50 more keys omitted")
})
it("should apply size limiting recursively in nested objects", () => {
const nested = {
level1: {
longString: "x".repeat(15_000),
level2: {
imageData: "data:image/png;base64,abc123",
},
},
}
ApiInferenceLogger.start({ provider: "test", operation: "test" }, nested)
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
level1: expect.objectContaining({
longString: "[Truncated len=15000]",
level2: expect.objectContaining({
imageData: expect.stringMatching(/^\[ImageData len=\d+\]$/),
}),
}),
}),
}),
)
})
it("should apply size limiting recursively in arrays", () => {
const messages = [
{ role: "user", content: "x".repeat(15_000) },
{ role: "assistant", content: "data:image/png;base64,abc123" },
]
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { messages })
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
messages: [
expect.objectContaining({
role: "user",
content: "[Truncated len=15000]",
}),
expect.objectContaining({
role: "assistant",
content: expect.stringMatching(/^\[ImageData len=\d+\]$/),
}),
],
}),
}),
)
})
})
describe("edge cases", () => {
beforeEach(() => {
ApiInferenceLogger.configure({ enabled: true, sink: mockSink })
})
it("should handle null and undefined values", () => {
expect(() => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, { value: null, other: undefined })
}).not.toThrow()
expect(mockSink).toHaveBeenCalledWith(
"[API][request]",
expect.objectContaining({
payload: expect.objectContaining({
value: null,
other: undefined,
}),
}),
)
})
it("should handle empty objects", () => {
expect(() => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
}).not.toThrow()
})
it("should not throw on circular references", () => {
const obj: any = { name: "test" }
obj.self = obj
expect(() => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, obj)
}).not.toThrow()
})
it("should handle primitive values in payload", () => {
expect(() => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, "string payload" as any)
}).not.toThrow()
})
it("should handle functions in payload without throwing", () => {
expect(() => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{ callback: () => console.log("test") },
)
}).not.toThrow()
})
it("should handle BigInt values without throwing", () => {
expect(() => {
ApiInferenceLogger.start(
{ provider: "test", operation: "test" },
{ bigValue: BigInt(9007199254740991) },
)
}).not.toThrow()
})
})
describe("sink error handling", () => {
it("should not throw if sink throws an error", () => {
const throwingSink = vi.fn(() => {
throw new Error("Sink error")
})
ApiInferenceLogger.configure({ enabled: true, sink: throwingSink })
expect(() => {
ApiInferenceLogger.start({ provider: "test", operation: "test" }, {})
}).not.toThrow()
})
})
})

View file

@ -30,6 +30,7 @@ import type { SingleCompletionHandler, ApiHandlerCreateMessageMetadata } from ".
// https://docs.anthropic.com/en/api/claude-on-vertex-ai
export class AnthropicVertexHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "Anthropic Vertex"
protected options: ApiHandlerOptions
private client: AnthropicVertex

View file

@ -33,7 +33,7 @@ import {
export class AnthropicHandler extends BaseProvider implements SingleCompletionHandler {
private options: ApiHandlerOptions
private client: Anthropic
private readonly providerName = "Anthropic"
protected readonly providerName = "Anthropic"
constructor(options: ApiHandlerOptions) {
super()
@ -63,6 +63,11 @@ export class AnthropicHandler extends BaseProvider implements SingleCompletionHa
reasoning: thinking,
} = this.getModel()
// Accumulators for final response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
// Filter out non-Anthropic blocks (reasoning, thoughtSignature, etc.) before sending to the API
const sanitizedMessages = filterNonAnthropicBlocks(messages)
@ -93,251 +98,291 @@ export class AnthropicHandler extends BaseProvider implements SingleCompletionHa
}
: {}
switch (modelId) {
case "claude-sonnet-4-5":
case "claude-sonnet-4-20250514":
case "claude-opus-4-5-20251101":
case "claude-opus-4-1-20250805":
case "claude-opus-4-20250514":
case "claude-3-7-sonnet-20250219":
case "claude-3-5-sonnet-20241022":
case "claude-3-5-haiku-20241022":
case "claude-3-opus-20240229":
case "claude-haiku-4-5-20251001":
case "claude-3-haiku-20240307": {
/**
* The latest message will be the new user message, one before
* will be the assistant message from a previous request, and
* the user message before that will be a previously cached user
* message. So we need to mark the latest user message as
* ephemeral to cache it for the next request, and mark the
* second to last user message as ephemeral to let the server
* know the last message to retrieve from the cache for the
* current request.
*/
const userMsgIndices = sanitizedMessages.reduce(
(acc, msg, index) => (msg.role === "user" ? [...acc, index] : acc),
[] as number[],
)
// Variable to hold the log handle - will be created after request body is built
let logHandle: ReturnType<typeof this.inferenceLogger.start>
const lastUserMsgIndex = userMsgIndices[userMsgIndices.length - 1] ?? -1
const secondLastMsgUserIndex = userMsgIndices[userMsgIndices.length - 2] ?? -1
try {
stream = await this.client.messages.create(
{
model: modelId,
max_tokens: maxTokens ?? ANTHROPIC_DEFAULT_MAX_TOKENS,
temperature,
thinking,
// Setting cache breakpoint for system prompt so new tasks can reuse it.
system: [{ text: systemPrompt, type: "text", cache_control: cacheControl }],
messages: sanitizedMessages.map((message, index) => {
if (index === lastUserMsgIndex || index === secondLastMsgUserIndex) {
return {
...message,
content:
typeof message.content === "string"
? [{ type: "text", text: message.content, cache_control: cacheControl }]
: message.content.map((content, contentIndex) =>
contentIndex === message.content.length - 1
? { ...content, cache_control: cacheControl }
: content,
),
}
}
return message
}),
stream: true,
...nativeToolParams,
},
(() => {
// prompt caching: https://x.com/alexalbert__/status/1823751995901272068
// https://github.com/anthropics/anthropic-sdk-typescript?tab=readme-ov-file#default-headers
// https://github.com/anthropics/anthropic-sdk-typescript/commit/c920b77fc67bd839bfeb6716ceab9d7c9bbe7393
// Then check for models that support prompt caching
switch (modelId) {
case "claude-sonnet-4-5":
case "claude-sonnet-4-20250514":
case "claude-opus-4-5-20251101":
case "claude-opus-4-1-20250805":
case "claude-opus-4-20250514":
case "claude-3-7-sonnet-20250219":
case "claude-3-5-sonnet-20241022":
case "claude-3-5-haiku-20241022":
case "claude-3-opus-20240229":
case "claude-haiku-4-5-20251001":
case "claude-3-haiku-20240307":
betas.push("prompt-caching-2024-07-31")
return { headers: { "anthropic-beta": betas.join(",") } }
default:
return undefined
}
})(),
try {
switch (modelId) {
case "claude-sonnet-4-5":
case "claude-sonnet-4-20250514":
case "claude-opus-4-5-20251101":
case "claude-opus-4-1-20250805":
case "claude-opus-4-20250514":
case "claude-3-7-sonnet-20250219":
case "claude-3-5-sonnet-20241022":
case "claude-3-5-haiku-20241022":
case "claude-3-opus-20240229":
case "claude-haiku-4-5-20251001":
case "claude-3-haiku-20240307": {
/**
* The latest message will be the new user message, one before
* will be the assistant message from a previous request, and
* the user message before that will be a previously cached user
* message. So we need to mark the latest user message as
* ephemeral to cache it for the next request, and mark the
* second to last user message as ephemeral to let the server
* know the last message to retrieve from the cache for the
* current request.
*/
const userMsgIndices = sanitizedMessages.reduce(
(acc, msg, index) => (msg.role === "user" ? [...acc, index] : acc),
[] as number[],
)
} catch (error) {
TelemetryService.instance.captureException(
new ApiProviderError(
error instanceof Error ? error.message : String(error),
this.providerName,
modelId,
"createMessage",
),
)
throw error
}
break
}
default: {
try {
stream = (await this.client.messages.create({
const lastUserMsgIndex = userMsgIndices[userMsgIndices.length - 1] ?? -1
const secondLastMsgUserIndex = userMsgIndices[userMsgIndices.length - 2] ?? -1
// Build the request body for logging and API call
const requestBody = {
model: modelId,
max_tokens: maxTokens ?? ANTHROPIC_DEFAULT_MAX_TOKENS,
temperature,
system: [{ text: systemPrompt, type: "text" }],
messages: sanitizedMessages,
stream: true,
thinking,
// Setting cache breakpoint for system prompt so new tasks can reuse it.
system: [{ text: systemPrompt, type: "text" as const, cache_control: cacheControl }],
messages: sanitizedMessages.map((message, index) => {
if (index === lastUserMsgIndex || index === secondLastMsgUserIndex) {
return {
...message,
content:
typeof message.content === "string"
? [
{
type: "text" as const,
text: message.content,
cache_control: cacheControl,
},
]
: message.content.map((content, contentIndex) =>
contentIndex === message.content.length - 1
? { ...content, cache_control: cacheControl }
: content,
),
}
}
return message
}),
stream: true as const,
...nativeToolParams,
})) as any
} catch (error) {
TelemetryService.instance.captureException(
new ApiProviderError(
error instanceof Error ? error.message : String(error),
this.providerName,
modelId,
"createMessage",
),
}
// Start inference logging with actual request body
logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestBody,
)
throw error
}
break
}
}
let inputTokens = 0
let outputTokens = 0
let cacheWriteTokens = 0
let cacheReadTokens = 0
try {
// Determine request options (beta headers)
betas.push("prompt-caching-2024-07-31")
const requestOptions = { headers: { "anthropic-beta": betas.join(",") } }
for await (const chunk of stream) {
switch (chunk.type) {
case "message_start": {
// Tells us cache reads/writes/input/output.
const {
input_tokens = 0,
output_tokens = 0,
cache_creation_input_tokens,
cache_read_input_tokens,
} = chunk.message.usage
yield {
type: "usage",
inputTokens: input_tokens,
outputTokens: output_tokens,
cacheWriteTokens: cache_creation_input_tokens || undefined,
cacheReadTokens: cache_read_input_tokens || undefined,
stream = await this.client.messages.create(requestBody, requestOptions)
} catch (error) {
TelemetryService.instance.captureException(
new ApiProviderError(
error instanceof Error ? error.message : String(error),
this.providerName,
modelId,
"createMessage",
),
)
throw error
}
inputTokens += input_tokens
outputTokens += output_tokens
cacheWriteTokens += cache_creation_input_tokens || 0
cacheReadTokens += cache_read_input_tokens || 0
break
}
case "message_delta":
// Tells us stop_reason, stop_sequence, and output tokens
// along the way and at the end of the message.
yield {
type: "usage",
inputTokens: 0,
outputTokens: chunk.usage.output_tokens || 0,
default: {
// Build the request body for logging and API call
const requestBody = {
model: modelId,
max_tokens: maxTokens ?? ANTHROPIC_DEFAULT_MAX_TOKENS,
temperature,
system: [{ text: systemPrompt, type: "text" as const }],
messages: sanitizedMessages,
stream: true as const,
...nativeToolParams,
}
break
case "message_stop":
// No usage data, just an indicator that the message is done.
break
case "content_block_start":
switch (chunk.content_block.type) {
case "thinking":
// We may receive multiple text blocks, in which
// case just insert a line break between them.
if (chunk.index > 0) {
yield { type: "reasoning", text: "\n" }
}
// Start inference logging with actual request body
logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestBody,
)
yield { type: "reasoning", text: chunk.content_block.thinking }
break
case "text":
// We may receive multiple text blocks, in which
// case just insert a line break between them.
if (chunk.index > 0) {
yield { type: "text", text: "\n" }
}
yield { type: "text", text: chunk.content_block.text }
break
case "tool_use": {
// Emit initial tool call partial with id and name
yield {
type: "tool_call_partial",
index: chunk.index,
id: chunk.content_block.id,
name: chunk.content_block.name,
arguments: undefined,
}
break
}
try {
stream = (await this.client.messages.create(requestBody)) as any
} catch (error) {
TelemetryService.instance.captureException(
new ApiProviderError(
error instanceof Error ? error.message : String(error),
this.providerName,
modelId,
"createMessage",
),
)
throw error
}
break
case "content_block_delta":
switch (chunk.delta.type) {
case "thinking_delta":
yield { type: "reasoning", text: chunk.delta.thinking }
break
case "text_delta":
yield { type: "text", text: chunk.delta.text }
break
case "input_json_delta": {
// Emit tool call partial chunks as arguments stream in
yield {
type: "tool_call_partial",
index: chunk.index,
id: undefined,
name: undefined,
arguments: chunk.delta.partial_json,
}
break
}
}
break
case "content_block_stop":
// Block complete - no action needed for now.
// NativeToolCallParser handles tool call completion
// Note: Signature for multi-turn thinking would require using stream.finalMessage()
// after iteration completes, which requires restructuring the streaming approach.
break
}
}
}
if (inputTokens > 0 || outputTokens > 0 || cacheWriteTokens > 0 || cacheReadTokens > 0) {
const { totalCost } = calculateApiCostAnthropic(
this.getModel().info,
inputTokens,
outputTokens,
cacheWriteTokens,
cacheReadTokens,
)
let inputTokens = 0
let outputTokens = 0
let cacheWriteTokens = 0
let cacheReadTokens = 0
yield {
type: "usage",
inputTokens: 0,
outputTokens: 0,
totalCost,
for await (const chunk of stream) {
switch (chunk.type) {
case "message_start": {
// Tells us cache reads/writes/input/output.
const {
input_tokens = 0,
output_tokens = 0,
cache_creation_input_tokens,
cache_read_input_tokens,
} = chunk.message.usage
yield {
type: "usage",
inputTokens: input_tokens,
outputTokens: output_tokens,
cacheWriteTokens: cache_creation_input_tokens || undefined,
cacheReadTokens: cache_read_input_tokens || undefined,
}
inputTokens += input_tokens
outputTokens += output_tokens
cacheWriteTokens += cache_creation_input_tokens || 0
cacheReadTokens += cache_read_input_tokens || 0
break
}
case "message_delta":
// Tells us stop_reason, stop_sequence, and output tokens
// along the way and at the end of the message.
yield {
type: "usage",
inputTokens: 0,
outputTokens: chunk.usage.output_tokens || 0,
}
break
case "message_stop":
// No usage data, just an indicator that the message is done.
break
case "content_block_start":
switch (chunk.content_block.type) {
case "thinking":
// We may receive multiple text blocks, in which
// case just insert a line break between them.
if (chunk.index > 0) {
accumulatedReasoning.push("\n")
yield { type: "reasoning", text: "\n" }
}
accumulatedReasoning.push(chunk.content_block.thinking)
yield { type: "reasoning", text: chunk.content_block.thinking }
break
case "text":
// We may receive multiple text blocks, in which
// case just insert a line break between them.
if (chunk.index > 0) {
accumulatedText.push("\n")
yield { type: "text", text: "\n" }
}
accumulatedText.push(chunk.content_block.text)
yield { type: "text", text: chunk.content_block.text }
break
case "tool_use": {
// Track tool call for logging
toolCalls.push({
id: chunk.content_block.id,
name: chunk.content_block.name,
})
// Emit initial tool call partial with id and name
yield {
type: "tool_call_partial",
index: chunk.index,
id: chunk.content_block.id,
name: chunk.content_block.name,
arguments: undefined,
}
break
}
}
break
case "content_block_delta":
switch (chunk.delta.type) {
case "thinking_delta":
accumulatedReasoning.push(chunk.delta.thinking)
yield { type: "reasoning", text: chunk.delta.thinking }
break
case "text_delta":
accumulatedText.push(chunk.delta.text)
yield { type: "text", text: chunk.delta.text }
break
case "input_json_delta": {
// Emit tool call partial chunks as arguments stream in
yield {
type: "tool_call_partial",
index: chunk.index,
id: undefined,
name: undefined,
arguments: chunk.delta.partial_json,
}
break
}
}
break
case "content_block_stop":
// Block complete - no action needed for now.
// NativeToolCallParser handles tool call completion
// Note: Signature for multi-turn thinking would require using stream.finalMessage()
// after iteration completes, which requires restructuring the streaming approach.
break
}
}
if (inputTokens > 0 || outputTokens > 0 || cacheWriteTokens > 0 || cacheReadTokens > 0) {
const { totalCost } = calculateApiCostAnthropic(
this.getModel().info,
inputTokens,
outputTokens,
cacheWriteTokens,
cacheReadTokens,
)
yield {
type: "usage",
inputTokens: 0,
outputTokens: 0,
totalCost,
}
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: { inputTokens, outputTokens, cacheWriteTokens, cacheReadTokens },
})
} catch (error) {
// logHandle may not be assigned if error occurs before request body is built
if (logHandle!) {
logHandle.error(error)
}
throw error
}
}

View file

@ -117,87 +117,167 @@ export abstract class BaseOpenAiCompatibleProvider<ModelName extends string>
messages: Anthropic.Messages.MessageParam[],
metadata?: ApiHandlerCreateMessageMetadata,
): ApiStream {
const stream = await this.createStream(systemPrompt, messages, metadata)
const { id: model, info } = this.getModel()
const matcher = new XmlMatcher(
"think",
(chunk) =>
({
type: chunk.matched ? "reasoning" : "text",
text: chunk.data,
}) as const,
// Build the actual params object that will be passed to the SDK
const max_tokens =
getModelMaxOutputTokens({
modelId: model,
model: info,
settings: this.options,
format: "openai",
}) ?? undefined
const temperature = this.options.modelTemperature ?? this.defaultTemperature
const requestParams: OpenAI.Chat.Completions.ChatCompletionCreateParamsStreaming = {
model,
max_tokens,
temperature,
messages: [{ role: "system", content: systemPrompt }, ...convertToOpenAiMessages(messages)],
stream: true,
stream_options: { include_usage: true },
...(metadata?.tools && { tools: this.convertToolsForOpenAI(metadata.tools) }),
...(metadata?.tool_choice && { tool_choice: metadata.tool_choice }),
...(metadata?.toolProtocol === "native" && {
parallel_tool_calls: metadata.parallelToolCalls ?? false,
}),
}
// Add thinking parameter if reasoning is enabled and model supports it
if (this.options.enableReasoningEffort && info.supportsReasoningBinary) {
;(requestParams as any).thinking = { type: "enabled" }
}
// Start inference logging with the actual request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model,
},
requestParams,
)
// Accumulators for final response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string; arguments?: string }> = []
let lastUsage: OpenAI.CompletionUsage | undefined
const activeToolCallIds = new Set<string>()
for await (const chunk of stream) {
// Check for provider-specific error responses (e.g., MiniMax base_resp)
const chunkAny = chunk as any
if (chunkAny.base_resp?.status_code && chunkAny.base_resp.status_code !== 0) {
throw new Error(
`${this.providerName} API Error (${chunkAny.base_resp.status_code}): ${chunkAny.base_resp.status_msg || "Unknown error"}`,
)
}
try {
const stream = await this.createStream(systemPrompt, messages, metadata)
const delta = chunk.choices?.[0]?.delta
const finishReason = chunk.choices?.[0]?.finish_reason
const matcher = new XmlMatcher(
"think",
(chunk) =>
({
type: chunk.matched ? "reasoning" : "text",
text: chunk.data,
}) as const,
)
if (delta?.content) {
for (const processedChunk of matcher.update(delta.content)) {
yield processedChunk
for await (const chunk of stream) {
// Check for provider-specific error responses (e.g., MiniMax base_resp)
const chunkAny = chunk as any
if (chunkAny.base_resp?.status_code && chunkAny.base_resp.status_code !== 0) {
const error = new Error(
`${this.providerName} API Error (${chunkAny.base_resp.status_code}): ${chunkAny.base_resp.status_msg || "Unknown error"}`,
)
logHandle.error(error)
throw error
}
}
if (delta) {
for (const key of ["reasoning_content", "reasoning"] as const) {
if (key in delta) {
const reasoning_content = ((delta as any)[key] as string | undefined) || ""
if (reasoning_content?.trim()) {
yield { type: "reasoning", text: reasoning_content }
const delta = chunk.choices?.[0]?.delta
const finishReason = chunk.choices?.[0]?.finish_reason
if (delta?.content) {
for (const processedChunk of matcher.update(delta.content)) {
if (processedChunk.type === "text") {
accumulatedText.push(processedChunk.text)
} else if (processedChunk.type === "reasoning") {
accumulatedReasoning.push(processedChunk.text)
}
break
yield processedChunk
}
}
if (delta) {
for (const key of ["reasoning_content", "reasoning"] as const) {
if (key in delta) {
const reasoning_content = ((delta as any)[key] as string | undefined) || ""
if (reasoning_content?.trim()) {
accumulatedReasoning.push(reasoning_content)
yield { type: "reasoning", text: reasoning_content }
}
break
}
}
}
// Emit raw tool call chunks - NativeToolCallParser handles state management
if (delta?.tool_calls) {
for (const toolCall of delta.tool_calls) {
// Track tool calls for logging
if (toolCall.id) {
activeToolCallIds.add(toolCall.id)
}
if (toolCall.id || toolCall.function?.name) {
toolCalls.push({
id: toolCall.id,
name: toolCall.function?.name,
arguments: toolCall.function?.arguments,
})
}
yield {
type: "tool_call_partial",
index: toolCall.index,
id: toolCall.id,
name: toolCall.function?.name,
arguments: toolCall.function?.arguments,
}
}
}
// Emit tool_call_end events when finish_reason is "tool_calls"
// This ensures tool calls are finalized even if the stream doesn't properly close
if (finishReason === "tool_calls" && activeToolCallIds.size > 0) {
for (const id of activeToolCallIds) {
yield { type: "tool_call_end", id }
}
activeToolCallIds.clear()
}
if (chunk.usage) {
lastUsage = chunk.usage
}
}
// Emit raw tool call chunks - NativeToolCallParser handles state management
if (delta?.tool_calls) {
for (const toolCall of delta.tool_calls) {
if (toolCall.id) {
activeToolCallIds.add(toolCall.id)
}
yield {
type: "tool_call_partial",
index: toolCall.index,
id: toolCall.id,
name: toolCall.function?.name,
arguments: toolCall.function?.arguments,
}
if (lastUsage) {
yield this.processUsageMetrics(lastUsage, this.getModel().info)
}
// Process any remaining content
for (const processedChunk of matcher.final()) {
if (processedChunk.type === "text") {
accumulatedText.push(processedChunk.text)
} else if (processedChunk.type === "reasoning") {
accumulatedReasoning.push(processedChunk.text)
}
yield processedChunk
}
// Emit tool_call_end events when finish_reason is "tool_calls"
// This ensures tool calls are finalized even if the stream doesn't properly close
if (finishReason === "tool_calls" && activeToolCallIds.size > 0) {
for (const id of activeToolCallIds) {
yield { type: "tool_call_end", id }
}
activeToolCallIds.clear()
}
if (chunk.usage) {
lastUsage = chunk.usage
}
}
if (lastUsage) {
yield this.processUsageMetrics(lastUsage, this.getModel().info)
}
// Process any remaining content
for (const processedChunk of matcher.final()) {
yield processedChunk
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
}

View file

@ -6,11 +6,23 @@ import type { ApiHandler, ApiHandlerCreateMessageMetadata } from "../index"
import { ApiStream } from "../transform/stream"
import { countTokens } from "../../utils/countTokens"
import { isMcpTool } from "../../utils/mcp-name"
import { ApiInferenceLogger } from "../logging/ApiInferenceLogger"
/**
* Base class for API providers that implements common functionality.
*/
export abstract class BaseProvider implements ApiHandler {
/**
* The name of this provider, used for logging and error reporting.
*/
protected abstract readonly providerName: string
/**
* Reference to the API inference logger singleton for logging requests/responses.
* Providers can use this to log inference calls when enabled.
*/
protected readonly inferenceLogger = ApiInferenceLogger
abstract createMessage(
systemPrompt: string,
messages: Anthropic.Messages.MessageParam[],

View file

@ -200,7 +200,7 @@ export class AwsBedrockHandler extends BaseProvider implements SingleCompletionH
protected options: ProviderSettings
private client: BedrockRuntimeClient
private arnInfo: any
private readonly providerName = "Bedrock"
protected readonly providerName = "Bedrock"
constructor(options: ProviderSettings) {
super()

View file

@ -20,6 +20,7 @@ const CEREBRAS_INTEGRATION_HEADER = "X-Cerebras-3rd-Party-Integration"
const CEREBRAS_INTEGRATION_NAME = "roocode"
export class CerebrasHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "Cerebras"
private apiKey: string
private providerModels: typeof cerebrasModels
private defaultProviderModelId: CerebrasModelId

View file

@ -20,6 +20,7 @@ import { t } from "../../i18n"
import { ApiHandlerOptions } from "../../shared/api"
import { countTokens } from "../../utils/countTokens"
import { convertOpenAIToolsToAnthropic } from "../../core/prompts/tools/native-tools/converters"
import { BaseProvider } from "./base-provider"
/**
* Converts OpenAI tool_choice to Anthropic ToolChoice format
@ -64,7 +65,8 @@ function convertOpenAIToolChoice(
return { type: "auto", disable_parallel_tool_use: disableParallelToolUse }
}
export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
export class ClaudeCodeHandler extends BaseProvider implements ApiHandler, SingleCompletionHandler {
protected readonly providerName = "Claude Code"
private options: ApiHandlerOptions
/**
* Store the last thinking block signature for interleaved thinking with tool use.
@ -75,6 +77,7 @@ export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
private lastThinkingSignature?: string
constructor(options: ApiHandlerOptions) {
super()
this.options = options
}
@ -114,7 +117,7 @@ export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
return null
}
async *createMessage(
override async *createMessage(
systemPrompt: string,
messages: Anthropic.Messages.MessageParam[],
metadata?: ApiHandlerCreateMessageMetadata,
@ -182,9 +185,8 @@ export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
thinking = { type: "disabled" }
}
// Create streaming request using OAuth
const stream = createStreamingMessage({
accessToken,
// Build request params for logging
const requestParams = {
model: modelId,
systemPrompt,
messages,
@ -195,79 +197,122 @@ export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
metadata: {
user_id: userId,
},
})
}
// Track usage for cost calculation
let inputTokens = 0
let outputTokens = 0
let cacheReadTokens = 0
let cacheWriteTokens = 0
// Start inference logging with request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestParams,
)
for await (const chunk of stream) {
switch (chunk.type) {
case "text":
yield {
type: "text",
text: chunk.text,
}
break
// Accumulators for response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
case "reasoning":
yield {
type: "reasoning",
text: chunk.text,
}
break
try {
// Create streaming request using OAuth
const stream = createStreamingMessage({
accessToken,
...requestParams,
})
case "thinking_complete":
// Capture the signature for persistence in api_conversation_history
// This enables tool use continuations where thinking blocks must be passed back
if (chunk.signature) {
this.lastThinkingSignature = chunk.signature
}
// Emit a complete thinking block with signature
// This is critical for interleaved thinking with tool use
// The signature must be included when passing thinking blocks back to the API
yield {
type: "reasoning",
text: chunk.thinking,
signature: chunk.signature,
}
break
// Track usage for cost calculation
let inputTokens = 0
let outputTokens = 0
let cacheReadTokens = 0
let cacheWriteTokens = 0
let lastUsage: any
case "tool_call_partial":
yield {
type: "tool_call_partial",
index: chunk.index,
id: chunk.id,
name: chunk.name,
arguments: chunk.arguments,
}
break
for await (const chunk of stream) {
switch (chunk.type) {
case "text":
accumulatedText.push(chunk.text)
yield {
type: "text",
text: chunk.text,
}
break
case "usage": {
inputTokens = chunk.inputTokens
outputTokens = chunk.outputTokens
cacheReadTokens = chunk.cacheReadTokens || 0
cacheWriteTokens = chunk.cacheWriteTokens || 0
case "reasoning":
accumulatedReasoning.push(chunk.text)
yield {
type: "reasoning",
text: chunk.text,
}
break
// Claude Code is subscription-based, no per-token cost
const usageChunk: ApiStreamUsageChunk = {
type: "usage",
inputTokens,
outputTokens,
cacheReadTokens: cacheReadTokens > 0 ? cacheReadTokens : undefined,
cacheWriteTokens: cacheWriteTokens > 0 ? cacheWriteTokens : undefined,
totalCost: 0,
case "thinking_complete":
// Capture the signature for persistence in api_conversation_history
// This enables tool use continuations where thinking blocks must be passed back
if (chunk.signature) {
this.lastThinkingSignature = chunk.signature
}
accumulatedReasoning.push(chunk.thinking)
// Emit a complete thinking block with signature
// This is critical for interleaved thinking with tool use
// The signature must be included when passing thinking blocks back to the API
yield {
type: "reasoning",
text: chunk.thinking,
signature: chunk.signature,
}
break
case "tool_call_partial":
if (chunk.id || chunk.name) {
toolCalls.push({ id: chunk.id, name: chunk.name })
}
yield {
type: "tool_call_partial",
index: chunk.index,
id: chunk.id,
name: chunk.name,
arguments: chunk.arguments,
}
break
case "usage": {
inputTokens = chunk.inputTokens
outputTokens = chunk.outputTokens
cacheReadTokens = chunk.cacheReadTokens || 0
cacheWriteTokens = chunk.cacheWriteTokens || 0
// Claude Code is subscription-based, no per-token cost
const usageChunk: ApiStreamUsageChunk = {
type: "usage",
inputTokens,
outputTokens,
cacheReadTokens: cacheReadTokens > 0 ? cacheReadTokens : undefined,
cacheWriteTokens: cacheWriteTokens > 0 ? cacheWriteTokens : undefined,
totalCost: 0,
}
lastUsage = usageChunk
yield usageChunk
break
}
yield usageChunk
break
case "error":
throw new Error(chunk.error)
}
case "error":
throw new Error(chunk.error)
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
}
@ -284,7 +329,7 @@ export class ClaudeCodeHandler implements ApiHandler, SingleCompletionHandler {
}
}
async countTokens(content: Anthropic.Messages.ContentBlockParam[]): Promise<number> {
override async countTokens(content: Anthropic.Messages.ContentBlockParam[]): Promise<number> {
if (content.length === 0) {
return 0
}

View file

@ -40,7 +40,7 @@ export class GeminiHandler extends BaseProvider implements SingleCompletionHandl
private client: GoogleGenAI
private lastThoughtSignature?: string
private lastResponseId?: string
private readonly providerName = "Gemini"
protected readonly providerName = "Gemini"
constructor({ isVertex, ...options }: GeminiHandlerOptions) {
super()

View file

@ -14,7 +14,7 @@ export class HuggingFaceHandler extends BaseProvider implements SingleCompletion
private client: OpenAI
private options: ApiHandlerOptions
private modelCache: ModelRecord | null = null
private readonly providerName = "HuggingFace"
protected readonly providerName = "HuggingFace"
constructor(options: ApiHandlerOptions) {
super()

View file

@ -21,7 +21,7 @@ import { handleOpenAIError } from "./utils/openai-error-handler"
export class LmStudioHandler extends BaseProvider implements SingleCompletionHandler {
protected options: ApiHandlerOptions
private client: OpenAI
private readonly providerName = "LM Studio"
protected readonly providerName = "LM Studio"
constructor(options: ApiHandlerOptions) {
super()

View file

@ -51,6 +51,7 @@ function convertOpenAIToolChoice(
}
export class MiniMaxHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "MiniMax"
private options: ApiHandlerOptions
private client: Anthropic

View file

@ -51,7 +51,7 @@ type MistralTool = {
export class MistralHandler extends BaseProvider implements SingleCompletionHandler {
protected options: ApiHandlerOptions
private client: Mistral
private readonly providerName = "Mistral"
protected readonly providerName = "Mistral"
constructor(options: ApiHandlerOptions) {
super()

View file

@ -146,6 +146,7 @@ function convertToOllamaMessages(anthropicMessages: Anthropic.Messages.MessagePa
}
export class NativeOllamaHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "Ollama"
protected options: ApiHandlerOptions
private client: Ollama | undefined
protected models: Record<string, ModelInfo> = {}

View file

@ -31,7 +31,7 @@ export type OpenAiNativeModel = ReturnType<OpenAiNativeHandler["getModel"]>
export class OpenAiNativeHandler extends BaseProvider implements SingleCompletionHandler {
protected options: ApiHandlerOptions
private client: OpenAI
private readonly providerName = "OpenAI Native"
protected readonly providerName = "OpenAI Native"
// Resolved service tier from Responses API (actual tier used by OpenAI)
private lastServiceTier: ServiceTier | undefined
// Complete response output array (includes reasoning items with encrypted_content)
@ -175,8 +175,54 @@ export class OpenAiNativeHandler extends BaseProvider implements SingleCompletio
metadata,
)
// Make the request (pass systemPrompt and messages for potential retry)
yield* this.executeRequest(requestBody, model, metadata, systemPrompt, messages)
// Start inference logging with actual request body
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: model.id,
},
requestBody,
)
// Accumulators for response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
let lastUsage: any
try {
// Make the request (pass systemPrompt and messages for potential retry)
for await (const chunk of this.executeRequest(requestBody, model, metadata, systemPrompt, messages)) {
// Accumulate for logging
if (chunk.type === "text") {
accumulatedText.push(chunk.text)
} else if (chunk.type === "reasoning") {
accumulatedReasoning.push(chunk.text)
} else if (chunk.type === "tool_call" || chunk.type === "tool_call_partial") {
if (chunk.id || chunk.name) {
toolCalls.push({ id: chunk.id, name: chunk.name })
}
} else if (chunk.type === "usage") {
lastUsage = chunk
}
yield chunk
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
responseId: this.lastResponseId,
responseOutput: this.lastResponseOutput,
})
} catch (error) {
logHandle.error(error)
throw error
}
}
private buildRequestBody(

View file

@ -33,7 +33,7 @@ import { handleOpenAIError } from "./utils/openai-error-handler"
export class OpenAiHandler extends BaseProvider implements SingleCompletionHandler {
protected options: ApiHandlerOptions
protected client: OpenAI
private readonly providerName = "OpenAI"
protected readonly providerName = "OpenAI"
constructor(options: ApiHandlerOptions) {
super()
@ -95,6 +95,13 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
const deepseekReasoner = modelId.includes("deepseek-reasoner") || enabledR1Format
const ark = modelUrl.includes(".volces.com")
// Accumulators for final response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
let lastUsage: any
// Handle O3 family models separately with their own logging
if (modelId.includes("o1") || modelId.includes("o3") || modelId.includes("o4")) {
yield* this.handleO3FamilyMessage(modelId, systemPrompt, messages, metadata)
return
@ -175,60 +182,102 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
// Add max_tokens if needed
this.addMaxTokensIfNeeded(requestOptions, modelInfo)
let stream
try {
stream = await this.client.chat.completions.create(
requestOptions,
isAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
const matcher = new XmlMatcher(
"think",
(chunk) =>
({
type: chunk.matched ? "reasoning" : "text",
text: chunk.data,
}) as const,
// Start inference logging with actual request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestOptions,
)
let lastUsage
const activeToolCallIds = new Set<string>()
try {
let stream
try {
stream = await this.client.chat.completions.create(
requestOptions,
isAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
for await (const chunk of stream) {
const delta = chunk.choices?.[0]?.delta ?? {}
const finishReason = chunk.choices?.[0]?.finish_reason
const matcher = new XmlMatcher(
"think",
(chunk) =>
({
type: chunk.matched ? "reasoning" : "text",
text: chunk.data,
}) as const,
)
if (delta.content) {
for (const chunk of matcher.update(delta.content)) {
yield chunk
const activeToolCallIds = new Set<string>()
for await (const chunk of stream) {
const delta = chunk.choices?.[0]?.delta ?? {}
const finishReason = chunk.choices?.[0]?.finish_reason
if (delta.content) {
for (const matchedChunk of matcher.update(delta.content)) {
if (matchedChunk.type === "text") {
accumulatedText.push(matchedChunk.text)
} else if (matchedChunk.type === "reasoning") {
accumulatedReasoning.push(matchedChunk.text)
}
yield matchedChunk
}
}
if ("reasoning_content" in delta && delta.reasoning_content) {
accumulatedReasoning.push((delta.reasoning_content as string | undefined) || "")
yield {
type: "reasoning",
text: (delta.reasoning_content as string | undefined) || "",
}
}
// Track tool calls for logging and use processToolCalls for proper tool_call_end events
if (delta.tool_calls) {
for (const toolCall of delta.tool_calls) {
if (toolCall.id || toolCall.function?.name) {
toolCalls.push({ id: toolCall.id, name: toolCall.function?.name })
}
}
}
yield* this.processToolCalls(delta, finishReason, activeToolCallIds)
if (chunk.usage) {
lastUsage = chunk.usage
}
}
if ("reasoning_content" in delta && delta.reasoning_content) {
yield {
type: "reasoning",
text: (delta.reasoning_content as string | undefined) || "",
for (const matchedChunk of matcher.final()) {
if (matchedChunk.type === "text") {
accumulatedText.push(matchedChunk.text)
} else if (matchedChunk.type === "reasoning") {
accumulatedReasoning.push(matchedChunk.text)
}
yield matchedChunk
}
yield* this.processToolCalls(delta, finishReason, activeToolCallIds)
if (chunk.usage) {
lastUsage = chunk.usage
if (lastUsage) {
yield this.processUsageMetrics(lastUsage, modelInfo)
}
}
for (const chunk of matcher.final()) {
yield chunk
}
if (lastUsage) {
yield this.processUsageMetrics(lastUsage, modelInfo)
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
} else {
// Non-streaming path
const requestOptions: OpenAI.Chat.Completions.ChatCompletionCreateParamsNonStreaming = {
model: modelId,
messages: deepseekReasoner
@ -246,37 +295,63 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
// Add max_tokens if needed
this.addMaxTokensIfNeeded(requestOptions, modelInfo)
let response
// Start inference logging with actual request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestOptions,
)
try {
response = await this.client.chat.completions.create(
requestOptions,
this._isAzureAiInference(modelUrl) ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
let response
try {
response = await this.client.chat.completions.create(
requestOptions,
this._isAzureAiInference(modelUrl) ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
const message = response.choices?.[0]?.message
const message = response.choices?.[0]?.message
if (message?.tool_calls) {
for (const toolCall of message.tool_calls) {
if (toolCall.type === "function") {
yield {
type: "tool_call",
id: toolCall.id,
name: toolCall.function.name,
arguments: toolCall.function.arguments,
if (message?.tool_calls) {
for (const toolCall of message.tool_calls) {
if (toolCall.type === "function") {
toolCalls.push({ id: toolCall.id, name: toolCall.function.name })
yield {
type: "tool_call",
id: toolCall.id,
name: toolCall.function.name,
arguments: toolCall.function.arguments,
}
}
}
}
}
yield {
type: "text",
text: message?.content || "",
}
accumulatedText.push(message?.content || "")
yield {
type: "text",
text: message?.content || "",
}
yield this.processUsageMetrics(response.usage, modelInfo)
lastUsage = response.usage
yield this.processUsageMetrics(response.usage, modelInfo)
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
}
}
@ -346,6 +421,11 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
const modelInfo = this.getModel().info
const methodIsAzureAiInference = this._isAzureAiInference(this.options.openAiBaseUrl)
// Accumulators for response logging
const accumulatedText: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
let lastUsage: any
if (this.options.openAiStreamingEnabled ?? true) {
const isGrokXAI = this._isGrokXAI(this.options.openAiBaseUrl)
@ -374,17 +454,73 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
// This allows O3 models to limit response length when includeMaxTokens is enabled
this.addMaxTokensIfNeeded(requestOptions, modelInfo)
let stream
try {
stream = await this.client.chat.completions.create(
requestOptions,
methodIsAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
// Start inference logging with actual request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestOptions,
)
yield* this.handleStreamResponse(stream)
try {
let stream
try {
stream = await this.client.chat.completions.create(
requestOptions,
methodIsAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
const activeToolCallIds = new Set<string>()
for await (const chunk of stream) {
const delta = chunk.choices?.[0]?.delta
const finishReason = chunk.choices?.[0]?.finish_reason
if (delta) {
if (delta.content) {
accumulatedText.push(delta.content)
yield {
type: "text",
text: delta.content,
}
}
// Track tool calls for logging and use processToolCalls for proper tool_call_end events
if (delta.tool_calls) {
for (const toolCall of delta.tool_calls) {
if (toolCall.id || toolCall.function?.name) {
toolCalls.push({ id: toolCall.id, name: toolCall.function?.name })
}
}
}
yield* this.processToolCalls(delta, finishReason, activeToolCallIds)
}
if (chunk.usage) {
lastUsage = chunk.usage
yield {
type: "usage",
inputTokens: chunk.usage.prompt_tokens || 0,
outputTokens: chunk.usage.completion_tokens || 0,
}
}
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
} else {
const requestOptions: OpenAI.Chat.Completions.ChatCompletionCreateParamsNonStreaming = {
model: modelId,
@ -409,35 +545,61 @@ export class OpenAiHandler extends BaseProvider implements SingleCompletionHandl
// This allows O3 models to limit response length when includeMaxTokens is enabled
this.addMaxTokensIfNeeded(requestOptions, modelInfo)
let response
try {
response = await this.client.chat.completions.create(
requestOptions,
methodIsAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
// Start inference logging with actual request params
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
requestOptions,
)
const message = response.choices?.[0]?.message
if (message?.tool_calls) {
for (const toolCall of message.tool_calls) {
if (toolCall.type === "function") {
yield {
type: "tool_call",
id: toolCall.id,
name: toolCall.function.name,
arguments: toolCall.function.arguments,
try {
let response
try {
response = await this.client.chat.completions.create(
requestOptions,
methodIsAzureAiInference ? { path: OPENAI_AZURE_AI_INFERENCE_PATH } : {},
)
} catch (error) {
throw handleOpenAIError(error, this.providerName)
}
const message = response.choices?.[0]?.message
if (message?.tool_calls) {
for (const toolCall of message.tool_calls) {
if (toolCall.type === "function") {
toolCalls.push({ id: toolCall.id, name: toolCall.function.name })
yield {
type: "tool_call",
id: toolCall.id,
name: toolCall.function.name,
arguments: toolCall.function.arguments,
}
}
}
}
}
yield {
type: "text",
text: message?.content || "",
accumulatedText.push(message?.content || "")
yield {
type: "text",
text: message?.content || "",
}
lastUsage = response.usage
yield this.processUsageMetrics(response.usage)
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
throw error
}
yield this.processUsageMetrics(response.usage)
}
}

View file

@ -140,7 +140,7 @@ export class OpenRouterHandler extends BaseProvider implements SingleCompletionH
private client: OpenAI
protected models: ModelRecord = {}
protected endpoints: ModelRecord = {}
private readonly providerName = "OpenRouter"
protected readonly providerName = "OpenRouter"
private currentReasoningDetails: any[] = []
constructor(options: ApiHandlerOptions) {
@ -319,198 +319,235 @@ export class OpenRouterHandler extends BaseProvider implements SingleCompletionH
? { headers: { "x-anthropic-beta": "fine-grained-tool-streaming-2025-05-14" } }
: undefined
let stream
try {
stream = await this.client.chat.completions.create(completionParams, requestOptions)
} catch (error) {
// Try to parse as OpenRouter error structure using Zod
const parseResult = OpenRouterErrorResponseSchema.safeParse(error)
if (parseResult.success && parseResult.data.error) {
const openRouterError = parseResult.data
const rawString = openRouterError.error?.metadata?.raw
const parsedError = extractErrorFromMetadataRaw(rawString)
const rawErrorMessage = parsedError || openRouterError.error?.message || "Unknown error"
const apiError = Object.assign(
new ApiProviderError(
rawErrorMessage,
this.providerName,
modelId,
"createMessage",
openRouterError.error?.code,
),
{
status: openRouterError.error?.code,
error: openRouterError.error,
},
)
TelemetryService.instance.captureException(apiError)
throw handleOpenAIError(error, this.providerName)
} else {
// Fallback for non-OpenRouter errors
const errorMessage = error instanceof Error ? error.message : String(error)
const apiError = new ApiProviderError(errorMessage, this.providerName, modelId, "createMessage")
TelemetryService.instance.captureException(apiError)
throw handleOpenAIError(error, this.providerName)
}
}
let lastUsage: CompletionUsage | undefined = undefined
// Accumulator for reasoning_details FROM the API.
// We preserve the original shape of reasoning_details to prevent malformed responses.
const reasoningDetailsAccumulator = new Map<
string,
// Start inference logging with actual request params
const logHandle = this.inferenceLogger.start(
{
type: string
text?: string
summary?: string
data?: string
id?: string | null
format?: string
signature?: string
index: number
}
>()
provider: this.providerName,
operation: "createMessage",
model: modelId,
},
{ completionParams, requestOptions },
)
// Track whether we've yielded displayable text from reasoning_details.
// When reasoning_details has displayable content (reasoning.text or reasoning.summary),
// we skip yielding the top-level reasoning field to avoid duplicate display.
let hasYieldedReasoningFromDetails = false
// Accumulators for response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
for await (const chunk of stream) {
// OpenRouter returns an error object instead of the OpenAI SDK throwing an error.
if ("error" in chunk) {
this.handleStreamingError(chunk.error as OpenRouterError, modelId, "createMessage")
try {
let stream
try {
stream = await this.client.chat.completions.create(completionParams, requestOptions)
} catch (error) {
// Try to parse as OpenRouter error structure using Zod
const parseResult = OpenRouterErrorResponseSchema.safeParse(error)
if (parseResult.success && parseResult.data.error) {
const openRouterError = parseResult.data
const rawString = openRouterError.error?.metadata?.raw
const parsedError = extractErrorFromMetadataRaw(rawString)
const rawErrorMessage = parsedError || openRouterError.error?.message || "Unknown error"
const apiError = Object.assign(
new ApiProviderError(
rawErrorMessage,
this.providerName,
modelId,
"createMessage",
openRouterError.error?.code,
),
{
status: openRouterError.error?.code,
error: openRouterError.error,
},
)
TelemetryService.instance.captureException(apiError)
throw handleOpenAIError(error, this.providerName)
} else {
// Fallback for non-OpenRouter errors
const errorMessage = error instanceof Error ? error.message : String(error)
const apiError = new ApiProviderError(errorMessage, this.providerName, modelId, "createMessage")
TelemetryService.instance.captureException(apiError)
throw handleOpenAIError(error, this.providerName)
}
}
const delta = chunk.choices[0]?.delta
const finishReason = chunk.choices[0]?.finish_reason
let lastUsage: CompletionUsage | undefined = undefined
// Accumulator for reasoning_details FROM the API.
// We preserve the original shape of reasoning_details to prevent malformed responses.
const reasoningDetailsAccumulator = new Map<
string,
{
type: string
text?: string
summary?: string
data?: string
id?: string | null
format?: string
signature?: string
index: number
}
>()
if (delta) {
// Handle reasoning_details array format (used by Gemini 3, Claude, OpenAI o-series, etc.)
// See: https://openrouter.ai/docs/use-cases/reasoning-tokens#preserving-reasoning-blocks
// Priority: Check for reasoning_details first, as it's the newer format
const deltaWithReasoning = delta as typeof delta & {
reasoning_details?: Array<{
type: string
text?: string
summary?: string
data?: string
id?: string | null
format?: string
signature?: string
index?: number
}>
// Track whether we've yielded displayable text from reasoning_details.
// When reasoning_details has displayable content (reasoning.text or reasoning.summary),
// we skip yielding the top-level reasoning field to avoid duplicate display.
let hasYieldedReasoningFromDetails = false
for await (const chunk of stream) {
// OpenRouter returns an error object instead of the OpenAI SDK throwing an error.
if ("error" in chunk) {
logHandle.error(chunk.error)
this.handleStreamingError(chunk.error as OpenRouterError, modelId, "createMessage")
}
if (deltaWithReasoning.reasoning_details && Array.isArray(deltaWithReasoning.reasoning_details)) {
for (const detail of deltaWithReasoning.reasoning_details) {
const index = detail.index ?? 0
const key = `${detail.type}-${index}`
const existing = reasoningDetailsAccumulator.get(key)
const delta = chunk.choices[0]?.delta
const finishReason = chunk.choices[0]?.finish_reason
if (existing) {
// Accumulate text/summary/data for existing reasoning detail
if (detail.text !== undefined) {
existing.text = (existing.text || "") + detail.text
}
if (detail.summary !== undefined) {
existing.summary = (existing.summary || "") + detail.summary
}
if (detail.data !== undefined) {
existing.data = (existing.data || "") + detail.data
}
// Update other fields if provided
if (detail.id !== undefined) existing.id = detail.id
if (detail.format !== undefined) existing.format = detail.format
if (detail.signature !== undefined) existing.signature = detail.signature
} else {
// Start new reasoning detail accumulation
reasoningDetailsAccumulator.set(key, {
type: detail.type,
text: detail.text,
summary: detail.summary,
data: detail.data,
id: detail.id,
format: detail.format,
signature: detail.signature,
index,
})
}
if (delta) {
// Handle reasoning_details array format (used by Gemini 3, Claude, OpenAI o-series, etc.)
// See: https://openrouter.ai/docs/use-cases/reasoning-tokens#preserving-reasoning-blocks
// Priority: Check for reasoning_details first, as it's the newer format
const deltaWithReasoning = delta as typeof delta & {
reasoning_details?: Array<{
type: string
text?: string
summary?: string
data?: string
id?: string | null
format?: string
signature?: string
index?: number
}>
}
// Yield text for display (still fragmented for live streaming)
// Only reasoning.text and reasoning.summary have displayable content
// reasoning.encrypted is intentionally skipped as it contains redacted content
let reasoningText: string | undefined
if (detail.type === "reasoning.text" && typeof detail.text === "string") {
reasoningText = detail.text
} else if (detail.type === "reasoning.summary" && typeof detail.summary === "string") {
reasoningText = detail.summary
}
if (deltaWithReasoning.reasoning_details && Array.isArray(deltaWithReasoning.reasoning_details)) {
for (const detail of deltaWithReasoning.reasoning_details) {
const index = detail.index ?? 0
const key = `${detail.type}-${index}`
const existing = reasoningDetailsAccumulator.get(key)
if (reasoningText) {
hasYieldedReasoningFromDetails = true
yield { type: "reasoning", text: reasoningText }
if (existing) {
// Accumulate text/summary/data for existing reasoning detail
if (detail.text !== undefined) {
existing.text = (existing.text || "") + detail.text
}
if (detail.summary !== undefined) {
existing.summary = (existing.summary || "") + detail.summary
}
if (detail.data !== undefined) {
existing.data = (existing.data || "") + detail.data
}
// Update other fields if provided
if (detail.id !== undefined) existing.id = detail.id
if (detail.format !== undefined) existing.format = detail.format
if (detail.signature !== undefined) existing.signature = detail.signature
} else {
// Start new reasoning detail accumulation
reasoningDetailsAccumulator.set(key, {
type: detail.type,
text: detail.text,
summary: detail.summary,
data: detail.data,
id: detail.id,
format: detail.format,
signature: detail.signature,
index,
})
}
// Yield text for display (still fragmented for live streaming)
// Only reasoning.text and reasoning.summary have displayable content
// reasoning.encrypted is intentionally skipped as it contains redacted content
let reasoningText: string | undefined
if (detail.type === "reasoning.text" && typeof detail.text === "string") {
reasoningText = detail.text
} else if (detail.type === "reasoning.summary" && typeof detail.summary === "string") {
reasoningText = detail.summary
}
if (reasoningText) {
hasYieldedReasoningFromDetails = true
accumulatedReasoning.push(reasoningText)
yield { type: "reasoning", text: reasoningText }
}
}
}
// Handle top-level reasoning field for UI display.
// Skip if we've already yielded from reasoning_details to avoid duplicate display.
if ("reasoning" in delta && delta.reasoning && typeof delta.reasoning === "string") {
if (!hasYieldedReasoningFromDetails) {
accumulatedReasoning.push(delta.reasoning)
yield { type: "reasoning", text: delta.reasoning }
}
}
// Emit raw tool call chunks - NativeToolCallParser handles state management
if ("tool_calls" in delta && Array.isArray(delta.tool_calls)) {
for (const toolCall of delta.tool_calls) {
if (toolCall.id || toolCall.function?.name) {
toolCalls.push({ id: toolCall.id, name: toolCall.function?.name })
}
yield {
type: "tool_call_partial",
index: toolCall.index,
id: toolCall.id,
name: toolCall.function?.name,
arguments: toolCall.function?.arguments,
}
}
}
if (delta.content) {
accumulatedText.push(delta.content)
yield { type: "text", text: delta.content }
}
}
// Handle top-level reasoning field for UI display.
// Skip if we've already yielded from reasoning_details to avoid duplicate display.
if ("reasoning" in delta && delta.reasoning && typeof delta.reasoning === "string") {
if (!hasYieldedReasoningFromDetails) {
yield { type: "reasoning", text: delta.reasoning }
// Process finish_reason to emit tool_call_end events
// This ensures tool calls are finalized even if the stream doesn't properly close
if (finishReason) {
const endEvents = NativeToolCallParser.processFinishReason(finishReason)
for (const event of endEvents) {
yield event
}
}
// Emit raw tool call chunks - NativeToolCallParser handles state management
if ("tool_calls" in delta && Array.isArray(delta.tool_calls)) {
for (const toolCall of delta.tool_calls) {
yield {
type: "tool_call_partial",
index: toolCall.index,
id: toolCall.id,
name: toolCall.function?.name,
arguments: toolCall.function?.arguments,
}
}
}
if (delta.content) {
yield { type: "text", text: delta.content }
if (chunk.usage) {
lastUsage = chunk.usage
}
}
// Process finish_reason to emit tool_call_end events
// This ensures tool calls are finalized even if the stream doesn't properly close
if (finishReason) {
const endEvents = NativeToolCallParser.processFinishReason(finishReason)
for (const event of endEvents) {
yield event
// After streaming completes, store ONLY the reasoning_details we received from the API.
if (reasoningDetailsAccumulator.size > 0) {
this.currentReasoningDetails = Array.from(reasoningDetailsAccumulator.values())
}
if (lastUsage) {
yield {
type: "usage",
inputTokens: lastUsage.prompt_tokens || 0,
outputTokens: lastUsage.completion_tokens || 0,
cacheReadTokens: lastUsage.prompt_tokens_details?.cached_tokens,
reasoningTokens: lastUsage.completion_tokens_details?.reasoning_tokens,
totalCost: (lastUsage.cost_details?.upstream_inference_cost || 0) + (lastUsage.cost || 0),
}
}
if (chunk.usage) {
lastUsage = chunk.usage
}
}
// After streaming completes, store ONLY the reasoning_details we received from the API.
if (reasoningDetailsAccumulator.size > 0) {
this.currentReasoningDetails = Array.from(reasoningDetailsAccumulator.values())
}
if (lastUsage) {
yield {
type: "usage",
inputTokens: lastUsage.prompt_tokens || 0,
outputTokens: lastUsage.completion_tokens || 0,
cacheReadTokens: lastUsage.prompt_tokens_details?.cached_tokens,
reasoningTokens: lastUsage.completion_tokens_details?.reasoning_tokens,
totalCost: (lastUsage.cost_details?.upstream_inference_cost || 0) + (lastUsage.cost || 0),
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
reasoningDetails:
reasoningDetailsAccumulator.size > 0 ? Array.from(reasoningDetailsAccumulator.values()) : undefined,
})
} catch (error) {
logHandle.error(error)
throw error
}
}

View file

@ -52,6 +52,7 @@ function objectToUrlEncoded(data: Record<string, string>): string {
}
export class QwenCodeHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "Qwen Code"
protected options: QwenCodeHandlerOptions
private credentials: QwenOAuthCredentials | null = null
private client: OpenAI | undefined

View file

@ -61,7 +61,7 @@ export class RequestyHandler extends BaseProvider implements SingleCompletionHan
protected models: ModelRecord = {}
private client: OpenAI
private baseURL: string
private readonly providerName = "Requesty"
protected readonly providerName = "Requesty"
constructor(options: ApiHandlerOptions) {
super()

View file

@ -125,6 +125,41 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
messages: Anthropic.Messages.MessageParam[],
metadata?: ApiHandlerCreateMessageMetadata,
): ApiStream {
const { id: model, info } = this.getModel()
// Get model parameters for logging
const params = getModelParams({
format: "openai",
modelId: model,
model: info,
settings: this.options,
defaultTemperature: this.defaultTemperature,
})
// Start inference logging
const logHandle = this.inferenceLogger.start(
{
provider: this.providerName,
operation: "createMessage",
model,
taskId: metadata?.taskId,
},
{
model,
maxTokens: params.maxTokens,
temperature: params.temperature,
messageCount: messages.length,
hasTools: !!metadata?.tools,
toolCount: metadata?.tools?.length ?? 0,
toolChoice: metadata?.tool_choice,
},
)
// Accumulators for final response logging
const accumulatedText: string[] = []
const accumulatedReasoning: string[] = []
const toolCalls: Array<{ id?: string; name?: string }> = []
try {
// Reset reasoning_details accumulator for this request
this.currentReasoningDetails = []
@ -234,6 +269,7 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
if (reasoningText) {
hasYieldedReasoningFromDetails = true
accumulatedReasoning.push(reasoningText)
yield { type: "reasoning", text: reasoningText }
}
}
@ -243,11 +279,13 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
// Skip if we've already yielded from reasoning_details to avoid duplicate display.
if ("reasoning" in delta && delta.reasoning && typeof delta.reasoning === "string") {
if (!hasYieldedReasoningFromDetails) {
accumulatedReasoning.push(delta.reasoning)
yield { type: "reasoning", text: delta.reasoning }
}
} else if ("reasoning_content" in delta && typeof delta.reasoning_content === "string") {
// Also check for reasoning_content for backward compatibility
if (!hasYieldedReasoningFromDetails) {
accumulatedReasoning.push(delta.reasoning_content)
yield { type: "reasoning", text: delta.reasoning_content }
}
}
@ -255,6 +293,10 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
// Emit raw tool call chunks - NativeToolCallParser handles state management
if ("tool_calls" in delta && Array.isArray(delta.tool_calls)) {
for (const toolCall of delta.tool_calls) {
// Track tool calls for logging
if (toolCall.id || toolCall.function?.name) {
toolCalls.push({ id: toolCall.id, name: toolCall.function?.name })
}
yield {
type: "tool_call_partial",
index: toolCall.index,
@ -266,6 +308,7 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
}
if (delta.content) {
accumulatedText.push(delta.content)
yield {
type: "text",
text: delta.content,
@ -317,7 +360,17 @@ export class RooHandler extends BaseOpenAiCompatibleProvider<string> {
totalCost: isFreeModel ? 0 : (lastUsage.cost ?? 0),
}
}
// Log successful response
logHandle.success({
text: accumulatedText.join(""),
reasoning: accumulatedReasoning.length > 0 ? accumulatedReasoning.join("") : undefined,
toolCalls: toolCalls.length > 0 ? toolCalls : undefined,
usage: lastUsage,
})
} catch (error) {
logHandle.error(error)
const errorContext = {
error: error instanceof Error ? error.message : String(error),
stack: error instanceof Error ? error.stack : undefined,

View file

@ -22,6 +22,10 @@ type RouterProviderOptions = {
export abstract class RouterProvider extends BaseProvider {
protected readonly options: ApiHandlerOptions
protected readonly name: RouterName
// Implement abstract providerName from BaseProvider using name
protected get providerName(): string {
return this.name
}
protected models: ModelRecord = {}
protected readonly modelId?: string
protected readonly defaultModelId: string

View file

@ -61,6 +61,7 @@ function convertToVsCodeLmTools(tools: OpenAI.Chat.ChatCompletionTool[]): vscode
* ```
*/
export class VsCodeLmHandler extends BaseProvider implements SingleCompletionHandler {
protected readonly providerName = "VS Code LM"
protected options: ApiHandlerOptions
private client: vscode.LanguageModelChat | null
private disposable: vscode.Disposable | null

View file

@ -21,7 +21,7 @@ const XAI_DEFAULT_TEMPERATURE = 0
export class XAIHandler extends BaseProvider implements SingleCompletionHandler {
protected options: ApiHandlerOptions
private client: OpenAI
private readonly providerName = "xAI"
protected readonly providerName = "xAI"
constructor(options: ApiHandlerOptions) {
super()

View file

@ -15,6 +15,7 @@ import {
// Create a mock ApiHandler for testing
class MockApiHandler extends BaseProvider {
protected readonly providerName = "Mock"
createMessage(): any {
// Mock implementation for testing - returns an async iterable stream
const mockStream = {

View file

@ -19,6 +19,7 @@ import {
// Create a mock ApiHandler for testing
class MockApiHandler extends BaseProvider {
protected readonly providerName = "Mock"
createMessage(): any {
// Mock implementation for testing - returns an async iterable stream
const mockStream = {

View file

@ -2,11 +2,14 @@ import * as vscode from "vscode"
import * as dotenvx from "@dotenvx/dotenvx"
import * as path from "path"
// Load environment variables from .env file
// Load environment variables from .env and .env.local files
try {
// Specify path to .env file in the project root directory
// Specify paths to .env and .env.local files in the project root directory
const envPath = path.join(__dirname, "..", ".env")
const envLocalPath = path.join(__dirname, "..", ".env.local")
// Load .env first, then .env.local (so .env.local can override)
dotenvx.config({ path: envPath })
dotenvx.config({ path: envLocalPath, override: true })
} catch (e) {
// Silently handle environment loading errors
console.warn("Failed to load environment variables:", e)
@ -19,6 +22,7 @@ import { customToolRegistry } from "@roo-code/core"
import "./utils/path" // Necessary to have access to String.prototype.toPosix.
import { createOutputChannelLogger, createDualLogger } from "./utils/outputChannelLogger"
import { ApiInferenceLogger } from "./api/logging/ApiInferenceLogger"
import { Package } from "./shared/package"
import { formatLanguage } from "./shared/language"
@ -68,6 +72,12 @@ export async function activate(context: vscode.ExtensionContext) {
context.subscriptions.push(outputChannel)
outputChannel.appendLine(`${Package.name} extension activated - ${JSON.stringify(Package)}`)
// Configure API inference logger for debugging (enable via ROO_CODE_API_LOGGING=true)
ApiInferenceLogger.configure({
enabled: process.env.ROO_CODE_API_LOGGING === "true",
sink: createOutputChannelLogger(outputChannel),
})
// Set extension path for custom tool registry to find bundled esbuild
customToolRegistry.setExtensionPath(context.extensionPath)