Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/fixes/11362-video-bridge-result-cache.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(video):** fingerprint protected Video Bridge bytes, coalesce concurrent work, and fail open when the bounded TTL/LRU result cache is unavailable or corrupt ([#11362](https://github.com/diegosouzapw/OmniRoute/pull/11362))
89 changes: 78 additions & 11 deletions src/lib/guardrails/modalityBridge/bridgeCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@ export function bridgeCacheKey(

export interface BridgeCacheOptions {
maxEntries: number;
/** Aggregate UTF-8 key/value/metadata budget; unlimited when omitted. */
maxBytes?: number;
ttlMs: number;
/** Injectable clock for tests. */
now?: () => number;
Expand All @@ -67,8 +69,37 @@ export interface BridgeCacheEntry {
metadata?: Record<string, unknown>;
}

export class BridgeCache {
private readonly entries = new Map<string, { entry: BridgeCacheEntry; expiresAt: number }>();
/** Minimal fail-open store contract accepted by complete-result bridge caches. */
export interface BridgeCacheStore {
delete(key: string): void;
getEntry(key: string): BridgeCacheEntry | undefined;
setEntry(key: string, entry: BridgeCacheEntry): void;
}

type StoredBridgeCacheEntry = {
bytes: number;
entry: BridgeCacheEntry;
expiresAt: number;
};

function cacheEntryBytes(entry: BridgeCacheEntry): number {
try {
const metadata = JSON.stringify({
metadata: entry.metadata,
producerModel: entry.producerModel,
});
return Buffer.byteLength(entry.value, "utf8") + Buffer.byteLength(metadata, "utf8");
} catch (error) {
console.debug("[MODALITY_BRIDGE_CACHE] Entry size calculation failed open", {
errorType: error instanceof Error ? error.name : typeof error,
});
return Number.POSITIVE_INFINITY;
}
}

export class BridgeCache implements BridgeCacheStore {
private readonly entries = new Map<string, StoredBridgeCacheEntry>();
private totalBytes = 0;

constructor(private readonly opts: BridgeCacheOptions) {}

Expand All @@ -81,7 +112,7 @@ export class BridgeCache {
if (!hit) return undefined;
const now = (this.opts.now ?? Date.now)();
if (hit.expiresAt <= now) {
this.entries.delete(key);
this.delete(key);
return undefined;
}
// Map preserves insertion order — re-insert to mark as most-recently-used.
Expand All @@ -96,34 +127,70 @@ export class BridgeCache {

setEntry(key: string, entry: BridgeCacheEntry): void {
const now = (this.opts.now ?? Date.now)();
this.entries.delete(key);
this.entries.set(key, { entry, expiresAt: now + this.opts.ttlMs });
while (this.entries.size > this.opts.maxEntries) {
const bytes = cacheEntryBytes(entry) + Buffer.byteLength(key, "utf8");
const maxBytes = Math.max(0, this.opts.maxBytes ?? Number.POSITIVE_INFINITY);
const maxEntries = Math.max(0, Math.floor(this.opts.maxEntries));
this.delete(key);
if (!Number.isFinite(bytes) || bytes > maxBytes || maxEntries === 0) return;
this.entries.set(key, { bytes, entry, expiresAt: now + this.opts.ttlMs });
this.totalBytes += bytes;
while (this.entries.size > maxEntries || this.totalBytes > maxBytes) {
const oldest = this.entries.keys().next().value;
if (oldest === undefined) break;
this.entries.delete(oldest);
this.delete(oldest);
}
}

get size(): number {
return this.entries.size;
}

/** Current aggregate UTF-8 bytes retained by this cache. */
get bytes(): number {
return this.totalBytes;
}

delete(key: string): void {
const existing = this.entries.get(key);
if (existing) this.totalBytes = Math.max(0, this.totalBytes - existing.bytes);
this.entries.delete(key);
}

clear(): void {
this.entries.clear();
this.totalBytes = 0;
}
}

/** Process-wide singleton used by the bridges; recreated when config changes. */
let shared: { cache: BridgeCache; ttlMs: number; maxEntries: number } | null = null;
let shared: { cache: BridgeCache; ttlMs: number; maxBytes: number; maxEntries: number } | null =
null;

export function getSharedBridgeCache(ttlMs: number, maxEntries: number): BridgeCache {
if (!shared || shared.ttlMs !== ttlMs || shared.maxEntries !== maxEntries) {
shared = { cache: new BridgeCache({ maxEntries, ttlMs }), ttlMs, maxEntries };
/**
* Resolve the process-wide bridge cache, recreating it when any bound changes.
*
* @param ttlMs - Entry lifetime in milliseconds.
* @param maxEntries - Maximum retained entry count.
* @param maxBytes - Aggregate UTF-8 storage budget.
* @returns The process-wide cache for these exact bounds.
*/
export function getSharedBridgeCache(
ttlMs: number,
maxEntries: number,
maxBytes = Number.POSITIVE_INFINITY
): BridgeCache {
if (
!shared ||
shared.ttlMs !== ttlMs ||
shared.maxEntries !== maxEntries ||
shared.maxBytes !== maxBytes
) {
shared = {
cache: new BridgeCache({ maxBytes, maxEntries, ttlMs }),
ttlMs,
maxBytes,
maxEntries,
};
}
return shared.cache;
}
Expand Down
6 changes: 6 additions & 0 deletions src/lib/guardrails/modalityBridge/bridgeStats.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ export interface BridgeModalityStats {
resultCacheBytes: number;
resultCacheHits: number;
resultCacheLatencyMs: number;
/** Requests that joined an in-flight complete result instead of hitting the persistent cache. */
resultSingleflightCoalesced: number;
failures: number;
/** Audio/video fusion runs (video bridge only; 0 for other modalities). */
fusionRuns: number;
Expand Down Expand Up @@ -47,6 +49,7 @@ function emptyStats(): BridgeModalityStats {
resultCacheBytes: 0,
resultCacheHits: 0,
resultCacheLatencyMs: 0,
resultSingleflightCoalesced: 0,
failures: 0,
fusionRuns: 0,
fusionPartials: 0,
Expand All @@ -69,6 +72,8 @@ export function recordBridgeUse(
resultCacheBytes?: number;
resultCacheHit?: boolean;
resultCacheLatencyMs?: number;
/** True only when this request joined existing in-flight result work. */
resultSingleflightCoalesced?: boolean;
} = {}
): void {
const s = stats[kind];
Expand Down Expand Up @@ -104,6 +109,7 @@ export function recordBridgeUse(
s.resultCacheLatencyMs += Math.max(0, opts.resultCacheLatencyMs);
}
}
if (opts.resultSingleflightCoalesced) s.resultSingleflightCoalesced += 1;
if (typeof opts.latencyMs === "number" && Number.isFinite(opts.latencyMs)) {
s.totalLatencyMs += Math.max(0, opts.latencyMs);
s.latencySamples += 1;
Expand Down
Loading
Loading