mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-26 00:52:18 +03:00
perf(proxy): implement non-blocking async proxy log batching and performance optimizations (A, B, C, D) (#11182)
Validated on the combined batch board + this branch: perf-a-b-c-d + account-fallback-service 93/93, gates + typecheck clean. Owner approved the full bundle including item A (async proxy-log batching, 1s/100-item flush with an unref'd timer — reviewed the implementation: flush helper exists for shutdown wiring, buffered logs are the accepted tradeoff). B (lazy modals), C (O(1) alias maps), D (pre-compiled regex — this is also the entire content of #11187, being closed as subsumed) ride along. Conflict with the tip was only stale provider-count docs. Thank you @rqzbeh!
This commit is contained in:
@@ -6,11 +6,10 @@ function isRecord(value: unknown): value is Record<string, unknown> {
|
|||||||
return !!value && typeof value === "object" && !Array.isArray(value);
|
return !!value && typeof value === "object" && !Array.isArray(value);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const PREFIX_TRIM_RE = /^(?:models\/|antigravity\/)+/i;
|
||||||
|
|
||||||
function normalizeCloudCodeModel(model: string): string {
|
function normalizeCloudCodeModel(model: string): string {
|
||||||
return String(model || "")
|
return String(model || "").trim().replace(PREFIX_TRIM_RE, "");
|
||||||
.trim()
|
|
||||||
.replace(/^models\//i, "")
|
|
||||||
.replace(/^antigravity\//i, "");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
function stripGeminiThinkingConfig(value: unknown): unknown {
|
function stripGeminiThinkingConfig(value: unknown): unknown {
|
||||||
|
|||||||
@@ -46,9 +46,19 @@ import {
|
|||||||
getCodexGlobalServiceMode,
|
getCodexGlobalServiceMode,
|
||||||
type CodexGlobalServiceMode,
|
type CodexGlobalServiceMode,
|
||||||
} from "@/lib/providers/codexFastTier";
|
} from "@/lib/providers/codexFastTier";
|
||||||
import AddCompatibleProviderModal from "./components/AddCompatibleProviderModal";
|
import dynamic from "next/dynamic";
|
||||||
|
const AddCompatibleProviderModal = dynamic(
|
||||||
|
() => import("./components/AddCompatibleProviderModal"),
|
||||||
|
{ ssr: false }
|
||||||
|
);
|
||||||
import { CategoryDot } from "./components/CategoryDot";
|
import { CategoryDot } from "./components/CategoryDot";
|
||||||
import { ImportProvidersFromFileModal } from "./components/ImportProvidersFromFileModal";
|
const ImportProvidersFromFileModal = dynamic(
|
||||||
|
() =>
|
||||||
|
import("./components/ImportProvidersFromFileModal").then(
|
||||||
|
(m) => m.ImportProvidersFromFileModal
|
||||||
|
),
|
||||||
|
{ ssr: false }
|
||||||
|
);
|
||||||
import NoAuthProvidersSection from "./components/NoAuthProvidersSection";
|
import NoAuthProvidersSection from "./components/NoAuthProvidersSection";
|
||||||
import HighlightableProviderCard from "./components/HighlightableProviderCard";
|
import HighlightableProviderCard from "./components/HighlightableProviderCard";
|
||||||
import ProviderCountBadge from "./components/ProviderCountBadge";
|
import ProviderCountBadge from "./components/ProviderCountBadge";
|
||||||
|
|||||||
@@ -194,45 +194,105 @@ export function logProxyEvent(entry: ProxyLogInput) {
|
|||||||
proxyLogs.length = MAX_IN_MEMORY_ENTRIES;
|
proxyLogs.length = MAX_IN_MEMORY_ENTRIES;
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. Persist to SQLite
|
// 2. Queue for background batch persistence (SQLite / Redis)
|
||||||
if (shouldPersistToDisk) {
|
if (shouldPersistToDisk) {
|
||||||
try {
|
enqueueProxyLog(log);
|
||||||
const db = getDbInstance();
|
|
||||||
db.prepare(
|
|
||||||
`INSERT INTO proxy_logs (id, timestamp, status, proxy_type, proxy_host, proxy_port,
|
|
||||||
level, level_id, provider, target_url, public_ip, egress_ip, latency_ms, error,
|
|
||||||
connection_id, combo_id, account, tls_fingerprint)
|
|
||||||
VALUES (@id, @timestamp, @status, @proxyType, @proxyHost, @proxyPort,
|
|
||||||
@level, @levelId, @provider, @targetUrl, @clientIp, @egressIp, @latencyMs, @error,
|
|
||||||
@connectionId, @comboId, @account, @tlsFingerprint)`
|
|
||||||
).run({
|
|
||||||
id: log.id,
|
|
||||||
timestamp: log.timestamp,
|
|
||||||
status: log.status,
|
|
||||||
proxyType: log.proxy?.type || null,
|
|
||||||
proxyHost: log.proxy?.host || null,
|
|
||||||
proxyPort: log.proxy?.port ? Number(log.proxy.port) : null,
|
|
||||||
level: log.level,
|
|
||||||
levelId: log.levelId,
|
|
||||||
provider: log.provider,
|
|
||||||
targetUrl: log.targetUrl,
|
|
||||||
clientIp: log.clientIp,
|
|
||||||
egressIp: log.egressIp,
|
|
||||||
latencyMs: log.latencyMs,
|
|
||||||
error: log.error,
|
|
||||||
connectionId: log.connectionId,
|
|
||||||
comboId: log.comboId,
|
|
||||||
account: log.account,
|
|
||||||
tlsFingerprint: log.tlsFingerprint ? 1 : 0,
|
|
||||||
});
|
|
||||||
} catch (err: any) {
|
|
||||||
console.warn("[proxyLogger] Failed to persist:", err.message);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return log;
|
return log;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ──────────────── Background Batch Persistence ────────────────
|
||||||
|
|
||||||
|
const BATCH_FLUSH_INTERVAL_MS = 1000;
|
||||||
|
const BATCH_SIZE_THRESHOLD = 100;
|
||||||
|
|
||||||
|
let pendingLogsQueue: ProxyLogEntry[] = [];
|
||||||
|
let batchTimer: NodeJS.Timeout | null = null;
|
||||||
|
|
||||||
|
function ensureBatchTimer() {
|
||||||
|
if (batchTimer) return;
|
||||||
|
batchTimer = setInterval(() => {
|
||||||
|
flushProxyLogsSync();
|
||||||
|
}, BATCH_FLUSH_INTERVAL_MS);
|
||||||
|
if (typeof batchTimer.unref === "function") {
|
||||||
|
batchTimer.unref();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function enqueueProxyLog(log: ProxyLogEntry) {
|
||||||
|
pendingLogsQueue.push(log);
|
||||||
|
ensureBatchTimer();
|
||||||
|
if (pendingLogsQueue.length >= BATCH_SIZE_THRESHOLD) {
|
||||||
|
flushProxyLogsSync();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function flushProxyLogsSync() {
|
||||||
|
if (pendingLogsQueue.length === 0) return;
|
||||||
|
const batch = pendingLogsQueue;
|
||||||
|
pendingLogsQueue = [];
|
||||||
|
|
||||||
|
// 1. If Redis driver is active, asynchronously publish batch to Redis Stream/Channel
|
||||||
|
if (process.env.QUOTA_STORE_DRIVER === "redis" || process.env.QUOTA_STORE_REDIS_URL) {
|
||||||
|
try {
|
||||||
|
import("@/lib/quota/redisQuotaStore").then(({ getRedisQuotaStore }) => {
|
||||||
|
const store = getRedisQuotaStore(process.env.QUOTA_STORE_REDIS_URL || "");
|
||||||
|
const client = (store as any)?.client;
|
||||||
|
if (client && typeof client.publish === "function") {
|
||||||
|
for (const entry of batch) {
|
||||||
|
client.publish("omniroute:proxy_logs", JSON.stringify(entry)).catch(() => {});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}).catch(() => {});
|
||||||
|
} catch {
|
||||||
|
/* ignore redis pub errors */
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. Persist to SQLite using a single transaction for high-performance non-blocking write
|
||||||
|
try {
|
||||||
|
const db = getDbInstance();
|
||||||
|
const insertStmt = db.prepare(
|
||||||
|
`INSERT INTO proxy_logs (id, timestamp, status, proxy_type, proxy_host, proxy_port,
|
||||||
|
level, level_id, provider, target_url, public_ip, egress_ip, latency_ms, error,
|
||||||
|
connection_id, combo_id, account, tls_fingerprint)
|
||||||
|
VALUES (@id, @timestamp, @status, @proxyType, @proxyHost, @proxyPort,
|
||||||
|
@level, @levelId, @provider, @targetUrl, @clientIp, @egressIp, @latencyMs, @error,
|
||||||
|
@connectionId, @comboId, @account, @tlsFingerprint)`
|
||||||
|
);
|
||||||
|
|
||||||
|
const transaction = db.transaction((entries: ProxyLogEntry[]) => {
|
||||||
|
for (const item of entries) {
|
||||||
|
insertStmt.run({
|
||||||
|
id: item.id,
|
||||||
|
timestamp: item.timestamp,
|
||||||
|
status: item.status,
|
||||||
|
proxyType: item.proxy?.type || null,
|
||||||
|
proxyHost: item.proxy?.host || null,
|
||||||
|
proxyPort: item.proxy?.port ? Number(item.proxy.port) : null,
|
||||||
|
level: item.level,
|
||||||
|
levelId: item.levelId,
|
||||||
|
provider: item.provider,
|
||||||
|
targetUrl: item.targetUrl,
|
||||||
|
clientIp: item.clientIp,
|
||||||
|
egressIp: item.egressIp,
|
||||||
|
latencyMs: item.latencyMs,
|
||||||
|
error: item.error,
|
||||||
|
connectionId: item.connectionId,
|
||||||
|
comboId: item.comboId,
|
||||||
|
account: item.account,
|
||||||
|
tlsFingerprint: item.tlsFingerprint ? 1 : 0,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
transaction(batch);
|
||||||
|
} catch (err: any) {
|
||||||
|
console.warn("[proxyLogger] Failed to write proxy log batch to disk:", err?.message || err);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ──────────────── Query ────────────────
|
// ──────────────── Query ────────────────
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
26
tests/unit/perf-a-b-c-d.test.ts
Normal file
26
tests/unit/perf-a-b-c-d.test.ts
Normal file
@@ -0,0 +1,26 @@
|
|||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
import { logProxyEvent, flushProxyLogsSync } from "../../src/lib/proxyLogger.ts";
|
||||||
|
import { shouldStripCloudCodeThinking } from "../../open-sse/services/cloudCodeThinking.ts";
|
||||||
|
|
||||||
|
test("Part A: async proxy log batching queues entries without synchronous failure", () => {
|
||||||
|
const sampleLog = {
|
||||||
|
status: "success",
|
||||||
|
provider: "test-provider",
|
||||||
|
latencyMs: 15,
|
||||||
|
};
|
||||||
|
|
||||||
|
const logged = logProxyEvent(sampleLog);
|
||||||
|
assert.equal(logged.provider, "test-provider");
|
||||||
|
assert.equal(typeof logged.id, "string");
|
||||||
|
|
||||||
|
// Ensure flush completes without throwing
|
||||||
|
assert.doesNotThrow(() => {
|
||||||
|
flushProxyLogsSync();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
test("Part D: pre-compiled regex in cloudCodeThinking model normalization", () => {
|
||||||
|
assert.equal(shouldStripCloudCodeThinking("antigravity", "antigravity/claude-3-7-sonnet"), true);
|
||||||
|
assert.equal(shouldStripCloudCodeThinking("antigravity", "models/gemini-2.5-pro"), false);
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user