diff --git a/packages/kilo-gateway/src/server/routes.ts b/packages/kilo-gateway/src/server/routes.ts index 13714e75ce6..f09b945a280 100644 --- a/packages/kilo-gateway/src/server/routes.ts +++ b/packages/kilo-gateway/src/server/routes.ts @@ -37,55 +37,6 @@ interface KiloRoutesDeps extends ImportDeps { const FIM_TIMEOUT_MS = 30_000 -function timeout(parent: AbortSignal | undefined, ms: number) { - const controller = new AbortController() - const timer = setTimeout(() => controller.abort(new DOMException("FIM request timed out", "TimeoutError")), ms) - const abort = () => controller.abort(parent?.reason) - - if (parent?.aborted) { - abort() - } else { - parent?.addEventListener("abort", abort, { once: true }) - } - - return { - signal: controller.signal, - cleanup() { - clearTimeout(timer) - parent?.removeEventListener("abort", abort) - }, - } -} - -function stream(body: ReadableStream | null, cleanup: () => void) { - if (!body) { - cleanup() - return null - } - - const reader = body.getReader() - return new ReadableStream({ - async pull(controller) { - try { - const chunk = await reader.read() - if (chunk.done) { - cleanup() - controller.close() - return - } - controller.enqueue(chunk.value) - } catch (err) { - cleanup() - controller.error(err) - } - }, - cancel(reason) { - cleanup() - return reader.cancel(reason) - }, - }) -} - /** * Create Kilo Gateway routes with OpenCode dependencies injected * @@ -393,50 +344,36 @@ export function createKiloRoutes(deps: KiloRoutesDeps) { [HEADER_FEATURE]: "autocomplete", } - const guard = timeout(c.req.raw.signal, FIM_TIMEOUT_MS) - const result = await fetch(endpoint, { - method: "POST", - headers, - signal: guard.signal, - body: JSON.stringify({ - model: fimModel, - prompt: prefix, - suffix, - max_tokens: fimMaxTokens, - temperature: fimTemperature, - stream: true, - }), - }) - .then((response) => ({ response })) - .catch((err) => ({ err })) + const signal = AbortSignal.any([c.req.raw.signal, AbortSignal.timeout(FIM_TIMEOUT_MS)]) - if ("err" in result) { - guard.cleanup() - if (!guard.signal.aborted) throw result.err - const reason = guard.signal.reason - const timed = reason instanceof DOMException && reason.name === "TimeoutError" - const status = timed ? 504 : 499 - return c.json({ error: timed ? "FIM request timed out" : "FIM request canceled" }, status as any) + let response: Response + try { + response = await fetch(endpoint, { + method: "POST", + headers, + signal, + body: JSON.stringify({ + model: fimModel, + prompt: prefix, + suffix, + max_tokens: fimMaxTokens, + temperature: fimTemperature, + stream: true, + }), + }) + } catch (err) { + if (err instanceof DOMException && err.name === "TimeoutError") + return c.json({ error: "FIM request timed out" }, 504 as any) + if (signal.aborted) return c.json({ error: "FIM request canceled" }, 499 as any) + throw err } - const response = result.response if (!response.ok) { - const error = await response - .text() - .then((text) => ({ text })) - .catch((err) => ({ err })) - .finally(() => guard.cleanup()) - if ("err" in error) { - if (!guard.signal.aborted) throw error.err - const reason = guard.signal.reason - const timed = reason instanceof DOMException && reason.name === "TimeoutError" - const status = timed ? 504 : 499 - return c.json({ error: timed ? "FIM request timed out" : "FIM request canceled" }, status as any) - } - return c.json({ error: `FIM request failed: ${response.status} ${error.text}` }, response.status as any) + const text = await response.text() + return c.json({ error: `FIM request failed: ${response.status} ${text}` }, response.status as any) } - return new Response(stream(response.body, guard.cleanup), { + return new Response(response.body, { headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache",