fix(qoder): address review findings

Correctness:
- testUtils: drop checkExpiry so the userinfo URL probe actually runs (revoked
  tokens used to look "active" until local 30-day expiry passed)
- auth.parseExpiry: handle numeric expiresAt, swap parseInt before Date.parse
  so "2026" doesn't get interpreted as year-2026, treat expires_in:0 as
  already-expired instead of fabricating a 30-day default
- providers.mapTokens: synthesize email from userId when fetchUserInfo fails
  so OAuth dedup works (re-logins no longer accumulate "Account N" rows)

SSE wrapper:
- wrapQoderSSE: add !doneEmitted guard on success branch (chunks could leak
  past [DONE] when an error envelope shared a TCP packet with a valid one)
- flush(): finalize TextDecoder + drain trailing buffer so the chunk carrying
  finish_reason is delivered when upstream closes without a final \n
- sanitize literal \n inside inner OpenAI body so SSE framing stays intact

Robustness:
- executor: wrap buildCosyHeaders in try/catch so a missing accessToken
  returns 401 (re-auth) instead of bubbling as 500
- executor: short-circuit on missing accessToken before signing
- executor: plumb proxyOptions/signal through buildQoderRequestBody so
  proxy-only networks can fetch the model_config catalog
- qoderModels: dedupe concurrent first-time misses with an in-flight Promise
  map (parallel chat windows now do 1 upstream fetch instead of N)
- qoderModels: check signal.aborted before addEventListener so a pre-aborted
  parent signal cancels the inner fetch immediately
- auth: AbortController + 15s timeout on pollDeviceToken / fetchUserInfo to
  prevent hung sockets when openapi.qoder.sh stalls mid-response

UX:
- OAuthModal: derive polling deadline from device-code expires_in (qoder
  publishes 300s; the previous fixed 120s caused timeouts when users took
  more than 2 minutes on the consent page)

Cleanup:
- delete src/lib/oauth/services/qoder.js — referenced removed config fields
  (clientId/clientSecret/tokenUrl/authorizeUrl) and was re-exported from
  services/index.js, so any future caller would TypeError on first use
This commit is contained in:
Simon Shi
2026-05-23 17:29:55 +09:00
committed by decolua
parent a6fd84691b
commit 620b59ca0b
8 changed files with 219 additions and 318 deletions

View File

