Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
c8f1bb6
feat: add versioned session binding compatibility layer
Brisbanehuang Jul 20, 2026
0e90937
fix: observe late capability probe failures
Brisbanehuang Jul 20, 2026
fa8342c
fix: guard versioned session termination conflicts
Brisbanehuang Jul 20, 2026
626f804
test(proxy): keep hedge lifecycle assertions current
Brisbanehuang Jul 20, 2026
488bc8c
fix: guard legacy binding fallback races
Brisbanehuang Jul 20, 2026
09fbb57
fix: preserve legacy mirror during binding races
Brisbanehuang Jul 20, 2026
45dac94
fix: validate canonical provider before mirror restore
Brisbanehuang Jul 20, 2026
9ac3d99
fix(binding): preserve imported legacy mirrors
Brisbanehuang Jul 20, 2026
8adf020
fix(binding): guard legacy owner refresh races
Brisbanehuang Jul 20, 2026
21c1219
fix(binding): close legacy session races
Brisbanehuang Jul 20, 2026
8ccaac3
fix(binding): preserve failover session metadata
Brisbanehuang Jul 20, 2026
c433d35
fix(binding): preserve scoped versioned session state
Brisbanehuang Jul 20, 2026
95a94e8
feat(binding): add snapshot-safe TTL touch
Brisbanehuang Jul 21, 2026
9391c30
fix(binding): close long-stream and terminate races
Brisbanehuang Jul 21, 2026
81d9b4f
fix(binding): fence failed hedge cleanup
Brisbanehuang Jul 21, 2026
684e334
fix(session): scope content hash mapping by API key
Brisbanehuang Jul 21, 2026
8b1b678
fix(binding): avoid reusing canonical legacy mirrors
Brisbanehuang Jul 21, 2026
cd9a40d
Merge remote-tracking branch 'origin/integration/discovery-stack-2026…
ding113 Jul 22, 2026
25f0fa9
test(discovery): update versioned cleanup expectation
ding113 Jul 22, 2026
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
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
"test:coverage:proxy-guard-pipeline": "vitest run --config tests/configs/proxy-guard-pipeline.config.ts --coverage",
"test:coverage:include-session-id-in-errors": "vitest run --config tests/configs/include-session-id-in-errors.config.ts --coverage",
"test:coverage:usage-logs-sessionid-search": "vitest run --config tests/configs/usage-logs-sessionid-search.config.ts --coverage",
"test:coverage:session-binding": "vitest run --config tests/configs/session-binding.config.ts --coverage",
"test:ci": "vitest run --reporter=default --reporter=junit --outputFile.junit=reports/vitest-junit.xml",
"test:v1": "vitest run --config tests/configs/v1.config.ts --coverage --reporter=verbose && bun scripts/check-v1-critical-coverage.ts",
"openapi:generate": "bun scripts/generate-v1-types.ts",
Expand Down
41 changes: 33 additions & 8 deletions src/app/v1/_lib/proxy/forwarder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,10 @@ import {
import { ProxyProviderResolver } from "./provider-selector";
import { finalizeHedgeLoserBilling } from "./response-handler";
import type { ProxySession } from "./session";
import { setDeferredStreamingFinalization } from "./stream-finalization";
import {
type DeferredStreamingHedgeBindingAuthority,
setDeferredStreamingFinalization,
} from "./stream-finalization";
import {
detectThinkingBudgetRectifierTrigger,
rectifyThinkingBudget,
Expand Down Expand Up @@ -4571,8 +4574,12 @@ export class ProxyForwarder {

abortAllAttempts(attempt, "hedge_loser");

if (session.sessionId) {
void (async () => {
// A non-hedged request is finalized through response-handler. Updating
// here as well would perform a duplicate binding read/CAS before the
// stream has passed its final validation.
let hedgeBindingAuthorityPromise: Promise<DeferredStreamingHedgeBindingAuthority> | undefined;
if (session.sessionId && isActualHedgeWin) {
hedgeBindingAuthorityPromise = (async () => {
const bindingResult = await SessionManager.updateSessionBindingSmart(
session.sessionId!,
attempt.provider.id,
Expand All @@ -4595,15 +4602,31 @@ export class ProxyForwarder {
}

if (session.shouldTrackSessionObservability()) {
await SessionManager.updateSessionProvider(session.sessionId!, {
void SessionManager.updateSessionProvider(session.sessionId!, {
providerId: attempt.provider.id,
providerName: attempt.provider.name,
}).catch((observabilityError) => {
logger.error(
"ProxyForwarder: Failed to update observable session provider for hedge winner",
{ error: observabilityError }
);
});
}

return {
// Only the exact snapshot returned by the first-byte CAS may keep
// a versioned binding alive or clear it after a failed stream.
snapshot: bindingResult.bindingSnapshot ?? null,
// A generic clear is safe only when this request demonstrably
// committed the legacy binding. A versioned CAS conflict must not
// clear a newer generation that happens to use the same Provider.
legacyClearAllowed: bindingResult.legacyBindingUpdated === true,
};
})().catch((bindingError) => {
logger.error("ProxyForwarder: Failed to update session provider info for hedge winner", {
error: bindingError,
});
return { snapshot: null, legacyClearAllowed: false };
});
}

Expand All @@ -4620,6 +4643,7 @@ export class ProxyForwarder {
upstreamStatusCode: attempt.response.status,
isHedgeWinner: isActualHedgeWin,
billHedgeLosers,
hedgeBindingAuthorityPromise,
});

const response = new Response(
Expand Down Expand Up @@ -5011,16 +5035,17 @@ export class ProxyForwarder {
expectedProviderId: number | null
): Promise<void> {
if (!session.sessionId) return;
await SessionManager.clearSessionProvider(session.sessionId, expectedProviderId);
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
await SessionManager.clearSessionProvider(session.sessionId, expectedProviderId, keyId);
}

private static async clearSessionProviderBindings(
session: ProxySession,
expectedProviderIds: Iterable<number>
): Promise<void> {
for (const providerId of new Set(expectedProviderIds)) {
await ProxyForwarder.clearSessionProviderBinding(session, providerId);
}
if (!session.sessionId) return;
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
await SessionManager.clearSessionProviders(session.sessionId, expectedProviderIds, keyId);
}

private static markProviderFailed(
Expand Down
18 changes: 8 additions & 10 deletions src/app/v1/_lib/proxy/provider-selector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -492,10 +492,8 @@ export class ProxyProviderResolver {
}

// 从 Redis 读取该 session 绑定的 provider
const providerId = await SessionManager.getSessionProvider(
session.sessionId,
session.authState?.key?.id ?? null
);
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
const providerId = await SessionManager.getSessionProvider(session.sessionId, keyId);
if (!providerId) {
logger.debug("ProviderSelector: Session has no bound provider", {
sessionId: session.sessionId,
Expand All @@ -510,7 +508,7 @@ export class ProxyProviderResolver {
sessionId: session.sessionId,
providerId,
});
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
return null;
}

Expand All @@ -520,7 +518,7 @@ export class ProxyProviderResolver {
providerId: provider.id,
providerName: provider.name,
});
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
return null;
}

Expand All @@ -534,7 +532,7 @@ export class ProxyProviderResolver {
activeTimeEnd: provider.activeTimeEnd,
timezone: systemTimezone,
});
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
return null;
}

Expand Down Expand Up @@ -575,7 +573,7 @@ export class ProxyProviderResolver {
providerType: provider.providerType,
originalFormat: session.originalFormat,
});
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
return null;
}

Expand All @@ -594,7 +592,7 @@ export class ProxyProviderResolver {
// 清除过时绑定,避免 SET NX 死锁
// 当 session 内请求模型发生变化时,旧绑定已无意义,
// 清除后新的成功请求可通过 SET NX 重新绑定匹配的 provider
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
logger.info("ProviderSelector: Cleared stale provider binding (model mismatch)", {
sessionId: session.sessionId,
staleProviderId: provider.id,
Expand Down Expand Up @@ -650,7 +648,7 @@ export class ProxyProviderResolver {
],
},
});
await SessionManager.clearSessionProvider(session.sessionId, providerId);
await SessionManager.clearSessionProvider(session.sessionId, providerId, keyId);
return null;
}

Expand Down
133 changes: 130 additions & 3 deletions src/app/v1/_lib/proxy/response-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ import { isClientAbortError, isTransportError } from "./errors";
import type { ProxySession } from "./session";
import {
consumeDeferredStreamingFinalization,
type DeferredStreamingBindingHeartbeat,
peekDeferredStreamingFinalization,
} from "./stream-finalization";

Expand Down Expand Up @@ -97,6 +98,96 @@ const STREAM_FINALIZATION_MAX_MS = 120_000;
const STREAM_FAILURE_PERSISTENCE_MAX_MS = 5_000;
const NON_STREAM_TERMINAL_PERSISTENCE_ERROR = Symbol("non_stream_terminal_persistence_error");

function startHedgeBindingHeartbeat(session: ProxySession): void {
const deferred = peekDeferredStreamingFinalization(session);
const authorityPromise = deferred?.hedgeBindingAuthorityPromise;
if (!deferred?.isHedgeWinner || !authorityPromise || deferred.hedgeBindingHeartbeat) return;

let periodicActive = true;
let authorityLost = false;
let timer: ReturnType<typeof setInterval> | null = null;
let touchInFlight: Promise<boolean> | null = null;
let completionPromise: Promise<void> | null = null;

const stopPeriodic = () => {
periodicActive = false;
if (timer) {
clearInterval(timer);
timer = null;
}
};

const loseAuthority = (status: string, reason?: string) => {
if (authorityLost) return;
authorityLost = true;
stopPeriodic();
logger.warn("[ResponseHandler] Hedge binding heartbeat stopped", {
sessionId: session.sessionId,
status,
reason,
});
};

const touch = (allowAfterStop = false): Promise<boolean> => {
if (authorityLost || (!periodicActive && !allowAfterStop)) return Promise.resolve(false);
if (touchInFlight) return touchInFlight;

const operation = (async () => {
const { snapshot } = await authorityPromise;
if (!snapshot) {
authorityLost = true;
stopPeriodic();
return false;
}
if (authorityLost || (!periodicActive && !allowAfterStop)) return false;

const touched = await SessionManager.touchVersionedSessionBinding(snapshot);
if (
touched.status !== "ok" ||
touched.snapshot.generation !== snapshot.generation ||
touched.snapshot.providerId !== snapshot.providerId
) {
loseAuthority(touched.status, "reason" in touched ? touched.reason : "snapshot_mismatch");
return false;
}
return true;
})()
.catch((error) => {
loseAuthority("error", error instanceof Error ? error.message : String(error));
return false;
})
.finally(() => {
if (touchInFlight === operation) touchInFlight = null;
});
touchInFlight = operation;
return operation;
};

const lifecycle: DeferredStreamingBindingHeartbeat = {
stop: stopPeriodic,
complete: () => {
if (completionPromise) return completionPromise;
stopPeriodic();
completionPromise = (async () => {
if (touchInFlight) await touchInFlight;
if (authorityLost) return;
await touch(true);
})();
return completionPromise;
},
};
deferred.hedgeBindingHeartbeat = lifecycle;

// The first touch validates that ownership really transferred with the
// first-byte CAS. Subsequent touches keep streams longer than SESSION_TTL alive.
void touch();
const intervalMs = Math.max(250, SessionManager.getVersionedSessionBindingRefreshIntervalMs());
timer = setInterval(() => {
void touch();
}, intervalMs);
timer.unref?.();
}

type MessageRequestTerminalDetails = Parameters<typeof updateMessageRequestDetailsDurably>[1];
type NonStreamTerminalPersistenceError = Error & {
[NON_STREAM_TERMINAL_PERSISTENCE_ERROR]: true;
Expand Down Expand Up @@ -1149,7 +1240,29 @@ function finalizeDeferredStreamingFinalizationIfNeeded(
const providerIdForPersistence = meta?.providerId ?? provider?.id ?? null;
const clearSessionBinding = async () => {
if (!session.sessionId) return;
await SessionManager.clearSessionProvider(session.sessionId, providerIdForPersistence);
const hedgeAuthority = meta?.isHedgeWinner
? await meta.hedgeBindingAuthorityPromise
: undefined;
if (hedgeAuthority?.snapshot) {
const hedgeSnapshot = hedgeAuthority.snapshot;
const cleared = await SessionManager.clearVersionedSessionProvider(
hedgeSnapshot,
providerIdForPersistence
);
if (cleared.status !== "ok") {
logger.warn("[ResponseHandler] Hedge winner binding clear stopped", {
sessionId: hedgeSnapshot.sessionId,
providerId: providerIdForPersistence,
reason: cleared.reason,
});
}
return;
}
if (meta?.isHedgeWinner && !hedgeAuthority?.legacyClearAllowed) {
return;
}
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
await SessionManager.clearSessionProvider(session.sessionId, providerIdForPersistence, keyId);
};

const isHedgeWinner = meta?.isHedgeWinner === true;
Expand Down Expand Up @@ -1239,11 +1352,15 @@ function finalizeDeferredStreamingFinalizationIfNeeded(
((clientAborted || !streamEndedNormally) && !clientAbortCompleteSuccess) ||
detected.isError ||
(upstreamStatusCode >= 400 && errorMessage !== null);
if (shouldClearSessionBindingOnFailure) {
meta?.hedgeBindingHeartbeat?.stop();
}

// 未启用延迟结算 / provider 缺失:
// - 只返回“内部状态码 + 错误原因”,由调用方写入统计;
// - 不在这里更新熔断/绑定(meta 缺失意味着 Forwarder 没有启用延迟结算;provider 缺失意味着无法归因)。
if (!meta || !provider) {
meta?.hedgeBindingHeartbeat?.stop();
return {
effectiveStatusCode,
errorMessage,
Expand Down Expand Up @@ -1461,7 +1578,13 @@ function finalizeDeferredStreamingFinalizationIfNeeded(
});
}

// Stop periodic refresh at the stream boundary and issue one final
// generation-safe touch so the next turn receives a full binding TTL.
// complete() is idempotent and permanently stops after any authority conflict.
const hedgeBindingCompletion = meta.hedgeBindingHeartbeat?.complete();
const commitSideEffects = async () => {
await hedgeBindingCompletion;

if (meta.endpointId != null) {
try {
const { recordEndpointSuccess } = await import("@/lib/endpoint-circuit-breaker");
Expand Down Expand Up @@ -1806,7 +1929,8 @@ export class ProxyResponseHandler {
if (session.sessionId) {
const sessionId = session.sessionId;
postTerminalSideEffects.push(async () => {
await SessionManager.clearSessionProvider(sessionId, provider.id);
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
await SessionManager.clearSessionProvider(sessionId, provider.id, keyId);
});
}
if (
Expand Down Expand Up @@ -1976,7 +2100,8 @@ export class ProxyResponseHandler {
if (session.sessionId) {
const sessionId = session.sessionId;
postTerminalSideEffects.push(async () => {
await SessionManager.clearSessionProvider(sessionId, provider.id);
const keyId = session.authState?.key?.id ?? session.messageContext?.key?.id ?? null;
await SessionManager.clearSessionProvider(sessionId, provider.id, keyId);

const sessionUsagePayload: SessionUsageUpdate = {
status:
Expand Down Expand Up @@ -2533,6 +2658,8 @@ export class ProxyResponseHandler {
return response;
}

startHedgeBindingHeartbeat(session);

let processedStream: ReadableStream<Uint8Array> = response.body;

// --- GEMINI STREAM HANDLING ---
Expand Down
17 changes: 17 additions & 0 deletions src/app/v1/_lib/proxy/stream-finalization.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,18 @@
import type { SessionBindingSnapshot } from "@/lib/redis/session-binding";
import type { ProxySession } from "./session";

export type DeferredStreamingBindingHeartbeat = {
stop: () => void;
complete: () => Promise<void>;
};

export type DeferredStreamingHedgeBindingAuthority = {
/** Exact versioned generation written by the first-byte winner, when available. */
snapshot: SessionBindingSnapshot | null;
/** Only a confirmed successful legacy write may use the non-versioned clear path. */
legacyClearAllowed: boolean;
};

/**
* 流式响应(SSE)在“收到响应头”时无法确定成功与否:
* - 上游可能返回 HTTP 200,但 body 是错误 JSON(假 200)
Expand Down Expand Up @@ -35,6 +48,10 @@ export type DeferredStreamingFinalization = {
* coexists with asynchronously accumulated loser costs without clobbering.
*/
billHedgeLosers?: boolean;
/** Binding authority established by the legacy Hedge winner's first-byte write. */
hedgeBindingAuthorityPromise?: Promise<DeferredStreamingHedgeBindingAuthority>;
/** ResponseHandler-owned runtime lifecycle; attached when streaming starts. */
hedgeBindingHeartbeat?: DeferredStreamingBindingHeartbeat;
};

const deferredMeta = new WeakMap<ProxySession, DeferredStreamingFinalization>();
Expand Down
Loading