From bd98e530ec316e77b9238ee4ac602827239e5ac7 Mon Sep 17 00:00:00 2001 From: Saul Carlin Date: Wed, 12 Aug 2026 00:11:01 -0700 Subject: [PATCH 1/3] feat: add Space upstream debug relay --- src/api/v2/agents/handlers/join-status.ts | 4 +- .../v2/conversations/conversations.router.ts | 14 +- .../conversations/handlers/space-upstream.ts | 242 ++++++++++ src/api/v2/index.ts | 6 +- src/middleware/rateLimit.ts | 14 + tests/space-upstream-production.test.ts | 33 ++ tests/space-upstream.test.ts | 439 ++++++++++++++++++ 7 files changed, 748 insertions(+), 4 deletions(-) create mode 100644 src/api/v2/conversations/handlers/space-upstream.ts create mode 100644 tests/space-upstream-production.test.ts create mode 100644 tests/space-upstream.test.ts diff --git a/src/api/v2/agents/handlers/join-status.ts b/src/api/v2/agents/handlers/join-status.ts index c75a9fdf..46746810 100644 --- a/src/api/v2/agents/handlers/join-status.ts +++ b/src/api/v2/agents/handlers/join-status.ts @@ -19,7 +19,7 @@ const paramsSchema = z.object({ // Optional dev-only variant routing hint. A malformed value (array, blank, // over-long) parses away to undefined and the poll falls back to the default // worker rather than 400ing — the status read still works. -const querySchema = z.object({ +export const joinStatusQuerySchema = z.object({ variantId: z.string().trim().min(1).max(64).optional(), }); @@ -85,7 +85,7 @@ export async function joinStatusHandler(req: Request, res: Response) { // the right runtime. Re-resolve the variant's ephemeral origin (dev-only, live + // allowlisted); anything else falls back to the default worker. let assistantBaseUrl = assistantApiUrl.replace(/\/+$/, ""); - const variantId = querySchema.safeParse(req.query).data?.variantId; + const variantId = joinStatusQuerySchema.safeParse(req.query).data?.variantId; if (variantId && XMTP_ENV !== "production") { const origin = await resolveVariantWorkerOrigin(variantId); if (origin) assistantBaseUrl = origin; diff --git a/src/api/v2/conversations/conversations.router.ts b/src/api/v2/conversations/conversations.router.ts index b4590429..3e8cc6e0 100644 --- a/src/api/v2/conversations/conversations.router.ts +++ b/src/api/v2/conversations/conversations.router.ts @@ -1,6 +1,9 @@ import { Router } from "express"; import { requireAccount } from "@/middleware/auth"; -import { agentParticipationLimiter } from "@/middleware/rateLimit"; +import { + agentParticipationLimiter, + spaceUpstreamLimiter, +} from "@/middleware/rateLimit"; import { conversationAbilitiesGetHandler } from "./handlers/abilities-get"; import { conversationAbilityDeleteHandler } from "./handlers/ability-delete"; import { conversationAbilityPutHandler } from "./handlers/ability-put"; @@ -8,12 +11,21 @@ import { getParticipationHandler, participationHandler, } from "./handlers/participation"; +import { spaceUpstreamHandler } from "./handlers/space-upstream"; // /v2/conversations — conversation-scoped surfaces. Mounted behind // authMiddleware in src/api/v2/index.ts; every route here applies // requireAccount itself. conversationId is the opaque XMTP string (no // Conversation table). export const conversationsRouter = Router(); +export const conversationsDebugRouter = Router(); + +conversationsDebugRouter.post( + "/:conversationId/debug/space-upstream", + spaceUpstreamLimiter, + requireAccount, + spaceUpstreamHandler, +); // How much the agents in this conversation may speak. `requireAccount` for the // same reason as /agents/join: an account-less JWT is an authorization failure, diff --git a/src/api/v2/conversations/handlers/space-upstream.ts b/src/api/v2/conversations/handlers/space-upstream.ts new file mode 100644 index 00000000..d09b3473 --- /dev/null +++ b/src/api/v2/conversations/handlers/space-upstream.ts @@ -0,0 +1,242 @@ +import type { Request, Response } from "express"; +import { z } from "zod"; +import { + getAssistantApiKey, + getAssistantApiUrl, +} from "@/api/v2/agents/handlers/assistant-config"; +import { joinStatusQuerySchema } from "@/api/v2/agents/handlers/join-status"; +import { resolveVariantWorkerOrigin } from "@/api/v2/agents/lib/variant-routing"; + +export const SPACE_UPSTREAM_FETCH_TIMEOUT_MS = 50_000; +const ERROR_BODY_LOG_LIMIT = 200; + +const conversationIdSchema = z + .string() + .regex(/^[0-9A-Za-z_-]{1,128}$/, "Invalid conversationId"); + +const paramsSchema = z.object({ + conversationId: conversationIdSchema, +}); + +const resultCountsSchema = { + wrote: z.number().int().nonnegative(), + deleted: z.number().int().nonnegative(), + refusedCount: z.number().int().nonnegative(), +}; + +export const spaceUpstreamResultSchema = z.discriminatedUnion("outcome", [ + z + .object({ + conversationId: conversationIdSchema, + outcome: z.literal("pull_request"), + prUrl: z.string().url(), + prNumber: z.number().int().positive(), + branch: z.string().min(1), + commitSha: z.string().min(1), + forkCommitSha: z.string().min(1), + ...resultCountsSchema, + }) + .strict(), + z + .object({ + conversationId: conversationIdSchema, + outcome: z.literal("unchanged"), + forkCommitSha: z.string().min(1), + ...resultCountsSchema, + }) + .strict(), +]); + +const upstreamErrorSchema = z + .object({ + error: z.string().min(1).max(500), + code: z.string().min(1).max(64), + }) + .strict(); + +type PublicError = { + status: number; + code: string; + error: string; +}; + +const ERRORS = { + INVALID_REQUEST: { + status: 400, + code: "INVALID_REQUEST", + error: "Invalid Space PR proposal request", + }, + VARIANT_UNAVAILABLE: { + status: 409, + code: "VARIANT_UNAVAILABLE", + error: "The selected agent variant is unavailable", + }, + SPACE_NOT_FOUND: { + status: 404, + code: "SPACE_NOT_FOUND", + error: "No Space was found for this conversation", + }, + SPACE_REPOSITORY_UNAVAILABLE: { + status: 409, + code: "SPACE_REPOSITORY_UNAVAILABLE", + error: "This Space does not have a repository", + }, + SPACE_UPSTREAM_NOT_ARMED: { + status: 503, + code: "SPACE_UPSTREAM_NOT_ARMED", + error: "The selected Space deployment is not armed for PR proposals", + }, + SPACE_UPSTREAM_UNAVAILABLE: { + status: 503, + code: "SPACE_UPSTREAM_UNAVAILABLE", + error: "Space PR proposals are unavailable", + }, + SPACE_UPSTREAM_REFUSED: { + status: 422, + code: "SPACE_UPSTREAM_REFUSED", + error: "The Space changes could not be proposed safely", + }, + SPACE_UPSTREAM_GITHUB_FAILED: { + status: 502, + code: "SPACE_UPSTREAM_GITHUB_FAILED", + error: "GitHub rejected the Space PR proposal; please try again", + }, + SPACE_UPSTREAM_FAILED: { + status: 502, + code: "SPACE_UPSTREAM_FAILED", + error: "The Space PR proposal failed", + }, + SPACE_UPSTREAM_TIMEOUT: { + status: 504, + code: "SPACE_UPSTREAM_TIMEOUT", + error: "The Space PR proposal timed out", + }, +} as const satisfies Record; + +function sendError(res: Response, value: PublicError): void { + const { status, ...body } = value; + res.status(status).json(body); +} + +function translateUpstreamError(status: number, raw: unknown): PublicError { + const parsed = upstreamErrorSchema.safeParse(raw); + if (!parsed.success) return ERRORS.SPACE_UPSTREAM_FAILED; + + const { code, error } = parsed.data; + if (status === 403 && code === "space_upstream_not_armed") { + return ERRORS.SPACE_UPSTREAM_NOT_ARMED; + } + if (status === 404 && code === "space_not_found") { + return ERRORS.SPACE_NOT_FOUND; + } + if (status === 409 && code === "space_repository_unavailable") { + return ERRORS.SPACE_REPOSITORY_UNAVAILABLE; + } + if (status === 503 && code === "space_repository_provider_unavailable") { + return ERRORS.SPACE_UPSTREAM_UNAVAILABLE; + } + if (status === 422 && code === "space_upstream_refused") { + return { ...ERRORS.SPACE_UPSTREAM_REFUSED, error }; + } + if (status === 502 && code === "space_upstream_github_failed") { + return ERRORS.SPACE_UPSTREAM_GITHUB_FAILED; + } + if (status === 502 && code === "space_upstream_failed") { + return ERRORS.SPACE_UPSTREAM_FAILED; + } + if (status === 504 && code === "space_upstream_timeout") { + return ERRORS.SPACE_UPSTREAM_TIMEOUT; + } + return ERRORS.SPACE_UPSTREAM_FAILED; +} + +export async function spaceUpstreamHandler(req: Request, res: Response) { + const parsedParams = paramsSchema.safeParse(req.params); + const parsedQuery = joinStatusQuerySchema.safeParse(req.query); + if (!parsedParams.success || !parsedQuery.success) { + sendError(res, ERRORS.INVALID_REQUEST); + return; + } + + const conversationId = parsedParams.data.conversationId.toLowerCase(); + const variantId = parsedQuery.data.variantId; + + let assistantOrigin: string; + if (variantId !== undefined) { + const resolvedOrigin = await resolveVariantWorkerOrigin(variantId); + if (!resolvedOrigin) { + sendError(res, ERRORS.VARIANT_UNAVAILABLE); + return; + } + assistantOrigin = resolvedOrigin; + } else { + assistantOrigin = getAssistantApiUrl(); + } + + const assistantApiKey = getAssistantApiKey().trim(); + const assistantBaseUrl = assistantOrigin.trim().replace(/\/+$/, ""); + if (!assistantApiKey || !assistantBaseUrl) { + req.log.error("Space upstream Worker is not configured"); + sendError(res, ERRORS.SPACE_UPSTREAM_UNAVAILABLE); + return; + } + + try { + const upstream = await fetch( + `${assistantBaseUrl}/api/conversations/${encodeURIComponent(conversationId)}/space-upstream`, + { + method: "POST", + headers: { Authorization: `Bearer ${assistantApiKey}` }, + signal: AbortSignal.timeout(SPACE_UPSTREAM_FETCH_TIMEOUT_MS), + }, + ); + + if (!upstream.ok) { + const text = await upstream.text(); + const bodyPreview = text.substring(0, ERROR_BODY_LOG_LIMIT); + req.log.error( + { status: upstream.status, bodyPreview }, + "Space upstream Worker request failed", + ); + + let raw: unknown; + try { + raw = JSON.parse(text); + } catch { + raw = null; + } + sendError(res, translateUpstreamError(upstream.status, raw)); + return; + } + + let raw: unknown; + try { + raw = await upstream.json(); + } catch { + raw = null; + } + const result = spaceUpstreamResultSchema.safeParse(raw); + if (!result.success) { + req.log.error( + { issues: result.error.issues }, + "Invalid Space upstream Worker response", + ); + sendError(res, ERRORS.SPACE_UPSTREAM_FAILED); + return; + } + + res.status(200).json(result.data); + } catch (error) { + if (error instanceof DOMException && error.name === "TimeoutError") { + req.log.error("Space upstream Worker request timed out"); + sendError(res, ERRORS.SPACE_UPSTREAM_TIMEOUT); + return; + } + + req.log.error( + { error, stack: error instanceof Error ? error.stack : undefined }, + "Space upstream Worker request failed", + ); + sendError(res, ERRORS.SPACE_UPSTREAM_FAILED); + } +} diff --git a/src/api/v2/index.ts b/src/api/v2/index.ts index fbb5f691..109b8349 100644 --- a/src/api/v2/index.ts +++ b/src/api/v2/index.ts @@ -40,7 +40,10 @@ import { composioRouter } from "./composio/composio.router"; import { connectionsRouter } from "./connections/connections.router"; import { actionsGetHandler } from "./connections/handlers/actions-get"; import { servicesGetHandler } from "./connections/handlers/services-get"; -import { conversationsRouter } from "./conversations/conversations.router"; +import { + conversationsDebugRouter, + conversationsRouter, +} from "./conversations/conversations.router"; import { creditsAdminRouter } from "./credits-admin/credits-admin.router"; import { dailyRefillRouter } from "./credits/daily.router"; import { devRouter } from "./dev/dev.router"; @@ -61,6 +64,7 @@ const v2Router = Router(); // /dev is a non-production test surface; keep it gated. if (process.env.XMTP_ENV !== "production") { v2Router.use("/dev", devAuthMiddleware, devRouter); + v2Router.use("/conversations", authMiddleware, conversationsDebugRouter); } v2Router.use("/agent-prompt-hints", agentPromptHintsRouter); diff --git a/src/middleware/rateLimit.ts b/src/middleware/rateLimit.ts index afd26705..f46f3995 100644 --- a/src/middleware/rateLimit.ts +++ b/src/middleware/rateLimit.ts @@ -57,6 +57,20 @@ export const agentParticipationLimiter = rateLimit({ }, }); +// Space-to-starter proposals can update a GitHub branch and draft pull request, +// so keep retries bounded independently of the cheaper participation controls. +export const spaceUpstreamLimiter = rateLimit({ + windowMs: 5 * 60 * 1000, + limit: 10, + keyGenerator: (req) => req.ip || "unknown", + legacyHeaders: false, + standardHeaders: "draft-8", + message: { + code: "RATE_LIMITED", + error: "Too many Space PR proposals; retry shortly", + }, +}); + // Rate limiting for asset renewal endpoint (10 batch requests per hour per device) export const assetRenewalLimiter = rateLimit({ windowMs: 60 * 60 * 1000, // 1 hour diff --git a/tests/space-upstream-production.test.ts b/tests/space-upstream-production.test.ts new file mode 100644 index 00000000..d24b19d3 --- /dev/null +++ b/tests/space-upstream-production.test.ts @@ -0,0 +1,33 @@ +import express from "express"; +import request from "supertest"; +import { afterEach, describe, expect, test, vi } from "vitest"; + +const originalXmtpEnv = process.env.XMTP_ENV; + +afterEach(() => { + process.env.XMTP_ENV = originalXmtpEnv; + vi.resetModules(); +}); + +describe("Space upstream production mount guard", () => { + test("does not mount the debug route in production", async () => { + process.env.XMTP_ENV = "production"; + vi.resetModules(); + const { default: v2Router } = await import("@/api/v2"); + const { pinoMiddleware } = await import("@/middleware/pino"); + const { createJwtToken, validateJWTKeys } = await import("@/utils/jwt"); + await validateJWTKeys(); + const token = await createJwtToken({ + deviceId: "production-mount-test", + accountId: "11111111-1111-4111-8111-111111111111", + }); + const app = express(); + app.use(pinoMiddleware); + app.use("/api/v2", v2Router); + + const response = await request(app) + .post("/api/v2/conversations/conversation_abc/debug/space-upstream") + .set("X-Convos-AuthToken", token); + expect(response.status).toBe(404); + }); +}); diff --git a/tests/space-upstream.test.ts b/tests/space-upstream.test.ts new file mode 100644 index 00000000..96099470 --- /dev/null +++ b/tests/space-upstream.test.ts @@ -0,0 +1,439 @@ +import express, { + type Request as ExpressRequest, + type Response as ExpressResponse, + type NextFunction, +} from "express"; +import request from "supertest"; +import { + afterAll, + beforeAll, + beforeEach, + describe, + expect, + test, + vi, +} from "vitest"; +import { __setAssistantConfigOverridesForTests } from "@/api/v2/agents/handlers/assistant-config"; +import { conversationsDebugRouter } from "@/api/v2/conversations/conversations.router"; +import { authMiddleware } from "@/middleware/auth"; +import { pinoMiddleware } from "@/middleware/pino"; +import { createJwtToken, validateJWTKeys } from "@/utils/jwt"; +import { prisma } from "@/utils/prisma"; + +const DEFAULT_URL = "https://assistants.test.local"; +const ASSISTANT_KEY = "test-space-upstream-key"; +const ACCOUNT_ID = "11111111-1111-4111-8111-111111111111"; +const GOOD_VARIANT = "pr-test-space-upstream"; +const GOOD_VARIANT_URL = `https://ephemeral-${GOOD_VARIANT}.convos.fun`; +const OFF_HOST_VARIANT = "pr-test-space-upstream-off-host"; +const MISMATCH_VARIANT = "pr-test-space-upstream-mismatch"; + +const pullRequestResult = { + conversationId: "conversation_abc", + outcome: "pull_request", + prUrl: "https://github.com/xmtplabs/convos-assistants/pull/123", + prNumber: 123, + branch: "space-upstream/conversation_abc", + commitSha: "commit-sha", + forkCommitSha: "fork-commit-sha", + wrote: 4, + deleted: 1, + refusedCount: 2, +} as const; + +const unchangedResult = { + conversationId: "conversation_abc", + outcome: "unchanged", + forkCommitSha: "fork-commit-sha", + wrote: 0, + deleted: 0, + refusedCount: 0, +} as const; + +type FetchCall = { url: string; init?: RequestInit }; +let fetchCalls: FetchCall[] = []; +let fetchImpl: (url: string, init?: RequestInit) => Promise; +const originalFetch = globalThis.fetch; +let nextIpOctet = 1; + +function jsonResponse(status: number, body: unknown): Response { + return new Response(JSON.stringify(body), { + status, + headers: { "Content-Type": "application/json" }, + }); +} + +function accountMiddleware( + _req: ExpressRequest, + res: ExpressResponse, + next: NextFunction, +) { + res.locals.accountId = ACCOUNT_ID; + next(); +} + +function buildApp(options?: { + auth?: boolean; + errorLog?: ReturnType; +}) { + const app = express(); + app.set("trust proxy", 1); + if (options?.errorLog) { + app.use((req, _res, next) => { + req.log = { + error: options.errorLog, + warn: vi.fn(), + info: vi.fn(), + } as unknown as ExpressRequest["log"]; + next(); + }); + } else { + app.use(pinoMiddleware); + } + app.use( + "/api/v2/conversations", + options?.auth ? authMiddleware : accountMiddleware, + conversationsDebugRouter, + ); + return app; +} + +function proposal( + app: ReturnType, + path = "/api/v2/conversations/CONVERSATION_ABC/debug/space-upstream", + ip?: string, +) { + const selectedIp = ip ?? `198.51.100.${nextIpOctet++}`; + return request(app).post(path).set("X-Forwarded-For", selectedIp); +} + +function responseBody(response: { body: unknown }): Record { + return response.body as Record; +} + +beforeAll(async () => { + await validateJWTKeys(); +}); + +afterAll(() => { + globalThis.fetch = originalFetch; + __setAssistantConfigOverridesForTests({}); +}); + +beforeEach(() => { + vi.restoreAllMocks(); + fetchCalls = []; + __setAssistantConfigOverridesForTests({ + assistantApiUrl: DEFAULT_URL, + assistantApiKey: ASSISTANT_KEY, + }); + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve(jsonResponse(200, pullRequestResult)); + }; + globalThis.fetch = (url: string | URL | Request, init?: RequestInit) => { + const urlString = + typeof url === "string" ? url : url instanceof URL ? url.href : url.url; + return fetchImpl(urlString, init); + }; + vi.spyOn(prisma.agentVariant, "findFirst").mockImplementation((args) => { + const slug = (args?.where as { slug?: string } | undefined)?.slug; + const assistantWorkerUrl = + slug === GOOD_VARIANT + ? GOOD_VARIANT_URL + : slug === OFF_HOST_VARIANT + ? "https://evil.example.com" + : slug === MISMATCH_VARIANT + ? "https://ephemeral-pr-test-space-upstream-other.convos.fun" + : null; + return Promise.resolve( + assistantWorkerUrl ? { assistantWorkerUrl } : null, + ) as never; + }); +}); + +describe("POST /conversations/:conversationId/debug/space-upstream", () => { + test("uses JWT auth and requires an account", async () => { + const app = buildApp({ auth: true }); + + const missing = await proposal(app); + expect(missing.status).toBe(401); + + const accountlessToken = await createJwtToken({ deviceId: "device-only" }); + const accountless = await proposal(app).set( + "X-Convos-AuthToken", + accountlessToken, + ); + expect(accountless.status).toBe(403); + expect(accountless.body).toEqual({ error: "Account required" }); + + const accountToken = await createJwtToken({ + deviceId: "device-account", + accountId: ACCOUNT_ID, + }); + const authenticated = await proposal(app).set( + "X-Convos-AuthToken", + accountToken, + ); + expect(authenticated.status).toBe(200); + expect(fetchCalls).toHaveLength(1); + }); + + test("normalizes the bounded conversation ID and sends only the shared key", async () => { + const res = await proposal(buildApp()); + expect(res.status).toBe(200); + expect(fetchCalls).toHaveLength(1); + expect(fetchCalls[0]?.url).toBe( + `${DEFAULT_URL}/api/conversations/conversation_abc/space-upstream`, + ); + expect(fetchCalls[0]?.init).toMatchObject({ + method: "POST", + headers: { Authorization: `Bearer ${ASSISTANT_KEY}` }, + }); + expect(fetchCalls[0]?.init?.body).toBeUndefined(); + expect(fetchCalls[0]?.init?.headers).toEqual({ + Authorization: `Bearer ${ASSISTANT_KEY}`, + }); + }); + + test.each([ + ["invalid characters", "bad%20conversation"], + ["overlong", "a".repeat(129)], + ])("rejects an %s conversation ID before fetch", async (_label, id) => { + const res = await proposal( + buildApp(), + `/api/v2/conversations/${id}/debug/space-upstream`, + ); + expect(res.status).toBe(400); + expect(res.body).toEqual({ + code: "INVALID_REQUEST", + error: "Invalid Space PR proposal request", + }); + expect(fetchCalls).toHaveLength(0); + }); + + test.each([ + ["empty", "variantId="], + ["array", "variantId=one&variantId=two"], + ["overlong", `variantId=${"a".repeat(65)}`], + ])("rejects an %s provided variant", async (_label, query) => { + const res = await proposal( + buildApp(), + `/api/v2/conversations/conversation_abc/debug/space-upstream?${query}`, + ); + expect(res.status).toBe(400); + expect(responseBody(res).code).toBe("INVALID_REQUEST"); + expect(fetchCalls).toHaveLength(0); + }); + + test("ignores unrelated query keys and uses the default Worker", async () => { + const res = await proposal( + buildApp(), + "/api/v2/conversations/conversation_abc/debug/space-upstream?future=value", + ); + expect(res.status).toBe(200); + expect(fetchCalls[0]?.url).toBe( + `${DEFAULT_URL}/api/conversations/conversation_abc/space-upstream`, + ); + }); + + test("passes the parsed variant slug to the registry and uses its exact origin", async () => { + const res = await proposal( + buildApp(), + `/api/v2/conversations/conversation_abc/debug/space-upstream?variantId=${GOOD_VARIANT}`, + ); + expect(res.status).toBe(200); + expect(fetchCalls[0]?.url).toBe( + `${GOOD_VARIANT_URL}/api/conversations/conversation_abc/space-upstream`, + ); + }); + + test.each([OFF_HOST_VARIANT, MISMATCH_VARIANT, "unknown-space-variant"])( + "fails a non-allowed variant closed without fetching (%s)", + async (variantId) => { + const res = await proposal( + buildApp(), + `/api/v2/conversations/conversation_abc/debug/space-upstream?variantId=${variantId}`, + ); + expect(res.status).toBe(409); + expect(res.body).toEqual({ + code: "VARIANT_UNAVAILABLE", + error: "The selected agent variant is unavailable", + }); + expect(fetchCalls).toHaveLength(0); + }, + ); + + test("uses a 50-second upstream AbortSignal", async () => { + const timeoutSpy = vi.spyOn(AbortSignal, "timeout"); + try { + const res = await proposal(buildApp()); + expect(res.status).toBe(200); + expect(timeoutSpy).toHaveBeenCalledWith(50_000); + } finally { + timeoutSpy.mockRestore(); + } + }); + + test.each([ + ["empty shared key", { assistantApiKey: "" }], + ["whitespace shared key", { assistantApiKey: " " }], + ["empty default origin", { assistantApiUrl: "" }], + ["whitespace default origin", { assistantApiUrl: " " }], + ])("returns unavailable for %s", async (_label, override) => { + __setAssistantConfigOverridesForTests({ + assistantApiUrl: DEFAULT_URL, + assistantApiKey: ASSISTANT_KEY, + ...override, + }); + const res = await proposal(buildApp()); + expect(res.status).toBe(503); + expect(responseBody(res).code).toBe("SPACE_UPSTREAM_UNAVAILABLE"); + expect(fetchCalls).toHaveLength(0); + }); + + test.each([ + [403, "space_upstream_not_armed", 503, "SPACE_UPSTREAM_NOT_ARMED"], + [404, "space_not_found", 404, "SPACE_NOT_FOUND"], + [409, "space_repository_unavailable", 409, "SPACE_REPOSITORY_UNAVAILABLE"], + [ + 503, + "space_repository_provider_unavailable", + 503, + "SPACE_UPSTREAM_UNAVAILABLE", + ], + [422, "space_upstream_refused", 422, "SPACE_UPSTREAM_REFUSED"], + [502, "space_upstream_github_failed", 502, "SPACE_UPSTREAM_GITHUB_FAILED"], + [502, "space_upstream_failed", 502, "SPACE_UPSTREAM_FAILED"], + [504, "space_upstream_timeout", 504, "SPACE_UPSTREAM_TIMEOUT"], + [401, "unauthorized", 502, "SPACE_UPSTREAM_FAILED"], + ])( + "maps Worker %i %s to %i %s", + async (workerStatus, workerCode, expectedStatus, expectedCode) => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve( + jsonResponse(workerStatus, { + error: "Safe upstream detail", + code: workerCode, + }), + ); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(expectedStatus); + expect(responseBody(res).code).toBe(expectedCode); + if (workerCode === "space_upstream_refused") { + expect(responseBody(res).error).toBe("Safe upstream detail"); + } + }, + ); + + test.each([ + ["uncoded old-route 404", 404, { error: "Not found" }], + ["unexpected coded status", 418, { error: "No", code: "unexpected" }], + [ + "extra error envelope fields", + 404, + { error: "Not found", code: "space_not_found", extra: true }, + ], + ])("maps %s to the generic failure", async (_label, status, body) => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve(jsonResponse(status, body)); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(502); + expect(responseBody(res).code).toBe("SPACE_UPSTREAM_FAILED"); + }); + + test("separates timeout and network failures", async () => { + fetchImpl = () => + Promise.reject(new DOMException("Timed out", "TimeoutError")); + const timedOut = await proposal(buildApp()); + expect(timedOut.status).toBe(504); + expect(responseBody(timedOut).code).toBe("SPACE_UPSTREAM_TIMEOUT"); + + fetchImpl = () => Promise.reject(new TypeError("network unavailable")); + const networkFailure = await proposal(buildApp()); + expect(networkFailure.status).toBe(502); + expect(responseBody(networkFailure).code).toBe("SPACE_UPSTREAM_FAILED"); + }); + + test.each([pullRequestResult, unchangedResult])( + "returns a valid $outcome result transparently", + async (result) => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve(jsonResponse(200, result)); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(200); + expect(res.body).toEqual(result); + }, + ); + + test.each([ + ["malformed JSON", "not-json"], + ["unknown outcome", { ...unchangedResult, outcome: "queued" }], + [ + "missing required field", + { ...unchangedResult, forkCommitSha: undefined }, + ], + ["extra success field", { ...unchangedResult, extra: true }], + [ + "mixed union fields", + { ...unchangedResult, prUrl: "https://example.com" }, + ], + ])("rejects a %s success response", async (_label, body) => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + if (typeof body === "string") { + return Promise.resolve( + new Response(body, { + status: 200, + headers: { "Content-Type": "application/json" }, + }), + ); + } + return Promise.resolve(jsonResponse(200, body)); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(502); + expect(responseBody(res).code).toBe("SPACE_UPSTREAM_FAILED"); + }); + + test("caps upstream body diagnostics", async () => { + const errorLog = vi.fn(); + fetchImpl = () => + Promise.resolve( + new Response("x".repeat(1_000), { + status: 500, + headers: { "Content-Type": "text/plain" }, + }), + ); + const res = await proposal(buildApp({ errorLog })); + expect(res.status).toBe(502); + expect(errorLog).toHaveBeenCalledWith( + { status: 500, bodyPreview: "x".repeat(200) }, + "Space upstream Worker request failed", + ); + }); + + test("returns the exact IP-keyed mutation limit response", async () => { + const app = buildApp(); + const ip = "203.0.113.77"; + for (let index = 0; index < 10; index += 1) { + const allowed = await proposal(app, undefined, ip); + expect(allowed.status).toBe(200); + } + const limited = await proposal(app, undefined, ip); + expect(limited.status).toBe(429); + expect(limited.body).toEqual({ + code: "RATE_LIMITED", + error: "Too many Space PR proposals; retry shortly", + }); + + const otherIp = await proposal(app, undefined, "203.0.113.78"); + expect(otherIp.status).toBe(200); + }); +}); From a90c53e74078d312467a411c9bff56c1fc6268ac Mon Sep 17 00:00:00 2001 From: Saul Carlin Date: Wed, 12 Aug 2026 00:36:54 -0700 Subject: [PATCH 2/3] refactor: align Space upstream relay with v2 idioms --- src/api/v2/agents/handlers/join-status.ts | 4 +- .../conversations/handlers/space-upstream.ts | 171 +++++++++--------- src/middleware/rateLimit.ts | 6 +- ...sations-space-upstream-production.test.ts} | 0 ...s => conversations-space-upstream.test.ts} | 132 +++++++------- 5 files changed, 161 insertions(+), 152 deletions(-) rename tests/{space-upstream-production.test.ts => conversations-space-upstream-production.test.ts} (100%) rename tests/{space-upstream.test.ts => conversations-space-upstream.test.ts} (80%) diff --git a/src/api/v2/agents/handlers/join-status.ts b/src/api/v2/agents/handlers/join-status.ts index 46746810..c75a9fdf 100644 --- a/src/api/v2/agents/handlers/join-status.ts +++ b/src/api/v2/agents/handlers/join-status.ts @@ -19,7 +19,7 @@ const paramsSchema = z.object({ // Optional dev-only variant routing hint. A malformed value (array, blank, // over-long) parses away to undefined and the poll falls back to the default // worker rather than 400ing — the status read still works. -export const joinStatusQuerySchema = z.object({ +const querySchema = z.object({ variantId: z.string().trim().min(1).max(64).optional(), }); @@ -85,7 +85,7 @@ export async function joinStatusHandler(req: Request, res: Response) { // the right runtime. Re-resolve the variant's ephemeral origin (dev-only, live + // allowlisted); anything else falls back to the default worker. let assistantBaseUrl = assistantApiUrl.replace(/\/+$/, ""); - const variantId = joinStatusQuerySchema.safeParse(req.query).data?.variantId; + const variantId = querySchema.safeParse(req.query).data?.variantId; if (variantId && XMTP_ENV !== "production") { const origin = await resolveVariantWorkerOrigin(variantId); if (origin) assistantBaseUrl = origin; diff --git a/src/api/v2/conversations/handlers/space-upstream.ts b/src/api/v2/conversations/handlers/space-upstream.ts index d09b3473..3d20cd37 100644 --- a/src/api/v2/conversations/handlers/space-upstream.ts +++ b/src/api/v2/conversations/handlers/space-upstream.ts @@ -4,165 +4,164 @@ import { getAssistantApiKey, getAssistantApiUrl, } from "@/api/v2/agents/handlers/assistant-config"; -import { joinStatusQuerySchema } from "@/api/v2/agents/handlers/join-status"; import { resolveVariantWorkerOrigin } from "@/api/v2/agents/lib/variant-routing"; -export const SPACE_UPSTREAM_FETCH_TIMEOUT_MS = 50_000; +const SPACE_UPSTREAM_FETCH_TIMEOUT_MS = 50_000; const ERROR_BODY_LOG_LIMIT = 200; const conversationIdSchema = z .string() - .regex(/^[0-9A-Za-z_-]{1,128}$/, "Invalid conversationId"); + .trim() + .min(1, "conversationId is required") + .max(256); const paramsSchema = z.object({ conversationId: conversationIdSchema, }); +const querySchema = z.object({ + variantId: z.string().trim().min(1).max(64).optional(), +}); + const resultCountsSchema = { wrote: z.number().int().nonnegative(), deleted: z.number().int().nonnegative(), refusedCount: z.number().int().nonnegative(), }; -export const spaceUpstreamResultSchema = z.discriminatedUnion("outcome", [ - z - .object({ - conversationId: conversationIdSchema, - outcome: z.literal("pull_request"), - prUrl: z.string().url(), - prNumber: z.number().int().positive(), - branch: z.string().min(1), - commitSha: z.string().min(1), - forkCommitSha: z.string().min(1), - ...resultCountsSchema, - }) - .strict(), - z - .object({ - conversationId: conversationIdSchema, - outcome: z.literal("unchanged"), - forkCommitSha: z.string().min(1), - ...resultCountsSchema, - }) - .strict(), +const spaceUpstreamResultSchema = z.discriminatedUnion("outcome", [ + z.object({ + conversationId: conversationIdSchema, + outcome: z.literal("pull_request"), + prUrl: z.string().url(), + prNumber: z.number().int().positive(), + branch: z.string().min(1), + commitSha: z.string().min(1), + forkCommitSha: z.string().min(1), + ...resultCountsSchema, + }), + z.object({ + conversationId: conversationIdSchema, + outcome: z.literal("unchanged"), + forkCommitSha: z.string().min(1), + ...resultCountsSchema, + }), ]); -const upstreamErrorSchema = z - .object({ - error: z.string().min(1).max(500), - code: z.string().min(1).max(64), - }) - .strict(); +const upstreamErrorSchema = z.object({ + error: z.string().min(1).max(500), + code: z.string().min(1).max(64), +}); type PublicError = { status: number; - code: string; error: string; + message: string; }; const ERRORS = { INVALID_REQUEST: { status: 400, - code: "INVALID_REQUEST", - error: "Invalid Space PR proposal request", + error: "INVALID_REQUEST", + message: "Invalid Space PR proposal request", }, VARIANT_UNAVAILABLE: { status: 409, - code: "VARIANT_UNAVAILABLE", - error: "The selected agent variant is unavailable", + error: "VARIANT_UNAVAILABLE", + message: "The selected agent variant is unavailable", }, SPACE_NOT_FOUND: { status: 404, - code: "SPACE_NOT_FOUND", - error: "No Space was found for this conversation", + error: "SPACE_NOT_FOUND", + message: "No Space was found for this conversation", }, SPACE_REPOSITORY_UNAVAILABLE: { status: 409, - code: "SPACE_REPOSITORY_UNAVAILABLE", - error: "This Space does not have a repository", + error: "SPACE_REPOSITORY_UNAVAILABLE", + message: "This Space does not have a repository", }, SPACE_UPSTREAM_NOT_ARMED: { status: 503, - code: "SPACE_UPSTREAM_NOT_ARMED", - error: "The selected Space deployment is not armed for PR proposals", + error: "SPACE_UPSTREAM_NOT_ARMED", + message: "The selected Space deployment is not armed for PR proposals", }, SPACE_UPSTREAM_UNAVAILABLE: { status: 503, - code: "SPACE_UPSTREAM_UNAVAILABLE", - error: "Space PR proposals are unavailable", + error: "SPACE_UPSTREAM_UNAVAILABLE", + message: "Space PR proposals are unavailable", }, SPACE_UPSTREAM_REFUSED: { status: 422, - code: "SPACE_UPSTREAM_REFUSED", - error: "The Space changes could not be proposed safely", + error: "SPACE_UPSTREAM_REFUSED", + message: "The Space changes could not be proposed safely", }, SPACE_UPSTREAM_GITHUB_FAILED: { status: 502, - code: "SPACE_UPSTREAM_GITHUB_FAILED", - error: "GitHub rejected the Space PR proposal; please try again", + error: "SPACE_UPSTREAM_GITHUB_FAILED", + message: "GitHub rejected the Space PR proposal; please try again", }, SPACE_UPSTREAM_FAILED: { status: 502, - code: "SPACE_UPSTREAM_FAILED", - error: "The Space PR proposal failed", + error: "SPACE_UPSTREAM_FAILED", + message: "The Space PR proposal failed", }, SPACE_UPSTREAM_TIMEOUT: { status: 504, - code: "SPACE_UPSTREAM_TIMEOUT", - error: "The Space PR proposal timed out", + error: "SPACE_UPSTREAM_TIMEOUT", + message: "The Space PR proposal timed out", }, } as const satisfies Record; +const UPSTREAM_ERRORS = { + space_upstream_not_armed: ERRORS.SPACE_UPSTREAM_NOT_ARMED, + space_not_found: ERRORS.SPACE_NOT_FOUND, + space_repository_unavailable: ERRORS.SPACE_REPOSITORY_UNAVAILABLE, + space_repository_provider_unavailable: ERRORS.SPACE_UPSTREAM_UNAVAILABLE, + space_upstream_refused: ERRORS.SPACE_UPSTREAM_REFUSED, + space_upstream_github_failed: ERRORS.SPACE_UPSTREAM_GITHUB_FAILED, + space_upstream_failed: ERRORS.SPACE_UPSTREAM_FAILED, + space_upstream_timeout: ERRORS.SPACE_UPSTREAM_TIMEOUT, +} as const satisfies Record; + function sendError(res: Response, value: PublicError): void { const { status, ...body } = value; - res.status(status).json(body); + res.status(status).json({ success: false, ...body }); } -function translateUpstreamError(status: number, raw: unknown): PublicError { +function translateUpstreamError(raw: unknown): PublicError { const parsed = upstreamErrorSchema.safeParse(raw); if (!parsed.success) return ERRORS.SPACE_UPSTREAM_FAILED; - const { code, error } = parsed.data; - if (status === 403 && code === "space_upstream_not_armed") { - return ERRORS.SPACE_UPSTREAM_NOT_ARMED; - } - if (status === 404 && code === "space_not_found") { - return ERRORS.SPACE_NOT_FOUND; - } - if (status === 409 && code === "space_repository_unavailable") { - return ERRORS.SPACE_REPOSITORY_UNAVAILABLE; - } - if (status === 503 && code === "space_repository_provider_unavailable") { - return ERRORS.SPACE_UPSTREAM_UNAVAILABLE; - } - if (status === 422 && code === "space_upstream_refused") { - return { ...ERRORS.SPACE_UPSTREAM_REFUSED, error }; - } - if (status === 502 && code === "space_upstream_github_failed") { - return ERRORS.SPACE_UPSTREAM_GITHUB_FAILED; - } - if (status === 502 && code === "space_upstream_failed") { - return ERRORS.SPACE_UPSTREAM_FAILED; - } - if (status === 504 && code === "space_upstream_timeout") { - return ERRORS.SPACE_UPSTREAM_TIMEOUT; - } - return ERRORS.SPACE_UPSTREAM_FAILED; + const { code, error: message } = parsed.data; + const publicError = UPSTREAM_ERRORS[code as keyof typeof UPSTREAM_ERRORS]; + if (!publicError) return ERRORS.SPACE_UPSTREAM_FAILED; + return code === "space_upstream_refused" + ? { ...publicError, message } + : publicError; } +/** + * Handler for POST /api/v2/conversations/:conversationId/debug/space-upstream + * + * Relays an authenticated, non-production Space PR proposal to the assistant + * Worker. The client never receives the shared Worker credential; it receives + * the standard v2 success or coded-error envelope instead. + */ export async function spaceUpstreamHandler(req: Request, res: Response) { const parsedParams = paramsSchema.safeParse(req.params); - const parsedQuery = joinStatusQuerySchema.safeParse(req.query); + const parsedQuery = querySchema.safeParse(req.query); if (!parsedParams.success || !parsedQuery.success) { sendError(res, ERRORS.INVALID_REQUEST); return; } - const conversationId = parsedParams.data.conversationId.toLowerCase(); + const conversationId = parsedParams.data.conversationId; const variantId = parsedQuery.data.variantId; let assistantOrigin: string; if (variantId !== undefined) { + // This mutation can create a GitHub branch and PR from variant-specific + // code, so it must not silently fall back to the default Worker. const resolvedOrigin = await resolveVariantWorkerOrigin(variantId); if (!resolvedOrigin) { sendError(res, ERRORS.VARIANT_UNAVAILABLE); @@ -173,9 +172,9 @@ export async function spaceUpstreamHandler(req: Request, res: Response) { assistantOrigin = getAssistantApiUrl(); } - const assistantApiKey = getAssistantApiKey().trim(); - const assistantBaseUrl = assistantOrigin.trim().replace(/\/+$/, ""); - if (!assistantApiKey || !assistantBaseUrl) { + const assistantApiKey = getAssistantApiKey(); + const assistantBaseUrl = assistantOrigin.replace(/\/+$/, ""); + if (!assistantApiKey) { req.log.error("Space upstream Worker is not configured"); sendError(res, ERRORS.SPACE_UPSTREAM_UNAVAILABLE); return; @@ -205,7 +204,7 @@ export async function spaceUpstreamHandler(req: Request, res: Response) { } catch { raw = null; } - sendError(res, translateUpstreamError(upstream.status, raw)); + sendError(res, translateUpstreamError(raw)); return; } @@ -225,7 +224,7 @@ export async function spaceUpstreamHandler(req: Request, res: Response) { return; } - res.status(200).json(result.data); + res.status(200).json({ success: true, ...result.data }); } catch (error) { if (error instanceof DOMException && error.name === "TimeoutError") { req.log.error("Space upstream Worker request timed out"); diff --git a/src/middleware/rateLimit.ts b/src/middleware/rateLimit.ts index f46f3995..b7b4550d 100644 --- a/src/middleware/rateLimit.ts +++ b/src/middleware/rateLimit.ts @@ -62,12 +62,12 @@ export const agentParticipationLimiter = rateLimit({ export const spaceUpstreamLimiter = rateLimit({ windowMs: 5 * 60 * 1000, limit: 10, - keyGenerator: (req) => req.ip || "unknown", legacyHeaders: false, standardHeaders: "draft-8", message: { - code: "RATE_LIMITED", - error: "Too many Space PR proposals; retry shortly", + success: false, + error: "RATE_LIMITED", + message: "Too many Space PR proposals; retry shortly", }, }); diff --git a/tests/space-upstream-production.test.ts b/tests/conversations-space-upstream-production.test.ts similarity index 100% rename from tests/space-upstream-production.test.ts rename to tests/conversations-space-upstream-production.test.ts diff --git a/tests/space-upstream.test.ts b/tests/conversations-space-upstream.test.ts similarity index 80% rename from tests/space-upstream.test.ts rename to tests/conversations-space-upstream.test.ts index 96099470..33bbd9a5 100644 --- a/tests/space-upstream.test.ts +++ b/tests/conversations-space-upstream.test.ts @@ -25,8 +25,6 @@ const ASSISTANT_KEY = "test-space-upstream-key"; const ACCOUNT_ID = "11111111-1111-4111-8111-111111111111"; const GOOD_VARIANT = "pr-test-space-upstream"; const GOOD_VARIANT_URL = `https://ephemeral-${GOOD_VARIANT}.convos.fun`; -const OFF_HOST_VARIANT = "pr-test-space-upstream-off-host"; -const MISMATCH_VARIANT = "pr-test-space-upstream-mismatch"; const pullRequestResult = { conversationId: "conversation_abc", @@ -138,14 +136,7 @@ beforeEach(() => { }; vi.spyOn(prisma.agentVariant, "findFirst").mockImplementation((args) => { const slug = (args?.where as { slug?: string } | undefined)?.slug; - const assistantWorkerUrl = - slug === GOOD_VARIANT - ? GOOD_VARIANT_URL - : slug === OFF_HOST_VARIANT - ? "https://evil.example.com" - : slug === MISMATCH_VARIANT - ? "https://ephemeral-pr-test-space-upstream-other.convos.fun" - : null; + const assistantWorkerUrl = slug === GOOD_VARIANT ? GOOD_VARIANT_URL : null; return Promise.resolve( assistantWorkerUrl ? { assistantWorkerUrl } : null, ) as never; @@ -179,12 +170,12 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { expect(fetchCalls).toHaveLength(1); }); - test("normalizes the bounded conversation ID and sends only the shared key", async () => { + test("forwards the bounded conversation ID verbatim and sends only the shared key", async () => { const res = await proposal(buildApp()); expect(res.status).toBe(200); expect(fetchCalls).toHaveLength(1); expect(fetchCalls[0]?.url).toBe( - `${DEFAULT_URL}/api/conversations/conversation_abc/space-upstream`, + `${DEFAULT_URL}/api/conversations/CONVERSATION_ABC/space-upstream`, ); expect(fetchCalls[0]?.init).toMatchObject({ method: "POST", @@ -197,8 +188,8 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { }); test.each([ - ["invalid characters", "bad%20conversation"], - ["overlong", "a".repeat(129)], + ["blank", "%20"], + ["overlong", "a".repeat(257)], ])("rejects an %s conversation ID before fetch", async (_label, id) => { const res = await proposal( buildApp(), @@ -206,8 +197,9 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { ); expect(res.status).toBe(400); expect(res.body).toEqual({ - code: "INVALID_REQUEST", - error: "Invalid Space PR proposal request", + success: false, + error: "INVALID_REQUEST", + message: "Invalid Space PR proposal request", }); expect(fetchCalls).toHaveLength(0); }); @@ -222,7 +214,7 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { `/api/v2/conversations/conversation_abc/debug/space-upstream?${query}`, ); expect(res.status).toBe(400); - expect(responseBody(res).code).toBe("INVALID_REQUEST"); + expect(responseBody(res).error).toBe("INVALID_REQUEST"); expect(fetchCalls).toHaveLength(0); }); @@ -248,21 +240,19 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { ); }); - test.each([OFF_HOST_VARIANT, MISMATCH_VARIANT, "unknown-space-variant"])( - "fails a non-allowed variant closed without fetching (%s)", - async (variantId) => { - const res = await proposal( - buildApp(), - `/api/v2/conversations/conversation_abc/debug/space-upstream?variantId=${variantId}`, - ); - expect(res.status).toBe(409); - expect(res.body).toEqual({ - code: "VARIANT_UNAVAILABLE", - error: "The selected agent variant is unavailable", - }); - expect(fetchCalls).toHaveLength(0); - }, - ); + test("fails a non-allowed variant closed without fetching", async () => { + const res = await proposal( + buildApp(), + "/api/v2/conversations/conversation_abc/debug/space-upstream?variantId=unknown-space-variant", + ); + expect(res.status).toBe(409); + expect(res.body).toEqual({ + success: false, + error: "VARIANT_UNAVAILABLE", + message: "The selected agent variant is unavailable", + }); + expect(fetchCalls).toHaveLength(0); + }); test("uses a 50-second upstream AbortSignal", async () => { const timeoutSpy = vi.spyOn(AbortSignal, "timeout"); @@ -275,20 +265,14 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { } }); - test.each([ - ["empty shared key", { assistantApiKey: "" }], - ["whitespace shared key", { assistantApiKey: " " }], - ["empty default origin", { assistantApiUrl: "" }], - ["whitespace default origin", { assistantApiUrl: " " }], - ])("returns unavailable for %s", async (_label, override) => { + test("returns unavailable without the optional shared key", async () => { __setAssistantConfigOverridesForTests({ assistantApiUrl: DEFAULT_URL, - assistantApiKey: ASSISTANT_KEY, - ...override, + assistantApiKey: "", }); const res = await proposal(buildApp()); expect(res.status).toBe(503); - expect(responseBody(res).code).toBe("SPACE_UPSTREAM_UNAVAILABLE"); + expect(responseBody(res).error).toBe("SPACE_UPSTREAM_UNAVAILABLE"); expect(fetchCalls).toHaveLength(0); }); @@ -321,21 +305,16 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { }; const res = await proposal(buildApp()); expect(res.status).toBe(expectedStatus); - expect(responseBody(res).code).toBe(expectedCode); + expect(responseBody(res).error).toBe(expectedCode); if (workerCode === "space_upstream_refused") { - expect(responseBody(res).error).toBe("Safe upstream detail"); + expect(responseBody(res).message).toBe("Safe upstream detail"); } }, ); test.each([ ["uncoded old-route 404", 404, { error: "Not found" }], - ["unexpected coded status", 418, { error: "No", code: "unexpected" }], - [ - "extra error envelope fields", - 404, - { error: "Not found", code: "space_not_found", extra: true }, - ], + ["unexpected code", 418, { error: "No", code: "unexpected" }], ])("maps %s to the generic failure", async (_label, status, body) => { fetchImpl = (url, init) => { fetchCalls.push({ url, init }); @@ -343,7 +322,27 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { }; const res = await proposal(buildApp()); expect(res.status).toBe(502); - expect(responseBody(res).code).toBe("SPACE_UPSTREAM_FAILED"); + expect(responseBody(res).error).toBe("SPACE_UPSTREAM_FAILED"); + }); + + test("accepts additive fields in a coded Worker error", async () => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve( + jsonResponse(404, { + error: "Not found", + code: "space_not_found", + extra: true, + }), + ); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(404); + expect(res.body).toEqual({ + success: false, + error: "SPACE_NOT_FOUND", + message: "No Space was found for this conversation", + }); }); test("separates timeout and network failures", async () => { @@ -351,12 +350,12 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { Promise.reject(new DOMException("Timed out", "TimeoutError")); const timedOut = await proposal(buildApp()); expect(timedOut.status).toBe(504); - expect(responseBody(timedOut).code).toBe("SPACE_UPSTREAM_TIMEOUT"); + expect(responseBody(timedOut).error).toBe("SPACE_UPSTREAM_TIMEOUT"); fetchImpl = () => Promise.reject(new TypeError("network unavailable")); const networkFailure = await proposal(buildApp()); expect(networkFailure.status).toBe(502); - expect(responseBody(networkFailure).code).toBe("SPACE_UPSTREAM_FAILED"); + expect(responseBody(networkFailure).error).toBe("SPACE_UPSTREAM_FAILED"); }); test.each([pullRequestResult, unchangedResult])( @@ -368,7 +367,7 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { }; const res = await proposal(buildApp()); expect(res.status).toBe(200); - expect(res.body).toEqual(result); + expect(res.body).toEqual({ success: true, ...result }); }, ); @@ -379,11 +378,6 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { "missing required field", { ...unchangedResult, forkCommitSha: undefined }, ], - ["extra success field", { ...unchangedResult, extra: true }], - [ - "mixed union fields", - { ...unchangedResult, prUrl: "https://example.com" }, - ], ])("rejects a %s success response", async (_label, body) => { fetchImpl = (url, init) => { fetchCalls.push({ url, init }); @@ -399,7 +393,22 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { }; const res = await proposal(buildApp()); expect(res.status).toBe(502); - expect(responseBody(res).code).toBe("SPACE_UPSTREAM_FAILED"); + expect(responseBody(res).error).toBe("SPACE_UPSTREAM_FAILED"); + }); + + test("accepts additive fields in a valid Worker result", async () => { + fetchImpl = (url, init) => { + fetchCalls.push({ url, init }); + return Promise.resolve( + jsonResponse(200, { + ...unchangedResult, + futureField: true, + }), + ); + }; + const res = await proposal(buildApp()); + expect(res.status).toBe(200); + expect(res.body).toEqual({ success: true, ...unchangedResult }); }); test("caps upstream body diagnostics", async () => { @@ -429,8 +438,9 @@ describe("POST /conversations/:conversationId/debug/space-upstream", () => { const limited = await proposal(app, undefined, ip); expect(limited.status).toBe(429); expect(limited.body).toEqual({ - code: "RATE_LIMITED", - error: "Too many Space PR proposals; retry shortly", + success: false, + error: "RATE_LIMITED", + message: "Too many Space PR proposals; retry shortly", }); const otherIp = await proposal(app, undefined, "203.0.113.78"); From ecfef4ba388b61d8b707ee9af9c6c44e5827a218 Mon Sep 17 00:00:00 2001 From: Saul Carlin Date: Wed, 12 Aug 2026 00:52:05 -0700 Subject: [PATCH 3/3] fix: mount Space upstream relay in production --- .../v2/conversations/conversations.router.ts | 3 +- .../conversations/handlers/space-upstream.ts | 8 ++--- src/api/v2/index.ts | 6 +--- ...rsations-space-upstream-production.test.ts | 33 ------------------- tests/conversations-space-upstream.test.ts | 4 +-- 5 files changed, 8 insertions(+), 46 deletions(-) delete mode 100644 tests/conversations-space-upstream-production.test.ts diff --git a/src/api/v2/conversations/conversations.router.ts b/src/api/v2/conversations/conversations.router.ts index 3e8cc6e0..69a4ee27 100644 --- a/src/api/v2/conversations/conversations.router.ts +++ b/src/api/v2/conversations/conversations.router.ts @@ -18,9 +18,8 @@ import { spaceUpstreamHandler } from "./handlers/space-upstream"; // requireAccount itself. conversationId is the opaque XMTP string (no // Conversation table). export const conversationsRouter = Router(); -export const conversationsDebugRouter = Router(); -conversationsDebugRouter.post( +conversationsRouter.post( "/:conversationId/debug/space-upstream", spaceUpstreamLimiter, requireAccount, diff --git a/src/api/v2/conversations/handlers/space-upstream.ts b/src/api/v2/conversations/handlers/space-upstream.ts index 3d20cd37..3bde9d32 100644 --- a/src/api/v2/conversations/handlers/space-upstream.ts +++ b/src/api/v2/conversations/handlers/space-upstream.ts @@ -133,8 +133,8 @@ function translateUpstreamError(raw: unknown): PublicError { if (!parsed.success) return ERRORS.SPACE_UPSTREAM_FAILED; const { code, error: message } = parsed.data; + if (!(code in UPSTREAM_ERRORS)) return ERRORS.SPACE_UPSTREAM_FAILED; const publicError = UPSTREAM_ERRORS[code as keyof typeof UPSTREAM_ERRORS]; - if (!publicError) return ERRORS.SPACE_UPSTREAM_FAILED; return code === "space_upstream_refused" ? { ...publicError, message } : publicError; @@ -143,9 +143,9 @@ function translateUpstreamError(raw: unknown): PublicError { /** * Handler for POST /api/v2/conversations/:conversationId/debug/space-upstream * - * Relays an authenticated, non-production Space PR proposal to the assistant - * Worker. The client never receives the shared Worker credential; it receives - * the standard v2 success or coded-error envelope instead. + * Relays an authenticated Space PR proposal to the assistant Worker. The + * client never receives the shared Worker credential; it receives the standard + * v2 success or coded-error envelope instead. */ export async function spaceUpstreamHandler(req: Request, res: Response) { const parsedParams = paramsSchema.safeParse(req.params); diff --git a/src/api/v2/index.ts b/src/api/v2/index.ts index 109b8349..fbb5f691 100644 --- a/src/api/v2/index.ts +++ b/src/api/v2/index.ts @@ -40,10 +40,7 @@ import { composioRouter } from "./composio/composio.router"; import { connectionsRouter } from "./connections/connections.router"; import { actionsGetHandler } from "./connections/handlers/actions-get"; import { servicesGetHandler } from "./connections/handlers/services-get"; -import { - conversationsDebugRouter, - conversationsRouter, -} from "./conversations/conversations.router"; +import { conversationsRouter } from "./conversations/conversations.router"; import { creditsAdminRouter } from "./credits-admin/credits-admin.router"; import { dailyRefillRouter } from "./credits/daily.router"; import { devRouter } from "./dev/dev.router"; @@ -64,7 +61,6 @@ const v2Router = Router(); // /dev is a non-production test surface; keep it gated. if (process.env.XMTP_ENV !== "production") { v2Router.use("/dev", devAuthMiddleware, devRouter); - v2Router.use("/conversations", authMiddleware, conversationsDebugRouter); } v2Router.use("/agent-prompt-hints", agentPromptHintsRouter); diff --git a/tests/conversations-space-upstream-production.test.ts b/tests/conversations-space-upstream-production.test.ts deleted file mode 100644 index d24b19d3..00000000 --- a/tests/conversations-space-upstream-production.test.ts +++ /dev/null @@ -1,33 +0,0 @@ -import express from "express"; -import request from "supertest"; -import { afterEach, describe, expect, test, vi } from "vitest"; - -const originalXmtpEnv = process.env.XMTP_ENV; - -afterEach(() => { - process.env.XMTP_ENV = originalXmtpEnv; - vi.resetModules(); -}); - -describe("Space upstream production mount guard", () => { - test("does not mount the debug route in production", async () => { - process.env.XMTP_ENV = "production"; - vi.resetModules(); - const { default: v2Router } = await import("@/api/v2"); - const { pinoMiddleware } = await import("@/middleware/pino"); - const { createJwtToken, validateJWTKeys } = await import("@/utils/jwt"); - await validateJWTKeys(); - const token = await createJwtToken({ - deviceId: "production-mount-test", - accountId: "11111111-1111-4111-8111-111111111111", - }); - const app = express(); - app.use(pinoMiddleware); - app.use("/api/v2", v2Router); - - const response = await request(app) - .post("/api/v2/conversations/conversation_abc/debug/space-upstream") - .set("X-Convos-AuthToken", token); - expect(response.status).toBe(404); - }); -}); diff --git a/tests/conversations-space-upstream.test.ts b/tests/conversations-space-upstream.test.ts index 33bbd9a5..3989b672 100644 --- a/tests/conversations-space-upstream.test.ts +++ b/tests/conversations-space-upstream.test.ts @@ -14,7 +14,7 @@ import { vi, } from "vitest"; import { __setAssistantConfigOverridesForTests } from "@/api/v2/agents/handlers/assistant-config"; -import { conversationsDebugRouter } from "@/api/v2/conversations/conversations.router"; +import { conversationsRouter } from "@/api/v2/conversations/conversations.router"; import { authMiddleware } from "@/middleware/auth"; import { pinoMiddleware } from "@/middleware/pino"; import { createJwtToken, validateJWTKeys } from "@/utils/jwt"; @@ -91,7 +91,7 @@ function buildApp(options?: { app.use( "/api/v2/conversations", options?.auth ? authMiddleware : accountMiddleware, - conversationsDebugRouter, + conversationsRouter, ); return app; }