From 3a646ffcb083ab3d069cf7a7bdd6cbca846d0368 Mon Sep 17 00:00:00 2001 From: Abdul Rahman <81605929+abdulrahmancodes@users.noreply.github.com> Date: Tue, 4 Aug 2026 08:17:23 +0530 Subject: [PATCH] feat(connections): add an onDisconnect lifecycle hook to connection providers (#23538) Platform half of the follow-up to https://github.com/twentyhq/twenty/pull/22984#discussion_r3673946334. The Slack app claims a `team_id` on connect and had no way to release it, because connection providers only had an on-connect hook. Nothing here is Slack-specific, so it targets `main`. The app side is #23540, on top of `feat/slack-bot`, and waits on this plus an SDK release. ## What changes `defineConnectionProvider` accepts `onDisconnectLogicFunction` alongside `onConnectLogicFunction`. It is stored on `connectionProvider.onDisconnectLogicFunctionUniversalIdentifier` (fast instance command `2.26.0_...1785350000000`) and enqueued right after the `ConnectedAccount` row is deleted, in the disconnecting workspace, with the same payload as on-connect: ```ts type OnDisconnectPayload = { connectionProviderId: string; connectionProviderName: string; connectedAccountId: string; }; ``` The `ConnectedAccount` is gone by the time the hook runs, so `getConnection` no longer resolves. Anything the cleanup needs has to be in the key-value store, written at connect time and keyed by `connectedAccountId`. The docs section spells that out, along with the fact that uninstalling an app drops its connections through a cascade that never reaches this hook, where `uninstallLogicFunction` is the right tool instead. Both dispatches moved into a new `ConnectionProviderLifecycleHookService`, so `ConnectionProviderOAuthFlowService` no longer owns hook plumbing and `ConnectedAccountMetadataService.delete` can reuse it. On-connect behaviour is unchanged: best effort, never blocks the caller, failures go to Sentry. ## Tests - `connection-provider-lifecycle-hook.service.spec.ts`: the on-connect cases moved over, plus on-disconnect dispatch, no-hook, and missing-provider cases - `connection-provider-oauth-flow.service.spec.ts`: now asserts delegation to the lifecycle hook service - SDK validation, manifest duplicate-identifier, and manifest to flat converter specs extended Server unit tests and typecheck for shared, sdk and server pass locally. --- .../extend/apps/logic/connections.mdx | 39 ++++ .../__tests__/manifest-validate.spec.ts | 63 +++-- .../build/manifest/manifest-validate.ts | 1 + .../define-connection-provider.spec.ts | 32 ++- .../define-connection-provider.ts | 26 ++- ...t-logic-function-to-connection-provider.ts | 21 ++ .../instance-commands.constant.ts | 2 + ...rsal-flat-connection-provider.util.spec.ts | 21 ++ ...universal-flat-connection-provider.util.ts | 3 + ...on-provider-lifecycle-hook.service.spec.ts | 219 ++++++++++++++++++ ...ection-provider-oauth-flow.service.spec.ts | 134 ++--------- ...connection-provider-exception-code.enum.ts | 1 + ...nection-provider-lifecycle-hook.service.ts | 157 +++++++++++++ .../connection-provider-oauth-flow.service.ts | 76 +----- .../connection-provider.entity.ts | 7 + .../connection-provider.exception.ts | 2 + .../connection-provider.module.ts | 8 +- .../connected-account-metadata.module.ts | 2 + .../connected-account-metadata.service.ts | 12 + ...ompare-and-stringify.constant.spec.ts.snap | 1 + ...configuration-by-metadata-name.constant.ts | 5 + ...te-connection-provider.integration-spec.ts | 1 + ...onnection-providers-by-application.util.ts | 4 +- .../connectionProviderManifestType.ts | 1 + 24 files changed, 594 insertions(+), 244 deletions(-) create mode 100644 packages/twenty-server/src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785810340935-add-on-disconnect-logic-function-to-connection-provider.ts create mode 100644 packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-lifecycle-hook.service.spec.ts create mode 100644 packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service.ts diff --git a/packages/twenty-docs/developers/extend/apps/logic/connections.mdx b/packages/twenty-docs/developers/extend/apps/logic/connections.mdx index 2d7093d393..5b573e07b3 100644 --- a/packages/twenty-docs/developers/extend/apps/logic/connections.mdx +++ b/packages/twenty-docs/developers/extend/apps/logic/connections.mdx @@ -45,6 +45,9 @@ export default defineConnectionProvider({ // Optional: a logic function in this app to run right after a connection is // established. See "Run a logic function on connect". // onConnectLogicFunction: { universalIdentifier: '3a2b1c0d-...-...' }, + // Optional: a logic function in this app to run right after a connection is + // removed. See "Run a logic function on disconnect". + // onDisconnectLogicFunction: { universalIdentifier: '4d5e6f70-...-...' }, }); ``` @@ -124,6 +127,42 @@ From there use `getConnection(connectedAccountId)` to read the fresh access toke + + +Anything an app claims at connect time has to be released when the connection goes away. A Slack integration that claims a `team_id` on connect, for instance, has to release that claim so another workspace can connect the same Slack team. Set `onDisconnectLogicFunction` to reference a logic function in the same app, and it runs right after the `ConnectedAccount` is deleted. + +```ts src/connection-providers/slack-connection.ts +export default defineConnectionProvider({ + universalIdentifier: '...', + name: 'slack', + displayName: 'Slack', + type: 'oauth', + oauth: { + /* ... */ + }, + // Runs releaseSlackTeam after every Slack disconnection. + onDisconnectLogicFunction: { + universalIdentifier: '4470aba8-5ff5-4800-88db-2a427cd8677c', + }, +}); +``` + +Like the on-connect hook it runs **asynchronously in the disconnecting workspace** and never blocks the disconnect. The handler receives the same payload shape: + +```ts +type OnDisconnectPayload = { + connectionProviderId: string; + connectionProviderName: string; // e.g. 'slack' + connectedAccountId: string; +}; +``` + +The `ConnectedAccount` is already gone when the hook runs, so `getConnection(connectedAccountId)` no longer resolves. Anything the cleanup needs (a `team_id`, an external subscription id) must have been written to the [key-value store](/developers/extend/apps/logic/key-value-store) at connect time, keyed by `connectedAccountId`. + +The hook fires when a connection is removed on its own. Uninstalling the app drops its connections through a database cascade instead, so the hook does not run there. Declare an `uninstallLogicFunction` on `defineApplication` for that path: it runs before the app's metadata is deleted, so it can still call `listConnections` and clean up whatever is left. + + + Inside a logic function handler, `listConnections({ providerName })` returns this app's `ConnectedAccount` rows for the given provider, with refreshed access tokens. diff --git a/packages/twenty-sdk/src/cli/utilities/build/manifest/__tests__/manifest-validate.spec.ts b/packages/twenty-sdk/src/cli/utilities/build/manifest/__tests__/manifest-validate.spec.ts index 83b913a80d..bb72dc46ba 100644 --- a/packages/twenty-sdk/src/cli/utilities/build/manifest/__tests__/manifest-validate.spec.ts +++ b/packages/twenty-sdk/src/cli/utilities/build/manifest/__tests__/manifest-validate.spec.ts @@ -197,36 +197,39 @@ describe('manifestValidate', () => { ); }); - it('should not flag a connection provider referencing a logic function via onConnectLogicFunction as a duplicate', () => { - const logicFunctionId = '550e8400-e29b-41d4-a716-446655440040'; + it.each(['onConnectLogicFunction', 'onDisconnectLogicFunction'] as const)( + 'should not flag a connection provider referencing a logic function via %s as a duplicate', + (lifecycleHookKey) => { + const logicFunctionId = '550e8400-e29b-41d4-a716-446655440040'; - const logicFunction = { - universalIdentifier: logicFunctionId, - name: 'onConnect', - sourceHandlerPath: 'src/logic-functions/on-connect.ts', - builtHandlerPath: 'dist/on-connect.js', - builtHandlerChecksum: '00000000-0000-4000-8000-000000000000', - handlerName: 'handler', - } as unknown as Manifest['logicFunctions'][number]; + const logicFunction = { + universalIdentifier: logicFunctionId, + name: lifecycleHookKey, + sourceHandlerPath: 'src/logic-functions/lifecycle-hook.ts', + builtHandlerPath: 'dist/lifecycle-hook.js', + builtHandlerChecksum: '00000000-0000-4000-8000-000000000000', + handlerName: 'handler', + } as unknown as Manifest['logicFunctions'][number]; - const connectionProvider = { - universalIdentifier: '550e8400-e29b-41d4-a716-446655440041', - name: 'slack', - displayName: 'Slack', - type: 'oauth', - oauth: {}, - onConnectLogicFunction: { universalIdentifier: logicFunctionId }, - } as unknown as NonNullable[number]; + const connectionProvider = { + universalIdentifier: '550e8400-e29b-41d4-a716-446655440041', + name: 'slack', + displayName: 'Slack', + type: 'oauth', + oauth: {}, + [lifecycleHookKey]: { universalIdentifier: logicFunctionId }, + } as unknown as NonNullable[number]; - const result = manifestValidate({ - ...validManifest, - logicFunctions: [logicFunction], - connectionProviders: [connectionProvider], - }); + const result = manifestValidate({ + ...validManifest, + logicFunctions: [logicFunction], + connectionProviders: [connectionProvider], + }); - expect(result.isValid).toBe(true); - expect(result.errors).toHaveLength(0); - }); + expect(result.isValid).toBe(true); + expect(result.errors).toHaveLength(0); + }, + ); it('should not flag a front component referenced via settingsFrontComponent as a duplicate', () => { const frontComponentId = '550e8400-e29b-41d4-a716-446655440050'; @@ -626,18 +629,14 @@ describe('manifestValidate', () => { expect(result.isValid).toBe(false); expect(result.errors).toContainEqual( - expect.stringContaining( - 'not "aggregateFieldMetadataId"', - ), + expect.stringContaining('not "aggregateFieldMetadataId"'), ); }); it('should ignore non-graph widgets that have no aggregate field', () => { const result = manifestValidate({ ...validManifest, - pageLayoutTabs: [ - makeGraphWidgetTab({ configurationType: 'TIMELINE' }), - ], + pageLayoutTabs: [makeGraphWidgetTab({ configurationType: 'TIMELINE' })], }); expect(result.isValid).toBe(true); diff --git a/packages/twenty-sdk/src/cli/utilities/build/manifest/manifest-validate.ts b/packages/twenty-sdk/src/cli/utilities/build/manifest/manifest-validate.ts index 498e45d837..1b84aaf29d 100644 --- a/packages/twenty-sdk/src/cli/utilities/build/manifest/manifest-validate.ts +++ b/packages/twenty-sdk/src/cli/utilities/build/manifest/manifest-validate.ts @@ -75,6 +75,7 @@ const findUniversalIdentifiers = (obj: object): string[] => { key === 'preInstallLogicFunction' || key === 'uninstallLogicFunction' || key === 'onConnectLogicFunction' || + key === 'onDisconnectLogicFunction' || key === 'settingsFrontComponent' ) { continue; diff --git a/packages/twenty-sdk/src/sdk/define/connection-providers/__tests__/define-connection-provider.spec.ts b/packages/twenty-sdk/src/sdk/define/connection-providers/__tests__/define-connection-provider.spec.ts index af72cf1fee..7c956ad805 100644 --- a/packages/twenty-sdk/src/sdk/define/connection-providers/__tests__/define-connection-provider.spec.ts +++ b/packages/twenty-sdk/src/sdk/define/connection-providers/__tests__/define-connection-provider.spec.ts @@ -83,10 +83,15 @@ describe('defineConnectionProvider', () => { expect(result.success).toBe(false); }); - it('accepts a valid onConnectLogicFunction', () => { + const lifecycleHookKeys = [ + 'onConnectLogicFunction', + 'onDisconnectLogicFunction', + ] as const; + + it.each(lifecycleHookKeys)('accepts a valid %s', (lifecycleHookKey) => { const result = defineConnectionProvider({ ...baseValidConfig, - onConnectLogicFunction: { + [lifecycleHookKey]: { universalIdentifier: 'b648f87b-1d26-4961-b974-0908fd991061', }, }); @@ -95,15 +100,18 @@ describe('defineConnectionProvider', () => { expect(result.errors).toEqual([]); }); - it('rejects a non-UUID onConnectLogicFunction universalIdentifier', () => { - const result = defineConnectionProvider({ - ...baseValidConfig, - onConnectLogicFunction: { universalIdentifier: 'claim-team-id' }, - }); + it.each(lifecycleHookKeys)( + 'rejects a non-UUID %s universalIdentifier', + (lifecycleHookKey) => { + const result = defineConnectionProvider({ + ...baseValidConfig, + [lifecycleHookKey]: { universalIdentifier: 'not-a-uuid' }, + }); - expect(result.success).toBe(false); - expect( - result.errors.some((error) => error.includes('onConnectLogicFunction')), - ).toBe(true); - }); + expect(result.success).toBe(false); + expect( + result.errors.some((error) => error.includes(lifecycleHookKey)), + ).toBe(true); + }, + ); }); diff --git a/packages/twenty-sdk/src/sdk/define/connection-providers/define-connection-provider.ts b/packages/twenty-sdk/src/sdk/define/connection-providers/define-connection-provider.ts index 9d760c31e4..216662e088 100644 --- a/packages/twenty-sdk/src/sdk/define/connection-providers/define-connection-provider.ts +++ b/packages/twenty-sdk/src/sdk/define/connection-providers/define-connection-provider.ts @@ -1,6 +1,7 @@ import { type DefineEntity } from '@/sdk/define/common/types/define-entity.type'; import { createValidationResult } from '@/sdk/define/common/utils/create-validation-result'; import { type ConnectionProviderManifest } from 'twenty-shared/application'; +import { isDefined } from 'twenty-shared/utils'; const PROVIDER_NAME_PATTERN = /^[a-z][a-z0-9-]*$/; // Matches UUID v1–v5 (and the `00000000-…` Nil UUID). Mirrors the server-side @@ -11,6 +12,16 @@ const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; const SUPPORTED_TYPES = ['oauth'] as const; +type ConnectionProviderLifecycleHookKey = Extract< + keyof ConnectionProviderManifest, + `on${string}LogicFunction` +>; + +const LIFECYCLE_HOOK_KEYS = [ + 'onConnectLogicFunction', + 'onDisconnectLogicFunction', +] as const satisfies readonly ConnectionProviderLifecycleHookKey[]; + export const defineConnectionProvider: DefineEntity< ConnectionProviderManifest > = (config) => { @@ -36,13 +47,14 @@ export const defineConnectionProvider: DefineEntity< errors.push('Connection provider must have a displayName'); } - if ( - config.onConnectLogicFunction && - !UUID_PATTERN.test(config.onConnectLogicFunction.universalIdentifier) - ) { - errors.push( - `Connection provider onConnectLogicFunction.universalIdentifier "${config.onConnectLogicFunction.universalIdentifier}" must be the UUID universalIdentifier of a logic function in this app.`, - ); + for (const hookKey of LIFECYCLE_HOOK_KEYS) { + const hook = config[hookKey]; + + if (isDefined(hook) && !UUID_PATTERN.test(hook.universalIdentifier)) { + errors.push( + `Connection provider ${hookKey}.universalIdentifier "${hook.universalIdentifier}" must be the UUID universalIdentifier of a logic function in this app.`, + ); + } } if (!config.type) { diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785810340935-add-on-disconnect-logic-function-to-connection-provider.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785810340935-add-on-disconnect-logic-function-to-connection-provider.ts new file mode 100644 index 0000000000..7ccdb4a2ab --- /dev/null +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785810340935-add-on-disconnect-logic-function-to-connection-provider.ts @@ -0,0 +1,21 @@ +import { type QueryRunner } from 'typeorm'; + +import { RegisteredInstanceCommand } from 'src/engine/core-modules/upgrade/decorators/registered-instance-command.decorator'; +import { type FastInstanceCommand } from 'src/engine/core-modules/upgrade/interfaces/fast-instance-command.interface'; + +@RegisteredInstanceCommand('2.27.0', 1785810340935) +export class AddOnDisconnectLogicFunctionToConnectionProviderFastInstanceCommand + implements FastInstanceCommand +{ + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query( + 'ALTER TABLE "core"."connectionProvider" ADD COLUMN IF NOT EXISTS "onDisconnectLogicFunctionUniversalIdentifier" uuid', + ); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query( + 'ALTER TABLE "core"."connectionProvider" DROP COLUMN IF EXISTS "onDisconnectLogicFunctionUniversalIdentifier"', + ); + } +} diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/instance-commands.constant.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/instance-commands.constant.ts index 32dab08a85..2f565fe0f7 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/instance-commands.constant.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/instance-commands.constant.ts @@ -130,6 +130,7 @@ import { AddAppTokenSsoExchangeIndexFastInstanceCommand } from './2-25/2-25-inst import { AddPageLayoutCascadeDeleteIndexesFastInstanceCommand } from './2-25/2-25-instance-command-fast-1784904030251-add-page-layout-cascade-delete-indexes'; import { AddChannelWebhookSubscriptionExternalIdIndexesFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-25/2-25-instance-command-fast-1785173910915-add-channel-webhook-subscription-external-id-indexes'; import { AddIsHiddenToAgentMessageFastInstanceCommand } from './2-25/2-25-instance-command-fast-1785230296000-add-is-hidden-to-agent-message'; +import { AddOnDisconnectLogicFunctionToConnectionProviderFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785810340935-add-on-disconnect-logic-function-to-connection-provider'; import { AddConnectedAccountHandleProviderIndexFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-26/2-26-instance-command-fast-1785420705255-add-connected-account-handle-provider-index'; import { AddOpenRecordInToObjectMetadataFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785504900000-add-open-record-in-to-object-metadata'; import { CreateUserSessionCoreTableFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-27/2-27-instance-command-fast-1785518325511-create-user-session-core-table'; @@ -266,6 +267,7 @@ export const INSTANCE_COMMANDS = [ AddPageLayoutCascadeDeleteIndexesFastInstanceCommand, AddChannelWebhookSubscriptionExternalIdIndexesFastInstanceCommand, AddIsHiddenToAgentMessageFastInstanceCommand, + AddOnDisconnectLogicFunctionToConnectionProviderFastInstanceCommand, AddConnectedAccountHandleProviderIndexFastInstanceCommand, AddOpenRecordInToObjectMetadataFastInstanceCommand, CreateUserSessionCoreTableFastInstanceCommand, diff --git a/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/__tests__/from-connection-provider-manifest-to-universal-flat-connection-provider.util.spec.ts b/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/__tests__/from-connection-provider-manifest-to-universal-flat-connection-provider.util.spec.ts index c9d7768571..5faf27cc68 100644 --- a/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/__tests__/from-connection-provider-manifest-to-universal-flat-connection-provider.util.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/__tests__/from-connection-provider-manifest-to-universal-flat-connection-provider.util.spec.ts @@ -51,6 +51,7 @@ describe('fromConnectionProviderManifestToUniversalFlatConnectionProvider', () = usePkce: true, }, onConnectLogicFunctionUniversalIdentifier: null, + onDisconnectLogicFunctionUniversalIdentifier: null, createdAt: NOW, updatedAt: NOW, }); @@ -76,6 +77,26 @@ describe('fromConnectionProviderManifestToUniversalFlatConnectionProvider', () = ); }); + it('resolves the onDisconnectLogicFunction universalIdentifier into the flat field when provided', () => { + const onDisconnectLogicFunctionUniversalIdentifier = + 'd2d2d2d2-d2d2-4d2d-d2d2-d2d2d2d2d2d2'; + + const result = + fromConnectionProviderManifestToUniversalFlatConnectionProvider({ + connectionProviderManifest: buildManifest({ + onDisconnectLogicFunction: { + universalIdentifier: onDisconnectLogicFunctionUniversalIdentifier, + }, + }), + applicationUniversalIdentifier: APP_UID, + now: NOW, + }); + + expect(result.onDisconnectLogicFunctionUniversalIdentifier).toBe( + onDisconnectLogicFunctionUniversalIdentifier, + ); + }); + it('passes through optional oauth config when provided', () => { const result = fromConnectionProviderManifestToUniversalFlatConnectionProvider({ diff --git a/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/from-connection-provider-manifest-to-universal-flat-connection-provider.util.ts b/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/from-connection-provider-manifest-to-universal-flat-connection-provider.util.ts index 599592a564..f7ff0c5c9c 100644 --- a/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/from-connection-provider-manifest-to-universal-flat-connection-provider.util.ts +++ b/packages/twenty-server/src/engine/core-modules/application/application-manifest/converters/from-connection-provider-manifest-to-universal-flat-connection-provider.util.ts @@ -46,6 +46,9 @@ export const fromConnectionProviderManifestToUniversalFlatConnectionProvider = onConnectLogicFunctionUniversalIdentifier: connectionProviderManifest.onConnectLogicFunction ?.universalIdentifier ?? null, + onDisconnectLogicFunctionUniversalIdentifier: + connectionProviderManifest.onDisconnectLogicFunction + ?.universalIdentifier ?? null, createdAt: now, updatedAt: now, }; diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-lifecycle-hook.service.spec.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-lifecycle-hook.service.spec.ts new file mode 100644 index 0000000000..c54a43e74f --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-lifecycle-hook.service.spec.ts @@ -0,0 +1,219 @@ +import { Test, type TestingModule } from '@nestjs/testing'; + +import { ConnectionProviderLifecycleHookService } from 'src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service'; +import { type ConnectionProviderEntity } from 'src/engine/core-modules/application/connection-provider/connection-provider.entity'; +import { ConnectionProviderService } from 'src/engine/core-modules/application/connection-provider/connection-provider.service'; +import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { getQueueToken } from 'src/engine/core-modules/message-queue/utils/get-queue-token.util'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; + +const ON_CONNECT_UID = 'c1c1c1c1-c1c1-4c1c-c1c1-c1c1c1c1c1c1'; +const ON_DISCONNECT_UID = 'd2d2d2d2-d2d2-4d2d-d2d2-d2d2d2d2d2d2'; + +describe('ConnectionProviderLifecycleHookService', () => { + let service: ConnectionProviderLifecycleHookService; + let connectionProviderService: { findOneByIdOrThrow: jest.Mock }; + let workspaceCacheService: { getOrRecompute: jest.Mock }; + let messageQueueService: { add: jest.Mock }; + let exceptionHandlerService: { captureExceptions: jest.Mock }; + + const baseProvider: Pick< + ConnectionProviderEntity, + | 'id' + | 'name' + | 'onConnectLogicFunctionUniversalIdentifier' + | 'onDisconnectLogicFunctionUniversalIdentifier' + > = { + id: 'provider-1', + name: 'linear', + onConnectLogicFunctionUniversalIdentifier: null, + onDisconnectLogicFunctionUniversalIdentifier: null, + }; + + beforeEach(async () => { + connectionProviderService = { + findOneByIdOrThrow: jest.fn(async () => baseProvider), + }; + workspaceCacheService = { + getOrRecompute: jest.fn(async () => ({ + flatLogicFunctionMaps: { byUniversalIdentifier: {} }, + })), + }; + messageQueueService = { add: jest.fn() }; + exceptionHandlerService = { captureExceptions: jest.fn() }; + + const module: TestingModule = await Test.createTestingModule({ + providers: [ + ConnectionProviderLifecycleHookService, + { + provide: ConnectionProviderService, + useValue: connectionProviderService, + }, + { provide: WorkspaceCacheService, useValue: workspaceCacheService }, + { + provide: getQueueToken(MessageQueue.logicFunctionQueue), + useValue: messageQueueService, + }, + { provide: ExceptionHandlerService, useValue: exceptionHandlerService }, + ], + }).compile(); + + service = module.get(ConnectionProviderLifecycleHookService); + }); + + afterEach(() => jest.clearAllMocks()); + + describe('dispatchOnConnect', () => { + it('does not dispatch a hook when the provider declares none', async () => { + await service.dispatchOnConnect({ + provider: baseProvider, + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(workspaceCacheService.getOrRecompute).not.toHaveBeenCalled(); + expect(messageQueueService.add).not.toHaveBeenCalled(); + }); + + it('enqueues the declared on-connect logic function in the connecting workspace', async () => { + workspaceCacheService.getOrRecompute.mockResolvedValue({ + flatLogicFunctionMaps: { + byUniversalIdentifier: { [ON_CONNECT_UID]: { id: 'function-1' } }, + }, + }); + + await service.dispatchOnConnect({ + provider: { + ...baseProvider, + onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, + }, + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(workspaceCacheService.getOrRecompute).toHaveBeenCalledWith( + 'workspace-1', + ['flatLogicFunctionMaps'], + ); + expect(messageQueueService.add).toHaveBeenCalledWith( + 'LogicFunctionTriggerJob', + { + logicFunctionId: 'function-1', + workspaceId: 'workspace-1', + payload: { + connectionProviderId: 'provider-1', + connectionProviderName: 'linear', + connectedAccountId: 'account-1', + }, + }, + { retryLimit: 3 }, + ); + }); + + it('reports to Sentry without throwing when the hook function is missing', async () => { + await service.dispatchOnConnect({ + provider: { + ...baseProvider, + onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, + }, + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(messageQueueService.add).not.toHaveBeenCalled(); + expect(exceptionHandlerService.captureExceptions).toHaveBeenCalledTimes( + 1, + ); + }); + + it('treats a soft-deleted hook function as missing instead of enqueuing it', async () => { + workspaceCacheService.getOrRecompute.mockResolvedValue({ + flatLogicFunctionMaps: { + byUniversalIdentifier: { + [ON_CONNECT_UID]: { + id: 'function-1', + deletedAt: new Date().toISOString(), + }, + }, + }, + }); + + await service.dispatchOnConnect({ + provider: { + ...baseProvider, + onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, + }, + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(messageQueueService.add).not.toHaveBeenCalled(); + expect(exceptionHandlerService.captureExceptions).toHaveBeenCalledTimes( + 1, + ); + }); + }); + + describe('dispatchOnDisconnect', () => { + it('does not dispatch a hook when the provider declares none', async () => { + await service.dispatchOnDisconnect({ + connectionProviderId: 'provider-1', + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(messageQueueService.add).not.toHaveBeenCalled(); + expect(exceptionHandlerService.captureExceptions).not.toHaveBeenCalled(); + }); + + it('enqueues the declared on-disconnect logic function in the disconnecting workspace', async () => { + connectionProviderService.findOneByIdOrThrow.mockResolvedValue({ + ...baseProvider, + onDisconnectLogicFunctionUniversalIdentifier: ON_DISCONNECT_UID, + }); + workspaceCacheService.getOrRecompute.mockResolvedValue({ + flatLogicFunctionMaps: { + byUniversalIdentifier: { [ON_DISCONNECT_UID]: { id: 'function-2' } }, + }, + }); + + await service.dispatchOnDisconnect({ + connectionProviderId: 'provider-1', + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(messageQueueService.add).toHaveBeenCalledWith( + 'LogicFunctionTriggerJob', + { + logicFunctionId: 'function-2', + workspaceId: 'workspace-1', + payload: { + connectionProviderId: 'provider-1', + connectionProviderName: 'linear', + connectedAccountId: 'account-1', + }, + }, + { retryLimit: 3 }, + ); + }); + + it('reports to Sentry without throwing when the provider no longer exists', async () => { + connectionProviderService.findOneByIdOrThrow.mockRejectedValue( + new Error('Provider not found'), + ); + + await service.dispatchOnDisconnect({ + connectionProviderId: 'provider-1', + workspaceId: 'workspace-1', + connectedAccountId: 'account-1', + }); + + expect(messageQueueService.add).not.toHaveBeenCalled(); + expect(exceptionHandlerService.captureExceptions).toHaveBeenCalledTimes( + 1, + ); + }); + }); +}); diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts index 2d1956493f..307cd51058 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts @@ -16,19 +16,16 @@ import { ConnectedAccountProvider } from 'twenty-shared/types'; import { isDefined } from 'twenty-shared/utils'; import { type ConnectionProviderEntity } from 'src/engine/core-modules/application/connection-provider/connection-provider.entity'; +import { ConnectionProviderLifecycleHookService } from 'src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service'; import { ConnectionProviderOAuthFlowService } from 'src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service'; import { ConnectionProviderService } from 'src/engine/core-modules/application/connection-provider/connection-provider.service'; import { JwtTokenTypeEnum } from 'src/engine/core-modules/auth/types/jwt-token-type.enum'; -import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { JwtWrapperService } from 'src/engine/core-modules/jwt/services/jwt-wrapper.service'; import { SecureHttpClientService } from 'src/engine/core-modules/secure-http-client/secure-http-client.service'; import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; -import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; -import { getQueueToken } from 'src/engine/core-modules/message-queue/utils/get-queue-token.util'; import { SECRET_ENCRYPTION_ENVELOPE_V2_PREFIX } from 'src/engine/core-modules/secret-encryption/constants/secret-encryption.constant'; import { ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity'; import { ConnectedAccountTokenEncryptionService } from 'src/engine/metadata-modules/connected-account/services/connected-account-token-encryption.service'; -import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; const FAKE_CIPHER_PREFIX = `${SECRET_ENCRYPTION_ENVELOPE_V2_PREFIX}keyid:`; @@ -51,9 +48,7 @@ describe('ConnectionProviderOAuthFlowService', () => { findOne: jest.Mock; findOneByOrFail: jest.Mock; }; - let workspaceCacheService: { getOrRecompute: jest.Mock }; - let messageQueueService: { add: jest.Mock }; - let exceptionHandlerService: { captureExceptions: jest.Mock }; + let connectionProviderLifecycleHookService: { dispatchOnConnect: jest.Mock }; const baseProvider: ConnectionProviderEntity = { id: 'provider-1', @@ -102,13 +97,9 @@ describe('ConnectionProviderOAuthFlowService', () => { provider: ConnectedAccountProvider.APP, })), }; - workspaceCacheService = { - getOrRecompute: jest.fn(async () => ({ - flatLogicFunctionMaps: { byUniversalIdentifier: {} }, - })), + connectionProviderLifecycleHookService = { + dispatchOnConnect: jest.fn(), }; - messageQueueService = { add: jest.fn() }; - exceptionHandlerService = { captureExceptions: jest.fn() }; const module: TestingModule = await Test.createTestingModule({ providers: [ @@ -128,16 +119,8 @@ describe('ConnectionProviderOAuthFlowService', () => { useValue: connectedAccountRepository, }, { - provide: WorkspaceCacheService, - useValue: workspaceCacheService, - }, - { - provide: getQueueToken(MessageQueue.logicFunctionQueue), - useValue: messageQueueService, - }, - { - provide: ExceptionHandlerService, - useValue: exceptionHandlerService, + provide: ConnectionProviderLifecycleHookService, + useValue: connectionProviderLifecycleHookService, }, { // Real prefix/round-trip behavior is asserted in @@ -464,103 +447,18 @@ describe('ConnectionProviderOAuthFlowService', () => { ).rejects.toThrow(/state/); }); - describe('on-connect hook', () => { - const ON_CONNECT_UID = 'c1c1c1c1-c1c1-4c1c-c1c1-c1c1c1c1c1c1'; - - it('does not dispatch a hook when the provider declares none', async () => { - await service.completeAuthorizationFlow({ - code: 'auth_code', - state: 'signed-state', - }); - - expect(workspaceCacheService.getOrRecompute).not.toHaveBeenCalled(); - expect(messageQueueService.add).not.toHaveBeenCalled(); + it('hands the created connection to the on-connect lifecycle hook', async () => { + const result = await service.completeAuthorizationFlow({ + code: 'auth_code', + state: 'signed-state', }); - it('enqueues the declared on-connect logic function in the connecting workspace', async () => { - connectionProviderService.findOneByIdOrThrow.mockResolvedValue({ - ...baseProvider, - onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, - }); - workspaceCacheService.getOrRecompute.mockResolvedValue({ - flatLogicFunctionMaps: { - byUniversalIdentifier: { - [ON_CONNECT_UID]: { id: 'logic-function-1' }, - }, - }, - }); - - const result = await service.completeAuthorizationFlow({ - code: 'auth_code', - state: 'signed-state', - }); - - expect(workspaceCacheService.getOrRecompute).toHaveBeenCalledWith( - 'workspace-1', - ['flatLogicFunctionMaps'], - ); - expect(messageQueueService.add).toHaveBeenCalledWith( - 'LogicFunctionTriggerJob', - { - logicFunctionId: 'logic-function-1', - workspaceId: 'workspace-1', - payload: { - connectionProviderId: 'provider-1', - connectionProviderName: 'linear', - connectedAccountId: result.connectedAccountId, - }, - }, - { retryLimit: 3 }, - ); - }); - - it('reports to Sentry without failing the connection when the hook function is missing', async () => { - connectionProviderService.findOneByIdOrThrow.mockResolvedValue({ - ...baseProvider, - onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, - }); - workspaceCacheService.getOrRecompute.mockResolvedValue({ - flatLogicFunctionMaps: { byUniversalIdentifier: {} }, - }); - - const result = await service.completeAuthorizationFlow({ - code: 'auth_code', - state: 'signed-state', - }); - - expect(result.connectedAccountId).toBe('new-account-id'); - expect(messageQueueService.add).not.toHaveBeenCalled(); - expect(exceptionHandlerService.captureExceptions).toHaveBeenCalledTimes( - 1, - ); - }); - - it('treats a soft-deleted hook function as missing instead of enqueuing it', async () => { - connectionProviderService.findOneByIdOrThrow.mockResolvedValue({ - ...baseProvider, - onConnectLogicFunctionUniversalIdentifier: ON_CONNECT_UID, - }); - workspaceCacheService.getOrRecompute.mockResolvedValue({ - flatLogicFunctionMaps: { - byUniversalIdentifier: { - [ON_CONNECT_UID]: { - id: 'logic-function-1', - deletedAt: new Date().toISOString(), - }, - }, - }, - }); - - const result = await service.completeAuthorizationFlow({ - code: 'auth_code', - state: 'signed-state', - }); - - expect(result.connectedAccountId).toBe('new-account-id'); - expect(messageQueueService.add).not.toHaveBeenCalled(); - expect(exceptionHandlerService.captureExceptions).toHaveBeenCalledTimes( - 1, - ); + expect( + connectionProviderLifecycleHookService.dispatchOnConnect, + ).toHaveBeenCalledWith({ + provider: baseProvider, + workspaceId: 'workspace-1', + connectedAccountId: result.connectedAccountId, }); }); }); diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum.ts index 1279f15a39..d39e72cdb7 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum.ts @@ -10,4 +10,5 @@ export enum ConnectionProviderExceptionCode { CONNECTION_PROVIDER_NOT_FOUND = 'CONNECTION_PROVIDER_NOT_FOUND', CONNECTION_PROVIDER_NAME_ALREADY_EXISTS = 'CONNECTION_PROVIDER_NAME_ALREADY_EXISTS', ON_CONNECT_LOGIC_FUNCTION_NOT_FOUND = 'ON_CONNECT_LOGIC_FUNCTION_NOT_FOUND', + ON_DISCONNECT_LOGIC_FUNCTION_NOT_FOUND = 'ON_DISCONNECT_LOGIC_FUNCTION_NOT_FOUND', } diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service.ts new file mode 100644 index 0000000000..5100500079 --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service.ts @@ -0,0 +1,157 @@ +import { Injectable } from '@nestjs/common'; + +import { isDefined } from 'twenty-shared/utils'; + +import { ConnectionProviderExceptionCode } from 'src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum'; +import { type ConnectionProviderEntity } from 'src/engine/core-modules/application/connection-provider/connection-provider.entity'; +import { ConnectionProviderException } from 'src/engine/core-modules/application/connection-provider/connection-provider.exception'; +import { ConnectionProviderService } from 'src/engine/core-modules/application/connection-provider/connection-provider.service'; +import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; +import { + LogicFunctionTriggerJob, + type LogicFunctionTriggerJobData, +} from 'src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job'; +import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; + +type ConnectionLifecycleHook = 'onConnect' | 'onDisconnect'; + +const MISSING_LOGIC_FUNCTION_EXCEPTION_CODE_BY_HOOK: Record< + ConnectionLifecycleHook, + ConnectionProviderExceptionCode +> = { + onConnect: + ConnectionProviderExceptionCode.ON_CONNECT_LOGIC_FUNCTION_NOT_FOUND, + onDisconnect: + ConnectionProviderExceptionCode.ON_DISCONNECT_LOGIC_FUNCTION_NOT_FOUND, +}; + +@Injectable() +export class ConnectionProviderLifecycleHookService { + constructor( + private readonly connectionProviderService: ConnectionProviderService, + @InjectMessageQueue(MessageQueue.logicFunctionQueue) + private readonly messageQueueService: MessageQueueService, + private readonly workspaceCacheService: WorkspaceCacheService, + private readonly exceptionHandlerService: ExceptionHandlerService, + ) {} + + async dispatchOnConnect({ + provider, + workspaceId, + connectedAccountId, + }: { + provider: Pick< + ConnectionProviderEntity, + 'id' | 'name' | 'onConnectLogicFunctionUniversalIdentifier' + >; + workspaceId: string; + connectedAccountId: string; + }): Promise { + await this.captureFailures(workspaceId, () => + this.enqueueLogicFunction({ + hook: 'onConnect', + logicFunctionUniversalIdentifier: + provider.onConnectLogicFunctionUniversalIdentifier, + provider, + workspaceId, + connectedAccountId, + }), + ); + } + + async dispatchOnDisconnect({ + connectionProviderId, + workspaceId, + connectedAccountId, + }: { + connectionProviderId: string; + workspaceId: string; + connectedAccountId: string; + }): Promise { + await this.captureFailures(workspaceId, async () => { + const provider = + await this.connectionProviderService.findOneByIdOrThrow( + connectionProviderId, + ); + + await this.enqueueLogicFunction({ + hook: 'onDisconnect', + logicFunctionUniversalIdentifier: + provider.onDisconnectLogicFunctionUniversalIdentifier, + provider, + workspaceId, + connectedAccountId, + }); + }); + } + + // Lifecycle hooks are best effort: a failing hook must never surface to the + // user connecting or disconnecting their account. + private async captureFailures( + workspaceId: string, + dispatch: () => Promise, + ): Promise { + try { + await dispatch(); + } catch (error) { + this.exceptionHandlerService.captureExceptions([error], { + workspace: { id: workspaceId }, + }); + } + } + + private async enqueueLogicFunction({ + hook, + logicFunctionUniversalIdentifier, + provider, + workspaceId, + connectedAccountId, + }: { + hook: ConnectionLifecycleHook; + logicFunctionUniversalIdentifier: string | null; + provider: Pick; + workspaceId: string; + connectedAccountId: string; + }): Promise { + if (!isDefined(logicFunctionUniversalIdentifier)) { + return; + } + + const { flatLogicFunctionMaps } = + await this.workspaceCacheService.getOrRecompute(workspaceId, [ + 'flatLogicFunctionMaps', + ]); + + const flatLogicFunction = + flatLogicFunctionMaps.byUniversalIdentifier[ + logicFunctionUniversalIdentifier + ]; + + if ( + !isDefined(flatLogicFunction) || + isDefined(flatLogicFunction.deletedAt) + ) { + throw new ConnectionProviderException( + `Connection provider ${provider.id} references ${hook} logic function ${logicFunctionUniversalIdentifier}, which was not found in workspace ${workspaceId}.`, + MISSING_LOGIC_FUNCTION_EXCEPTION_CODE_BY_HOOK[hook], + ); + } + + await this.messageQueueService.add( + LogicFunctionTriggerJob.name, + { + logicFunctionId: flatLogicFunction.id, + workspaceId, + payload: { + connectionProviderId: provider.id, + connectionProviderName: provider.name, + connectedAccountId, + }, + }, + { retryLimit: 3 }, + ); + } +} diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts index efce22153e..a9cb9d2768 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts @@ -7,6 +7,7 @@ import { ConnectedAccountProvider } from 'twenty-shared/types'; import { isDefined } from 'twenty-shared/utils'; import { ConnectionProviderExceptionCode } from 'src/engine/core-modules/application/connection-provider/connection-provider-exception-code.enum'; +import { ConnectionProviderLifecycleHookService } from 'src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service'; import { type ConnectionProviderEntity } from 'src/engine/core-modules/application/connection-provider/connection-provider.entity'; import { ConnectionProviderException } from 'src/engine/core-modules/application/connection-provider/connection-provider.exception'; import { ConnectionProviderService } from 'src/engine/core-modules/application/connection-provider/connection-provider.service'; @@ -21,20 +22,11 @@ import { exchangeCodeForToken } from 'src/engine/core-modules/application/connec import { generatePkceVerifier } from 'src/engine/core-modules/application/connection-provider/utils/generate-pkce-verifier.util'; import { type AppOAuthStateJwtPayload } from 'src/engine/core-modules/auth/types/app-oauth-state-jwt-payload.type'; import { JwtTokenTypeEnum } from 'src/engine/core-modules/auth/types/jwt-token-type.enum'; -import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { JwtWrapperService } from 'src/engine/core-modules/jwt/services/jwt-wrapper.service'; -import { - LogicFunctionTriggerJob, - type LogicFunctionTriggerJobData, -} from 'src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job'; -import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; -import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; -import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { SecureHttpClientService } from 'src/engine/core-modules/secure-http-client/secure-http-client.service'; import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; import { ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity'; import { ConnectedAccountTokenEncryptionService } from 'src/engine/metadata-modules/connected-account/services/connected-account-token-encryption.service'; -import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; const STATE_JWT_EXPIRES_IN = '10m'; @@ -70,12 +62,9 @@ export class ConnectionProviderOAuthFlowService { private readonly secureHttpClientService: SecureHttpClientService, private readonly twentyConfigService: TwentyConfigService, private readonly connectedAccountTokenEncryptionService: ConnectedAccountTokenEncryptionService, + private readonly connectionProviderLifecycleHookService: ConnectionProviderLifecycleHookService, @InjectRepository(ConnectedAccountEntity) private readonly connectedAccountRepository: Repository, - @InjectMessageQueue(MessageQueue.logicFunctionQueue) - private readonly messageQueueService: MessageQueueService, - private readonly workspaceCacheService: WorkspaceCacheService, - private readonly exceptionHandlerService: ExceptionHandlerService, ) {} async startAuthorizationFlow( @@ -199,7 +188,7 @@ export class ConnectionProviderOAuthFlowService { statePayload.reconnectingConnectedAccountId, }); - await this.dispatchOnConnectHook({ + await this.connectionProviderLifecycleHookService.dispatchOnConnect({ provider, workspaceId: statePayload.workspaceId, connectedAccountId: connectedAccount.id, @@ -213,65 +202,6 @@ export class ConnectionProviderOAuthFlowService { }; } - private async dispatchOnConnectHook({ - provider, - workspaceId, - connectedAccountId, - }: { - provider: OAuthConnectionProvider; - workspaceId: string; - connectedAccountId: string; - }): Promise { - const { onConnectLogicFunctionUniversalIdentifier } = provider; - - if (!isDefined(onConnectLogicFunctionUniversalIdentifier)) { - return; - } - - // The on-connect hook is best-effort: the ConnectedAccount is already - // persisted, so a misconfigured or failing hook must not break the OAuth - // callback. We still report failures to Sentry so they don't go unnoticed. - try { - const { flatLogicFunctionMaps } = - await this.workspaceCacheService.getOrRecompute(workspaceId, [ - 'flatLogicFunctionMaps', - ]); - - const flatLogicFunction = - flatLogicFunctionMaps.byUniversalIdentifier[ - onConnectLogicFunctionUniversalIdentifier - ]; - - if ( - !isDefined(flatLogicFunction) || - isDefined(flatLogicFunction.deletedAt) - ) { - throw new ConnectionProviderException( - `Connection provider ${provider.id} references on-connect logic function ${onConnectLogicFunctionUniversalIdentifier}, which was not found in workspace ${workspaceId}.`, - ConnectionProviderExceptionCode.ON_CONNECT_LOGIC_FUNCTION_NOT_FOUND, - ); - } - - await this.messageQueueService.add( - LogicFunctionTriggerJob.name, - { - logicFunctionId: flatLogicFunction.id, - workspaceId, - payload: { - connectionProviderId: provider.id, - connectionProviderName: provider.name, - connectedAccountId, - }, - }, - { retryLimit: 3 }, - ); - } catch (error) { - this.exceptionHandlerService.captureExceptions([error], { - workspace: { id: workspaceId }, - }); - } - } - private async signState(payload: AppOAuthStateJwtPayload): Promise { return this.jwtWrapperService.signAsyncOrThrow(payload, { expiresIn: STATE_JWT_EXPIRES_IN, diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.entity.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.entity.ts index 6d6f06a284..cbc0e2331b 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.entity.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.entity.ts @@ -47,6 +47,13 @@ export class ConnectionProviderEntity }) onConnectLogicFunctionUniversalIdentifier: string | null; + @Column({ nullable: true, type: 'uuid' }) + @WasIntroducedInUpgrade({ + upgradeCommandName: + '2.27.0_AddOnDisconnectLogicFunctionToConnectionProviderFastInstanceCommand_1785810340935', + }) + onDisconnectLogicFunctionUniversalIdentifier: string | null; + @CreateDateColumn({ type: 'timestamptz' }) createdAt: Date; diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.exception.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.exception.ts index 0081737744..2d786c4cec 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.exception.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.exception.ts @@ -31,6 +31,8 @@ const getConnectionProviderExceptionUserFriendlyMessage = ( return msg`A connection provider with this name already exists for this application.`; case ConnectionProviderExceptionCode.ON_CONNECT_LOGIC_FUNCTION_NOT_FOUND: return msg`The logic function to run on connect was not found.`; + case ConnectionProviderExceptionCode.ON_DISCONNECT_LOGIC_FUNCTION_NOT_FOUND: + return msg`The logic function to run on disconnect was not found.`; default: assertUnreachable(code); } diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.module.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.module.ts index fe15e96d0b..c9715b9669 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.module.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider.module.ts @@ -3,6 +3,7 @@ import { TypeOrmModule } from '@nestjs/typeorm'; import { ApplicationConnectionProviderResolver } from 'src/engine/core-modules/application/connection-provider/application-connection-provider.resolver'; import { ConnectionProviderEntity } from 'src/engine/core-modules/application/connection-provider/connection-provider.entity'; +import { ConnectionProviderLifecycleHookService } from 'src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service'; import { ConnectionProviderOAuthFlowService } from 'src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service'; import { ConnectionProviderService } from 'src/engine/core-modules/application/connection-provider/connection-provider.service'; import { ApplicationRegistrationVariableEntity } from 'src/engine/core-modules/application/application-registration-variable/application-registration-variable.entity'; @@ -35,8 +36,13 @@ import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache providers: [ ConnectionProviderService, ConnectionProviderOAuthFlowService, + ConnectionProviderLifecycleHookService, ApplicationConnectionProviderResolver, ], - exports: [ConnectionProviderService, ConnectionProviderOAuthFlowService], + exports: [ + ConnectionProviderService, + ConnectionProviderOAuthFlowService, + ConnectionProviderLifecycleHookService, + ], }) export class ConnectionProviderModule {} diff --git a/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.module.ts b/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.module.ts index 75cead6837..f449b38edd 100644 --- a/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.module.ts +++ b/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.module.ts @@ -1,6 +1,7 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; +import { ConnectionProviderModule } from 'src/engine/core-modules/application/connection-provider/connection-provider.module'; import { AppOAuthRefreshModule } from 'src/engine/core-modules/application/connection-provider/refresh/app-oauth-refresh.module'; import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; import { CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity'; @@ -21,6 +22,7 @@ import { WorkspaceEventEmitterModule } from 'src/engine/workspace-event-emitter/ MessageChannelEntity, ]), AppOAuthRefreshModule, + ConnectionProviderModule, FeatureFlagModule, PermissionsModule, WorkspaceEventEmitterModule, diff --git a/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.service.ts b/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.service.ts index 8548756c64..70e3aeeccf 100644 --- a/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/connected-account/connected-account-metadata.service.ts @@ -3,6 +3,9 @@ import { InjectRepository } from '@nestjs/typeorm'; import { In, IsNull, Repository } from 'typeorm'; +import { isDefined } from 'twenty-shared/utils'; + +import { ConnectionProviderLifecycleHookService } from 'src/engine/core-modules/application/connection-provider/connection-provider-lifecycle-hook.service'; import { AppOAuthRevokeService } from 'src/engine/core-modules/application/connection-provider/refresh/services/app-oauth-revoke.service'; import { CALENDAR_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/calendar-channel/constants/calendar-channel-deleted.constant'; import { CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity'; @@ -31,6 +34,7 @@ export class ConnectedAccountMetadataService { @InjectRepository(MessageChannelEntity) private readonly messageChannelRepository: Repository, private readonly appOAuthRevokeService: AppOAuthRevokeService, + private readonly connectionProviderLifecycleHookService: ConnectionProviderLifecycleHookService, private readonly workspaceEventEmitter: WorkspaceEventEmitter, ) {} @@ -246,6 +250,14 @@ export class ConnectedAccountMetadataService { await this.repository.delete({ id, workspaceId }); + if (isDefined(connectedAccount.connectionProviderId)) { + await this.connectionProviderLifecycleHookService.dispatchOnDisconnect({ + connectionProviderId: connectedAccount.connectionProviderId, + workspaceId, + connectedAccountId: id, + }); + } + this.workspaceEventEmitter.emitCustomBatchEvent( MESSAGE_CHANNEL_DELETED_EVENT, messageChannels.map((messageChannel) => ({ diff --git a/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/__tests__/__snapshots__/all-universal-flat-entity-properties-to-compare-and-stringify.constant.spec.ts.snap b/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/__tests__/__snapshots__/all-universal-flat-entity-properties-to-compare-and-stringify.constant.spec.ts.snap index ff31f1ce92..2d0eda97e8 100644 --- a/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/__tests__/__snapshots__/all-universal-flat-entity-properties-to-compare-and-stringify.constant.spec.ts.snap +++ b/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/__tests__/__snapshots__/all-universal-flat-entity-properties-to-compare-and-stringify.constant.spec.ts.snap @@ -61,6 +61,7 @@ exports[`ALL_UNIVERSAL_FLAT_ENTITY_PROPERTIES_TO_COMPARE_AND_STRINGIFY should ma "type", "oauthConfig", "onConnectLogicFunctionUniversalIdentifier", + "onDisconnectLogicFunctionUniversalIdentifier", ], "propertiesToStringify": [ "oauthConfig", diff --git a/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/all-entity-properties-configuration-by-metadata-name.constant.ts b/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/all-entity-properties-configuration-by-metadata-name.constant.ts index 1ced675da8..17508683dd 100644 --- a/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/all-entity-properties-configuration-by-metadata-name.constant.ts +++ b/packages/twenty-server/src/engine/metadata-modules/flat-entity/constant/all-entity-properties-configuration-by-metadata-name.constant.ts @@ -1813,6 +1813,11 @@ export const ALL_ENTITY_PROPERTIES_CONFIGURATION_BY_METADATA_NAME = { toStringify: false, universalProperty: undefined, }, + onDisconnectLogicFunctionUniversalIdentifier: { + toCompare: true, + toStringify: false, + universalProperty: undefined, + }, createdAt: { toCompare: false, toStringify: false, diff --git a/packages/twenty-server/test/integration/metadata/suites/application/successful-manifest-update-connection-provider.integration-spec.ts b/packages/twenty-server/test/integration/metadata/suites/application/successful-manifest-update-connection-provider.integration-spec.ts index 610f2ba413..1815778bf0 100644 --- a/packages/twenty-server/test/integration/metadata/suites/application/successful-manifest-update-connection-provider.integration-spec.ts +++ b/packages/twenty-server/test/integration/metadata/suites/application/successful-manifest-update-connection-provider.integration-spec.ts @@ -84,6 +84,7 @@ describe('Manifest update - connection providers', () => { usePkce: true, }, onConnectLogicFunctionUniversalIdentifier: null, + onDisconnectLogicFunctionUniversalIdentifier: null, }); }, 60000); diff --git a/packages/twenty-server/test/integration/metadata/suites/connection-provider/utils/find-connection-providers-by-application.util.ts b/packages/twenty-server/test/integration/metadata/suites/connection-provider/utils/find-connection-providers-by-application.util.ts index b5d3a5925d..5d7817c70e 100644 --- a/packages/twenty-server/test/integration/metadata/suites/connection-provider/utils/find-connection-providers-by-application.util.ts +++ b/packages/twenty-server/test/integration/metadata/suites/connection-provider/utils/find-connection-providers-by-application.util.ts @@ -10,6 +10,7 @@ type ConnectionProviderRow = { type: string; oauthConfig: StoredOAuthConnectionProviderConfig | null; onConnectLogicFunctionUniversalIdentifier: string | null; + onDisconnectLogicFunctionUniversalIdentifier: string | null; }; export const findConnectionProvidersByApplication = async ( @@ -18,7 +19,8 @@ export const findConnectionProvidersByApplication = async ( return globalThis.testDataSource.query( `SELECT cp.id, cp."universalIdentifier", cp."applicationId", cp."workspaceId", cp.name, cp."displayName", cp.type, - cp."oauthConfig", cp."onConnectLogicFunctionUniversalIdentifier" + cp."oauthConfig", cp."onConnectLogicFunctionUniversalIdentifier", + cp."onDisconnectLogicFunctionUniversalIdentifier" FROM core."connectionProvider" cp JOIN core."application" app ON app.id = cp."applicationId" WHERE app."universalIdentifier" = $1 diff --git a/packages/twenty-shared/src/application/connectionProviderManifestType.ts b/packages/twenty-shared/src/application/connectionProviderManifestType.ts index d7a37be78a..d1fcfcaab8 100644 --- a/packages/twenty-shared/src/application/connectionProviderManifestType.ts +++ b/packages/twenty-shared/src/application/connectionProviderManifestType.ts @@ -7,4 +7,5 @@ export type ConnectionProviderManifest = SyncableEntityOptions & { type: 'oauth'; oauth: OAuthConnectionProviderConfig; onConnectLogicFunction?: SyncableEntityOptions; + onDisconnectLogicFunction?: SyncableEntityOptions; };