From 5da36bf23e8861db33af995eaceb62372d2ac9b5 Mon Sep 17 00:00:00 2001 From: Adam Chmara Date: Fri, 21 Aug 2026 12:24:17 +0200 Subject: [PATCH 1/2] fix(dashboard,novu): improve novu connect dashboard commands fixes NV-8636 (#12413) --- .../agents/agent-chat-setup-content.tsx | 1 + .../agents/agent-chat-setup-guide.tsx | 1 + .../agents/agent-chat-setup-steps.tsx | 13 ++- .../agents/agent-code-setup-section.tsx | 18 +++- .../agent-chat-agent-integration-guide.tsx | 2 +- .../pipeline/resolve-existing-agent.ts | 40 +++++++++ .../src/commands/connect/pipeline/runner.ts | 34 +++++++- packages/novu/src/commands/connect/types.ts | 2 + packages/novu/src/index.ts | 5 +- .../utils/agent-chat-connect-prompt.spec.ts | 27 +++++- .../src/utils/agent-chat-connect-prompt.ts | 36 +++++++- .../shared/src/utils/novu-connect-cli.spec.ts | 28 ++++++ packages/shared/src/utils/novu-connect-cli.ts | 86 +++++++++++++++++-- 13 files changed, 269 insertions(+), 24 deletions(-) create mode 100644 packages/novu/src/commands/connect/pipeline/resolve-existing-agent.ts 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 638880d1063..13d30f9bad9 100644 --- a/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx +++ b/apps/dashboard/src/components/agents/agent-chat-setup-content.tsx @@ -3,6 +3,7 @@ export { APPLICATION_IDENTIFIER_PLACEHOLDER, buildAgentChatPrompt, buildAgentChatTuiCommand, + buildAgentChatTuiCommandForDisplay, NOVU_CONNECT_AGENT_CHAT_TUI_COMMAND, SUBSCRIBER_ID_PLACEHOLDER, } from '@novu/shared'; diff --git a/apps/dashboard/src/components/agents/agent-chat-setup-guide.tsx b/apps/dashboard/src/components/agents/agent-chat-setup-guide.tsx index 40162c9812a..82b0075b47d 100644 --- a/apps/dashboard/src/components/agents/agent-chat-setup-guide.tsx +++ b/apps/dashboard/src/components/agents/agent-chat-setup-guide.tsx @@ -43,6 +43,7 @@ export function AgentChatSetupGuide({ const stepsColumn = ( openPreview(agent.identifier)} diff --git a/apps/dashboard/src/components/agents/agent-chat-setup-steps.tsx b/apps/dashboard/src/components/agents/agent-chat-setup-steps.tsx index 7c6e2f9904f..fcf8e5e00d6 100644 --- a/apps/dashboard/src/components/agents/agent-chat-setup-steps.tsx +++ b/apps/dashboard/src/components/agents/agent-chat-setup-steps.tsx @@ -3,6 +3,7 @@ import { AGENT_CHAT_DOCS_URL, AgentChatEmbedResources, buildAgentChatTuiCommand, + buildAgentChatTuiCommandForDisplay, } from '@/components/agents/agent-chat-setup-content'; import { CopyableTerminalBlock } from '@/components/primitives/copyable-terminal-block'; import { ExternalLink } from '@/components/shared/external-link'; @@ -12,6 +13,7 @@ import { deriveStepStatus } from './setup-guide-step-utils'; type AgentChatSetupStepsProps = { prompt: string; + agentIdentifier: string; stepOffset?: number; /** Omit to show all steps as completed (connected recap). */ firstIncompleteStep?: number; @@ -20,12 +22,19 @@ type AgentChatSetupStepsProps = { export function AgentChatSetupSteps({ prompt, + agentIdentifier, stepOffset = 1, firstIncompleteStep, onOpenChat, }: AgentChatSetupStepsProps) { const base = stepOffset; - const tuiCommand = buildAgentChatTuiCommand(apiHostnameManager.getHostname()); + const tuiCommandOptions = { + apiUrl: apiHostnameManager.getHostname(), + agentIdentifier, + connectDashboardUrl: window.location.origin, + }; + const tuiDisplayCommand = buildAgentChatTuiCommandForDisplay(tuiCommandOptions); + const tuiCopyCommand = buildAgentChatTuiCommand(tuiCommandOptions); return ( <> @@ -66,7 +75,7 @@ export function AgentChatSetupSteps({ } - extraContent={} + extraContent={} rightContent={} /> diff --git a/apps/dashboard/src/components/agents/agent-code-setup-section.tsx b/apps/dashboard/src/components/agents/agent-code-setup-section.tsx index dcc59536a69..ac3255a16b9 100644 --- a/apps/dashboard/src/components/agents/agent-code-setup-section.tsx +++ b/apps/dashboard/src/components/agents/agent-code-setup-section.tsx @@ -44,16 +44,21 @@ function resolveConnectRuntime(connectorId: ConnectorId | undefined): ConnectRun function buildConnectScaffoldParts({ secretKey, apiUrl, + connectDashboardUrl, runtime, }: { secretKey: string; apiUrl: string; + connectDashboardUrl: string; runtime: ConnectRuntimeFlag; }): string[] { return [ `${getNovuConnectInvocation(apiUrl)} --runtime ${runtime}`, `--secret-key ${secretKey}`, - ...getNovuConnectTargetFlags(apiUrl), + ...getNovuConnectTargetFlags({ + apiUrl, + connectDashboardUrl, + }), '--channel skip', ]; } @@ -61,29 +66,33 @@ function buildConnectScaffoldParts({ function buildConnectScaffoldCommand({ secretKey, apiUrl, + connectDashboardUrl, runtime, masked, }: { secretKey: string; apiUrl: string; + connectDashboardUrl: string; runtime: ConnectRuntimeFlag; masked: boolean; }): string { const key = masked ? maskSecretKey(secretKey) : secretKey; - return buildConnectScaffoldParts({ secretKey: key, apiUrl, runtime }).join(' \\\n '); + return buildConnectScaffoldParts({ secretKey: key, apiUrl, connectDashboardUrl, runtime }).join(' \\\n '); } function buildConnectScaffoldCopyCommand({ secretKey, apiUrl, + connectDashboardUrl, runtime, }: { secretKey: string; apiUrl: string; + connectDashboardUrl: string; runtime: ConnectRuntimeFlag; }): string { - return buildConnectScaffoldParts({ secretKey, apiUrl, runtime }).join(' '); + return buildConnectScaffoldParts({ secretKey, apiUrl, connectDashboardUrl, runtime }).join(' '); } function getProviderSlackMessage(agentName: string): string { @@ -361,6 +370,7 @@ export function AgentCodeSetupSection({ const connectRuntime = resolveConnectRuntime(connectorId); const currentApiUrl = apiHostnameManager.getHostname(); + const connectDashboardUrl = window.location.origin; const bridgeConnected = useBridgeConnectionPolling(agent, onBridgeConnected); @@ -402,12 +412,14 @@ export function AgentCodeSetupSection({ displayCommand={buildConnectScaffoldCommand({ secretKey, apiUrl: currentApiUrl, + connectDashboardUrl, runtime: connectRuntime, masked: true, })} copyCommand={buildConnectScaffoldCopyCommand({ secretKey, apiUrl: currentApiUrl, + connectDashboardUrl, runtime: connectRuntime, })} /> diff --git a/apps/dashboard/src/components/agents/agent-integration-guides/agent-chat-agent-integration-guide.tsx b/apps/dashboard/src/components/agents/agent-integration-guides/agent-chat-agent-integration-guide.tsx index 392cf11eb0f..ed937f79ffb 100644 --- a/apps/dashboard/src/components/agents/agent-integration-guides/agent-chat-agent-integration-guide.tsx +++ b/apps/dashboard/src/components/agents/agent-integration-guides/agent-chat-agent-integration-guide.tsx @@ -168,7 +168,7 @@ function AgentChatConnectedRecap({ - + diff --git a/packages/novu/src/commands/connect/pipeline/resolve-existing-agent.ts b/packages/novu/src/commands/connect/pipeline/resolve-existing-agent.ts new file mode 100644 index 00000000000..d7346ac6ef5 --- /dev/null +++ b/packages/novu/src/commands/connect/pipeline/resolve-existing-agent.ts @@ -0,0 +1,40 @@ +import type { AgentRecord } from '../api/agents'; +import type { AgentConnectMode, AgentSummary, ConnectCommandOptions } from '../types'; + +export function shouldSkipAgentConnectModePicker(options: ConnectCommandOptions): boolean { + return Boolean(options.agentIdentifier?.trim()); +} + +export function resolveConnectModeForExistingAgent(agent: Pick): AgentConnectMode { + return agent.runtime === 'self-hosted' ? 'custom-code' : 'demo'; +} + +export type ExistingAgentContext = { + summary: AgentSummary; + connectMode: AgentConnectMode; +}; + +export function resolveExistingAgentContext( + existingAgents: AgentRecord[], + agentIdentifier: string +): ExistingAgentContext { + const normalized = agentIdentifier.trim(); + const record = existingAgents.find((agent) => agent.identifier === normalized); + + if (!record) { + throw new Error(`No agent found with identifier "${normalized}" in this environment.`); + } + + return { + summary: { + id: record._id, + identifier: record.identifier, + name: record.name, + }, + connectMode: resolveConnectModeForExistingAgent(record), + }; +} + +export function resolveExistingAgentByIdentifier(existingAgents: AgentRecord[], agentIdentifier: string): AgentSummary { + return resolveExistingAgentContext(existingAgents, agentIdentifier).summary; +} diff --git a/packages/novu/src/commands/connect/pipeline/runner.ts b/packages/novu/src/commands/connect/pipeline/runner.ts index 36b4234d713..112d8a0fb0f 100644 --- a/packages/novu/src/commands/connect/pipeline/runner.ts +++ b/packages/novu/src/commands/connect/pipeline/runner.ts @@ -54,6 +54,11 @@ import { maybeRunChatSdkTunnel, runChatSdkProjectSetup } from './chat-sdk'; import { runCustomCodeProjectSetup } from './custom-code'; import { maybeRunLangChainTunnel, runLangChainProjectSetup } from './langchain'; import { resolveAgentRuntimeIntegration, resolveRuntimeFromOptions } from './resolve-agent-runtime-integration'; +import { + type ExistingAgentContext, + resolveExistingAgentContext, + shouldSkipAgentConnectModePicker, +} from './resolve-existing-agent'; export interface ConnectPipelineInput { options: ConnectCommandOptions; @@ -151,7 +156,11 @@ export async function runConnectPipeline(input: ConnectPipelineInput): Promise 0 && !options.prompt) { const pick = await ui.pickExistingOrCreate(existingAgents.map(toSummary)); if (pick.action === 'use') { @@ -550,7 +566,10 @@ function resolveBridgeProject(outcomes: { return outcomes.chatSdkOutcome ?? outcomes.aiSdkOutcome ?? outcomes.langChainOutcome ?? outcomes.customCodeOutcome; } -async function resolveAgentConnectMode(ctx: PipelineContext): Promise { +async function resolveAgentConnectMode( + ctx: PipelineContext, + preselectedAgent?: ExistingAgentContext +): Promise { const { options, ui, track, sessionProps } = ctx; if (options.runtime) { @@ -562,6 +581,17 @@ async function resolveAgentConnectMode(ctx: PipelineContext): Promise', 'Use an existing agent-runtime integration (skips credential setup for BYOK runtimes)' ) + .option('--agent-identifier ', 'Use an existing agent by identifier (skips the agent picker)') .option('--anthropic-api-key ', 'Anthropic API key for --runtime claude non-interactive runs') .option( '--llm-auth ', @@ -276,9 +277,9 @@ program const channel = options.skipSlack ? 'skip' : options.channel; const connectMode = options.chatSdk ? 'chat-sdk' : options.brain === 'chat-sdk' ? 'chat-sdk' : options.runtime; - if (!prompt && (!connectMode || !isBridgeConnectMode(connectMode))) { + if (!prompt && !options.agentIdentifier?.trim() && (!connectMode || !isBridgeConnectMode(connectMode))) { console.error( - 'Non-interactive mode requires a prompt (positional or --prompt), unless --runtime is a bridge mode (ai-sdk, langchain, custom-code, chat-sdk).\n(run `novu connect --help` for the non-interactive contract and examples)' + 'Non-interactive mode requires a prompt (positional or --prompt), --agent-identifier, or --runtime as a bridge mode (ai-sdk, langchain, custom-code, chat-sdk).\n(run `novu connect --help` for the non-interactive contract and examples)' ); process.exit(1); } diff --git a/packages/shared/src/utils/agent-chat-connect-prompt.spec.ts b/packages/shared/src/utils/agent-chat-connect-prompt.spec.ts index 669300c4dfb..06be5fc27da 100644 --- a/packages/shared/src/utils/agent-chat-connect-prompt.spec.ts +++ b/packages/shared/src/utils/agent-chat-connect-prompt.spec.ts @@ -2,6 +2,7 @@ import { describe, expect, it } from 'vitest'; import { buildAgentChatPrompt, buildAgentChatTuiCommand, + buildAgentChatTuiCommandForDisplay, buildOnboardingAgentPrompt, } from './agent-chat-connect-prompt'; import { NOVU_STAGING_API_URL } from './novu-connect-cli'; @@ -22,7 +23,31 @@ describe('agent-chat-connect-prompt', () => { 'npx novu@latest connect --channel agent-chat --api-url https://eu.api.novu.co' ); expect(buildAgentChatTuiCommand('http://localhost:3000')).toBe( - 'npx novu@latest connect --channel agent-chat --api-url http://localhost:3000' + 'npx novu@rc connect --channel agent-chat --api-url http://localhost:3000' + ); + }); + + it('includes agent identifier and local dashboard URLs when provided', () => { + expect( + buildAgentChatTuiCommand({ + apiUrl: 'http://localhost:3000', + agentIdentifier: 'support-agent', + connectDashboardUrl: 'http://localhost:4201', + }) + ).toBe( + 'npx novu@rc connect --channel agent-chat --api-url http://localhost:3000 --connect-dashboard-url http://localhost:4201 --dashboard-url http://localhost:4201 --agent-identifier support-agent' + ); + }); + + it('formats the TUI command for terminal display with line continuations', () => { + expect( + buildAgentChatTuiCommandForDisplay({ + apiUrl: 'http://localhost:3000', + agentIdentifier: 'support-agent', + connectDashboardUrl: 'http://localhost:4201', + }) + ).toBe( + 'npx novu@rc connect --channel agent-chat \\\n --api-url http://localhost:3000 \\\n --connect-dashboard-url http://localhost:4201 \\\n --dashboard-url http://localhost:4201 \\\n --agent-identifier support-agent' ); }); diff --git a/packages/shared/src/utils/agent-chat-connect-prompt.ts b/packages/shared/src/utils/agent-chat-connect-prompt.ts index 3a902d53dde..5ed564ea038 100644 --- a/packages/shared/src/utils/agent-chat-connect-prompt.ts +++ b/packages/shared/src/utils/agent-chat-connect-prompt.ts @@ -1,17 +1,45 @@ import { buildNovuConnectStagingHint, + formatNovuConnectCommandForDisplay, getNovuConnectInvocation, getNovuConnectTargetFlags, isNovuStagingApiUrl, + type NovuConnectTargetOptions, + normalizeConnectTargetOptions, } from './novu-connect-cli'; export const AGENT_CHAT_DOCS_URL = 'https://docs.novu.co/agents/channels/agent-chat'; export const AGENT_ONBOARDING_PLAYBOOK_URL = 'https://novu.co/agents.md'; -export function buildAgentChatTuiCommand(apiUrl?: string | null): string { - const parts = [`${getNovuConnectInvocation(apiUrl)} --channel agent-chat`, ...getNovuConnectTargetFlags(apiUrl)]; +export type BuildAgentChatTuiCommandOptions = NovuConnectTargetOptions & { + agentIdentifier?: string | null; +}; - return parts.join(' '); +export function buildAgentChatTuiCommandParts( + apiUrlOrOptions?: string | null | BuildAgentChatTuiCommandOptions +): string[] { + const options = normalizeConnectTargetOptions(apiUrlOrOptions); + const parts = [ + `${getNovuConnectInvocation(options.apiUrl)} --channel agent-chat`, + ...getNovuConnectTargetFlags(options), + ]; + const agentIdentifier = options.agentIdentifier?.trim(); + + if (agentIdentifier) { + parts.push(`--agent-identifier ${agentIdentifier}`); + } + + return parts; +} + +export function buildAgentChatTuiCommand(apiUrlOrOptions?: string | null | BuildAgentChatTuiCommandOptions): string { + return buildAgentChatTuiCommandParts(apiUrlOrOptions).join(' '); +} + +export function buildAgentChatTuiCommandForDisplay( + apiUrlOrOptions?: string | null | BuildAgentChatTuiCommandOptions +): string { + return formatNovuConnectCommandForDisplay(buildAgentChatTuiCommandParts(apiUrlOrOptions)); } /** Production TUI command. Prefer `buildAgentChatTuiCommand(apiUrl)` in the dashboard. */ @@ -33,7 +61,7 @@ function signedInDashboardLine(apiUrl?: string | null): string { /** * Dashboard Copy prompt / Open in Cursor text. Same shape as the onboarding * prompt: signed-in line + intent + agents.md. The playbook owns `--ci`, - * questions, and flags. Staging adds `novu@rc` + `--region staging`. + * questions, and flags. Staging and local dev use `novu@rc` (staging also passes `--region staging`). */ export function buildAgentChatPrompt(agentName: string, agentIdentifier: string, apiUrl?: string | null): string { const lines = [ diff --git a/packages/shared/src/utils/novu-connect-cli.spec.ts b/packages/shared/src/utils/novu-connect-cli.spec.ts index b95c746f425..81d7edb726b 100644 --- a/packages/shared/src/utils/novu-connect-cli.spec.ts +++ b/packages/shared/src/utils/novu-connect-cli.spec.ts @@ -28,4 +28,32 @@ describe('novu-connect-cli', () => { expect(getNovuConnectTargetFlags('https://eu.api.novu.co')).toEqual(['--api-url https://eu.api.novu.co']); expect(getNovuConnectTargetFlags('http://localhost:3000')).toEqual(['--api-url http://localhost:3000']); }); + + it('uses rc on local API hosts', () => { + expect(getNovuConnectPackageTag('http://localhost:3000')).toBe('rc'); + expect(getNovuConnectPackageTag('http://127.0.0.1:3000')).toBe('rc'); + expect(getNovuConnectInvocation('http://localhost:3000')).toBe('npx novu@rc connect'); + }); + + it('adds dashboard URLs for local connect commands', () => { + expect( + getNovuConnectTargetFlags({ + apiUrl: 'http://localhost:3000', + connectDashboardUrl: 'http://localhost:4201', + }) + ).toEqual([ + '--api-url http://localhost:3000', + '--connect-dashboard-url http://localhost:4201', + '--dashboard-url http://localhost:4201', + ]); + }); + + it('does not add dashboard URLs for cloud dashboards', () => { + expect( + getNovuConnectTargetFlags({ + apiUrl: 'http://localhost:3000', + connectDashboardUrl: 'https://dashboard.novu.co', + }) + ).toEqual(['--api-url http://localhost:3000']); + }); }); diff --git a/packages/shared/src/utils/novu-connect-cli.ts b/packages/shared/src/utils/novu-connect-cli.ts index 19f0f5e53e7..0455064d262 100644 --- a/packages/shared/src/utils/novu-connect-cli.ts +++ b/packages/shared/src/utils/novu-connect-cli.ts @@ -3,16 +3,58 @@ export const NOVU_STAGING_API_URL = 'https://api.novu-staging.co'; export type NovuConnectPackageTag = 'latest' | 'rc'; +export type NovuConnectTargetOptions = { + apiUrl?: string | null; + connectDashboardUrl?: string | null; + dashboardUrl?: string | null; +}; + +const NOVU_CLOUD_DASHBOARD_URLS = new Set([ + 'https://dashboard.novu.co', + 'https://eu.dashboard.novu.co', + 'https://dashboard.novu-staging.co', + 'https://dashboard.novu.localhost', +]); + +function normalizeUrl(url: string | null | undefined): string { + return (url ?? '').replace(/\/$/, ''); +} + function normalizeApiUrl(apiUrl: string | null | undefined): string { - return (apiUrl ?? '').replace(/\/$/, ''); + return normalizeUrl(apiUrl); +} + +export function normalizeConnectTargetOptions( + apiUrlOrOptions?: string | null | T +): T { + if (typeof apiUrlOrOptions === 'string' || apiUrlOrOptions == null) { + return { apiUrl: apiUrlOrOptions } as T; + } + + return apiUrlOrOptions; } export function isNovuStagingApiUrl(apiUrl: string | null | undefined): boolean { return normalizeApiUrl(apiUrl) === NOVU_STAGING_API_URL; } +export function isNovuLocalApiUrl(apiUrl: string | null | undefined): boolean { + const normalized = normalizeApiUrl(apiUrl); + if (!normalized) { + return false; + } + + try { + const hostname = new URL(normalized).hostname; + + return hostname === 'localhost' || hostname === '127.0.0.1'; + } catch { + return false; + } +} + export function getNovuConnectPackageTag(apiUrl?: string | null): NovuConnectPackageTag { - return isNovuStagingApiUrl(apiUrl) ? 'rc' : 'latest'; + return isNovuStagingApiUrl(apiUrl) || isNovuLocalApiUrl(apiUrl) ? 'rc' : 'latest'; } export function getNovuConnectRegionFlag(apiUrl?: string | null): '--region staging' | undefined { @@ -27,22 +69,48 @@ export function getNovuConnectInvocation(apiUrl?: string | null): string { return `npx novu@${getNovuConnectPackageTag(apiUrl)} connect`; } +function shouldEmitConnectDashboardFlags(connectDashboardUrl?: string | null): boolean { + const normalized = normalizeUrl(connectDashboardUrl); + + if (!normalized) { + return false; + } + + return !NOVU_CLOUD_DASHBOARD_URLS.has(normalized); +} + +export function formatNovuConnectCommandForDisplay(parts: readonly string[]): string { + return parts.join(' \\\n '); +} + /** * Staging uses `--region staging` so OAuth hits dashboard.novu-staging.co. - * Other non-US Cloud APIs keep `--api-url`. + * Other non-US Cloud APIs keep `--api-url`. Local dev also needs dashboard URLs + * so browser OAuth opens the same dashboard the user copied the command from. */ -export function getNovuConnectTargetFlags(apiUrl?: string | null): string[] { - const regionFlag = getNovuConnectRegionFlag(apiUrl); +export function getNovuConnectTargetFlags(apiUrlOrOptions?: string | null | NovuConnectTargetOptions): string[] { + const options = normalizeConnectTargetOptions(apiUrlOrOptions); + const regionFlag = getNovuConnectRegionFlag(options.apiUrl); if (regionFlag) { return [regionFlag]; } - const normalized = normalizeApiUrl(apiUrl); - if (normalized && normalized !== NOVU_CLOUD_API_URL) { - return [`--api-url ${normalized}`]; + const flags: string[] = []; + const normalizedApiUrl = normalizeApiUrl(options.apiUrl); + + if (normalizedApiUrl && normalizedApiUrl !== NOVU_CLOUD_API_URL) { + flags.push(`--api-url ${normalizedApiUrl}`); + } + + if (shouldEmitConnectDashboardFlags(options.connectDashboardUrl)) { + const connectDashboardUrl = normalizeUrl(options.connectDashboardUrl); + const dashboardUrl = normalizeUrl(options.dashboardUrl) || connectDashboardUrl; + + flags.push(`--connect-dashboard-url ${connectDashboardUrl}`); + flags.push(`--dashboard-url ${dashboardUrl}`); } - return []; + return flags; } export function buildNovuConnectStagingHint(apiUrl?: string | null): string | undefined { From 093c3a2b625b0378a9f7116752a1b61b1ef2e5a3 Mon Sep 17 00:00:00 2001 From: Adam Chmara Date: Fri, 21 Aug 2026 15:01:40 +0200 Subject: [PATCH 2/2] feat(js): add agent conversation runtime with immutable snapshots fixes NV-8640 (#12415) --- packages/js/scripts/size-limit.mjs | 6 +- packages/js/src/agent-chat/agent-chat.ts | 58 +++ .../agent-conversation-runtime.test.ts | 319 ++++++++++++++ .../agent-chat/agent-conversation-runtime.ts | 400 ++++++++++++++++++ .../agent-chat/conversation-runtime.types.ts | 80 ++++ packages/js/src/agent-chat/index.ts | 13 + .../js/src/agent-chat/runtime-cache-key.ts | 4 + packages/js/src/agent-chat/types.ts | 1 + packages/js/src/api/agent-chat-service.ts | 4 +- packages/js/src/index.ts | 12 +- 10 files changed, 892 insertions(+), 5 deletions(-) create mode 100644 packages/js/src/agent-chat/agent-conversation-runtime.test.ts create mode 100644 packages/js/src/agent-chat/agent-conversation-runtime.ts create mode 100644 packages/js/src/agent-chat/conversation-runtime.types.ts create mode 100644 packages/js/src/agent-chat/runtime-cache-key.ts diff --git a/packages/js/scripts/size-limit.mjs b/packages/js/scripts/size-limit.mjs index bc0a6acc00d..d934f59dce8 100644 --- a/packages/js/scripts/size-limit.mjs +++ b/packages/js/scripts/size-limit.mjs @@ -15,13 +15,13 @@ const modules = [ { name: 'UMD minified', filePath: umdPath, - // Raised for headless agentChat on Novu. Split to ./agent-chat later if needed. - limitInBytes: 225_000, + // Raised for agent conversation runtime (NV-8640). Split to ./agent-chat later if needed. + limitInBytes: 235_000, }, { name: 'UMD gzip', filePath: umdGzipPath, - limitInBytes: 62_000, + limitInBytes: 64_000, }, ]; diff --git a/packages/js/src/agent-chat/agent-chat.ts b/packages/js/src/agent-chat/agent-chat.ts index 698030f83b5..8cb7981767c 100644 --- a/packages/js/src/agent-chat/agent-chat.ts +++ b/packages/js/src/agent-chat/agent-chat.ts @@ -6,8 +6,12 @@ import type { Result } from '../types'; import { NovuError } from '../utils/errors'; import type { BaseSocketInterface } from '../ws/base-socket'; import { AgentChatStore, type ConversationEntry, createLocalConversationKey } from './agent-chat-store'; +import { AgentConversationRuntime } from './agent-conversation-runtime'; import { type AgentMessage, derivePendingActions } from './agent-message.types'; +import type { ConversationArgs, ConversationResult } from './conversation-runtime.types'; +import { runtimeCacheKey } from './runtime-cache-key'; import type { + AgentChatMessagesUpdated, FetchMoreArgs, FetchMoreResult, LoadConversationArgs, @@ -25,6 +29,7 @@ export class AgentChat extends BaseModule { #store: AgentChatStore; #socket: Pick; #liveSubscriberCount = 0; + #runtimes = new Map(); /** * Non-null while a reconnect catch-up is in flight: live envelopes are buffered here * and applied after the HTTP page is absorbed. Serialized via `#catchUpChain`. @@ -93,8 +98,60 @@ export class AgentChat extends BaseModule { } clearCache(): void { + for (const runtime of [...this.#runtimes.values()]) { + runtime.dispose(); + } + this.#store.clear(); this.#catchUpBuffer = null; + this.#runtimes.clear(); + } + + /** + * Return a stable conversation runtime for one agent thread. + * Resume sessions (`conversationId` set) are reused across calls with the same identity. + */ + conversation(args: ConversationArgs): ConversationResult { + const cacheKey = args.conversationId ? runtimeCacheKey(args.agentId, args.conversationId) : undefined; + if (cacheKey) { + const existing = this.#runtimes.get(cacheKey); + if (existing) { + return { ok: true, data: existing }; + } + } + + const runtime = new AgentConversationRuntime(this, args); + if (cacheKey) { + this.#runtimes.set(cacheKey, runtime); + } + + return { ok: true, data: runtime }; + } + + /** @internal */ + onMessagesUpdated(listener: (data: AgentChatMessagesUpdated) => void): () => void { + return this._emitter.on('agent_chat.messages.updated', ({ data }) => { + listener(data); + }); + } + + /** @internal */ + registerRuntime(cacheKey: string, runtime: AgentConversationRuntime): void { + const existing = this.#runtimes.get(cacheKey); + if (existing && existing !== runtime) { + existing.dispose(); + } + + this.#runtimes.set(cacheKey, runtime); + } + + /** @internal */ + unregisterRuntime(runtime: AgentConversationRuntime): void { + for (const [key, value] of this.#runtimes.entries()) { + if (value === runtime) { + this.#runtimes.delete(key); + } + } } getConversation({ agentId, conversationId, key }: { agentId: string; conversationId?: string; key?: string }): @@ -276,6 +333,7 @@ export class AgentChat extends BaseModule { text: args.text, conversationId, agentHash: args.agentHash, + metadata: args.metadata, }); this.#store.markSent(entry, { diff --git a/packages/js/src/agent-chat/agent-conversation-runtime.test.ts b/packages/js/src/agent-chat/agent-conversation-runtime.test.ts new file mode 100644 index 00000000000..34a5bb08edd --- /dev/null +++ b/packages/js/src/agent-chat/agent-conversation-runtime.test.ts @@ -0,0 +1,319 @@ +import { AgentChatService } from '../api'; +import { NovuEventEmitter } from '../event-emitter'; +import { AgentChat } from './agent-chat'; +import type { AgentConversationSnapshot } from './conversation-runtime.types'; + +describe('AgentConversationRuntime', () => { + const inboxServiceInstance = { isSessionInitialized: true } as any; + let emitter: NovuEventEmitter; + let sendMessage: jest.Mock; + let getEvents: jest.Mock; + let connect: jest.Mock; + let agentChat: AgentChat; + + beforeEach(() => { + emitter = new NovuEventEmitter(); + sendMessage = jest.fn(); + getEvents = jest.fn(); + connect = jest.fn().mockResolvedValue({ data: undefined }); + const agentChatService = { + sendMessage, + respondToAction: jest.fn(), + sendAction: jest.fn(), + getEvents, + } as unknown as AgentChatService; + agentChat = new AgentChat({ + inboxServiceInstance, + eventEmitterInstance: emitter, + agentChatService, + socket: { connect }, + }); + }); + + afterEach(() => { + for (const result of [ + agentChat.conversation({ agentId: 'agent_1', conversationId: 'conv_abcdefghijkl' }), + agentChat.conversation({ agentId: 'agent_1' }), + ]) { + if (result.ok) { + result.data.dispose(); + } + } + }); + + it('reuses a runtime keyed by agent and conversation id', () => { + const first = agentChat.conversation({ agentId: 'agent_1', conversationId: 'conv_abcdefghijkl' }); + const second = agentChat.conversation({ agentId: 'agent_1', conversationId: 'conv_abcdefghijkl' }); + + expect(first.ok).toBe(true); + expect(second.ok).toBe(true); + if (!first.ok || !second.ok) { + return; + } + + expect(first.data).toBe(second.data); + first.data.dispose(); + }); + + it('returns the same frozen snapshot reference until the next publication', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + expect(created.ok).toBe(true); + if (!created.ok) { + return; + } + + const runtime = created.data; + const before = runtime.getSnapshot(); + + await runtime.sendMessage('hello'); + await runtime.sendMessage({ text: 'with metadata', metadata: { source: 'test' } }); + + const afterSend = runtime.getSnapshot(); + expect(afterSend).not.toBe(before); + expect(afterSend.messages[0]).toMatchObject({ + role: 'user', + status: 'sent', + parts: [{ type: 'text', text: 'hello', state: 'done' }], + }); + expect(Object.isFrozen(afterSend)).toBe(true); + expect(Object.isFrozen(afterSend.messages)).toBe(true); + expect(Object.isFrozen(afterSend.pendingActions)).toBe(true); + expect(Object.isFrozen(afterSend.run)).toBe(true); + expect(Object.isFrozen(afterSend.pagination)).toBe(true); + expect(Object.isFrozen(afterSend.messages[0])).toBe(true); + expect(Object.isFrozen(afterSend.messages[0]?.parts[0])).toBe(true); + + const store = agentChat.getConversation({ agentId: 'agent_1', key: runtime.key }); + expect(afterSend.messages[0]).not.toBe(store?.messages[0]); + + const again = runtime.getSnapshot(); + expect(again).toBe(afterSend); + + expect(sendMessage).toHaveBeenNthCalledWith(1, expect.objectContaining({ text: 'hello' })); + expect(sendMessage).toHaveBeenNthCalledWith( + 2, + expect.objectContaining({ text: 'with metadata', metadata: { source: 'test' } }) + ); + + runtime.dispose(); + }); + + it('notifies subscribers only when the snapshot reference changes', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + if (!created.ok) { + return; + } + + const runtime = created.data; + const seen: AgentConversationSnapshot[] = []; + + const unsubscribe = runtime.subscribe((snapshot) => { + seen.push(snapshot); + }); + + expect(seen).toHaveLength(1); + const initial = seen[0]; + + await runtime.sendMessage('hello'); + + // Optimistic send publishes twice: sending, then sent. + expect(seen).toHaveLength(3); + expect(seen[0]).toBe(initial); + expect(seen[1]).not.toBe(initial); + expect(seen[2]).not.toBe(seen[1]); + + const current = runtime.getSnapshot(); + runtime.getSnapshot(); + expect(seen).toHaveLength(3); + expect(seen.every((snapshot, index) => snapshot === seen[index])).toBe(true); + expect(current).toBe(seen[2]); + + unsubscribe(); + runtime.dispose(); + }); + + it('separates run, conversationStatus, pagination, and session status in the snapshot', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_user0000001' }); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + if (!created.ok) { + return; + } + + await created.data.sendMessage('hello'); + + emitter.emit('agent_chat.agent_event', { + result: { + version: 1, + conversationId: 'internal', + conversationIdentifier: 'conv_abcdefghijkl', + agentId: 'agent_1', + runId: 'run_1', + turnId: 'turn_1', + sequence: 1, + timestamp: '2026-08-07T12:00:00.000Z', + event: { type: 'run-start' }, + }, + }); + + emitter.emit('agent_chat.agent_event', { + result: { + version: 1, + conversationId: 'internal', + conversationIdentifier: 'conv_abcdefghijkl', + agentId: 'agent_1', + runId: 'run_1', + turnId: 'turn_1', + sequence: 2, + timestamp: '2026-08-07T12:00:01.000Z', + event: { type: 'channel.typing', state: 'on', status: 'Thinking…' }, + }, + }); + + const snapshot = created.data.getSnapshot(); + expect(snapshot.status).toBe('ready'); + expect(snapshot.run.isRunning).toBe(true); + expect(snapshot.run.typing?.status).toBe('Thinking…'); + expect(snapshot.conversationStatus).toBe('active'); + expect(snapshot.pagination.hasMore).toBe(false); + + created.data.dispose(); + }); + + it('replaces a stale resume runtime when the create-flow runtime registers', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + + const createFlow = agentChat.conversation({ agentId: 'agent_1' }); + if (!createFlow.ok) { + return; + } + + const staleResume = agentChat.conversation({ + agentId: 'agent_1', + conversationId: 'conv_abcdefghijkl', + }); + expect(staleResume.ok).toBe(true); + if (!staleResume.ok) { + return; + } + + expect(staleResume.data).not.toBe(createFlow.data); + + await createFlow.data.sendMessage('hello'); + + const resumed = agentChat.conversation({ + agentId: 'agent_1', + conversationId: 'conv_abcdefghijkl', + }); + expect(resumed.ok).toBe(true); + if (!resumed.ok) { + return; + } + + expect(resumed.data).toBe(createFlow.data); + expect(resumed.data).not.toBe(staleResume.data); + expect(resumed.data.getSnapshot().messages[0]).toMatchObject({ + role: 'user', + parts: [{ type: 'text', text: 'hello' }], + }); + + createFlow.data.dispose(); + }); + + it('clearCache disposes runtimes so a later resume gets a fresh instance', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + if (!created.ok) { + return; + } + + await created.data.sendMessage('hello'); + const disposedRuntime = created.data; + + agentChat.clearCache(); + + const resumed = agentChat.conversation({ + agentId: 'agent_1', + conversationId: 'conv_abcdefghijkl', + }); + expect(resumed.ok).toBe(true); + if (!resumed.ok) { + return; + } + + expect(resumed.data).not.toBe(disposedRuntime); + + resumed.data.dispose(); + }); + + it('does not register a disposed runtime after an in-flight send completes', async () => { + let resolveSend!: (value: { identifier: string; messageId: string }) => void; + sendMessage.mockImplementation( + () => + new Promise((resolve) => { + resolveSend = resolve; + }) + ); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + if (!created.ok) { + return; + } + + const sendPromise = created.data.sendMessage('hello'); + await Promise.resolve(); + expect(sendMessage).toHaveBeenCalledTimes(1); + + created.data.dispose(); + + resolveSend({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + await sendPromise; + + const resumed = agentChat.conversation({ + agentId: 'agent_1', + conversationId: 'conv_abcdefghijkl', + }); + expect(resumed.ok).toBe(true); + if (!resumed.ok) { + return; + } + + expect(resumed.data).not.toBe(created.data); + + resumed.data.dispose(); + }); + + it('isolates snapshot messages from store mutations', async () => { + sendMessage.mockResolvedValue({ identifier: 'conv_abcdefghijkl', messageId: 'msg_abcdefghijkl' }); + + const created = agentChat.conversation({ agentId: 'agent_1' }); + if (!created.ok) { + return; + } + + await created.data.sendMessage('hello'); + + const snapshot = created.data.getSnapshot(); + const storeBefore = agentChat.getConversation({ agentId: 'agent_1', key: created.data.key }); + + expect(snapshot.messages[0]?.parts[0]).toMatchObject({ type: 'text', text: 'hello' }); + expect(Object.isFrozen(snapshot.messages[0])).toBe(true); + expect(Object.isFrozen(snapshot.messages[0]?.parts[0])).toBe(true); + if (snapshot.messages[0]?.parts[0]?.type !== 'text') { + return; + } + + expect(() => { + (snapshot.messages[0].parts[0] as { text: string }).text = 'mutated'; + }).toThrow(TypeError); + + expect(storeBefore?.messages[0]?.parts[0]).toMatchObject({ type: 'text', text: 'hello' }); + + created.data.dispose(); + }); +}); diff --git a/packages/js/src/agent-chat/agent-conversation-runtime.ts b/packages/js/src/agent-chat/agent-conversation-runtime.ts new file mode 100644 index 00000000000..c64284744ea --- /dev/null +++ b/packages/js/src/agent-chat/agent-conversation-runtime.ts @@ -0,0 +1,400 @@ +import type { AgentChatPlanLimitError } from '../api'; +import { NovuError } from '../utils/errors'; +import type { AgentChat } from './agent-chat'; +import { createLocalConversationKey } from './agent-chat-store'; +import type { AgentToolApprovalDecision } from './agent-message.types'; +import { derivePendingActions } from './agent-message.types'; +import type { + AgentConversationRunSnapshot, + AgentConversationSessionStatus, + AgentConversationSnapshot, + ConversationArgs, + ConversationResult, + SendMessageInput, +} from './conversation-runtime.types'; +import { runtimeCacheKey } from './runtime-cache-key'; + +const EMPTY_RUN: AgentConversationRunSnapshot = Object.freeze({ isRunning: false }); + +function cloneSnapshot(snapshot: AgentConversationSnapshot): AgentConversationSnapshot { + return { + ...snapshot, + run: { + ...snapshot.run, + typing: snapshot.run.typing ? { ...snapshot.run.typing } : undefined, + }, + pagination: { ...snapshot.pagination }, + messages: structuredClone(snapshot.messages), + pendingActions: structuredClone(snapshot.pendingActions), + }; +} + +function deepFreeze(value: T): T { + if (value === null || typeof value !== 'object') { + return value; + } + + Object.freeze(value); + + for (const nested of Object.values(value)) { + deepFreeze(nested); + } + + return value; +} + +function freezeSnapshot(snapshot: AgentConversationSnapshot): AgentConversationSnapshot { + return deepFreeze({ + ...snapshot, + run: { + ...snapshot.run, + typing: snapshot.run.typing ? { ...snapshot.run.typing } : undefined, + }, + pagination: { ...snapshot.pagination }, + messages: snapshot.messages, + pendingActions: snapshot.pendingActions, + }) as AgentConversationSnapshot; +} + +function createEmptySnapshot(key: string, conversationId?: string): AgentConversationSnapshot { + return freezeSnapshot({ + key, + conversationId, + status: conversationId ? 'loading' : 'ready', + run: EMPTY_RUN, + conversationStatus: 'active', + pagination: { hasMore: false }, + messages: [], + pendingActions: [], + }); +} + +function normalizeSendMessageInput(input: SendMessageInput): { text: string; metadata?: Record } { + if (typeof input === 'string') { + return { text: input }; + } + + return { text: input.text, metadata: input.metadata }; +} + +/** + * Framework-independent conversation runtime for one agent thread. + * Owns identity, immutable snapshots, and bound actions. + */ +export class AgentConversationRuntime { + readonly agentId: string; + readonly key: string; + + #agentChat: AgentChat; + #agentHash?: string; + #conversationId?: string; + #snapshot: AgentConversationSnapshot; + #listeners = new Set<(snapshot: AgentConversationSnapshot) => void>(); + #stopListening?: () => void; + #disposed = false; + #registeredConversationKey?: string; + + constructor(agentChat: AgentChat, args: ConversationArgs) { + this.#agentChat = agentChat; + this.agentId = args.agentId; + this.#agentHash = args.agentHash; + this.#conversationId = args.conversationId; + this.key = args.conversationId ?? createLocalConversationKey(); + this.#snapshot = createEmptySnapshot(this.key, args.conversationId); + + this.#agentChat.subscribe(); + this.#stopListening = this.#agentChat.onMessagesUpdated((data) => { + if (this.#disposed || data.key !== this.key) { + return; + } + + if (data.conversationId && !this.#conversationId) { + this.#conversationId = data.conversationId; + this.#registerByConversationId(data.conversationId); + } + + this.#publishFromStore({ + messages: data.messages, + isRunning: data.isRunning, + typing: data.typing, + status: data.status, + hasMore: data.hasMore, + conversationId: data.conversationId, + sessionStatus: + this.#snapshot.status === 'loading' || this.#snapshot.status === 'fetching' ? this.#snapshot.status : 'ready', + }); + }); + + if (args.conversationId) { + void this.load(); + } + } + + getSnapshot(): AgentConversationSnapshot { + return this.#snapshot; + } + + subscribe(listener: (snapshot: AgentConversationSnapshot) => void): () => void { + this.#listeners.add(listener); + listener(this.#snapshot); + + return () => { + this.#listeners.delete(listener); + }; + } + + dispose(): void { + if (this.#disposed) { + return; + } + + this.#disposed = true; + this.#stopListening?.(); + this.#stopListening = undefined; + this.#agentChat.unregisterRuntime(this); + this.#agentChat.unsubscribe(); + this.#listeners.clear(); + } + + async load(): Promise<{ + data?: { conversationId: string; messages: AgentConversationSnapshot['messages']; hasMore: boolean }; + error?: NovuError; + }> { + const conversationId = this.#conversationId; + if (!conversationId) { + return { + error: new NovuError( + 'Cannot load conversation without a conversation id', + new Error('missing conversation id') + ), + }; + } + + this.#publishSessionStatus('loading'); + + const response = await this.#agentChat.loadConversation({ + agentId: this.agentId, + conversationId, + }); + + if (response.error) { + this.#publishError(response.error); + this.#publishSessionStatus('ready'); + + return response; + } + + if (response.data) { + this.#publishFromStore({ + messages: response.data.messages, + isRunning: this.#snapshot.run.isRunning, + typing: this.#snapshot.run.typing, + status: this.#snapshot.conversationStatus, + hasMore: response.data.hasMore, + conversationId: response.data.conversationId, + sessionStatus: 'ready', + }); + } + + return response; + } + + async fetchMore(): Promise<{ + data?: { messages: AgentConversationSnapshot['messages']; hasMore: boolean }; + error?: NovuError; + }> { + this.#publishSessionStatus('fetching'); + + const response = await this.#agentChat.fetchMore({ + agentId: this.agentId, + key: this.key, + conversationId: this.#conversationId, + }); + + if (response.error) { + this.#publishError(response.error); + this.#publishSessionStatus('ready'); + + return response; + } + + if (response.data) { + const store = this.#agentChat.getConversation({ + agentId: this.agentId, + key: this.key, + conversationId: this.#conversationId, + }); + + this.#publishFromStore({ + messages: response.data.messages, + isRunning: store?.isRunning ?? this.#snapshot.run.isRunning, + typing: store?.typing ?? this.#snapshot.run.typing, + status: store?.status ?? this.#snapshot.conversationStatus, + hasMore: response.data.hasMore, + conversationId: store?.conversationId ?? this.#conversationId, + sessionStatus: 'ready', + }); + } + + return response; + } + + async sendMessage( + input: SendMessageInput + ): Promise<{ data?: { conversationId: string; messageId: string }; error?: NovuError | AgentChatPlanLimitError }> { + const { text, metadata } = normalizeSendMessageInput(input); + + const response = await this.#agentChat.sendMessage({ + agentId: this.agentId, + agentHash: this.#agentHash, + text, + metadata, + key: this.key, + conversationId: this.#conversationId, + }); + + if (response.error) { + this.#publishError(response.error); + + return response; + } + + if (response.data?.conversationId && !this.#disposed) { + this.#conversationId = response.data.conversationId; + this.#registerByConversationId(response.data.conversationId); + } + + return response; + } + + async respondToAction(args: { + actionId: string; + decision: AgentToolApprovalDecision; + }): Promise<{ data?: { conversationId: string }; error?: NovuError | AgentChatPlanLimitError }> { + const response = await this.#agentChat.respondToAction({ + agentId: this.agentId, + agentHash: this.#agentHash, + key: this.key, + conversationId: this.#conversationId, + actionId: args.actionId, + decision: args.decision, + }); + + if (response.error) { + this.#publishError(response.error); + } + + return response; + } + + async sendAction(args: { + actionId: string; + sourceMessageId: string; + value?: string; + }): Promise<{ data?: { conversationId: string }; error?: NovuError | AgentChatPlanLimitError }> { + const response = await this.#agentChat.sendAction({ + agentId: this.agentId, + agentHash: this.#agentHash, + key: this.key, + conversationId: this.#conversationId, + actionId: args.actionId, + sourceMessageId: args.sourceMessageId, + value: args.value, + }); + + if (response.error) { + this.#publishError(response.error); + } + + return response; + } + + cancelRun(): ConversationResult { + return { + ok: false, + error: new NovuError('Run cancellation is not supported yet', new Error('cancelRun is not implemented')), + }; + } + + /** @internal */ + get conversationId(): string | undefined { + return this.#conversationId; + } + + /** @internal */ + getRuntimeCacheKey(): string | undefined { + return this.#conversationId ? runtimeCacheKey(this.agentId, this.#conversationId) : undefined; + } + + #registerByConversationId(conversationId: string): void { + if (this.#disposed) { + return; + } + + const cacheKey = runtimeCacheKey(this.agentId, conversationId); + if (this.#registeredConversationKey === cacheKey) { + return; + } + + this.#registeredConversationKey = cacheKey; + this.#agentChat.registerRuntime(cacheKey, this); + } + + #publishSessionStatus(status: AgentConversationSessionStatus): void { + if (this.#snapshot.status === status) { + return; + } + + this.#publish({ + ...this.#snapshot, + status, + }); + } + + #publishError(error: NovuError | AgentChatPlanLimitError): void { + this.#publish({ + ...this.#snapshot, + error, + }); + } + + #publishFromStore(args: { + messages: AgentConversationSnapshot['messages']; + isRunning: boolean; + typing?: AgentConversationRunSnapshot['typing']; + status: AgentConversationSnapshot['conversationStatus']; + hasMore: boolean; + conversationId?: string; + sessionStatus: AgentConversationSessionStatus; + }): void { + this.#publish({ + key: this.key, + conversationId: args.conversationId ?? this.#conversationId, + status: args.sessionStatus, + run: { + isRunning: args.isRunning, + typing: args.typing, + }, + conversationStatus: args.status, + pagination: { + hasMore: args.hasMore, + }, + messages: args.messages, + pendingActions: derivePendingActions([...args.messages]), + error: undefined, + }); + } + + #publish(next: AgentConversationSnapshot): void { + if (this.#disposed) { + return; + } + + const frozen = freezeSnapshot(cloneSnapshot(next)); + this.#snapshot = frozen; + + for (const listener of this.#listeners) { + listener(frozen); + } + } +} diff --git a/packages/js/src/agent-chat/conversation-runtime.types.ts b/packages/js/src/agent-chat/conversation-runtime.types.ts new file mode 100644 index 00000000000..d0968abc69c --- /dev/null +++ b/packages/js/src/agent-chat/conversation-runtime.types.ts @@ -0,0 +1,80 @@ +import type { AgentChatPlanLimitError } from '../api'; +import type { NovuError } from '../utils/errors'; +import type { + AgentConversationError, + AgentConversationStatus, + AgentConversationTyping, + AgentMessage, + AgentPendingAction, + AgentToolApprovalDecision, +} from './agent-message.types'; +import type { + FetchMoreResult, + LoadConversationResult, + RespondToActionResult, + SendActionResult, + SendMessageResult, +} from './types'; + +/** Session lifecycle for history loads on this runtime. */ +export type AgentConversationSessionStatus = 'ready' | 'loading' | 'fetching'; + +export type AgentConversationRunSnapshot = { + isRunning: boolean; + typing?: AgentConversationTyping; +}; + +export type AgentConversationPaginationSnapshot = { + hasMore: boolean; +}; + +/** + * Immutable published view of one agent conversation thread. + * `getSnapshot()` returns the same object reference until the next publication. + */ +export type AgentConversationSnapshot = { + /** Holder key for this runtime session. Stable for the life of the runtime. */ + key: string; + conversationId?: string; + /** History load / pagination state for this runtime. */ + status: AgentConversationSessionStatus; + run: AgentConversationRunSnapshot; + conversationStatus: AgentConversationStatus; + pagination: AgentConversationPaginationSnapshot; + messages: readonly AgentMessage[]; + pendingActions: readonly AgentPendingAction[]; + error?: NovuError | AgentChatPlanLimitError | AgentConversationError; +}; + +export type ConversationOk = { ok: true; data: T }; +export type ConversationErr = { ok: false; error: NovuError }; +export type ConversationResult = ConversationOk | ConversationErr; + +export type ConversationArgs = { + agentId: string; + conversationId?: string; + agentHash?: string; +}; + +export type SendMessageInput = string | { text: string; metadata?: Record }; + +export type AgentConversationRuntimeActions = { + getSnapshot(): AgentConversationSnapshot; + subscribe(listener: (snapshot: AgentConversationSnapshot) => void): () => void; + dispose(): void; + load(): Promise<{ data?: LoadConversationResult; error?: NovuError }>; + fetchMore(): Promise<{ data?: FetchMoreResult; error?: NovuError }>; + sendMessage( + input: SendMessageInput + ): Promise<{ data?: SendMessageResult; error?: NovuError | AgentChatPlanLimitError }>; + respondToAction(args: { + actionId: string; + decision: AgentToolApprovalDecision; + }): Promise<{ data?: RespondToActionResult; error?: NovuError | AgentChatPlanLimitError }>; + sendAction(args: { + actionId: string; + sourceMessageId: string; + value?: string; + }): Promise<{ data?: SendActionResult; error?: NovuError | AgentChatPlanLimitError }>; + cancelRun(): ConversationResult; +}; diff --git a/packages/js/src/agent-chat/index.ts b/packages/js/src/agent-chat/index.ts index 22f5a7fdbd8..14fadece595 100644 --- a/packages/js/src/agent-chat/index.ts +++ b/packages/js/src/agent-chat/index.ts @@ -1,5 +1,18 @@ export { AgentChat } from './agent-chat'; +export { AgentConversationRuntime } from './agent-conversation-runtime'; export { derivePendingActions } from './agent-message.types'; +export type { + AgentConversationPaginationSnapshot, + AgentConversationRunSnapshot, + AgentConversationRuntimeActions, + AgentConversationSessionStatus, + AgentConversationSnapshot, + ConversationArgs, + ConversationErr, + ConversationOk, + ConversationResult, + SendMessageInput, +} from './conversation-runtime.types'; export type { AgentChatChange, AgentChatMessagesUpdated, diff --git a/packages/js/src/agent-chat/runtime-cache-key.ts b/packages/js/src/agent-chat/runtime-cache-key.ts new file mode 100644 index 00000000000..c98521ed832 --- /dev/null +++ b/packages/js/src/agent-chat/runtime-cache-key.ts @@ -0,0 +1,4 @@ +/** Stable registry key for one agent + conversation thread. */ +export function runtimeCacheKey(agentId: string, conversationId: string): string { + return `${agentId}::${conversationId}`; +} diff --git a/packages/js/src/agent-chat/types.ts b/packages/js/src/agent-chat/types.ts index 00898b10327..6aba5b0f70f 100644 --- a/packages/js/src/agent-chat/types.ts +++ b/packages/js/src/agent-chat/types.ts @@ -33,6 +33,7 @@ export type AgentHashFields = { export type SendMessageArgs = AgentHashFields & { agentId: string; text: string; + metadata?: Record; /** * Existing conversation to append to. * Omit this field to create a new conversation. The client does not reuse a prior chat. diff --git a/packages/js/src/api/agent-chat-service.ts b/packages/js/src/api/agent-chat-service.ts index e077916626b..3cfcffa7103 100644 --- a/packages/js/src/api/agent-chat-service.ts +++ b/packages/js/src/api/agent-chat-service.ts @@ -19,6 +19,7 @@ export class AgentChatPlanLimitError extends Error { export type AgentChatSendMessageArgs = AgentHashFields & { agentId: string; text: string; + metadata?: Record; /** Existing conversation id. Omit this field to create a new conversation. */ conversationId?: string; }; @@ -76,6 +77,7 @@ export class AgentChatService { text: args.text, ...(args.conversationId ? { conversationIdentifier: args.conversationId } : {}), ...(args.agentHash ? { agentHash: args.agentHash } : {}), + ...(args.metadata ? { metadata: args.metadata } : {}), }); } @@ -100,7 +102,7 @@ export class AgentChatService { } async #postAccept( - body: Record + body: Record ): Promise { try { return await this.#httpClient.post(AGENT_CHAT_CONVERSATIONS_ROUTE, body); diff --git a/packages/js/src/index.ts b/packages/js/src/index.ts index 5ef9d0c4f96..deb4b701459 100644 --- a/packages/js/src/index.ts +++ b/packages/js/src/index.ts @@ -1,6 +1,11 @@ export type * from 'json-logic-js'; export type { AgentChatChange, + AgentConversationPaginationSnapshot, + AgentConversationRunSnapshot, + AgentConversationRuntimeActions, + AgentConversationSessionStatus, + AgentConversationSnapshot, AgentConversationStatus, AgentConversationTyping, AgentEventEnvelope, @@ -11,6 +16,10 @@ export type { AgentPendingAction, AgentToolApprovalAction, AgentToolApprovalDecision, + ConversationArgs, + ConversationErr, + ConversationOk, + ConversationResult, FetchMoreArgs, FetchMoreResult, LoadConversationArgs, @@ -20,9 +29,10 @@ export type { SendActionArgs, SendActionResult, SendMessageArgs, + SendMessageInput, SendMessageResult, } from './agent-chat'; -export { derivePendingActions } from './agent-chat'; +export { AgentConversationRuntime, derivePendingActions } from './agent-chat'; export type { AgentChatPlanLimitReason } from './api/agent-chat-service'; export { AgentChatPlanLimitError } from './api/agent-chat-service'; export type {