From 413dc384817aba5ad68fd66b723d3fa62db793a0 Mon Sep 17 00:00:00 2001 From: Himanshu Garg Date: Thu, 20 Aug 2026 14:22:14 +0530 Subject: [PATCH 1/7] fix(root): harden release-packages.yml against shell injection (#12398) --- .github/workflows/release-packages.yml | 88 ++++++++++++++++---------- 1 file changed, 56 insertions(+), 32 deletions(-) diff --git a/.github/workflows/release-packages.yml b/.github/workflows/release-packages.yml index 24435054a00..518c6466529 100644 --- a/.github/workflows/release-packages.yml +++ b/.github/workflows/release-packages.yml @@ -82,14 +82,17 @@ jobs: - name: Set version for nightly or rc if: github.event.inputs.release_type != 'stable' + env: + INPUT_RELEASE_TYPE: ${{ github.event.inputs.release_type }} + INPUT_VERSION: ${{ github.event.inputs.version }} run: | COMMIT_SHA=$(git rev-parse --short HEAD) - if [ "${{ github.event.inputs.release_type }}" = "nightly" ]; then + if [ "$INPUT_RELEASE_TYPE" = "nightly" ]; then DATE=$(date +'%Y%m%d') - echo "RELEASE_VERSION=${{ github.event.inputs.version }}-nightly.${DATE}.${COMMIT_SHA}" >> $GITHUB_ENV + echo "RELEASE_VERSION=${INPUT_VERSION}-nightly.${DATE}.${COMMIT_SHA}" >> $GITHUB_ENV echo "Using nightly version: $RELEASE_VERSION" - elif [ "${{ github.event.inputs.release_type }}" = "rc" ]; then - echo "RELEASE_VERSION=${{ github.event.inputs.version }}-rc.${COMMIT_SHA}" >> $GITHUB_ENV + elif [ "$INPUT_RELEASE_TYPE" = "rc" ]; then + echo "RELEASE_VERSION=${INPUT_VERSION}-rc.${COMMIT_SHA}" >> $GITHUB_ENV echo "Using rc version: $RELEASE_VERSION" fi @@ -99,31 +102,41 @@ jobs: git config --global user.name "github-actions[bot]" - name: Release version (without commit) + env: + INPUT_RELEASE_TYPE: ${{ github.event.inputs.release_type }} + INPUT_VERSION: ${{ github.event.inputs.version }} + INPUT_PACKAGES: ${{ github.event.inputs.packages }} run: | - if [ "${{ github.event.inputs.release_type }}" = "nightly" ]; then - echo "Running nightly release with version: ${{ env.RELEASE_VERSION }}" - pnpm nx release version ${{ env.RELEASE_VERSION }} --projects=${{ github.event.inputs.packages }} --preid nightly --git-commit=false --verbose - elif [ "${{ github.event.inputs.release_type }}" = "rc" ]; then - echo "Running rc release with version: ${{ env.RELEASE_VERSION }}" - pnpm nx release version ${{ env.RELEASE_VERSION }} --projects=${{ github.event.inputs.packages }} --preid rc --git-commit=false --verbose + if [ "$INPUT_RELEASE_TYPE" = "nightly" ]; then + echo "Running nightly release with version: $RELEASE_VERSION" + pnpm nx release version "$RELEASE_VERSION" --projects="$INPUT_PACKAGES" --preid nightly --git-commit=false --verbose + elif [ "$INPUT_RELEASE_TYPE" = "rc" ]; then + echo "Running rc release with version: $RELEASE_VERSION" + pnpm nx release version "$RELEASE_VERSION" --projects="$INPUT_PACKAGES" --preid rc --git-commit=false --verbose else - echo "Running stable release with version: ${{ github.event.inputs.version }}" - pnpm nx release version ${{ github.event.inputs.version }} --projects=${{ github.event.inputs.packages }} --git-commit=false --verbose + echo "Running stable release with version: $INPUT_VERSION" + pnpm nx release version "$INPUT_VERSION" --projects="$INPUT_PACKAGES" --git-commit=false --verbose fi - name: Generate changelog (without commit) + env: + INPUT_RELEASE_TYPE: ${{ github.event.inputs.release_type }} + INPUT_VERSION: ${{ github.event.inputs.version }} + INPUT_PACKAGES: ${{ github.event.inputs.packages }} + INPUT_PREVIOUS_TAG: ${{ github.event.inputs.previous_tag }} run: | - if [ "${{ github.event.inputs.release_type }}" = "stable" ]; then - pnpm nx release changelog ${{ github.event.inputs.version }} --projects=${{ github.event.inputs.packages }} --from=${{ github.event.inputs.previous_tag }} --git-commit=false + if [ "$INPUT_RELEASE_TYPE" = "stable" ]; then + pnpm nx release changelog "$INPUT_VERSION" --projects="$INPUT_PACKAGES" --from="$INPUT_PREVIOUS_TAG" --git-commit=false else - pnpm nx release changelog ${{ env.RELEASE_VERSION }} --projects=${{ github.event.inputs.packages }} --from=${{ github.event.inputs.previous_tag }} --git-commit=false + pnpm nx release changelog "$RELEASE_VERSION" --projects="$INPUT_PACKAGES" --from="$INPUT_PREVIOUS_TAG" --git-commit=false fi - name: Build packages env: NX_NO_CLOUD: ${{ secrets.NX_CLOUD_ACCESS_TOKEN == '' && 'true' || 'false' }} + INPUT_RELEASE_TYPE: ${{ github.event.inputs.release_type }} run: | - if [ "${{ github.event.inputs.release_type }}" != "stable" ]; then + if [ "$INPUT_RELEASE_TYPE" != "stable" ]; then pnpm run build:packages else echo "Skipping build for stable release (will be built after PR merge)" @@ -131,13 +144,16 @@ jobs: - name: Publish packages (nightly and rc only) if: github.event.inputs.release_type != 'stable' + env: + INPUT_RELEASE_TYPE: ${{ github.event.inputs.release_type }} + INPUT_PACKAGES: ${{ github.event.inputs.packages }} run: | - if [ "${{ github.event.inputs.release_type }}" = "nightly" ]; then + if [ "$INPUT_RELEASE_TYPE" = "nightly" ]; then echo "πŸ“¦ Publishing nightly release with OIDC trusted publishing..." - pnpm nx run-many -t nx-release-publish --projects=${{ github.event.inputs.packages }} -- --tag=nightly --provenance - elif [ "${{ github.event.inputs.release_type }}" = "rc" ]; then + pnpm nx run-many -t nx-release-publish --projects="$INPUT_PACKAGES" -- --tag=nightly --provenance + elif [ "$INPUT_RELEASE_TYPE" = "rc" ]; then echo "πŸ“¦ Publishing RC release with OIDC trusted publishing..." - pnpm nx run-many -t nx-release-publish --projects=${{ github.event.inputs.packages }} -- --tag=next --provenance + pnpm nx run-many -t nx-release-publish --projects="$INPUT_PACKAGES" -- --tag=next --provenance fi - name: Create Pull Request if: github.event.inputs.release_type == 'stable' @@ -187,9 +203,9 @@ jobs: - name: Extract version and packages from PR id: extract-info + env: + PR_TITLE: ${{ github.event.pull_request.title }} run: | - PR_TITLE='${{ github.event.pull_request.title }}' - echo "=== PR Title ===" echo "$PR_TITLE" echo "" @@ -243,20 +259,25 @@ jobs: run: pnpm run build:packages - name: Publish packages to NPM + env: + RELEASE_PACKAGES: ${{ steps.extract-info.outputs.packages }} run: | echo "πŸ“¦ Publishing stable release with OIDC trusted publishing..." - pnpm nx run-many -t nx-release-publish --projects=${{ steps.extract-info.outputs.packages }} -- --tag=latest --provenance + pnpm nx run-many -t nx-release-publish --projects="$RELEASE_PACKAGES" -- --tag=latest --provenance - name: Create and push tags + env: + RELEASE_PACKAGES: ${{ steps.extract-info.outputs.packages }} + RELEASE_VERSION: ${{ steps.extract-info.outputs.version }} run: | git config --global user.email "github-actions[bot]@users.noreply.github.com" git config --global user.name "github-actions[bot]" - echo "Creating tags for packages: ${{ steps.extract-info.outputs.packages }}" - IFS=',' read -ra PACKAGES <<< "${{ steps.extract-info.outputs.packages }}" + echo "Creating tags for packages: $RELEASE_PACKAGES" + IFS=',' read -ra PACKAGES <<< "$RELEASE_PACKAGES" for package in "${PACKAGES[@]}"; do package=$(echo "$package" | sed 's/^[[:space:]]*//;s/[[:space:]]*$//') - tag_name="${package}@${{ steps.extract-info.outputs.version }}" + tag_name="${package}@${RELEASE_VERSION}" echo "Creating tag: $tag_name" git tag "$tag_name" echo "Pushing tag: $tag_name" @@ -266,18 +287,21 @@ jobs: - name: Create GitHub Releases env: GH_TOKEN: ${{ secrets.GITHUB_TOKEN }} + RELEASE_PACKAGES: ${{ steps.extract-info.outputs.packages }} + RELEASE_VERSION: ${{ steps.extract-info.outputs.version }} + PR_BASE_REF: ${{ github.event.pull_request.base.ref }} run: | - echo "Creating GitHub releases for packages: ${{ steps.extract-info.outputs.packages }}" + echo "Creating GitHub releases for packages: $RELEASE_PACKAGES" - IFS=',' read -ra PACKAGES <<< "${{ steps.extract-info.outputs.packages }}" + IFS=',' read -ra PACKAGES <<< "$RELEASE_PACKAGES" for package in "${PACKAGES[@]}"; do package=$(echo "$package" | sed 's/^[[:space:]]*//;s/[[:space:]]*$//') - tag_name="${package}@${{ steps.extract-info.outputs.version }}" + tag_name="${package}@${RELEASE_VERSION}" echo "Creating GitHub release for $package with tag: $tag_name" package_name=$(echo $package | sed 's/@novu\///') - release_body="Release of ${package} version ${{ steps.extract-info.outputs.version }}" + release_body="Release of ${package} version ${RELEASE_VERSION}" possible_paths=( "packages/${package_name}/CHANGELOG.md" @@ -288,7 +312,7 @@ jobs: for changelog_path in "${possible_paths[@]}"; do if [ -f "$changelog_path" ]; then echo "Found changelog at: $changelog_path" - changelog_content=$(awk -v version="${{ steps.extract-info.outputs.version }}" ' + changelog_content=$(awk -v version="$RELEASE_VERSION" ' BEGIN { capture=0 } /^## / || /^# / || /^v[0-9]+\./ { if (capture) exit @@ -310,7 +334,7 @@ jobs: gh release create "$tag_name" \ --title "$tag_name" \ --notes "$release_body" \ - --target ${{ github.event.pull_request.base.ref }} + --target "$PR_BASE_REF" echo "βœ… Created GitHub release for $tag_name" done From 8c1c9e193c90e257197315b35a43f492d85d0084 Mon Sep 17 00:00:00 2001 From: Nikita Grossman <98285260+nikitagrossman@users.noreply.github.com> Date: Thu, 20 Aug 2026 11:55:59 +0300 Subject: [PATCH 2/7] feat(api-service,worker): MS Teams workflow-origin DM catch-up hydration fixes NV-8591 (#12381) Co-authored-by: Cursor --- .../ingress/inbound-turn.handler.ts | 2 + .../ingress/workflow-origin.helpers.ts | 35 ++- .../ingress/workflow-origin.service.spec.ts | 273 +++++++++++++++++- .../ingress/workflow-origin.service.ts | 16 +- .../send-message-chat.usecase.spec.ts | 86 ++++-- .../send-message/send-message-chat.usecase.ts | 6 +- 6 files changed, 395 insertions(+), 23 deletions(-) diff --git a/apps/api/src/app/agents/conversation-runtime/ingress/inbound-turn.handler.ts b/apps/api/src/app/agents/conversation-runtime/ingress/inbound-turn.handler.ts index f484e528081..6c7e24e5884 100644 --- a/apps/api/src/app/agents/conversation-runtime/ingress/inbound-turn.handler.ts +++ b/apps/api/src/app/agents/conversation-runtime/ingress/inbound-turn.handler.ts @@ -436,6 +436,7 @@ export class AgentInboundHandler implements OnModuleInit { subscriberId, message, existingConversation, + isDirectMessage: thread.isDM, }); const conversation = await this.conversationService.createOrGetConversation({ @@ -1154,6 +1155,7 @@ export class AgentInboundHandler implements OnModuleInit { subscriberId, message: null, existingConversation, + isDirectMessage: thread.isDM, }); const conversation = await this.conversationService.createOrGetConversation({ diff --git a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.helpers.ts b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.helpers.ts index 65c8cecb66a..6bc8c98dfff 100644 --- a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.helpers.ts +++ b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.helpers.ts @@ -11,6 +11,7 @@ export const RECHECK_WORKFLOW_ORIGIN_PLATFORMS: ReadonlySet = AgentPlatformEnum.WHATSAPP, AgentPlatformEnum.TELEGRAM, AgentPlatformEnum.SENDBLUE, + AgentPlatformEnum.TEAMS, ]); export function buildWorkflowOriginSummary( @@ -76,6 +77,31 @@ export function extractTelegramQuotedMessageId(message: Message | null): string return typeof quotedId === 'string' && quotedId.length > 0 ? quotedId : null; } +/** Teams quote-reply activity id from `message.raw.entities[].quotedReply.messageId`, else `replyToId`. */ +export function extractTeamsQuotedActivityId(message: Message | null): string | null { + if (!message) { + return null; + } + + const raw = asRecord(message.raw); + const entities = Array.isArray(raw?.entities) ? raw.entities : []; + for (const entity of entities) { + const record = asRecord(entity); + if (record?.type !== 'quotedReply') { + continue; + } + + const messageId = asRecord(record.quotedReply)?.messageId; + if (typeof messageId === 'string' && messageId.length > 0) { + return messageId; + } + } + + const replyToId = raw?.replyToId; + + return typeof replyToId === 'string' && replyToId.length > 0 ? replyToId : null; +} + /** * Bare chat id from `telegram:{chatId}` or `telegram:{chatId}:{messageThreadId}` (forum topics). * Returns null when the prefix is absent or the segment is empty. @@ -101,7 +127,7 @@ export function isSendblueDirectThreadId(platformThreadId: string): boolean { return segments.length === 3 && segments[0] === 'sendblue' && segments[1].length > 0 && segments[2].length > 0; } -/** Email β†’ Message._id; WhatsApp β†’ wamid; Sendblue β†’ message_handle; Slack `{channel}:{ts}` β†’ `ts`; Telegram β†’ `{chatId}:{message_id}`. */ +/** Email β†’ Message._id; WhatsApp β†’ wamid; Sendblue β†’ message_handle; Teams β†’ activity id; Slack `{channel}:{ts}` β†’ `ts`; Telegram β†’ `{chatId}:{message_id}`. */ export function resolvePlatformMessageId( platform: AgentPlatformEnum, originMessage: MessageEntity, @@ -115,7 +141,11 @@ export function resolvePlatformMessageId( return undefined; } - if (platform === AgentPlatformEnum.WHATSAPP || platform === AgentPlatformEnum.SENDBLUE) { + if ( + platform === AgentPlatformEnum.WHATSAPP || + platform === AgentPlatformEnum.SENDBLUE || + platform === AgentPlatformEnum.TEAMS + ) { return originMessage.identifier; } @@ -132,6 +162,7 @@ export function resolvePlatformMessageId( return `${chatId}:${originMessage.identifier}`; } + // Slack-only: Message.identifier is `{channel}:{ts}`; hydration keys off the bare `ts`. const colon = originMessage.identifier.indexOf(':'); if (colon <= 0 || colon === originMessage.identifier.length - 1) { return undefined; diff --git a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.spec.ts b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.spec.ts index fac5baf3966..10b348257af 100644 --- a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.spec.ts +++ b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.spec.ts @@ -1,4 +1,4 @@ -import { ChannelTypeEnum } from '@novu/shared'; +import { ChannelTypeEnum, ENDPOINT_TYPES } from '@novu/shared'; import { expect } from 'chai'; import sinon from 'sinon'; import { AgentPlatformEnum } from '../../shared/enums/agent-platform.enum'; @@ -993,4 +993,275 @@ describe('WorkflowOriginService', () => { expect(result?.origin.identifier).to.equal(sendblueOrigin.identifier); }); }); + + describe('Teams', () => { + const teamsDmThreadId = 'teams:YToxMjM:aHR0cHM6Ly9zbWJhLnRyYWZmaWNtYW5hZ2VyLm5ldC90ZWFtcw'; + const teamsConfig = { + environmentId: 'env1', + organizationId: 'org1', + platform: AgentPlatformEnum.TEAMS, + agentIdentifier: 'support-agent', + providerId: 'msteams', + }; + + const teamsOrigin = { + _id: 'teams-msg1', + _notificationId: 'teams-notif1', + _jobId: 'teams-job1', + transactionId: 'teams-txn1', + templateIdentifier: 'order-alerts', + stepId: 'chat-1', + content: 'Your order shipped', + identifier: 'activity-abc123', + createdAt: new Date(Date.now() - 60_000).toISOString(), + providerId: 'msteams', + }; + + it('hydrates the latest ms_teams_user origin on a new DM conversation', async () => { + const { service, messageRepository } = makeService({ + find: sinon.stub().resolves([teamsOrigin]), + }); + + const before = Date.now(); + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { id: 'inbound', text: 'hi', author: { userId: '29:user1' }, raw: {} } as any, + existingConversation: null, + isDirectMessage: true, + }); + const after = Date.now(); + + expect(messageRepository.find.calledOnce).to.equal(true); + const [query, , options] = messageRepository.find.firstCall.args; + expect(query).to.include({ + _environmentId: 'env1', + _agentId: 'agent1', + _subscriberId: 'subscriber-mongo-1', + providerId: 'msteams', + channel: ChannelTypeEnum.CHAT, + 'channelData.type': ENDPOINT_TYPES.MS_TEAMS_USER, + }); + expect(query._notificationId).to.deep.equal({ $exists: true, $ne: null }); + expect(query.createdAt.$gt.getTime()).to.be.at.least(before - WORKFLOW_ORIGIN_LOOKBACK_MS - 5); + expect(query.createdAt.$gt.getTime()).to.be.at.most(after - WORKFLOW_ORIGIN_LOOKBACK_MS + 5); + expect(options).to.deep.equal({ sort: { createdAt: -1 }, limit: 1 }); + expect(result).to.deep.equal({ origin: teamsOrigin, notificationId: 'teams-notif1' }); + }); + + it('prefers a quotedReply entity messageId over latest-by-subscriber and bypasses the catch-up window', async () => { + const lastActivityAt = new Date().toISOString(); + const existingConversation = { + ...conversation, + lastActivityAt, + participants: [{ type: 'subscriber', id: 'sub1' }], + }; + const { service, messageRepository } = makeService({ + findOne: sinon.stub().resolves({ + ...teamsOrigin, + createdAt: new Date(Date.now() - 86_400_000).toISOString(), + }), + }); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { + id: 'inbound', + text: 'hi', + author: { userId: '29:user1' }, + raw: { + entities: [ + { type: 'clientInfo' }, + { type: 'quotedReply', quotedReply: { messageId: teamsOrigin.identifier } }, + ], + }, + } as any, + existingConversation: existingConversation as any, + isDirectMessage: true, + }); + + expect(messageRepository.findOne.calledOnce).to.equal(true); + expect(messageRepository.findOne.firstCall.args[0]).to.deep.equal({ + _environmentId: 'env1', + _agentId: 'agent1', + _subscriberId: 'subscriber-mongo-1', + providerId: 'msteams', + channel: ChannelTypeEnum.CHAT, + _notificationId: { $exists: true, $ne: null }, + 'channelData.type': ENDPOINT_TYPES.MS_TEAMS_USER, + identifier: teamsOrigin.identifier, + }); + expect(messageRepository.find.called).to.equal(false); + expect(result?.origin.identifier).to.equal(teamsOrigin.identifier); + expect(result?.notificationId).to.equal(undefined); + }); + + it('falls back to replyToId when no quotedReply entity is present', async () => { + const { service, messageRepository } = makeService({ + findOne: sinon.stub().resolves(teamsOrigin), + }); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { + id: 'inbound', + text: 'hi', + author: { userId: '29:user1' }, + raw: { replyToId: teamsOrigin.identifier }, + } as any, + existingConversation: null, + isDirectMessage: true, + }); + + expect(messageRepository.findOne.calledOnce).to.equal(true); + expect(messageRepository.findOne.firstCall.args[0].identifier).to.equal(teamsOrigin.identifier); + expect(messageRepository.find.called).to.equal(false); + expect(result?.origin.identifier).to.equal(teamsOrigin.identifier); + }); + + it('catch-up hydrates an existing DM conversation when the origin is not hydrated yet', async () => { + const existingConversation = { + ...conversation, + lastActivityAt: new Date(Date.now() - 3_600_000).toISOString(), + participants: [{ type: 'subscriber', id: 'sub1' }], + }; + const { service, messageRepository, conversationService } = makeService({ + find: sinon.stub().resolves([teamsOrigin]), + }); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { id: 'inbound', text: 'hi', author: { userId: '29:user1' }, raw: {} } as any, + existingConversation: existingConversation as any, + isDirectMessage: true, + }); + + expect(messageRepository.find.calledOnce).to.equal(true); + expect(conversationService.isWorkflowOriginHydrated.firstCall.args).to.deep.equal([ + 'env1', + 'conversation1', + teamsOrigin.identifier, + ]); + expect(result).to.deep.equal({ origin: teamsOrigin, notificationId: undefined }); + }); + + it('skips an origin already hydrated into the conversation', async () => { + const existingConversation = { + ...conversation, + lastActivityAt: new Date(Date.now() - 3_600_000).toISOString(), + participants: [{ type: 'subscriber', id: 'sub1' }], + }; + const { service } = makeService({ + find: sinon.stub().resolves([teamsOrigin]), + isWorkflowOriginHydrated: sinon.stub().resolves(true), + }); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { id: 'inbound', text: 'hi', author: { userId: '29:user1' }, raw: {} } as any, + existingConversation: existingConversation as any, + isDirectMessage: true, + }); + + expect(result).to.equal(null); + }); + + it('fails closed for non-DM turns without looking up an origin', async () => { + const { service, messageRepository } = makeService({ + find: sinon.stub().resolves([teamsOrigin]), + }); + + for (const isDirectMessage of [false, undefined] as const) { + messageRepository.find.resetHistory(); + messageRepository.findOne.resetHistory(); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: { id: 'inbound', text: 'hi', author: { userId: '29:user1' }, raw: {} } as any, + existingConversation: null, + ...(isDirectMessage === undefined ? {} : { isDirectMessage }), + }); + + expect(result).to.equal(null); + expect(messageRepository.find.called).to.equal(false); + expect(messageRepository.findOne.called).to.equal(false); + } + }); + + it('resolves on action turns without a message using latest-by-subscriber', async () => { + const { service, messageRepository } = makeService({ + find: sinon.stub().resolves([teamsOrigin]), + }); + + const result = await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'sub1', + message: null, + existingConversation: null, + isDirectMessage: true, + }); + + expect(messageRepository.find.calledOnce).to.equal(true); + expect(result?.notificationId).to.equal('teams-notif1'); + }); + + it('scopes the origin lookup to the resolved subscriber', async () => { + const { service, messageRepository } = makeService({ + findBySubscriberId: sinon.stub().resolves({ _id: 'attacker-mongo' }), + find: sinon.stub().resolves([]), + }); + + await service.resolve({ + agentId: 'agent1', + config: teamsConfig as any, + platformThreadId: teamsDmThreadId, + subscriberId: 'attacker-subscriber', + message: { id: 'inbound', text: 'hi', author: { userId: '29:user1' }, raw: {} } as any, + existingConversation: null, + isDirectMessage: true, + }); + + expect(messageRepository.find.firstCall.args[0]._subscriberId).to.equal('attacker-mongo'); + }); + + it('hydrates using the bare activity id as platformMessageId', async () => { + const { service, conversationService } = makeService({ + notificationFindOne: sinon.stub().resolves({ payload: { orderId: 'ORD-9' } }), + }); + + await service.hydrate({ + agentId: 'agent1', + config: teamsConfig as any, + conversation: conversation as any, + platformThreadId: teamsDmThreadId, + origin: teamsOrigin as any, + }); + + expect(conversationService.persistWorkflowOriginHydration.firstCall.args[0].platformMessageId).to.equal( + teamsOrigin.identifier + ); + expect( + conversationService.persistWorkflowOriginHydration.firstCall.args[0].signalData.workflowIdentifier + ).to.equal('order-alerts'); + }); + }); }); diff --git a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.ts b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.ts index c70f90d2623..a5403a9320b 100644 --- a/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.ts +++ b/apps/api/src/app/agents/conversation-runtime/ingress/workflow-origin.service.ts @@ -8,7 +8,7 @@ import { NotificationRepository, SubscriberRepository, } from '@novu/dal'; -import { ChannelTypeEnum } from '@novu/shared'; +import { ChannelTypeEnum, ENDPOINT_TYPES } from '@novu/shared'; import type { Message } from 'chat'; import { ResolvedAgentConfig } from '../../channels/agent-config-resolver.service'; import { AgentPlatformEnum } from '../../shared/enums/agent-platform.enum'; @@ -19,6 +19,7 @@ import { extractAgentEmailOriginToken, extractTelegramChatIdFromThreadId, extractTelegramQuotedMessageId, + extractTeamsQuotedActivityId, extractWhatsAppQuotedWamid, isSendblueDirectThreadId, RECHECK_WORKFLOW_ORIGIN_PLATFORMS, @@ -50,8 +51,9 @@ export class WorkflowOriginService { subscriberId: string | null; message: Message | null; existingConversation: ConversationEntity | null; + isDirectMessage?: boolean; }): Promise { - const { agentId, config, platformThreadId, subscriberId, message, existingConversation } = params; + const { agentId, config, platformThreadId, subscriberId, message, existingConversation, isDirectMessage } = params; if (!subscriberId) { return null; @@ -112,6 +114,16 @@ export class WorkflowOriginService { : null; break; case AgentPlatformEnum.TEAMS: + origin = isDirectMessage + ? await this.findRecentChatWorkflowOriginMessage( + agentId, + config, + subscriber._id, + extractTeamsQuotedActivityId(message), + { 'channelData.type': ENDPOINT_TYPES.MS_TEAMS_USER } + ) + : null; + break; case AgentPlatformEnum.AGENT_CHAT: break; default: { diff --git a/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.spec.ts b/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.spec.ts index 92baa697926..08c1bcb0489 100644 --- a/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.spec.ts +++ b/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.spec.ts @@ -336,17 +336,28 @@ describe('SendMessageChat - agent assigned path', () => { sinon.restore(); }); - function buildAgentUsecase(options: { - channelData?: unknown[]; - linked?: boolean; - chatHandlerSend?: sinon.SinonStub; - updateMessage?: sinon.SinonStub; - jobAgentId?: string | null; - workflowAgentIdentifier?: string; - } = {}) { + function buildAgentUsecase( + options: { + channelData?: unknown[]; + linked?: boolean; + linkedRefs?: Array<{ identifier: string; providerId: string }>; + integration?: { + _id: string; + identifier: string; + providerId: ChatProviderIdEnum; + channel: ChannelTypeEnum; + credentials: Record; + }; + chatHandlerSend?: sinon.SinonStub; + updateMessage?: sinon.SinonStub; + jobAgentId?: string | null; + workflowAgentIdentifier?: string; + } = {} + ) { const { channelData = [slackUserData], linked = true, + integration = slackIntegration, chatHandlerSend = sinon.stub().resolves({ id: 'D123:1777837477.371619', date: new Date().toISOString(), @@ -355,6 +366,9 @@ describe('SendMessageChat - agent assigned path', () => { jobAgentId, workflowAgentIdentifier, } = options; + const linkedRefs = + options.linkedRefs ?? + (linked ? [{ identifier: integration.identifier, providerId: integration.providerId }] : []); sinon.stub(ChatFactory.prototype, 'getHandler').returns({ send: chatHandlerSend, @@ -365,9 +379,7 @@ describe('SendMessageChat - agent assigned path', () => { findOne: sinon.stub().resolves(workflowAgentIdentifier ? { _id: 'agent_from_workflow' } : null), }; const agentIntegrationRepository = { - listLinkedIntegrationRefs: sinon - .stub() - .resolves(linked ? [{ identifier: 'slack-main', providerId: ChatProviderIdEnum.Slack }] : []), + listLinkedIntegrationRefs: sinon.stub().resolves(linkedRefs), }; const createExecutionDetails = { execute: sinon.stub().resolves(undefined) }; const sendWebhookMessage = { execute: sinon.stub().resolves(undefined) }; @@ -377,7 +389,7 @@ describe('SendMessageChat - agent assigned path', () => { update: updateMessage, }; const selectIntegration = { - execute: sinon.stub().resolves(slackIntegration), + execute: sinon.stub().resolves(integration), }; const featureFlagsService = { getFlag: sinon.stub().resolves(false), @@ -400,8 +412,8 @@ describe('SendMessageChat - agent assigned path', () => { { execute: sinon.stub().resolves([ { - integrationIdentifier: 'slack-main', - providerId: ChatProviderIdEnum.Slack, + integrationIdentifier: integration.identifier, + providerId: integration.providerId, channelData, }, ]), @@ -454,9 +466,7 @@ describe('SendMessageChat - agent assigned path', () => { content: 'agent hello', }, } as never, - workflow: workflowAgentIdentifier - ? ({ agent: { identifier: workflowAgentIdentifier } } as never) - : undefined, + workflow: workflowAgentIdentifier ? ({ agent: { identifier: workflowAgentIdentifier } } as never) : undefined, job: { _id: 'job_1', _environmentId: 'env_1', @@ -620,6 +630,48 @@ describe('SendMessageChat - agent assigned path', () => { }); }); + it('routes agent-assigned MS Teams user endpoints without the chat-channels fallback', async () => { + const teamsActivityId = 'activity-abc123'; + const chatHandlerSend = sinon.stub().resolves({ id: teamsActivityId, date: new Date().toISOString() }); + const { usecase, updateMessage, createExecutionDetails } = buildAgentUsecase({ + jobAgentId: 'agent_1', + chatHandlerSend, + integration: { + _id: 'integration_teams', + identifier: 'msteams-main', + providerId: ChatProviderIdEnum.MsTeams, + channel: ChannelTypeEnum.CHAT, + credentials: {}, + }, + channelData: [ + { + type: ENDPOINT_TYPES.MS_TEAMS_USER, + identifier: 'ep_teams_user', + token: 'bot-framework-token', + endpoint: { userId: '29:user1' }, + subscriberTenantId: 'tenant-1', + clientId: 'client-1', + }, + ], + }); + + const result = await usecase.execute(buildAgentCommand({ jobAgentId: 'agent_1' })); + + expect(result.status).to.equal(SendMessageStatus.SUCCESS); + sinon.assert.calledOnce(chatHandlerSend); + sinon.assert.calledOnce(updateMessage); + expect(updateMessage.firstCall.args[1]).to.deep.equal({ + $set: { + identifier: teamsActivityId, + _agentId: 'agent_1', + }, + }); + const fallbackWarning = createExecutionDetails.execute + .getCalls() + .find((call) => call.args[0]?.detail === DetailEnum.CHAT_AGENT_CHANNELS_FALLBACK); + expect(fallbackWarning).to.equal(undefined); + }); + it('resolves workflow.agent when job._agentId is unset and stamps with that agent', async () => { const { usecase, chatHandlerSend, updateMessage } = buildAgentUsecase({ workflowAgentIdentifier: 'support-agent', diff --git a/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.ts b/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.ts index 1c1bcaf7cab..5f5419c6130 100644 --- a/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.ts +++ b/apps/worker/src/app/workflow/usecases/send-message/send-message-chat.usecase.ts @@ -89,7 +89,11 @@ type MessageContext = { assignedAgentId: string | null; }; -const AGENT_SUPPORTED_ENDPOINT_TYPES = new Set([ENDPOINT_TYPES.SLACK_USER, ENDPOINT_TYPES.SLACK_CHANNEL]); +const AGENT_SUPPORTED_ENDPOINT_TYPES = new Set([ + ENDPOINT_TYPES.SLACK_USER, + ENDPOINT_TYPES.SLACK_CHANNEL, + ENDPOINT_TYPES.MS_TEAMS_USER, +]); function filterAgentSupportedEndpoints(endpoints: ChannelData[]): ChannelData[] { return endpoints.filter((endpoint) => { From 6921b223372a8daee9572640b1399afcb4b573d6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Tymczuk?= Date: Thu, 20 Aug 2026 12:15:44 +0200 Subject: [PATCH 3/7] fix(dashboard): agents list page - single row name click and copy identifier button (#12368) --- .../src/components/agents/agents-table.tsx | 22 +++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/apps/dashboard/src/components/agents/agents-table.tsx b/apps/dashboard/src/components/agents/agents-table.tsx index 07b796db369..3141e7a6c60 100644 --- a/apps/dashboard/src/components/agents/agents-table.tsx +++ b/apps/dashboard/src/components/agents/agents-table.tsx @@ -16,6 +16,7 @@ import type { AgentResponse } from '@/api/agents'; import { ExceedsPlanIndicator } from '@/components/agents/exceeds-plan-indicator'; import { ProviderIcon } from '@/components/integrations/components/provider-icon'; import { CompactButton } from '@/components/primitives/button-compact'; +import { CopyButton } from '@/components/primitives/copy-button'; import { DropdownMenu, DropdownMenuContent, @@ -270,12 +271,25 @@ export function AgentsTable({
- + {agent.name} - - {agent.identifier} - +
+ + {agent.identifier} + + +
From d05484daef20b04ab09d756c8e60ccdf700536d6 Mon Sep 17 00:00:00 2001 From: Dima Grossman Date: Thu, 20 Aug 2026 14:54:40 +0300 Subject: [PATCH 4/7] fix(dashboard): silence Sanity changelog fetch error toasts fixes NV-8615 (#12399) Co-authored-by: Cursor Agent --- .../components/side-navigation/changelog-cards.tsx | 14 +++++++++++--- apps/dashboard/src/hooks/use-agent-templates.ts | 12 ++++++++++-- 2 files changed, 21 insertions(+), 5 deletions(-) diff --git a/apps/dashboard/src/components/side-navigation/changelog-cards.tsx b/apps/dashboard/src/components/side-navigation/changelog-cards.tsx index 135c1d13cb7..8599d5b7865 100644 --- a/apps/dashboard/src/components/side-navigation/changelog-cards.tsx +++ b/apps/dashboard/src/components/side-navigation/changelog-cards.tsx @@ -141,10 +141,16 @@ export function ChangelogStack() { } `; - const sanityPosts = await fetchSanity(query); + try { + const sanityPosts = await fetchSanity(query); + const transformedData = transformSanityData(sanityPosts ?? []); - const transformedData = transformSanityData(sanityPosts ?? []); - return filterChangelogs(transformedData, getDismissedChangelogs()); + return filterChangelogs(transformedData, getDismissedChangelogs()); + } catch (error) { + console.error('Failed to fetch changelogs from Sanity', error); + + return []; + } }; const { data: changelogs = [] } = useQuery({ @@ -152,6 +158,8 @@ export function ChangelogStack() { queryFn: fetchChangelogs, // Refetch every hour to ensure users see new changelogs staleTime: 60 * 60 * 1000, + // Non-critical marketing content β€” never surface Sanity outages to users + meta: { showError: false }, }); const handleChangelogClick = async (changelog: Changelog) => { diff --git a/apps/dashboard/src/hooks/use-agent-templates.ts b/apps/dashboard/src/hooks/use-agent-templates.ts index 5840993d061..842f8ed77a9 100644 --- a/apps/dashboard/src/hooks/use-agent-templates.ts +++ b/apps/dashboard/src/hooks/use-agent-templates.ts @@ -59,12 +59,20 @@ export function useAgentTemplates() { const { data, isLoading } = useQuery({ queryKey: QUERY_KEY, queryFn: async ({ signal }) => { - const result = await fetchSanity(agentTemplatesQuery, { signal }); + try { + const result = await fetchSanity(agentTemplatesQuery, { signal }); - return (result ?? []).map(mapSanityTemplate).filter((template): template is AgentTemplate => template !== null); + return (result ?? []).map(mapSanityTemplate).filter((template): template is AgentTemplate => template !== null); + } catch (error) { + console.error('Failed to fetch agent templates from Sanity', error); + + return []; + } }, // Templates rarely change β€” keep them fresh for an hour like the changelog query. staleTime: 60 * 60 * 1000, + // Non-critical Sanity content β€” never surface fetch failures as toasts + meta: { showError: false }, }); const templates = data && data.length > 0 ? data : AGENT_TEMPLATES; From 8e3ab6b09f02a798622a2c97e94efdc4ad733f01 Mon Sep 17 00:00:00 2001 From: Dima Grossman Date: Thu, 20 Aug 2026 14:56:50 +0300 Subject: [PATCH 5/7] feat(api-service): add topic.data metafield on topics API fixes NV-8462 (#12158) Co-authored-by: Cursor Agent Co-authored-by: Pawan Jain --- .../create-subscriptions-response.dto.ts | 9 ++ .../create-subscriptions.usecase.ts | 3 + .../update-subscription.usecase.ts | 1 + .../topics-v2/dtos/create-update-topic.dto.ts | 21 ++++- ...delete-topic-subscriptions-response.dto.ts | 9 ++ .../app/topics-v2/dtos/topic-response.dto.ts | 13 ++- .../app/topics-v2/dtos/update-topic.dto.ts | 28 +++++-- .../src/app/topics-v2/e2e/update-topic.e2e.ts | 38 ++++++++- .../src/app/topics-v2/e2e/upsert-topic.e2e.ts | 62 +++++++++++++- .../src/app/topics-v2/topics.controller.ts | 4 +- .../delete-topic-subscriptions.usecase.ts | 1 + .../list-topics/map-topic-entity-to.dto.ts | 1 + .../update-topic/update-topic.command.ts | 16 +++- .../update-topic/update-topic.usecase.ts | 26 +++++- .../upsert-topic/upsert-topic.command.ts | 12 ++- .../upsert-topic/upsert-topic.usecase.ts | 83 ++++++++++++++----- docs/api-reference/topics/topic-schema.mdx | 1 + .../src/decorators/index.ts | 1 + .../is-valid-custom-data.decorator.ts | 62 ++++++++++++++ .../src/repositories/topic/topic.entity.ts | 2 + .../repositories/topic/topic.repository.ts | 4 +- .../src/repositories/topic/topic.schema.ts | 1 + packages/shared/src/types/topic.ts | 3 + 23 files changed, 355 insertions(+), 46 deletions(-) create mode 100644 libs/application-generic/src/decorators/is-valid-custom-data.decorator.ts diff --git a/apps/api/src/app/shared/dtos/subscriptions/create-subscriptions-response.dto.ts b/apps/api/src/app/shared/dtos/subscriptions/create-subscriptions-response.dto.ts index caf1fdc5d92..eba289b65b8 100644 --- a/apps/api/src/app/shared/dtos/subscriptions/create-subscriptions-response.dto.ts +++ b/apps/api/src/app/shared/dtos/subscriptions/create-subscriptions-response.dto.ts @@ -26,6 +26,15 @@ export class TopicDto { @IsString() @IsOptional() name?: string; + + @ApiPropertyOptional({ + description: 'Additional custom data associated with the topic', + type: Object, + additionalProperties: true, + example: { category: 'product', priority: 1 }, + }) + @IsOptional() + data?: Record; } export class SubscriberDto { diff --git a/apps/api/src/app/subscriptions/usecases/create-subscriptions/create-subscriptions.usecase.ts b/apps/api/src/app/subscriptions/usecases/create-subscriptions/create-subscriptions.usecase.ts index b22db145073..1c865f60275 100644 --- a/apps/api/src/app/subscriptions/usecases/create-subscriptions/create-subscriptions.usecase.ts +++ b/apps/api/src/app/subscriptions/usecases/create-subscriptions/create-subscriptions.usecase.ts @@ -171,6 +171,7 @@ export class CreateSubscriptionsUsecase { _id: topic._id, key: topic.key, name: topic.name, + data: topic.data, }, subscriber: subscriber ? { @@ -238,6 +239,7 @@ export class CreateSubscriptionsUsecase { _id: topic._id, key: topic.key, name: topic.name, + data: topic.data, }, subscriber: subscriber ? { @@ -276,6 +278,7 @@ export class CreateSubscriptionsUsecase { _id: topic._id, key: topic.key, name: topic.name, + data: topic.data, }, subscriber: subscriber ? { diff --git a/apps/api/src/app/subscriptions/usecases/update-subscription/update-subscription.usecase.ts b/apps/api/src/app/subscriptions/usecases/update-subscription/update-subscription.usecase.ts index 287b8905c88..ca537734553 100644 --- a/apps/api/src/app/subscriptions/usecases/update-subscription/update-subscription.usecase.ts +++ b/apps/api/src/app/subscriptions/usecases/update-subscription/update-subscription.usecase.ts @@ -386,6 +386,7 @@ export class UpdateSubscriptionUsecase { _id: topic._id, key: topic.key, name: topic.name, + data: topic.data, }, subscriber: subscriber ? { diff --git a/apps/api/src/app/topics-v2/dtos/create-update-topic.dto.ts b/apps/api/src/app/topics-v2/dtos/create-update-topic.dto.ts index c6985087b74..c4cf475e316 100644 --- a/apps/api/src/app/topics-v2/dtos/create-update-topic.dto.ts +++ b/apps/api/src/app/topics-v2/dtos/create-update-topic.dto.ts @@ -1,5 +1,7 @@ import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; -import { IsNotEmpty, IsOptional, IsString, Length } from 'class-validator'; +import { IsValidContextData, IsValidCustomData } from '@novu/application-generic'; +import { TopicCustomData } from '@novu/shared'; +import { IsNotEmpty, IsObject, IsOptional, IsString, Length, ValidateIf } from 'class-validator'; export class CreateUpdateTopicRequestDto { @ApiProperty({ @@ -19,5 +21,20 @@ export class CreateUpdateTopicRequestDto { @IsString() @IsOptional() @Length(0, 100) - name: string; + name?: string; + + @ApiPropertyOptional({ + description: + 'Additional custom data associated with the topic. Flat key-value pairs of scalars (string, number, boolean, string[]). Maximum size: 64KB.', + type: Object, + nullable: true, + additionalProperties: true, + example: { category: 'product', priority: 1 }, + }) + @IsOptional() + @ValidateIf((obj) => obj.data !== null) + @IsObject() + @IsValidCustomData() + @IsValidContextData() + data?: TopicCustomData | null; } diff --git a/apps/api/src/app/topics-v2/dtos/delete-topic-subscriptions-response.dto.ts b/apps/api/src/app/topics-v2/dtos/delete-topic-subscriptions-response.dto.ts index 4939aaa3a64..ff14f5fedc1 100644 --- a/apps/api/src/app/topics-v2/dtos/delete-topic-subscriptions-response.dto.ts +++ b/apps/api/src/app/topics-v2/dtos/delete-topic-subscriptions-response.dto.ts @@ -20,6 +20,15 @@ export class TopicDto { required: false, }) name?: string; + + @ApiProperty({ + description: 'Additional custom data associated with the topic', + example: { category: 'product', priority: 1 }, + required: false, + type: Object, + additionalProperties: true, + }) + data?: Record; } export class SubscriberDto { diff --git a/apps/api/src/app/topics-v2/dtos/topic-response.dto.ts b/apps/api/src/app/topics-v2/dtos/topic-response.dto.ts index 00d2818c17a..b216be211bd 100644 --- a/apps/api/src/app/topics-v2/dtos/topic-response.dto.ts +++ b/apps/api/src/app/topics-v2/dtos/topic-response.dto.ts @@ -1,5 +1,6 @@ import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; -import { IsOptional, IsString } from 'class-validator'; +import { TopicCustomData } from '@novu/shared'; +import { IsObject, IsOptional, IsString } from 'class-validator'; export class TopicResponseDto { @ApiProperty({ @@ -27,6 +28,16 @@ export class TopicResponseDto { @IsOptional() name?: string; + @ApiPropertyOptional({ + description: 'Additional custom data associated with the topic', + type: Object, + additionalProperties: true, + example: { category: 'product', priority: 1 }, + }) + @IsObject() + @IsOptional() + data?: TopicCustomData; + @ApiPropertyOptional({ description: 'The date the topic was created', type: String, diff --git a/apps/api/src/app/topics-v2/dtos/update-topic.dto.ts b/apps/api/src/app/topics-v2/dtos/update-topic.dto.ts index 42ff06c060c..fea791a8635 100644 --- a/apps/api/src/app/topics-v2/dtos/update-topic.dto.ts +++ b/apps/api/src/app/topics-v2/dtos/update-topic.dto.ts @@ -1,12 +1,30 @@ -import { ApiProperty } from '@nestjs/swagger'; -import { IsNotEmpty, IsString } from 'class-validator'; +import { ApiPropertyOptional } from '@nestjs/swagger'; +import { IsValidContextData, IsValidCustomData } from '@novu/application-generic'; +import { TopicCustomData } from '@novu/shared'; +import { IsObject, IsOptional, IsString, Length, ValidateIf } from 'class-validator'; export class UpdateTopicRequestDto { - @ApiProperty({ + @ApiPropertyOptional({ description: 'The display name for the topic', example: 'Updated Topic Name', }) @IsString() - @IsNotEmpty() - name: string; + @IsOptional() + @Length(0, 100) + name?: string; + + @ApiPropertyOptional({ + description: + 'Additional custom data associated with the topic. Flat key-value pairs of scalars (string, number, boolean, string[]). Maximum size: 64KB. Pass null to clear.', + type: Object, + nullable: true, + additionalProperties: true, + example: { category: 'product', priority: 1 }, + }) + @IsOptional() + @ValidateIf((obj) => obj.data !== null) + @IsObject() + @IsValidCustomData() + @IsValidContextData() + data?: TopicCustomData | null; } diff --git a/apps/api/src/app/topics-v2/e2e/update-topic.e2e.ts b/apps/api/src/app/topics-v2/e2e/update-topic.e2e.ts index 580fc079043..64eca8016cd 100644 --- a/apps/api/src/app/topics-v2/e2e/update-topic.e2e.ts +++ b/apps/api/src/app/topics-v2/e2e/update-topic.e2e.ts @@ -15,7 +15,6 @@ describe('Update topic by key - /v2/topics/:topicKey (PATCH) #novu-v2', async () await session.initialize(); novuClient = initNovuClassSdk(session); - // Create a topic to update later const createResponse = await novuClient.topics.create({ key: topicKey, name: initialName, @@ -39,12 +38,46 @@ describe('Update topic by key - /v2/topics/:topicKey (PATCH) #novu-v2', async () expect(response.result).to.have.property('createdAt'); expect(response.result).to.have.property('updatedAt'); - // Verify the update persisted by fetching the topic const getResponse = await novuClient.topics.get(topicKey); expect(getResponse.result).to.exist; expect(getResponse.result.name).to.equal(updatedName); }); + it('should update topic data without changing name', async () => { + const dataKey = `topic-key-patch-data-${Date.now()}`; + await session.testAgent.post('/v2/topics').send({ + key: dataKey, + name: 'Keep Name', + data: { category: 'old' }, + }); + + const { body } = await session.testAgent.patch(`/v2/topics/${dataKey}`).send({ + data: { category: 'new', featured: true }, + }); + + expect(body.data.name).to.equal('Keep Name'); + expect(body.data.data).to.deep.equal({ category: 'new', featured: true }); + }); + + it('should clear topic data when patching with null', async () => { + const dataKey = `topic-key-clear-data-${Date.now()}`; + await session.testAgent.post('/v2/topics').send({ + key: dataKey, + name: 'Clear Data', + data: { category: 'old' }, + }); + + const { body } = await session.testAgent.patch(`/v2/topics/${dataKey}`).send({ + data: null, + }); + + expect(body.data.name).to.equal('Clear Data'); + expect(body.data.data).to.be.undefined; + + const getResponse = await session.testAgent.get(`/v2/topics/${dataKey}`); + expect(getResponse.body.data.data).to.be.undefined; + }); + it('should return 404 for updating a non-existent topic key', async () => { const nonExistentKey = 'non-existent-topic-key'; try { @@ -55,7 +88,6 @@ describe('Update topic by key - /v2/topics/:topicKey (PATCH) #novu-v2', async () nonExistentKey ); - /* If we reach here, the test failed */ expect.fail('Should have thrown an error for non-existent topic'); } catch (error) { expect(error.statusCode).to.equal(404); diff --git a/apps/api/src/app/topics-v2/e2e/upsert-topic.e2e.ts b/apps/api/src/app/topics-v2/e2e/upsert-topic.e2e.ts index 6c44fd068e9..70461654b71 100644 --- a/apps/api/src/app/topics-v2/e2e/upsert-topic.e2e.ts +++ b/apps/api/src/app/topics-v2/e2e/upsert-topic.e2e.ts @@ -31,8 +31,47 @@ describe('Upsert topic - /v2/topics (POST) #novu-v2', async () => { expect(response.result).to.have.property('updatedAt'); }); + it('should create a topic with custom data', async () => { + const key = `topic-key-data-${Date.now()}`; + const name = 'Topic With Data'; + const data = { category: 'product', priority: 1, tags: ['a', 'b'] }; + + const { body } = await session.testAgent.post('/v2/topics').send({ key, name, data }); + + expect(body.data).to.exist; + expect(body.data.key).to.equal(key); + expect(body.data.name).to.equal(name); + expect(body.data.data).to.deep.equal(data); + }); + + it('should reject topic data larger than 64KB', async () => { + const key = `topic-key-oversized-${Date.now()}`; + const oversizedValue = 'x'.repeat(65 * 1024); + + const { body } = await session.testAgent.post('/v2/topics').send({ + key, + name: 'Oversized', + data: { big: oversizedValue }, + }); + + expect(body.statusCode).to.equal(422); + expect(JSON.stringify(body)).to.match(/too large|Data is too large|Validation Error/i); + }); + + it('should reject nested objects in topic data', async () => { + const key = `topic-key-nested-${Date.now()}`; + + const { body } = await session.testAgent.post('/v2/topics').send({ + key, + name: 'Nested', + data: { nested: { not: 'allowed' } }, + }); + + expect(body.statusCode).to.equal(422); + expect(JSON.stringify(body)).to.match(/must be a string, number, boolean, or string\[\]/i); + }); + it('should update an existing topic when it already exists', async () => { - // First create a topic const key = `topic-key-${Date.now()}`; const originalName = 'Original Name'; @@ -44,7 +83,6 @@ describe('Upsert topic - /v2/topics (POST) #novu-v2', async () => { expect(createResponse.result).to.exist; const originalId = createResponse.result.id; - // Now update the same topic by creating with the same key const updatedName = 'Updated Name'; const updateResponse = await novuClient.topics.update( { @@ -57,8 +95,26 @@ describe('Upsert topic - /v2/topics (POST) #novu-v2', async () => { expect(updateResponse.result.id).to.equal(originalId); expect(updateResponse.result.key).to.equal(key); expect(updateResponse.result.name).to.equal(updatedName); - // Verify the update persisted by fetching the topic const getResponse = await novuClient.topics.get(key); expect(getResponse.result.name).to.equal(updatedName); }); + + it('should update topic data on upsert when topic already exists', async () => { + const key = `topic-key-upsert-data-${Date.now()}`; + + await session.testAgent.post('/v2/topics').send({ + key, + name: 'Original', + data: { category: 'old' }, + }); + + const { body } = await session.testAgent.post('/v2/topics').send({ + key, + name: 'Updated', + data: { category: 'new', priority: 2 }, + }); + + expect(body.data.name).to.equal('Updated'); + expect(body.data.data).to.deep.equal({ category: 'new', priority: 2 }); + }); }); diff --git a/apps/api/src/app/topics-v2/topics.controller.ts b/apps/api/src/app/topics-v2/topics.controller.ts index 8a23d63abb2..c272591546f 100644 --- a/apps/api/src/app/topics-v2/topics.controller.ts +++ b/apps/api/src/app/topics-v2/topics.controller.ts @@ -148,6 +148,7 @@ export class TopicsController { organizationId: user.organizationId, key: body.key, name: body.name, + data: body.data, failIfExists, }) ); @@ -184,7 +185,7 @@ export class TopicsController { @SdkMethodName('update') @ApiOperation({ summary: 'Update a topic', - description: `Update a topic name by its unique key identifier **topicKey**`, + description: `Update a topic name or data by its unique key identifier **topicKey**`, }) @ApiParam({ name: 'topicKey', description: 'The key identifier of the topic', type: String }) @ApiResponse(TopicResponseDto, 200) @@ -201,6 +202,7 @@ export class TopicsController { userId: user._id, topicKey, name: body.name, + data: body.data, }) ); } diff --git a/apps/api/src/app/topics-v2/usecases/delete-topic-subscriptions/delete-topic-subscriptions.usecase.ts b/apps/api/src/app/topics-v2/usecases/delete-topic-subscriptions/delete-topic-subscriptions.usecase.ts index 7b28ab48f04..82d9e782add 100644 --- a/apps/api/src/app/topics-v2/usecases/delete-topic-subscriptions/delete-topic-subscriptions.usecase.ts +++ b/apps/api/src/app/topics-v2/usecases/delete-topic-subscriptions/delete-topic-subscriptions.usecase.ts @@ -301,6 +301,7 @@ export class DeleteTopicSubscriptionsUsecase { _id: topic._id, key: topic.key, name: topic.name, + data: topic.data, }, subscriber: subscriber ? { diff --git a/apps/api/src/app/topics-v2/usecases/list-topics/map-topic-entity-to.dto.ts b/apps/api/src/app/topics-v2/usecases/list-topics/map-topic-entity-to.dto.ts index a00d41212f7..7913b33510c 100644 --- a/apps/api/src/app/topics-v2/usecases/list-topics/map-topic-entity-to.dto.ts +++ b/apps/api/src/app/topics-v2/usecases/list-topics/map-topic-entity-to.dto.ts @@ -8,6 +8,7 @@ export function mapTopicEntityToDto(topicEntity: TopicEntity): TopicResponseDto _id: String(topicEntity._id), name: topicEntity.name, key: topicEntity.key, + ...(topicEntity.data != null ? { data: topicEntity.data } : {}), createdAt: topicEntity.createdAt, updatedAt: topicEntity.updatedAt, }; diff --git a/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.command.ts b/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.command.ts index c2ea26183aa..345cbfc212e 100644 --- a/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.command.ts +++ b/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.command.ts @@ -1,4 +1,6 @@ -import { IsNotEmpty, IsString } from 'class-validator'; +import { IsValidContextData, IsValidCustomData } from '@novu/application-generic'; +import { TopicCustomData } from '@novu/shared'; +import { IsNotEmpty, IsObject, IsOptional, IsString, Length, ValidateIf } from 'class-validator'; import { EnvironmentWithUserCommand } from '../../../shared/commands/project.command'; export class UpdateTopicCommand extends EnvironmentWithUserCommand { @@ -7,6 +9,14 @@ export class UpdateTopicCommand extends EnvironmentWithUserCommand { topicKey: string; @IsString() - @IsNotEmpty() - name: string; + @IsOptional() + @Length(0, 100) + name?: string; + + @IsOptional() + @ValidateIf((obj) => obj.data !== null) + @IsObject() + @IsValidCustomData() + @IsValidContextData() + data?: TopicCustomData | null; } diff --git a/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.usecase.ts b/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.usecase.ts index e0762974630..daa9590d568 100644 --- a/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.usecase.ts +++ b/apps/api/src/app/topics-v2/usecases/update-topic/update-topic.usecase.ts @@ -1,4 +1,4 @@ -import { Injectable, NotFoundException } from '@nestjs/common'; +import { BadRequestException, Injectable, NotFoundException } from '@nestjs/common'; import { InstrumentUsecase } from '@novu/application-generic'; import { TopicRepository } from '@novu/dal'; import { TopicResponseDto } from '../../dtos/topic-response.dto'; @@ -21,6 +21,25 @@ export class UpdateTopicUseCase { throw new NotFoundException(`Topic with key ${command.topicKey} not found`); } + if (command.name === undefined && command.data === undefined) { + throw new BadRequestException('At least one of name or data must be provided'); + } + + const setBody: Record = {}; + const unsetBody: Record = {}; + + if (command.name !== undefined) { + setBody.name = command.name; + } + + if (command.data !== undefined) { + if (command.data === null) { + unsetBody.data = ''; + } else { + setBody.data = command.data; + } + } + const updatedTopic = await this.topicRepository.findOneAndUpdate( { _id: existingTopic._id, @@ -28,9 +47,8 @@ export class UpdateTopicUseCase { _organizationId: command.organizationId, }, { - $set: { - name: command.name, - }, + ...(Object.keys(setBody).length > 0 ? { $set: setBody } : {}), + ...(Object.keys(unsetBody).length > 0 ? { $unset: unsetBody } : {}), }, { new: true } ); diff --git a/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.command.ts b/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.command.ts index 6001a584eb4..2a6eb831941 100644 --- a/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.command.ts +++ b/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.command.ts @@ -1,5 +1,6 @@ -import { EnvironmentCommand } from '@novu/application-generic'; -import { IsBoolean, IsNotEmpty, IsOptional, IsString, Length } from 'class-validator'; +import { EnvironmentCommand, IsValidContextData, IsValidCustomData } from '@novu/application-generic'; +import { TopicCustomData } from '@novu/shared'; +import { IsBoolean, IsNotEmpty, IsObject, IsOptional, IsString, Length, ValidateIf } from 'class-validator'; export class UpsertTopicCommand extends EnvironmentCommand { @IsString() @@ -12,6 +13,13 @@ export class UpsertTopicCommand extends EnvironmentCommand { @Length(0, 100) name?: string; + @IsOptional() + @ValidateIf((obj) => obj.data !== null) + @IsObject() + @IsValidCustomData() + @IsValidContextData() + data?: TopicCustomData | null; + @IsBoolean() @IsOptional() failIfExists?: boolean; diff --git a/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.usecase.ts b/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.usecase.ts index ee77b99796b..8fa4c118556 100644 --- a/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.usecase.ts +++ b/apps/api/src/app/topics-v2/usecases/upsert-topic/upsert-topic.usecase.ts @@ -1,4 +1,4 @@ -import { BadRequestException, ConflictException, Injectable } from '@nestjs/common'; +import { BadRequestException, ConflictException, Injectable, NotFoundException } from '@nestjs/common'; import { InstrumentUsecase } from '@novu/application-generic'; import { ErrorCodesEnum, TopicRepository } from '@novu/dal'; import { VALID_ID_REGEX } from '@novu/shared'; @@ -17,6 +17,8 @@ export class UpsertTopicUseCase { throw new ConflictException(`Topic with key "${command.key}" already exists`); } + let created = !topic; + if (!topic) { this.isValidTopicKey(command.key); @@ -26,39 +28,78 @@ export class UpsertTopicUseCase { _organizationId: command.organizationId, key: command.key, name: command.name, + data: command.data ?? undefined, }); } catch (error: unknown) { - if (this.isDuplicateKeyError(error)) { - topic = await this.topicRepository.findTopicByKey(command.key, command.organizationId, command.environmentId); - } else { + if (!this.isDuplicateKeyError(error)) { throw error; } - } - } else { - const updateBody: Record = {}; - if (command.name) { - updateBody.name = command.name; - } + if (command.failIfExists) { + throw new ConflictException(`Topic with key "${command.key}" already exists`); + } - topic = await this.topicRepository.findOneAndUpdate( - { - _id: topic._id, - _environmentId: command.environmentId, - _organizationId: command.organizationId, - }, - { - $set: updateBody, + const winningTopic = await this.topicRepository.findTopicByKey( + command.key, + command.organizationId, + command.environmentId + ); + + if (!winningTopic) { + throw error; } - ); + + created = false; + topic = await this.applyTopicUpdate(winningTopic._id, command); + } + } else { + topic = await this.applyTopicUpdate(topic._id, command); + } + + if (!topic) { + throw new NotFoundException(`Topic with key "${command.key}" not found`); } return { - topic: mapTopicEntityToDto(topic!), - created: !topic, + topic: mapTopicEntityToDto(topic), + created, }; } + private async applyTopicUpdate(topicId: string, command: UpsertTopicCommand) { + const setBody: Record = {}; + const unsetBody: Record = {}; + + if (command.name !== undefined) { + setBody.name = command.name; + } + + if (command.data !== undefined) { + if (command.data === null) { + unsetBody.data = ''; + } else { + setBody.data = command.data; + } + } + + if (Object.keys(setBody).length === 0 && Object.keys(unsetBody).length === 0) { + return this.topicRepository.findTopicByKey(command.key, command.organizationId, command.environmentId); + } + + return await this.topicRepository.findOneAndUpdate( + { + _id: topicId, + _environmentId: command.environmentId, + _organizationId: command.organizationId, + }, + { + ...(Object.keys(setBody).length > 0 ? { $set: setBody } : {}), + ...(Object.keys(unsetBody).length > 0 ? { $unset: unsetBody } : {}), + }, + { new: true } + ); + } + private isValidTopicKey(key: string): void { if (VALID_ID_REGEX.test(key)) { return; diff --git a/docs/api-reference/topics/topic-schema.mdx b/docs/api-reference/topics/topic-schema.mdx index 5b8b333ce6c..0249125fb0c 100644 --- a/docs/api-reference/topics/topic-schema.mdx +++ b/docs/api-reference/topics/topic-schema.mdx @@ -12,6 +12,7 @@ Topic is a collection of subscribers that share a common interest. Subscriber ca | `id` | `string` | The identifier of the topic | | `key` | `string` | The unique key of the topic | | `name` | `string` | The name of the topic | +| `data` | `object` | Optional flat custom data (string, number, boolean, or string[] values). Maximum serialized size: 64KB. Returned on topic and topic-subscription API responses. | | `createdAt` | `string` | The date the topic was created | | `updatedAt` | `string` | The date the topic was last updated | diff --git a/libs/application-generic/src/decorators/index.ts b/libs/application-generic/src/decorators/index.ts index 86717da4dc8..2a177e6ef08 100644 --- a/libs/application-generic/src/decorators/index.ts +++ b/libs/application-generic/src/decorators/index.ts @@ -1,6 +1,7 @@ export * from './context-payload.decorator'; export * from './external-api.decorator'; export * from './is-valid-context-payload.decorator'; +export * from './is-valid-custom-data.decorator'; export * from './is-valid-locale.decorator'; export * from './json-schema.validator'; export * from './permissions.decorator'; diff --git a/libs/application-generic/src/decorators/is-valid-custom-data.decorator.ts b/libs/application-generic/src/decorators/is-valid-custom-data.decorator.ts new file mode 100644 index 00000000000..ef708c5262b --- /dev/null +++ b/libs/application-generic/src/decorators/is-valid-custom-data.decorator.ts @@ -0,0 +1,62 @@ +import { registerDecorator, ValidationOptions } from 'class-validator'; + +function isScalarCustomDataValue(value: unknown): boolean { + if (value === undefined) { + return true; + } + + if (typeof value === 'string' || typeof value === 'number' || typeof value === 'boolean') { + return true; + } + + if (Array.isArray(value)) { + return value.every((item) => typeof item === 'string'); + } + + return false; +} + +export function validateCustomDataWithDetails(value: unknown): { isValid: boolean; error?: string } { + if (value == null) { + return { isValid: true }; + } + + if (typeof value !== 'object' || Array.isArray(value)) { + return { isValid: false, error: 'Custom data must be an object' }; + } + + for (const [key, entry] of Object.entries(value as Record)) { + if (!isScalarCustomDataValue(entry)) { + return { + isValid: false, + error: `Custom data field "${key}" must be a string, number, boolean, or string[]`, + }; + } + } + + return { isValid: true }; +} + +export function IsValidCustomData(validationOptions?: ValidationOptions) { + return (object: object, propertyName: string) => { + let lastValidationError: string | undefined; + + registerDecorator({ + name: 'isValidCustomData', + target: object.constructor, + propertyName, + options: validationOptions, + validator: { + validate(value: unknown) { + const result = validateCustomDataWithDetails(value); + lastValidationError = result.error; + + return result.isValid; + }, + defaultMessage() { + return lastValidationError || 'Invalid custom data'; + }, + }, + }); + }; +} diff --git a/libs/dal/src/repositories/topic/topic.entity.ts b/libs/dal/src/repositories/topic/topic.entity.ts index e69be47b56d..ce3ea518f77 100644 --- a/libs/dal/src/repositories/topic/topic.entity.ts +++ b/libs/dal/src/repositories/topic/topic.entity.ts @@ -1,3 +1,4 @@ +import { CustomDataType } from '@novu/shared'; import { Types } from 'mongoose'; import { EnvironmentId, OrganizationId, TopicId, TopicKey, TopicName } from './types'; @@ -8,6 +9,7 @@ export class TopicEntity { _organizationId: OrganizationId; key: TopicKey; name?: TopicName; + data?: CustomDataType; createdAt?: string; updatedAt?: string; diff --git a/libs/dal/src/repositories/topic/topic.repository.ts b/libs/dal/src/repositories/topic/topic.repository.ts index 5e0991b1579..1f793bf6a6e 100644 --- a/libs/dal/src/repositories/topic/topic.repository.ts +++ b/libs/dal/src/repositories/topic/topic.repository.ts @@ -19,6 +19,7 @@ const topicWithSubscribersProjection = { updatedAt: 1, key: 1, name: 1, + data: 1, subscribers: '$topicSubscribers.externalSubscriberId', }, }; @@ -38,12 +39,13 @@ export class TopicRepository extends BaseRepository): Promise { - const { key, name, _environmentId, _organizationId } = entity; + const { key, name, data, _environmentId, _organizationId } = entity; return await this.create({ _environmentId, key, name, + data, _organizationId, }); } diff --git a/libs/dal/src/repositories/topic/topic.schema.ts b/libs/dal/src/repositories/topic/topic.schema.ts index f384ccda86e..d44ea203aad 100644 --- a/libs/dal/src/repositories/topic/topic.schema.ts +++ b/libs/dal/src/repositories/topic/topic.schema.ts @@ -24,6 +24,7 @@ const topicSchema = new Schema( name: { type: Schema.Types.String, }, + data: Schema.Types.Mixed, }, schemaOptions ); diff --git a/packages/shared/src/types/topic.ts b/packages/shared/src/types/topic.ts index a89f460bd55..b12e32749df 100644 --- a/packages/shared/src/types/topic.ts +++ b/packages/shared/src/types/topic.ts @@ -1,3 +1,6 @@ +import type { CustomDataType } from './utils'; + export type TopicId = string; export type TopicKey = string; export type TopicName = string; +export type TopicCustomData = CustomDataType; From 83f5edf9970eac4cdabe6fcee4f30456b6bbd7f5 Mon Sep 17 00:00:00 2001 From: Adam Chmara Date: Thu, 20 Aug 2026 14:07:27 +0200 Subject: [PATCH 6/7] fix(agent-chat): await ingress processing and fix Agent Chat runtime bugs fixes NV-8593 (#12392) --- packages/chat-adapter-agent-chat/src/adapter.ts | 4 ++-- packages/js/src/agent-chat/apply-envelope.ts | 5 +++-- packages/js/src/ws/socket-factory.ts | 2 ++ 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/packages/chat-adapter-agent-chat/src/adapter.ts b/packages/chat-adapter-agent-chat/src/adapter.ts index fd93c566b66..fc95abd2d66 100644 --- a/packages/chat-adapter-agent-chat/src/adapter.ts +++ b/packages/chat-adapter-agent-chat/src/adapter.ts @@ -240,7 +240,7 @@ export class NovuAgentChatAdapterImpl implements Adapter= 0) { - next[streamingIndex] = { type: 'text', text: content.markdown, state: 'done' }; + const textIndex = streamingIndex >= 0 ? streamingIndex : next.findIndex((part) => part.type === 'text'); + if (textIndex >= 0) { + next[textIndex] = { type: 'text', text: content.markdown, state: 'done' }; } else { next = [...next, { type: 'text', text: content.markdown, state: 'done' }]; } diff --git a/packages/js/src/ws/socket-factory.ts b/packages/js/src/ws/socket-factory.ts index 1bf07d6df24..b6dc7750b6f 100644 --- a/packages/js/src/ws/socket-factory.ts +++ b/packages/js/src/ws/socket-factory.ts @@ -13,6 +13,8 @@ const PARTY_SOCKET_URLS = [ 'wss://socket-worker-local.cli-shortener.workers.dev', 'ws://127.0.0.1:8787', 'http://127.0.0.1:8787', + 'ws://localhost:8787', + 'http://localhost:8787', ]; const URL_TRANSFORMATIONS: Record = { From b6027facff239888ef353cfdc38edfb71ea0a5e3 Mon Sep 17 00:00:00 2001 From: Adam Chmara Date: Thu, 20 Aug 2026 14:34:39 +0200 Subject: [PATCH 7/7] feat(novu): add Agent Chat channel to npx novu connect fixes NV-8593 (#12393) --- .../add-agent-integration.usecase.ts | 8 +- .../agents/e2e/agent-chat-conversation.e2e.ts | 22 - .../shared/assert-agent-chat-enabled.ts | 37 +- apps/api/src/app/agents/shared/index.ts | 2 +- .../agents/agent-chat-setup-content.tsx | 33 +- packages/novu/package.json | 2 +- .../src/commands/connect/analytics/events.ts | 1 + .../novu/src/commands/connect/api/agents.ts | 17 + .../resolve-connect-application-identifier.ts | 25 + .../connect/connect-channel-picker-options.ts | 19 + .../src/commands/connect/dashboard-urls.ts | 2 + .../novu/src/commands/connect/help-text.ts | 3 +- .../agent-chat/detect-agent-chat-ui.ts | 74 ++ .../offer-post-connect-bridge-tunnel.ts | 73 ++ .../resolve-agent-chat-wiring-state.spec.ts | 56 ++ .../resolve-agent-chat-wiring-state.ts | 71 ++ .../agent-chat/run-agent-chat-setup.ts | 183 ++++ .../agent-chat/scaffold-agent-chat.spec.ts | 39 + .../agent-chat/scaffold-agent-chat.ts | 483 ++++++++++ .../agent-chat/wire-agent-chat-env.ts | 165 ++++ .../wrap-ui-for-agent-chat-handoff.ts | 39 + .../connect/pipeline/bridge-adapter/engine.ts | 18 +- .../bridge/print-bridge-dev-next-steps.ts | 15 +- .../connect/pipeline/channels/agent-chat.ts | 75 ++ .../connect/pipeline/chat-sdk/index.ts | 2 +- .../llm-auth/ensure-subscription-auth.ts | 4 + .../pipeline/llm-auth/llm-auth-options.ts | 17 +- .../pipeline/llm-auth/llm-auth.spec.ts | 7 +- .../pipeline/llm-auth/resolve-llm-auth.ts | 4 + .../src/commands/connect/pipeline/runner.ts | 72 +- .../templates/agent-chat/ts/agent-chat.tsx | 42 + .../templates/agent-chat/ts/chat-panel.tsx | 83 ++ .../templates/agent-chat/ts/chat-thread.tsx | 132 +++ .../templates/agent-chat/ts/composer.tsx | 80 ++ .../templates/agent-chat/ts/globals.css | 902 ++++++++++++++++++ .../connect/templates/agent-chat/ts/icons.tsx | 247 +++++ .../templates/agent-chat/ts/markdown.tsx | 18 + .../agent-chat/ts/message-bubble.tsx | 236 +++++ .../templates/agent-chat/ts/message-utils.ts | 17 + .../agent-chat/ts/pending-action-card.tsx | 146 +++ packages/novu/src/commands/connect/types.ts | 28 +- .../ui/agent-chat-embed-finish-content.tsx | 94 ++ packages/novu/src/commands/connect/ui/app.tsx | 4 +- .../ui/bridge-reconcile-phase-content.tsx | 1 + .../connect/ui/confirm-scaffold-content.tsx | 40 +- .../ui/console-bridge-scaffold-prompts.ts | 4 + .../src/commands/connect/ui/copyable-link.tsx | 8 +- .../connect/ui/format-agent-chat-success.ts | 153 +++ .../src/commands/connect/ui/handoff-events.ts | 16 + .../novu/src/commands/connect/ui/index.tsx | 117 ++- .../commands/connect/ui/llm-auth-picker.tsx | 14 +- .../src/commands/connect/ui/logging-ui.ts | 29 +- .../src/commands/connect/ui/orb/orb-tint.ts | 12 + .../src/commands/connect/ui/phase-content.tsx | 152 ++- .../ui/phase-content/agent-chat-phases.tsx | 172 ++++ .../connect/ui/phase-has-copyable-url.ts | 11 +- .../connect/ui/print-connect-success.ts | 109 ++- .../novu/src/commands/connect/ui/store.ts | 20 + packages/novu/src/commands/connect/ui/ui.ts | 18 +- packages/novu/src/index.ts | 17 +- packages/novu/tsconfig.json | 1 + .../src/utils/agent-chat-connect-prompt.ts | 35 + .../utils/connect-agent-chat-dashboard-url.ts | 15 + .../utils/connect-embed-prompt-constants.ts | 2 + .../shared/src/utils/connect-embed-prompt.ts | 325 +++++++ packages/shared/src/utils/index.ts | 3 + 66 files changed, 4716 insertions(+), 155 deletions(-) create mode 100644 packages/novu/src/commands/connect/auth/resolve-connect-application-identifier.ts create mode 100644 packages/novu/src/commands/connect/connect-channel-picker-options.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/detect-agent-chat-ui.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/offer-post-connect-bridge-tunnel.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.spec.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/run-agent-chat-setup.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.spec.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/wire-agent-chat-env.ts create mode 100644 packages/novu/src/commands/connect/pipeline/agent-chat/wrap-ui-for-agent-chat-handoff.ts create mode 100644 packages/novu/src/commands/connect/pipeline/channels/agent-chat.ts create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/agent-chat.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/chat-panel.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/chat-thread.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/composer.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/globals.css create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/icons.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/markdown.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/message-bubble.tsx create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/message-utils.ts create mode 100644 packages/novu/src/commands/connect/templates/agent-chat/ts/pending-action-card.tsx create mode 100644 packages/novu/src/commands/connect/ui/agent-chat-embed-finish-content.tsx create mode 100644 packages/novu/src/commands/connect/ui/format-agent-chat-success.ts create mode 100644 packages/novu/src/commands/connect/ui/phase-content/agent-chat-phases.tsx create mode 100644 packages/shared/src/utils/agent-chat-connect-prompt.ts create mode 100644 packages/shared/src/utils/connect-agent-chat-dashboard-url.ts create mode 100644 packages/shared/src/utils/connect-embed-prompt-constants.ts create mode 100644 packages/shared/src/utils/connect-embed-prompt.ts diff --git a/apps/api/src/app/agents/channels/integrations/add-agent-integration/add-agent-integration.usecase.ts b/apps/api/src/app/agents/channels/integrations/add-agent-integration/add-agent-integration.usecase.ts index d28e8da52de..ed5fa0767ca 100644 --- a/apps/api/src/app/agents/channels/integrations/add-agent-integration/add-agent-integration.usecase.ts +++ b/apps/api/src/app/agents/channels/integrations/add-agent-integration/add-agent-integration.usecase.ts @@ -8,7 +8,7 @@ import { Injectable, NotFoundException, } from '@nestjs/common'; -import { AnalyticsService, encryptSecret, isAgentEmailEnabled } from '@novu/application-generic'; +import { AnalyticsService, encryptSecret, FeatureFlagsService, isAgentEmailEnabled } from '@novu/application-generic'; import { type AgentEntity, AgentIntegrationRepository, @@ -28,6 +28,7 @@ import { } from '@novu/shared'; import { NovuEmailProvisioningService } from '../../../email/novu-email/find-or-create-novu-email/find-or-create-novu-email.service'; import { trackAgentIntegrationConnected } from '../../../shared/analytics/agent-analytics'; +import { assertAgentChatEnabledForConnect } from '../../../shared/assert-agent-chat-enabled'; import type { AgentIntegrationResponseDto } from '../../../shared/dtos'; import { toAgentIntegrationResponse } from '../../../shared/mappers/agent-response.mapper'; import { NovuAgentChatProvisioningService } from '../../agent-chat/find-or-create-novu-agent-chat/find-or-create-novu-agent-chat.service'; @@ -43,7 +44,8 @@ export class AddAgentIntegration { private readonly environmentRepository: EnvironmentRepository, private readonly findOrCreateNovuEmail: NovuEmailProvisioningService, private readonly findOrCreateNovuAgentChat: NovuAgentChatProvisioningService, - private readonly analyticsService: AnalyticsService + private readonly analyticsService: AnalyticsService, + private readonly featureFlagsService: FeatureFlagsService ) {} async execute(command: AddAgentIntegrationCommand): Promise { @@ -96,6 +98,8 @@ export class AddAgentIntegration { } if (command.providerId === ChatProviderIdEnum.NovuAgentChat) { + await assertAgentChatEnabledForConnect(this.featureFlagsService, command.organizationId, command.environmentId); + const { response, provisionedNewLink } = await this.findOrCreateNovuAgentChat.execute( agent._id, command.environmentId, diff --git a/apps/api/src/app/agents/e2e/agent-chat-conversation.e2e.ts b/apps/api/src/app/agents/e2e/agent-chat-conversation.e2e.ts index e40f0713b9f..1c81c340b09 100644 --- a/apps/api/src/app/agents/e2e/agent-chat-conversation.e2e.ts +++ b/apps/api/src/app/agents/e2e/agent-chat-conversation.e2e.ts @@ -412,28 +412,6 @@ describe('Agent Chat - /agent-chat/conversations #novu-v2', () => { expect(res.body.message).to.equal('This agent is not available on agent chat'); }); - it('should return 404 when IS_AGENT_WEB_CHAT_ENABLED is off', async () => { - const previousFlag = process.env.IS_AGENT_WEB_CHAT_ENABLED; - delete process.env.IS_AGENT_WEB_CHAT_ENABLED; - - try { - await linkAgentChat(); - - const res = await createConversation({ - agentId: ctx.agentIdentifier, - text: 'Should be hidden', - }); - - expect(res.status).to.equal(404); - } finally { - if (previousFlag === undefined) { - delete process.env.IS_AGENT_WEB_CHAT_ENABLED; - } else { - process.env.IS_AGENT_WEB_CHAT_ENABLED = previousFlag; - } - } - }); - it('should return 401 when Authorization is missing on POST', async () => { await linkAgentChat(); diff --git a/apps/api/src/app/agents/shared/assert-agent-chat-enabled.ts b/apps/api/src/app/agents/shared/assert-agent-chat-enabled.ts index e3afe7b5754..fd75931ed37 100644 --- a/apps/api/src/app/agents/shared/assert-agent-chat-enabled.ts +++ b/apps/api/src/app/agents/shared/assert-agent-chat-enabled.ts @@ -1,7 +1,22 @@ -import { NotFoundException } from '@nestjs/common'; +import { ForbiddenException, NotFoundException } from '@nestjs/common'; import { FeatureFlagsService } from '@novu/application-generic'; import { FeatureFlagsKeysEnum } from '@novu/shared'; +const AGENT_CHAT_DISABLED_CONNECT_MESSAGE = 'Agent Chat is not enabled for this workspace. Contact support@novu.co'; + +async function isAgentChatFeatureEnabled( + featureFlagsService: FeatureFlagsService, + organizationId: string, + environmentId: string +): Promise { + return featureFlagsService.getFlag({ + key: FeatureFlagsKeysEnum.IS_AGENT_WEB_CHAT_ENABLED, + defaultValue: false, + organization: { _id: organizationId }, + environment: { _id: environmentId }, + }); +} + /** * Gates subscriber agent-chat HTTP (`/v1/agent-chat/*`) on `IS_AGENT_WEB_CHAT_ENABLED`. * Shared by Nest `AgentChatEnabledGuard` (GET) and the POST create path (no AuthGuard). @@ -11,14 +26,22 @@ export async function assertAgentChatEnabled( organizationId: string, environmentId: string ): Promise { - const isEnabled = await featureFlagsService.getFlag({ - key: FeatureFlagsKeysEnum.IS_AGENT_WEB_CHAT_ENABLED, - defaultValue: false, - organization: { _id: organizationId }, - environment: { _id: environmentId }, - }); + const isEnabled = await isAgentChatFeatureEnabled(featureFlagsService, organizationId, environmentId); if (!isEnabled) { throw new NotFoundException(); } } + +/** Connect CLI / dashboard provisioning β€” clear error instead of a silent 404. */ +export async function assertAgentChatEnabledForConnect( + featureFlagsService: FeatureFlagsService, + organizationId: string, + environmentId: string +): Promise { + const isEnabled = await isAgentChatFeatureEnabled(featureFlagsService, organizationId, environmentId); + + if (!isEnabled) { + throw new ForbiddenException(AGENT_CHAT_DISABLED_CONNECT_MESSAGE); + } +} diff --git a/apps/api/src/app/agents/shared/index.ts b/apps/api/src/app/agents/shared/index.ts index 4b6b0bcfb25..015a4b98f68 100644 --- a/apps/api/src/app/agents/shared/index.ts +++ b/apps/api/src/app/agents/shared/index.ts @@ -1,4 +1,4 @@ export { AgentChatEnabledGuard } from './agent-chat-enabled.guard'; export { AgentConversationEnabledGuard } from './agent-conversation-enabled.guard'; export { AgentRuntimeExceptionFilter } from './agent-runtime-exception.filter'; -export { assertAgentChatEnabled } from './assert-agent-chat-enabled'; +export { assertAgentChatEnabled, assertAgentChatEnabledForConnect } from './assert-agent-chat-enabled'; diff --git a/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx b/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx index e75fd6ad645..2be351a4fc9 100644 --- a/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx +++ b/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx @@ -1,31 +1,14 @@ +export { + AGENT_CHAT_DOCS_URL, + APPLICATION_IDENTIFIER_PLACEHOLDER, + buildAgentChatPrompt, + SUBSCRIBER_ID_PLACEHOLDER, +} from '@novu/shared'; + +import { AGENT_CHAT_DOCS_URL } from '@novu/shared'; import { PrebuiltPromptBanner } from '@/components/onboarding/connect-agent/prebuilt-prompt-banner'; import { ExternalLink } from '@/components/shared/external-link'; -export const AGENT_CHAT_DOCS_URL = 'https://docs.novu.co/agents/channels/agent-chat'; -export const APPLICATION_IDENTIFIER_PLACEHOLDER = ''; -export const SUBSCRIBER_ID_PLACEHOLDER = 'YOUR_SUBSCRIBER_ID'; - -export function buildAgentChatPrompt( - agentName: string, - agentIdentifier: string, - applicationIdentifier: string, - subscriberId: string -): string { - return `Add Novu Agent Chat to my app with useAgentChat from @novu/react so end users can chat with the "${agentName}" agent in-product. - -Context: I'm already signed in to the Novu dashboard and the "${agentName}" Agent Chat integration already exists (agent id: ${agentIdentifier}). This is purely a frontend code integration: do NOT run the Novu CLI, the agent-onboarding flow, or keyless mode. - -Requirements: -- Install @novu/react with my project's package manager. -- Follow the Agent Chat docs: ${AGENT_CHAT_DOCS_URL} -- Render a production-quality chat UI (message list from message.parts, composer, tool approvals via respondToAction). Match my app's styling β€” do not dump raw JSON. -- Wrap that UI in configured for the currently signed-in end user. -- Use applicationIdentifier="${applicationIdentifier}". Store applicationIdentifier in an environment variable rather than hardcoding it. -- Pass the authenticated user's id as subscriberId: source it from my app's existing auth (Clerk, NextAuth, Firebase, Supabase, or custom). If no auth system exists yet, use "${subscriberId}" for a quick smoke test. -- If my app enables Novu subscriber HMAC, pass the matching subscriber hash into NovuProvider (same pattern as Inbox). -- Follow my app's existing framework, routing, styling, and TypeScript conventions, place the chat in a sensible spot in the UI, and add no unnecessary wrappers.`; -} - export function AgentChatEmbedResources({ prompt }: { prompt: string }) { return (
diff --git a/packages/novu/package.json b/packages/novu/package.json index cd2dbfca01a..e7630f5c4d7 100644 --- a/packages/novu/package.json +++ b/packages/novu/package.json @@ -18,7 +18,7 @@ ], "scripts": { "prebuild": "rimraf dist", - "build": "pnpm prebuild && tsc -p tsconfig.json && tsc -p tsconfig.ui.json && node scripts/build-ui.mjs && cp -r src/commands/init/templates/app* dist/src/commands/init/templates && cp -r src/commands/init/templates/github dist/src/commands/init/templates && cp -r src/commands/wizard/skills/content dist/src/commands/wizard/skills", + "build": "pnpm prebuild && tsc -p tsconfig.json && tsc -p tsconfig.ui.json && node scripts/build-ui.mjs && cp -r src/commands/init/templates/app* dist/src/commands/init/templates && cp -r src/commands/init/templates/github dist/src/commands/init/templates && cp -r src/commands/connect/templates dist/src/commands/connect/templates && cp -r src/commands/wizard/skills/content dist/src/commands/wizard/skills", "build:prod": "pnpm prebuild && pnpm build", "precommit": "lint-staged", "start": "pnpm start:dev", diff --git a/packages/novu/src/commands/connect/analytics/events.ts b/packages/novu/src/commands/connect/analytics/events.ts index cc1f2dad7d2..a3c71a3bd1c 100644 --- a/packages/novu/src/commands/connect/analytics/events.ts +++ b/packages/novu/src/commands/connect/analytics/events.ts @@ -37,6 +37,7 @@ export const CONNECT_EVENTS = { TELEGRAM_CONNECTED: 'Connect Telegram Connected', EMAIL_CONNECTED: 'Connect Email Connected', SENDBLUE_CONNECTED: 'Connect Sendblue Connected', + AGENT_CHAT_LINKED: 'Connect Agent Chat Linked', WELCOME_SENT: 'Connect Welcome Sent', COMPLETED: 'Connect Completed', ERROR: 'Connect Error', diff --git a/packages/novu/src/commands/connect/api/agents.ts b/packages/novu/src/commands/connect/api/agents.ts index d10c05774cd..2764bc54b22 100644 --- a/packages/novu/src/commands/connect/api/agents.ts +++ b/packages/novu/src/commands/connect/api/agents.ts @@ -156,6 +156,23 @@ export async function addAgentEmailIntegration( return 'data' in body && body.data ? body.data : (body as AgentIntegrationLink); } +/** + * `POST /v1/agents/:id/integrations` with `providerId: 'novu-agent-chat'` + * auto-provisions the Agent Chat integration and links it to the agent. + */ +export async function addAgentChatIntegration( + client: ConnectApiClient, + agentIdentifier: string +): Promise { + const res = await client.axios.post<{ data?: AgentIntegrationLink } | AgentIntegrationLink>( + `/v1/agents/${encodeURIComponent(agentIdentifier)}/integrations`, + { providerId: 'novu-agent-chat' } + ); + const body = res.data; + + return 'data' in body && body.data ? body.data : (body as AgentIntegrationLink); +} + export async function listAgentIntegrations( client: ConnectApiClient, agentIdentifier: string, diff --git a/packages/novu/src/commands/connect/auth/resolve-connect-application-identifier.ts b/packages/novu/src/commands/connect/auth/resolve-connect-application-identifier.ts new file mode 100644 index 00000000000..83773118219 --- /dev/null +++ b/packages/novu/src/commands/connect/auth/resolve-connect-application-identifier.ts @@ -0,0 +1,25 @@ +import { APPLICATION_IDENTIFIER_PLACEHOLDER } from '@novu/shared'; +import { requestApiJson } from '../../shared/novu-http'; +import type { ResolvedConnectAuth } from './resolve-connect-auth'; + +export async function resolveConnectApplicationIdentifier(auth: ResolvedConnectAuth): Promise { + const keyless = auth.keylessApplicationIdentifier?.trim(); + if (keyless) { + return keyless; + } + + const secretKey = auth.secretKey?.trim(); + if (!secretKey) { + return APPLICATION_IDENTIFIER_PLACEHOLDER; + } + + try { + const environment = await requestApiJson<{ identifier: string }>(auth.apiUrl, '/environments/me', { + headers: { Authorization: `ApiKey ${secretKey}` }, + }); + + return environment.identifier?.trim() || APPLICATION_IDENTIFIER_PLACEHOLDER; + } catch { + return APPLICATION_IDENTIFIER_PLACEHOLDER; + } +} diff --git a/packages/novu/src/commands/connect/connect-channel-picker-options.ts b/packages/novu/src/commands/connect/connect-channel-picker-options.ts new file mode 100644 index 00000000000..92732c63656 --- /dev/null +++ b/packages/novu/src/commands/connect/connect-channel-picker-options.ts @@ -0,0 +1,19 @@ +import { channelDisplayName } from './dashboard-urls'; +import type { ChannelChoice } from './types'; + +export type ConnectChannelPickerOption = { + label: string; + value: ChannelChoice; +}; + +/** Single source for the interactive channel picker β€” keep values unique. */ +export const CONNECT_CHANNEL_PICKER_OPTIONS: readonly ConnectChannelPickerOption[] = [ + { label: 'Slack (recommended)', value: 'slack' }, + { label: 'Telegram', value: 'telegram' }, + { label: channelDisplayName('email'), value: 'email' }, + { label: channelDisplayName('sendblue'), value: 'sendblue' }, + { label: channelDisplayName('whatsapp'), value: 'whatsapp' }, + { label: channelDisplayName('agent-chat'), value: 'agent-chat' }, + { label: channelDisplayName('teams'), value: 'teams' }, + { label: 'Skip β€” set up later in dashboard', value: 'skip' }, +]; diff --git a/packages/novu/src/commands/connect/dashboard-urls.ts b/packages/novu/src/commands/connect/dashboard-urls.ts index a1522b31c71..e30922545f8 100644 --- a/packages/novu/src/commands/connect/dashboard-urls.ts +++ b/packages/novu/src/commands/connect/dashboard-urls.ts @@ -86,6 +86,8 @@ export function channelDisplayName(channel: ChannelChoice): string { return 'Email'; case 'sendblue': return 'iMessage (Sendblue)'; + case 'agent-chat': + return 'Agent Chat'; default: return channel; } diff --git a/packages/novu/src/commands/connect/help-text.ts b/packages/novu/src/commands/connect/help-text.ts index ba2563a9146..935ce49f101 100644 --- a/packages/novu/src/commands/connect/help-text.ts +++ b/packages/novu/src/commands/connect/help-text.ts @@ -119,7 +119,7 @@ Non-interactive (agent / CI) contract: Required for --ci mode: - Pass the agent description as the positional argument or --prompt (managed path only). - - Pass --channel (teams only without --keyless). + - Pass --channel (teams only without --keyless). - Bridge path: pass --runtime ai-sdk or --runtime langchain (omit the positional description). Authentication (pick one): @@ -134,6 +134,7 @@ Non-interactive (agent / CI) contract: - --channel sendblue β†’ requires --sendblue-api-key, --sendblue-secret-key, --sendblue-from (E.164 agent/sender number), and --sendblue-test-phone (E.164 recipient phone); no secure setup page - --channel whatsapp β†’ no extra flags; works with --keyless. When Meta Embedded Signup is available the CLI prints a public tokenized signup URL (valid ~30 min, no dashboard login needed), polls until signup completes in the browser (~15 min budget), then waits for an inbound WhatsApp message. When unavailable (flag off / self-hosted without Meta Tech Provider credentials) a keyless run prompts dashboard sign-in and retries; if still unavailable the CLI opens the dashboard integrations tab and exits. - --channel skip β†’ no extra flags (agent only, no channel) + - --channel agent-chat β†’ no extra flags. Links Agent Chat, prints dashboard Chat tab URL. Embed path writes env vars; scaffold path creates an example app. Use --agent-chat-setup scaffold|embed|skip in --ci (auto-detect when omitted). Optional CI-only escape hatches (secrets injected via env β€” never paste in chat): - --slack-config-token "xoxe.xoxp-…" β†’ skip the setup page; pass token directly diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/detect-agent-chat-ui.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/detect-agent-chat-ui.ts new file mode 100644 index 00000000000..1fd6498e3bf --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/detect-agent-chat-ui.ts @@ -0,0 +1,74 @@ +import fs from 'node:fs'; +import path from 'node:path'; + +const SOURCE_DIRS = ['lib', 'src', 'app', 'components', 'pages'] as const; + +const NOVU_REACT_IMPORT = /@novu\/react/; +const USE_AGENT_CHAT = /\buseAgentChat\b/; + +export type AgentChatUiWiringDetection = { + hasNovuReact: boolean; + hasUseAgentChat: boolean; + isWired: boolean; +}; + +function listSourceFiles(projectDir: string): string[] { + const files: string[] = []; + + for (const dirName of SOURCE_DIRS) { + const dir = path.join(projectDir, dirName); + if (!fs.existsSync(dir)) { + continue; + } + + walkDir(dir, files); + } + + return files; +} + +function walkDir(dir: string, files: string[]): void { + for (const entry of fs.readdirSync(dir, { withFileTypes: true })) { + const fullPath = path.join(dir, entry.name); + if (entry.isDirectory()) { + walkDir(fullPath, files); + continue; + } + + if (/\.(tsx?|jsx?|mjs)$/.test(entry.name)) { + files.push(fullPath); + } + } +} + +export function detectAgentChatUiWiring(projectDir: string): AgentChatUiWiringDetection { + const resolvedDir = path.resolve(projectDir); + const files = listSourceFiles(resolvedDir); + + let hasNovuReact = false; + let hasUseAgentChat = false; + + for (const filePath of files) { + let contents: string; + + try { + contents = fs.readFileSync(filePath, 'utf8'); + } catch { + continue; + } + + if (NOVU_REACT_IMPORT.test(contents)) { + hasNovuReact = true; + } + + if (USE_AGENT_CHAT.test(contents)) { + hasUseAgentChat = true; + } + } + + return { + hasNovuReact, + hasUseAgentChat, + isWired: hasNovuReact && hasUseAgentChat, + }; +} diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/offer-post-connect-bridge-tunnel.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/offer-post-connect-bridge-tunnel.ts new file mode 100644 index 00000000000..febbe585e51 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/offer-post-connect-bridge-tunnel.ts @@ -0,0 +1,73 @@ +import type { + AgentConnectMode, + AiSdkConnectOutcome, + ChatSdkConnectOutcome, + ConnectAgentChatHandoff, + LangChainConnectOutcome, +} from '../../types'; +import { isBridgeConnectMode } from '../../types'; +import { promptBridgeTunnelInConsole } from '../../ui/console-bridge-reconcile-prompts'; +import { aiSdkAdapter } from '../ai-sdk/adapter'; +import { buildDevNovuScript as buildAiSdkDevNovuScript } from '../ai-sdk/dev-script'; +import { isBridgeAdapterWiringReadyForTunnel } from '../bridge-adapter/engine'; +import { isChatSdkWiringReadyForTunnel } from '../chat-sdk'; +import { buildDevNovuScript as buildChatSdkDevNovuScript } from '../chat-sdk/dev-script'; +import { langChainAdapter } from '../langchain/adapter'; + +type DeferredBridgeOutcome = AiSdkConnectOutcome | LangChainConnectOutcome | ChatSdkConnectOutcome; + +export async function offerPostConnectBridgeTunnel(input: { + connectMode: AgentConnectMode; + chatSdkOutcome?: ChatSdkConnectOutcome; + aiSdkOutcome?: AiSdkConnectOutcome; + langChainOutcome?: LangChainConnectOutcome; + agentChatHandoff?: ConnectAgentChatHandoff; + agentChatProjectDir?: string; + ci?: boolean; +}): Promise { + if (input.ci || !input.agentChatHandoff || !isBridgeConnectMode(input.connectMode)) { + return false; + } + + const bridgeOutcome = (input.chatSdkOutcome ?? input.aiSdkOutcome ?? input.langChainOutcome) as + | DeferredBridgeOutcome + | undefined; + + if (!bridgeOutcome || !('coreReady' in bridgeOutcome)) { + return false; + } + + if (!bridgeOutcome.coreReady || bridgeOutcome.tunnelAccepted || bridgeOutcome.skippedInstall) { + return false; + } + + const projectDir = input.agentChatProjectDir ?? bridgeOutcome.projectDir; + const adapter = + input.connectMode === 'langchain' ? langChainAdapter : input.connectMode === 'ai-sdk' ? aiSdkAdapter : null; + + if (input.connectMode === 'chat-sdk') { + if (!isChatSdkWiringReadyForTunnel(bridgeOutcome.requirements, projectDir, bridgeOutcome.scaffolded)) { + return false; + } + } else if (adapter) { + if ( + !isBridgeAdapterWiringReadyForTunnel(adapter, bridgeOutcome.requirements, projectDir, bridgeOutcome.scaffolded) + ) { + return false; + } + } else { + return false; + } + + const devCommand = + input.connectMode === 'chat-sdk' ? buildChatSdkDevNovuScript(projectDir) : buildAiSdkDevNovuScript(projectDir); + + const choice = await promptBridgeTunnelInConsole({ projectDir, devCommand }); + if (choice === 'accept') { + bridgeOutcome.tunnelAccepted = true; + + return true; + } + + return false; +} diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.spec.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.spec.ts new file mode 100644 index 00000000000..c129de36d01 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.spec.ts @@ -0,0 +1,56 @@ +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { describe, expect, it } from 'vitest'; +import { resolveAgentChatProjectWiringState } from './resolve-agent-chat-wiring-state'; + +function writeFile(dir: string, rel: string, contents: string) { + const full = path.join(dir, rel); + fs.mkdirSync(path.dirname(full), { recursive: true }); + fs.writeFileSync(full, contents); +} + +describe('resolveAgentChatProjectWiringState', () => { + it('returns wired when handler, UI, and env are ready for bridge runtimes', () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'novu-agent-chat-wire-')); + writeFile( + dir, + 'app/page.tsx', + "import { NovuProvider, useAgentChat } from '@novu/react';\nexport default function Page() { useAgentChat({ agentId: 'a' }); }" + ); + writeFile(dir, '.env.local', 'NEXT_PUBLIC_NOVU_APP_ID=app\nNEXT_PUBLIC_NOVU_AGENT_ID=agent\n'); + + expect( + resolveAgentChatProjectWiringState(dir, 'ai-sdk', { + projectDir: dir, + scaffolded: false, + requirements: [{ id: 'code-wiring', status: 'ok', detail: 'Wired' }], + }) + ).toBe('wired'); + }); + + it('returns partial when only env vars exist', () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'novu-agent-chat-wire-')); + writeFile(dir, '.env.local', 'NEXT_PUBLIC_NOVU_APP_ID=app\nNEXT_PUBLIC_NOVU_AGENT_ID=agent\n'); + + expect(resolveAgentChatProjectWiringState(dir, 'ai-sdk')).toBe('partial'); + }); + + it('returns unwired for managed demo runtime without UI or env', () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'novu-agent-chat-wire-')); + + expect(resolveAgentChatProjectWiringState(dir, 'demo')).toBe('unwired'); + }); + + it('returns wired for managed runtime when UI and env are both present', () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'novu-agent-chat-wire-')); + writeFile( + dir, + 'app/page.tsx', + "import { useAgentChat } from '@novu/react';\nexport default function Page() { useAgentChat({ agentId: 'a' }); }" + ); + writeFile(dir, '.env.local', 'NEXT_PUBLIC_NOVU_APP_ID=app\nNEXT_PUBLIC_NOVU_AGENT_ID=agent\n'); + + expect(resolveAgentChatProjectWiringState(dir, 'demo')).toBe('wired'); + }); +}); diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.ts new file mode 100644 index 00000000000..8f521193e74 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/resolve-agent-chat-wiring-state.ts @@ -0,0 +1,71 @@ +import type { + AgentConnectMode, + AiSdkConnectOutcome, + ChatSdkConnectOutcome, + CustomCodeConnectOutcome, + LangChainConnectOutcome, +} from '../../types'; +import { detectAgentChatUiWiring } from './detect-agent-chat-ui'; +import { hasAgentChatEnvConfigured } from './wire-agent-chat-env'; + +export type BridgeSetupSnapshot = + | AiSdkConnectOutcome + | LangChainConnectOutcome + | ChatSdkConnectOutcome + | CustomCodeConnectOutcome; + +export type AgentChatProjectWiringState = 'unwired' | 'partial' | 'wired'; + +const MANAGED_RUNTIMES = new Set(['demo', 'claude', 'claude-aws']); + +export function resolveHandlerWired(bridgeOutcome?: BridgeSetupSnapshot): boolean { + if (!bridgeOutcome) { + return false; + } + + if ('requirements' in bridgeOutcome && bridgeOutcome.requirements) { + const wiring = bridgeOutcome.requirements.find((req) => req.id === 'code-wiring'); + if (wiring) { + return wiring.status === 'ok'; + } + } + + if ('agentFilePath' in bridgeOutcome && bridgeOutcome.agentFilePath) { + return true; + } + + return false; +} + +export function resolveAgentChatProjectWiringState( + projectDir: string, + connectMode: AgentConnectMode, + bridgeOutcome?: BridgeSetupSnapshot +): AgentChatProjectWiringState { + const uiWired = detectAgentChatUiWiring(projectDir).isWired; + const envReady = hasAgentChatEnvConfigured(projectDir); + + if (MANAGED_RUNTIMES.has(connectMode)) { + if (uiWired && envReady) { + return 'wired'; + } + + if (uiWired || envReady) { + return 'partial'; + } + + return 'unwired'; + } + + const handlerWired = resolveHandlerWired(bridgeOutcome); + + if (handlerWired && uiWired && envReady) { + return 'wired'; + } + + if (handlerWired || uiWired || envReady) { + return 'partial'; + } + + return 'unwired'; +} diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/run-agent-chat-setup.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/run-agent-chat-setup.ts new file mode 100644 index 00000000000..fef52d0e673 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/run-agent-chat-setup.ts @@ -0,0 +1,183 @@ +import path from 'node:path'; +import { buildConnectEmbedPrompt, type ConnectEmbedRuntime, SUBSCRIBER_ID_PLACEHOLDER } from '@novu/shared'; +import chalk from 'chalk'; +import { resolveConnectApplicationIdentifier } from '../../auth/resolve-connect-application-identifier'; +import type { ResolvedConnectAuth } from '../../auth/resolve-connect-auth'; +import type { + AgentChatConnectOutcome, + AgentChatSetupMode, + AgentConnectMode, + AgentSummary, + ConnectAgentChatHandoff, + ConnectCommandOptions, +} from '../../types'; +import { logAgentChatEmbedPromptFileHandoffEvent, writeAgentChatEmbedPromptHandoffFile } from '../../ui/handoff-events'; +import type { ConnectUI } from '../../ui/ui'; +import { + type BridgeSetupSnapshot, + resolveAgentChatProjectWiringState, + resolveHandlerWired, +} from './resolve-agent-chat-wiring-state'; +import { + defaultAgentChatScaffoldDirName, + detectAgentChatProjectKind, + scaffoldAgentChatProject, +} from './scaffold-agent-chat'; +import { mergeAgentChatEnv } from './wire-agent-chat-env'; + +export type { BridgeSetupSnapshot } from './resolve-agent-chat-wiring-state'; +export { resolveHandlerWired } from './resolve-agent-chat-wiring-state'; + +export async function runAgentChatProjectSetup(input: { + options: ConnectCommandOptions; + ui: ConnectUI; + auth: ResolvedConnectAuth; + agent: AgentSummary; + handoff: ConnectAgentChatHandoff; + connectMode: AgentConnectMode; + bridgeOutcome?: BridgeSetupSnapshot; + bridgeProjectDir?: string; + autoMergeIntoBridge?: boolean; +}): Promise { + const projectDir = path.resolve(input.options.projectDir ?? process.cwd()); + const projectKind = detectAgentChatProjectKind(projectDir); + const explicitMode = resolveExplicitAgentChatSetupMode(input.options); + const wiringState = + projectKind === 'project' + ? resolveAgentChatProjectWiringState(projectDir, input.connectMode, input.bridgeOutcome) + : 'unwired'; + + if (!explicitMode && wiringState === 'wired') { + return runAgentChatEmbedSetup(input, projectDir, { alreadyWired: true }); + } + + const mode = + explicitMode ?? + (input.autoMergeIntoBridge && input.bridgeProjectDir + ? 'scaffold' + : input.ui.interactive + ? await input.ui.pickAgentChatSetup({ projectKind }) + : resolveAgentChatSetupMode(input.options, projectKind)); + + if (mode === 'skip') { + return { mode }; + } + + if (mode === 'embed') { + return runAgentChatEmbedSetup(input, projectDir, { alreadyWired: false }); + } + + const applicationIdentifier = await resolveConnectApplicationIdentifier(input.auth); + const subscriberId = input.auth.user?.id ?? SUBSCRIBER_ID_PLACEHOLDER; + + if (input.ui.interactive) { + await input.ui.releaseTerminal(); + console.log(chalk.cyan('Scaffolding your Agent Chat example app…')); + console.log(`${chalk.gray('Installing dependencies β€” this may take a minute.')}\n`); + } else { + input.ui.scaffoldingAgentChat(); + } + + const result = await scaffoldAgentChatProject({ + parentDir: projectDir, + appName: defaultAgentChatScaffoldDirName(input.agent.identifier), + agentIdentifier: input.agent.identifier, + applicationIdentifier, + subscriberId, + apiUrl: input.auth.apiUrl, + mergeIntoProjectDir: input.bridgeProjectDir, + mergeAtRoot: input.autoMergeIntoBridge, + }); + + return { + mode: 'scaffold', + projectDir: result.projectDir, + scaffolded: result.scaffolded, + mergedIntoBridge: result.mergedIntoBridge, + chatPath: result.chatPath, + }; +} + +async function runAgentChatEmbedSetup( + input: { + options: ConnectCommandOptions; + ui: ConnectUI; + auth: ResolvedConnectAuth; + agent: AgentSummary; + handoff: ConnectAgentChatHandoff; + connectMode: AgentConnectMode; + bridgeOutcome?: BridgeSetupSnapshot; + }, + projectDir: string, + opts: { alreadyWired: boolean } +): Promise { + const applicationIdentifier = await resolveConnectApplicationIdentifier(input.auth); + const subscriberId = input.auth.user?.id ?? SUBSCRIBER_ID_PLACEHOLDER; + const envResult = mergeAgentChatEnv({ + projectDir, + applicationIdentifier, + subscriberId, + agentIdentifier: input.agent.identifier, + apiUrl: input.auth.apiUrl, + }); + const handlerWired = opts.alreadyWired ? true : resolveHandlerWired(input.bridgeOutcome); + const embedPrompt = buildConnectEmbedPrompt({ + agentName: input.agent.name, + agentIdentifier: input.agent.identifier, + applicationIdentifier, + subscriberId, + envPaths: envResult.envPaths, + connectMode: resolveConnectEmbedRuntime(input.connectMode), + handlerWired, + }); + input.handoff.embedPrompt = embedPrompt; + + const embedPromptFile = + !opts.alreadyWired && !input.ui.interactive ? await writeAgentChatEmbedPromptHandoffFile(embedPrompt) : undefined; + + input.handoff.embedPromptFile = embedPromptFile; + + if (!input.ui.interactive && embedPromptFile) { + logAgentChatEmbedPromptFileHandoffEvent({ embedPromptFile }); + } + + return { + mode: 'embed', + projectDir, + embedPromptFile, + envPaths: envResult.envPaths, + alreadyWired: opts.alreadyWired, + }; +} + +export function resolveConnectEmbedRuntime(connectMode: AgentConnectMode): ConnectEmbedRuntime { + if ( + connectMode === 'ai-sdk' || + connectMode === 'langchain' || + connectMode === 'custom-code' || + connectMode === 'chat-sdk' || + connectMode === 'demo' || + connectMode === 'claude' || + connectMode === 'claude-aws' + ) { + return connectMode; + } + + return 'ai-sdk'; +} + +export function resolveAgentChatSetupMode( + options: ConnectCommandOptions, + projectKind: 'empty' | 'project' +): AgentChatSetupMode { + return resolveExplicitAgentChatSetupMode(options) ?? (projectKind === 'empty' ? 'scaffold' : 'embed'); +} + +function resolveExplicitAgentChatSetupMode(options: ConnectCommandOptions): AgentChatSetupMode | undefined { + const explicit = options.agentChatSetup?.trim().toLowerCase(); + if (explicit === 'scaffold' || explicit === 'embed' || explicit === 'skip') { + return explicit; + } + + return undefined; +} diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.spec.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.spec.ts new file mode 100644 index 00000000000..565f5d8e58c --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.spec.ts @@ -0,0 +1,39 @@ +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { describe, expect, it } from 'vitest'; +import { assertSafeScaffoldDirectoryName, scaffoldAgentChatProject } from './scaffold-agent-chat'; + +describe('assertSafeScaffoldDirectoryName', () => { + it('accepts a simple directory name', () => { + expect(() => assertSafeScaffoldDirectoryName('support-agent-agent-chat')).not.toThrow(); + }); + + it('rejects path traversal segments', () => { + expect(() => assertSafeScaffoldDirectoryName('../../../../tmp/malicious-agent-chat')).toThrow( + /Invalid scaffold directory name/ + ); + }); + + it('rejects absolute paths', () => { + expect(() => assertSafeScaffoldDirectoryName('/tmp/malicious-agent-chat')).toThrow( + /Invalid scaffold directory name/ + ); + }); +}); + +describe('scaffoldAgentChatProject', () => { + it('rejects unsafe agent identifiers before writing outside the parent directory', async () => { + const parentDir = fs.mkdtempSync(path.join(os.tmpdir(), 'novu-agent-chat-scaffold-')); + + await expect( + scaffoldAgentChatProject({ + parentDir, + agentIdentifier: '../../../../tmp/malicious', + applicationIdentifier: 'app-id', + subscriberId: 'subscriber-id', + apiUrl: 'http://localhost:3000', + }) + ).rejects.toThrow(/Invalid scaffold directory name/); + }); +}); diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.ts new file mode 100644 index 00000000000..8217e9096b4 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/scaffold-agent-chat.ts @@ -0,0 +1,483 @@ +import { execFileSync } from 'node:child_process'; +import fs from 'node:fs'; +import path from 'node:path'; +import { tryGitInit } from '../../../init/helpers/git'; +import { isFolderEmpty } from '../../../init/helpers/is-folder-empty'; +import { getOnline } from '../../../init/helpers/is-online'; +import { detectBridgeProject } from '../bridge/detect-project'; +import { buildAgentChatEnvEntries, mergeAgentChatEnv } from './wire-agent-chat-env'; + +export type ScaffoldAgentChatProjectInput = { + parentDir: string; + appName?: string; + agentIdentifier: string; + applicationIdentifier: string; + subscriberId: string; + apiUrl: string; + /** When set, merge chat UI into an existing bridge scaffold instead of creating a sibling app. */ + mergeIntoProjectDir?: string; + /** Replace the root page when merging into a project scaffolded during this connect run. */ + mergeAtRoot?: boolean; +}; + +export type ScaffoldAgentChatProjectResult = { + projectDir: string; + appName: string; + scaffolded: boolean; + mergedIntoBridge: boolean; + /** Browser path where Agent Chat is served inside the app. */ + chatPath: string; +}; + +const TEMPLATE_ROOT = path.join(__dirname, '../../templates/agent-chat/ts'); + +export function defaultAgentChatScaffoldDirName(agentIdentifier: string): string { + return `${agentIdentifier}-agent-chat`; +} + +export function assertSafeScaffoldDirectoryName(name: string): void { + if (path.isAbsolute(name) || path.basename(name) !== name || name === '.' || name === '..') { + throw new Error(`Invalid scaffold directory name "${name}". Use a single relative directory name.`); + } +} + +function resolveScaffoldAppName(input: ScaffoldAgentChatProjectInput): string { + const appName = input.appName?.trim() || defaultAgentChatScaffoldDirName(input.agentIdentifier); + assertSafeScaffoldDirectoryName(appName); + + return appName; +} + +export async function scaffoldAgentChatProject( + input: ScaffoldAgentChatProjectInput +): Promise { + const mergeTarget = input.mergeIntoProjectDir?.trim(); + if (mergeTarget) { + await mergeAgentChatIntoProject(mergeTarget, input); + const chatPath = input.mergeAtRoot ? '/' : '/agent-chat'; + + return { + projectDir: mergeTarget, + appName: path.basename(mergeTarget), + scaffolded: true, + mergedIntoBridge: true, + chatPath, + }; + } + + const parentDir = path.resolve(input.parentDir); + const appName = resolveScaffoldAppName(input); + const root = path.join(parentDir, appName); + + if (fs.existsSync(root) && !isFolderEmpty(root, appName)) { + throw new Error(`Cannot scaffold Agent Chat into "${root}" β€” the directory is not empty.`); + } + + fs.mkdirSync(root, { recursive: true }); + await writeStandaloneAgentChatApp(root, input); + tryGitInit(root); + + return { projectDir: root, appName, scaffolded: true, mergedIntoBridge: false, chatPath: '/' }; +} + +async function mergeAgentChatIntoProject(projectDir: string, input: ScaffoldAgentChatProjectInput): Promise { + const resolved = path.resolve(projectDir); + if (!fs.existsSync(path.join(resolved, 'package.json'))) { + throw new Error(`Cannot merge Agent Chat into "${resolved}" β€” no package.json found.`); + } + + const dependenciesChanged = ensureAgentChatDependencies(resolved, findLocalNovuDeps()); + const componentsDir = path.join(resolved, 'components', 'agent-chat'); + fs.mkdirSync(componentsDir, { recursive: true }); + copyTemplateComponents(componentsDir); + + const chatPageDir = input.mergeAtRoot ? path.join(resolved, 'app') : path.join(resolved, 'app', 'agent-chat'); + fs.mkdirSync(chatPageDir, { recursive: true }); + fs.writeFileSync( + path.join(chatPageDir, 'page.tsx'), + renderChatPage({ standalone: false, configImport: '../config' }), + 'utf8' + ); + + appendEnvExample(resolved, input); + + if (dependenciesChanged && (await getOnline())) { + execFileSync(resolvePackageManager(resolved), ['install'], { cwd: resolved, stdio: 'inherit' }); + } +} + +function ensureAgentChatDependencies(projectDir: string, localNovuDeps: LocalNovuDeps | undefined): boolean { + const packageJsonPath = path.join(projectDir, 'package.json'); + const packageJson = JSON.parse(fs.readFileSync(packageJsonPath, 'utf8')) as { + dependencies?: Record; + scripts?: Record; + }; + const dependencies = packageJson.dependencies ?? {}; + const required = { + '@novu/react': resolveNovuReactDependency(localNovuDeps), + ...(localNovuDeps ? { '@novu/js': `file:${localNovuDeps.jsDir}` } : {}), + 'react-markdown': '^10.1.0', + 'remark-gfm': '^4.0.1', + }; + let changed = false; + + for (const [name, version] of Object.entries(required)) { + if (dependencies[name] !== version) { + dependencies[name] = version; + changed = true; + } + } + + if (localNovuDeps && packageJson.scripts) { + for (const [name, script] of Object.entries(packageJson.scripts)) { + const usesNextDev = script.includes('next dev') && !script.includes('next dev --webpack'); + const usesNextBuild = script.includes('next build') && !script.includes('next build --webpack'); + if (!usesNextDev && !usesNextBuild) continue; + + packageJson.scripts[name] = script + .replace(/next dev(?=\s|$)/g, 'next dev --webpack') + .replace(/next build(?=\s|$)/g, 'next build --webpack'); + changed = true; + } + } + + if (changed) { + packageJson.dependencies = dependencies; + fs.writeFileSync(packageJsonPath, `${JSON.stringify(packageJson, null, 2)}\n`, 'utf8'); + } + + return changed; +} + +function resolvePackageManager(projectDir: string): 'npm' | 'pnpm' | 'yarn' | 'bun' { + if (fs.existsSync(path.join(projectDir, 'pnpm-lock.yaml'))) return 'pnpm'; + if (fs.existsSync(path.join(projectDir, 'yarn.lock'))) return 'yarn'; + if (fs.existsSync(path.join(projectDir, 'bun.lock')) || fs.existsSync(path.join(projectDir, 'bun.lockb'))) { + return 'bun'; + } + + return 'npm'; +} + +async function writeStandaloneAgentChatApp(root: string, input: ScaffoldAgentChatProjectInput): Promise { + fs.mkdirSync(path.join(root, 'app'), { recursive: true }); + fs.mkdirSync(path.join(root, 'components', 'agent-chat'), { recursive: true }); + + const localNovuDeps = findLocalNovuDeps(); + + copyTemplateComponents(path.join(root, 'components', 'agent-chat')); + + fs.writeFileSync(path.join(root, 'app', 'layout.tsx'), STANDALONE_LAYOUT, 'utf8'); + fs.writeFileSync( + path.join(root, 'app', 'page.tsx'), + renderChatPage({ standalone: true, configImport: '../config' }), + 'utf8' + ); + fs.writeFileSync( + path.join(root, 'config.ts'), + renderConfigModule({ + applicationIdentifier: input.applicationIdentifier, + subscriberId: input.subscriberId, + agentIdentifier: input.agentIdentifier, + apiUrl: input.apiUrl, + }), + 'utf8' + ); + fs.writeFileSync( + path.join(root, 'package.json'), + renderPackageJson(path.basename(root), root, localNovuDeps), + 'utf8' + ); + fs.writeFileSync(path.join(root, 'tsconfig.json'), STANDALONE_TSCONFIG, 'utf8'); + fs.writeFileSync(path.join(root, 'next.config.mjs'), renderNextConfig(root, localNovuDeps), 'utf8'); + fs.writeFileSync(path.join(root, '.env.local'), renderEnvLocal(input), 'utf8'); + fs.writeFileSync(path.join(root, '.env.example'), renderEnvExample(input), 'utf8'); + fs.writeFileSync(path.join(root, '.gitignore'), STANDALONE_GITIGNORE, 'utf8'); + + const isOnline = await getOnline(); + if (isOnline) { + const { execSync } = await import('node:child_process'); + execSync('npm install', { cwd: root, stdio: 'inherit' }); + } +} + +function copyTemplateComponents(targetDir: string): void { + for (const file of fs.readdirSync(TEMPLATE_ROOT)) { + if (!file.endsWith('.tsx') && !file.endsWith('.css') && !file.endsWith('.ts')) continue; + fs.copyFileSync(path.join(TEMPLATE_ROOT, file), path.join(targetDir, file)); + } +} + +function renderChatPage(opts: { standalone: boolean; configImport: string }): string { + if (opts.standalone) { + return `'use client'; + +import { NovuProvider } from '@novu/react'; +import { AgentChat } from '@/components/agent-chat/agent-chat'; +import { config } from '@/config'; +import '@/components/agent-chat/globals.css'; + +export default function Page() { + return ( + +
+ +
+
+ ); +} +`; + } + + return `'use client'; + +import { Inter } from 'next/font/google'; +import { NovuProvider } from '@novu/react'; +import { AgentChat } from '@/components/agent-chat/agent-chat'; +import '@/components/agent-chat/globals.css'; + +const inter = Inter({ subsets: ['latin'], display: 'swap' }); + +export default function AgentChatPage() { + const applicationIdentifier = process.env.NEXT_PUBLIC_NOVU_APP_ID ?? ''; + const subscriberId = process.env.NEXT_PUBLIC_NOVU_SUBSCRIBER_ID ?? ''; + const apiUrl = process.env.NEXT_PUBLIC_NOVU_BACKEND_URL ?? 'http://localhost:3000'; + const socketUrl = process.env.NEXT_PUBLIC_NOVU_SOCKET_URL; + + return ( + +
+ +
+
+ ); +} +`; +} + +function renderConfigModule(input: { + applicationIdentifier: string; + subscriberId: string; + agentIdentifier: string; + apiUrl: string; +}): string { + return `export const config = { + applicationIdentifier: process.env.NEXT_PUBLIC_NOVU_APP_ID ?? '${escapeJs(input.applicationIdentifier)}', + subscriberId: process.env.NEXT_PUBLIC_NOVU_SUBSCRIBER_ID ?? '${escapeJs(input.subscriberId)}', + agentId: process.env.NEXT_PUBLIC_NOVU_AGENT_ID ?? '${escapeJs(input.agentIdentifier)}', + backendUrl: process.env.NEXT_PUBLIC_NOVU_BACKEND_URL ?? '${escapeJs(input.apiUrl)}', + socketUrl: process.env.NEXT_PUBLIC_NOVU_SOCKET_URL, +} as const; +`; +} + +function renderEnvExample(input: ScaffoldAgentChatProjectInput): string { + const entries = buildAgentChatEnvEntries({ + projectDir: '', + applicationIdentifier: input.applicationIdentifier, + subscriberId: input.subscriberId, + agentIdentifier: input.agentIdentifier, + apiUrl: input.apiUrl, + }); + const lines = ['# Copy to .env.local', '# Agent Chat (added by npx novu connect)']; + for (const [key, value] of Object.entries(entries)) { + lines.push(`${key}=${value}`); + } + + return `${lines.join('\n')}\n`; +} + +function renderEnvLocal(input: ScaffoldAgentChatProjectInput): string { + return renderEnvExample(input).replace('# Copy to .env.local\n', ''); +} + +function appendEnvExample(projectDir: string, input: ScaffoldAgentChatProjectInput): void { + mergeAgentChatEnv({ + projectDir, + applicationIdentifier: input.applicationIdentifier, + subscriberId: input.subscriberId, + agentIdentifier: input.agentIdentifier, + apiUrl: input.apiUrl, + }); +} + +function findMonorepoReactPackageDir(): string | undefined { + let dir = __dirname; + for (let depth = 0; depth < 12; depth++) { + const candidate = path.join(dir, 'packages', 'react', 'package.json'); + if (fs.existsSync(candidate)) { + return path.dirname(candidate); + } + + const parent = path.dirname(dir); + if (parent === dir) break; + dir = parent; + } + + return undefined; +} + +type LocalNovuDeps = { + reactDir: string; + jsDir: string; +}; + +function findLocalNovuDeps(): LocalNovuDeps | undefined { + const reactDir = findMonorepoReactPackageDir(); + if (!reactDir) { + return undefined; + } + + const jsDir = path.join(path.dirname(reactDir), 'js'); + if (!fs.existsSync(path.join(jsDir, 'package.json'))) { + return undefined; + } + + return { reactDir, jsDir }; +} + +function toPosixRelative(from: string, to: string): string { + const rel = path.relative(from, to).split(path.sep).join('/'); + return rel.startsWith('.') ? rel : `./${rel}`; +} + +function resolveNovuReactDependency(localNovuDeps: LocalNovuDeps | undefined): string { + if (localNovuDeps) { + return `file:${localNovuDeps.reactDir}`; + } + + return 'latest'; +} + +function renderPackageJson(name: string, scaffoldRoot: string, localNovuDeps: LocalNovuDeps | undefined): string { + const usesLocalNovu = Boolean(localNovuDeps); + + return JSON.stringify( + { + name, + private: true, + scripts: { + dev: usesLocalNovu ? 'next dev -p 4012 --webpack' : 'next dev -p 4012', + build: usesLocalNovu ? 'next build --webpack' : 'next build', + start: 'next start -p 4012', + }, + dependencies: { + '@novu/react': resolveNovuReactDependency(localNovuDeps), + next: '^16.2.11', + react: '^18.3.1', + 'react-dom': '^18.3.1', + 'react-markdown': '^10.1.0', + 'remark-gfm': '^4.0.1', + }, + devDependencies: { + '@types/node': '^22.0.0', + '@types/react': '^19.0.0', + '@types/react-dom': '^19.0.0', + typescript: '5.6.2', + }, + }, + null, + 2 + ); +} + +function renderNextConfig(scaffoldRoot: string, localNovuDeps: LocalNovuDeps | undefined): string { + if (!localNovuDeps) { + return STANDALONE_NEXT_CONFIG; + } + + const reactRel = toPosixRelative(scaffoldRoot, localNovuDeps.reactDir); + const jsRel = toPosixRelative(scaffoldRoot, localNovuDeps.jsDir); + + return `import path from 'node:path'; +import { fileURLToPath } from 'node:url'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +const novuReact = path.resolve(__dirname, '${reactRel}'); +const novuJs = path.resolve(__dirname, '${jsRel}'); + +/** @type {import('next').NextConfig} */ +const nextConfig = { + transpilePackages: ['@novu/react', '@novu/js'], + // Turbopack cannot resolve file: symlinks outside the app root β€” webpack aliases + // point at the local monorepo packages when connect scaffolds from a dev checkout. + webpack: (config) => { + config.resolve.alias = { + ...config.resolve.alias, + '@novu/react': novuReact, + '@novu/js': novuJs, + }; + + return config; + }, +}; + +export default nextConfig; +`; +} + +function escapeJs(value: string): string { + return value.replace(/\\/g, '\\\\').replace(/'/g, "\\'"); +} + +export function detectAgentChatProjectKind(projectDir: string): 'empty' | 'project' { + return detectBridgeProject(projectDir).kind; +} + +const STANDALONE_LAYOUT = `import { Inter } from 'next/font/google'; + +const inter = Inter({ subsets: ['latin'], display: 'swap' }); + +export const metadata = { + title: 'Novu Agent Chat', + description: 'A standalone Agent Chat example powered by Novu.', +}; + +export default function RootLayout({ children }: { children: React.ReactNode }) { + return ( + + + {children} + + + ); +} +`; + +const STANDALONE_TSCONFIG = JSON.stringify( + { + compilerOptions: { + target: 'ES2017', + lib: ['dom', 'dom.iterable', 'esnext'], + allowJs: true, + skipLibCheck: true, + strict: true, + noEmit: true, + esModuleInterop: true, + module: 'esnext', + moduleResolution: 'bundler', + resolveJsonModule: true, + isolatedModules: true, + jsx: 'preserve', + incremental: true, + plugins: [{ name: 'next' }], + paths: { '@/*': ['./*'] }, + }, + include: ['next-env.d.ts', '**/*.ts', '**/*.tsx'], + exclude: ['node_modules'], + }, + null, + 2 +); + +const STANDALONE_NEXT_CONFIG = `/** @type {import('next').NextConfig} */ +const nextConfig = {}; +export default nextConfig; +`; + +const STANDALONE_GITIGNORE = `.next\nnode_modules\n.env.local\n`; diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/wire-agent-chat-env.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/wire-agent-chat-env.ts new file mode 100644 index 00000000000..47aed686ea9 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/wire-agent-chat-env.ts @@ -0,0 +1,165 @@ +import fs from 'node:fs'; +import { parseEnvFile, resolveProjectEnvPaths } from '../bridge/wire-env'; + +export type AgentChatEnvInput = { + projectDir: string; + applicationIdentifier: string; + subscriberId: string; + agentIdentifier: string; + apiUrl: string; +}; + +export type AgentChatEnvMergeResult = { + envPaths: string[]; + updatedKeys: string[]; +}; + +const AGENT_CHAT_ENV_KEYS = [ + 'NEXT_PUBLIC_NOVU_APP_ID', + 'NEXT_PUBLIC_NOVU_SUBSCRIBER_ID', + 'NEXT_PUBLIC_NOVU_AGENT_ID', + 'NEXT_PUBLIC_NOVU_BACKEND_URL', + 'NEXT_PUBLIC_NOVU_SOCKET_URL', +] as const; + +const DEFAULT_US_API_URL = 'https://api.novu.co'; + +const API_HOST_SOCKET_URL: Record = { + 'eu.api.novu.co': 'wss://eu.socket.novu.co', + 'api.novu-staging.co': 'wss://socket.novu-staging.co', +}; + +export function resolveScaffoldSocketUrl(apiUrl: string): string | undefined { + try { + const hostname = new URL(apiUrl).hostname; + if (hostname === 'localhost' || hostname === '127.0.0.1' || hostname.endsWith('.novu.localhost')) { + return 'ws://127.0.0.1:8787'; + } + + return API_HOST_SOCKET_URL[hostname]; + } catch { + return undefined; + } +} + +export function buildAgentChatEnvEntries(input: AgentChatEnvInput): Record { + const backendUrl = input.apiUrl.replace(/\/$/, ''); + const socketUrl = resolveScaffoldSocketUrl(input.apiUrl); + const entries: Record = { + NEXT_PUBLIC_NOVU_APP_ID: input.applicationIdentifier, + NEXT_PUBLIC_NOVU_SUBSCRIBER_ID: input.subscriberId, + NEXT_PUBLIC_NOVU_AGENT_ID: input.agentIdentifier, + }; + + // US Cloud defaults to https://api.novu.co β€” omit so NovuProvider uses SDK defaults. + if (backendUrl !== DEFAULT_US_API_URL) { + entries.NEXT_PUBLIC_NOVU_BACKEND_URL = backendUrl; + } + + // Only non-default sockets (local, EU, staging). US Cloud defaults to wss://socket.novu.co. + if (socketUrl) { + entries.NEXT_PUBLIC_NOVU_SOCKET_URL = socketUrl; + } + + return entries; +} + +function mergeEnvFileAtPath(envPath: string, entriesToSet: Record): string[] { + const created = !fs.existsSync(envPath); + const existingContents = created ? '' : fs.readFileSync(envPath, 'utf8'); + const entries = parseEnvFile(existingContents); + const updatedKeys: string[] = []; + + for (const [key, value] of Object.entries(entriesToSet)) { + if (entries.get(key) === value) { + continue; + } + + entries.set(key, value); + updatedKeys.push(key); + } + + if (updatedKeys.length === 0 && !created) { + return updatedKeys; + } + + fs.writeFileSync(envPath, serializeEnvWithComment(existingContents, entries, created)); + + return updatedKeys; +} + +function serializeEnvWithComment(existingContents: string, entries: Map, created: boolean): string { + if (created || !existingContents.trim()) { + const lines = ['# Agent Chat (added by npx novu connect)']; + for (const key of AGENT_CHAT_ENV_KEYS) { + const value = entries.get(key); + if (value) { + lines.push(`${key}=${value}`); + } + } + + return `${lines.join('\n')}\n`; + } + + const lines = existingContents.replace(/\s+$/, '').split(/\r?\n/); + const preserved = lines.filter((line) => { + const trimmed = line.trim(); + if (!trimmed || trimmed.startsWith('#')) { + return true; + } + + const eq = trimmed.indexOf('='); + if (eq <= 0) { + return true; + } + + const key = trimmed.slice(0, eq).trim(); + + return !(AGENT_CHAT_ENV_KEYS as readonly string[]).includes(key); + }); + + preserved.push('', '# Agent Chat (added by npx novu connect)'); + for (const key of AGENT_CHAT_ENV_KEYS) { + const value = entries.get(key); + if (value) { + preserved.push(`${key}=${value}`); + } + } + + return `${preserved.join('\n')}\n`; +} + +export function mergeAgentChatEnv(input: AgentChatEnvInput): AgentChatEnvMergeResult { + const envPaths = resolveProjectEnvPaths(input.projectDir); + const entriesToSet = buildAgentChatEnvEntries(input); + const updatedKeys = new Set(); + + for (const envPath of envPaths) { + for (const key of mergeEnvFileAtPath(envPath, entriesToSet)) { + updatedKeys.add(key); + } + } + + return { + envPaths, + updatedKeys: [...updatedKeys], + }; +} + +export function hasAgentChatEnvConfigured(projectDir: string): boolean { + for (const envPath of resolveProjectEnvPaths(projectDir)) { + if (!fs.existsSync(envPath)) { + continue; + } + + const entries = parseEnvFile(fs.readFileSync(envPath, 'utf8')); + const appId = entries.get('NEXT_PUBLIC_NOVU_APP_ID')?.trim(); + const agentId = entries.get('NEXT_PUBLIC_NOVU_AGENT_ID')?.trim(); + + if (appId && agentId) { + return true; + } + } + + return false; +} diff --git a/packages/novu/src/commands/connect/pipeline/agent-chat/wrap-ui-for-agent-chat-handoff.ts b/packages/novu/src/commands/connect/pipeline/agent-chat/wrap-ui-for-agent-chat-handoff.ts new file mode 100644 index 00000000000..d11126b611e --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/agent-chat/wrap-ui-for-agent-chat-handoff.ts @@ -0,0 +1,39 @@ +import type { ConnectUI } from '../../ui/ui'; + +export type AgentChatHandoffUiPolicy = { + suppressBridgeScaffoldSummary: boolean; + suppressBridgeReconcilePlan: boolean; + suppressBridgeTunnel: boolean; +}; + +export function wrapUiForAgentChatHandoff(ui: ConnectUI, policy: AgentChatHandoffUiPolicy): ConnectUI { + if (!policy.suppressBridgeScaffoldSummary && !policy.suppressBridgeReconcilePlan && !policy.suppressBridgeTunnel) { + return ui; + } + + return { + ...ui, + bridgeScaffolded: policy.suppressBridgeScaffoldSummary ? () => {} : ui.bridgeScaffolded.bind(ui), + showBridgeReconcilePlan: policy.suppressBridgeReconcilePlan ? async () => {} : ui.showBridgeReconcilePlan.bind(ui), + offerBridgeTunnel: policy.suppressBridgeTunnel ? async () => 'decline' as const : ui.offerBridgeTunnel.bind(ui), + }; +} + +export function resolveAgentChatHandoffUiPolicy(input: { + channel: string; + agentChatHandoff: boolean; + agentChatSetup?: string; +}): AgentChatHandoffUiPolicy | null { + if (!input.agentChatHandoff) { + return null; + } + + const setup = input.agentChatSetup?.trim().toLowerCase(); + const skipCombinedScaffoldSummary = input.channel === 'agent-chat' && setup !== 'embed' && setup !== 'skip'; + + return { + suppressBridgeScaffoldSummary: skipCombinedScaffoldSummary, + suppressBridgeReconcilePlan: true, + suppressBridgeTunnel: true, + }; +} diff --git a/packages/novu/src/commands/connect/pipeline/bridge-adapter/engine.ts b/packages/novu/src/commands/connect/pipeline/bridge-adapter/engine.ts index 07f5e790858..55f8e3a6125 100644 --- a/packages/novu/src/commands/connect/pipeline/bridge-adapter/engine.ts +++ b/packages/novu/src/commands/connect/pipeline/bridge-adapter/engine.ts @@ -5,6 +5,7 @@ import { defaultCustomCodeScaffoldDirName, resolveAgentHandlerPathIfExists } fro import { resolveEmptyDirScaffoldTarget } from '../bridge/confirm-empty-dir-scaffold'; import { requireConnectSecretKey } from '../bridge/require-secret-key'; import { runScaffoldWithConsole } from '../bridge/run-scaffold-with-console'; +import { describeLlmAuthChoice } from '../llm-auth/llm-auth-options'; import { resolveLlmAuthChoice } from '../llm-auth/resolve-llm-auth'; import type { LlmAuthChoice } from '../llm-auth/types'; import type { @@ -20,6 +21,12 @@ type ReconcileOptions = { agentFilePath?: string; }; +function deferInkInput(): Promise { + return new Promise((resolve) => { + setTimeout(resolve, 0); + }); +} + export async function runBridgeAdapterProjectSetup( input: BridgeAdapterSetupInput, adapter: BridgeAdapter @@ -56,10 +63,13 @@ export async function runBridgeAdapterProjectSetup( ui: input.ui, }); + await deferInkInput(); + const confirmed = await input.ui.confirmScaffold({ projectDir: target.projectDir, appName: target.appName, variant: adapter.variant, + llmAuthLabel: describeLlmAuthChoice(llmAuth), }); if (!confirmed) { @@ -352,7 +362,9 @@ async function promptTunnelIfReady(opts: { return false; } - if (!isWiringReadyForTunnel(opts.adapter, opts.reconcilePlan.requirements, opts.reconcilePlan.projectDir)) { + if ( + !isBridgeAdapterWiringReadyForTunnel(opts.adapter, opts.reconcilePlan.requirements, opts.reconcilePlan.projectDir) + ) { return false; } @@ -403,14 +415,14 @@ function shouldRunTunnel( return false; } - if (!isWiringReadyForTunnel(adapter, outcome.requirements, outcome.projectDir, outcome.scaffolded)) { + if (!isBridgeAdapterWiringReadyForTunnel(adapter, outcome.requirements, outcome.projectDir, outcome.scaffolded)) { return false; } return outcome.tunnelAccepted === true; } -function isWiringReadyForTunnel( +export function isBridgeAdapterWiringReadyForTunnel( adapter: BridgeAdapter, requirements: BridgeRequirement[] | undefined, projectDir: string, diff --git a/packages/novu/src/commands/connect/pipeline/bridge/print-bridge-dev-next-steps.ts b/packages/novu/src/commands/connect/pipeline/bridge/print-bridge-dev-next-steps.ts index 6d5926159a9..3341cb413b5 100644 --- a/packages/novu/src/commands/connect/pipeline/bridge/print-bridge-dev-next-steps.ts +++ b/packages/novu/src/commands/connect/pipeline/bridge/print-bridge-dev-next-steps.ts @@ -1,10 +1,7 @@ import { cyan, dim } from 'picocolors'; -export function printBridgeDevNextSteps(opts: { projectDir: string; skippedInstall?: boolean }): void { - const cmd = opts.skippedInstall - ? `cd ${opts.projectDir} && npm install && npm run dev:novu` - : `cd ${opts.projectDir} && npm run dev:novu`; - const cmdLine = `$ ${cmd}`; +export function printDevCommandBox(command: string): void { + const cmdLine = `$ ${command}`; const innerWidth = Math.max(cmdLine.length + 4, 50); console.log(); @@ -14,6 +11,14 @@ export function printBridgeDevNextSteps(opts: { projectDir: string; skippedInsta console.log(dim(` β”‚${' '.repeat(innerWidth)}β”‚`)); console.log(dim(` β•°${'─'.repeat(innerWidth)}β•―`)); console.log(); +} + +export function printBridgeDevNextSteps(opts: { projectDir: string; skippedInstall?: boolean }): void { + const command = opts.skippedInstall + ? `cd ${opts.projectDir} && npm install && npm run dev:novu` + : `cd ${opts.projectDir} && npm run dev:novu`; + + printDevCommandBox(command); console.log(` ${dim('npm run dev')} ${dim('Start app without tunnel')}`); console.log(` ${dim('npm run dev:novu')} ${dim('Start app + dev tunnel')}`); console.log(); diff --git a/packages/novu/src/commands/connect/pipeline/channels/agent-chat.ts b/packages/novu/src/commands/connect/pipeline/channels/agent-chat.ts new file mode 100644 index 00000000000..78f0e0d1ebe --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/channels/agent-chat.ts @@ -0,0 +1,75 @@ +import { buildConnectAgentChatDashboardUrl } from '@novu/shared'; +import { CONNECT_EVENTS } from '../../analytics/events'; +import { addAgentChatIntegration } from '../../api/agents'; +import type { ConnectApiClient } from '../../api/client'; +import { NovuApiError } from '../../api/client'; +import type { IntegrationRecord } from '../../api/integrations'; +import type { ResolvedConnectAuth } from '../../auth/resolve-connect-auth'; +import type { AgentSummary, ConnectAgentChatHandoff, ConnectCommandOptions } from '../../types'; +import { logAgentChatDashboardUrlHandoffEvent } from '../../ui/handoff-events'; +import type { ConnectUI } from '../../ui/ui'; + +export async function connectAgentChatForAgent( + client: ConnectApiClient, + agent: AgentSummary, + ui: ConnectUI, + options: ConnectCommandOptions, + auth: ResolvedConnectAuth, + track: (event: string, data?: Record) => void +): Promise<{ + integration: IntegrationRecord; + handoff: ConnectAgentChatHandoff; +}> { + ui.addingAgentChatIntegration(); + + let link; + try { + link = await addAgentChatIntegration(client, agent.identifier); + } catch (err) { + throw normalizeAgentChatProvisionError(err); + } + + const integration: IntegrationRecord = { + _id: link.integration._id, + identifier: link.integration.identifier, + name: link.integration.name, + providerId: link.integration.providerId, + channel: 'chat', + active: link.integration.active !== false, + }; + + track(CONNECT_EVENTS.AGENT_CHAT_LINKED, { + agent: agent.identifier, + alreadyLinked: Boolean(link.connectedAt), + }); + + const dashboardUrl = buildConnectAgentChatDashboardUrl({ + connectDashboardUrl: options.connectDashboardUrl, + environmentSlug: auth.environmentSlug ?? null, + agentIdentifier: agent.identifier, + }); + + const handoff: ConnectAgentChatHandoff = { + dashboardUrl, + embedPrompt: '', + }; + + await ui.awaitAgentChatHandoff({ + dashboardUrl, + embedPrompt: handoff.embedPrompt, + }); + + if (!ui.interactive) { + logAgentChatDashboardUrlHandoffEvent({ dashboardUrl }); + } + + return { integration, handoff }; +} + +function normalizeAgentChatProvisionError(err: unknown): Error { + if (err instanceof NovuApiError && err.status === 403) { + return new Error(err.message); + } + + return err instanceof Error ? err : new Error(String(err)); +} diff --git a/packages/novu/src/commands/connect/pipeline/chat-sdk/index.ts b/packages/novu/src/commands/connect/pipeline/chat-sdk/index.ts index 7c0b3ce1da7..1a498ef67af 100644 --- a/packages/novu/src/commands/connect/pipeline/chat-sdk/index.ts +++ b/packages/novu/src/commands/connect/pipeline/chat-sdk/index.ts @@ -390,7 +390,7 @@ function shouldRunChatSdkTunnel( return outcome.tunnelAccepted === true; } -function isChatSdkWiringReadyForTunnel( +export function isChatSdkWiringReadyForTunnel( requirements: BridgeRequirement[] | undefined, projectDir: string, scaffolded = false diff --git a/packages/novu/src/commands/connect/pipeline/llm-auth/ensure-subscription-auth.ts b/packages/novu/src/commands/connect/pipeline/llm-auth/ensure-subscription-auth.ts index a6930d09280..8738bcf6273 100644 --- a/packages/novu/src/commands/connect/pipeline/llm-auth/ensure-subscription-auth.ts +++ b/packages/novu/src/commands/connect/pipeline/llm-auth/ensure-subscription-auth.ts @@ -24,6 +24,8 @@ function printBrowserAuthHint(): void { async function ensureCodexCliAuth(ui: ConnectUI, ci?: boolean): Promise { if (hasCodexCliAuth()) { + console.log(chalk.green('βœ“ Codex CLI login already on this machine β€” using your existing session.')); + return; } @@ -79,6 +81,8 @@ async function ensureLangchainCodexOauthAuth(ui: ConnectUI, ci?: boolean): Promi async function ensureClaudeCodeAuth(ui: ConnectUI, ci?: boolean): Promise { if (hasClaudeCodeAuth()) { + console.log(chalk.green('βœ“ Claude Code login already on this machine β€” using your existing session.')); + return; } diff --git a/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth-options.ts b/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth-options.ts index c99280b4920..dd4c707f063 100644 --- a/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth-options.ts +++ b/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth-options.ts @@ -1,5 +1,5 @@ import type { BridgeAdapterVariant } from '../bridge-adapter/types'; -import type { LlmAuthPickerOption } from './types'; +import type { LlmAuthChoice, LlmAuthPickerOption } from './types'; const OPENAI_API_KEY_OPTION: LlmAuthPickerOption = { kind: 'openai-api-key', @@ -46,3 +46,18 @@ export const LLM_AUTH_PICKER_TITLE = 'How do you want to power your agent?'; export const LLM_AUTH_PICKER_SUBTITLE = 'Optional for local dev. API keys are saved to .env.local; subscriptions use CLI OAuth on your machine.'; + +export function describeLlmAuthChoice(llmAuth: LlmAuthChoice): string { + switch (llmAuth.kind) { + case 'openai-api-key': + return 'OpenAI API key'; + case 'anthropic-api-key': + return 'Anthropic API key'; + case 'codex-subscription': + return 'ChatGPT subscription (Codex CLI)'; + case 'claude-subscription': + return 'Claude subscription (Claude Code)'; + case 'skip': + return 'Demo echo (no LLM wired yet)'; + } +} diff --git a/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth.spec.ts b/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth.spec.ts index 2f2fb3fa779..f6a764e5845 100644 --- a/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth.spec.ts +++ b/packages/novu/src/commands/connect/pipeline/llm-auth/llm-auth.spec.ts @@ -1,7 +1,7 @@ import { describe, expect, it } from 'vitest'; import { generateAgentNextConfigSource } from './codegen/generate-agent-next-config'; import { generateSupportAgentSource } from './codegen/generate-support-agent'; -import { getLlmAuthPickerOptions } from './llm-auth-options'; +import { describeLlmAuthChoice, getLlmAuthPickerOptions } from './llm-auth-options'; import { resolveLlmAuthEnvVars, resolveLlmAuthPackageDependencies, resolveLlmAuthPackages } from './registry'; describe('llm-auth registry', () => { @@ -87,6 +87,11 @@ describe('llm-auth picker options', () => { expect(aiSdkKinds).toContain('claude-subscription'); expect(langChainKinds).not.toContain('claude-subscription'); }); + + it('describes scaffold wiring choices for confirm screens', () => { + expect(describeLlmAuthChoice({ kind: 'claude-subscription' })).toContain('Claude Code'); + expect(describeLlmAuthChoice({ kind: 'skip' })).toContain('echo'); + }); }); describe('generateSupportAgentSource', () => { diff --git a/packages/novu/src/commands/connect/pipeline/llm-auth/resolve-llm-auth.ts b/packages/novu/src/commands/connect/pipeline/llm-auth/resolve-llm-auth.ts index fa357954d8d..e132b7ffa2a 100644 --- a/packages/novu/src/commands/connect/pipeline/llm-auth/resolve-llm-auth.ts +++ b/packages/novu/src/commands/connect/pipeline/llm-auth/resolve-llm-auth.ts @@ -98,6 +98,10 @@ async function resolveFromCliFlags(input: ResolveLlmAuthInput): Promise { const kind = await input.ui.pickLlmAuthKind({ connectMode: input.connectMode }); + await new Promise((resolve) => { + setTimeout(resolve, 0); + }); + if (kind === 'skip') { return { kind: 'skip' }; } diff --git a/packages/novu/src/commands/connect/pipeline/runner.ts b/packages/novu/src/commands/connect/pipeline/runner.ts index bb912135c4c..36b4234d713 100644 --- a/packages/novu/src/commands/connect/pipeline/runner.ts +++ b/packages/novu/src/commands/connect/pipeline/runner.ts @@ -18,11 +18,13 @@ import { buildConnectAgentDetailsUrl, buildConnectClaimUrl, channelDisplayName } import { ConnectChannelBackError } from '../errors'; import { shouldUpgradeFromKeylessGenerateLimit } from '../keyless-limit-errors'; import type { + AgentChatConnectOutcome, AgentConnectMode, AgentSummary, AiSdkConnectOutcome, ChannelChoice, ChatSdkConnectOutcome, + ConnectAgentChatHandoff, ConnectCommandOptions, CustomCodeConnectOutcome, LangChainConnectOutcome, @@ -34,8 +36,15 @@ import { isVanillaCustomCodeConnectMode, } from '../types'; import type { ConnectUI } from '../ui/ui'; +import { offerPostConnectBridgeTunnel } from './agent-chat/offer-post-connect-bridge-tunnel'; +import { runAgentChatProjectSetup } from './agent-chat/run-agent-chat-setup'; +import { + resolveAgentChatHandoffUiPolicy, + wrapUiForAgentChatHandoff, +} from './agent-chat/wrap-ui-for-agent-chat-handoff'; import { maybeRunAiSdkTunnel, runAiSdkProjectSetup } from './ai-sdk'; import { createBridgeAgentFlow } from './bridge/create-bridge-agent'; +import { connectAgentChatForAgent } from './channels/agent-chat'; import { connectEmailForAgent } from './channels/email'; import { connectSendblueForAgent } from './channels/sendblue'; import { connectSlackForAgent } from './channels/slack'; @@ -158,6 +167,8 @@ export async function runConnectPipeline(input: ConnectPipelineInput): Promise { const { options, ui, track, sessionProps } = ctx; diff --git a/packages/novu/src/commands/connect/templates/agent-chat/ts/agent-chat.tsx b/packages/novu/src/commands/connect/templates/agent-chat/ts/agent-chat.tsx new file mode 100644 index 00000000000..553a6bda9b9 --- /dev/null +++ b/packages/novu/src/commands/connect/templates/agent-chat/ts/agent-chat.tsx @@ -0,0 +1,42 @@ +'use client'; + +import { useAgentChat } from '@novu/react'; +import { ChatPanel } from './chat-panel'; + +export function AgentChat() { + const agentId = process.env.NEXT_PUBLIC_NOVU_AGENT_ID ?? ''; + const subscriberId = process.env.NEXT_PUBLIC_NOVU_SUBSCRIBER_ID ?? ''; + + const { + messages, + pendingActions = [], + sendMessage, + sendAction, + respondToAction, + error, + isRunning, + isLoading, + typing, + hasMore, + isFetching, + fetchMore, + } = useAgentChat({ agentId }); + + return ( + + ); +} diff --git a/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-panel.tsx b/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-panel.tsx new file mode 100644 index 00000000000..2032911258b --- /dev/null +++ b/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-panel.tsx @@ -0,0 +1,83 @@ +'use client'; + +import type { AgentConversationTyping, AgentMessage, AgentPendingAction, UseAgentChatResult } from '@novu/react'; +import { ChatThread } from './chat-thread'; +import { Composer } from './composer'; +import { PendingActionCard } from './pending-action-card'; + +/** + * Presentational shell. Swap this (and the components it uses) for your own UI β€” + * it only consumes values from `useAgentChat`. + */ +export type ChatPanelProps = { + subscriberId: string; + error?: { message: string }; + messages: AgentMessage[]; + pendingActions: AgentPendingAction[]; + isRunning: boolean; + isLoading: boolean; + typing?: AgentConversationTyping; + hasMore: boolean; + isFetching: boolean; + onFetchMore: () => Promise; + onRespond: UseAgentChatResult['respondToAction']; + onCardAction: UseAgentChatResult['sendAction']; + onSend: UseAgentChatResult['sendMessage']; +}; + +export function ChatPanel({ + subscriberId, + error, + messages, + pendingActions, + isRunning, + isLoading, + typing, + hasMore, + isFetching, + onFetchMore, + onRespond, + onCardAction, + onSend, +}: ChatPanelProps) { + const interactionDisabled = isRunning || isLoading; + + return ( +
+
+

+ Chatting as {subscriberId} +

+
+ + + +
+
+ {pendingActions.map((action) => ( + + ))} + + {error ? ( +
+ {error.message} +
+ ) : null} + + +
+
+
+ ); +} diff --git a/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-thread.tsx b/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-thread.tsx new file mode 100644 index 00000000000..27b4428b32a --- /dev/null +++ b/packages/novu/src/commands/connect/templates/agent-chat/ts/chat-thread.tsx @@ -0,0 +1,132 @@ +'use client'; + +import type { AgentMessage, UseAgentChatResult } from '@novu/react'; +import { useCallback, useEffect, useRef } from 'react'; +import { ChatIcon } from './icons'; +import { MessageRow } from './message-bubble'; + +/** Matches `AgentConversationTyping` from `@novu/react` / the event protocol. */ +type TypingState = { status?: string }; + +const STARTER_PROMPTS = ['Hello', 'What can you do?', 'List my MCP tools'] as const; + +type ChatThreadProps = { + messages: AgentMessage[]; + /** Durable run state: set by live run lifecycle events and by history replay after a reload. */ + isRunning: boolean; + isLoading: boolean; + /** Live `channel.typing` from `useAgentChat`. Absent when the agent is idle. */ + typing?: TypingState; + hasMore: boolean; + isFetching: boolean; + onFetchMore: () => Promise; + onCardAction: UseAgentChatResult['sendAction']; + cardActionsDisabled: boolean; + onSend: UseAgentChatResult['sendMessage']; +}; + +function AgentStatusRow({ status }: { status?: string }) { + // Server statuses often arrive with their own trailing ellipsis or dots. + const label = status?.trim().replace(/[.\u2026]+$/, ''); + + return ( + + + + + + {label ? ( + {label}… + ) : ( + + + + + + )} + + + ); +} + +export function ChatThread({ + messages, + isRunning, + isLoading, + typing, + hasMore, + isFetching, + onFetchMore, + onCardAction, + cardActionsDisabled, + onSend, +}: ChatThreadProps) { + const scrollRef = useRef(null); + const lastMessage = messages[messages.length - 1]; + + // Keyed on the tail, not the count: an older page prepends and must not scroll. + useEffect(() => { + const container = scrollRef.current; + if (!container) return; + + container.scrollTop = container.scrollHeight; + }, [lastMessage?.id, lastMessage?.parts, typing, isRunning]); + + const loadOlder = useCallback(async () => { + const container = scrollRef.current; + const heightBefore = container?.scrollHeight ?? 0; + + await onFetchMore(); + + // Hold the reading position: the prepended page grows the thread upwards. + requestAnimationFrame(() => { + if (!container) return; + container.scrollTop += container.scrollHeight - heightBefore; + }); + }, [onFetchMore]); + + return ( +
+
+ {hasMore ? ( +
+ +
+ ) : null} + + {messages.length === 0 && !isRunning && !isLoading ? ( +
+
+ +
+

Your agent is ready

+

Send a message to see how it replies.

+
+ {STARTER_PROMPTS.map((prompt) => ( + + ))} +
+
+ ) : ( + messages.map((message, index) => ( + + )) + )} + {typing || isRunning || isLoading ? ( + + ) : null} +
+
+ ); +} diff --git a/packages/novu/src/commands/connect/templates/agent-chat/ts/composer.tsx b/packages/novu/src/commands/connect/templates/agent-chat/ts/composer.tsx new file mode 100644 index 00000000000..3792407ad8f --- /dev/null +++ b/packages/novu/src/commands/connect/templates/agent-chat/ts/composer.tsx @@ -0,0 +1,80 @@ +'use client'; + +import type { UseAgentChatResult } from '@novu/react'; +import { FormEvent, useEffect, useLayoutEffect, useRef, useState } from 'react'; +import { SendIcon } from './icons'; + +type ComposerProps = { + isLoading: boolean; + isRunning: boolean; + onSend: UseAgentChatResult['sendMessage']; +}; + +const MAX_HEIGHT_PX = 128; + +export function Composer({ isLoading, isRunning, onSend }: ComposerProps) { + const [draft, setDraft] = useState(''); + const inputRef = useRef(null); + const disabled = isLoading || isRunning; + + useEffect(() => { + if (!disabled) { + inputRef.current?.focus(); + } + }, [disabled]); + + useLayoutEffect(() => { + const input = inputRef.current; + if (!input) return; + + input.style.height = '0px'; + input.style.height = `${Math.min(input.scrollHeight, MAX_HEIGHT_PX)}px`; + }, [draft]); + + function submit(event: FormEvent) { + event.preventDefault(); + const text = draft.trim(); + if (!text || disabled) return; + setDraft(''); + void onSend(text); + requestAnimationFrame(() => inputRef.current?.focus()); + } + + return ( +
+
+