diff --git a/src/server.mjs b/src/server.mjs index 5cac260..77b174b 100644 --- a/src/server.mjs +++ b/src/server.mjs @@ -521,6 +521,13 @@ function isClientClosedError(error) { return error instanceof ClientClosedError || error?.name === 'ClientClosedError' || error?.code === 'client_closed'; } +const QUEUE_BACKPRESSURE_CODES = new Set(['RUNTIME_QUEUE_TIMEOUT', 'RUNTIME_QUEUE_FULL']); + +/** Queue saturation (local, or relayed by a federated node) is load, not target failure. */ +function isQueueBackpressureError(error) { + return QUEUE_BACKPRESSURE_CODES.has(error?.code) || QUEUE_BACKPRESSURE_CODES.has(error?.upstreamCode); +} + function clientClosedStatus(error) { return isClientClosedError(error) ? 499 : 0; } @@ -554,8 +561,11 @@ export function shouldFailoverModelRequest(error, res = null) { async function upstreamStatusError(upstream) { const text = await readErrorDiagnostic(upstream); let message = text; + let upstreamCode = null; try { - message = JSON.parse(text)?.error?.message ?? text; + const parsed = JSON.parse(text)?.error; + message = parsed?.message ?? text; + upstreamCode = typeof parsed?.code === 'string' ? parsed.code : null; } catch { // Keep the raw upstream response as the diagnostic message. } @@ -568,10 +578,12 @@ async function upstreamStatusError(upstream) { if (Number.isFinite(retryAt)) retryAfterSeconds = Math.max(1, Math.ceil((retryAt - Date.now()) / 1000)); } return Object.assign(new Error(message || `upstream status ${upstream.status}`), { - code: 'upstream_error', + // Relay queue saturation by its own code so every federated hop sees backpressure. + code: QUEUE_BACKPRESSURE_CODES.has(upstreamCode) ? upstreamCode : 'upstream_error', statusCode: upstream.status, upstreamGenerationId: upstream.headers.get('x-generation-id'), upstreamHeadersReceived: true, + ...(upstreamCode == null ? {} : { upstreamCode }), ...(retryAfterSeconds == null ? {} : { retryAfterSeconds }) }); } @@ -2860,7 +2872,8 @@ export function createLloomServer( releaseTargetProbe(resolved); throw error; } - if (shouldFailoverModelRequest(error, res)) noteTargetFailure(resolved, error); + if (isQueueBackpressureError(error)) releaseTargetProbe(resolved); + else if (shouldFailoverModelRequest(error, res)) noteTargetFailure(resolved, error); else if (!isClientClosedError(error)) noteTargetSuccess(resolved); else releaseTargetProbe(resolved); if (!hasNext || !shouldFailoverModelRequest(error, res)) throw error; diff --git a/test/model-failover.test.mjs b/test/model-failover.test.mjs index 4f3b06b..95f4689 100644 --- a/test/model-failover.test.mjs +++ b/test/model-failover.test.mjs @@ -40,6 +40,8 @@ function embedding(model, value) { async function createFixture({ primaryStatus = 503, primaryHeaders = {}, + primaryError = null, + slotError = null, cloudStatus = 200, runtimeStatus = 'running', preserveResident = false, @@ -54,7 +56,7 @@ async function createFixture({ const body = primaryStatus === 200 ? completion('local-upstream', 'local') - : JSON.stringify({ error: { message: `primary status ${primaryStatus}` } }); + : JSON.stringify({ error: primaryError ?? { message: `primary status ${primaryStatus}` } }); res.writeHead(primaryStatus, { 'content-type': 'application/json', ...primaryHeaders }); res.end(body); }); @@ -126,6 +128,15 @@ async function createFixture({ contextWindow: 8192, maxOutputTokens: 1024 }, + { + id: 'primary-direct', + kind: 'chat', + backend: 'primary', + runtime: 'primary-runtime', + upstreamModel: 'local-upstream', + contextWindow: 8192, + maxOutputTokens: 1024 + }, { id: 'cloud-chat', kind: 'chat', @@ -148,6 +159,7 @@ async function createFixture({ }, withSlot: async (runtimeId, fn) => { if (runtimeId) operations.push(`slot:${runtimeId}`); + if (slotError && runtimeId === 'primary-runtime') throw slotError(); return fn(); }, noteRequestOutcome() {}, @@ -621,6 +633,79 @@ async function embed(url, body = {}) { await fixture.close(); } +// Queue saturation is backpressure, not target failure. A federated node's +// queue-timeout 429 may fail over, but must not open the target circuit and +// fast-fail every other caller as "temporarily unavailable". +for (const code of ['RUNTIME_QUEUE_TIMEOUT', 'RUNTIME_QUEUE_FULL']) { + const fixture = await createFixture({ + primaryStatus: 429, + primaryHeaders: { 'retry-after': '30' }, + primaryError: { message: 'runtime primary request queue wait timed out', type: 'runtime_queue_error', code } + }); + try { + assert.equal((await chat(fixture.url)).status, 200); + assert.equal((await chat(fixture.url)).status, 200); + assert.deepEqual(fixture.hits, { primary: 2, cloud: 2, ensures: 2 }); + const routing = await fetch(`${fixture.url}/gateway/routing`).then((response) => response.json()); + assert.equal( + routing.targetBackoffs.find((entry) => entry.model === 'stable-chat'), + undefined + ); + } finally { + await fixture.close(); + } +} + +// A relaying gateway keeps the queue code, so backpressure survives any number +// of federated hops instead of degrading to a generic upstream error. +{ + const fixture = await createFixture({ + primaryStatus: 429, + primaryHeaders: { 'retry-after': '3' }, + primaryError: { message: 'runtime primary request queue wait timed out', code: 'RUNTIME_QUEUE_TIMEOUT' } + }); + try { + const response = await chat(fixture.url, { model: 'primary-direct' }); + assert.equal(response.status, 429); + assert.equal(response.headers.get('retry-after'), '3'); + assert.equal((await response.json()).error.code, 'RUNTIME_QUEUE_TIMEOUT'); + const routing = await fetch(`${fixture.url}/gateway/routing`).then((value) => value.json()); + assert.equal( + routing.targetBackoffs.find((entry) => entry.model === 'primary-direct'), + undefined + ); + } finally { + await fixture.close(); + } +} + +// The same holds for this gateway's own runtime queue. +{ + const fixture = await createFixture({ + primaryStatus: 200, + slotError: () => + Object.assign(new Error('runtime primary-runtime request queue is full; retry after 2 seconds'), { + name: 'RuntimeQueueError', + code: 'RUNTIME_QUEUE_FULL', + type: 'runtime_queue_error', + statusCode: 429, + retryAfterSeconds: 2 + }) + }); + try { + assert.equal((await chat(fixture.url)).status, 200); + assert.equal((await chat(fixture.url)).status, 200); + assert.deepEqual(fixture.hits, { primary: 0, cloud: 2, ensures: 2 }); + const routing = await fetch(`${fixture.url}/gateway/routing`).then((response) => response.json()); + assert.equal( + routing.targetBackoffs.find((entry) => entry.model === 'stable-chat'), + undefined + ); + } finally { + await fixture.close(); + } +} + // A worker/control-plane LLooM remains healthy and administrable but refuses // direct inference, making an accidental client connection fail visibly. {