@@ -123,17 +123,17 @@ function truncate(s, n) {
/**
* Map the OpenAI-style request body into the exact shape Qoder expects.
*/
async function buildQoderRequestBody({ model, body, credentials, log }) {
async function buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal }) {
const qoderKey = String(model || "").replace(/^qoder\//, "");
if (!QODER_MODEL_MAP[qoderKey]) {
throw new Error(`Unsupported qoder model: "${qoderKey}" (received "${model}")`);
}
let modelConfig = await getQoderModelConfig(credentials, qoderKey, { log });
let modelConfig = await getQoderModelConfig(credentials, qoderKey, { log, proxyOptions, signal });
if (!modelConfig) {
// Try a forced refresh once before giving up — the cache may simply
// not be populated yet on first ever call for this credential.
const refreshed = await resolveQoderModels(credentials, { forceRefresh: true, log });
const refreshed = await resolveQoderModels(credentials, { forceRefresh: true, log, proxyOptions, signal });
const retried = refreshed?.rawConfigs.get(qoderKey);
if (!retried) {
throw new Error(
@@ -230,6 +230,54 @@ function wrapQoderSSE(response, model) {
let buffer = "";
let doneEmitted = false;
// Process one already-extracted SSE line (no trailing newline). Returns
// false when the line indicated end-of-stream so the caller can stop
// forwarding any remaining chunks after [DONE].
const processLine = (line, controller) => {
const trimmed = line.replace(/\r$/, "").trim();
if (!trimmed) return;
if (!trimmed.startsWith("data:")) return;
if (doneEmitted) return; // never forward chunks past stream end
const data = trimmed.slice(5).trimStart();
if (data === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
return;
}
let envelope;
try { envelope = JSON.parse(data); } catch { return; }
const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200;
const inner = typeof envelope.body === "string" ? envelope.body : "";
if (statusVal !== 200) {
const msg = inner || `upstream status ${statusVal}`;
const errChunk = JSON.stringify({
id: `qoder-error-${Date.now()}`,
object: "chat.completion.chunk",
created: Math.floor(Date.now() / 1000),
model,
choices: [{ index: 0, delta: { content: `\n[qoder error ${statusVal}: ${truncate(msg, 200)}]` }, finish_reason: "stop" }],
});
controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
return;
}
if (!inner) return;
if (inner === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
return;
}
// Inner is an OpenAI-shaped chunk. Strip any embedded newlines so the
// SSE frame stays a single event (a literal "\n" inside `inner` would
// otherwise split the frame across multiple data: lines and downstream
// parsers would reassemble them as separate events).
const sanitized = inner.replace(/\r?\n/g, "");
controller.enqueue(encoder.encode(`data: ${sanitized}\n\n`));
};
const transform = new TransformStream({
transform(chunk, controller) {
buffer += decoder.decode(chunk, { stream: true });
@@ -237,54 +285,24 @@ function wrapQoderSSE(response, model) {
while ((nl = buffer.indexOf("\n")) !== -1) {
const line = buffer.slice(0, nl);
buffer = buffer.slice(nl + 1);
const trimmed = line.replace(/\r$/, "").trim();
if (!trimmed) continue;
if (!trimmed.startsWith("data:")) continue;
let data = trimmed.slice(5).trimStart();
if (data === "[DONE]") {
if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
}
continue;
}
let envelope;
try { envelope = JSON.parse(data); } catch { continue; }
const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200;
const inner = typeof envelope.body === "string" ? envelope.body : "";
if (statusVal !== 200) {
const msg = inner || `upstream status ${statusVal}`;
const errChunk = JSON.stringify({
id: `qoder-error-${Date.now()}`,
object: "chat.completion.chunk",
created: Math.floor(Date.now() / 1000),
model,
choices: [{ index: 0, delta: { content: `\n[qoder error ${statusVal}: ${truncate(msg, 200)}]` }, finish_reason: "stop" }],
});
controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`));
if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
}
continue;
}
if (!inner) continue;
if (inner === "[DONE]") {
if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
}
continue;
}
// Inner is already an OpenAI-shaped chunk; forward as-is.
controller.enqueue(encoder.encode(`data: ${inner}\n\n`));
processLine(line, controller);
}
},
flush(controller) {
// Finalize the decoder so any pending multi-byte sequence is
// released into `buffer` instead of being silently dropped.
buffer += decoder.decode();
// Drain any trailing line that arrived without a terminating newline
// (e.g. upstream closed the socket immediately after the last write,
// or a CDN stripped the final CRLF). Without this, the chunk that
// carries finish_reason is silently lost.
if (buffer.length > 0) {
processLine(buffer, controller);
buffer = "";
}
if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
doneEmitted = true;
}
},
});
@@ -329,11 +347,20 @@ export class QoderExecutor extends BaseExecutor {
);
return { response: fakeResp, url, headers: {}, transformedBody: body };
}
if (!credentials?.accessToken) {
// Same shape as the userId guard — clean 401 so chatCore reports
// "reconnect" rather than bubbling cosy.js's synchronous throw as 500.
const fakeResp = new Response(
JSON.stringify({ error: { message: "qoder credential is missing accessToken; reconnect the account" } }),
{ status: 401, headers: { "Content-Type": "application/json" } },
);
return { response: fakeResp, url, headers: {}, transformedBody: body };
}
let qoderKey;
let payload;
try {
({ qoderKey, payload } = await buildQoderRequestBody({ model, body, credentials, log }));
({ qoderKey, payload } = await buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal }));
} catch (err) {
const fakeResp = new Response(
JSON.stringify({ error: { message: err.message } }),
@@ -346,17 +373,28 @@ export class QoderExecutor extends BaseExecutor {
const encodedBodyStr = qoderEncodeBody(plainBody);
const encodedBodyBuf = Buffer.from(encodedBodyStr, "latin1");
const cosyHeaders = buildCosyHeaders(
encodedBodyBuf,
url,
{
userId: psd.userId,
authToken: credentials.accessToken,
name: credentials.displayName || "",
email: credentials.email || "",
machineId: psd.machineId || "",
},
);
let cosyHeaders;
try {
cosyHeaders = buildCosyHeaders(
encodedBodyBuf,
url,
{
userId: psd.userId,
authToken: credentials.accessToken,
name: credentials.displayName || "",
email: credentials.email || "",
machineId: psd.machineId || "",
},
);
} catch (err) {
// cosy.js throws synchronously on missing userId/authToken — surface
// as 401 so chatCore prompts re-auth instead of returning a 500.
const fakeResp = new Response(
JSON.stringify({ error: { message: `qoder cosy signing failed: ${err.message}` } }),
{ status: 401, headers: { "Content-Type": "application/json" } },
);
return { response: fakeResp, url, headers: {}, transformedBody: body };
}
const modelSource = (payload.model_config && payload.model_config.source) || "system";
const headers = {

View File

@@ -26,6 +26,14 @@ const CACHE_TTL_MS = 60 * 60 * 1000; // 1h, same as the Kiro catalog
/** @type {Map<string, { expiresAt: number, models: any[], rawConfigs: Map<string, object>, fetched: boolean }>} */
const catalogCache = new Map();
/**
* In-flight fetch promises keyed by cacheKey. Concurrent first-time
* callers (parallel chat windows) all observe the same Promise so we
* fan-out exactly one upstream request per credential per miss.
* @type {Map<string, Promise<{ expiresAt: number, models: any[], rawConfigs: Map<string, object>, fetched: boolean } | null>>}
*/
const inflight = new Map();
/**
* Stable cache key per credential (so different login sessions for the same
* account share an entry).
@@ -73,8 +81,15 @@ async function fetchQoderCatalogRaw(credentials, signal, proxyOptions = null) {
try {
timer = setTimeout(() => controller.abort("timeout"), FETCH_TIMEOUT_MS);
if (signal && typeof signal.addEventListener === "function") {
abortListener = () => controller.abort(signal.reason);
signal.addEventListener("abort", abortListener);
// If the parent signal already aborted before we got here, the
// 'abort' event has already fired and addEventListener won't
// re-trigger it. Propagate the cancellation immediately.
if (signal.aborted) {
controller.abort(signal.reason);
} else {
abortListener = () => controller.abort(signal.reason);
signal.addEventListener("abort", abortListener);
}
}
response = await proxyAwareFetch(
QODER_MODEL_LIST_URL,
@@ -137,7 +152,9 @@ export async function getQoderModelConfig(credentials, modelKey, options = {}) {
/**
* Resolve the live model catalog + raw configs for a credential. Caches
* results for CACHE_TTL_MS so repeated chat requests don't re-fetch.
* results for CACHE_TTL_MS so repeated chat requests don't re-fetch, and
* deduplicates concurrent misses so parallel chat windows fan-out exactly
* one upstream request per credential.
*/
export async function resolveQoderModels(credentials, options = {}) {
if (!credentials?.accessToken) return null;
@@ -153,17 +170,36 @@ export async function resolveQoderModels(credentials, options = {}) {
}
}
const fetched = await fetchQoderCatalogRaw(credentials, options.signal, options.proxyOptions);
if (!fetched) return null;
// Coalesce concurrent misses on the same credential into one upstream call.
// forceRefresh callers still get their own fetch (they wanted fresh data).
const existing = inflight.get(key);
if (existing && !options.forceRefresh) {
return existing;
}
const entry = {
expiresAt: now + CACHE_TTL_MS,
models: fetched.models,
rawConfigs: fetched.rawConfigs,
fetched: true,
};
catalogCache.set(key, entry);
return entry;
const fetchPromise = (async () => {
const fetched = await fetchQoderCatalogRaw(credentials, options.signal, options.proxyOptions);
if (!fetched) return null;
const entry = {
expiresAt: Date.now() + CACHE_TTL_MS,
models: fetched.models,
rawConfigs: fetched.rawConfigs,
fetched: true,
};
catalogCache.set(key, entry);
return entry;
})();
inflight.set(key, fetchPromise);
try {
return await fetchPromise;
} finally {
// Clear only if this is still the in-flight entry — a forceRefresh
// call that started later may have replaced it.
if (inflight.get(key) === fetchPromise) {
inflight.delete(key);
}
}
}
export function invalidateQoderCatalog(credentials) {