feat(ai): add ask_questions interactive clarifying-question tool (#22346)
## What & why Adds an `ask_questions` tool that lets the in-app **Ask AI** assistant **pause a turn to ask the user one or more multiple-choice questions** (per the [Figma design](https://www.figma.com/design/xt8O9mFeLl46C5InWwoMrN/Twenty?node-id=105959-117153)) and resume once answered — instead of guessing on ambiguous/consequential decisions. The tool is **harness-only**: an interactive question UI is meaningless without a user to answer it, so it must be absent from MCP and from head-less workflow agents. ## Design — true tool-result resume (not a synthetic user message) The user's answer is a **structured tool result bound to the `toolCallId`**, and the **same agent turn resumes** — exactly how Anthropic (`tool_result` by `tool_use_id`) and OpenAI (`function_call_output`) model human-in-the-loop. The naive form of this (leave the tool call in `input-available` to mean "pending") is **impossible** here: `finalizeDanglingToolParts` rewrites `input-available` → `output-error` ("Tool execution was interrupted") on both the persist path (`addMessage`) and the model-reload path (`chat-execution.service.ts`). That util is a load-bearing safety net, so weakening it is the wrong move. Instead: - `ask_questions` is an **inline, chat-only tool with an `execute` that returns a `status: 'pending'` result immediately**, so the tool part is always `output-available` and **immune to `finalizeDanglingToolParts`**. `stopWhen(hasToolCall('ask_questions'))` halts the turn right after the call (the model never sees the placeholder). - A nullable **`thread.pendingQuestionMessageId`** marker records that a turn is awaiting an answer. - The new **`answerAgentChatQuestion`** mutation atomically *claims* the question (clears the marker, marks the thread streaming), **writes the answer onto the same tool part** (`status: 'answered'`), and **re-enqueues the turn via the existing `existingTurnId` plumbing** (`isResume` bypasses the per-turn dedup guard). On resume `finalizeDanglingToolParts` leaves the `output-available` part untouched and `convertToModelMessages` emits `assistant(tool_use)` + `tool_result(answers)`, so the model continues. This achieves the platform-aligned semantics **without** weakening the finalize safety net or inventing a fragile new part state. ### Meets the two requirements - **Survives refresh, scoped per-thread** — the pending state is a normal persisted `output-available` part + the thread marker; the frontend card is derived per-thread from the loaded messages, so it re-appears on reload and only on its own thread. - **Takes priority over the queue** — a unified `isBlocked = activeStreamId || pendingQuestionMessageId` gate is applied in both `sendChatMessage` (new messages queue) and `flushNextQueuedMessage` (the drain). The queue cannot unpile until the question is answered and the resumed turn completes. ### Harness-only by construction `ask_questions` is added **only** to the chat's inline `activeTools` (like `learn_tools`/`execute_tool`/`load_skills`). It never enters the tool registry/catalog, so it is invisible to MCP and to workflow agents — no `MCP_EXCLUDED_TOOL_NAMES` entry needed. ## UX While a question is pending, the **composer is replaced by the question card** (matching the Figma): question title + pager (`1/2`), numbered option rows (`IconSquareNumber*`) with per-option info-icon descriptions and a "Recommended" badge, and the normal composer as the free-text fallback ("Type anything to do differently."). The transcript shows a compact "Asking questions…" status line that becomes an answered summary. ## Changes **twenty-shared** - `ai/types/AskQuestionsToolTypes.ts` — `AskQuestionItem/Option/Answer/Result`, `ASK_QUESTIONS_TOOL_NAME`. **twenty-server** - `ai-chat/tools/ask-questions.tool.ts` — inline tool factory (pending-result `execute`, zod schema, 1–4 questions × 2–4 options). - `chat-execution.service.ts` — add to `activeTools` + `preloadedToolNames`; `hasToolCall` in `stopWhen`. - `chat-system-prompts.const.ts` — when-to-use guidance. - `entities/agent-chat-thread.entity.ts` — `pendingQuestionMessageId` column. - `stream-agent-chat.job.ts` — set the marker on a question pause; bypass the dedup guard on resume; suppress the no-text warning for question pauses. - `agent-chat-streaming.service.ts` — gate `flushNextQueuedMessage`; `enqueueResumeStream`. - `agent-chat.resolver.ts` — gate `sendChatMessage`; `answerAgentChatQuestion` mutation. - `agent-chat.service.ts` — `resolvePendingQuestion` (atomic claim + write answer). - `dtos/agent-chat-question-answer.input.ts`, `ai.exception.ts` (`QUESTION_NOT_PENDING`), `utils/find-pending-question-part.util.ts`. **twenty-front** - `components/AiChatQuestionCard.tsx` — the interactive card (matches Figma tokens) + `__stories__/AiChatQuestionCard.stories.tsx`. - `components/AiChatEditorSection.tsx` — swap the composer for the card while pending. - `components/AiChatQuestionStatusRenderer.tsx` + branch in `AiChatAssistantMessageRenderer.tsx`. - `states/selectors/agentChatPendingQuestionComponentSelector.ts`, `types/AgentChatPendingQuestion.ts`. - `hooks/useSubmitQuestionAnswer.ts` + `utils/markQuestionAnswered.ts` (optimistic) + `graphql/mutations/answerAgentChatQuestion.ts`. A design doc lives at `packages/twenty-server/docs/ASK_USER_QUESTION_TOOL_PLAN.md`. ## Migration Adds a nullable `pendingQuestionMessageId` (uuid) column to `core.agentChatThread`. Needs a generated **fast instance command** (`database:migrate:generate --name addThreadPendingQuestion --type fast`) — see "Verification status". ## Tests - Server: `ask-questions.tool.spec.ts` (pending echo + schema bounds), `find-pending-question-part.util.spec.ts`. - Front: `markQuestionAnswered.test.ts`, plus the Storybook story. ## Verification status (please read) This branch was authored in an environment where the monorepo `yarn install` repeatedly failed on transient TLS resets from the package registry, so I could **not** locally run the mechanical gates. The logic was reviewed by hand and the `ai@6.0.97` exports used (`hasToolCall`, `stepCountIs`, `generateId`) were confirmed against the package's type defs. Still **TODO** (will rely on CI / a follow-up once deps install): - [ ] `nx run twenty-shared:generateBarrels` (the `ai/index.ts` export was added by hand; regen to reconcile) - [ ] `nx run twenty-front:graphql:generate` (new mutation + input type) - [ ] generate the fast instance command (migration) for the new column - [ ] `typecheck` + `lint:diff-with-main` (front + server) — expect minor import-ordering autofixes - [ ] run the unit tests **Screenshots:** reproducing the live flow needs an AI provider API key (to get the model to actually call `ask_questions`), which isn't available here. The card can be screenshotted from its **Storybook story** (`AiChatQuestionCard.stories.tsx`) with no API key — I'll add that image once deps install, or a reviewer can run `nx storybook twenty-front`. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01AArS8H3y3Z1Qwm763xhPLB --- _Generated by [Claude Code](https://claude.ai/code/session_01AArS8H3y3Z1Qwm763xhPLB)_ <!-- This is an auto-generated description by cubic. --> <a href="https://cubic.dev/pr/twentyhq/twenty/pull/22346?utm_source=github" target="_blank" rel="noopener noreferrer" data-no-image-dialog="true"><picture><source media="(prefers-color-scheme: dark)" srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source media="(prefers-color-scheme: light)" srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img alt="Review in cubic" src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a> <!-- End of auto-generated description by cubic. -->
This commit is contained in:
+6
@@ -59,6 +59,12 @@ Building or editing dashboards through the AI is not available yet — it is a c
|
||||
|
||||
- **Favorites are navigation menu items.** Twenty has no separate "Favorites" concept. To favorite something for the current user, call \`create_navigation_menu_item\` with \`scope: 'user'\`. Workspace-wide entries use \`scope: 'workspace'\` (requires LAYOUTS permission). Both are the same primitive — do not look for a separate favorites tool.
|
||||
- **A default OBJECT navigation menu item is auto-created with \`create_object_metadata\`.** Don't immediately create another OBJECT item for the new object — only add a follow-up navigation item when the user is asking to pin a *different* view, folder, link, record, or page layout.
|
||||
|
||||
## Asking the user questions
|
||||
|
||||
- When a decision is genuinely ambiguous or consequential and you cannot infer it from the request or context, call \`ask_questions\` to ask the user one or more multiple-choice questions instead of guessing. The conversation pauses until they answer.
|
||||
- Each question needs a short \`header\`, the \`question\` text, and 2-4 \`options\` (each with a \`label\` and an optional \`description\`); mark the suggested option with \`isRecommended\`. The user can always type a free-form answer instead of picking an option.
|
||||
- Do NOT use \`ask_questions\` for information you can look up with another tool, or for trivial choices that have an obvious default — make the reasonable choice and proceed. Ask at most a few focused questions at once.
|
||||
`,
|
||||
|
||||
// Browsing context hint
|
||||
|
||||
+13
@@ -0,0 +1,13 @@
|
||||
import { Field, InputType, Int } from '@nestjs/graphql';
|
||||
|
||||
@InputType()
|
||||
export class AgentChatQuestionAnswerInput {
|
||||
@Field(() => Int)
|
||||
questionIndex: number;
|
||||
|
||||
@Field(() => [Int])
|
||||
selectedOptionIndices: number[];
|
||||
|
||||
@Field(() => String, { nullable: true })
|
||||
freeText?: string;
|
||||
}
|
||||
+8
@@ -11,6 +11,7 @@ import {
|
||||
} from 'typeorm';
|
||||
|
||||
import { ADD_LAST_STREAM_ERROR_TO_AGENT_CHAT_THREAD_UPGRADE_COMMAND_NAME } from 'src/database/commands/upgrade-version-command/2-19/add-last-stream-error-to-agent-chat-thread-upgrade-command-name.constant';
|
||||
import { ADD_PENDING_QUESTION_MESSAGE_ID_TO_AGENT_CHAT_THREAD_UPGRADE_COMMAND_NAME } from 'src/database/commands/upgrade-version-command/2-19/add-pending-question-message-id-to-agent-chat-thread-upgrade-command-name.constant';
|
||||
import { WasIntroducedInUpgrade } from 'src/engine/core-modules/upgrade/decorators/was-introduced-in-upgrade.decorator';
|
||||
import { UserWorkspaceEntity } from 'src/engine/core-modules/user-workspace/user-workspace.entity';
|
||||
import { AgentMessageEntity } from 'src/engine/metadata-modules/ai/ai-agent-execution/entities/agent-message.entity';
|
||||
@@ -73,6 +74,13 @@ export class AgentChatThreadEntity {
|
||||
@Column({ type: 'varchar', nullable: true })
|
||||
activeStreamId: string | null;
|
||||
|
||||
@WasIntroducedInUpgrade({
|
||||
upgradeCommandName:
|
||||
ADD_PENDING_QUESTION_MESSAGE_ID_TO_AGENT_CHAT_THREAD_UPGRADE_COMMAND_NAME,
|
||||
})
|
||||
@Column({ type: 'uuid', nullable: true })
|
||||
pendingQuestionMessageId: string | null;
|
||||
|
||||
@WasIntroducedInUpgrade({
|
||||
upgradeCommandName:
|
||||
ADD_LAST_STREAM_ERROR_TO_AGENT_CHAT_THREAD_UPGRADE_COMMAND_NAME,
|
||||
|
||||
+1
@@ -18,4 +18,5 @@ export type StreamAgentChatJobData = {
|
||||
hasTitle: boolean;
|
||||
existingTurnId?: string;
|
||||
conversationSizeTokens: number;
|
||||
isResume?: boolean;
|
||||
};
|
||||
|
||||
+26
-9
@@ -10,7 +10,7 @@ import type {
|
||||
} from 'twenty-shared/ai';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { Repository } from 'typeorm';
|
||||
import { v4 } from 'uuid';
|
||||
import { v5 as uuidv5 } from 'uuid';
|
||||
|
||||
import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator';
|
||||
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
|
||||
@@ -27,6 +27,7 @@ import { AgentChatEventPublisherService } from 'src/engine/metadata-modules/ai/a
|
||||
import { AgentChatStreamingService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service';
|
||||
import { AgentChatService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service';
|
||||
import { ChatExecutionService } from 'src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service';
|
||||
import { findPendingQuestionPart } from 'src/engine/metadata-modules/ai/ai-chat/utils/find-pending-question-part.util';
|
||||
import { getCancelChannel } from 'src/engine/metadata-modules/ai/ai-chat/utils/get-cancel-channel.util';
|
||||
import { mapErrorToStreamError } from 'src/engine/metadata-modules/ai/ai-chat/utils/map-error-to-stream-error.util';
|
||||
import type { AiModelConfig } from 'src/engine/metadata-modules/ai/ai-models/types/ai-model-config.type';
|
||||
@@ -38,6 +39,11 @@ import { type StreamAgentChatJobData } from './stream-agent-chat-job.types';
|
||||
|
||||
export { STREAM_AGENT_CHAT_JOB_NAME, type StreamAgentChatJobData };
|
||||
|
||||
// Derive assistantMessageId deterministically from streamId so assistant-message
|
||||
// persistence is idempotent per stream: a retried job for the stream is skipped,
|
||||
// while each distinct resume in a turn persists its own message.
|
||||
const ASSISTANT_MESSAGE_ID_NAMESPACE = '0b9c2a3d-4e5f-4a1b-8c2d-3e4f5a6b7c8d';
|
||||
|
||||
@Processor({ queueName: MessageQueue.aiStreamQueue, scope: Scope.REQUEST })
|
||||
export class StreamAgentChatJob {
|
||||
private readonly logger = new Logger(StreamAgentChatJob.name);
|
||||
@@ -204,7 +210,10 @@ export class StreamAgentChatJob {
|
||||
abortSignal: AbortSignal;
|
||||
}): Promise<void> {
|
||||
return new Promise<void>((resolve, reject) => {
|
||||
const assistantMessageId = v4();
|
||||
const assistantMessageId = uuidv5(
|
||||
data.streamId,
|
||||
ASSISTANT_MESSAGE_ID_NAMESPACE,
|
||||
);
|
||||
|
||||
let streamUsage = {
|
||||
inputTokens: 0,
|
||||
@@ -516,7 +525,9 @@ export class StreamAgentChatJob {
|
||||
(part) => part.type === 'text' && isNonEmptyString(part.text),
|
||||
);
|
||||
|
||||
if (isAborted || !hasText) {
|
||||
const pendingQuestionPart = findPendingQuestionPart(responseMessage.parts);
|
||||
|
||||
if ((isAborted || !hasText) && !isDefined(pendingQuestionPart)) {
|
||||
this.logAssistantTurnWithoutText({
|
||||
responseMessage,
|
||||
isAborted,
|
||||
@@ -544,13 +555,16 @@ export class StreamAgentChatJob {
|
||||
|
||||
const userMessage = await userMessagePromise;
|
||||
|
||||
if (
|
||||
isDefined(userMessage.turnId) &&
|
||||
(await this.agentChatService.hasAssistantMessageForTurn({
|
||||
turnId: userMessage.turnId,
|
||||
// Idempotent per stream: assistantMessageId is derived from the streamId,
|
||||
// so a retried job for this stream is skipped here while each distinct
|
||||
// resume in the turn persists its own message.
|
||||
const assistantMessageAlreadyPersisted =
|
||||
await this.agentChatService.hasMessageById({
|
||||
id: assistantMessageId,
|
||||
workspaceId,
|
||||
}))
|
||||
) {
|
||||
});
|
||||
|
||||
if (assistantMessageAlreadyPersisted) {
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -580,6 +594,9 @@ export class StreamAgentChatJob {
|
||||
`"totalCacheCreationTokens" + ${totalCacheCreationTokens}`,
|
||||
contextWindowTokens: modelConfig.contextWindowTokens,
|
||||
conversationSize: lastStepConversationSize,
|
||||
pendingQuestionMessageId: isDefined(pendingQuestionPart)
|
||||
? assistantMessageId
|
||||
: null,
|
||||
lastStreamError: null,
|
||||
},
|
||||
);
|
||||
|
||||
+78
-1
@@ -8,6 +8,7 @@ import {
|
||||
ResolveField,
|
||||
} from '@nestjs/graphql';
|
||||
|
||||
import { generateId } from 'ai';
|
||||
import GraphQLJSON from 'graphql-type-json';
|
||||
import { PermissionFlagType } from 'twenty-shared/constants';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
@@ -24,6 +25,7 @@ import { SettingsPermissionGuard } from 'src/engine/guards/settings-permission.g
|
||||
import { WorkspaceAuthGuard } from 'src/engine/guards/workspace-auth.guard';
|
||||
import { AgentMessageDTO } from 'src/engine/metadata-modules/ai/ai-agent-execution/dtos/agent-message.dto';
|
||||
import { type BrowsingContextType } from 'src/engine/metadata-modules/ai/ai-agent/types/browsingContext.type';
|
||||
import { AgentChatQuestionAnswerInput } from 'src/engine/metadata-modules/ai/ai-chat/dtos/agent-chat-question-answer.input';
|
||||
import { AgentChatThreadDTO } from 'src/engine/metadata-modules/ai/ai-chat/dtos/agent-chat-thread.dto';
|
||||
import { FileAttachmentInput } from 'src/engine/metadata-modules/ai/ai-chat/dtos/file-attachment.input';
|
||||
import { AiSystemPromptPreviewDTO } from 'src/engine/metadata-modules/ai/ai-chat/dtos/ai-system-prompt-preview.dto';
|
||||
@@ -186,7 +188,10 @@ export class AgentChatResolver {
|
||||
});
|
||||
}
|
||||
|
||||
if (isDefined(thread.activeStreamId)) {
|
||||
if (
|
||||
isDefined(thread.activeStreamId) ||
|
||||
isDefined(thread.pendingQuestionMessageId)
|
||||
) {
|
||||
const queuedMessage = await this.agentChatService.queueMessage({
|
||||
threadId,
|
||||
text,
|
||||
@@ -259,6 +264,78 @@ export class AgentChatResolver {
|
||||
};
|
||||
}
|
||||
|
||||
@Mutation(() => SendChatMessageResultDTO)
|
||||
async answerAgentChatQuestion(
|
||||
@Args('threadId', { type: () => UUIDScalarType }) threadId: string,
|
||||
@Args('messageId', { type: () => UUIDScalarType }) messageId: string,
|
||||
@Args('answers', { type: () => [AgentChatQuestionAnswerInput] })
|
||||
answers: AgentChatQuestionAnswerInput[],
|
||||
@Args('modelId', { type: () => String, nullable: true })
|
||||
modelId: string | undefined,
|
||||
@AuthUserWorkspaceId() userWorkspaceId: string,
|
||||
@AuthWorkspace() workspace: WorkspaceEntity,
|
||||
): Promise<SendChatMessageResultDTO> {
|
||||
if (this.aiModelRegistryService.getAvailableModels().length === 0) {
|
||||
throw new AiException(
|
||||
'No AI models are available. Configure at least one AI provider.',
|
||||
AiExceptionCode.API_KEY_NOT_CONFIGURED,
|
||||
);
|
||||
}
|
||||
|
||||
const resolvedModelId = modelId ?? workspace.smartModel;
|
||||
|
||||
this.aiModelRegistryService.validateModelAvailability(
|
||||
resolvedModelId,
|
||||
workspace,
|
||||
);
|
||||
|
||||
await this.billingUsageService.hasAvailableCreditsOrThrow(workspace.id);
|
||||
|
||||
const thread = await this.threadRepository.findOne(workspace.id, {
|
||||
where: { id: threadId, userWorkspaceId },
|
||||
});
|
||||
|
||||
if (!isDefined(thread)) {
|
||||
throw new AiException(
|
||||
'Thread not found',
|
||||
AiExceptionCode.THREAD_NOT_FOUND,
|
||||
);
|
||||
}
|
||||
|
||||
const streamId = generateId();
|
||||
|
||||
const { turnId } = await this.agentChatService.resolvePendingQuestion({
|
||||
threadId,
|
||||
messageId,
|
||||
answers,
|
||||
streamId,
|
||||
workspaceId: workspace.id,
|
||||
});
|
||||
|
||||
try {
|
||||
await this.agentChatStreamingService.enqueueResumeStream({
|
||||
threadId,
|
||||
userWorkspaceId,
|
||||
workspace,
|
||||
turnId,
|
||||
streamId,
|
||||
modelId,
|
||||
});
|
||||
} catch (error) {
|
||||
// Roll back the streaming claim so the thread isn't stuck "streaming".
|
||||
await this.threadRepository
|
||||
.update(
|
||||
workspace.id,
|
||||
{ id: threadId, activeStreamId: streamId },
|
||||
{ activeStreamId: null },
|
||||
)
|
||||
.catch(() => {});
|
||||
throw error;
|
||||
}
|
||||
|
||||
return { messageId, queued: false, streamId };
|
||||
}
|
||||
|
||||
@Mutation(() => Boolean)
|
||||
async stopAgentChatStream(
|
||||
@Args('threadId', { type: () => UUIDScalarType }) threadId: string,
|
||||
|
||||
+50
-1
@@ -242,6 +242,51 @@ export class AgentChatStreamingService {
|
||||
return { streamId, messageId: lastUserMessage.id };
|
||||
}
|
||||
|
||||
async enqueueResumeStream({
|
||||
threadId,
|
||||
userWorkspaceId,
|
||||
workspace,
|
||||
turnId,
|
||||
streamId,
|
||||
modelId,
|
||||
}: {
|
||||
threadId: string;
|
||||
userWorkspaceId: string;
|
||||
workspace: WorkspaceEntity;
|
||||
turnId: string | null;
|
||||
streamId: string;
|
||||
modelId?: string;
|
||||
}): Promise<void> {
|
||||
const thread = await this.threadRepository.findOneOrFail(workspace.id, {
|
||||
where: { id: threadId },
|
||||
});
|
||||
|
||||
const messages = await this.loadMessagesFromDB(
|
||||
threadId,
|
||||
userWorkspaceId,
|
||||
workspace.id,
|
||||
);
|
||||
|
||||
await this.messageQueueService.add<StreamAgentChatJobData>(
|
||||
STREAM_AGENT_CHAT_JOB_NAME,
|
||||
{
|
||||
threadId,
|
||||
streamId,
|
||||
userWorkspaceId,
|
||||
workspaceId: workspace.id,
|
||||
messages,
|
||||
browsingContext: null,
|
||||
modelId,
|
||||
lastUserMessageText: '',
|
||||
lastUserMessageParts: [],
|
||||
hasTitle: !!thread.title,
|
||||
conversationSizeTokens: thread.conversationSize,
|
||||
existingTurnId: turnId ?? undefined,
|
||||
isResume: true,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
async flushNextQueuedMessage(
|
||||
threadId: string,
|
||||
userWorkspaceId: string,
|
||||
@@ -250,13 +295,17 @@ export class AgentChatStreamingService {
|
||||
): Promise<void> {
|
||||
const threadStatus = await this.threadRepository.findOne(workspaceId, {
|
||||
where: { id: threadId },
|
||||
select: ['id', 'deletedAt'],
|
||||
select: ['id', 'deletedAt', 'pendingQuestionMessageId'],
|
||||
});
|
||||
|
||||
if (!threadStatus || threadStatus.deletedAt) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (isDefined(threadStatus.pendingQuestionMessageId)) {
|
||||
return;
|
||||
}
|
||||
|
||||
const queuedMessages = await this.agentChatService.getQueuedMessages({
|
||||
threadId,
|
||||
workspaceId,
|
||||
|
||||
+141
-5
@@ -1,6 +1,12 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
|
||||
import { ExtendedUIMessage } from 'twenty-shared/ai';
|
||||
import {
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
type AskQuestionAnswer,
|
||||
type AskQuestionItem,
|
||||
type AskQuestionsToolResult,
|
||||
ExtendedUIMessage,
|
||||
} from 'twenty-shared/ai';
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { In, IsNull, Not } from 'typeorm';
|
||||
import type { QueryDeepPartialEntity } from 'typeorm/query-builder/QueryPartialEntity';
|
||||
@@ -291,15 +297,15 @@ export class AgentChatService {
|
||||
});
|
||||
}
|
||||
|
||||
async hasAssistantMessageForTurn({
|
||||
turnId,
|
||||
async hasMessageById({
|
||||
id,
|
||||
workspaceId,
|
||||
}: {
|
||||
turnId: string;
|
||||
id: string;
|
||||
workspaceId: string;
|
||||
}): Promise<boolean> {
|
||||
const existingMessage = await this.messageRepository.findOne(workspaceId, {
|
||||
where: { turnId, role: AgentMessageRole.ASSISTANT },
|
||||
where: { id },
|
||||
select: ['id'],
|
||||
});
|
||||
|
||||
@@ -481,6 +487,136 @@ export class AgentChatService {
|
||||
return savedTurnId;
|
||||
}
|
||||
|
||||
async resolvePendingQuestion({
|
||||
threadId,
|
||||
messageId,
|
||||
answers,
|
||||
streamId,
|
||||
workspaceId,
|
||||
}: {
|
||||
threadId: string;
|
||||
messageId: string;
|
||||
answers: AskQuestionAnswer[];
|
||||
streamId: string;
|
||||
workspaceId: string;
|
||||
}): Promise<{ turnId: string | null }> {
|
||||
const message = await this.messageRepository.findOne(workspaceId, {
|
||||
where: { id: messageId, threadId },
|
||||
relations: ['parts'],
|
||||
});
|
||||
|
||||
if (!message) {
|
||||
throw new AiException(
|
||||
'Question message not found',
|
||||
AiExceptionCode.MESSAGE_NOT_FOUND,
|
||||
);
|
||||
}
|
||||
|
||||
const pendingPart = (message.parts ?? []).find(
|
||||
(part) =>
|
||||
part.toolName === ASK_QUESTIONS_TOOL_NAME &&
|
||||
(part.toolOutput as { result?: AskQuestionsToolResult } | null)?.result
|
||||
?.status === 'pending',
|
||||
);
|
||||
|
||||
if (!pendingPart) {
|
||||
throw new AiException(
|
||||
'No pending question to answer',
|
||||
AiExceptionCode.QUESTION_NOT_PENDING,
|
||||
);
|
||||
}
|
||||
|
||||
const previousOutput =
|
||||
(pendingPart.toolOutput as Record<string, unknown> | null) ?? {};
|
||||
const previousResult = previousOutput.result as
|
||||
| AskQuestionsToolResult
|
||||
| undefined;
|
||||
const questions = previousResult?.questions ?? [];
|
||||
|
||||
this.validateQuestionAnswers(answers, questions);
|
||||
|
||||
const claim = await this.threadRepository.update(
|
||||
workspaceId,
|
||||
{ id: threadId, pendingQuestionMessageId: messageId },
|
||||
{ pendingQuestionMessageId: null, activeStreamId: streamId },
|
||||
);
|
||||
|
||||
if ((claim.affected ?? 0) === 0) {
|
||||
throw new AiException(
|
||||
'No pending question to answer',
|
||||
AiExceptionCode.QUESTION_NOT_PENDING,
|
||||
);
|
||||
}
|
||||
|
||||
try {
|
||||
await this.messagePartRepository.update(
|
||||
workspaceId,
|
||||
{ id: pendingPart.id },
|
||||
{
|
||||
toolOutput: {
|
||||
...previousOutput,
|
||||
success: true,
|
||||
message: 'User answered the questions.',
|
||||
result: {
|
||||
questions,
|
||||
status: 'answered',
|
||||
answers,
|
||||
},
|
||||
},
|
||||
},
|
||||
);
|
||||
} catch (error) {
|
||||
await this.threadRepository
|
||||
.update(
|
||||
workspaceId,
|
||||
{ id: threadId, activeStreamId: streamId },
|
||||
{ pendingQuestionMessageId: messageId, activeStreamId: null },
|
||||
)
|
||||
.catch(() => {});
|
||||
throw error;
|
||||
}
|
||||
|
||||
return { turnId: message.turnId };
|
||||
}
|
||||
|
||||
private validateQuestionAnswers(
|
||||
answers: AskQuestionAnswer[],
|
||||
questions: AskQuestionItem[],
|
||||
): void {
|
||||
for (const answer of answers) {
|
||||
const question = questions[answer.questionIndex];
|
||||
|
||||
if (!isDefined(question)) {
|
||||
throw new AiException(
|
||||
'Answer references an unknown question.',
|
||||
AiExceptionCode.INVALID_QUESTION_ANSWER,
|
||||
);
|
||||
}
|
||||
|
||||
const hasInvalidOption = answer.selectedOptionIndices.some(
|
||||
(optionIndex) =>
|
||||
optionIndex < 0 || optionIndex >= question.options.length,
|
||||
);
|
||||
|
||||
if (hasInvalidOption) {
|
||||
throw new AiException(
|
||||
'Answer references an unknown option.',
|
||||
AiExceptionCode.INVALID_QUESTION_ANSWER,
|
||||
);
|
||||
}
|
||||
|
||||
if (
|
||||
question.allowMultiSelect !== true &&
|
||||
answer.selectedOptionIndices.length > 1
|
||||
) {
|
||||
throw new AiException(
|
||||
'This question allows only one selection.',
|
||||
AiExceptionCode.INVALID_QUESTION_ANSWER,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async updateThreadTitle({
|
||||
threadId,
|
||||
userWorkspaceId,
|
||||
|
||||
+10
-1
@@ -3,6 +3,7 @@ import { Injectable, Logger } from '@nestjs/common';
|
||||
import { isNonEmptyString, isObject } from '@sniptt/guards';
|
||||
import {
|
||||
convertToModelMessages,
|
||||
hasToolCall,
|
||||
type LanguageModelUsage,
|
||||
stepCountIs,
|
||||
type StepResult,
|
||||
@@ -53,6 +54,10 @@ import {
|
||||
extractCacheCreationTokensFromSteps,
|
||||
} from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util';
|
||||
import { AI_CHAT_TOOL_NAMES_TO_PRELOAD } from 'src/engine/metadata-modules/ai/ai-chat/constants/ai-chat-tool-names-to-preload.const';
|
||||
import {
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
createAskQuestionsTool,
|
||||
} from 'src/engine/metadata-modules/ai/ai-chat/tools/ask-questions.tool';
|
||||
import { MessagePruningService } from 'src/engine/metadata-modules/ai/ai-chat/services/message-pruning.service';
|
||||
import { SystemPromptBuilderService } from 'src/engine/metadata-modules/ai/ai-chat/services/system-prompt-builder.service';
|
||||
import { type ExtractedFile } from 'src/engine/metadata-modules/ai/ai-chat/types/extracted-file.type';
|
||||
@@ -195,12 +200,14 @@ export class ChatExecutionService {
|
||||
const preloadedToolNames = [
|
||||
...Object.keys(preloadedTools),
|
||||
...Object.keys(nativeTools),
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
];
|
||||
|
||||
// ToolSet is constant for the entire conversation — no mutation.
|
||||
// learn_tools returns schemas as text; execute_tool dispatches via the registry.
|
||||
const activeTools: ToolSet = {
|
||||
...directTools,
|
||||
[ASK_QUESTIONS_TOOL_NAME]: createAskQuestionsTool(),
|
||||
[LEARN_TOOLS_TOOL_NAME]: createLearnToolsTool(
|
||||
this.toolRegistry,
|
||||
toolContext,
|
||||
@@ -420,7 +427,9 @@ export class ChatExecutionService {
|
||||
tools: activeTools,
|
||||
abortSignal,
|
||||
stopWhen: (step) =>
|
||||
stepCountIs(AGENT_CONFIG.MAX_STEPS)(step) || hasNoMoreAvailableCredits,
|
||||
stepCountIs(AGENT_CONFIG.MAX_STEPS)(step) ||
|
||||
hasToolCall(ASK_QUESTIONS_TOOL_NAME)(step) ||
|
||||
hasNoMoreAvailableCredits,
|
||||
experimental_telemetry: AI_TELEMETRY_CONFIG,
|
||||
providerOptions: getCallLevelProviderOptions({
|
||||
sdkPackage: registeredModel.sdkPackage,
|
||||
|
||||
+119
@@ -0,0 +1,119 @@
|
||||
import {
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
askQuestionsInputSchema,
|
||||
createAskQuestionsTool,
|
||||
} from 'src/engine/metadata-modules/ai/ai-chat/tools/ask-questions.tool';
|
||||
|
||||
describe('ask_questions tool', () => {
|
||||
it('is named ask_questions (plural)', () => {
|
||||
expect(ASK_QUESTIONS_TOOL_NAME).toBe('ask_questions');
|
||||
});
|
||||
|
||||
it('execute echoes the questions with a pending status', async () => {
|
||||
const tool = createAskQuestionsTool();
|
||||
const questions = [
|
||||
{
|
||||
header: 'Email type',
|
||||
question: 'What type of email?',
|
||||
options: [{ label: 'Welcome' }, { label: 'Offer' }],
|
||||
},
|
||||
];
|
||||
|
||||
const output = await tool.execute({ questions });
|
||||
|
||||
expect(output).toEqual({
|
||||
success: true,
|
||||
message: expect.any(String),
|
||||
result: { questions, status: 'pending' },
|
||||
});
|
||||
});
|
||||
|
||||
it('rejects fewer than two options', () => {
|
||||
const result = askQuestionsInputSchema.safeParse({
|
||||
questions: [
|
||||
{ header: 'h', question: 'q', options: [{ label: 'only one' }] },
|
||||
],
|
||||
});
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
});
|
||||
|
||||
it('rejects zero questions', () => {
|
||||
const result = askQuestionsInputSchema.safeParse({ questions: [] });
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
});
|
||||
|
||||
it('rejects more than four questions', () => {
|
||||
const question = {
|
||||
header: 'h',
|
||||
question: 'q',
|
||||
options: [{ label: 'a' }, { label: 'b' }],
|
||||
};
|
||||
const result = askQuestionsInputSchema.safeParse({
|
||||
questions: [question, question, question, question, question],
|
||||
});
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
});
|
||||
|
||||
it('rejects more than four options', () => {
|
||||
const result = askQuestionsInputSchema.safeParse({
|
||||
questions: [
|
||||
{
|
||||
header: 'h',
|
||||
question: 'q',
|
||||
options: [
|
||||
{ label: 'a' },
|
||||
{ label: 'b' },
|
||||
{ label: 'c' },
|
||||
{ label: 'd' },
|
||||
{ label: 'e' },
|
||||
],
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
});
|
||||
|
||||
it('rejects more than one recommended option', () => {
|
||||
const result = askQuestionsInputSchema.safeParse({
|
||||
questions: [
|
||||
{
|
||||
header: 'h',
|
||||
question: 'q',
|
||||
options: [
|
||||
{ label: 'a', isRecommended: true },
|
||||
{ label: 'b', isRecommended: true },
|
||||
],
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
});
|
||||
|
||||
it('accepts a valid multi-question payload', () => {
|
||||
const result = askQuestionsInputSchema.safeParse({
|
||||
questions: [
|
||||
{
|
||||
header: 'h1',
|
||||
question: 'q1',
|
||||
options: [{ label: 'a' }, { label: 'b' }],
|
||||
},
|
||||
{
|
||||
header: 'h2',
|
||||
question: 'q2',
|
||||
options: [
|
||||
{ label: 'c', description: 'desc', isRecommended: true },
|
||||
{ label: 'd' },
|
||||
],
|
||||
allowMultiSelect: true,
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
expect(result.success).toBe(true);
|
||||
});
|
||||
});
|
||||
+85
@@ -0,0 +1,85 @@
|
||||
import { z } from 'zod';
|
||||
|
||||
import {
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
type AskQuestionsToolInput,
|
||||
type AskQuestionsToolResult,
|
||||
} from 'twenty-shared/ai';
|
||||
|
||||
export { ASK_QUESTIONS_TOOL_NAME };
|
||||
|
||||
export const askQuestionsInputSchema = z.object({
|
||||
questions: z
|
||||
.array(
|
||||
z.object({
|
||||
header: z
|
||||
.string()
|
||||
.describe(
|
||||
'Very short label/tag for the question (≤ ~32 chars), e.g. "Email type".',
|
||||
),
|
||||
question: z
|
||||
.string()
|
||||
.describe(
|
||||
'The full question to ask the user. Be clear and specific.',
|
||||
),
|
||||
options: z
|
||||
.array(
|
||||
z.object({
|
||||
label: z
|
||||
.string()
|
||||
.describe('Concise option the user can pick (1-5 words).'),
|
||||
description: z
|
||||
.string()
|
||||
.optional()
|
||||
.describe(
|
||||
'Longer explanation shown when the user opens the option info icon.',
|
||||
),
|
||||
isRecommended: z
|
||||
.boolean()
|
||||
.optional()
|
||||
.describe('Mark the single suggested option, if any.'),
|
||||
}),
|
||||
)
|
||||
.min(2)
|
||||
.max(4)
|
||||
.refine(
|
||||
(options) =>
|
||||
options.filter((option) => option.isRecommended === true)
|
||||
.length <= 1,
|
||||
{ message: 'At most one option can be marked as recommended.' },
|
||||
)
|
||||
.describe('2-4 mutually exclusive options.'),
|
||||
allowMultiSelect: z
|
||||
.boolean()
|
||||
.optional()
|
||||
.describe('Allow the user to select more than one option.'),
|
||||
}),
|
||||
)
|
||||
.min(1)
|
||||
.max(4)
|
||||
.describe('One to four questions to ask the user.'),
|
||||
});
|
||||
|
||||
type AskQuestionsPendingOutput = {
|
||||
success: true;
|
||||
message: string;
|
||||
result: AskQuestionsToolResult;
|
||||
};
|
||||
|
||||
export const createAskQuestionsTool = () => ({
|
||||
description:
|
||||
'Ask the user one or more multiple-choice questions when you need a decision you cannot ' +
|
||||
'infer from the request or context and that has no obvious default. The conversation ' +
|
||||
'pauses until the user answers, then continues with their choice in mind. Prefer this ' +
|
||||
'over guessing on consequential or ambiguous decisions. Do NOT use it for information you ' +
|
||||
'could look up with another tool, or for trivial choices with an obvious default. The ' +
|
||||
'user can always type a free-form answer instead of picking an option.',
|
||||
inputSchema: askQuestionsInputSchema,
|
||||
execute: async (
|
||||
input: AskQuestionsToolInput,
|
||||
): Promise<AskQuestionsPendingOutput> => ({
|
||||
success: true,
|
||||
message: 'Questions presented to the user; awaiting their answer.',
|
||||
result: { questions: input.questions, status: 'pending' },
|
||||
}),
|
||||
});
|
||||
+58
@@ -0,0 +1,58 @@
|
||||
import { type ExtendedUIMessagePart } from 'twenty-shared/ai';
|
||||
|
||||
import { findPendingQuestionPart } from 'src/engine/metadata-modules/ai/ai-chat/utils/find-pending-question-part.util';
|
||||
|
||||
const askQuestionsPart = (
|
||||
status: 'pending' | 'answered',
|
||||
): ExtendedUIMessagePart =>
|
||||
({
|
||||
type: 'tool-ask_questions',
|
||||
toolCallId: 'call-1',
|
||||
state: 'output-available',
|
||||
input: { questions: [] },
|
||||
output: {
|
||||
success: true,
|
||||
message: 'x',
|
||||
result: {
|
||||
questions: [{ header: 'h', question: 'q', options: [] }],
|
||||
status,
|
||||
},
|
||||
},
|
||||
}) as unknown as ExtendedUIMessagePart;
|
||||
|
||||
const textPart = (text: string): ExtendedUIMessagePart =>
|
||||
({ type: 'text', text }) as ExtendedUIMessagePart;
|
||||
|
||||
describe('findPendingQuestionPart', () => {
|
||||
it('returns the ask_questions part when status is pending', () => {
|
||||
const part = findPendingQuestionPart([
|
||||
textPart('hello'),
|
||||
askQuestionsPart('pending'),
|
||||
]);
|
||||
|
||||
expect(part).toBeDefined();
|
||||
expect(part?.toolCallId).toBe('call-1');
|
||||
});
|
||||
|
||||
it('returns undefined when the question has been answered', () => {
|
||||
expect(
|
||||
findPendingQuestionPart([askQuestionsPart('answered')]),
|
||||
).toBeUndefined();
|
||||
});
|
||||
|
||||
it('returns undefined when there is no ask_questions part', () => {
|
||||
expect(findPendingQuestionPart([textPart('hello')])).toBeUndefined();
|
||||
});
|
||||
|
||||
it('ignores other tool parts', () => {
|
||||
const otherTool = {
|
||||
type: 'tool-search_help_center',
|
||||
toolCallId: 'call-2',
|
||||
state: 'output-available',
|
||||
input: {},
|
||||
output: { success: true },
|
||||
} as unknown as ExtendedUIMessagePart;
|
||||
|
||||
expect(findPendingQuestionPart([otherTool])).toBeUndefined();
|
||||
});
|
||||
});
|
||||
+29
@@ -0,0 +1,29 @@
|
||||
import { getToolName, isToolUIPart } from 'ai';
|
||||
import {
|
||||
ASK_QUESTIONS_TOOL_NAME,
|
||||
type AskQuestionsToolResult,
|
||||
type ExtendedUIMessagePart,
|
||||
} from 'twenty-shared/ai';
|
||||
|
||||
type ToolPartWithOutput = ExtendedUIMessagePart & {
|
||||
toolCallId: string;
|
||||
output?: { result?: AskQuestionsToolResult };
|
||||
};
|
||||
|
||||
export const findPendingQuestionPart = (
|
||||
parts: ExtendedUIMessagePart[],
|
||||
): ToolPartWithOutput | undefined => {
|
||||
for (const part of parts) {
|
||||
if (!isToolUIPart(part) || getToolName(part) !== ASK_QUESTIONS_TOOL_NAME) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const output = (part as ToolPartWithOutput).output;
|
||||
|
||||
if (output?.result?.status === 'pending') {
|
||||
return part as ToolPartWithOutput;
|
||||
}
|
||||
}
|
||||
|
||||
return undefined;
|
||||
};
|
||||
@@ -13,6 +13,8 @@ export enum AiExceptionCode {
|
||||
THREAD_NOT_FOUND = 'THREAD_NOT_FOUND',
|
||||
INVALID_CHAT_THREAD_TITLE = 'INVALID_CHAT_THREAD_TITLE',
|
||||
MESSAGE_NOT_FOUND = 'MESSAGE_NOT_FOUND',
|
||||
QUESTION_NOT_PENDING = 'QUESTION_NOT_PENDING',
|
||||
INVALID_QUESTION_ANSWER = 'INVALID_QUESTION_ANSWER',
|
||||
API_KEY_NOT_CONFIGURED = 'API_KEY_NOT_CONFIGURED',
|
||||
USER_WORKSPACE_ID_NOT_FOUND = 'USER_WORKSPACE_ID_NOT_FOUND',
|
||||
ROLE_NOT_FOUND = 'ROLE_NOT_FOUND',
|
||||
@@ -38,6 +40,10 @@ const getAiExceptionUserFriendlyMessage = (code: AiExceptionCode) => {
|
||||
return msg`Chat thread title cannot be empty.`;
|
||||
case AiExceptionCode.MESSAGE_NOT_FOUND:
|
||||
return msg`Chat message not found.`;
|
||||
case AiExceptionCode.QUESTION_NOT_PENDING:
|
||||
return msg`This question has already been answered.`;
|
||||
case AiExceptionCode.INVALID_QUESTION_ANSWER:
|
||||
return msg`Invalid answer for this question.`;
|
||||
case AiExceptionCode.API_KEY_NOT_CONFIGURED:
|
||||
return msg`API key is not configured.`;
|
||||
case AiExceptionCode.USER_WORKSPACE_ID_NOT_FOUND:
|
||||
|
||||
+2
@@ -28,6 +28,8 @@ export const aiGraphqlApiExceptionHandler = (error: Error) => {
|
||||
throw new NotFoundError(error);
|
||||
case AiExceptionCode.INVALID_AGENT_INPUT:
|
||||
case AiExceptionCode.INVALID_CHAT_THREAD_TITLE:
|
||||
case AiExceptionCode.QUESTION_NOT_PENDING:
|
||||
case AiExceptionCode.INVALID_QUESTION_ANSWER:
|
||||
throw new UserInputError(error);
|
||||
case AiExceptionCode.AGENT_ALREADY_EXISTS:
|
||||
case AiExceptionCode.NO_FAILED_TURN_TO_RETRY:
|
||||
|
||||
Reference in New Issue
Block a user