From a51c37dae55f76c4287bcd0435441a46e87ad0d8 Mon Sep 17 00:00:00 2001
From: Etienne <45695613+etiennejouan@users.noreply.github.com>
Date: Thu, 9 Jul 2026 10:24:27 +0200
Subject: [PATCH] feat(ai) - add light AI chat turn instrumentation (metrics +
Sentry correlation) (#22692)
## Summary
Adds minimal server-side observability for AI chat turns: lifecycle
counters to measure success/failure rates, Sentry scope tags to
correlate API and worker traces, and per-LLM-call telemetry metadata for
turn/stream correlation.
- Add turn lifecycle metrics: `ai-chat/turn-started`,
`ai-chat/turn-completed`, `ai-chat/turn-failed` (with `failure_phase`
and `error_code` attributes)
- Emit counters at key points: job start, clean completion, execution
failures, enqueue failures, interrupted streams, and empty completions
- Tag Sentry scope with `streamId`, `turnId`, `threadId`, and
`workspaceId` at API entry points (`sendChatMessage`,
`retryChatMessage`, `answerAgentChatQuestion`) and worker entry
(`StreamAgentChatJob`)
- Enrich LLM `experimental_telemetry` metadata with
stream/turn/thread/workspace IDs
- Return `turnId` from streaming service methods so resolvers can tag
the scope
- Remove granular tool-learned/skill-loaded metrics in favor of the
turn-level counters
## Test plan
- [ ] Send a chat message and verify `ai-chat/turn-started` and
`ai-chat/turn-completed` increment
- [ ] Trigger a stream failure (e.g. interrupted/dead stream) and verify
`ai-chat/turn-failed` with correct `failure_phase`
- [ ] Retry a failed turn and confirm a new `turn-started` is emitted
for the retry attempt
- [ ] Answer an `ask_questions` prompt and confirm Sentry tags include
`streamId` and `turnId`
- [ ] Check Sentry spans for LLM calls include `streamId`, `turnId`,
`threadId`, `workspaceId` in telemetry metadata
- [ ] Run unit tests:
- `npx jest
packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts`
- `npx jest
packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts`
- `npx jest
packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.retry.spec.ts`
---
.../metrics/types/metrics-keys.type.ts | 7 +-
.../__tests__/stream-agent-chat.job.spec.ts | 1 +
.../ai/ai-chat/jobs/stream-agent-chat.job.ts | 48 ++++++++++++
.../ai-chat/resolvers/agent-chat.resolver.ts | 23 ++++++
...agent-chat-streaming.service.claim.spec.ts | 7 ++
...agent-chat-streaming.service.retry.spec.ts | 2 +
.../services/agent-chat-streaming.service.ts | 68 ++++++++++++++++-
.../services/chat-execution.service.ts | 76 +++++--------------
.../utils/tag-ai-chat-stream-scope.util.ts | 20 +++++
9 files changed, 189 insertions(+), 63 deletions(-)
create mode 100644 packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util.ts
diff --git a/packages/twenty-server/src/engine/core-modules/metrics/types/metrics-keys.type.ts b/packages/twenty-server/src/engine/core-modules/metrics/types/metrics-keys.type.ts
index ac7c4f150b..5760c83a56 100644
--- a/packages/twenty-server/src/engine/core-modules/metrics/types/metrics-keys.type.ts
+++ b/packages/twenty-server/src/engine/core-modules/metrics/types/metrics-keys.type.ts
@@ -25,10 +25,6 @@ export enum MetricsKeys {
WorkflowRunSystemError = 'workflow-run/system-error',
AiChatToolExecutionSucceeded = 'ai-chat/tool-execution-succeeded',
AiChatToolExecutionFailed = 'ai-chat/tool-execution-failed',
- AiChatToolLearnedSucceeded = 'ai-chat/tool-learned-succeeded',
- AiChatToolLearnedFailed = 'ai-chat/tool-learned-failed',
- AiChatSkillLoadedSucceeded = 'ai-chat/skill-loaded-succeeded',
- AiChatSkillLoadedFailed = 'ai-chat/skill-loaded-failed',
WorkflowAgentToolExecutionSucceeded = 'workflow-agent/tool-execution-succeeded',
WorkflowAgentToolExecutionFailed = 'workflow-agent/tool-execution-failed',
McpToolExecutionSucceeded = 'mcp/tool-execution-succeeded',
@@ -53,5 +49,8 @@ export enum MetricsKeys {
AiChatTurnLatencyMs = 'ai-chat/turn-latency-ms',
AiChatStepLatencyMs = 'ai-chat/step-latency-ms',
AiChatTtftMs = 'ai-chat/ttft-ms',
+ AiChatTurnStarted = 'ai-chat/turn-started',
+ AiChatTurnCompleted = 'ai-chat/turn-completed',
+ AiChatTurnFailed = 'ai-chat/turn-failed',
WorkspaceMetadataCacheLocalEviction = 'workspace-metadata-cache/local-eviction',
}
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts
index 6a8d5e2ce1..bcdac582f1 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts
@@ -173,6 +173,7 @@ describe('StreamAgentChatJob', () => {
cancelSubscriberService as never,
agentChatStreamingService as never,
streamHeartbeatService as never,
+ { incrementCounterBy: jest.fn() } as never,
);
return {
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts
index 54259f8c8a..b8215f9f77 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts
@@ -17,6 +17,8 @@ import { type MessageQueueJobContext } from 'src/engine/core-modules/message-que
import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator';
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
+import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
+import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type';
import { toDisplayCredits } from 'src/engine/core-modules/usage/utils/to-display-credits.util';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { AgentMessageRole } from 'src/engine/metadata-modules/ai/ai-agent-execution/entities/agent-message.entity';
@@ -38,6 +40,7 @@ import { findPendingQuestionPart } from 'src/engine/metadata-modules/ai/ai-chat/
import { AGENT_CHAT_CHECKPOINT_INTERVAL_MS } from 'src/engine/metadata-modules/ai/ai-chat/constants/agent-chat-checkpoint-interval-ms.constant';
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 { tagAiChatStreamScope } from 'src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util';
import type { AiModelConfig } from 'src/engine/metadata-modules/ai/ai-models/types/ai-model-config.type';
import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator';
import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
@@ -67,6 +70,7 @@ export class StreamAgentChatJob {
private readonly cancelSubscriberService: AgentChatCancelSubscriberService,
private readonly agentChatStreamingService: AgentChatStreamingService,
private readonly streamHeartbeatService: AgentChatStreamHeartbeatService,
+ private readonly metricsService: MetricsService,
) {}
@Process(STREAM_AGENT_CHAT_JOB_NAME)
@@ -74,6 +78,13 @@ export class StreamAgentChatJob {
data: StreamAgentChatJobData,
context?: MessageQueueJobContext,
): Promise {
+ tagAiChatStreamScope({
+ streamId: data.streamId,
+ turnId: data.existingTurnId,
+ threadId: data.threadId,
+ workspaceId: data.workspaceId,
+ });
+
const thread = await this.threadRepository.findOne(data.workspaceId, {
where: { id: data.threadId },
select: ['id', 'activeStreamId'],
@@ -87,6 +98,12 @@ export class StreamAgentChatJob {
return;
}
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnStarted,
+ amount: 1,
+ attributes: { model: data.modelId ?? 'unknown' },
+ });
+
await this.eventPublisherService.resetStreamState(data.threadId);
const abortController = new AbortController();
@@ -132,6 +149,16 @@ export class StreamAgentChatJob {
);
const streamError = mapErrorToStreamError(error);
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: {
+ model: data.modelId ?? 'unknown',
+ failure_phase: 'execution',
+ error_code: streamError.code,
+ },
+ });
+
await this.threadRepository
.update(
data.workspaceId,
@@ -337,6 +364,8 @@ export class StreamAgentChatJob {
workspace,
userWorkspaceId: data.userWorkspaceId,
threadId: data.threadId,
+ streamId: data.streamId,
+ turnId: data.existingTurnId,
messages: data.messages,
browsingContext: data.browsingContext,
modelId: data.modelId,
@@ -662,6 +691,7 @@ export class StreamAgentChatJob {
threadId,
workspaceId,
streamUsage,
+ modelId: modelConfig.modelId,
});
}
@@ -726,6 +756,14 @@ export class StreamAgentChatJob {
return;
}
+ if (hasText && !isAborted && !isDefined(pendingQuestionPart)) {
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnCompleted,
+ amount: 1,
+ attributes: { model: modelConfig.modelId },
+ });
+ }
+
await this.agentChatService.notifyThreadUsageUpdated({
threadId,
userWorkspaceId,
@@ -742,6 +780,7 @@ export class StreamAgentChatJob {
threadId,
workspaceId,
streamUsage,
+ modelId,
}: {
responseMessage: Omit;
isAborted: boolean;
@@ -754,6 +793,7 @@ export class StreamAgentChatJob {
inputTokens: number;
outputTokens: number;
};
+ modelId: string;
}): void {
const reason = isAborted
? 'user-cancelled'
@@ -763,6 +803,14 @@ export class StreamAgentChatJob {
? 'credits-exhausted'
: 'empty-completion';
+ if (reason === 'empty-completion') {
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: { model: modelId, failure_phase: 'no_text' },
+ });
+ }
+
const errorDetail =
streamError instanceof Error
? `${streamError.name}: ${streamError.message}`
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat.resolver.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat.resolver.ts
index 6841384edf..7088240a6d 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat.resolver.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat.resolver.ts
@@ -37,6 +37,7 @@ import { AgentChatStreamingService } from 'src/engine/metadata-modules/ai/ai-cha
import { AgentChatService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service';
import { SystemPromptBuilderService } from 'src/engine/metadata-modules/ai/ai-chat/services/system-prompt-builder.service';
import { getCancelChannel } from 'src/engine/metadata-modules/ai/ai-chat/utils/get-cancel-channel.util';
+import { tagAiChatStreamScope } from 'src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util';
import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models/services/ai-model-registry.service';
import {
AiException,
@@ -45,6 +46,7 @@ import {
import { AiGraphqlApiExceptionInterceptor } from 'src/engine/metadata-modules/ai/interceptors/ai-graphql-api-exception.interceptor';
import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator';
import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
+
@UseGuards(WorkspaceAuthGuard, SettingsPermissionGuard(PermissionFlagType.AI))
@UseInterceptors(AiGraphqlApiExceptionInterceptor)
@MetadataResolver(() => AgentChatThreadDTO)
@@ -255,6 +257,13 @@ export class AgentChatResolver {
return { messageId: result.messageId, queued: true };
}
+ tagAiChatStreamScope({
+ streamId: result.streamId,
+ turnId: result.turnId,
+ threadId,
+ workspaceId: workspace.id,
+ });
+
return {
messageId: result.messageId,
queued: false,
@@ -291,6 +300,13 @@ export class AgentChatResolver {
modelId,
});
+ tagAiChatStreamScope({
+ streamId: result.streamId,
+ turnId: result.turnId,
+ threadId,
+ workspaceId: workspace.id,
+ });
+
return {
messageId: result.messageId,
queued: false,
@@ -375,6 +391,13 @@ export class AgentChatResolver {
throw error;
}
+ tagAiChatStreamScope({
+ streamId,
+ turnId,
+ threadId,
+ workspaceId: workspace.id,
+ });
+
return { messageId, queued: false, streamId };
}
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts
index adfacf8489..65c6d48d6c 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts
@@ -63,6 +63,7 @@ describe('AgentChatStreamingService claim & reap', () => {
eventPublisherService as never,
{ signFileByIdUrl: jest.fn() } as never,
streamHeartbeatService as never,
+ { incrementCounterBy: jest.fn() } as never,
);
return {
@@ -92,6 +93,12 @@ describe('AgentChatStreamingService claim & reap', () => {
const result = await service.streamAgentChat(sendArguments);
expect(result.queued).toBe(false);
+ expect(result).toEqual(
+ expect.objectContaining({
+ messageId: 'user-message-id',
+ turnId: 'turn-id',
+ }),
+ );
expect(threadRepository.update).toHaveBeenCalledWith(
'workspace-id',
expect.objectContaining({ id: 'thread-id' }),
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.retry.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.retry.spec.ts
index 791eef09ec..438f373242 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.retry.spec.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.retry.spec.ts
@@ -61,6 +61,7 @@ describe('AgentChatStreamingService.retryLastFailedTurn', () => {
{ publish: jest.fn() } as never,
{ signFileByIdUrl: jest.fn() } as never,
streamHeartbeatService as never,
+ { incrementCounterBy: jest.fn() } as never,
);
return { service, threadRepository, messageQueueService, agentChatService };
@@ -153,5 +154,6 @@ describe('AgentChatStreamingService.retryLastFailedTurn', () => {
{ activeStreamId: result.streamId, lastStreamError: null },
);
expect(result.messageId).toBe('user-message-id');
+ expect(result.turnId).toBe('turn-id');
});
});
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts
index fc4fd0708a..dfc3a3f2c7 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts
@@ -15,6 +15,8 @@ import { FileUrlService } from 'src/engine/core-modules/file/file-url/file-url.s
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
+import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service';
+import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type';
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import {
AgentMessageRole,
@@ -30,6 +32,7 @@ import { AgentChatEventPublisherService } from 'src/engine/metadata-modules/ai/a
import { AgentChatStreamHeartbeatService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-stream-heartbeat.service';
import { AgentChatService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service';
import { AiChatFileAttachment } from 'src/engine/metadata-modules/ai/ai-chat/types/ai-chat-file-attachment.type';
+import { mapErrorToStreamError } from 'src/engine/metadata-modules/ai/ai-chat/utils/map-error-to-stream-error.util';
import {
AiException,
AiExceptionCode,
@@ -62,6 +65,7 @@ export class AgentChatStreamingService {
private readonly eventPublisherService: AgentChatEventPublisherService,
private readonly fileUrlService: FileUrlService,
private readonly streamHeartbeatService: AgentChatStreamHeartbeatService,
+ private readonly metricsService: MetricsService,
) {}
async reapDeadStream({
@@ -95,6 +99,15 @@ export class AgentChatStreamingService {
return null;
}
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: {
+ failure_phase: 'interrupted',
+ error_code: interruptedError.code,
+ },
+ });
+
await this.eventPublisherService.resetStreamState(thread.id);
await this.eventPublisherService
.publish({
@@ -147,7 +160,12 @@ export class AgentChatStreamingService {
messageId,
fileAttachments,
}: StreamAgentChatOptions): Promise<
- | { queued: false; streamId: string; messageId: string }
+ | {
+ queued: false;
+ streamId: string;
+ messageId: string;
+ turnId: string | null;
+ }
| { queued: true; messageId: string }
> {
const thread = await this.threadRepository.findOne(workspace.id, {
@@ -253,9 +271,25 @@ export class AgentChatStreamingService {
},
);
- return { queued: false, streamId, messageId: savedUserMessage.id };
+ return {
+ queued: false,
+ streamId,
+ messageId: savedUserMessage.id,
+ turnId: savedUserMessage.turnId,
+ };
} catch (error) {
await this.releaseStreamClaim(threadId, workspace.id, streamId);
+ const streamError = mapErrorToStreamError(error);
+
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: {
+ model: modelId ?? 'unknown',
+ failure_phase: 'enqueue',
+ error_code: streamError.code,
+ },
+ });
throw error;
}
}
@@ -270,7 +304,7 @@ export class AgentChatStreamingService {
userWorkspaceId: string;
workspace: WorkspaceEntity;
modelId?: string;
- }): Promise<{ streamId: string; messageId: string }> {
+ }): Promise<{ streamId: string; messageId: string; turnId: string }> {
const thread = await this.threadRepository.findOne(workspace.id, {
where: { id: threadId, userWorkspaceId },
});
@@ -364,11 +398,26 @@ export class AgentChatStreamingService {
},
);
- return { streamId, messageId: lastUserMessage.id };
+ return {
+ streamId,
+ messageId: lastUserMessage.id,
+ turnId: lastUserMessage.turnId,
+ };
} catch (error) {
await this.releaseStreamClaim(threadId, workspace.id, streamId, {
lastStreamError: thread.lastStreamError,
});
+ const streamError = mapErrorToStreamError(error);
+
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: {
+ model: modelId ?? 'unknown',
+ failure_phase: 'enqueue',
+ error_code: streamError.code,
+ },
+ });
throw error;
}
}
@@ -543,6 +592,17 @@ export class AgentChatStreamingService {
);
} catch (error) {
await this.releaseStreamClaim(threadId, workspaceId, streamId);
+ const streamError = mapErrorToStreamError(error);
+
+ this.metricsService.incrementCounterBy({
+ key: MetricsKeys.AiChatTurnFailed,
+ amount: 1,
+ attributes: {
+ model: 'unknown',
+ failure_phase: 'enqueue',
+ error_code: streamError.code,
+ },
+ });
throw error;
}
}
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts
index d28d7a4f02..4ddc781bb3 100644
--- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/chat-execution.service.ts
@@ -1,6 +1,5 @@
import { Injectable, Logger } from '@nestjs/common';
-import { isNonEmptyString, isObject } from '@sniptt/guards';
import {
convertToModelMessages,
hasToolCall,
@@ -39,14 +38,10 @@ import { getToolMetricName } from 'src/engine/core-modules/tool-provider/utils/g
import { isToolOutputSuccessful } from 'src/engine/core-modules/tool-provider/utils/is-tool-output-successful.util';
import { resolveToolName } from 'src/engine/core-modules/tool-provider/utils/resolve-tool-name.util';
import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
-import {
- AiException,
- AiExceptionCode,
-} from 'src/engine/metadata-modules/ai/ai.exception';
import { AgentActorContextService } from 'src/engine/metadata-modules/ai/ai-agent-execution/services/agent-actor-context.service';
import { finalizeDanglingToolParts } from 'src/engine/metadata-modules/ai/ai-agent-execution/utils/finalize-dangling-tool-parts.util';
import { AGENT_CONFIG } from 'src/engine/metadata-modules/ai/ai-agent/constants/agent-config.const';
-import { type BrowsingContextType } from 'src/engine/metadata-modules/ai/ai-agent/types/browsingContext.type';
+import { BrowsingContextType } from 'src/engine/metadata-modules/ai/ai-agent/types/browsingContext.type';
import { repairToolCall } from 'src/engine/metadata-modules/ai/ai-agent/utils/repair-tool-call.util';
import { AiBillingService } from 'src/engine/metadata-modules/ai/ai-billing/services/ai-billing.service';
import { convertDollarsToBillingCredits } from 'src/engine/metadata-modules/ai/ai-billing/utils/convert-dollars-to-billing-credits.util';
@@ -56,12 +51,12 @@ 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 { 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 {
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';
import { extractCodeInterpreterFiles } from 'src/engine/metadata-modules/ai/ai-chat/utils/extract-code-interpreter-files.util';
import { injectMessageTimestamps } from 'src/engine/metadata-modules/ai/ai-chat/utils/inject-message-timestamps.util';
@@ -76,12 +71,18 @@ import { AiModelRegistryService } from 'src/engine/metadata-modules/ai/ai-models
import { NativeToolBinderService } from 'src/engine/metadata-modules/ai/ai-models/services/native-tool-binder.service';
import { type AiModelConfig } from 'src/engine/metadata-modules/ai/ai-models/types/ai-model-config.type';
import { getNativeModelCapabilities } from 'src/engine/metadata-modules/ai/ai-models/utils/get-native-model-capabilities.util';
+import {
+ AiException,
+ AiExceptionCode,
+} from 'src/engine/metadata-modules/ai/ai.exception';
import { SkillService } from 'src/engine/metadata-modules/skill/skill.service';
export type ChatExecutionOptions = {
workspace: WorkspaceEntity;
userWorkspaceId: string;
threadId?: string;
+ streamId?: string;
+ turnId?: string;
messages: ExtendedUIMessage[];
browsingContext: BrowsingContextType | null;
onCodeExecutionUpdate?: CodeExecutionStreamEmitter;
@@ -120,6 +121,8 @@ export class ChatExecutionService {
workspace,
userWorkspaceId,
threadId,
+ streamId,
+ turnId,
messages,
browsingContext,
onCodeExecutionUpdate,
@@ -439,7 +442,16 @@ export class ChatExecutionService {
stepCountIs(AGENT_CONFIG.MAX_STEPS)(step) ||
hasToolCall(ASK_QUESTIONS_TOOL_NAME)(step) ||
hasNoMoreAvailableCredits,
- experimental_telemetry: AI_TELEMETRY_CONFIG,
+ experimental_telemetry: {
+ ...AI_TELEMETRY_CONFIG,
+ functionId: 'ai-chat-stream',
+ metadata: {
+ streamId: streamId ?? '',
+ turnId: turnId ?? '',
+ threadId: threadId ?? '',
+ workspaceId: workspace.id,
+ },
+ },
providerOptions: getCallLevelProviderOptions({
sdkPackage: registeredModel.sdkPackage,
providerOptions: undefined,
@@ -533,52 +545,6 @@ export class ChatExecutionService {
unit: 'token',
attributes: executionAttributes,
});
-
- const { input } = part;
-
- if (part.toolName === LEARN_TOOLS_TOOL_NAME) {
- const learntToolNames =
- isObject(input) && 'toolNames' in input
- ? input.toolNames
- : undefined;
-
- for (const learntToolName of Array.isArray(learntToolNames)
- ? learntToolNames.filter(isNonEmptyString)
- : []) {
- this.metricsService.incrementCounterBy({
- key: succeeded
- ? MetricsKeys.AiChatToolLearnedSucceeded
- : MetricsKeys.AiChatToolLearnedFailed,
- amount: 1,
- attributes: {
- model: registeredModel.modelId,
- tool: getToolMetricName(learntToolName),
- },
- });
- }
- }
-
- if (part.toolName === LOAD_SKILL_TOOL_NAME) {
- const loadedSkillNames =
- isObject(input) && 'skillNames' in input
- ? input.skillNames
- : undefined;
-
- for (const loadedSkillName of Array.isArray(loadedSkillNames)
- ? loadedSkillNames.filter(isNonEmptyString)
- : []) {
- this.metricsService.incrementCounterBy({
- key: succeeded
- ? MetricsKeys.AiChatSkillLoadedSucceeded
- : MetricsKeys.AiChatSkillLoadedFailed,
- amount: 1,
- attributes: {
- model: registeredModel.modelId,
- skill: loadedSkillName,
- },
- });
- }
- }
}
},
onAbort: async ({ steps }) => {
diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util.ts
new file mode 100644
index 0000000000..2931baae9e
--- /dev/null
+++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/utils/tag-ai-chat-stream-scope.util.ts
@@ -0,0 +1,20 @@
+import * as Sentry from '@sentry/node';
+
+export const tagAiChatStreamScope = ({
+ streamId,
+ turnId,
+ threadId,
+ workspaceId,
+}: {
+ streamId: string;
+ turnId?: string | null;
+ threadId: string;
+ workspaceId: string;
+}) => {
+ Sentry.getCurrentScope().setTags({
+ streamId,
+ turnId: turnId ?? undefined,
+ threadId,
+ workspaceId,
+ });
+};