From 148fc40bde2c736c4135229d4b62bc4ded1c062d Mon Sep 17 00:00:00 2001 From: Dotta Date: Thu, 10 Sep 2026 16:55:57 -0500 Subject: [PATCH] fix(avatars): await stream disposal before cache cleanup Co-Authored-By: Paperclip --- server/src/__tests__/agent-avatars.test.ts | 27 +++++++++++++++++----- server/src/routes/agent-avatars.ts | 10 ++++++-- 2 files changed, 29 insertions(+), 8 deletions(-) diff --git a/server/src/__tests__/agent-avatars.test.ts b/server/src/__tests__/agent-avatars.test.ts index d1cbb41136..9e5e9f3892 100644 --- a/server/src/__tests__/agent-avatars.test.ts +++ b/server/src/__tests__/agent-avatars.test.ts @@ -17,6 +17,9 @@ import { agentAvatarRoutes } from "../routes/agent-avatars.js"; const request: AgentAvatarRequest = { appearance: appearanceForPalette("arctic-blue"), size: 24, scale: 2, pose: "rest", muted: false }; const cleanups: Array<() => Promise> = []; afterEach(async () => { await Promise.all(cleanups.splice(0).map(fn => fn())); }); +async function consume(stream: Readable) { + for await (const _ of stream) { /* Finish reads before deleting their cache files. */ } +} async function storage() { const dir = await mkdtemp(path.join(os.tmpdir(), "agent-avatar-test-")); cleanups.push(() => rm(dir, { recursive: true, force: true })); @@ -37,10 +40,10 @@ describe("on-demand agent avatars", () => { expect(render).toHaveBeenCalledTimes(1); expect(new Set(results.map(result => result.etag)).size).toBe(1); await Promise.all(results.map(async result => { for await (const _ of result.stream) { /* consume */ } })); - (await createAgentAvatarService(provider, render).get(request)).stream.destroy(); + await consume((await createAgentAvatarService(provider, render).get(request)).stream); expect(render).toHaveBeenCalledTimes(1); await provider.deleteObject({ objectKey: avatarCacheKey(request) }); - (await service.get(request)).stream.destroy(); + await consume((await service.get(request)).stream); expect(render).toHaveBeenCalledTimes(2); }); it("limits cold keys per client while admitting warm hits, joiners and other clients", async () => { @@ -58,10 +61,10 @@ describe("on-demand agent avatars", () => { pending.push(service.get(keys[32], "two")); // Another client still has room. await vi.waitFor(() => expect(render).toHaveBeenCalledTimes(33)); } finally { release(); } - for (const result of await Promise.all(pending)) result.stream.destroy(); - (await service.get(keys[0], "one")).stream.destroy(); + for (const result of await Promise.all(pending)) await consume(result.stream); + await consume((await service.get(keys[0], "one")).stream); expect(render).toHaveBeenCalledTimes(33); - (await service.get(keys[33], "one")).stream.destroy(); // Completed renders release slots. + await consume((await service.get(keys[33], "one")).stream); // Completed renders release slots. expect(render).toHaveBeenCalledTimes(34); }); it("returns retryable admission errors without caching them", async () => { @@ -102,7 +105,7 @@ describe("on-demand agent avatars", () => { expect(Buffer.concat(secondBytes)).toEqual(Buffer.concat(firstBytes)); expect(render).toHaveBeenCalledTimes(1); await provider.deleteObject({ objectKey: avatarCacheKey(request) }); - (await createAgentAvatarService(provider, render).get(request)).stream.destroy(); + await consume((await createAgentAvatarService(provider, render).get(request)).stream); expect(render).toHaveBeenCalledTimes(2); } finally { send.mockRestore(); } }); @@ -126,6 +129,18 @@ describe("on-demand agent avatars", () => { } expect(render).toHaveBeenCalledTimes(1); }); + it("settles stream disposal on a 304 even if the cached file disappears during open", async () => { + const service = createAgentAvatarService(await storage(), async () => Buffer.from("png")); + const stream = new Readable({ + read() {}, + destroy(_error, callback) { setImmediate(() => callback(new Error("cached file removed during open"))); }, + }); + vi.spyOn(service, "get").mockResolvedValueOnce({ stream, byteSize: 3, etag: '"cached"' }); + const url = await serve(service); + const response = await fetch(url, { headers: { "If-None-Match": '"cached"' } }); + expect(response.status).toBe(304); + expect(stream.closed).toBe(true); + }); it("does not cache rendering failures and permits retry", async () => { const render = vi.fn().mockRejectedValueOnce(new Error("unavailable")).mockResolvedValue(Buffer.from("png")); const url = await serve(createAgentAvatarService(await storage(), render)); diff --git a/server/src/routes/agent-avatars.ts b/server/src/routes/agent-avatars.ts index 1e62473be6..3e5c66ffc6 100644 --- a/server/src/routes/agent-avatars.ts +++ b/server/src/routes/agent-avatars.ts @@ -1,4 +1,4 @@ -import { pipeline } from "node:stream/promises"; +import { finished, pipeline } from "node:stream/promises"; import { Router } from "express"; import { z } from "zod"; import { AGENT_PALETTE_IDS, AGENT_AVATAR_SIZES, CHARACTER_STATES, appearanceForPalette, type AgentAvatarSize } from "@paperclipai/shared"; @@ -35,7 +35,13 @@ export function agentAvatarRoutes(injected?: ReturnType value.trim().replace(/^W\//, "")); - if (validators?.some(value => value === "*" || value === etag)) { stream.destroy(); res.status(304).end(); return; } + if (validators?.some(value => value === "*" || value === etag)) { + // ReadStream opens asynchronously. Await disposal so a concurrent cache + // deletion cannot emit an unhandled open error after the 304 is sent. + stream.destroy(); + await finished(stream, { cleanup: true }).catch(() => {}); + res.status(304).end(); return; + } res.setHeader("Content-Length", byteSize); await pipeline(stream, res); } catch (error) {