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
1 change: 1 addition & 0 deletions changelog.d/features/14810-call-resilience-actions.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **feat(call-logs):** show per-request recovery actions in the request log ([#14810](https://github.com/diegosouzapw/OmniRoute/pull/14810)) — thanks @maxmad64bis
9 changes: 5 additions & 4 deletions config/quality/file-size-baseline.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"_rebaseline_2026_09_25_14810_resilience_actions_wiring": "PR #14810 own growth (resilience-actions observability wiring, irreducible call-site lines after extraction of the note logic into chatCore/*_notes modules + opencodeResilienceNotes + requestLoggerResilience, all under cap): open-sse/handlers/chatCore.ts 6426->6441 (+15 = wrapper handleChatCore/handleChatCoreInner + resumed-flag param/note/strip + 3 one-line note calls, logic in chatCore/resilienceAttemptContext + resumedResilienceNotes + emptyTurnResilienceNotes + continuationResilienceNotes + recoveryTraceLogging); open-sse/executors/opencode.ts 1338->1354, re-measured 1354->1370 after merging release/v3.8.51 @86870a82 (the tip itself had grown to 1354 via the applied-egress tracker) (+16 = served-account tracker + park/replay/stored notes wiring, logic in opencodeResilienceNotes); src/sse/handlers/chat.ts 2567->2582 (+15 on the reconciled tip = resumed-flag threading via body marker + dispatch param + 2 combo forwards); src/shared/components/RequestLoggerDetail.tsx 1210->1216, re-measured on release/v3.8.51 @4b9388a1 (tip file 1205 after #14795; +11 = one import + badge computation + badge render loop, badge logic in requestLoggerResilience); src/lib/usage/callLogs.ts stays under the 1200 cap (resilience-actions parser moved to src/lib/usage/resilienceActionsParse.ts). Measured with split(\"\\n\").length. Covered by tests/unit/resilience-actions-context (11) + sink (5) + notes (9) + badges (5) + db/migration-191 (2).",
"_rebaseline_2026_09_24_14802_transport_setaside": "PR #14802 own growth: open-sse/utils/proxyFetch.ts 1334->1341 (+7 = two-symbol import + proxied-success hook at the single dispatch return + tagged-failure guard at the single final throw; full evidence helpers live in open-sse/utils/proxyRefusalMemory.ts (not frozen), irreducible call-site wiring at the two existing chokepoints). Covered by tests/unit/proxy-transport-setaside-cross-evidence.test.ts (10/10) + tests/unit/proxy-transport-traffic-simulation.test.ts (4/4).",
"_rebaseline_2026_09_24_14708_persisted_cooldown_precedence": "PR #14708 own growth: open-sse/services/combo/roundRobinCombo.ts 1263->1285 (+22 = the persisted-cooldown pre-check at the only round-robin dispatch point with no persisted check (import + async pre-check + skip/continue), plus the sticky-target reflow). Irreducible: the check must live at the dispatch site that bypasses executeTargetAttempt. Covered by tests/unit/combo-transient-persisted-precedence.test.ts.",
"_rebaseline_2026_09_24_14367_synthetic_key_budget_orphan": "PR #14367 (original fix; #14661 re-landed the same diff): src/lib/db/apiKeys.ts 1718->1719 (+1 = the single import of SYNTHETIC_ENV_API_KEY_ID threading the canonical synthetic-identity constant into getApiKeyMetadata; the exclusion list lives in src/shared/constants/apiKeyIdentities.ts). Irreducible. Covered by tests/unit/db/health-check-orphan-budgets.test.ts.",
Expand Down Expand Up @@ -500,7 +501,7 @@
"open-sse/executors/codex.ts": 1584,
"open-sse/executors/cursor.ts": 1868,
"open-sse/executors/muse-spark-web.ts": 1405,
"open-sse/handlers/chatCore.ts": 6426,
"open-sse/handlers/chatCore.ts": 6441,
"open-sse/handlers/imageGeneration.ts": 3334,
"open-sse/handlers/search.ts": 1789,
"open-sse/mcp-server/schemas/tools.ts": 1621,
Expand Down Expand Up @@ -540,18 +541,18 @@
"src/shared/components/RequestLoggerV2.tsx": 1748,
"src/shared/constants/providers/apikey/gateways.ts": 1544,
"src/shared/services/cliRuntime.ts": 1296,
"src/sse/handlers/chat.ts": 2567,
"src/sse/handlers/chat.ts": 2582,
"src/sse/services/auth.ts": 3602,
"tests/unit/account-fallback-service.test.ts": 2453,
"tests/unit/provider-validation-specialty.test.ts": 4656,
"open-sse/services/autoCombo/virtualFactory.ts": 1258,
"open-sse/services/combo/roundRobinCombo.ts": 1285,
"src/shared/components/RequestLoggerDetail.tsx": 1210,
"src/shared/components/RequestLoggerDetail.tsx": 1216,
"src/shared/middleware/chatBodyAdmission.ts": 1206,
"open-sse/executors/deepseek-web.ts": 1224,
"open-sse/executors/default.ts": 1207,
"open-sse/services/rateLimitManager.ts": 1329,
"open-sse/executors/opencode.ts": 1338
"open-sse/executors/opencode.ts": 1370
},
"_rebaseline_2026_09_15_roundrobin_dashboard_events": "Fix #13089 (Combo Studio Live dashboard shows an empty backlog for round-robin combos): open-sse/services/combo/roundRobinCombo.ts 1205->1213. Round-robin is the only combo strategy that bypasses handleComboChat/executeTargetAttempt.ts, the path that publishes the combo.target.attempt/succeeded/failed EventBus events the Live dashboard listens for — so round-robin completions never showed up. The new call-site wiring (createRRDashboardEvents(...) instantiated once per target, one-line .attempt()/.succeeded()/.failed() calls at the 6 existing dispatch/outcome points) is the emitter logic actually extracted into a new module, open-sse/services/combo/rrDashboardEvents.ts — this is the minimum irreducible footprint for wiring 6 required call sites into 6 fixed control-flow points of the frozen file. Covered by tests/unit/issue-13089-roundrobin-live-ws-events.test.ts (2 tests: success + failure paths).",
"_rebaseline_base_2026_08_10_proxyfetch": "Base-red fix (green-prs sweep, issue #9985): open-sse/utils/proxyFetch.ts 1207 > cap 1000 — new proxied-TLS fetch helper introduced by the Fal reference-image work. Owner-authorized quick rebaseline to green; structural slim tracked for v3.9.0.",
Expand Down
16 changes: 16 additions & 0 deletions open-sse/executors/opencode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,12 @@ import {
noteRotationAccount,
resolveProxyForRequest,
} from "../utils/proxyFetch.ts";
import {
createServedAccountTracker,
noteParkWait,
noteReplayed,
noteStoredFallback,
} from "./opencodeResilienceNotes.ts";
import {
clientSuppliedOpencodeSession,
forwardOpencodeClientHeaders,
Expand Down Expand Up @@ -645,6 +651,8 @@ export class OpencodeExecutor extends BaseExecutor {
let burstStreak = 0,
parked = false;
const requestPacing = egressPacing.initEgressPacingForRequest(); // Off by default.
// served-account changes (effective-change counting) live in the leaf tracker.
const noteServedAccount = createServedAccountTracker();
const appliedEgress = egressPacing.createAppliedEgressTracker(
this.buildUrl(String(input.model ?? ""), Boolean(input.stream)),
resolveProxyForRequest
Expand Down Expand Up @@ -749,6 +757,9 @@ export class OpencodeExecutor extends BaseExecutor {
if (attributionOn && (accounts.length > 1 || account.fingerprint !== "")) {
noteRotationAccount(masked);
}
// effective-change counting on the masked id (the raw
// fingerprint never reaches the log, just the counter).
noteServedAccount(masked);

// Pin egress to this account's proxy for the whole BaseExecutor dispatch
// (incl. its intra-URL 429 retries). skipUpstreamRetry lets THIS loop own
Expand Down Expand Up @@ -905,6 +916,8 @@ export class OpencodeExecutor extends BaseExecutor {
"OPENCODE",
`${cid}burstStreak=${burstStreak} freshD2=${marker.fresh} park`
);
// local monotone park measure (Date.now diff, integer ms).
const parkStartMs = Date.now();
const p = await runParkAndReplay(
{
execute: (i: ExecuteInput) =>
Expand All @@ -920,17 +933,20 @@ export class OpencodeExecutor extends BaseExecutor {
log,
cid
);
noteParkWait(Date.now() - parkStartMs);
if (p && p !== result) {
if (attributionOn && skippedCooldown.size > 0) {
this.logSkippedCooldownAccounts(log, cid, skippedCooldown);
}
noteReplayed();
return this.normalizeMuseSparkResponse(input, p);
}
if (p) {
discardResponseBody(abandonedResponse);
if (attributionOn && skippedCooldown.size > 0) {
this.logSkippedCooldownAccounts(log, cid, skippedCooldown);
}
noteStoredFallback();
return this.normalizeMuseSparkResponse(input, result);
}
}
Expand Down
15 changes: 14 additions & 1 deletion open-sse/executors/opencodeParkResume.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import { PARKED_STREAM_HEADER, PARKED_STREAM_VALUE } from "../utils/streamReadin
import { isProxyAvoided, proxyEgressKey, proxySetAsideSeq } from "../utils/proxyRefusalMemory.ts";
import { maskAccountId, type RotatableAccount } from "./accountRotation.ts";
import { runWithProxyContext } from "../utils/proxyFetch.ts";
import { noteResilienceAction } from "@/lib/usage/resilienceActionsContext.ts";
import type { ExecuteInput, ExecutorExecuteResult } from "./base.ts";

/** Consecutive transient 429s before a request parks. */
Expand Down Expand Up @@ -251,10 +252,22 @@ export async function runParkAndReplay<TAccount extends RotatableAccount>(
}
const probe = await replayOneLeg(driver, input, driver.accounts, log, cid);
const finalBody = probe?.result.response ?? fallback.response;
// stream note: note only AFTER the recopy outcome is known. A stored 429
// fallback recopied into the 200 SSE envelope is still a stored
// error replay — flag alone decides at read time.
let recopied = true;
try {
await copyFinalBodyAsValidFrames(controller, encoder, finalBody);
} catch {
/* unreadable body — close with the pings already sent */
// Unreadable body — the client got pings only, not the replay.
recopied = false;
}
if (probe == null && fallback.response.status === 429) {
noteResilienceAction({ stored429: true, replayed: false });
} else if (probe != null && recopied) {
noteResilienceAction({ replayed: true });
} else if (probe != null) {
noteResilienceAction({ replayed: false });
}
try {
controller.close();
Expand Down
31 changes: 31 additions & 0 deletions open-sse/executors/opencodeResilienceNotes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import { noteResilienceAction } from "@/lib/usage/resilienceActionsContext.ts";

/**
* resilience effective served-account changes (effective-change counting): the first served dispatch is 0,
* each later dispatch served by a distinct account is +1. Abort and
* shared-egress exclusion arms never reach the call site, so they never
* count. Best-effort notes only.
*/
export function createServedAccountTracker(): (fingerprint: string) => void {
let lastServedFingerprint: string | null = null;
return (fingerprint: string): void => {
if (lastServedFingerprint !== null && lastServedFingerprint !== fingerprint) {
noteResilienceAction({ rotations: 1 });
}
lastServedFingerprint = fingerprint;
};
}

/** resilience park outcome notes (local monotone wait + replayed/stored flags). */
export function noteParkWait(elapsedMs: number): void {
noteResilienceAction({ parked: true, parkMs: Math.max(0, Math.round(elapsedMs)) });
}

export function noteReplayed(): void {
noteResilienceAction({ replayed: true });
}

export function noteStoredFallback(): void {
// The stored 429 fallback is served without a new send.
noteResilienceAction({ replayed: false, stored429: true });
}
18 changes: 17 additions & 1 deletion open-sse/handlers/chatCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,9 @@ import {
readBoundedResponseOutcome,
FLUSH_EMPTY_RETRY_MAX_BYTES,
} from "../utils/emptyTurnRetry.ts";
import { withResilienceActionsContext } from "./chatCore/resilienceAttemptContext.ts";
import { noteBufferedVerdictOutcome } from "./chatCore/emptyTurnResilienceNotes.ts";
import { notePreviousResponseResumed } from "./chatCore/resumedResilienceNotes.ts";
import { assembleStreamingResponseHeaders } from "./chatCore/streamingResponseHeaders.ts";
import { storeStreamingSemanticCacheResponse } from "./chatCore/streamingSemanticCacheStore.ts";
import { captureStreamReasoningForReplay } from "./chatCore/streamReasoningCapture.ts";
Expand Down Expand Up @@ -474,7 +477,13 @@ type VideoBridgeLogParam = { observed: boolean; redaction: VideoBridgeLogRedacti
*/
// extractSystemRoleMessages extracted to chatCore/claudeSystemRole.ts (#3501); re-exported above so
// existing importers (e.g. tests/unit/system-role-extraction.test.ts) keep resolving it from here.
export async function handleChatCore({
export async function handleChatCore(args: Parameters<typeof handleChatCoreInner>[0]) {
// one implicit resilience store per attempt (combo legs each run
// handleChatCore, so each leg gets its own isolated store).
return withResilienceActionsContext([args], (forwarded) => handleChatCoreInner(forwarded));
}

async function handleChatCoreInner({
body,
modelInfo,
credentials,
Expand Down Expand Up @@ -514,7 +523,9 @@ export async function handleChatCore({
videoBridgeLog = undefined,
fallbackAttempts = undefined,
forcedConnectionId = null, // #14116: caller's pinned/requested connection, vs credentials.connectionId below
previousResponseResumed = undefined, // rehydrated-continuation flag from chat.ts; noted below, no semantics.
}) {
delete (body as Record<string, unknown>)._omniroutePreviousResponseResumed;
const {
model: originModel,
resolvedThinkingEffort,
Expand All @@ -540,6 +551,8 @@ export async function handleChatCore({
// (chatCore/memoryExtraction.ts::runMemoryExtractionGate).
const videoBridgeObserved: boolean =
(videoBridgeLog as VideoBridgeLogParam | undefined)?.observed === true;
// resume flag from chat.ts, noted under the attempt store opened above.
notePreviousResponseResumed(previousResponseResumed);
const resilienceSettings = resolveResilienceSettings(cachedSettings);
if (!skipResourcePressureGuard) {
try {
Expand Down Expand Up @@ -5905,19 +5918,22 @@ export async function handleChatCore({
if (verdict.kind === "pass") {
const v = formatBufferedVerdictLog(verdict, correlationId, traceId);
log?.[v.level]?.("FLUSH_EMPTY_RETRY", v.line);
noteBufferedVerdictOutcome(verdict, null, false);
break;
}
if (emptyTurnRetries >= STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX) {
log?.warn?.(
"FLUSH_EMPTY_RETRY",
"retry budget exhausted, falling back to current behavior"
);
noteBufferedVerdictOutcome(verdict, null, true);
break;
}
log?.warn?.(
"FLUSH_EMPTY_RETRY",
`${verdict.reason}, bounded retry through the normal credential path`
);
noteBufferedVerdictOutcome(verdict, { level: "warn", line: verdict.reason }, false);
const nextCreds = await getProviderCredentials(
provider,
null,
Expand Down
6 changes: 6 additions & 0 deletions open-sse/handlers/chatCore/continuationResilienceNotes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
import { noteResilienceAction } from "@/lib/usage/resilienceActionsContext.ts";

/** only a stitched suffix counts as a resume of the turn. */
export function noteContinuedSuffix(): void {
noteResilienceAction({ continued: 1 });
}
26 changes: 26 additions & 0 deletions open-sse/handlers/chatCore/emptyTurnResilienceNotes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
import { judgeBufferedTurn } from "../../utils/emptyTurnRetry.ts";
import { noteResilienceAction } from "@/lib/usage/resilienceActionsContext.ts";

type Verdict = ReturnType<typeof judgeBufferedTurn>;

/**
* record the buffered-turn verdict outcome (pass/retry) and the
* empty-retry count into the attempt resilience store. A pass or an
* exhausted budget closes the loop with a terminal verdict note; each
* scheduled retry adds one empty-retry.
*/
export function noteBufferedVerdictOutcome(
verdict: Verdict,
logLine: { level: "warn" | "debug" | "info"; line: string } | null,
budgetExhausted: boolean
): void {
if (verdict.kind === "pass") {
noteResilienceAction({ buffered: "pass" });
return;
}
if (budgetExhausted) {
noteResilienceAction({ buffered: "retry" });
return;
}
if (logLine) noteResilienceAction({ emptyRetries: 1 });
}
14 changes: 13 additions & 1 deletion open-sse/handlers/chatCore/recoveryTraceLogging.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
* joins the attempt line.
*/
import { STREAM_RECOVERY } from "../../config/constants.ts";
import { noteContinuedSuffix } from "./continuationResilienceNotes.ts";
import type {
ContinuationOutcome,
RecoverableStreamOptions,
Expand Down Expand Up @@ -54,8 +55,13 @@ export function isContinuationGiveUp(event: ContinuationOutcome): boolean {

export function buildContinuationLogHooks(
log: RecoveryLogger,
correlationId?: CorrelationId
correlationId?: CorrelationId,
resilience?: { onSuffix?: (event: ContinuationOutcome) => void }
): Pick<RecoverableStreamOptions, "onContinue" | "onContinueOutcome"> {
// Default resilience note: only a stitched suffix counts as a resume of
// the turn. Terminal/empty/overlap-reject/refused/no-stream outcomes
// record nothing. Kept here so the call site stays one line.
const onSuffix = resilience?.onSuffix ?? noteContinuedSuffix;
// Same empty-string-to-"none" fallback as formatBufferedVerdictLog, without trim.
const cid = correlationId && correlationId.length > 0 ? correlationId : "none";
return {
Expand All @@ -66,6 +72,12 @@ export function buildContinuationLogHooks(
if (isContinuationGiveUp(event)) log?.warn?.(TAG, line);
else if (event.outcome === "refused" && event.reason === "tool-call") log?.debug?.(TAG, line);
else log?.info?.(TAG, line);
// best-effort resilience note; never throws into the stream path.
try {
if (event.outcome === "suffix") onSuffix(event);
} catch {
/* observability only */
}
},
};
}
13 changes: 13 additions & 0 deletions open-sse/handlers/chatCore/resilienceAttemptContext.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { runWithResilienceActionsContext } from "@/lib/usage/resilienceActionsContext.ts";

/**
* one implicit resilience store per attempt. Combo legs each run
* handleChatCore, so each leg gets its own isolated store. The public
* signature stays single-arg (forwarded verbatim); no per-caller change.
*/
export function withResilienceActionsContext<TArgs extends unknown[], TResult>(
args: TArgs,
inner: (...innerArgs: TArgs) => TResult
): TResult {
return runWithResilienceActionsContext(() => inner(...args));
}
15 changes: 15 additions & 0 deletions open-sse/handlers/chatCore/resumedResilienceNotes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
import { noteResilienceAction } from "@/lib/usage/resilienceActionsContext.ts";

/**
* record a rehydrated previous_response_id continuation as a Responses
* resume of this attempt. Called inside handleChatCoreInner, under the
* attempt store — best-effort, never throws.
*/
export function notePreviousResponseResumed(resumed: unknown): void {
if (resumed !== true) return;
try {
noteResilienceAction({ resumed: true });
} catch {
/* observability only */
}
}
8 changes: 7 additions & 1 deletion src/i18n/messages/am.json
Original file line number Diff line number Diff line change
Expand Up @@ -11509,7 +11509,13 @@
"correlationIdValue": "የትስስር መታወቂያ፦ {id}",
"detailedPayloadInfo": "ለአዳዲስ ጥያቄዎች ባለአራት-ደረጃ የደንበኛ/አቅራቢ ውሂብ እይታን ከፈለጉ፣ መጀመሪያ ዝርዝር ምዝገባን ያንቁ።",
"copyAll": "ሁሉንም ቅዳ",
"copiedAll": "ሁሉም ተቀድቷል"
"copiedAll": "ሁሉም ተቀድቷል",
"resilienceStoredError": "ከተቀመጠ ስህተት ቀርቧል",
"resilienceParked": "ከመላኩ በፊት {ms} ቆሟል",
"resilienceRotations": "{count} የመለያ ሽክርክሪቶች",
"resilienceRetried": "ባዶ ዙር እንደገና ተሞክሯል",
"resilienceContinued": "በዥረቱ መሃል ቀጥሏል",
"resilienceRecovered": "በመቋቋም እርምጃ ተመልሷል"
}
},
"proxyLogger": {
Expand Down
Loading
Loading