Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
522 changes: 218 additions & 304 deletions open-sse/executors/index.ts

Large diffs are not rendered by default.

41 changes: 38 additions & 3 deletions open-sse/executors/registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,46 @@ export function getRegisteredExecutor(alias: string): BaseExecutor | undefined {
return registry.get(alias);
}

// ── #11220: lazy registration ───────────────────────────────────────────────
// Aliases may register a deferred loader instead of an instance. The alias and
// its registration ORDER are declared eagerly — hasRegisteredExecutor() and
// listExecutorAliases() stay synchronous and the golden snapshot keeps its
// shape — while the class import + construction happen on first use. A
// completed load caches into `registry`, so later resolution is identical to a
// static registration.
const lazyLoaders = new Map<string, () => Promise<BaseExecutor>>();
const lazyInFlight = new Map<string, Promise<BaseExecutor>>();

export function registerLazyExecutor(alias: string, load: () => Promise<BaseExecutor>): void {
if (registry.has(alias) || lazyLoaders.has(alias)) {
throw new Error(`executor alias already registered: "${alias}"`);
}
lazyLoaders.set(alias, load);
}

export function loadRegisteredExecutor(alias: string): Promise<BaseExecutor> | undefined {
const cached = registry.get(alias);
if (cached) return Promise.resolve(cached);
const load = lazyLoaders.get(alias);
if (!load) return undefined;
let inFlight = lazyInFlight.get(alias);
if (!inFlight) {
inFlight = load().then((executor) => {
registerExecutor(alias, executor);
lazyLoaders.delete(alias);
lazyInFlight.delete(alias);
return executor;
});
lazyInFlight.set(alias, inFlight);
}
return inFlight;
}

export function hasRegisteredExecutor(alias: string): boolean {
return registry.has(alias);
return registry.has(alias) || lazyLoaders.has(alias);
}

/** All registered aliases, in registration order. */
/** All registered aliases — static and lazy — in registration order. */
export function listExecutorAliases(): string[] {
return [...registry.keys()];
return [...registry.keys(), ...lazyLoaders.keys()];
}
4 changes: 2 additions & 2 deletions open-sse/handlers/chatCore/cliproxyModelMapping.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,12 @@
type ExecutorInput = {
model: string;
body: unknown;
[key: string]: unknown;
};

// No index signature: executors (BaseExecutor subclasses) must satisfy this
// structurally, and class instances don't carry index signatures.
type ExecutorLike = {
execute: (input: ExecutorInput) => Promise<unknown>;
[key: string]: unknown;
};

export type CliproxyapiModelMapping = Record<string, unknown> | null | undefined;
Expand Down
4 changes: 2 additions & 2 deletions open-sse/handlers/chatCore/cliproxyapiCredentials.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,12 @@ import type { ProviderCredentials } from "../../executors/base.ts";

type ExecutorInput = {
credentials: ProviderCredentials;
[key: string]: unknown;
};

// No index signature: executors (BaseExecutor subclasses) must satisfy this
// structurally, and class instances don't carry index signatures.
type ExecutorLike = {
execute: (input: ExecutorInput) => Promise<unknown>;
[key: string]: unknown;
};

/**
Expand Down
13 changes: 8 additions & 5 deletions open-sse/handlers/chatCore/executorProxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,12 +76,15 @@ async function loadCliproxyapiSettings(): Promise<{
* dedicated-credential wrappers applied. Used by the direct `cliproxyapi` leg
* and the CLIProxyAPI branch of `fallback`.
*/
function resolveCliproxyapiExecutor(
async function resolveCliproxyapiExecutor(
cliproxyapiModelMapping: Record<string, unknown> | null,
dedicatedApiKey: string | null
) {
return wrapExecutorWithCliproxyapiCredentials(
wrapExecutorWithCliproxyapiModelMapping(getExecutor("cliproxyapi"), cliproxyapiModelMapping),
wrapExecutorWithCliproxyapiModelMapping(
await getExecutor("cliproxyapi"),
cliproxyapiModelMapping
),
dedicatedApiKey
);
}
Expand Down Expand Up @@ -138,16 +141,16 @@ export async function resolveExecutorWithProxy(
// backend on specific failures. The backend defaults to CLIProxyAPI so every
// pre-existing fallback config behaves exactly as before; fallbackBackend
// === "dario" opts the retry leg over to Dario instead.
const nativeExec = getExecutor(prov);
const nativeExec = await getExecutor(prov);
const fallbackBackend: FallbackBackend = cfg.fallbackBackend;
const { fallbackCodes, dedicatedApiKey } = await loadCliproxyapiSettings();

// The model mapping applies only to the CLIProxyAPI retry leg (proxyExec) —
// the native leg must keep seeing the original, unmapped model.
const proxyExec =
fallbackBackend === "dario"
? getExecutor("dario")
: resolveCliproxyapiExecutor(cfg.cliproxyapiModelMapping, dedicatedApiKey);
? await getExecutor("dario")
: await resolveCliproxyapiExecutor(cfg.cliproxyapiModelMapping, dedicatedApiKey);
const backendLabel = fallbackBackend === "dario" ? "Dario" : "CLIProxyAPI";
const isRetryableStatus = (s: number) => fallbackCodes.includes(s) || s === 0;

Expand Down
2 changes: 1 addition & 1 deletion open-sse/handlers/videoGeneration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -374,7 +374,7 @@ async function handleVertexVeoGeneration({ model, body, credentials, log }) {
* Submits an AnimateDiff or SVD workflow, polls for completion, fetches output video
*/
async function handleVeoAiFreeVideoGeneration({ model, provider, body, credentials, log }) {
const executor = getExecutor(provider);
const executor = await getExecutor(provider);
if (!executor) {
return { success: false, status: 400, error: `Unknown video provider: ${provider}` };
}
Expand Down
3 changes: 2 additions & 1 deletion open-sse/services/compression/eval/executorModelClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,10 @@ export function createExecutorModelClient(
credentials: ProviderCredentials,
costPerKTokenOut?: number
): ModelClient {
const executor = getExecutor(provider);
return {
async complete(model: string, messages: ChatTurn[]): Promise<ModelCallResult> {
// #11220: getExecutor is async (lazy registry) — resolve per call.
const executor = await getExecutor(provider);
const body = { model, messages, stream: false };
const input: ExecuteInput = {
model,
Expand Down
4 changes: 2 additions & 2 deletions scripts/check/check-known-symbols.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,8 +142,8 @@ export function diffComboStrategies(
* getExecutor() na função main().
*/
export function extractExecutorAliases(indexSource: string): string[] {
const start = indexSource.indexOf("const executors = {");
if (start < 0) throw new Error("could not find `const executors = {` in executors/index.ts");
const start = indexSource.indexOf("const lazyExecutors");
if (start < 0) throw new Error("could not find 'const lazyExecutors' in executors/index.ts");
const end = indexSource.indexOf("\n};", start);
if (end < 0) throw new Error("could not find end of executors map (`\\n};`)");
const block = indexSource.slice(start, end);
Expand Down
3 changes: 2 additions & 1 deletion src/lib/compression/judgeModelClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,10 @@ export function createPricedJudgeClient(
provider: string,
credentials: ProviderCredentials
): ModelClient {
const executor = getExecutor(provider);
return {
async complete(model: string, messages: ChatTurn[]): Promise<ModelCallResult> {
// #11220: getExecutor is async (lazy registry) — resolve per call.
const executor = await getExecutor(provider);
const input: ExecuteInput = {
model,
body: { model, messages, stream: false },
Expand Down
8 changes: 5 additions & 3 deletions src/lib/providers/validation/anthropicFormat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -132,12 +132,13 @@ export async function validateClaudeOAuthInline({
modelId: string | null | undefined;
providerSpecificData?: Record<string, unknown>;
}) {
const testModelId =
providerSpecificData?.validationModelId || modelId || "claude-haiku-4-5-20251001";
const override = providerSpecificData?.validationModelId;
const testModelId: string =
typeof override === "string" && override ? override : modelId || "claude-haiku-4-5-20251001";

try {
const { getExecutor } = await import("@omniroute/open-sse/executors/index.ts");
const { response } = await getExecutor("claude").execute({
const executed = await (await getExecutor("claude")).execute({
model: testModelId,
body: {
model: testModelId,
Expand All @@ -148,6 +149,7 @@ export async function validateClaudeOAuthInline({
credentials: { accessToken: apiKey, providerSpecificData },
});

const response = executed instanceof Response ? executed : executed.response;
if (response.status === 401 || response.status === 403) {
return { valid: false, error: "Invalid OAuth token" };
}
Expand Down
5 changes: 3 additions & 2 deletions src/lib/services/quotaAutoPing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import { logger } from "@omniroute/open-sse/utils/logger.ts";
import { sanitizeErrorMessage } from "@omniroute/open-sse/utils/error.ts";
import { getExecutor } from "@omniroute/open-sse/executors/index.ts";
import type { BaseExecutor } from "@omniroute/open-sse/executors/base";
import { getCodexUsage } from "@omniroute/open-sse/services/usage/codex.ts";
import { getSettings, getProviderConnections, updateProviderConnection } from "@/lib/localDb";
import { isConnectionUnavailableToAuxiliaryActivity } from "@/lib/exclusiveLeaseIsolation";
Expand Down Expand Up @@ -67,7 +68,7 @@ export interface QuotaAutoPingDeps {
accessToken?: string,
providerSpecificData?: JsonRecord
) => Promise<JsonRecord>;
getExecutor: (provider: string) => { execute: (input: JsonRecord) => Promise<JsonRecord> };
getExecutor: (provider: string) => Promise<BaseExecutor>;
canExecuteProvider: (provider: string) => boolean;
isConnectionUnavailableToAuxiliaryActivity: (connectionId: string) => Promise<boolean>;
}
Expand Down Expand Up @@ -208,7 +209,7 @@ async function sendCodexPing(
providerConfig: QuotaAutoPingProviderConfig,
deps: QuotaAutoPingDeps
): Promise<boolean> {
const executor = deps.getExecutor("codex");
const executor = await deps.getExecutor("codex");
const result = await executor.execute({
model: providerConfig.pingModel,
stream: true,
Expand Down
9 changes: 5 additions & 4 deletions tests/unit/adobe-firefly.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ test("adobe-firefly is registered in VIDEO_PROVIDERS with adobe-firefly-video fo
});

test("getExecutor(adobe-firefly) rejects chat completions", async () => {
const executor = getExecutor("adobe-firefly");
const executor = await getExecutor("adobe-firefly");
assert.ok(executor);
const result = await executor.execute({
model: "adobe-firefly/nano-banana-pro",
Expand All @@ -95,9 +95,10 @@ test("getExecutor(adobe-firefly) rejects chat completions", async () => {
stream: false,
credentials: { apiKey: "tok" },
});
assert.ok(result.response, "executor must return a Response wrapper");
assert.equal(result.response.status, 400);
const bodyText = await result.response.text();
const response = result instanceof Response ? result : result.response;
assert.ok(response, "executor must return a Response wrapper");
assert.equal(response.status, 400);
const bodyText = await response.text();
assert.match(bodyText, /images\/generations|videos\/generations|media-generation/i);
});

Expand Down
7 changes: 4 additions & 3 deletions tests/unit/azure-param-rules.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@ import {
applyAzureParamRules,
AZURE_COMPLETION_TOKEN_DEPLOYMENT,
} from "../../open-sse/executors/azureParamRules.ts";
import { getExecutor, AzureAiExecutor } from "../../open-sse/executors/index.ts";
import { getExecutor } from "../../open-sse/executors/index.ts";
import { AzureAiExecutor } from "../../open-sse/executors/azure-ai.ts";

/**
* Regression guards for two Azure 400s observed against a live Azure AI Foundry
Expand Down Expand Up @@ -87,8 +88,8 @@ test("the regex does not match unrelated names by accident", () => {
assert.equal(AZURE_COMPLETION_TOKEN_DEPLOYMENT.test("Kimi-K2.7-Code"), false);
});

test("azure-ai resolves to AzureAiExecutor, not the bare DefaultExecutor", () => {
const executor = getExecutor("azure-ai");
test("azure-ai resolves to AzureAiExecutor, not the bare DefaultExecutor", async () => {
const executor = await getExecutor("azure-ai");
assert.ok(
executor instanceof AzureAiExecutor,
"azure-ai must have its own executor so it inherits the Azure param rules"
Expand Down
6 changes: 3 additions & 3 deletions tests/unit/blackbox-web.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,11 +59,11 @@ function mockFetchCapture(status = 200, text = "Hello from Blackbox") {
};
}

test("BlackboxWebExecutor is registered in executor index", () => {
test("BlackboxWebExecutor is registered in executor index", async () => {
assert.ok(hasSpecializedExecutor("blackbox-web"));
assert.ok(hasSpecializedExecutor("bb-web"));
const executor = getExecutor("blackbox-web");
const alias = getExecutor("bb-web");
const executor = await getExecutor("blackbox-web");
const alias = await getExecutor("bb-web");
assert.ok(executor instanceof BlackboxWebExecutor);
assert.ok(alias instanceof BlackboxWebExecutor);
});
Expand Down
20 changes: 10 additions & 10 deletions tests/unit/chatcore-executor-proxy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ after(() => {
test("no config (disabled by default) returns the provider's own executor", async () => {
clearUpstreamProxyConfigCache("openai");
const exec = await resolveExecutorWithProxy("openai");
assert.equal(exec, getExecutor("openai"));
assert.equal(exec, await getExecutor("openai"));
});

test("mode 'native' returns the provider's own executor", async () => {
Expand All @@ -52,7 +52,7 @@ test("mode 'native' returns the provider's own executor", async () => {
});
clearUpstreamProxyConfigCache("openai");
const exec = await resolveExecutorWithProxy("openai");
assert.equal(exec, getExecutor("openai"));
assert.equal(exec, await getExecutor("openai"));
});

test("mode 'cliproxyapi' returns the CLIProxyAPI passthrough executor", async () => {
Expand All @@ -63,7 +63,7 @@ test("mode 'cliproxyapi' returns the CLIProxyAPI passthrough executor", async ()
});
clearUpstreamProxyConfigCache("anthropic");
const exec = await resolveExecutorWithProxy("anthropic");
assert.equal(exec, getExecutor("cliproxyapi"));
assert.equal(exec, await getExecutor("cliproxyapi"));
});

test("mode 'fallback' returns a distinct wrapper owning its own execute()", async () => {
Expand All @@ -74,8 +74,8 @@ test("mode 'fallback' returns a distinct wrapper owning its own execute()", asyn
});
clearUpstreamProxyConfigCache("openai");
const exec = await resolveExecutorWithProxy("openai");
assert.notEqual(exec, getExecutor("openai"));
assert.notEqual(exec, getExecutor("cliproxyapi"));
assert.notEqual(exec, await getExecutor("openai"));
assert.notEqual(exec, await getExecutor("cliproxyapi"));
assert.equal(typeof exec.execute, "function");
});

Expand All @@ -94,15 +94,15 @@ test("connection override 'claude-native' selects CLIProxyAPI even when provider
const exec = await resolveExecutorWithProxy("openai", undefined, {
cliproxyapiMode: "claude-native",
});
assert.equal(exec, getExecutor("cliproxyapi"));
assert.equal(exec, await getExecutor("cliproxyapi"));
});

test("connection override 'claude-native' selects CLIProxyAPI even with no provider config (default)", async () => {
clearUpstreamProxyConfigCache("anthropic");
const exec = await resolveExecutorWithProxy("anthropic", undefined, {
cliproxyapiMode: "claude-native",
});
assert.equal(exec, getExecutor("cliproxyapi"));
assert.equal(exec, await getExecutor("cliproxyapi"));
});

test("no connection override + provider mode native → native executor (unchanged)", async () => {
Expand All @@ -115,13 +115,13 @@ test("no connection override + provider mode native → native executor (unchang
const exec = await resolveExecutorWithProxy("openai", undefined, {
someOtherField: "x",
});
assert.equal(exec, getExecutor("openai"));
assert.equal(exec, await getExecutor("openai"));
});

test("connection override absent (undefined providerSpecificData) preserves default behaviour", async () => {
clearUpstreamProxyConfigCache("openai");
const exec = await resolveExecutorWithProxy("openai");
assert.equal(exec, getExecutor("openai"));
assert.equal(exec, await getExecutor("openai"));
});

test("connection override wins over provider mode 'fallback'", async () => {
Expand All @@ -135,5 +135,5 @@ test("connection override wins over provider mode 'fallback'", async () => {
cliproxyapiMode: "claude-native",
});
// Connection override short-circuits to the passthrough executor, not the fallback wrapper.
assert.equal(exec, getExecutor("cliproxyapi"));
assert.equal(exec, await getExecutor("cliproxyapi"));
});
2 changes: 1 addition & 1 deletion tests/unit/chatcore-translation-paths.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -456,7 +456,7 @@ test("chatCore times out upstream execution before provider response headers", a
// (fresh-DB default leaves it off → the waitFor below would never resolve;
// failed deterministically on CI and on an isolated run, incl. at v3.8.18).
await settingsDb.updateSettings({ call_log_pipeline_enabled: true });
const executor = getExecutor("openai");
const executor = await getExecutor("openai");
const originalGetTimeoutMs = executor.getTimeoutMs?.bind(executor);
executor.getTimeoutMs = () => 200;

Expand Down
10 changes: 5 additions & 5 deletions tests/unit/chatgpt-web.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -370,16 +370,16 @@ function reset() {

// ─── Registration ───────────────────────────────────────────────────────────

test("ChatGptWebExecutor is registered in executor index", () => {
test("ChatGptWebExecutor is registered in executor index", async () => {
assert.ok(hasSpecializedExecutor("chatgpt-web"));
assert.ok(hasSpecializedExecutor("cgpt-web"));
const executor = getExecutor("chatgpt-web");
const executor = await getExecutor("chatgpt-web");
assert.ok(executor instanceof ChatGptWebExecutor);
});

test("ChatGptWebExecutor alias resolves to same type", () => {
const a = getExecutor("chatgpt-web");
const b = getExecutor("cgpt-web");
test("ChatGptWebExecutor alias resolves to same type", async () => {
const a = await getExecutor("chatgpt-web");
const b = await getExecutor("cgpt-web");
assert.ok(a instanceof ChatGptWebExecutor);
assert.ok(b instanceof ChatGptWebExecutor);
});
Expand Down
4 changes: 2 additions & 2 deletions tests/unit/chipotle-executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,14 +63,14 @@ describe("ChipotleExecutor", () => {

it("is registered in executor index", async () => {
const { getExecutor } = await import("../../open-sse/executors/index.ts");
const exec = getExecutor("chipotle");
const exec = await getExecutor("chipotle");
assert.ok(exec, "chipotle executor should be registered");
assert.ok(exec instanceof ChipotleExecutor);
});

it("pepper alias works", async () => {
const { getExecutor } = await import("../../open-sse/executors/index.ts");
const exec = getExecutor("pepper");
const exec = await getExecutor("pepper");
assert.ok(exec, "pepper alias should be registered");
assert.ok(exec instanceof ChipotleExecutor);
});
Expand Down
Loading
Loading