From 2fa6f18059a76e96cac48182877ce10e236564dd Mon Sep 17 00:00:00 2001 From: Ramiro Rivera Date: Tue, 29 Sep 2026 12:46:25 +0200 Subject: [PATCH] Add an on-demand OAuth refresh for connections OAuth access tokens are refreshed lazily: when a token is within the expiry skew, or after an upstream 401. A connection nobody uses never spends its refresh token, so providers that expire idle refresh tokens kill it and the user has to reconnect in a browser. Nothing lets an operator run the grant ahead of time. Add connections.refreshOAuthToken(ref) and POST /connections/:owner/:integration/:name/oauth/refresh. The grant runs through refreshConnectionToken, so it joins an in-flight refresh of the same connection and persists a rotated refresh token before the access token, like the lazy paths. It returns the new expiry and a health verdict. A refused grant is reported as a verdict (a dead grant reads expired and is not re-sent); a connection with no OAuth grant is InvalidConnectionInputError (400). The org-write guard matches connections.refresh. --- .changeset/connection-oauth-force-refresh.md | 5 + .../connection-oauth-force-refresh.test.ts | 291 ++++++++++++++++++ packages/core/api/src/connections/api.ts | 29 ++ .../api/src/connections/oauth-refresh.test.ts | 228 ++++++++++++++ packages/core/api/src/handlers/connections.ts | 17 + packages/core/sdk/src/connection.ts | 16 + packages/core/sdk/src/errors.ts | 5 +- packages/core/sdk/src/executor.ts | 115 ++++++- packages/core/sdk/src/index.ts | 1 + packages/core/sdk/src/oauth-flow.test.ts | 210 ++++++++++++- packages/core/sdk/src/org-writes.test.ts | 6 + packages/core/sdk/src/shared.ts | 1 + 12 files changed, 917 insertions(+), 7 deletions(-) create mode 100644 .changeset/connection-oauth-force-refresh.md create mode 100644 e2e/scenarios/connection-oauth-force-refresh.test.ts create mode 100644 packages/core/api/src/connections/oauth-refresh.test.ts diff --git a/.changeset/connection-oauth-force-refresh.md b/.changeset/connection-oauth-force-refresh.md new file mode 100644 index 0000000000..faac4fd215 --- /dev/null +++ b/.changeset/connection-oauth-force-refresh.md @@ -0,0 +1,5 @@ +--- +"executor": patch +--- + +Add `connections.refreshOAuthToken` and `POST /connections/:owner/:integration/:name/oauth/refresh`, which run a connection's OAuth refresh grant on demand, even when the access token is not yet due, and return the new expiry with a health verdict. Calling it on a schedule keeps an idle connection's refresh token inside a provider's inactivity window, so it does not die with `invalid_grant`. diff --git a/e2e/scenarios/connection-oauth-force-refresh.test.ts b/e2e/scenarios/connection-oauth-force-refresh.test.ts new file mode 100644 index 0000000000..7487e25817 --- /dev/null +++ b/e2e/scenarios/connection-oauth-force-refresh.test.ts @@ -0,0 +1,291 @@ +// Cross-target: an OAuth connection's refresh grant can be run on demand. +// +// Refresh is lazy. A token is re-minted only when it is due or an upstream +// rejects it, so a connection nobody uses never spends its refresh token, and +// providers that expire refresh tokens after a period of inactivity kill it: +// the next call meets invalid_grant and the user has to reconnect in a +// browser. `POST /connections/:owner/:integration/:name/oauth/refresh` runs +// the grant now, however long the access token has left, which is what lets +// an operator keep an idle connection's refresh token exercised. +// +// The journey: an OpenAPI integration with a declared health check completes a +// real authorization-code flow against a live test AS that mints hour-long +// access tokens and rotates refresh tokens on every grant (forgetting the spent +// one). A health check leaves the token alone; the forced refresh reaches the +// token endpoint, reports the new expiry and a probed verdict, and a second +// forced refresh succeeds only because the rotated refresh token was stored. +// A connection whose grant the AS refuses reads `expired`, and is not re-sent. +import { randomBytes } from "node:crypto"; +import { createServer } from "node:http"; + +import { expect } from "@effect/vitest"; +import { Effect } from "effect"; +import type { HttpApiClient } from "effect/unstable/httpapi"; +import { composePluginApi } from "@executor-js/api/server"; +import { openApiHttpPlugin } from "@executor-js/plugin-openapi/api"; +import { + AuthTemplateSlug, + ConnectionName, + IntegrationSlug, + OAuthClientSlug, +} from "@executor-js/sdk/shared"; +import { serveOAuthTestServer, type OAuthTestServerOptions } from "@executor-js/sdk/testing"; + +import { scenario } from "../src/scenario"; +import { Api, Target } from "../src/services"; + +const api = composePluginApi([openApiHttpPlugin()] as const); +type Client = HttpApiClient.ForApi; + +const unique = (prefix: string) => `${prefix}${randomBytes(4).toString("hex")}`; + +/** Upstream on 127.0.0.1 whose `GET /me` accepts any bearer the test AS + * minted, so a `healthy` verdict means the probe carried a real token. */ +const serveUpstream = () => + Effect.acquireRelease( + Effect.callback<{ readonly url: string; readonly close: () => void }>((resume) => { + const server = createServer((request, response) => { + if (request.method === "GET" && (request.url ?? "").startsWith("/me")) { + const authorized = (request.headers["authorization"] ?? "").startsWith("Bearer at_"); + response.writeHead(authorized ? 200 : 401, { + "content-type": "application/json", + }); + response.end(JSON.stringify(authorized ? { email: "probe@example.test" } : {})); + return; + } + response.writeHead(404, { "content-type": "application/json" }); + response.end(JSON.stringify({ error: "not_found" })); + }); + server.listen(0, "127.0.0.1", () => { + const address = server.address(); + const port = typeof address === "object" && address ? address.port : 0; + resume( + Effect.succeed({ + url: `http://127.0.0.1:${port}`, + close: () => { + server.close(); + server.closeAllConnections(); + }, + }), + ); + }); + }), + (server) => Effect.sync(server.close), + ); + +const spec = ( + baseUrl: string, + oauth: { + readonly authorizationEndpoint: string; + readonly tokenEndpoint: string; + }, +): string => + JSON.stringify({ + openapi: "3.0.3", + info: { title: "Identity API", version: "1.0.0" }, + servers: [{ url: baseUrl }], + paths: { + "/me": { + get: { + operationId: "getMe", + summary: "The current account", + security: [{ oauth: ["identity.read"] }], + responses: { "200": { description: "account" } }, + }, + }, + }, + components: { + securitySchemes: { + oauth: { + type: "oauth2", + flows: { + authorizationCode: { + authorizationUrl: oauth.authorizationEndpoint, + tokenUrl: oauth.tokenEndpoint, + scopes: { "identity.read": "Read the account" }, + }, + }, + }, + }, + }, + }); + +/** One integration with a declared health check, one OAuth client, and one + * connection completed through a real authorization-code flow. */ +const connect = ( + client: Client, + upstream: { readonly url: string }, + name: ConnectionName, + serverOptions: OAuthTestServerOptions, +) => + Effect.gen(function* () { + const oauth = yield* serveOAuthTestServer({ + scopes: ["identity.read"], + ...serverOptions, + }); + const slug = IntegrationSlug.make(unique("forcerefresh")); + const clientSlug = OAuthClientSlug.make(unique("forcerefreshc")); + + yield* Effect.addFinalizer(() => + Effect.all( + [ + client.connections + .remove({ params: { owner: "org", integration: slug, name } }) + .pipe(Effect.ignore), + client.oauth + .removeClient({ params: { slug: clientSlug }, payload: { owner: "org" } }) + .pipe(Effect.ignore), + client.openapi.removeSpec({ params: { slug } }).pipe(Effect.ignore), + ], + { discard: true }, + ), + ); + + yield* client.openapi.addSpec({ + payload: { + spec: { kind: "blob", value: spec(upstream.url, oauth) }, + slug, + baseUrl: upstream.url, + authenticationTemplate: [ + { + slug: "oauth", + kind: "oauth2", + authorizationUrl: oauth.authorizationEndpoint, + tokenUrl: oauth.tokenEndpoint, + scopes: ["identity.read"], + }, + ], + }, + }); + const candidates = yield* client.integrations.healthCheckCandidates({ + params: { slug }, + }); + const getMe = candidates.find((candidate) => candidate.method === "get"); + if (!getMe) return yield* Effect.die("the identity spec exposed no GET candidate"); + yield* client.integrations.healthCheckSet({ + params: { slug }, + payload: { spec: { operation: getMe.operation, identityField: "email" } }, + }); + + yield* client.oauth.createClient({ + payload: { + owner: "org", + slug: clientSlug, + grant: "authorization_code", + authorizationUrl: oauth.authorizationEndpoint, + tokenUrl: oauth.tokenEndpoint, + clientId: "test-client", + clientSecret: "test-secret", + originIntegration: slug, + }, + }); + const started = yield* client.oauth.start({ + payload: { + client: clientSlug, + clientOwner: "org", + owner: "org", + name, + integration: slug, + template: AuthTemplateSlug.make("oauth"), + }, + }); + expect(started.status, "oauth.start redirects to the authorization server").toBe("redirect"); + if (started.status !== "redirect") return yield* Effect.die("no redirect"); + + // Drive the test IdP's consent by hand (authorize → login → code). + const code = yield* Effect.promise(async () => { + const authorize = await fetch(started.authorizationUrl, { redirect: "manual" }); + const loginUrl = authorize.headers.get("location"); + if (!loginUrl) throw new Error(`authorize did not redirect: ${authorize.status}`); + const login = await fetch(loginUrl, { + method: "POST", + headers: { + authorization: `Basic ${Buffer.from("alice:password").toString("base64")}`, + }, + redirect: "manual", + }); + const callbackUrl = login.headers.get("location"); + if (!callbackUrl) throw new Error(`login did not redirect: ${login.status}`); + const minted = new URL(callbackUrl).searchParams.get("code"); + if (!minted) throw new Error("callback carried no authorization code"); + return minted; + }); + yield* client.oauth.complete({ payload: { state: started.state, code } }); + yield* oauth.clearRequests; + + return { + params: { owner: "org" as const, integration: slug, name }, + /** The refresh grants the authorization server actually received. */ + refreshGrants: oauth.requests.pipe( + Effect.map((all) => + all.filter( + (request) => + request.path === "/token" && + request.method === "POST" && + request.body.includes("grant_type=refresh_token"), + ), + ), + ), + }; + }); + +scenario( + "Connections · an OAuth refresh can be forced before the access token is due", + {}, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const { client: makeClient } = yield* Api; + const identity = yield* target.newIdentity(); + const client = yield* makeClient(api, identity); + const upstream = yield* serveUpstream(); + + const live = yield* connect(client, upstream, ConnectionName.make("forcerefreshlive"), { + tokenExpiresInSeconds: 3600, + }); + + // Baseline: with an hour left, nothing touches the refresh token. + const before = yield* client.connections.checkHealth({ + params: live.params, + query: {}, + }); + expect(before.status, "the fresh connection is healthy").toBe("healthy"); + expect(yield* live.refreshGrants, "a health check does not refresh a token not due").toEqual( + [], + ); + + const forced = yield* client.connections.refreshOAuthToken({ params: live.params }); + expect(forced.refreshed, "the authorization server issued a new token").toBe(true); + expect(forced.health.status, "the new token passes the declared probe").toBe("healthy"); + expect(forced.expiresAt ?? 0, "the new expiry is about an hour out").toBeGreaterThan( + Date.now() + 30 * 60_000, + ); + expect(yield* live.refreshGrants, "exactly one refresh grant was sent").toHaveLength(1); + + // The AS forgets a spent refresh token, so this grant succeeds only if + // the rotated one from the first grant was stored. + const again = yield* client.connections.refreshOAuthToken({ params: live.params }); + expect(again.refreshed, "the rotated refresh token was persisted").toBe(true); + const grants = yield* live.refreshGrants; + expect(grants).toHaveLength(2); + const sent = grants.map((grant) => new URLSearchParams(grant.body).get("refresh_token")); + expect(sent[1], "the second grant spent the rotated refresh token").not.toBe(sent[0]); + + // A grant the AS refuses is a verdict, and a dead grant is not re-sent. + const refused = yield* connect(client, upstream, ConnectionName.make("forcerefreshdead"), { + tokenExpiresInSeconds: 3600, + supportRefresh: false, + invalidRefreshTokenDescription: "Grant revoked", + }); + const dead = yield* client.connections.refreshOAuthToken({ params: refused.params }); + expect(dead.refreshed, "no token was issued").toBe(false); + expect(dead.health.status, "a refused grant reads expired, so the UI offers reconnect").toBe( + "expired", + ); + expect(dead.health.detail ?? "").toContain("Grant revoked"); + const deadAgain = yield* client.connections.refreshOAuthToken({ params: refused.params }); + expect(deadAgain.health.status).toBe("expired"); + expect(yield* refused.refreshGrants, "the dead grant was sent once").toHaveLength(1); + }), + ), +); diff --git a/packages/core/api/src/connections/api.ts b/packages/core/api/src/connections/api.ts index ca70b50ad5..3a91865ce3 100644 --- a/packages/core/api/src/connections/api.ts +++ b/packages/core/api/src/connections/api.ts @@ -66,6 +66,12 @@ const ConnectionResponse = Schema.Struct({ lastHealth: Schema.NullOr(HealthCheckResult), }); +const OAuthRefreshResponse = Schema.Struct({ + refreshed: Schema.Boolean, + expiresAt: Schema.NullOr(Schema.Number), + health: HealthCheckResult, +}); + const ToolResponse = Schema.Struct({ address: Schema.String, owner: Owner, @@ -249,6 +255,29 @@ export const ConnectionsApi = HttpApiGroup.make("connections") error: [InternalError, ConnectionNotFound, IntegrationNotFound], }), ) + // Run the OAuth refresh grant NOW, however far the access token is from + // expiry, then report health. Refresh is otherwise lazy, so this is what + // keeps an idle connection's refresh token inside the provider's inactivity + // window. Distinct from `refresh` above, which re-syncs the tool catalog. A + // refused grant is a verdict (`health.status: "expired"`), not an error; a + // connection with no OAuth grant is a 400. + .add( + HttpApiEndpoint.post( + "refreshOAuthToken", + "/connections/:owner/:integration/:name/oauth/refresh", + { + params: ConnectionParams, + success: OAuthRefreshResponse, + error: [ + InternalError, + ConnectionNotFound, + IntegrationNotFound, + InvalidConnectionInput, + OrgWriteDeniedError, + ], + }, + ), + ) // Run the health check against an IN-FLIGHT credential without saving it (the // key-first connect flow): confirm the pasted key works and surface the // identity the UI derives a connection name from before anything persists. diff --git a/packages/core/api/src/connections/oauth-refresh.test.ts b/packages/core/api/src/connections/oauth-refresh.test.ts new file mode 100644 index 0000000000..d414ec1acc --- /dev/null +++ b/packages/core/api/src/connections/oauth-refresh.test.ts @@ -0,0 +1,228 @@ +import { HttpApiBuilder } from "effect/unstable/httpapi"; +import { HttpRouter, HttpServer } from "effect/unstable/http"; +import { describe, expect, it } from "@effect/vitest"; +import { Context, Effect, Layer } from "effect"; + +import { + AuthTemplateSlug, + ConnectionName, + IntegrationSlug, + OAuthClientSlug, + ToolName, + createExecutor, + definePlugin, + type Executor, +} from "@executor-js/sdk"; +import { + makeTestConfig, + memoryCredentialsPlugin, + serveOAuthTestServer, + type OAuthTestServerShape, +} from "@executor-js/sdk/testing"; + +import { ExecutorApi } from "../api"; +import { observabilityMiddleware } from "../observability"; +import { CoreHandlers, ExecutionEngineService, ExecutorService } from "../server"; + +// `POST /connections/:owner/:integration/:name/oauth/refresh` — the on-demand +// OAuth refresh grant, driven over HTTP against a live test authorization +// server. + +const webHandlerFor = (executor: Executor) => + Effect.acquireRelease( + Effect.sync(() => + HttpRouter.toWebHandler( + HttpApiBuilder.layer(ExecutorApi).pipe( + Layer.provide(CoreHandlers), + Layer.provide(observabilityMiddleware(ExecutorApi)), + Layer.provide(Layer.succeed(ExecutorService)(executor)), + Layer.provide( + Layer.succeed(ExecutionEngineService)({} as ExecutionEngineService["Service"]), + ), + Layer.provideMerge(HttpServer.layerServices), + Layer.provideMerge(Layer.succeed(HttpRouter.RouterConfig)({ maxParamLength: 1000 })), + ), + { disableLogger: true }, + ), + ), + (web) => Effect.promise(() => web.dispose()), + ); + +const handlerContextFor = (executor: Executor) => + Context.make(ExecutorService, executor).pipe( + Context.add(ExecutionEngineService, {} as ExecutionEngineService["Service"]), + ); + +const INTEGRATION = IntegrationSlug.make("acme"); +const TEMPLATE = AuthTemplateSlug.make("oauth"); +const CLIENT = OAuthClientSlug.make("acme-app"); + +const acmePlugin = definePlugin(() => ({ + id: "acme" as const, + storage: () => ({}), + resolveTools: () => + Effect.succeed({ + tools: [{ name: ToolName.make("whoami"), description: "whoami" }], + }), + describeAuthMethods: () => [ + { + id: "oauth", + label: "OAuth2", + kind: "oauth" as const, + template: String(TEMPLATE), + oauth: { scopes: ["read"] }, + }, + ], + invokeTool: ({ credential }) => Effect.succeed({ token: credential.value }), + checkHealth: ({ credential }) => + Effect.succeed({ + status: credential.value === null ? "expired" : "healthy", + checkedAt: Date.now(), + }), + extension: (ctx) => ({ + seed: () => + ctx.core.integrations.register({ + slug: INTEGRATION, + description: "Acme", + config: {}, + }), + }), +}))(); + +const plugins = [memoryCredentialsPlugin(), acmePlugin] as const; + +const connectOAuth = (executor: Executor, server: OAuthTestServerShape) => + Effect.gen(function* () { + yield* executor.acme.seed(); + yield* executor.oauth.createClient({ + owner: "org", + slug: CLIENT, + authorizationUrl: server.authorizationEndpoint, + tokenUrl: server.tokenEndpoint, + grant: "authorization_code", + clientId: "test-client", + clientSecret: "test-secret", + }); + const started = yield* executor.oauth.start({ + owner: "org", + client: CLIENT, + clientOwner: "org", + name: ConnectionName.make("main"), + integration: INTEGRATION, + template: TEMPLATE, + }); + expect(started.status).toBe("redirect"); + if (started.status !== "redirect") return; + const callback = yield* server.completeAuthorizationCodeFlow({ + authorizationUrl: started.authorizationUrl, + }); + yield* executor.oauth.complete({ state: started.state, code: callback.code }); + }); + +const refreshRequest = (owner: string, name: string) => + new Request(`http://localhost/connections/${owner}/${INTEGRATION}/${name}/oauth/refresh`, { + method: "POST", + }); + +const refreshGrantCount = (server: OAuthTestServerShape) => + server.requests.pipe( + Effect.map( + (requests) => + requests.filter( + (request) => + request.path === "/token" && request.body.includes("grant_type=refresh_token"), + ).length, + ), + ); + +describe("POST /connections/:owner/:integration/:name/oauth/refresh", () => { + it.effect("refreshes a token that is nowhere near expiry", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveOAuthTestServer({ + scopes: ["read"], + tokenExpiresInSeconds: 3600, + }); + const config = makeTestConfig({ plugins }); + const executor = yield* createExecutor(config); + yield* Effect.addFinalizer(() => executor.close().pipe(Effect.ignore)); + yield* connectOAuth(executor, server); + yield* server.clearRequests; + const web = yield* webHandlerFor(executor); + + const response = yield* Effect.promise(() => + web.handler(refreshRequest("org", "main"), handlerContextFor(executor)), + ); + expect(response.status).toBe(200); + const body = (yield* Effect.promise(() => response.json())) as { + readonly refreshed: boolean; + readonly expiresAt: number | null; + readonly health: { readonly status: string }; + }; + expect(body.refreshed).toBe(true); + expect(body.expiresAt).toBeGreaterThan(Date.now() + 30 * 60_000); + expect(body.health.status).toBe("healthy"); + expect(yield* refreshGrantCount(server)).toBe(1); + }), + ), + ); + + it.effect("maps non-OAuth and unknown connections to 400 and 404", () => + Effect.scoped( + Effect.gen(function* () { + const config = makeTestConfig({ plugins }); + const executor = yield* createExecutor(config); + yield* Effect.addFinalizer(() => executor.close().pipe(Effect.ignore)); + yield* executor.acme.seed(); + yield* executor.connections.create({ + owner: "org", + name: ConnectionName.make("pasted"), + integration: INTEGRATION, + template: AuthTemplateSlug.make("apiKey"), + value: "static-token", + }); + const web = yield* webHandlerFor(executor); + const context = handlerContextFor(executor); + + const notOAuth = yield* Effect.promise(() => + web.handler(refreshRequest("org", "pasted"), context), + ); + expect(notOAuth.status).toBe(400); + expect(yield* Effect.promise(() => notOAuth.json())).toMatchObject({ + _tag: "InvalidConnectionInputError", + }); + + const missing = yield* Effect.promise(() => + web.handler(refreshRequest("org", "nope"), context), + ); + expect(missing.status).toBe(404); + }), + ), + ); + + it.effect("refuses a workspace connection for a member without workspace writes", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveOAuthTestServer({ scopes: ["read"] }); + const config = makeTestConfig({ plugins }); + const admin = yield* createExecutor(config); + const member = yield* createExecutor({ ...config, orgWrites: "denied" }); + yield* Effect.addFinalizer(() => + admin.close().pipe(Effect.andThen(member.close()), Effect.ignore), + ); + yield* connectOAuth(admin, server); + yield* server.clearRequests; + const web = yield* webHandlerFor(member); + + const response = yield* Effect.promise(() => + web.handler(refreshRequest("org", "main"), handlerContextFor(member)), + ); + expect(response.status).toBe(403); + expect(yield* Effect.promise(() => response.json())).toMatchObject({ + _tag: "OrgWriteDeniedError", + }); + expect(yield* refreshGrantCount(server), "nothing reached the token endpoint").toBe(0); + }), + ), + ); +}); diff --git a/packages/core/api/src/handlers/connections.ts b/packages/core/api/src/handlers/connections.ts index 9ae476a97f..d99f6b700d 100644 --- a/packages/core/api/src/handlers/connections.ts +++ b/packages/core/api/src/handlers/connections.ts @@ -160,6 +160,23 @@ export const ConnectionsHandlers = HttpApiBuilder.group(ExecutorApi, "connection }), ), ) + .handle("refreshOAuthToken", ({ params: path }) => + capture( + Effect.gen(function* () { + const executor = yield* ExecutorService; + const result = yield* executor.connections.refreshOAuthToken({ + owner: path.owner, + integration: path.integration, + name: path.name, + }); + return { + refreshed: result.refreshed, + expiresAt: result.expiresAt, + health: toHealthResponse(result.health), + }; + }), + ), + ) .handle("validate", ({ payload }) => capture( Effect.gen(function* () { diff --git a/packages/core/sdk/src/connection.ts b/packages/core/sdk/src/connection.ts index 41fbd8b24c..9c124145f7 100644 --- a/packages/core/sdk/src/connection.ts +++ b/packages/core/sdk/src/connection.ts @@ -123,3 +123,19 @@ export interface UpdateConnectionInput { readonly description?: string | null; readonly identityLabel?: string | null; } + +/** Outcome of `connections.refreshOAuthToken`: a refresh-token grant run on + * demand, however far the stored access token is from expiry. */ +export interface ConnectionOAuthRefreshResult { + /** True when the authorization server issued a new access token. False when + * it refused the grant, or when the grant was already recorded dead and + * nothing was sent; `health` then says why. */ + readonly refreshed: boolean; + /** Epoch ms when the stored access token expires after this call; null when + * the authorization server advertised no lifetime. */ + readonly expiresAt: number | null; + /** The connection's health after the refresh, in the same shape and with the + * same persistence as `checkHealth`. A refused grant reads `expired` (dead + * grant) or `degraded`. */ + readonly health: HealthCheckResult; +} diff --git a/packages/core/sdk/src/errors.ts b/packages/core/sdk/src/errors.ts index d0f16e7d12..527ae49859 100644 --- a/packages/core/sdk/src/errors.ts +++ b/packages/core/sdk/src/errors.ts @@ -211,8 +211,9 @@ export class ConnectionAlreadyExistsError extends Schema.TaggedErrorClass()( "InvalidConnectionInputError", { message: Schema.String }, diff --git a/packages/core/sdk/src/executor.ts b/packages/core/sdk/src/executor.ts index 0f338405cd..61f7e68fbc 100644 --- a/packages/core/sdk/src/executor.ts +++ b/packages/core/sdk/src/executor.ts @@ -37,6 +37,7 @@ import { coreToolsPlugin } from "./core-tools"; import type { Connection, ConnectionInputOrigin, + ConnectionOAuthRefreshResult, ConnectionRef, CreateConnectionInput, ConnectionValueInput, @@ -456,6 +457,25 @@ export type Executor = { HealthCheckResult, ConnectionNotFoundError | IntegrationNotFoundError | StorageFailure >; + /** Run the OAuth refresh grant for a saved connection NOW, even while its + * access token is far from expiry, then report its health. Refresh is + * otherwise lazy (due within the expiry skew, or after an upstream 401), + * so an idle connection never exercises its refresh token and can outlive + * the provider's refresh-token inactivity window; calling this on a + * schedule keeps it alive. Shares the in-flight gate and rotation-safe + * persistence with every other refresh. A refused grant is a verdict, not + * an error; a connection with no OAuth grant is `InvalidConnectionInputError`. + * (`refresh` above re-syncs the tool catalog; it does not touch tokens.) */ + readonly refreshOAuthToken: ( + ref: ConnectionRef, + ) => Effect.Effect< + ConnectionOAuthRefreshResult, + | ConnectionNotFoundError + | IntegrationNotFoundError + | InvalidConnectionInputError + | OrgWriteDeniedError + | StorageFailure + >; /** Validate an in-flight credential WITHOUT saving it (key-first connect): * resolve the pasted value(s), run the health check, and return the result * so the caller can confirm the key works and derive a name from the @@ -2265,9 +2285,10 @@ export const createExecutor = => + Effect.gen(function* () { + yield* guardOrgWrite(ref.owner); + const row = yield* findConnectionRow(ref); + if (!row) { + return yield* new ConnectionNotFoundError({ + owner: ref.owner, + integration: ref.integration, + name: ref.name, + }); + } + if (row.oauth_client == null) { + return yield* new InvalidConnectionInputError({ + message: `Connection "${ref.name}" does not use OAuth, so it has no token to refresh.`, + }); + } + const expiresAtOf = (current: ConnectionRow | null): number | null => + current?.expires_at == null ? null : Number(current.expires_at); + // A recorded dead grant answers as `checkHealth` does, without a token + // request. The refresh path would refuse to re-send it anyway; this + // keeps the verdict identical to the one every other surface serves. + const reauthState = oauthReauthRequiredFromProviderState(row.provider_state); + if (reauthState !== null) { + return { + refreshed: false, + expiresAt: expiresAtOf(row), + health: deadGrantVerdict(reauthState, row), + }; + } + const provider = credentialProviders.get(row.provider); + if (!provider) { + return yield* new StorageError({ + message: `Credential provider "${row.provider}" is not registered.`, + cause: undefined, + }); + } + const refusal = yield* refreshConnectionToken(row, provider, "manual").pipe( + Effect.as(null), + Effect.catchTag("CredentialResolutionError", (failure) => Effect.succeed(failure)), + ); + if (refusal !== null) { + // Fold the refusal the way a probe would. A rejected grant has + // already recorded itself dead (which makes this persist a no-op); + // anything else lands as the connection's verdict. + const health = healthFromCredentialResolutionFailure(refusal); + yield* persistProbeHealthResult(ref, health); + return { + refreshed: false, + expiresAt: expiresAtOf(yield* findConnectionRow(ref)), + health, + }; + } + // A minted token proves only that the authorization server still + // honours the grant. Probe so the verdict (and the persisted one) + // reflects whether the upstream accepts it. + const health = yield* connectionCheckHealth(ref); + return { + refreshed: true, + expiresAt: expiresAtOf(yield* findConnectionRow(ref)), + health, + }; + }).pipe( + Effect.withSpan("executor.connection.oauth.refresh", { + attributes: { + "executor.tenant": tenant, + ...(subject != null ? { "executor.subject": subject } : {}), + "executor.integration": String(ref.integration), + "executor.connection": String(ref.name), + }, + }), + ); + const connectionValidate = ( input: ValidateConnectionInput, ): Effect.Effect => @@ -7280,6 +7388,7 @@ export const createExecutor = { ), ); }); + +// Refresh is lazy: a token is re-minted only when it is due or an upstream +// rejects it, so a connection nobody uses never spends its refresh token and +// can outlive the provider's refresh-token inactivity window. +// `connections.refreshOAuthToken` runs the grant on demand through the same +// gate and persistence as the lazy paths. +describe("connections.refreshOAuthToken", () => { + const MAIN = ConnectionName.make("main"); + const mainRef = { owner: "org" as const, integration: INTEG, name: MAIN }; + const address = ToolAddress.make("tools.acme.org.main.whoami"); + + const connect = (executor: Executor, server: OAuthTestServerShape) => + Effect.gen(function* () { + yield* executor.acme.seed(); + yield* executor.oauth.createClient({ + owner: "org", + slug: CLIENT, + authorizationUrl: server.authorizationEndpoint, + tokenUrl: server.tokenEndpoint, + grant: "authorization_code", + clientId: "test-client", + clientSecret: "test-secret", + }); + const started = yield* executor.oauth.start({ + owner: "org", + client: CLIENT, + clientOwner: "org", + name: MAIN, + integration: INTEG, + template: TEMPLATE, + }); + if (started.status !== "redirect") { + return yield* Effect.die(`expected a redirect, got ${started.status}`); + } + const callback = yield* server.completeAuthorizationCodeFlow({ + authorizationUrl: started.authorizationUrl, + }); + return yield* executor.oauth.complete({ state: started.state, code: callback.code }); + }); + + it.effect("refreshes a token far from expiry and persists the rotated refresh token", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveOAuthTestServer({ + scopes: ["read"], + tokenExpiresInSeconds: 3600, + }); + const { executor } = yield* makeTestWorkspaceHarness({ plugins }); + const connection = yield* connect(executor, server); + const originalExpiry = connection.expiresAt ?? 0; + expect(originalExpiry).toBeGreaterThan(Date.now() + 30 * 60_000); + + // Baseline: the lazy path leaves a token with an hour to go alone. + yield* server.clearRequests; + const original = (yield* executor.execute(address, {})) as { token: string }; + expect(refreshGrantsIn(yield* server.requests)).toHaveLength(0); + + const first = yield* executor.connections.refreshOAuthToken(mainRef); + expect(first.refreshed).toBe(true); + expect(first.health.status).toBe("healthy"); + expect(first.expiresAt).toBeGreaterThanOrEqual(originalExpiry); + expect(refreshGrantsIn(yield* server.requests)).toHaveLength(1); + + const after = (yield* executor.execute(address, {})) as { token: string }; + expect(after.token, "tool calls use the newly minted access token").not.toBe( + original.token, + ); + expect(yield* server.acceptsAccessToken(after.token)).toBe(true); + + // The test authorization server rotates refresh tokens and forgets + // the spent one, so a second grant succeeds only if the rotated + // token was stored. + const second = yield* executor.connections.refreshOAuthToken(mainRef); + expect(second.refreshed).toBe(true); + const grants = refreshGrantsIn(yield* server.requests); + expect(grants).toHaveLength(2); + const sentRefreshTokens = grants.map((grant) => + new URLSearchParams(grant.body).get("refresh_token"), + ); + expect(sentRefreshTokens[1]).not.toBe(sentRefreshTokens[0]); + + const persisted = yield* executor.connections.get(mainRef); + expect(persisted?.lastHealth?.status).toBe("healthy"); + }), + ), + ); + + it.effect("joins a refresh already in flight instead of spending the token twice", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveOAuthTestServer({ scopes: ["read"] }); + const park = makeTokenRequestPark(); + const config = { ...makeTestConfig({ plugins }), fetch: park.fetch }; + const executor = yield* createExecutor(config); + yield* Effect.addFinalizer(() => executor.close().pipe(Effect.ignore)); + yield* Effect.addFinalizer(() => + Effect.promise(() => config.testDb.close()).pipe(Effect.ignore), + ); + yield* connect(executor, server); + + // A tool call that must refresh, racing a keepalive caller. + yield* Effect.promise(() => + config.db.updateMany("connection", { + where: (b) => b("name", "=", "main"), + set: { expires_at: Date.now() - 60_000 }, + }), + ); + yield* server.clearRequests; + park.arm(); + + const call = yield* Effect.forkChild(executor.execute(address, {})); + const forced = yield* Effect.forkChild(executor.connections.refreshOAuthToken(mainRef)); + yield* Effect.promise(() => park.seen); + park.release(); + + const called = (yield* Fiber.join(call)) as { token: string }; + const result = yield* Fiber.join(forced); + expect(result.refreshed).toBe(true); + expect( + refreshGrantsIn(yield* server.requests), + "the forced refresh joined the tool call's grant", + ).toHaveLength(1); + const next = (yield* executor.execute(address, {})) as { token: string }; + expect(next.token).toBe(called.token); + + // The rotated refresh token survived the race: another grant works. + const again = yield* executor.connections.refreshOAuthToken(mainRef); + expect(again.refreshed).toBe(true); + expect(again.health.status).toBe("healthy"); + }), + ), + ); + + it.effect("reports a dead grant as expired and does not re-send it", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveOAuthTestServer({ + scopes: ["read"], + supportRefresh: false, + invalidRefreshTokenDescription: "Grant revoked", + }); + const { executor, config } = yield* makeTestWorkspaceHarness({ plugins }); + yield* connect(executor, server); + yield* server.clearRequests; + + const first = yield* executor.connections.refreshOAuthToken(mainRef); + expect(first.refreshed).toBe(false); + expect(first.health).toMatchObject({ + status: "expired", + reason: "credential_refresh_rejected", + }); + expect(first.health.detail).toContain("invalid_grant"); + expect(refreshGrantsIn(yield* server.requests)).toHaveLength(1); + + const row = yield* Effect.promise(() => + config.db.findFirst("connection", { where: (b) => b("name", "=", "main") }), + ); + expect( + (row?.provider_state as { oauthReauthRequiredAt?: number } | null)?.oauthReauthRequiredAt, + ).toEqual(expect.any(Number)); + + const second = yield* executor.connections.refreshOAuthToken(mainRef); + expect(second.refreshed).toBe(false); + expect(second.health.status).toBe("expired"); + expect( + refreshGrantsIn(yield* server.requests), + "the dead grant is not re-sent", + ).toHaveLength(1); + }), + ), + ); + + it.effect("rejects connections without an OAuth grant and unknown connections", () => + Effect.scoped( + Effect.gen(function* () { + const { executor } = yield* makeTestWorkspaceHarness({ plugins }); + yield* executor.acme.seed(); + const pasted = yield* executor.connections.create({ + owner: "org", + name: ConnectionName.make("pasted"), + integration: INTEG, + template: AuthTemplateSlug.make("apiKey"), + value: "static-token", + }); + + const notOAuth = yield* Effect.flip( + executor.connections.refreshOAuthToken({ + owner: pasted.owner, + integration: pasted.integration, + name: pasted.name, + }), + ); + expect(notOAuth).toMatchObject({ _tag: "InvalidConnectionInputError" }); + + const missing = yield* Effect.flip( + executor.connections.refreshOAuthToken({ + owner: "org", + integration: INTEG, + name: ConnectionName.make("nope"), + }), + ); + expect(missing).toMatchObject({ _tag: "ConnectionNotFoundError" }); + }), + ), + ); +}); diff --git a/packages/core/sdk/src/org-writes.test.ts b/packages/core/sdk/src/org-writes.test.ts index f30c5d14e7..63cc042435 100644 --- a/packages/core/sdk/src/org-writes.test.ts +++ b/packages/core/sdk/src/org-writes.test.ts @@ -185,6 +185,11 @@ describe("orgWrites: denied", () => { description: "my credential", }); expect(yield* member.connections.refresh(mineRef)).toHaveLength(1); + // Past the gate for a Personal connection: refused only because a + // pasted credential has no OAuth grant to refresh. + expect(yield* Effect.flip(member.connections.refreshOAuthToken(mineRef))).toMatchObject({ + _tag: "InvalidConnectionInputError", + }); const shared = yield* admin.connections.create({ owner: "org", @@ -200,6 +205,7 @@ describe("orgWrites: denied", () => { }; yield* expectOrgWriteDenied(member.connections.update(ref, { description: "renamed" })); yield* expectOrgWriteDenied(member.connections.refresh(ref)); + yield* expectOrgWriteDenied(member.connections.refreshOAuthToken(ref)); yield* expectOrgWriteDenied(member.connections.remove(ref)); expect(yield* admin.connections.refresh(ref)).toHaveLength(1); yield* member.connections.remove(mineRef); diff --git a/packages/core/sdk/src/shared.ts b/packages/core/sdk/src/shared.ts index 391c12e3fa..87bb9a703d 100644 --- a/packages/core/sdk/src/shared.ts +++ b/packages/core/sdk/src/shared.ts @@ -40,6 +40,7 @@ export type { } from "./integration"; export type { Connection, + ConnectionOAuthRefreshResult, ConnectionRef, ConnectionValueInput, CreateConnectionInput,