From fa0215f98719168c59fc5e470cfac2c2b5cc2f76 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Wed, 16 Sep 2026 10:13:28 -0700 Subject: [PATCH] Keep request resources alive for background OAuth sync --- .../oauth-background-resource-lifetime.md | 5 + e2e/cloud/oauth-background-catalog.test.ts | 178 ++++++++++++++++++ .../api/src/server/request-scoped.test.ts | 144 ++++++++++++++ .../core/api/src/server/request-scoped.ts | 77 ++++++-- .../core/api/src/server/scoped-executor.ts | 15 +- 5 files changed, 404 insertions(+), 15 deletions(-) create mode 100644 .changeset/oauth-background-resource-lifetime.md create mode 100644 e2e/cloud/oauth-background-catalog.test.ts create mode 100644 packages/core/api/src/server/request-scoped.test.ts diff --git a/.changeset/oauth-background-resource-lifetime.md b/.changeset/oauth-background-resource-lifetime.md new file mode 100644 index 0000000000..2680db0ed6 --- /dev/null +++ b/.changeset/oauth-background-resource-lifetime.md @@ -0,0 +1,5 @@ +--- +"@executor-js/api": patch +--- + +Keep request resources alive until background OAuth tool discovery finishes, so slow cloud connections can publish their tools after the callback returns. diff --git a/e2e/cloud/oauth-background-catalog.test.ts b/e2e/cloud/oauth-background-catalog.test.ts new file mode 100644 index 0000000000..c6c630d943 --- /dev/null +++ b/e2e/cloud/oauth-background-catalog.test.ts @@ -0,0 +1,178 @@ +import { randomBytes } from "node:crypto"; + +import { expect } from "@effect/vitest"; +import { Effect, Option, Schema } from "effect"; +import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"; +import { composePluginApi } from "@executor-js/api/server"; +import { mcpHttpPlugin } from "@executor-js/plugin-mcp/api"; +import { + AuthTemplateSlug, + ConnectionName, + IntegrationSlug, + OAuthClientSlug, +} from "@executor-js/sdk/shared"; +import { serveTestHttpApp } from "@executor-js/sdk/testing"; + +import { createEmulatorInstance } from "../src/emulator-instance"; +import { scenario } from "../src/scenario"; +import { Api, Browser, Target } from "../src/services"; +import { visit } from "../src/surfaces/browser"; + +const api = composePluginApi([mcpHttpPlugin()] as const); +const decodeRpc = Schema.decodeUnknownOption( + Schema.fromJsonString(Schema.Struct({ method: Schema.String })), +); + +scenario( + "OAuth ยท slow catalog discovery persists after the cloud callback returns", + { timeout: 180_000 }, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const browser = yield* Browser; + const { client: makeApiClient } = yield* Api; + const identity = yield* target.newIdentity(); + const client = yield* makeApiClient(api, identity); + const upstream = yield* createEmulatorInstance("mcp", "oauth-background-catalog"); + const firstListing = Promise.withResolvers(); + const releaseCallbackListing = Promise.withResolvers(); + const releaseOtherListings = Promise.withResolvers(); + let listings = 0; + + // Only delay traffic. OAuth, credentials and MCP responses all come from + // the published emulator. Hold later listings separately so a UI/API read + // cannot repair the catalog and hide failure of the callback's own sync. + const proxy = yield* serveTestHttpApp((request) => + Effect.promise(async () => { + const web = await Effect.runPromise(HttpServerRequest.toWeb(request)); + const body = web.method === "GET" || web.method === "HEAD" ? undefined : await web.text(); + const path = new URL(web.url); + const headers = new Headers(web.headers); + headers.delete("host"); + const response = await fetch(`${upstream}${path.pathname}${path.search}`, { + method: web.method, + headers, + body, + redirect: "manual", + }); + const rpc = body === undefined ? Option.none() : decodeRpc(body); + if (response.ok && Option.isSome(rpc) && rpc.value.method === "tools/list") { + listings += 1; + if (listings === 1) { + firstListing.resolve(); + await releaseCallbackListing.promise; + } else { + await releaseOtherListings.promise; + } + } + return HttpServerResponse.fromWeb(response); + }), + ); + yield* Effect.addFinalizer(() => + Effect.sync(() => { + releaseCallbackListing.resolve(); + releaseOtherListings.resolve(); + }), + ); + + const slug = IntegrationSlug.make(`slow_oauth_${randomBytes(4).toString("hex")}`); + yield* client.mcp.addServer({ + payload: { + transport: "remote", + name: "Slow OAuth catalog", + endpoint: `${proxy.baseUrl}/mcp`, + slug, + authenticationTemplate: [{ kind: "oauth2" }], + }, + }); + yield* Effect.addFinalizer(() => + Effect.gen(function* () { + releaseCallbackListing.resolve(); + releaseOtherListings.resolve(); + yield* client.mcp.removeServer({ params: { slug } }).pipe(Effect.ignore); + }), + ); + + const probe = yield* client.oauth.probe({ payload: { url: `${proxy.baseUrl}/mcp` } }); + if (!probe.registrationEndpoint || !probe.authorizationUrl || !probe.tokenUrl) { + return yield* Effect.die("Emulator did not advertise OAuth registration"); + } + const { client: oauthClient } = yield* client.oauth.registerDynamic({ + payload: { + owner: "org", + slug: OAuthClientSlug.make(`${slug}_client`), + registrationEndpoint: probe.registrationEndpoint, + authorizationUrl: probe.authorizationUrl, + tokenUrl: probe.tokenUrl, + resource: probe.resource, + scopes: probe.scopesSupported ?? [], + originIntegration: slug, + }, + }); + yield* Effect.addFinalizer(() => + client.oauth + .removeClient({ + params: { slug: oauthClient }, + payload: { owner: "org" }, + }) + .pipe(Effect.ignore), + ); + const started = yield* client.oauth.start({ + payload: { + owner: "org", + client: oauthClient, + clientOwner: "org", + name: ConnectionName.make("main"), + integration: slug, + template: AuthTemplateSlug.make("oauth2"), + }, + }); + if (started.status !== "redirect") return yield* Effect.die("Expected OAuth authorization"); + yield* Effect.addFinalizer(() => + client.oauth.cancel({ payload: { state: started.state } }).pipe(Effect.ignore), + ); + + yield* browser.session(identity, async ({ page, step }) => { + await step("Authorize the OAuth connection", async () => { + // Stay off the integration screen until discovery is verified: its + // refresh after the callback could otherwise repair a failed sync. + await page.goto(started.authorizationUrl); + const authorization = new URL(page.url()); + const approved = await page.request.post(`${upstream}/authorize/approve`, { + form: { ...Object.fromEntries(authorization.searchParams), login: "admin" }, + maxRedirects: 0, + }); + expect(approved.status()).toBe(302); + const callback = approved.headers()["location"]; + if (!callback) throw new Error("Emulator did not return the OAuth callback"); + await page.goto(callback, { waitUntil: "domcontentloaded" }); + await page.getByRole("heading", { name: "Connected" }).waitFor(); + }); + await step("Keep discovery blocked after the callback returns", async () => { + await expect.poll(() => listings, { timeout: 15_000 }).toBeGreaterThan(0); + await firstListing.promise; + const connections = await Effect.runPromise( + client.connections.list({ query: { integration: slug } }), + ); + expect(connections).toHaveLength(1); + }); + await step("Receive the tools from the completed background discovery", async () => { + releaseCallbackListing.resolve(); + await expect + .poll( + async () => { + const tools = await Effect.runPromise( + client.tools.list({ query: { integration: slug } }), + ); + return tools.length; + }, + { timeout: 30_000, interval: 500 }, + ) + .toBeGreaterThan(0); + await visit(page, `/integrations/${slug}`); + await page.getByRole("button", { name: "Add connection" }).waitFor(); + }); + }); + }), + ), +); diff --git a/packages/core/api/src/server/request-scoped.test.ts b/packages/core/api/src/server/request-scoped.test.ts new file mode 100644 index 0000000000..98c18226d0 --- /dev/null +++ b/packages/core/api/src/server/request-scoped.test.ts @@ -0,0 +1,144 @@ +import { describe, expect, it, onTestFinished } from "@effect/vitest"; +import { Context, Effect, Layer } from "effect"; +import { HttpRouter, HttpServer, HttpServerResponse } from "effect/unstable/http"; + +import { RequestBackgroundTasks, requestScopedMiddleware } from "./request-scoped"; + +class Resource extends Context.Service()( + "test/RequestResource", +) {} + +const fixture = ( + options: { + background?: boolean; + failTask?: boolean; + failRequest?: boolean; + blockRequest?: boolean; + } = {}, +) => { + const resources: Resource["Service"][] = []; + const gates: ReturnType>[] = []; + const keptAlive: Promise[] = []; + const writes: number[] = []; + const entered = Promise.withResolvers(); + const resource = Layer.effect(Resource)( + Effect.acquireRelease( + Effect.sync(() => { + const acquired = { id: resources.length, closed: false }; + resources.push(acquired); + gates.push(Promise.withResolvers()); + return acquired; + }), + (acquired) => + Effect.sync(() => { + acquired.closed = true; + }), + ), + ); + const routes = HttpRouter.add( + "GET", + "/", + Effect.gen(function* () { + const acquired = yield* Resource; + if (options.background !== false) { + const tasks = yield* RequestBackgroundTasks; + const gate = gates[acquired.id]; + if (!gate) return yield* Effect.die("Missing request gate"); + const task = Effect.runPromise( + Effect.gen(function* () { + yield* Effect.promise(() => gate.promise); + if (options.failTask) return yield* Effect.fail("Background failure"); + if (acquired.closed) + return yield* Effect.fail("Resource closed before background write"); + writes.push(acquired.id); + }), + ); + keptAlive.push(tasks.retain(task)); + } + entered.resolve(); + if (options.blockRequest) return yield* Effect.never; + if (options.failRequest) return yield* Effect.die("Request failure"); + return HttpServerResponse.empty(); + }), + ); + const app = HttpRouter.toWebHandler( + routes.pipe( + Layer.provide(requestScopedMiddleware(resource).layer), + Layer.provideMerge(HttpServer.layerServices), + ), + { disableLogger: true }, + ); + onTestFinished(async () => { + for (const gate of gates) gate.resolve(); + await Promise.all(keptAlive); + await app.dispose(); + }); + return { + resources, + gates, + keptAlive, + writes, + entered, + request: (signal?: AbortSignal) => + app.handler(new Request("http://test.local/", { signal }), Context.empty()), + }; +}; + +describe("request background resource ownership", () => { + it("returns the response before background work, and releases after its write", async () => { + const test = fixture(); + expect((await test.request()).status).toBe(204); + expect(test.resources.map((item) => item.closed)).toEqual([false]); + expect(test.writes).toEqual([]); + test.gates[0]?.resolve(); + await Promise.all(test.keptAlive); + expect(test.writes).toEqual([0]); + expect(test.resources.map((item) => item.closed)).toEqual([true]); + }); + + it("keeps concurrent requests isolated and releases each independently", async () => { + const test = fixture(); + const responses = await Promise.all([test.request(), test.request()]); + expect(responses.map((response) => response.status)).toEqual([204, 204]); + test.gates[0]?.resolve(); + await test.keptAlive[0]; + expect(test.resources.map((item) => item.closed)).toEqual([true, false]); + test.gates[1]?.resolve(); + await test.keptAlive[1]; + expect(test.writes).toEqual([0, 1]); + expect(test.resources.map((item) => item.closed)).toEqual([true, true]); + }); + + for (const failure of ["task", "request"] as const) { + it(`releases resources after a ${failure} failure`, async () => { + const test = fixture({ failTask: failure === "task", failRequest: failure === "request" }); + expect((await test.request()).status).toBe(failure === "request" ? 500 : 204); + expect(test.resources.map((item) => item.closed)).toEqual([false]); + test.gates[0]?.resolve(); + await Promise.all(test.keptAlive); + expect(test.writes).toEqual(failure === "task" ? [] : [0]); + expect(test.resources.map((item) => item.closed)).toEqual([true]); + }); + } + + it("closes before responding when no background work was registered", async () => { + const test = fixture({ background: false }); + expect((await test.request()).status).toBe(204); + expect(test.resources.map((item) => item.closed)).toEqual([true]); + expect(test.keptAlive).toEqual([]); + }); + + it("finishes retained work and cleanup after the client cancels the request", async () => { + const test = fixture({ blockRequest: true }); + const controller = new AbortController(); + const response = test.request(controller.signal); + await test.entered.promise; + controller.abort(); + expect((await response).status).toBe(499); + expect(test.resources.map((item) => item.closed)).toEqual([false]); + test.gates[0]?.resolve(); + await Promise.all(test.keptAlive); + expect(test.writes).toEqual([0]); + expect(test.resources.map((item) => item.closed)).toEqual([true]); + }); +}); diff --git a/packages/core/api/src/server/request-scoped.ts b/packages/core/api/src/server/request-scoped.ts index 1436d851bb..6f0aa05da8 100644 --- a/packages/core/api/src/server/request-scoped.ts +++ b/packages/core/api/src/server/request-scoped.ts @@ -13,8 +13,9 @@ // both build the inner layer at construction time. The only primitive // that actually rebuilds per request is a router middleware whose // per-request handler builds the layer with a *fresh* `MemoMap` and a -// per-request scope, so `acquireRelease` fires per request and finalizers -// run when the request fiber's scope closes. +// per-request scope, so `acquireRelease` fires per request. Background work +// retains that scope until its database writes finish; the response need not +// wait, but the platform keep-alive must include resource cleanup too. // // The fresh `MemoMap` matters: `Layer.build` would otherwise inherit // `CurrentMemoMap` from the boot context (`HttpRouter.toWebHandler` @@ -29,14 +30,24 @@ // coverage that pins this rule down (sequential AND concurrent cases). // --------------------------------------------------------------------------- -import { Effect, Layer } from "effect"; +import { Context, Effect, Layer, Scope } from "effect"; import { HttpRouter } from "effect/unstable/http"; +/** + * Retain this request's resources for already-started background work. The + * caller owns task error reporting; the returned promise settles only after + * all retained tasks AND resource finalizers finish, for the host's waitUntil. + */ +export class RequestBackgroundTasks extends Context.Service< + RequestBackgroundTasks, + { readonly retain: (task: Promise) => Promise } +>()("@executor-js/api/RequestBackgroundTasks") {} + /** * Build an `HttpRouter.middleware` that provides `layer`'s services to * each request. The layer is rebuilt per HTTP request so - * `Effect.acquireRelease` fires per request and is released when the - * request fiber's scope closes. + * `Effect.acquireRelease` fires per request. Resources close after the handler + * and its registered background work settle, without delaying the response. * * The returned value is a `Middleware`. Use `.layer` to apply it as a * standalone layer; use `.combine(other)` to fold it into another @@ -46,15 +57,55 @@ import { HttpRouter } from "effect/unstable/http"; * outer middleware's `requires`). */ export const requestScopedMiddleware = (layer: Layer.Layer) => - HttpRouter.middleware<{ provides: A }>()((httpEffect) => - Effect.scoped( + HttpRouter.middleware<{ provides: A | RequestBackgroundTasks }>()((httpEffect) => + Effect.uninterruptibleMask((restore) => Effect.gen(function* () { - // Fresh MemoMap per request โ€” see file-level note for why we - // must NOT inherit `CurrentMemoMap` from the boot context. - const memoMap = yield* Layer.makeMemoMap; - const scope = yield* Effect.scope; - const services = yield* Layer.buildWithMemoMap(layer, memoMap, scope); - return yield* Effect.provideContext(httpEffect, services); + const scope = yield* Scope.make(); + const pending = new Set>(); + const released = Promise.withResolvers(); + const background = RequestBackgroundTasks.of({ + retain: (task) => { + // SDK tasks report their own failures. Both outcomes release the + // resource lease; a rejected task must not leak its database. + const settled = task.then( + () => { + pending.delete(settled); + }, + () => { + pending.delete(settled); + }, + ); + pending.add(settled); + return released.promise; + }, + }); + return yield* restore( + Effect.gen(function* () { + // Never inherit the boot MemoMap: concurrent requests each own + // their socket, including after either response has been sent. + const memoMap = yield* Layer.makeMemoMap; + const services = yield* Layer.buildWithMemoMap(layer, memoMap, scope); + return yield* Effect.provideContext(httpEffect, services); + }).pipe( + Effect.provideService(Scope.Scope, scope), + Effect.provideService(RequestBackgroundTasks, background), + ), + ).pipe( + Effect.onExit((exit) => { + const release = Effect.gen(function* () { + // A retained task may start another task before it settles. + while (pending.size > 0) { + yield* Effect.promise(() => Promise.all(pending)); + } + }).pipe( + Effect.ensuring(Scope.close(scope, exit)), + Effect.ensuring(Effect.sync(() => released.resolve())), + ); + // With no background work, preserve synchronous teardown. Otherwise + // the resource owner, including cleanup, is kept alive by the host. + return pending.size === 0 ? release : release.pipe(Effect.forkDetach, Effect.asVoid); + }), + ); }), ), ); diff --git a/packages/core/api/src/server/scoped-executor.ts b/packages/core/api/src/server/scoped-executor.ts index 5e10bbbfd1..749839be3e 100644 --- a/packages/core/api/src/server/scoped-executor.ts +++ b/packages/core/api/src/server/scoped-executor.ts @@ -51,6 +51,7 @@ import { } from "@executor-js/sdk/host-internal"; import { DbProvider } from "./executor-fuma-db"; +import { RequestBackgroundTasks } from "./request-scoped"; // --------------------------------------------------------------------------- // HostConfig seam โ€” the two host scalars that vary the `createExecutor` options. @@ -126,12 +127,14 @@ export interface HostConfigShape { */ readonly toolsSyncTtlMs?: number | null; /** - * Forwarded verbatim to `ExecutorConfig.waitUntil`: the host's keep-alive + * Forwarded to `ExecutorConfig.waitUntil`: the host's keep-alive * for background work that outlives a request (stale tool-catalog rebuilds * that keep running after a read stops waiting). Cloud supplies the * platform `waitUntil` from `cloudflare:workers`, which binds to the * in-flight invocation ambiently; long-lived hosts (self-host, local, * tests) omit it and detached fibers simply run to completion in-process. + * Under requestScopedMiddleware, the promise also covers releasing that + * request's database after background work finishes. */ readonly waitUntil?: (promise: Promise) => void; } @@ -267,6 +270,14 @@ export const makeScopedExecutor = < const { db, blobs } = yield* DbProvider.asEffect(); const { plugins: pluginsFactory } = yield* PluginsProvider.asEffect(); const config = yield* HostConfig.asEffect(); + const background = yield* Effect.serviceOption(RequestBackgroundTasks); + const waitUntil = Option.match(background, { + onNone: () => config.waitUntil, + onSome: (tasks) => (task: Promise) => { + const released = tasks.retain(task); + config.waitUntil?.(released); + }, + }); // Explicit config wins; otherwise fall back to the request origin if a host // provided one (HTTP middleware / MCP session DO). Stays `undefined` for // non-request callers โ€” `coreTools.webBaseUrl` is optional and only the @@ -323,7 +334,7 @@ export const makeScopedExecutor = < fetch: hostedFetch, onIntegrationChange: config.onIntegrationChange, ...(config.toolsSyncTtlMs !== undefined ? { toolsSyncTtlMs: config.toolsSyncTtlMs } : {}), - ...(config.waitUntil !== undefined ? { waitUntil: config.waitUntil } : {}), + ...(waitUntil !== undefined ? { waitUntil } : {}), onElicitation: "accept-all", ...(options?.orgWrites === undefined ? {} : { orgWrites: options.orgWrites }), redirectUri,