fix(agent): handled ollama-cloud task backoff
Added ollama-cloud subagent concurrency limiting, role fallback-chain inheritance, and visible empty length errors for native Ollama responses. Fixes #3464
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed Ollama/Ollama Cloud native chat responses that finish with `done_reason: "length"` and no assistant content surfacing as a normal empty stop; they now become a context-window error instead of entering empty-stop retry recovery. ([#3464](https://github.com/can1357/oh-my-pi/issues/3464))
|
||||
|
||||
## [16.1.19] - 2026-06-25
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -437,6 +437,17 @@ function mapDoneReason(doneReason: string | undefined, output: AssistantMessage)
|
||||
return "stop";
|
||||
}
|
||||
|
||||
const EMPTY_OLLAMA_LENGTH_COMPLETION_MESSAGE =
|
||||
"Model returned no content: prompt filled the context window; raise Ollama num_ctx or shorten the prompt.";
|
||||
|
||||
function hasVisibleAssistantContent(output: AssistantMessage): boolean {
|
||||
return output.content.some(block => {
|
||||
if (block.type === "text") return block.text.trim().length > 0;
|
||||
if (block.type === "thinking") return block.thinking.trim().length > 0;
|
||||
return block.type === "toolCall";
|
||||
});
|
||||
}
|
||||
|
||||
const OLLAMA_RETRY_DELAYS_MS = [2_000, 5_000, 10_000];
|
||||
|
||||
export const streamOllama: StreamFunction<"ollama-chat"> = (
|
||||
@@ -702,6 +713,10 @@ export const streamOllama: StreamFunction<"ollama-chat"> = (
|
||||
}
|
||||
endActiveThinkingBlock();
|
||||
endActiveTextBlock();
|
||||
if (output.stopReason === "length" && !hasVisibleAssistantContent(output)) {
|
||||
output.stopReason = "error";
|
||||
output.errorMessage = EMPTY_OLLAMA_LENGTH_COMPLETION_MESSAGE;
|
||||
}
|
||||
// Tool calls always mean "execute and continue" in the OpenAI/Ollama contract.
|
||||
// If the turn produced tool-call blocks but reported a natural `stop`, promote
|
||||
// to `toolUse` so the agent loop runs them (it gates execution on the stop
|
||||
@@ -713,6 +728,11 @@ export const streamOllama: StreamFunction<"ollama-chat"> = (
|
||||
if (firstTokenTime) {
|
||||
output.ttft = firstTokenTime - startTime;
|
||||
}
|
||||
if (output.stopReason === "error") {
|
||||
stream.push({ type: "error", reason: "error", error: output });
|
||||
stream.end();
|
||||
return;
|
||||
}
|
||||
const doneReason =
|
||||
output.stopReason === "length" ? "length" : output.stopReason === "toolUse" ? "toolUse" : "stop";
|
||||
stream.push({ type: "done", reason: doneReason, message: output });
|
||||
|
||||
@@ -267,6 +267,25 @@ describe("ollama-cloud provider support", () => {
|
||||
]);
|
||||
});
|
||||
|
||||
test("surfaces empty length completions as context-window errors", async () => {
|
||||
const fetchMock: FetchImpl = vi.fn(async () =>
|
||||
createNdjsonResponse([
|
||||
{ model: "gpt-oss:120b", done: true, done_reason: "length", prompt_eval_count: 1000, eval_count: 0 },
|
||||
]),
|
||||
);
|
||||
|
||||
const result = await stream(
|
||||
cloudModel,
|
||||
{
|
||||
messages: [{ role: "user", content: "Large task prompt", timestamp: Date.now() }],
|
||||
},
|
||||
{ apiKey: "cloud-test-key", fetch: fetchMock },
|
||||
).result();
|
||||
|
||||
expect(result.stopReason).toBe("error");
|
||||
expect(result.errorMessage).toContain("prompt filled the context window");
|
||||
});
|
||||
|
||||
test("sends max for GLM-5.2 xhigh reasoning on Ollama Cloud", async () => {
|
||||
let requestBody: Record<string, unknown> | undefined;
|
||||
const fetchMock: FetchImpl = vi.fn(async (_input, init) => {
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed ollama-cloud task/subagent fan-out exceeding the provider's three-request concurrency cap by adding a provider-specific subagent limiter, and let configured task/smol/advisor model roles inherit the default retry fallback chain when they do not define their own chain. ([#3464](https://github.com/can1357/oh-my-pi/issues/3464))
|
||||
- Fixed `omp install <plugin>` failing extension validation in compiled-binary mode with `Cannot find module '@(scope)/pi-ai/oauth' from '<plugin>/src/oauth.ts'` (and any other non-wildcard pi-* subpath import like `@oh-my-pi/pi-coding-agent/tools`). The bundled-registry override map seeded by `__buildLegacyPiPackageRootOverrides` only covered the bare package roots, so `rewriteLegacyPiImports` rewrote `@(scope)/pi-ai/oauth` to `@oh-my-pi/pi-ai/oauth`, fell through to `Bun.resolveSync` (which bunfs can't satisfy on Bun 1.3.14+), then left the original specifier alone — at which point Bun's native resolver failed because most plugins declare `@(scope)/pi-ai` as a `peerDependency` only and never materialize a real install. The new `scripts/generate-legacy-pi-bundled-registry.ts` reads every bundled pi-* package's non-wildcard `exports` field and emits both the heavy `legacy-pi-bundled-registry.ts` (static imports + map) and a light `legacy-pi-bundled-keys.ts` (statically imported by `legacy-pi-compat.ts` to seed the override map without the cascade through `legacy-pi-coding-agent-shim → ../index → export/html/...`). `scripts/build-binary.ts` now runs the generator before `bun build --compile`. ([#3442](https://github.com/can1357/oh-my-pi/issues/3442))
|
||||
- Fixed `skill://` tool resolution losing loaded session skills when a tool runs outside the session-initialization module state. Internal URL resolution now prefers the caller's `session.skills` snapshot before falling back to the process-global skill list, so `read skill://<name>` works across tool execution boundaries. ([#3436](https://github.com/can1357/oh-my-pi/issues/3436))
|
||||
- Fixed `@image` mentions on OpenAI Codex Responses (chatgpt.com `gpt-5.5` and siblings) failing with `Codex error event: [OneOfParam] [input[N].content[…]] [invalid_enum_value] Invalid value: 'input_image'. Supported values are: 'input_text'.`. `convertToLlm` for `fileMention` always emitted a `developer`-role message, so the auto-attached image landed in a Responses content array that the Codex backend (and OpenAI Responses generally) only allows to carry `input_text`. #3421's previous fix only stopped the Codex Responses Lite header from going out on image-bearing turns; the full transport kept rejecting the same payload. `convertToLlm` now splits a mixed-content `fileMention` into two messages — text-only files stay on `developer` (so the auto-read context keeps instruction priority), while image-bearing files ride on `user` (the only Responses content slot that accepts `input_image`). ([#3443](https://github.com/can1357/oh-my-pi/issues/3443))
|
||||
|
||||
@@ -4041,6 +4041,17 @@ export const SETTINGS_SCHEMA = {
|
||||
},
|
||||
|
||||
// Provider selection
|
||||
"providers.ollama-cloud.maxConcurrency": {
|
||||
type: "number",
|
||||
default: 3,
|
||||
ui: {
|
||||
tab: "providers",
|
||||
group: "Services",
|
||||
label: "Ollama Cloud Max Concurrency",
|
||||
description:
|
||||
"Maximum concurrent Ollama Cloud subagent runs per process; 0 disables the provider-specific limit",
|
||||
},
|
||||
},
|
||||
"providers.webSearch": {
|
||||
type: "enum",
|
||||
values: SEARCH_PROVIDER_PREFERENCES,
|
||||
|
||||
@@ -10619,7 +10619,16 @@ export class AgentSession {
|
||||
#getRetryFallbackChains(): RetryFallbackChains {
|
||||
const configuredChains = this.settings.get("retry.fallbackChains");
|
||||
if (!configuredChains || typeof configuredChains !== "object") return {};
|
||||
return configuredChains as RetryFallbackChains;
|
||||
const chains: RetryFallbackChains = { ...(configuredChains as RetryFallbackChains) };
|
||||
const defaultChain = chains.default;
|
||||
if (Array.isArray(defaultChain)) {
|
||||
for (const role of Object.keys(this.settings.getModelRoles())) {
|
||||
if (role !== "default" && chains[role] === undefined) {
|
||||
chains[role] = defaultChain;
|
||||
}
|
||||
}
|
||||
}
|
||||
return chains;
|
||||
}
|
||||
|
||||
#validateRetryFallbackChains(): void {
|
||||
|
||||
@@ -50,12 +50,12 @@ import {
|
||||
type OutputValidator,
|
||||
summarizeValidationFailure,
|
||||
} from "../tools/output-schema-validator";
|
||||
|
||||
import { type ReportFindingDetails, toReviewFinding } from "../tools/review";
|
||||
import { ToolAbortError } from "../tools/tool-errors";
|
||||
import type { EventBus } from "../utils/event-bus";
|
||||
import { buildNamedToolChoice } from "../utils/tool-choice";
|
||||
import type { WorkspaceTree } from "../workspace-tree";
|
||||
import { Semaphore } from "./parallel";
|
||||
import { subprocessToolRegistry } from "./subprocess-tool-registry";
|
||||
import {
|
||||
type AgentDefinition,
|
||||
@@ -194,6 +194,35 @@ function installSubagentRetryFallbackChain(args: {
|
||||
return role;
|
||||
}
|
||||
|
||||
const PROVIDER_MAX_CONCURRENCY_SETTINGS: Record<string, SettingPath> = {
|
||||
"ollama-cloud": "providers.ollama-cloud.maxConcurrency",
|
||||
};
|
||||
|
||||
interface ProviderSemaphoreEntry {
|
||||
limit: number;
|
||||
semaphore: Semaphore;
|
||||
}
|
||||
|
||||
const providerSemaphores = new Map<string, ProviderSemaphoreEntry>();
|
||||
|
||||
function getProviderConcurrencyLimit(settings: Settings, provider: string): number {
|
||||
const settingPath = PROVIDER_MAX_CONCURRENCY_SETTINGS[provider];
|
||||
if (!settingPath) return 0;
|
||||
const raw = settings.get(settingPath);
|
||||
const limit = Number.isFinite(raw) ? Math.trunc(raw) : 0;
|
||||
return limit > 0 ? limit : 0;
|
||||
}
|
||||
|
||||
function getProviderSemaphore(settings: Settings, provider: string): Semaphore | undefined {
|
||||
const limit = getProviderConcurrencyLimit(settings, provider);
|
||||
if (limit <= 0) return undefined;
|
||||
const existing = providerSemaphores.get(provider);
|
||||
if (existing?.limit === limit) return existing.semaphore;
|
||||
const semaphore = new Semaphore(limit);
|
||||
providerSemaphores.set(provider, { limit, semaphore });
|
||||
return semaphore;
|
||||
}
|
||||
|
||||
function renderIrcPeerRoster(selfId: string): string {
|
||||
const peers = AgentRegistry.global()
|
||||
.list()
|
||||
@@ -1889,6 +1918,8 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
let sessionOpenedAt: number | undefined;
|
||||
let sessionCreatedAt: number | undefined;
|
||||
let readyAt: number | undefined;
|
||||
let providerSemaphore: Semaphore | undefined;
|
||||
let providerSemaphoreAcquired = false;
|
||||
|
||||
try {
|
||||
checkAbort();
|
||||
@@ -1957,6 +1988,14 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
? resolvedThinkingLevel
|
||||
: (thinkingLevel ?? resolvedThinkingLevel);
|
||||
resolvedAt = performance.now();
|
||||
if (model) {
|
||||
providerSemaphore = getProviderSemaphore(settings, model.provider);
|
||||
if (providerSemaphore) {
|
||||
await awaitAbortable(providerSemaphore.acquire());
|
||||
providerSemaphoreAcquired = true;
|
||||
checkAbort();
|
||||
}
|
||||
}
|
||||
|
||||
const effectiveCwd = worktree ?? cwd;
|
||||
const sessionManager = sessionFile
|
||||
@@ -2247,6 +2286,10 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
}
|
||||
if (exitCode === 0) exitCode = 1;
|
||||
}
|
||||
if (providerSemaphoreAcquired) {
|
||||
providerSemaphore?.release();
|
||||
providerSemaphoreAcquired = false;
|
||||
}
|
||||
sessionAbortController.abort();
|
||||
if (unsubscribe) {
|
||||
try {
|
||||
|
||||
@@ -0,0 +1,220 @@
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test";
|
||||
import * as path from "node:path";
|
||||
import { Agent } from "@oh-my-pi/pi-agent-core";
|
||||
import type { Model } from "@oh-my-pi/pi-ai";
|
||||
import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock";
|
||||
import { type GeneratedProvider, getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types";
|
||||
import type { CreateAgentSessionResult } from "@oh-my-pi/pi-coding-agent/sdk";
|
||||
import * as sdkModule from "@oh-my-pi/pi-coding-agent/sdk";
|
||||
import {
|
||||
AgentSession,
|
||||
type AgentSessionEvent,
|
||||
type PromptOptions,
|
||||
} from "@oh-my-pi/pi-coding-agent/session/agent-session";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
import { runSubprocess } from "@oh-my-pi/pi-coding-agent/task/executor";
|
||||
import type { AgentDefinition } from "@oh-my-pi/pi-coding-agent/task/types";
|
||||
import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
type MockPromptSession = AgentSession & {
|
||||
emit(event: AgentSessionEvent): void;
|
||||
};
|
||||
|
||||
interface Deferred {
|
||||
promise: Promise<void>;
|
||||
resolve: () => void;
|
||||
}
|
||||
|
||||
function deferred(): Deferred {
|
||||
const { promise, resolve } = Promise.withResolvers<void>();
|
||||
return { promise, resolve };
|
||||
}
|
||||
|
||||
function createSessionResult(session: AgentSession): CreateAgentSessionResult {
|
||||
return {
|
||||
session,
|
||||
extensionsResult: { extensions: [], errors: [], runtime: {} as unknown } as LoadExtensionsResult,
|
||||
setToolUIContext: () => {},
|
||||
eventBus: new EventBus(),
|
||||
};
|
||||
}
|
||||
|
||||
function createGateSession(onPrompt: () => Promise<void>): MockPromptSession {
|
||||
const listeners: Array<(event: AgentSessionEvent) => void> = [];
|
||||
const session = {
|
||||
agent: { state: { systemPrompt: ["test"] } },
|
||||
state: { messages: [] },
|
||||
extensionRunner: undefined,
|
||||
sessionManager: { appendSessionInit: () => {} },
|
||||
getActiveToolNames: () => ["yield"],
|
||||
setActiveToolsByName: async () => {},
|
||||
subscribe: (listener: (event: AgentSessionEvent) => void) => {
|
||||
listeners.push(listener);
|
||||
return () => {};
|
||||
},
|
||||
prompt: async (_text: string, _options?: PromptOptions) => {
|
||||
await onPrompt();
|
||||
for (const listener of listeners) {
|
||||
listener({
|
||||
type: "tool_execution_end",
|
||||
toolCallId: "tool-yield",
|
||||
toolName: "yield",
|
||||
result: { content: [{ type: "text", text: "Result submitted." }], details: { status: "success" } },
|
||||
isError: false,
|
||||
});
|
||||
}
|
||||
},
|
||||
waitForIdle: async () => {},
|
||||
getLastAssistantMessage: () => undefined,
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
emit: (event: AgentSessionEvent) => {
|
||||
for (const listener of listeners) listener(event);
|
||||
},
|
||||
};
|
||||
return session as unknown as MockPromptSession;
|
||||
}
|
||||
|
||||
function requireModel(provider: GeneratedProvider, id: string): Model {
|
||||
const model = getBundledModel(provider, id);
|
||||
if (!model) throw new Error(`Expected bundled model ${provider}/${id}`);
|
||||
return model;
|
||||
}
|
||||
|
||||
const taskAgent: AgentDefinition = {
|
||||
name: "task",
|
||||
description: "General task agent",
|
||||
systemPrompt: "test",
|
||||
source: "bundled",
|
||||
};
|
||||
|
||||
describe("issue #3464: ollama-cloud task backoff", () => {
|
||||
let tempDir: TempDir;
|
||||
let authStorage: AuthStorage;
|
||||
let modelRegistry: ModelRegistry;
|
||||
let session: AgentSession | undefined;
|
||||
|
||||
beforeAll(async () => {
|
||||
tempDir = TempDir.createSync("@omp-issue-3464-");
|
||||
authStorage = await AuthStorage.create(path.join(tempDir.path(), "auth.db"));
|
||||
authStorage.setRuntimeApiKey("anthropic", "anthropic-test-key");
|
||||
authStorage.setRuntimeApiKey("openai", "openai-test-key");
|
||||
authStorage.setRuntimeApiKey("ollama-cloud", "ollama-cloud-test-key");
|
||||
modelRegistry = new ModelRegistry(authStorage);
|
||||
});
|
||||
|
||||
afterAll(() => {
|
||||
authStorage.close();
|
||||
tempDir.removeSync();
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
if (session) {
|
||||
await session.dispose();
|
||||
session = undefined;
|
||||
}
|
||||
modelRegistry.clearSuppressedSelectors();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("uses the default fallback chain for a configured task role with no task chain", async () => {
|
||||
const primary = requireModel("anthropic", "claude-sonnet-4-5");
|
||||
const fallback = requireModel("openai", "gpt-4o-mini");
|
||||
const requestedModels: string[] = [];
|
||||
const mock = createMockModel();
|
||||
let primaryAttempts = 0;
|
||||
const agent = new Agent({
|
||||
getApiKey: model => `${model.provider}-test-key`,
|
||||
initialState: { model: primary, systemPrompt: ["Test"], tools: [], messages: [] },
|
||||
streamFn: (model, context, options) => {
|
||||
requestedModels.push(`${model.provider}/${model.id}`);
|
||||
if (model.provider === primary.provider && model.id === primary.id && primaryAttempts === 0) {
|
||||
primaryAttempts += 1;
|
||||
mock.push({ throw: "rate limit exceeded retry-after-ms=200" });
|
||||
} else {
|
||||
mock.push({ content: [`ok:${model.provider}/${model.id}`] });
|
||||
}
|
||||
return mock.stream(model, context, options);
|
||||
},
|
||||
});
|
||||
const settings = Settings.isolated({
|
||||
"compaction.enabled": false,
|
||||
"retry.baseDelayMs": 5,
|
||||
"retry.maxRetries": 1,
|
||||
"retry.fallbackChains": { default: [`${fallback.provider}/${fallback.id}`] },
|
||||
});
|
||||
settings.setModelRole("task", `${primary.provider}/${primary.id}`);
|
||||
|
||||
session = new AgentSession({ agent, sessionManager: SessionManager.inMemory(), settings, modelRegistry });
|
||||
|
||||
await session.prompt("Task role should inherit the default fallback chain");
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(requestedModels).toEqual([`${primary.provider}/${primary.id}`, `${fallback.provider}/${fallback.id}`]);
|
||||
expect(session.model?.provider).toBe(fallback.provider);
|
||||
expect(session.model?.id).toBe(fallback.id);
|
||||
});
|
||||
|
||||
it("bounds concurrent subagent runs by the resolved ollama-cloud provider limit", async () => {
|
||||
const cloudModel = requireModel("ollama-cloud", "gpt-oss:120b");
|
||||
const started: string[] = [];
|
||||
const gates = new Map<string, Deferred>();
|
||||
const firstStarted = deferred();
|
||||
const secondStarted = deferred();
|
||||
vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async options => {
|
||||
const id = options?.agentId ?? "unknown";
|
||||
const gate = deferred();
|
||||
gates.set(id, gate);
|
||||
return createSessionResult(
|
||||
createGateSession(async () => {
|
||||
started.push(id);
|
||||
if (id === "CloudOne") firstStarted.resolve();
|
||||
if (id === "CloudTwo") secondStarted.resolve();
|
||||
await gate.promise;
|
||||
}),
|
||||
);
|
||||
});
|
||||
const settings = Settings.isolated({
|
||||
"providers.ollama-cloud.maxConcurrency": 1,
|
||||
});
|
||||
|
||||
const first = runSubprocess({
|
||||
cwd: "/tmp",
|
||||
agent: taskAgent,
|
||||
task: "first",
|
||||
index: 0,
|
||||
id: "CloudOne",
|
||||
modelOverride: `${cloudModel.provider}/${cloudModel.id}`,
|
||||
settings,
|
||||
modelRegistry,
|
||||
enableLsp: false,
|
||||
});
|
||||
const second = runSubprocess({
|
||||
cwd: "/tmp",
|
||||
agent: taskAgent,
|
||||
task: "second",
|
||||
index: 1,
|
||||
id: "CloudTwo",
|
||||
modelOverride: `${cloudModel.provider}/${cloudModel.id}`,
|
||||
settings,
|
||||
modelRegistry,
|
||||
enableLsp: false,
|
||||
});
|
||||
|
||||
await firstStarted.promise;
|
||||
expect(started).toEqual(["CloudOne"]);
|
||||
expect(gates.has("CloudTwo")).toBe(false);
|
||||
|
||||
gates.get("CloudOne")?.resolve();
|
||||
await first;
|
||||
await secondStarted.promise;
|
||||
expect(started).toEqual(["CloudOne", "CloudTwo"]);
|
||||
gates.get("CloudTwo")?.resolve();
|
||||
await second;
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user