import type { AgentIntegrationConfig, ChatHubMessageStatus, InstanceAiEvent, PushMessage, PushPayload, WorkerStatus, WorkflowPublicationStatusMessage, } from '@n8n/api-types'; import type { IWorkflowBase, RelatedAgentRun, WorkflowActivateMode } from 'n8n-workflow'; export type PubSubCommandMap = { // #region Lifecycle 'reload-license': never; 'restart-event-bus': never; 'reload-external-secrets-providers': never; // #endregion // # region Credentials 'reload-overwrite-credentials': never; // #endregion // # region SSO 'reload-oidc-config': never; 'reload-saml-config': never; // # sso provisioning 'reload-sso-provisioning-configuration': never; // #endregion 'reload-source-control-config': never; 'reload-mcp-registry': never; 'reload-otel-config': never; 'reload-instance-ai-settings': never; // #region Community packages 'community-package-install': { packageName: string; packageVersion: string; checksum?: string; }; 'community-package-update': { packageName: string; packageVersion: string; checksum?: string; }; 'community-package-uninstall': { packageName: string; }; // #endregion // #region Worker view 'get-worker-id': never; 'get-worker-status': { requestingUserId: string; }; // #endregion // #region Execution control /** * Stop a specific in-memory execution on whichever worker is running it. Used in queue * mode for subworkflow executions, which run inline in the parent's worker process and * therefore have no Bull job to abort. Each worker checks its own `ActiveExecutions` and * cancels the execution if it holds it. */ 'stop-execution': { executionId: string; }; // #endregion // #region Multi-main setup 'add-webhooks-triggers-and-pollers': { workflowId: string; activeVersionId: string; activationMode: WorkflowActivateMode; }; 'remove-triggers-and-pollers': { workflowId: string; }; 'display-workflow-activation': { workflowId: string; activeVersionId: string; }; 'display-workflow-deactivation': { workflowId: string; }; /** Wake the leader's publication outbox consumer to drain pending records now. * The consumer polls as a backup, but this lets publication feel fast without frequent polling. */ 'workflow-publish-wake-up': never; 'display-workflow-activation-error': { workflowId: string; errorMessage: string; errorDescription?: string; nodeId?: string; }; /** Relay a publication status push from the leader (which drains the publication * outbox) to the other mains, so clients connected to followers get the push too. */ 'display-workflow-publication-status': WorkflowPublicationStatusMessage; 'relay-execution-lifecycle-event': PushMessage & { pushRef: string; asBinary: boolean; }; 'relay-agent-execution-update': { data: PushPayload<'agentExecutionUpdated'>; userIds: string[]; }; 'relay-agent-background-tasks-update': { data: PushPayload<'agentBackgroundTasksUpdated'>; userIds: string[]; }; 'relay-agent-update': { data: PushPayload<'agentUpdated'>; userIds: string[]; excludePushRef?: string; }; /** Ask mains to wake the agent run a finished sub-execution was parked on. */ 'resume-agent-workflow-tool': { agentRun: RelatedAgentRun; status: string; }; /** * Ask mains to abort a background job's live run. The job row is already * claimed as cancelled by the publisher; only the main holding the * in-process abort handle acts on this. */ 'cancel-agent-background-job': { jobId: string; }; /** Ask main instances to deliver background job results to the parent thread. */ 'wake-agent-background-job': { threadId: string; }; 'clear-test-webhooks': { webhookKey: string; workflowEntity: IWorkflowBase; pushRef: string; }; /** * Relay chat stream events between main instances. * Used when the main handling the workflow execution is different from * the main that has the WebSocket connection to the client. */ 'relay-chat-stream-event': { /** Event type: execution-level or message-level */ eventType: 'execution-begin' | 'execution-end' | 'begin' | 'chunk' | 'end' | 'error'; /** User ID - sends to all user connections */ userId: string; /** Chat session ID */ sessionId: string; /** Message ID being streamed (empty for execution-level events) */ messageId: string; /** Sequence number for ordering (0 for execution-level events) */ sequenceNumber: number; /** Event-specific payload */ payload: { previousMessageId?: string | null; retryOfMessageId?: string | null; executionId?: number | null; content?: string; status?: ChatHubMessageStatus; error?: string; }; }; /** * Relay an Instance AI stream event to sibling mains. * * The agent runs on whichever main received `POST /chat/:threadId` and emits * events into that main's in-process bus, but the client's `GET /events/:threadId` * SSE connection may be held by a different main. The producing main relays each * event; the main holding the SSE subscription re-emits it locally to its client. */ 'relay-instance-ai-event': { threadId: string; /** * Producer-assigned stored event. The id comes from the shared per-thread * sequence, so every main stores and serves identical event ids and the * frontend's replay cursor is valid against any main. With the durable * log enabled ids are DB-assigned seqs and ephemeral events (deltas, * status) carry no id at all: they are live-only. */ storedEvent: { id?: number; event: InstanceAiEvent }; }; /** * Relay an Instance AI task-control action (correction / cancel / clear) to * sibling mains. The action may target a background task or run held * in another main's in-memory state. Broadcast + local-gate: every main applies * it only to its own local slice of the thread. */ 'relay-instance-ai-task-control': { threadId: string; userId?: string; taskId?: string; action: 'correct' | 'cancel-task' | 'cancel-thread' | 'clear-thread'; correction?: string; }; /** * Relay human message events between main instances. * Used for cross-client synchronization when a user sends a message * from one browser window and other windows need to be updated. */ 'relay-chat-human-message': { /** User ID - sends to all user connections */ userId: string; /** Chat session ID */ sessionId: string; /** Human message ID */ messageId: string; /** ID of the previous message in the chain */ previousMessageId: string | null; /** Message content */ content: string; /** Attachments on the message */ attachments: Array<{ id: string; fileName: string; mimeType: string }>; }; /** * Relay message edit events between main instances. * Used for cross-client synchronization when a user edits a message * from one browser window and other windows need to be updated. */ 'relay-chat-message-edit': { /** User ID - sends to all user connections */ userId: string; /** Chat session ID */ sessionId: string; /** ID of the message being revised */ revisionOfMessageId: string; /** ID of this message (the revised version) */ messageId: string; /** New message content */ content: string; /** Attachments on the new message */ attachments: Array<{ id: string; fileName: string; mimeType: string }>; }; // #endregion // #region Evaluation /** * Cancel a test run across all main instances. * Used in multi-main mode to signal cancellation to the instance running the test. */ 'cancel-test-run': { testRunId: string; }; /** * Cancel every running test run inside an evaluation collection across all * main instances. Used when a user cancels a collection — each main checks * its in-flight runs and aborts those that belong to the collection. */ 'cancel-collection': { collectionId: string; }; // #endregion // #region Agents /** * Reconcile a single agent chat integration across main instances. * Published by the main that handled the user's connect/disconnect request * after the change is persisted; every main applies the same connect or * disconnect locally so the in-memory `connections` map stays in sync. */ 'agent-chat-integration-changed': { agentId: string; integration: AgentIntegrationConfig; action: 'connect' | 'disconnect'; }; /** * Ask the current leader to start or stop a leader-only channel runtime * (e.g. Telegram in polling mode) on behalf of the main that handled the * user's request. See `LeaderChannelRelayService` for the mechanism. * * Broadcast to every main rather than targeted at the leader's host: the * leader key can be stale by the time we read it, and leadership can change * between publish and receive. Followers ignore this command, so whoever * holds leadership when it lands is the one that acts. * * `replyTo` carries the requester's host ID because the subscriber emits * only the payload — `senderId` is stripped before handlers see the message. * * A connect carries the whole config: it runs before the config is persisted, * so the leader cannot read it back. A teardown only needs the connection key, * which every caller already holds. */ 'agent-chat-leader-channel-request': { requestId: string; replyTo: string; agentId: string; } & ( | { action: 'connect'; integration: AgentIntegrationConfig } | { action: 'disconnect'; integration: { type: string; credentialId: string } } ); /** * Acknowledge an `agent-chat-leader-channel-request`, targeted at the * requester's host ID. Published by the leader once the operation has fully * completed, so the requesting main can only report success after the * leader's startup or teardown finished. */ 'agent-chat-leader-channel-result': { requestId: string; ok: boolean; error?: string; }; /** * Keep per-main thread subscription state in sync across the cluster. The * originating main persists the subscription change before publishing; peers * update their local in-memory subscription state so load-balanced follow-up * messages can route to `onSubscribedMessage` without re-mentioning the bot. * * Subscriptions are currently backed by the Vercel Chat SDK memory adapter, * but this event describes the intent (thread subscription changed) rather * than that implementation detail. */ 'agent-chat-subscription-changed': { agentId: string; integration: AgentIntegrationConfig; threadId: string; action: 'subscribe' | 'unsubscribe'; }; /** * Drop the cached agent runtime in `AgentsService.runtimes` across mains. * Published by the main that handled an agent mutation (publish, unpublish, * config update, tool/skill change, delete) after the change is persisted. * Every main drops its cache entry so the next request rebuilds the runtime * from the current DB state, picking up the new model/credential/tools/skills. * * Without this, peer mains keep serving webhook traffic from a stale * compiled runtime — including stale embedded credentials — until the * 30-minute TTL evicts the entry. */ 'agent-config-changed': { agentId: string; }; /** * Reconcile an agent's scheduled task cron jobs across main instances. * Published by the main that handled a publish/unpublish/delete after the * change is persisted. Only the leader owns task crons, so every main runs * the same reconcile but registration is a no-op on followers. */ 'agent-tasks-changed': { agentId: string; }; // #endregion // #region Redaction /** * Drop the cached instance redaction floor across main instances. * Published by the main that handled a redaction-floor update after the new * value is persisted; every other main clears its local cache key so the next * read re-loads the current value from the DB. * * Must NOT be added to SELF_SEND_COMMANDS: the originating main already * updates its own cache synchronously in set(). */ 'redaction-floor-changed': never; // #endregion }; export type PubSubWorkerResponseMap = { 'response-to-get-worker-status': WorkerStatus & { requestingUserId: string; }; }; export type PubSubEventMap = PubSubCommandMap & PubSubWorkerResponseMap;