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` <!-- This is an auto-generated description by cubic. --> <a href="https://cubic.dev/pr/twentyhq/twenty/pull/22692?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:
@@ -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',
|
||||
}
|
||||
|
||||
+1
@@ -173,6 +173,7 @@ describe('StreamAgentChatJob', () => {
|
||||
cancelSubscriberService as never,
|
||||
agentChatStreamingService as never,
|
||||
streamHeartbeatService as never,
|
||||
{ incrementCounterBy: jest.fn() } as never,
|
||||
);
|
||||
|
||||
return {
|
||||
|
||||
+48
@@ -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<void> {
|
||||
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<ExtendedUIMessage, 'id'>;
|
||||
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}`
|
||||
|
||||
+23
@@ -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 };
|
||||
}
|
||||
|
||||
|
||||
+7
@@ -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' }),
|
||||
|
||||
+2
@@ -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');
|
||||
});
|
||||
});
|
||||
|
||||
+64
-4
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
+21
-55
@@ -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 }) => {
|
||||
|
||||
+20
@@ -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,
|
||||
});
|
||||
};
|
||||
Reference in New Issue
Block a user