mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 21:07:01 -06:00
5cabd2b72e
This reverts commit f5e9622fc9.
3165 lines
112 KiB
TypeScript
3165 lines
112 KiB
TypeScript
import { expectDefined } from "@openclaw/normalization-core";
|
||
// Ollama tests cover stream runtime plugin behavior.
|
||
import { createRequireRecord } from "openclaw/plugin-sdk/test-fixtures";
|
||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||
|
||
const { fetchWithSsrFGuardMock, ollamaStreamWarnMock } = vi.hoisted(() => ({
|
||
fetchWithSsrFGuardMock: vi.fn(),
|
||
ollamaStreamWarnMock: vi.fn(),
|
||
}));
|
||
|
||
vi.mock("openclaw/plugin-sdk/ssrf-runtime", () => ({
|
||
fetchWithSsrFGuard: fetchWithSsrFGuardMock,
|
||
}));
|
||
|
||
vi.mock("openclaw/plugin-sdk/runtime-env", async (importOriginal) => {
|
||
const actual = await importOriginal<typeof import("openclaw/plugin-sdk/runtime-env")>();
|
||
return {
|
||
...actual,
|
||
createSubsystemLogger: () => ({ warn: ollamaStreamWarnMock }),
|
||
};
|
||
});
|
||
|
||
import { OLLAMA_INCOMPLETE_STREAM_ERROR } from "./stream-contract.js";
|
||
import {
|
||
buildOllamaChatRequest,
|
||
createConfiguredOllamaCompatStreamWrapper,
|
||
createConfiguredOllamaStreamFn,
|
||
createOllamaStreamFn,
|
||
convertToOllamaMessages,
|
||
buildAssistantMessage,
|
||
parseNdjsonStream,
|
||
resolveOllamaBaseUrlForRun,
|
||
} from "./stream.runtime.js";
|
||
|
||
type GuardedFetchCall = {
|
||
url: string;
|
||
init?: RequestInit;
|
||
policy?: unknown;
|
||
signal?: AbortSignal;
|
||
timeoutMs?: number;
|
||
auditContext?: string;
|
||
};
|
||
|
||
function requireEntry<T>(entries: readonly T[], index: number, context: string): T {
|
||
return expectDefined(entries[index], context);
|
||
}
|
||
|
||
const requireRecord = createRequireRecord("object", "expected-label");
|
||
|
||
function requireOptionalRecord(value: unknown): Record<string, unknown> | undefined {
|
||
return value === undefined ? undefined : requireRecord(value, "request options");
|
||
}
|
||
|
||
function requireHeaders(value: unknown): Record<string, string> {
|
||
return requireRecord(value, "request headers") as Record<string, string>;
|
||
}
|
||
|
||
function expectToolCallContent(
|
||
value: unknown,
|
||
expected: { name: string; arguments: Record<string, unknown> },
|
||
) {
|
||
const content = requireRecord(value, "tool call content");
|
||
expect(content.type).toBe("toolCall");
|
||
expect(content.name).toBe(expected.name);
|
||
expect(content.arguments).toEqual(expected.arguments);
|
||
}
|
||
|
||
function expectIteratorEvent(
|
||
value: unknown,
|
||
expected: { type?: string; delta?: string; content?: string; done: boolean },
|
||
) {
|
||
const result = requireRecord(value, "iterator result");
|
||
expect(result.done).toBe(expected.done);
|
||
if (expected.type !== undefined) {
|
||
const event = requireRecord(result.value, "iterator result value");
|
||
expect(event.type).toBe(expected.type);
|
||
if (expected.delta !== undefined) {
|
||
expect(event.delta).toBe(expected.delta);
|
||
}
|
||
if (expected.content !== undefined) {
|
||
expect(event.content).toBe(expected.content);
|
||
}
|
||
} else {
|
||
expect(result.value).toBeUndefined();
|
||
}
|
||
}
|
||
|
||
function convertAssistantContent(
|
||
content: Array<Record<string, unknown>>,
|
||
options?: Parameters<typeof convertToOllamaMessages>[2],
|
||
) {
|
||
return convertToOllamaMessages([{ role: "assistant", content }] as never, undefined, options);
|
||
}
|
||
|
||
type AssistantResponse = Parameters<typeof buildAssistantMessage>[0];
|
||
type AssistantResponseMessage = AssistantResponse["message"];
|
||
|
||
function createAssistantResponse(
|
||
message: Omit<AssistantResponseMessage, "role">,
|
||
overrides: Partial<Omit<AssistantResponse, "message">> = {},
|
||
): AssistantResponse {
|
||
return {
|
||
model: "qwen3:32b",
|
||
created_at: "2026-01-01T00:00:00Z",
|
||
message: { role: "assistant", ...message },
|
||
done: true,
|
||
...overrides,
|
||
};
|
||
}
|
||
|
||
function createToolCallResponse(
|
||
toolCalls: NonNullable<AssistantResponseMessage["tool_calls"]>,
|
||
overrides: Partial<Omit<AssistantResponse, "message">> = {},
|
||
): AssistantResponse {
|
||
return createAssistantResponse({ content: "", tool_calls: toolCalls }, overrides);
|
||
}
|
||
|
||
afterEach(() => {
|
||
fetchWithSsrFGuardMock.mockReset();
|
||
ollamaStreamWarnMock.mockReset();
|
||
});
|
||
|
||
describe("buildOllamaChatRequest", () => {
|
||
it("omits tools when none are provided", () => {
|
||
expect(
|
||
buildOllamaChatRequest({
|
||
modelId: "qwen3.5:9b",
|
||
messages: [{ role: "user", content: "hello" }],
|
||
options: { num_ctx: 65536 },
|
||
}),
|
||
).toEqual({
|
||
model: "qwen3.5:9b",
|
||
messages: [{ role: "user", content: "hello" }],
|
||
stream: true,
|
||
options: { num_ctx: 65536 },
|
||
});
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "strips the ollama/ prefix from chat model ids",
|
||
modelId: "ollama/qwen3:14b-q8_0",
|
||
expected: "qwen3:14b-q8_0",
|
||
},
|
||
{
|
||
name: "strips the active custom provider prefix from chat model ids",
|
||
modelId: "ollama-spark/qwen3:32b",
|
||
providerId: "ollama-spark",
|
||
expected: "qwen3:32b",
|
||
},
|
||
{
|
||
name: "keeps unrelated slash-containing Ollama model ids intact",
|
||
modelId: "library/qwen3:32b",
|
||
providerId: "ollama-spark",
|
||
expected: "library/qwen3:32b",
|
||
},
|
||
])("$name", ({ modelId, providerId, expected }) => {
|
||
expect(
|
||
buildOllamaChatRequest({
|
||
modelId,
|
||
providerId,
|
||
messages: [{ role: "user", content: "hello" }],
|
||
}).model,
|
||
).toBe(expected);
|
||
});
|
||
|
||
it("keeps native Ollama replay tool arguments as objects", () => {
|
||
const messages = convertToOllamaMessages([
|
||
{
|
||
role: "assistant",
|
||
content: [
|
||
{
|
||
type: "toolCall",
|
||
name: "gateway",
|
||
arguments: '{"action":"config.get","path":"gateway.port"}',
|
||
},
|
||
],
|
||
},
|
||
]);
|
||
|
||
expect(messages[0]?.tool_calls?.[0]?.function.arguments).toEqual({
|
||
action: "config.get",
|
||
path: "gateway.port",
|
||
});
|
||
});
|
||
});
|
||
|
||
describe("createConfiguredOllamaCompatStreamWrapper", () => {
|
||
it("adds Moonshot thinking config for Ollama cloud Kimi compat requests", async () => {
|
||
let patchedPayload: Record<string, unknown> | undefined;
|
||
const baseStreamFn = vi.fn((_model, _context, options) => {
|
||
options?.onPayload?.({ tool_choice: "auto" });
|
||
return (async function* () {})();
|
||
});
|
||
const model = {
|
||
api: "openai-completions",
|
||
provider: "ollama",
|
||
id: "kimi-k2.5:cloud",
|
||
contextWindow: 262144,
|
||
params: { num_ctx: 65536 },
|
||
};
|
||
|
||
const wrapped = createConfiguredOllamaCompatStreamWrapper({
|
||
provider: "ollama",
|
||
modelId: "kimi-k2.5:cloud",
|
||
model,
|
||
streamFn: baseStreamFn,
|
||
thinkingLevel: "high",
|
||
extraParams: {},
|
||
} as never);
|
||
|
||
await wrapped?.(
|
||
model as never,
|
||
{ messages: [] } as never,
|
||
{
|
||
onPayload: (payload: unknown) => {
|
||
patchedPayload = payload as Record<string, unknown>;
|
||
},
|
||
} as never,
|
||
);
|
||
|
||
const payload = requireRecord(patchedPayload, "patched payload");
|
||
expect(payload.thinking).toEqual({ type: "enabled" });
|
||
expect(payload.options).toEqual({ num_ctx: 65536 });
|
||
});
|
||
|
||
it("preserves OpenAI-compatible replay tool arguments as strings", async () => {
|
||
let patchedPayload: Record<string, unknown> | undefined;
|
||
const baseStreamFn = vi.fn((_model, _context, options) => {
|
||
options?.onPayload?.({
|
||
messages: [
|
||
{
|
||
role: "assistant",
|
||
function_call: {
|
||
name: "legacy_gateway",
|
||
arguments: '{"action":"config.get"}',
|
||
},
|
||
tool_calls: [
|
||
{
|
||
id: "call_gateway",
|
||
type: "function",
|
||
function: {
|
||
name: "gateway",
|
||
arguments: '{"action":"config.get","path":"gateway.port"}',
|
||
},
|
||
},
|
||
],
|
||
},
|
||
],
|
||
});
|
||
return (async function* () {})();
|
||
});
|
||
const model = {
|
||
api: "openai-completions",
|
||
provider: "ollama",
|
||
id: "glm-5.2:cloud",
|
||
contextWindow: 262144,
|
||
};
|
||
|
||
const wrapped = createConfiguredOllamaCompatStreamWrapper({
|
||
provider: "ollama",
|
||
modelId: "glm-5.2:cloud",
|
||
model,
|
||
streamFn: baseStreamFn,
|
||
} as never);
|
||
|
||
await wrapped?.(
|
||
model as never,
|
||
{ messages: [] } as never,
|
||
{
|
||
onPayload: (payload: unknown) => {
|
||
patchedPayload = payload as Record<string, unknown>;
|
||
},
|
||
} as never,
|
||
);
|
||
|
||
const payload = requireRecord(patchedPayload, "patched payload");
|
||
const messages = payload.messages as Array<Record<string, unknown>>;
|
||
const assistantMessage = requireRecord(messages[0], "assistant message");
|
||
const functionCall = requireRecord(assistantMessage.function_call, "function call");
|
||
const toolCalls = assistantMessage.tool_calls as Array<Record<string, unknown>>;
|
||
const toolCallFunction = requireRecord(toolCalls[0]?.function, "tool call function");
|
||
expect(functionCall.arguments).toBe('{"action":"config.get"}');
|
||
expect(toolCallFunction.arguments).toBe('{"action":"config.get","path":"gateway.port"}');
|
||
expect(payload.options).toEqual({ num_ctx: 262144 });
|
||
});
|
||
|
||
it("falls back to contextWindow when configured num_ctx is invalid", async () => {
|
||
let patchedPayload: Record<string, unknown> | undefined;
|
||
const baseStreamFn = vi.fn((_model, _context, options) => {
|
||
options?.onPayload?.({});
|
||
return (async function* () {})();
|
||
});
|
||
const model = {
|
||
api: "openai-completions",
|
||
provider: "ollama",
|
||
id: "qwen3:32b",
|
||
contextWindow: 131072,
|
||
params: { num_ctx: 0 },
|
||
};
|
||
|
||
const wrapped = createConfiguredOllamaCompatStreamWrapper({
|
||
provider: "ollama",
|
||
modelId: "qwen3:32b",
|
||
model,
|
||
streamFn: baseStreamFn,
|
||
} as never);
|
||
|
||
await wrapped?.(
|
||
model as never,
|
||
{ messages: [] } as never,
|
||
{
|
||
onPayload: (payload: unknown) => {
|
||
patchedPayload = payload as Record<string, unknown>;
|
||
},
|
||
} as never,
|
||
);
|
||
|
||
const payload = requireRecord(patchedPayload, "patched payload");
|
||
expect(payload.options).toEqual({ num_ctx: 131072 });
|
||
});
|
||
|
||
it.each<{
|
||
name: string;
|
||
id: string;
|
||
contextWindow: number;
|
||
provider?: string;
|
||
reasoning?: boolean;
|
||
thinkingLevel: string;
|
||
params?: Record<string, unknown>;
|
||
expectedThink: boolean | string | undefined;
|
||
}>([
|
||
{
|
||
name: "forwards think=false on native Ollama chat requests when thinking is off",
|
||
id: "qwen3:32b",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "off",
|
||
expectedThink: false,
|
||
},
|
||
{
|
||
name: "does not overwrite configured native Ollama params.thinking with implicit off",
|
||
id: "qwen3:32b",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "off",
|
||
params: { thinking: "medium" },
|
||
expectedThink: "medium",
|
||
},
|
||
{
|
||
name: "does not forward truthy configured native Ollama thinking for non-reasoning models",
|
||
id: "llama3.2:latest",
|
||
contextWindow: 8192,
|
||
reasoning: false,
|
||
thinkingLevel: "off",
|
||
params: { thinking: "medium" },
|
||
expectedThink: undefined,
|
||
},
|
||
{
|
||
name: "does not forward runtime native Ollama thinking for non-reasoning models",
|
||
id: "llama3.2:latest",
|
||
contextWindow: 8192,
|
||
reasoning: false,
|
||
thinkingLevel: "low",
|
||
expectedThink: undefined,
|
||
},
|
||
...(["low", "medium", "high"] as const).map((thinkingLevel) => ({
|
||
name: `preserves native Ollama ${thinkingLevel} thinking on the wire`,
|
||
id: "gpt-oss:20b",
|
||
contextWindow: 131072,
|
||
thinkingLevel,
|
||
expectedThink: thinkingLevel,
|
||
})),
|
||
{
|
||
name: "keeps the compatible local Ollama max mapping",
|
||
id: "gpt-oss:20b",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "max",
|
||
expectedThink: "high",
|
||
},
|
||
{
|
||
name: "does not infer native max support from a local cloud model alias",
|
||
id: "glm-5.2:cloud",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "max",
|
||
expectedThink: "high",
|
||
},
|
||
{
|
||
name: "preserves native Ollama Cloud max thinking on the wire",
|
||
id: "glm-5.2",
|
||
provider: "ollama-cloud",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "max",
|
||
expectedThink: "max",
|
||
},
|
||
{
|
||
name: "keeps the high fallback for Ollama Cloud GPT-OSS",
|
||
id: "gpt-oss:120b",
|
||
provider: "ollama-cloud",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "max",
|
||
expectedThink: "high",
|
||
},
|
||
{
|
||
name: "keeps the high fallback for Cloud models without a verified max tier",
|
||
id: "kimi-k2.5",
|
||
provider: "ollama-cloud",
|
||
contextWindow: 131072,
|
||
thinkingLevel: "max",
|
||
expectedThink: "high",
|
||
},
|
||
])(
|
||
"$name",
|
||
async ({
|
||
id,
|
||
provider = "ollama",
|
||
contextWindow,
|
||
reasoning,
|
||
thinkingLevel,
|
||
params,
|
||
expectedThink,
|
||
}) => {
|
||
await withSuccessfulOllamaFetch(async (fetchMock) => {
|
||
const model = {
|
||
api: "ollama",
|
||
provider,
|
||
id,
|
||
contextWindow,
|
||
...(reasoning === undefined ? {} : { reasoning }),
|
||
...(params ? { params } : {}),
|
||
};
|
||
const wrapped = expectDefined(
|
||
createConfiguredOllamaCompatStreamWrapper({
|
||
provider,
|
||
modelId: id,
|
||
model,
|
||
streamFn: createOllamaStreamFn("http://ollama-host:11434"),
|
||
thinkingLevel,
|
||
} as never),
|
||
"wrapped Ollama stream function",
|
||
);
|
||
const stream = await Promise.resolve(
|
||
wrapped(
|
||
model as never,
|
||
{ messages: [{ role: "user", content: "hello" }] } as never,
|
||
{} as never,
|
||
),
|
||
);
|
||
await collectStreamEvents(stream);
|
||
|
||
const requestBody = getGuardedFetchJsonBody(fetchMock);
|
||
expect(requestBody.think).toBe(expectedThink);
|
||
expect(requireOptionalRecord(requestBody.options)?.think).toBeUndefined();
|
||
if (reasoning !== false) {
|
||
expect(requireOptionalRecord(requestBody.options)?.num_ctx).toBeUndefined();
|
||
}
|
||
});
|
||
},
|
||
);
|
||
|
||
it("passes resolved provider request timeouts to native Ollama chat fetches", async () => {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
async (fetchMock) => {
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { requestTimeoutMs: 450_000 },
|
||
});
|
||
|
||
await collectStreamEvents(stream);
|
||
|
||
expect(getGuardedFetchCall(fetchMock).timeoutMs).toBe(450_000);
|
||
},
|
||
);
|
||
});
|
||
|
||
it("passes caller abort signals at guard level when a timeout is present", async () => {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
async (fetchMock) => {
|
||
const signal = new AbortController().signal;
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: { signal, timeoutMs: 123_456 },
|
||
});
|
||
|
||
await collectStreamEvents(stream);
|
||
|
||
const request = getGuardedFetchCall(fetchMock);
|
||
expect(request.timeoutMs).toBe(123_456);
|
||
expect(request.signal).toBe(signal);
|
||
expect(request.init?.signal).toBeUndefined();
|
||
},
|
||
);
|
||
});
|
||
|
||
it("sends custom-provider Ollama chat requests with the bare Ollama model id", async () => {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
async (fetchMock) => {
|
||
const streamFn = createOllamaStreamFn("http://ollama-host:11434");
|
||
const model = {
|
||
api: "ollama",
|
||
provider: "ollama-spark",
|
||
id: "ollama-spark/qwen3:32b",
|
||
contextWindow: 131072,
|
||
};
|
||
|
||
const stream = await Promise.resolve(
|
||
streamFn(
|
||
model as never,
|
||
{
|
||
messages: [{ role: "user", content: "hello" }],
|
||
} as never,
|
||
{} as never,
|
||
),
|
||
);
|
||
|
||
await collectStreamEvents(stream);
|
||
|
||
const requestInit = getGuardedFetchCall(fetchMock).init ?? {};
|
||
if (typeof requestInit.body !== "string") {
|
||
throw new Error("Expected string request body");
|
||
}
|
||
const requestBody = JSON.parse(requestInit.body) as { model?: string };
|
||
expect(requestBody.model).toBe("qwen3:32b");
|
||
},
|
||
);
|
||
});
|
||
|
||
it("adds direct type hints to native Ollama tool schemas before sending them", async () => {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
async (fetchMock) => {
|
||
const streamFn = createOllamaStreamFn("http://ollama-host:11434");
|
||
const model = {
|
||
api: "ollama",
|
||
provider: "ollama",
|
||
id: "qwen3:32b",
|
||
contextWindow: 131072,
|
||
};
|
||
|
||
const stream = await Promise.resolve(
|
||
streamFn(
|
||
model as never,
|
||
{
|
||
messages: [{ role: "user", content: "hello" }],
|
||
tools: [
|
||
{
|
||
name: "search",
|
||
description: "search",
|
||
parameters: {
|
||
properties: {
|
||
query: {
|
||
anyOf: [{ type: "string" }, { type: "null" }],
|
||
},
|
||
tags: {
|
||
items: { type: "string" },
|
||
},
|
||
},
|
||
required: ["query"],
|
||
},
|
||
},
|
||
],
|
||
} as never,
|
||
{} as never,
|
||
),
|
||
);
|
||
|
||
await collectStreamEvents(stream);
|
||
|
||
const requestInit = getGuardedFetchCall(fetchMock).init ?? {};
|
||
if (typeof requestInit.body !== "string") {
|
||
throw new Error("Expected string request body");
|
||
}
|
||
const requestBody = JSON.parse(requestInit.body) as {
|
||
tools?: Array<{
|
||
function?: {
|
||
parameters?: {
|
||
type?: string;
|
||
properties?: Record<string, { type?: string }>;
|
||
};
|
||
};
|
||
}>;
|
||
};
|
||
const parameters = requestBody.tools?.[0]?.function?.parameters;
|
||
expect(parameters?.type).toBe("object");
|
||
expect(parameters?.properties?.query?.type).toBe("string");
|
||
expect(parameters?.properties?.tags?.type).toBe("array");
|
||
},
|
||
);
|
||
});
|
||
});
|
||
|
||
describe("convertToOllamaMessages", () => {
|
||
it("converts user text messages", () => {
|
||
const messages = [{ role: "user", content: "hello" }];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "user", content: "hello" }]);
|
||
});
|
||
|
||
it("converts user messages with content parts", () => {
|
||
const messages = [
|
||
{
|
||
role: "user",
|
||
content: [
|
||
{ type: "text", text: "describe this" },
|
||
{ type: "image", data: "base64data" },
|
||
],
|
||
},
|
||
];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "user", content: "describe this", images: ["base64data"] }]);
|
||
});
|
||
|
||
it("prepends system message when provided", () => {
|
||
const messages = [{ role: "user", content: "hello" }];
|
||
const result = convertToOllamaMessages(messages, "You are helpful.");
|
||
expect(requireEntry(result, 0, "first converted Ollama message")).toEqual({
|
||
role: "system",
|
||
content: "You are helpful.",
|
||
});
|
||
expect(requireEntry(result, 1, "second converted Ollama message")).toEqual({
|
||
role: "user",
|
||
content: "hello",
|
||
});
|
||
});
|
||
|
||
it("converts assistant messages with toolCall content blocks", () => {
|
||
const result = convertAssistantContent([
|
||
{ type: "text", text: "Let me check." },
|
||
{ type: "toolCall", id: "call_1", name: "bash", arguments: { command: "ls" } },
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").role).toBe("assistant");
|
||
expect(requireEntry(result, 0, "first converted Ollama message").content).toBe("Let me check.");
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_1", function: { name: "bash", arguments: { command: "ls" } } },
|
||
]);
|
||
});
|
||
|
||
it("preserves assistant tool-call ids before Ollama replay", () => {
|
||
const result = convertAssistantContent([
|
||
{
|
||
type: "toolCall",
|
||
id: "fc_ollama_123",
|
||
name: "bash",
|
||
arguments: { command: "pwd" },
|
||
},
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "fc_ollama_123", function: { name: "bash", arguments: { command: "pwd" } } },
|
||
]);
|
||
});
|
||
|
||
it("normalizes provider-prefixed tool-call names before Ollama replay", () => {
|
||
const result = convertAssistantContent([
|
||
{ type: "toolCall", id: "call_1", name: "functions.exec", arguments: { command: "pwd" } },
|
||
{ type: "tool_use", id: "call_2", name: "tools/read", input: { path: "README.md" } },
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_1", function: { name: "exec", arguments: { command: "pwd" } } },
|
||
{ id: "call_2", function: { name: "read", arguments: { path: "README.md" } } },
|
||
]);
|
||
});
|
||
|
||
it("preserves exact allowlisted tool-prefix names before Ollama replay", () => {
|
||
const result = convertAssistantContent(
|
||
[
|
||
{ type: "toolCall", id: "call_1", name: "tool_a", arguments: { value: 1 } },
|
||
{ type: "tool_use", id: "call_2", name: "tools_invoke_test", input: { value: 2 } },
|
||
{ type: "toolCall", id: "call_3", name: "function-run", arguments: { value: 3 } },
|
||
],
|
||
{
|
||
availableToolNames: new Set(["tool_a", "tools_invoke_test", "function-run"]),
|
||
},
|
||
);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_1", function: { name: "tool_a", arguments: { value: 1 } } },
|
||
{ id: "call_2", function: { name: "tools_invoke_test", arguments: { value: 2 } } },
|
||
{ id: "call_3", function: { name: "function-run", arguments: { value: 3 } } },
|
||
]);
|
||
});
|
||
|
||
it("strips underscore and dash provider prefixes only when the suffix is allowlisted", () => {
|
||
const result = convertAssistantContent(
|
||
[
|
||
{ type: "toolCall", id: "call_1", name: "tools_exec", arguments: { command: "pwd" } },
|
||
{ type: "tool_use", id: "call_2", name: "function-read", input: { path: "." } },
|
||
{ type: "toolCall", id: "call_3", name: "tool_missing", arguments: {} },
|
||
],
|
||
{
|
||
availableToolNames: new Set(["exec", "read"]),
|
||
},
|
||
);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_1", function: { name: "exec", arguments: { command: "pwd" } } },
|
||
{ id: "call_2", function: { name: "read", arguments: { path: "." } } },
|
||
{ id: "call_3", function: { name: "tool_missing", arguments: {} } },
|
||
]);
|
||
});
|
||
|
||
it("keeps non-prefixed Ollama replay tool names intact", () => {
|
||
const result = convertAssistantContent([
|
||
{ type: "toolCall", id: "call_1", name: "functionshell", arguments: {} },
|
||
{ type: "toolCall", id: "call_2", name: "tooling", arguments: {} },
|
||
{ type: "toolCall", id: "call_3", name: "tools", arguments: {} },
|
||
{ type: "toolCall", id: "call_4", name: "tool_a", arguments: {} },
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_1", function: { name: "functionshell", arguments: {} } },
|
||
{ id: "call_2", function: { name: "tooling", arguments: {} } },
|
||
{ id: "call_3", function: { name: "tools", arguments: {} } },
|
||
{ id: "call_4", function: { name: "tool_a", arguments: {} } },
|
||
]);
|
||
});
|
||
|
||
it("deserializes string arguments back to objects for Ollama (round-trip fix)", () => {
|
||
// When tool calls round-trip through OpenAI-format storage, arguments
|
||
// are serialized as a JSON string. Ollama expects an object.
|
||
const result = convertAssistantContent([
|
||
{
|
||
type: "toolCall",
|
||
id: "call_2",
|
||
name: "Read",
|
||
arguments: '{"file_path":"/tmp/test.txt"}',
|
||
},
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "call_2", function: { name: "Read", arguments: { file_path: "/tmp/test.txt" } } },
|
||
]);
|
||
});
|
||
|
||
it("handles tool_use blocks with string input (Anthropic format round-trip)", () => {
|
||
const result = convertAssistantContent([
|
||
{ type: "tool_use", id: "toolu_1", name: "exec", input: '{"command":"echo hello"}' },
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{ id: "toolu_1", function: { name: "exec", arguments: { command: "echo hello" } } },
|
||
]);
|
||
});
|
||
|
||
it("preserves unsafe integers as strings when replay args are deserialized", () => {
|
||
const result = convertAssistantContent([
|
||
{
|
||
type: "toolCall",
|
||
id: "call_3",
|
||
name: "read",
|
||
arguments: '{"path":9223372036854775807,"nested":{"thread":1234567890123456789}}',
|
||
},
|
||
]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message").tool_calls).toEqual([
|
||
{
|
||
id: "call_3",
|
||
function: {
|
||
name: "read",
|
||
arguments: {
|
||
path: "9223372036854775807",
|
||
nested: { thread: "1234567890123456789" },
|
||
},
|
||
},
|
||
},
|
||
]);
|
||
});
|
||
it("converts tool result messages with 'tool' role", () => {
|
||
const messages = [{ role: "tool", content: "file1.txt\nfile2.txt" }];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "tool", content: "file1.txt\nfile2.txt" }]);
|
||
});
|
||
|
||
it("converts SDK 'toolResult' role to Ollama 'tool' role", () => {
|
||
const messages = [{ role: "toolResult", content: "command output here" }];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "tool", content: "command output here" }]);
|
||
});
|
||
|
||
it("includes tool_name from SDK toolResult messages", () => {
|
||
const messages = [{ role: "toolResult", content: "file contents here", toolName: "read" }];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "tool", content: "file contents here", tool_name: "read" }]);
|
||
});
|
||
|
||
it("omits tool_name when not provided in toolResult", () => {
|
||
const messages = [{ role: "toolResult", content: "output" }];
|
||
const result = convertToOllamaMessages(messages);
|
||
expect(result).toEqual([{ role: "tool", content: "output" }]);
|
||
expect(requireEntry(result, 0, "first converted Ollama message")).not.toHaveProperty(
|
||
"tool_name",
|
||
);
|
||
});
|
||
|
||
it("handles empty messages array", () => {
|
||
const result = convertToOllamaMessages([]);
|
||
expect(result).toStrictEqual([]);
|
||
});
|
||
});
|
||
|
||
describe("buildAssistantMessage", () => {
|
||
const modelInfo = { api: "ollama", provider: "ollama", id: "qwen3:32b" };
|
||
|
||
it("builds text-only response", () => {
|
||
const response = createAssistantResponse(
|
||
{ content: "Hello!" },
|
||
{
|
||
prompt_eval_count: 10,
|
||
eval_count: 5,
|
||
},
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.role).toBe("assistant");
|
||
expect(result.content).toEqual([{ type: "text", text: "Hello!" }]);
|
||
expect(result.stopReason).toBe("stop");
|
||
expect(result.usage.input).toBe(10);
|
||
expect(result.usage.output).toBe(5);
|
||
expect(result.usage.totalTokens).toBe(15);
|
||
});
|
||
|
||
it("keeps thinking-only output when content is empty", () => {
|
||
const response = createAssistantResponse({ content: "", thinking: "Thinking output" });
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.stopReason).toBe("stop");
|
||
expect(result.content).toEqual([{ type: "thinking", thinking: "Thinking output" }]);
|
||
});
|
||
|
||
it("keeps reasoning-only output when content and thinking are empty", () => {
|
||
const response = createAssistantResponse({ content: "", reasoning: "Reasoning output" });
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.stopReason).toBe("stop");
|
||
expect(result.content).toEqual([{ type: "thinking", thinking: "Reasoning output" }]);
|
||
});
|
||
|
||
it("drops provider-returned thinking for non-reasoning models", () => {
|
||
const response = createAssistantResponse(
|
||
{ content: "", thinking: "Thinking output" },
|
||
{
|
||
model: "minimax-m2.7:cloud",
|
||
prompt_eval_count: 10,
|
||
eval_count: 6,
|
||
},
|
||
);
|
||
const result = buildAssistantMessage(response, {
|
||
...modelInfo,
|
||
id: "minimax-m2.7:cloud",
|
||
reasoning: false,
|
||
});
|
||
expect(result.stopReason).toBe("stop");
|
||
expect(result.content).toEqual([]);
|
||
expect(result.usage.output).toBe(6);
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "strips inline reasoning prefix from kimi cloud visible text",
|
||
id: "kimi-k2.6:cloud",
|
||
answer: "Final answer only.",
|
||
separator: "️",
|
||
},
|
||
{
|
||
name: "strips inline reasoning for provider-qualified Kimi cloud refs",
|
||
id: "ollama/kimi-k2.6:cloud",
|
||
answer: "Final answer only.",
|
||
separator: "️",
|
||
},
|
||
{
|
||
name: "strips inline reasoning when the Kimi boundary is followed by whitespace",
|
||
id: "kimi-k2.6:cloud",
|
||
answer: "Final answer only.",
|
||
separator: "️ ",
|
||
},
|
||
{
|
||
name: "strips inline reasoning before short Kimi cloud answers",
|
||
id: "kimi-k2.6:cloud",
|
||
answer: "OK.",
|
||
separator: "️ ",
|
||
},
|
||
{
|
||
name: "strips inline reasoning when the Kimi boundary has no visible answer",
|
||
id: "kimi-k2.6:cloud",
|
||
answer: "",
|
||
separator: "️ ",
|
||
},
|
||
])("$name", ({ id, answer, separator }) => {
|
||
const response = createAssistantResponse(
|
||
{
|
||
content:
|
||
"I should think privately and not leak this planning text in the answer. " +
|
||
`I need to keep deciding what to say next. ${separator}${answer}`,
|
||
},
|
||
{ model: "kimi-k2.6:cloud" },
|
||
);
|
||
expect(
|
||
buildAssistantMessage(response, { api: "ollama", provider: "ollama", id }).content,
|
||
).toEqual(answer ? [{ type: "text", text: answer }] : []);
|
||
});
|
||
|
||
it("does not strip inline boundary marker on non-kimi models", () => {
|
||
const response = createAssistantResponse({ content: "intro ️keep this intact" });
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.content).toEqual([{ type: "text", text: "intro ️keep this intact" }]);
|
||
});
|
||
|
||
it("does not treat emoji variation selectors as Kimi inline-reasoning boundaries", () => {
|
||
const response = createAssistantResponse(
|
||
{
|
||
content:
|
||
"This is a normal Kimi cloud answer with enough length to cross the prefix threshold and no hidden reasoning leak. ☀️sunshine should remain visible to the user.",
|
||
},
|
||
{ model: "kimi-k2.6:cloud" },
|
||
);
|
||
const result = buildAssistantMessage(response, {
|
||
api: "ollama",
|
||
provider: "ollama",
|
||
id: "kimi-k2.6:cloud",
|
||
});
|
||
expect(result.content).toEqual([
|
||
{
|
||
type: "text",
|
||
text: "This is a normal Kimi cloud answer with enough length to cross the prefix threshold and no hidden reasoning leak. ☀️sunshine should remain visible to the user.",
|
||
},
|
||
]);
|
||
});
|
||
|
||
it("estimates usage when Ollama omits eval counters", () => {
|
||
const response = createAssistantResponse({ content: "Estimated output" });
|
||
const result = buildAssistantMessage(response, modelInfo, { input: 11, output: 4 });
|
||
expect(result.usage.input).toBe(11);
|
||
expect(result.usage.output).toBe(4);
|
||
expect(result.usage.totalTokens).toBe(15);
|
||
});
|
||
|
||
it("preserves explicit zero usage counters from Ollama", () => {
|
||
const response = createAssistantResponse(
|
||
{ content: "" },
|
||
{
|
||
prompt_eval_count: 0,
|
||
eval_count: 0,
|
||
},
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo, { input: 11, output: 4 });
|
||
expect(result.usage.input).toBe(0);
|
||
expect(result.usage.output).toBe(0);
|
||
expect(result.usage.totalTokens).toBe(0);
|
||
});
|
||
|
||
it("builds response with tool calls", () => {
|
||
const response = createToolCallResponse(
|
||
[{ function: { name: "bash", arguments: { command: "ls -la" } } }],
|
||
{
|
||
prompt_eval_count: 20,
|
||
eval_count: 10,
|
||
},
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.stopReason).toBe("toolUse");
|
||
expect(result.content.length).toBe(1); // toolCall only (empty content is skipped)
|
||
expect(requireEntry(result.content, 0, "Ollama tool-call content").type).toBe("toolCall");
|
||
const toolCall = requireEntry(result.content, 0, "Ollama tool-call content") as {
|
||
type: "toolCall";
|
||
id: string;
|
||
name: string;
|
||
arguments: Record<string, unknown>;
|
||
};
|
||
expect(toolCall.name).toBe("bash");
|
||
expect(toolCall.arguments).toEqual({ command: "ls -la" });
|
||
expect(toolCall.id).toMatch(/^ollama_call_[0-9a-f-]{36}$/);
|
||
});
|
||
|
||
it("preserves Ollama response tool-call ids", () => {
|
||
const response = createToolCallResponse(
|
||
[{ id: "fc_ollama_real_1", function: { name: "bash", arguments: { command: "pwd" } } }],
|
||
{ model: "gemini-3-flash-preview:cloud" },
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "bash",
|
||
arguments: { command: "pwd" },
|
||
});
|
||
expect(
|
||
(requireEntry(result.content, 0, "Ollama tool-call content") as { id?: string }).id,
|
||
).toBe("fc_ollama_real_1");
|
||
});
|
||
|
||
it("preserves parallel Ollama response tool-call ids independently", () => {
|
||
const response = createToolCallResponse(
|
||
[
|
||
{ id: "fc_ollama_real_1", function: { name: "read", arguments: { path: "a.txt" } } },
|
||
{ id: "fc_ollama_real_2", function: { name: "exec", arguments: { command: "date" } } },
|
||
],
|
||
{ model: "gemini-3-flash-preview:cloud" },
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.content.map((part) => (part as { id?: string }).id)).toEqual([
|
||
"fc_ollama_real_1",
|
||
"fc_ollama_real_2",
|
||
]);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "read",
|
||
arguments: { path: "a.txt" },
|
||
});
|
||
expectToolCallContent(result.content[1], { name: "exec", arguments: { command: "date" } });
|
||
});
|
||
|
||
it("normalizes provider-prefixed tool-call names in Ollama responses", () => {
|
||
const response = createToolCallResponse([
|
||
{ function: { name: "functions.exec", arguments: { command: "pwd" } } },
|
||
{ function: { name: "tools/read", arguments: { path: "README.md" } } },
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.content).toHaveLength(2);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "exec",
|
||
arguments: { command: "pwd" },
|
||
});
|
||
expectToolCallContent(result.content[1], { name: "read", arguments: { path: "README.md" } });
|
||
});
|
||
|
||
it("preserves exact allowlisted tool-prefix names in Ollama responses", () => {
|
||
const response = createToolCallResponse([
|
||
{ function: { name: "tool_a", arguments: { value: 1 } } },
|
||
{ function: { name: "tools_invoke_test", arguments: { value: 2 } } },
|
||
{ function: { name: "function-run", arguments: { value: 3 } } },
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo, undefined, {
|
||
availableToolNames: new Set(["tool_a", "tools_invoke_test", "function-run"]),
|
||
});
|
||
expect(result.content).toHaveLength(3);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "tool_a",
|
||
arguments: { value: 1 },
|
||
});
|
||
expectToolCallContent(result.content[1], {
|
||
name: "tools_invoke_test",
|
||
arguments: { value: 2 },
|
||
});
|
||
expectToolCallContent(result.content[2], { name: "function-run", arguments: { value: 3 } });
|
||
});
|
||
|
||
it("keeps non-prefixed Ollama response tool names intact", () => {
|
||
const response = createToolCallResponse([
|
||
{ function: { name: "functionshell", arguments: {} } },
|
||
{ function: { name: "tooling", arguments: {} } },
|
||
{ function: { name: "tools", arguments: {} } },
|
||
{ function: { name: "tool_a", arguments: {} } },
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.content).toHaveLength(4);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "functionshell",
|
||
arguments: {},
|
||
});
|
||
expectToolCallContent(result.content[1], { name: "tooling", arguments: {} });
|
||
expectToolCallContent(result.content[2], { name: "tools", arguments: {} });
|
||
expectToolCallContent(result.content[3], { name: "tool_a", arguments: {} });
|
||
});
|
||
|
||
it("parses stringified tool call arguments from Ollama responses", () => {
|
||
const response = createToolCallResponse([
|
||
{ function: { name: "bash", arguments: '{"command":"ls","path":"/tmp"}' } },
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "bash",
|
||
arguments: { command: "ls", path: "/tmp" },
|
||
});
|
||
});
|
||
|
||
it("preserves unsafe integers in stringified tool call arguments", () => {
|
||
const response = createToolCallResponse([
|
||
{
|
||
function: {
|
||
name: "send",
|
||
arguments: '{"target":9223372036854775807,"nested":{"thread":1234567890123456789}}',
|
||
},
|
||
},
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "send",
|
||
arguments: {
|
||
target: "9223372036854775807",
|
||
nested: { thread: "1234567890123456789" },
|
||
},
|
||
});
|
||
});
|
||
|
||
it("falls back to empty arguments for malformed stringified tool call arguments", () => {
|
||
const response = createToolCallResponse([
|
||
{ function: { name: "bash", arguments: '{"command":"ls"' } },
|
||
]);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expectToolCallContent(requireEntry(result.content, 0, "Ollama tool-call content"), {
|
||
name: "bash",
|
||
arguments: {},
|
||
});
|
||
});
|
||
|
||
it("sets all costs to zero for local models", () => {
|
||
const response = createAssistantResponse({ content: "ok" });
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.usage.cost).toEqual({
|
||
input: 0,
|
||
output: 0,
|
||
cacheRead: 0,
|
||
cacheWrite: 0,
|
||
total: 0,
|
||
});
|
||
});
|
||
|
||
it("records unavailable cache telemetry when Ollama omits the cache split", () => {
|
||
const response = createAssistantResponse(
|
||
{ content: "ok" },
|
||
{
|
||
prompt_eval_count: 10,
|
||
eval_count: 2,
|
||
},
|
||
);
|
||
const result = buildAssistantMessage(response, modelInfo);
|
||
expect(result.usage.cacheTelemetry).toEqual({ state: "unavailable" });
|
||
});
|
||
});
|
||
|
||
// Helper: build a ReadableStreamDefaultReader from NDJSON lines
|
||
function mockNdjsonReader(
|
||
lines: string[],
|
||
options: { trailingNewline?: boolean } = {},
|
||
): ReadableStreamDefaultReader<Uint8Array> {
|
||
const encoder = new TextEncoder();
|
||
const payload = lines.join("\n") + (options.trailingNewline === false ? "" : "\n");
|
||
let consumed = false;
|
||
return {
|
||
read: async () => {
|
||
if (consumed) {
|
||
return { done: true as const, value: undefined };
|
||
}
|
||
consumed = true;
|
||
return { done: false as const, value: encoder.encode(payload) };
|
||
},
|
||
releaseLock: () => {},
|
||
cancel: async () => {},
|
||
closed: Promise.resolve(undefined),
|
||
} as unknown as ReadableStreamDefaultReader<Uint8Array>;
|
||
}
|
||
|
||
function createPendingCancelNdjsonStream(lines: string[]) {
|
||
const encoder = new TextEncoder();
|
||
let markCancelStarted!: () => void;
|
||
let settleCancel!: () => void;
|
||
const cancelStarted = new Promise<void>((resolve) => {
|
||
markCancelStarted = resolve;
|
||
});
|
||
const cancelPending = new Promise<void>((resolve) => {
|
||
settleCancel = resolve;
|
||
});
|
||
const stream = new ReadableStream<Uint8Array>({
|
||
start(controller) {
|
||
controller.enqueue(encoder.encode(`${lines.join("\n")}\n`));
|
||
},
|
||
cancel() {
|
||
markCancelStarted();
|
||
return cancelPending;
|
||
},
|
||
});
|
||
return {
|
||
cancelPending,
|
||
cancelStarted,
|
||
reader: stream.getReader(),
|
||
settleCancel,
|
||
stream,
|
||
};
|
||
}
|
||
|
||
function createClosedNdjsonStream(lines: string[]) {
|
||
const encoder = new TextEncoder();
|
||
return new ReadableStream<Uint8Array>({
|
||
start(controller) {
|
||
controller.enqueue(encoder.encode(`${lines.join("\n")}\n`));
|
||
controller.close();
|
||
},
|
||
});
|
||
}
|
||
|
||
async function expectDoneEventContent(lines: string[], expectedContent: unknown) {
|
||
const events = await collectMockedOllamaEvents(lines);
|
||
const doneEvent = events.at(-1);
|
||
if (!doneEvent || doneEvent.type !== "done") {
|
||
throw new Error("Expected done event");
|
||
}
|
||
expect(doneEvent.message.content).toEqual(expectedContent);
|
||
}
|
||
|
||
async function expectNoParsedChunks(reader: ReadableStreamDefaultReader<Uint8Array>) {
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
}
|
||
expect(chunks).toEqual([]);
|
||
}
|
||
|
||
describe("parseNdjsonStream", () => {
|
||
it("cancels an oversized unterminated record", async () => {
|
||
const oversizedRecord = new Uint8Array(16 * 1024 * 1024 + 1).fill(0x20);
|
||
let canceled = false;
|
||
const stream = new ReadableStream<Uint8Array>({
|
||
start(controller) {
|
||
controller.enqueue(oversizedRecord);
|
||
},
|
||
cancel() {
|
||
canceled = true;
|
||
},
|
||
});
|
||
const reader = stream.getReader();
|
||
|
||
await expect(expectNoParsedChunks(reader)).rejects.toThrow(
|
||
"Ollama NDJSON record exceeds 16777216 bytes",
|
||
);
|
||
expect(canceled).toBe(true);
|
||
expect(stream.locked).toBe(false);
|
||
});
|
||
|
||
it("resets the record limit after each newline", async () => {
|
||
const legalRecord = new Uint8Array(9 * 1024 * 1024 + 1).fill(0x20);
|
||
legalRecord[legalRecord.length - 1] = 0x0a;
|
||
const reader = new ReadableStream<Uint8Array>({
|
||
start(controller) {
|
||
controller.enqueue(legalRecord);
|
||
controller.enqueue(legalRecord);
|
||
controller.close();
|
||
},
|
||
}).getReader();
|
||
|
||
await expectNoParsedChunks(reader);
|
||
});
|
||
|
||
it("does not log a dangling surrogate for a malformed complete line", async () => {
|
||
const prefix = "x".repeat(119);
|
||
const reader = mockNdjsonReader([`${prefix}😀tail`]);
|
||
|
||
await expectNoParsedChunks(reader);
|
||
|
||
expect(ollamaStreamWarnMock).toHaveBeenCalledExactlyOnceWith(
|
||
`Skipping malformed NDJSON line: ${prefix}`,
|
||
);
|
||
});
|
||
|
||
it("does not log a dangling surrogate for malformed trailing data", async () => {
|
||
const prefix = "x".repeat(119);
|
||
const reader = mockNdjsonReader([`${prefix}😀tail`], { trailingNewline: false });
|
||
|
||
await expectNoParsedChunks(reader);
|
||
|
||
expect(ollamaStreamWarnMock).toHaveBeenCalledExactlyOnceWith(
|
||
`Skipping malformed trailing data: ${prefix}`,
|
||
);
|
||
});
|
||
|
||
it("parses text-only streaming chunks", async () => {
|
||
const reader = mockNdjsonReader([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"Hello"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":" world"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":5,"eval_count":2}',
|
||
]);
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
}
|
||
expect(chunks).toHaveLength(3);
|
||
expect(requireEntry(chunks, 0, "first parsed Ollama chunk").message.content).toBe("Hello");
|
||
expect(requireEntry(chunks, 1, "second parsed Ollama chunk").message.content).toBe(" world");
|
||
expect(requireEntry(chunks, 2, "final parsed Ollama chunk").done).toBe(true);
|
||
});
|
||
|
||
it("parses tool_calls from intermediate chunk (not final)", async () => {
|
||
// Ollama sends tool_calls in done:false chunk, final done:true has no tool_calls
|
||
const reader = mockNdjsonReader([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":10,"eval_count":5}',
|
||
]);
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
}
|
||
expect(chunks).toHaveLength(2);
|
||
const firstChunk = requireEntry(chunks, 0, "first parsed Ollama chunk");
|
||
const firstToolCall = requireEntry(
|
||
expectDefined(firstChunk.message.tool_calls, "first parsed Ollama chunk tool calls"),
|
||
0,
|
||
"first parsed Ollama tool call",
|
||
);
|
||
expect(firstChunk.done).toBe(false);
|
||
expect(firstChunk.message.tool_calls).toHaveLength(1);
|
||
expect(firstToolCall.function.name).toBe("bash");
|
||
expect(requireEntry(chunks, 1, "second parsed Ollama chunk").done).toBe(true);
|
||
expect(
|
||
requireEntry(chunks, 1, "second parsed Ollama chunk").message.tool_calls,
|
||
).toBeUndefined();
|
||
});
|
||
|
||
it("accumulates tool_calls across multiple intermediate chunks", async () => {
|
||
const reader = mockNdjsonReader([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"read","arguments":{"path":"/tmp/a"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
]);
|
||
|
||
// Simulate the accumulation logic from createOllamaStreamFn
|
||
const accumulatedToolCalls: Array<{
|
||
function: { name: string; arguments: unknown };
|
||
}> = [];
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
if (chunk.message?.tool_calls) {
|
||
accumulatedToolCalls.push(...chunk.message.tool_calls);
|
||
}
|
||
}
|
||
expect(accumulatedToolCalls).toHaveLength(2);
|
||
expect(requireEntry(accumulatedToolCalls, 0, "first accumulated tool call").function.name).toBe(
|
||
"read",
|
||
);
|
||
expect(
|
||
requireEntry(accumulatedToolCalls, 1, "second accumulated tool call").function.name,
|
||
).toBe("bash");
|
||
// Final done:true chunk has no tool_calls
|
||
expect(requireEntry(chunks, 2, "final parsed Ollama chunk").message.tool_calls).toBeUndefined();
|
||
});
|
||
|
||
it("preserves unsafe integer tool arguments as exact strings", async () => {
|
||
const reader = mockNdjsonReader([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"send","arguments":{"target":1234567890123456789,"nested":{"thread":9223372036854775807}}}}]},"done":false}',
|
||
]);
|
||
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
}
|
||
|
||
const args = chunks[0]?.message.tool_calls?.[0]?.function.arguments as
|
||
| { target?: unknown; nested?: { thread?: unknown } }
|
||
| undefined;
|
||
expect(args?.target).toBe("1234567890123456789");
|
||
expect(args?.nested?.thread).toBe("9223372036854775807");
|
||
});
|
||
|
||
it("keeps safe integer tool arguments as numbers", async () => {
|
||
const reader = mockNdjsonReader([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"send","arguments":{"retries":3,"delayMs":2500}}}]},"done":false}',
|
||
]);
|
||
|
||
const chunks = [];
|
||
for await (const chunk of parseNdjsonStream(reader)) {
|
||
chunks.push(chunk);
|
||
}
|
||
|
||
const args = chunks[0]?.message.tool_calls?.[0]?.function.arguments as
|
||
| { retries?: unknown; delayMs?: unknown }
|
||
| undefined;
|
||
expect(args?.retries).toBe(3);
|
||
expect(args?.delayMs).toBe(2500);
|
||
});
|
||
|
||
it("unlocks a real stream before pending cancellation settles on early break", async () => {
|
||
const source = createPendingCancelNdjsonStream([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"one"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"two"},"done":true}',
|
||
]);
|
||
let iterationFinished = false;
|
||
|
||
const iteration = (async () => {
|
||
for await (const chunk of parseNdjsonStream(source.reader)) {
|
||
expect(chunk.message.content).toBe("one");
|
||
break;
|
||
}
|
||
iterationFinished = true;
|
||
})();
|
||
|
||
await source.cancelStarted;
|
||
await iteration;
|
||
expect(iterationFinished).toBe(true);
|
||
expect(source.stream.locked).toBe(false);
|
||
|
||
source.settleCancel();
|
||
await source.cancelPending;
|
||
});
|
||
|
||
it("preserves a consumer error while unlocking a real stream", async () => {
|
||
const source = createPendingCancelNdjsonStream([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"one"},"done":false}',
|
||
]);
|
||
const testError = new Error("consumer abort");
|
||
|
||
const iteration = (async () => {
|
||
for await (const chunk of parseNdjsonStream(source.reader)) {
|
||
void chunk;
|
||
throw testError;
|
||
}
|
||
})();
|
||
|
||
await source.cancelStarted;
|
||
await expect(iteration).rejects.toBe(testError);
|
||
expect(source.stream.locked).toBe(false);
|
||
|
||
source.settleCancel();
|
||
await source.cancelPending;
|
||
});
|
||
|
||
it("skips malformed NDJSON and unlocks after a valid terminal record", async () => {
|
||
const stream = createClosedNdjsonStream([
|
||
"not-json",
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"done"},"done":true}',
|
||
]);
|
||
const chunks = [];
|
||
|
||
for await (const chunk of parseNdjsonStream(stream.getReader())) {
|
||
chunks.push(chunk);
|
||
}
|
||
|
||
expect(chunks).toHaveLength(1);
|
||
expect(chunks[0]?.done).toBe(true);
|
||
expect(ollamaStreamWarnMock).toHaveBeenCalledWith("Skipping malformed NDJSON line: not-json");
|
||
expect(stream.locked).toBe(false);
|
||
});
|
||
});
|
||
|
||
async function withMockNdjsonFetch(
|
||
lines: string[],
|
||
run: (fetchMock: typeof fetchWithSsrFGuardMock) => Promise<void>,
|
||
): Promise<void> {
|
||
fetchWithSsrFGuardMock.mockImplementation(async () => {
|
||
const payload = lines.join("\n");
|
||
return {
|
||
response: new Response(`${payload}\n`, {
|
||
status: 200,
|
||
headers: { "Content-Type": "application/x-ndjson" },
|
||
}),
|
||
release: vi.fn(async () => undefined),
|
||
};
|
||
});
|
||
await run(fetchWithSsrFGuardMock);
|
||
}
|
||
|
||
async function withSuccessfulOllamaFetch(
|
||
run: (fetchMock: typeof fetchWithSsrFGuardMock) => Promise<void>,
|
||
): Promise<void> {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
run,
|
||
);
|
||
}
|
||
|
||
function createControlledNdjsonFetch(): {
|
||
fetchImpl: () => Promise<{ response: Response; release: () => Promise<void> }>;
|
||
pushLine: (line: string) => void;
|
||
close: () => void;
|
||
} {
|
||
const encoder = new TextEncoder();
|
||
let controller: ReadableStreamDefaultController<Uint8Array> | undefined;
|
||
const body = new ReadableStream<Uint8Array>({
|
||
start(streamController) {
|
||
controller = streamController;
|
||
},
|
||
});
|
||
return {
|
||
fetchImpl: async () => ({
|
||
response: new Response(body, {
|
||
status: 200,
|
||
headers: { "Content-Type": "application/x-ndjson" },
|
||
}),
|
||
release: vi.fn(async () => undefined),
|
||
}),
|
||
pushLine(line: string) {
|
||
if (!controller) {
|
||
throw new Error("NDJSON controller not initialized");
|
||
}
|
||
controller.enqueue(encoder.encode(`${line}\n`));
|
||
},
|
||
close() {
|
||
if (!controller) {
|
||
throw new Error("NDJSON controller not initialized");
|
||
}
|
||
controller.close();
|
||
},
|
||
};
|
||
}
|
||
|
||
function getGuardedFetchCall(fetchMock: typeof fetchWithSsrFGuardMock): GuardedFetchCall {
|
||
return (fetchMock.mock.calls.at(0)?.[0] as GuardedFetchCall | undefined) ?? { url: "" };
|
||
}
|
||
|
||
function getGuardedFetchJsonBody(
|
||
fetchMock: typeof fetchWithSsrFGuardMock,
|
||
): Record<string, unknown> {
|
||
const body = getGuardedFetchCall(fetchMock).init?.body;
|
||
if (typeof body !== "string") {
|
||
throw new Error("Expected string request body");
|
||
}
|
||
return requireRecord(JSON.parse(body), "Ollama request body");
|
||
}
|
||
|
||
function cancelTrackedResponse(
|
||
text: string,
|
||
init: ResponseInit,
|
||
): {
|
||
response: Response;
|
||
wasCanceled: () => boolean;
|
||
} {
|
||
let canceled = false;
|
||
const stream = new ReadableStream<Uint8Array>({
|
||
start(controller) {
|
||
controller.enqueue(new TextEncoder().encode(text));
|
||
},
|
||
cancel() {
|
||
canceled = true;
|
||
},
|
||
});
|
||
return {
|
||
response: new Response(stream, init),
|
||
wasCanceled: () => canceled,
|
||
};
|
||
}
|
||
|
||
async function createOllamaTestStream(params: {
|
||
baseUrl: string;
|
||
defaultHeaders?: Record<string, string>;
|
||
model?: Record<string, unknown>;
|
||
context?: Record<string, unknown>;
|
||
options?: Parameters<ReturnType<typeof createOllamaStreamFn>>[2];
|
||
}) {
|
||
const streamFn = createOllamaStreamFn(params.baseUrl, params.defaultHeaders);
|
||
return streamFn(
|
||
{
|
||
id: "qwen3:32b",
|
||
api: "ollama",
|
||
provider: "custom-ollama",
|
||
contextWindow: 131072,
|
||
...params.model,
|
||
} as unknown as Parameters<typeof streamFn>[0],
|
||
(params.context ?? {
|
||
messages: [{ role: "user", content: "hello" }],
|
||
}) as unknown as Parameters<typeof streamFn>[1],
|
||
(params.options ?? {}) as unknown as Parameters<typeof streamFn>[2],
|
||
);
|
||
}
|
||
|
||
async function collectStreamEvents<T>(stream: AsyncIterable<T>): Promise<T[]> {
|
||
const events: T[] = [];
|
||
for await (const event of stream) {
|
||
events.push(event);
|
||
}
|
||
return events;
|
||
}
|
||
|
||
type OllamaStreamEvent =
|
||
Awaited<ReturnType<typeof createOllamaTestStream>> extends AsyncIterable<infer Event>
|
||
? Event
|
||
: never;
|
||
|
||
async function collectMockedOllamaEvents(
|
||
lines: string[],
|
||
params: Parameters<typeof createOllamaTestStream>[0] = {
|
||
baseUrl: "http://ollama-host:11434",
|
||
},
|
||
): Promise<OllamaStreamEvent[]> {
|
||
let events: OllamaStreamEvent[] | undefined;
|
||
await withMockNdjsonFetch(lines, async () => {
|
||
events = await collectStreamEvents(await createOllamaTestStream(params));
|
||
});
|
||
return expectDefined(events, "mocked Ollama stream events");
|
||
}
|
||
|
||
async function expectSuccessfulOllamaRequest(
|
||
params: Parameters<typeof createOllamaTestStream>[0],
|
||
verify: (observation: {
|
||
body: Record<string, unknown>;
|
||
fetchMock: typeof fetchWithSsrFGuardMock;
|
||
request: GuardedFetchCall;
|
||
}) => void | Promise<void>,
|
||
): Promise<void> {
|
||
await withSuccessfulOllamaFetch(async (fetchMock) => {
|
||
const events = await collectStreamEvents(await createOllamaTestStream(params));
|
||
expect(events.at(-1)?.type).toBe("done");
|
||
await verify({
|
||
body: getGuardedFetchJsonBody(fetchMock),
|
||
fetchMock,
|
||
request: getGuardedFetchCall(fetchMock),
|
||
});
|
||
});
|
||
}
|
||
|
||
async function nextEventWithin<T>(
|
||
iterator: AsyncIterator<T>,
|
||
timeoutMs = 100,
|
||
): Promise<IteratorResult<T> | "timeout"> {
|
||
let timer: NodeJS.Timeout | undefined;
|
||
try {
|
||
return await Promise.race([
|
||
iterator.next(),
|
||
new Promise<"timeout">((resolve) => {
|
||
timer = setTimeout(() => resolve("timeout"), timeoutMs);
|
||
}),
|
||
]);
|
||
} finally {
|
||
if (timer) {
|
||
clearTimeout(timer);
|
||
}
|
||
}
|
||
}
|
||
|
||
describe("createOllamaStreamFn streaming events", () => {
|
||
it("reports the successful HTTP response before streaming events", async () => {
|
||
const timeline: string[] = [];
|
||
const onResponse = vi.fn((response, callbackModel) => {
|
||
timeline.push("response");
|
||
expect(response).toEqual({
|
||
status: 200,
|
||
headers: {
|
||
"content-type": "application/x-ndjson",
|
||
"x-ollama-request-id": "req-1",
|
||
},
|
||
});
|
||
expect(callbackModel.id).toBe("qwen3:32b");
|
||
});
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: new Response(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
].join("\n"),
|
||
{
|
||
status: 200,
|
||
headers: {
|
||
"Content-Type": "application/x-ndjson",
|
||
"X-Ollama-Request-Id": "req-1",
|
||
},
|
||
},
|
||
),
|
||
release: vi.fn(async () => undefined),
|
||
});
|
||
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: { onResponse },
|
||
});
|
||
for await (const event of stream) {
|
||
timeline.push(event.type);
|
||
}
|
||
|
||
expect(onResponse).toHaveBeenCalledTimes(1);
|
||
expect(timeline).toEqual(["response", "start", "text_start", "text_delta", "text_end", "done"]);
|
||
});
|
||
|
||
it("reports failed HTTP responses before the stream error", async () => {
|
||
const timeline: string[] = [];
|
||
const onResponse = vi.fn(() => {
|
||
timeline.push("response");
|
||
});
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: new Response("rate limited", {
|
||
status: 429,
|
||
headers: { "Retry-After": "30" },
|
||
}),
|
||
release: vi.fn(async () => undefined),
|
||
});
|
||
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: { onResponse },
|
||
});
|
||
for await (const event of stream) {
|
||
timeline.push(event.type);
|
||
}
|
||
|
||
expect(onResponse).toHaveBeenCalledWith(
|
||
{ status: 429, headers: { "content-type": "text/plain;charset=UTF-8", "retry-after": "30" } },
|
||
expect.objectContaining({ id: "qwen3:32b" }),
|
||
);
|
||
expect(timeline).toEqual(["response", "error"]);
|
||
});
|
||
|
||
it("does not wait for unread response cancellation when the response hook fails", async () => {
|
||
const source = createPendingCancelNdjsonStream([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":true}',
|
||
]);
|
||
source.reader.releaseLock();
|
||
const release = vi.fn(async () => undefined);
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: new Response(source.stream, {
|
||
status: 200,
|
||
headers: { "Content-Type": "application/x-ndjson" },
|
||
}),
|
||
release,
|
||
});
|
||
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: {
|
||
onResponse: () => {
|
||
throw new Error("response hook failed");
|
||
},
|
||
},
|
||
});
|
||
const event = await nextEventWithin(stream[Symbol.asyncIterator]());
|
||
await source.cancelStarted;
|
||
source.settleCancel();
|
||
await source.cancelPending;
|
||
|
||
expect(event).not.toBe("timeout");
|
||
if (event !== "timeout") {
|
||
expect(event.done).toBe(false);
|
||
expect(event.value).toMatchObject({ type: "error", reason: "error" });
|
||
}
|
||
expect(release).toHaveBeenCalledTimes(1);
|
||
});
|
||
|
||
it("stops waiting for the response hook when the request is aborted", async () => {
|
||
let markHookStarted: () => void = () => undefined;
|
||
const hookStarted = new Promise<void>((resolve) => {
|
||
markHookStarted = resolve;
|
||
});
|
||
const onResponse = vi.fn(async () => {
|
||
markHookStarted();
|
||
await new Promise<void>(() => {
|
||
// Keep the hook pending so the request signal must end the stream.
|
||
});
|
||
});
|
||
const release = vi.fn(async () => undefined);
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: new Response(
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":true}',
|
||
{
|
||
status: 200,
|
||
headers: { "Content-Type": "application/x-ndjson" },
|
||
},
|
||
),
|
||
release,
|
||
});
|
||
const abortController = new AbortController();
|
||
|
||
const eventsPromise = collectStreamEvents(
|
||
await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: { onResponse, signal: abortController.signal },
|
||
}),
|
||
);
|
||
await hookStarted;
|
||
abortController.abort();
|
||
const events = await eventsPromise;
|
||
|
||
expect(events).toHaveLength(1);
|
||
expect(events[0]).toMatchObject({ type: "error", reason: "aborted" });
|
||
expect(release).toHaveBeenCalledTimes(1);
|
||
});
|
||
|
||
it("does not resume response handling after the hook resolves concurrently with abort", async () => {
|
||
let markHookStarted: () => void = () => undefined;
|
||
const hookStarted = new Promise<void>((resolve) => {
|
||
markHookStarted = resolve;
|
||
});
|
||
let settleHook: () => void = () => undefined;
|
||
const hookPending = new Promise<void>((resolve) => {
|
||
settleHook = resolve;
|
||
});
|
||
const getReader = vi.fn(() => ({
|
||
read: vi.fn(async () => ({ done: true as const, value: undefined })),
|
||
cancel: vi.fn(async () => undefined),
|
||
releaseLock: vi.fn(),
|
||
}));
|
||
const cancel = vi.fn(async () => undefined);
|
||
const release = vi.fn(async () => undefined);
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: {
|
||
status: 200,
|
||
ok: true,
|
||
headers: new Headers({ "Content-Type": "application/x-ndjson" }),
|
||
body: { getReader, cancel },
|
||
} as unknown as Response,
|
||
release,
|
||
});
|
||
const abortController = new AbortController();
|
||
|
||
const eventsPromise = collectStreamEvents(
|
||
await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: {
|
||
onResponse: () => {
|
||
markHookStarted();
|
||
return hookPending;
|
||
},
|
||
signal: abortController.signal,
|
||
},
|
||
}),
|
||
);
|
||
await hookStarted;
|
||
await Promise.resolve();
|
||
void hookPending.then(() => abortController.abort());
|
||
settleHook();
|
||
const events = await eventsPromise;
|
||
|
||
expect(events).toHaveLength(1);
|
||
expect(events[0]).toMatchObject({ type: "error", reason: "aborted" });
|
||
expect(getReader).not.toHaveBeenCalled();
|
||
expect(cancel).toHaveBeenCalledTimes(1);
|
||
expect(release).toHaveBeenCalledTimes(1);
|
||
});
|
||
|
||
it("emits start, text_start, text_delta, text_end, done for text responses", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"Hello"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":" world"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":5,"eval_count":2}',
|
||
]);
|
||
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual(["start", "text_start", "text_delta", "text_delta", "text_end", "done"]);
|
||
|
||
// text_delta events carry incremental deltas
|
||
const deltas = events.filter((e) => e.type === "text_delta");
|
||
expect(deltas[0]?.contentIndex).toBe(0);
|
||
expect(deltas[0]?.delta).toBe("Hello");
|
||
expect(deltas[1]?.contentIndex).toBe(0);
|
||
expect(deltas[1]?.delta).toBe(" world");
|
||
|
||
// text_end carries the full accumulated content
|
||
const textEnd = events.find((e) => e.type === "text_end");
|
||
expect(textEnd?.contentIndex).toBe(0);
|
||
expect(textEnd?.content).toBe("Hello world");
|
||
|
||
// start/text_start carry empty partials (before any content accumulates)
|
||
const startEvent = events.find((e) => e.type === "start");
|
||
expect(startEvent?.partial.content).toStrictEqual([]);
|
||
const textStartEvent = events.find((e) => e.type === "text_start");
|
||
expect(textStartEvent?.partial.content).toStrictEqual([]);
|
||
|
||
// text_delta events stay lightweight; text_end/done carry the full snapshot.
|
||
expect(deltas[0]).not.toHaveProperty("partial");
|
||
expect(deltas[1]).not.toHaveProperty("partial");
|
||
|
||
// done event contains the final message
|
||
const doneEvent = events.at(-1);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.content).toEqual([{ type: "text", text: "Hello world" }]);
|
||
}
|
||
});
|
||
|
||
it("streams the complete lifecycle for tool-call-only responses", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":10,"eval_count":5}',
|
||
]);
|
||
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual(["start", "toolcall_start", "toolcall_delta", "toolcall_end", "done"]);
|
||
expect(events[1]).toMatchObject({
|
||
type: "toolcall_start",
|
||
contentIndex: 0,
|
||
partial: { content: [{ type: "toolCall", name: "bash", arguments: {} }] },
|
||
});
|
||
expect(events[2]).toMatchObject({
|
||
type: "toolcall_delta",
|
||
contentIndex: 0,
|
||
delta: '{"command":"ls"}',
|
||
});
|
||
expect(events[3]).toMatchObject({
|
||
type: "toolcall_end",
|
||
contentIndex: 0,
|
||
toolCall: { name: "bash", arguments: { command: "ls" } },
|
||
});
|
||
const doneEvent = requireEntry(events, 4, "tool-call-only done event");
|
||
if (doneEvent.type === "done") {
|
||
expect(doneEvent.reason).toBe("toolUse");
|
||
expect(doneEvent.message.content[0]).toMatchObject({
|
||
type: "toolCall",
|
||
id: events[3]?.type === "toolcall_end" ? events[3].toolCall.id : undefined,
|
||
});
|
||
}
|
||
});
|
||
|
||
it("estimates usage when the final Ollama chunk omits counters", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"Estimated answer"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
]);
|
||
|
||
const doneEvent = events.at(-1);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.usage.input).toBeGreaterThan(0);
|
||
expect(doneEvent.message.usage.output).toBeGreaterThan(0);
|
||
expect(doneEvent.message.usage.totalTokens).toBeGreaterThan(0);
|
||
}
|
||
});
|
||
|
||
it("counts image payloads in prompt usage estimates when Ollama omits counters", async () => {
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"vision answer"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "llava" },
|
||
context: {
|
||
messages: [{ role: "user", content: [{ type: "image", data: "a".repeat(400) }] }],
|
||
},
|
||
},
|
||
);
|
||
const doneEvent = events.at(-1);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.usage.input).toBeGreaterThan(50);
|
||
}
|
||
});
|
||
|
||
it("emits text streaming events before done for mixed text + tool responses", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"Let me check."},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":10,"eval_count":5}',
|
||
]);
|
||
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual([
|
||
"start",
|
||
"text_start",
|
||
"text_delta",
|
||
"text_end",
|
||
"toolcall_start",
|
||
"toolcall_delta",
|
||
"toolcall_end",
|
||
"done",
|
||
]);
|
||
expect(events[5]).toMatchObject({
|
||
type: "toolcall_delta",
|
||
contentIndex: 1,
|
||
delta: '{"command":"ls"}',
|
||
});
|
||
const doneEvent = events.at(-1);
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.reason).toBe("toolUse");
|
||
}
|
||
});
|
||
|
||
it("streams multiple native calls with stable provider ids across chunks", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"id":"call-read","function":{"name":"read","arguments":{"path":"/tmp/a"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"id":"call-bash","function":{"name":"bash","arguments":"{\\"command\\":\\"ls\\"}"}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
]);
|
||
|
||
expect(events.map((event) => event.type)).toEqual([
|
||
"start",
|
||
"toolcall_start",
|
||
"toolcall_delta",
|
||
"toolcall_end",
|
||
"toolcall_start",
|
||
"toolcall_delta",
|
||
"toolcall_end",
|
||
"done",
|
||
]);
|
||
const toolCallEnds = events.filter((event) => event.type === "toolcall_end");
|
||
expect(toolCallEnds).toMatchObject([
|
||
{
|
||
contentIndex: 0,
|
||
toolCall: { id: "call-read", name: "read", arguments: { path: "/tmp/a" } },
|
||
},
|
||
{
|
||
contentIndex: 1,
|
||
toolCall: { id: "call-bash", name: "bash", arguments: { command: "ls" } },
|
||
},
|
||
]);
|
||
expect(events.filter((event) => event.type === "toolcall_delta")).toMatchObject([
|
||
{ contentIndex: 0, delta: '{"path":"/tmp/a"}' },
|
||
{ contentIndex: 1, delta: '{"command":"ls"}' },
|
||
]);
|
||
expect(events.filter((event) => event.type === "toolcall_start")).toMatchObject([
|
||
{ partial: { content: [{ arguments: {} }] } },
|
||
{
|
||
partial: {
|
||
content: [{ arguments: { path: "/tmp/a" } }, { arguments: {} }],
|
||
},
|
||
},
|
||
]);
|
||
const done = events.at(-1);
|
||
if (done?.type !== "done") {
|
||
throw new Error("missing terminal Ollama message");
|
||
}
|
||
expect(done.message.content).toMatchObject([
|
||
{ type: "toolCall", id: "call-read" },
|
||
{ type: "toolCall", id: "call-bash" },
|
||
]);
|
||
});
|
||
|
||
it("does not stream non-executable calls from a token-limited final chunk", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":true,"done_reason":"length"}',
|
||
]);
|
||
|
||
expect(events.map((event) => event.type)).toEqual(["done"]);
|
||
expect(events[0]).toMatchObject({
|
||
type: "done",
|
||
reason: "length",
|
||
message: { content: [], stopReason: "length" },
|
||
});
|
||
});
|
||
|
||
it("never exposes an intermediate native call invalidated by a later length terminal", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"done_reason":"length"}',
|
||
]);
|
||
|
||
expect(events.map((event) => event.type)).toEqual(["done"]);
|
||
expect(events[0]).toMatchObject({
|
||
type: "done",
|
||
reason: "length",
|
||
message: { content: [], stopReason: "length" },
|
||
});
|
||
});
|
||
|
||
it("emits text_end as soon as Ollama switches from text to tool calls", async () => {
|
||
const controlledFetch = createControlledNdjsonFetch();
|
||
fetchWithSsrFGuardMock.mockImplementation(controlledFetch.fetchImpl);
|
||
|
||
try {
|
||
const stream = await createOllamaTestStream({ baseUrl: "http://ollama-host:11434" });
|
||
const iterator = stream[Symbol.asyncIterator]();
|
||
|
||
controlledFetch.pushLine(
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"Let me check."},"done":false}',
|
||
);
|
||
|
||
const startEvent = await nextEventWithin(iterator);
|
||
const textStartEvent = await nextEventWithin(iterator);
|
||
const textDeltaEvent = await nextEventWithin(iterator);
|
||
|
||
expect(startEvent).not.toBe("timeout");
|
||
expect(textStartEvent).not.toBe("timeout");
|
||
expect(textDeltaEvent).not.toBe("timeout");
|
||
expectIteratorEvent(startEvent, { type: "start", done: false });
|
||
expectIteratorEvent(textStartEvent, { type: "text_start", done: false });
|
||
expectIteratorEvent(textDeltaEvent, {
|
||
type: "text_delta",
|
||
delta: "Let me check.",
|
||
done: false,
|
||
});
|
||
|
||
controlledFetch.pushLine(
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
);
|
||
|
||
const textEndEvent = await nextEventWithin(iterator);
|
||
expect(textEndEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textEndEvent, {
|
||
type: "text_end",
|
||
content: "Let me check.",
|
||
done: false,
|
||
});
|
||
if (textEndEvent !== "timeout") {
|
||
const textEndValue = requireRecord(textEndEvent.value, "text_end value");
|
||
expect(textEndValue.contentIndex).toBe(0);
|
||
expect(requireRecord(textEndValue.partial, "text_end partial").content).toEqual([
|
||
{ type: "text", text: "Let me check." },
|
||
]);
|
||
}
|
||
|
||
controlledFetch.pushLine(
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":10,"eval_count":5}',
|
||
);
|
||
controlledFetch.close();
|
||
|
||
const toolCallStartEvent = await nextEventWithin(iterator);
|
||
const toolCallDeltaEvent = await nextEventWithin(iterator);
|
||
const toolCallEndEvent = await nextEventWithin(iterator);
|
||
expect(toolCallStartEvent).not.toBe("timeout");
|
||
expect(toolCallDeltaEvent).not.toBe("timeout");
|
||
expect(toolCallEndEvent).not.toBe("timeout");
|
||
expectIteratorEvent(toolCallStartEvent, { type: "toolcall_start", done: false });
|
||
expectIteratorEvent(toolCallDeltaEvent, {
|
||
type: "toolcall_delta",
|
||
delta: '{"command":"ls"}',
|
||
done: false,
|
||
});
|
||
expectIteratorEvent(toolCallEndEvent, { type: "toolcall_end", done: false });
|
||
|
||
const doneEvent = await nextEventWithin(iterator);
|
||
expect(doneEvent).not.toBe("timeout");
|
||
if (doneEvent !== "timeout" && doneEvent.done === false) {
|
||
expectIteratorEvent(doneEvent, { type: "done", done: false });
|
||
expect(requireRecord(doneEvent.value, "done value").reason).toBe("toolUse");
|
||
|
||
const streamEnd = await nextEventWithin(iterator);
|
||
expect(streamEnd).not.toBe("timeout");
|
||
expectIteratorEvent(streamEnd, { done: true });
|
||
} else {
|
||
expectIteratorEvent(doneEvent, { done: true });
|
||
}
|
||
} finally {
|
||
fetchWithSsrFGuardMock.mockReset();
|
||
}
|
||
});
|
||
|
||
it("emits error without text_end when stream fails mid-response", async () => {
|
||
// Simulate a stream that sends one content chunk then ends without done:true.
|
||
// The stream function throws "Ollama API stream ended without a final response".
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"partial"},"done":false}',
|
||
]);
|
||
|
||
const types = events.map((e) => e.type);
|
||
// Should have streaming events for the partial content, then error (no text_end).
|
||
expect(types).toEqual(["start", "text_start", "text_delta", "error"]);
|
||
const errorEvent = events.at(-1);
|
||
expect(errorEvent?.type).toBe("error");
|
||
if (errorEvent?.type === "error") {
|
||
expect(errorEvent.error.errorMessage).toBe(OLLAMA_INCOMPLETE_STREAM_ERROR);
|
||
}
|
||
});
|
||
|
||
it("emits an error instead of accepting garbled Kimi visible text", async () => {
|
||
const garbled =
|
||
'$$"##"%#"##"####""$""""##""$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$' +
|
||
'#"$"$"""$""""#$"""$"""%"%###"""#%""""&"#"""$"""#"#""""%#""""&"#"""$"""$"""#%"""';
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.5:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: garbled },
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.5:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.5:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual(["error"]);
|
||
const errorEvent = events.at(-1);
|
||
expect(errorEvent?.type).toBe("error");
|
||
if (errorEvent?.type === "error") {
|
||
expect(errorEvent.error.errorMessage).toContain("garbled visible text");
|
||
}
|
||
});
|
||
|
||
it("buffers Kimi inline reasoning until the streaming boundary is safe", async () => {
|
||
const controlledFetch = createControlledNdjsonFetch();
|
||
fetchWithSsrFGuardMock.mockImplementation(controlledFetch.fetchImpl);
|
||
|
||
try {
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
});
|
||
const iterator = stream[Symbol.asyncIterator]();
|
||
|
||
controlledFetch.pushLine(
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: {
|
||
role: "assistant",
|
||
content:
|
||
"The user is asking for a short answer. I should reason privately before answering. I need to avoid showing this planning text to the user.",
|
||
},
|
||
done: false,
|
||
}),
|
||
);
|
||
|
||
const pendingStartEvent = iterator.next();
|
||
expect(
|
||
await Promise.race([
|
||
pendingStartEvent.then(() => "event" as const),
|
||
new Promise<"timeout">((resolve) => {
|
||
setTimeout(() => resolve("timeout"), 100);
|
||
}),
|
||
]),
|
||
).toBe("timeout");
|
||
|
||
controlledFetch.pushLine(
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: " ️ OK." },
|
||
done: false,
|
||
}),
|
||
);
|
||
|
||
const startEvent = await pendingStartEvent;
|
||
expectIteratorEvent(startEvent, { type: "start", done: false });
|
||
|
||
const textStartEvent = await nextEventWithin(iterator);
|
||
expect(textStartEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textStartEvent, { type: "text_start", done: false });
|
||
|
||
const textDeltaEvent = await nextEventWithin(iterator);
|
||
expect(textDeltaEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textDeltaEvent, {
|
||
type: "text_delta",
|
||
delta: "OK.",
|
||
done: false,
|
||
});
|
||
if (textDeltaEvent !== "timeout" && textDeltaEvent.done === false) {
|
||
const value = requireRecord(textDeltaEvent.value, "text_delta value");
|
||
expect(JSON.stringify(value)).not.toContain("The user is asking");
|
||
}
|
||
|
||
controlledFetch.pushLine(
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
);
|
||
controlledFetch.close();
|
||
|
||
const textEndEvent = await nextEventWithin(iterator);
|
||
expect(textEndEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textEndEvent, {
|
||
type: "text_end",
|
||
content: "OK.",
|
||
done: false,
|
||
});
|
||
|
||
const doneEvent = await nextEventWithin(iterator);
|
||
expect(doneEvent).not.toBe("timeout");
|
||
expectIteratorEvent(doneEvent, { type: "done", done: false });
|
||
if (doneEvent !== "timeout" && doneEvent.done === false) {
|
||
const value = requireRecord(doneEvent.value, "done value");
|
||
expect(JSON.stringify(value)).not.toContain("The user is asking");
|
||
}
|
||
|
||
const streamEnd = await nextEventWithin(iterator);
|
||
expect(streamEnd).not.toBe("timeout");
|
||
expectIteratorEvent(streamEnd, { done: true });
|
||
} finally {
|
||
fetchWithSsrFGuardMock.mockReset();
|
||
}
|
||
});
|
||
|
||
it("streams marker-less Kimi answers after the bounded sanitizer window", async () => {
|
||
const controlledFetch = createControlledNdjsonFetch();
|
||
fetchWithSsrFGuardMock.mockImplementation(controlledFetch.fetchImpl);
|
||
|
||
try {
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
});
|
||
const iterator = stream[Symbol.asyncIterator]();
|
||
const visibleAnswer = "This is a normal Kimi cloud answer without hidden reasoning. ".repeat(
|
||
12,
|
||
);
|
||
|
||
controlledFetch.pushLine(
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: visibleAnswer },
|
||
done: false,
|
||
}),
|
||
);
|
||
|
||
const startEvent = await nextEventWithin(iterator);
|
||
expect(startEvent).not.toBe("timeout");
|
||
expectIteratorEvent(startEvent, { type: "start", done: false });
|
||
|
||
const textStartEvent = await nextEventWithin(iterator);
|
||
expect(textStartEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textStartEvent, { type: "text_start", done: false });
|
||
|
||
const textDeltaEvent = await nextEventWithin(iterator);
|
||
expect(textDeltaEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textDeltaEvent, {
|
||
type: "text_delta",
|
||
delta: visibleAnswer,
|
||
done: false,
|
||
});
|
||
|
||
controlledFetch.pushLine(
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
);
|
||
controlledFetch.close();
|
||
|
||
const textEndEvent = await nextEventWithin(iterator);
|
||
expect(textEndEvent).not.toBe("timeout");
|
||
expectIteratorEvent(textEndEvent, {
|
||
type: "text_end",
|
||
content: visibleAnswer,
|
||
done: false,
|
||
});
|
||
|
||
const doneEvent = await nextEventWithin(iterator);
|
||
expect(doneEvent).not.toBe("timeout");
|
||
expectIteratorEvent(doneEvent, { type: "done", done: false });
|
||
|
||
const streamEnd = await nextEventWithin(iterator);
|
||
expect(streamEnd).not.toBe("timeout");
|
||
expectIteratorEvent(streamEnd, { done: true });
|
||
} finally {
|
||
fetchWithSsrFGuardMock.mockReset();
|
||
}
|
||
});
|
||
|
||
it("keeps Kimi deltas append-only after the bounded sanitizer window is bypassed", async () => {
|
||
const longPrefix = "This Kimi cloud output has streamed past the sanitizer window. ".repeat(10);
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: longPrefix },
|
||
done: false,
|
||
}),
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: " ️ OK." },
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
const deltas = events.filter((event) => event.type === "text_delta");
|
||
expect(deltas.map((event) => event.delta)).toEqual([longPrefix, " ️ OK."]);
|
||
|
||
const rawText = `${longPrefix} ️ OK.`;
|
||
const textEnd = events.find((event) => event.type === "text_end");
|
||
expect(textEnd?.content).toBe(rawText);
|
||
const doneEvent = events.find((event) => event.type === "done");
|
||
expect(doneEvent?.message.content).toEqual([{ type: "text", text: rawText }]);
|
||
});
|
||
|
||
it("does not reject punctuation-heavy text from unrelated Ollama models", async () => {
|
||
const punctuationHeavy =
|
||
'$$"##"%#"##"####""$""""##""$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$' +
|
||
'#"$"$"""$""""#$"""$"""%"%###"""#%""""&"#"""$"""#"#""""%#""""&"#"""$"""$"""#%"""';
|
||
const events = await collectMockedOllamaEvents([
|
||
JSON.stringify({
|
||
model: "qwen3:32b",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: punctuationHeavy },
|
||
done: false,
|
||
}),
|
||
'{"model":"qwen3:32b","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
]);
|
||
|
||
expect(events.map((e) => e.type)).toEqual([
|
||
"start",
|
||
"text_start",
|
||
"text_delta",
|
||
"text_end",
|
||
"done",
|
||
]);
|
||
});
|
||
|
||
it("emits a single text_delta for single-chunk responses", async () => {
|
||
const events = await collectMockedOllamaEvents([
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"one shot"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
]);
|
||
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual(["start", "text_start", "text_delta", "text_end", "done"]);
|
||
|
||
const textStart = events.find((e) => e.type === "text_start");
|
||
expect(textStart?.partial.content).toEqual([]);
|
||
const delta = events.find((e) => e.type === "text_delta");
|
||
expect(delta?.delta).toBe("one shot");
|
||
});
|
||
|
||
it("sanitizes Kimi inline reasoning in text_delta, text_end, and done output", async () => {
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: {
|
||
role: "assistant",
|
||
content:
|
||
"I should think privately and not leak this planning text in the answer. I need to keep deciding what to say next. ️Final answer",
|
||
},
|
||
done: false,
|
||
}),
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: " only." },
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":3,"eval_count":4}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
const types = events.map((e) => e.type);
|
||
expect(types).toEqual(["start", "text_start", "text_delta", "text_delta", "text_end", "done"]);
|
||
|
||
const textStart = events.find((e) => e.type === "text_start");
|
||
expect(textStart?.partial.content).toEqual([]);
|
||
const deltas = events.filter((e) => e.type === "text_delta");
|
||
expect(deltas).toHaveLength(2);
|
||
expect(deltas[0]?.delta).toBe("Final answer");
|
||
expect(deltas[1]?.delta).toBe(" only.");
|
||
expect(deltas[0]).not.toHaveProperty("partial");
|
||
expect(deltas[1]).not.toHaveProperty("partial");
|
||
|
||
const textEnd = events.find((e) => e.type === "text_end");
|
||
expect(textEnd?.content).toBe("Final answer only.");
|
||
expect(textEnd?.partial.content).toEqual([{ type: "text", text: "Final answer only." }]);
|
||
|
||
const doneEvent = events.at(-1);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.content).toEqual([{ type: "text", text: "Final answer only." }]);
|
||
}
|
||
});
|
||
|
||
it("does not re-sanitize visible Kimi stream output before done", async () => {
|
||
const visibleAnswer =
|
||
"This visible answer is intentionally long enough to look like a reasoning prefix if it is sanitized a second time. ️ keep this marker visible.";
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: {
|
||
role: "assistant",
|
||
content:
|
||
"I should think privately and not leak this planning text in the answer. I need to keep deciding what to say next. ️",
|
||
},
|
||
done: false,
|
||
}),
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: visibleAnswer },
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":3,"eval_count":4}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
const textEnd = events.find((event) => event.type === "text_end");
|
||
const doneEvent = events.at(-1);
|
||
|
||
expect(textEnd?.content).toBe(visibleAnswer);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.content).toEqual([{ type: "text", text: visibleAnswer }]);
|
||
}
|
||
});
|
||
|
||
it("does not leak Kimi inline reasoning when a boundary is followed by tool calls only", async () => {
|
||
const hiddenPrefix =
|
||
"I should think privately and not leak this planning text in the answer. " +
|
||
"I need to keep deciding what tool to call before showing any visible text.";
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: `${hiddenPrefix} ️ ` },
|
||
done: false,
|
||
}),
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: {
|
||
role: "assistant",
|
||
content: "",
|
||
tool_calls: [{ function: { name: "bash", arguments: { command: "ls" } } }],
|
||
},
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":20,"eval_count":40}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
|
||
expect(events.map((event) => event.type)).toEqual([
|
||
"start",
|
||
"toolcall_start",
|
||
"toolcall_delta",
|
||
"toolcall_end",
|
||
"done",
|
||
]);
|
||
expect(JSON.stringify(events)).not.toContain("I should think privately");
|
||
const doneEvent = events.at(-1);
|
||
expect(doneEvent?.type).toBe("done");
|
||
if (doneEvent?.type === "done") {
|
||
expect(doneEvent.message.content).toEqual([
|
||
{
|
||
type: "toolCall",
|
||
id: expect.any(String),
|
||
name: "bash",
|
||
arguments: { command: "ls" },
|
||
},
|
||
]);
|
||
expect(JSON.stringify(doneEvent)).not.toContain("I should think privately");
|
||
}
|
||
});
|
||
|
||
it("flushes buffered visible Kimi text before streaming its native tool call", async () => {
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":"Visible answer"},"done":false}',
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[{"function":{"name":"bash","arguments":{"command":"ls"}}}]},"done":false}',
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
|
||
expect(events.map((event) => event.type)).toEqual([
|
||
"start",
|
||
"text_start",
|
||
"text_delta",
|
||
"text_end",
|
||
"toolcall_start",
|
||
"toolcall_delta",
|
||
"toolcall_end",
|
||
"done",
|
||
]);
|
||
expect(events[2]).toMatchObject({ type: "text_delta", delta: "Visible answer" });
|
||
expect(events[4]).toMatchObject({ type: "toolcall_start", contentIndex: 1 });
|
||
expect(events[6]).toMatchObject({ type: "toolcall_end", contentIndex: 1 });
|
||
expect(events.at(-1)).toMatchObject({
|
||
type: "done",
|
||
message: {
|
||
content: [
|
||
{ type: "text", text: "Visible answer" },
|
||
{ type: "toolCall", name: "bash" },
|
||
],
|
||
},
|
||
});
|
||
});
|
||
|
||
it("does not reveal buffered Kimi reasoning for an empty tool-call chunk", async () => {
|
||
const hiddenPrefix =
|
||
"I should think privately and not leak this planning text in the answer. " +
|
||
"I need to keep deciding what to say next.";
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
JSON.stringify({
|
||
model: "kimi-k2.6:cloud",
|
||
created_at: "t",
|
||
message: { role: "assistant", content: hiddenPrefix },
|
||
done: false,
|
||
}),
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":"","tool_calls":[]},"done":false}',
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":" ️ Visible answer"},"done":false}',
|
||
'{"model":"kimi-k2.6:cloud","created_at":"t","message":{"role":"assistant","content":""},"done":true}',
|
||
],
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: { id: "kimi-k2.6:cloud", provider: "ollama" },
|
||
},
|
||
);
|
||
|
||
expect(events.map((event) => event.type)).toEqual([
|
||
"start",
|
||
"text_start",
|
||
"text_delta",
|
||
"text_end",
|
||
"done",
|
||
]);
|
||
expect(events[2]).toMatchObject({ type: "text_delta", delta: "Visible answer" });
|
||
expect(JSON.stringify(events)).not.toContain("I should think privately");
|
||
});
|
||
});
|
||
|
||
describe("createOllamaStreamFn", () => {
|
||
it("normalizes /v1 baseUrl and maps maxTokens + signal", async () => {
|
||
const signal = new AbortController().signal;
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434/v1/", options: { maxTokens: 123, signal } },
|
||
({ body, fetchMock, request }) => {
|
||
expect(fetchMock).toHaveBeenCalledTimes(1);
|
||
expect(request.url).toBe("http://ollama-host:11434/api/chat");
|
||
expect(request.auditContext).toBe("ollama-stream.chat");
|
||
expect(request.signal).toBe(signal);
|
||
expect(request.init?.signal).toBeUndefined();
|
||
const options = requireRecord(body.options, "Ollama request options");
|
||
expect(options.num_ctx).toBeUndefined();
|
||
expect(options.num_predict).toBe(123);
|
||
},
|
||
);
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "forwards request stop sequences as native Ollama options",
|
||
model: {},
|
||
stop: ["END", "DONE"],
|
||
expectedStop: ["END", "DONE"],
|
||
},
|
||
{
|
||
name: "lets request stop sequences override configured model stop sequences",
|
||
model: { params: { stop: ["MODEL"] } },
|
||
stop: ["REQUEST"],
|
||
expectedStop: ["REQUEST"],
|
||
},
|
||
{
|
||
name: "keeps configured model stop sequences when request stops are empty",
|
||
model: { params: { stop: ["MODEL"] } },
|
||
stop: [],
|
||
expectedStop: ["MODEL"],
|
||
},
|
||
])("$name", async ({ model, stop, expectedStop }) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", model, options: { stop } },
|
||
({ body }) =>
|
||
expect(requireRecord(body.options, "Ollama request options").stop).toEqual(expectedStop),
|
||
);
|
||
});
|
||
|
||
it("awaits asynchronous payload mutations before dispatching native Ollama requests", async () => {
|
||
await withSuccessfulOllamaFetch(async (fetchMock) => {
|
||
let releasePayload: (() => void) | undefined;
|
||
const payloadGate = new Promise<void>((resolve) => {
|
||
releasePayload = resolve;
|
||
});
|
||
const onPayload = vi.fn(async (payload: unknown) => {
|
||
await payloadGate;
|
||
requireRecord(payload, "Ollama request payload").model = "patched-model";
|
||
});
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: { onPayload },
|
||
});
|
||
|
||
await vi.waitFor(() => expect(onPayload).toHaveBeenCalledTimes(1));
|
||
expect(fetchMock).not.toHaveBeenCalled();
|
||
expectDefined(releasePayload, "pending Ollama payload hook")();
|
||
const events = await collectStreamEvents(stream);
|
||
|
||
expect(events.at(-1)?.type).toBe("done");
|
||
expect(getGuardedFetchJsonBody(fetchMock).model).toBe("patched-model");
|
||
});
|
||
});
|
||
|
||
it("dispatches asynchronous payload replacements for native Ollama requests", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: {
|
||
onPayload: async (payload) => {
|
||
await Promise.resolve();
|
||
const current = requireRecord(payload, "Ollama request payload");
|
||
return {
|
||
...current,
|
||
model: "replacement-model",
|
||
options: {
|
||
...requireRecord(current.options, "Ollama request options"),
|
||
stop: ["REPLACEMENT"],
|
||
},
|
||
};
|
||
},
|
||
},
|
||
},
|
||
({ body }) => {
|
||
expect(body.model).toBe("replacement-model");
|
||
expect(requireRecord(body.options, "Ollama request options").stop).toEqual(["REPLACEMENT"]);
|
||
},
|
||
);
|
||
});
|
||
|
||
it("surfaces asynchronous payload hook rejection without dispatching a request", async () => {
|
||
await withSuccessfulOllamaFetch(async (fetchMock) => {
|
||
const stream = await createOllamaTestStream({
|
||
baseUrl: "http://ollama-host:11434",
|
||
options: {
|
||
onPayload: async () => {
|
||
await Promise.resolve();
|
||
throw new Error("payload admission rejected");
|
||
},
|
||
},
|
||
});
|
||
|
||
const events = await collectStreamEvents(stream);
|
||
|
||
expect(events).toMatchObject([
|
||
{
|
||
type: "error",
|
||
reason: "error",
|
||
error: {
|
||
stopReason: "error",
|
||
errorMessage: "payload admission rejected",
|
||
},
|
||
},
|
||
]);
|
||
expect(fetchMock).not.toHaveBeenCalled();
|
||
});
|
||
});
|
||
|
||
it("maps responseFormat JSON Schema to native Ollama format", async () => {
|
||
const schema = {
|
||
type: "object",
|
||
properties: { reply: { type: "string" } },
|
||
required: ["reply"],
|
||
additionalProperties: false,
|
||
};
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", options: { responseFormat: schema } },
|
||
({ body }) => expect(body).toMatchObject({ format: schema }),
|
||
);
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "omits native Ollama format when responseFormat is absent",
|
||
baseUrl: "http://ollama-host:11434",
|
||
},
|
||
{
|
||
name: "keeps provider-shaped text response formats off the native Ollama wire",
|
||
baseUrl: "http://ollama-host:11434",
|
||
responseFormat: { type: "text" },
|
||
},
|
||
{
|
||
name: "omits native Ollama format for cloud model through a local daemon",
|
||
baseUrl: "http://ollama-host:11434",
|
||
id: "gemma4:cloud",
|
||
responseFormat: {
|
||
type: "object",
|
||
properties: { reply: { type: "string" } },
|
||
required: ["reply"],
|
||
additionalProperties: false,
|
||
},
|
||
},
|
||
{
|
||
name: "omits native Ollama format for hosted Ollama Cloud",
|
||
baseUrl: "https://ollama.com/v1",
|
||
id: "gemma4",
|
||
responseFormat: {
|
||
type: "object",
|
||
properties: { reply: { type: "string" } },
|
||
required: ["reply"],
|
||
additionalProperties: false,
|
||
},
|
||
},
|
||
])("$name", async ({ baseUrl, id, responseFormat }) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl,
|
||
...(id ? { model: { id } } : {}),
|
||
...(responseFormat ? { options: { responseFormat } } : {}),
|
||
},
|
||
({ body }) => expect(body).not.toHaveProperty("format"),
|
||
);
|
||
});
|
||
|
||
it("lets native Ollama tools win over responseFormat", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
context: {
|
||
messages: [{ role: "user", content: "weather" }],
|
||
tools: [
|
||
{
|
||
name: "weather",
|
||
description: "Get weather",
|
||
parameters: { type: "object", properties: {} },
|
||
},
|
||
],
|
||
},
|
||
options: {
|
||
responseFormat: {
|
||
type: "object",
|
||
properties: { reply: { type: "string" } },
|
||
required: ["reply"],
|
||
additionalProperties: false,
|
||
},
|
||
},
|
||
},
|
||
({ body }) => {
|
||
expect(body.tools).toHaveLength(1);
|
||
expect(body).not.toHaveProperty("format");
|
||
},
|
||
);
|
||
});
|
||
|
||
it("uses configured params.num_ctx for native Ollama chat options", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
model: {
|
||
params: {
|
||
num_ctx: 32768,
|
||
temperature: 0.2,
|
||
top_p: 0.9,
|
||
thinking: false,
|
||
streaming: false,
|
||
},
|
||
contextWindow: 131072,
|
||
contextTokens: 16384,
|
||
},
|
||
options: { temperature: 0.7, maxTokens: 55 },
|
||
},
|
||
({ body }) => {
|
||
const options = requireRecord(body.options, "Ollama request options");
|
||
expect(options).toMatchObject({
|
||
num_ctx: 32768,
|
||
num_predict: 55,
|
||
temperature: 0.7,
|
||
top_p: 0.9,
|
||
});
|
||
expect(options.streaming).toBeUndefined();
|
||
expect(body.think).toBe(false);
|
||
},
|
||
);
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "uses effective contextTokens for native Ollama chat options",
|
||
model: { contextWindow: 262_144, contextTokens: 32_768 },
|
||
expected: 32_768,
|
||
},
|
||
{
|
||
name: "omits num_ctx when the model has no params.num_ctx and no catalog window",
|
||
model: { contextWindow: undefined },
|
||
expected: undefined,
|
||
},
|
||
{
|
||
name: "does not fall back to catalog contextWindow as native Ollama num_ctx",
|
||
model: { contextWindow: 32_768 },
|
||
expected: undefined,
|
||
},
|
||
{
|
||
name: "does not fall back to catalog maxTokens as native Ollama num_ctx",
|
||
model: { contextWindow: undefined, maxTokens: 65_536 },
|
||
expected: undefined,
|
||
},
|
||
])("$name", async ({ model, expected }) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", model },
|
||
({ body }) => expect(requireOptionalRecord(body.options)?.num_ctx).toBe(expected),
|
||
);
|
||
});
|
||
|
||
it.each([
|
||
{
|
||
name: "sets top_p=1 for native Ollama greedy sampling requests",
|
||
params: { num_ctx: 4096, top_p: 0.9, thinking: false },
|
||
temperature: 0,
|
||
expectedTopP: 1,
|
||
},
|
||
{
|
||
name: "sets top_p=1 for native Ollama greedy requests without configured top_p",
|
||
params: { num_ctx: 4096, thinking: false },
|
||
temperature: 0,
|
||
expectedTopP: 1,
|
||
},
|
||
{
|
||
name: "preserves configured top_p for native Ollama non-greedy sampling requests",
|
||
params: { top_p: 0.9 },
|
||
temperature: 0.2,
|
||
expectedTopP: 0.9,
|
||
},
|
||
])("$name", async ({ params, temperature, expectedTopP }) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", model: { params }, options: { temperature } },
|
||
({ body }) => {
|
||
const options = requireRecord(body.options, "Ollama sampling options");
|
||
expect(options.temperature).toBe(temperature);
|
||
expect(options.top_p).toBe(expectedTopP);
|
||
},
|
||
);
|
||
});
|
||
|
||
it.each(["low", "medium", "high"] as const)(
|
||
"preserves configured native Ollama params.thinking=%s",
|
||
async (thinking) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", model: { params: { thinking } } },
|
||
({ body }) => {
|
||
expect(body.think).toBe(thinking);
|
||
expect(requireOptionalRecord(body.options)?.think).toBeUndefined();
|
||
},
|
||
);
|
||
},
|
||
);
|
||
|
||
it("keeps configured local Ollama params.thinking=max compatible", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{ baseUrl: "http://ollama-host:11434", model: { params: { thinking: "max" } } },
|
||
({ body }) => {
|
||
expect(body.think).toBe("high");
|
||
},
|
||
);
|
||
});
|
||
|
||
it("preserves configured Ollama Cloud params.thinking=max", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "https://ollama.com",
|
||
model: { provider: "ollama-cloud", id: "glm-5.2", params: { thinking: "max" } },
|
||
},
|
||
({ body }) => {
|
||
expect(body.think).toBe("max");
|
||
},
|
||
);
|
||
});
|
||
|
||
it.each(["gpt-oss:120b", "kimi-k2.5", "custom-thinking-model"])(
|
||
"keeps configured Ollama Cloud %s params.thinking=max compatible",
|
||
async (id) => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "https://ollama.com",
|
||
model: { provider: "ollama-cloud", id, params: { thinking: "max" } },
|
||
},
|
||
({ body }) => {
|
||
expect(body.think).toBe("high");
|
||
},
|
||
);
|
||
},
|
||
);
|
||
|
||
it("uses the default loopback policy when baseUrl is empty", async () => {
|
||
await expectSuccessfulOllamaRequest({ baseUrl: "" }, ({ request }) => {
|
||
expect(request.url).toBe("http://127.0.0.1:11434/api/chat");
|
||
const policy = requireRecord(request.policy, "ssrf policy");
|
||
expect(policy.hostnameAllowlist).toEqual(["127.0.0.1"]);
|
||
expect(policy.allowPrivateNetwork).toBe(true);
|
||
});
|
||
});
|
||
|
||
it("merges default headers and allows request headers to override them", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
defaultHeaders: { "X-OLLAMA-KEY": "provider-secret", "X-Trace": "default" },
|
||
options: { headers: { "X-Trace": "request", "X-Request-Only": "1" } },
|
||
},
|
||
({ request }) =>
|
||
expect(requireHeaders(request.init?.headers)).toMatchObject({
|
||
"Content-Type": "application/json",
|
||
"X-OLLAMA-KEY": "provider-secret",
|
||
"X-Trace": "request",
|
||
"X-Request-Only": "1",
|
||
}),
|
||
);
|
||
});
|
||
|
||
it("preserves an explicit Authorization header when apiKey is a local marker", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
defaultHeaders: { Authorization: "Bearer proxy-token" },
|
||
options: {
|
||
apiKey: "ollama-local", // pragma: allowlist secret
|
||
headers: { Authorization: "Bearer proxy-token" },
|
||
},
|
||
},
|
||
({ request }) =>
|
||
expect(requireHeaders(request.init?.headers).Authorization).toBe("Bearer proxy-token"),
|
||
);
|
||
});
|
||
|
||
it("allows a real apiKey to override an explicit Authorization header", async () => {
|
||
await expectSuccessfulOllamaRequest(
|
||
{
|
||
baseUrl: "http://ollama-host:11434",
|
||
defaultHeaders: { Authorization: "Bearer proxy-token" },
|
||
options: { apiKey: "real-token" }, // pragma: allowlist secret
|
||
},
|
||
({ request }) =>
|
||
expect(requireHeaders(request.init?.headers).Authorization).toBe("Bearer real-token"),
|
||
);
|
||
});
|
||
|
||
it("surfaces bounded non-2xx HTTP response text as a status-prefixed error", async () => {
|
||
const tracked = cancelTrackedResponse(`${"Service Unavailable ".repeat(1024)}tail`, {
|
||
status: 503,
|
||
statusText: "Service Unavailable",
|
||
});
|
||
const textSpy = vi.spyOn(tracked.response, "text").mockRejectedValue(new Error("unbounded"));
|
||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||
response: tracked.response,
|
||
release: vi.fn(async () => undefined),
|
||
});
|
||
try {
|
||
const stream = await createOllamaTestStream({ baseUrl: "http://ollama-host:11434" });
|
||
const events = await collectStreamEvents(stream);
|
||
|
||
const errorEvent = events.find((e) => e.type === "error") as
|
||
| { type: "error"; error: { errorMessage?: string } }
|
||
| undefined;
|
||
if (!errorEvent) {
|
||
throw new Error("expected Ollama stream error event");
|
||
}
|
||
// The error message must start with the HTTP status code so that
|
||
// extractLeadingHttpStatus can parse it for failover/retry logic.
|
||
expect(errorEvent.error.errorMessage).toMatch(/^503\b/);
|
||
expect(errorEvent.error.errorMessage).toContain("Service Unavailable");
|
||
expect(errorEvent.error.errorMessage).not.toContain("tail");
|
||
expect(tracked.wasCanceled()).toBe(true);
|
||
expect(textSpy).not.toHaveBeenCalled();
|
||
} finally {
|
||
fetchWithSsrFGuardMock.mockReset();
|
||
}
|
||
});
|
||
|
||
it("keeps thinking chunks when no final content is emitted", async () => {
|
||
await expectDoneEventContent(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","thinking":"reasoned"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","thinking":" output"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":2}',
|
||
],
|
||
[{ type: "thinking", thinking: "reasoned output" }],
|
||
);
|
||
});
|
||
|
||
it("keeps streamed content after earlier thinking chunks", async () => {
|
||
await expectDoneEventContent(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","thinking":"internal"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"final"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":" answer"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":2}',
|
||
],
|
||
[
|
||
{ type: "thinking", thinking: "internal" },
|
||
{ type: "text", text: "final answer" },
|
||
],
|
||
);
|
||
});
|
||
|
||
it("keeps reasoning chunks when no final content is emitted", async () => {
|
||
await expectDoneEventContent(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","reasoning":"reasoned"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","reasoning":" output"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1}',
|
||
],
|
||
[{ type: "thinking", thinking: "reasoned output" }],
|
||
);
|
||
});
|
||
|
||
it("drops streamed reasoning chunks for non-reasoning models", async () => {
|
||
const events = await collectMockedOllamaEvents(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","reasoning":"reasoned"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","reasoning":" output"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":2}',
|
||
],
|
||
{ baseUrl: "http://ollama-host:11434", model: { reasoning: false } },
|
||
);
|
||
const doneEvent = events.at(-1);
|
||
if (!doneEvent || doneEvent.type !== "done") {
|
||
throw new Error("Expected done event");
|
||
}
|
||
|
||
expect(doneEvent.message.content).toEqual([]);
|
||
expect(doneEvent.message.usage.output).toBeGreaterThan(0);
|
||
expect(events.some((event) => event.type === "thinking_delta")).toBe(false);
|
||
});
|
||
|
||
it("keeps streamed content after earlier reasoning chunks", async () => {
|
||
await expectDoneEventContent(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"","reasoning":"internal"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"final"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":" answer"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":2}',
|
||
],
|
||
[
|
||
{ type: "thinking", thinking: "internal" },
|
||
{ type: "text", text: "final answer" },
|
||
],
|
||
);
|
||
});
|
||
});
|
||
|
||
describe("resolveOllamaBaseUrlForRun", () => {
|
||
it("prefers provider baseUrl over model baseUrl", () => {
|
||
expect(
|
||
resolveOllamaBaseUrlForRun({
|
||
modelBaseUrl: "http://model-host:11434",
|
||
providerBaseUrl: "http://provider-host:11434",
|
||
}),
|
||
).toBe("http://provider-host:11434");
|
||
});
|
||
|
||
it("falls back to model baseUrl when provider baseUrl is missing", () => {
|
||
expect(
|
||
resolveOllamaBaseUrlForRun({
|
||
modelBaseUrl: "http://model-host:11434",
|
||
}),
|
||
).toBe("http://model-host:11434");
|
||
});
|
||
|
||
it("falls back to native default when neither baseUrl is configured", () => {
|
||
expect(resolveOllamaBaseUrlForRun({})).toBe("http://127.0.0.1:11434");
|
||
});
|
||
});
|
||
|
||
describe("createConfiguredOllamaStreamFn", () => {
|
||
it("uses provider-level baseUrl when model baseUrl is absent", async () => {
|
||
await withMockNdjsonFetch(
|
||
[
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":"ok"},"done":false}',
|
||
'{"model":"m","created_at":"t","message":{"role":"assistant","content":""},"done":true,"prompt_eval_count":1,"eval_count":1}',
|
||
],
|
||
async (fetchMock) => {
|
||
const streamFn = createConfiguredOllamaStreamFn({
|
||
model: {
|
||
headers: { Authorization: "Bearer proxy-token" },
|
||
},
|
||
providerBaseUrl: "http://provider-host:11434/v1",
|
||
});
|
||
const stream = await Promise.resolve(
|
||
streamFn(
|
||
{
|
||
id: "qwen3:32b",
|
||
api: "ollama",
|
||
provider: "custom-ollama",
|
||
contextWindow: 131072,
|
||
} as never,
|
||
{
|
||
messages: [{ role: "user", content: "hello" }],
|
||
} as never,
|
||
{
|
||
apiKey: "ollama-local", // pragma: allowlist secret
|
||
} as never,
|
||
),
|
||
);
|
||
|
||
await collectStreamEvents(stream);
|
||
const request = getGuardedFetchCall(fetchMock);
|
||
expect(request.url).toBe("http://provider-host:11434/api/chat");
|
||
const requestInit = request.init ?? {};
|
||
expect(requireHeaders(requestInit.headers).Authorization).toBe("Bearer proxy-token");
|
||
},
|
||
);
|
||
});
|
||
});
|
||
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */
|