mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-09-22 06:42:19 +03:00
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.
47 lines
1.4 KiB
TypeScript
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 });
|
|
},
|
|
};
|
|
}
|
|
}
|