diff --git a/package-lock.json b/package-lock.json index 79750b77..80ff3ec8 100644 --- a/package-lock.json +++ b/package-lock.json @@ -3514,6 +3514,7 @@ "integrity": "sha512-6/cmF2piao+f6wSxUsJLZjck7OQsYyRtcOZS02k7XINSNlz93v6emM8WutDQSXnroG2xwYlEVHJI+cPA7CPM3Q==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "@typescript-eslint/scope-manager": "8.50.0", "@typescript-eslint/types": "8.50.0", @@ -3822,6 +3823,7 @@ "integrity": "sha512-oukfKT9Mk41LreEW09vt45f8wx7DordoWUZMYdY/cyAk7w5TWkTRCNZYF7sX7n2wB7jyGAl74OxgwhPgKaqDMQ==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "@vitest/utils": "3.2.4", "pathe": "^2.0.3", @@ -3837,6 +3839,7 @@ "integrity": "sha512-dEYtS7qQP2CjU27QBC5oUOxLE/v5eLkGqPE0ZKEIDGMs4vKWe7IjgLOeauHsR0D5YuuycGRO5oSRXnwnmA78fQ==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "@vitest/pretty-format": "3.2.4", "magic-string": "^0.30.17", @@ -3894,6 +3897,7 @@ "integrity": "sha512-NZyJarBfL7nWwIq+FDL6Zp/yHEhePMNnnJ0y3qfieCrmNvYct8uvtiV41UvlSe6apAfk0fY1FbWx+NwfmpvtTg==", "dev": true, "license": "MIT", + "peer": true, "bin": { "acorn": "bin/acorn" }, @@ -4655,6 +4659,7 @@ "dev": true, "hasInstallScript": true, "license": "MIT", + "peer": true, "bin": { "esbuild": "bin/esbuild" }, @@ -4726,6 +4731,7 @@ "integrity": "sha512-LEyamqS7W5HB3ujJyvi0HQK/dtVINZvd5mAAp9eT5S/ujByGjiZLCzPcHVzuXbpJDJF/cxwHlfceVUDZ2lnSTw==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "@eslint-community/eslint-utils": "^4.8.0", "@eslint-community/regexpp": "^4.12.1", @@ -4786,6 +4792,7 @@ "integrity": "sha512-82GZUjRS0p/jganf6q1rEO25VSoHH0hKPCTrgillPjdI/3bgBhAE1QzHrHTizjpRvy6pGAvKjDJtk2pF9NDq8w==", "dev": true, "license": "MIT", + "peer": true, "bin": { "eslint-config-prettier": "bin/cli.js" }, @@ -6398,6 +6405,7 @@ "integrity": "sha512-5gTmgEY/sqK6gFXLIsQNH19lWb4ebPDLA4SdLP7dsWkIXHWlG66oPuVvXSGFPppYZz8ZDZq0dYYrbHfBCVUb1Q==", "dev": true, "license": "MIT", + "peer": true, "engines": { "node": ">=12" }, @@ -6447,6 +6455,7 @@ } ], "license": "MIT", + "peer": true, "dependencies": { "nanoid": "^3.3.11", "picocolors": "^1.1.1", @@ -6515,6 +6524,7 @@ "integrity": "sha512-v6UNi1+3hSlVvv8fSaoUbggEM5VErKmmpGA7Pl3HF8V6uKY7rvClBOJlH6yNwQtfTueNkGVpOv/mtWL9L4bgRA==", "dev": true, "license": "MIT", + "peer": true, "bin": { "prettier": "bin/prettier.cjs" }, @@ -7502,6 +7512,7 @@ "integrity": "sha512-5C1sg4USs1lfG0GFb2RLXsdpXqBSEhAaA/0kPL01wxzpMqLILNxIxIOKiILz+cdg/pLnOUxFYOR5yhHU666wbw==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "esbuild": "~0.27.0", "get-tsconfig": "^4.7.5" @@ -7550,6 +7561,7 @@ "integrity": "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw==", "dev": true, "license": "Apache-2.0", + "peer": true, "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" @@ -7612,6 +7624,7 @@ "integrity": "sha512-i7qRCmY42zmCwnYlh9H2SvLEypEFGye5iRmEMKjcGi7zk9UquigRjFtTLz0TYqr0ZGLZhaMHl/foy1bZR+Cwlw==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "pathe": "^2.0.3" } @@ -7678,6 +7691,7 @@ "integrity": "sha512-dZwN5L1VlUBewiP6H9s2+B3e3Jg96D0vzN+Ry73sOefebhYr9f94wwkMNN/9ouoU8pV1BqA1d1zGk8928cx0rg==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "esbuild": "^0.27.0", "fdir": "^6.5.0", @@ -7776,6 +7790,7 @@ "integrity": "sha512-LUCP5ev3GURDysTWiP47wRRUpLKMOfPh+yKTx3kVIEiu5KOMeqzpnYNsKyOoVrULivR8tLcks4+lga33Whn90A==", "dev": true, "license": "MIT", + "peer": true, "dependencies": { "@types/chai": "^5.2.2", "@vitest/expect": "3.2.4", @@ -7893,6 +7908,7 @@ "dev": true, "hasInstallScript": true, "license": "Apache-2.0", + "peer": true, "bin": { "workerd": "bin/workerd" }, diff --git a/src/client/index.ts b/src/client/index.ts index a9b4298f..ada95783 100644 --- a/src/client/index.ts +++ b/src/client/index.ts @@ -30,6 +30,7 @@ export type { ClientCallInput, ClientCallResult, } from './interceptors.js'; +export { TransportStats, AdaptiveTransportInterceptor } from './transports/adaptive_transport.js'; export { ServiceParameters, type ServiceParametersUpdate, diff --git a/src/client/transports/adaptive_transport.ts b/src/client/transports/adaptive_transport.ts new file mode 100644 index 00000000..0e101a03 --- /dev/null +++ b/src/client/transports/adaptive_transport.ts @@ -0,0 +1,207 @@ +/** + * Adaptive transport selection — bio-inspired strategy. + * + * Inspired by cellular signal pathway selection: cells activate different + * signaling cascades (cAMP, Ca²⁺, MAPK) depending on recent efficacy and + * response speed. This module applies the same principle to A2A transport + * selection: transports that succeed more often and respond faster are + * preferred, while failing transports are deprioritized and allowed to + * recover over time. + * + * Usage: + * ```ts + * const stats = new TransportStats(); + * const options = ClientFactoryOptions.createFrom(ClientFactoryOptions.default, { + * clientConfig: { interceptors: [new AdaptiveTransportInterceptor(stats)] }, + * }); + * // Before each request, read preferred order: + * const preferred = stats.preferredOrder(); + * ``` + * + * @see Kholodenko, B.N. (2006) "Cell-signalling dynamics in time and space." + * Nat Rev Mol Cell Biol 7:165-176. + * + * @module + */ + +import { CallInterceptor, BeforeArgs, AfterArgs } from '../interceptors.js'; + +// --------------------------------------------------------------------------- +// Transport statistics +// --------------------------------------------------------------------------- + +interface TransportRecord { + ok: boolean; + latencyMs: number; + timestamp: number; +} + +/** + * Tracks per-transport success rates and latencies within a sliding window. + * + * Analogous to how cells monitor pathway efficacy: each signaling cascade + * has an implicit "success rate" (fraction of signals that reach the + * nucleus) and a "latency" (time from receptor activation to gene + * expression). Pathways with higher efficacy and lower latency are + * preferentially activated. + */ +export class TransportStats { + private readonly windowSize: number; + private readonly latencyNormalizer: number; + private readonly records = new Map(); + + /** + * @param windowSize Number of recent outcomes to keep per transport. + * @param latencyNormalizer When average latency equals this value (ms), + * the latency factor is 0.5. Analogous to the Km in Michaelis-Menten + * kinetics: half-maximal response at this "concentration". + */ + constructor(windowSize = 20, latencyNormalizer = 1000) { + this.windowSize = Math.max(1, windowSize); + this.latencyNormalizer = Math.max(1, latencyNormalizer); + } + + /** Record a transport attempt outcome. */ + record(transport: string, ok: boolean, latencyMs: number): void { + let list = this.records.get(transport); + if (!list) { + list = []; + this.records.set(transport, list); + } + list.push({ ok, latencyMs: Math.max(0, latencyMs), timestamp: Date.now() }); + if (list.length > this.windowSize) { + list.splice(0, list.length - this.windowSize); + } + } + + /** Success rate in [0, 1]. Returns 1 for unknown transports (explore-first). */ + successRate(transport: string): number { + const list = this.records.get(transport); + if (!list || list.length === 0) return 1; + return list.filter((r) => r.ok).length / list.length; + } + + /** Average latency (ms) of successful calls. Returns 0 if unknown. */ + avgLatency(transport: string): number { + const list = this.records.get(transport); + if (!list) return 0; + const ok = list.filter((r) => r.ok); + if (ok.length === 0) return 0; + return ok.reduce((s, r) => s + r.latencyMs, 0) / ok.length; + } + + /** + * Composite score: `successRate × latencyFactor`. + * + * `latencyFactor = 1 / (1 + avgLatency / normalizer)` + * + * Higher score = preferred transport. Unknown transports score 1.0 + * (explore-first, analogous to immune system sampling novel antigens). + */ + getScore(transport: string): number { + const sr = this.successRate(transport); + const avg = this.avgLatency(transport); + return sr * (1 / (1 + avg / this.latencyNormalizer)); + } + + /** Number of recorded outcomes for a transport. */ + count(transport: string): number { + return this.records.get(transport)?.length ?? 0; + } + + /** + * Returns transport names ordered by descending score. + * Only includes transports that have at least one recorded outcome. + */ + preferredOrder(): string[] { + return Array.from(this.records.keys()).sort((a, b) => this.getScore(b) - this.getScore(a)); + } + + /** Clear all recorded data. */ + clear(): void { + this.records.clear(); + } +} + +// --------------------------------------------------------------------------- +// Adaptive transport interceptor +// --------------------------------------------------------------------------- + +/** + * A {@link CallInterceptor} that records transport outcomes into + * {@link TransportStats} for adaptive transport selection. + * + * Attach to `ClientConfig.interceptors` to automatically track success/failure + * and latency of every transport call. + * + * @example + * ```ts + * const stats = new TransportStats(); + * const interceptor = new AdaptiveTransportInterceptor(stats); + * + * const factory = new ClientFactory( + * ClientFactoryOptions.createFrom(ClientFactoryOptions.default, { + * clientConfig: { interceptors: [interceptor] }, + * }) + * ); + * ``` + */ +export class AdaptiveTransportInterceptor implements CallInterceptor { + private readonly stats: TransportStats; + private readonly timers = new Map(); + + constructor(stats: TransportStats) { + this.stats = stats; + } + + async before(args: BeforeArgs): Promise { + if (!args.input) return; + // Record start time keyed by a unique request identifier. + // We use the method name + timestamp as a simple key since + // interceptors are invoked synchronously per request. + const key = `${args.input.method}:${Date.now()}`; + this.timers.set(key, performance.now()); + // Store key in options for retrieval in after() + if (!args.options) { + args.options = {}; + } + (args.options as Record)['_adaptiveTimerKey'] = key; + } + + async after(args: AfterArgs): Promise { + if (!args.result) return; + const key = (args.options as Record | undefined)?.['_adaptiveTimerKey'] as + | string + | undefined; + if (!key) return; + + const startTime = this.timers.get(key); + this.timers.delete(key); + if (startTime === undefined) return; + + const latencyMs = performance.now() - startTime; + + // Determine transport name from agent card's supported protocols. + // The first matching protocol is the one being used. + const protocols = + args.agentCard?.additionalInterfaces?.map((i: { transport?: string }) => i.transport) ?? []; + const transport = protocols[0] ?? 'unknown'; + + // Determine success based on whether result contains an error-like value. + const ok = !isErrorResult(args.result); + this.stats.record(transport, ok, latencyMs); + } + + /** Access the underlying stats for reading preferred order or scores. */ + getStats(): TransportStats { + return this.stats; + } +} + +function isErrorResult(result: { value?: unknown }): boolean { + const val = result?.value; + if (val == null) return false; + if (val instanceof Error) return true; + if (typeof val === 'object' && 'error' in val) return true; + return false; +} diff --git a/test/client/adaptive_transport.spec.ts b/test/client/adaptive_transport.spec.ts new file mode 100644 index 00000000..f5423123 --- /dev/null +++ b/test/client/adaptive_transport.spec.ts @@ -0,0 +1,104 @@ +import { describe, it, expect, beforeEach } from 'vitest'; +import { + TransportStats, + AdaptiveTransportInterceptor, +} from '../../src/client/transports/adaptive_transport.js'; + +// --------------------------------------------------------------------------- +// TransportStats +// --------------------------------------------------------------------------- + +describe('TransportStats', () => { + let stats: TransportStats; + + beforeEach(() => { + stats = new TransportStats(10, 1000); + }); + + it('returns score 1.0 for unknown transports (explore-first)', () => { + expect(stats.getScore('JSONRPC')).toBe(1); + expect(stats.successRate('JSONRPC')).toBe(1); + expect(stats.avgLatency('JSONRPC')).toBe(0); + }); + + it('tracks success rate correctly', () => { + stats.record('JSONRPC', true, 50); + stats.record('JSONRPC', true, 50); + stats.record('JSONRPC', false, 5000); + expect(stats.successRate('JSONRPC')).toBeCloseTo(2 / 3); + }); + + it('tracks average latency of successful calls only', () => { + stats.record('JSONRPC', true, 100); + stats.record('JSONRPC', true, 200); + stats.record('JSONRPC', false, 9999); // failure excluded + expect(stats.avgLatency('JSONRPC')).toBe(150); + }); + + it('composite score penalizes high latency', () => { + stats.record('FAST', true, 50); + stats.record('SLOW', true, 2000); + expect(stats.getScore('FAST')).toBeGreaterThan(stats.getScore('SLOW')); + }); + + it('composite score penalizes low success rate', () => { + for (let i = 0; i < 10; i++) stats.record('RELIABLE', true, 100); + for (let i = 0; i < 10; i++) stats.record('FLAKY', i < 5, 100); + expect(stats.getScore('RELIABLE')).toBeGreaterThan(stats.getScore('FLAKY')); + }); + + it('reliable beats flaky even with higher latency', () => { + // 50% success, 30ms + for (let i = 0; i < 10; i++) stats.record('FAST_FLAKY', i < 5, 30); + // 100% success, 800ms + for (let i = 0; i < 10; i++) stats.record('SLOW_RELIABLE', true, 800); + + expect(stats.getScore('SLOW_RELIABLE')).toBeGreaterThan(stats.getScore('FAST_FLAKY')); + }); + + it('sliding window evicts old records', () => { + // Fill window with failures + for (let i = 0; i < 10; i++) stats.record('JSONRPC', false, 5000); + expect(stats.successRate('JSONRPC')).toBe(0); + + // Now record 10 successes — old failures should be evicted + for (let i = 0; i < 10; i++) stats.record('JSONRPC', true, 50); + expect(stats.successRate('JSONRPC')).toBe(1); + }); + + it('preferredOrder sorts by descending score', () => { + stats.record('WORST', false, 5000); + stats.record('MIDDLE', true, 500); + stats.record('BEST', true, 50); + + const order = stats.preferredOrder(); + expect(order[0]).toBe('BEST'); + expect(order[order.length - 1]).toBe('WORST'); + }); + + it('count returns number of records', () => { + expect(stats.count('JSONRPC')).toBe(0); + stats.record('JSONRPC', true, 50); + stats.record('JSONRPC', true, 60); + expect(stats.count('JSONRPC')).toBe(2); + }); + + it('clear removes all data', () => { + stats.record('JSONRPC', true, 50); + stats.clear(); + expect(stats.count('JSONRPC')).toBe(0); + expect(stats.getScore('JSONRPC')).toBe(1); // back to unknown default + }); +}); + +// --------------------------------------------------------------------------- +// AdaptiveTransportInterceptor +// --------------------------------------------------------------------------- + +describe('AdaptiveTransportInterceptor', () => { + it('exposes stats via getStats()', () => { + const stats = new TransportStats(); + const interceptor = new AdaptiveTransportInterceptor(stats); + expect(interceptor.getStats()).toBe(stats); + }); +});