From a13ce70a7ae33e263d1a9ed58d161758fbcc69fa Mon Sep 17 00:00:00 2001 From: Dotta Date: Thu, 10 Sep 2026 16:37:08 -0500 Subject: [PATCH] fix(agents): bound cold avatar admission and retain company drafts Co-Authored-By: Paperclip --- server/src/__tests__/agent-avatars.test.ts | 35 +++++++++++++++++- server/src/routes/agent-avatars.ts | 10 +++--- server/src/routes/openapi.ts | 1 + server/src/services/agent-avatar-pool.ts | 2 +- server/src/services/agent-avatars.ts | 41 ++++++++++++++++------ ui/src/components/AgentPersona.test.tsx | 11 ++++++ ui/src/hooks/useAgentAppearanceDraft.ts | 10 +++--- 7 files changed, 88 insertions(+), 22 deletions(-) diff --git a/server/src/__tests__/agent-avatars.test.ts b/server/src/__tests__/agent-avatars.test.ts index 4911e47031..d1cbb41136 100644 --- a/server/src/__tests__/agent-avatars.test.ts +++ b/server/src/__tests__/agent-avatars.test.ts @@ -8,7 +8,7 @@ import path from "node:path"; import express from "express"; import type { Server } from "node:http"; import sharp from "sharp"; -import { appearanceForPalette } from "@paperclipai/shared"; +import { AGENT_PALETTE_IDS, appearanceForPalette } from "@paperclipai/shared"; import { createLocalDiskStorageProvider } from "../storage/local-disk-provider.js"; import { createAgentAvatarService, avatarCacheKey, type AgentAvatarRequest } from "../services/agent-avatars.js"; import { createAgentAvatarPool } from "../services/agent-avatar-pool.js"; @@ -43,6 +43,39 @@ describe("on-demand agent avatars", () => { (await service.get(request)).stream.destroy(); expect(render).toHaveBeenCalledTimes(2); }); + it("limits cold keys per client while admitting warm hits, joiners and other clients", async () => { + let release!: () => void; + const blocked = new Promise(resolve => { release = resolve; }); + const render = vi.fn(async () => { await blocked; return Buffer.from("png"); }); + const provider = await storage(); + const service = createAgentAvatarService(provider, render); + const keys = AGENT_PALETTE_IDS.flatMap(palette => ([16, 20, 24] as const).map(size => ({ ...request, appearance: appearanceForPalette(palette), size }))); + const pending = keys.slice(0, 32).map(key => service.get(key, "one")); + try { + await vi.waitFor(() => expect(render).toHaveBeenCalledTimes(32)); + await expect(service.get(keys[32], "one")).rejects.toThrow("Too many cold avatar requests"); + pending.push(service.get(keys[0], "one")); // Same cold key is free. + 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(); + expect(render).toHaveBeenCalledTimes(33); + (await service.get(keys[33], "one")).stream.destroy(); // Completed renders release slots. + expect(render).toHaveBeenCalledTimes(34); + }); + it("returns retryable admission errors without caching them", async () => { + const service = createAgentAvatarService(await storage(), async () => Buffer.from("png")); + const { AvatarAdmissionError } = await import("../services/agent-avatars.js"); + vi.spyOn(service, "get").mockRejectedValueOnce(new AvatarAdmissionError(12)); + const url = await serve(service); + const denied = await fetch(url); + expect(denied.status).toBe(429); + expect(denied.headers.get("retry-after")).toBe("12"); + expect(denied.headers.get("cache-control")).toBe("no-store"); + await denied.text(); + const retry = await fetch(url); expect(retry.status).toBe(200); await retry.arrayBuffer(); + }); it("uses the configured S3 prefix and reuses bytes across service instances", async () => { const objects = new Map(); const send = vi.spyOn(S3Client.prototype, "send").mockImplementation(async (command: any) => { diff --git a/server/src/routes/agent-avatars.ts b/server/src/routes/agent-avatars.ts index 0b885f62d8..1e62473be6 100644 --- a/server/src/routes/agent-avatars.ts +++ b/server/src/routes/agent-avatars.ts @@ -2,7 +2,7 @@ import { 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"; -import { createAgentAvatarService } from "../services/agent-avatars.js"; +import { AvatarAdmissionError, createAgentAvatarService } from "../services/agent-avatars.js"; import { createStorageProviderFromConfig } from "../storage/provider-registry.js"; import { loadConfig } from "../config.js"; import { logger } from "../middleware/logger.js"; @@ -29,7 +29,7 @@ export function agentAvatarRoutes(injected?: ReturnType(); let closed = false; function spawn() { - const source = import.meta.url.endsWith(".ts"); + const source = new URL(import.meta.url).pathname.endsWith(".ts"); const url = new URL(source ? "./agent-avatar-worker.ts" : "./agent-avatar-worker.js", import.meta.url); const worker = source ? new Worker(`import(${JSON.stringify(import.meta.resolve('tsx/esm/api'))}).then(({tsImport}) => tsImport(${JSON.stringify(url.href)}, ${JSON.stringify(import.meta.url)}));`, { eval: true }) diff --git a/server/src/services/agent-avatars.ts b/server/src/services/agent-avatars.ts index 6e2f385d38..66715f0023 100644 --- a/server/src/services/agent-avatars.ts +++ b/server/src/services/agent-avatars.ts @@ -1,6 +1,7 @@ import { createHash } from "node:crypto"; import type { AgentAppearance, AgentAvatarSize, CharacterState } from "@paperclipai/shared"; import type { StorageProvider } from "../storage/types.js"; +import { createInviteRateLimiter } from "./invite-rate-limit.js"; import { createAgentAvatarPool } from "./agent-avatar-pool.js"; export interface AgentAvatarRequest { @@ -14,11 +15,18 @@ export function avatarCacheKey(request: AgentAvatarRequest) { const { appearance, size, scale, pose, muted } = request; return `generated-agent-avatars/${appearance.characterVersion}/${muted ? "muted-dream" : appearance.paletteId}/${pose}-${size}-${scale}.png`; } +export class AvatarAdmissionError extends Error { + constructor(public readonly retryAfterSeconds: number) { super("Too many cold avatar requests"); } +} type CacheMetadata = { sha256: string; byteSize: number }; export function createAgentAvatarService(storage: StorageProvider, render?: (request: AgentAvatarRequest) => Promise) { const pool = render ? undefined : createAgentAvatarPool(); const pending = new Map>(); - async function ensure(request: AgentAvatarRequest, key: string): Promise { + // Only cold keys consume admission; warm images and single-flight joiners + // remain available. Leave at least half the pool queue for other clients. + const activeByClient = new Map(); + const limiter = createInviteRateLimiter({ maxRequests: 256 }); + async function ensure(request: AgentAvatarRequest, key: string, client: string): Promise { const metadataKey = `${key}.json`; const [image, metadata] = await Promise.all([ storage.headObject({ objectKey: key }), storage.headObject({ objectKey: metadataKey }), @@ -32,21 +40,32 @@ export function createAgentAvatarService(storage: StorageProvider, render?: (req if (/^[a-f0-9]{64}$/.test(cached.sha256) && cached.byteSize > 0 && cached.byteSize === image.contentLength) return cached; } catch { /* Disposable metadata: regenerate a corrupt or old cache entry. */ } } - const bytes = await (render ?? pool!.render)(request); - const result = { sha256: createHash("sha256").update(bytes).digest("hex"), byteSize: bytes.length }; - // Both providers publish whole objects atomically. Publish metadata last so - // readers never consider an unfinished image a completed cache entry. - await storage.putObject({ objectKey: key, body: bytes, contentLength: bytes.length, contentType: "image/png" }); - const encoded = Buffer.from(JSON.stringify(result)); - await storage.putObject({ objectKey: metadataKey, body: encoded, contentLength: encoded.length, contentType: "application/json" }); - return result; + const active = activeByClient.get(client) ?? 0; + if (active >= 32) throw new AvatarAdmissionError(5); + const admission = limiter.consume(client); + if (!admission.allowed) throw new AvatarAdmissionError(admission.retryAfterSeconds); + activeByClient.set(client, active + 1); + try { + const bytes = await (render ?? pool!.render)(request); + const result = { sha256: createHash("sha256").update(bytes).digest("hex"), byteSize: bytes.length }; + // Both providers publish whole objects atomically. Publish metadata last so + // readers never consider an unfinished image a completed cache entry. + await storage.putObject({ objectKey: key, body: bytes, contentLength: bytes.length, contentType: "image/png" }); + const encoded = Buffer.from(JSON.stringify(result)); + await storage.putObject({ objectKey: metadataKey, body: encoded, contentLength: encoded.length, contentType: "application/json" }); + return result; + } finally { + const remaining = (activeByClient.get(client) ?? 1) - 1; + if (remaining) activeByClient.set(client, remaining); + else activeByClient.delete(client); + } } return { - async get(request: AgentAvatarRequest) { + async get(request: AgentAvatarRequest, client = "unknown") { const key = avatarCacheKey(request); let result = pending.get(key); if (!result) { - result = ensure(request, key).finally(() => pending.delete(key)); + result = ensure(request, key, client).finally(() => pending.delete(key)); pending.set(key, result); } const metadata = await result; diff --git a/ui/src/components/AgentPersona.test.tsx b/ui/src/components/AgentPersona.test.tsx index 5509930cef..340166f9a7 100644 --- a/ui/src/components/AgentPersona.test.tsx +++ b/ui/src/components/AgentPersona.test.tsx @@ -63,6 +63,17 @@ describe("agent persona presentation", () => { expect(createCharacter).not.toHaveBeenCalled(); expect(host.querySelectorAll("img")).toHaveLength(2); }); + it("selects the correct saved draft when the company changes without remounting", async () => { + let draft!: ReturnType; + function Draft({ company }: { company: string }) { draft = useAgentAppearanceDraft(`${company}:new-agent`); return null; } + sessionStorage.setItem("paperclip.agent-appearance.two:new-agent", JSON.stringify(appearance)); + await act(async () => root.render()); + const first = draft.appearance; + await act(async () => root.render()); + expect(draft.appearance).toEqual(appearance); + await act(async () => root.render()); + expect(draft.appearance).toEqual(first); + }); it("retains the draft assignment across remounts and clears it only after creation", async () => { let draft!: ReturnType; function Draft() { draft = useAgentAppearanceDraft("company:new-agent"); return null; } diff --git a/ui/src/hooks/useAgentAppearanceDraft.ts b/ui/src/hooks/useAgentAppearanceDraft.ts index e876e2db5f..8cbd6bc85a 100644 --- a/ui/src/hooks/useAgentAppearanceDraft.ts +++ b/ui/src/hooks/useAgentAppearanceDraft.ts @@ -1,10 +1,10 @@ import { useState } from "react"; import { agentAppearanceSchema, randomAgentAppearance } from "@paperclipai/shared"; -/** Non-secret visual identity only. The caller remounts when its draft key changes. */ +/** Non-secret visual identity, retained across navigation and company changes. */ export function useAgentAppearanceDraft(draftKey: string) { const key = `paperclip.agent-appearance.${draftKey}`; - const [appearance] = useState(() => { + function readDraft() { try { const stored = agentAppearanceSchema.safeParse(JSON.parse(sessionStorage.getItem(key) ?? "null")); if (stored.success) return stored.data; @@ -12,6 +12,8 @@ export function useAgentAppearanceDraft(draftKey: string) { const value = randomAgentAppearance(); try { sessionStorage.setItem(key, JSON.stringify(value)); } catch { /* In-memory draft still works. */ } return value; - }); - return { appearance, clear() { try { sessionStorage.removeItem(key); } catch { /* Best effort. */ } } }; + } + const [draft, setDraft] = useState(() => ({ key, appearance: readDraft() })); + if (draft.key !== key) setDraft({ key, appearance: readDraft() }); + return { appearance: draft.appearance, clear() { try { sessionStorage.removeItem(key); } catch { /* Best effort. */ } } }; }