From 5f19b3868cf5694ea49d291384a70f92076322cf Mon Sep 17 00:00:00 2001 From: Kieran Klukas Date: Sun, 2 Aug 2026 22:43:53 -0400 Subject: [PATCH] Replace Elysia with framework-free Bun routes and add the web UI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The server now runs on Bun's native per-method routes with Standard Schema (valibot) body validation, dropping the Elysia dependency. Tests exercise a real ephemeral server instead of app.handle. Adds the vanilla-JS chat client: SSE streaming, optimistic send with retry, and the model curation settings page. Also folds in the review cleanup: the server and worker drive loops are unified in a shared JobDriver, dead store/event/sse code is deleted, idle actor eviction no longer orphans live streams, Last-Event-ID validates its conversation scope, and conversations list by most recent activity rather than creation time. πŸ’˜ Generated with Crush Assisted-by: Crush:qwen3.8-max-preview --- bun.lock | 40 +--- package.json | 4 +- scripts/smoke.ts | 12 +- server.ts | 431 +++--------------------------------- src/actor.ts | 6 +- src/client/app.css | 207 ++++++++++++++++++ src/client/app.js | 460 +++++++++++++++++++++++++++++++++++++++ src/client/index.html | 63 ++++++ src/client/settings.html | 20 ++ src/client/settings.js | 125 +++++++++++ src/drive.ts | 82 +++++++ src/events.ts | 14 +- src/http.ts | 301 +++++++++++++++++++++++++ src/schemas.ts | 36 +++ src/sse.ts | 13 -- src/store.ts | 161 ++++++++------ src/validate.ts | 58 +++++ tests/curation.test.ts | 124 ++++------- tests/server.test.ts | 274 +++++++++++++++-------- tests/validation.test.ts | 69 ++++++ worker.ts | 139 ++---------- 21 files changed, 1817 insertions(+), 822 deletions(-) create mode 100644 src/client/app.css create mode 100644 src/client/app.js create mode 100644 src/client/index.html create mode 100644 src/client/settings.html create mode 100644 src/client/settings.js create mode 100644 src/drive.ts create mode 100644 src/http.ts create mode 100644 src/schemas.ts create mode 100644 src/validate.ts create mode 100644 tests/validation.test.ts diff --git a/bun.lock b/bun.lock index 5839294..fc00a32 100644 --- a/bun.lock +++ b/bun.lock @@ -8,8 +8,10 @@ "@ai-sdk/anthropic": "^4.0.27", "@ai-sdk/openai": "^4.0.27", "@ai-sdk/openai-compatible": "^3.0.20", + "@standard-schema/spec": "^1.1.0", "ai": "^7.0.48", - "elysia": "^1.4.29", + "streaming-markdown": "^0.2.15", + "valibot": "^1.4.2", }, "devDependencies": { "@types/bun": "^1.2.2", @@ -29,16 +31,8 @@ "@ai-sdk/provider-utils": ["@ai-sdk/provider-utils@5.0.18", "", { "dependencies": { "@ai-sdk/provider": "4.0.4", "@standard-schema/spec": "^1.1.0", "@workflow/serde": "4.1.0", "eventsource-parser": "^3.0.8", "undici": "^7.28.0" }, "peerDependencies": { "zod": "^3.25.76 || ^4.1.8" } }, "sha512-UBNCrkxS5llgN2/RXLBRkjCFTIuL6YB1Goq5c+yRecnznhtX4XPjTwLHz2hrsTltTVZ0rqVegtnwFXdPBkjDHA=="], - "@borewit/text-codec": ["@borewit/text-codec@0.2.2", "", {}, "sha512-DDaRehssg1aNrH4+2hnj1B7vnUGEjU6OIlyRdkMd0aUdIUvKXrJfXsy8LVtXAy7DRvYVluWbMspsRhz2lcW0mQ=="], - - "@sinclair/typebox": ["@sinclair/typebox@0.34.52", "", {}, "sha512-XiMQh7qqVlxZzcVD+kkGMNGMzcTrDMLWI7S4x7z1MkCkbDPrekpZXEUK0eZqZFMuHQg2a2DZOcDIh9o5v3Gonw=="], - "@standard-schema/spec": ["@standard-schema/spec@1.1.0", "", {}, "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w=="], - "@tokenizer/inflate": ["@tokenizer/inflate@0.4.1", "", { "dependencies": { "debug": "^4.4.3", "token-types": "^6.1.1" } }, "sha512-2mAv+8pkG6GIZiF1kNg1jAjh27IDxEPKwdGul3snfztFerfPGI1LjDezZp3i7BElXompqEtPmoPx6c2wgtWsOA=="], - - "@tokenizer/token": ["@tokenizer/token@0.3.0", "", {}, "sha512-OvjF+z51L3ov0OyAU0duzsYuvO01PH7x4t6DJx+guahgTnBHkhJdG7soQeTSFLWN3efnHyibZ4Z8l2EuWwJN3A=="], - "@types/bun": ["@types/bun@1.3.14", "", { "dependencies": { "bun-types": "1.3.14" } }, "sha512-h1hFqFVcvAvD9j9K7ZW7vd82aSA+rTdznZa+5bwvCwqSB1jmmfLcbIWhOLx1/+boy/xmjgCs/OMUL8hRJSmnPw=="], "@types/node": ["@types/node@26.1.2", "", { "dependencies": { "undici-types": "~8.3.0" } }, "sha512-Vu4a5UFA9rIIFJ7rB/Vaafh9lrCQszopTCx6KjFboXTGQbPNasehVR5TEiithSDGyd1DEiUByggTZsg8jukeIg=="], @@ -51,40 +45,18 @@ "bun-types": ["bun-types@1.3.14", "", { "dependencies": { "@types/node": "*" } }, "sha512-4N0ig0fEomHt5R0KCFWjovxow98rIoRwKolrYdCcknNwMekCXRnWEUvgu5soYV8QXtVsrUD8B95MBOZGPvr6KQ=="], - "cookie": ["cookie@1.1.1", "", {}, "sha512-ei8Aos7ja0weRpFzJnEA9UHJ/7XQmqglbRwnf2ATjcB9Wq874VKH9kfjjirM6UhU2/E5fFYadylyhFldcqSidQ=="], - - "debug": ["debug@4.4.3", "", { "dependencies": { "ms": "^2.1.3" }, "peerDependencies": { "supports-color": "*" }, "optionalPeers": ["supports-color"] }, "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA=="], - - "elysia": ["elysia@1.4.29", "", { "dependencies": { "cookie": "^1.1.1", "exact-mirror": "^0.2.7", "fast-decode-uri-component": "^1.0.1", "memoirist": "^0.4.0" }, "peerDependencies": { "@sinclair/typebox": ">= 0.34.0 < 1", "@types/bun": ">= 1.2.0", "file-type": ">= 20.0.0", "openapi-types": ">= 12.0.0", "typescript": ">= 5.0.0" }, "optionalPeers": ["@types/bun", "typescript"] }, "sha512-GwMRGGwSdjfPt+w3LA0fqTuYJtS8uVRJicvoar98/HrO5qdFKDc9CwjIb6Kja+v39lkY+58hr2JvdR9jQzlUuA=="], - "eventsource-parser": ["eventsource-parser@3.1.0", "", {}, "sha512-kJezFj9YFAMLeORyi7aCLxLbD5/qWMQnoMVlVPyHIll7lgRJCc3JVln9Vgl9nwQi0YkMnhdGTMNn7CkRRAptMg=="], - "exact-mirror": ["exact-mirror@0.2.7", "", { "peerDependencies": { "@sinclair/typebox": "^0.34.15" }, "optionalPeers": ["@sinclair/typebox"] }, "sha512-+MeEmDcLA4o/vjK2zujgk+1VTxPR4hdp23qLqkWfStbECtAq9gmsvQa3LW6z/0GXZyHJobrCnmy1cdeE7BjsYg=="], - - "fast-decode-uri-component": ["fast-decode-uri-component@1.0.1", "", {}, "sha512-WKgKWg5eUxvRZGwW8FvfbaH7AXSh2cL+3j5fMGzUMCxWBJ3dV3a7Wz8y2f/uQ0e3B6WmodD3oS54jTQ9HVTIIg=="], - - "file-type": ["file-type@22.0.1", "", { "dependencies": { "@tokenizer/inflate": "^0.4.1", "strtok3": "^10.3.5", "token-types": "^6.1.2", "uint8array-extras": "^1.5.0" } }, "sha512-ww5Mhre0EE+jmBvOXTmXAbEMuZE7uX4a3+oRCQFNj8w++g3ev913N6tXQz0XTXbueQ5TWQfm6BdaViEHHn8bhA=="], - - "ieee754": ["ieee754@1.2.1", "", {}, "sha512-dcyqhDvX1C46lXZcVqCpK+FtMRQVdIMN6/Df5js2zouUsqG7I6sFxitIC+7KYK29KdXOLHdu9zL4sFnoVQnqaA=="], - "json-schema": ["json-schema@0.4.0", "", {}, "sha512-es94M3nTIfsEPisRafak+HDLfHXnKBhV3vU5eqPcS3flIWqcxJWgXHXiey3YrpaNsanY5ei1VoYEbOzijuq9BA=="], - "memoirist": ["memoirist@0.4.0", "", {}, "sha512-zxTgA0mSYELa66DimuNQDvyLq36AwDlTuVRbnQtB+VuTcKWm5Qc4z3WkSpgsFWHNhexqkIooqpv4hdcqrX5Nmg=="], - - "ms": ["ms@2.1.3", "", {}, "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA=="], - - "openapi-types": ["openapi-types@12.1.3", "", {}, "sha512-N4YtSYJqghVu4iek2ZUvcN/0aqH1kRDuNqzcycDxhOUpg7GdvLa2F3DgS6yBNhInhv2r/6I0Flkn7CqL8+nIcw=="], - - "strtok3": ["strtok3@10.3.5", "", { "dependencies": { "@tokenizer/token": "^0.3.0" } }, "sha512-ki4hZQfh5rX0QDLLkOCj+h+CVNkqmp/CMf8v8kZpkNVK6jGQooMytqzLZYUVYIZcFZ6yDB70EfD8POcFXiF5oA=="], - - "token-types": ["token-types@6.1.2", "", { "dependencies": { "@borewit/text-codec": "^0.2.1", "@tokenizer/token": "^0.3.0", "ieee754": "^1.2.1" } }, "sha512-dRXchy+C0IgK8WPC6xvCHFRIWYUbqqdEIKPaKo/AcTUNzwLTK6AH7RjdLWsEZcAN/TBdtfUw3PYEgPr5VPr6ww=="], - - "uint8array-extras": ["uint8array-extras@1.5.0", "", {}, "sha512-rvKSBiC5zqCCiDZ9kAOszZcDvdAHwwIKJG33Ykj43OKcWsnmcBRL09YTU4nOeHZ8Y2a7l1MgTd08SBe9A8Qj6A=="], + "streaming-markdown": ["streaming-markdown@0.2.15", "", {}, "sha512-eSWjdvLSTznl2+FIzuxw1TRXwhjZiY1RfSygGCMpcNY8CrsaUM2zHtxe/JFcPoXCQxeO3TSJ72QMDIN5B/Ui1w=="], "undici": ["undici@7.29.0", "", {}, "sha512-IDxfleLmmbSskfWSUATiN1nfn2rDuvnMOqb5CWR92iIfojA0Ud+ulOAAEQ57LPr9rWmsreUyf5lwyao+7GNNVw=="], "undici-types": ["undici-types@8.3.0", "", {}, "sha512-j375ScV60dom+YkPFIfTLcOiPxkN/buHz5GobjLhixFuANaNs3C9l4GmrWqejgXWJ7BbJcFYpTEUkS1Ge8bpZQ=="], + "valibot": ["valibot@1.4.2", "", { "peerDependencies": { "typescript": ">=5" }, "optionalPeers": ["typescript"] }, "sha512-gjdCvJ6d3RyHAneqxMYMW9QMCwYMb3jpOO0IyHZV1bnRHFBHrX3VkIILt5XYR0WhwHiH7Mty8ovuPZ/O3gamrg=="], + "zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], } } diff --git a/package.json b/package.json index 9b203c8..86e66f4 100644 --- a/package.json +++ b/package.json @@ -12,8 +12,10 @@ "@ai-sdk/anthropic": "^4.0.27", "@ai-sdk/openai": "^4.0.27", "@ai-sdk/openai-compatible": "^3.0.20", + "@standard-schema/spec": "^1.1.0", "ai": "^7.0.48", - "elysia": "^1.4.29" + "streaming-markdown": "^0.2.15", + "valibot": "^1.4.2" }, "devDependencies": { "@types/bun": "^1.2.2" diff --git a/scripts/smoke.ts b/scripts/smoke.ts index 86a686d..5ad4045 100644 --- a/scripts/smoke.ts +++ b/scripts/smoke.ts @@ -88,7 +88,7 @@ async function readUntil( } // 1. POST a prompt -const res = await fetch(`${base}/conversations/${conv}/prompt`, { +const res = await fetch(`${base}/api/conversations/${conv}/prompt`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ content: "hello", model: "echo" }), @@ -96,7 +96,7 @@ const res = await fetch(`${base}/conversations/${conv}/prompt`, { console.log("prompt status:", res.status, await res.text()); // 2. Open SSE and read until message-end -const es = await fetch(`${base}/conversations/${conv}/stream`); +const es = await fetch(`${base}/api/conversations/${conv}/stream`); const frames = await readUntil(es, (f) => f.some((x) => x.event === "message-end")); console.log("stream events:", frames.map((f) => f.event).join(", ")); const delta = frames.find((f) => f.event === "text-delta")!.data.delta; @@ -109,7 +109,7 @@ console.log("last id:", lastId); await Bun.sleep(100); // let the previous connection fully close let resumed: Frame[] = []; try { - const es2 = await fetch(`${base}/conversations/${conv}/stream`, { + const es2 = await fetch(`${base}/api/conversations/${conv}/stream`, { headers: { "last-event-id": lastId }, }); resumed = await readUntil(es2, () => false, 500); @@ -120,18 +120,18 @@ try { console.log("resume-at-tail frames:", resumed.length, "(expect 0)"); // 4. Cancel path: prompt a second run then cancel before it finishes. -const res2 = await fetch(`${base}/conversations/${conv}/prompt`, { +const res2 = await fetch(`${base}/api/conversations/${conv}/prompt`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ content: "cancel me", model: "echo" }), }); console.log("prompt2 status:", res2.status); await Bun.sleep(1200); // let the drive loop claim + start the run -const res3 = await fetch(`${base}/conversations/${conv}/cancel`, { method: "POST" }); +const res3 = await fetch(`${base}/api/conversations/${conv}/cancel`, { method: "POST" }); console.log("cancel status:", res3.status, await res3.text()); // Re-read stream to see cancellation -const es3 = await fetch(`${base}/conversations/${conv}/stream`); +const es3 = await fetch(`${base}/api/conversations/${conv}/stream`); const frames3 = await readUntil(es3, (f) => f.some((x) => x.event === "cancelled")); console.log("post-cancel stream:", frames3.map((f) => f.event).join(", ")); diff --git a/server.ts b/server.ts index bd92fb8..9cc7770 100644 --- a/server.ts +++ b/server.ts @@ -1,415 +1,54 @@ -import { Elysia, t } from "elysia"; -import { randomUUID } from "node:crypto"; +import indexHTML from "./src/client/index.html"; +import settingsHTML from "./src/client/settings.html"; import { Store } from "./src/store"; -import { ConversationActor, type Subscriber, type WireEvent } from "./src/actor"; -import { run, initInference, getRegistry } from "./src/inference"; -import { sseBlock } from "./src/sse"; -import { parseEventId } from "./src/events"; -import { - ACTOR_IDLE_TTL_MS, - LEASE_GRACE_MS, - REAP_INTERVAL_MS, - SUBSCRIBER_HEARTBEAT_MS, -} from "./src/config"; +import { initInference } from "./src/inference"; +import { apiRoutes, getActor, evictIdleActors } from "./src/http"; +import { JobDriver } from "./src/drive"; +import { REAP_INTERVAL_MS } from "./src/config"; /** - * Parses the Last-Event-ID header. The browser sends back the full event id - * (`:`); we extract the seq. Returns 0 (replay all) on - * any parse failure. + * Web entrypoint. Bun's native `routes` serve the HTML pages (transpiled, + * bundled, content-hashed, HMR in dev β€” no bundler config, no Vite) alongside + * the framework-free API routes from src/http. Tests import `apiRoutes` directly + * and never touch this file, so importing it never triggers frontend bundling. + * + * The inline drive loop shares its implementation with worker.ts via JobDriver; + * both entrypoints run the same claim/run/checkpoint logic. */ -function readLastEventId(req: Request): number { - const raw = req.headers.get("last-event-id"); - if (raw === null) return 0; - try { - return parseEventId(raw).seq; - } catch { - return 0; - } -} - -/** - * Subscriber backed by a ReadableStream: Bun's stream layer pulls strictly - * serially, so a double-delivery is structurally impossible (unlike an async - * iterator whose concurrent pulls race on a shared waiter). The actor `push`es - * each event straight into the controller; backpressure is the platform's job. - */ -class MySubscriber implements Subscriber { - closed = false; - readonly stream: ReadableStream; - private controller!: ReadableStreamDefaultController; - private readonly onClose: () => void; - - constructor(onClose: () => void) { - this.onClose = onClose; - this.stream = new ReadableStream({ - start: (controller) => { - this.controller = controller; - }, - cancel: () => { - this.close(); - }, - }); - } - - push(e: WireEvent): void { - if (this.closed) return; - this.controller.enqueue(e); - } - - close(): void { - if (this.closed) return; - this.closed = true; - this.onClose(); - try { - this.controller.close(); - } catch { - // already closed by the consumer - } - } -} - -/** - * Per-conversation actors with idle eviction. Actors with no subscribers and - * no recent activity are removed so the map doesn't grow without bound. - */ -const actors = new Map(); -export function getActor(id: string, store: Store): ConversationActor { - let a = actors.get(id); - if (!a) { - a = new ConversationActor(id, store); - actors.set(id, a); - } - return a; -} - -function evictIdleActors(): void { - const cutoff = Date.now() - ACTOR_IDLE_TTL_MS; - for (const [id, actor] of actors) { - if (actor.lastActivity < cutoff) { - actors.delete(id); - } - } -} - -/** Refs of every model this deployment can run (enabled providers + echo). */ -function knownModelRefs(): Set { - return new Set(getRegistry().listModels().map((m) => m.ref)); -} - -/** - * Every available model joined to its curation state, for the settings UI. - * Models with no curation row default to hidden with their catalog name. - */ -function adminModels(store: Store) { - const settings = new Map(store.listModelSettings().map((s) => [s.ref, s])); - return getRegistry() - .listModels() - .map((m) => { - const s = settings.get(m.ref); - return { - ...m, - visible: s?.visible ?? false, - displayName: s?.displayName ?? null, - sortOrder: s?.sortOrder ?? 0, - }; - }); -} - -/** - * The curated subset shown in the chat picker: opt-in (visible only), with - * displayName applied and ordered by sortOrder then name. - */ -function chatModels(store: Store) { - const settings = new Map(store.listModelSettings().map((s) => [s.ref, s])); - return getRegistry() - .listModels() - .filter((m) => settings.get(m.ref)?.visible) - .map((m) => { - const s = settings.get(m.ref)!; - return { - ref: m.ref, - name: s.displayName ?? m.name, - contextWindow: m.contextWindow, - reasoningLevels: m.reasoningLevels, - supportsImages: m.supportsImages, - sortOrder: s.sortOrder, - }; - }) - .sort((a, b) => a.sortOrder - b.sortOrder || a.name.localeCompare(b.name)); -} - -export function buildApp(deps: { - store: Store; -}): Elysia { - const { store } = deps; - const app = new Elysia(); - - app.get("/health", () => ({ ok: true })); - - // Settings/admin view: every model available to this deployment (enabled - // providers Γ— catalog) with its catalog metadata joined to its curation - // state (visible / displayName / sortOrder). - app.get("/models", () => ({ models: adminModels(store) })); - - // Chat view: the curated, opt-in subset β€” only models explicitly marked - // visible, with pretty names applied, ordered for the picker. - app.get("/models/chat", () => ({ models: chatModels(store) })); - - // Curation mutation for the settings menu. Partial: omitted fields keep - // their current value; `displayName: null` clears an override. - app.patch( - "/models", - ({ body, set }) => { - if (!knownModelRefs().has(body.ref)) { - set.status = 422; - return { error: `unknown model "${body.ref}"` }; - } - const prev = store.getModelSetting(body.ref); - const merged = { - ref: body.ref, - visible: body.visible ?? prev?.visible ?? false, - displayName: - body.displayName !== undefined - ? body.displayName - : (prev?.displayName ?? null), - sortOrder: body.sortOrder ?? prev?.sortOrder ?? 0, - }; - store.setModelSetting(merged); - return merged; - }, - { - body: t.Object({ - ref: t.String(), - visible: t.Optional(t.Boolean()), - displayName: t.Optional(t.Union([t.String(), t.Null()])), - sortOrder: t.Optional(t.Number()), - }), - }, - ); - - app.get( - "/conversations/:id/stream", - ({ params, request }) => { - // TODO(auth): cookie/session auth so native EventSource works unmodified. - const conversationId = params.id; - const actor = getActor(conversationId, store); - const after = readLastEventId(request); - - let unsub: () => void = () => {}; - const sub = new MySubscriber(() => unsub()); - unsub = actor.subscribe(sub, after); - - // NOTE: this stream is intentionally long-lived. It does not close when - // a run ends; the next run's events arrive on the same connection. A - // dropped connection is what triggers replay via Last-Event-ID. - // - // Keepalive: emit SSE comments periodically so proxies and load - // balancers don't close the idle connection. - const KEEPALIVE = Symbol("keepalive"); - const sseTransform = new TransformStream({ - transform: (e, controller) => { - if (e === KEEPALIVE) { - controller.enqueue(": keepalive\n\n"); - } else { - controller.enqueue(sseBlock(e)); - } - }, - }); - const body = sub.stream.pipeThrough(sseTransform); - - const keepAlive = setInterval(() => { - if (sub.closed) { - clearInterval(keepAlive); - return; - } - sub.push(KEEPALIVE as unknown as WireEvent); - }, SUBSCRIBER_HEARTBEAT_MS); - - return new Response(body, { - headers: { - "Content-Type": "text/event-stream", - "Cache-Control": "no-cache, no-transform", - Connection: "keep-alive", - "X-Accel-Buffering": "no", - }, - }); - }, - { - params: t.Object({ id: t.String() }), - }, - ); - - app.post( - "/conversations/:id/prompt", - async ({ params, body, set }) => { - const conversationId = params.id; - - // Reject unknown models up front: otherwise the job is accepted (202) - // and fails silently in the worker at resolve time. - if (!knownModelRefs().has(body.model)) { - set.status = 422; - return { error: `unknown model "${body.model}"` }; - } - - const actor = getActor(conversationId, store); - const runId = body.runId ?? randomUUID(); - const messageId = randomUUID(); - - // Use the SAME runId for the user message and the run so clients can - // correlate them. - actor.appendUser(body.content, runId); - - const jobId = `${conversationId}:${randomUUID()}`; - store.enqueue(jobId, conversationId, { - conversationId, - runId, - messageId, - prompt: body.content, - model: body.model, - }); - - // Enqueue decouples the response from the run: the client opens /stream - // separately and receives message-start/text-deltas/message-end as the - // worker (possibly in another process) executes the job. - set.status = 202; - return { jobId, runId, messageId }; - }, - { - body: t.Object({ - content: t.String(), - runId: t.Optional(t.String()), - model: t.String(), - }), - }, - ); - - app.post( - "/conversations/:id/cancel", - ({ params }) => { - const actor = getActor(params.id, store); - actor.requestCancel(); - return { ok: true }; - }, - { - params: t.Object({ id: t.String() }), - }, - ); - - app.post( - "/conversations/:id/steer", - async ({ params, body, set }) => { - const conversationId = params.id; - - if (!knownModelRefs().has(body.model)) { - set.status = 422; - return { error: `unknown model "${body.model}"` }; - } - - const actor = getActor(conversationId, store); - // Hard steer: abort the current run, then queue a new one with the - // redirect message. - actor.requestCancel(); - - const runId = randomUUID(); - const messageId = randomUUID(); - actor.appendUser(body.content, runId); - - const jobId = `${conversationId}:${randomUUID()}`; - store.enqueue(jobId, conversationId, { - conversationId, - runId, - messageId, - prompt: body.content, - model: body.model, - }); - - set.status = 202; - return { ok: true, jobId, runId, messageId }; - }, - { - body: t.Object({ content: t.String(), model: t.String() }), - }, - ); - - app.get("/conversations/:id/events", ({ params }) => { - const actor = getActor(params.id, store); - return actor.replay(0).map((e) => e); - }); - - return app; -} - -// Side effects happen only when run as the entry point (`bun server.ts`), so -// tests can import `buildApp` and use their own store without touching `data/`. -const isEntryPoint = import.meta.main; - -let store: Store; -let app: Elysia; -if (isEntryPoint) { - // Load the catalog and build the provider registry before serving, so the - // first request already has models resolvable. +if (import.meta.main) { + // Load the catalog + provider registry before serving, so the first request + // already has models resolvable. await initInference(); - store = new Store(); - app = buildApp({ store }); - - // Track in-flight runs per conversation so the drive loop never claims two - // jobs for the same conversation concurrently (single-writer invariant). - const activeRuns = new Set(); - - async function driveOnce(): Promise { - const row = store.claimExpiredExclusive(Date.now()); - if (!row) return; - if (activeRuns.has(row.conversation_id)) return; - - const params = JSON.parse(row.params) as { - runId: string; - messageId: string; - prompt: string; - model: string; - }; - const actor = getActor(row.conversation_id, store); - activeRuns.add(row.conversation_id); - try { - await actor.runText( - params.runId, - params.messageId, - async function* (signal) { - for await (const step of run(params.prompt, { - runId: params.runId, - model: params.model, - abortSignal: signal, - })) { - yield step; - } - }, - (seq) => { - // Advance the job's durable checkpoint + lease on each flush so a - // crash mid-run is re-claimed from the last flushed seq. - store.checkpoint(row.id, seq); - store.heartbeat(row.id, Date.now() + LEASE_GRACE_MS); - }, - ); - store.markDone(row.id); - } catch (err) { - store.markFailed(row.id); - } finally { - activeRuns.delete(row.conversation_id); - } - } + const store = new Store(); + const driver = new JobDriver(store, (id) => getActor(id, store)); setInterval(() => { - void driveOnce(); + void driver.driveOnce(); }, 1000); // Reaper: re-queue jobs whose lease expired (worker died mid-run) so any - // process (this one or a peer) can claim them again from checkpoint_seq. + // process can re-claim them from checkpoint_seq; also evict idle actors. setInterval(() => { store.reap(Date.now()); evictIdleActors(); }, REAP_INTERVAL_MS); const port = Number(process.env.PORT ?? 3000); - Bun.serve({ port, fetch: app.fetch }); + // The SSE stream is intentionally long-lived and can sit idle between + // generations. Bun's default idleTimeout is 10s β€” shorter than our 15s + // keepalive β€” so an idle stream would be killed before the first keepalive + // fires. Raise it to Bun's max (255s); the keepalive resets the idle clock + // well inside that window, so streams stay open indefinitely. + Bun.serve({ + port, + idleTimeout: 255, + development: process.env.NODE_ENV !== "production", + routes: { + "/": indexHTML, + "/settings": settingsHTML, + ...apiRoutes({ store }), + }, + }); console.log(`kloe listening on http://localhost:${port}`); } diff --git a/src/actor.ts b/src/actor.ts index 066e001..aaeb8f0 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -115,9 +115,9 @@ export class ConversationActor { })); } - /** Current tail seq, used for Last-Event-ID resume cursor arithmetic. */ - lastCommittedSeq(): number { - return this.seq; + /** True while at least one SSE stream is attached to this actor. */ + hasSubscribers(): boolean { + return this.subscribers.size > 0; } /** diff --git a/src/client/app.css b/src/client/app.css new file mode 100644 index 0000000..fa64c1d --- /dev/null +++ b/src/client/app.css @@ -0,0 +1,207 @@ +:root { + --sans: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, "Helvetica Neue", Arial, sans-serif; + --mono: ui-monospace, "SF Mono", "Cascadia Code", "Roboto Mono", Menlo, Consolas, monospace; + --bg:#ffffff; --bg-sunk:#f6f6f7; --bg-raise:#ffffff; --ink:#17181a; + --ink-muted:#6a6d73; --ink-faint:#9a9da3; --rule:#e6e7e9; --rule-strong:#d9dade; + --accent:#0a66d6; --ok:#3f8f5f; --warn:#b07d1a; --err:#c0392b; --measure:72ch; +} +@media (prefers-color-scheme: dark) { + :root { --bg:#121315; --bg-sunk:#1b1c1f; --bg-raise:#1a1b1e; --ink:#e8e9ea; + --ink-muted:#9a9da3; --ink-faint:#6a6d73; --rule:#2a2c2f; --rule-strong:#3a3d42; + --accent:#4c9aff; --ok:#5bb37f; --warn:#d0a044; --err:#e06c5e; } +} +[data-theme="light"] { --bg:#ffffff; --bg-sunk:#f6f6f7; --bg-raise:#ffffff; --ink:#17181a; + --ink-muted:#6a6d73; --ink-faint:#9a9da3; --rule:#e6e7e9; --rule-strong:#d9dade; + --accent:#0a66d6; --ok:#3f8f5f; --warn:#b07d1a; --err:#c0392b; } +[data-theme="dark"] { --bg:#121315; --bg-sunk:#1b1c1f; --bg-raise:#1a1b1e; --ink:#e8e9ea; + --ink-muted:#9a9da3; --ink-faint:#6a6d73; --rule:#2a2c2f; --rule-strong:#3a3d42; + --accent:#4c9aff; --ok:#5bb37f; --warn:#d0a044; --err:#e06c5e; } + +* { box-sizing: border-box; } +html, body { height: 100%; } +body { margin: 0; background: var(--bg); color: var(--ink); font-family: var(--sans); + font-size: 15.5px; line-height: 1.6; -webkit-font-smoothing: antialiased; + text-rendering: optimizeLegibility; display: flex; flex-direction: column; } +:focus-visible { outline: 2px solid var(--accent); outline-offset: 2px; border-radius: 3px; } +a { color: var(--accent); } + +.app { flex: 1; display: grid; grid-template-columns: 232px 1fr; min-height: 0; } +@media (max-width: 720px) { .app { grid-template-columns: 1fr; } } + +.rail { border-right: 1px solid var(--rule); min-height: 0; overflow-y: auto; + padding: 10px 8px; display: flex; flex-direction: column; gap: 2px; } +@media (max-width: 720px) { .rail { display: none; } .rail.open { display: flex; + position: fixed; inset: 0 30% 0 0; z-index: 20; background: var(--bg); } } +.rail .new { text-align: left; font: inherit; color: var(--ink); background: none; + border: 1px solid var(--rule-strong); border-radius: 6px; padding: 7px 10px; cursor: pointer; margin-bottom: 8px; } +.rail .new:hover { background: var(--bg-sunk); } +.rail .conv { text-align: left; font: inherit; font-size: 14px; color: var(--ink-muted); + background: none; border: 0; border-radius: 6px; padding: 7px 10px; cursor: pointer; + width: 100%; white-space: nowrap; overflow: hidden; text-overflow: ellipsis; } +.rail .conv:hover { background: var(--bg-sunk); color: var(--ink); } +.rail .conv[aria-current="true"] { background: var(--bg-sunk); color: var(--ink); font-weight: 500; } +.rail .group { font-family: var(--mono); font-size: 10.5px; color: var(--ink-faint); + text-transform: uppercase; letter-spacing: .08em; padding: 12px 10px 4px; } +.rail .spacer { flex: 1; } +.rail .railfoot { padding: 8px 10px; } +.rail .railfoot a { font-family: var(--mono); font-size: 12px; color: var(--ink-muted); text-decoration: none; } +.rail .railfoot a:hover { color: var(--ink); } + +.main { display: flex; flex-direction: column; min-height: 0; min-width: 0; } +.head { border-bottom: 1px solid var(--rule); padding: 10px 16px; display: flex; align-items: baseline; gap: 10px; } +.head .title { font-size: 14.5px; font-weight: 600; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.head .menu { display: none; } +@media (max-width: 720px) { .head .menu { display: inline-flex; font: inherit; background: none; + border: 1px solid var(--rule-strong); border-radius: 6px; padding: 4px 9px; cursor: pointer; color: var(--ink); } } +.head .connstat { margin-left: auto; font-family: var(--mono); font-size: 11.5px; } +.head .connstat .conn { display: none; } +.head .connstat[data-state="reconnecting"] .conn { display: inline; color: var(--warn); } +.head .connstat[data-state="offline"] .conn { display: inline; color: var(--err); } + +.scroll { flex: 1; overflow-y: auto; min-height: 0; position: relative; } +.thread { max-width: var(--measure); margin: 0 auto; padding: 22px 20px 32px; } +.turn { padding: 18px 0; border-top: 1px solid var(--rule); } +.turn:first-child { border-top: 0; } +.turn .label { display: flex; align-items: baseline; gap: 10px; margin-bottom: 8px; } +.turn .who { font-size: 13px; font-weight: 600; } +.turn .who.user { color: var(--ink-muted); } +.turn .meta { font-family: var(--mono); font-size: 11px; color: var(--ink-faint); + margin-left: auto; text-align: right; white-space: nowrap; + opacity: .42; transition: opacity .12s ease; } +.turn:hover .meta, .turn:focus-within .meta, .turn.generating .meta { opacity: 1; } +.turn .body > :first-child { margin-top: 0; } +.turn .body > :last-child { margin-bottom: 0; } +.turn .body p { margin: 0 0 12px; } +.turn .body code { font-family: var(--mono); font-size: 0.88em; background: var(--bg-sunk); + border: 1px solid var(--rule); border-radius: 4px; padding: 0.05em 0.35em; } +.turn .body pre { font-family: var(--mono); font-size: 12.5px; line-height: 1.55; + background: var(--bg-sunk); border: 1px solid var(--rule); border-radius: 7px; + padding: 12px 14px; overflow-x: auto; margin: 0 0 12px; } +.turn .body pre code { background: none; border: 0; padding: 0; font-size: inherit; } +.turn .body ul, .turn .body ol { margin: 0 0 12px; padding-left: 1.4em; } +.turn .body li { margin: 2px 0; } +.turn .body blockquote { margin: 0 0 12px; padding-left: 12px; border-left: 3px solid var(--rule-strong); color: var(--ink-muted); } +.turn .body h1, .turn .body h2, .turn .body h3 { margin: 16px 0 8px; line-height: 1.3; } +.turn .body h1 { font-size: 1.3em; } .turn .body h2 { font-size: 1.18em; } .turn .body h3 { font-size: 1.05em; } +.turn .body a { color: var(--accent); text-decoration: underline; text-underline-offset: 2px; } +.turn .body table { border-collapse: collapse; margin: 0 0 12px; display: block; overflow-x: auto; } +.turn .body th, .turn .body td { border: 1px solid var(--rule); padding: 4px 9px; text-align: left; } +.turn .body img { max-width: 100%; border-radius: 6px; } +/* streaming caret (spec: blinking caret at the stream head), no injected node */ +.turn.generating .body::after { content: ""; display: inline-block; width: 7px; height: 1.05em; + vertical-align: text-bottom; margin-left: 1px; background: var(--accent); border-radius: 1px; + animation: blink 1s steps(2) infinite; } +@keyframes blink { 50% { opacity: 0; } } + +.turn.pending { opacity: .6; } +.turn.failed .body { color: var(--ink-muted); } +.failbar { margin-top: 10px; font-family: var(--mono); font-size: 12px; color: var(--err); + display: flex; align-items: center; gap: 10px; } +.failbar button { font-family: var(--mono); font-size: 12px; color: var(--ink); + background: var(--bg); border: 1px solid var(--rule-strong); border-radius: 5px; padding: 2px 8px; cursor: pointer; } + +.empty { max-width: var(--measure); margin: 0 auto; padding: 60px 20px; color: var(--ink-muted); + text-align: center; font-size: 14px; } +.empty code { font-family: var(--mono); font-size: 12.5px; background: var(--bg-sunk); + border: 1px solid var(--rule); border-radius: 4px; padding: 1px 5px; } + +.jump { position: absolute; left: 50%; transform: translateX(-50%); bottom: 12px; + font-family: var(--mono); font-size: 12px; color: var(--ink-muted); background: var(--bg); + border: 1px solid var(--rule-strong); border-radius: 999px; padding: 4px 12px; cursor: pointer; display: none; } + +/* ---- composer dock ----------------------------------------------------- */ +.dock { border-top: 0; background: var(--bg); } +.dockinner { max-width: var(--measure); margin: 0 auto; padding: 8px 16px 16px; } + +.banner { display: none; font-family: var(--mono); font-size: 12px; color: var(--warn); + padding: 8px 12px; margin-bottom: 10px; border: 1px solid color-mix(in srgb, var(--warn) 40%, var(--rule)); + border-radius: 9px; background: color-mix(in srgb, var(--warn) 8%, var(--bg)); } +body:has(.connstat[data-state="offline"]) .banner { display: block; } + +.card { border: 1px solid var(--rule-strong); border-radius: 16px; background: var(--bg-raise); + padding: 0 12px 8px; display: flex; flex-direction: column; transition: border-color .12s ease; } +.card:hover, .card:focus-within { border-color: color-mix(in srgb, var(--accent) 15%, var(--rule-strong)); } +.card textarea { width: 100%; resize: none; border: 0; background: none; color: var(--ink); + font: inherit; font-size: 16px; line-height: 1.5; padding: 13px 8px 8px; min-height: 45px; max-height: 44vh; overflow-y: auto; } +.card textarea:focus { outline: none; } +.card textarea::placeholder { color: var(--ink-faint); } + +.tools { display: flex; align-items: center; gap: 8px; margin-top: 0; } +.tools .right { margin-left: auto; display: flex; align-items: center; gap: 4px; } +.icon { display: inline-flex; align-items: center; justify-content: center; width: 34px; height: 34px; + border-radius: 8px; border: 0; background: none; color: var(--ink-muted); cursor: pointer; text-decoration: none; } +.icon:hover { background: color-mix(in srgb, var(--ink) 9%, transparent); color: var(--ink); } +.icon svg { width: 18px; height: 18px; } + +/* context window β€” unicode shaded bar, next to the model selector (~approximate) */ +.ctx { font-family: var(--mono); font-size: 12px; color: var(--ink-muted); background: none; + border: 0; border-radius: 8px; padding: 6px 8px; cursor: default; display: inline-flex; + align-items: center; gap: 7px; white-space: nowrap; } +.ctx.hidden { display: none; } +.ctx .bar { letter-spacing: -0.5px; color: color-mix(in srgb, var(--accent) 55%, var(--ink-muted)); } +.ctx .pct { color: var(--ink-faint); } + +.pickwrap { position: relative; display: inline-flex; } +.pill { font: inherit; font-size: 13.5px; color: var(--ink); background: none; border: 0; + border-radius: 8px; padding: 6px 8px; cursor: pointer; display: inline-flex; align-items: center; gap: 4px; } +.pill:hover { background: color-mix(in srgb, var(--ink) 9%, transparent); } +.pill:disabled { color: var(--ink-faint); cursor: default; } +.pill .m { font-weight: 500; } +.pill .e { color: var(--ink-muted); } +.pill svg { width: 13px; height: 13px; opacity: .6; margin-left: 2px; } + +.picker { position: absolute; bottom: calc(100% + 6px); right: 0; z-index: 30; min-width: 260px; + max-width: 340px; max-height: 50vh; overflow-y: auto; background: var(--bg-raise); + border: 1px solid var(--rule-strong); border-radius: 12px; padding: 6px; + box-shadow: 0 8px 30px rgba(0,0,0,.18); } +.picker[hidden] { display: none; } +.picker .opt { display: block; width: 100%; text-align: left; font: inherit; color: var(--ink); + background: none; border: 0; border-radius: 8px; padding: 8px 10px; cursor: pointer; } +.picker .opt:hover { background: var(--bg-sunk); } +.picker .opt[aria-selected="true"] { background: color-mix(in srgb, var(--accent) 12%, var(--bg-raise)); } +.picker .opt .name { font-weight: 500; font-size: 14px; } +.picker .opt .sub { font-family: var(--mono); font-size: 11px; color: var(--ink-faint); margin-top: 2px; } +.picker .none { padding: 10px; font-size: 13px; color: var(--ink-muted); } +.picker .none a { white-space: nowrap; } + +.send { display: inline-flex; align-items: center; justify-content: center; width: 36px; height: 36px; + border-radius: 10px; border: 1px solid transparent; background: var(--bg-sunk); color: var(--ink-faint); + cursor: default; margin-left: 4px; padding: 0; font-family: var(--mono); font-size: 16px; line-height: 1; } +.send svg { width: 18px; height: 18px; } +.send.ready { background: var(--accent); color: #fff; cursor: pointer; } +.send.ready:hover { filter: brightness(1.08); } +.send.stop { background: var(--bg); color: var(--ink); border-color: var(--rule-strong); cursor: pointer; } +.send.stop:hover { background: color-mix(in srgb, var(--ink) 9%, transparent); } + +/* ---- settings ---------------------------------------------------------- */ +.settings { max-width: 860px; margin: 0 auto; padding: 28px 24px 60px; width: 100%; } +.settings h1 { font-size: 20px; margin: 0 0 4px; } +.settings .lede { color: var(--ink-muted); font-size: 14px; margin: 0 0 24px; } +.settings .back { font-family: var(--mono); font-size: 12px; color: var(--ink-muted); text-decoration: none; } +.settings .back:hover { color: var(--ink); } +.mtable { width: 100%; border-collapse: collapse; font-size: 14px; } +.mtable th { text-align: left; font-family: var(--mono); font-size: 11px; text-transform: uppercase; + letter-spacing: .06em; color: var(--ink-faint); font-weight: 500; padding: 8px 10px; border-bottom: 1px solid var(--rule); } +.mtable td { padding: 10px; border-bottom: 1px solid var(--rule); vertical-align: middle; } +.mtable tr:hover td { background: var(--bg-sunk); } +.mtable .ref { font-family: var(--mono); font-size: 12px; color: var(--ink-muted); } +.mtable .cap { font-family: var(--mono); font-size: 11px; color: var(--ink-faint); } +.mtable input[type="text"] { font: inherit; font-size: 13px; color: var(--ink); background: var(--bg); + border: 1px solid var(--rule-strong); border-radius: 6px; padding: 4px 8px; width: 140px; } +.mtable input[type="number"] { font: inherit; font-size: 13px; color: var(--ink); background: var(--bg); + border: 1px solid var(--rule-strong); border-radius: 6px; padding: 4px 6px; width: 56px; } +.mtable .prov { font-weight: 600; font-size: 13px; padding-top: 20px; color: var(--ink-muted); } +.toggle { appearance: none; width: 34px; height: 20px; border-radius: 999px; background: var(--rule-strong); + position: relative; cursor: pointer; transition: background .12s ease; flex: none; } +.toggle:checked { background: var(--accent); } +.toggle::after { content: ""; position: absolute; top: 2px; left: 2px; width: 16px; height: 16px; + border-radius: 50%; background: #fff; transition: transform .12s ease; } +.toggle:checked::after { transform: translateX(14px); } +.saved { font-family: var(--mono); font-size: 11px; color: var(--ok); opacity: 0; transition: opacity .2s; } +.saved.show { opacity: 1; } + +@media (prefers-reduced-motion: reduce) { + .turn .meta { transition: none; } + .cursor { animation: none; } + * { scroll-behavior: auto !important; } +} diff --git a/src/client/app.js b/src/client/app.js new file mode 100644 index 0000000..8d168ad --- /dev/null +++ b/src/client/app.js @@ -0,0 +1,460 @@ +/* + * kloe chat frontend β€” aligned with spec.md. + * + * - Eagerness (spec "Eagerness"): the user's own message is appended + * OPTIMISTICALLY the instant they submit β€” network is off the feedback path. + * The message carries a client-generated runId; the server echoes that same + * runId back over SSE, and we dedupe against it. POST failure marks the turn + * failed and offers retry (the spec's rollback). + * - Rendering (spec "Rendering"): assistant deltas are streamed into the DOM + * with streaming-markdown (smd) β€” append-only, so already-streamed text is + * never re-parsed (no O(nΒ²) re-render). Writes are batched to a rAF (~60fps). + * Sanitization is structural: smd never emits raw HTML (verified β€” raw tags + * become text nodes), so the only XSS vector is link/image URLs, which the + * wrapped renderer neutralizes as they render. That makes a heavyweight + * HTML sanitizer (DOMPurify, ~118KB) unnecessary β€” it's deliberately absent. + * + * No framework, no build step: an ES module importing one small vendored, + * zero-dep library (streaming-markdown, ~3KB brotli). The durable event log + * stays authoritative β€” optimism is a local projection reconciled by the echo. + */ +import * as smd from "streaming-markdown"; + +(function () { + "use strict"; + + var SEND = ''; + var STOP = ''; + + var $ = function (id) { return document.getElementById(id); }; + var thread = $("thread"), scroll = $("scroll"), jump = $("jump"); + var input = $("input"), send = $("send"), composer = $("composer"); + var title = $("title"), status = $("status"), conn = $("conn"); + var railList = $("railList"), rail = $("rail"); + var pill = $("pill"), pillModel = $("pillModel"), picker = $("picker"); + var ctx = $("ctx"), ctxbar = $("ctxbar"), ctxpct = $("ctxpct"); + + // ---- state ------------------------------------------------------------- + var convId = null; // current conversation id + var source = null; // active EventSource + var streaming = false; // a run is in flight for the current conversation + var atBottom = true; + var models = [], selected = null; + var msgs = Object.create(null); // messageId -> assistant render record + var pending = Object.create(null); // runId -> optimistic user turn awaiting echo + var flushHandle = null; + var lastUsage = null; // real token usage from the last completed turn + + // ---- streaming-markdown rendering -------------------------------------- + // Wrap smd's default renderer to harden URLs. smd never emits raw HTML tags + // (model text lands in text nodes), so href/src are the only injection + // vector β€” we neutralize dangerous schemes and reveal external link targets. + function makeRenderer(root) { + var r = smd.default_renderer(root); + r._stripped = false; + var base = r.set_attr; + r.set_attr = function (data, type, value) { + var out = value; + if (type === smd.HREF || type === smd.SRC) { + if (/^\s*(javascript|vbscript|file):/i.test(value)) out = "#"; + else if (/^\s*data:/i.test(value) && + !(type === smd.SRC && /^\s*data:image\//i.test(value))) out = "#"; + if (out !== value) r._stripped = true; + } + base(data, type, out); + if (type === smd.HREF) { + var node = data.nodes[data.index]; + if (node && node.tagName === "A") { + node.setAttribute("target", "_blank"); + node.setAttribute("rel", "noopener noreferrer nofollow"); + node.setAttribute("title", out); // spec: show the full URL before navigating + } + } + }; + return r; + } + function newParser(root) { + var renderer = makeRenderer(root); + return { renderer: renderer, parser: smd.parser(renderer) }; + } + // Reports whether the renderer had to neutralize a dangerous URL. No HTML + // sanitizer needed: smd builds the DOM node-by-node and never emits raw tags, + // so there's no untrusted HTML string to purify β€” only the href/src the + // wrapped renderer already guarded. + function finalize(renderer) { return !!(renderer && renderer._stripped); } + function renderStaticMd(el, text) { + var np = newParser(el); + smd.parser_write(np.parser, text); + smd.parser_end(np.parser); + return finalize(np.renderer); + } + + // rAF-batched delta flush: models emit faster than the eye needs (spec). + function scheduleFlush() { if (!flushHandle) flushHandle = requestAnimationFrame(flush); } + function flush() { + flushHandle = null; + var painted = false; + for (var id in msgs) { + var r = msgs[id]; + if (r.buf) { smd.parser_write(r.parser, r.buf); r.buf = ""; painted = true; liveMeta(r); } + } + if (painted) autoScroll(); + } + + // ---- helpers ----------------------------------------------------------- + function autoScroll() { if (atBottom) scroll.scrollTop = scroll.scrollHeight; } + // Compact context-window label: 1000000 -> "1M", 1048576 -> "1M", 1500000 -> "1.5M", else "Nk". + function fmtCtx(n) { + if (n >= 1e6) return (n / 1e6).toFixed(1).replace(/\.0$/, "") + "M"; + return Math.round(n / 1000) + "k"; + } + // Real context fill: the last completed turn's input + output tokens (that's + // what the provider actually counted, and the baseline the next turn sends). + function usedTokens(u) { + if (!u) return null; + if (u.inputTokens != null || u.outputTokens != null) return (u.inputTokens || 0) + (u.outputTokens || 0); + if (u.totalTokens != null) return u.totalTokens; + return null; + } + function updateCtx() { + var used = usedTokens(lastUsage); + if (!selected || !selected.contextWindow || used == null) { ctx.classList.add("hidden"); return; } + var pct = Math.max(0, Math.min(100, Math.round((used / selected.contextWindow) * 100))); + var n = 12, f = Math.round((pct / 100) * n); + ctxbar.textContent = "β–“".repeat(f) + "β–‘".repeat(n - f); + ctxpct.textContent = pct + "%"; + ctx.classList.remove("hidden"); + ctx.title = used.toLocaleString() + " / " + selected.contextWindow.toLocaleString() + " tokens"; + } + function updateSend() { + if (streaming) { + send.innerHTML = STOP; send.className = "send stop"; send.disabled = false; + send.setAttribute("aria-label", "Stop"); return; + } + send.innerHTML = SEND; + var has = input.value.trim().length > 0 && selected; + send.className = "send" + (has ? " ready" : ""); + send.disabled = !has; + send.setAttribute("aria-label", "Send"); + } + // While streaming we don't have real token counts yet (the provider reports + // usage only at the end), so show only measured wall-clock: ttft + elapsed. + // No estimated token rate β€” real counts land on message-end. + function liveMeta(rec) { + var elapsed = (Date.now() - rec.startedAt) / 1000; + var ttft = rec.firstDeltaAt ? ((rec.firstDeltaAt - rec.startedAt) / 1000).toFixed(1) + "s ttft Β· " : ""; + rec.meta.textContent = ttft + elapsed.toFixed(1) + "s"; + } + + // ---- thread rendering -------------------------------------------------- + function clearThread() { + thread.innerHTML = ""; + msgs = Object.create(null); + pending = Object.create(null); + streaming = false; + lastUsage = null; + if (flushHandle) { cancelAnimationFrame(flushHandle); flushHandle = null; } + } + function makeTurn(who, cls) { + var t = document.createElement("article"); + t.className = "turn" + (cls ? " " + cls : ""); + t.innerHTML = + '
' + + '
'; + t.querySelector(".who").textContent = who; + thread.appendChild(t); + return t; + } + + function optimisticUser(content, runId) { + var t = makeTurn("You", "pending"); + t.dataset.runId = runId; + renderStaticMd(t.querySelector(".body"), content); + autoScroll(); + pending[runId] = { turn: t, content: content }; + } + function confirmUser(runId, content) { + var p = pending[runId]; + if (p) { p.turn.classList.remove("pending", "failed"); delete pending[runId]; return; } + // Not ours (history, or another device): render fresh. + var t = makeTurn("You"); + renderStaticMd(t.querySelector(".body"), content); + autoScroll(); + } + function failUser(runId) { + var p = pending[runId]; + if (!p) return; + p.turn.classList.remove("pending"); + p.turn.classList.add("failed"); + if (p.turn.querySelector(".failbar")) return; + var fb = document.createElement("div"); + fb.className = "failbar"; + fb.textContent = "not sent β€” "; + var btn = document.createElement("button"); + btn.type = "button"; btn.textContent = "Retry"; + btn.onclick = function () { + p.turn.remove(); delete pending[runId]; + doSend(p.content, runId); + }; + fb.appendChild(btn); + p.turn.querySelector(".body").appendChild(fb); + } + + function assistantTurn(messageId) { + if (msgs[messageId]) return msgs[messageId]; + var t = makeTurn("Assistant", "generating"); + var body = t.querySelector(".body"); + var np = newParser(body); + var rec = { + turn: t, body: body, meta: t.querySelector(".meta"), + parser: np.parser, renderer: np.renderer, buf: "", + startedAt: Date.now(), firstDeltaAt: 0, + }; + msgs[messageId] = rec; + autoScroll(); + return rec; + } + function endAssistant(rec, finishReason, usage) { + if (rec.buf) { smd.parser_write(rec.parser, rec.buf); rec.buf = ""; } + smd.parser_end(rec.parser); + var stripped = finalize(rec.renderer); + rec.turn.classList.remove("generating"); + if (finishReason === "aborted") { + rec.meta.textContent = "stopped"; + } else if (finishReason === "error") { + rec.turn.classList.add("failed"); failbar(rec.body, "generation failed"); rec.meta.textContent = "error"; + } else { + var elapsed = (Date.now() - rec.startedAt) / 1000; + var live = elapsed >= 0.3; // actually streamed in this session vs replayed from the log + var out = usage && usage.outputTokens != null ? usage.outputTokens : null; + if (out != null && live) { + var ttft = rec.firstDeltaAt ? ((rec.firstDeltaAt - rec.startedAt) / 1000).toFixed(1) + "s ttft Β· " : ""; + rec.meta.textContent = ttft + out.toLocaleString() + " tok Β· " + Math.round(out / elapsed) + " tok/s"; + } else if (out != null) { + rec.meta.textContent = out.toLocaleString() + " tok"; // replayed: real count, no synthetic rate + } else if (!live) { + rec.meta.textContent = ""; + } else { + liveMeta(rec); // provider reported no usage: fall back to wall-clock + } + } + if (usage) { lastUsage = usage; updateCtx(); } + if (stripped) failbar(rec.body, "some content was removed by the sanitizer"); + } + function failbar(bodyEl, msg) { + if (bodyEl.querySelector(".failbar")) return; + var fb = document.createElement("div"); + fb.className = "failbar"; fb.textContent = msg; + bodyEl.appendChild(fb); + } + function lastAssistant() { + var keys = Object.keys(msgs); + return keys.length ? msgs[keys[keys.length - 1]] : null; + } + + function applyEvent(name, data) { + switch (name) { + case "user-message": + confirmUser(data.runId, data.content); + break; + case "message-start": + streaming = true; + assistantTurn(data.messageId); + updateSend(); + break; + case "text-delta": { + var rec = assistantTurn(data.messageId); + if (!rec.firstDeltaAt) rec.firstDeltaAt = Date.now(); + rec.buf += data.delta; + scheduleFlush(); + break; + } + case "message-end": { + streaming = false; + var r = msgs[data.messageId]; + if (r) endAssistant(r, data.finishReason, data.usage); + updateSend(); + break; + } + case "run-error": { + streaming = false; + var la = lastAssistant(); + if (la) { la.turn.classList.add("failed"); failbar(la.body, String(data.error || "run error")); la.meta.textContent = "error"; } + updateSend(); + break; + } + case "run-started": + case "cancelled": + break; + } + } + + // ---- SSE stream -------------------------------------------------------- + var connTimer = null; + function openStream(id) { + if (source) { source.close(); source = null; } + if (connTimer) { clearTimeout(connTimer); connTimer = null; } + clearThread(); + convId = id; + var es = new EventSource("/api/conversations/" + encodeURIComponent(id) + "/stream"); + source = es; + ["user-message", "run-started", "message-start", "text-delta", + "message-end", "run-error", "cancelled"].forEach(function (nm) { + es.addEventListener(nm, function (ev) { + var data; try { data = JSON.parse(ev.data); } catch (_) { return; } + applyEvent(nm, data); + }); + }); + es.onopen = function () { setConn("connected"); }; + es.onerror = function () { + // readyState 2 (CLOSED) is terminal; 0 (CONNECTING) is the browser already + // auto-reconnecting β€” usually done within a second. Only surface + // "reconnecting" if the gap actually lingers, so a quick reconnect doesn't + // flash the header. + if (es.readyState === 2) { setConn("offline"); return; } + if (!connTimer) connTimer = setTimeout(function () { connTimer = null; setConn("reconnecting"); }, 1500); + }; + } + function setConn(s) { + if (s === "connected" && connTimer) { clearTimeout(connTimer); connTimer = null; } + status.dataset.state = s; + conn.textContent = s === "reconnecting" ? "reconnecting…" : s; + } + + // ---- composer / sending ------------------------------------------------ + function autosize() { + input.style.height = "auto"; + input.style.height = Math.min(input.scrollHeight, window.innerHeight * 0.44) + "px"; + } + function submit() { + var content = input.value.trim(); + if (!content || !selected || streaming) return; + input.value = ""; autosize(); + var runId = (crypto.randomUUID ? crypto.randomUUID() : String(Date.now()) + Math.random()); + doSend(content, runId); + } + // Optimistic apply now, POST after: the user's turn is on screen before the + // request leaves. Reconciled by the server's user-message echo (same runId). + async function doSend(content, runId) { + var wasNew = !hasConversation(convId); + optimisticUser(content, runId); + streaming = true; updateSend(); // job is queued + cancellable even pre-first-token + try { + var res = await fetch("/api/conversations/" + encodeURIComponent(convId) + "/prompt", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: content, model: selected.ref, runId: runId }), + }); + if (!res.ok) { + streaming = false; updateSend(); + var err = await res.json().catch(function () { return {}; }); + failUser(runId); + if (err.error) console.warn("prompt rejected:", err.error); + return; + } + if (wasNew) { title.textContent = content.slice(0, 80); setTimeout(loadConversations, 400); } + } catch (e) { + streaming = false; updateSend(); + failUser(runId); + } + } + function stop() { + fetch("/api/conversations/" + encodeURIComponent(convId) + "/cancel", { method: "POST" }).catch(function () {}); + } + + // ---- conversation rail ------------------------------------------------- + var conversations = []; + function hasConversation(id) { return conversations.some(function (c) { return c.id === id; }); } + async function loadConversations() { + try { + var res = await fetch("/api/conversations"); + conversations = (await res.json()).conversations || []; + } catch (_) { conversations = []; } + renderRail(); + } + function renderRail() { + railList.innerHTML = ""; + conversations.forEach(function (c) { + var b = document.createElement("button"); + b.className = "conv"; + b.textContent = c.title || "Untitled"; + if (c.id === convId) b.setAttribute("aria-current", "true"); + b.onclick = function () { selectConversation(c.id, c.title); rail.classList.remove("open"); }; + railList.appendChild(b); + }); + } + function selectConversation(id, t) { title.textContent = t || "Conversation"; openStream(id); renderRail(); } + function newConversation() { + var id = (crypto.randomUUID ? crypto.randomUUID() : String(Date.now()) + Math.random()); + title.textContent = "New conversation"; + openStream(id); renderRail(); input.focus(); + } + + // ---- model picker ------------------------------------------------------ + async function loadModels() { + try { models = (await (await fetch("/api/models/chat")).json()).models || []; } + catch (_) { models = []; } + var saved = localStorage.getItem("kloe.model"); + selected = models.find(function (m) { return m.ref === saved; }) || models[0] || null; + renderPicker(); renderPill(); updateSend(); updateCtx(); + } + function renderPill() { + pill.disabled = models.length === 0; + pillModel.textContent = selected ? selected.name : (models.length ? "Select model" : "No models"); + } + function renderPicker() { + picker.innerHTML = ""; + if (models.length === 0) { + var none = document.createElement("div"); + none.className = "none"; + none.innerHTML = 'No models enabled. Turn some on in Settings.'; + picker.appendChild(none); + return; + } + models.forEach(function (m) { + var b = document.createElement("button"); + b.className = "opt"; b.type = "button"; b.setAttribute("role", "option"); + if (selected && m.ref === selected.ref) b.setAttribute("aria-selected", "true"); + var sub = []; + if (m.contextWindow) sub.push(fmtCtx(m.contextWindow) + " ctx"); + if (m.reasoningLevels && m.reasoningLevels.length) sub.push("reasoning"); + if (m.supportsImages) sub.push("images"); + b.innerHTML = '
' + (sub.length ? '
' + sub.join(" Β· ") + "
" : ""); + b.querySelector(".name").textContent = m.name; + b.onclick = function () { + selected = m; localStorage.setItem("kloe.model", m.ref); + renderPill(); renderPicker(); updateSend(); updateCtx(); closePicker(); + }; + picker.appendChild(b); + }); + } + function openPicker() { picker.hidden = false; pill.setAttribute("aria-expanded", "true"); } + function closePicker() { picker.hidden = true; pill.setAttribute("aria-expanded", "false"); } + + // ---- wiring ------------------------------------------------------------ + composer.addEventListener("submit", function (e) { e.preventDefault(); if (streaming) stop(); else submit(); }); + input.addEventListener("input", function () { autosize(); updateSend(); }); + input.addEventListener("keydown", function (e) { + if (e.key === "Enter" && !e.shiftKey) { e.preventDefault(); if (!streaming) submit(); } + }); + $("new").addEventListener("click", newConversation); + $("menu").addEventListener("click", function () { rail.classList.toggle("open"); }); + pill.addEventListener("click", function () { if (picker.hidden) openPicker(); else closePicker(); }); + document.addEventListener("click", function (e) { + if (!picker.hidden && !picker.contains(e.target) && e.target !== pill && !pill.contains(e.target)) closePicker(); + }); + document.addEventListener("keydown", function (e) { if (e.key === "Escape") closePicker(); }); + scroll.addEventListener("scroll", function () { + atBottom = scroll.scrollHeight - scroll.scrollTop - scroll.clientHeight < 40; + jump.style.display = atBottom ? "none" : "block"; + }); + jump.addEventListener("click", function () { atBottom = true; jump.style.display = "none"; scroll.scrollTop = scroll.scrollHeight; }); + + // ---- boot -------------------------------------------------------------- + (async function init() { + await Promise.all([loadModels(), loadConversations()]); + if (conversations.length) selectConversation(conversations[0].id, conversations[0].title); + else newConversation(); + updateSend(); + })(); +})(); diff --git a/src/client/index.html b/src/client/index.html new file mode 100644 index 0000000..86d3445 --- /dev/null +++ b/src/client/index.html @@ -0,0 +1,63 @@ + + + + + + +kloe + + + +
+ + +
+
+ + New conversation + +
+ +
+
+ +
+ +
+
+ + +
+ +
+ + + +
+ + + + + + +
+
+
+
+
+
+
+ + + diff --git a/src/client/settings.html b/src/client/settings.html new file mode 100644 index 0000000..15dd3b4 --- /dev/null +++ b/src/client/settings.html @@ -0,0 +1,20 @@ + + + + + + +kloe β€” settings + + + +
+ ← Back to chat +

Models

+

Choose which models appear in the chat picker, and give them display names. + Chat is opt-in β€” nothing shows until you enable it here.

+
Loading…
+
+ + + diff --git a/src/client/settings.js b/src/client/settings.js new file mode 100644 index 0000000..be174ee --- /dev/null +++ b/src/client/settings.js @@ -0,0 +1,125 @@ +/* + * Model curation. Lists every model this deployment can run (GET /models, joined + * to catalog metadata + curation state) and lets the operator toggle chat + * visibility, set a display name, and order the picker (PATCH /models). Each + * edit is a partial PATCH β€” omitted fields keep their stored value. + */ +(function () { + "use strict"; + var content = document.getElementById("content"); + + function group(models) { + // Bucket by provider (the ref prefix) for a readable table. + var by = {}; + models.forEach(function (m) { + var prov = m.ref.split("/")[0]; + (by[prov] = by[prov] || []).push(m); + }); + return by; + } + + async function patch(ref, field, value, savedEl) { + var body = { ref: ref }; + body[field] = value; + try { + var res = await fetch("/api/models", { + method: "PATCH", + headers: { "content-type": "application/json" }, + body: JSON.stringify(body), + }); + if (!res.ok) throw new Error(String(res.status)); + if (savedEl) { savedEl.classList.add("show"); setTimeout(function () { savedEl.classList.remove("show"); }, 900); } + } catch (e) { + alert("Save failed: " + e.message); + } + } + + function fmtCtx(n) { + if (n >= 1e6) return (n / 1e6).toFixed(1).replace(/\.0$/, "") + "M"; + return Math.round(n / 1000) + "k"; + } + function cap(m) { + var c = []; + if (m.contextWindow) c.push(fmtCtx(m.contextWindow)); + if (m.reasoningLevels && m.reasoningLevels.length) c.push("reasoning"); + if (m.supportsImages) c.push("images"); + return c.join(" Β· "); + } + + function render(models) { + if (!models.length) { + content.innerHTML = '

No models available. Enable providers in providers.json ' + + 'and restart the server.

'; + return; + } + var by = group(models); + var table = document.createElement("table"); + table.className = "mtable"; + table.innerHTML = "ShowModelDisplay nameOrder"; + var tbody = document.createElement("tbody"); + + Object.keys(by).sort().forEach(function (prov) { + var head = document.createElement("tr"); + head.innerHTML = '' + prov + ""; + tbody.appendChild(head); + + by[prov].forEach(function (m) { + var tr = document.createElement("tr"); + + var tdToggle = document.createElement("td"); + var toggle = document.createElement("input"); + toggle.type = "checkbox"; toggle.className = "toggle"; toggle.checked = !!m.visible; + tdToggle.appendChild(toggle); + + var tdName = document.createElement("td"); + tdName.innerHTML = '
'; + tdName.children[0].textContent = m.name; + tdName.children[1].textContent = m.ref; + + var tdDisplay = document.createElement("td"); + var nameInput = document.createElement("input"); + nameInput.type = "text"; nameInput.placeholder = m.name; + nameInput.value = m.displayName || ""; + tdDisplay.appendChild(nameInput); + + var tdOrder = document.createElement("td"); + var orderInput = document.createElement("input"); + orderInput.type = "number"; orderInput.value = m.sortOrder || 0; + tdOrder.appendChild(orderInput); + + var tdSaved = document.createElement("td"); + var saved = document.createElement("span"); + saved.className = "saved"; saved.textContent = "saved"; + var capEl = document.createElement("span"); + capEl.className = "cap"; capEl.textContent = cap(m); + tdSaved.appendChild(capEl); tdSaved.appendChild(document.createTextNode(" ")); tdSaved.appendChild(saved); + + toggle.addEventListener("change", function () { patch(m.ref, "visible", toggle.checked, saved); }); + nameInput.addEventListener("change", function () { + patch(m.ref, "displayName", nameInput.value.trim() === "" ? null : nameInput.value.trim(), saved); + }); + orderInput.addEventListener("change", function () { + patch(m.ref, "sortOrder", Number(orderInput.value) || 0, saved); + }); + + tr.appendChild(tdToggle); tr.appendChild(tdName); tr.appendChild(tdDisplay); + tr.appendChild(tdOrder); tr.appendChild(tdSaved); + tbody.appendChild(tr); + }); + }); + + table.appendChild(tbody); + content.innerHTML = ""; + content.appendChild(table); + } + + (async function () { + try { + var res = await fetch("/api/models"); + var body = await res.json(); + render(body.models || []); + } catch (e) { + content.innerHTML = '

Failed to load models: ' + e.message + "

"; + } + })(); +})(); diff --git a/src/drive.ts b/src/drive.ts new file mode 100644 index 0000000..c96d685 --- /dev/null +++ b/src/drive.ts @@ -0,0 +1,82 @@ +import { Store, parseJobParams, type EnqueueParams } from "./store"; +import { ConversationActor } from "./actor"; +import { run } from "./inference"; +import { LEASE_GRACE_MS, HEARTBEAT_INTERVAL_MS } from "./config"; + +/** + * The job drive loop, shared by the server (inline driver) and worker.ts + * (standalone process). Both used to copy-paste this body, so claim/run/ + * checkpoint logic could drift between the two; this is the single canonical + * implementation. + * + * Single-writer per conversation is enforced twice: in SQL + * (`claimExpiredExclusive` refuses a second claim while a lease is live) and + * in-process (`activeRuns` here), since the poll timer can fire again while a + * run is still finishing. + */ +export class JobDriver { + private readonly store: Store; + private readonly getActor: (conversationId: string) => ConversationActor; + /** Conversations with a run in flight in this process. */ + private readonly activeRuns = new Set(); + + constructor(store: Store, getActor: (conversationId: string) => ConversationActor) { + this.store = store; + this.getActor = getActor; + } + + /** + * Claim one job and run it to completion: stream provider steps through the + * conversation actor, advancing the durable checkpoint and lease on each + * flush so a crash mid-run is re-claimed from the last flushed seq. + * Corrupt job params mark the job failed immediately instead of letting it + * sit claimed until the lease expires. + */ + async driveOnce(): Promise { + const row = this.store.claimExpiredExclusive(Date.now()); + if (!row) return; + if (this.activeRuns.has(row.conversation_id)) { + // The previous run is still finishing up; hand the job back so the + // next poll picks it up immediately instead of waiting for the lease + // to expire and the reaper to reclaim it. + this.store.requeue(row.id); + return; + } + + let params: EnqueueParams; + try { + params = parseJobParams(row.params); + } catch { + this.store.markFailed(row.id); + return; + } + + const actor = this.getActor(row.conversation_id); + this.activeRuns.add(row.conversation_id); + // Keep the lease alive independent of delta flushes: a slow first token + // (or a long quiet stretch) must not let the reaper steal a healthy run. + const leaseRefresh = setInterval(() => { + this.store.heartbeat(row.id, Date.now() + LEASE_GRACE_MS); + }, HEARTBEAT_INTERVAL_MS); + try { + await actor.runText( + params.runId, + params.messageId, + (signal) => + run(params.prompt, { runId: params.runId, model: params.model, abortSignal: signal }), + (seq) => { + // Advance the job's durable checkpoint + lease on each flush so a + // crash mid-run is re-claimed from the last flushed seq. + this.store.checkpoint(row.id, seq); + this.store.heartbeat(row.id, Date.now() + LEASE_GRACE_MS); + }, + ); + this.store.markDone(row.id); + } catch { + this.store.markFailed(row.id); + } finally { + clearInterval(leaseRefresh); + this.activeRuns.delete(row.conversation_id); + } + } +} diff --git a/src/events.ts b/src/events.ts index 88ec804..196c3f7 100644 --- a/src/events.ts +++ b/src/events.ts @@ -15,7 +15,6 @@ export const Event = { RunStart: "run-started", RunErr: "run-error", Cancelled: "cancelled", - Steer: "steer", } as const; export type EventName = (typeof Event)[keyof typeof Event]; @@ -87,12 +86,6 @@ export interface CancelledData { runId: string; } -export interface SteerData { - threadId: string; - runId: string; - message: string; -} - export type EventData = | UserMessageData | RunStartedData @@ -102,8 +95,7 @@ export type EventData = | ToolResultData | MessageEndData | RunErrorData - | CancelledData - | SteerData; + | CancelledData; export interface ParsedEvent { id: string; @@ -115,10 +107,6 @@ export function makeId(conversationId: string, seq: number): string { return `${conversationId}:${seq}`; } -export function isEventName(v: unknown): v is EventName { - return typeof v === "string" && (Object.values(Event) as string[]).includes(v); -} - export function parseEventId(id: string): { conversationId: string; seq: number; diff --git a/src/http.ts b/src/http.ts new file mode 100644 index 0000000..86c9334 --- /dev/null +++ b/src/http.ts @@ -0,0 +1,301 @@ +import { randomUUID } from "node:crypto"; +import { Store } from "./store"; +import { ConversationActor, type Subscriber, type WireEvent } from "./actor"; +import { getRegistry } from "./inference"; +import { sseBlock } from "./sse"; +import { parseEventId } from "./events"; +import { withBody } from "./validate"; +import { PromptBody, SteerBody, ModelPatchBody } from "./schemas"; +import { ACTOR_IDLE_TTL_MS, SUBSCRIBER_HEARTBEAT_MS } from "./config"; + +/** + * The web layer as a plain data structure: `apiRoutes(deps)` returns a Bun + * `routes` object, and the entrypoint (server.ts) merges it with the HTML page + * routes and hands it to `Bun.serve`. Keeping this framework-free (Bun's native + * per-method routes + a Standard-Schema `withBody`, no Elysia) means it's the + * same shape production runs and tests exercise against a real ephemeral server. + */ + +/** + * Parses the Last-Event-ID header. The browser sends back the full event id + * (`:`); we extract the seq. Returns 0 (replay all) on + * any parse failure, or when the id belongs to a different conversation β€” an + * id is scoped to its own conversation, so a stale or foreign cursor must + * never skip events here. + */ +function readLastEventId(req: Request, conversationId: string): number { + const raw = req.headers.get("last-event-id"); + if (raw === null) return 0; + try { + const parsed = parseEventId(raw); + return parsed.conversationId === conversationId ? parsed.seq : 0; + } catch { + return 0; + } +} + +/** Sentinel enqueued into an idle stream so proxies keep the connection open. */ +const KEEPALIVE = Symbol("keepalive"); +type StreamItem = WireEvent | typeof KEEPALIVE; + +/** + * Subscriber backed by a ReadableStream: Bun's stream layer pulls strictly + * serially, so a double-delivery is structurally impossible (unlike an async + * iterator whose concurrent pulls race on a shared waiter). The actor `push`es + * each event straight into the controller; backpressure is the platform's job. + */ +class StreamSubscriber implements Subscriber { + closed = false; + readonly stream: ReadableStream; + private controller!: ReadableStreamDefaultController; + private readonly onClose: () => void; + + constructor(onClose: () => void) { + this.onClose = onClose; + this.stream = new ReadableStream({ + start: (controller) => { + this.controller = controller; + }, + cancel: () => { + this.close(); + }, + }); + } + + push(e: WireEvent): void { + if (this.closed) return; + this.controller.enqueue(e); + } + + keepalive(): void { + if (this.closed) return; + this.controller.enqueue(KEEPALIVE); + } + + close(): void { + if (this.closed) return; + this.closed = true; + this.onClose(); + try { + this.controller.close(); + } catch { + // already closed by the consumer + } + } +} + +/** + * Per-conversation actors with idle eviction. Actors with no subscribers and + * no recent activity are removed so the map doesn't grow without bound. + * Subscribed actors are never evicted: an SSE stream is pinned to its actor + * instance, and evicting it would orphan the stream from every future run. + */ +const actors = new Map(); +export function getActor(id: string, store: Store): ConversationActor { + let a = actors.get(id); + if (!a) { + a = new ConversationActor(id, store); + actors.set(id, a); + } + return a; +} + +export function evictIdleActors(): void { + const cutoff = Date.now() - ACTOR_IDLE_TTL_MS; + for (const [id, actor] of actors) { + if (!actor.hasSubscribers() && actor.lastActivity < cutoff) { + actors.delete(id); + } + } +} + +/** Refs of every model this deployment can run (enabled providers + echo). */ +function knownModelRefs(): Set { + return new Set(getRegistry().listModels().map((m) => m.ref)); +} + +/** 422 for unknown model refs, null when known. The gate every model-taking endpoint starts with. */ +function requireKnownModel(ref: string): Response | null { + if (knownModelRefs().has(ref)) return null; + return Response.json({ error: `unknown model "${ref}"` }, { status: 422 }); +} + +/** + * Every available model joined to its curation state, for the settings UI. + * Models with no curation row default to hidden with their catalog name. + */ +function adminModels(store: Store) { + const settings = new Map(store.listModelSettings().map((s) => [s.ref, s])); + return getRegistry() + .listModels() + .map((m) => { + const s = settings.get(m.ref); + return { + ...m, + visible: s?.visible ?? false, + displayName: s?.displayName ?? null, + sortOrder: s?.sortOrder ?? 0, + }; + }); +} + +/** + * The curated subset shown in the chat picker: opt-in (visible only), with + * displayName applied and ordered by sortOrder then name. + */ +function chatModels(store: Store) { + const settings = new Map(store.listModelSettings().map((s) => [s.ref, s])); + return getRegistry() + .listModels() + .flatMap((m) => { + const s = settings.get(m.ref); + if (!s?.visible) return []; + return [{ + ref: m.ref, + name: s.displayName ?? m.name, + contextWindow: m.contextWindow, + reasoningLevels: m.reasoningLevels, + supportsImages: m.supportsImages, + sortOrder: s.sortOrder, + }]; + }) + .sort((a, b) => a.sortOrder - b.sortOrder || a.name.localeCompare(b.name)); +} + +/** The long-lived SSE stream for a conversation (down channel). */ +function openStream(conversationId: string, req: Request, store: Store): Response { + // TODO(auth): cookie/session auth so native EventSource works unmodified. + const actor = getActor(conversationId, store); + const after = readLastEventId(req, conversationId); + + let unsub: () => void = () => {}; + const sub = new StreamSubscriber(() => unsub()); + unsub = actor.subscribe(sub, after); + + // NOTE: this stream is intentionally long-lived. It does not close when a run + // ends; the next run's events arrive on the same connection. A dropped + // connection is what triggers replay via Last-Event-ID. Keepalive comments + // keep proxies (and Bun's idle timeout) from closing the idle connection. + const sseTransform = new TransformStream({ + transform: (item, controller) => { + if (item === KEEPALIVE) controller.enqueue(": keepalive\n\n"); + else controller.enqueue(sseBlock(item)); + }, + }); + const body = sub.stream.pipeThrough(sseTransform); + + const keepAlive = setInterval(() => { + if (sub.closed) { + clearInterval(keepAlive); + return; + } + sub.keepalive(); + }, SUBSCRIBER_HEARTBEAT_MS); + + return new Response(body, { + headers: { + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache, no-transform", + Connection: "keep-alive", + "X-Accel-Buffering": "no", + }, + }); +} + +/** + * Starts a generation run: appends the user message through the actor and + * enqueues a job (with `cancel: true` for /steer, any current run is aborted + * first). Rejects unknown models up front (422): otherwise the job is + * accepted and fails silently in the drive loop at resolve time. The user + * message and the run share `runId` so the client can correlate the + * optimistic echo. + */ +function startRun( + conversationId: string, + data: PromptBody, + opts: { cancel: boolean }, + store: Store, +): Response { + const rejected = requireKnownModel(data.model); + if (rejected) return rejected; + + const actor = getActor(conversationId, store); + if (opts.cancel) actor.requestCancel(); + const runId = data.runId ?? randomUUID(); + const messageId = randomUUID(); + actor.appendUser(data.content, runId); + + const jobId = `${conversationId}:${randomUUID()}`; + store.enqueue(jobId, conversationId, { conversationId, runId, messageId, prompt: data.content, model: data.model }); + + // Enqueue decouples the response from the run: the client opens the stream + // separately and receives message-start/text-deltas/message-end as the drive + // loop (possibly in another process) executes the job. + return Response.json({ jobId, runId, messageId }, { status: 202 }); +} + +/** Partial curation update; `displayName: null` clears an override. */ +function patchModel(data: ModelPatchBody, store: Store): Response { + const rejected = requireKnownModel(data.ref); + if (rejected) return rejected; + const prev = store.getModelSetting(data.ref); + const merged = { + ref: data.ref, + visible: data.visible ?? prev?.visible ?? false, + displayName: data.displayName !== undefined ? data.displayName : (prev?.displayName ?? null), + sortOrder: data.sortOrder ?? prev?.sortOrder ?? 0, + }; + store.setModelSetting(merged); + return Response.json(merged); +} + +/** + * The API routes as a Bun `routes` object. Per-method handlers give free method + * dispatch + `req.params`; `withBody` validates + types JSON bodies before the + * handler runs. Everything dynamic lives under `/api/` so a future service + * worker can bypass it cleanly and cache only the static shell. + */ +export function apiRoutes(deps: { store: Store }) { + const { store } = deps; + return { + "/health": { GET: () => Response.json({ ok: true }) }, + + "/api/conversations": { + GET: () => Response.json({ conversations: store.listConversations() }), + }, + + // Settings/admin view: every available model + its curation state. + "/api/models": { + GET: () => Response.json({ models: adminModels(store) }), + PATCH: withBody(ModelPatchBody, (data) => patchModel(data, store)), + }, + + // Chat view: the curated, opt-in subset, ordered for the picker. + "/api/models/chat": { + GET: () => Response.json({ models: chatModels(store) }), + }, + + "/api/conversations/:id/stream": { + GET: (req: Bun.BunRequest<"/api/conversations/:id/stream">) => + openStream(req.params.id, req, store), + }, + "/api/conversations/:id/prompt": { + POST: withBody(PromptBody, (data, req: Bun.BunRequest<"/api/conversations/:id/prompt">) => + startRun(req.params.id, data, { cancel: false }, store)), + }, + "/api/conversations/:id/cancel": { + POST: (req: Bun.BunRequest<"/api/conversations/:id/cancel">) => { + getActor(req.params.id, store).requestCancel(); + return Response.json({ ok: true }); + }, + }, + "/api/conversations/:id/steer": { + POST: withBody(SteerBody, (data, req: Bun.BunRequest<"/api/conversations/:id/steer">) => + startRun(req.params.id, data, { cancel: true }, store)), + }, + "/api/conversations/:id/events": { + GET: (req: Bun.BunRequest<"/api/conversations/:id/events">) => + Response.json(getActor(req.params.id, store).replay(0)), + }, + }; +} diff --git a/src/schemas.ts b/src/schemas.ts new file mode 100644 index 0000000..92b8ffa --- /dev/null +++ b/src/schemas.ts @@ -0,0 +1,36 @@ +import * as v from "valibot"; + +/** + * Request-body schemas (valibot). + * + * `v.object` ignores unknown keys (lenient like the current API); + * `v.minLength(1)` on `model`/`ref` reproduces the "required, non-empty" rule + * that returns 422 today. + */ + +/** POST /api/conversations/:id/prompt */ +export const PromptBody = v.object({ + content: v.string(), + model: v.pipe(v.string(), v.minLength(1, "model is required")), + runId: v.optional(v.string()), +}); +export type PromptBody = v.InferOutput; + +/** POST /api/conversations/:id/steer */ +export const SteerBody = v.object({ + content: v.string(), + model: v.pipe(v.string(), v.minLength(1, "model is required")), +}); +export type SteerBody = v.InferOutput; + +/** + * PATCH /api/models β€” partial curation update. `displayName: null` clears the + * override; omitted fields keep their stored value (merge happens in the handler). + */ +export const ModelPatchBody = v.object({ + ref: v.pipe(v.string(), v.minLength(1)), + visible: v.optional(v.boolean()), + displayName: v.optional(v.nullable(v.string())), + sortOrder: v.optional(v.number()), +}); +export type ModelPatchBody = v.InferOutput; diff --git a/src/sse.ts b/src/sse.ts index 38d2b03..9dcedd2 100644 --- a/src/sse.ts +++ b/src/sse.ts @@ -16,16 +16,3 @@ export function truncateUtf8(s: string, n: number = MAX_SSE_FIELD_BYTES): string export function sseBlock(e: { id: string; event: string; data: unknown }): string { return `event: ${e.event}\nid: ${e.id}\ndata: ${JSON.stringify(e.data)}\n\n`; } - -/** - * Wraps a feed of typed events as SSE-formatted strings. The feed must - * implement the async-iterator protocol with strict sequential `next()` calls - * (e.g. a ReadableStream's reader guarantees this). - */ -export async function* sseBlocks( - events: AsyncIterable<{ id: string; event: string; data: unknown }>, -): AsyncGenerator { - for await (const e of events) { - yield sseBlock(e); - } -} diff --git a/src/store.ts b/src/store.ts index 6d74a26..ec45c22 100644 --- a/src/store.ts +++ b/src/store.ts @@ -19,12 +19,36 @@ export interface EnqueueParams { model: string; } +/** + * Decodes a job's params blob β€” the inverse of `enqueue`, and the only place + * job params are parsed. Throws on a corrupt or incomplete row; callers mark + * the job failed rather than re-claiming it forever. + */ +export function parseJobParams(params: string): EnqueueParams { + const p = JSON.parse(params) as Record; + for (const key of ["conversationId", "runId", "messageId", "prompt", "model"] as const) { + if (typeof p[key] !== "string") { + throw new Error(`malformed job params: bad or missing "${key}"`); + } + } + return p as unknown as EnqueueParams; +} + export interface StoredEvent { seq: number; event: string; data: unknown; } +/** A conversation as shown in the chat rail: id, age, and a derived title. */ +export interface ConversationSummary { + id: string; + createdAt: number; + lastSeq: number; + /** First user message, truncated β€” null for a conversation with no prompt yet. */ + title: string | null; +} + /** * Curation of a model for the chat UI. Opt-in: a model with no row (or * `visible: false`) is hidden from the chat picker. `displayName` overrides the @@ -105,11 +129,13 @@ export class Store { private insertConversationStmt: ReturnType; private upsertConversationStmt: ReturnType; private enqueueStmt: ReturnType; - private claimStmt: ReturnType; + private claimExclusiveStmt: ReturnType; private heartbeatStmt: ReturnType; private checkpointStmt: ReturnType; private reapStmt: ReturnType; + private requeueStmt: ReturnType; private finishStmt: ReturnType; + private listConversationsStmt: ReturnType; private listSettingsStmt: ReturnType; private getSettingStmt: ReturnType; private upsertSettingStmt: ReturnType; @@ -145,10 +171,24 @@ export class Store { this.enqueueStmt = this.db.prepare( `INSERT INTO jobs (id, conversation_id, status, params) VALUES (?, ?, 'queued', ?)`, ); - this.claimStmt = this.db.prepare( + // The hot claim query: queued or expired-lease jobs, but only for + // conversations with no other live running job. Enforces the single-writer + // invariant (one active run per conversation) atomically in SQL. + this.claimExclusiveStmt = this.db.prepare( `UPDATE jobs SET status = 'running', lease_until = ? - WHERE id = (SELECT id FROM jobs WHERE status = 'queued' ORDER BY id LIMIT 1) + WHERE id = ( + SELECT j.id FROM jobs j + WHERE (j.status = 'queued' + OR (j.status = 'running' AND j.lease_until < ?)) + AND NOT EXISTS ( + SELECT 1 FROM jobs j2 + WHERE j2.conversation_id = j.conversation_id + AND j2.status = 'running' + AND j2.lease_until >= ? + ) + ORDER BY j.id LIMIT 1 + ) RETURNING id, conversation_id, status, lease_until, checkpoint_seq, params`, ); this.heartbeatStmt = this.db.prepare( @@ -160,10 +200,24 @@ export class Store { this.reapStmt = this.db.prepare( `UPDATE jobs SET status = 'queued' WHERE status = 'running' AND lease_until < ?`, ); + this.requeueStmt = this.db.prepare( + `UPDATE jobs SET status = 'queued', lease_until = 0 WHERE id = ?`, + ); this.finishStmt = this.db.prepare( `UPDATE jobs SET status = ?, lease_until = 0 WHERE id = ?`, ); + this.listConversationsStmt = this.db.prepare( + `SELECT c.id AS id, c.created_at AS created_at, c.last_seq AS last_seq, + (SELECT e.data FROM events e + WHERE e.conversation_id = c.id AND e.event = 'user-message' + ORDER BY e.seq ASC LIMIT 1) AS first_user, + COALESCE((SELECT MAX(e.created_at) FROM events e + WHERE e.conversation_id = c.id), c.created_at) AS last_activity + FROM conversations c + ORDER BY last_activity DESC`, + ); + this.listSettingsStmt = this.db.prepare( `SELECT model_ref, visible, display_name, sort_order FROM model_settings`, ); @@ -181,6 +235,33 @@ export class Store { ); } + /** + * All conversations, most recently active first (by newest event, so a + * conversation that gets activity again rises to the top), each with a title + * derived from its first user message. Used by the chat rail. The title + * subquery pulls the earliest `user-message` event's content per conversation. + */ + listConversations(): ConversationSummary[] { + const rows = this.listConversationsStmt.all() as Array<{ + id: string; + created_at: number; + last_seq: number; + first_user: string | null; + }>; + return rows.map((r) => { + let title: string | null = null; + if (r.first_user) { + try { + const content = (JSON.parse(r.first_user) as { content?: string }).content; + if (typeof content === "string") title = content.slice(0, 80); + } catch { + // malformed row: leave title null rather than crash the list + } + } + return { id: r.id, createdAt: r.created_at, lastSeq: r.last_seq, title }; + }); + } + /** All curation rows (models with no row are hidden by default). */ listModelSettings(): ModelSetting[] { return (this.listSettingsStmt.all() as ModelSettingRow[]).map(rowToSetting); @@ -221,28 +302,6 @@ export class Store { })(); } - append( - conversationId: string, - seq: number, - eventName: string, - data: unknown, - ): void { - const id = `${conversationId}:${seq}`; - this.insertEventStmt.run( - id, - conversationId, - seq, - eventName, - JSON.stringify(data), - Date.now(), - ); - } - - /** Marks a delta against `conversation_id`; called when a delta batch flushes. */ - bumpSeq(conversationId: string, seq: number): void { - this.upsertConversationStmt.run(conversationId, Date.now(), seq); - } - /** The last seq durable for a conversation (0 if none). */ lastSeq(conversationId: string): number { const row = this.db @@ -265,53 +324,14 @@ export class Store { this.enqueueStmt.run(jobId, conversationId, JSON.stringify(params)); } - /** Atomic race-free claim; returns the job or null. */ - claim(now: number): JobRow | null { - return this.claimStmt.get(now) as JobRow | null; - } - /** * Returns a job that needs running (queued, or an expired running lease the - * original worker may have died with), claiming it atomically. - */ - claimExpired(now: number): JobRow | null { - return this.db - .prepare( - `UPDATE jobs - SET status = 'running', lease_until = ? - WHERE id = (SELECT id FROM jobs - WHERE status = 'queued' - OR (status = 'running' AND lease_until < ?) - ORDER BY id LIMIT 1) - RETURNING id, conversation_id, status, lease_until, checkpoint_seq, params`, - ) - .get(now, now) as JobRow | null; - } - - /** - * Like claimExpired, but only for conversations with no other running job. - * Enforces the single-writer invariant: one active run per conversation. + * original worker may have died with) for a conversation with no other live + * run, claiming it atomically. The returned row is claimed (status + * 'running') β€” the caller owns it until markDone/markFailed. */ claimExpiredExclusive(now: number): JobRow | null { - return this.db - .prepare( - `UPDATE jobs - SET status = 'running', lease_until = ? - WHERE id = ( - SELECT j.id FROM jobs j - WHERE (j.status = 'queued' - OR (j.status = 'running' AND j.lease_until < ?)) - AND NOT EXISTS ( - SELECT 1 FROM jobs j2 - WHERE j2.conversation_id = j.conversation_id - AND j2.status = 'running' - AND j2.lease_until >= ? - ) - ORDER BY j.id LIMIT 1 - ) - RETURNING id, conversation_id, status, lease_until, checkpoint_seq, params`, - ) - .get(now, now, now) as JobRow | null; + return this.claimExclusiveStmt.get(now, now, now) as JobRow | null; } heartbeat(id: string, leaseUntil: number): void { @@ -327,6 +347,11 @@ export class Store { return this.reapStmt.run(now).changes; } + /** Voluntarily hand a claimed job back without running it. */ + requeue(id: string): void { + this.requeueStmt.run(id); + } + markDone(id: string): void { this.finishStmt.run("done", id); } diff --git a/src/validate.ts b/src/validate.ts new file mode 100644 index 0000000..70af87a --- /dev/null +++ b/src/validate.ts @@ -0,0 +1,58 @@ +import type { StandardSchemaV1 } from "@standard-schema/spec"; + +/** + * Framework-agnostic request-body validation, built on the Standard Schema + * interface (implemented by valibot, zod v4, arktype, typebox...). Nothing here + * imports valibot directly β€” swapping validators never touches this file. + * + * Designed for Bun's per-method route handlers (`{ POST: withBody(Schema, fn) }`), + * which already give method dispatch, `req.params`, and automatic 405s β€” so all + * that's left is shaping + typing the JSON body, which is what this does. + */ + +export type Validated = + | { ok: true; value: T } + | { ok: false; status: number; error: string; issues: readonly StandardSchemaV1.Issue[] }; + +/** Validate an already-parsed value against a schema. 422 with issues on failure. */ +export async function validate( + schema: S, + input: unknown, +): Promise>> { + let result = schema["~standard"].validate(input); + if (result instanceof Promise) result = await result; + if (result.issues) { + return { ok: false, status: 422, error: "validation failed", issues: result.issues }; + } + return { ok: true, value: result.value }; +} + +/** The minimal request surface we need: a JSON body reader. */ +interface JsonRequest { + json(): Promise; +} + +/** + * Wraps a route handler so the JSON body is read and validated first: 400 on + * unparseable JSON, 422 (with Standard Schema issues) on a shape mismatch, and + * only then the handler runs with fully-typed, validated data. `req` is passed + * through untouched so the handler still sees `req.params`, headers, etc. + */ +export function withBody( + schema: S, + handler: (data: StandardSchemaV1.InferOutput, req: R) => Response | Promise, +): (req: R) => Promise { + return async (req) => { + let raw: unknown; + try { + raw = await req.json(); + } catch { + return Response.json({ error: "invalid JSON body" }, { status: 400 }); + } + const result = await validate(schema, raw); + if (!result.ok) { + return Response.json({ error: result.error, issues: result.issues }, { status: result.status }); + } + return handler(result.value, req); + }; +} diff --git a/tests/curation.test.ts b/tests/curation.test.ts index 7613a76..99fd3ed 100644 --- a/tests/curation.test.ts +++ b/tests/curation.test.ts @@ -1,15 +1,13 @@ -import { test, expect, beforeEach } from "bun:test"; -import { mkdtempSync } from "node:fs"; +import { test, expect, beforeEach, afterAll } from "bun:test"; +import { mkdtempSync, rmSync } from "node:fs"; import { join } from "node:path"; import { tmpdir } from "node:os"; import { Store } from "../src/store"; -import { buildApp } from "../server"; +import { apiRoutes } from "../src/http"; import { ProviderRegistry } from "../src/providers"; import { setRegistry } from "../src/inference"; import { Catalog } from "../src/catalog"; -const base = "http://localhost"; - function fixtureRegistry(): ProviderRegistry { const catalog = Catalog.fromRaw([ { @@ -28,52 +26,57 @@ function fixtureRegistry(): ProviderRegistry { }); } +// Each test gets a fresh store behind a real ephemeral-port server. +const servers: Array> = []; +const tmpDirs: string[] = []; function freshApp() { const tmp = mkdtempSync(join(tmpdir(), "kloe-cur-")); + tmpDirs.push(tmp); const store = new Store(join(tmp, "test.db")); - const app = buildApp({ store }); - return { app, store }; + const server = Bun.serve({ port: 0, routes: apiRoutes({ store }) }); + servers.push(server); + return { base: server.url.origin, store }; } +afterAll(() => { + for (const s of servers) s.stop(true); + for (const d of tmpDirs) rmSync(d, { recursive: true, force: true }); +}); async function json(res: Response): Promise { return JSON.parse(await res.text()); } +const patchModels = (base: string, bodyObj: unknown) => + fetch(`${base}/api/models`, { + method: "PATCH", + headers: { "content-type": "application/json" }, + body: JSON.stringify(bodyObj), + }); + beforeEach(() => { setRegistry(fixtureRegistry()); }); -test("GET /models returns all models hidden by default (opt-in)", async () => { - const { app } = freshApp(); - const body = await json(await app.handle(new Request(`${base}/models`))); +test("GET /api/models returns all models hidden by default (opt-in)", async () => { + const { base } = freshApp(); + const body = await json(await fetch(`${base}/api/models`)); const refs = body.models.map((m: any) => m.ref); expect(refs).toContain("acme/acme-1"); expect(refs).toContain("echo"); - // Nothing curated yet β†’ everything hidden. expect(body.models.every((m: any) => m.visible === false)).toBe(true); }); -test("GET /models/chat is empty until a model is made visible", async () => { - const { app } = freshApp(); - const body = await json(await app.handle(new Request(`${base}/models/chat`))); +test("GET /api/models/chat is empty until a model is made visible", async () => { + const { base } = freshApp(); + const body = await json(await fetch(`${base}/api/models/chat`)); expect(body.models).toEqual([]); }); -test("PATCH /models makes a model visible and renames it; chat reflects it", async () => { - const { app } = freshApp(); +test("PATCH /api/models makes a model visible and renames it; chat reflects it", async () => { + const { base } = freshApp(); const patched = await json( - await app.handle( - new Request(`${base}/models`, { - method: "PATCH", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ - ref: "acme/acme-1", - visible: true, - displayName: "Acme (fast)", - }), - }), - ), + await patchModels(base, { ref: "acme/acme-1", visible: true, displayName: "Acme (fast)" }), ); expect(patched).toMatchObject({ ref: "acme/acme-1", @@ -81,7 +84,7 @@ test("PATCH /models makes a model visible and renames it; chat reflects it", asy displayName: "Acme (fast)", }); - const chat = await json(await app.handle(new Request(`${base}/models/chat`))); + const chat = await json(await fetch(`${base}/api/models/chat`)); expect(chat.models).toHaveLength(1); expect(chat.models[0]).toMatchObject({ ref: "acme/acme-1", @@ -91,74 +94,39 @@ test("PATCH /models makes a model visible and renames it; chat reflects it", asy }); test("PATCH is partial: a second patch keeps prior fields", async () => { - const { app } = freshApp(); - const patch = (bodyObj: unknown) => - app.handle( - new Request(`${base}/models`, { - method: "PATCH", - headers: { "content-type": "application/json" }, - body: JSON.stringify(bodyObj), - }), - ); - - await patch({ ref: "acme/acme-1", visible: true, displayName: "Kept" }); - // Only change sortOrder; visibility + name must persist. - const result = await json(await patch({ ref: "acme/acme-1", sortOrder: 5 })); - expect(result).toMatchObject({ - visible: true, - displayName: "Kept", - sortOrder: 5, - }); + const { base } = freshApp(); + await patchModels(base, { ref: "acme/acme-1", visible: true, displayName: "Kept" }); + const result = await json(await patchModels(base, { ref: "acme/acme-1", sortOrder: 5 })); + expect(result).toMatchObject({ visible: true, displayName: "Kept", sortOrder: 5 }); }); test("PATCH displayName:null clears the override", async () => { - const { app } = freshApp(); - const patch = (bodyObj: unknown) => - app.handle( - new Request(`${base}/models`, { - method: "PATCH", - headers: { "content-type": "application/json" }, - body: JSON.stringify(bodyObj), - }), - ); - await patch({ ref: "acme/acme-1", visible: true, displayName: "Temp" }); - const cleared = await json(await patch({ ref: "acme/acme-1", displayName: null })); + const { base } = freshApp(); + await patchModels(base, { ref: "acme/acme-1", visible: true, displayName: "Temp" }); + const cleared = await json(await patchModels(base, { ref: "acme/acme-1", displayName: null })); expect(cleared.displayName).toBeNull(); - const chat = await json(await app.handle(new Request(`${base}/models/chat`))); + const chat = await json(await fetch(`${base}/api/models/chat`)); expect(chat.models[0].name).toBe("Acme One"); // back to catalog name }); test("chat models are ordered by sortOrder then name", async () => { - const { app } = freshApp(); - const patch = (bodyObj: unknown) => - app.handle( - new Request(`${base}/models`, { - method: "PATCH", - headers: { "content-type": "application/json" }, - body: JSON.stringify(bodyObj), - }), - ); - await patch({ ref: "acme/acme-1", visible: true, sortOrder: 10 }); - await patch({ ref: "acme/acme-2", visible: true, sortOrder: 1 }); - const chat = await json(await app.handle(new Request(`${base}/models/chat`))); + const { base } = freshApp(); + await patchModels(base, { ref: "acme/acme-1", visible: true, sortOrder: 10 }); + await patchModels(base, { ref: "acme/acme-2", visible: true, sortOrder: 1 }); + const chat = await json(await fetch(`${base}/api/models/chat`)); expect(chat.models.map((m: any) => m.ref)).toEqual(["acme/acme-2", "acme/acme-1"]); }); test("PATCH rejects an unknown model ref with 422", async () => { - const { app } = freshApp(); - const res = await app.handle( - new Request(`${base}/models`, { - method: "PATCH", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ ref: "acme/ghost", visible: true }), - }), - ); + const { base } = freshApp(); + const res = await patchModels(base, { ref: "acme/ghost", visible: true }); expect(res.status).toBe(422); }); test("curation persists across Store re-open", async () => { const tmp = mkdtempSync(join(tmpdir(), "kloe-cur-persist-")); + tmpDirs.push(tmp); const dbPath = join(tmp, "test.db"); const store1 = new Store(dbPath); store1.setModelSetting({ diff --git a/tests/server.test.ts b/tests/server.test.ts index f254c4e..9839639 100644 --- a/tests/server.test.ts +++ b/tests/server.test.ts @@ -1,17 +1,31 @@ -import { test, expect, beforeEach } from "bun:test"; +import { test, expect, beforeAll, afterAll, beforeEach } from "bun:test"; import { mkdtempSync, rmSync } from "node:fs"; import { join } from "node:path"; import { tmpdir } from "node:os"; -import { Store } from "../src/store"; -import { buildApp, getActor } from "../server"; +import { Store, parseJobParams } from "../src/store"; +import { apiRoutes, getActor, evictIdleActors } from "../src/http"; +import { JobDriver } from "../src/drive"; import { setRegistry } from "../src/inference"; import { ProviderRegistry } from "../src/providers"; import { Catalog } from "../src/catalog"; const tmp = mkdtempSync(join(tmpdir(), "kloe-srv-")); const store = new Store(join(tmp, "test.db")); -const app = buildApp({ store }); -const base = "http://localhost"; + +// A real Bun server on an ephemeral port β€” the same routing/validation path +// production uses. `apiRoutes` carries no HTML routes, so starting it never +// triggers frontend bundling. +let server: ReturnType; +let base: string; +beforeAll(() => { + server = Bun.serve({ port: 0, routes: apiRoutes({ store }) }); + base = server.url.origin; +}); +afterAll(() => { + server.stop(true); + store.db.close(); + rmSync(tmp, { recursive: true, force: true }); +}); // A minimal registry so model validation resolves (echo is always known). // Set per-test to survive interleaving with other files' registry mutations. @@ -27,28 +41,26 @@ interface Frame { /** * Reads an SSE response until `until` returns true (or EOF), returning the - * frames seen. Leaves the reader open past the stop point, then cancels. + * frames seen. Cancels the reader as soon as the condition is met. */ async function readSse( res: Response, until: (frames: Frame[]) => boolean, ): Promise { const reader = res.body!.getReader(); + const decoder = new TextDecoder(); let buffer = ""; const frames: Frame[] = []; while (true) { const { done, value } = await reader.read(); if (done) break; - // The body is a TransformStream; chunks are strings. - buffer += typeof value === "string" ? value : new TextDecoder().decode(value); - // Parse complete SSE blocks (each ends with a blank line). + buffer += typeof value === "string" ? value : decoder.decode(value); let idx; while ((idx = buffer.indexOf("\n\n")) !== -1) { const block = buffer.slice(0, idx); buffer = buffer.slice(idx + 2); if (!block.trim()) continue; - // Skip SSE comments (keepalive). - if (block.startsWith(":")) continue; + if (block.startsWith(":")) continue; // keepalive comment let event = "message"; let id = ""; const dataLines: string[] = []; @@ -68,59 +80,71 @@ async function readSse( } test("prompt rejects an unknown model with 422 (no silent async failure)", async () => { - const res = await app.handle( - new Request(`${base}/conversations/badmodel/prompt`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ content: "hi", model: "openai/gpt-4" }), - }), - ); + const res = await fetch(`${base}/api/conversations/badmodel/prompt`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: "hi", model: "openai/gpt-4" }), + }); expect(res.status).toBe(422); - // Nothing should have been enqueued for that conversation. + // Nothing should have been enqueued/appended for that conversation. expect(getActor("badmodel", store).replay(0).length).toBe(0); }); -test("prompt rejects an empty model string with 422", async () => { - const res = await app.handle( - new Request(`${base}/conversations/emptymodel/prompt`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ content: "hi", model: "" }), - }), - ); +test("prompt rejects an empty model string with 422 (schema minLength)", async () => { + const res = await fetch(`${base}/api/conversations/emptymodel/prompt`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: "hi", model: "" }), + }); expect(res.status).toBe(422); }); +test("prompt rejects a malformed JSON body with 400", async () => { + const res = await fetch(`${base}/api/conversations/badjson/prompt`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{not json", + }); + expect(res.status).toBe(400); +}); + test("steer rejects an unknown model with 422", async () => { - const res = await app.handle( - new Request(`${base}/conversations/steerbad/steer`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ content: "go", model: "nope/nope" }), - }), - ); + const res = await fetch(`${base}/api/conversations/steerbad/steer`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: "go", model: "nope/nope" }), + }); expect(res.status).toBe(422); }); +test("GET /api/conversations lists conversations newest-first with a derived title", async () => { + const actor = getActor("conv-list-1", store); + actor.appendUser("What is the meaning of it all?", "r-list-1"); + + const res = await fetch(`${base}/api/conversations`); + expect(res.status).toBe(200); + const body = (await res.json()) as { + conversations: Array<{ id: string; title: string | null; lastSeq: number }>; + }; + const row = body.conversations.find((c) => c.id === "conv-list-1"); + expect(row).toBeDefined(); + expect(row!.title).toBe("What is the meaning of it all?"); + expect(row!.lastSeq).toBeGreaterThan(0); +}); + test("prompt β†’ SSE stream emits user-message, message-start, deltas, message-end", async () => { const conv = "s1"; - const res = await app.handle( - new Request(`${base}/conversations/${conv}/prompt`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ content: "hello kloe", model: "echo" }), - }), - ); + const res = await fetch(`${base}/api/conversations/${conv}/prompt`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: "hello kloe", model: "echo" }), + }); expect(res.status).toBe(202); // Claim the queued job (as the drive loop would) and run it inline. const row = store.claimExpiredExclusive(Date.now()); expect(row).not.toBeNull(); - const params = JSON.parse(row!.params) as { - runId: string; - messageId: string; - prompt: string; - }; + const params = parseJobParams(row!.params); const actor = getActor(conv, store); await actor.runText(params.runId, params.messageId, async function* (_signal) { yield { kind: "text", chunk: `echo: ${params.prompt}` }; @@ -128,11 +152,9 @@ test("prompt β†’ SSE stream emits user-message, message-start, deltas, message-e store.markDone(row!.id); // Open the stream after the run: it replays the durable log. - const streamRes = await app.handle(new Request(`${base}/conversations/${conv}/stream`)); + const streamRes = await fetch(`${base}/api/conversations/${conv}/stream`); expect(streamRes.headers.get("content-type")).toBe("text/event-stream"); - const frames = await readSse(streamRes, (f) => - f.some((x) => x.event === "message-end"), - ); + const frames = await readSse(streamRes, (f) => f.some((x) => x.event === "message-end")); const events = frames.map((f) => f.event); expect(events).toContain("user-message"); expect(events).toContain("message-start"); @@ -140,7 +162,6 @@ test("prompt β†’ SSE stream emits user-message, message-start, deltas, message-e expect(events).toContain("message-end"); const delta = frames.find((f) => f.event === "text-delta"); expect((delta!.data as { delta: string }).delta).toBe("echo: hello kloe"); - // Ids monotonic. const seqs = frames.map((f) => Number(f.id.split(":")[1]!)); for (let i = 1; i < seqs.length; i++) expect(seqs[i]!).toBeGreaterThan(seqs[i - 1]!); }); @@ -154,61 +175,128 @@ test("resume via HTTP Last-Event-ID replays only the gap", async () => { // Cursor 0: full replay over HTTP. const full = await readSse( - await app.handle(new Request(`${base}/conversations/${conv}/stream`)), + await fetch(`${base}/api/conversations/${conv}/stream`), (f) => f.some((x) => x.event === "message-end"), ); const lastSeq = Math.max(...full.map((f) => Number(f.id.split(":")[1]!))); - // Resume over HTTP with Last-Event-ID = lastSeq. The server must parse the - // `:` header and replay strictly after it. With nothing new, - // the stream stays open (no message-end arrives), so we read with a short - // timeout and assert we got zero real frames. - const resumeRes = await app.handle( - new Request(`${base}/conversations/${conv}/stream`, { + // Append a new event after the full replay, then resume from lastSeq: the + // stream must deliver ONLY that gap event, never re-sending the earlier ones. + // (Deterministic β€” no timing race on an idle connection.) + actor.appendUser("follow-up", "r3"); + const gap = await readSse( + await fetch(`${base}/api/conversations/${conv}/stream`, { headers: { "last-event-id": `${conv}:${lastSeq}` }, }), + (f) => f.length >= 1, ); - const reader = resumeRes.body!.getReader(); - let buffer = ""; - const resumed: Frame[] = []; - const deadline = Bun.sleep(300); - const readLoop = (async () => { - while (true) { - const { done, value } = await reader.read(); - if (done) break; - buffer += typeof value === "string" ? value : new TextDecoder().decode(value); - let idx; - while ((idx = buffer.indexOf("\n\n")) !== -1) { - const block = buffer.slice(0, idx); - buffer = buffer.slice(idx + 2); - if (!block.trim() || block.startsWith(":")) continue; - let event = "message"; - let id = ""; - const dataLines: string[] = []; - for (const line of block.split("\n")) { - if (line.startsWith("event:")) event = line.slice(6).trim(); - else if (line.startsWith("id:")) id = line.slice(3).trim(); - else if (line.startsWith("data:")) dataLines.push(line.slice(5).trim()); - } - resumed.push({ event, id, data: JSON.parse(dataLines.join("\n")) }); - } - } - })(); - await Promise.race([readLoop, deadline]); - await reader.cancel(); - // Nothing new to replay: no real frames (only keepalive comments, skipped). - expect(resumed.length).toBe(0); + expect(gap.length).toBeGreaterThan(0); + expect(gap.every((f) => Number(f.id.split(":")[1]!) > lastSeq)).toBe(true); + expect(gap.some((f) => f.event === "user-message")).toBe(true); // An intermediate cursor must replay the missing middle strictly. const firstSeq = Math.min(...full.map((f) => Number(f.id.split(":")[1]!))); - const midRes = await app.handle( - new Request(`${base}/conversations/${conv}/stream`, { - headers: { "last-event-id": `${conv}:${firstSeq}` }, - }), - ); + const midRes = await fetch(`${base}/api/conversations/${conv}/stream`, { + headers: { "last-event-id": `${conv}:${firstSeq}` }, + }); const mid = await readSse(midRes, (f) => f.some((x) => x.event === "message-end")); expect(mid.length).toBeGreaterThan(0); for (const e of mid) { expect(Number(e.id.split(":")[1]!)).toBeGreaterThan(firstSeq); } }); + +test("a Last-Event-ID from a different conversation replays from the start", async () => { + const conv = "s3"; + const actor = getActor(conv, store); + actor.appendUser("first turn", "r4"); + + // A cursor claiming to belong to another conversation must be treated as + // "no cursor" β€” replaying strictly after a foreign seq could skip events. + const frames = await readSse( + await fetch(`${base}/api/conversations/${conv}/stream`, { + headers: { "last-event-id": "some-other-conv:42" }, + }), + (f) => f.length >= 1, + ); + expect(frames.length).toBeGreaterThan(0); + expect(frames[0]!.id).toBe(`${conv}:1`); // full replay from the start +}); + +test("eviction skips actors with live subscribers", async () => { + const conv = "evict-me"; + const actor = getActor(conv, store); + + // Idle + no subscribers: evictable. + actor.lastActivity = 0; + evictIdleActors(); + expect(getActor(conv, store)).not.toBe(actor); + + // Idle + a live subscriber: pinned. Evicting would orphan the stream, so + // the map must keep the instance and getActor returns the same object. + const actor2 = getActor(conv, store); + actor2.lastActivity = 0; + const unsub = actor2.follow({ push: () => {}, closed: false }); + evictIdleActors(); + expect(getActor(conv, store)).toBe(actor2); + unsub(); +}); + +test("JobDriver runs a queued job end-to-end (claim β†’ run β†’ done)", async () => { + const conv = "drive-1"; + const res = await fetch(`${base}/api/conversations/${conv}/prompt`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ content: "drive me", model: "echo" }), + }); + expect(res.status).toBe(202); + + const driver = new JobDriver(store, (id) => getActor(id, store)); + await driver.driveOnce(); + + const replay = getActor(conv, store).replay(0); + expect(replay.some((e) => e.event === "message-end")).toBe(true); + // Second call finds nothing claimable. + await driver.driveOnce(); +}); + +test("JobDriver marks a job with corrupt params failed instead of re-claiming it forever", async () => { + const conv = "drive-corrupt"; + const jobId = `${conv}:bad-params`; + store.enqueue(jobId, conv, { + conversationId: conv, + runId: "r-c", + messageId: "m-c", + prompt: "hi", + model: "echo", + }); + // Clobber the durable row so it looks like a torn/half-written insert. + store.db + .prepare("UPDATE jobs SET params = ? WHERE id = ?") + .run(JSON.stringify({ conversationId: conv }), jobId); + + const driver = new JobDriver(store, (id) => getActor(id, store)); + await driver.driveOnce(); + + const row = store.db.prepare("SELECT status FROM jobs WHERE id = ?").get(jobId) as { status: string }; + expect(row.status).toBe("failed"); +}); + +test("GET /api/conversations orders by most recent activity, not creation", async () => { + // Two conversations; the older one gets new activity afterwards and must + // rise above the newer one. + const oldActor = getActor("conv-old", store); + oldActor.appendUser("created first", "r-old"); + const newerActor = getActor("conv-newer", store); + newerActor.appendUser("created second", "r-newer"); + + // Wait a tick so created_at differs, then bump the older conversation. + await Bun.sleep(5); + oldActor.appendUser("still going", "r-old-2"); + + const body = (await (await fetch(`${base}/api/conversations`)).json()) as { + conversations: Array<{ id: string }>; + }; + const ids = body.conversations.map((c) => c.id); + expect(ids.indexOf("conv-old")).toBeLessThan(ids.indexOf("conv-newer")); +}); diff --git a/tests/validation.test.ts b/tests/validation.test.ts new file mode 100644 index 0000000..b7b0d0b --- /dev/null +++ b/tests/validation.test.ts @@ -0,0 +1,69 @@ +import { test, expect } from "bun:test"; +import { validate, withBody } from "../src/validate"; +import { PromptBody, ModelPatchBody } from "../src/schemas"; + +test("validate accepts a well-formed prompt body and infers the type", async () => { + const r = await validate(PromptBody, { content: "hi", model: "openai/gpt-4" }); + expect(r.ok).toBe(true); + if (r.ok) { + expect(r.value.model).toBe("openai/gpt-4"); + expect(r.value.runId).toBeUndefined(); + } +}); + +test("validate rejects a missing model with 422 + issues", async () => { + const r = await validate(PromptBody, { content: "hi" }); + expect(r.ok).toBe(false); + if (!r.ok) { + expect(r.status).toBe(422); + expect(r.issues.length).toBeGreaterThan(0); + } +}); + +test("validate rejects an empty model string (minLength)", async () => { + const r = await validate(PromptBody, { content: "hi", model: "" }); + expect(r.ok).toBe(false); +}); + +test("validate rejects a wrong-typed field", async () => { + const r = await validate(ModelPatchBody, { ref: "echo", visible: "yes" }); + expect(r.ok).toBe(false); +}); + +test("ModelPatchBody allows partial fields and an explicit null displayName", async () => { + const r = await validate(ModelPatchBody, { ref: "echo", displayName: null }); + expect(r.ok).toBe(true); + if (r.ok) { + expect(r.value.displayName).toBeNull(); + expect(r.value.visible).toBeUndefined(); + } +}); + +// --- withBody: the Bun per-method route wrapper ------------------------------- +function fakeReq(body: unknown, badJson = false) { + return { json: async () => { if (badJson) throw new Error("bad json"); return body; } }; +} + +test("withBody passes validated, typed data to the handler (200)", async () => { + const handler = withBody(PromptBody, (data) => Response.json({ model: data.model })); + const res = await handler(fakeReq({ content: "hi", model: "echo" })); + expect(res.status).toBe(200); + expect(await res.json()).toEqual({ model: "echo" }); +}); + +test("withBody returns 400 on unparseable JSON", async () => { + const handler = withBody(PromptBody, () => new Response("unreached")); + const res = await handler(fakeReq(null, true)); + expect(res.status).toBe(400); +}); + +test("withBody returns 422 and does NOT call the handler on a schema failure", async () => { + let called = false; + const handler = withBody(PromptBody, () => { called = true; return new Response("x"); }); + const res = await handler(fakeReq({ content: "hi" })); // no model + expect(res.status).toBe(422); + expect(called).toBe(false); + const body = (await res.json()) as { error: string; issues: unknown[] }; + expect(body.error).toBe("validation failed"); + expect(body.issues.length).toBeGreaterThan(0); +}); diff --git a/worker.ts b/worker.ts index 9c41202..122f41b 100644 --- a/worker.ts +++ b/worker.ts @@ -1,135 +1,40 @@ import { Store } from "./src/store"; import { ConversationActor } from "./src/actor"; -import { run, initInference } from "./src/inference"; -import { LEASE_GRACE_MS, HEARTBEAT_INTERVAL_MS, REAP_INTERVAL_MS } from "./src/config"; - -interface CurrentJob { - jobId: string; - conversationId: string; - checkpointSeq: number; - runId: string; - messageId: string; - prompt: string; - model: string; -} +import { initInference } from "./src/inference"; +import { JobDriver } from "./src/drive"; +import { REAP_INTERVAL_MS } from "./src/config"; /** - * Standalone worker process: claim β†’ run β†’ heartbeat β†’ checkpoint β†’ done. - * Reclaims expired leases on startup; `bun:sqlite`'s atomic UPDATE...RETURNING - * does the race-free claim. Run with `bun worker.ts` alongside the server for - * crash isolation (a crash in either tier doesn't take the other down). + * Standalone worker process: claim β†’ run β†’ heartbeat β†’ checkpoint β†’ done, via + * the shared JobDriver (same code the server's inline loop runs). Run with + * `bun worker.ts` alongside the server for crash isolation (a crash in either + * tier doesn't take the other down). Both processes may poll the same job + * table; the SQL claim is atomic and single-writer per conversation. */ -export class Worker { - private readonly store: Store; - private heartbeatTimer: ReturnType | null = null; - private current: CurrentJob | null = null; - - constructor(store: Store) { - this.store = store; - } - - start(): void { - // Reclaim any jobs whose lease expired while we were down. - this.store.reap(Date.now()); - if (!this.heartbeatTimer) { - this.heartbeatTimer = setInterval( - () => this.heartbeat(), - HEARTBEAT_INTERVAL_MS, - ); - } - } - - private heartbeat(): void { - if (this.current) { - this.store.heartbeat(this.current.jobId, Date.now() + LEASE_GRACE_MS); - } - } - - /** - * One claim-attempt: pulls a queued (or expired-lease) job for a - * conversation with no other active run, and returns true if work was - * claimed. Returns false when idle. - */ - async maybeClaim(): Promise { - const row = this.store.claimExpiredExclusive(Date.now()); - if (!row) return false; - const params = JSON.parse(row.params) as { - runId: string; - messageId: string; - prompt: string; - model: string; - }; - this.current = { - jobId: row.id, - conversationId: row.conversation_id, - checkpointSeq: row.checkpoint_seq, - runId: params.runId, - messageId: params.messageId, - prompt: params.prompt, - model: params.model, - }; - return true; - } - - async runClaimed(): Promise { - const j = this.current; - if (!j) return; - const actor = new ConversationActor(j.conversationId, this.store); - try { - await actor.runText( - j.runId, - j.messageId, - async function* (signal) { - for await (const step of run(j.prompt, { - runId: j.runId, - model: j.model, - abortSignal: signal, - })) { - yield step; - } - }, - (seq) => { - // Durable progress: checkpoint + lease advance on each delta flush. - this.store.checkpoint(j.jobId, seq); - this.store.heartbeat(j.jobId, Date.now() + LEASE_GRACE_MS); - }, - ); - this.store.markDone(j.jobId); - } catch (err) { - this.store.markFailed(j.jobId); - } finally { - this.current = null; - } - } - - stop(): void { - if (this.heartbeatTimer) clearInterval(this.heartbeatTimer); - this.heartbeatTimer = null; - } -} - -// Entry point: run the claim loop when invoked directly. if (import.meta.main) { // Build the provider registry (catalog + ops config) before claiming work. await initInference(); const store = new Store(); - const worker = new Worker(store); - worker.start(); + // One actor per claimed conversation. Each job is fully executed before the + // next is claimed, so a plain map keyed by conversation is sufficient. + const actors = new Map(); + const driver = new JobDriver(store, (id) => { + let a = actors.get(id); + if (!a) { + a = new ConversationActor(id, store); + actors.set(id, a); + } + return a; + }); - // Reap expired leases periodically. + // Reclaim any jobs whose lease expired while we were down, then poll. + store.reap(Date.now()); setInterval(() => { store.reap(Date.now()); }, REAP_INTERVAL_MS); - - // Claim loop: poll for work every second. - const drive = async () => { - if (await worker.maybeClaim()) { - await worker.runClaimed(); - } - }; setInterval(() => { - void drive(); + void driver.driveOnce(); }, 1000); console.log("kloe worker started"); -- 2.51.2