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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(proxy-logs):** search proxy-log rows keep the HTTP status the provider actually returned (`upstream_status`, null when no response arrived), matching the chat writer ([#14220](https://github.com/diegosouzapw/OmniRoute/pull/14220)) — thanks @maxmad64bis
28 changes: 23 additions & 5 deletions open-sse/handlers/search/searchProxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,13 @@ export async function fetchWithSearchProxy(
/**
* Emit a sanitized proxy event for a search provider attempt.
* Never includes query, API key, proxy username, or proxy password.
*
* `upstreamStatus` carries the HTTP status the provider actually returned for
* this attempt (null when no response arrived). The other proxy-log writers
* correctly keep null: the history websocket writer derives its status
* post-hoc (fallbacks instead of a received response), and the provider-test
* writer sometimes synthesizes its status code (network failure, refresh
* failure) — copying either number would fabricate a status.
*/
export async function emitSearchProxyEvent(
provider: string,
Expand All @@ -88,7 +95,8 @@ export async function emitSearchProxyEvent(
proxyLevel: string,
targetUrl: string,
startTime: number,
status: string
status: string,
upstreamStatus: number | null = null
): Promise<void> {
try {
const { logProxyEvent } = await import("@/lib/proxyLogger");
Expand All @@ -112,6 +120,7 @@ export async function emitSearchProxyEvent(
: null;
logProxyEvent({
status,
upstreamStatus,
proxy: proxyInfo,
level: proxyLevel,
levelId: connectionId || null,
Expand Down Expand Up @@ -187,8 +196,17 @@ export async function executeProviderFetch(
): Promise<ProviderFetchResult> {
const { config, url, init, controller, timer, query, searchType, maxResults, startTime } = p;
const { connectionId, proxy, proxyLevel, log, normalize } = p;
const emitEvent = (status: string) =>
emitSearchProxyEvent(config.id, connectionId, proxy, proxyLevel, url, startTime, status);
const emitEvent = (status: string, upstreamStatus: number | null = null) =>
emitSearchProxyEvent(
config.id,
connectionId,
proxy,
proxyLevel,
url,
startTime,
status,
upstreamStatus
);
const logCall = (fields: Record<string, unknown>) =>
saveCallLog({
method: config.method,
Expand Down Expand Up @@ -227,7 +245,7 @@ export async function executeProviderFetch(
duration: Date.now() - startTime,
error: errorText.slice(0, 500),
});
await emitEvent("error");
await emitEvent("error", response.status);
return {
success: false,
status: response.status,
Expand All @@ -246,7 +264,7 @@ export async function executeProviderFetch(
tokens: { prompt_tokens: 0, completion_tokens: 0 },
responseBody: { results_count: results.length, cached: false },
});
await emitEvent("success");
await emitEvent("success", response.status);

return {
success: true,
Expand Down
157 changes: 157 additions & 0 deletions tests/unit/search-proxy-upstream-status.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
import test from "node:test";
import assert from "node:assert/strict";
import fs from "node:fs";
import http from "node:http";
import os from "node:os";
import path from "node:path";
import type { AddressInfo } from "node:net";

// Search proxy-log rows keep the HTTP status the provider actually
// returned. Real codes (incl. 429/500 from the wire) land on the journal line;
// locally synthesized codes (transport/timeout, envelope errors) stay null.

const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-search-upstream-"));
process.env.DATA_DIR = TEST_DATA_DIR;
process.env.API_KEY_SECRET = process.env.API_KEY_SECRET || "search-upstream-test-secret";
process.env.DISABLE_SQLITE_AUTO_BACKUP = "true";

const core = await import("../../src/lib/db/core.ts");
const proxyLogger = await import("../../src/lib/proxyLogger.ts");
const searchProxy = await import("../../open-sse/handlers/search/searchProxy.ts");
const { SEARCH_PROVIDERS } = await import("../../open-sse/config/searchRegistry.ts");
const { closeCallLogSaves } = await import("../../src/lib/usage/callLogs.ts");

let server: http.Server;
let baseUrl = "";

test.before(async () => {
server = http.createServer((req, res) => {
const url = new URL(req.url ?? "/", "http://local");
const code = Number(url.searchParams.get("code") ?? 200);
if (url.searchParams.get("hang") === "1") return; // never respond -> caller aborts
res.writeHead(code, { "Content-Type": "application/json", Connection: "close" });
if (code >= 200 && code < 300) {
res.end(JSON.stringify({ results: [], total: 0 }));
} else {
res.end(JSON.stringify({ error: `provider error ${code}` }));
}
});
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", () => resolve()));
baseUrl = `http://127.0.0.1:${(server.address() as AddressInfo).port}`;
});

test.after(async () => {
await new Promise<void>((resolve) => server.close(() => resolve()));
proxyLogger.clearProxyLogs();
await closeCallLogSaves(500).catch(() => {});
try {
core.resetDbInstance();
} catch {}
try {
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 });
} catch {}
});

test.beforeEach(() => {
proxyLogger.clearProxyLogs();
});

function fetchParams(pathSuffix: string, timeoutMs = 5000) {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), timeoutMs);
timer.unref?.();
return {
config: SEARCH_PROVIDERS["tavily-search"],
url: `${baseUrl}${pathSuffix}`,
init: { method: "POST", headers: { "Content-Type": "application/json" } },
controller,
timer,
query: "test query",
searchType: "web",
maxResults: 5,
startTime: Date.now(),
connectionId: undefined,
proxy: null,
proxyLevel: "direct",
normalize: () => ({ results: [], totalResults: 0 }),
};
}

async function latestUpstreamStatus() {
const [entry] = proxyLogger.getProxyLogs();
assert.ok(entry, "a proxy log line must have been emitted");
return entry.upstreamStatus;
}

test("a real 429 from the wire lands on the journal line", async () => {
const p = fetchParams("/search?code=429");
try {
const result = await searchProxy.executeProviderFetch(p);
assert.equal(result.success, false);
assert.equal(result.status, 429);
assert.equal(await latestUpstreamStatus(), 429);
} finally {
clearTimeout(p.timer);
}
});

test("a real 200 from the wire lands on the journal line", async () => {
const p = fetchParams("/search?code=200");
try {
const result = await searchProxy.executeProviderFetch(p);
assert.equal(result.success, true);
assert.equal(await latestUpstreamStatus(), 200);
} finally {
clearTimeout(p.timer);
}
});

test("a real 500 from the wire lands on the journal line (not the 502 transport code)", async () => {
const p = fetchParams("/search?code=500");
try {
const result = await searchProxy.executeProviderFetch(p);
assert.equal(result.success, false);
assert.equal(result.status, 500);
assert.equal(await latestUpstreamStatus(), 500);
} finally {
clearTimeout(p.timer);
}
});

test("a caller abort (no response received) keeps null", async () => {
const p = fetchParams("/search?hang=1", 200);
try {
const result = await searchProxy.executeProviderFetch(p);
assert.equal(result.success, false);
assert.equal(await latestUpstreamStatus(), null);
} finally {
clearTimeout(p.timer);
}
});

for (const quota of [true, false]) {
test(`an AnySearch envelope error (quota=${quota}) keeps null despite result.status`, async () => {
const p = fetchParams("/search?code=200");
const envelope = new Error("AnySearch envelope error") as Error & { quota?: boolean };
envelope.name = "AnysearchSearchEnvelopeError";
envelope.quota = quota;
const failingNormalize = () => {
throw envelope;
};
try {
const result = await searchProxy.executeProviderFetch({
...p,
normalize: failingNormalize,
});
assert.equal(result.success, false);
assert.equal(result.status, quota ? 402 : 502);
assert.equal(
await latestUpstreamStatus(),
null,
"synthesized envelope codes must never reach upstream_status"
);
} finally {
clearTimeout(p.timer);
}
});
}
Loading