From a63722e2adf3f9c158fda5858c20e4bd1e89db55 Mon Sep 17 00:00:00 2001 From: Drew Stone Date: Fri, 21 Aug 2026 03:59:26 -0700 Subject: [PATCH] refactor(http): one drain loop and one closed-response guard for the local servers --- CHANGELOG.md | 8 ++++++ clients/python/pyproject.toml | 2 +- clients/python/src/agent_eval_rpc/__init__.py | 2 +- clients/python/uv.lock | 2 +- package.json | 2 +- src/analyst/benchmark-implementation.ts | 4 +-- src/analyst/trace-tool-callback.ts | 18 +++++-------- src/campaign/external-optimizer-callback.ts | 18 +++++-------- src/campaign/external-optimizer-http.ts | 25 +++++++++++++++++++ .../external-optimizer-model-proxy.ts | 18 +++++-------- 10 files changed, 57 insertions(+), 42 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4d2fd976..8bdacf3c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,14 @@ All notable changes to `@tangle-network/agent-eval` and its sibling `agent-eval- --- +## [0.167.1] — 2026-08-21 + +### Changed + +- `waitForActiveHandlers` and `sendJsonIfOpen` had three byte-identical private copies each, one per local HTTP server (`src/analyst/trace-tool-callback.ts`, `src/campaign/external-optimizer-callback.ts`, `src/campaign/external-optimizer-model-proxy.ts`). Both now live in `src/campaign/external-optimizer-http.ts`, the module all three already imported `closeServer`, `listenLocal`, and `sendJson` from. The drain loop re-reads the handler set on every pass because a running handler can register another; a copy that awaited one snapshot would let the caller close the server with work outstanding. No behavior change and no export change: the helpers are internal to that module's consumers. + +--- + ## [0.167.0] — 2026-08-21 ### Changed diff --git a/clients/python/pyproject.toml b/clients/python/pyproject.toml index 7d295a15..fd686781 100644 --- a/clients/python/pyproject.toml +++ b/clients/python/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "agent-eval-rpc" -version = "0.167.0" +version = "0.167.1" description = "Python RPC client, official optimizer bridge, and DSPy metric adapter for @tangle-network/agent-eval." readme = "README.md" requires-python = ">=3.10" diff --git a/clients/python/src/agent_eval_rpc/__init__.py b/clients/python/src/agent_eval_rpc/__init__.py index ffa697da..0e5d021b 100644 --- a/clients/python/src/agent_eval_rpc/__init__.py +++ b/clients/python/src/agent_eval_rpc/__init__.py @@ -53,7 +53,7 @@ try: __version__ = version("agent-eval-rpc") except PackageNotFoundError: - __version__ = "0.167.0" + __version__ = "0.167.1" __all__ = [ "Client", diff --git a/clients/python/uv.lock b/clients/python/uv.lock index 77c5179a..ac6e1522 100644 --- a/clients/python/uv.lock +++ b/clients/python/uv.lock @@ -34,7 +34,7 @@ conflicts = [[ [[package]] name = "agent-eval-rpc" -version = "0.167.0" +version = "0.167.1" source = { editable = "." } dependencies = [ { name = "filelock" }, diff --git a/package.json b/package.json index 0be2596c..b5a1665f 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-eval", - "version": "0.167.0", + "version": "0.167.1", "description": "Evaluate and improve AI agents from runs, traces, judges, and feedback. Compare candidates, cluster failures, measure lift, and gate releases.", "homepage": "https://github.com/tangle-network/agent-eval#readme", "repository": { diff --git a/src/analyst/benchmark-implementation.ts b/src/analyst/benchmark-implementation.ts index 16b3e42f..9ae9cba3 100644 --- a/src/analyst/benchmark-implementation.ts +++ b/src/analyst/benchmark-implementation.ts @@ -10,7 +10,7 @@ export const ANALYST_BENCHMARK_DEPENDENCY_LOCK_FILES = Object.freeze([ ]) export const ANALYST_BENCHMARK_DEPENDENCY_LOCK_SHA256 = - '5487ec58c4a939ece38536a982a9a72e3d59363a2ed8677c6ad2aaa9699c06f1' + 'a372aba7edb147c768813488fa7c69395f093f4c2b9bb9862aff2d3b468be47a' /** The published benchmark evidence was produced at this package version, by * the retired one-shot direct runner, before trace analysts moved to the @@ -139,7 +139,7 @@ export const ANALYST_BENCHMARK_IMPLEMENTATION_FILES = Object.freeze([ ]) export const ANALYST_BENCHMARK_IMPLEMENTATION_SHA256 = - 'c3c6778b7102a635706b76a8bd5a1f0375cd84e30f900a4026c7939cb6fb04b9' + '2ab66a98fac14974e8e3fde86600899a38a1cbba2957bed3afad44261108c43e' export function analystBenchmarkImplementationDigest() { return ANALYST_BENCHMARK_IMPLEMENTATION_SHA256 diff --git a/src/analyst/trace-tool-callback.ts b/src/analyst/trace-tool-callback.ts index 581e1d4f..8ce8e66b 100644 --- a/src/analyst/trace-tool-callback.ts +++ b/src/analyst/trace-tool-callback.ts @@ -4,7 +4,12 @@ import { type ExternalOptimizerCallbackLimits, resolveExternalOptimizerCallbackLimits, } from '../campaign/external-optimizer-contracts' -import { closeServer, listenLocal, sendJson } from '../campaign/external-optimizer-http' +import { + closeServer, + listenLocal, + sendJsonIfOpen, + waitForActiveHandlers, +} from '../campaign/external-optimizer-http' import type { TraceAnalysisToolDescriptor } from '../trace-analyst/tools' export interface TraceToolCallback { @@ -144,12 +149,6 @@ export async function startTraceToolCallback(args: { } } -async function waitForActiveHandlers(activeHandlers: Set>): Promise { - while (activeHandlers.size > 0) { - await Promise.allSettled([...activeHandlers]) - } -} - function readJson(request: IncomingMessage, maxRequestBytes: number): Promise { return new Promise((resolve, reject) => { let size = 0 @@ -174,11 +173,6 @@ function readJson(request: IncomingMessage, maxRequestBytes: number): Promise { return typeof value === 'object' && value !== null && !Array.isArray(value) } diff --git a/src/campaign/external-optimizer-callback.ts b/src/campaign/external-optimizer-callback.ts index deb5eb38..15d756f7 100644 --- a/src/campaign/external-optimizer-callback.ts +++ b/src/campaign/external-optimizer-callback.ts @@ -10,7 +10,12 @@ import { isRecord, resolveExternalOptimizerCallbackLimits, } from './external-optimizer-contracts' -import { closeServer, listenLocal, sendJson } from './external-optimizer-http' +import { + closeServer, + listenLocal, + sendJsonIfOpen, + waitForActiveHandlers, +} from './external-optimizer-http' type UnsequencedObservation = ExternalOptimizerEvaluationObservation extends infer T ? T extends ExternalOptimizerEvaluationObservation @@ -233,17 +238,6 @@ async function handleCallback( } } -async function waitForActiveHandlers(activeHandlers: Set>): Promise { - while (activeHandlers.size > 0) { - await Promise.allSettled([...activeHandlers]) - } -} - -function sendJsonIfOpen(response: ServerResponse, status: number, body: unknown): void { - if (response.destroyed || response.writableEnded) return - sendJson(response, status, body) -} - function readJson(request: IncomingMessage, maxRequestBytes: number): Promise { return new Promise((resolvePromise, reject) => { let size = 0 diff --git a/src/campaign/external-optimizer-http.ts b/src/campaign/external-optimizer-http.ts index f2fc74f4..d065ac24 100644 --- a/src/campaign/external-optimizer-http.ts +++ b/src/campaign/external-optimizer-http.ts @@ -25,3 +25,28 @@ export function sendJson(response: ServerResponse, status: number, body: unknown response.writeHead(status, { 'content-type': 'application/json; charset=utf-8' }) response.end(JSON.stringify(body)) } + +/** + * Send a JSON body only while the response can still take one. + * + * A handler that lost its client — the request was aborted, or the response + * already ended — must not write again; Node throws `ERR_STREAM_WRITE_AFTER_END` + * and the throw escapes into the server's error path rather than the caller's. + */ +export function sendJsonIfOpen(response: ServerResponse, status: number, body: unknown): void { + if (response.destroyed || response.writableEnded) return + sendJson(response, status, body) +} + +/** + * Wait until every in-flight handler has settled. + * + * Re-reads the set on each pass: a handler that is still running can register + * another, so awaiting one snapshot would return while work is outstanding and + * the caller would close the server under it. + */ +export async function waitForActiveHandlers(activeHandlers: Set>): Promise { + while (activeHandlers.size > 0) { + await Promise.allSettled([...activeHandlers]) + } +} diff --git a/src/campaign/external-optimizer-model-proxy.ts b/src/campaign/external-optimizer-model-proxy.ts index 2756d642..e009de74 100644 --- a/src/campaign/external-optimizer-model-proxy.ts +++ b/src/campaign/external-optimizer-model-proxy.ts @@ -35,7 +35,12 @@ import { type ExternalOptimizerWireCounts, isRecord, } from './external-optimizer-contracts' -import { closeServer, listenLocal, sendJson } from './external-optimizer-http' +import { + closeServer, + listenLocal, + sendJsonIfOpen, + waitForActiveHandlers, +} from './external-optimizer-http' const MODEL_PROXY_PATHS = new Set(['/v1/chat/completions', '/v1/responses']) type ModelProxyPath = '/v1/chat/completions' | '/v1/responses' | '/v1/messages' @@ -474,17 +479,6 @@ async function handleModelProxyRequest(args: { } } -async function waitForActiveHandlers(activeHandlers: Set>): Promise { - while (activeHandlers.size > 0) { - await Promise.allSettled([...activeHandlers]) - } -} - -function sendJsonIfOpen(response: ServerResponse, status: number, body: unknown): void { - if (response.destroyed || response.writableEnded) return - sendJson(response, status, body) -} - async function forwardModelProxyRequest(args: { call: ExternalOptimizerModelCall callId: string