mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-09-15 03:12:36 +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.
59 lines
2.1 KiB
TypeScript
59 lines
2.1 KiB
TypeScript
/**
|
|
* @file earlyKeepaliveByteBuffer.ts
|
|
* @description Bridges bytes withEarlyStreamKeepalive writes directly to the
|
|
* client (outside the request handler's own reqLogger) back into that same
|
|
* request's persisted call-log artifact.
|
|
*
|
|
* withEarlyStreamKeepalive wraps a route's handler Promise from OUTSIDE the
|
|
* handler's own call tree — it has no reference to the reqLogger the handler
|
|
* creates deep inside chatCore.ts, and by the time that reqLogger exists the
|
|
* keepalive/startup frames may already be written. A shared correlationId
|
|
* (threaded by the route as handleChat's 4th positional arg, and separately
|
|
* into withEarlyStreamKeepalive's options) is the only thing both sides
|
|
* share, so recordEarlyKeepaliveBytes/takeEarlyKeepaliveBytes key on that
|
|
* instead of trying to pass a live object reference across the boundary.
|
|
*
|
|
* Entries are consumed once (chatCore/attemptLogging.ts calls
|
|
* takeEarlyKeepaliveBytes exactly when it assembles the final call-log
|
|
* payload) and swept on a TTL so a request that never reaches that point
|
|
* (aborted, detailed logging disabled, a route that never wires this up)
|
|
* cannot leak buffered bytes forever.
|
|
*/
|
|
|
|
const MAX_ITEMS_PER_CORRELATION = 200;
|
|
const ENTRY_TTL_MS = 10 * 60 * 1000;
|
|
|
|
type BufferEntry = { chunks: string[]; createdAt: number };
|
|
|
|
const buffers = new Map<string, BufferEntry>();
|
|
|
|
function sweepExpired(): void {
|
|
const cutoff = Date.now() - ENTRY_TTL_MS;
|
|
for (const [correlationId, entry] of buffers) {
|
|
if (entry.createdAt < cutoff) {
|
|
buffers.delete(correlationId);
|
|
}
|
|
}
|
|
}
|
|
|
|
export function recordEarlyKeepaliveBytes(correlationId: string, chunk: string): void {
|
|
if (!correlationId || !chunk) return;
|
|
sweepExpired();
|
|
let entry = buffers.get(correlationId);
|
|
if (!entry) {
|
|
entry = { chunks: [], createdAt: Date.now() };
|
|
buffers.set(correlationId, entry);
|
|
}
|
|
if (entry.chunks.length < MAX_ITEMS_PER_CORRELATION) {
|
|
entry.chunks.push(chunk);
|
|
}
|
|
}
|
|
|
|
export function takeEarlyKeepaliveBytes(correlationId: string): string[] {
|
|
sweepExpired();
|
|
const entry = buffers.get(correlationId);
|
|
if (!entry) return [];
|
|
buffers.delete(correlationId);
|
|
return entry.chunks;
|
|
}
|