Files
OmniRoute/open-sse/vendor/codex-chatgpt-web/event-queue.ts
Praveen K Palaniswamy 65e81158ab fix(ollama): route models by advertised capability (#11088)
Landed with the design call resolved per the owner's pick — **option 1**: the synced store is now endpoint-agnostic (persistDiscoveredModels and managedModelImport no longer drop non-chat models at write time), and chat selectability moved to read time (auto-pool expansion in autoStrategy applies filterChatSelectableModels; the models-route projection already had its chatOnly filter). Your discovery test now passes end-to-end (3/3): /api/show capabilities persist per connection and image/embedding requests route through the advertising host.

Reconciliation notes: conflicted areas merged onto the current tip (adobe discovery import, requestedModel preflight signature, resolvedProvider fast-path coexists with the synced-route override — explicit resolution wins); carried base-red drains (#10055 memoization, #11071 test variants) dropped as already-landed; the managed-model-import exclusion test was propagated to the new contract (image/video models persist; the read filter still hides them from chat pickers — pinned by a new assertion). Full battery: 205/206 focused (the one red is a confirmed periodic-timer timing flake on the loaded devbox — 20/20 isolated), autoCombo vitest 30/30, combo suites 46/46, gates + typecheck clean.

Thank you @yourspraveen — the capability probe + routing design was right; it just needed the store contract opened up. Fixes #11087.
2026-08-23 11:45:01 -03:00

47 lines
1.4 KiB
TypeScript

/* Adapted from miuuyy/codex-chatgpt-web commit 55592fca0ba19a27f1b769cec8fff61ff340a785 (MIT). */
export class AsyncEventQueue<T> implements AsyncIterable<T> {
private readonly buffered: T[] = [];
private readonly waiters: Array<(result: IteratorResult<T>) => void> = [];
private closed = false;
constructor(private readonly maxBuffered = 10_000) {}
push(value: T): void {
if (this.closed) return;
const waiter = this.waiters.shift();
if (waiter) {
waiter({ value, done: false });
return;
}
if (this.buffered.length >= this.maxBuffered) throw new Error("Adapter event backlog exceeded");
this.buffered.push(value);
}
close(): void {
if (this.closed) return;
this.closed = true;
while (this.waiters.length > 0) this.waiters.shift()!({ value: undefined, done: true });
}
async collect(): Promise<T[]> {
const values: T[] = [];
for await (const value of this) values.push(value);
return values;
}
[Symbol.asyncIterator](): AsyncIterator<T> {
return {
next: () => {
const value = this.buffered.shift();
if (value !== undefined) return Promise.resolve({ value, done: false });
if (this.closed) return Promise.resolve({ value: undefined, done: true });
return new Promise((resolve) => this.waiters.push(resolve));
},
return: () => {
this.close();
return Promise.resolve({ value: undefined, done: true });
},
};
}
}