diff --git a/docs/generated/postgres-schema/README.md b/docs/generated/postgres-schema/README.md index e0091b4052d..88d31b873e4 100644 --- a/docs/generated/postgres-schema/README.md +++ b/docs/generated/postgres-schema/README.md @@ -87,7 +87,7 @@ Auto-generated from the PostgreSQL migrations in @n8n/db. Do not edit by hand. | [public.mcp_registry_server](public.mcp_registry_server.md) | 7 | | BASE TABLE | | [public.oauth_access_tokens](public.oauth_access_tokens.md) | 3 | | BASE TABLE | | [public.oauth_authorization_codes](public.oauth_authorization_codes.md) | 13 | | BASE TABLE | -| [public.oauth_clients](public.oauth_clients.md) | 9 | | BASE TABLE | +| [public.oauth_clients](public.oauth_clients.md) | 10 | | BASE TABLE | | [public.oauth_refresh_tokens](public.oauth_refresh_tokens.md) | 7 | | BASE TABLE | | [public.oauth_user_consents](public.oauth_user_consents.md) | 5 | | BASE TABLE | | [public.processed_data](public.processed_data.md) | 5 | | BASE TABLE | @@ -1121,6 +1121,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/postgres-schema/public.oauth_access_tokens.md b/docs/generated/postgres-schema/public.oauth_access_tokens.md index fca472ccbd9..4a3c6947130 100644 --- a/docs/generated/postgres-schema/public.oauth_access_tokens.md +++ b/docs/generated/postgres-schema/public.oauth_access_tokens.md @@ -44,6 +44,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/postgres-schema/public.oauth_authorization_codes.md b/docs/generated/postgres-schema/public.oauth_authorization_codes.md index 6ee13ccbf37..454fb742b73 100644 --- a/docs/generated/postgres-schema/public.oauth_authorization_codes.md +++ b/docs/generated/postgres-schema/public.oauth_authorization_codes.md @@ -72,6 +72,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/postgres-schema/public.oauth_clients.md b/docs/generated/postgres-schema/public.oauth_clients.md index df8e99d1fce..60895179557 100644 --- a/docs/generated/postgres-schema/public.oauth_clients.md +++ b/docs/generated/postgres-schema/public.oauth_clients.md @@ -9,6 +9,7 @@ | createdAt | timestamp(3) with time zone | CURRENT_TIMESTAMP(3) | false | | | | | grantTypes | json | | false | | | | | id | varchar | | false | [public.oauth_access_tokens](public.oauth_access_tokens.md) [public.oauth_authorization_codes](public.oauth_authorization_codes.md) [public.oauth_refresh_tokens](public.oauth_refresh_tokens.md) [public.oauth_user_consents](public.oauth_user_consents.md) | | | +| isFirstParty | boolean | false | false | | | | | name | varchar(255) | | false | | | | | redirectUris | json | | false | | | | | tokenEndpointAuthMethod | varchar(255) | 'none'::character varying | false | | | Possible values: none, client_secret_basic or client_secret_post | @@ -22,6 +23,7 @@ | oauth_clients_createdAt_not_null | n | NOT NULL "createdAt" | | oauth_clients_grantTypes_not_null | n | NOT NULL "grantTypes" | | oauth_clients_id_not_null | n | NOT NULL id | +| oauth_clients_isFirstParty_not_null | n | NOT NULL "isFirstParty" | | oauth_clients_name_not_null | n | NOT NULL name | | oauth_clients_redirectUris_not_null | n | NOT NULL "redirectUris" | | oauth_clients_tokenEndpointAuthMethod_not_null | n | NOT NULL "tokenEndpointAuthMethod" | @@ -49,6 +51,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/postgres-schema/public.oauth_refresh_tokens.md b/docs/generated/postgres-schema/public.oauth_refresh_tokens.md index cafa6fe353d..edde80d3771 100644 --- a/docs/generated/postgres-schema/public.oauth_refresh_tokens.md +++ b/docs/generated/postgres-schema/public.oauth_refresh_tokens.md @@ -56,6 +56,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/postgres-schema/public.oauth_user_consents.md b/docs/generated/postgres-schema/public.oauth_user_consents.md index 78451d89ecc..00b7ab42fb6 100644 --- a/docs/generated/postgres-schema/public.oauth_user_consents.md +++ b/docs/generated/postgres-schema/public.oauth_user_consents.md @@ -52,6 +52,7 @@ erDiagram timestamp_3__with_time_zone createdAt json grantTypes varchar id + boolean isFirstParty varchar_255_ name json redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/README.md b/docs/generated/sqlite-schema/README.md index fbc98ff3fcf..056cd791652 100644 --- a/docs/generated/sqlite-schema/README.md +++ b/docs/generated/sqlite-schema/README.md @@ -87,7 +87,7 @@ Auto-generated from the SQLite migrations in @n8n/db. Do not edit by hand. | [mcp_registry_server](mcp_registry_server.md) | 7 | | table | | [oauth_access_tokens](oauth_access_tokens.md) | 3 | | table | | [oauth_authorization_codes](oauth_authorization_codes.md) | 13 | | table | -| [oauth_clients](oauth_clients.md) | 9 | | table | +| [oauth_clients](oauth_clients.md) | 10 | | table | | [oauth_refresh_tokens](oauth_refresh_tokens.md) | 7 | | table | | [oauth_user_consents](oauth_user_consents.md) | 5 | | table | | [processed_data](processed_data.md) | 5 | | table | @@ -1108,6 +1108,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/oauth_access_tokens.md b/docs/generated/sqlite-schema/oauth_access_tokens.md index b8e4d97757b..904d880aa81 100644 --- a/docs/generated/sqlite-schema/oauth_access_tokens.md +++ b/docs/generated/sqlite-schema/oauth_access_tokens.md @@ -53,6 +53,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/oauth_authorization_codes.md b/docs/generated/sqlite-schema/oauth_authorization_codes.md index 446e13cb6ae..5bc4b4947f3 100644 --- a/docs/generated/sqlite-schema/oauth_authorization_codes.md +++ b/docs/generated/sqlite-schema/oauth_authorization_codes.md @@ -73,6 +73,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/oauth_clients.md b/docs/generated/sqlite-schema/oauth_clients.md index aff09badeb2..69a9df541a1 100644 --- a/docs/generated/sqlite-schema/oauth_clients.md +++ b/docs/generated/sqlite-schema/oauth_clients.md @@ -6,7 +6,7 @@ Table Definition ```sql -CREATE TABLE "oauth_clients" ("id" varchar PRIMARY KEY NOT NULL, "name" varchar(255) NOT NULL, "redirectUris" text NOT NULL, "grantTypes" text NOT NULL, "clientSecret" varchar(255), "clientSecretExpiresAt" bigint, "tokenEndpointAuthMethod" varchar(255) NOT NULL DEFAULT ('none'), "createdAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW')), "updatedAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW'))) +CREATE TABLE "oauth_clients" ("id" varchar PRIMARY KEY NOT NULL, "name" varchar(255) NOT NULL, "redirectUris" text NOT NULL, "grantTypes" text NOT NULL, "clientSecret" varchar(255), "clientSecretExpiresAt" bigint, "tokenEndpointAuthMethod" varchar(255) NOT NULL DEFAULT ('none'), "createdAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW')), "updatedAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW')), "isFirstParty" boolean NOT NULL DEFAULT (false)) ``` @@ -20,6 +20,7 @@ CREATE TABLE "oauth_clients" ("id" varchar PRIMARY KEY NOT NULL, "name" varchar( | createdAt | datetime(3) | STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW') | false | | | | | grantTypes | TEXT | | false | | | | | id | varchar | | false | [oauth_access_tokens](oauth_access_tokens.md) [oauth_authorization_codes](oauth_authorization_codes.md) [oauth_refresh_tokens](oauth_refresh_tokens.md) [oauth_user_consents](oauth_user_consents.md) | | | +| isFirstParty | boolean | false | false | | | | | name | varchar(255) | | false | | | | | redirectUris | TEXT | | false | | | | | tokenEndpointAuthMethod | varchar(255) | 'none' | false | | | | @@ -54,6 +55,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/oauth_refresh_tokens.md b/docs/generated/sqlite-schema/oauth_refresh_tokens.md index 075a5050cfd..b4dc6c8a780 100644 --- a/docs/generated/sqlite-schema/oauth_refresh_tokens.md +++ b/docs/generated/sqlite-schema/oauth_refresh_tokens.md @@ -61,6 +61,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/docs/generated/sqlite-schema/oauth_user_consents.md b/docs/generated/sqlite-schema/oauth_user_consents.md index 61a4d9e7df2..02ab80bae06 100644 --- a/docs/generated/sqlite-schema/oauth_user_consents.md +++ b/docs/generated/sqlite-schema/oauth_user_consents.md @@ -57,6 +57,7 @@ erDiagram datetime_3_ createdAt TEXT grantTypes varchar id PK + boolean isFirstParty varchar_255_ name TEXT redirectUris varchar_255_ tokenEndpointAuthMethod diff --git a/packages/@n8n/db/src/migrations/common/1785162364001-AddIsFirstPartyToOAuthClients.ts b/packages/@n8n/db/src/migrations/common/1785162364001-AddIsFirstPartyToOAuthClients.ts new file mode 100644 index 00000000000..5ccb896575c --- /dev/null +++ b/packages/@n8n/db/src/migrations/common/1785162364001-AddIsFirstPartyToOAuthClients.ts @@ -0,0 +1,22 @@ +import type { MigrationContext, ReversibleMigration } from '../migration-types'; + +const OAUTH_CLIENTS_TABLE = 'oauth_clients'; +const OAUTH_IS_FIRST_PARTY_COLUMN = 'isFirstParty'; + +export class AddIsFirstPartyToOAuthClients1785162364001 implements ReversibleMigration { + async up({ schemaBuilder: { addColumns, column } }: MigrationContext) { + await addColumns( + OAUTH_CLIENTS_TABLE, + [column(OAUTH_IS_FIRST_PARTY_COLUMN).bool.notNull.default(false)], + { + recreatesOnSqlite: true, + }, + ); + } + + async down({ schemaBuilder: { dropColumns } }: MigrationContext) { + await dropColumns(OAUTH_CLIENTS_TABLE, [OAUTH_IS_FIRST_PARTY_COLUMN], { + recreatesOnSqlite: true, + }); + } +} diff --git a/packages/@n8n/db/src/migrations/postgresdb/index.ts b/packages/@n8n/db/src/migrations/postgresdb/index.ts index dd6b75e6d80..2409fda64fb 100644 --- a/packages/@n8n/db/src/migrations/postgresdb/index.ts +++ b/packages/@n8n/db/src/migrations/postgresdb/index.ts @@ -231,6 +231,7 @@ import { AddStoredAtToAgentExecution1784815940110 } from '../common/178481594011 import { AddInstanceCredentials1784815940111 } from '../common/1784815940111-AddInstanceCredentials'; import { CreateAgentEvalTables1784815940112 } from '../common/1784815940112-CreateAgentEvalTables'; import { AddAvailableInMcpToAgents1784897791636 } from '../common/1784897791636-AddAvailableInMcpToAgents'; +import { AddIsFirstPartyToOAuthClients1785162364001 } from '../common/1785162364001-AddIsFirstPartyToOAuthClients'; import type { Migration } from '../migration-types'; export const postgresMigrations: Migration[] = [ @@ -467,4 +468,5 @@ export const postgresMigrations: Migration[] = [ CreateAgentEvalTables1784815940112, AddAvailableInMcpToAgents1784897791636, ChangeInstalledNodeVersionType1785162364000, + AddIsFirstPartyToOAuthClients1785162364001, ]; diff --git a/packages/@n8n/db/src/migrations/sqlite/1785162364001-AddIsFirstPartyToOAuthClients.ts b/packages/@n8n/db/src/migrations/sqlite/1785162364001-AddIsFirstPartyToOAuthClients.ts new file mode 100644 index 00000000000..370f7249ed5 --- /dev/null +++ b/packages/@n8n/db/src/migrations/sqlite/1785162364001-AddIsFirstPartyToOAuthClients.ts @@ -0,0 +1,27 @@ +import type { MigrationContext, ReversibleMigration } from '../migration-types'; + +const OAUTH_CLIENTS_TABLE = 'oauth_clients'; +const OAUTH_IS_FIRST_PARTY_COLUMN = 'isFirstParty'; + +export class AddIsFirstPartyToOAuthClients1785162364001 implements ReversibleMigration { + // oauth_clients has inbound ON DELETE CASCADE FKs (access/refresh tokens, + // authorization codes, user consents). Recreating the table on SQLite would + // fire those cascades and wipe the referencing rows, so disable FKs. + withFKsDisabled = true as const; + + async up({ schemaBuilder: { addColumns, column } }: MigrationContext) { + await addColumns( + OAUTH_CLIENTS_TABLE, + [column(OAUTH_IS_FIRST_PARTY_COLUMN).bool.notNull.default(false)], + { + recreatesOnSqlite: true, + }, + ); + } + + async down({ schemaBuilder: { dropColumns } }: MigrationContext) { + await dropColumns(OAUTH_CLIENTS_TABLE, [OAUTH_IS_FIRST_PARTY_COLUMN], { + recreatesOnSqlite: true, + }); + } +} diff --git a/packages/@n8n/db/src/migrations/sqlite/index.ts b/packages/@n8n/db/src/migrations/sqlite/index.ts index 5d300b8fac1..d4bb95768af 100644 --- a/packages/@n8n/db/src/migrations/sqlite/index.ts +++ b/packages/@n8n/db/src/migrations/sqlite/index.ts @@ -60,6 +60,7 @@ import { DropAgentDescriptionFromAgents1784000000037 } from './1784000000037-Dro import { AddRecurringCronScheduleKind1784000000045 } from './1784000000045-AddRecurringCronScheduleKind'; import { AddAvailableInMcpToAgents1784897791636 } from './1784897791636-AddAvailableInMcpToAgents'; import { ChangeInstalledNodeVersionType1785162364000 } from './1785162364000-ChangeInstalledNodeVersionType'; +import { AddIsFirstPartyToOAuthClients1785162364001 } from './1785162364001-AddIsFirstPartyToOAuthClients'; import { UniqueWorkflowNames1620821879465 } from '../common/1620821879465-UniqueWorkflowNames'; import { UpdateWorkflowCredentials1630330987096 } from '../common/1630330987096-UpdateWorkflowCredentials'; import { AddNodeIds1658930531669 } from '../common/1658930531669-AddNodeIds'; @@ -449,6 +450,7 @@ const sqliteMigrations: Migration[] = [ CreateAgentEvalTables1784815940112, AddAvailableInMcpToAgents1784897791636, ChangeInstalledNodeVersionType1785162364000, + AddIsFirstPartyToOAuthClients1785162364001, ]; export { sqliteMigrations }; diff --git a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/__tests__/n8n-oauth.test.ts b/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/__tests__/n8n-oauth.test.ts deleted file mode 100644 index 09d6e9d2ab5..00000000000 --- a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/__tests__/n8n-oauth.test.ts +++ /dev/null @@ -1,154 +0,0 @@ -import type { Mocked } from 'vitest'; -import type { Logger } from '@n8n/backend-common'; -import type { User } from '@n8n/db'; -import type { ContextEstablishmentOptions } from '@n8n/decorators'; -import { mock } from 'vitest-mock-extended'; -import type { Cipher } from 'n8n-core'; -import type { ICredentialContext, INode, IRunExecutionData } from 'n8n-workflow'; - -import type { AuthService } from '@/auth/auth.service'; -import type { OAuthTokenVerifierProxy } from '@/services/oauth-token-verifier-proxy.service'; - -import { N8NIdentifier } from '../../../credential-resolvers/identifiers/n8n-identifier'; -import { N8NOAuth2Extractor, N8N_OAUTH_EXTRACTOR_NAME } from '../n8n-oauth-extractor'; -import { N8nOAuthIdentitySeeder } from '../n8n-oauth-seeder'; - -const TOKEN = 'super-secret-oauth-token'; -const RESOURCE = 'https://host/mcp/workflow-a'; - -const scopedLogger = (): Mocked => { - const logger = mock(); - logger.scoped.mockReturnValue(logger); - return logger; -}; - -// Faithful in-memory stand-in for Cipher: encryptV2 returns an opaque handle (never the -// plaintext), decryptV2 resolves it back — so round-trips work while the raw token never -// appears in the serialized run data. -const createVaultCipher = (): Mocked => { - const vault = new Map(); - let counter = 0; - const cipher = mock(); - cipher.encryptV2.mockImplementation(async (data) => { - const plaintext = typeof data === 'string' ? data : JSON.stringify(data); - const handle = `enc:${counter++}`; - vault.set(handle, plaintext); - return handle; - }); - cipher.decryptV2.mockImplementation(async (handle) => vault.get(handle) ?? 'not-json'); - return cipher; -}; - -const makeTriggerNode = (): INode => ({ - id: 'node-1', - name: 'MCP Trigger', - type: '@n8n/n8n-nodes-langchain.mcpTrigger', - typeVersion: 2, - position: [0, 0], - parameters: { authentication: 'n8nOAuth2' }, -}); - -const makeRunExecutionData = (triggerNode: INode): IRunExecutionData => - ({ - resultData: { runData: {} }, - executionData: { - nodeExecutionStack: [{ node: triggerNode, data: { main: [] }, source: null }], - contextData: {}, - waitingExecution: {}, - waitingExecutionSource: {}, - metadata: {}, - }, - }) as unknown as IRunExecutionData; - -describe('n8n-oauth identity seeding', () => { - describe('round-trip: seeder → extractor → identifier', () => { - it('propagates token + resource end-to-end and resolves the user', async () => { - const cipher = createVaultCipher(); - const runExecutionData = makeRunExecutionData(makeTriggerNode()); - - // 1. Seed - await new N8nOAuthIdentitySeeder(scopedLogger(), cipher).seed( - runExecutionData, - TOKEN, - RESOURCE, - ); - - // 2. Extract (reads the item the seeder placed on the stack) - const triggerItems = runExecutionData.executionData!.nodeExecutionStack[0].data.main[0]!; - const result = await new N8NOAuth2Extractor(scopedLogger(), cipher).execute({ - triggerItems, - } as ContextEstablishmentOptions); - - expect(result.contextUpdate?.credentials).toEqual({ - version: 1, - identity: TOKEN, - metadata: { source: 'n8n-oauth', resource: RESOURCE }, - }); - expect(triggerItems[0]).not.toHaveProperty('encryptedMetadata'); - - // 3. Identify (the extractor's metadata must satisfy the identifier's discriminated union; - // this is what catches a source-literal contract drift between the two) - const oauthVerifier = mock(); - oauthVerifier.verifyOAuthAccessToken.mockResolvedValue({ - user: mock({ id: 'user-9' }), - }); - const identifier = new N8NIdentifier(mock(), oauthVerifier); - - const userId = await identifier.resolve( - result.contextUpdate!.credentials as ICredentialContext, - {}, - ); - - expect(userId).toBe('user-9'); - expect(oauthVerifier.verifyOAuthAccessToken).toHaveBeenCalledWith(TOKEN, RESOURCE); - }); - - it('throws and still deletes encryptedMetadata when the blob is invalid', async () => { - const extractor = new N8NOAuth2Extractor(scopedLogger(), createVaultCipher()); - const item = { json: {}, encryptedMetadata: 'unknown-handle' }; - - await expect( - extractor.execute({ triggerItems: [item] } as unknown as ContextEstablishmentOptions), - ).rejects.toThrow('No valid n8n OAuth authentication metadata could be extracted.'); - expect(item).not.toHaveProperty('encryptedMetadata'); - }); - }); - - describe('N8nOAuthIdentitySeeder', () => { - it('persists only the encrypted blob — never the raw token', async () => { - const runExecutionData = makeRunExecutionData(makeTriggerNode()); - - await new N8nOAuthIdentitySeeder(scopedLogger(), createVaultCipher()).seed( - runExecutionData, - TOKEN, - RESOURCE, - ); - - expect(JSON.stringify(runExecutionData)).not.toContain(TOKEN); - expect(runExecutionData.executionData!.nodeExecutionStack[0].data.main[0]![0]).toHaveProperty( - 'encryptedMetadata', - ); - }); - - it('injects the hook on a clone without mutating the shared trigger node', async () => { - const triggerNode = makeTriggerNode(); - const runExecutionData = makeRunExecutionData(triggerNode); - - await new N8nOAuthIdentitySeeder(scopedLogger(), createVaultCipher()).seed( - runExecutionData, - TOKEN, - RESOURCE, - ); - - // Original node untouched (would otherwise leak the injected hook into the persisted snapshot) - expect(triggerNode.parameters).toEqual({ authentication: 'n8nOAuth2' }); - - const stackNode = runExecutionData.executionData!.nodeExecutionStack[0].node; - expect(stackNode).not.toBe(triggerNode); - expect(stackNode.parameters.executionsHooksVersion).toBe(1); - expect(stackNode.parameters.contextEstablishmentHooks).toEqual({ - hooks: [{ hookName: N8N_OAUTH_EXTRACTOR_NAME, isAllowedToFail: true }], - }); - }); - }); -}); diff --git a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/index.ts b/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/index.ts deleted file mode 100644 index 936f681244f..00000000000 --- a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/index.ts +++ /dev/null @@ -1,2 +0,0 @@ -export * from './n8n-oauth-seeder'; -export * from './n8n-oauth-extractor'; diff --git a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/metadata.ts b/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/metadata.ts deleted file mode 100644 index c3c3a8cef99..00000000000 --- a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/metadata.ts +++ /dev/null @@ -1,13 +0,0 @@ -import { z } from 'zod'; - -export const EncryptedMetadataSchema = z.object({ - encryptedMetadata: z.string(), -}); - -export const N8NOAuth2ExtractorMetadataSchema = z.object({ - authToken: z.string(), - resource: z.string(), -}); - -export type EncryptedMetadata = z.infer; -export type N8NOAuth2ExtractorMetadata = z.infer; diff --git a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-extractor.ts b/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-extractor.ts deleted file mode 100644 index a941486cd55..00000000000 --- a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-extractor.ts +++ /dev/null @@ -1,85 +0,0 @@ -import { Logger } from '@n8n/backend-common'; -import { - ContextEstablishmentHook, - ContextEstablishmentOptions, - ContextEstablishmentResult, - HookDescription, - IContextEstablishmentHook, -} from '@n8n/decorators'; -import { Cipher } from 'n8n-core'; -import { EncryptedMetadataSchema, N8NOAuth2ExtractorMetadataSchema } from './metadata'; -import { ensureError } from '@n8n/utils/errors/ensure-error'; -import { jsonParse } from 'n8n-workflow'; - -export const N8N_OAUTH_EXTRACTOR_NAME = 'N8nOAuthExtractor'; - -@ContextEstablishmentHook() -export class N8NOAuth2Extractor implements IContextEstablishmentHook { - constructor( - private readonly logger: Logger, - private readonly cipher: Cipher, - ) { - this.logger = this.logger.scoped('dynamic-credentials'); - } - - hookDescription: HookDescription = { - name: N8N_OAUTH_EXTRACTOR_NAME, - displayName: 'N8N OAuth2 Extractor', - options: [], - }; - - isApplicableToTriggerNode(_nodeType: string): boolean { - // This extractor is not showing up in any UI for selection, it can only be used - // when referenced directly - return false; - } - - async execute(options: ContextEstablishmentOptions): Promise { - if (!options.triggerItems || options.triggerItems.length === 0) { - this.logger.debug('No trigger items found, skipping n8n OAuth extractor hook.'); - throw new Error('No trigger items found, skipping n8n OAuth extractor hook.'); - } - const [triggerItem] = options.triggerItems; - - const encryptedMetadataResult = EncryptedMetadataSchema.safeParse(triggerItem); - - // Always delete encryptedMetadata from the item - delete triggerItem.encryptedMetadata; - - if (encryptedMetadataResult.success) { - try { - const decrypted = await this.cipher.decryptV2( - encryptedMetadataResult.data.encryptedMetadata, - ); - const parsed = jsonParse(decrypted); - const n8nOAuthInformation = N8NOAuth2ExtractorMetadataSchema.safeParse(parsed); - if (n8nOAuthInformation.success) { - return { - triggerItems: options.triggerItems, - contextUpdate: { - credentials: { - version: 1, - identity: n8nOAuthInformation.data.authToken, - metadata: { - source: 'n8n-oauth', - resource: n8nOAuthInformation.data.resource, - }, - }, - }, - }; - } else { - this.logger.warn('Invalid format for encryptedMetadata in n8n OAuth extractor', { - errors: n8nOAuthInformation.error.errors, - }); - } - } catch (error) { - this.logger.error('Failed to decrypt/parse encrypted n8n OAuth metadata', { - error: ensureError(error), - }); - } - } else { - this.logger.warn('No encryptedMetadata found in trigger item for n8n OAuth extractor.'); - } - throw new Error('No valid n8n OAuth authentication metadata could be extracted.'); - } -} diff --git a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-seeder.ts b/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-seeder.ts deleted file mode 100644 index ce7c79083b1..00000000000 --- a/packages/cli/src/modules/dynamic-credentials.ee/context-establishment-hooks/n8n-oauth/n8n-oauth-seeder.ts +++ /dev/null @@ -1,66 +0,0 @@ -import { Logger } from '@n8n/backend-common'; -import { Service } from '@n8n/di'; -import { Cipher } from 'n8n-core'; -import { IRunExecutionData } from 'n8n-workflow'; -import { N8NOAuth2ExtractorMetadata } from './metadata'; -import { N8N_OAUTH_EXTRACTOR_NAME } from './n8n-oauth-extractor'; -import { TriggerAuthIdentitySeeder } from '@/services/trigger-auth-identity-seeder-proxy.service'; - -@Service() -export class N8nOAuthIdentitySeeder implements TriggerAuthIdentitySeeder { - constructor( - private readonly logger: Logger, - private readonly cipher: Cipher, - ) { - this.logger = this.logger.scoped('dynamic-credentials'); - } - - async seed(runExecutionData: IRunExecutionData, token: string, resource: string): Promise { - const item = runExecutionData.executionData?.nodeExecutionStack?.[0]?.data?.main; - const metadata: N8NOAuth2ExtractorMetadata = { - resource, - authToken: token, - }; - - const encryptedMetadata = await this.cipher.encryptV2(metadata); - if (item) { - this.logger.debug('Seeding n8n OAuth identity into trigger execution data.'); - if (item?.[0]?.[0]) { - item[0][0].encryptedMetadata = encryptedMetadata; - } else if (item?.[0]) { - item[0].push({ - encryptedMetadata, - json: {}, - }); - } else { - runExecutionData.executionData!.nodeExecutionStack[0].data.main = [ - [ - { - encryptedMetadata, - json: {}, - }, - ], - ]; - } - } - - if (runExecutionData.executionData?.nodeExecutionStack?.[0]) { - this.logger.debug('Seeding n8n OAuth identity into trigger node parameters.'); - runExecutionData.executionData.nodeExecutionStack[0].node = { - ...runExecutionData.executionData?.nodeExecutionStack?.[0].node, - parameters: { - ...runExecutionData.executionData?.nodeExecutionStack?.[0].node.parameters, - executionsHooksVersion: 1, - contextEstablishmentHooks: { - hooks: [ - { - hookName: N8N_OAUTH_EXTRACTOR_NAME, - isAllowedToFail: true, - }, - ], - }, - }, - }; - } - } -} diff --git a/packages/cli/src/modules/dynamic-credentials.ee/dynamic-credentials.module.ts b/packages/cli/src/modules/dynamic-credentials.ee/dynamic-credentials.module.ts index 409e6041576..2c0d3edda6d 100644 --- a/packages/cli/src/modules/dynamic-credentials.ee/dynamic-credentials.module.ts +++ b/packages/cli/src/modules/dynamic-credentials.ee/dynamic-credentials.module.ts @@ -3,8 +3,6 @@ import type { ModuleInterface } from '@n8n/decorators'; import { BackendModule, OnShutdown } from '@n8n/decorators'; import { Container } from '@n8n/di'; -import { TriggerAuthIdentitySeederProxy } from '@/services/trigger-auth-identity-seeder-proxy.service'; - /** * Superset capability: external/custom credential resolvers (OAuth/Slack) plus * their management surfaces and identity-extractor hooks. The base "private @@ -19,15 +17,6 @@ export class DynamicCredentialsModule implements ModuleInterface { async init() { await import('./dynamic-credentials.controller.js'); - // Import the n8n oauth extractor and seeder - const { N8nOAuthIdentitySeeder } = await import( - './context-establishment-hooks/n8n-oauth/index.js' - ); - - Container.get(TriggerAuthIdentitySeederProxy).registerSeeder( - Container.get(N8nOAuthIdentitySeeder), - ); - // System resolver powers private credentials; OAuth/Slack resolvers and // their management/identity-extractor surfaces are external-only. await import('./credential-resolvers/n8n-credential-resolver.js'); diff --git a/packages/cli/src/modules/dynamic-credentials.ee/services/__tests__/credential-check-proxy.service.test.ts b/packages/cli/src/modules/dynamic-credentials.ee/services/__tests__/credential-check-proxy.service.test.ts index 5d8ed2412c1..46d5435536f 100644 --- a/packages/cli/src/modules/dynamic-credentials.ee/services/__tests__/credential-check-proxy.service.test.ts +++ b/packages/cli/src/modules/dynamic-credentials.ee/services/__tests__/credential-check-proxy.service.test.ts @@ -1,6 +1,10 @@ import type { Mocked } from 'vitest'; import type { GlobalConfig } from '@n8n/config'; -import type { IExecutionContext, PlaintextExecutionContext } from 'n8n-workflow'; +import type { + ICredentialContext, + IExecutionContext, + PlaintextExecutionContext, +} from 'n8n-workflow'; import type { EnterpriseCredentialsService } from '@/credentials/credentials.service.ee'; import type { UrlService } from '@/services/url.service'; @@ -67,7 +71,7 @@ describe('CredentialCheckProxyService', () => { } as unknown as Mocked; mockExecutionContextService = { - decryptExecutionContext: vi.fn().mockResolvedValue(plaintextContext), + decryptCredentialContext: vi.fn().mockResolvedValue(plaintextContext.credentials), } as unknown as Mocked; mockEnterpriseCredentialsService = { @@ -178,12 +182,9 @@ describe('CredentialCheckProxyService', () => { }); it('should throw when no credential context in execution context', async () => { - mockExecutionContextService.decryptExecutionContext.mockResolvedValue({ - version: 1, - establishedAt: Date.now(), - source: 'webhook', - credentials: undefined, - } as PlaintextExecutionContext); + mockExecutionContextService.decryptCredentialContext.mockResolvedValue( + undefined as unknown as ICredentialContext, + ); await expect(service.checkCredentialStatus('workflow-1', executionContext)).rejects.toThrow( 'Execution context is present but contains no credential context. Ensure credential context establishment hooks are configured for this workflow.', @@ -244,15 +245,10 @@ describe('CredentialCheckProxyService', () => { }); it('should capture an empty identity in the intent when identity is missing', async () => { - mockExecutionContextService.decryptExecutionContext.mockResolvedValue({ + mockExecutionContextService.decryptCredentialContext.mockResolvedValue({ version: 1, - establishedAt: Date.now(), - source: 'webhook', - credentials: { - version: 1, - metadata: {}, - }, - } as PlaintextExecutionContext); + metadata: {}, + } as unknown as ICredentialContext); mockCredentialResolverWorkflowService.getWorkflowStatus.mockResolvedValue([ { diff --git a/packages/cli/src/modules/dynamic-credentials.ee/services/credential-check-proxy.service.ts b/packages/cli/src/modules/dynamic-credentials.ee/services/credential-check-proxy.service.ts index 323cf25fd20..f05f9b1cd4f 100644 --- a/packages/cli/src/modules/dynamic-credentials.ee/services/credential-check-proxy.service.ts +++ b/packages/cli/src/modules/dynamic-credentials.ee/services/credential-check-proxy.service.ts @@ -5,7 +5,6 @@ import type { CredentialCheckStatus, DynamicCredentialCheckProxyProvider, ICredentialContext, - IExecutionContext, } from 'n8n-workflow'; import { EnterpriseCredentialsService } from '@/credentials/credentials.service.ee'; @@ -30,11 +29,20 @@ export class CredentialCheckProxyService implements DynamicCredentialCheckProxyP async checkCredentialStatus( workflowId: string, - executionContext: IExecutionContext, + executionContext: { + credentials?: string; + }, ): Promise { - const plaintext = await this.executionContextService.decryptExecutionContext(executionContext); + if (!executionContext.credentials) { + throw new Error( + 'Execution context is present but contains no credential context. Ensure credential context establishment hooks are configured for this workflow.', + ); + } + const plaintext = await this.executionContextService.decryptCredentialContext( + executionContext.credentials, + ); - if (!plaintext.credentials) { + if (!plaintext) { throw new Error( 'Execution context is present but contains no credential context. Ensure credential context establishment hooks are configured for this workflow.', ); @@ -42,7 +50,7 @@ export class CredentialCheckProxyService implements DynamicCredentialCheckProxyP const statuses = await this.credentialResolverWorkflowService.getWorkflowStatus( workflowId, - plaintext.credentials, + plaintext, ); const credentials: CredentialCheckStatus[] = await Promise.all( @@ -59,7 +67,7 @@ export class CredentialCheckProxyService implements DynamicCredentialCheckProxyP checkStatus.authorizationUrl = await this.generateAuthorizationUrl( status.credentialId, status.resolverId, - plaintext.credentials!, + plaintext, ); } diff --git a/packages/cli/src/modules/oauth-server/__tests__/form-trigger-resource.resolver.api.test.ts b/packages/cli/src/modules/oauth-server/__tests__/form-trigger-resource.resolver.api.test.ts new file mode 100644 index 00000000000..6c629a84e7c --- /dev/null +++ b/packages/cli/src/modules/oauth-server/__tests__/form-trigger-resource.resolver.api.test.ts @@ -0,0 +1,297 @@ +import { + createWorkflowWithHistory, + setActiveVersion, + shareWorkflowWithUsers, + testDb, +} from '@n8n/backend-test-utils'; +import { GlobalConfig } from '@n8n/config'; +import type { User } from '@n8n/db'; +import { WebhookRepository, WorkflowRepository } from '@n8n/db'; +import { Container } from '@n8n/di'; +import type { INode } from 'n8n-workflow'; +import { FORM_TRIGGER_NODE_TYPE } from 'n8n-workflow'; +import { randomUUID } from 'node:crypto'; + +import { createOwner, createMember } from '@test-integration/db/users'; +import { setupTestServer } from '@test-integration/utils'; + +import { CacheService } from '@/services/cache/cache.service'; +import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; + +const testServer = setupTestServer({ modules: ['oauth-server', 'mcp'], endpointGroups: ['mcp'] }); + +let owner: User; +let member: User; +let formEndpoint: string; + +const webhookBaseUrl = () => Container.get(UrlService).getWebhookBaseUrl().replace(/\/$/, ''); +const resourceUrlFor = (webhookPath: string) => + `${webhookBaseUrl()}/${formEndpoint}/${webhookPath}`; +const prmPathFor = (webhookPath: string) => + `/.well-known/oauth-protected-resource/${formEndpoint}/${webhookPath}`; + +const formTriggerNode = ({ + name = 'On form submission', + authentication = 'n8nUserAuth', + disabled = false, + requireExecuteAccess, +}: { + name?: string; + authentication?: string; + disabled?: boolean; + requireExecuteAccess?: boolean; +} = {}): INode => ({ + id: randomUUID(), + name, + type: FORM_TRIGGER_NODE_TYPE, + typeVersion: 2, + position: [0, 0], + disabled, + parameters: { + path: 'unused', + authentication, + ...(requireExecuteAccess === undefined ? {} : { requireExecuteAccess }), + }, +}); + +/** Mirrors what `ActiveWorkflowManager.addWebhooks` persists on activation. */ +const insertWebhookRow = async (workflowId: string, webhookPath: string, node: string) => { + await Container.get(WebhookRepository).insert({ workflowId, webhookPath, method: 'POST', node }); +}; + +/** Active workflow whose published version contains the given trigger node. */ +const createPublishedFormWorkflow = async (webhookPath: string, node: INode, ownedBy = owner) => { + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, ownedBy); + await setActiveVersion(workflow.id, workflow.versionId); + await insertWebhookRow(workflow.id, webhookPath, node.name); + return workflow; +}; + +/** Overwrite the draft nodes without touching the published (active) version. */ +const updateDraftNodes = async (workflowId: string, nodes: INode[]) => { + await Container.get(WorkflowRepository).update(workflowId, { nodes, versionId: randomUUID() }); +}; + +const resolveResource = async (webhookPath: string) => + await Container.get(ProtectedResourceRegistry).getByResourcePath( + `/${formEndpoint}/${webhookPath}`, + ); + +beforeAll(async () => { + process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 = 'true'; // gates the form-trigger resolver + owner = await createOwner(); + member = await createMember(); + formEndpoint = Container.get(GlobalConfig).endpoints.form; +}); + +afterAll(() => { + delete process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2; +}); + +afterEach(async () => { + await Container.get(CacheService).reset(); // WebhookService caches static webhook lookups + await testDb.truncate([ + 'AccessToken', + 'RefreshToken', + 'AuthorizationCode', + 'OAuthClient', + 'WebhookEntity', + 'SharedWorkflow', + 'WorkflowEntity', + 'WorkflowHistory', + ]); +}); + +describe('protected resource metadata for form triggers', () => { + test('should serve the metadata document for an active n8nUserAuth form trigger', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(200); + // exact match: `scopes_supported` must be absent (the resource advertises no scopes) + expect(response.body).toEqual({ + resource: resourceUrlFor(webhookPath), + bearer_methods_supported: ['header'], + authorization_servers: [expect.any(String)], + }); + }); + + test('should resolve as a first-party resource whose only redirect URI is the trigger URL', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + const resource = await resolveResource(webhookPath); + + expect(resource?.isFirstParty).toBe(true); + expect(resource?.getResourceUrl()).toBe(resourceUrlFor(webhookPath)); + await expect(resource?.getAllowedRedirectUris?.()).resolves.toEqual([ + resourceUrlFor(webhookPath), + ]); + }); + + test('should expose the workflow name for the consent screen', async () => { + const webhookPath = randomUUID(); + const workflow = await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + const resource = await resolveResource(webhookPath); + + expect(resource?.displayName).toBe(workflow.name); + }); + + test('should not resolve an unknown path', async () => { + const response = await testServer.restlessAgent.get(prmPathFor(randomUUID())); + + expect(response.statusCode).toBe(404); + }); + + test('should not resolve when the feature flag is disabled', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + delete process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2; + try { + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + expect(response.statusCode).toBe(404); + } finally { + process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 = 'true'; + } + }); + + test('should not resolve a non-form path even if the webhook exists', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + const response = await testServer.restlessAgent.get( + `/.well-known/oauth-protected-resource/webhook/${webhookPath}`, + ); + + expect(response.statusCode).toBe(404); + }); + + test.each([ + ['authentication is none', formTriggerNode({ authentication: 'none' })], + ['authentication is basicAuth', formTriggerNode({ authentication: 'basicAuth' })], + ['authentication is an expression', formTriggerNode({ authentication: '={{ $json.auth }}' })], + ['the node is disabled', formTriggerNode({ disabled: true })], + ])('should not resolve when %s', async (_, node) => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, node); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(404); + }); + + test('should not resolve a workflow without a published version', async () => { + const node = formTriggerNode(); + const webhookPath = randomUUID(); + const workflow = await createWorkflowWithHistory({ active: false, nodes: [node] }, owner); + await insertWebhookRow(workflow.id, webhookPath, node.name); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(404); + }); + + test('should not resolve when the webhook node is missing from the active version', async () => { + const node = formTriggerNode(); + const webhookPath = randomUUID(); + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, owner); + await setActiveVersion(workflow.id, workflow.versionId); + await insertWebhookRow(workflow.id, webhookPath, 'Ghost node'); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(404); + }); + + test('should not resolve a dynamic webhook path', async () => { + const node = formTriggerNode(); + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, owner); + await setActiveVersion(workflow.id, workflow.versionId); + const webhookId = randomUUID(); + await Container.get(WebhookRepository).insert({ + workflowId: workflow.id, + webhookPath: ':param', + method: 'POST', + node: node.name, + webhookId, + pathLength: 1, + }); + + const response = await testServer.restlessAgent.get(prmPathFor(`${webhookId}/anything`)); + + expect(response.statusCode).toBe(404); + }); + + test('should follow the published version, not the draft', async () => { + // published n8nOAuth2, draft switched to none -> resource stays + const protectedPath = randomUUID(); + const protectedWorkflow = await createPublishedFormWorkflow(protectedPath, formTriggerNode()); + await updateDraftNodes(protectedWorkflow.id, [formTriggerNode({ authentication: 'none' })]); + + const stillProtected = await testServer.restlessAgent.get(prmPathFor(protectedPath)); + expect(stillProtected.statusCode).toBe(200); + expect(stillProtected.body.resource).toBe(resourceUrlFor(protectedPath)); + + // published none, draft switched to n8nOAuth2 -> no resource + const unprotectedPath = randomUUID(); + const unprotectedWorkflow = await createPublishedFormWorkflow( + unprotectedPath, + formTriggerNode({ authentication: 'none' }), + ); + await updateDraftNodes(unprotectedWorkflow.id, [formTriggerNode()]); + + const stillUnprotected = await testServer.restlessAgent.get(prmPathFor(unprotectedPath)); + expect(stillUnprotected.statusCode).toBe(404); + }); + + test('should stop resolving once the webhook is deregistered', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + expect((await testServer.restlessAgent.get(prmPathFor(webhookPath))).statusCode).toBe(200); + + await Container.get(WebhookRepository).delete({ webhookPath }); + await Container.get(CacheService).reset(); + + expect((await testServer.restlessAgent.get(prmPathFor(webhookPath))).statusCode).toBe(404); + }); +}); + +describe('authorize gate (workflow:execute)', () => { + test('authorizes the owner but denies a user without execute access', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + + const resource = await resolveResource(webhookPath); + + await expect(resource?.authorize(owner)).resolves.toBe(true); + await expect(resource?.authorize(member)).resolves.toBe(false); + }); + + test('authorizes a user granted execute via a project role', async () => { + const webhookPath = randomUUID(); + const workflow = await createPublishedFormWorkflow(webhookPath, formTriggerNode()); + await shareWorkflowWithUsers(workflow, [member]); + + const resource = await resolveResource(webhookPath); + + await expect(resource?.authorize(member)).resolves.toBe(true); + }); + + test('authorizes any authenticated user when require-execute is turned off', async () => { + const webhookPath = randomUUID(); + await createPublishedFormWorkflow( + webhookPath, + formTriggerNode({ requireExecuteAccess: false }), + ); + + const resource = await resolveResource(webhookPath); + + await expect(resource?.authorize(member)).resolves.toBe(true); + }); +}); diff --git a/packages/cli/src/modules/oauth-server/__tests__/form-trigger-test-resource.resolver.api.test.ts b/packages/cli/src/modules/oauth-server/__tests__/form-trigger-test-resource.resolver.api.test.ts new file mode 100644 index 00000000000..44a0ec4c20e --- /dev/null +++ b/packages/cli/src/modules/oauth-server/__tests__/form-trigger-test-resource.resolver.api.test.ts @@ -0,0 +1,276 @@ +import { createWorkflowWithHistory, setActiveVersion, testDb } from '@n8n/backend-test-utils'; +import { GlobalConfig } from '@n8n/config'; +import type { User } from '@n8n/db'; +import { WebhookRepository } from '@n8n/db'; +import { Container } from '@n8n/di'; +import type { INode, IWebhookData, IWorkflowBase } from 'n8n-workflow'; +import { FORM_TRIGGER_NODE_TYPE } from 'n8n-workflow'; +import { randomUUID } from 'node:crypto'; + +import { createOwner } from '@test-integration/db/users'; +import { setupTestServer } from '@test-integration/utils'; + +import { OAuthClientRepository } from '../database/repositories/oauth-client.repository'; +import { OAuthTokenService } from '@/modules/oauth-server/oauth-token.service'; +import { CacheService } from '@/services/cache/cache.service'; +import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; +import { TestWebhookRegistrationsService } from '@/webhooks/test-webhook-registrations.service'; + +const testServer = setupTestServer({ modules: ['oauth-server', 'mcp'], endpointGroups: ['mcp'] }); + +let owner: User; +let formEndpoint: string; +let formTestEndpoint: string; +let registrations: TestWebhookRegistrationsService; + +const webhookBaseUrl = () => Container.get(UrlService).getWebhookBaseUrl().replace(/\/$/, ''); +const testWebhookBaseUrl = () => + Container.get(UrlService).getTestWebhookBaseUrl().replace(/\/$/, ''); +const testResourceUrlFor = (webhookPath: string) => + `${testWebhookBaseUrl()}/${formTestEndpoint}/${webhookPath}`; +const prmPathFor = (webhookPath: string) => + `/.well-known/oauth-protected-resource/${formTestEndpoint}/${webhookPath}`; + +const formTriggerNode = ({ + name = 'On form submission', + authentication = 'n8nUserAuth', + disabled = false, +}: { name?: string; authentication?: string; disabled?: boolean } = {}): INode => ({ + id: randomUUID(), + name, + type: FORM_TRIGGER_NODE_TYPE, + typeVersion: 2, + position: [0, 0], + disabled, + parameters: { path: 'unused', authentication }, +}); + +/** Mirrors what `TestWebhooks.needsWebhook` registers when the user tests a form trigger. */ +const registerTestWebhook = async ( + webhookPath: string, + node: INode, + { + workflowId = randomUUID(), + workflowName = 'My test workflow', + }: { workflowId?: string; workflowName?: string } = {}, +) => { + await registrations.register({ + version: 1, + workflowEntity: { + id: workflowId, + name: workflowName, + active: false, + nodes: [node], + connections: {}, + } as IWorkflowBase, + webhook: { + httpMethod: 'POST', + path: webhookPath, + node: node.name, + workflowId, + } as IWebhookData, + }); + return { workflowId, workflowName }; +}; + +const resolveResource = async (webhookPath: string) => + await Container.get(ProtectedResourceRegistry).getByResourcePath( + `/${formTestEndpoint}/${webhookPath}`, + ); + +beforeAll(async () => { + process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 = 'true'; // gates the form-trigger resolver + owner = await createOwner(); + const { endpoints } = Container.get(GlobalConfig); + formEndpoint = endpoints.form; + formTestEndpoint = endpoints.formTest; + registrations = Container.get(TestWebhookRegistrationsService); +}); + +afterAll(() => { + delete process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2; +}); + +afterEach(async () => { + await Container.get(CacheService).reset(); // test webhook registrations live in the cache + await testDb.truncate([ + 'AccessToken', + 'RefreshToken', + 'AuthorizationCode', + 'OAuthClient', + 'WebhookEntity', + 'SharedWorkflow', + 'WorkflowEntity', + 'WorkflowHistory', + ]); +}); + +describe('protected resource metadata for test form triggers', () => { + test('should serve the metadata document while a test registration exists', async () => { + const webhookPath = randomUUID(); + await registerTestWebhook(webhookPath, formTriggerNode()); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(200); + expect(response.body).toEqual({ + resource: testResourceUrlFor(webhookPath), + bearer_methods_supported: ['header'], + authorization_servers: [expect.any(String)], + }); + }); + + test('should resolve as a first-party resource whose only redirect URI is the trigger URL', async () => { + const webhookPath = randomUUID(); + await registerTestWebhook(webhookPath, formTriggerNode()); + + const resource = await resolveResource(webhookPath); + + expect(resource?.isFirstParty).toBe(true); + await expect(resource?.getAllowedRedirectUris?.()).resolves.toEqual([ + testResourceUrlFor(webhookPath), + ]); + }); + + test('should resolve from the registration alone, without the workflow in the DB', async () => { + const webhookPath = randomUUID(); + await registerTestWebhook(webhookPath, formTriggerNode(), { workflowName: 'Unsaved workflow' }); + + const resource = await resolveResource(webhookPath); + + expect(resource?.displayName).toBe('Unsaved workflow'); + }); + + test('should not resolve an unknown test path', async () => { + const response = await testServer.restlessAgent.get(prmPathFor(randomUUID())); + + expect(response.statusCode).toBe(404); + }); + + test('should not resolve when the registration node name does not match', async () => { + const webhookPath = randomUUID(); + await registrations.register({ + version: 1, + workflowEntity: { + id: randomUUID(), + name: 'My test workflow', + active: false, + nodes: [formTriggerNode()], + connections: {}, + } as IWorkflowBase, + webhook: { + httpMethod: 'POST', + path: webhookPath, + node: 'Ghost node', + workflowId: randomUUID(), + } as IWebhookData, + }); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(404); + }); + + test.each([ + ['authentication is none', formTriggerNode({ authentication: 'none' })], + ['authentication is basicAuth', formTriggerNode({ authentication: 'basicAuth' })], + ['the node is disabled', formTriggerNode({ disabled: true })], + ])('should not resolve when %s', async (_, node) => { + const webhookPath = randomUUID(); + await registerTestWebhook(webhookPath, node); + + const response = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(response.statusCode).toBe(404); + }); + + test('should stop resolving as soon as the registration is removed', async () => { + const webhookPath = randomUUID(); + await registerTestWebhook(webhookPath, formTriggerNode()); + + expect((await testServer.restlessAgent.get(prmPathFor(webhookPath))).statusCode).toBe(200); + + await registrations.deregister(registrations.toKey({ httpMethod: 'POST', path: webhookPath })); + + expect((await testServer.restlessAgent.get(prmPathFor(webhookPath))).statusCode).toBe(404); + }); +}); + +describe('test vs production form resources', () => { + test('should serve the same trigger path as two distinct resources', async () => { + const webhookPath = randomUUID(); + const node = formTriggerNode(); + + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, owner); + await setActiveVersion(workflow.id, workflow.versionId); + await Container.get(WebhookRepository).insert({ + workflowId: workflow.id, + webhookPath, + method: 'POST', + node: node.name, + }); + await registerTestWebhook(webhookPath, node, { workflowId: workflow.id }); + + const production = await testServer.restlessAgent.get( + `/.well-known/oauth-protected-resource/${formEndpoint}/${webhookPath}`, + ); + const test = await testServer.restlessAgent.get(prmPathFor(webhookPath)); + + expect(production.statusCode).toBe(200); + expect(test.statusCode).toBe(200); + expect(production.body.resource).not.toBe(test.body.resource); + }); + + test('should reject a test-resource token at the production resource and vice versa', async () => { + const webhookPath = randomUUID(); + const node = formTriggerNode(); + + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, owner); + await setActiveVersion(workflow.id, workflow.versionId); + await Container.get(WebhookRepository).insert({ + workflowId: workflow.id, + webhookPath, + method: 'POST', + node: node.name, + }); + await registerTestWebhook(webhookPath, node, { workflowId: workflow.id }); + + const tokenService = Container.get(OAuthTokenService); + const productionResourceUrl = `${webhookBaseUrl()}/${formEndpoint}/${webhookPath}`; + const testResourceUrl = testResourceUrlFor(webhookPath); + + // A registered client is needed only to satisfy the token rows' FK. + const clientId = `client-${randomUUID()}`; + await Container.get(OAuthClientRepository).save({ + id: clientId, + name: 'Form resolver tests', + redirectUris: ['https://example.com/callback'], + grantTypes: ['authorization_code'], + tokenEndpointAuthMethod: 'none', + }); + + const mint = async (resourceUrl: string) => { + const pair = tokenService.generateTokenPair(owner.id, clientId, resourceUrl, []); + await tokenService.saveTokenPair(pair.accessToken, pair.refreshToken, clientId, owner.id, []); + return pair.accessToken; + }; + + const testToken = await mint(testResourceUrl); + const productionToken = await mint(productionResourceUrl); + + await expect(tokenService.verifyAccessToken(testToken, testResourceUrl)).resolves.toMatchObject( + { clientId }, + ); + await expect( + tokenService.verifyAccessToken(productionToken, productionResourceUrl), + ).resolves.toMatchObject({ clientId }); + + await expect( + tokenService.verifyAccessToken(testToken, productionResourceUrl), + ).rejects.toThrow(); + await expect( + tokenService.verifyAccessToken(productionToken, testResourceUrl), + ).rejects.toThrow(); + }); +}); diff --git a/packages/cli/src/modules/oauth-server/__tests__/oauth-clients.controller.api.test.ts b/packages/cli/src/modules/oauth-server/__tests__/oauth-clients.controller.api.test.ts index 3ad7a1a0355..b4624fb8b18 100644 --- a/packages/cli/src/modules/oauth-server/__tests__/oauth-clients.controller.api.test.ts +++ b/packages/cli/src/modules/oauth-server/__tests__/oauth-clients.controller.api.test.ts @@ -105,6 +105,45 @@ describe('GET /rest/mcp/oauth-clients', () => { expect(memberResponse.statusCode).toBe(200); expect(memberResponse.body.data.totals).toEqual({ mine: 1 }); }); + + test('should exclude first-party (internal) clients from the list and the totals', async () => { + const normalClient = await oauthClientRepository.save({ + id: 'normal-client', + name: 'Normal Client', + redirectUris: ['https://example.com/callback'], + grantTypes: ['authorization_code'], + tokenEndpointAuthMethod: 'none', + }); + const firstPartyClient = await oauthClientRepository.save({ + id: 'https://n8n.example.com/form/abc', + name: 'My Form', + redirectUris: ['https://n8n.example.com/form/abc'], + grantTypes: ['authorization_code'], + tokenEndpointAuthMethod: 'none', + isFirstParty: true, + }); + + await userConsentRepository.save({ + userId: owner.id, + clientId: normalClient.id, + grantedAt: Date.now(), + scope: [], + }); + await userConsentRepository.save({ + userId: owner.id, + clientId: firstPartyClient.id, + grantedAt: Date.now(), + scope: [], + }); + + const response = await testServer.authAgentFor(owner).get('/mcp/oauth-clients'); + + expect(response.statusCode).toBe(200); + expect(response.body.data.count).toBe(1); + expect(response.body.data.data).toHaveLength(1); + expect(response.body.data.data[0].id).toBe(normalClient.id); + expect(response.body.data.totals).toEqual({ mine: 1, all: 1 }); + }); }); describe('GET /rest/mcp/oauth-clients?ownership=all', () => { diff --git a/packages/cli/src/modules/oauth-server/__tests__/oauth-flow.service.api.test.ts b/packages/cli/src/modules/oauth-server/__tests__/oauth-flow.service.api.test.ts new file mode 100644 index 00000000000..9800ebff656 --- /dev/null +++ b/packages/cli/src/modules/oauth-server/__tests__/oauth-flow.service.api.test.ts @@ -0,0 +1,243 @@ +import { createWorkflowWithHistory, setActiveVersion, testDb } from '@n8n/backend-test-utils'; +import { GlobalConfig } from '@n8n/config'; +import type { User } from '@n8n/db'; +import { WebhookRepository } from '@n8n/db'; +import { Container } from '@n8n/di'; +import type { INode } from 'n8n-workflow'; +import { FORM_TRIGGER_NODE_TYPE, UserError } from 'n8n-workflow'; +import { randomUUID } from 'node:crypto'; + +import { CacheService } from '@/services/cache/cache.service'; +import { UrlService } from '@/services/url.service'; +import { createOwner, createMember } from '@test-integration/db/users'; +import { setupTestServer } from '@test-integration/utils'; + +import { OAuthAuthorizationCodeService } from '../oauth-authorization-code.service'; +import { OAuth2FlowService } from '../oauth-flow.service'; +import { OAuthServerService } from '../oauth-server.service'; +import { OAuthTokenService } from '../oauth-token.service'; + +// The flow service is driven directly via the DI container; the test server is +// set up only for the real DB + module registration (resolvers, token service). +setupTestServer({ modules: ['oauth-server', 'mcp'], endpointGroups: ['mcp'] }); + +let owner: User; +let member: User; +let formEndpoint: string; + +let flow: OAuth2FlowService; +let codes: OAuthAuthorizationCodeService; +let oauthServer: OAuthServerService; +let tokenService: OAuthTokenService; + +const webhookBaseUrl = () => Container.get(UrlService).getWebhookBaseUrl().replace(/\/$/, ''); +const resourceUrlFor = (webhookPath: string) => + `${webhookBaseUrl()}/${formEndpoint}/${webhookPath}`; + +const decodeJwtPayload = (token: string): Record => + JSON.parse(Buffer.from(token.split('.')[1], 'base64url').toString()) as Record; + +const formTriggerNode = (): INode => ({ + id: randomUUID(), + name: 'On form submission', + type: FORM_TRIGGER_NODE_TYPE, + typeVersion: 2, + position: [0, 0], + parameters: { path: 'unused', authentication: 'n8nUserAuth' }, +}); + +/** Active form workflow + production webhook row; returns the canonical resource URL. */ +const createProtectedFormWorkflow = async (ownedBy = owner) => { + const node = formTriggerNode(); + const webhookPath = randomUUID(); + const workflow = await createWorkflowWithHistory({ active: true, nodes: [node] }, ownedBy); + await setActiveVersion(workflow.id, workflow.versionId); + await Container.get(WebhookRepository).insert({ + workflowId: workflow.id, + webhookPath, + method: 'POST', + node: node.name, + }); + return resourceUrlFor(webhookPath); +}; + +/** + * Drive the browser legs the backend never performs itself in a test: pull the + * PKCE challenge + state out of the authorize URL, materialize the virtual client + * row, and mint the authorization code the AS would issue after consent. Returns + * the code + state to hand to `complete`. + */ +const authorizeAndMintCode = async ( + resourceUrl: string, + userId: string, + metadata?: Record, +) => { + const url = new URL(await flow.begin(resourceUrl, metadata)); + const state = url.searchParams.get('state')!; + const codeChallenge = url.searchParams.get('code_challenge')!; + + await oauthServer.clientsStore.getClient(resourceUrl); // lazy-upsert the virtual client row + const code = await codes.createAuthorizationCode( + resourceUrl, + userId, + resourceUrl, + codeChallenge, + state, + resourceUrl, + [], + ); + return { code, state }; +}; + +beforeAll(async () => { + process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 = 'true'; // gates the form-trigger resolver + owner = await createOwner(); + member = await createMember(); + formEndpoint = Container.get(GlobalConfig).endpoints.form; + flow = Container.get(OAuth2FlowService); + codes = Container.get(OAuthAuthorizationCodeService); + oauthServer = Container.get(OAuthServerService); + tokenService = Container.get(OAuthTokenService); +}); + +afterAll(() => { + delete process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2; +}); + +afterEach(async () => { + await Container.get(CacheService).reset(); + await testDb.truncate([ + 'AccessToken', + 'RefreshToken', + 'AuthorizationCode', + 'OAuthClient', + 'WebhookEntity', + 'SharedWorkflow', + 'WorkflowEntity', + 'WorkflowHistory', + ]); +}); + +describe('begin', () => { + test('builds an /oauth/authorize URL with coincident client_id, redirect_uri and resource', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + + const url = new URL(await flow.begin(resourceUrl)); + + expect(url.pathname).toBe('/oauth/authorize'); + expect(Object.fromEntries(url.searchParams)).toMatchObject({ + response_type: 'code', + client_id: resourceUrl, + redirect_uri: resourceUrl, + resource: resourceUrl, + code_challenge_method: 'S256', + }); + expect(url.searchParams.get('code_challenge')).toBeTruthy(); + expect(url.searchParams.get('state')).toBeTruthy(); + }); + + test('rejects a resource URL that is not a first-party protected resource', async () => { + await expect(flow.begin(resourceUrlFor(randomUUID()))).rejects.toThrow(UserError); + }); +}); + +describe('complete', () => { + test('exchanges the code for a validated token (sub=submitter, aud=form resource)', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + const { code, state } = await authorizeAndMintCode(resourceUrl, owner.id); + + const result = await flow.complete(code, state); + + expect(result).toMatchObject({ valid: true, user: { id: owner.id } }); + if (result.valid) { + expect(decodeJwtPayload(result.token).sub).toBe(owner.id); + expect(decodeJwtPayload(result.token).aud).toBe(resourceUrl); + } + }); + + test('returns metadata stashed at begin, and undefined when none was stashed', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + + const withMeta = await authorizeAndMintCode(resourceUrl, owner.id, { query: 'foo=bar' }); + const result = await flow.complete(withMeta.code, withMeta.state); + expect(result).toMatchObject({ valid: true, metadata: { query: 'foo=bar' } }); + + const withoutMeta = await authorizeAndMintCode(resourceUrl, owner.id); + const bareResult = await flow.complete(withoutMeta.code, withoutMeta.state); + expect(bareResult.valid).toBe(true); + if (bareResult.valid) expect(bareResult.metadata).toBeUndefined(); + }); + + test('rejects an unknown state', async () => { + const result = await flow.complete('some-code', 'unknown-state'); + + expect(result).toEqual({ valid: false, reason: 'invalid_state' }); + }); + + test('consumes the state so a replay is rejected', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + const { code, state } = await authorizeAndMintCode(resourceUrl, owner.id); + + await flow.complete(code, state); + const replay = await flow.complete(code, state); + + expect(replay).toEqual({ valid: false, reason: 'invalid_state' }); + }); + + test('rejects when the PKCE verifier does not match the code challenge', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + const url = new URL(await flow.begin(resourceUrl)); + const state = url.searchParams.get('state')!; + + await oauthServer.clientsStore.getClient(resourceUrl); + // Mint the code with a challenge that does NOT correspond to the cached verifier. + const code = await codes.createAuthorizationCode( + resourceUrl, + owner.id, + resourceUrl, + 'a-different-but-well-formed-code-challenge-value', + state, + resourceUrl, + [], + ); + + const result = await flow.complete(code, state); + + expect(result).toEqual({ valid: false, reason: 'invalid_grant' }); + }); + + test('rejects when the submitter lacks execute access on the workflow', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + const { code, state } = await authorizeAndMintCode(resourceUrl, member.id); + + const result = await flow.complete(code, state); + + expect(result).toEqual({ valid: false, reason: 'insufficient_scope' }); + }); + + test('produces a token that a different form resource rejects', async () => { + const resourceUrlA = await createProtectedFormWorkflow(); + const resourceUrlB = await createProtectedFormWorkflow(); + const { code, state } = await authorizeAndMintCode(resourceUrlA, owner.id); + + const result = await flow.complete(code, state); + expect(result.valid).toBe(true); + + if (result.valid) { + const crossResource = await tokenService.verifyOAuthAccessToken(result.token, resourceUrlB); + expect(crossResource.user).toBeNull(); + } + }); + + test('maps an already-consumed authorization code to invalid_grant instead of throwing', async () => { + const resourceUrl = await createProtectedFormWorkflow(); + const { code, state } = await authorizeAndMintCode(resourceUrl, owner.id); + // The loser of a concurrent completion (double-submitted callback) hits an + // already-used code; it must surface as a graceful result, not a thrown error. + await codes.markAuthorizationCodeAsUsed(code); + + const result = await flow.complete(code, state); + + expect(result).toEqual({ valid: false, reason: 'invalid_grant' }); + }); +}); diff --git a/packages/cli/src/modules/oauth-server/__tests__/oauth-server.service.test.ts b/packages/cli/src/modules/oauth-server/__tests__/oauth-server.service.test.ts index fb2893f5abc..66bec0e1d74 100644 --- a/packages/cli/src/modules/oauth-server/__tests__/oauth-server.service.test.ts +++ b/packages/cli/src/modules/oauth-server/__tests__/oauth-server.service.test.ts @@ -1,4 +1,3 @@ -import type { Mock, Mocked } from 'vitest'; import { InvalidGrantError, InvalidTargetError, @@ -7,8 +6,16 @@ import { Logger } from '@n8n/backend-common'; import { mockInstance } from '@n8n/backend-test-utils'; import { GlobalConfig } from '@n8n/config'; import type { Response } from 'express'; +import type { Mock, Mocked } from 'vitest'; import { mock } from 'vitest-mock-extended'; +import { McpProtectedResource } from '@/modules/mcp/mcp-protected-resource'; +import type { McpConfig } from '@/modules/mcp/mcp.config'; +import type { McpSettingsService } from '@/modules/mcp/mcp.settings.service'; +import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import type { UrlService } from '@/services/url.service'; +import { UserManagementMailer } from '@/user-management/email'; + import type { AuthorizationCode } from '../database/entities/oauth-authorization-code.entity'; import type { OAuthClient } from '../database/entities/oauth-client.entity'; import { OAuthClientRepository } from '../database/repositories/oauth-client.repository'; @@ -17,12 +24,6 @@ import { OAuthAuthorizationCodeService } from '../oauth-authorization-code.servi import { OAuthServerService } from '../oauth-server.service'; import { OAuthSessionService } from '../oauth-session.service'; import { OAuthTokenService } from '../oauth-token.service'; -import { McpProtectedResource } from '@/modules/mcp/mcp-protected-resource'; -import type { McpConfig } from '@/modules/mcp/mcp.config'; -import type { McpSettingsService } from '@/modules/mcp/mcp.settings.service'; -import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; -import type { UrlService } from '@/services/url.service'; -import { UserManagementMailer } from '@/user-management/email'; const SUPPORTED_SCOPES = ['tool:listWorkflows', 'tool:getWorkflowDetails']; const TEST_RESOURCE_URL = 'https://n8n.example.com/mcp-server/http'; @@ -37,6 +38,9 @@ let userConsentRepository: Mocked; let mailer: Mocked; let getAllowedRedirectUris: Mock<() => Promise>; +// Shared, immutable across tests: the base URLs gate the first-party client_id guard. +const urlServiceMock = mock(); + describe('OAuthServerService', () => { beforeAll(() => { logger = mockInstance(Logger); @@ -46,6 +50,8 @@ describe('OAuthServerService', () => { authorizationCodeService = mockInstance(OAuthAuthorizationCodeService); userConsentRepository = mockInstance(UserConsentRepository); mailer = mockInstance(UserManagementMailer); + urlServiceMock.getWebhookBaseUrl.mockReturnValue('https://n8n.example.com/'); + urlServiceMock.getTestWebhookBaseUrl.mockReturnValue('https://n8n.example.com/'); getAllowedRedirectUris = vi.fn<(...args: []) => Promise>().mockResolvedValue([]); const resourceRegistry = new ProtectedResourceRegistry(mock()); @@ -69,6 +75,7 @@ describe('OAuthServerService', () => { userConsentRepository, resourceRegistry, mailer, + urlServiceMock, ); }); @@ -137,6 +144,138 @@ describe('OAuthServerService', () => { }); }); + describe('getClient — virtual first-party client', () => { + const FIRST_PARTY_URL = 'https://n8n.example.com/form/abc'; + const NON_FIRST_PARTY_URL = 'https://n8n.example.com/mcp-server/http'; + let firstPartyService: OAuthServerService; + + beforeAll(() => { + const registry = new ProtectedResourceRegistry(mock()); + registry.register({ + id: 'form-abc', + isFirstParty: true, + displayName: 'My Form', + getResourceUrl: () => FIRST_PARTY_URL, + getAudiences: () => [FIRST_PARTY_URL], + getAllowedRedirectUris: async () => [FIRST_PARTY_URL], + scopes: [], + authorize: async () => true, + }); + // A resource that exists but is not first-party (mirror of an MCP resource). + registry.register({ + id: 'mcp-x', + getResourceUrl: () => NON_FIRST_PARTY_URL, + getAudiences: () => [NON_FIRST_PARTY_URL], + scopes: [], + authorize: async () => true, + }); + firstPartyService = new OAuthServerService( + logger, + mockInstance(GlobalConfig), + oauthSessionService, + oauthClientRepository, + tokenService, + authorizationCodeService, + userConsentRepository, + registry, + mailer, + urlServiceMock, + ); + }); + + it('lazily upserts and returns a virtual client on a DB miss for a first-party resource', async () => { + oauthClientRepository.findOneBy.mockResolvedValue(null); + + const result = await firstPartyService.clientsStore.getClient(FIRST_PARTY_URL); + + expect(oauthClientRepository.upsert).toHaveBeenCalledWith( + { + id: FIRST_PARTY_URL, + name: 'My Form', + redirectUris: [FIRST_PARTY_URL], + grantTypes: ['authorization_code'], + tokenEndpointAuthMethod: 'none', + clientSecret: null, + clientSecretExpiresAt: null, + isFirstParty: true, + }, + ['id'], + ); + expect(result).toEqual({ + client_id: FIRST_PARTY_URL, + client_name: 'My Form', + redirect_uris: [FIRST_PARTY_URL], + grant_types: ['authorization_code'], + token_endpoint_auth_method: 'none', + response_types: ['code'], + logo_uri: undefined, + tos_uri: undefined, + }); + }); + + it('returns undefined and does not upsert when the resolved resource is not first-party', async () => { + oauthClientRepository.findOneBy.mockResolvedValue(null); + + const result = await firstPartyService.clientsStore.getClient(NON_FIRST_PARTY_URL); + + expect(result).toBeUndefined(); + expect(oauthClientRepository.upsert).not.toHaveBeenCalled(); + }); + + it('returns undefined and does not upsert when no resource resolves', async () => { + oauthClientRepository.findOneBy.mockResolvedValue(null); + + const result = await firstPartyService.clientsStore.getClient( + 'https://n8n.example.com/form/unknown', + ); + + expect(result).toBeUndefined(); + expect(oauthClientRepository.upsert).not.toHaveBeenCalled(); + }); + + it('short-circuits without consulting the resolver registry for a non-webhook client_id', async () => { + oauthClientRepository.findOneBy.mockResolvedValue(null); + const registry = new ProtectedResourceRegistry(mock()); + const getByResourceUrl = vi.spyOn(registry, 'getByResourceUrl'); + const svc = new OAuthServerService( + logger, + mockInstance(GlobalConfig), + oauthSessionService, + oauthClientRepository, + tokenService, + authorizationCodeService, + userConsentRepository, + registry, + mailer, + urlServiceMock, + ); + + const result = await svc.clientsStore.getClient('https://evil.example.com/form/abc'); + + expect(result).toBeUndefined(); + expect(getByResourceUrl).not.toHaveBeenCalled(); + expect(oauthClientRepository.upsert).not.toHaveBeenCalled(); + }); + }); + + describe('registered-client cap excludes first-party clients', () => { + it('counts only non-first-party clients for the limit check', async () => { + oauthClientRepository.countBy.mockResolvedValue(0); + + await service.isClientLimitReached(); + + expect(oauthClientRepository.countBy).toHaveBeenCalledWith({ isFirstParty: false }); + }); + + it('counts only non-first-party clients for the instance stats', async () => { + oauthClientRepository.countBy.mockResolvedValue(0); + + await service.getInstanceClientStats(); + + expect(oauthClientRepository.countBy).toHaveBeenCalledWith({ isFirstParty: false }); + }); + }); + describe('registerClient', () => { it('should save client with all required fields', async () => { const clientInfo = { @@ -163,6 +302,7 @@ describe('OAuthServerService', () => { clientSecret: null, clientSecretExpiresAt: null, tokenEndpointAuthMethod: 'none', + isFirstParty: false, }); expect(result).toEqual(clientInfo); }); @@ -194,6 +334,7 @@ describe('OAuthServerService', () => { clientSecret: 'secret-123', clientSecretExpiresAt: 1234567890, tokenEndpointAuthMethod: 'client_secret_post', + isFirstParty: false, }); }); @@ -1177,6 +1318,7 @@ describe('OAuthServerService', () => { userConsentRepository, multiRegistry, mailer, + urlServiceMock, ); expect( @@ -1222,6 +1364,7 @@ describe('OAuthServerService', () => { userConsentRepository, configuredRegistry, mailer, + urlService, ); }; diff --git a/packages/cli/src/modules/oauth-server/database/entities/oauth-client.entity.ts b/packages/cli/src/modules/oauth-server/database/entities/oauth-client.entity.ts index 47aa98d2101..26b2e83609f 100644 --- a/packages/cli/src/modules/oauth-server/database/entities/oauth-client.entity.ts +++ b/packages/cli/src/modules/oauth-server/database/entities/oauth-client.entity.ts @@ -23,6 +23,9 @@ export class OAuthClient extends WithTimestamps { @Column({ type: String, default: 'none' }) tokenEndpointAuthMethod: string; + @Column({ type: Boolean, default: false }) + isFirstParty: boolean; + @OneToMany('AuthorizationCode', 'client') authorizationCodes: AuthorizationCode[]; diff --git a/packages/cli/src/modules/oauth-server/database/repositories/oauth-user-consent.repository.ts b/packages/cli/src/modules/oauth-server/database/repositories/oauth-user-consent.repository.ts index 1cabc26c996..f1f7503b17a 100644 --- a/packages/cli/src/modules/oauth-server/database/repositories/oauth-user-consent.repository.ts +++ b/packages/cli/src/modules/oauth-server/database/repositories/oauth-user-consent.repository.ts @@ -53,7 +53,11 @@ export class UserConsentRepository extends Repository { async findConnectedClients( options: FindConnectedClientsOptions, ): Promise<{ rows: UserConsent[]; total: number }> { - const qb = this.createQueryBuilder('consent').leftJoinAndSelect('consent.client', 'client'); + const qb = this.createQueryBuilder('consent') + .leftJoinAndSelect('consent.client', 'client') + // First-party clients are internal per-resource virtual clients (e.g. form + // triggers), not user-manageable OAuth clients — keep them out of the list. + .andWhere('client.isFirstParty = :isFirstParty', { isFirstParty: false }); if (options.withOwner) qb.leftJoinAndSelect('consent.user', 'user'); if (options.userId) qb.andWhere('consent.userId = :userId', { userId: options.userId }); @@ -87,10 +91,13 @@ export class UserConsentRepository extends Repository { /** * Distinct owners across every consent, for the "Connected by" filter. Not * scoped by the current filters, so the dropdown always lists all owners. + * First-party (internal) client consents are excluded, matching the list. */ async findConsentOwners(): Promise { return await this.createQueryBuilder('consent') .innerJoin('consent.user', 'user') + .innerJoin('consent.client', 'client') + .where('client.isFirstParty = :isFirstParty', { isFirstParty: false }) .select('user.id', 'id') .addSelect('user.firstName', 'firstName') .addSelect('user.lastName', 'lastName') @@ -98,4 +105,18 @@ export class UserConsentRepository extends Repository { .distinct(true) .getRawMany(); } + + /** + * Count connected-client consents, optionally scoped to a single user. Backs + * the "mine"/"all" tab totals. Excludes first-party (internal) client consents + * so the totals match the listing. + */ + async countConnectedConsents(userId?: string): Promise { + return await this.count({ + where: { + ...(userId ? { userId } : {}), + client: { isFirstParty: false }, + }, + }); + } } diff --git a/packages/cli/src/modules/oauth-server/oauth-flow.service.ts b/packages/cli/src/modules/oauth-server/oauth-flow.service.ts new file mode 100644 index 00000000000..229959f518d --- /dev/null +++ b/packages/cli/src/modules/oauth-server/oauth-flow.service.ts @@ -0,0 +1,125 @@ +import { InvalidGrantError } from '@modelcontextprotocol/sdk/server/auth/errors.js'; +import { Time } from '@n8n/constants'; +import { Service } from '@n8n/di'; +import { UserError, type N8nOAuth2FlowResult } from 'n8n-workflow'; +import { createHash, randomBytes } from 'node:crypto'; +import pkceChallenge from 'pkce-challenge'; + +import { CacheService } from '@/services/cache/cache.service'; +import { OAuthTokenVerifierProxy } from '@/services/oauth-token-verifier-proxy.service'; +import type { N8nOAuth2Flow } from '@/services/oauth2-flow-proxy.service'; +import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; + +import { OAuthServerService } from './oauth-server.service'; + +const FLOW_STATE_PREFIX = 'oauth-flow:'; +const FLOW_STATE_TTL = 5 * Time.minutes.toMilliseconds; + +type FlowState = { codeVerifier: string; resourceUrl: string; metadata?: Record }; + +@Service() +export class OAuth2FlowService implements N8nOAuth2Flow { + constructor( + private readonly oauthServerService: OAuthServerService, + private readonly resourceRegistry: ProtectedResourceRegistry, + private readonly urlService: UrlService, + private readonly cacheService: CacheService, + private readonly tokenVerifier: OAuthTokenVerifierProxy, + ) {} + + /** + * Begin authorization-code + PKCE for a form trigger: validate the resource is + * first-party, stash the PKCE verifier under an unguessable single-use `state`, + * and return the `/oauth/authorize` URL to redirect the browser to. For form + * triggers client_id = redirect_uri = resource = the trigger URL. + * + * Optional `metadata` is stashed against the `state` (server-side, never sent to the + * browser) and handed back by `complete` on success. + */ + async begin(resourceUrl: string, metadata?: Record): Promise { + const resource = await this.resourceRegistry.getByResourceUrl(resourceUrl); + if (!resource?.isFirstParty) { + throw new UserError(`Not a first-party protected resource: ${resourceUrl}`); + } + + const { code_verifier, code_challenge } = await pkceChallenge(); + const state = randomBytes(32).toString('hex'); + await this.cacheService.set( + FLOW_STATE_PREFIX + state, + { codeVerifier: code_verifier, resourceUrl, metadata } satisfies FlowState, + FLOW_STATE_TTL, + ); + + const url = new URL(`${this.urlService.getInstanceBaseUrl()}/oauth/authorize`); + url.searchParams.set('response_type', 'code'); + url.searchParams.set('client_id', resourceUrl); + url.searchParams.set('redirect_uri', resourceUrl); + url.searchParams.set('resource', resourceUrl); + url.searchParams.set('code_challenge', code_challenge); + url.searchParams.set('code_challenge_method', 'S256'); + url.searchParams.set('state', state); + return url.toString(); + } + + /** + * Complete the flow: consume the cached state once, verify PKCE (the SDK token + * handler normally does this — we exchange in-process, so replicate the S256 + * check), exchange the code for an AS token, and validate it against the form's + * resource. Ends with a validated token (sub=submitter, aud=form resource). + */ + async complete(code: string, state: string): Promise { + const cacheKey = FLOW_STATE_PREFIX + state; + const flow = await this.cacheService.get(cacheKey); + if (!flow) return { valid: false, reason: 'invalid_state' }; + await this.cacheService.delete(cacheKey); // consume-once + + const { codeVerifier, resourceUrl, metadata } = flow; + + const client = await this.oauthServerService.clientsStore.getClient(resourceUrl); + if (!client) return { valid: false, reason: 'invalid_client' }; + + try { + const challenge = await this.oauthServerService.challengeForAuthorizationCode(client, code); + if (createHash('sha256').update(codeVerifier).digest('base64url') !== challenge) { + return { valid: false, reason: 'invalid_grant' }; + } + + const tokens = await this.oauthServerService.exchangeAuthorizationCode( + client, + code, + codeVerifier, + resourceUrl, + new URL(resourceUrl), + ); + + const result = await this.tokenVerifier.verifyOAuthAccessToken( + tokens.access_token, + resourceUrl, + ); + if (!result.user) { + return { valid: false, reason: result.context?.reason ?? 'invalid_token' }; + } + return { + valid: true, + token: tokens.access_token, + user: { + id: result.user.id, + email: result.user.email, + firstName: result.user.firstName, + lastName: result.user.lastName, + }, + metadata, + }; + } catch (error) { + // A concurrent completion (double-submitted callback) loses the atomic + // `markAuthorizationCodeAsUsed` race; a missing/used/expired code throws the + // same way. Surface it as a graceful invalid_grant rather than a thrown 500 — + // the return contract is a discriminated union, not an exception channel. + if (error instanceof InvalidGrantError) { + return { valid: false, reason: 'invalid_grant' }; + } + throw error; + } + } +} diff --git a/packages/cli/src/modules/oauth-server/oauth-server.module.ts b/packages/cli/src/modules/oauth-server/oauth-server.module.ts index d4ecb01f0ed..7ac271ab4a1 100644 --- a/packages/cli/src/modules/oauth-server/oauth-server.module.ts +++ b/packages/cli/src/modules/oauth-server/oauth-server.module.ts @@ -1,4 +1,3 @@ -import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; import type { ModuleInterface } from '@n8n/decorators'; import { BackendModule } from '@n8n/decorators'; import { Container } from '@n8n/di'; @@ -35,19 +34,14 @@ export class OAuthServerModule implements ModuleInterface { const { OAuthTokenService } = await import('./oauth-token.service.js'); Container.get(OAuthTokenVerifierProxy).registerProvider(Container.get(OAuthTokenService)); - const { WorkflowMcpTriggerResourceResolver } = await import( - './protected-resource-resolvers/workflow-mcp-trigger-resource.resolver.js' - ); - Container.get(ProtectedResourceRegistry).registerResolver( - Container.get(WorkflowMcpTriggerResourceResolver), + const { registerProtectedResourceResolvers } = await import( + './protected-resource-resolvers/index.js' ); + registerProtectedResourceResolvers(); - const { WorkflowMcpTestTriggerResourceResolver } = await import( - './protected-resource-resolvers/workflow-mcp-test-trigger-resource.resolver.js' - ); - Container.get(ProtectedResourceRegistry).registerResolver( - Container.get(WorkflowMcpTestTriggerResourceResolver), - ); + const { OAuth2FlowProxy } = await import('@/services/oauth2-flow-proxy.service.js'); + const { OAuth2FlowService } = await import('./oauth-flow.service.js'); + Container.get(OAuth2FlowProxy).registerProvider(Container.get(OAuth2FlowService)); } async entities() { diff --git a/packages/cli/src/modules/oauth-server/oauth-server.service.ts b/packages/cli/src/modules/oauth-server/oauth-server.service.ts index 9347cd834bc..8d66083c160 100644 --- a/packages/cli/src/modules/oauth-server/oauth-server.service.ts +++ b/packages/cli/src/modules/oauth-server/oauth-server.service.ts @@ -22,7 +22,10 @@ import { Service } from '@n8n/di'; import { hasGlobalScope } from '@n8n/permissions'; import type { Response } from 'express'; +import { ForbiddenError } from '@/errors/response-errors/forbidden.error'; import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; +import { UserManagementMailer } from '@/user-management/email'; import { OAuthClient } from './database/entities/oauth-client.entity'; import { OAuthClientRepository } from './database/repositories/oauth-client.repository'; @@ -31,8 +34,6 @@ import { OAuthAuthorizationCodeService } from './oauth-authorization-code.servic import { OAuthSessionService } from './oauth-session.service'; import { OAuthTokenService } from './oauth-token.service'; import { OAuthClientLimitReachedError } from './oauth.errors'; -import { ForbiddenError } from '@/errors/response-errors/forbidden.error'; -import { UserManagementMailer } from '@/user-management/email'; /** Maximum number of redirect URIs per client */ const MAX_REDIRECT_URIS = 10; @@ -102,6 +103,7 @@ export class OAuthServerService implements OAuthServerProvider { private readonly userConsentRepository: UserConsentRepository, private readonly resourceRegistry: ProtectedResourceRegistry, private readonly mailer: UserManagementMailer, + private readonly urlService: UrlService, ) {} get clientsStore(): OAuthRegisteredClientsStore { @@ -109,7 +111,7 @@ export class OAuthServerService implements OAuthServerProvider { getClient: async (clientId: string): Promise => { const client = await this.oauthClientRepository.findOneBy({ id: clientId }); if (!client) { - return undefined; + return await this.resolveVirtualClient(clientId); } // Some clients echo back the `scope` they saw on registration and @@ -146,6 +148,7 @@ export class OAuthServerService implements OAuthServerProvider { clientSecret: client.client_secret ?? null, clientSecretExpiresAt: client.client_secret_expires_at ?? null, tokenEndpointAuthMethod: client.token_endpoint_auth_method ?? 'none', + isFirstParty: false, }); await this.enforceClientLimit(client.client_id); @@ -157,7 +160,7 @@ export class OAuthServerService implements OAuthServerProvider { /** Returns true when the instance is already at or above the registered-client cap. */ async isClientLimitReached(): Promise { - const clientCount = await this.oauthClientRepository.count(); + const clientCount = await this.oauthClientRepository.countBy({ isFirstParty: false }); return clientCount >= this.globalConfig.endpoints.mcpMaxRegisteredClients; } @@ -166,7 +169,7 @@ export class OAuthServerService implements OAuthServerProvider { limit: number; atCapacity: boolean; }> { - const count = await this.oauthClientRepository.count(); + const count = await this.oauthClientRepository.countBy({ isFirstParty: false }); const limit = this.globalConfig.endpoints.mcpMaxRegisteredClients; return { count, limit, atCapacity: count >= limit }; } @@ -180,7 +183,7 @@ export class OAuthServerService implements OAuthServerProvider { * — matching the response shape of the pre-check guard at the route layer. */ private async enforceClientLimit(clientId: string): Promise { - const clientCount = await this.oauthClientRepository.count(); + const clientCount = await this.oauthClientRepository.countBy({ isFirstParty: false }); const limit = this.globalConfig.endpoints.mcpMaxRegisteredClients; if (clientCount > limit) { await this.oauthClientRepository.delete({ id: clientId }); @@ -192,6 +195,62 @@ export class OAuthServerService implements OAuthServerProvider { } } + /** + * On-demand per-trigger virtual client for a first-party protected resource + * (form trigger). Public + PKCE, single redirect_uri = the trigger URL (which + * equals the client_id and the resource URL). The row is persisted lazily only + * to satisfy the FKs from auth codes / tokens; it is never a DCR client and is + * excluded from the registered-client cap. + */ + private async resolveVirtualClient( + clientId: string, + ): Promise { + // First-party resources are form triggers served under the (test) webhook base + // URL, so a client_id that isn't can never resolve to one. Skip the resolver + // sweep + lazy upsert for anything else, so the unauthenticated /authorize path + // can't be used to fan out DB lookups on arbitrary client_ids. + if (!this.isFormTriggerClientId(clientId)) { + return undefined; + } + + const resource = await this.resourceRegistry.getByResourceUrl(clientId); + if (!resource?.isFirstParty) { + return undefined; + } + + await this.oauthClientRepository.upsert( + { + id: clientId, + name: resource.displayName ?? clientId, + redirectUris: [clientId], + grantTypes: ['authorization_code'], + tokenEndpointAuthMethod: 'none', + clientSecret: null, + clientSecretExpiresAt: null, + isFirstParty: true, + }, + ['id'], + ); + + return { + client_id: clientId, + client_name: resource.displayName ?? clientId, + redirect_uris: [clientId], + grant_types: ['authorization_code'], + token_endpoint_auth_method: 'none', + response_types: ['code'], + logo_uri: undefined, + tos_uri: undefined, + }; + } + + /** Whether a client_id could be a form-trigger resource URL (served under a webhook base URL). */ + private isFormTriggerClientId(clientId: string): boolean { + return [this.urlService.getWebhookBaseUrl(), this.urlService.getTestWebhookBaseUrl()] + .map((base) => (base.endsWith('/') ? base : `${base}/`)) + .some((base) => clientId.startsWith(base)); + } + private validateClientRegistration(client: OAuthClientInformationFull): void { if (!client.client_name) { throw new Error('client_name is required'); @@ -506,6 +565,7 @@ export class OAuthServerService implements OAuthServerProvider { if (options.type) { const registered = await this.oauthClientRepository.find({ select: { id: true, name: true }, + where: { isFirstParty: false }, }); clientIds = registered .filter((client) => matchesTypeFilter(client.name, options.type!)) @@ -552,8 +612,8 @@ export class OAuthServerService implements OAuthServerProvider { // dedicated counts rather than the filtered page above. const [consentOwners, mineCount, allCount] = await Promise.all([ listAll ? this.userConsentRepository.findConsentOwners() : undefined, - this.userConsentRepository.countBy({ userId: user.id }), - canSeeAll ? this.userConsentRepository.count() : undefined, + this.userConsentRepository.countConnectedConsents(user.id), + canSeeAll ? this.userConsentRepository.countConnectedConsents() : undefined, ]); const owners = consentOwners ? sortOwners(consentOwners) : undefined; const totals: ConnectedOAuthClientTotals = { mine: mineCount }; diff --git a/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-resource.resolver.ts b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-resource.resolver.ts new file mode 100644 index 00000000000..198eed82076 --- /dev/null +++ b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-resource.resolver.ts @@ -0,0 +1,105 @@ +import type { ProtectedResourceResolver } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; +import { WebhookService } from '@/webhooks/webhook.service'; +import { WorkflowFinderService } from '@/workflows/workflow-finder.service'; +import { Logger } from '@n8n/backend-common'; +import { GlobalConfig } from '@n8n/config'; +import { User, WorkflowRepository } from '@n8n/db'; +import { Service } from '@n8n/di'; +import { FORM_TRIGGER_NODE_TYPE } from 'n8n-workflow'; + +import { + FORM_TRIGGER_SCOPES, + isFormOAuth2Enabled, + resourceUrlToWebhookPath, + trimSlashes, + trimTrailingSlash, +} from './utils'; + +@Service() +export class FormTriggerResourceResolver implements ProtectedResourceResolver { + constructor( + private readonly config: GlobalConfig, + private readonly webhookService: WebhookService, + private readonly workflowRepository: WorkflowRepository, + private readonly urlService: UrlService, + private readonly logger: Logger, + private readonly workflowFinderService: WorkflowFinderService, + ) {} + + readonly id = 'form-trigger'; + readonly scopes = FORM_TRIGGER_SCOPES; + + async resolveByUrl(resourceUrl: string) { + const pathname = resourceUrlToWebhookPath(resourceUrl, this.urlService.getWebhookBaseUrl()); + if (pathname === undefined) { + this.logger.debug(`Resource URL is not under the webhook base URL: ${resourceUrl}`); + return undefined; + } + return await this.resolveByPath(pathname); + } + + async resolveByPath(pathname: string) { + if (!isFormOAuth2Enabled()) { + return undefined; + } + + if (!pathname.startsWith(`/${this.config.endpoints.form}/`)) { + return undefined; + } + + const path = trimSlashes(pathname.slice(this.config.endpoints.form.length + 1)); + + const webhook = await this.webhookService.findStaticWebhook('POST', path); + if (!webhook || webhook.isDynamic) { + return undefined; + } + + const { workflowId, node: nodeName } = webhook; + + const workflow = await this.workflowRepository.findOne({ + where: { id: workflowId }, + relations: { activeVersion: true }, + }); + if (!workflow?.activeVersion) { + return undefined; + } + + const node = workflow.activeVersion.nodes.find((n) => n.name === nodeName); + if (!node) { + return undefined; + } + + if ( + node.type === FORM_TRIGGER_NODE_TYPE && + !node.disabled && + node.parameters.authentication === 'n8nUserAuth' + ) { + const resourceUrl = `${trimTrailingSlash(this.urlService.getWebhookBaseUrl())}/${this.config.endpoints.form}/${path}`; + const requireExecute = node.parameters.requireExecuteAccess !== false; + return { + id: 'workflow-form:' + workflow.id, + isFirstParty: true, + getResourceUrl: () => resourceUrl, + getAudiences: () => [resourceUrl], + getAllowedRedirectUris: async () => [resourceUrl], + scopes: FORM_TRIGGER_SCOPES, + displayName: workflow.name, + authorize: async (user: User) => { + if (requireExecute) { + return ( + await this.workflowFinderService.findWorkflowIdsWithScopeForUser( + [workflow.id], + user, + ['workflow:execute'], + ) + ).has(workflow.id); + } + return true; + }, + }; + } + + return undefined; + } +} diff --git a/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-test-resource.resolver.ts b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-test-resource.resolver.ts new file mode 100644 index 00000000000..c3882b983ea --- /dev/null +++ b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/form-trigger-test-resource.resolver.ts @@ -0,0 +1,108 @@ +import type { ProtectedResourceResolver } from '@/services/protected-resource.registry'; +import { UrlService } from '@/services/url.service'; +import { WorkflowFinderService } from '@/workflows/workflow-finder.service'; +import { Logger } from '@n8n/backend-common'; +import { GlobalConfig } from '@n8n/config'; +import { User } from '@n8n/db'; +import { Service } from '@n8n/di'; +import { FORM_TRIGGER_NODE_TYPE } from 'n8n-workflow'; + +import { + FORM_TRIGGER_SCOPES, + isFormOAuth2Enabled, + resourceUrlToWebhookPath, + trimSlashes, + trimTrailingSlash, +} from './utils'; +import { TestWebhookRegistrationsService } from '@/webhooks/test-webhook-registrations.service'; + +@Service() +export class FormTriggerTestResourceResolver implements ProtectedResourceResolver { + constructor( + private readonly config: GlobalConfig, + private readonly registrations: TestWebhookRegistrationsService, + private readonly urlService: UrlService, + private readonly logger: Logger, + private readonly workflowFinderService: WorkflowFinderService, + ) {} + + readonly id = 'form-trigger-test'; + readonly scopes = FORM_TRIGGER_SCOPES; + + async resolveByUrl(resourceUrl: string) { + const pathname = resourceUrlToWebhookPath(resourceUrl, this.urlService.getTestWebhookBaseUrl()); + if (pathname === undefined) { + this.logger.debug(`Resource URL is not under the webhook base URL: ${resourceUrl}`); + return undefined; + } + return await this.resolveByPath(pathname); + } + + async resolveByPath(pathname: string) { + if (!isFormOAuth2Enabled()) { + return undefined; + } + + if (!pathname.startsWith(`/${this.config.endpoints.formTest}/`)) { + return undefined; + } + + const path = trimSlashes(pathname.slice(this.config.endpoints.formTest.length + 1)); + + // The registration holds the workflow exactly as the editor is testing it + // (including unsaved changes), so it is the source of truth here — not the + // DB draft. Only the static key shape is looked up: dynamic registrations + // are keyed differently and are not protectable resources. + const registration = await this.registrations.get( + this.registrations.toKey({ httpMethod: 'POST', path }), + ); + + if (!registration) { + this.logger.debug(`No test webhook registration found for path: ${path}`); + return undefined; + } + + const { workflowEntity, webhook } = registration; + + const node = workflowEntity.nodes.find((n) => n.name === webhook.node); + + if (!node) { + this.logger.debug( + `No node found with name ${webhook.node} in test registration for workflow with ID: ${workflowEntity.id}`, + ); + return undefined; + } + + if ( + node.type === FORM_TRIGGER_NODE_TYPE && + !node.disabled && + node.parameters.authentication === 'n8nUserAuth' + ) { + const resourceUrl = `${trimTrailingSlash(this.urlService.getTestWebhookBaseUrl())}/${this.config.endpoints.formTest}/${path}`; + const requireExecute = node.parameters.requireExecuteAccess !== false; + return { + id: 'workflow-form:' + workflowEntity.id, + isFirstParty: true, + getResourceUrl: () => resourceUrl, + getAudiences: () => [resourceUrl], + getAllowedRedirectUris: async () => [resourceUrl], + scopes: FORM_TRIGGER_SCOPES, + displayName: workflowEntity.name, + authorize: async (user: User) => { + if (requireExecute) { + return ( + await this.workflowFinderService.findWorkflowIdsWithScopeForUser( + [workflowEntity.id], + user, + ['workflow:execute'], + ) + ).has(workflowEntity.id); + } + return true; + }, + }; + } + + return undefined; + } +} diff --git a/packages/cli/src/modules/oauth-server/protected-resource-resolvers/index.ts b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/index.ts new file mode 100644 index 00000000000..f64ac1ad142 --- /dev/null +++ b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/index.ts @@ -0,0 +1,21 @@ +import { Container } from '@n8n/di'; +import { ProtectedResourceRegistry } from '@/services/protected-resource.registry'; +import { FormTriggerTestResourceResolver } from './form-trigger-test-resource.resolver'; +import { FormTriggerResourceResolver } from './form-trigger-resource.resolver'; +import { WorkflowMcpTestTriggerResourceResolver } from './workflow-mcp-test-trigger-resource.resolver'; +import { WorkflowMcpTriggerResourceResolver } from './workflow-mcp-trigger-resource.resolver'; + +export function registerProtectedResourceResolvers() { + Container.get(ProtectedResourceRegistry).registerResolver( + Container.get(WorkflowMcpTriggerResourceResolver), + ); + Container.get(ProtectedResourceRegistry).registerResolver( + Container.get(WorkflowMcpTestTriggerResourceResolver), + ); + Container.get(ProtectedResourceRegistry).registerResolver( + Container.get(FormTriggerResourceResolver), + ); + Container.get(ProtectedResourceRegistry).registerResolver( + Container.get(FormTriggerTestResourceResolver), + ); +} diff --git a/packages/cli/src/modules/oauth-server/protected-resource-resolvers/utils.ts b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/utils.ts index 77a2af7df80..f0911f8a682 100644 --- a/packages/cli/src/modules/oauth-server/protected-resource-resolvers/utils.ts +++ b/packages/cli/src/modules/oauth-server/protected-resource-resolvers/utils.ts @@ -5,6 +5,18 @@ */ export const WORKFLOW_MCP_TRIGGER_SCOPES: string[] = []; +/** Scopes advertised for per-workflow Form trigger resources. Empty, like MCP triggers. */ +export const FORM_TRIGGER_SCOPES: string[] = []; + +/** + * Form-trigger OAuth2 is opt-in. When the flag is off, `n8nUserAuth` form triggers + * keep their existing cookie/HMAC auth and must not be exposed as OAuth protected + * resources, so the resolvers short-circuit. + */ +export function isFormOAuth2Enabled(): boolean { + return process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 === 'true'; +} + export function trimTrailingSlash(path: string): string { if (path.endsWith('/')) { path = path.slice(0, -1); diff --git a/packages/cli/src/services/oauth2-flow-proxy.service.ts b/packages/cli/src/services/oauth2-flow-proxy.service.ts new file mode 100644 index 00000000000..9fc04a7fd40 --- /dev/null +++ b/packages/cli/src/services/oauth2-flow-proxy.service.ts @@ -0,0 +1,26 @@ +import { Service } from '@n8n/di'; +import { UnexpectedError, type N8nOAuth2FlowResult } from 'n8n-workflow'; + +export interface N8nOAuth2Flow { + begin(resourceUrl: string, metadata?: Record): Promise; + complete(code: string, state: string): Promise; +} + +@Service() +export class OAuth2FlowProxy implements N8nOAuth2Flow { + private provider: N8nOAuth2Flow | null = null; + + registerProvider(provider: N8nOAuth2Flow): void { + this.provider = provider; + } + + async begin(resourceUrl: string, metadata?: Record): Promise { + if (!this.provider) throw new UnexpectedError('OAuth2 form flow is not available'); + return await this.provider.begin(resourceUrl, metadata); + } + + async complete(code: string, state: string): Promise { + if (!this.provider) throw new UnexpectedError('OAuth2 form flow is not available'); + return await this.provider.complete(code, state); + } +} diff --git a/packages/cli/src/services/protected-resource.registry.ts b/packages/cli/src/services/protected-resource.registry.ts index ef316559b19..82636084d83 100644 --- a/packages/cli/src/services/protected-resource.registry.ts +++ b/packages/cli/src/services/protected-resource.registry.ts @@ -65,6 +65,8 @@ export interface ProtectedResource { */ getAllowedRedirectUris?(): Promise; + isFirstParty?: boolean; + /** * Determine whether the given user is authorized to access this resource. * Called during the consent flow to gate access to the resource. diff --git a/packages/cli/src/services/trigger-auth-identity-seeder-proxy.service.ts b/packages/cli/src/services/trigger-auth-identity-seeder-proxy.service.ts deleted file mode 100644 index 6f3535b15f7..00000000000 --- a/packages/cli/src/services/trigger-auth-identity-seeder-proxy.service.ts +++ /dev/null @@ -1,24 +0,0 @@ -import { Service } from '@n8n/di'; -import type { IRunExecutionData } from 'n8n-workflow'; - -export interface TriggerAuthIdentitySeeder { - seed(runExecutionData: IRunExecutionData, token: string, resource: string): Promise; -} - -@Service() -export class TriggerAuthIdentitySeederProxy implements TriggerAuthIdentitySeeder { - private seeder: TriggerAuthIdentitySeeder | null = null; - - constructor() {} - - registerSeeder(seeder: TriggerAuthIdentitySeeder): void { - this.seeder = seeder; - } - - async seed(runExecutionData: IRunExecutionData, token: string, resource: string): Promise { - if (this.seeder) { - return await this.seeder.seed(runExecutionData, token, resource); - } - return; - } -} diff --git a/packages/cli/src/webhooks/webhook-helpers.ts b/packages/cli/src/webhooks/webhook-helpers.ts index a615b7c6cb4..3b0331402fa 100644 --- a/packages/cli/src/webhooks/webhook-helpers.ts +++ b/packages/cli/src/webhooks/webhook-helpers.ts @@ -15,6 +15,7 @@ import { BinaryDataService, ErrorReporter, establishExecutionContext, + ExecutionContextService, WAITING_TOKEN_QUERY_PARAM, } from 'n8n-core'; import type { @@ -67,8 +68,8 @@ import { type AuthFailureReason, OAuthTokenVerifierProxy, } from '@/services/oauth-token-verifier-proxy.service'; +import { OAuth2FlowProxy } from '@/services/oauth2-flow-proxy.service'; import { OwnershipService } from '@/services/ownership.service'; -import { TriggerAuthIdentitySeederProxy } from '@/services/trigger-auth-identity-seeder-proxy.service'; import { WorkflowStatisticsService } from '@/services/workflow-statistics.service'; import { WaitTracker } from '@/wait-tracker'; import { WebhookExecutionContext } from '@/webhooks/webhook-execution-context'; @@ -535,6 +536,14 @@ export async function executeWebhook( } }; + additionalData.beginN8nOAuth2Flow = async ( + resourceUrl: string, + metadata?: Record, + ) => await Container.get(OAuth2FlowProxy).begin(resourceUrl, metadata); + + additionalData.completeN8nOAuth2Flow = async (code: string, state: string) => + await Container.get(OAuth2FlowProxy).complete(code, state); + additionalData.validateN8nOAuth2Token = async (token: string, resourceUrl: string) => { const oauthTokenVerifierProxy = Container.get(OAuthTokenVerifierProxy); const result = await oauthTokenVerifierProxy.verifyOAuthAccessToken(token, resourceUrl); @@ -557,12 +566,12 @@ export async function executeWebhook( }; additionalData.establishTriggerIdentity = async (token: string, resource: string) => { - if (runExecutionData === undefined) { - throw new UnexpectedError('Execution data is not available to establish trigger identity'); + additionalData.encryptedRunnerIdentity = await Container.get( + ExecutionContextService, + ).buildTriggerIdentityCredentials(token, resource); + if (runExecutionData) { + await establishExecutionContext(workflow, runExecutionData, additionalData, executionMode); } - await Container.get(TriggerAuthIdentitySeederProxy).seed(runExecutionData, token, resource); - - await establishExecutionContext(workflow, runExecutionData, additionalData, executionMode); }; // Eager pre-execution credential-status gate. Uses the execution context that @@ -572,10 +581,26 @@ export async function executeWebhook( // established, in which case the caller proceeds to execute normally. additionalData.checkTriggerCredentialStatus = async () => { const credentialCheckProxy = additionalData['dynamic-credentials']?.credentialCheckProxy; - const executionContext = runExecutionData?.executionData?.runtimeData; - if (!credentialCheckProxy || !workflow.id || !executionContext?.credentials) { + + if (!credentialCheckProxy || !workflow.id) { return undefined; } + const executionContext = + runExecutionData?.executionData?.runtimeData ?? + (additionalData.encryptedRunnerIdentity + ? { + credentials: additionalData.encryptedRunnerIdentity, + } + : undefined); + + if (!executionContext) { + return undefined; + } + + if (!executionContext.credentials) { + return undefined; + } + return await credentialCheckProxy.checkCredentialStatus(workflow.id, executionContext); }; @@ -739,6 +764,7 @@ export async function executeWebhook( projectId: project?.id, projectName: project?.name, userId: webhookData.userId, + encryptedRunnerIdentity: additionalData.encryptedRunnerIdentity, }; // When resuming from a wait node, copy over the pushRef from the execution-data diff --git a/packages/cli/src/workflow-execute-additional-data.ts b/packages/cli/src/workflow-execute-additional-data.ts index 8c0b52d06c4..fc08eef1209 100644 --- a/packages/cli/src/workflow-execute-additional-data.ts +++ b/packages/cli/src/workflow-execute-additional-data.ts @@ -768,6 +768,8 @@ export async function getBase({ restApiUrl: urlBaseWebhook + globalConfig.endpoints.rest, instanceBaseUrl: `${instanceBaseUrl}/`, formWaitingBaseUrl: urlBaseWebhook + globalConfig.endpoints.formWaiting, + formBaseUrl: urlBaseWebhook + globalConfig.endpoints.form, + formTestBaseUrl: urlBaseTestWebhook + globalConfig.endpoints.formTest, webhookBaseUrl: urlBaseWebhook + globalConfig.endpoints.webhook, webhookWaitingBaseUrl: urlBaseWebhook + globalConfig.endpoints.webhookWaiting, webhookTestBaseUrl: urlBaseTestWebhook + globalConfig.endpoints.webhookTest, diff --git a/packages/cli/test/migration/1785162364001-add-is-first-party-to-oauth-clients.test.ts b/packages/cli/test/migration/1785162364001-add-is-first-party-to-oauth-clients.test.ts new file mode 100644 index 00000000000..37c7a22dee3 --- /dev/null +++ b/packages/cli/test/migration/1785162364001-add-is-first-party-to-oauth-clients.test.ts @@ -0,0 +1,173 @@ +import { + createTestMigrationContext, + initDbUpToMigration, + runSingleMigration, + undoLastSingleMigration, + type TestMigrationContext, +} from '@n8n/backend-test-utils'; +import { DbConnection } from '@n8n/db'; +import { Container } from '@n8n/di'; +import { DataSource } from '@n8n/typeorm'; +import { nanoid } from 'nanoid'; +import { randomUUID } from 'node:crypto'; + +const MIGRATION_NAME = 'AddIsFirstPartyToOAuthClients1785162364001'; + +interface SqliteColumnInfo { + name: string; + notnull: number; +} + +interface PgColumnInfo { + column_name: string; + is_nullable: string; +} + +describe('AddIsFirstPartyToOAuthClients Migration', () => { + let dataSource: DataSource; + + beforeAll(async () => { + const dbConnection = Container.get(DbConnection); + await dbConnection.init(); + dataSource = Container.get(DataSource); + }); + + beforeEach(async () => { + const context = createTestMigrationContext(dataSource); + await context.queryRunner.clearDatabase(); + await context.queryRunner.release(); + await initDbUpToMigration(MIGRATION_NAME); + }); + + afterAll(async () => { + const dbConnection = Container.get(DbConnection); + await dbConnection.close(); + }); + + async function insertUser(context: TestMigrationContext, id: string) { + const table = context.escape.tableName('user'); + await context.runQuery( + `INSERT INTO ${table} ("id", "email", "firstName", "lastName", "password", "roleSlug", "createdAt", "updatedAt") + VALUES (:id, :email, :firstName, :lastName, :password, :roleSlug, :createdAt, :updatedAt)`, + { + id, + email: `${id}@test.com`, + firstName: 'Test', + lastName: 'User', + password: 'hashed', + roleSlug: 'global:member', + createdAt: new Date(), + updatedAt: new Date(), + }, + ); + } + + async function insertClient(context: TestMigrationContext, id: string) { + const table = context.escape.tableName('oauth_clients'); + await context.runQuery( + `INSERT INTO ${table} ("id", "name", "redirectUris", "grantTypes", "createdAt", "updatedAt") + VALUES (:id, :name, :redirectUris, :grantTypes, :createdAt, :updatedAt)`, + { + id, + name: 'Test Client', + redirectUris: JSON.stringify(['https://example.com/callback']), + grantTypes: JSON.stringify(['authorization_code']), + createdAt: new Date(), + updatedAt: new Date(), + }, + ); + } + + async function insertRefreshToken( + context: TestMigrationContext, + token: string, + clientId: string, + userId: string, + ) { + const table = context.escape.tableName('oauth_refresh_tokens'); + await context.runQuery( + `INSERT INTO ${table} ("token", "clientId", "userId", "expiresAt", "createdAt", "updatedAt") + VALUES (:token, :clientId, :userId, :expiresAt, :createdAt, :updatedAt)`, + { + token, + clientId, + userId, + expiresAt: Date.now() + 60_000, + createdAt: new Date(), + updatedAt: new Date(), + }, + ); + } + + async function getColumnMeta( + context: TestMigrationContext, + table: string, + columnName: string, + ): Promise { + if (context.isSqlite) { + const rows = (await context.queryRunner.query( + `PRAGMA table_info(${context.escape.tableName(table)})`, + )) as SqliteColumnInfo[]; + return rows.find((r) => r.name === columnName); + } + const rows = (await context.queryRunner.query( + 'SELECT column_name, is_nullable FROM information_schema.columns WHERE table_name = $1 AND column_name = $2', + [`${context.tablePrefix}${table}`, columnName], + )) as PgColumnInfo[]; + return rows[0]; + } + + describe('up', () => { + it('should add a NOT NULL isFirstParty column to oauth_clients', async () => { + await runSingleMigration(MIGRATION_NAME); + const context = createTestMigrationContext(dataSource); + + const col = await getColumnMeta(context, 'oauth_clients', 'isFirstParty'); + expect(col).toBeDefined(); + if (context.isSqlite) { + expect((col as SqliteColumnInfo).notnull).toBe(1); + } else { + expect((col as PgColumnInfo).is_nullable).toBe('NO'); + } + + await context.queryRunner.release(); + }); + + it('should preserve rows in tables that cascade from oauth_clients', async () => { + // oauth_clients has inbound ON DELETE CASCADE FKs. On SQLite the column + // add recreates the table (drop + rename); without disabling foreign + // keys the drop would cascade and wipe referencing rows. + const context = createTestMigrationContext(dataSource); + const userId = randomUUID(); + const clientId = nanoid(16); + const token = nanoid(); + await insertUser(context, userId); + await insertClient(context, clientId); + await insertRefreshToken(context, token, clientId, userId); + await context.queryRunner.release(); + + await runSingleMigration(MIGRATION_NAME); + + const postContext = createTestMigrationContext(dataSource); + const table = postContext.escape.tableName('oauth_refresh_tokens'); + const rows: Array<{ token: string }> = await postContext.runQuery( + `SELECT ${postContext.escape.columnName('token')} FROM ${table} WHERE ${postContext.escape.columnName('token')} = :token`, + { token }, + ); + expect(rows).toHaveLength(1); + + await postContext.queryRunner.release(); + }); + }); + + describe('down', () => { + it('should remove the isFirstParty column from oauth_clients', async () => { + await runSingleMigration(MIGRATION_NAME); + await undoLastSingleMigration(); + + const context = createTestMigrationContext(dataSource); + expect(await getColumnMeta(context, 'oauth_clients', 'isFirstParty')).toBeUndefined(); + await context.queryRunner.release(); + }); + }); +}); diff --git a/packages/core/src/execution-engine/__tests__/execution-context.service.test.ts b/packages/core/src/execution-engine/__tests__/execution-context.service.test.ts index 7714da2cedb..5f10b02c24d 100644 --- a/packages/core/src/execution-engine/__tests__/execution-context.service.test.ts +++ b/packages/core/src/execution-engine/__tests__/execution-context.service.test.ts @@ -329,6 +329,32 @@ describe('ExecutionContextService', () => { }); }); + describe('buildTriggerIdentityCredentials()', () => { + it('should encrypt the credential context with the token as identity and resource in metadata', async () => { + mockCipher.encryptV2.mockResolvedValue('encrypted-credential-blob'); + + const result = await service.buildTriggerIdentityCredentials( + 'oauth-token-jwt', + 'https://api.example.com/resource', + ); + + expect(mockCipher.encryptV2).toHaveBeenCalledWith({ + version: 1, + identity: 'oauth-token-jwt', + metadata: { source: 'n8n-oauth', resource: 'https://api.example.com/resource' }, + }); + expect(result).toBe('encrypted-credential-blob'); + }); + + it('should propagate errors raised by the cipher', async () => { + mockCipher.encryptV2.mockRejectedValue(new Error('encryption key missing')); + + await expect(service.buildTriggerIdentityCredentials('token', 'resource')).rejects.toThrow( + 'encryption key missing', + ); + }); + }); + describe('encrypt → decrypt round-trip', () => { it('should preserve secureArtifacts through a full round-trip', async () => { // JSON-stringify on encrypt, identity on decrypt — simulates a symmetric cipher diff --git a/packages/core/src/execution-engine/__tests__/execution-context.test.ts b/packages/core/src/execution-engine/__tests__/execution-context.test.ts index aa951c84ccb..3dd716e7970 100644 --- a/packages/core/src/execution-engine/__tests__/execution-context.test.ts +++ b/packages/core/src/execution-engine/__tests__/execution-context.test.ts @@ -1265,29 +1265,24 @@ describe('establishExecutionContext', () => { expect(runExecutionData.executionData!.runtimeData!.credentials).toBeUndefined(); }); - it('should NOT inject credentials for webhook mode even when ciphertext is present', async () => { - const runExecutionData = buildRunDataWithManualTrigger(); - const additionalData = mock({ - encryptedRunnerIdentity: 'encrypted-credential-blob', - }); + // The identity channel is not manual-only: identity-bearing triggers (Form, MCP) + // run in webhook/trigger mode and must resolve the submitter's credentials too. + it.each(['webhook', 'trigger'] as const)( + 'should inject credentials for %s mode when ciphertext is present', + async (mode) => { + const runExecutionData = buildRunDataWithManualTrigger(); + const additionalData = mock({ + encryptedRunnerIdentity: 'encrypted-credential-blob', + }); - await establishExecutionContext(mockWorkflow, runExecutionData, additionalData, 'webhook'); + await establishExecutionContext(mockWorkflow, runExecutionData, additionalData, mode); - expect(mockExecutionContextService.buildManualExecutionCredentials).not.toHaveBeenCalled(); - expect(runExecutionData.executionData!.runtimeData!.credentials).toBeUndefined(); - }); - - it('should NOT inject credentials for trigger mode even when ciphertext is present', async () => { - const runExecutionData = buildRunDataWithManualTrigger(); - const additionalData = mock({ - encryptedRunnerIdentity: 'encrypted-credential-blob', - }); - - await establishExecutionContext(mockWorkflow, runExecutionData, additionalData, 'trigger'); - - expect(mockExecutionContextService.buildManualExecutionCredentials).not.toHaveBeenCalled(); - expect(runExecutionData.executionData!.runtimeData!.credentials).toBeUndefined(); - }); + expect(mockExecutionContextService.buildManualExecutionCredentials).not.toHaveBeenCalled(); + expect(runExecutionData.executionData!.runtimeData!.credentials).toBe( + 'encrypted-credential-blob', + ); + }, + ); it('should not overwrite existing runtimeData when it is already established', async () => { const runExecutionData = buildRunDataWithManualTrigger(); diff --git a/packages/core/src/execution-engine/__tests__/workflow-execute-process-process-run-execution-data.test.ts b/packages/core/src/execution-engine/__tests__/workflow-execute-process-process-run-execution-data.test.ts index 358b8013e84..be4df709ef5 100644 --- a/packages/core/src/execution-engine/__tests__/workflow-execute-process-process-run-execution-data.test.ts +++ b/packages/core/src/execution-engine/__tests__/workflow-execute-process-process-run-execution-data.test.ts @@ -40,6 +40,7 @@ describe('processRunExecutionData', () => { const additionalData = mock({ hooks: { runHook }, restartExecutionId: undefined, + encryptedRunnerIdentity: undefined, webhookWaitingBaseUrl: 'http://localhost:5678/webhook-waiting', formWaitingBaseUrl: 'http://localhost:5678/form-waiting', }); diff --git a/packages/core/src/execution-engine/execution-context.service.ts b/packages/core/src/execution-engine/execution-context.service.ts index 0447a4f6727..4510707cd03 100644 --- a/packages/core/src/execution-engine/execution-context.service.ts +++ b/packages/core/src/execution-engine/execution-context.service.ts @@ -25,12 +25,16 @@ export class ExecutionContextService { private readonly cipher: Cipher, ) {} + async decryptCredentialContext(encrypted: string): Promise { + const decrypted = await this.cipher.decryptV2(encrypted); + return toCredentialContext(decrypted); + } + async decryptExecutionContext(context: IExecutionContext): Promise { const { credentials: encCredentials, secureArtifacts: encSecureArtifacts, ...rest } = context; const result: PlaintextExecutionContext = { ...rest }; if (encCredentials) { - const decrypted = await this.cipher.decryptV2(encCredentials); - result.credentials = toCredentialContext(decrypted); + result.credentials = await this.decryptCredentialContext(encCredentials); } if (encSecureArtifacts) { const decrypted = await this.cipher.decryptV2(encSecureArtifacts); @@ -54,6 +58,15 @@ export class ExecutionContextService { return await this.cipher.encryptV2(payload); } + async buildTriggerIdentityCredentials(token: string, resource: string): Promise { + const payload: ICredentialContext = { + version: 1, + identity: token, + metadata: { source: 'n8n-oauth', resource }, + }; + return await this.cipher.encryptV2(payload); + } + async encryptExecutionContext(context: PlaintextExecutionContext): Promise { const { credentials, secureArtifacts, ...rest } = context; const result: IExecutionContext = { ...rest }; diff --git a/packages/core/src/execution-engine/execution-context.ts b/packages/core/src/execution-engine/execution-context.ts index fecd7e8a538..cc9d5c389c9 100644 --- a/packages/core/src/execution-engine/execution-context.ts +++ b/packages/core/src/execution-engine/execution-context.ts @@ -126,7 +126,7 @@ export const establishExecutionContext = async ( source: mode, }; - if (mode === 'manual' && additionalData?.encryptedRunnerIdentity) { + if (additionalData?.encryptedRunnerIdentity) { executionData.runtimeData.credentials = additionalData.encryptedRunnerIdentity; } diff --git a/packages/core/src/execution-engine/node-execution-context/__tests__/webhook-context.test.ts b/packages/core/src/execution-engine/node-execution-context/__tests__/webhook-context.test.ts index 8a83752fb7c..35bdfbfc34d 100644 --- a/packages/core/src/execution-engine/node-execution-context/__tests__/webhook-context.test.ts +++ b/packages/core/src/execution-engine/node-execution-context/__tests__/webhook-context.test.ts @@ -187,6 +187,56 @@ describe('WebhookContext', () => { }); }); + describe('getNodeWebhookUrl', () => { + const buildUrlContext = (webhookNodeType: 'form' | 'mcp' | undefined, isTest: boolean) => { + const urlNodeType = mock({ + description: { + webhooks: [ + { name: 'default', nodeType: webhookNodeType, path: 'my-path', isFullPath: false }, + ], + }, + }); + nodeTypes.getByNameAndVersion.mockReturnValue(urlNodeType); + expression.getSimpleParameterValue.mockImplementation((_node, value) => value); + + const urlAdditionalData = mock({ + formBaseUrl: 'http://localhost/prod-webhook', + formTestBaseUrl: 'http://localhost/test-webhook', + webhookBaseUrl: 'http://localhost/prod-webhook', + webhookTestBaseUrl: 'http://localhost/test-webhook', + }); + const urlWebhookData = mock({ + webhookDescription: { name: 'default', nodeType: webhookNodeType }, + isTest, + }); + + return new WebhookContext( + workflow, + node, + urlAdditionalData, + mode, + urlWebhookData, + [], + runExecutionData, + ); + }; + + it('should use the test base URL for a form webhook running as a test', () => { + const context = buildUrlContext('form', true); + expect(context.getNodeWebhookUrl('default')).toContain('test-webhook'); + }); + + it('should use the production base URL for a form webhook running in production', () => { + const context = buildUrlContext('form', false); + expect(context.getNodeWebhookUrl('default')).toContain('prod-webhook'); + }); + + it('should ignore isTest for non-form/non-mcp webhooks (production base)', () => { + const context = buildUrlContext(undefined, true); + expect(context.getNodeWebhookUrl('default')).toContain('prod-webhook'); + }); + }); + describe('getNodeParameter', () => { beforeEach(() => { nodeTypes.getByNameAndVersion.mockReturnValue(nodeType); diff --git a/packages/core/src/execution-engine/node-execution-context/utils/__tests__/credential-check-helper-functions.test.ts b/packages/core/src/execution-engine/node-execution-context/utils/__tests__/credential-check-helper-functions.test.ts index eb8a9a8d500..37dad4905b9 100644 --- a/packages/core/src/execution-engine/node-execution-context/utils/__tests__/credential-check-helper-functions.test.ts +++ b/packages/core/src/execution-engine/node-execution-context/utils/__tests__/credential-check-helper-functions.test.ts @@ -29,9 +29,7 @@ describe('getCredentialCheckHelperFunctions', () => { expect(result.checkCredentialStatus).toBeDefined(); const status = await result.checkCredentialStatus!('wf-1', { - version: 1, - establishedAt: Date.now(), - source: 'webhook', + credentials: 'encrypted-credentials', }); expect(status).toEqual(mockResult); diff --git a/packages/core/src/execution-engine/node-execution-context/utils/webhook-helper-functions.ts b/packages/core/src/execution-engine/node-execution-context/utils/webhook-helper-functions.ts index 2bfacebbe0b..d26fa55d490 100644 --- a/packages/core/src/execution-engine/node-execution-context/utils/webhook-helper-functions.ts +++ b/packages/core/src/execution-engine/node-execution-context/utils/webhook-helper-functions.ts @@ -45,6 +45,8 @@ export function getNodeWebhookUrl( let baseUrl: string; if (webhookDescription.nodeType === 'mcp') { baseUrl = isTest === true ? additionalData.mcpTestBaseUrl : additionalData.mcpBaseUrl; + } else if (webhookDescription.nodeType === 'form') { + baseUrl = isTest === true ? additionalData.formTestBaseUrl : additionalData.formBaseUrl; } else { baseUrl = isTest === true ? additionalData.webhookTestBaseUrl : additionalData.webhookBaseUrl; } diff --git a/packages/core/src/execution-engine/node-execution-context/webhook-context.ts b/packages/core/src/execution-engine/node-execution-context/webhook-context.ts index e42873a9eef..c668133e23e 100644 --- a/packages/core/src/execution-engine/node-execution-context/webhook-context.ts +++ b/packages/core/src/execution-engine/node-execution-context/webhook-context.ts @@ -19,6 +19,7 @@ import type { WebhookType, Workflow, WorkflowExecuteMode, + N8nOAuth2FlowResult, } from 'n8n-workflow'; import { UnexpectedError, createEmptyRunExecutionData } from 'n8n-workflow'; @@ -136,11 +137,12 @@ export class WebhookContext extends NodeExecutionContext implements IWebhookFunc } getNodeWebhookUrl(name: WebhookType): string | undefined { - // MCP webhooks are served under dedicated /mcp and /mcp-test endpoints; the OAuth - // resource URL must match the endpoint the request actually arrived on. Other webhook - // types keep their existing behaviour (production base) here. - const isTest = - this.webhookData.webhookDescription.nodeType === 'mcp' ? this.webhookData.isTest : undefined; + // MCP and form webhooks are served under dedicated /mcp(+/mcp-test) and + // /form(+/form-test) endpoints; the OAuth resource URL must match the endpoint the + // request actually arrived on. Other webhook types keep their existing behaviour + // (production base) here. + const { nodeType } = this.webhookData.webhookDescription; + const isTest = nodeType === 'mcp' || nodeType === 'form' ? this.webhookData.isTest : undefined; return getNodeWebhookUrl( name, @@ -173,6 +175,23 @@ export class WebhookContext extends NodeExecutionContext implements IWebhookFunc return await this.additionalData.validateCookieAuth(cookieValue); } + async beginN8nOAuth2Flow( + resourceUrl: string, + metadata?: Record, + ): Promise { + if (!this.additionalData.beginN8nOAuth2Flow) { + throw new UnexpectedError('OAuth2 flow is not available'); + } + return await this.additionalData.beginN8nOAuth2Flow(resourceUrl, metadata); + } + + async completeN8nOAuth2Flow(code: string, state: string): Promise { + if (!this.additionalData.completeN8nOAuth2Flow) { + throw new UnexpectedError('OAuth2 flow is not available'); + } + return await this.additionalData.completeN8nOAuth2Flow(code, state); + } + async validateN8nOAuth2Token( token: string, resourceUrl: string, diff --git a/packages/nodes-base/nodes/Form/test/utils.test.ts b/packages/nodes-base/nodes/Form/test/utils.test.ts index 3d651c72415..116b1bc540e 100644 --- a/packages/nodes-base/nodes/Form/test/utils.test.ts +++ b/packages/nodes-base/nodes/Form/test/utils.test.ts @@ -1026,6 +1026,296 @@ describe('FormTrigger, formWebhook', () => { expect(json.user).toBeUndefined(); }); }); + + describe('n8nUserAuth with OAuth2 flow (N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2)', () => { + const authedUser = { + id: 'user-1', + email: 'user@example.com', + firstName: 'Test', + lastName: 'User', + }; + const formFields: FormFieldsParameter = [ + { fieldLabel: 'Name', fieldType: 'text', requiredField: true }, + ]; + const resourceUrl = 'http://localhost:5678/form/test'; + + const setupContext = ( + ctx: ReturnType>, + overrides: { + method: 'GET' | 'POST'; + query?: IDataObject; + headers?: Record; + originalUrl?: string; + } = { method: 'GET' }, + ) => { + const send = vi.fn(); + const status = vi.fn(() => ({ send })) as any; + const writeHead = vi.fn(); + const end = vi.fn(); + const setHeader = vi.fn(); + const render = vi.fn(); + const cookie = vi.fn(); + const clearCookie = vi.fn(); + const request = { + method: overrides.method, + originalUrl: overrides.originalUrl ?? '/form/test', + query: overrides.query ?? {}, + headers: { host: 'localhost:5678', ...(overrides.headers ?? {}) }, + protocol: 'http', + contentType: overrides.method === 'POST' ? 'multipart/form-data' : undefined, + }; + + ctx.getNode.mockReturnValue({ typeVersion: 2.6 } as INode); + ctx.getNodeParameter.calledWith('options').mockReturnValue({}); + ctx.getNodeParameter.calledWith('formTitle').mockReturnValue('Test Form'); + ctx.getNodeParameter.calledWith('formDescription').mockReturnValue('Test Description'); + ctx.getNodeParameter.calledWith('responseMode').mockReturnValue('onReceived'); + ctx.getNodeParameter.calledWith('authentication', 'none').mockReturnValue('n8nUserAuth'); + ctx.getNodeParameter.calledWith('formFields.values').mockReturnValue(formFields); + ctx.getNodeWebhookUrl.mockReturnValue(resourceUrl); + ctx.getRequestObject.mockReturnValue(request as any); + ctx.getHeaderData.mockReturnValue(request.headers); + ctx.getResponseObject.mockReturnValue({ + status, + writeHead, + end, + setHeader, + render, + cookie, + clearCookie, + } as any); + ctx.getMode.mockReturnValue('manual'); + ctx.getInstanceId.mockReturnValue('instanceId'); + ctx.getBodyData.mockReturnValue({ data: { 'field-0': 'John' }, files: {} }); + ctx.getWorkflowSettings.mockReturnValue(mock({})); + ctx.getChildNodes.mockReturnValue([]); + (ctx as any).logger = { warn: vi.fn(), error: vi.fn(), debug: vi.fn(), info: vi.fn() }; + + return { status, send, writeHead, end, setHeader, render, cookie, clearCookie }; + }; + + beforeEach(() => { + vi.clearAllMocks(); + vi.stubEnv('N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2', 'true'); + }); + + afterEach(() => { + vi.unstubAllEnvs(); + }); + + it('redirects to the authorization URL on GET without a code', async () => { + const ctx = mock(); + const { writeHead, end } = setupContext(ctx, { method: 'GET' }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=abc'); + + const result = await formWebhook(ctx); + + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, undefined); + expect(writeHead).toHaveBeenCalledWith(302, { + Location: 'http://localhost:5678/oauth/authorize?state=abc', + }); + expect(end).toHaveBeenCalled(); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('responds 403 without restarting the flow when consent is denied', async () => { + const ctx = mock(); + const { status, send } = setupContext(ctx, { + method: 'GET', + query: { error: 'access_denied', error_description: 'User denied', state: 'the-state' }, + }); + + const result = await formWebhook(ctx); + + expect(status).toHaveBeenCalledWith(403); + expect(send).toHaveBeenCalled(); + expect(ctx.beginN8nOAuth2Flow).not.toHaveBeenCalled(); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('restarts the flow when the callback fails validation', async () => { + const ctx = mock(); + const { writeHead } = setupContext(ctx, { + method: 'GET', + query: { code: 'the-code', state: 'the-state' }, + }); + ctx.completeN8nOAuth2Flow.mockResolvedValue({ valid: false, reason: 'invalid_state' }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=fresh'); + + const result = await formWebhook(ctx); + + expect(ctx.completeN8nOAuth2Flow).toHaveBeenCalledWith('the-code', 'the-state'); + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, undefined); + expect(writeHead).toHaveBeenCalledWith(302, { + Location: 'http://localhost:5678/oauth/authorize?state=fresh', + }); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('redirects to a clean URL with the token in a cookie on a valid callback', async () => { + const ctx = mock(); + const { render, writeHead, cookie } = setupContext(ctx, { + method: 'GET', + query: { code: 'the-code', state: 'the-state' }, + }); + ctx.completeN8nOAuth2Flow.mockResolvedValue({ + valid: true, + token: 'as-token', + user: authedUser, + }); + + const result = await formWebhook(ctx); + + expect(ctx.completeN8nOAuth2Flow).toHaveBeenCalledWith('the-code', 'the-state'); + expect(ctx.beginN8nOAuth2Flow).not.toHaveBeenCalled(); + // The code/state must not reach the sandboxed form page: redirect to the + // clean resource URL instead of rendering here. + expect(render).not.toHaveBeenCalled(); + expect(writeHead).toHaveBeenCalledWith(302, { Location: resourceUrl }); + expect(cookie).toHaveBeenCalledWith( + 'n8n-form-oauth', + 'as-token', + expect.objectContaining({ httpOnly: true, sameSite: 'lax' }), + ); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('stashes the original query params as flow metadata on a fresh GET before redirecting', async () => { + const ctx = mock(); + const { writeHead } = setupContext(ctx, { + method: 'GET', + originalUrl: '/form/test?foo=bar', + }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=abc'); + + const result = await formWebhook(ctx); + + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, { query: 'foo=bar' }); + expect(writeHead).toHaveBeenCalledWith(302, { + Location: 'http://localhost:5678/oauth/authorize?state=abc', + }); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('preserves a lone code query param as flow metadata on a fresh GET', async () => { + // A form field literally named `code` (or `state`) is not a provider callback + // (which needs both), so it is a genuine fresh GET and must be preserved. + const ctx = mock(); + setupContext(ctx, { + method: 'GET', + query: { foo: 'bar', code: 'x' }, + originalUrl: '/form/test?foo=bar&code=x', + }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=abc'); + + await formWebhook(ctx); + + expect(ctx.completeN8nOAuth2Flow).not.toHaveBeenCalled(); + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, { query: 'foo=bar&code=x' }); + }); + + it('re-appends the query stashed as flow metadata on a valid callback', async () => { + const ctx = mock(); + const { writeHead } = setupContext(ctx, { + method: 'GET', + query: { code: 'the-code', state: 'the-state' }, + }); + ctx.completeN8nOAuth2Flow.mockResolvedValue({ + valid: true, + token: 'as-token', + user: authedUser, + metadata: { query: 'foo=bar' }, + }); + + const result = await formWebhook(ctx); + + expect(writeHead).toHaveBeenCalledWith(302, { Location: `${resourceUrl}?foo=bar` }); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('does not stash code/state as flow metadata on a callback fall-through', async () => { + const ctx = mock(); + setupContext(ctx, { + method: 'GET', + query: { code: 'the-code', state: 'the-state' }, + originalUrl: '/form/test?code=the-code&state=the-state', + }); + ctx.completeN8nOAuth2Flow.mockResolvedValue({ valid: false, reason: 'invalid_state' }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=fresh'); + + await formWebhook(ctx); + + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, undefined); + }); + + it('renders the form on the clean GET carrying the oauth cookie', async () => { + const ctx = mock(); + const { render, clearCookie } = setupContext(ctx, { + method: 'GET', + headers: { cookie: 'n8n-form-oauth=as-token' }, + }); + ctx.validateN8nOAuth2Token.mockResolvedValue({ valid: true, user: authedUser }); + + const result = await formWebhook(ctx); + + expect(ctx.validateN8nOAuth2Token).toHaveBeenCalledWith('as-token', resourceUrl); + expect(ctx.beginN8nOAuth2Flow).not.toHaveBeenCalled(); + expect(clearCookie).toHaveBeenCalledWith('n8n-form-oauth', expect.any(Object)); + expect(render).toHaveBeenCalledWith( + 'form-trigger', + expect.objectContaining({ authToken: 'as-token' }), + ); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('restarts the flow when the cookie token is invalid', async () => { + const ctx = mock(); + const { writeHead, render } = setupContext(ctx, { + method: 'GET', + headers: { cookie: 'n8n-form-oauth=stale-token' }, + }); + ctx.validateN8nOAuth2Token.mockResolvedValue({ valid: false, reason: 'invalid_token' }); + ctx.beginN8nOAuth2Flow.mockResolvedValue('http://localhost:5678/oauth/authorize?state=fresh'); + + const result = await formWebhook(ctx); + + expect(ctx.validateN8nOAuth2Token).toHaveBeenCalledWith('stale-token', resourceUrl); + expect(ctx.beginN8nOAuth2Flow).toHaveBeenCalledWith(resourceUrl, undefined); + expect(writeHead).toHaveBeenCalledWith(302, { + Location: 'http://localhost:5678/oauth/authorize?state=fresh', + }); + expect(render).not.toHaveBeenCalled(); + expect(result).toEqual({ noWebhookResponse: true }); + }); + + it('establishes the submitter identity on POST with a valid token', async () => { + const ctx = mock(); + setupContext(ctx, { method: 'POST', headers: { 'x-auth-token': 'as-token' } }); + ctx.validateN8nOAuth2Token.mockResolvedValue({ valid: true, user: authedUser }); + + const result = await formWebhook(ctx); + + expect(ctx.validateN8nOAuth2Token).toHaveBeenCalledWith('as-token', resourceUrl); + expect(ctx.establishTriggerIdentity).toHaveBeenCalledWith('as-token', resourceUrl); + expect(result).toMatchObject({ webhookResponse: { status: 200 } }); + }); + + it('returns 401 on POST with an invalid token', async () => { + const ctx = mock(); + const { status, send } = setupContext(ctx, { + method: 'POST', + headers: { 'x-auth-token': 'bad-token' }, + }); + ctx.validateN8nOAuth2Token.mockResolvedValue({ valid: false, reason: 'invalid_token' }); + + const result = await formWebhook(ctx); + + expect(ctx.establishTriggerIdentity).not.toHaveBeenCalled(); + expect(status).toHaveBeenCalledWith(401); + expect(send).toHaveBeenCalled(); + expect(result).toEqual({ noWebhookResponse: true }); + }); + }); }); describe('FormTrigger, prepareFormData', () => { diff --git a/packages/nodes-base/nodes/Form/utils/utils.ts b/packages/nodes-base/nodes/Form/utils/utils.ts index 8c24ed9055c..da53b6819e0 100644 --- a/packages/nodes-base/nodes/Form/utils/utils.ts +++ b/packages/nodes-base/nodes/Form/utils/utils.ts @@ -25,6 +25,7 @@ import { tryToParseUrl, BINARY_MODE_COMBINED, tryToParseJsonToFormFields, + UnexpectedError, } from 'n8n-workflow'; import * as a from 'node:assert'; import sanitize from 'sanitize-html'; @@ -70,6 +71,10 @@ function isFormUserAuthClaims(value: unknown): value is FormUserAuthClaims { ); } +export function isFormOAuth2Enabled(): boolean { + return process.env.N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2 === 'true'; +} + export function sanitizeHtml(text: string) { return sanitize(text, { allowedTags: [ @@ -705,6 +710,46 @@ export function verifyFormUserAuthToken(token: string, node: INode): IUser | nul }; } +function trimTrailingSlash(url: string): string { + return url.endsWith('/') ? url.slice(0, -1) : url; +} + +// Carries the OAuth2 access token across the single same-site redirect from the +// provider callback to the clean form URL, so `code`/`state` never reach the +// sandboxed form page. The token is otherwise already embedded in the form HTML +// (the page sends it back as `x-auth-token` on POST), so this is not a new exposure. +const FORM_OAUTH_COOKIE_NAME = 'n8n-form-oauth'; + +function formOAuthCookieOptions(req: Request, resourceUrl: string) { + // Derive `secure` from the request scheme (honouring x-forwarded-proto, as + // buildAbsoluteFormUrl does) rather than config, so the cookie is actually sent + // back on the follow-up GET over http in dev while staying Secure over https. + const forwardedProto = req.headers['x-forwarded-proto']; + const proto = (typeof forwardedProto === 'string' ? forwardedProto.trim() : '') || req.protocol; + return { + httpOnly: true, + sameSite: 'lax' as const, // must be Lax: sent on our own top-level 302 → GET + secure: proto === 'https', + path: new URL(resourceUrl).pathname, // scope the bearer token to this form + }; +} + +function setFormOAuthToken(res: Response, req: Request, resourceUrl: string, token: string): void { + res.cookie(FORM_OAUTH_COOKIE_NAME, token, { + ...formOAuthCookieOptions(req, resourceUrl), + maxAge: 60_000, // one redirect hop; short by design + }); +} + +function readFormOAuthToken(req: Request): string | null { + const match = (req.headers.cookie ?? '').match(/(?:^|;\s*)n8n-form-oauth=([^;]+)/); + return match ? decodeURIComponent(match[1].trim()) : null; +} + +function clearFormOAuthToken(res: Response, req: Request, resourceUrl: string): void { + res.clearCookie(FORM_OAUTH_COOKIE_NAME, formOAuthCookieOptions(req, resourceUrl)); +} + /** * Authenticate an `n8nUserAuth` request via: * 1. the `n8n-auth` cookie (sent on top-level GET when the user is logged in), or @@ -715,22 +760,139 @@ export function verifyFormUserAuthToken(token: string, node: INode): IUser | nul * (302 to `/signin` on GET, 401 on POST) and returns `null` — the caller * must abort with `noWebhookResponse`. */ -async function authenticateFormUserOrRespond(context: IWebhookFunctions): Promise { +async function authenticateFormUserOrRespond( + context: IWebhookFunctions, + oauth2Enabled: boolean = false, +): Promise<{ user: IUser; token: string | null } | null> { const req = context.getRequestObject(); + if (oauth2Enabled) { + const res = context.getResponseObject(); + const url = context.getNodeWebhookUrl('default'); + if (!url) { + throw new UnexpectedError('Webhook URL not found for the node'); + } + const resourceUrl = trimTrailingSlash(url); + if (req.method === 'GET') { + const { code, state } = req.query; + + if (typeof req.query.error === 'string') { + // The provider returned an error (e.g. the user denied consent). Restarting the + // flow here would loop straight back to the same denial, so stop and report. + context.logger.warn('Form OAuth2 authorization was denied or failed', { + error: req.query.error, + error_description: req.query.error_description, + }); + res.status(403).send('Access denied'); + res.end(); + return null; + } + + // handle OAuth2 callback from the provider (code + state query params) + if (typeof code === 'string' && typeof state === 'string') { + try { + const authorizationResult = await context.completeN8nOAuth2Flow(code, state); + if (authorizationResult.valid) { + // TODO: exchange the OAuth2 token to a specficly scoped token. + // Don't render the form here: the callback URL still carries `code`/`state`, + // which must never reach the sandboxed form page. Stash the token in a + // one-hop cookie and redirect to the clean resource URL — the follow-up GET + // (below) picks up the cookie and renders the form. + setFormOAuthToken(res, req, resourceUrl, authorizationResult.token); + // Re-append the original query params (dropped by the first provider + // redirect, stashed against this flow's `state` on the fresh GET) so the + // follow-up GET restores field prefill and `formQueryParameters`. + const preservedQuery = authorizationResult.metadata?.query; + res.writeHead(302, { + Location: preservedQuery ? `${resourceUrl}?${preservedQuery}` : resourceUrl, + }); + res.end(); + return null; + } + // We fall through to the redirect below to restart the OAuth2 flow if the callback is invalid. + context.logger.warn('Form OAuth2 flow failed, restarting', { + reason: authorizationResult.reason, + }); + } catch (error) { + // Ignore errors and fall through to the redirect below + context.logger.warn('Form OAuth2 flow failed, restarting', { + error, + }); + } + } else { + // Not a provider callback. If we just completed the flow, the token rides in a + // one-hop cookie set on the redirect above. Consume it once and render. + const cookieToken = readFormOAuthToken(req); + if (cookieToken) { + clearFormOAuthToken(res, req, resourceUrl); + const validation = await context.validateN8nOAuth2Token(cookieToken, resourceUrl); + if (validation.valid) { + return { user: validation.user, token: cookieToken }; + } + // Stale/invalid cookie — fall through to restart the OAuth2 flow. + } + } + + // Stash the original query params against this flow's OAuth `state` so we can + // restore them after the bounce. Only on a genuine fresh GET — a callback + // fall-through (invalid completion) carries `code`/`state`, which we must not + // preserve as the form's query, and whose original query was consumed with the + // now-invalid state, so a restart begins clean. + const isCallback = typeof code === 'string' && typeof state === 'string'; + const originalQuery = isCallback ? '' : req.originalUrl.split('?').slice(1).join('?'); + + try { + // start authentication flow by redirecting to the OAuth2 provider's authorization URL + const authorizationUrl = await context.beginN8nOAuth2Flow( + resourceUrl, + originalQuery ? { query: originalQuery } : undefined, + ); + res.writeHead(302, { + Location: authorizationUrl, + }); + res.end(); + } catch (error) { + // Can't build the authorization URL — nothing to redirect to, so abort. + context.logger.warn('Form OAuth2 flow failed', { + error, + }); + throw new UnexpectedError('Form OAuth2 flow failed'); + } + return null; + } else { + // For POST requests, the OAuth2 flow is not applicable. We fall through to the cookie and token checks below. + const formToken = req.headers['x-auth-token']; + if (typeof formToken === 'string' && formToken) { + const validation = await context.validateN8nOAuth2Token(formToken, resourceUrl); + if (validation.valid) { + await context.establishTriggerIdentity(formToken, resourceUrl); // seeds the run + return { + user: validation.user, + token: null, + }; + } + } + res.status(401).send(); + return null; + } + } + // Parse the raw Cookie header rather than `req.cookies` because the webhook // path may bypass cookie-parser middleware in some deployments. const cookieMatch = (req.headers.cookie ?? '').match(/(?:^|;\s*)n8n-auth=([^;]+)/); if (cookieMatch) { try { - return await context.validateCookieAuth(cookieMatch[1].trim()); + return { + user: await context.validateCookieAuth(cookieMatch[1].trim()), + token: null, + }; } catch {} } const formToken = req.headers['x-auth-token']; if (typeof formToken === 'string' && formToken) { const user = verifyFormUserAuthToken(formToken, context.getNode()); - if (user) return user; + if (user) return { user, token: null }; } const res = context.getResponseObject(); @@ -758,7 +920,7 @@ export async function validateFormPageAuth( triggerAuthentication: string, ): Promise<{ authedUser?: IUser; responded?: boolean }> { if (triggerAuthentication !== 'n8nUserAuth') return {}; - const user = await authenticateFormUserOrRespond(context); + const { user } = (await authenticateFormUserOrRespond(context, false)) ?? {}; return user ? { authedUser: user } : { responded: true }; } @@ -802,11 +964,14 @@ export async function formWebhook( const authentication = context.getNodeParameter(authProperty, 'none') as string; let authedUser: IUser | undefined; + let oAuth2Token: string | undefined; if (node.typeVersion > 1) { if (authentication === 'n8nUserAuth') { - const user = await authenticateFormUserOrRespond(context); + const { user, token } = + (await authenticateFormUserOrRespond(context, isFormOAuth2Enabled())) ?? {}; if (!user) return { noWebhookResponse: true }; authedUser = user; + oAuth2Token = token ?? undefined; } else { try { await validateWebhookAuthentication(context, authProperty); @@ -885,10 +1050,14 @@ export async function formWebhook( let authToken: string | undefined; if (node.typeVersion > 1) { if (authentication === 'n8nUserAuth' && authedUser) { - // Cookies aren't sent on POST from the sandboxed form page - // (null origin + SameSite=Lax). Embed an HMAC token so the - // POST handler can re-authenticate the user. - authToken = generateFormUserAuthToken(node, authedUser); + if (!isFormOAuth2Enabled()) { + // Cookies aren't sent on POST from the sandboxed form page + // (null origin + SameSite=Lax). Embed an HMAC token so the + // POST handler can re-authenticate the user. + authToken = generateFormUserAuthToken(node, authedUser); + } else { + authToken = oAuth2Token; + } } else { authToken = await generateFormPostBasicAuthToken(context, authProperty); } diff --git a/packages/nodes-base/nodes/Form/v2/FormTriggerV2.node.ts b/packages/nodes-base/nodes/Form/v2/FormTriggerV2.node.ts index 1295db28396..34017b9cd8f 100644 --- a/packages/nodes-base/nodes/Form/v2/FormTriggerV2.node.ts +++ b/packages/nodes-base/nodes/Form/v2/FormTriggerV2.node.ts @@ -143,6 +143,18 @@ const descriptionV2: INodeTypeDescription = { "Default to 'none'. n8n exposes inbound trigger URLs publicly by design. Only select an authentication method when the user explicitly asks to authenticate inbound traffic.", }, }, + { + displayName: 'Require Workflow Execute Permission', + name: 'requireExecuteAccess', + type: 'boolean', + default: true, + envFeatureFlag: 'FORM_TRIGGER_OAUTH2', + displayOptions: { + show: { authentication: ['n8nUserAuth'], '@version': [{ _cnd: { gte: 2.6 } }] }, + }, + description: + 'Whether the triggering user must also have permission to execute the workflow in the project it belongs to', + }, { ...webhookPath, displayOptions: { show: { '@version': [{ _cnd: { lte: 2.1 } }] } } }, formTitle, formDescription, diff --git a/packages/workflow/src/interfaces.ts b/packages/workflow/src/interfaces.ts index c9349adca53..19228eeb144 100644 --- a/packages/workflow/src/interfaces.ts +++ b/packages/workflow/src/interfaces.ts @@ -154,6 +154,10 @@ export type N8nOAuth2ValidationResult = | { valid: true; user: IUser } | { valid: false; reason: OAuth2FailureReason }; +export type N8nOAuth2FlowResult = + | { valid: true; token: string; user: IUser; metadata?: Record } + | { valid: false; reason: string }; + export type ProjectSharingData = { id: string; name: string | null; @@ -1166,7 +1170,9 @@ export type CredentialCheckResult = { export type DynamicCredentialCheckProxyProvider = { checkCredentialStatus( workflowId: string, - executionContext: IExecutionContext, + executionContext: { + credentials?: string; + }, ): Promise; }; @@ -1174,7 +1180,9 @@ export type CredentialCheckProxyFunctions = { // Optional to account for situations where the dynamic-credentials module is disabled checkCredentialStatus?( workflowId: string, - executionContext: IExecutionContext, + executionContext: { + credentials?: string; + }, ): Promise; }; @@ -1452,7 +1460,58 @@ export interface IHookFunctions export interface IWebhookFunctions extends FunctionsBaseWithRequiredKeys<'getMode'> { getBodyData(): IDataObject; getHeaderData(): IncomingHttpHeaders; + /** + * Identity pipeline for identity-bearing triggers (Form, MCP). + * + * These let a trigger run a workflow as the *caller* — resolving that user's own + * private (per-user) credentials instead of a shared/static credential — by + * proving the caller's n8n identity against the internal Authorization Server (AS) + * and binding it to the execution. The shape is acquire → verify → bind: + * + * - Browser-facing triggers (Form) acquire a token interactively: + * `beginN8nOAuth2Flow` → (AS redirect) → `completeN8nOAuth2Flow`. + * - Resource-server triggers (MCP) receive a bearer token directly and only + * `validateN8nOAuth2Token` it. + * - Either way, `establishTriggerIdentity` binds the verified token to the run. + */ + + /** + * Starts the interactive authorization-code + PKCE flow against n8n's internal AS + * for `resourceUrl` (the trigger's own protected-resource URL). Returns the + * `/oauth/authorize` URL to redirect the browser to; the AS identifies the + * already-logged-in user from their n8n session (no login prompt for the + * first-party trigger client) and redirects back to the trigger URL with a code. + * Used on the initial GET of a browser-facing trigger. Pair with + * `completeN8nOAuth2Flow`. + * + * Optional `metadata` is stashed server-side against this flow's one-time `state` + * (never sent to the browser) and returned by `completeN8nOAuth2Flow` on success — + * a per-flow slot for carrying data (e.g. the original request query) across the bounce. + */ + beginN8nOAuth2Flow(resourceUrl: string, metadata?: Record): Promise; + /** + * Completes the flow started by `beginN8nOAuth2Flow` once the AS redirects back to + * the trigger URL with `?code&state`. Consumes the one-time `state`, verifies PKCE, + * exchanges the code for an access token **server-side** (the code never reaches the + * sandboxed form page), and validates it. Returns the token and resolved user on + * success, or a failure reason. + */ + completeN8nOAuth2Flow(code: string, state: string): Promise; + /** + * Verifies an AS access token against `resourceUrl` (the expected audience) without + * running a redirect flow. Used by resource-server triggers (MCP) that receive a + * bearer token directly, and by browser triggers on the POST leg to re-check the + * token the page presents. Returns validity + resolved user, or a failure reason + * (`invalid_token`, `insufficient_scope`, …). + */ validateN8nOAuth2Token(token: string, resourceUrl: string): Promise; + /** + * Binds the verified submitter to the current execution: builds an encrypted + * credential context from `token`/`resource` and threads it into the execution's + * runtime context, so downstream nodes resolve *that user's* per-user credentials. + * Call only after the token is validated. The identity persists for the whole + * execution (including across a Wait), within the token's validity window. + */ establishTriggerIdentity(token: string, resource: string): Promise; /** * Checks the status of the triggering identity's resolvable (private) credentials @@ -3452,7 +3511,7 @@ export interface IWorkflowExecutionDataProcess { deduplicationKey?: string; /** W3C trace context extracted from inbound webhook headers. */ tracingContext?: { traceparent: string; tracestate?: string }; - /** Encrypted credential context for a manual editor-triggered execution. */ + /** Encrypted credential context for a triggered execution. */ encryptedRunnerIdentity?: string; /** Parent evaluation TestRun.id, exposed to expressions as `$evaluation.runId`. */ evaluationRunId?: string; @@ -3543,6 +3602,14 @@ export interface IWorkflowExecuteAdditionalData { runExecutionData: IRunExecutionData, alias: string, ): Promise; + /** + * Backing implementations for the trigger identity pipeline exposed to nodes on + * `IWebhookFunctions` (see the docs there). Optional here because they are only + * wired in by the CLI webhook layer when the OAuth server is available; the + * context methods throw if a node reaches them while unset. + */ + beginN8nOAuth2Flow?: (resourceUrl: string, metadata?: Record) => Promise; + completeN8nOAuth2Flow?: (code: string, state: string) => Promise; validateN8nOAuth2Token?: ( token: string, resourceUrl: string, @@ -3557,7 +3624,9 @@ export interface IWorkflowExecuteAdditionalData { instanceBaseUrl: string; setExecutionStatus?: (status: ExecutionStatus) => void; sendDataToUI?: (type: string, data: IDataObject | IDataObject[]) => void; + formBaseUrl: string; formWaitingBaseUrl: string; + formTestBaseUrl: string; webhookBaseUrl: string; webhookWaitingBaseUrl: string; webhookTestBaseUrl: string; diff --git a/packages/workflow/src/trigger-identity.ts b/packages/workflow/src/trigger-identity.ts index 445d68484af..fd734e52447 100644 --- a/packages/workflow/src/trigger-identity.ts +++ b/packages/workflow/src/trigger-identity.ts @@ -1,6 +1,7 @@ import { CHAT_TRIGGER_NODE_TYPE, EXECUTE_WORKFLOW_TRIGGER_NODE_TYPE, + FORM_TRIGGER_NODE_TYPE, MANUAL_TRIGGER_NODE_TYPES, MCP_TRIGGER_NODE_TYPE, } from './constants'; @@ -55,7 +56,13 @@ export function classifyTriggerIdentity( nodeType === CHAT_TRIGGER_NODE_TYPE && parameters?.availableInChat === true; const isMcpTrigger = nodeType === MCP_TRIGGER_NODE_TYPE && parameters?.authentication === 'n8nOAuth2'; - if (isSubWorkflowTrigger || isChatHubTrigger || isMcpTrigger) { + // Not gated by N8N_ENV_FEAT_FORM_TRIGGER_OAUTH2: this shared classification describes the + // trigger's capability when the OAuth feature is enabled, and reading the env flag here + // would mean threading it through both callers of a low-level FE/BE-shared function. The + // only inconsistency is the narrow dynamic-credentials + form + flag-off combination. + const isFormTrigger = + nodeType === FORM_TRIGGER_NODE_TYPE && parameters?.authentication === 'n8nUserAuth'; + if (isSubWorkflowTrigger || isChatHubTrigger || isMcpTrigger || isFormTrigger) { return { providesN8nIdentity: true, providesExternalIdentity: true }; }